Skip to content

Commit f2e2f97

Browse files
authored
fix(syncing): clean up already processed DA data (#3413)
fix: clean up processed DA data
1 parent 1f8b332 commit f2e2f97

2 files changed

Lines changed: 57 additions & 5 deletions

File tree

block/internal/syncing/syncer.go

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -589,6 +589,7 @@ func (s *Syncer) processHeightEvent(ctx context.Context, event *common.DAHeightE
589589

590590
// Skip if already processed
591591
if height <= currentHeight || s.cache.IsHeaderSeen(headerHash) {
592+
s.removePendingDAData(event)
592593
s.logger.Debug().
593594
Uint64("height", height).
594595
Str("source", string(event.Source)).
@@ -844,15 +845,21 @@ func (s *Syncer) trySyncNextBlockWithState(ctx context.Context, event *common.DA
844845
s.p2pHandler.SetProcessedHeight(newState.LastBlockHeight)
845846
}
846847

847-
if event.Source == common.SourceDA {
848-
if cleaner, ok := s.daRetriever.(pendingDataCleaner); ok {
849-
cleaner.removePendingData(nextHeight)
850-
}
851-
}
848+
s.removePendingDAData(event)
852849

853850
return nil
854851
}
855852

853+
func (s *Syncer) removePendingDAData(event *common.DAHeightEvent) {
854+
if event.Source != common.SourceDA {
855+
return
856+
}
857+
858+
if cleaner, ok := s.daRetriever.(pendingDataCleaner); ok {
859+
cleaner.removePendingData(event.Header.Height())
860+
}
861+
}
862+
856863
// ApplyBlock applies a block to get the new state
857864
func (s *Syncer) ApplyBlock(ctx context.Context, header types.Header, data *types.Data, currentState types.State) (types.State, error) {
858865
// Prepare transactions

block/internal/syncing/syncer_test.go

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -622,6 +622,51 @@ func TestProcessHeightEvent_SyncsAndUpdatesState(t *testing.T) {
622622
assert.Equal(t, uint64(1), st1.LastBlockHeight)
623623
}
624624

625+
func TestProcessHeightEvent_AlreadyProcessedDADataDoesNotAccumulate(t *testing.T) {
626+
const alreadyProcessedBlocks = 8
627+
628+
ds := dssync.MutexWrap(datastore.NewMapDatastore())
629+
st := store.New(ds)
630+
631+
batch, err := st.NewBatch(t.Context())
632+
require.NoError(t, err)
633+
require.NoError(t, batch.SetHeight(alreadyProcessedBlocks))
634+
require.NoError(t, batch.Commit())
635+
636+
addr, pub, signer := buildSyncTestSigner(t)
637+
gen := genesis.Genesis{
638+
ChainID: "tchain",
639+
InitialHeight: 1,
640+
StartTime: time.Now().Add(-time.Second),
641+
ProposerAddress: addr,
642+
}
643+
daRetriever := newTestDARetriever(t, nil, config.DefaultConfig(), gen)
644+
s := &Syncer{
645+
store: st,
646+
daRetriever: daRetriever,
647+
logger: zerolog.Nop(),
648+
}
649+
650+
futureHeight := uint64(alreadyProcessedBlocks + 1)
651+
futureDataBin, _ := makeSignedDataBytes(t, gen.ChainID, futureHeight, addr, pub, signer, 8)
652+
require.Empty(t, daRetriever.processBlobs(t.Context(), [][]byte{futureDataBin}, futureHeight))
653+
require.Contains(t, daRetriever.pendingData, futureHeight)
654+
655+
for height := uint64(1); height <= alreadyProcessedBlocks; height++ {
656+
dataBin, data := makeSignedDataBytes(t, gen.ChainID, height, addr, pub, signer, 8)
657+
headerBin, _ := makeSignedHeaderBytes(t, gen.ChainID, height, addr, pub, signer, nil, &data.Data, nil)
658+
659+
events := daRetriever.processBlobs(t.Context(), [][]byte{headerBin, dataBin}, height)
660+
require.Len(t, events, 1)
661+
662+
s.processHeightEvent(t.Context(), &events[0])
663+
require.NotContains(t, daRetriever.pendingData, height, "processed DA data must be removed immediately")
664+
}
665+
666+
require.Len(t, daRetriever.pendingData, 1)
667+
require.Contains(t, daRetriever.pendingData, futureHeight, "DA data above the current height must remain pending")
668+
}
669+
625670
func TestProcessHeightEvent_UnexpectedProposerFromDAIsNotCriticalStateError(t *testing.T) {
626671
ds := dssync.MutexWrap(datastore.NewMapDatastore())
627672
st := store.New(ds)

0 commit comments

Comments
 (0)