From 5d454ae5424ea9f900a4ce9f0313a819cf03d891 Mon Sep 17 00:00:00 2001 From: HAOYUatHZ <37070449+HAOYUatHZ@users.noreply.github.com> Date: Mon, 20 Nov 2023 13:06:18 +0800 Subject: [PATCH] add `rollup/sync_service` package (#565) --- rollup/sync_service/bindings.go | 150 ++++++++++++++++++ rollup/sync_service/bridge_client.go | 129 +++++++++++++++ rollup/sync_service/sync_service.go | 227 +++++++++++++++++++++++++++ rollup/sync_service/types.go | 21 +++ 4 files changed, 527 insertions(+) create mode 100644 rollup/sync_service/bindings.go create mode 100644 rollup/sync_service/bridge_client.go create mode 100644 rollup/sync_service/sync_service.go create mode 100644 rollup/sync_service/types.go diff --git a/rollup/sync_service/bindings.go b/rollup/sync_service/bindings.go new file mode 100644 index 0000000000..1c2b1393e9 --- /dev/null +++ b/rollup/sync_service/bindings.go @@ -0,0 +1,150 @@ +// Code generated - DO NOT EDIT. +// This file is a generated binding and any manual changes will be lost. + +// generated using: +// forge flatten src/L1/rollup/L1MessageQueue.sol > flatten.sol +// go run github.com/scroll-tech/go-ethereum/cmd/abigen@develop --sol flatten.sol --pkg rollup --out ./L1MessageQueue.go --contract L1MessageQueue + +package sync_service + +import ( + "math/big" + "strings" + + ethereum "github.com/scroll-tech/go-ethereum" + "github.com/scroll-tech/go-ethereum/accounts/abi" + "github.com/scroll-tech/go-ethereum/accounts/abi/bind" + "github.com/scroll-tech/go-ethereum/common" + "github.com/scroll-tech/go-ethereum/core/types" +) + +// L1MessageQueueMetaData contains all meta data concerning the L1MessageQueue contract. +var L1MessageQueueMetaData = &bind.MetaData{ + ABI: "[{\"anonymous\":false,\"inputs\":[{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"startIndex\",\"type\":\"uint256\"},{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"count\",\"type\":\"uint256\"},{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"skippedBitmap\",\"type\":\"uint256\"}],\"name\":\"DequeueTransaction\",\"type\":\"event\"},{\"anonymous\":false,\"inputs\":[{\"indexed\":true,\"internalType\":\"address\",\"name\":\"previousOwner\",\"type\":\"address\"},{\"indexed\":true,\"internalType\":\"address\",\"name\":\"newOwner\",\"type\":\"address\"}],\"name\":\"OwnershipTransferred\",\"type\":\"event\"},{\"anonymous\":false,\"inputs\":[{\"indexed\":true,\"internalType\":\"address\",\"name\":\"sender\",\"type\":\"address\"},{\"indexed\":true,\"internalType\":\"address\",\"name\":\"target\",\"type\":\"address\"},{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"value\",\"type\":\"uint256\"},{\"indexed\":false,\"internalType\":\"uint64\",\"name\":\"queueIndex\",\"type\":\"uint64\"},{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"gasLimit\",\"type\":\"uint256\"},{\"indexed\":false,\"internalType\":\"bytes\",\"name\":\"data\",\"type\":\"bytes\"}],\"name\":\"QueueTransaction\",\"type\":\"event\"},{\"anonymous\":false,\"inputs\":[{\"indexed\":false,\"internalType\":\"address\",\"name\":\"_oldGateway\",\"type\":\"address\"},{\"indexed\":false,\"internalType\":\"address\",\"name\":\"_newGateway\",\"type\":\"address\"}],\"name\":\"UpdateEnforcedTxGateway\",\"type\":\"event\"},{\"anonymous\":false,\"inputs\":[{\"indexed\":false,\"internalType\":\"address\",\"name\":\"_oldGasOracle\",\"type\":\"address\"},{\"indexed\":false,\"internalType\":\"address\",\"name\":\"_newGasOracle\",\"type\":\"address\"}],\"name\":\"UpdateGasOracle\",\"type\":\"event\"},{\"anonymous\":false,\"inputs\":[{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"_oldMaxGasLimit\",\"type\":\"uint256\"},{\"indexed\":false,\"internalType\":\"uint256\",\"name\":\"_newMaxGasLimit\",\"type\":\"uint256\"}],\"name\":\"UpdateMaxGasLimit\",\"type\":\"event\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"_target\",\"type\":\"address\"},{\"internalType\":\"uint256\",\"name\":\"_gasLimit\",\"type\":\"uint256\"},{\"internalType\":\"bytes\",\"name\":\"_data\",\"type\":\"bytes\"}],\"name\":\"appendCrossDomainMessage\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"_sender\",\"type\":\"address\"},{\"internalType\":\"address\",\"name\":\"_target\",\"type\":\"address\"},{\"internalType\":\"uint256\",\"name\":\"_value\",\"type\":\"uint256\"},{\"internalType\":\"uint256\",\"name\":\"_gasLimit\",\"type\":\"uint256\"},{\"internalType\":\"bytes\",\"name\":\"_data\",\"type\":\"bytes\"}],\"name\":\"appendEnforcedTransaction\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"bytes\",\"name\":\"_calldata\",\"type\":\"bytes\"}],\"name\":\"calculateIntrinsicGasFee\",\"outputs\":[{\"internalType\":\"uint256\",\"name\":\"\",\"type\":\"uint256\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"_sender\",\"type\":\"address\"},{\"internalType\":\"uint256\",\"name\":\"_queueIndex\",\"type\":\"uint256\"},{\"internalType\":\"uint256\",\"name\":\"_value\",\"type\":\"uint256\"},{\"internalType\":\"address\",\"name\":\"_target\",\"type\":\"address\"},{\"internalType\":\"uint256\",\"name\":\"_gasLimit\",\"type\":\"uint256\"},{\"internalType\":\"bytes\",\"name\":\"_data\",\"type\":\"bytes\"}],\"name\":\"computeTransactionHash\",\"outputs\":[{\"internalType\":\"bytes32\",\"name\":\"\",\"type\":\"bytes32\"}],\"stateMutability\":\"pure\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"enforcedTxGateway\",\"outputs\":[{\"internalType\":\"address\",\"name\":\"\",\"type\":\"address\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"uint256\",\"name\":\"_gasLimit\",\"type\":\"uint256\"}],\"name\":\"estimateCrossDomainMessageFee\",\"outputs\":[{\"internalType\":\"uint256\",\"name\":\"\",\"type\":\"uint256\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"gasOracle\",\"outputs\":[{\"internalType\":\"address\",\"name\":\"\",\"type\":\"address\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"uint256\",\"name\":\"_queueIndex\",\"type\":\"uint256\"}],\"name\":\"getCrossDomainMessage\",\"outputs\":[{\"internalType\":\"bytes32\",\"name\":\"\",\"type\":\"bytes32\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"_messenger\",\"type\":\"address\"},{\"internalType\":\"address\",\"name\":\"_scrollChain\",\"type\":\"address\"},{\"internalType\":\"address\",\"name\":\"_enforcedTxGateway\",\"type\":\"address\"},{\"internalType\":\"address\",\"name\":\"_gasOracle\",\"type\":\"address\"},{\"internalType\":\"uint256\",\"name\":\"_maxGasLimit\",\"type\":\"uint256\"}],\"name\":\"initialize\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"maxGasLimit\",\"outputs\":[{\"internalType\":\"uint256\",\"name\":\"\",\"type\":\"uint256\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"uint256\",\"name\":\"\",\"type\":\"uint256\"}],\"name\":\"messageQueue\",\"outputs\":[{\"internalType\":\"bytes32\",\"name\":\"\",\"type\":\"bytes32\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"messenger\",\"outputs\":[{\"internalType\":\"address\",\"name\":\"\",\"type\":\"address\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"nextCrossDomainMessageIndex\",\"outputs\":[{\"internalType\":\"uint256\",\"name\":\"\",\"type\":\"uint256\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"owner\",\"outputs\":[{\"internalType\":\"address\",\"name\":\"\",\"type\":\"address\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"pendingQueueIndex\",\"outputs\":[{\"internalType\":\"uint256\",\"name\":\"\",\"type\":\"uint256\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"uint256\",\"name\":\"_startIndex\",\"type\":\"uint256\"},{\"internalType\":\"uint256\",\"name\":\"_count\",\"type\":\"uint256\"},{\"internalType\":\"uint256\",\"name\":\"_skippedBitmap\",\"type\":\"uint256\"}],\"name\":\"popCrossDomainMessage\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"renounceOwnership\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[],\"name\":\"scrollChain\",\"outputs\":[{\"internalType\":\"address\",\"name\":\"\",\"type\":\"address\"}],\"stateMutability\":\"view\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"newOwner\",\"type\":\"address\"}],\"name\":\"transferOwnership\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"_newGateway\",\"type\":\"address\"}],\"name\":\"updateEnforcedTxGateway\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"address\",\"name\":\"_newGasOracle\",\"type\":\"address\"}],\"name\":\"updateGasOracle\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"},{\"inputs\":[{\"internalType\":\"uint256\",\"name\":\"_newMaxGasLimit\",\"type\":\"uint256\"}],\"name\":\"updateMaxGasLimit\",\"outputs\":[],\"stateMutability\":\"nonpayable\",\"type\":\"function\"}]", +} + +// L1MessageQueueABI is the input ABI used to generate the binding from. +// Deprecated: Use L1MessageQueueMetaData.ABI instead. +var L1MessageQueueABI = L1MessageQueueMetaData.ABI + +// L1MessageQueueFilterer is an auto generated log filtering Go binding around an Ethereum contract events. +type L1MessageQueueFilterer struct { + contract *bind.BoundContract // Generic contract wrapper for the low level calls +} + +// NewL1MessageQueueFilterer creates a new log filterer instance of L1MessageQueue, bound to a specific deployed contract. +func NewL1MessageQueueFilterer(address common.Address, filterer bind.ContractFilterer) (*L1MessageQueueFilterer, error) { + contract, err := bindL1MessageQueue(address, nil, nil, filterer) + if err != nil { + return nil, err + } + return &L1MessageQueueFilterer{contract: contract}, nil +} + +// bindL1MessageQueue binds a generic wrapper to an already deployed contract. +func bindL1MessageQueue(address common.Address, caller bind.ContractCaller, transactor bind.ContractTransactor, filterer bind.ContractFilterer) (*bind.BoundContract, error) { + parsed, err := abi.JSON(strings.NewReader(L1MessageQueueABI)) + if err != nil { + return nil, err + } + return bind.NewBoundContract(address, parsed, caller, transactor, filterer), nil +} + +// L1MessageQueueQueueTransactionIterator is returned from FilterQueueTransaction and is used to iterate over the raw logs and unpacked data for QueueTransaction events raised by the L1MessageQueue contract. +type L1MessageQueueQueueTransactionIterator struct { + Event *L1MessageQueueQueueTransaction // Event containing the contract specifics and raw log + + contract *bind.BoundContract // Generic contract to use for unpacking event data + event string // Event name to use for unpacking event data + + logs chan types.Log // Log channel receiving the found contract events + sub ethereum.Subscription // Subscription for errors, completion and termination + done bool // Whether the subscription completed delivering logs + fail error // Occurred error to stop iteration +} + +// Next advances the iterator to the subsequent event, returning whether there +// are any more events found. In case of a retrieval or parsing error, false is +// returned and Error() can be queried for the exact failure. +func (it *L1MessageQueueQueueTransactionIterator) Next() bool { + // If the iterator failed, stop iterating + if it.fail != nil { + return false + } + // If the iterator completed, deliver directly whatever's available + if it.done { + select { + case log := <-it.logs: + it.Event = new(L1MessageQueueQueueTransaction) + if err := it.contract.UnpackLog(it.Event, it.event, log); err != nil { + it.fail = err + return false + } + it.Event.Raw = log + return true + + default: + return false + } + } + // Iterator still in progress, wait for either a data or an error event + select { + case log := <-it.logs: + it.Event = new(L1MessageQueueQueueTransaction) + if err := it.contract.UnpackLog(it.Event, it.event, log); err != nil { + it.fail = err + return false + } + it.Event.Raw = log + return true + + case err := <-it.sub.Err(): + it.done = true + it.fail = err + return it.Next() + } +} + +// Error returns any retrieval or parsing error occurred during filtering. +func (it *L1MessageQueueQueueTransactionIterator) Error() error { + return it.fail +} + +// Close terminates the iteration process, releasing any pending underlying +// resources. +func (it *L1MessageQueueQueueTransactionIterator) Close() error { + it.sub.Unsubscribe() + return nil +} + +// L1MessageQueueQueueTransaction represents a QueueTransaction event raised by the L1MessageQueue contract. +type L1MessageQueueQueueTransaction struct { + Sender common.Address + Target common.Address + Value *big.Int + QueueIndex uint64 + GasLimit *big.Int + Data []byte + Raw types.Log // Blockchain specific contextual infos +} + +// FilterQueueTransaction is a free log retrieval operation binding the contract event 0x69cfcb8e6d4192b8aba9902243912587f37e550d75c1fa801491fce26717f37e. +// +// Solidity: event QueueTransaction(address indexed sender, address indexed target, uint256 value, uint64 queueIndex, uint256 gasLimit, bytes data) +func (_L1MessageQueue *L1MessageQueueFilterer) FilterQueueTransaction(opts *bind.FilterOpts, sender []common.Address, target []common.Address) (*L1MessageQueueQueueTransactionIterator, error) { + + var senderRule []interface{} + for _, senderItem := range sender { + senderRule = append(senderRule, senderItem) + } + var targetRule []interface{} + for _, targetItem := range target { + targetRule = append(targetRule, targetItem) + } + + logs, sub, err := _L1MessageQueue.contract.FilterLogs(opts, "QueueTransaction", senderRule, targetRule) + if err != nil { + return nil, err + } + return &L1MessageQueueQueueTransactionIterator{contract: _L1MessageQueue.contract, event: "QueueTransaction", logs: logs, sub: sub}, nil +} diff --git a/rollup/sync_service/bridge_client.go b/rollup/sync_service/bridge_client.go new file mode 100644 index 0000000000..ea46fdcb8b --- /dev/null +++ b/rollup/sync_service/bridge_client.go @@ -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) +} diff --git a/rollup/sync_service/sync_service.go b/rollup/sync_service/sync_service.go new file mode 100644 index 0000000000..2720aec76b --- /dev/null +++ b/rollup/sync_service/sync_service.go @@ -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) + } + } +} diff --git a/rollup/sync_service/types.go b/rollup/sync_service/types.go new file mode 100644 index 0000000000..0bc062772a --- /dev/null +++ b/rollup/sync_service/types.go @@ -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) +}