fix: sync l1msg fix from develop branch (#826)

* fix: ensure L1 messages are stored in db consistently (#679)

* fix: ensure L1 messages are stored in db consistently

* check db iterator errors

* fix(sync-service): only add queue index when message index is not zero (#682)

* fix(sync-service): increase queue index only when L1 message index is not zero

---------

Co-authored-by: Péter Garamvölgyi <peter@scroll.io>

fix

---------

Co-authored-by: Péter Garamvölgyi <peter@scroll.io>
Co-authored-by: colin <102356659+colinlyguo@users.noreply.github.com>
This commit is contained in:
HAOYUatHZ 2024-06-13 20:15:21 +08:00 committed by GitHub
parent 1ca3b74ec1
commit bed81b02f1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 35 additions and 1 deletions

View file

@ -226,6 +226,9 @@ func (v *BlockValidator) ValidateL1Messages(block *types.Block) error {
// TODO: consider verifying that skipped messages overflow // TODO: consider verifying that skipped messages overflow
for index := queueIndex; index < txQueueIndex; index++ { for index := queueIndex; index < txQueueIndex; index++ {
if exists := it.Next(); !exists { if exists := it.Next(); !exists {
if err := it.Error(); err != nil {
log.Error("Unexpected DB error in ValidateL1Messages", "err", err, "queueIndex", queueIndex)
}
// the message in this block is not available in our local db. // the message in this block is not available in our local db.
// we'll reprocess this block at a later time. // we'll reprocess this block at a later time.
return consensus.ErrMissingL1MessageData return consensus.ErrMissingL1MessageData
@ -240,6 +243,9 @@ func (v *BlockValidator) ValidateL1Messages(block *types.Block) error {
queueIndex = txQueueIndex + 1 queueIndex = txQueueIndex + 1
if exists := it.Next(); !exists { if exists := it.Next(); !exists {
if err := it.Error(); err != nil {
log.Error("Unexpected DB error in ValidateL1Messages", "err", err, "queueIndex", txQueueIndex)
}
// the message in this block is not available in our local db. // the message in this block is not available in our local db.
// we'll reprocess this block at a later time. // we'll reprocess this block at a later time.
return consensus.ErrMissingL1MessageData return consensus.ErrMissingL1MessageData

View file

@ -202,6 +202,12 @@ func (it *L1MessageIterator) Release() {
it.inner.Release() it.inner.Release()
} }
// Error returns any accumulated error.
// Exhausting all the key/value pairs is not considered to be an error.
func (it *L1MessageIterator) Error() error {
return it.inner.Error()
}
// ReadL1MessagesFrom retrieves up to `maxCount` L1 messages starting at `startIndex`. // ReadL1MessagesFrom retrieves up to `maxCount` L1 messages starting at `startIndex`.
func ReadL1MessagesFrom(db ethdb.Database, startIndex, maxCount uint64) []types.L1MessageTx { func ReadL1MessagesFrom(db ethdb.Database, startIndex, maxCount uint64) []types.L1MessageTx {
msgs := make([]types.L1MessageTx, 0, maxCount) msgs := make([]types.L1MessageTx, 0, maxCount)
@ -236,6 +242,10 @@ func ReadL1MessagesFrom(db ethdb.Database, startIndex, maxCount uint64) []types.
} }
} }
if err := it.Error(); err != nil {
log.Crit("Failed to read L1 messages", "err", err)
}
return msgs return msgs
} }

View file

@ -86,6 +86,10 @@ func (c *BridgeClient) fetchMessagesInRange(ctx context.Context, from, to uint64
}) })
} }
if err := it.Error(); err != nil {
return nil, err
}
return msgs, nil return msgs, nil
} }

View file

@ -149,6 +149,9 @@ func (s *SyncService) fetchMessages() {
log.Trace("Sync service fetchMessages", "latestProcessedBlock", s.latestProcessedBlock, "latestConfirmed", latestConfirmed) log.Trace("Sync service fetchMessages", "latestProcessedBlock", s.latestProcessedBlock, "latestConfirmed", latestConfirmed)
// keep track of next queue index we're expecting to see
queueIndex := rawdb.ReadHighestSyncedQueueIndex(s.db)
batchWriter := s.db.NewBatch() batchWriter := s.db.NewBatch()
numBlocksPendingDbWrite := uint64(0) numBlocksPendingDbWrite := uint64(0)
numMessagesPendingDbWrite := 0 numMessagesPendingDbWrite := 0
@ -216,7 +219,18 @@ func (s *SyncService) fetchMessages() {
numMsgsCollected += len(msgs) numMsgsCollected += len(msgs)
} }
numBlocksPendingDbWrite += to - from for _, msg := range msgs {
if msg.QueueIndex > 0 {
queueIndex++
}
// check if received queue index matches expected queue index
if msg.QueueIndex != queueIndex {
log.Error("Unexpected queue index in SyncService", "expected", queueIndex, "got", msg.QueueIndex, "msg", msg)
return // do not flush inconsistent data to disk
}
}
numBlocksPendingDbWrite += to - from + 1
numMessagesPendingDbWrite += len(msgs) numMessagesPendingDbWrite += len(msgs)
// flush new messages to database periodically // flush new messages to database periodically