diff --git a/p2p/message.go b/p2p/message.go index 1b44a9326d..d7b5b5335d 100644 --- a/p2p/message.go +++ b/p2p/message.go @@ -30,8 +30,6 @@ import ( "github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/rlp" - - ) // Msg defines the structure of a p2p message. @@ -276,56 +274,54 @@ func ExpectMsg(r MsgReader, code uint64, content interface{}) error { return nil } - -// wraps a msgreadwriter to allow for emitting message events upon -// send or receive -type MsgReporterRW struct { +// MsgEventer wraps a MsgReadWriter and sends events whenever a message is sent +// or received +type MsgEventer struct { MsgReadWriter - feed *event.Feed - peerid discover.NodeID - closefunc func() error + + feed *event.Feed + peerID discover.NodeID } -func NewMsgReporterRW(feed *event.Feed, rw MsgReadWriter, id discover.NodeID, closefunc func() error) *MsgReporterRW { - return &MsgReporterRW{ +func NewMsgEventer(rw MsgReadWriter, feed *event.Feed, peerID discover.NodeID) *MsgEventer { + return &MsgEventer{ MsgReadWriter: rw, - feed: feed, - peerid: id, - closefunc: closefunc, + feed: feed, + peerID: peerID, } } -func (self *MsgReporterRW) ReadMsg() (Msg, error) { +func (self *MsgEventer) ReadMsg() (Msg, error) { msg, err := self.MsgReadWriter.ReadMsg() if err != nil { return msg, err } - event := PeerEvent{ - Type: PeerEventTypeMsgRecv, - Peer: self.peerid, - Label: fmt.Sprintf("%d,%d", msg.Code, msg.Size), - } - self.feed.Send(event) + self.feed.Send(&PeerEvent{ + Type: PeerEventTypeMsgRecv, + Peer: self.peerID, + MsgCode: &msg.Code, + MsgSize: &msg.Size, + }) return msg, nil } -func (self *MsgReporterRW) WriteMsg(msg Msg) error { +func (self *MsgEventer) WriteMsg(msg Msg) error { err := self.MsgReadWriter.WriteMsg(msg) if err != nil { return err } - event := PeerEvent{ - Type: PeerEventTypeMsgSend, - Peer: self.peerid, - Label: fmt.Sprintf("%d,%d", msg.Code, msg.Size), - } - self.feed.Send(event) + self.feed.Send(&PeerEvent{ + Type: PeerEventTypeMsgSend, + Peer: self.peerID, + MsgCode: &msg.Code, + MsgSize: &msg.Size, + }) return nil } -func (self *MsgReporterRW) Close() error { - if self.closefunc != nil { - return self.closefunc() +func (self *MsgEventer) Close() error { + if v, ok := self.MsgReadWriter.(io.Closer); ok { + return v.Close() } return nil } diff --git a/p2p/peer.go b/p2p/peer.go index c17606f03b..94b5b3fad8 100644 --- a/p2p/peer.go +++ b/p2p/peer.go @@ -25,6 +25,7 @@ import ( "time" "github.com/ethereum/go-ethereum/common/mclock" + "github.com/ethereum/go-ethereum/event" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/rlp" @@ -82,12 +83,13 @@ const ( ) // 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 PeerEventType - Peer discover.NodeID - Error string - Label string + Type PeerEventType `json:"type"` + Peer discover.NodeID `json:"peer"` + Error string `json:"error,omitempty"` + MsgCode *uint64 `json:"msg_code,omitempty"` + MsgSize *uint32 `json:"msg_size,omitempty"` } // Peer represents a connected remote node. @@ -101,6 +103,9 @@ type Peer struct { protoErr chan error closed chan struct{} disc chan DiscReason + + // events receives message send / receive events if set + events *event.Feed } // 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.wstart = writeStart 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)) go func() { - err := proto.Run(p, proto) + err := proto.Run(p, rw) if err == nil { p.log.Trace(fmt.Sprintf("Protocol %s/%d returned", proto.Name, proto.Version)) 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 // peer. Sub-protocol independent fields are contained and initialized here, with // protocol specifics delegated to all connected sub-protocols. diff --git a/p2p/server.go b/p2p/server.go index 0f21642ee8..c641e3a628 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -135,6 +135,10 @@ type Config struct { // If NoDial is true, the server will not dial any peers. 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 { @@ -567,6 +571,11 @@ running: if err == nil { // The handshakes are done and it passed all checks. 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) log.Debug("Adding p2p peer", "id", c.id, "name", name, "addr", c.fd.RemoteAddr(), "peers", len(peers)+1) peers[c.id] = p diff --git a/p2p/simulations/adapters/docker.go b/p2p/simulations/adapters/docker.go index 7c768d17a0..bc801a8a94 100644 --- a/p2p/simulations/adapters/docker.go +++ b/p2p/simulations/adapters/docker.go @@ -49,6 +49,7 @@ func (d *DockerAdapter) NewNode(config *NodeConfig) (Node, error) { // generate the config conf := node.DefaultConfig conf.DataDir = "/data" + conf.P2P.EnableMsgEvents = true conf.P2P.NoDiscovery = true conf.P2P.NAT = nil diff --git a/p2p/simulations/adapters/exec.go b/p2p/simulations/adapters/exec.go index 12167a271d..bebf4ea853 100644 --- a/p2p/simulations/adapters/exec.go +++ b/p2p/simulations/adapters/exec.go @@ -62,6 +62,7 @@ func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) { // generate the config conf := node.DefaultConfig conf.DataDir = filepath.Join(dir, "data") + conf.P2P.EnableMsgEvents = true conf.P2P.NoDiscovery = true conf.P2P.NAT = nil diff --git a/p2p/simulations/adapters/inproc.go b/p2p/simulations/adapters/inproc.go index ea6285ddf6..bfbba74eeb 100644 --- a/p2p/simulations/adapters/inproc.go +++ b/p2p/simulations/adapters/inproc.go @@ -234,7 +234,9 @@ func (self *SimNode) AddPeer(peer *discover.Node) { if !exists { 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 peerNode.RunProtocol(self, peerRW) self.RunProtocol(peerNode, localRW) diff --git a/p2p/simulations/network.go b/p2p/simulations/network.go index e4953e8f83..ef408121c1 100644 --- a/p2p/simulations/network.go +++ b/p2p/simulations/network.go @@ -400,7 +400,8 @@ func (self *Conn) String() string { type Msg struct { One *adapters.NodeId `json:"one"` Other *adapters.NodeId `json:"other"` - Code uint64 `json:"conn"` + Code uint64 `json:"code"` + Received bool `json:"received"` 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 { 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(): if err != nil { @@ -768,6 +777,28 @@ func (self *Network) Send(senderid, receiverid *adapters.NodeId, msgcode uint64, 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 // returns nil if the node does not exist func (self *Network) GetNode(id *adapters.NodeId) *Node {