mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-26 14:46:42 +00:00
fix: ensure L1 messages are stored in db consistently (#679)
* fix: ensure L1 messages are stored in db consistently * check db iterator errors
This commit is contained in:
parent
e9246f1cc2
commit
411889ef92
5 changed files with 34 additions and 2 deletions
|
|
@ -178,6 +178,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
|
||||||
|
|
@ -192,6 +195,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
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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 = 1 // Minor version component of the current release
|
VersionMinor = 1 // Minor version component of the current release
|
||||||
VersionPatch = 24 // Patch version component of the current release
|
VersionPatch = 25 // 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
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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,16 @@ func (s *SyncService) fetchMessages() {
|
||||||
numMsgsCollected += len(msgs)
|
numMsgsCollected += len(msgs)
|
||||||
}
|
}
|
||||||
|
|
||||||
numBlocksPendingDbWrite += to - from
|
for _, msg := range msgs {
|
||||||
|
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
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue