From e113804bfc9d985b7fbd1eb32344c164d774e348 Mon Sep 17 00:00:00 2001 From: HAOYUatHZ <37070449+HAOYUatHZ@users.noreply.github.com> Date: Wed, 24 Jul 2024 21:31:25 +0800 Subject: [PATCH] feat: pipeline worker (#907) * comment out miner/worker.go * add miner/scroll_worker.go * rename package * add rollup/pipeline/pipeline.go * rename package * some renamings * some fixes * some fixes * some fixes * some fixes * some fixes * some fixes * some fixes * some fixes * some cleaning * some fixes * some fixes * remove `RelaxedPeriod` * some fixes * try `orderedTransactionSet` * try lazy * try lazy * update miner/worker_test.go * update miner/worker_test.go * feat: implement relaxed period in clique * fix `parent.Number` * fix coinbase * fix "create block" * fix "updateSnapshot" * make consensus engine wait for a bit before giving up on worker * fix error handling --- consensus/clique/clique.go | 10 +- core/blockchain.go | 5 + miner/payload_building.go | 7 + miner/scroll_worker.go | 751 ++++++++++++++++++++++++++++++++++++ miner/worker.go | 35 +- miner/worker_test.go | 17 +- params/config.go | 2 + rollup/pipeline/pipeline.go | 559 +++++++++++++++++++++++++++ 8 files changed, 1356 insertions(+), 30 deletions(-) create mode 100644 miner/scroll_worker.go create mode 100644 rollup/pipeline/pipeline.go diff --git a/consensus/clique/clique.go b/consensus/clique/clique.go index 8e6da2ddfa..6cb40fc9f8 100644 --- a/consensus/clique/clique.go +++ b/consensus/clique/clique.go @@ -329,7 +329,7 @@ func (c *Clique) verifyCascadingFields(chain consensus.ChainHeaderReader, header if parent == nil || parent.Number.Uint64() != number-1 || parent.Hash() != header.ParentHash { return consensus.ErrUnknownAncestor } - if parent.Time+c.config.Period > header.Time { + if header.Time < parent.Time { return errInvalidTimestamp } // Verify that the gasUsed is <= gasLimit @@ -559,7 +559,9 @@ func (c *Clique) Prepare(chain consensus.ChainHeaderReader, header *types.Header return consensus.ErrUnknownAncestor } header.Time = parent.Time + c.config.Period - if header.Time < uint64(time.Now().Unix()) { + // If RelaxedPeriod is enabled, always set the header timestamp to now (ie the time we start building it) as + // we don't know when it will be sealed + if c.config.RelaxedPeriod || header.Time < uint64(time.Now().Unix()) { header.Time = uint64(time.Now().Unix()) } return nil @@ -651,6 +653,8 @@ func (c *Clique) Seal(chain consensus.ChainHeaderReader, block *types.Block, res // Wait until sealing is terminated or delay timeout. log.Trace("Waiting for slot to sign and propagate", "delay", common.PrettyDuration(delay)) go func() { + defer close(results) + select { case <-stop: return @@ -659,7 +663,7 @@ func (c *Clique) Seal(chain consensus.ChainHeaderReader, block *types.Block, res select { case results <- block.WithSeal(header): - default: + case <-time.After(time.Second): log.Warn("Sealing result is not read by miner", "sealhash", SealHash(header)) } }() diff --git a/core/blockchain.go b/core/blockchain.go index 8fd9b303cc..6cf9839fca 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -2669,3 +2669,8 @@ func (bc *BlockChain) SetTrieFlushInterval(interval time.Duration) { func (bc *BlockChain) GetTrieFlushInterval() time.Duration { return time.Duration(bc.flushInterval.Load()) } + +// Database gives access to the underlying database for convenience +func (bc *BlockChain) Database() ethdb.Database { + return bc.db +} diff --git a/miner/payload_building.go b/miner/payload_building.go index 69ffab75b5..2ece05d0ee 100644 --- a/miner/payload_building.go +++ b/miner/payload_building.go @@ -19,6 +19,7 @@ package miner import ( "crypto/sha256" "encoding/binary" + "errors" "math/big" "sync" "time" @@ -175,6 +176,11 @@ func (payload *Payload) ResolveFull() *engine.ExecutionPayloadEnvelope { } // buildPayload builds the payload according to the provided parameters. +func (w *worker) buildPayload(args *BuildPayloadArgs) (*Payload, error) { + return nil, errors.New("not implemented") +} + +/* func (w *worker) buildPayload(args *BuildPayloadArgs) (*Payload, error) { // Build the initial version with no transaction included. It should be fast // enough to run. The empty payload can at least make sure there is something @@ -241,3 +247,4 @@ func (w *worker) buildPayload(args *BuildPayloadArgs) (*Payload, error) { }() return payload, nil } +*/ diff --git a/miner/scroll_worker.go b/miner/scroll_worker.go new file mode 100644 index 0000000000..221bfa2de2 --- /dev/null +++ b/miner/scroll_worker.go @@ -0,0 +1,751 @@ +// Copyright 2015 The go-ethereum Authors +// This file is part of the go-ethereum library. +// +// The go-ethereum library is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// The go-ethereum library is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with the go-ethereum library. If not, see . + +package miner + +import ( + "bytes" + "errors" + "math/big" + "sync" + "sync/atomic" + "time" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/consensus" + "github.com/ethereum/go-ethereum/consensus/misc" + "github.com/ethereum/go-ethereum/consensus/misc/eip1559" + "github.com/ethereum/go-ethereum/core" + "github.com/ethereum/go-ethereum/core/rawdb" + "github.com/ethereum/go-ethereum/core/state" + "github.com/ethereum/go-ethereum/core/txpool" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/event" + "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/metrics" + "github.com/ethereum/go-ethereum/params" + "github.com/ethereum/go-ethereum/rollup/circuitcapacitychecker" + "github.com/ethereum/go-ethereum/rollup/fees" + "github.com/ethereum/go-ethereum/rollup/pipeline" + "github.com/ethereum/go-ethereum/trie" +) + +const ( + // txChanSize is the size of channel listening to NewTxsEvent. + // The number is referenced from the size of tx pool. + txChanSize = 4096 + + // chainHeadChanSize is the size of channel listening to ChainHeadEvent. + chainHeadChanSize = 10 +) + +var ( + // Metrics for the skipped txs + l1TxGasLimitExceededCounter = metrics.NewRegisteredCounter("miner/skipped_txs/l1/gas_limit_exceeded", nil) + l1TxRowConsumptionOverflowCounter = metrics.NewRegisteredCounter("miner/skipped_txs/l1/row_consumption_overflow", nil) + l2TxRowConsumptionOverflowCounter = metrics.NewRegisteredCounter("miner/skipped_txs/l2/row_consumption_overflow", nil) + l1TxCccUnknownErrCounter = metrics.NewRegisteredCounter("miner/skipped_txs/l1/ccc_unknown_err", nil) + l2TxCccUnknownErrCounter = metrics.NewRegisteredCounter("miner/skipped_txs/l2/ccc_unknown_err", nil) + l1TxStrangeErrCounter = metrics.NewRegisteredCounter("miner/skipped_txs/l1/strange_err", nil) + + collectL1MsgsTimer = metrics.NewRegisteredTimer("miner/collect_l1_msgs", nil) + prepareTimer = metrics.NewRegisteredTimer("miner/prepare", nil) + collectL2Timer = metrics.NewRegisteredTimer("miner/collect_l2_txns", nil) + l2CommitTimer = metrics.NewRegisteredTimer("miner/commit", nil) + + commitReasonCCCCounter = metrics.NewRegisteredCounter("miner/commit_reason_ccc", nil) + commitReasonDeadlineCounter = metrics.NewRegisteredCounter("miner/commit_reason_deadline", nil) + commitGasCounter = metrics.NewRegisteredCounter("miner/commit_gas", nil) +) + +// prioritizedTransaction represents a single transaction that +// should be processed as the first transaction in the next block. +type prioritizedTransaction struct { + blockNumber uint64 + tx *types.Transaction +} + +// worker is the main object which takes care of submitting new work to consensus engine +// and gathering the sealing result. +type worker struct { + config *Config + chainConfig *params.ChainConfig + engine consensus.Engine + eth Backend + chain *core.BlockChain + + // Feeds + pendingLogsFeed event.Feed + + // Subscriptions + mux *event.TypeMux + txsCh chan core.NewTxsEvent + txsSub event.Subscription + chainHeadCh chan core.ChainHeadEvent + chainHeadSub event.Subscription + + // Channels + startCh chan struct{} + exitCh chan struct{} + + wg sync.WaitGroup + + currentPipelineStart time.Time + currentPipeline *pipeline.Pipeline + + mu sync.RWMutex // The lock used to protect the coinbase and extra fields + coinbase common.Address + extra []byte + + snapshotMu sync.RWMutex // The lock used to protect the snapshots below + snapshotBlock *types.Block + snapshotReceipts types.Receipts + snapshotState *state.StateDB + + // atomic status counters + running atomic.Bool // The indicator whether the consensus engine is running or not. + newTxs atomic.Int32 // New arrival transaction count since last sealing work submitting. + syncing atomic.Bool // The indicator whether the node is still syncing. + + // noempty is the flag used to control whether the feature of pre-seal empty + // block is enabled. The default value is false(pre-seal is enabled by default). + // But in some special scenario the consensus engine will seal blocks instantaneously, + // in this case this feature will add all empty blocks into canonical chain + // non-stop and no real transaction will be included. + noempty uint32 + + // External functions + isLocalBlock func(block *types.Header) bool // Function used to determine whether the specified block is mined by local miner. + + circuitCapacityChecker *circuitcapacitychecker.CircuitCapacityChecker + prioritizedTx *prioritizedTransaction + + // Test hooks + beforeTxHook func() // Method to call before processing a transaction. +} + +func newWorker(config *Config, chainConfig *params.ChainConfig, engine consensus.Engine, eth Backend, mux *event.TypeMux, isLocalBlock func(*types.Header) bool, init bool) *worker { + worker := &worker{ + config: config, + chainConfig: chainConfig, + engine: engine, + eth: eth, + chain: eth.BlockChain(), + mux: mux, + isLocalBlock: isLocalBlock, + coinbase: config.Etherbase, + extra: config.ExtraData, + txsCh: make(chan core.NewTxsEvent, txChanSize), + chainHeadCh: make(chan core.ChainHeadEvent, chainHeadChanSize), + exitCh: make(chan struct{}), + startCh: make(chan struct{}, 1), + circuitCapacityChecker: circuitcapacitychecker.NewCircuitCapacityChecker(true), + } + log.Info("created new worker", "CircuitCapacityChecker ID", worker.circuitCapacityChecker.ID) + + // Subscribe NewTxsEvent for tx pool + worker.txsSub = eth.TxPool().SubscribeTransactions(worker.txsCh, true) + + // Subscribe events for blockchain + worker.chainHeadSub = eth.BlockChain().SubscribeChainHeadEvent(worker.chainHeadCh) + + worker.wg.Add(1) + go worker.mainLoop() + + // Submit first work to initialize pending state. + if init { + worker.startCh <- struct{}{} + } + return worker +} + +// getCCC returns a pointer to this worker's CCC instance. +// Only used in tests. +func (w *worker) getCCC() *circuitcapacitychecker.CircuitCapacityChecker { + return w.circuitCapacityChecker +} + +// disablePreseal disables pre-sealing mining feature +func (w *worker) disablePreseal() { + atomic.StoreUint32(&w.noempty, 1) +} + +// enablePreseal enables pre-sealing mining feature +func (w *worker) enablePreseal() { + atomic.StoreUint32(&w.noempty, 0) +} + +// mainLoop is a standalone goroutine to regenerate the sealing task based on the received event. +func (w *worker) mainLoop() { + defer w.wg.Done() + defer w.txsSub.Unsubscribe() + defer w.chainHeadSub.Unsubscribe() + + deadCh := make(chan *pipeline.Result) + pipelineResultCh := func() <-chan *pipeline.Result { + if w.currentPipeline == nil { + return deadCh + } + return w.currentPipeline.ResultCh + } + + for { + select { + case <-w.startCh: + w.startNewPipeline(time.Now().Unix()) + case <-w.chainHeadCh: + w.startNewPipeline(time.Now().Unix()) + case result := <-pipelineResultCh(): + w.handlePipelineResult(result) + case ev := <-w.txsCh: + // Apply transactions to the pending state + // + // Note all transactions received may not be continuous with transactions + // already included in the current mining block. These transactions will + // be automatically eliminated. + if w.currentPipeline != nil { + txs := make(map[common.Address][]*txpool.LazyTransaction) + signer := types.MakeSigner(w.chainConfig, w.currentPipeline.Header.Number, w.currentPipeline.Header.Time) + for _, tx := range ev.Txs { + acc, _ := types.Sender(signer, tx) + txs[acc] = append(txs[acc], &txpool.LazyTransaction{ + Pool: w.eth.TxPool(), // We don't know where this came from, yolo resolve from everywhere + Hash: tx.Hash(), + Tx: nil, // Do *not* set this! We need to resolve it later to pull blobs in + Time: tx.Time(), + GasFeeCap: tx.GasFeeCap(), + GasTipCap: tx.GasTipCap(), + Gas: tx.Gas(), + BlobGas: tx.BlobGas(), + }) + } + txset := newTransactionsByPriceAndNonce(signer, txs, w.currentPipeline.Header.BaseFee) + if result := w.currentPipeline.TryPushTxns(txset, w.onTxFailingInPipeline); result != nil { + w.handlePipelineResult(result) + } + } + w.newTxs.Add(int32(len(ev.Txs))) + + // System stopped + case <-w.exitCh: + return + case <-w.txsSub.Err(): + return + case <-w.chainHeadSub.Err(): + return + } + } +} + +// updateSnapshot updates pending snapshot block and state. +// Note this function assumes the current variable is thread safe. +func (w *worker) updateSnapshot(current *pipeline.BlockCandidate) { + w.snapshotMu.Lock() + defer w.snapshotMu.Unlock() + + w.snapshotBlock = types.NewBlock( + current.Header, + current.Txs, + nil, + current.Receipts, + trie.NewStackTrie(nil), + ) + w.snapshotReceipts = copyReceipts(current.Receipts) + w.snapshotState = current.State.Copy() +} + +func (w *worker) collectPendingL1Messages(startIndex uint64) []types.L1MessageTx { + maxCount := w.chainConfig.Scroll.L1Config.NumL1MessagesPerBlock + return rawdb.ReadL1MessagesFrom(w.eth.ChainDb(), startIndex, maxCount) +} + +// startNewPipeline generates several new sealing tasks based on the parent block. +func (w *worker) startNewPipeline(timestamp int64) { + // Abort if node is still syncing + if w.syncing.Load() { + return + } + + if w.currentPipeline != nil { + w.currentPipeline.Release() + w.currentPipeline = nil + } + + parent := w.chain.CurrentBlock() + + num := parent.Number + header := &types.Header{ + ParentHash: parent.Hash(), + Number: new(big.Int).Add(num, common.Big1), + GasLimit: core.CalcGasLimit(parent.GasLimit, w.config.GasCeil), + Extra: w.extra, + Time: uint64(timestamp), + } + // Set baseFee if we are on an EIP-1559 chain + if w.chainConfig.IsCurie(header.Number) { + state, err := w.chain.StateAt(parent.Root) + if err != nil { + log.Error("Failed to create mining context", "err", err) + return + } + parentL1BaseFee := fees.GetL1BaseFee(state) + header.BaseFee = eip1559.CalcBaseFee(w.chainConfig, parent, parentL1BaseFee) + } + // Only set the coinbase if our consensus engine is running (avoid spurious block rewards) + if w.isRunning() { + if w.coinbase == (common.Address{}) { + log.Error("Refusing to mine without etherbase") + return + } + header.Coinbase = w.coinbase + } + + prepareStart := time.Now() + if err := w.engine.Prepare(w.chain, header); err != nil { + log.Error("Failed to prepare header for mining", "err", err) + return + } + prepareTimer.UpdateSince(prepareStart) + + // If we are care about TheDAO hard-fork check whether to override the extra-data or not + if daoBlock := w.chainConfig.DAOForkBlock; daoBlock != nil { + // Check whether the block is among the fork extra-override range + limit := new(big.Int).Add(daoBlock, params.DAOForkExtraRange) + if header.Number.Cmp(daoBlock) >= 0 && header.Number.Cmp(limit) < 0 { + // Depending whether we support or oppose the fork, override differently + if w.chainConfig.DAOForkSupport { + header.Extra = common.CopyBytes(params.DAOForkBlockExtra) + } else if bytes.Equal(header.Extra, params.DAOForkBlockExtra) { + header.Extra = []byte{} // If miner opposes, don't let it use the reserved extra-data + } + } + } + + parentState, err := w.chain.StateAt(parent.Root) + if err != nil { + log.Error("failed to fetch parent state", "err", err) + return + } + + // Apply special state transition at Curie block + if w.chainConfig.CurieBlock != nil && w.chainConfig.CurieBlock.Cmp(header.Number) == 0 { + misc.ApplyCurieHardFork(parentState) + + var nextL1MsgIndex uint64 + if dbVal := rawdb.ReadFirstQueueIndexNotInL2Block(w.eth.ChainDb(), header.ParentHash); dbVal != nil { + nextL1MsgIndex = *dbVal + } + + // zkEVM requirement: Curie transition block contains 0 transactions, bypass pipeline. + err = w.commit(&pipeline.Result{ + // Note: Signer nodes will not store CCC results for empty blocks in their database. + // In practice, this is acceptable, since this block will never overflow, and follower + // nodes will still store CCC results. + Rows: &types.RowConsumption{}, + FinalBlock: &pipeline.BlockCandidate{ + Header: header, + State: parentState, + Txs: types.Transactions{}, + Receipts: types.Receipts{}, + CoalescedLogs: []*types.Log{}, + NextL1MsgIndex: nextL1MsgIndex, + }, + }) + + if err != nil { + log.Error("failed to commit Curie fork block", "reason", err) + } + + return + } + + // fetch l1Txs + var l1Messages []types.L1MessageTx + if w.chainConfig.Scroll.ShouldIncludeL1Messages() { + common.WithTimer(collectL1MsgsTimer, func() { + l1Messages = w.collectPendingL1Messages(*rawdb.ReadFirstQueueIndexNotInL2Block(w.eth.ChainDb(), parent.Hash())) + }) + } + + tidyPendingStart := time.Now() + // Fill the block with all available pending transactions. + pending := w.eth.TxPool().Pending(false) + // Split the pending transactions into locals and remotes + localTxs, remoteTxs := make(map[common.Address][]*txpool.LazyTransaction), pending + for _, account := range w.eth.TxPool().Locals() { + if txs := remoteTxs[account]; len(txs) > 0 { + delete(remoteTxs, account) + localTxs[account] = txs + } + } + collectL2Timer.UpdateSince(tidyPendingStart) + + var nextL1MsgIndex uint64 + if dbIndex := rawdb.ReadFirstQueueIndexNotInL2Block(w.chain.Database(), parent.Hash()); dbIndex != nil { + nextL1MsgIndex = *dbIndex + } else { + log.Error("failed to read nextL1MsgIndex", "parent", parent.Hash()) + return + } + + w.currentPipelineStart = time.Now() + pipelineCCC := w.getCCC() + if !w.isRunning() { + pipelineCCC = nil + } + w.currentPipeline = pipeline.NewPipeline(w.chain, *w.chain.GetVMConfig(), parentState, header, nextL1MsgIndex, pipelineCCC).WithBeforeTxHook(w.beforeTxHook) + + deadline := time.Unix(int64(header.Time), 0) + if w.chainConfig.Clique != nil && w.chainConfig.Clique.RelaxedPeriod { + // clique with relaxed period uses time.Now() as the header.Time, calculate the deadline + deadline = time.Unix(int64(header.Time+w.chainConfig.Clique.Period), 0) + } + + if err := w.currentPipeline.Start(deadline); err != nil { + log.Error("failed to start pipeline", "err", err) + return + } + + // Short circuit if there is no available pending transactions. + // But if we disable empty precommit already, ignore it. Since + // empty block is necessary to keep the liveness of the network. + if len(localTxs) == 0 && len(remoteTxs) == 0 && len(l1Messages) == 0 && atomic.LoadUint32(&w.noempty) == 0 { + return + } + + if w.chainConfig.Scroll.ShouldIncludeL1Messages() && len(l1Messages) > 0 { + log.Trace("Processing L1 messages for inclusion", "count", len(l1Messages)) + txs, err := newL1MessagesByQueueIndex(w.eth.TxPool(), l1Messages) + if err != nil { + log.Error("Failed to create L1 message set", "l1Messages", l1Messages, "err", err) + return + } + + if result := w.currentPipeline.TryPushTxns(txs, w.onTxFailingInPipeline); result != nil { + w.handlePipelineResult(result) + return + } + } + signer := types.MakeSigner(w.chainConfig, header.Number, header.Time) + + if w.prioritizedTx != nil && w.currentPipeline.Header.Number.Uint64() > w.prioritizedTx.blockNumber { + w.prioritizedTx = nil + } + if w.prioritizedTx != nil { + from, _ := types.Sender(signer, w.prioritizedTx.tx) // error already checked before + txList := map[common.Address][]*txpool.LazyTransaction{from: {txToLazyTx(w.eth.TxPool(), w.prioritizedTx.tx)}} + txs := newTransactionsByPriceAndNonce(signer, txList, header.BaseFee) + if result := w.currentPipeline.TryPushTxns(txs, w.onTxFailingInPipeline); result != nil { + w.handlePipelineResult(result) + return + } + } + + if len(localTxs) > 0 { + txs := newTransactionsByPriceAndNonce(signer, localTxs, header.BaseFee) + if result := w.currentPipeline.TryPushTxns(txs, w.onTxFailingInPipeline); result != nil { + w.handlePipelineResult(result) + return + } + } + if len(remoteTxs) > 0 { + txs := newTransactionsByPriceAndNonce(signer, remoteTxs, header.BaseFee) + if result := w.currentPipeline.TryPushTxns(txs, w.onTxFailingInPipeline); result != nil { + w.handlePipelineResult(result) + return + } + } + + // pipelineCCC was nil, so the block was built for RPC purposes only. Stop the pipeline immediately + // and update the pending block. + if pipelineCCC == nil { + w.currentPipeline.Stop() + } +} + +func (w *worker) handlePipelineResult(res *pipeline.Result) error { + startingHeader := w.currentPipeline.Header + w.currentPipeline.Release() + w.currentPipeline = nil + + if res.FinalBlock != nil { + w.updateSnapshot(res.FinalBlock) + } + + // Rows being nil without an OverflowingTx means that block didn't go thru CCC, + // which means that we are not the sequencer. Do not attempt to commit. + if res.Rows == nil && res.OverflowingTx == nil { + return nil + } + + if res.OverflowingTx != nil { + if res.FinalBlock == nil { + // first txn overflowed the circuit, skip + log.Info("Circuit capacity limit reached for a single tx", "tx", res.OverflowingTx.Hash().String(), + "isL1Message", res.OverflowingTx.IsL1MessageTx(), "reason", res.CCCErr.Error()) + + // Store skipped transaction in local db + overflowingTrace := res.OverflowingTrace + if !w.config.StoreSkippedTxTraces { + overflowingTrace = nil + } + rawdb.WriteSkippedTransaction(w.eth.ChainDb(), res.OverflowingTx, overflowingTrace, res.CCCErr.Error(), + startingHeader.Number.Uint64(), nil) + + if overflowingL1MsgTx := res.OverflowingTx.AsL1MessageTx(); overflowingL1MsgTx != nil { + rawdb.WriteFirstQueueIndexNotInL2Block(w.eth.ChainDb(), startingHeader.ParentHash, overflowingL1MsgTx.QueueIndex+1) + } else { + w.prioritizedTx = nil + w.eth.TxPool().RemoveTx(res.OverflowingTx.Hash(), true, true) + } + } else if !res.OverflowingTx.IsL1MessageTx() { + // prioritize overflowing L2 message as the first txn next block + // no need to prioritize L1 messages, they are fetched in order + // and processed first in every block anyways + w.prioritizedTx = &prioritizedTransaction{ + blockNumber: startingHeader.Number.Uint64() + 1, + tx: res.OverflowingTx, + } + } + + switch { + case res.OverflowingTx.IsL1MessageTx() && + errors.Is(res.CCCErr, circuitcapacitychecker.ErrBlockRowConsumptionOverflow): + l1TxRowConsumptionOverflowCounter.Inc(1) + case !res.OverflowingTx.IsL1MessageTx() && + errors.Is(res.CCCErr, circuitcapacitychecker.ErrBlockRowConsumptionOverflow): + l2TxRowConsumptionOverflowCounter.Inc(1) + case res.OverflowingTx.IsL1MessageTx() && + errors.Is(res.CCCErr, circuitcapacitychecker.ErrUnknown): + l1TxCccUnknownErrCounter.Inc(1) + case !res.OverflowingTx.IsL1MessageTx() && + errors.Is(res.CCCErr, circuitcapacitychecker.ErrUnknown): + l2TxCccUnknownErrCounter.Inc(1) + } + } + + var commitError error + if res.FinalBlock != nil { + if commitError = w.commit(res); commitError == nil { + return nil + } + log.Error("Commit failed", "header", res.FinalBlock.Header, "reason", commitError) + if _, isRetryable := commitError.(retryableCommitError); !isRetryable { + return commitError + } + } + w.startNewPipeline(time.Now().Unix()) + return nil +} + +// retryableCommitError wraps an error that happened during commit phase and indicates that worker can retry to build a new block +type retryableCommitError struct { + inner error +} + +func (e retryableCommitError) Error() string { + return e.inner.Error() +} + +func (e retryableCommitError) Unwrap() error { + return e.inner +} + +// commit runs any post-transaction state modifications, assembles the final block +// and commits new work if consensus engine is running. +func (w *worker) commit(res *pipeline.Result) error { + sealDelay := time.Duration(0) + defer func(t0 time.Time) { + l2CommitTimer.Update(time.Since(t0) - sealDelay) + }(time.Now()) + + if res.CCCErr != nil { + commitReasonCCCCounter.Inc(1) + } else { + commitReasonDeadlineCounter.Inc(1) + } + commitGasCounter.Inc(int64(res.FinalBlock.Header.GasUsed)) + + block, err := w.engine.FinalizeAndAssemble(w.chain, res.FinalBlock.Header, res.FinalBlock.State, + res.FinalBlock.Txs, nil, res.FinalBlock.Receipts, nil) + if err != nil { + return err + } + + sealHash := w.engine.SealHash(block.Header()) + log.Info("Committing new mining work", "number", block.Number(), "sealhash", sealHash, + "txs", res.FinalBlock.Txs.Len(), + "gas", block.GasUsed(), "fees", totalFees(block, res.FinalBlock.Receipts), + "elapsed", common.PrettyDuration(time.Since(w.currentPipelineStart))) + + resultCh, stopCh := make(chan *types.Block), make(chan struct{}) + if err := w.engine.Seal(w.chain, block, resultCh, stopCh); err != nil { + return err + } + // Clique.Seal() will only wait for a second before giving up on us. So make sure there is nothing computational heavy + // or a call that blocks between the call to Seal and the line below. Seal might introduce some delay, so we keep track of + // that artificially added delay and subtract it from overall runtime of commit(). + sealStart := time.Now() + block = <-resultCh + sealDelay = time.Since(sealStart) + if block == nil { + return errors.New("missed seal response from consensus engine") + } + + // verify the generated block with local consensus engine to make sure everything is as expected + if err = w.engine.VerifyHeader(w.chain, block.Header()); err != nil { + return retryableCommitError{inner: err} + } + + blockHash := block.Hash() + var logs []*types.Log + for i, receipt := range res.FinalBlock.Receipts { + // add block location fields + receipt.BlockHash = blockHash + receipt.BlockNumber = block.Number() + receipt.TransactionIndex = uint(i) + + for _, log := range receipt.Logs { + log.BlockHash = blockHash + } + + logs = append(logs, receipt.Logs...) + } + + for _, log := range res.FinalBlock.CoalescedLogs { + log.BlockHash = blockHash + } + + // It's possible that we've stored L1 queue index for this block previously, + // in this case do not overwrite it. + if index := rawdb.ReadFirstQueueIndexNotInL2Block(w.eth.ChainDb(), blockHash); index == nil { + // Store first L1 queue index not processed by this block. + // Note: This accounts for both included and skipped messages. This + // way, if a block only skips messages, we won't reprocess the same + // messages from the next block. + log.Trace( + "Worker WriteFirstQueueIndexNotInL2Block", + "number", block.Number(), + "hash", blockHash.String(), + "nextL1MsgIndex", res.FinalBlock.NextL1MsgIndex, + ) + rawdb.WriteFirstQueueIndexNotInL2Block(w.eth.ChainDb(), blockHash, res.FinalBlock.NextL1MsgIndex) + } else { + log.Trace( + "Worker WriteFirstQueueIndexNotInL2Block: not overwriting existing index", + "number", block.Number(), + "hash", blockHash.String(), + "index", *index, + "nextL1MsgIndex", res.FinalBlock.NextL1MsgIndex, + ) + } + // Store circuit row consumption. + log.Trace( + "Worker write block row consumption", + "id", w.circuitCapacityChecker.ID, + "number", block.Number(), + "hash", blockHash.String(), + "accRows", res.Rows, + ) + + rawdb.WriteBlockRowConsumption(w.eth.ChainDb(), blockHash, res.Rows) + // Commit block and state to database. + _, err = w.chain.WriteBlockAndSetHead(block, res.FinalBlock.Receipts, logs, res.FinalBlock.State, true) + if err != nil { + log.Error("Failed writing block to chain", "err", err) + return err + } + + log.Info("Successfully sealed new block", "number", block.Number(), "sealhash", sealHash, "hash", blockHash) + + // Broadcast the block and announce chain insertion event + w.mux.Post(core.NewMinedBlockEvent{Block: block}) + + return nil +} + +// copyReceipts makes a deep copy of the given receipts. +func copyReceipts(receipts []*types.Receipt) []*types.Receipt { + result := make([]*types.Receipt, len(receipts)) + for i, l := range receipts { + cpy := *l + result[i] = &cpy + } + return result +} + +func (w *worker) onTxFailingInPipeline(txIndex int, tx *types.Transaction, err error) bool { + if !w.isRunning() { + return false + } + + writeTrace := func() { + var trace *types.BlockTrace + var errWithTrace *pipeline.ErrorWithTrace + if w.config.StoreSkippedTxTraces && errors.As(err, &errWithTrace) { + trace = errWithTrace.Trace + } + rawdb.WriteSkippedTransaction(w.eth.ChainDb(), tx, trace, err.Error(), + w.currentPipeline.Header.Number.Uint64(), nil) + } + + switch { + case errors.Is(err, core.ErrGasLimitReached) && tx.IsL1MessageTx(): + // If this block already contains some L1 messages try again in the next block. + if txIndex > 0 { + break + } + // A single L1 message leads to out-of-gas. Skip it. + queueIndex := tx.AsL1MessageTx().QueueIndex + log.Info("Skipping L1 message", "queueIndex", queueIndex, "tx", tx.Hash().String(), "block", + w.currentPipeline.Header.Number, "reason", "gas limit exceeded") + writeTrace() + l1TxGasLimitExceededCounter.Inc(1) + + case errors.Is(err, core.ErrInsufficientFunds): + log.Trace("Skipping tx with insufficient funds", "tx", tx.Hash().String()) + w.eth.TxPool().RemoveTx(tx.Hash(), true, true) + + case errors.Is(err, pipeline.ErrUnexpectedL1MessageIndex): + log.Warn( + "Unexpected L1 message queue index in worker", + "got", tx.AsL1MessageTx().QueueIndex, + ) + case errors.Is(err, core.ErrGasLimitReached), errors.Is(err, core.ErrNonceTooLow), errors.Is(err, core.ErrNonceTooHigh), errors.Is(err, core.ErrTxTypeNotSupported): + break + default: + // Strange error + log.Debug("Transaction failed, account skipped", "hash", tx.Hash().String(), "err", err) + if tx.IsL1MessageTx() { + queueIndex := tx.AsL1MessageTx().QueueIndex + log.Info("Skipping L1 message", "queueIndex", queueIndex, "tx", tx.Hash().String(), "block", + w.currentPipeline.Header.Number, "reason", "strange error", "err", err) + writeTrace() + l1TxStrangeErrCounter.Inc(1) + } + } + return false +} + +// totalFees computes total consumed miner fees in ETH. Block transactions and receipts have to have the same order. +func totalFees(block *types.Block, receipts []*types.Receipt) *big.Float { + feesWei := new(big.Int) + for i, tx := range block.Transactions() { + minerFee, _ := tx.EffectiveGasTip(block.BaseFee()) + feesWei.Add(feesWei, new(big.Int).Mul(new(big.Int).SetUint64(receipts[i].GasUsed), minerFee)) + } + return new(big.Float).Quo(new(big.Float).SetInt(feesWei), new(big.Float).SetInt(big.NewInt(params.Ether))) +} diff --git a/miner/worker.go b/miner/worker.go index 52d2b3178c..b0061645aa 100644 --- a/miner/worker.go +++ b/miner/worker.go @@ -17,34 +17,16 @@ package miner import ( - "errors" - "fmt" "math/big" - "sync" - "sync/atomic" - "time" "github.com/ethereum/go-ethereum/common" - "github.com/ethereum/go-ethereum/consensus" - "github.com/ethereum/go-ethereum/consensus/misc" - "github.com/ethereum/go-ethereum/consensus/misc/eip1559" - "github.com/ethereum/go-ethereum/consensus/misc/eip4844" - "github.com/ethereum/go-ethereum/core" - "github.com/ethereum/go-ethereum/core/rawdb" "github.com/ethereum/go-ethereum/core/state" "github.com/ethereum/go-ethereum/core/txpool" "github.com/ethereum/go-ethereum/core/types" - "github.com/ethereum/go-ethereum/core/vm" - "github.com/ethereum/go-ethereum/event" - "github.com/ethereum/go-ethereum/log" - "github.com/ethereum/go-ethereum/metrics" - "github.com/ethereum/go-ethereum/params" - "github.com/ethereum/go-ethereum/rollup/circuitcapacitychecker" - "github.com/ethereum/go-ethereum/rollup/fees" - "github.com/ethereum/go-ethereum/rollup/tracing" - "github.com/ethereum/go-ethereum/trie" ) +/* + const ( // resultQueueSize is the size of channel listening to sealing result. resultQueueSize = 10 @@ -199,6 +181,7 @@ type newWorkReq struct { interrupt *atomic.Int32 timestamp int64 } +*/ // newPayloadResult is the result of payload generation. type newPayloadResult struct { @@ -208,6 +191,8 @@ type newPayloadResult struct { sidecars []*types.BlobTxSidecar // collected blobs of blob transactions } +/* + // getWorkReq represents a request for getting a new sealing work with provided parameters. type getWorkReq struct { params *generateParams @@ -387,6 +372,8 @@ func (w *worker) getCCC() *circuitcapacitychecker.CircuitCapacityChecker { return w.circuitCapacityChecker } +*/ + // setEtherbase sets the etherbase used to initialize the block coinbase field. func (w *worker) setEtherbase(addr common.Address) { w.mu.Lock() @@ -465,6 +452,8 @@ func (w *worker) close() { w.wg.Wait() } +/* + // recalcRecommit recalculates the resubmitting interval upon feedback. func recalcRecommit(minRecommit, prev time.Duration, target float64, inc bool) time.Duration { var ( @@ -1388,6 +1377,8 @@ func (w *worker) prepareWork(genParams *generateParams) (*environment, error) { return env, nil } +*/ + func txToLazyTx(txPool *txpool.TxPool, tx *types.Transaction) *txpool.LazyTransaction { if tx.IsL1MessageTx() { return &txpool.LazyTransaction{ @@ -1414,6 +1405,8 @@ func txToLazyTx(txPool *txpool.TxPool, tx *types.Transaction) *txpool.LazyTransa } } +/* + // fillTransactions retrieves the pending transactions from the txpool and fills them // into the given sealing block. The transaction selection and ordering strategy can // be customized with the plugin in the future. @@ -1780,3 +1773,5 @@ func signalToErr(signal int32) error { panic(fmt.Errorf("undefined signal %d", signal)) } } + +*/ diff --git a/miner/worker_test.go b/miner/worker_test.go index ab8e455493..4e7935f5d8 100644 --- a/miner/worker_test.go +++ b/miner/worker_test.go @@ -18,7 +18,6 @@ package miner import ( "math/big" - "sync/atomic" "testing" "time" @@ -185,10 +184,10 @@ func TestGenerateAndImportBlock(t *testing.T) { chain, _ := core.NewBlockChain(rawdb.NewMemoryDatabase(), nil, b.genesis, nil, engine, vm.Config{}, nil, nil) defer chain.Stop() - // Ignore empty commit here for less noise. - w.skipSealHook = func(task *task) bool { - return len(task.receipts) == 0 - } + // // Ignore empty commit here for less noise. + // w.skipSealHook = func(task *task) bool { + // return len(task.receipts) == 0 + // } // Wait for mined blocks. sub := w.mux.Subscribe(core.NewMinedBlockEvent{}) @@ -213,6 +212,8 @@ func TestGenerateAndImportBlock(t *testing.T) { } } +/* + func TestEmptyWorkEthash(t *testing.T) { testEmptyWork(t, ethashChainConfig, ethash.NewFaker()) } @@ -322,7 +323,7 @@ func testAdjustInterval(t *testing.T, chainConfig *params.ChainConfig, engine co time.Sleep(time.Second) // Ensure two tasks have been submitted due to start opt start.Store(true) - // w.setRecommitInterval(3 * time.Second) + w.setRecommitInterval(3 * time.Second) select { case <-progress: case <-time.NewTimer(time.Second).C: @@ -343,7 +344,7 @@ func testAdjustInterval(t *testing.T, chainConfig *params.ChainConfig, engine co t.Error("interval reset timeout") } - // w.setRecommitInterval(500 * time.Millisecond) + w.setRecommitInterval(500 * time.Millisecond) select { case <-progress: case <-time.NewTimer(time.Second).C: @@ -505,3 +506,5 @@ func testGetSealingWork(t *testing.T, chainConfig *params.ChainConfig, engine co } } } + +*/ diff --git a/params/config.go b/params/config.go index 5e0947f607..97d698fb7c 100644 --- a/params/config.go +++ b/params/config.go @@ -643,6 +643,8 @@ func (c *EthashConfig) String() string { type CliqueConfig struct { Period uint64 `json:"period"` // Number of seconds between blocks to enforce Epoch uint64 `json:"epoch"` // Epoch length to reset votes and checkpoint + + RelaxedPeriod bool `json:"relaxed_period"` // Relaxes the period to be just an upper bound } // String implements the stringer interface, returning the consensus engine details. diff --git a/rollup/pipeline/pipeline.go b/rollup/pipeline/pipeline.go new file mode 100644 index 0000000000..abf4858a6e --- /dev/null +++ b/rollup/pipeline/pipeline.go @@ -0,0 +1,559 @@ +package pipeline + +import ( + "bytes" + "context" + "errors" + "sync" + "time" + "unsafe" + + "github.com/ethereum/go-ethereum/core" + "github.com/ethereum/go-ethereum/core/state" + "github.com/ethereum/go-ethereum/core/txpool" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/core/vm" + "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/metrics" + "github.com/ethereum/go-ethereum/params" + "github.com/ethereum/go-ethereum/rollup/circuitcapacitychecker" + "github.com/ethereum/go-ethereum/rollup/tracing" +) + +type ErrorWithTrace struct { + Trace *types.BlockTrace + err error +} + +func (e *ErrorWithTrace) Error() string { + return e.err.Error() +} + +func (e *ErrorWithTrace) Unwrap() error { + return e.err +} + +var ( + ErrPipelineDone = errors.New("pipeline is done") + ErrUnexpectedL1MessageIndex = errors.New("unexpected L1 message index") + + lifetimeTimer = func() metrics.Timer { + t := metrics.NewCustomTimer(metrics.NewHistogram(metrics.NewExpDecaySample(128, 0.015)), metrics.NewMeter()) + metrics.DefaultRegistry.Register("pipeline/lifetime", t) + return t + }() + applyTimer = metrics.NewRegisteredTimer("pipeline/apply", nil) + applyIdleTimer = metrics.NewRegisteredTimer("pipeline/apply_idle", nil) + applyStallTimer = metrics.NewRegisteredTimer("pipeline/apply_stall", nil) + encodeTimer = metrics.NewRegisteredTimer("pipeline/encode", nil) + encodeIdleTimer = metrics.NewRegisteredTimer("pipeline/encode_idle", nil) + encodeStallTimer = metrics.NewRegisteredTimer("pipeline/encode_stall", nil) + cccTimer = metrics.NewRegisteredTimer("pipeline/ccc", nil) + cccIdleTimer = metrics.NewRegisteredTimer("pipeline/ccc_idle", nil) +) + +type Pipeline struct { + chain *core.BlockChain + vmConfig vm.Config + parent *types.Block + start time.Time + wg sync.WaitGroup + ctx context.Context + cancelCtx context.CancelFunc + + // accumulators + ccc *circuitcapacitychecker.CircuitCapacityChecker + Header types.Header + state *state.StateDB + nextL1MsgIndex uint64 + blockSize uint64 + txs types.Transactions + coalescedLogs []*types.Log + receipts types.Receipts + gasPool *core.GasPool + + // com channels + txnQueue chan *txpool.LazyTransaction + applyStageRespCh <-chan error + ResultCh <-chan *Result + + // Test hooks + beforeTxHook func() // Method to call before processing a transaction. +} + +func NewPipeline( + chain *core.BlockChain, + vmConfig vm.Config, + state *state.StateDB, + + header *types.Header, + nextL1MsgIndex uint64, + ccc *circuitcapacitychecker.CircuitCapacityChecker, +) *Pipeline { + // make sure we are not sharing a tracer with the caller and not in debug mode + vmConfig.Tracer = nil + + ctx, cancel := context.WithCancel(context.Background()) + return &Pipeline{ + chain: chain, + vmConfig: vmConfig, + parent: chain.GetBlock(header.ParentHash, header.Number.Uint64()-1), + nextL1MsgIndex: nextL1MsgIndex, + Header: *header, + ccc: ccc, + state: state, + gasPool: new(core.GasPool).AddGas(header.GasLimit), + ctx: ctx, + cancelCtx: cancel, + } +} + +func (p *Pipeline) WithBeforeTxHook(beforeTxHook func()) *Pipeline { + p.beforeTxHook = beforeTxHook + return p +} + +func (p *Pipeline) Start(deadline time.Time) error { + p.start = time.Now() + p.txnQueue = make(chan *txpool.LazyTransaction) + applyStageRespCh, applyToEncodeCh, err := p.traceAndApplyStage(p.txnQueue) + if err != nil { + log.Error("Failed starting traceAndApplyStage", "err", err) + return err + } + p.applyStageRespCh = applyStageRespCh + encodeToCccCh := p.encodeStage(applyToEncodeCh) + p.ResultCh = p.cccStage(encodeToCccCh, deadline) + return nil +} + +// Stop forces pipeline to stop its operation and return whatever progress it has so far +func (p *Pipeline) Stop() { + if p.txnQueue != nil { + close(p.txnQueue) + p.txnQueue = nil + } +} + +// orderedTransactionSet represents a set of transactions and some ordering on top of this set. +type orderedTransactionSet interface { + // Peek returns the next transaction. + Peek() *txpool.LazyTransaction + + // Shift removes the next transaction. + Shift() + + // Pop removes all transactions from the current account. + Pop() +} + +func (p *Pipeline) TryPushTxns(txs orderedTransactionSet, onFailingTxn func(txnIndex int, tx *types.Transaction, err error) bool) *Result { + for { + ltx := txs.Peek() + if ltx == nil { + break + } + + result, err := p.TryPushTxn(ltx) + if result != nil { + return result + } + + // TODO: return tx via `TryPushTxn` so that we don't need to resolve it here again + tx := ltx.Resolve() + if tx == nil { + txs.Shift() + continue + } + + switch { + case err == nil, errors.Is(err, core.ErrNonceTooLow): + txs.Shift() + default: + if errors.Is(err, ErrPipelineDone) || onFailingTxn(p.txs.Len(), tx, err) { + p.Stop() + return nil + } + + if tx.IsL1MessageTx() { + txs.Shift() + } else { + txs.Pop() + } + } + } + + return nil +} + +func (p *Pipeline) TryPushTxn(tx *txpool.LazyTransaction) (*Result, error) { + if p.txnQueue == nil { + return nil, ErrPipelineDone + } + + select { + case p.txnQueue <- tx: + case <-p.ctx.Done(): + return nil, ErrPipelineDone + case res := <-p.ResultCh: + return res, nil + } + + select { + case err, valid := <-p.applyStageRespCh: + if !valid { + return nil, ErrPipelineDone + } + return nil, err + case res := <-p.ResultCh: + return res, nil + } +} + +// Release releases all resources related to the pipeline +func (p *Pipeline) Release() { + p.cancelCtx() + p.wg.Wait() +} + +type BlockCandidate struct { + LastTrace *types.BlockTrace + RustTrace unsafe.Pointer + NextL1MsgIndex uint64 + + // accumulated state + Header *types.Header + State *state.StateDB + Txs types.Transactions + Receipts types.Receipts + CoalescedLogs []*types.Log +} + +// sendCancellable tries to send msg to resCh but allows send operation to be cancelled +// by closing cancelCh. Returns true if cancelled. +func sendCancellable[T any, C comparable](resCh chan T, msg T, cancelCh <-chan C) bool { + var zeroC C + + select { + case resCh <- msg: + return false + case cancelSignal := <-cancelCh: + if cancelSignal != zeroC { + panic("shouldn't have happened") + } + return true + } +} + +func (p *Pipeline) traceAndApplyStage(txsIn <-chan *txpool.LazyTransaction) (<-chan error, <-chan *BlockCandidate, error) { + p.state.StartPrefetcher("miner") + downstreamCh := make(chan *BlockCandidate, p.downstreamChCapacity()) + resCh := make(chan error) + p.wg.Add(1) + go func() { + defer func() { + close(downstreamCh) + close(resCh) + p.state.StopPrefetcher() + p.wg.Done() + }() + + var ltx *txpool.LazyTransaction + for { + idleStart := time.Now() + select { + case ltx = <-txsIn: + if ltx == nil { + return + } + case <-p.ctx.Done(): + return + } + applyIdleTimer.UpdateSince(idleStart) + + applyStart := time.Now() + + // If we don't have enough gas for any further transactions then we're done + if p.gasPool.Gas() < params.TxGas { + return + } + + // If we have collected enough transactions then we're done + // Originally we only limit l2txs count, but now strictly limit total txs number. + if !p.chain.Config().Scroll.IsValidTxCount(p.txs.Len() + 1) { + return + } + + if p.gasPool.Gas() < ltx.Gas { + // we don't have enough space for the next transaction, skip the account and continue looking for more txns + sendCancellable(resCh, core.ErrGasLimitReached, p.ctx.Done()) + continue + } + + // TODO: blob gas check + + tx := ltx.Resolve() + if tx == nil { + log.Trace("Ignoring evicted transaction", "hash", ltx.Hash) + // can't resolve the tx, silently ignore and continue looking for more txns + sendCancellable(resCh, errors.New("cannot resolve evicted tx"), p.ctx.Done()) + continue + } + + if tx.IsL1MessageTx() && tx.AsL1MessageTx().QueueIndex != p.nextL1MsgIndex { + // Continue, we might still be able to include some L2 messages + sendCancellable(resCh, ErrUnexpectedL1MessageIndex, p.ctx.Done()) + continue + } + + if !tx.IsL1MessageTx() && !p.chain.Config().Scroll.IsValidBlockSize(p.blockSize+tx.Size()) { + // can't fit this txn in this block, silently ignore and continue looking for more txns + sendCancellable(resCh, nil, p.ctx.Done()) + continue + } + + // Start executing the transaction + p.state.SetTxContext(tx.Hash(), p.txs.Len()) + receipt, trace, err := p.traceAndApply(tx) + + if p.txs.Len() == 0 && tx.IsL1MessageTx() && err != nil { + // L1 message errored as the first txn, skip + p.nextL1MsgIndex = tx.AsL1MessageTx().QueueIndex + 1 + } + + if err == nil { + // Everything ok, collect the logs and shift in the next transaction from the same account + p.coalescedLogs = append(p.coalescedLogs, receipt.Logs...) + p.txs = append(p.txs, tx) + p.receipts = append(p.receipts, receipt) + + if !tx.IsL1MessageTx() { + // only consider block size limit for L2 transactions + p.blockSize += tx.Size() + } else { + p.nextL1MsgIndex = tx.AsL1MessageTx().QueueIndex + 1 + } + + stallStart := time.Now() + if sendCancellable(downstreamCh, &BlockCandidate{ + NextL1MsgIndex: p.nextL1MsgIndex, + LastTrace: trace, + + Header: types.CopyHeader(&p.Header), + State: p.state.Copy(), + Txs: p.txs, + Receipts: p.receipts, + CoalescedLogs: p.coalescedLogs, + }, p.ctx.Done()) { + // next stage terminated and caller terminated us as well + return + } + applyStallTimer.UpdateSince(stallStart) + } + if err != nil && trace != nil { + err = &ErrorWithTrace{ + Trace: trace, + err: err, + } + } + applyTimer.UpdateSince(applyStart) + sendCancellable(resCh, err, p.ctx.Done()) + } + }() + return resCh, downstreamCh, nil +} + +type Result struct { + OverflowingTx *types.Transaction + OverflowingTrace *types.BlockTrace + CCCErr error + + Rows *types.RowConsumption + FinalBlock *BlockCandidate +} + +func (p *Pipeline) encodeStage(traces <-chan *BlockCandidate) <-chan *BlockCandidate { + downstreamCh := make(chan *BlockCandidate, p.downstreamChCapacity()) + p.wg.Add(1) + + go func() { + defer func() { + close(downstreamCh) + p.wg.Done() + }() + buffer := new(bytes.Buffer) + for { + idleStart := time.Now() + select { + case trace := <-traces: + if trace == nil { + return + } + encodeIdleTimer.UpdateSince(idleStart) + + encodeStart := time.Now() + if p.ccc != nil { + trace.RustTrace = circuitcapacitychecker.MakeRustTrace(trace.LastTrace, buffer) + if trace.RustTrace == nil { + log.Error("making rust trace", "txHash", trace.LastTrace.Transactions[0].TxHash) + return + } + } + encodeTimer.UpdateSince(encodeStart) + + stallStart := time.Now() + if sendCancellable(downstreamCh, trace, p.ctx.Done()) && trace.RustTrace != nil { + // failed to send the trace downstream, free it here. + circuitcapacitychecker.FreeRustTrace(trace.RustTrace) + } + encodeStallTimer.UpdateSince(stallStart) + case <-p.ctx.Done(): + return + } + + } + }() + return downstreamCh +} + +func (p *Pipeline) cccStage(candidates <-chan *BlockCandidate, deadline time.Time) <-chan *Result { + if p.ccc != nil { + p.ccc.Reset() + } + resultCh := make(chan *Result, 1) + var lastCandidate *BlockCandidate + var lastAccRows *types.RowConsumption + var deadlineReached bool + + p.wg.Add(1) + go func() { + deadlineTimer := time.NewTimer(time.Until(deadline)) + defer func() { + close(resultCh) + deadlineTimer.Stop() + lifetimeTimer.UpdateSince(p.start) + // consume candidates and free all rust traces + for candidate := range candidates { + if candidate == nil { + break + } + if candidate.RustTrace != nil { + circuitcapacitychecker.FreeRustTrace(candidate.RustTrace) + } + } + p.wg.Done() + }() + for { + idleStart := time.Now() + select { + case <-p.ctx.Done(): + return + case <-deadlineTimer.C: + cccIdleTimer.UpdateSince(idleStart) + // note: currently we don't allow empty blocks, but if we ever do; make sure to CCC check it first + if lastCandidate != nil { + resultCh <- &Result{ + Rows: lastAccRows, + FinalBlock: lastCandidate, + } + return + } + deadlineReached = true + case candidate := <-candidates: + cccIdleTimer.UpdateSince(idleStart) + cccStart := time.Now() + var accRows *types.RowConsumption + var err error + if candidate != nil && p.ccc != nil { + accRows, err = p.ccc.ApplyTransactionRustTrace(candidate.RustTrace) + lastTxn := candidate.Txs[candidate.Txs.Len()-1] + cccTimer.UpdateSince(cccStart) + if err != nil { + resultCh <- &Result{ + OverflowingTx: lastTxn, + OverflowingTrace: candidate.LastTrace, + CCCErr: err, + Rows: lastAccRows, + FinalBlock: lastCandidate, + } + return + } + + lastCandidate = candidate + lastAccRows = accRows + } else if candidate != nil && p.ccc == nil { + lastCandidate = candidate + } + + // immediately close the block if deadline reached or apply stage is done + if candidate == nil || deadlineReached { + resultCh <- &Result{ + Rows: lastAccRows, + FinalBlock: lastCandidate, + } + return + } + } + } + }() + return resultCh +} + +func (p *Pipeline) traceAndApply(tx *types.Transaction) (*types.Receipt, *types.BlockTrace, error) { + var trace *types.BlockTrace + var err error + + if p.beforeTxHook != nil { + p.beforeTxHook() + } + + // do gas limit check up-front and do not run CCC if it fails + if p.gasPool.Gas() < tx.Gas() { + return nil, nil, core.ErrGasLimitReached + } + + if p.ccc != nil { + // don't commit the state during tracing for circuit capacity checker, otherwise we cannot revert. + // and even if we don't commit the state, the `refund` value will still be correct, as explained in `CommitTransaction` + finaliseStateAfterApply := false + snap := p.state.Snapshot() + + // 1. we have to check circuit capacity before `core.ApplyTransaction`, + // because if the tx can be successfully executed but circuit capacity overflows, it will be inconvenient to revert. + // 2. even if we don't commit to the state during the tracing (which means `clearJournalAndRefund` is not called during the tracing), + // the `refund` value will still be correct, because: + // 2.1 when starting handling the first tx, `state.refund` is 0 by default, + // 2.2 after tracing, the state is either committed in `core.ApplyTransaction`, or reverted, so the `state.refund` can be cleared, + // 2.3 when starting handling the following txs, `state.refund` comes as 0 + trace, err = tracing.NewTracerWrapper().CreateTraceEnvAndGetBlockTrace(p.chain.Config(), p.chain, p.chain.Engine(), p.chain.Database(), + p.state, p.parent.Header(), types.NewBlockWithHeader(&p.Header).WithBody([]*types.Transaction{tx}, nil), finaliseStateAfterApply) + // `w.current.traceEnv.State` & `w.current.state` share a same pointer to the state, so only need to revert `w.current.state` + // revert to snapshot for calling `core.ApplyMessage` again, (both `traceEnv.GetBlockTrace` & `core.ApplyTransaction` will call `core.ApplyMessage`) + p.state.RevertToSnapshot(snap) + if err != nil { + return nil, nil, err + } + } + + // create new snapshot for `core.ApplyTransaction` + snap := p.state.Snapshot() + + var receipt *types.Receipt + receipt, err = core.ApplyTransaction(p.chain.Config(), p.chain, nil /* coinbase will default to chainConfig.Scroll.FeeVaultAddress */, p.gasPool, + p.state, &p.Header, tx, &p.Header.GasUsed, p.vmConfig) + if err != nil { + p.state.RevertToSnapshot(snap) + return nil, trace, err + } + return receipt, trace, nil +} + +// downstreamChCapacity returns the channel capacity that should be used for downstream channels. +// It aims to minimize stalls caused by different computational costs of different transactions +func (p *Pipeline) downstreamChCapacity() int { + cap := 1 + if p.chain.Config().Scroll.MaxTxPerBlock != nil { + cap = *p.chain.Config().Scroll.MaxTxPerBlock + } + return cap +}