From 3d40d32e41d8e552c1cb97f0d5c0142102265829 Mon Sep 17 00:00:00 2001 From: lightclient Date: Mon, 22 Sep 2025 15:36:55 -0600 Subject: [PATCH] core/txpool/blobpool: swap converted tx atomically --- core/txpool/blobpool/blobpool.go | 152 +++++++++++++++++-------------- 1 file changed, 85 insertions(+), 67 deletions(-) diff --git a/core/txpool/blobpool/blobpool.go b/core/txpool/blobpool/blobpool.go index a4b9722587..c9f9595bbe 100644 --- a/core/txpool/blobpool/blobpool.go +++ b/core/txpool/blobpool/blobpool.go @@ -477,7 +477,7 @@ func (p *BlobPool) Init(gasTip uint64, head *types.Header, reserver txpool.Reser // Since the user might have modified their pool's capacity, evict anything // above the current allowance for p.stored > p.config.Datacap { - p.drop(true) + p.drop() } // Update the metrics and return the constructed pool datacapGauge.Update(int64(p.config.Datacap)) @@ -893,22 +893,21 @@ func (p *BlobPool) Reset(oldHead, newHead *types.Header) { // Evicts all transactions from the pool. The pool should, at this stage, // contain only legacy blob transactions. Evicting any new-format blob // transactions is also safe, as they will be re-injected later. - txs := make(map[common.Address]map[uint64]uint64) - for p.stored > 0 { - sender, nonce, id := p.drop(false) - if txs[sender] == nil { - txs[sender] = make(map[uint64]uint64) + all := make(map[common.Address]map[uint64]*blobTxMeta) + for sender, txs := range p.index { + all[sender] = make(map[uint64]*blobTxMeta) + for _, m := range txs { + all[sender][m.nonce] = m } - txs[sender][nonce] = id } // Initiate the background conversion thread, the mutex is not required // and won't block any pool operation. - go p.convertLegacySidecars(txs) + go p.convertLegacySidecars(all) } } // convertLegacySidecars converts all given transactions to sidecar version 1. -func (p *BlobPool) convertLegacySidecars(txs map[common.Address]map[uint64]uint64) { +func (p *BlobPool) convertLegacySidecars(txs map[common.Address]map[uint64]*blobTxMeta) { var ( start = time.Now() success int @@ -922,35 +921,52 @@ func (p *BlobPool) convertLegacySidecars(txs map[common.Address]map[uint64]uint6 // Convert and insert the txs in order for _, nonce := range nonces { - id := list[nonce] - blob, err := p.store.Get(id) + m := list[nonce] + p.lock.RLock() + blob, err := p.store.Get(m.id) + p.lock.RUnlock() if err != nil { fail++ continue } - // Remove the blob transaction regardless of conversion succeeds or not - if err := p.store.Delete(id); err != nil { - log.Error("Failed to delete blob transaction", "from", addr, "id", id, "err", err) + + // Use this to remove the blob transaction regardless of if the conversion + // succeeds or not. + delete := func() { + p.lock.Lock() + defer p.lock.Unlock() + if err := p.store.Delete(m.id); err != nil { + log.Error("Failed to delete blob transaction", "from", addr, "id", m.id, "err", err) + } } + + // Decode the transaction and perform the migration. var tx types.Transaction if err = rlp.DecodeBytes(blob, &tx); err != nil { - log.Error("Blob transaction is corrupted", "id", id, "err", err) + log.Error("Blob transaction is corrupted", "id", m.id, "err", err) + delete() fail++ continue } if err := tx.BlobTxSidecar().ToV1(); err != nil { log.Error("Failed to convert blob transaction", "hash", tx.Hash(), "err", err) + delete() fail++ continue } - errs := p.Add([]*types.Transaction{&tx}, true) - if errs[0] != nil { + + // Now atomically swap the old sidecar with the converted sidecar. + p.lock.Lock() + p.remove(m, addr) + err = p.add(&tx) + p.lock.Unlock() + if err != nil { fail++ - log.Error("Failed to re-inject the tx", "hash", tx.Hash(), "err", errs[0]) - } else { - success++ - log.Debug("Reinjected the converted tx", "hash", tx.Hash(), "err", errs[0]) + log.Error("Failed to re-inject the tx", "hash", tx.Hash(), "err", err) + continue } + success++ + log.Debug("Reinjected the converted tx", "hash", tx.Hash(), "err", err) } } @@ -1593,10 +1609,20 @@ func (p *BlobPool) Add(txs []*types.Transaction, sync bool) []error { if errs[i] != nil { continue } - errs[i] = p.add(tx) - if errs[i] == nil { - adds = append(adds, tx.WithoutBlobTxSidecar()) - } + func() { + // The blob pool blocks on adding a transaction. This is because blob txs are + // only even pulled from the network, so this method will act as the overload + // protection for fetches. + waitStart := time.Now() + p.lock.Lock() + addwaitHist.Update(time.Since(waitStart).Nanoseconds()) + defer p.lock.Unlock() + defer func(start time.Time) { addtimeHist.Update(time.Since(start).Nanoseconds()) }(time.Now()) + errs[i] = p.add(tx) + if errs[i] == nil { + adds = append(adds, tx.WithoutBlobTxSidecar()) + } + }() } if len(adds) > 0 { p.discoverFeed.Send(core.NewTxsEvent{Txs: adds}) @@ -1606,20 +1632,9 @@ func (p *BlobPool) Add(txs []*types.Transaction, sync bool) []error { } // add inserts a new blob transaction into the pool if it passes validation (both -// consensus validity and pool restrictions). +// consensus validity and pool restrictions). It assume the pool lock is already +// held by the caller. func (p *BlobPool) add(tx *types.Transaction) (err error) { - // The blob pool blocks on adding a transaction. This is because blob txs are - // only even pulled from the network, so this method will act as the overload - // protection for fetches. - waitStart := time.Now() - p.lock.Lock() - addwaitHist.Update(time.Since(waitStart).Nanoseconds()) - defer p.lock.Unlock() - - defer func(start time.Time) { - addtimeHist.Update(time.Since(start).Nanoseconds()) - }(time.Now()) - // Ensure the transaction is valid from all perspectives if err := p.validateTx(tx); err != nil { log.Trace("Transaction validation failed", "hash", tx.Hash(), "err", err) @@ -1764,7 +1779,7 @@ func (p *BlobPool) add(tx *types.Transaction) (err error) { // If the pool went over the allowed data limit, evict transactions until // we're again below the threshold for p.stored > p.config.Datacap { - p.drop(true) + p.drop() } p.updateStorageMetrics() @@ -1772,60 +1787,63 @@ func (p *BlobPool) add(tx *types.Transaction) (err error) { return nil } -// drop removes the worst transaction from the pool. It is primarily used when a -// freshly added transaction overflows the pool and needs to evict something. The -// method is also called on startup if the user resizes their storage, might be an -// expensive run but it should be fine-ish. -func (p *BlobPool) drop(deleteTxBlob bool) (common.Address, uint64, uint64) { - // Peek at the account with the worse transaction set to evict from (Go's heap - // stores the minimum at index zero of the heap slice) and retrieve it's last - // transaction. +// remove will remove the specified transaction from the pool and delete it from +// the store. It expects the caller to hold the pool lock. +func (p *BlobPool) remove(tx *blobTxMeta, from common.Address) { + // Remove the transaction from the pool's index var ( - from = p.evict.addrs[0] // cannot call drop on empty pool - txs = p.index[from] - drop = txs[len(txs)-1] last = len(txs) == 1 ) - // Remove the transaction from the pool's index - if last { + if len(p.index[from]) == 1 { delete(p.index, from) delete(p.spent, from) p.reserver.Release(from) } else { txs[len(txs)-1] = nil txs = txs[:len(txs)-1] - p.index[from] = txs - p.spent[from] = new(uint256.Int).Sub(p.spent[from], drop.costCap) + p.spent[from] = new(uint256.Int).Sub(p.spent[from], tx.costCap) } - p.stored -= uint64(drop.storageSize) - p.lookup.untrack(drop) + p.stored -= uint64(tx.storageSize) + p.lookup.untrack(tx) // Remove the transaction from the pool's eviction heap: // - If the entire account was dropped, pop off the address // - Otherwise, if the new tail has better eviction caps, fix the heap if last { - heap.Pop(p.evict) + heap.Remove(p.evict, p.evict.index[from]) } else { tail := txs[len(txs)-1] // new tail, surely exists - evictionExecFeeDiff := tail.evictionExecFeeJumps - drop.evictionExecFeeJumps - evictionBlobFeeDiff := tail.evictionBlobFeeJumps - drop.evictionBlobFeeJumps + evictionExecFeeDiff := tail.evictionExecFeeJumps - tx.evictionExecFeeJumps + evictionBlobFeeDiff := tail.evictionBlobFeeJumps - tx.evictionBlobFeeJumps if evictionExecFeeDiff > 0.001 || evictionBlobFeeDiff > 0.001 { // no need for math.Abs, monotonic decreasing - heap.Fix(p.evict, 0) + heap.Fix(p.evict, p.evict.index[from]) } } - // Remove the transaction from the data store + if err := p.store.Delete(tx.id); err != nil { + log.Error("Failed to drop evicted transaction", "id", tx.id, "err", err) + } +} + +// drop removes the worst transaction from the pool. It is primarily used when a +// freshly added transaction overflows the pool and needs to evict something. The +// method is also called on startup if the user resizes their storage, might be an +// expensive run but it should be fine-ish. +func (p *BlobPool) drop() (common.Address, uint64, uint64) { + // Peek at the account with the worse transaction set to evict from (Go's heap + // stores the minimum at index zero of the heap slice) and retrieve it's last + // transaction. + var ( + from = p.evict.addrs[0] // cannot call drop on empty pool + txs = p.index[from] + drop = txs[len(txs)-1] + ) + p.remove(drop, from) log.Debug("Evicting overflown blob transaction", "from", from, "evicted", drop.nonce, "id", drop.id) dropOverflownMeter.Mark(1) - - if deleteTxBlob { - if err := p.store.Delete(drop.id); err != nil { - log.Error("Failed to drop evicted transaction", "id", drop.id, "err", err) - } - } return from, drop.nonce, drop.id }