Skip to content

Commit c4104c6

Browse files
committed
x
1 parent 0804010 commit c4104c6

9 files changed

Lines changed: 305 additions & 94 deletions

File tree

block/components.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,11 @@ func (bc *Components) Stop() error {
119119
errs = errors.Join(errs, fmt.Errorf("failed to stop submitter: %w", err))
120120
}
121121
}
122+
if bc.Cache != nil {
123+
if err := bc.Cache.SaveToDisk(); err != nil {
124+
errs = errors.Join(errs, fmt.Errorf("failed to save caches: %w", err))
125+
}
126+
}
122127

123128
return errs
124129
}

block/internal/executing/executor.go

Lines changed: 24 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -408,6 +408,30 @@ func (e *Executor) produceBlock() error {
408408
return fmt.Errorf("failed to validate block: %w", err)
409409
}
410410

411+
// Propose block to raft to share state in the cluster
412+
if e.raftNode != nil {
413+
headerBytes, err := header.MarshalBinary()
414+
if err != nil {
415+
return fmt.Errorf("failed to marshal header: %w", err)
416+
}
417+
dataBytes, err := data.MarshalBinary()
418+
if err != nil {
419+
return fmt.Errorf("failed to marshal data: %w", err)
420+
}
421+
422+
raftState := &raft.RaftBlockState{
423+
Height: newHeight,
424+
Hash: header.Hash(),
425+
Timestamp: header.BaseHeader.Time,
426+
Header: headerBytes,
427+
Data: dataBytes,
428+
}
429+
if err := e.raftNode.Broadcast(e.ctx, raftState); err != nil {
430+
return fmt.Errorf("failed to propose block to raft: %w", err)
431+
}
432+
e.logger.Debug().Uint64("height", newHeight).Msg("proposed block to raft")
433+
}
434+
411435
batch, err := e.store.NewBatch(e.ctx)
412436
if err != nil {
413437
return fmt.Errorf("failed to create batch: %w", err)
@@ -433,30 +457,6 @@ func (e *Executor) produceBlock() error {
433457
e.setLastState(newState)
434458
e.metrics.Height.Set(float64(newState.LastBlockHeight))
435459

436-
// Propose block to raft before p2p broadcast if raft is enabled
437-
if e.raftNode != nil {
438-
headerBytes, err := header.MarshalBinary()
439-
if err != nil {
440-
return fmt.Errorf("failed to marshal header: %w", err)
441-
}
442-
dataBytes, err := data.MarshalBinary()
443-
if err != nil {
444-
return fmt.Errorf("failed to marshal data: %w", err)
445-
}
446-
447-
raftState := &raft.RaftBlockState{
448-
Height: newHeight,
449-
Hash: header.Hash(),
450-
Timestamp: header.BaseHeader.Time,
451-
Header: headerBytes,
452-
Data: dataBytes,
453-
}
454-
if err := e.raftNode.Broadcast(e.ctx, raftState); err != nil {
455-
return fmt.Errorf("failed to propose block to raft: %w", err)
456-
}
457-
e.logger.Debug().Uint64("height", newHeight).Msg("proposed block to raft")
458-
}
459-
460460
// broadcast header and data to P2P network
461461
g, ctx := errgroup.WithContext(e.ctx)
462462
g.Go(func() error { return e.headerBroadcaster.WriteToStoreAndBroadcast(ctx, header) })
@@ -465,7 +465,6 @@ func (e *Executor) produceBlock() error {
465465
e.logger.Error().Err(err).Msg("failed to broadcast header and/data")
466466
// don't fail block production on broadcast error
467467
}
468-
469468
e.recordBlockMetrics(data)
470469

471470
e.logger.Info().

block/internal/submitting/da_submitter.go

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -191,7 +191,15 @@ func (s *DASubmitter) SubmitHeaders(ctx context.Context, cache cache.Manager) er
191191
return nil
192192
}
193193

194-
s.logger.Info().Int("count", len(headers)).Msg("submitting headers to DA")
194+
heights := make([]uint64, len(headers))
195+
for i, sd := range headers {
196+
heights[i] = sd.Height()
197+
if len(sd.Header.ProposerAddress) == 0 || bytes.Equal(sd.Header.ProposerAddress, make([]byte, 32)) {
198+
panic("empty proposer address")
199+
}
200+
}
201+
202+
s.logger.Info().Uints64("heights", heights).Int("count", len(headers)).Msg("submitting headers to DA")
195203

196204
return submitToDA(s, ctx, headers,
197205
func(header *types.SignedHeader) ([]byte, error) {
@@ -239,7 +247,11 @@ func (s *DASubmitter) SubmitData(ctx context.Context, cache cache.Manager, signe
239247
return nil // No non-empty data to submit
240248
}
241249

242-
s.logger.Info().Int("count", len(signedDataList)).Msg("submitting data to DA")
250+
heights := make([]uint64, len(signedDataList))
251+
for i, sd := range signedDataList {
252+
heights[i] = sd.Height()
253+
}
254+
s.logger.Info().Uints64("heights", heights).Int("count", len(signedDataList)).Msg("submitting data to DA")
243255

244256
return submitToDA(s, ctx, signedDataList,
245257
func(signedData *types.SignedData) ([]byte, error) {

block/internal/submitting/submitter.go

Lines changed: 60 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,7 @@ func (s *Submitter) Start(ctx context.Context) error {
103103

104104
// Start DA submission loop if signer is available (aggregator nodes only)
105105
if s.signer != nil {
106+
s.logger.Info().Msg("starting DA submission loop")
106107
s.wg.Add(1)
107108
go func() {
108109
defer s.wg.Done()
@@ -113,7 +114,14 @@ func (s *Submitter) Start(ctx context.Context) error {
113114
// Start DA inclusion processing loop (both sync and aggregator nodes)
114115
s.wg.Add(1)
115116
go func() {
116-
defer s.wg.Done()
117+
defer func() {
118+
// fetch all from DA to ensure final height is set
119+
// todo: revisit this for a graceful shutdown scenario
120+
tCtx, cancelFunc := context.WithTimeout(context.Background(), 2*time.Second)
121+
defer cancelFunc()
122+
s.processAllFromDA(tCtx)
123+
s.wg.Done()
124+
}()
117125
s.processDAInclusionLoop()
118126
}()
119127

@@ -193,48 +201,63 @@ func (s *Submitter) processDAInclusionLoop() {
193201
case <-s.ctx.Done():
194202
return
195203
case <-ticker.C:
196-
currentDAIncluded := s.GetDAIncludedHeight()
204+
s.processAllFromDA(s.ctx)
205+
}
206+
}
207+
}
197208

198-
for {
199-
nextHeight := currentDAIncluded + 1
209+
func (s *Submitter) processAllFromDA(ctx context.Context) {
210+
currentDAIncluded := s.GetDAIncludedHeight()
200211

201-
// Get block data first
202-
header, data, err := s.store.GetBlockData(s.ctx, nextHeight)
203-
if err != nil {
204-
break
205-
}
212+
for {
213+
select {
214+
case <-ctx.Done():
215+
return
216+
default:
217+
}
206218

207-
// Check if this height is DA included
208-
if included, err := s.IsHeightDAIncluded(nextHeight, header, data); err != nil || !included {
209-
break
210-
}
219+
nextHeight := currentDAIncluded + 1
220+
221+
// Get block data first
222+
header, data, err := s.store.GetBlockData(s.ctx, nextHeight)
223+
if err != nil {
224+
return
225+
}
211226

212-
s.logger.Debug().Uint64("height", nextHeight).Msg("advancing DA included height")
227+
// Check if this height is DA included
228+
if included, err := s.IsHeightDAIncluded(nextHeight, header, data); err != nil || !included {
229+
return
230+
}
213231

214-
// Set sequencer height to DA height mapping using already retrieved data
215-
if err := s.setSequencerHeightToDAHeight(s.ctx, nextHeight, header, data, currentDAIncluded == 0); err != nil {
216-
s.logger.Error().Err(err).Uint64("height", nextHeight).Msg("failed to set sequencer height to DA height mapping")
217-
break
218-
}
232+
s.logger.Debug().Uint64("height", nextHeight).Msg("advancing DA included height")
219233

220-
// Set final height in executor
221-
if err := s.setFinalWithRetry(nextHeight); err != nil {
222-
s.sendCriticalError(fmt.Errorf("failed to set final height: %w", err))
223-
s.logger.Error().Err(err).Uint64("height", nextHeight).Msg("failed to set final height")
224-
break
225-
}
234+
// Set sequencer height to DA height mapping using already retrieved data
235+
if err := s.setSequencerHeightToDAHeight(s.ctx, nextHeight, header, data, currentDAIncluded == 0); err != nil {
236+
s.logger.Error().Err(err).Uint64("height", nextHeight).Msg("failed to set sequencer height to DA height mapping")
237+
return
238+
}
226239

227-
// Update DA included height
228-
s.SetDAIncludedHeight(nextHeight)
229-
currentDAIncluded = nextHeight
240+
// Set final height in executor
241+
if err := s.setFinalWithRetry(nextHeight); err != nil {
242+
s.sendCriticalError(fmt.Errorf("failed to set final height: %w", err))
243+
s.logger.Error().Err(err).Uint64("height", nextHeight).Msg("failed to set final height")
244+
return
245+
}
230246

231-
// Persist DA included height
232-
bz := make([]byte, 8)
233-
binary.LittleEndian.PutUint64(bz, nextHeight)
234-
if err := s.store.SetMetadata(s.ctx, store.DAIncludedHeightKey, bz); err != nil {
235-
s.logger.Error().Err(err).Uint64("height", nextHeight).Msg("failed to persist DA included height")
236-
}
237-
}
247+
// Update DA included height
248+
s.SetDAIncludedHeight(nextHeight)
249+
currentDAIncluded = nextHeight
250+
251+
// Persist DA included height
252+
bz := make([]byte, 8)
253+
binary.LittleEndian.PutUint64(bz, nextHeight)
254+
if err := s.store.SetMetadata(s.ctx, store.DAIncludedHeightKey, bz); err != nil {
255+
s.logger.Error().Err(err).Uint64("height", nextHeight).Msg("failed to persist DA included height")
256+
}
257+
if s.signer == nil {
258+
// bump height to keep cache in sync when submission loop is not active
259+
s.cache.SetLastSubmittedDataHeight(s.ctx, currentDAIncluded)
260+
s.cache.SetLastSubmittedHeaderHeight(s.ctx, currentDAIncluded)
238261
}
239262
}
240263
}
@@ -280,6 +303,7 @@ func (s *Submitter) SetDAIncludedHeight(height uint64) {
280303
// initializeDAIncludedHeight loads the DA included height from store
281304
func (s *Submitter) initializeDAIncludedHeight(ctx context.Context) error {
282305
if height, err := s.store.GetMetadata(ctx, store.DAIncludedHeightKey); err == nil && len(height) == 8 {
306+
s.logger.Debug().Uint64("height", binary.LittleEndian.Uint64(height)).Msg("Initial DA included height loaded from store")
283307
s.SetDAIncludedHeight(binary.LittleEndian.Uint64(height))
284308
}
285309
return nil
@@ -367,7 +391,7 @@ func (s *Submitter) IsHeightDAIncluded(height uint64, header *types.SignedHeader
367391
_, headerIncluded := s.cache.GetHeaderDAIncluded(headerHash)
368392
_, dataIncluded := s.cache.GetDataDAIncluded(dataHash)
369393

370-
dataIncluded = bytes.Equal(data.DACommitment(), common.DataHashForEmptyTxs) || dataIncluded
394+
dataIncluded = dataIncluded || bytes.Equal(data.DACommitment(), common.DataHashForEmptyTxs)
371395

372396
return headerIncluded && dataIncluded, nil
373397
}

block/internal/syncing/syncer.go

Lines changed: 26 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -421,7 +421,12 @@ func (s *Syncer) processHeightEvent(event *common.DAHeightEvent) {
421421
if err := s.trySyncNextBlock(event); err != nil {
422422
s.logger.Error().Err(err).Msg("failed to sync next block")
423423
// If the error is not due to an validation error, re-store the event as pending
424-
if !errors.Is(err, errInvalidBlock) {
424+
switch {
425+
case errors.Is(err, errInvalidBlock):
426+
// do not reschedule
427+
case errors.Is(err, errInvalidState):
428+
s.logger.Fatal().Err(err).Msg("Invalid state, shutting down")
429+
default:
425430
s.cache.SetPendingEvent(height, event)
426431
}
427432
return
@@ -450,7 +455,7 @@ func (s *Syncer) trySyncNextBlock(event *common.DAHeightEvent) error {
450455
// Compared to the executor logic where the current block needs to be applied first,
451456
// here only the previous block needs to be applied to proceed to the verification.
452457
// The header validation must be done before applying the block to avoid executing gibberish
453-
if err := s.validateBlock(header, data); err != nil {
458+
if err := s.validateBlock(header, data, currentState); err != nil {
454459
// remove header as da included (not per se needed, but keep cache clean)
455460
s.cache.RemoveHeaderDAIncluded(header.Hash().String())
456461
return errors.Join(errInvalidBlock, fmt.Errorf("failed to validate block: %w", err))
@@ -555,14 +560,13 @@ func (s *Syncer) executeTxsWithRetry(ctx context.Context, rawTxs [][]byte, heade
555560
return nil, nil
556561
}
557562

563+
var errInvalidState = errors.New("invalid state")
564+
558565
// validateBlock validates a synced block
559566
// NOTE: if the header was gibberish and somehow passed all validation prior but the data was correct
560567
// or if the data was gibberish and somehow passed all validation prior but the header was correct
561568
// we are still losing both in the pending event. This should never happen.
562-
func (s *Syncer) validateBlock(
563-
header *types.SignedHeader,
564-
data *types.Data,
565-
) error {
569+
func (s *Syncer) validateBlock(header *types.SignedHeader, data *types.Data, state types.State, ) error {
566570
// Set custom verifier for aggregator node signature
567571
header.SetCustomVerifierForSyncNode(s.options.SyncNodeSignatureBytesProvider)
568572

@@ -575,7 +579,23 @@ func (s *Syncer) validateBlock(
575579
if err := types.Validate(header, data); err != nil {
576580
return fmt.Errorf("header-data validation failed: %w", err)
577581
}
582+
if state.LastBlockHeight < s.genesis.InitialHeight {
583+
return nil
584+
}
585+
// Validate header against state
586+
if header.Height() != state.LastBlockHeight+1 {
587+
return fmt.Errorf("%w: invalid block height - got: %d, want: %d", errInvalidState, header.Height(), state.LastBlockHeight+1)
588+
}
578589

590+
if !header.Time().After(state.LastBlockTime) {
591+
return fmt.Errorf("%w: invalid block time - got: %v, last: %v", errInvalidState, header.Time(), state.LastBlockTime)
592+
}
593+
if !bytes.Equal(header.LastHeaderHash, state.LastHeaderHash) {
594+
return fmt.Errorf("%w: invalid last header hash - got: %x, want: %x", errInvalidState, header.LastHeaderHash, header.LastHeaderHash)
595+
}
596+
if !bytes.Equal(header.AppHash, state.AppHash) {
597+
return fmt.Errorf("%w: invalid last appHash hash - got: %x, want: %x", errInvalidState, header.AppHash, header.AppHash)
598+
}
579599
return nil
580600
}
581601

node/full.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -105,7 +105,7 @@ func newFullNode(
105105
nodeConfig.P2P.Peers = "" // raft cluster provides most up to date blocks. No need for p2p peers.
106106
case nodeConfig.Node.Aggregator && !nodeConfig.Raft.Enable:
107107
leaderElection = AlwaysLeader{}
108-
case !nodeConfig.Node.Aggregator:
108+
case !nodeConfig.Node.Aggregator && !nodeConfig.Raft.Enable:
109109
leaderElection = AlwaysFollower{}
110110
default:
111111
return nil, fmt.Errorf("raft config must be used in sequencer setup only")
@@ -120,7 +120,7 @@ func newFullNode(
120120
aggregatorComponentFunc: func(ctx context.Context) error {
121121
logger.Info().Msg("Starting aggregator-MODE")
122122
nodeConfig.Node.Aggregator = true
123-
nodeConfig.P2P.Peers = ""
123+
nodeConfig.P2P.Peers = "" // peers are not supported in aggregator mode
124124
m, err := newAggregatorMode(nodeConfig, nodeKey, signer, genesis, database, exec, sequencer, da, logger, rktStore, mainKV, blockMetrics, nodeOpts, raftNode)
125125
if err != nil {
126126
return err

0 commit comments

Comments
 (0)