mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-21 12:16:44 +00:00
This is the implementation of EIP-8070 Sparse Blobpool. It introduces protocol version eth/72 which relays blob transaction cells instead of full blobs. The blobpool now store 'incomplete' transactions, where only some of the cells are provided. The stored cell indexes are taken from the custody bitmap, which is provided by the consensus layer in forkchoiceUpdatedV4. This method will be called once Glamsterdam activates, and the default custody is full custody, so for now there is no change in the amount of stored cells for now. The main entities added are the BlobBuffer and BlobFetcher, which work together to track and fetch missing cells from connected peers. The partial transactions become available for inclusion in blocks when they are covered by enough peers that hold all cells. This change also introduces engine_getBlobsV4, which allows for cell-based responses (and partial blobs with only some of the cells). We maintain backward compatibility with getBlobsV3 which expects full blob responses by recovering the blob from available cells. Since this process is resource-intensive, we proactively cache the conversion so it is ready in time for getBlobsV3 calls. This mechanism will be removed once support for getBlobsV4 is universal across all consensus layer implementations. devp2p tests for eth/72 are not part of this initial change. This is to avoid breaking test success status for execution clients that do not have eth/72 implemented yet. The tests will be added in a subsequent change. --------- Co-authored-by: Felix Lange <fjl@twurst.com>
151 lines
5.1 KiB
Go
151 lines
5.1 KiB
Go
// Copyright 2020 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 eth
|
|
|
|
import (
|
|
"errors"
|
|
"fmt"
|
|
|
|
"github.com/ethereum/go-ethereum/common"
|
|
"github.com/ethereum/go-ethereum/core"
|
|
"github.com/ethereum/go-ethereum/core/types"
|
|
"github.com/ethereum/go-ethereum/crypto/kzg4844"
|
|
"github.com/ethereum/go-ethereum/eth/protocols/eth"
|
|
"github.com/ethereum/go-ethereum/p2p/enode"
|
|
"github.com/ethereum/go-ethereum/params"
|
|
)
|
|
|
|
// ethHandler implements the eth.Backend interface to handle the various network
|
|
// packets that are sent as replies or broadcasts.
|
|
type ethHandler handler
|
|
|
|
func (h *ethHandler) Chain() *core.BlockChain { return h.chain }
|
|
func (h *ethHandler) TxPool() eth.TxPool { return h.txpool }
|
|
func (h *ethHandler) BlobPool() eth.BlobPool { return h.blobpool }
|
|
|
|
// RunPeer is invoked when a peer joins on the `eth` protocol.
|
|
func (h *ethHandler) RunPeer(peer *eth.Peer, hand eth.Handler) error {
|
|
return (*handler)(h).runEthPeer(peer, hand)
|
|
}
|
|
|
|
// PeerInfo retrieves all known `eth` information about a peer.
|
|
func (h *ethHandler) PeerInfo(id enode.ID) interface{} {
|
|
if p := h.peers.peer(id.String()); p != nil {
|
|
return p.info()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// AcceptTxs retrieves whether transaction processing is enabled on the node
|
|
// or if inbound transactions should simply be dropped.
|
|
func (h *ethHandler) AcceptTxs() bool {
|
|
return h.synced.Load()
|
|
}
|
|
|
|
// Handle is invoked from a peer's message handler when it receives a new remote
|
|
// message that the handler couldn't consume and serve itself.
|
|
func (h *ethHandler) Handle(peer *eth.Peer, packet eth.Packet) error {
|
|
// Consume any broadcasts and announces, forwarding the rest to the downloader
|
|
switch packet := packet.(type) {
|
|
case *eth.NewPooledTransactionHashesPacket72:
|
|
hashes, err := h.txFetcher.Notify(peer.ID(), packet.Types, packet.Sizes, packet.Hashes)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if len(hashes) != 0 {
|
|
return h.blobFetcher.Notify(peer.ID(), hashes, packet.Mask)
|
|
}
|
|
return nil
|
|
|
|
case *eth.NewPooledTransactionHashesPacket71:
|
|
_, err := h.txFetcher.Notify(peer.ID(), packet.Types, packet.Sizes, packet.Hashes)
|
|
return err
|
|
|
|
case *eth.TransactionsPacket:
|
|
txs, err := packet.Items()
|
|
if err != nil {
|
|
return fmt.Errorf("Transactions: %v", err)
|
|
}
|
|
if err := handleTransactions(peer, txs, true); err != nil {
|
|
return fmt.Errorf("Transactions: %v", err)
|
|
}
|
|
return h.txFetcher.Enqueue(peer.ID(), peer.Version(), txs, false)
|
|
|
|
case *eth.PooledTransactionsPacket:
|
|
txs, err := packet.List.Items()
|
|
if err != nil {
|
|
return fmt.Errorf("PooledTransactions: %v", err)
|
|
}
|
|
if err := handleTransactions(peer, txs, false); err != nil {
|
|
return fmt.Errorf("PooledTransactions: %v", err)
|
|
}
|
|
return h.txFetcher.Enqueue(peer.ID(), peer.Version(), txs, true)
|
|
|
|
case *eth.CellsResponse:
|
|
outer, err := packet.Cells.Items()
|
|
if err != nil {
|
|
return fmt.Errorf("Cells: %v", err)
|
|
}
|
|
cells := make([][]kzg4844.Cell, len(outer))
|
|
for i := range outer {
|
|
if outer[i].Len() > params.BlobTxMaxBlobs*kzg4844.CellsPerBlob {
|
|
return fmt.Errorf("Cells: cells per tx exceeded the possible maximum")
|
|
}
|
|
if cells[i], err = outer[i].Items(); err != nil {
|
|
return fmt.Errorf("Cells: %v", err)
|
|
}
|
|
}
|
|
return h.blobFetcher.Enqueue(peer.ID(), packet.Hashes, cells, packet.Mask)
|
|
|
|
default:
|
|
return fmt.Errorf("unexpected eth packet type: %T", packet)
|
|
}
|
|
}
|
|
|
|
// handleTransactions marks all given transactions as known to the peer
|
|
// and performs basic validations.
|
|
func handleTransactions(peer *eth.Peer, list []*types.Transaction, directBroadcast bool) error {
|
|
seen := make(map[common.Hash]struct{}, len(list))
|
|
for _, tx := range list {
|
|
if tx.Type() == types.BlobTxType {
|
|
if directBroadcast {
|
|
return errors.New("disallowed broadcast blob transaction")
|
|
} else {
|
|
// If we receive any blob transactions missing sidecars, or with
|
|
// sidecars that don't correspond to the versioned hashes reported
|
|
// in the header, disconnect from the sending peer.
|
|
if tx.BlobTxSidecar() == nil {
|
|
return errors.New("received sidecar-less blob transaction")
|
|
}
|
|
if err := tx.BlobTxSidecar().ValidateBlobCommitmentHashes(tx.BlobHashes()); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
|
|
// Check for duplicates.
|
|
hash := tx.Hash()
|
|
if _, exists := seen[hash]; exists {
|
|
return fmt.Errorf("multiple copies of the same hash %v", hash)
|
|
}
|
|
seen[hash] = struct{}{}
|
|
|
|
// Mark as known.
|
|
peer.MarkTransaction(hash)
|
|
}
|
|
return nil
|
|
}
|