mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-25 22:26:42 +00:00
eth/protocols/eth: update eth/69 for status changes and BlockRangeUpdate
This commit is contained in:
parent
ff3cb690b5
commit
c3960e0a07
8 changed files with 171 additions and 71 deletions
|
|
@ -26,7 +26,6 @@ import (
|
|||
|
||||
"github.com/ethereum/go-ethereum/common"
|
||||
"github.com/ethereum/go-ethereum/core"
|
||||
"github.com/ethereum/go-ethereum/core/forkid"
|
||||
"github.com/ethereum/go-ethereum/core/txpool"
|
||||
"github.com/ethereum/go-ethereum/core/types"
|
||||
"github.com/ethereum/go-ethereum/crypto"
|
||||
|
|
@ -105,7 +104,6 @@ type handlerConfig struct {
|
|||
type handler struct {
|
||||
nodeID enode.ID
|
||||
networkID uint64
|
||||
forkFilter forkid.Filter // Fork ID filter, constant across the lifetime of the node
|
||||
|
||||
snapSync atomic.Bool // Flag whether snap sync is enabled (gets disabled if we already have blocks)
|
||||
synced atomic.Bool // Flag whether we're considered synchronised (enables transaction processing)
|
||||
|
|
@ -143,7 +141,6 @@ func newHandler(config *handlerConfig) (*handler, error) {
|
|||
h := &handler{
|
||||
nodeID: config.NodeID,
|
||||
networkID: config.Network,
|
||||
forkFilter: forkid.NewFilter(config.Chain),
|
||||
eventMux: config.EventMux,
|
||||
database: config.Database,
|
||||
txpool: config.TxPool,
|
||||
|
|
@ -256,14 +253,7 @@ func (h *handler) runEthPeer(peer *eth.Peer, handler eth.Handler) error {
|
|||
}
|
||||
|
||||
// Execute the Ethereum handshake
|
||||
var (
|
||||
genesis = h.chain.Genesis()
|
||||
head = h.chain.CurrentHeader()
|
||||
hash = head.Hash()
|
||||
number = head.Number.Uint64()
|
||||
)
|
||||
forkID := forkid.NewID(h.chain.Config(), genesis, number, head.Time)
|
||||
if err := peer.Handshake(h.networkID, hash, genesis.Hash(), forkID, h.forkFilter); err != nil {
|
||||
if err := peer.Handshake(h.networkID, h.chain); err != nil {
|
||||
peer.Log().Debug("Ethereum handshake failed", "err", err)
|
||||
return err
|
||||
}
|
||||
|
|
|
|||
|
|
@ -25,7 +25,6 @@ import (
|
|||
"github.com/ethereum/go-ethereum/common"
|
||||
"github.com/ethereum/go-ethereum/consensus/ethash"
|
||||
"github.com/ethereum/go-ethereum/core"
|
||||
"github.com/ethereum/go-ethereum/core/forkid"
|
||||
"github.com/ethereum/go-ethereum/core/rawdb"
|
||||
"github.com/ethereum/go-ethereum/core/types"
|
||||
"github.com/ethereum/go-ethereum/core/vm"
|
||||
|
|
@ -257,11 +256,7 @@ func testRecvTransactions(t *testing.T, protocol uint) {
|
|||
return eth.Handle((*ethHandler)(handler.handler), peer)
|
||||
})
|
||||
// Run the handshake locally to avoid spinning up a source handler
|
||||
var (
|
||||
genesis = handler.chain.Genesis()
|
||||
head = handler.chain.CurrentBlock()
|
||||
)
|
||||
if err := src.Handshake(1, head.Hash(), genesis.Hash(), forkid.NewIDWithChain(handler.chain), forkid.NewFilter(handler.chain)); err != nil {
|
||||
if err := src.Handshake(1, handler.chain); err != nil {
|
||||
t.Fatalf("failed to run protocol handshake")
|
||||
}
|
||||
// Send the transaction to the sink and verify that it's added to the tx pool
|
||||
|
|
@ -316,11 +311,7 @@ func testSendTransactions(t *testing.T, protocol uint) {
|
|||
return eth.Handle((*ethHandler)(handler.handler), peer)
|
||||
})
|
||||
// Run the handshake locally to avoid spinning up a source handler
|
||||
var (
|
||||
genesis = handler.chain.Genesis()
|
||||
head = handler.chain.CurrentBlock()
|
||||
)
|
||||
if err := sink.Handshake(1, head.Hash(), genesis.Hash(), forkid.NewIDWithChain(handler.chain), forkid.NewFilter(handler.chain)); err != nil {
|
||||
if err := sink.Handshake(1, handler.chain); err != nil {
|
||||
t.Fatalf("failed to run protocol handshake")
|
||||
}
|
||||
// After the handshake completes, the source handler should stream the sink
|
||||
|
|
|
|||
|
|
@ -192,6 +192,7 @@ var eth69 = map[uint64]msgHandler{
|
|||
ReceiptsMsg: handleReceipts69,
|
||||
GetPooledTransactionsMsg: handleGetPooledTransactions,
|
||||
PooledTransactionsMsg: handlePooledTransactions,
|
||||
BlockRangeUpdateMsg: handleBlockRangeUpdate,
|
||||
}
|
||||
|
||||
// handleMessage is invoked whenever an inbound message is received from a remote
|
||||
|
|
|
|||
|
|
@ -532,3 +532,12 @@ func handlePooledTransactions(backend Backend, msg Decoder, peer *Peer) error {
|
|||
|
||||
return backend.Handle(peer, &txs.PooledTransactionsResponse)
|
||||
}
|
||||
|
||||
func handleBlockRangeUpdate(backend Backend, msg Decoder, peer *Peer) error {
|
||||
var update BlockRangeUpdatePacket
|
||||
if err := msg.Decode(&update); err != nil {
|
||||
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
||||
}
|
||||
// We don't do anything with these messages for now.
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,6 +22,7 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/ethereum/go-ethereum/common"
|
||||
"github.com/ethereum/go-ethereum/core"
|
||||
"github.com/ethereum/go-ethereum/core/forkid"
|
||||
"github.com/ethereum/go-ethereum/metrics"
|
||||
"github.com/ethereum/go-ethereum/p2p"
|
||||
|
|
@ -35,44 +36,112 @@ const (
|
|||
|
||||
// Handshake executes the eth protocol handshake, negotiating version number,
|
||||
// network IDs, difficulties, head and genesis blocks.
|
||||
func (p *Peer) Handshake(network uint64, head common.Hash, genesis common.Hash, forkID forkid.ID, forkFilter forkid.Filter) error {
|
||||
// Send out own handshake in a new thread
|
||||
func (p *Peer) Handshake(networkID uint64, chain *core.BlockChain) error {
|
||||
switch p.version {
|
||||
case ETH69:
|
||||
return p.handshake69(networkID, chain)
|
||||
case ETH68:
|
||||
return p.handshake68(networkID, chain)
|
||||
default:
|
||||
return errors.New("unsupported protocol version")
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Peer) handshake68(networkID uint64, chain *core.BlockChain) error {
|
||||
var (
|
||||
genesis = chain.Genesis()
|
||||
latest = chain.CurrentBlock()
|
||||
forkID = forkid.NewID(chain.Config(), genesis, latest.Number.Uint64(), latest.Time)
|
||||
forkFilter = forkid.NewFilter(chain)
|
||||
)
|
||||
errc := make(chan error, 2)
|
||||
|
||||
var status StatusPacket // safe to read after two values have been received from errc
|
||||
|
||||
go func() {
|
||||
pkt := &StatusPacket{
|
||||
pkt := &StatusPacket68{
|
||||
ProtocolVersion: uint32(p.version),
|
||||
NetworkID: network,
|
||||
Head: head,
|
||||
Genesis: genesis,
|
||||
NetworkID: networkID,
|
||||
Head: latest.Hash(),
|
||||
Genesis: genesis.Hash(),
|
||||
ForkID: forkID,
|
||||
}
|
||||
errc <- p2p.Send(p.rw, StatusMsg, pkt)
|
||||
}()
|
||||
var status StatusPacket68 // safe to read after two values have been received from errc
|
||||
go func() {
|
||||
errc <- p.readStatus(network, &status, genesis, forkFilter)
|
||||
errc <- p.readStatus68(networkID, &status, genesis.Hash(), forkFilter)
|
||||
}()
|
||||
timeout := time.NewTimer(handshakeTimeout)
|
||||
defer timeout.Stop()
|
||||
for i := 0; i < 2; i++ {
|
||||
select {
|
||||
case err := <-errc:
|
||||
if err != nil {
|
||||
markError(p, err)
|
||||
|
||||
return waitForHandshake(errc, p)
|
||||
}
|
||||
|
||||
func (p *Peer) readStatus68(networkID uint64, status *StatusPacket68, genesis common.Hash, forkFilter forkid.Filter) error {
|
||||
if err := p.readStatusMsg(status); err != nil {
|
||||
return err
|
||||
}
|
||||
case <-timeout.C:
|
||||
markError(p, p2p.DiscReadTimeout)
|
||||
return p2p.DiscReadTimeout
|
||||
if status.NetworkID != networkID {
|
||||
return fmt.Errorf("%w: %d (!= %d)", errNetworkIDMismatch, status.NetworkID, networkID)
|
||||
}
|
||||
if uint(status.ProtocolVersion) != p.version {
|
||||
return fmt.Errorf("%w: %d (!= %d)", errProtocolVersionMismatch, status.ProtocolVersion, p.version)
|
||||
}
|
||||
if status.Genesis != genesis {
|
||||
return fmt.Errorf("%w: %x (!= %x)", errGenesisMismatch, status.Genesis, genesis)
|
||||
}
|
||||
if err := forkFilter(status.ForkID); err != nil {
|
||||
return fmt.Errorf("%w: %v", errForkIDRejected, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// readStatus reads the remote handshake message.
|
||||
func (p *Peer) readStatus(network uint64, status *StatusPacket, genesis common.Hash, forkFilter forkid.Filter) error {
|
||||
func (p *Peer) handshake69(networkID uint64, chain *core.BlockChain) error {
|
||||
var (
|
||||
genesis = chain.Genesis()
|
||||
latest = chain.CurrentBlock()
|
||||
forkID = forkid.NewID(chain.Config(), genesis, latest.Number.Uint64(), latest.Time)
|
||||
forkFilter = forkid.NewFilter(chain)
|
||||
earliest, _ = chain.HistoryPruningCutoff()
|
||||
)
|
||||
errc := make(chan error, 2)
|
||||
go func() {
|
||||
pkt := &StatusPacket69{
|
||||
ProtocolVersion: uint32(p.version),
|
||||
NetworkID: networkID,
|
||||
Genesis: genesis.Hash(),
|
||||
ForkID: forkID,
|
||||
EarliestBlock: earliest,
|
||||
LatestBlock: latest.Number.Uint64(),
|
||||
LatestBlockHash: latest.Hash(),
|
||||
}
|
||||
errc <- p2p.Send(p.rw, StatusMsg, pkt)
|
||||
}()
|
||||
var status StatusPacket69 // safe to read after two values have been received from errc
|
||||
go func() {
|
||||
errc <- p.readStatus69(networkID, &status, genesis.Hash(), forkFilter)
|
||||
}()
|
||||
|
||||
return waitForHandshake(errc, p)
|
||||
}
|
||||
|
||||
func (p *Peer) readStatus69(networkID uint64, status *StatusPacket69, genesis common.Hash, forkFilter forkid.Filter) error {
|
||||
if err := p.readStatusMsg(status); err != nil {
|
||||
return err
|
||||
}
|
||||
if status.NetworkID != networkID {
|
||||
return fmt.Errorf("%w: %d (!= %d)", errNetworkIDMismatch, status.NetworkID, networkID)
|
||||
}
|
||||
if uint(status.ProtocolVersion) != p.version {
|
||||
return fmt.Errorf("%w: %d (!= %d)", errProtocolVersionMismatch, status.ProtocolVersion, p.version)
|
||||
}
|
||||
if status.Genesis != genesis {
|
||||
return fmt.Errorf("%w: %x (!= %x)", errGenesisMismatch, status.Genesis, genesis)
|
||||
}
|
||||
if err := forkFilter(status.ForkID); err != nil {
|
||||
return fmt.Errorf("%w: %v", errForkIDRejected, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// readStatusMsg reads the first message on the connection.
|
||||
func (p *Peer) readStatusMsg(dst any) error {
|
||||
msg, err := p.rw.ReadMsg()
|
||||
if err != nil {
|
||||
return err
|
||||
|
|
@ -83,21 +152,26 @@ func (p *Peer) readStatus(network uint64, status *StatusPacket, genesis common.H
|
|||
if msg.Size > maxMessageSize {
|
||||
return fmt.Errorf("%w: %v > %v", errMsgTooLarge, msg.Size, maxMessageSize)
|
||||
}
|
||||
// Decode the handshake and make sure everything matches
|
||||
if err := msg.Decode(&status); err != nil {
|
||||
if err := msg.Decode(dst); err != nil {
|
||||
return fmt.Errorf("%w: message %v: %v", errDecode, msg, err)
|
||||
}
|
||||
if status.NetworkID != network {
|
||||
return fmt.Errorf("%w: %d (!= %d)", errNetworkIDMismatch, status.NetworkID, network)
|
||||
return nil
|
||||
}
|
||||
|
||||
func waitForHandshake(errc <-chan error, p *Peer) error {
|
||||
timeout := time.NewTimer(handshakeTimeout)
|
||||
defer timeout.Stop()
|
||||
for range 2 {
|
||||
select {
|
||||
case err := <-errc:
|
||||
if err != nil {
|
||||
markError(p, err)
|
||||
return err
|
||||
}
|
||||
if uint(status.ProtocolVersion) != p.version {
|
||||
return fmt.Errorf("%w: %d (!= %d)", errProtocolVersionMismatch, status.ProtocolVersion, p.version)
|
||||
case <-timeout.C:
|
||||
markError(p, p2p.DiscReadTimeout)
|
||||
return p2p.DiscReadTimeout
|
||||
}
|
||||
if status.Genesis != genesis {
|
||||
return fmt.Errorf("%w: %x (!= %x)", errGenesisMismatch, status.Genesis, genesis)
|
||||
}
|
||||
if err := forkFilter(status.ForkID); err != nil {
|
||||
return fmt.Errorf("%w: %v", errForkIDRejected, err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -52,19 +52,19 @@ func testHandshake(t *testing.T, protocol uint) {
|
|||
want: errNoStatusMsg,
|
||||
},
|
||||
{
|
||||
code: StatusMsg, data: StatusPacket{10, 1, new(big.Int), head.Hash(), genesis.Hash(), forkID},
|
||||
code: StatusMsg, data: StatusPacket68{10, 1, new(big.Int), head.Hash(), genesis.Hash(), forkID},
|
||||
want: errProtocolVersionMismatch,
|
||||
},
|
||||
{
|
||||
code: StatusMsg, data: StatusPacket{uint32(protocol), 999, new(big.Int), head.Hash(), genesis.Hash(), forkID},
|
||||
code: StatusMsg, data: StatusPacket68{uint32(protocol), 999, new(big.Int), head.Hash(), genesis.Hash(), forkID},
|
||||
want: errNetworkIDMismatch,
|
||||
},
|
||||
{
|
||||
code: StatusMsg, data: StatusPacket{uint32(protocol), 1, new(big.Int), head.Hash(), common.Hash{3}, forkID},
|
||||
code: StatusMsg, data: StatusPacket68{uint32(protocol), 1, new(big.Int), head.Hash(), common.Hash{3}, forkID},
|
||||
want: errGenesisMismatch,
|
||||
},
|
||||
{
|
||||
code: StatusMsg, data: StatusPacket{uint32(protocol), 1, new(big.Int), head.Hash(), genesis.Hash(), forkid.ID{Hash: [4]byte{0x00, 0x01, 0x02, 0x03}}},
|
||||
code: StatusMsg, data: StatusPacket68{uint32(protocol), 1, new(big.Int), head.Hash(), genesis.Hash(), forkid.ID{Hash: [4]byte{0x00, 0x01, 0x02, 0x03}}},
|
||||
want: errForkIDRejected,
|
||||
},
|
||||
}
|
||||
|
|
@ -80,7 +80,7 @@ func testHandshake(t *testing.T, protocol uint) {
|
|||
// Send the junk test with one peer, check the handshake failure
|
||||
go p2p.Send(app, test.code, test.data)
|
||||
|
||||
err := peer.Handshake(1, head.Hash(), genesis.Hash(), forkID, forkid.NewFilter(backend.chain))
|
||||
err := peer.Handshake(1, backend.chain)
|
||||
if err == nil {
|
||||
t.Errorf("test %d: protocol returned nil error, want %q", i, test.want)
|
||||
} else if !errors.Is(err, test.want) {
|
||||
|
|
|
|||
|
|
@ -343,6 +343,15 @@ func (p *Peer) RequestTxs(hashes []common.Hash) error {
|
|||
})
|
||||
}
|
||||
|
||||
// SendBlockRangeUpdate sends a notification about our available block range to the peer.
|
||||
func (p *Peer) SendBlockRangeUpdate(earliestBlock, latestBlock uint64, latestBlockHash common.Hash) error {
|
||||
return p2p.Send(p.rw, BlockRangeUpdateMsg, &BlockRangeUpdatePacket{
|
||||
EarliestBlock: earliestBlock,
|
||||
LatestBlock: latestBlock,
|
||||
LatestBlockHash: latestBlockHash,
|
||||
})
|
||||
}
|
||||
|
||||
// knownCache is a cache for known hashes.
|
||||
type knownCache struct {
|
||||
hashes mapset.Set[common.Hash]
|
||||
|
|
|
|||
|
|
@ -44,7 +44,7 @@ var ProtocolVersions = []uint{ETH69, ETH68}
|
|||
|
||||
// protocolLengths are the number of implemented message corresponding to
|
||||
// different protocol versions.
|
||||
var protocolLengths = map[uint]uint64{ETH68: 17, ETH69: 17}
|
||||
var protocolLengths = map[uint]uint64{ETH68: 17, ETH69: 18}
|
||||
|
||||
// maxMessageSize is the maximum cap on the size of a protocol message.
|
||||
const maxMessageSize = 10 * 1024 * 1024
|
||||
|
|
@ -63,6 +63,7 @@ const (
|
|||
PooledTransactionsMsg = 0x0a
|
||||
GetReceiptsMsg = 0x0f
|
||||
ReceiptsMsg = 0x10
|
||||
BlockRangeUpdateMsg = 0x11
|
||||
)
|
||||
|
||||
var (
|
||||
|
|
@ -83,7 +84,7 @@ type Packet interface {
|
|||
}
|
||||
|
||||
// StatusPacket is the network packet for the status message.
|
||||
type StatusPacket struct {
|
||||
type StatusPacket68 struct {
|
||||
ProtocolVersion uint32
|
||||
NetworkID uint64
|
||||
TD *big.Int
|
||||
|
|
@ -92,6 +93,18 @@ type StatusPacket struct {
|
|||
ForkID forkid.ID
|
||||
}
|
||||
|
||||
// StatusPacket69 is the network packet for the status message.
|
||||
type StatusPacket69 struct {
|
||||
ProtocolVersion uint32
|
||||
NetworkID uint64
|
||||
Genesis common.Hash
|
||||
ForkID forkid.ID
|
||||
// initial available block range
|
||||
EarliestBlock uint64
|
||||
LatestBlock uint64
|
||||
LatestBlockHash common.Hash
|
||||
}
|
||||
|
||||
// NewBlockHashesPacket is the network packet for the block announcements.
|
||||
type NewBlockHashesPacket []struct {
|
||||
Hash common.Hash // Hash of one particular block being announced
|
||||
|
|
@ -312,8 +325,18 @@ type PooledTransactionsRLPPacket struct {
|
|||
PooledTransactionsRLPResponse
|
||||
}
|
||||
|
||||
func (*StatusPacket) Name() string { return "Status" }
|
||||
func (*StatusPacket) Kind() byte { return StatusMsg }
|
||||
// BlockRangeUpdatePacket is an announcement of the node's available block range.
|
||||
type BlockRangeUpdatePacket struct {
|
||||
EarliestBlock uint64
|
||||
LatestBlock uint64
|
||||
LatestBlockHash common.Hash
|
||||
}
|
||||
|
||||
func (*StatusPacket68) Name() string { return "Status" }
|
||||
func (*StatusPacket68) Kind() byte { return StatusMsg }
|
||||
|
||||
func (*StatusPacket69) Name() string { return "Status" }
|
||||
func (*StatusPacket69) Kind() byte { return StatusMsg }
|
||||
|
||||
func (*NewBlockHashesPacket) Name() string { return "NewBlockHashes" }
|
||||
func (*NewBlockHashesPacket) Kind() byte { return NewBlockHashesMsg }
|
||||
|
|
@ -353,3 +376,6 @@ func (*ReceiptsResponse) Kind() byte { return ReceiptsMsg }
|
|||
|
||||
func (*ReceiptsRLPResponse) Name() string { return "Receipts" }
|
||||
func (*ReceiptsRLPResponse) Kind() byte { return ReceiptsMsg }
|
||||
|
||||
func (*BlockRangeUpdatePacket) Name() string { return "BlockRangeUpdate" }
|
||||
func (*BlockRangeUpdatePacket) Kind() byte { return BlockRangeUpdateMsg }
|
||||
|
|
|
|||
Loading…
Reference in a new issue