From c5055a88b3ad02b206ed906b470ad7dc2adaa9f6 Mon Sep 17 00:00:00 2001 From: Jared Wasinger Date: Fri, 4 Apr 2025 12:35:48 +0200 Subject: [PATCH] attempt at making txpool Add sync param work --- core/txpool/legacypool/legacypool.go | 76 +++++++++++++++++++++++----- core/txpool/txpool.go | 1 + eth/catalyst/api_test.go | 1 + tests/testdata | 2 +- 4 files changed, 66 insertions(+), 14 deletions(-) diff --git a/core/txpool/legacypool/legacypool.go b/core/txpool/legacypool/legacypool.go index 278ad0791f..2871c600d6 100644 --- a/core/txpool/legacypool/legacypool.go +++ b/core/txpool/legacypool/legacypool.go @@ -19,6 +19,7 @@ package legacypool import ( "errors" + "fmt" "maps" "math" "math/big" @@ -246,7 +247,7 @@ type LegacyPool struct { priced *pricedList // All transactions sorted by price reqResetCh chan *txpoolResetRequest - reqPromoteCh chan *accountSet + reqPromoteCh chan *reqPromote queueTxEventCh chan *types.Transaction reorgDoneCh chan chan struct{} reorgShutdownCh chan struct{} // requests shutdown of scheduleReorgLoop @@ -258,6 +259,7 @@ type LegacyPool struct { type txpoolResetRequest struct { oldHead, newHead *types.Header + done chan struct{} } // New creates a new transaction pool to gather, sort and filter inbound @@ -277,7 +279,7 @@ func New(config Config, chain BlockChain) *LegacyPool { beats: make(map[common.Address]time.Time), all: newLookup(), reqResetCh: make(chan *txpoolResetRequest), - reqPromoteCh: make(chan *accountSet), + reqPromoteCh: make(chan *reqPromote), queueTxEventCh: make(chan *types.Transaction), reorgDoneCh: make(chan chan struct{}), reorgShutdownCh: make(chan struct{}), @@ -962,6 +964,9 @@ func (pool *LegacyPool) Add(txs []*types.Transaction, sync bool) []error { pool.mu.Lock() newErrs, dirtyAddrs := pool.addTxsLocked(news) pool.mu.Unlock() + for _, new := range news { + fmt.Printf("prepare tx for queue %x\n", new.Hash()) + } var nilSlot = 0 for _, err := range newErrs { @@ -972,9 +977,11 @@ func (pool *LegacyPool) Add(txs []*types.Transaction, sync bool) []error { nilSlot++ } // Reorg the pool internals if needed and return - done := pool.requestPromoteExecutables(dirtyAddrs) + req := &reqPromote{dirtyAddrs, make(chan struct{})} + pool.requestPromoteExecutables(req) if sync { - <-done + <-req.done + fmt.Printf("done adding txs %v\n", txs[0]) } return errs } @@ -1140,9 +1147,10 @@ func (pool *LegacyPool) removeTx(hash common.Hash, outofbound bool, unreserve bo // requestReset requests a pool reset to the new head block. // The returned channel is closed when the reset has occurred. func (pool *LegacyPool) requestReset(oldHead *types.Header, newHead *types.Header) chan struct{} { + req := &txpoolResetRequest{oldHead, newHead, make(chan struct{})} select { - case pool.reqResetCh <- &txpoolResetRequest{oldHead, newHead}: - return <-pool.reorgDoneCh + case pool.reqResetCh <- req: + return req.done case <-pool.reorgShutdownCh: return pool.reorgShutdownCh } @@ -1150,10 +1158,10 @@ func (pool *LegacyPool) requestReset(oldHead *types.Header, newHead *types.Heade // requestPromoteExecutables requests transaction promotion checks for the given addresses. // The returned channel is closed when the promotion checks have occurred. -func (pool *LegacyPool) requestPromoteExecutables(set *accountSet) chan struct{} { +func (pool *LegacyPool) requestPromoteExecutables(req *reqPromote) chan struct{} { select { - case pool.reqPromoteCh <- set: - return <-pool.reorgDoneCh + case pool.reqPromoteCh <- req: + return req.done case <-pool.reorgShutdownCh: return pool.reorgShutdownCh } @@ -1174,6 +1182,8 @@ func (pool *LegacyPool) scheduleReorgLoop() { defer pool.wg.Done() var ( + curDones []chan struct{} + nextDones []chan struct{} curDone chan struct{} // non-nil while runReorg is active nextDone = make(chan struct{}) launchNextRun bool @@ -1184,11 +1194,15 @@ func (pool *LegacyPool) scheduleReorgLoop() { for { // Launch next background reorg if needed if curDone == nil && launchNextRun { + fmt.Printf("initiating reorg with dones %v\n", nextDones) + fmt.Printf("initiating reorg with dirty accounts: %v\n", dirtyAccounts) + // Run the background reorg and announcements go pool.runReorg(nextDone, reset, dirtyAccounts, queuedEvents) // Prepare everything for the next round of reorg curDone, nextDone = nextDone, make(chan struct{}) + curDones, nextDones = nextDones, []chan struct{}{} launchNextRun = false reset, dirtyAccounts = nil, nil @@ -1203,20 +1217,37 @@ func (pool *LegacyPool) scheduleReorgLoop() { } else { reset.newHead = req.newHead } + fmt.Printf("reset: appending to next dones %v\n", req.done) + nextDones = append(nextDones, req.done) launchNextRun = true - pool.reorgDoneCh <- nextDone + //pool.reorgDoneCh <- nextDone case req := <-pool.reqPromoteCh: + select { + case tx := <-pool.queueTxEventCh: + fmt.Printf("queue tx %x\n", tx.Hash()) + // Queue up the event, but don't schedule a reorg. It's up to the caller to + // request one later if they want the events sent. + addr, _ := types.Sender(pool.signer, tx) + if _, ok := queuedEvents[addr]; !ok { + queuedEvents[addr] = NewSortedMap() + } + queuedEvents[addr].Put(tx) + default: + } // Promote request: update address set if request is already pending. if dirtyAccounts == nil { - dirtyAccounts = req + dirtyAccounts = req.set } else { - dirtyAccounts.merge(req) + dirtyAccounts.merge(req.set) } + fmt.Printf("promote executables: appending to next dones %v\n", req.done) + nextDones = append(nextDones, req.done) launchNextRun = true - pool.reorgDoneCh <- nextDone + //pool.reorgDoneCh <- nextDone case tx := <-pool.queueTxEventCh: + fmt.Printf("queue tx %x\n", tx.Hash()) // Queue up the event, but don't schedule a reorg. It's up to the caller to // request one later if they want the events sent. addr, _ := types.Sender(pool.signer, tx) @@ -1226,12 +1257,24 @@ func (pool *LegacyPool) scheduleReorgLoop() { queuedEvents[addr].Put(tx) case <-curDone: + fmt.Println("reorg done") + for _, ch := range curDones { + fmt.Println(ch) + close(ch) + } + fmt.Println("end reorg done") curDone = nil case <-pool.reorgShutdownCh: // Wait for current run to finish. if curDone != nil { + fmt.Println("shutdown reorg done") <-curDone + for _, ch := range curDones { + fmt.Println(ch) + close(ch) + } + fmt.Println("end reorg done") } close(nextDone) return @@ -1309,6 +1352,7 @@ func (pool *LegacyPool) runReorg(done chan struct{}, reset *txpoolResetRequest, if _, ok := events[addr]; !ok { events[addr] = NewSortedMap() } + fmt.Printf("promoted tx %v\n", tx.Hash()) events[addr].Put(tx) } if len(events) > 0 { @@ -1429,6 +1473,7 @@ func (pool *LegacyPool) promoteExecutables(accounts []common.Address) []*types.T // Iterate over all accounts and promote any executable transactions gasLimit := pool.currentHead.Load().GasLimit for _, addr := range accounts { + fmt.Printf("promoteExecutables %x\n", addr) list := pool.queue[addr] if list == nil { continue // Just in case someone calls with a non existing account @@ -1687,6 +1732,11 @@ type accountSet struct { cache []common.Address } +type reqPromote struct { + set *accountSet + done chan struct{} +} + // newAccountSet creates a new address set with an associated signer for sender // derivations. func newAccountSet(signer types.Signer, addrs ...common.Address) *accountSet { diff --git a/core/txpool/txpool.go b/core/txpool/txpool.go index 47d83e03d4..c1c8078094 100644 --- a/core/txpool/txpool.go +++ b/core/txpool/txpool.go @@ -361,6 +361,7 @@ func (p *TxPool) Add(txs []*types.Transaction, sync bool) []error { splits := make([]int, len(txs)) for i, tx := range txs { + fmt.Printf("txpool attempting to add %v\n", tx.Hash()) // Mark this transaction belonging to no-subpool splits[i] = -1 diff --git a/eth/catalyst/api_test.go b/eth/catalyst/api_test.go index e91e07d05d..8c0e8290bf 100644 --- a/eth/catalyst/api_test.go +++ b/eth/catalyst/api_test.go @@ -119,6 +119,7 @@ func TestEth2AssembleBlock(t *testing.T) { t.Fatalf("error signing transaction, err=%v", err) } ethservice.TxPool().Add([]*types.Transaction{tx}, true) + fmt.Println("in test: added txs") blockParams := engine.PayloadAttributes{ Timestamp: blocks[9].Time() + 5, } diff --git a/tests/testdata b/tests/testdata index 81862e4848..faf33b4714 160000 --- a/tests/testdata +++ b/tests/testdata @@ -1 +1 @@ -Subproject commit 81862e4848585a438d64f911a19b3825f0f4cd95 +Subproject commit faf33b471465d3c6cdc3d04fbd690895f78d33f2