From 17acae881a415b2e09c540ec5175188f51cb9ee4 Mon Sep 17 00:00:00 2001 From: HAOYUatHZ <37070449+HAOYUatHZ@users.noreply.github.com> Date: Thu, 1 Aug 2024 14:19:44 +0800 Subject: [PATCH] defer txpool reorg until worker fetches txns for the next block (#944) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: defer txpool reorg until worker fetches txns for the next block (#905) * fix --------- Co-authored-by: Ă–mer Faruk Irmak --- core/txpool/blobpool/blobpool.go | 7 +++++++ core/txpool/legacypool/legacypool.go | 24 +++++++++++++++++++++++- core/txpool/subpool.go | 3 +++ core/txpool/txpool.go | 12 ++++++++++++ miner/scroll_worker.go | 7 +++++++ 5 files changed, 52 insertions(+), 1 deletion(-) diff --git a/core/txpool/blobpool/blobpool.go b/core/txpool/blobpool/blobpool.go index 84c62b3401..6e64adcec0 100644 --- a/core/txpool/blobpool/blobpool.go +++ b/core/txpool/blobpool/blobpool.go @@ -1576,3 +1576,10 @@ func (p *BlobPool) Status(hash common.Hash) txpool.TxStatus { func (p *BlobPool) RemoveTx(hash common.Hash, outofbound bool, unreserve bool) int { return 0 } + +func (p *BlobPool) PauseReorgs() { + log.Debug("skip BlobPool `PauseReorgs`") +} +func (p *BlobPool) ResumeReorgs() { + log.Debug("skip BlobPool `ResumeReorgs`") +} diff --git a/core/txpool/legacypool/legacypool.go b/core/txpool/legacypool/legacypool.go index 97f820ed75..9bd7aa4116 100644 --- a/core/txpool/legacypool/legacypool.go +++ b/core/txpool/legacypool/legacypool.go @@ -227,6 +227,7 @@ type LegacyPool struct { queueTxEventCh chan *types.Transaction reorgDoneCh chan chan struct{} reorgShutdownCh chan struct{} // requests shutdown of scheduleReorgLoop + reorgPauseCh chan bool // requests to pause scheduleReorgLoop wg sync.WaitGroup // tracks loop, scheduleReorgLoop initDoneCh chan struct{} // is closed once the pool is initialized (for tests) @@ -258,6 +259,7 @@ func New(config Config, chain BlockChain) *LegacyPool { queueTxEventCh: make(chan *types.Transaction), reorgDoneCh: make(chan chan struct{}), reorgShutdownCh: make(chan struct{}), + reorgPauseCh: make(chan bool), initDoneCh: make(chan struct{}), } pool.locals = newAccountSet(pool.signer) @@ -1198,13 +1200,14 @@ func (pool *LegacyPool) scheduleReorgLoop() { curDone chan struct{} // non-nil while runReorg is active nextDone = make(chan struct{}) launchNextRun bool + reorgsPaused bool reset *txpoolResetRequest dirtyAccounts *accountSet queuedEvents = make(map[common.Address]*sortedMap) ) for { // Launch next background reorg if needed - if curDone == nil && launchNextRun { + if curDone == nil && launchNextRun && !reorgsPaused { // Run the background reorg and announcements go pool.runReorg(nextDone, reset, dirtyAccounts, queuedEvents) @@ -1256,6 +1259,7 @@ func (pool *LegacyPool) scheduleReorgLoop() { } close(nextDone) return + case reorgsPaused = <-pool.reorgPauseCh: } } } @@ -1705,6 +1709,24 @@ func (pool *LegacyPool) demoteUnexecutables() { } } +// PauseReorgs stops any new reorg jobs to be started but doesn't interrupt any existing ones that are in flight +// Keep in mind this function might block, although it is not expected to block for any significant amount of time +func (pool *LegacyPool) PauseReorgs() { + select { + case pool.reorgPauseCh <- true: + case <-pool.reorgShutdownCh: + } +} + +// ResumeReorgs allows new reorg jobs to be started. +// Keep in mind this function might block, although it is not expected to block for any significant amount of time +func (pool *LegacyPool) ResumeReorgs() { + select { + case pool.reorgPauseCh <- false: + case <-pool.reorgShutdownCh: + } +} + // addressByHeartbeat is an account address tagged with its last activity timestamp. type addressByHeartbeat struct { address common.Address diff --git a/core/txpool/subpool.go b/core/txpool/subpool.go index 4843b97e87..9fd5e40d30 100644 --- a/core/txpool/subpool.go +++ b/core/txpool/subpool.go @@ -142,4 +142,7 @@ type SubPool interface { // RemoveTx removes a transaction from the pool, returning the number of transactions removed. RemoveTx(hash common.Hash, outofbound bool, unreserve bool) int + + PauseReorgs() + ResumeReorgs() } diff --git a/core/txpool/txpool.go b/core/txpool/txpool.go index b132b8b98d..04e9d6b088 100644 --- a/core/txpool/txpool.go +++ b/core/txpool/txpool.go @@ -434,3 +434,15 @@ func (p *TxPool) RemoveTx(hash common.Hash, outofbound bool, unreserve bool) int } return ret } + +func (pool *TxPool) PauseReorgs() { + for _, subpool := range pool.subpools { + subpool.PauseReorgs() + } +} + +func (pool *TxPool) ResumeReorgs() { + for _, subpool := range pool.subpools { + subpool.ResumeReorgs() + } +} diff --git a/miner/scroll_worker.go b/miner/scroll_worker.go index c7e6249ccd..9eac3a0167 100644 --- a/miner/scroll_worker.go +++ b/miner/scroll_worker.go @@ -400,6 +400,9 @@ func (w *worker) startNewPipeline(timestamp int64) { } collectL2Timer.UpdateSince(tidyPendingStart) + // Allow txpool to be reorged as we build current block + w.eth.TxPool().ResumeReorgs() + var nextL1MsgIndex uint64 if dbIndex := rawdb.ReadFirstQueueIndexNotInL2Block(w.chain.Database(), parent.Hash()); dbIndex != nil { nextL1MsgIndex = *dbIndex @@ -668,6 +671,10 @@ func (w *worker) commit(res *pipeline.Result) error { "accRows", res.Rows, ) + // A new block event will trigger a reorg in the txpool, pause reorgs to defer this until we fetch txns for next block. + // We may end up trying to process txns that we already included in the previous block, but they will all fail the nonce check + w.eth.TxPool().PauseReorgs() + 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)