From d35d52b0ca43e8e23fd605c8d8e873b0c1e131a6 Mon Sep 17 00:00:00 2001 From: healthykim Date: Mon, 8 Sep 2025 22:04:03 +0900 Subject: [PATCH] add: implement pooledTx type --- core/txpool/blobpool/blobpool.go | 136 +++++++++++++++++++------- core/txpool/blobpool/blobpool_test.go | 17 +++- core/txpool/blobpool/limbo.go | 8 +- core/txpool/blobpool/lookup.go | 2 +- core/types/tx_blob.go | 9 +- 5 files changed, 122 insertions(+), 50 deletions(-) diff --git a/core/txpool/blobpool/blobpool.go b/core/txpool/blobpool/blobpool.go index d92cfc84dc..c3fe2fd05d 100644 --- a/core/txpool/blobpool/blobpool.go +++ b/core/txpool/blobpool/blobpool.go @@ -137,8 +137,8 @@ func newBlobTxMeta(id uint64, size uint64, storageSize uint32, tx *types.Transac vhashes: tx.BlobHashes(), version: tx.BlobTxSidecar().Version, id: id, - storageSize: storageSize, - size: size, + storageSize: storageSize, // size of tx including cells + size: size, // size of tx only nonce: tx.Nonce(), costCap: uint256.MustFromBig(tx.Cost()), execTipCap: uint256.MustFromBig(tx.GasTipCap()), @@ -153,6 +153,26 @@ func newBlobTxMeta(id uint64, size uint64, storageSize uint32, tx *types.Transac return meta } +type PooledBlobTx struct { + Transaction *types.Transaction + Sidecar *types.BlobTxCellSidecar +} + +func NewPooledBlobTx(tx *types.Transaction, sidecar *types.BlobTxCellSidecar) *PooledBlobTx { + return &PooledBlobTx{ + Transaction: tx, + Sidecar: sidecar, + } +} + +func (ptx *PooledBlobTx) Hash() common.Hash { + return ptx.Transaction.Hash() +} + +func (ptx *PooledBlobTx) Convert() *types.Transaction { + return ptx.Transaction.WithBlobTxSidecar(ptx.Sidecar.ToBlobTxSidecar()) +} + // BlobPool is the transaction pool dedicated to EIP-4844 blob transactions. // // Blob transactions are special snowflakes that are designed for a very specific @@ -497,19 +517,31 @@ func (p *BlobPool) Close() error { // each transaction on disk to create the in-memory metadata index. func (p *BlobPool) parseTransaction(id uint64, size uint32, blob []byte) error { tx := new(types.Transaction) - if err := rlp.DecodeBytes(blob, tx); err != nil { + + pooledTx := new(PooledBlobTx) + if err := rlp.DecodeBytes(blob, pooledTx); err != nil { // This path is impossible unless the disk data representation changes // across restarts. For that ever improbable case, recover gracefully // by ignoring this data entry. - log.Error("Failed to decode blob pool entry", "id", id, "err", err) - return err - } - if tx.BlobTxSidecar() == nil { - log.Error("Missing sidecar in blob pool entry", "id", id, "hash", tx.Hash()) - return errors.New("missing blob sidecar") - } + if err := rlp.DecodeBytes(blob, tx); err != nil { + log.Error("Failed to decode blob pool entry", "id", id, "err", err) + return errors.New("unknown tx type") + } + if tx.BlobTxSidecar() == nil { + log.Error("Missing sidecar in blob pool entry", "id", id, "hash", tx.Hash()) + return errors.New("missing sidecar") + } + } else { + sidecar := pooledTx.Sidecar.ToBlobTxSidecar() + if sidecar == nil { + log.Error("Missing sidecar in blob pool entry", "id", id, "hash", pooledTx.Transaction.Hash()) + return errors.New("missing sidecar") + } + tx = pooledTx.Convert() + } meta := newBlobTxMeta(id, tx.Size(), size, tx) + if p.lookup.exists(meta.hash) { // This path is only possible after a crash, where deleted items are not // removed via the normal shutdown-startup procedure and thus may get @@ -799,17 +831,17 @@ func (p *BlobPool) offload(addr common.Address, nonce uint64, id uint64, inclusi log.Error("Blobs missing for included transaction", "from", addr, "nonce", nonce, "id", id, "err", err) return } - var tx types.Transaction - if err = rlp.DecodeBytes(data, &tx); err != nil { + var pooledTx PooledBlobTx + if err = rlp.DecodeBytes(data, &pooledTx); err != nil { log.Error("Blobs corrupted for included transaction", "from", addr, "nonce", nonce, "id", id, "err", err) return } - block, ok := inclusions[tx.Hash()] + block, ok := inclusions[pooledTx.Transaction.Hash()] if !ok { log.Warn("Blob transaction swapped out by signer", "from", addr, "nonce", nonce, "id", id) return } - if err := p.limbo.push(&tx, block); err != nil { + if err := p.limbo.push(&pooledTx, block); err != nil { log.Warn("Failed to offload blob tx into limbo", "err", err) return } @@ -1021,20 +1053,20 @@ func (p *BlobPool) reinject(addr common.Address, txhash common.Hash) error { // Serialize the transaction back into the primary datastore. blob, err := rlp.EncodeToBytes(tx) if err != nil { - log.Error("Failed to encode transaction for storage", "hash", tx.Hash(), "err", err) + log.Error("Failed to encode transaction for storage", "hash", tx.Transaction.Hash(), "err", err) return err } id, err := p.store.Put(blob) if err != nil { - log.Error("Failed to write transaction into storage", "hash", tx.Hash(), "err", err) + log.Error("Failed to write transaction into storage", "hash", tx.Transaction.Hash(), "err", err) return err } // Update the indices and metrics - meta := newBlobTxMeta(id, tx.Size(), p.store.Size(id), tx) + meta := newBlobTxMeta(id, tx.Transaction.Size(), p.store.Size(id), tx.Transaction) if _, ok := p.index[addr]; !ok { if err := p.reserver.Hold(addr); err != nil { - log.Warn("Failed to reserve account for blob pool", "tx", tx.Hash(), "from", addr, "err", err) + log.Warn("Failed to reserve account for blob pool", "tx", tx.Transaction.Hash(), "from", addr, "err", err) return err } p.index[addr] = []*blobTxMeta{meta} @@ -1282,10 +1314,26 @@ func (p *BlobPool) getRLP(hash common.Hash) []byte { } data, err := p.store.Get(id) if err != nil { - log.Error("Tracked blob transaction missing from store", "hash", hash, "id", id, "err", err) + log.Error("Failed to get transaction in blobpool", "hash", hash, "id", id, "err", err) return nil } - return data + tx := new(types.Transaction) + pooledTx := new(PooledBlobTx) + if err := rlp.DecodeBytes(data, pooledTx); err != nil { + if err := rlp.DecodeBytes(data, tx); err != nil { + log.Error("Failed to decode transaction in blobpool", "hash", hash, "id", id, "err", err) + return nil + } + return data + } + + encoded, err := rlp.EncodeToBytes(pooledTx.Convert()) + if err != nil { + log.Error("Failed to encode transaction in blobpool", "hash", hash, "id", id, "err", err) + return nil + } + + return encoded } // Get returns a transaction if it is contained in the pool, or nil otherwise. @@ -1294,15 +1342,19 @@ func (p *BlobPool) Get(hash common.Hash) *types.Transaction { if len(data) == 0 { return nil } - item := new(types.Transaction) - if err := rlp.DecodeBytes(data, item); err != nil { - id, _ := p.lookup.storeidOfTx(hash) + tx := new(types.Transaction) + pooledTx := new(PooledBlobTx) + if err := rlp.DecodeBytes(data, pooledTx); err != nil { + if err := rlp.DecodeBytes(data, tx); err != nil { + id, _ := p.lookup.storeidOfTx(hash) - log.Error("Blobs corrupted for traced transaction", - "hash", hash, "id", id, "err", err) - return nil + log.Error("Blobs corrupted for traced transaction", + "hash", hash, "id", id, "err", err) + return nil + } + return tx } - return item + return pooledTx.Convert() } // GetRLP returns a RLP-encoded transaction if it is contained in the pool. @@ -1314,7 +1366,7 @@ func (p *BlobPool) GetRLP(hash common.Hash) []byte { // given transaction hash. // // The size refers the length of the 'rlp encoding' of a blob transaction -// including the attached blobs. +// excluding the attached blobs. func (p *BlobPool) GetMetadata(hash common.Hash) *txpool.TxMetadata { p.lock.RLock() defer p.lock.RUnlock() @@ -1372,11 +1424,18 @@ func (p *BlobPool) GetBlobs(vhashes []common.Hash, version byte) ([]*kzg4844.Blo // Decode the blob transaction tx := new(types.Transaction) - if err := rlp.DecodeBytes(data, tx); err != nil { - log.Error("Blobs corrupted for traced transaction", "id", txID, "err", err) - continue + var sidecar *types.BlobTxSidecar + pooledTx := new(PooledBlobTx) + if err := rlp.DecodeBytes(data, pooledTx); err != nil { + if err := rlp.DecodeBytes(data, tx); err != nil { + log.Error("Blobs corrupted for traced transaction", "id", txID, "err", err) + continue + } + sidecar = tx.BlobTxSidecar() + } else { + tx = pooledTx.Transaction + sidecar = pooledTx.Sidecar.ToBlobTxSidecar() } - sidecar := tx.BlobTxSidecar() if sidecar == nil { log.Error("Blob tx without sidecar", "hash", tx.Hash(), "id", txID) continue @@ -1454,7 +1513,11 @@ func (p *BlobPool) Add(txs []*types.Transaction, sync bool) []error { if errs[i] != nil { continue } - cellSidecar := tx.BlobTxSidecar().ToBlobTxCellSidecar() + cellSidecar, err := tx.BlobTxSidecar().ToBlobTxCellSidecar() + if err != nil { + errs[i] = err + continue + } errs[i] = p.add(tx.WithoutBlobTxSidecar(), cellSidecar) if errs[i] == nil { adds = append(adds, tx.WithoutBlobTxSidecar()) @@ -1526,12 +1589,11 @@ func (p *BlobPool) add(tx *types.Transaction, cellSidecar *types.BlobTxCellSidec }() } - //todo(healthykim) remove this and seperate database or introduce list - blobSidecar := cellSidecar.ToBlobTxSidecar() - tx = tx.WithBlobTxSidecar(blobSidecar) // Transaction permitted into the pool from a nonce and cost perspective, // insert it into the database and update the indices - blob, err := rlp.EncodeToBytes(tx) + pooledTx := NewPooledBlobTx(tx, cellSidecar) + + blob, err := rlp.EncodeToBytes(pooledTx) if err != nil { log.Error("Failed to encode transaction for storage", "hash", tx.Hash(), "err", err) return err diff --git a/core/txpool/blobpool/blobpool_test.go b/core/txpool/blobpool/blobpool_test.go index 1016461080..b6d232f06e 100644 --- a/core/txpool/blobpool/blobpool_test.go +++ b/core/txpool/blobpool/blobpool_test.go @@ -452,8 +452,12 @@ func verifyBlobRetrievals(t *testing.T, pool *BlobPool) { } // Item retrieved, make sure it matches the expectation index := testBlobIndices[hash] - if *blobs1[i] != *testBlobs[index] || proofs1[i][0] != testBlobProofs[index] { - t.Errorf("retrieved blob or proof mismatch: item %d, hash %x", i, hash) + if *blobs1[i] != *testBlobs[index] { + t.Errorf("retrieved blob mismatch: item %d, hash %x", i, hash) + continue + } + if proofs1[i][0] != testBlobProofs[index] { + t.Errorf("retrieved proof mismatch: item %d, hash %x", i, hash) continue } if *blobs2[i] != *testBlobs[index] || !slices.Equal(proofs2[i], testBlobCellProofs[index]) { @@ -1752,7 +1756,9 @@ func TestAdd(t *testing.T) { // Add each transaction one by one, verifying the pool internals in between for j, add := range tt.adds { signed, _ := types.SignNewTx(keys[add.from], types.LatestSigner(params.MainnetChainConfig), add.tx) - if err := pool.add(signed.WithoutBlobTxSidecar(), signed.BlobTxSidecar().ToBlobTxCellSidecar()); !errors.Is(err, add.err) { + sidecar, _ := signed.BlobTxSidecar().ToBlobTxCellSidecar() + + if err := pool.add(signed.WithoutBlobTxSidecar(), sidecar); !errors.Is(err, add.err) { t.Errorf("test %d, tx %d: adding transaction error mismatch: have %v, want %v", i, j, err, add.err) } if add.err == nil { @@ -1760,7 +1766,7 @@ func TestAdd(t *testing.T) { if !exist { t.Errorf("test %d, tx %d: failed to lookup transaction's size", i, j) } - if size != signed.Size() { + if size != signed.WithoutBlobTxSidecar().Size() { t.Errorf("test %d, tx %d: transaction's size mismatches: have %v, want %v", i, j, size, signed.Size()) } @@ -2124,7 +2130,8 @@ func benchmarkPoolPending(b *testing.B, datacap uint64) { b.Fatal(err) } statedb.AddBalance(addr, uint256.NewInt(1_000_000_000), tracing.BalanceChangeUnspecified) - pool.add(tx.WithoutBlobTxSidecar(), tx.BlobTxSidecar().ToBlobTxCellSidecar()) + sidecar, _ := tx.BlobTxSidecar().ToBlobTxCellSidecar() + pool.add(tx.WithoutBlobTxSidecar(), sidecar) } statedb.Commit(0, true, false) defer pool.Close() diff --git a/core/txpool/blobpool/limbo.go b/core/txpool/blobpool/limbo.go index 50c40c9d83..7b2bf600d5 100644 --- a/core/txpool/blobpool/limbo.go +++ b/core/txpool/blobpool/limbo.go @@ -34,7 +34,7 @@ import ( type limboBlob struct { TxHash common.Hash // Owner transaction's hash to support resurrecting reorged txs Block uint64 // Block in which the blob transaction was included - Tx *types.Transaction + Tx *PooledBlobTx } // limbo is a light, indexed database to temporarily store recently included @@ -147,7 +147,7 @@ func (l *limbo) finalize(final *types.Header) { // push stores a new blob transaction into the limbo, waiting until finality for // it to be automatically evicted. -func (l *limbo) push(tx *types.Transaction, block uint64) error { +func (l *limbo) push(tx *PooledBlobTx, block uint64) error { // If the blobs are already tracked by the limbo, consider it a programming // error. There's not much to do against it, but be loud. if _, ok := l.index[tx.Hash()]; ok { @@ -164,7 +164,7 @@ func (l *limbo) push(tx *types.Transaction, block uint64) error { // pull retrieves a previously pushed set of blobs back from the limbo, removing // it at the same time. This method should be used when a previously included blob // transaction gets reorged out. -func (l *limbo) pull(tx common.Hash) (*types.Transaction, error) { +func (l *limbo) pull(tx common.Hash) (*PooledBlobTx, error) { // If the blobs are not tracked by the limbo, there's not much to do. This // can happen for example if a blob transaction is mined without pushing it // into the network first. @@ -241,7 +241,7 @@ func (l *limbo) getAndDrop(id uint64) (*limboBlob, error) { // setAndIndex assembles a limbo blob database entry and stores it, also updating // the in-memory indices. -func (l *limbo) setAndIndex(tx *types.Transaction, block uint64) error { +func (l *limbo) setAndIndex(tx *PooledBlobTx, block uint64) error { txhash := tx.Hash() item := &limboBlob{ TxHash: txhash, diff --git a/core/txpool/blobpool/lookup.go b/core/txpool/blobpool/lookup.go index 7607cd487a..68817fa99d 100644 --- a/core/txpool/blobpool/lookup.go +++ b/core/txpool/blobpool/lookup.go @@ -22,7 +22,7 @@ import ( type txMetadata struct { id uint64 // the billy id of transction - size uint64 // the RLP encoded size of transaction (blobs are included) + size uint64 // the RLP encoded size of transaction (blobs are excluded) } // lookup maps blob versioned hashes to transaction hashes that include them, diff --git a/core/types/tx_blob.go b/core/types/tx_blob.go index 8242cd6146..27abdd808a 100644 --- a/core/types/tx_blob.go +++ b/core/types/tx_blob.go @@ -176,15 +176,18 @@ func (sc *BlobTxSidecar) Copy() *BlobTxSidecar { } } -func (sc *BlobTxSidecar) ToBlobTxCellSidecar() *BlobTxCellSidecar { - cells, _ := kzg4844.ComputeCells(sc.Blobs) +func (sc *BlobTxSidecar) ToBlobTxCellSidecar() (*BlobTxCellSidecar, error) { + cells, err := kzg4844.ComputeCells(sc.Blobs) + if err != nil { + return nil, err + } return &BlobTxCellSidecar{ Version: sc.Version, Cells: cells, Commitments: sc.Commitments, Proofs: sc.Proofs, CellIndices: CustodyBitmap{}.SetAll(), - } + }, nil } type BlobTxCellSidecar struct {