mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 10:22:23 +00:00
add rollup/sync_service package (#565)
This commit is contained in:
parent
261983dbdb
commit
5d454ae542
4 changed files with 527 additions and 0 deletions
150
rollup/sync_service/bindings.go
Normal file
150
rollup/sync_service/bindings.go
Normal file
File diff suppressed because one or more lines are too long
129
rollup/sync_service/bridge_client.go
Normal file
129
rollup/sync_service/bridge_client.go
Normal file
|
|
@ -0,0 +1,129 @@
|
|||
package sync_service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math/big"
|
||||
|
||||
"github.com/scroll-tech/go-ethereum/accounts/abi/bind"
|
||||
"github.com/scroll-tech/go-ethereum/common"
|
||||
"github.com/scroll-tech/go-ethereum/core/types"
|
||||
"github.com/scroll-tech/go-ethereum/log"
|
||||
"github.com/scroll-tech/go-ethereum/rpc"
|
||||
)
|
||||
|
||||
// BridgeClient is a wrapper around EthClient that adds
|
||||
// methods for conveniently collecting L1 messages.
|
||||
type BridgeClient struct {
|
||||
client EthClient
|
||||
confirmations rpc.BlockNumber
|
||||
l1MessageQueueAddress common.Address
|
||||
filterer *L1MessageQueueFilterer
|
||||
}
|
||||
|
||||
func newBridgeClient(ctx context.Context, l1Client EthClient, l1ChainId uint64, confirmations rpc.BlockNumber, l1MessageQueueAddress common.Address) (*BridgeClient, error) {
|
||||
if l1MessageQueueAddress == (common.Address{}) {
|
||||
return nil, errors.New("must pass non-zero l1MessageQueueAddress to BridgeClient")
|
||||
}
|
||||
|
||||
// sanity check: compare chain IDs
|
||||
got, err := l1Client.ChainID(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to query L1 chain ID, err = %w", err)
|
||||
}
|
||||
if got.Cmp(big.NewInt(0).SetUint64(l1ChainId)) != 0 {
|
||||
return nil, fmt.Errorf("unexpected chain ID, expected = %v, got = %v", l1ChainId, got)
|
||||
}
|
||||
|
||||
filterer, err := NewL1MessageQueueFilterer(l1MessageQueueAddress, l1Client)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to initialize L1MessageQueueFilterer, err = %w", err)
|
||||
}
|
||||
|
||||
client := BridgeClient{
|
||||
client: l1Client,
|
||||
confirmations: confirmations,
|
||||
l1MessageQueueAddress: l1MessageQueueAddress,
|
||||
filterer: filterer,
|
||||
}
|
||||
|
||||
return &client, nil
|
||||
}
|
||||
|
||||
// fetchMessagesInRange retrieves and parses all L1 messages between the
|
||||
// provided from and to L1 block numbers (inclusive).
|
||||
func (c *BridgeClient) fetchMessagesInRange(ctx context.Context, from, to uint64) ([]types.L1MessageTx, error) {
|
||||
log.Trace("BridgeClient fetchMessagesInRange", "fromBlock", from, "toBlock", to)
|
||||
|
||||
opts := bind.FilterOpts{
|
||||
Start: from,
|
||||
End: &to,
|
||||
Context: ctx,
|
||||
}
|
||||
it, err := c.filterer.FilterQueueTransaction(&opts, nil, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var msgs []types.L1MessageTx
|
||||
|
||||
for it.Next() {
|
||||
event := it.Event
|
||||
log.Trace("Received new L1 QueueTransaction event", "event", event)
|
||||
|
||||
if !event.GasLimit.IsUint64() {
|
||||
return nil, fmt.Errorf("invalid QueueTransaction event: QueueIndex = %v, GasLimit = %v", event.QueueIndex, event.GasLimit)
|
||||
}
|
||||
|
||||
msgs = append(msgs, types.L1MessageTx{
|
||||
QueueIndex: event.QueueIndex,
|
||||
Gas: event.GasLimit.Uint64(),
|
||||
To: &event.Target,
|
||||
Value: event.Value,
|
||||
Data: event.Data,
|
||||
Sender: event.Sender,
|
||||
})
|
||||
}
|
||||
|
||||
return msgs, nil
|
||||
}
|
||||
|
||||
func (c *BridgeClient) getLatestConfirmedBlockNumber(ctx context.Context) (uint64, error) {
|
||||
// confirmation based on "safe" or "finalized" block tag
|
||||
if c.confirmations == rpc.SafeBlockNumber || c.confirmations == rpc.FinalizedBlockNumber {
|
||||
tag := big.NewInt(int64(c.confirmations))
|
||||
header, err := c.client.HeaderByNumber(ctx, tag)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if !header.Number.IsInt64() {
|
||||
return 0, fmt.Errorf("received unexpected block number in BridgeClient: %v", header.Number)
|
||||
}
|
||||
return header.Number.Uint64(), nil
|
||||
}
|
||||
|
||||
// confirmation based on latest block number
|
||||
if c.confirmations == rpc.LatestBlockNumber {
|
||||
number, err := c.client.BlockNumber(ctx)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return number, nil
|
||||
}
|
||||
|
||||
// confirmation based on a certain number of blocks
|
||||
if c.confirmations.Int64() >= 0 {
|
||||
number, err := c.client.BlockNumber(ctx)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
confirmations := uint64(c.confirmations.Int64())
|
||||
if number >= confirmations {
|
||||
return number - confirmations, nil
|
||||
}
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
return 0, fmt.Errorf("unknown confirmation type: %v", c.confirmations)
|
||||
}
|
||||
227
rollup/sync_service/sync_service.go
Normal file
227
rollup/sync_service/sync_service.go
Normal file
|
|
@ -0,0 +1,227 @@
|
|||
package sync_service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"reflect"
|
||||
"time"
|
||||
|
||||
"github.com/scroll-tech/go-ethereum/core"
|
||||
"github.com/scroll-tech/go-ethereum/core/rawdb"
|
||||
"github.com/scroll-tech/go-ethereum/ethdb"
|
||||
"github.com/scroll-tech/go-ethereum/event"
|
||||
"github.com/scroll-tech/go-ethereum/log"
|
||||
"github.com/scroll-tech/go-ethereum/node"
|
||||
"github.com/scroll-tech/go-ethereum/params"
|
||||
)
|
||||
|
||||
const (
|
||||
// DefaultFetchBlockRange is the number of blocks that we collect in a single eth_getLogs query.
|
||||
DefaultFetchBlockRange = uint64(100)
|
||||
|
||||
// DefaultPollInterval is the frequency at which we query for new L1 messages.
|
||||
DefaultPollInterval = time.Second * 10
|
||||
|
||||
// LogProgressInterval is the frequency at which we log progress.
|
||||
LogProgressInterval = time.Second * 10
|
||||
|
||||
// DbWriteThresholdBytes is the size of batched database writes in bytes.
|
||||
DbWriteThresholdBytes = 10 * 1024
|
||||
|
||||
// DbWriteThresholdBlocks is the number of blocks scanned after which we write to the database
|
||||
// even if we have not collected DbWriteThresholdBytes bytes of data yet. This way, if there is
|
||||
// a long section of L1 blocks with no messages and we stop or crash, we will not need to re-scan
|
||||
// this secion.
|
||||
DbWriteThresholdBlocks = 1000
|
||||
)
|
||||
|
||||
// SyncService collects all L1 messages and stores them in a local database.
|
||||
type SyncService struct {
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
client *BridgeClient
|
||||
db ethdb.Database
|
||||
msgCountFeed event.Feed
|
||||
pollInterval time.Duration
|
||||
latestProcessedBlock uint64
|
||||
scope event.SubscriptionScope
|
||||
}
|
||||
|
||||
func NewSyncService(ctx context.Context, genesisConfig *params.ChainConfig, nodeConfig *node.Config, db ethdb.Database, l1Client EthClient) (*SyncService, error) {
|
||||
// terminate if the caller does not provide an L1 client (e.g. in tests)
|
||||
if l1Client == nil || (reflect.ValueOf(l1Client).Kind() == reflect.Ptr && reflect.ValueOf(l1Client).IsNil()) {
|
||||
log.Warn("No L1 client provided, L1 sync service will not run")
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
if genesisConfig.Scroll.L1Config == nil {
|
||||
return nil, fmt.Errorf("missing L1 config in genesis")
|
||||
}
|
||||
|
||||
client, err := newBridgeClient(ctx, l1Client, genesisConfig.Scroll.L1Config.L1ChainId, nodeConfig.L1Confirmations, genesisConfig.Scroll.L1Config.L1MessageQueueAddress)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to initialize bridge client: %w", err)
|
||||
}
|
||||
|
||||
// assume deployment block has 0 messages
|
||||
latestProcessedBlock := nodeConfig.L1DeploymentBlock
|
||||
block := rawdb.ReadSyncedL1BlockNumber(db)
|
||||
if block != nil {
|
||||
// restart from latest synced block number
|
||||
latestProcessedBlock = *block
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
|
||||
service := SyncService{
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
client: client,
|
||||
db: db,
|
||||
pollInterval: DefaultPollInterval,
|
||||
latestProcessedBlock: latestProcessedBlock,
|
||||
}
|
||||
|
||||
return &service, nil
|
||||
}
|
||||
|
||||
func (s *SyncService) Start() {
|
||||
if s == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// wait for initial sync before starting node
|
||||
log.Info("Starting L1 message sync service", "latestProcessedBlock", s.latestProcessedBlock)
|
||||
|
||||
// block node startup during initial sync and print some helpful logs
|
||||
latestConfirmed, err := s.client.getLatestConfirmedBlockNumber(s.ctx)
|
||||
if err == nil && latestConfirmed > s.latestProcessedBlock+1000 {
|
||||
log.Warn("Running initial sync of L1 messages before starting l2geth, this might take a while...")
|
||||
s.fetchMessages()
|
||||
log.Info("L1 message initial sync completed", "latestProcessedBlock", s.latestProcessedBlock)
|
||||
}
|
||||
|
||||
go func() {
|
||||
t := time.NewTicker(s.pollInterval)
|
||||
defer t.Stop()
|
||||
|
||||
for {
|
||||
// don't wait for ticker during startup
|
||||
s.fetchMessages()
|
||||
|
||||
select {
|
||||
case <-s.ctx.Done():
|
||||
return
|
||||
case <-t.C:
|
||||
continue
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func (s *SyncService) Stop() {
|
||||
if s == nil {
|
||||
return
|
||||
}
|
||||
|
||||
log.Info("Stopping sync service")
|
||||
|
||||
// Unsubscribe all subscriptions registered
|
||||
s.scope.Close()
|
||||
|
||||
if s.cancel != nil {
|
||||
s.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
// SubscribeNewL1MsgsEvent registers a subscription of NewL1MsgsEvent and
|
||||
// starts sending event to the given channel.
|
||||
func (s *SyncService) SubscribeNewL1MsgsEvent(ch chan<- core.NewL1MsgsEvent) event.Subscription {
|
||||
return s.scope.Track(s.msgCountFeed.Subscribe(ch))
|
||||
}
|
||||
|
||||
func (s *SyncService) fetchMessages() {
|
||||
latestConfirmed, err := s.client.getLatestConfirmedBlockNumber(s.ctx)
|
||||
if err != nil {
|
||||
log.Warn("Failed to get latest confirmed block number", "err", err)
|
||||
return
|
||||
}
|
||||
|
||||
log.Trace("Sync service fetchMessages", "latestProcessedBlock", s.latestProcessedBlock, "latestConfirmed", latestConfirmed)
|
||||
|
||||
batchWriter := s.db.NewBatch()
|
||||
numBlocksPendingDbWrite := uint64(0)
|
||||
numMessagesPendingDbWrite := 0
|
||||
|
||||
// helper function to flush database writes cached in memory
|
||||
flush := func(lastBlock uint64) {
|
||||
// update sync progress
|
||||
rawdb.WriteSyncedL1BlockNumber(batchWriter, lastBlock)
|
||||
|
||||
// write batch in a single transaction
|
||||
err := batchWriter.Write()
|
||||
if err != nil {
|
||||
// crash on database error, no risk of inconsistency here
|
||||
log.Crit("Failed to write L1 messages to database", "err", err)
|
||||
}
|
||||
|
||||
batchWriter.Reset()
|
||||
numBlocksPendingDbWrite = 0
|
||||
|
||||
if numMessagesPendingDbWrite > 0 {
|
||||
s.msgCountFeed.Send(core.NewL1MsgsEvent{Count: numMessagesPendingDbWrite})
|
||||
numMessagesPendingDbWrite = 0
|
||||
}
|
||||
|
||||
s.latestProcessedBlock = lastBlock
|
||||
}
|
||||
|
||||
// ticker for logging progress
|
||||
t := time.NewTicker(LogProgressInterval)
|
||||
numMsgsCollected := 0
|
||||
|
||||
// query in batches
|
||||
for from := s.latestProcessedBlock + 1; from <= latestConfirmed; from += DefaultFetchBlockRange {
|
||||
select {
|
||||
case <-s.ctx.Done():
|
||||
// flush pending writes to database
|
||||
if from > 0 {
|
||||
flush(from - 1)
|
||||
}
|
||||
return
|
||||
case <-t.C:
|
||||
progress := 100 * float64(s.latestProcessedBlock) / float64(latestConfirmed)
|
||||
log.Info("Syncing L1 messages", "processed", s.latestProcessedBlock, "confirmed", latestConfirmed, "collected", numMsgsCollected, "progress(%)", progress)
|
||||
default:
|
||||
}
|
||||
|
||||
to := from + DefaultFetchBlockRange - 1
|
||||
if to > latestConfirmed {
|
||||
to = latestConfirmed
|
||||
}
|
||||
|
||||
msgs, err := s.client.fetchMessagesInRange(s.ctx, from, to)
|
||||
if err != nil {
|
||||
// flush pending writes to database
|
||||
if from > 0 {
|
||||
flush(from - 1)
|
||||
}
|
||||
log.Warn("Failed to fetch L1 messages in range", "fromBlock", from, "toBlock", to, "err", err)
|
||||
return
|
||||
}
|
||||
|
||||
if len(msgs) > 0 {
|
||||
log.Debug("Received new L1 events", "fromBlock", from, "toBlock", to, "count", len(msgs))
|
||||
rawdb.WriteL1Messages(batchWriter, msgs) // collect messages in memory
|
||||
numMsgsCollected += len(msgs)
|
||||
}
|
||||
|
||||
numBlocksPendingDbWrite += to - from
|
||||
numMessagesPendingDbWrite += len(msgs)
|
||||
|
||||
// flush new messages to database periodically
|
||||
if to == latestConfirmed || batchWriter.ValueSize() >= DbWriteThresholdBytes || numBlocksPendingDbWrite >= DbWriteThresholdBlocks {
|
||||
flush(to)
|
||||
}
|
||||
}
|
||||
}
|
||||
21
rollup/sync_service/types.go
Normal file
21
rollup/sync_service/types.go
Normal file
|
|
@ -0,0 +1,21 @@
|
|||
package sync_service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"math/big"
|
||||
|
||||
"github.com/scroll-tech/go-ethereum"
|
||||
"github.com/scroll-tech/go-ethereum/common"
|
||||
"github.com/scroll-tech/go-ethereum/core/types"
|
||||
)
|
||||
|
||||
// We cannot use ethclient.Client directly as that would lead
|
||||
// to circular dependency between eth, rollup, and ethclient.
|
||||
type EthClient interface {
|
||||
BlockNumber(ctx context.Context) (uint64, error)
|
||||
ChainID(ctx context.Context) (*big.Int, error)
|
||||
FilterLogs(ctx context.Context, q ethereum.FilterQuery) ([]types.Log, error)
|
||||
HeaderByNumber(ctx context.Context, number *big.Int) (*types.Header, error)
|
||||
SubscribeFilterLogs(ctx context.Context, query ethereum.FilterQuery, ch chan<- types.Log) (ethereum.Subscription, error)
|
||||
TransactionByHash(ctx context.Context, txHash common.Hash) (tx *types.Transaction, isPending bool, err error)
|
||||
}
|
||||
Loading…
Reference in a new issue