cmd, consensus, eth, ethstats: add protocol interface into consensus to support custom messages

This commit is contained in:
mark.lin 2017-09-12 11:29:12 +08:00
parent 933972d139
commit e22d86f6f9
12 changed files with 169 additions and 36 deletions

View file

@ -113,7 +113,6 @@ func version(ctx *cli.Context) error {
fmt.Println("Git Commit:", gitCommit) fmt.Println("Git Commit:", gitCommit)
} }
fmt.Println("Architecture:", runtime.GOARCH) fmt.Println("Architecture:", runtime.GOARCH)
fmt.Println("Protocol Versions:", eth.ProtocolVersions)
fmt.Println("Network Id:", eth.DefaultConfig.NetworkId) fmt.Println("Network Id:", eth.DefaultConfig.NetworkId)
fmt.Println("Go Version:", runtime.Version()) fmt.Println("Go Version:", runtime.Version())
fmt.Println("Operating System:", runtime.GOOS) fmt.Println("Operating System:", runtime.GOOS)

View file

@ -683,3 +683,8 @@ func (c *Clique) APIs(chain consensus.ChainReader) []rpc.API {
Public: false, Public: false,
}} }}
} }
// Protocol implements consensus.Engine.Protocol
func (c *Clique) Protocol() consensus.Protocol {
return consensus.EthProtocol
}

View file

@ -21,6 +21,7 @@ import (
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/state" "github.com/ethereum/go-ethereum/core/state"
"github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/params" "github.com/ethereum/go-ethereum/params"
"github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/rpc"
"math/big" "math/big"
@ -95,6 +96,21 @@ type Engine interface {
// APIs returns the RPC APIs this consensus engine provides. // APIs returns the RPC APIs this consensus engine provides.
APIs(chain ChainReader) []rpc.API APIs(chain ChainReader) []rpc.API
// Protocol returns the protocol for this consensus
Protocol() Protocol
}
// Handler should be implemented is the consensus needs to handle and send peer's message
type Handler interface {
// NewChainHead handles a new head block comes
NewChainHead() error
// HandleMsg handles a message from peer
HandleMsg(address common.Address, data p2p.Msg) (bool, error)
// SetBroadcaster sets the broadcaster to send message to peers
SetBroadcaster(Broadcaster)
} }
// PoW is a consensus engine based on proof-of-work. // PoW is a consensus engine based on proof-of-work.

View file

@ -576,3 +576,8 @@ func (ethash *Ethash) APIs(chain consensus.ChainReader) []rpc.API {
func SeedHash(block uint64) []byte { func SeedHash(block uint64) []byte {
return seedHash(block) return seedHash(block)
} }
// Protocol implements consensus.Engine.Protocol
func (ethash *Ethash) Protocol() consensus.Protocol {
return consensus.EthProtocol
}

61
consensus/protocol.go Normal file
View file

@ -0,0 +1,61 @@
// Copyright 2017 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 consensus implements different Ethereum consensus engines.
package consensus
import (
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/core/types"
)
// Constants to match up protocol versions and messages
const (
Eth62 = 62
Eth63 = 63
)
var (
EthProtocol = Protocol{
Name: "eth",
Versions: []uint{Eth62, Eth63},
Lengths: []uint64{17, 8},
}
)
// Protocol defines the protocol of the consensus
type Protocol struct {
// Official short name of the protocol used during capability negotiation.
Name string
// Supported versions of the eth protocol (first is primary).
Versions []uint
// Number of implemented message corresponding to different protocol versions.
Lengths []uint64
}
// Broadcaster defines the interface to enqueue blocks to fetcher and find peer
type Broadcaster interface {
// Enqueue add a block into fetcher queue
Enqueue(id string, block *types.Block)
// FindPeers retrives peers by addresses
FindPeers(map[common.Address]bool) map[common.Address]Peer
}
// Peer defines the interface to communicate with peer
type Peer interface {
// Send sends the message to this peer
Send(msgcode uint64, data interface{}) error
}

View file

@ -135,7 +135,7 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) {
bloomIndexer: NewBloomIndexer(chainDb, params.BloomBitsBlocks), bloomIndexer: NewBloomIndexer(chainDb, params.BloomBitsBlocks),
} }
log.Info("Initialising Ethereum protocol", "versions", ProtocolVersions, "network", config.NetworkId) log.Info("Initialising Ethereum protocol", "versions", eth.engine.Protocol().Versions, "network", config.NetworkId)
if !config.SkipBcVersionCheck { if !config.SkipBcVersionCheck {
bcVersion := core.GetBlockChainVersion(chainDb) bcVersion := core.GetBlockChainVersion(chainDb)

View file

@ -31,6 +31,7 @@ import (
"github.com/ethereum/go-ethereum/consensus/misc" "github.com/ethereum/go-ethereum/consensus/misc"
"github.com/ethereum/go-ethereum/core" "github.com/ethereum/go-ethereum/core"
"github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/eth/downloader" "github.com/ethereum/go-ethereum/eth/downloader"
"github.com/ethereum/go-ethereum/eth/fetcher" "github.com/ethereum/go-ethereum/eth/fetcher"
"github.com/ethereum/go-ethereum/ethdb" "github.com/ethereum/go-ethereum/ethdb"
@ -94,6 +95,8 @@ type ProtocolManager struct {
// wait group is used for graceful shutdowns during downloading // wait group is used for graceful shutdowns during downloading
// and processing // and processing
wg sync.WaitGroup wg sync.WaitGroup
engine consensus.Engine
} }
// NewProtocolManager returns a new ethereum sub protocol manager. The Ethereum sub protocol manages peers capable // NewProtocolManager returns a new ethereum sub protocol manager. The Ethereum sub protocol manages peers capable
@ -111,7 +114,13 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne
noMorePeers: make(chan struct{}), noMorePeers: make(chan struct{}),
txsyncCh: make(chan *txsync), txsyncCh: make(chan *txsync),
quitSync: make(chan struct{}), quitSync: make(chan struct{}),
engine: engine,
} }
if handler, ok := manager.engine.(consensus.Handler); ok {
handler.SetBroadcaster(manager)
}
// Figure out whether to allow fast sync or not // Figure out whether to allow fast sync or not
if mode == downloader.FastSync && blockchain.CurrentBlock().NumberU64() > 0 { if mode == downloader.FastSync && blockchain.CurrentBlock().NumberU64() > 0 {
log.Warn("Blockchain not empty, fast sync disabled") log.Warn("Blockchain not empty, fast sync disabled")
@ -120,19 +129,20 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne
if mode == downloader.FastSync { if mode == downloader.FastSync {
manager.fastSync = uint32(1) manager.fastSync = uint32(1)
} }
protocol := engine.Protocol()
// Initiate a sub-protocol for every implemented version we can handle // Initiate a sub-protocol for every implemented version we can handle
manager.SubProtocols = make([]p2p.Protocol, 0, len(ProtocolVersions)) manager.SubProtocols = make([]p2p.Protocol, 0, len(protocol.Versions))
for i, version := range ProtocolVersions { for i, version := range protocol.Versions {
// Skip protocol version if incompatible with the mode of operation // Skip protocol version if incompatible with the mode of operation
if mode == downloader.FastSync && version < eth63 { if mode == downloader.FastSync && version < consensus.Eth63 {
continue continue
} }
// Compatible; initialise the sub-protocol // Compatible; initialise the sub-protocol
version := version // Closure for the run version := version // Closure for the run
manager.SubProtocols = append(manager.SubProtocols, p2p.Protocol{ manager.SubProtocols = append(manager.SubProtocols, p2p.Protocol{
Name: ProtocolName, Name: protocol.Name,
Version: version, Version: version,
Length: ProtocolLengths[i], Length: protocol.Lengths[i],
Run: func(p *p2p.Peer, rw p2p.MsgReadWriter) error { Run: func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
peer := manager.newPeer(int(version), p, rw) peer := manager.newPeer(int(version), p, rw)
select { select {
@ -326,6 +336,18 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
} }
defer msg.Discard() defer msg.Discard()
if handler, ok := pm.engine.(consensus.Handler); ok {
pubKey, err := p.ID().Pubkey()
if err != nil {
return err
}
addr := crypto.PubkeyToAddress(*pubKey)
handled, err := handler.HandleMsg(addr, msg)
if handled {
return err
}
}
// Handle the message depending on its contents // Handle the message depending on its contents
switch { switch {
case msg.Code == StatusMsg: case msg.Code == StatusMsg:
@ -517,7 +539,7 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
} }
} }
case p.version >= eth63 && msg.Code == GetNodeDataMsg: case p.version >= consensus.Eth63 && msg.Code == GetNodeDataMsg:
// Decode the retrieval message // Decode the retrieval message
msgStream := rlp.NewStream(msg.Payload, uint64(msg.Size)) msgStream := rlp.NewStream(msg.Payload, uint64(msg.Size))
if _, err := msgStream.List(); err != nil { if _, err := msgStream.List(); err != nil {
@ -544,7 +566,7 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
} }
return p.SendNodeData(data) return p.SendNodeData(data)
case p.version >= eth63 && msg.Code == NodeDataMsg: case p.version >= consensus.Eth63 && msg.Code == NodeDataMsg:
// A batch of node state data arrived to one of our previous requests // A batch of node state data arrived to one of our previous requests
var data [][]byte var data [][]byte
if err := msg.Decode(&data); err != nil { if err := msg.Decode(&data); err != nil {
@ -555,7 +577,7 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
log.Debug("Failed to deliver node state data", "err", err) log.Debug("Failed to deliver node state data", "err", err)
} }
case p.version >= eth63 && msg.Code == GetReceiptsMsg: case p.version >= consensus.Eth63 && msg.Code == GetReceiptsMsg:
// Decode the retrieval message // Decode the retrieval message
msgStream := rlp.NewStream(msg.Payload, uint64(msg.Size)) msgStream := rlp.NewStream(msg.Payload, uint64(msg.Size))
if _, err := msgStream.List(); err != nil { if _, err := msgStream.List(); err != nil {
@ -591,7 +613,7 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
} }
return p.SendReceiptsRLP(receipts) return p.SendReceiptsRLP(receipts)
case p.version >= eth63 && msg.Code == ReceiptsMsg: case p.version >= consensus.Eth63 && msg.Code == ReceiptsMsg:
// A batch of receipts arrived to one of our previous requests // A batch of receipts arrived to one of our previous requests
var receipts [][]*types.Receipt var receipts [][]*types.Receipt
if err := msg.Decode(&receipts); err != nil { if err := msg.Decode(&receipts); err != nil {
@ -679,6 +701,10 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
return nil return nil
} }
func (pm *ProtocolManager) Enqueue(id string, block *types.Block) {
pm.fetcher.Enqueue(id, block)
}
// BroadcastBlock will either propagate a block to a subset of it's peers, or // BroadcastBlock will either propagate a block to a subset of it's peers, or
// will only announce it's availability (depending what's requested). // will only announce it's availability (depending what's requested).
func (pm *ProtocolManager) BroadcastBlock(block *types.Block, propagate bool) { func (pm *ProtocolManager) BroadcastBlock(block *types.Block, propagate bool) {
@ -770,3 +796,18 @@ func (self *ProtocolManager) NodeInfo() *NodeInfo {
Head: currentBlock.Hash(), Head: currentBlock.Hash(),
} }
} }
func (self *ProtocolManager) FindPeers(targets map[common.Address]bool) map[common.Address]consensus.Peer {
m := make(map[common.Address]consensus.Peer)
for _, p := range self.peers.Peers() {
pubKey, err := p.ID().Pubkey()
if err != nil {
continue
}
addr := crypto.PubkeyToAddress(*pubKey)
if targets[addr] {
m[addr] = p
}
}
return m
}

View file

@ -24,6 +24,7 @@ import (
"time" "time"
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/consensus"
"github.com/ethereum/go-ethereum/consensus/ethash" "github.com/ethereum/go-ethereum/consensus/ethash"
"github.com/ethereum/go-ethereum/core" "github.com/ethereum/go-ethereum/core"
"github.com/ethereum/go-ethereum/core/state" "github.com/ethereum/go-ethereum/core/state"
@ -49,12 +50,12 @@ func TestProtocolCompatibility(t *testing.T) {
{61, downloader.FastSync, false}, {62, downloader.FastSync, false}, {63, downloader.FastSync, true}, {61, downloader.FastSync, false}, {62, downloader.FastSync, false}, {63, downloader.FastSync, true},
} }
// Make sure anything we screw up is restored // Make sure anything we screw up is restored
backup := ProtocolVersions backup := consensus.EthProtocol.Versions
defer func() { ProtocolVersions = backup }() defer func() { consensus.EthProtocol.Versions = backup }()
// Try all available compatibility configs and check for errors // Try all available compatibility configs and check for errors
for i, tt := range tests { for i, tt := range tests {
ProtocolVersions = []uint{tt.version} consensus.EthProtocol.Versions = []uint{tt.version}
pm, _, err := newTestProtocolManager(tt.mode, 0, nil, nil) pm, _, err := newTestProtocolManager(tt.mode, 0, nil, nil)
if pm != nil { if pm != nil {
@ -482,7 +483,7 @@ func testDAOChallenge(t *testing.T, localForked, remoteForked bool, timeout bool
defer pm.Stop() defer pm.Stop()
// Connect a new peer and check that we receive the DAO challenge // Connect a new peer and check that we receive the DAO challenge
peer, _ := newTestPeer("peer", eth63, pm, true) peer, _ := newTestPeer("peer", consensus.Eth63, pm, true)
defer peer.close() defer peer.close()
challenge := &getBlockHeadersData{ challenge := &getBlockHeadersData{

View file

@ -17,6 +17,7 @@
package eth package eth
import ( import (
"github.com/ethereum/go-ethereum/consensus"
"github.com/ethereum/go-ethereum/metrics" "github.com/ethereum/go-ethereum/metrics"
"github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p"
) )
@ -92,9 +93,9 @@ func (rw *meteredMsgReadWriter) ReadMsg() (p2p.Msg, error) {
case msg.Code == BlockBodiesMsg: case msg.Code == BlockBodiesMsg:
packets, traffic = reqBodyInPacketsMeter, reqBodyInTrafficMeter packets, traffic = reqBodyInPacketsMeter, reqBodyInTrafficMeter
case rw.version >= eth63 && msg.Code == NodeDataMsg: case rw.version >= consensus.Eth63 && msg.Code == NodeDataMsg:
packets, traffic = reqStateInPacketsMeter, reqStateInTrafficMeter packets, traffic = reqStateInPacketsMeter, reqStateInTrafficMeter
case rw.version >= eth63 && msg.Code == ReceiptsMsg: case rw.version >= consensus.Eth63 && msg.Code == ReceiptsMsg:
packets, traffic = reqReceiptInPacketsMeter, reqReceiptInTrafficMeter packets, traffic = reqReceiptInPacketsMeter, reqReceiptInTrafficMeter
case msg.Code == NewBlockHashesMsg: case msg.Code == NewBlockHashesMsg:
@ -119,9 +120,9 @@ func (rw *meteredMsgReadWriter) WriteMsg(msg p2p.Msg) error {
case msg.Code == BlockBodiesMsg: case msg.Code == BlockBodiesMsg:
packets, traffic = reqBodyOutPacketsMeter, reqBodyOutTrafficMeter packets, traffic = reqBodyOutPacketsMeter, reqBodyOutTrafficMeter
case rw.version >= eth63 && msg.Code == NodeDataMsg: case rw.version >= consensus.Eth63 && msg.Code == NodeDataMsg:
packets, traffic = reqStateOutPacketsMeter, reqStateOutTrafficMeter packets, traffic = reqStateOutPacketsMeter, reqStateOutTrafficMeter
case rw.version >= eth63 && msg.Code == ReceiptsMsg: case rw.version >= consensus.Eth63 && msg.Code == ReceiptsMsg:
packets, traffic = reqReceiptOutPacketsMeter, reqReceiptOutTrafficMeter packets, traffic = reqReceiptOutPacketsMeter, reqReceiptOutTrafficMeter
case msg.Code == NewBlockHashesMsg: case msg.Code == NewBlockHashesMsg:

View file

@ -130,6 +130,12 @@ func (p *peer) MarkTransaction(hash common.Hash) {
p.knownTxs.Add(hash) p.knownTxs.Add(hash)
} }
// Send writes an RLP-encoded message with the given code.
// data should encode as an RLP list.
func (p *peer) Send(msgcode uint64, data interface{}) error {
return p2p.Send(p.rw, msgcode, data)
}
// SendTransactions sends transactions to the peer and includes the hashes // SendTransactions sends transactions to the peer and includes the hashes
// in its transaction hash set for future reference. // in its transaction hash set for future reference.
func (p *peer) SendTransactions(txs types.Transactions) error { func (p *peer) SendTransactions(txs types.Transactions) error {
@ -341,6 +347,18 @@ func (ps *peerSet) Unregister(id string) error {
return nil return nil
} }
// Peers returns all registered peers
func (ps *peerSet) Peers() map[string]*peer {
ps.lock.RLock()
defer ps.lock.RUnlock()
set := make(map[string]*peer)
for id, p := range ps.peers {
set[id] = p
}
return set
}
// Peer retrieves the registered peer with the given id. // Peer retrieves the registered peer with the given id.
func (ps *peerSet) Peer(id string) *peer { func (ps *peerSet) Peer(id string) *peer {
ps.lock.RLock() ps.lock.RLock()

View file

@ -28,21 +28,6 @@ import (
"github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rlp"
) )
// Constants to match up protocol versions and messages
const (
eth62 = 62
eth63 = 63
)
// Official short name of the protocol used during capability negotiation.
var ProtocolName = "eth"
// Supported versions of the eth protocol (first is primary).
var ProtocolVersions = []uint{eth63, eth62}
// Number of implemented message corresponding to different protocol versions.
var ProtocolLengths = []uint64{17, 8}
const ProtocolMaxMsgSize = 10 * 1024 * 1024 // Maximum cap on the size of a protocol message const ProtocolMaxMsgSize = 10 * 1024 * 1024 // Maximum cap on the size of a protocol message
// eth protocol message codes // eth protocol message codes

View file

@ -373,9 +373,10 @@ func (s *Service) login(conn *websocket.Conn) error {
infos := s.server.NodeInfo() infos := s.server.NodeInfo()
var network, protocol string var network, protocol string
if info := infos.Protocols["eth"]; info != nil { p := s.engine.Protocol()
if info := infos.Protocols[p.Name]; info != nil {
network = fmt.Sprintf("%d", info.(*eth.NodeInfo).Network) network = fmt.Sprintf("%d", info.(*eth.NodeInfo).Network)
protocol = fmt.Sprintf("eth/%d", eth.ProtocolVersions[0]) protocol = fmt.Sprintf("%s/%d", p.Name, p.Versions[0])
} else { } else {
network = fmt.Sprintf("%d", infos.Protocols["les"].(*les.NodeInfo).Network) network = fmt.Sprintf("%d", infos.Protocols["les"].(*les.NodeInfo).Network)
protocol = fmt.Sprintf("les/%d", les.ClientProtocolVersions[0]) protocol = fmt.Sprintf("les/%d", les.ClientProtocolVersions[0])