Skip to content

Commit 044903f

Browse files
committed
x
1 parent 0b04da3 commit 044903f

15 files changed

Lines changed: 374 additions & 219 deletions

block/components.go

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"errors"
66
"fmt"
77

8+
"github.com/evstack/ev-node/pkg/sync"
89
"github.com/rs/zerolog"
910

1011
"github.com/evstack/ev-node/block/internal/cache"
@@ -132,8 +133,8 @@ func NewSyncComponents(
132133
store store.Store,
133134
exec coreexecutor.Executor,
134135
da coreda.DA,
135-
headerStore common.Broadcaster[*types.SignedHeaderWithDAHint],
136-
dataStore common.Broadcaster[*types.Data],
136+
headerStore *sync.HeaderSyncService,
137+
dataStore *sync.DataSyncService,
137138
logger zerolog.Logger,
138139
metrics *Metrics,
139140
blockOpts BlockOptions,
@@ -157,15 +158,15 @@ func NewSyncComponents(
157158
metrics,
158159
config,
159160
genesis,
160-
headerStore,
161-
dataStore,
161+
common.NewDecorator[*types.SignedHeader](headerStore),
162+
common.NewDecorator[*types.Data](dataStore),
162163
logger,
163164
blockOpts,
164165
errorCh,
165166
)
166167

167168
// Create submitter for sync nodes (no signer, only DA inclusion processing)
168-
daSubmitter := submitting.NewDASubmitter(daClient, config, genesis, blockOpts, metrics, logger, nil) // todo (Alex): use a noop
169+
daSubmitter := submitting.NewDASubmitter(daClient, config, genesis, blockOpts, metrics, logger, headerStore, dataStore)
169170
submitter := submitting.NewSubmitter(
170171
store,
171172
exec,
@@ -198,8 +199,8 @@ func NewAggregatorComponents(
198199
sequencer coresequencer.Sequencer,
199200
da coreda.DA,
200201
signer signer.Signer,
201-
headerBroadcaster common.Broadcaster[*types.SignedHeaderWithDAHint],
202-
dataBroadcaster common.Broadcaster[*types.Data],
202+
headerBroadcaster *sync.HeaderSyncService,
203+
dataBroadcaster *sync.DataSyncService,
203204
logger zerolog.Logger,
204205
metrics *Metrics,
205206
blockOpts BlockOptions,
@@ -222,8 +223,8 @@ func NewAggregatorComponents(
222223
metrics,
223224
config,
224225
genesis,
225-
headerBroadcaster,
226-
dataBroadcaster,
226+
common.NewDecorator[*types.SignedHeader](headerBroadcaster),
227+
common.NewDecorator[*types.Data](dataBroadcaster),
227228
logger,
228229
blockOpts,
229230
errorCh,
@@ -247,7 +248,7 @@ func NewAggregatorComponents(
247248

248249
// Create DA client and submitter for aggregator nodes (with signer for submission)
249250
daClient := NewDAClient(da, config, logger)
250-
daSubmitter := submitting.NewDASubmitter(daClient, config, genesis, blockOpts, metrics, logger, headerBroadcaster)
251+
daSubmitter := submitting.NewDASubmitter(daClient, config, genesis, blockOpts, metrics, logger, headerBroadcaster, dataBroadcaster)
251252
submitter := submitting.NewSubmitter(
252253
store,
253254
exec,

block/internal/common/broadcaster_mock.go

Lines changed: 77 additions & 57 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

block/internal/common/event.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,6 @@ type DAHeightEvent struct {
2121
// Source indicates where this event originated from (DA or P2P)
2222
Source EventSource
2323

24-
// Optional DA height hint from P2P
25-
DaHeightHint uint64
24+
// Optional DA height hints from P2P. first is the DA height hint for the header, second is the DA height hint for the data
25+
DaHeightHints [2]uint64
2626
}

block/internal/common/expected_interfaces.go

Lines changed: 51 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,14 +3,63 @@ package common
33
import (
44
"context"
55

6+
"github.com/evstack/ev-node/types"
67
pubsub "github.com/libp2p/go-libp2p-pubsub"
78

89
"github.com/celestiaorg/go-header"
910
)
1011

11-
// broadcaster interface for P2P broadcasting
12+
type (
13+
HeaderP2PBroadcaster = Decorator[*types.SignedHeader]
14+
DataP2PBroadcaster = Decorator[*types.Data]
15+
)
16+
17+
// Broadcaster interface for P2P broadcasting
1218
type Broadcaster[H header.Header[H]] interface {
1319
WriteToStoreAndBroadcast(ctx context.Context, payload H, opts ...pubsub.PubOpt) error
1420
Store() header.Store[H]
15-
XXX(ctx context.Context, headerOrData H) error
21+
AppendDAHint(ctx context.Context, daHeight uint64, hashes ...types.Hash) error
22+
}
23+
24+
// Decorator to access the the payload type without the container
25+
type Decorator[H header.Header[H]] struct {
26+
nested Broadcaster[*types.DAHeightHintContainer[H]]
27+
}
28+
29+
func NewDecorator[H header.Header[H]](nested Broadcaster[*types.DAHeightHintContainer[H]]) Decorator[H] {
30+
return Decorator[H]{nested: nested}
31+
}
32+
33+
func (d Decorator[H]) WriteToStoreAndBroadcast(ctx context.Context, payload H, opts ...pubsub.PubOpt) error {
34+
return d.nested.WriteToStoreAndBroadcast(ctx, &types.DAHeightHintContainer[H]{Entry: payload}, opts...)
35+
}
36+
37+
func (d Decorator[H]) Store() HeightStore[H] {
38+
return HeightStoreImpl[H]{store: d.nested.Store()}
39+
}
40+
func (d Decorator[H]) XStore() header.Store[*types.DAHeightHintContainer[H]] {
41+
return d.nested.Store()
42+
}
43+
44+
func (d Decorator[H]) AppendDAHint(ctx context.Context, daHeight uint64, hashes ...types.Hash) error {
45+
return d.nested.AppendDAHint(ctx, daHeight, hashes...)
46+
}
47+
48+
// HeightStore is a subset of goheader.Store
49+
type HeightStore[H header.Header[H]] interface {
50+
GetByHeight(context.Context, uint64) (H, error)
51+
}
52+
53+
type HeightStoreImpl[H header.Header[H]] struct {
54+
store header.Store[*types.DAHeightHintContainer[H]]
55+
}
56+
57+
func (s HeightStoreImpl[H]) GetByHeight(ctx context.Context, height uint64) (H, error) {
58+
var zero H
59+
v, err := s.store.GetByHeight(ctx, height)
60+
if err != nil {
61+
return zero, err
62+
}
63+
return v.Entry, nil
64+
1665
}

block/internal/executing/executor.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -37,8 +37,8 @@ type Executor struct {
3737
metrics *common.Metrics
3838

3939
// Broadcasting
40-
headerBroadcaster common.Broadcaster[*types.SignedHeaderWithDAHint]
41-
dataBroadcaster common.Broadcaster[*types.Data]
40+
headerBroadcaster common.HeaderP2PBroadcaster
41+
dataBroadcaster common.DataP2PBroadcaster
4242

4343
// Configuration
4444
config config.Config
@@ -76,8 +76,8 @@ func NewExecutor(
7676
metrics *common.Metrics,
7777
config config.Config,
7878
genesis genesis.Genesis,
79-
headerBroadcaster common.Broadcaster[*types.SignedHeaderWithDAHint],
80-
dataBroadcaster common.Broadcaster[*types.Data],
79+
headerBroadcaster common.HeaderP2PBroadcaster,
80+
dataBroadcaster common.DataP2PBroadcaster,
8181
logger zerolog.Logger,
8282
options common.BlockOptions,
8383
errorCh chan<- error,
@@ -420,7 +420,7 @@ func (e *Executor) produceBlock() error {
420420

421421
// broadcast header and data to P2P network
422422
g, ctx := errgroup.WithContext(e.ctx)
423-
g.Go(func() error { return e.headerBroadcaster.WriteToStoreAndBroadcast(ctx, &types.SignedHeaderWithDAHint{SignedHeader: header}) })
423+
g.Go(func() error { return e.headerBroadcaster.WriteToStoreAndBroadcast(ctx, header) })
424424
g.Go(func() error { return e.dataBroadcaster.WriteToStoreAndBroadcast(ctx, data) })
425425
if err := g.Wait(); err != nil {
426426
e.logger.Error().Err(err).Msg("failed to broadcast header and/data")

0 commit comments

Comments
 (0)