core/txpool/blobpool: add fork boundary conversion

This commit is contained in:
Gary Rong 2025-09-18 21:22:14 +08:00
parent d54a983742
commit de75c79fe4
2 changed files with 119 additions and 1 deletions

View file

@ -883,6 +883,26 @@ func (p *BlobPool) Reset(oldHead, newHead *types.Header) {
basefeeGauge.Update(int64(basefee.Uint64())) basefeeGauge.Update(int64(basefee.Uint64()))
blobfeeGauge.Update(int64(blobfee.Uint64())) blobfeeGauge.Update(int64(blobfee.Uint64()))
p.updateStorageMetrics() p.updateStorageMetrics()
if p.chain.Config().IsOsaka(newHead.Number, newHead.Time) {
txs := make(map[common.Address]map[common.Hash]uint64)
for address, list := range p.index {
txs[address] = make(map[common.Hash]uint64)
for _, meta := range list {
if meta.version == types.BlobSidecarVersion1 {
continue
}
txs[address][meta.hash] = meta.id
}
}
go func() {
for address, list := range txs {
for hash, id := range list {
p.convert(address, hash, id)
}
}
}()
}
} }
// reorg assembles all the transactors and missing transactions between an old // reorg assembles all the transactors and missing transactions between an old
@ -1061,6 +1081,93 @@ func (p *BlobPool) reinject(addr common.Address, txhash common.Hash) error {
return nil return nil
} }
// convert fetches transaction data from the store, performs an on-the-fly conversion,
// and persists the result. This function is intended for use only during the
// Osaka fork transition period.
func (p *BlobPool) convert(sender common.Address, hash common.Hash, id uint64) (success bool) {
// Remove the legacy blob tx data if conversion succeeds
defer func() {
if success {
if err := p.store.Delete(id); err != nil {
log.Error("Failed to delete old blob transaction", "hash", hash, "id", id, "err", err)
}
}
}()
// Retrieves the legacy blob transaction from the underlying store without
// holding the pool mutex. The transaction may have been evicted simultaneously,
// if so, it is safely ignored.
start := time.Now()
data, err := p.store.Get(id)
if err != nil || len(data) == 0 {
log.Debug("Blob transaction is missing", "hash", hash, "id", id, "err", err)
return false
}
oldStorageSize := p.store.Size(id)
// Convert the sidecar without the lock, as this procedure can take a while
var tx types.Transaction
if err = rlp.DecodeBytes(data, &tx); err != nil {
log.Error("Blob transaction is corrupted", "hash", hash, "id", id, "err", err)
return false
}
if err := tx.BlobTxSidecar().ToV1(); err != nil {
log.Error("Failed to convert blob transaction", "hash", hash, "err", err)
return false
}
// 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)
return false
}
newSize := uint64(len(blob))
newId, err := p.store.Put(blob)
if err != nil {
log.Error("Failed to store transaction", "hash", hash, "err", err)
return false
}
newStorageSize := p.store.Size(newId)
// Remove the newly added blob transaction if the transaction is not held
// by the pool anymore.
defer func() {
if !success {
if err := p.store.Delete(newId); err != nil {
log.Error("Failed to delete outdated transaction", "hash", hash, "err", err)
}
log.Info("Blob transaction is outdated", "hash", hash)
}
}()
// Update the metadata in the blob pool by holding the mutex
p.lock.Lock()
defer p.lock.Unlock()
// The transaction may have been evicted simultaneously, if so, it is
// safely ignored.
if !p.lookup.exists(hash) {
return false
}
for _, meta := range p.index[sender] {
if meta.hash == hash {
meta.id = newId
meta.version = types.BlobSidecarVersion1
meta.storageSize = newStorageSize
meta.size = newSize
p.lookup.updateTxIndex(hash, newId, newSize)
p.stored += uint64(newStorageSize)
p.stored -= uint64(oldStorageSize)
break
}
}
log.Debug("Converted the legacy transaction", "hash", hash, "elapsed", common.PrettyDuration(time.Since(start)))
return true
}
// SetGasTip implements txpool.SubPool, allowing the blob pool's gas requirements // SetGasTip implements txpool.SubPool, allowing the blob pool's gas requirements
// to be kept in sync with the main transaction pool's gas requirements. // to be kept in sync with the main transaction pool's gas requirements.
func (p *BlobPool) SetGasTip(tip *big.Int) { func (p *BlobPool) SetGasTip(tip *big.Int) {

View file

@ -21,7 +21,7 @@ import (
) )
type txMetadata struct { type txMetadata struct {
id uint64 // the billy id of transction id uint64 // the billy id of transaction
size uint64 // the RLP encoded size of transaction (blobs are included) size uint64 // the RLP encoded size of transaction (blobs are included)
} }
@ -110,3 +110,14 @@ func (l *lookup) untrack(tx *blobTxMeta) {
} }
} }
} }
// updateTxIndex updates the transaction index
func (l *lookup) updateTxIndex(hash common.Hash, id uint64, size uint64) bool {
meta, exists := l.txIndex[hash]
if !exists {
return false
}
meta.id = id
meta.size = size
return true
}