p2p/simulations: Track message events

Signed-off-by: Lewis Marshall <lewis@lmars.net>
This commit is contained in:
Lewis Marshall 2017-05-03 15:13:45 +01:00
parent 12eb088e5b
commit d3543fb458
7 changed files with 88 additions and 43 deletions

View file

@ -30,8 +30,6 @@ import (
"github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rlp"
) )
// Msg defines the structure of a p2p message. // Msg defines the structure of a p2p message.
@ -276,56 +274,54 @@ func ExpectMsg(r MsgReader, code uint64, content interface{}) error {
return nil return nil
} }
// MsgEventer wraps a MsgReadWriter and sends events whenever a message is sent
// wraps a msgreadwriter to allow for emitting message events upon // or received
// send or receive type MsgEventer struct {
type MsgReporterRW struct {
MsgReadWriter MsgReadWriter
feed *event.Feed
peerid discover.NodeID feed *event.Feed
closefunc func() error peerID discover.NodeID
} }
func NewMsgReporterRW(feed *event.Feed, rw MsgReadWriter, id discover.NodeID, closefunc func() error) *MsgReporterRW { func NewMsgEventer(rw MsgReadWriter, feed *event.Feed, peerID discover.NodeID) *MsgEventer {
return &MsgReporterRW{ return &MsgEventer{
MsgReadWriter: rw, MsgReadWriter: rw,
feed: feed, feed: feed,
peerid: id, peerID: peerID,
closefunc: closefunc,
} }
} }
func (self *MsgReporterRW) ReadMsg() (Msg, error) { func (self *MsgEventer) ReadMsg() (Msg, error) {
msg, err := self.MsgReadWriter.ReadMsg() msg, err := self.MsgReadWriter.ReadMsg()
if err != nil { if err != nil {
return msg, err return msg, err
} }
event := PeerEvent{ self.feed.Send(&PeerEvent{
Type: PeerEventTypeMsgRecv, Type: PeerEventTypeMsgRecv,
Peer: self.peerid, Peer: self.peerID,
Label: fmt.Sprintf("%d,%d", msg.Code, msg.Size), MsgCode: &msg.Code,
} MsgSize: &msg.Size,
self.feed.Send(event) })
return msg, nil return msg, nil
} }
func (self *MsgReporterRW) WriteMsg(msg Msg) error { func (self *MsgEventer) WriteMsg(msg Msg) error {
err := self.MsgReadWriter.WriteMsg(msg) err := self.MsgReadWriter.WriteMsg(msg)
if err != nil { if err != nil {
return err return err
} }
event := PeerEvent{ self.feed.Send(&PeerEvent{
Type: PeerEventTypeMsgSend, Type: PeerEventTypeMsgSend,
Peer: self.peerid, Peer: self.peerID,
Label: fmt.Sprintf("%d,%d", msg.Code, msg.Size), MsgCode: &msg.Code,
} MsgSize: &msg.Size,
self.feed.Send(event) })
return nil return nil
} }
func (self *MsgReporterRW) Close() error { func (self *MsgEventer) Close() error {
if self.closefunc != nil { if v, ok := self.MsgReadWriter.(io.Closer); ok {
return self.closefunc() return v.Close()
} }
return nil return nil
} }

View file

@ -25,6 +25,7 @@ import (
"time" "time"
"github.com/ethereum/go-ethereum/common/mclock" "github.com/ethereum/go-ethereum/common/mclock"
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rlp"
@ -82,12 +83,13 @@ const (
) )
// PeerEvent is an event emitted when peers are either added or dropped from // PeerEvent is an event emitted when peers are either added or dropped from
// a p2p.Server // a p2p.Server or when a message is sent or received on a peer connection
type PeerEvent struct { type PeerEvent struct {
Type PeerEventType Type PeerEventType `json:"type"`
Peer discover.NodeID Peer discover.NodeID `json:"peer"`
Error string Error string `json:"error,omitempty"`
Label string MsgCode *uint64 `json:"msg_code,omitempty"`
MsgSize *uint32 `json:"msg_size,omitempty"`
} }
// Peer represents a connected remote node. // Peer represents a connected remote node.
@ -101,6 +103,9 @@ type Peer struct {
protoErr chan error protoErr chan error
closed chan struct{} closed chan struct{}
disc chan DiscReason disc chan DiscReason
// events receives message send / receive events if set
events *event.Feed
} }
// NewPeer returns a peer for testing purposes. // NewPeer returns a peer for testing purposes.
@ -326,9 +331,13 @@ func (p *Peer) startProtocols(writeStart <-chan struct{}, writeErr chan<- error)
proto.closed = p.closed proto.closed = p.closed
proto.wstart = writeStart proto.wstart = writeStart
proto.werr = writeErr proto.werr = writeErr
var rw MsgReadWriter = proto
if p.events != nil {
rw = NewMsgEventer(rw, p.events, p.ID())
}
p.log.Trace(fmt.Sprintf("Starting protocol %s/%d", proto.Name, proto.Version)) p.log.Trace(fmt.Sprintf("Starting protocol %s/%d", proto.Name, proto.Version))
go func() { go func() {
err := proto.Run(p, proto) err := proto.Run(p, rw)
if err == nil { if err == nil {
p.log.Trace(fmt.Sprintf("Protocol %s/%d returned", proto.Name, proto.Version)) p.log.Trace(fmt.Sprintf("Protocol %s/%d returned", proto.Name, proto.Version))
err = errProtocolReturned err = errProtocolReturned
@ -391,10 +400,6 @@ func (rw *protoRW) ReadMsg() (Msg, error) {
} }
} }
func (rw *protoRW) Close() error {
return nil
}
// PeerInfo represents a short summary of the information known about a connected // PeerInfo represents a short summary of the information known about a connected
// peer. Sub-protocol independent fields are contained and initialized here, with // peer. Sub-protocol independent fields are contained and initialized here, with
// protocol specifics delegated to all connected sub-protocols. // protocol specifics delegated to all connected sub-protocols.

View file

@ -135,6 +135,10 @@ type Config struct {
// If NoDial is true, the server will not dial any peers. // If NoDial is true, the server will not dial any peers.
NoDial bool `toml:",omitempty"` NoDial bool `toml:",omitempty"`
// If EnableMsgEvents is set then the server will emit PeerEvents
// whenever a message is sent to or received from a peer
EnableMsgEvents bool
} }
type Server interface { type Server interface {
@ -567,6 +571,11 @@ running:
if err == nil { if err == nil {
// The handshakes are done and it passed all checks. // The handshakes are done and it passed all checks.
p := newPeer(c, srv.Protocols) p := newPeer(c, srv.Protocols)
// If message events are enabled, pass the peerFeed
// to the peer
if srv.EnableMsgEvents {
p.events = &srv.peerFeed
}
name := truncateName(c.name) name := truncateName(c.name)
log.Debug("Adding p2p peer", "id", c.id, "name", name, "addr", c.fd.RemoteAddr(), "peers", len(peers)+1) log.Debug("Adding p2p peer", "id", c.id, "name", name, "addr", c.fd.RemoteAddr(), "peers", len(peers)+1)
peers[c.id] = p peers[c.id] = p

View file

@ -49,6 +49,7 @@ func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) {
// generate the config // generate the config
conf := node.DefaultConfig conf := node.DefaultConfig
conf.DataDir = "/data" conf.DataDir = "/data"
conf.P2P.EnableMsgEvents = true
conf.P2P.NoDiscovery = true conf.P2P.NoDiscovery = true
conf.P2P.NAT = nil conf.P2P.NAT = nil

View file

@ -62,6 +62,7 @@ func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) {
// generate the config // generate the config
conf := node.DefaultConfig conf := node.DefaultConfig
conf.DataDir = filepath.Join(dir, "data") conf.DataDir = filepath.Join(dir, "data")
conf.P2P.EnableMsgEvents = true
conf.P2P.NoDiscovery = true conf.P2P.NoDiscovery = true
conf.P2P.NAT = nil conf.P2P.NAT = nil

View file

@ -234,7 +234,9 @@ func (self *SimNode) AddPeer(peer *discover.Node) {
if !exists { if !exists {
panic(fmt.Sprintf("unknown peer: %s", peer.ID)) panic(fmt.Sprintf("unknown peer: %s", peer.ID))
} }
localRW, peerRW := p2p.MsgPipe() p1, p2 := p2p.MsgPipe()
localRW := NewMsgReporter(p1, &self.peerFeed, peer.ID)
peerRW := NewMsgReporter(p2, &self.peerFeed, self.Id.NodeID)
self.peers[peer.ID] = peerRW self.peers[peer.ID] = peerRW
peerNode.RunProtocol(self, peerRW) peerNode.RunProtocol(self, peerRW)
self.RunProtocol(peerNode, localRW) self.RunProtocol(peerNode, localRW)

View file

@ -400,7 +400,8 @@ func (self *Conn) String() string {
type Msg struct { type Msg struct {
One *adapters.NodeId `json:"one"` One *adapters.NodeId `json:"one"`
Other *adapters.NodeId `json:"other"` Other *adapters.NodeId `json:"other"`
Code uint64 `json:"conn"` Code uint64 `json:"code"`
Received bool `json:"received"`
controlFired bool controlFired bool
} }
@ -624,6 +625,14 @@ func (self *Network) watchPeerEvents(id *adapters.NodeId, events chan *p2p.PeerE
if err := self.DidDisconnect(id, peer); err != nil { if err := self.DidDisconnect(id, peer); err != nil {
log.Error(fmt.Sprintf("error generating connection down event %s => %s", id.Label(), peer.Label()), "err", err) log.Error(fmt.Sprintf("error generating connection down event %s => %s", id.Label(), peer.Label()), "err", err)
} }
case p2p.PeerEventTypeMsgSend:
if err := self.DidSend(id, peer, *event.MsgCode); err != nil {
log.Error(fmt.Sprintf("error generating msg send event %s => %s", id.Label(), peer.Label()), "err", err)
}
case p2p.PeerEventTypeMsgRecv:
if err := self.DidReceive(peer, id, *event.MsgCode); err != nil {
log.Error(fmt.Sprintf("error generating msg receive event %s => %s", peer.Label(), id.Label()), "err", err)
}
} }
case err := <-sub.Err(): case err := <-sub.Err():
if err != nil { if err != nil {
@ -768,6 +777,28 @@ func (self *Network) Send(senderid, receiverid *adapters.NodeId, msgcode uint64,
self.events.Post(msg.EmitEvent(ControlEvent)) self.events.Post(msg.EmitEvent(ControlEvent))
} }
func (self *Network) DidSend(sender, receiver *adapters.NodeId, msgcode uint64) error {
msg := &Msg{
One: sender,
Other: receiver,
Code: msgcode,
Received: false,
}
self.events.Post(msg.EmitEvent(LiveEvent))
return nil
}
func (self *Network) DidReceive(sender, receiver *adapters.NodeId, msgcode uint64) error {
msg := &Msg{
One: sender,
Other: receiver,
Code: msgcode,
Received: true,
}
self.events.Post(msg.EmitEvent(LiveEvent))
return nil
}
// GetNode retrieves the node model for the id given as arg // GetNode retrieves the node model for the id given as arg
// returns nil if the node does not exist // returns nil if the node does not exist
func (self *Network) GetNode(id *adapters.NodeId) *Node { func (self *Network) GetNode(id *adapters.NodeId) *Node {