fix(sync service): skip messages with lower index (#1134)

* skip messages with lower index

* address review comments

* chore: auto version bump [bot]

---------

Co-authored-by: Thegaram <Thegaram@users.noreply.github.com>
This commit is contained in:
Jonas Theis 2025-03-05 18:00:31 +08:00 committed by GitHub
parent 3c454e7101
commit 6954c73459
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
2 changed files with 30 additions and 4 deletions

View file

@ -24,7 +24,7 @@ import (
const ( const (
VersionMajor = 5 // Major version component of the current release VersionMajor = 5 // Major version component of the current release
VersionMinor = 8 // Minor version component of the current release VersionMinor = 8 // Minor version component of the current release
VersionPatch = 20 // Patch version component of the current release VersionPatch = 21 // Patch version component of the current release
VersionMeta = "mainnet" // Version metadata to append to the version string VersionMeta = "mainnet" // Version metadata to append to the version string
) )

View file

@ -1,6 +1,7 @@
package sync_service package sync_service
import ( import (
"bytes"
"context" "context"
"fmt" "fmt"
"reflect" "reflect"
@ -15,6 +16,7 @@ import (
"github.com/scroll-tech/go-ethereum/metrics" "github.com/scroll-tech/go-ethereum/metrics"
"github.com/scroll-tech/go-ethereum/node" "github.com/scroll-tech/go-ethereum/node"
"github.com/scroll-tech/go-ethereum/params" "github.com/scroll-tech/go-ethereum/params"
"github.com/scroll-tech/go-ethereum/rlp"
) )
const ( const (
@ -272,19 +274,43 @@ func (s *SyncService) fetchMessages() {
if len(msgs) > 0 { if len(msgs) > 0 {
log.Debug("Received new L1 events", "fromBlock", from, "toBlock", to, "count", len(msgs)) log.Debug("Received new L1 events", "fromBlock", from, "toBlock", to, "count", len(msgs))
rawdb.WriteL1Messages(batchWriter, msgs) // collect messages in memory
numMsgsCollected += len(msgs)
} }
for _, msg := range msgs { for _, msg := range msgs {
if msg.QueueIndex > 0 { if msg.QueueIndex > 0 {
queueIndex++ queueIndex++
} }
// check if received queue index matches expected queue index // check if received queue index matches expected queue index
if msg.QueueIndex != queueIndex { if msg.QueueIndex > queueIndex {
log.Error("Unexpected queue index in SyncService", "expected", queueIndex, "got", msg.QueueIndex, "msg", msg) log.Error("Unexpected queue index in SyncService", "expected", queueIndex, "got", msg.QueueIndex, "msg", msg)
return // do not flush inconsistent data to disk return // do not flush inconsistent data to disk
} }
// compare with stored message in database, abort if not equal, ignore if already exists
if msg.QueueIndex < queueIndex {
log.Warn("Duplicate queue index in SyncService", "expected", queueIndex, "got", msg.QueueIndex)
receivedMsgBytes, err := rlp.EncodeToBytes(msg)
if err != nil {
log.Error("Failed to encode message", "err", err)
return
}
storedMsgBytes := rawdb.ReadL1MessageRLP(s.db, msg.QueueIndex)
if !bytes.Equal(storedMsgBytes, receivedMsgBytes) {
storedL1Message := rawdb.ReadL1Message(s.db, msg.QueueIndex)
log.Error("Stored message at same queue index does not match received message", "queueIndex", msg.QueueIndex, "expected", storedL1Message, "got", msg)
return
}
// already exists, ignore
queueIndex--
continue
}
// store message to database (collected in memory and flushed periodically)
rawdb.WriteL1Message(batchWriter, msg)
numMsgsCollected++
} }
numBlocksPendingDbWrite += to - from + 1 numBlocksPendingDbWrite += to - from + 1