core/txpool/blobpool: introduce sidecar conversion

This commit is contained in:
Gary Rong 2025-09-17 14:19:07 +08:00
parent 8a171dce1f
commit e1f3232f88
2 changed files with 166 additions and 14 deletions

View file

@ -61,11 +61,11 @@ const (
// small buffer is added to the proof overhead. // small buffer is added to the proof overhead.
txBlobOverhead = uint32(kzg4844.CellProofsPerBlob*len(kzg4844.Proof{}) + 64) txBlobOverhead = uint32(kzg4844.CellProofsPerBlob*len(kzg4844.Proof{}) + 64)
// txMaxSize is the maximum size a single transaction can have, outside // txMaxSize is the maximum size a single transaction can have, including the
// the included blobs. Since blob transactions are pulled instead of pushed, // blobs. Since blob transactions are pulled instead of pushed, and only a
// and only a small metadata is kept in ram, the rest is on disk, there is // small metadata is kept in ram, the rest is on disk, there is no critical
// no critical limit that should be enforced. Still, capping it to some sane // limit that should be enforced. Still, capping it to some sane limit can
// limit can never hurt. // never hurt, which is aligned with maxBlobsPerTx constraint enforced internally.
txMaxSize = 1024 * 1024 txMaxSize = 1024 * 1024
// maxBlobsPerTx is the maximum number of blobs that a single transaction can // maxBlobsPerTx is the maximum number of blobs that a single transaction can
@ -331,6 +331,7 @@ type BlobPool struct {
signer types.Signer // Transaction signer to use for sender recovery signer types.Signer // Transaction signer to use for sender recovery
chain BlockChain // Chain object to access the state through chain BlockChain // Chain object to access the state through
cQueue *conversionQueue // The queue for performing legacy sidecar conversion
head *types.Header // Current head of the chain head *types.Header // Current head of the chain
state *state.StateDB // Current state at the head of the chain state *state.StateDB // Current state at the head of the chain
@ -359,6 +360,7 @@ func New(config Config, chain BlockChain, hasPendingAuth func(common.Address) bo
hasPendingAuth: hasPendingAuth, hasPendingAuth: hasPendingAuth,
signer: types.LatestSigner(chain.Config()), signer: types.LatestSigner(chain.Config()),
chain: chain, chain: chain,
cQueue: newConversionQueue(),
lookup: newLookup(), lookup: newLookup(),
index: make(map[common.Address][]*blobTxMeta), index: make(map[common.Address][]*blobTxMeta),
spent: make(map[common.Address]*uint256.Int), spent: make(map[common.Address]*uint256.Int),
@ -475,6 +477,7 @@ func (p *BlobPool) Init(gasTip uint64, head *types.Header, reserver txpool.Reser
// Close closes down the underlying persistent store. // Close closes down the underlying persistent store.
func (p *BlobPool) Close() error { func (p *BlobPool) Close() error {
var errs []error var errs []error
p.cQueue.close()
if p.limbo != nil { // Close might be invoked due to error in constructor, before p,limbo is set if p.limbo != nil { // Close might be invoked due to error in constructor, before p,limbo is set
if err := p.limbo.Close(); err != nil { if err := p.limbo.Close(); err != nil {
errs = append(errs, err) errs = append(errs, err)
@ -1164,10 +1167,10 @@ func (p *BlobPool) checkDelegationLimit(tx *types.Transaction) error {
// validateTx checks whether a transaction is valid according to the consensus // validateTx checks whether a transaction is valid according to the consensus
// rules and adheres to some heuristic limits of the local node (price and size). // rules and adheres to some heuristic limits of the local node (price and size).
//
// This function assumes the static validation has been performed already and
// only runs the stateful checks with lock protection.
func (p *BlobPool) validateTx(tx *types.Transaction) error { func (p *BlobPool) validateTx(tx *types.Transaction) error {
if err := p.ValidateTxBasics(tx); err != nil {
return err
}
// Ensure the transaction adheres to the stateful pool filters (nonce, balance) // Ensure the transaction adheres to the stateful pool filters (nonce, balance)
stateOpts := &txpool.ValidationOptionsWithState{ stateOpts := &txpool.ValidationOptionsWithState{
State: p.state, State: p.state,
@ -1432,17 +1435,50 @@ func (p *BlobPool) AvailableBlobs(vhashes []common.Hash) int {
return available return available
} }
// preprocess performs the static validation upon the provided txs and converts
// the legacy sidecars if Osaka fork has been activated.
//
// This function is pure static and lock free.
func (p *BlobPool) preprocess(txs []*types.Transaction) ([]*types.Transaction, []error) {
head := p.chain.CurrentBlock()
if !p.chain.Config().IsOsaka(head.Number, head.Time) {
return txs, make([]error, len(txs))
}
var errs []error
for _, tx := range txs {
// Validate the transaction statically at first to avoid unnecessary
// conversion. This step doesn't require lock.
if err := p.ValidateTxBasics(tx); err != nil {
errs = append(errs, err)
continue
}
// Convert the legacy sidecar after Osaka fork. This could be a long
// procedure which takes a few seconds, even minutes. Fortunately it
// will only block the routine of the source peer without affecting
// other parts.
if tx.BlobTxSidecar().Version == types.BlobSidecarVersion0 {
if err := p.cQueue.convert(tx); err != nil {
errs = append(errs, err)
continue
}
}
errs = append(errs, nil)
}
return txs, errs
}
// Add inserts a set of blob transactions into the pool if they pass validation (both // Add inserts a set of blob transactions into the pool if they pass validation (both
// consensus validity and pool restrictions). // consensus validity and pool restrictions).
//
// Note, if sync is set the method will block until all internal maintenance
// related to the add is finished. Only use this during tests for determinism.
func (p *BlobPool) Add(txs []*types.Transaction, sync bool) []error { func (p *BlobPool) Add(txs []*types.Transaction, sync bool) []error {
var ( var (
errs = make([]error, len(txs)) errs []error
adds = make([]*types.Transaction, 0, len(txs)) adds = make([]*types.Transaction, 0, len(txs))
) )
txs, errs = p.preprocess(txs)
for i, tx := range txs { for i, tx := range txs {
if errs[i] != nil {
continue
}
errs[i] = p.add(tx) errs[i] = p.add(tx)
if errs[i] == nil { if errs[i] == nil {
adds = append(adds, tx.WithoutBlobTxSidecar()) adds = append(adds, tx.WithoutBlobTxSidecar())

View file

@ -0,0 +1,116 @@
// Copyright 2025 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 <http://www.gnu.org/licenses/>.
package blobpool
import (
"errors"
"github.com/ethereum/go-ethereum/core/types"
)
// cTask represents a conversion task with an attached result channel.
type cTask struct {
tx *types.Transaction // Blob transaction, sidecar is expected
done chan error // Channel for signaling back if the conversion succeeds
}
// conversionQueue is a dedicated queue for converting legacy blob transactions
// received from the network after the Osaka fork. Since conversion is expensive,
// it is performed in the background by a single thread, ensuring the main Geth
// process is not overloaded.
type conversionQueue struct {
tasks chan *cTask
quit chan struct{}
closed chan struct{}
}
// newConversionQueue constructs the conversion queue.
func newConversionQueue() *conversionQueue {
q := &conversionQueue{
tasks: make(chan *cTask),
quit: make(chan struct{}),
closed: make(chan struct{}),
}
go q.loop()
return q
}
// convert accepts a legacy blob transaction with version-0 blobs and queues it
// for conversion.
//
// This function may block for a long time until the transaction is processed.
func (q *conversionQueue) convert(tx *types.Transaction) error {
done := make(chan error, 1)
select {
case q.tasks <- &cTask{tx: tx, done: done}:
return <-done
case <-q.closed:
return errors.New("conversion queue closed")
}
}
// close terminates the conversion queue.
func (q *conversionQueue) close() {
select {
case <-q.closed:
return
default:
close(q.quit)
<-q.closed
}
}
func (q *conversionQueue) run(tasks []*cTask, done chan struct{}) {
defer close(done)
for _, t := range tasks {
sidecar := t.tx.BlobTxSidecar()
if sidecar == nil {
t.done <- errors.New("tx without sidecar")
continue
}
// Run the conversion, the original sidecar will be mutated in place
t.done <- sidecar.ToV1()
}
}
func (q *conversionQueue) loop() {
defer close(q.closed)
var (
done = make(chan struct{})
cTasks []*cTask
)
for {
select {
case t := <-q.tasks:
cTasks = append(cTasks, t)
if done == nil {
done = make(chan struct{})
tasks := cTasks
cTasks = cTasks[:0]
go q.run(tasks, done)
}
case <-done:
done = nil
case <-q.quit:
return
}
}
}