diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index b68a4da5e6..a298c3ae84 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -23,6 +23,7 @@ import ( "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/hexutil" "github.com/ethereum/go-ethereum/common/mclock" + "github.com/ethereum/go-ethereum/metrics" "github.com/ethereum/go-ethereum/p2p/discover/v5wire" "github.com/VictoriaMetrics/fastcache" @@ -43,6 +44,26 @@ import ( "go.uber.org/zap" ) +var ( + messagesReceivedAccept metrics.Meter + messagesReceivedNodes metrics.Meter + messagesReceivedFindNodes metrics.Meter + messagesReceivedFindContent metrics.Meter + messagesReceivedContent metrics.Meter + messagesReceivedOffer metrics.Meter + messagesReceivedPing metrics.Meter + messagesReceivedPong metrics.Meter + + messagesSentAccept metrics.Meter + messagesSentNodes metrics.Meter + messagesSentFindNodes metrics.Meter + messagesSentFindContent metrics.Meter + messagesSentContent metrics.Meter + messagesSentOffer metrics.Meter + messagesSentPing metrics.Meter + messagesSentPong metrics.Meter +) + const ( // TalkResp message is a response message so the session is established and a @@ -237,6 +258,26 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId portalwire.Proto opt(protocol) } + if metrics.Enabled { + messagesReceivedAccept = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/accept", nil) + messagesReceivedNodes = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/nodes", nil) + messagesReceivedFindNodes = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/find_nodes", nil) + messagesReceivedFindContent = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/find_content", nil) + messagesReceivedContent = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/content", nil) + messagesReceivedOffer = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/offer", nil) + messagesReceivedPing = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/ping", nil) + messagesReceivedPong = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/received/pong", nil) + + messagesSentAccept = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/accept", nil) + messagesSentNodes = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/nodes", nil) + messagesSentFindNodes = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/find_nodes", nil) + messagesSentFindContent = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/find_content", nil) + messagesSentContent = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/content", nil) + messagesSentOffer = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/offer", nil) + messagesSentPing = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/ping", nil) + messagesSentPong = metrics.NewRegisteredMeter("portal/"+protocolId.Name()+"/sent/pong", nil) + } + return protocol, nil } @@ -428,6 +469,9 @@ func (p *PortalProtocol) pingInner(node *enode.Node) (*portalwire.Pong, error) { } p.Log.Trace(">> PING/"+p.protocolName, "protocol", p.protocolName, "ip", p.Self().IP().String(), "source", p.Self().ID(), "target", node.ID(), "ping", pingRequest) + if metrics.Enabled { + messagesSentPing.Mark(1) + } pingRequestBytes, err := pingRequest.MarshalSSZ() if err != nil { return nil, err @@ -444,6 +488,9 @@ func (p *PortalProtocol) pingInner(node *enode.Node) (*portalwire.Pong, error) { } p.Log.Trace("<< PONG/"+p.protocolName, "source", p.Self().ID(), "target", node.ID(), "res", talkResp) + if metrics.Enabled { + messagesReceivedPong.Mark(1) + } return p.processPong(node, talkResp) } @@ -463,6 +510,9 @@ func (p *PortalProtocol) findNodes(node *enode.Node, distances []uint) ([]*enode } p.Log.Trace(">> FIND_NODES/"+p.protocolName, "id", node.ID(), "findNodes", findNodes) + if metrics.Enabled { + messagesSentFindNodes.Mark(1) + } findNodesBytes, err := findNodes.MarshalSSZ() if err != nil { p.Log.Error("failed to marshal find nodes request", "err", err) @@ -488,6 +538,9 @@ func (p *PortalProtocol) findContent(node *enode.Node, contentKey []byte) (byte, } p.Log.Trace(">> FIND_CONTENT/"+p.protocolName, "id", node.ID(), "findContent", findContent) + if metrics.Enabled { + messagesSentFindContent.Mark(1) + } findContentBytes, err := findContent.MarshalSSZ() if err != nil { p.Log.Error("failed to marshal find content request", "err", err) @@ -515,6 +568,9 @@ func (p *PortalProtocol) offer(node *enode.Node, offerRequest *OfferRequest) ([] } p.Log.Trace(">> OFFER/"+p.protocolName, "offer", offer) + if metrics.Enabled { + messagesSentOffer.Mark(1) + } offerBytes, err := offer.MarshalSSZ() if err != nil { p.Log.Error("failed to marshal offer request", "err", err) @@ -553,6 +609,9 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * } p.Log.Trace("<< ACCEPT/"+p.protocolName, "id", target.ID(), "accept", accept) + if metrics.Enabled { + messagesReceivedAccept.Mark(1) + } isAdded := p.table.addFoundNode(target, true) if isAdded { log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) @@ -648,6 +707,9 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * return } p.Log.Trace(">> CONTENT/"+p.protocolName, "id", target.ID(), "contents", contents, "size", written) + if metrics.Enabled { + messagesSentContent.Mark(1) + } return } } @@ -677,6 +739,9 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, } p.Log.Trace("<< CONTENT/"+p.protocolName, "id", target.ID(), "content", content) + if metrics.Enabled { + messagesReceivedContent.Mark(1) + } isAdded := p.table.addFoundNode(target, true) if isAdded { log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) @@ -692,6 +757,9 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, } p.Log.Trace("<< CONTENT_CONNECTION_ID/"+p.protocolName, "id", target.ID(), "resp", common.Bytes2Hex(resp), "connIdMsg", connIdMsg) + if metrics.Enabled { + messagesReceivedContent.Mark(1) + } isAdded := p.table.addFoundNode(target, true) if isAdded { log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) @@ -729,6 +797,9 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, return 0xff, nil, err } p.Log.Trace("<< CONTENT/"+p.protocolName, "id", target.ID(), "size", len(data), "data", data) + if metrics.Enabled { + messagesReceivedContent.Mark(1) + } return resp[1], data, nil case portalwire.ContentEnrsSelector: enrs := &portalwire.Enrs{} @@ -739,6 +810,9 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, } p.Log.Trace("<< CONTENT_ENRS/"+p.protocolName, "id", target.ID(), "enrs", enrs) + if metrics.Enabled { + messagesReceivedContent.Mark(1) + } isAdded := p.table.addFoundNode(target, true) if isAdded { log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) @@ -804,6 +878,9 @@ func (p *PortalProtocol) filterNodes(target *enode.Node, enrs [][]byte, distance } p.Log.Trace("<< NODES/"+p.protocolName, "id", target.ID(), "total", len(enrs), "verified", verified, "nodes", nodes) + if metrics.Enabled { + messagesReceivedNodes.Mark(1) + } return nodes } @@ -821,6 +898,9 @@ func (p *PortalProtocol) processPong(target *enode.Node, resp []byte) (*portalwi } p.Log.Trace("<< PONG_RESPONSE/"+p.protocolName, "id", target.ID(), "pong", pong) + if metrics.Enabled { + messagesReceivedPong.Mark(1) + } customPayload := &portalwire.PingPongCustomData{} err = customPayload.UnmarshalSSZ(pong.CustomPayload) @@ -829,6 +909,9 @@ func (p *PortalProtocol) processPong(target *enode.Node, resp []byte) (*portalwi } p.Log.Trace("<< PONG_RESPONSE/"+p.protocolName, "id", target.ID(), "pong", pong, "customPayload", customPayload) + if metrics.Enabled { + messagesReceivedPong.Mark(1) + } isAdded := p.table.addFoundNode(target, true) if isAdded { log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) @@ -869,6 +952,9 @@ func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg [ } p.Log.Trace("<< PING/"+p.protocolName, "protocol", p.protocolName, "source", id, "pingRequest", pingRequest) + if metrics.Enabled { + messagesReceivedPing.Mark(1) + } resp, err := p.handlePing(id, pingRequest) if err != nil { p.Log.Error("failed to handle ping request", "err", err) @@ -885,6 +971,9 @@ func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg [ } p.Log.Trace("<< FIND_NODES/"+p.protocolName, "protocol", p.protocolName, "source", id, "findNodesRequest", findNodesRequest) + if metrics.Enabled { + messagesReceivedFindNodes.Mark(1) + } resp, err := p.handleFindNodes(addr, findNodesRequest) if err != nil { p.Log.Error("failed to handle find nodes request", "err", err) @@ -900,7 +989,10 @@ func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg [ return nil } - p.Log.Trace("<< FIND_NODES/"+p.protocolName, "protocol", p.protocolName, "source", id, "findContentRequest", findContentRequest) + p.Log.Trace("<< FIND_CONTENT/"+p.protocolName, "protocol", p.protocolName, "source", id, "findContentRequest", findContentRequest) + if metrics.Enabled { + messagesReceivedFindContent.Mark(1) + } resp, err := p.handleFindContent(id, addr, findContentRequest) if err != nil { p.Log.Error("failed to handle find content request", "err", err) @@ -917,6 +1009,9 @@ func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg [ } p.Log.Trace("<< OFFER/"+p.protocolName, "protocol", p.protocolName, "source", id, "offerRequest", offerRequest) + if metrics.Enabled { + messagesReceivedOffer.Mark(1) + } resp, err := p.handleOffer(id, addr, offerRequest) if err != nil { p.Log.Error("failed to handle offer request", "err", err) @@ -958,6 +1053,9 @@ func (p *PortalProtocol) handlePing(id enode.ID, ping *portalwire.Ping) ([]byte, } p.Log.Trace(">> PONG/"+p.protocolName, "protocol", p.protocolName, "source", id, "pong", pong) + if metrics.Enabled { + messagesSentPong.Mark(1) + } pongBytes, err := pong.MarshalSSZ() if err != nil { @@ -991,6 +1089,9 @@ func (p *PortalProtocol) handleFindNodes(fromAddr *net.UDPAddr, request *portalw } p.Log.Trace(">> NODES/"+p.protocolName, "protocol", p.protocolName, "source", fromAddr, "nodes", nodesMsg) + if metrics.Enabled { + messagesSentNodes.Mark(1) + } nodesMsgBytes, err := nodesMsg.MarshalSSZ() if err != nil { return nil, err @@ -1042,6 +1143,9 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque } p.Log.Trace(">> CONTENT_ENRS/"+p.protocolName, "protocol", p.protocolName, "source", addr, "enrs", enrsMsg) + if metrics.Enabled { + messagesSentContent.Mark(1) + } var enrsMsgBytes []byte enrsMsgBytes, err = enrsMsg.MarshalSSZ() if err != nil { @@ -1063,6 +1167,9 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque } p.Log.Trace(">> CONTENT_RAW/"+p.protocolName, "protocol", p.protocolName, "source", addr, "content", rawContentMsg) + if metrics.Enabled { + messagesSentContent.Mark(1) + } var rawContentMsgBytes []byte rawContentMsgBytes, err = rawContentMsg.MarshalSSZ() @@ -1136,6 +1243,9 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque } p.Log.Trace(">> CONTENT_CONNECTION_ID/"+p.protocolName, "protocol", p.protocolName, "source", addr, "connId", connIdMsg) + if metrics.Enabled { + messagesSentContent.Mark(1) + } var connIdMsgBytes []byte connIdMsgBytes, err = connIdMsg.MarshalSSZ() if err != nil { @@ -1164,6 +1274,9 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po } p.Log.Trace(">> ACCEPT/"+p.protocolName, "protocol", p.protocolName, "source", addr, "accept", acceptMsg) + if metrics.Enabled { + messagesSentAccept.Mark(1) + } var acceptMsgBytes []byte acceptMsgBytes, err = acceptMsg.MarshalSSZ() if err != nil { @@ -1239,6 +1352,9 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po return } p.Log.Trace("<< OFFER_CONTENT/"+p.protocolName, "id", id, "size", len(data), "data", data) + if metrics.Enabled { + messagesReceivedContent.Mark(1) + } err = p.handleOfferedContents(id, contentKeys, data) if err != nil { @@ -1262,6 +1378,9 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po } p.Log.Trace(">> ACCEPT/"+p.protocolName, "protocol", p.protocolName, "source", addr, "accept", acceptMsg) + if metrics.Enabled { + messagesSentAccept.Mark(1) + } var acceptMsgBytes []byte acceptMsgBytes, err = acceptMsg.MarshalSSZ() if err != nil {