portal_protocol.go metrics

This commit is contained in:
Rafael Sampaio 2024-10-17 10:45:44 -03:00
parent c5236fc9a9
commit 9b9d666dc4

View file

@ -23,6 +23,7 @@ import (
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil" "github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/common/mclock" "github.com/ethereum/go-ethereum/common/mclock"
"github.com/ethereum/go-ethereum/metrics"
"github.com/ethereum/go-ethereum/p2p/discover/v5wire" "github.com/ethereum/go-ethereum/p2p/discover/v5wire"
"github.com/VictoriaMetrics/fastcache" "github.com/VictoriaMetrics/fastcache"
@ -43,6 +44,26 @@ import (
"go.uber.org/zap" "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 ( const (
// TalkResp message is a response message so the session is established and a // 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) 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 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) 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() pingRequestBytes, err := pingRequest.MarshalSSZ()
if err != nil { if err != nil {
return nil, err 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) 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) 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) p.Log.Trace(">> FIND_NODES/"+p.protocolName, "id", node.ID(), "findNodes", findNodes)
if metrics.Enabled {
messagesSentFindNodes.Mark(1)
}
findNodesBytes, err := findNodes.MarshalSSZ() findNodesBytes, err := findNodes.MarshalSSZ()
if err != nil { if err != nil {
p.Log.Error("failed to marshal find nodes request", "err", err) 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) p.Log.Trace(">> FIND_CONTENT/"+p.protocolName, "id", node.ID(), "findContent", findContent)
if metrics.Enabled {
messagesSentFindContent.Mark(1)
}
findContentBytes, err := findContent.MarshalSSZ() findContentBytes, err := findContent.MarshalSSZ()
if err != nil { if err != nil {
p.Log.Error("failed to marshal find content request", "err", err) 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) p.Log.Trace(">> OFFER/"+p.protocolName, "offer", offer)
if metrics.Enabled {
messagesSentOffer.Mark(1)
}
offerBytes, err := offer.MarshalSSZ() offerBytes, err := offer.MarshalSSZ()
if err != nil { if err != nil {
p.Log.Error("failed to marshal offer request", "err", err) 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) p.Log.Trace("<< ACCEPT/"+p.protocolName, "id", target.ID(), "accept", accept)
if metrics.Enabled {
messagesReceivedAccept.Mark(1)
}
isAdded := p.table.addFoundNode(target, true) isAdded := p.table.addFoundNode(target, true)
if isAdded { if isAdded {
log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) 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 return
} }
p.Log.Trace(">> CONTENT/"+p.protocolName, "id", target.ID(), "contents", contents, "size", written) p.Log.Trace(">> CONTENT/"+p.protocolName, "id", target.ID(), "contents", contents, "size", written)
if metrics.Enabled {
messagesSentContent.Mark(1)
}
return 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) p.Log.Trace("<< CONTENT/"+p.protocolName, "id", target.ID(), "content", content)
if metrics.Enabled {
messagesReceivedContent.Mark(1)
}
isAdded := p.table.addFoundNode(target, true) isAdded := p.table.addFoundNode(target, true)
if isAdded { if isAdded {
log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) 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) 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) isAdded := p.table.addFoundNode(target, true)
if isAdded { if isAdded {
log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) 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 return 0xff, nil, err
} }
p.Log.Trace("<< CONTENT/"+p.protocolName, "id", target.ID(), "size", len(data), "data", data) 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 return resp[1], data, nil
case portalwire.ContentEnrsSelector: case portalwire.ContentEnrsSelector:
enrs := &portalwire.Enrs{} 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) p.Log.Trace("<< CONTENT_ENRS/"+p.protocolName, "id", target.ID(), "enrs", enrs)
if metrics.Enabled {
messagesReceivedContent.Mark(1)
}
isAdded := p.table.addFoundNode(target, true) isAdded := p.table.addFoundNode(target, true)
if isAdded { if isAdded {
log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) 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) p.Log.Trace("<< NODES/"+p.protocolName, "id", target.ID(), "total", len(enrs), "verified", verified, "nodes", nodes)
if metrics.Enabled {
messagesReceivedNodes.Mark(1)
}
return nodes 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) p.Log.Trace("<< PONG_RESPONSE/"+p.protocolName, "id", target.ID(), "pong", pong)
if metrics.Enabled {
messagesReceivedPong.Mark(1)
}
customPayload := &portalwire.PingPongCustomData{} customPayload := &portalwire.PingPongCustomData{}
err = customPayload.UnmarshalSSZ(pong.CustomPayload) 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) 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) isAdded := p.table.addFoundNode(target, true)
if isAdded { if isAdded {
log.Debug("Node added to bucket", "protocol", p.protocolName, "node", target.IP(), "port", target.UDP()) 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) 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) resp, err := p.handlePing(id, pingRequest)
if err != nil { if err != nil {
p.Log.Error("failed to handle ping request", "err", err) 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) 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) resp, err := p.handleFindNodes(addr, findNodesRequest)
if err != nil { if err != nil {
p.Log.Error("failed to handle find nodes request", "err", err) 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 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) resp, err := p.handleFindContent(id, addr, findContentRequest)
if err != nil { if err != nil {
p.Log.Error("failed to handle find content request", "err", err) 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) 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) resp, err := p.handleOffer(id, addr, offerRequest)
if err != nil { if err != nil {
p.Log.Error("failed to handle offer request", "err", err) 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) p.Log.Trace(">> PONG/"+p.protocolName, "protocol", p.protocolName, "source", id, "pong", pong)
if metrics.Enabled {
messagesSentPong.Mark(1)
}
pongBytes, err := pong.MarshalSSZ() pongBytes, err := pong.MarshalSSZ()
if err != nil { 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) p.Log.Trace(">> NODES/"+p.protocolName, "protocol", p.protocolName, "source", fromAddr, "nodes", nodesMsg)
if metrics.Enabled {
messagesSentNodes.Mark(1)
}
nodesMsgBytes, err := nodesMsg.MarshalSSZ() nodesMsgBytes, err := nodesMsg.MarshalSSZ()
if err != nil { if err != nil {
return nil, err 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) p.Log.Trace(">> CONTENT_ENRS/"+p.protocolName, "protocol", p.protocolName, "source", addr, "enrs", enrsMsg)
if metrics.Enabled {
messagesSentContent.Mark(1)
}
var enrsMsgBytes []byte var enrsMsgBytes []byte
enrsMsgBytes, err = enrsMsg.MarshalSSZ() enrsMsgBytes, err = enrsMsg.MarshalSSZ()
if err != nil { 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) p.Log.Trace(">> CONTENT_RAW/"+p.protocolName, "protocol", p.protocolName, "source", addr, "content", rawContentMsg)
if metrics.Enabled {
messagesSentContent.Mark(1)
}
var rawContentMsgBytes []byte var rawContentMsgBytes []byte
rawContentMsgBytes, err = rawContentMsg.MarshalSSZ() 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) p.Log.Trace(">> CONTENT_CONNECTION_ID/"+p.protocolName, "protocol", p.protocolName, "source", addr, "connId", connIdMsg)
if metrics.Enabled {
messagesSentContent.Mark(1)
}
var connIdMsgBytes []byte var connIdMsgBytes []byte
connIdMsgBytes, err = connIdMsg.MarshalSSZ() connIdMsgBytes, err = connIdMsg.MarshalSSZ()
if err != nil { 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) p.Log.Trace(">> ACCEPT/"+p.protocolName, "protocol", p.protocolName, "source", addr, "accept", acceptMsg)
if metrics.Enabled {
messagesSentAccept.Mark(1)
}
var acceptMsgBytes []byte var acceptMsgBytes []byte
acceptMsgBytes, err = acceptMsg.MarshalSSZ() acceptMsgBytes, err = acceptMsg.MarshalSSZ()
if err != nil { if err != nil {
@ -1239,6 +1352,9 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
return return
} }
p.Log.Trace("<< OFFER_CONTENT/"+p.protocolName, "id", id, "size", len(data), "data", data) 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) err = p.handleOfferedContents(id, contentKeys, data)
if err != nil { 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) p.Log.Trace(">> ACCEPT/"+p.protocolName, "protocol", p.protocolName, "source", addr, "accept", acceptMsg)
if metrics.Enabled {
messagesSentAccept.Mark(1)
}
var acceptMsgBytes []byte var acceptMsgBytes []byte
acceptMsgBytes, err = acceptMsg.MarshalSSZ() acceptMsgBytes, err = acceptMsg.MarshalSSZ()
if err != nil { if err != nil {