From 9b9d666dc43593b194685d33e9d46433e2e864fe Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Thu, 17 Oct 2024 10:45:44 -0300 Subject: [PATCH 1/7] portal_protocol.go metrics --- p2p/discover/portal_protocol.go | 121 +++++++++++++++++++++++++++++++- 1 file changed, 120 insertions(+), 1 deletion(-) 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 { From c5d13e4b74a2b5d7e3b02547b51418880eef02d2 Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Thu, 17 Oct 2024 12:02:47 -0300 Subject: [PATCH 2/7] metric radius rate --- portalnetwork/history/storage.go | 20 ++++++++++++++++++++ 1 file changed, 20 insertions(+) diff --git a/portalnetwork/history/storage.go b/portalnetwork/history/storage.go index 9844f96f89..7878093f6b 100644 --- a/portalnetwork/history/storage.go +++ b/portalnetwork/history/storage.go @@ -14,12 +14,17 @@ import ( "github.com/ethereum/go-ethereum/common/hexutil" "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/metrics" "github.com/ethereum/go-ethereum/p2p/enode" "github.com/ethereum/go-ethereum/portalnetwork/storage" "github.com/holiman/uint256" "github.com/mattn/go-sqlite3" ) +var ( + radiusRatio metrics.GaugeFloat64 +) + const ( sqliteName = "history.sqlite" contentDeletionFraction = 0.05 // 5% of the content will be deleted when the storage capacity is hit and radius gets adjusted. @@ -107,6 +112,10 @@ func NewHistoryStorage(config storage.PortalStorageConfig) (storage.ContentStora log: log.New("storage", config.NetworkName), } hs.radius.Store(storage.MaxDistance) + if metrics.Enabled { + radiusRatio = metrics.NewRegisteredGaugeFloat64("portal/radius_ratio", nil) + radiusRatio.Update(1) + } err := hs.createTable() if err != nil { return nil, err @@ -332,6 +341,12 @@ func (p *ContentStorage) EstimateNewRadius(currentRadius *uint256.Int) (*uint256 sizeRatio := currrentSize / p.storageCapacityInBytes if sizeRatio > 0 { bigFormat := new(big.Int).SetUint64(sizeRatio) + if metrics.Enabled { + newRadius := new(uint256.Int).Div(currentRadius, uint256.MustFromBig(bigFormat)) + newRadius.Mul(newRadius, uint256.NewInt(100)) + newRadius.Mod(newRadius, storage.MaxDistance) + radiusRatio.Update(newRadius.Float64() / 100) + } return new(uint256.Int).Div(currentRadius, uint256.MustFromBig(bigFormat)), nil } return currentRadius, nil @@ -391,6 +406,11 @@ func (p *ContentStorage) deleteContentFraction(fraction float64) (deleteCount in return 0, err } p.radius.Store(dis) + if metrics.Enabled { + dis.Mul(dis, uint256.NewInt(100)) + dis.Mod(dis, storage.MaxDistance) + radiusRatio.Update(dis.Float64() / 100) + } } // row must close first, or database is locked // rows.Close() can call multi times From d1b5606b0004b4b7342bcd26b948701b52c51da0 Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Thu, 17 Oct 2024 18:41:12 -0300 Subject: [PATCH 3/7] metric storage_capacity --- cmd/shisui/main.go | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/cmd/shisui/main.go b/cmd/shisui/main.go index 6260029552..6977180a4d 100644 --- a/cmd/shisui/main.go +++ b/cmd/shisui/main.go @@ -42,6 +42,10 @@ import ( "github.com/urfave/cli/v2" ) +var ( + storageCapacity metrics.Gauge +) + const ( privateKeyFileName = "clientKey" ) @@ -124,6 +128,11 @@ func shisui(ctx *cli.Context) error { // Start system runtime metrics collection go metrics.CollectProcessMetrics(3 * time.Second) + if metrics.Enabled { + storageCapacity = metrics.NewRegisteredGauge("portal/storage_capacity", nil) + storageCapacity.Update(ctx.Int64(utils.PortalDataCapacityFlag.Name)) + } + config, err := getPortalConfig(ctx) if err != nil { return nil From efa2e9002a5f7d3a804b44adb22cc11f15fee849 Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Thu, 17 Oct 2024 21:46:22 -0300 Subject: [PATCH 4/7] refactor: portal protocol metrics --- p2p/discover/portal_protocol.go | 92 ++++++++----------------- p2p/discover/portal_protocol_metrics.go | 44 ++++++++++++ 2 files changed, 73 insertions(+), 63 deletions(-) create mode 100644 p2p/discover/portal_protocol_metrics.go diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index a298c3ae84..b5f116880a 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -44,26 +44,6 @@ 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 @@ -220,6 +200,8 @@ type PortalProtocol struct { portMappingRegister chan *portMapping clock mclock.Clock NAT nat.Interface + + portalMetrics *portalMetrics } func defaultContentIdFunc(contentKey []byte) []byte { @@ -259,23 +241,7 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId portalwire.Proto } 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) + protocol.portalMetrics = newPortalMetrics(protocolId.Name()) } return protocol, nil @@ -470,7 +436,7 @@ 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) + p.portalMetrics.messagesSentPing.Mark(1) } pingRequestBytes, err := pingRequest.MarshalSSZ() if err != nil { @@ -489,7 +455,7 @@ 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) + p.portalMetrics.messagesReceivedPong.Mark(1) } return p.processPong(node, talkResp) @@ -511,7 +477,7 @@ 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) + p.portalMetrics.messagesSentFindNodes.Mark(1) } findNodesBytes, err := findNodes.MarshalSSZ() if err != nil { @@ -539,7 +505,7 @@ 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) + p.portalMetrics.messagesSentFindContent.Mark(1) } findContentBytes, err := findContent.MarshalSSZ() if err != nil { @@ -569,7 +535,7 @@ func (p *PortalProtocol) offer(node *enode.Node, offerRequest *OfferRequest) ([] p.Log.Trace(">> OFFER/"+p.protocolName, "offer", offer) if metrics.Enabled { - messagesSentOffer.Mark(1) + p.portalMetrics.messagesSentOffer.Mark(1) } offerBytes, err := offer.MarshalSSZ() if err != nil { @@ -610,7 +576,7 @@ 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) + p.portalMetrics.messagesReceivedAccept.Mark(1) } isAdded := p.table.addFoundNode(target, true) if isAdded { @@ -708,7 +674,7 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * } p.Log.Trace(">> CONTENT/"+p.protocolName, "id", target.ID(), "contents", contents, "size", written) if metrics.Enabled { - messagesSentContent.Mark(1) + p.portalMetrics.messagesSentContent.Mark(1) } return } @@ -740,7 +706,7 @@ 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) + p.portalMetrics.messagesReceivedContent.Mark(1) } isAdded := p.table.addFoundNode(target, true) if isAdded { @@ -758,7 +724,7 @@ 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) + p.portalMetrics.messagesReceivedContent.Mark(1) } isAdded := p.table.addFoundNode(target, true) if isAdded { @@ -798,7 +764,7 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, } p.Log.Trace("<< CONTENT/"+p.protocolName, "id", target.ID(), "size", len(data), "data", data) if metrics.Enabled { - messagesReceivedContent.Mark(1) + p.portalMetrics.messagesReceivedContent.Mark(1) } return resp[1], data, nil case portalwire.ContentEnrsSelector: @@ -811,7 +777,7 @@ 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) + p.portalMetrics.messagesReceivedContent.Mark(1) } isAdded := p.table.addFoundNode(target, true) if isAdded { @@ -879,7 +845,7 @@ 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) + p.portalMetrics.messagesReceivedNodes.Mark(1) } return nodes } @@ -899,7 +865,7 @@ 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) + p.portalMetrics.messagesReceivedPong.Mark(1) } customPayload := &portalwire.PingPongCustomData{} @@ -910,7 +876,7 @@ 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) + p.portalMetrics.messagesReceivedPong.Mark(1) } isAdded := p.table.addFoundNode(target, true) if isAdded { @@ -953,7 +919,7 @@ 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) + p.portalMetrics.messagesReceivedPing.Mark(1) } resp, err := p.handlePing(id, pingRequest) if err != nil { @@ -972,7 +938,7 @@ 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) + p.portalMetrics.messagesReceivedFindNodes.Mark(1) } resp, err := p.handleFindNodes(addr, findNodesRequest) if err != nil { @@ -991,7 +957,7 @@ func (p *PortalProtocol) handleTalkRequest(id enode.ID, addr *net.UDPAddr, msg [ p.Log.Trace("<< FIND_CONTENT/"+p.protocolName, "protocol", p.protocolName, "source", id, "findContentRequest", findContentRequest) if metrics.Enabled { - messagesReceivedFindContent.Mark(1) + p.portalMetrics.messagesReceivedFindContent.Mark(1) } resp, err := p.handleFindContent(id, addr, findContentRequest) if err != nil { @@ -1010,7 +976,7 @@ 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) + p.portalMetrics.messagesReceivedOffer.Mark(1) } resp, err := p.handleOffer(id, addr, offerRequest) if err != nil { @@ -1054,7 +1020,7 @@ 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) + p.portalMetrics.messagesSentPong.Mark(1) } pongBytes, err := pong.MarshalSSZ() @@ -1090,7 +1056,7 @@ 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) + p.portalMetrics.messagesSentNodes.Mark(1) } nodesMsgBytes, err := nodesMsg.MarshalSSZ() if err != nil { @@ -1144,7 +1110,7 @@ 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) + p.portalMetrics.messagesSentContent.Mark(1) } var enrsMsgBytes []byte enrsMsgBytes, err = enrsMsg.MarshalSSZ() @@ -1168,7 +1134,7 @@ 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) + p.portalMetrics.messagesSentContent.Mark(1) } var rawContentMsgBytes []byte @@ -1244,7 +1210,7 @@ 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) + p.portalMetrics.messagesSentContent.Mark(1) } var connIdMsgBytes []byte connIdMsgBytes, err = connIdMsg.MarshalSSZ() @@ -1275,7 +1241,7 @@ 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) + p.portalMetrics.messagesSentAccept.Mark(1) } var acceptMsgBytes []byte acceptMsgBytes, err = acceptMsg.MarshalSSZ() @@ -1353,7 +1319,7 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po } p.Log.Trace("<< OFFER_CONTENT/"+p.protocolName, "id", id, "size", len(data), "data", data) if metrics.Enabled { - messagesReceivedContent.Mark(1) + p.portalMetrics.messagesReceivedContent.Mark(1) } err = p.handleOfferedContents(id, contentKeys, data) @@ -1379,7 +1345,7 @@ 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) + p.portalMetrics.messagesSentAccept.Mark(1) } var acceptMsgBytes []byte acceptMsgBytes, err = acceptMsg.MarshalSSZ() diff --git a/p2p/discover/portal_protocol_metrics.go b/p2p/discover/portal_protocol_metrics.go new file mode 100644 index 0000000000..6dc7ca2358 --- /dev/null +++ b/p2p/discover/portal_protocol_metrics.go @@ -0,0 +1,44 @@ +package discover + +import "github.com/ethereum/go-ethereum/metrics" + +type portalMetrics struct { + 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 +} + +func newPortalMetrics(protocolName string) *portalMetrics { + return &portalMetrics{ + messagesReceivedAccept: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/accept", nil), + messagesReceivedNodes: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/nodes", nil), + messagesReceivedFindNodes: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/find_nodes", nil), + messagesReceivedFindContent: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/find_content", nil), + messagesReceivedContent: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/content", nil), + messagesReceivedOffer: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/offer", nil), + messagesReceivedPing: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/ping", nil), + messagesReceivedPong: metrics.NewRegisteredMeter("portal/"+protocolName+"/received/pong", nil), + messagesSentAccept: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/accept", nil), + messagesSentNodes: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/nodes", nil), + messagesSentFindNodes: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/find_nodes", nil), + messagesSentFindContent: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/find_content", nil), + messagesSentContent: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/content", nil), + messagesSentOffer: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/offer", nil), + messagesSentPing: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/ping", nil), + messagesSentPong: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/pong", nil), + } +} From e85b17cf37c5e7bcad6e404d816fe5ebf9f50671 Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Thu, 17 Oct 2024 23:21:04 -0300 Subject: [PATCH 5/7] utp protocol metrics --- p2p/discover/portal_protocol.go | 47 +++++++++++++++++++++++++ p2p/discover/portal_protocol_metrics.go | 18 ++++++++++ 2 files changed, 65 insertions(+) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index b5f116880a..0b6136611d 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -656,12 +656,18 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * conn, err = utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId))) conncancel() if err != nil { + if metrics.Enabled { + p.portalMetrics.utpOutFailConn.Inc(1) + } p.Log.Error("failed to dial utp connection", "err", err) return } err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpOutFailShutdown.Inc(1) + } p.Log.Error("failed to set write deadline", "err", err) return } @@ -669,12 +675,16 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * var written int written, err = conn.Write(contentsPayload) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpOutFailTx.Inc(1) + } p.Log.Error("failed to write to utp connection", "err", err) return } p.Log.Trace(">> CONTENT/"+p.protocolName, "id", target.ID(), "contents", contents, "size", written) if metrics.Enabled { p.portalMetrics.messagesSentContent.Mark(1) + p.portalMetrics.utpOutSuccess.Inc(1) } return } @@ -740,6 +750,9 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, conn, err := utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId))) defer func() { if conn == nil { + if metrics.Enabled { + p.portalMetrics.utpInFailConn.Inc(1) + } return } err := conn.Close() @@ -754,17 +767,24 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout)) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpInFailShutdown.Inc(1) + } return 0xff, nil, err } // Read ALL the data from the connection until EOF and return it data, err := io.ReadAll(conn) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpInFailTx.Inc(1) + } p.Log.Error("failed to read from utp connection", "err", err) return 0xff, nil, err } p.Log.Trace("<< CONTENT/"+p.protocolName, "id", target.ID(), "size", len(data), "data", data) if metrics.Enabled { p.portalMetrics.messagesReceivedContent.Mark(1) + p.portalMetrics.utpInSuccess.Inc(1) } return resp[1], data, nil case portalwire.ContentEnrsSelector: @@ -1179,12 +1199,18 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque conn, err = p.utp.AcceptUTPContext(ctx, connIdSend) cancel() if err != nil { + if metrics.Enabled { + p.portalMetrics.utpOutFailConn.Inc(1) + } p.Log.Error("failed to accept utp connection for handle find content", "connId", connIdSend, "err", err) return } err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpOutFailShutdown.Inc(1) + } p.Log.Error("failed to set write deadline", "err", err) return } @@ -1192,10 +1218,16 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque var n int n, err = conn.Write(content) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpOutFailTx.Inc(1) + } p.Log.Error("failed to write content to utp connection", "err", err) return } + if metrics.Enabled { + p.portalMetrics.utpOutSuccess.Inc(1) + } p.Log.Trace("wrote content size to utp connection", "n", n) return } @@ -1301,12 +1333,18 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po conn, err = p.utp.AcceptUTPContext(ctx, connIdSend) cancel() if err != nil { + if metrics.Enabled { + p.portalMetrics.utpInFailConn.Inc(1) + } p.Log.Error("failed to accept utp connection for handle offer", "connId", connIdSend, "err", err) return } err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout)) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpInFailShutdown.Inc(1) + } p.Log.Error("failed to set read deadline", "err", err) return } @@ -1314,6 +1352,9 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po var data []byte data, err = io.ReadAll(conn) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpInFailTx.Inc(1) + } p.Log.Error("failed to read from utp connection", "err", err) return } @@ -1324,10 +1365,16 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po err = p.handleOfferedContents(id, contentKeys, data) if err != nil { + if metrics.Enabled { + p.portalMetrics.utpInFailTx.Inc(1) + } p.Log.Error("failed to handle offered Contents", "err", err) return } + if metrics.Enabled { + p.portalMetrics.utpInSuccess.Inc(1) + } return } } diff --git a/p2p/discover/portal_protocol_metrics.go b/p2p/discover/portal_protocol_metrics.go index 6dc7ca2358..1fc2d13bf1 100644 --- a/p2p/discover/portal_protocol_metrics.go +++ b/p2p/discover/portal_protocol_metrics.go @@ -20,6 +20,16 @@ type portalMetrics struct { messagesSentOffer metrics.Meter messagesSentPing metrics.Meter messagesSentPong metrics.Meter + + utpInFailConn metrics.Counter + utpInFailTx metrics.Counter + utpInFailShutdown metrics.Counter + utpInSuccess metrics.Counter + + utpOutFailConn metrics.Counter + utpOutFailTx metrics.Counter + utpOutFailShutdown metrics.Counter + utpOutSuccess metrics.Counter } func newPortalMetrics(protocolName string) *portalMetrics { @@ -40,5 +50,13 @@ func newPortalMetrics(protocolName string) *portalMetrics { messagesSentOffer: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/offer", nil), messagesSentPing: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/ping", nil), messagesSentPong: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/pong", nil), + utpInFailConn: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_conn", nil), + utpInFailTx: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_tx", nil), + utpInFailShutdown: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_shutdown", nil), + utpInSuccess: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/success", nil), + utpOutFailConn: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_conn", nil), + utpOutFailTx: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_tx", nil), + utpOutFailShutdown: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_shutdown", nil), + utpOutSuccess: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/success", nil), } } From c9ddfbf290b80a63e2084bf301b4fae5af8bb1cb Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Fri, 18 Oct 2024 00:19:22 -0300 Subject: [PATCH 6/7] content validation metrics --- p2p/discover/portal_protocol.go | 9 +++++++++ p2p/discover/portal_protocol_metrics.go | 5 +++++ 2 files changed, 14 insertions(+) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index 0b6136611d..8cb95ec77e 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -1410,12 +1410,18 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, payload []byte) error { contents, err := decodeContents(payload) if err != nil { + if metrics.Enabled { + p.portalMetrics.contentInvalidated.Inc(1) + } return err } keyLen := len(keys) contentLen := len(contents) if keyLen != contentLen { + if metrics.Enabled { + p.portalMetrics.contentInvalidated.Inc(1) + } return fmt.Errorf("content keys len %d doesn't match content values len %d", keyLen, contentLen) } @@ -1427,6 +1433,9 @@ func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, paylo p.contentQueue <- contentElement + if metrics.Enabled { + p.portalMetrics.contentValidated.Inc(1) + } return nil } diff --git a/p2p/discover/portal_protocol_metrics.go b/p2p/discover/portal_protocol_metrics.go index 1fc2d13bf1..86c02663a2 100644 --- a/p2p/discover/portal_protocol_metrics.go +++ b/p2p/discover/portal_protocol_metrics.go @@ -30,6 +30,9 @@ type portalMetrics struct { utpOutFailTx metrics.Counter utpOutFailShutdown metrics.Counter utpOutSuccess metrics.Counter + + contentValidated metrics.Counter + contentInvalidated metrics.Counter } func newPortalMetrics(protocolName string) *portalMetrics { @@ -58,5 +61,7 @@ func newPortalMetrics(protocolName string) *portalMetrics { utpOutFailTx: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_tx", nil), utpOutFailShutdown: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_shutdown", nil), utpOutSuccess: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/success", nil), + contentValidated: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/validated", nil), + contentInvalidated: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/invalidated", nil), } } From 9d046eab4223c8ec4ebfe402b9a0cf3df9e53ea5 Mon Sep 17 00:00:00 2001 From: Rafael Sampaio <5679073+r4f4ss@users.noreply.github.com> Date: Fri, 18 Oct 2024 02:11:13 -0300 Subject: [PATCH 7/7] improves metric names --- p2p/discover/portal_protocol.go | 25 +++++++++++-------------- p2p/discover/portal_protocol_metrics.go | 24 ++++++++++++------------ 2 files changed, 23 insertions(+), 26 deletions(-) diff --git a/p2p/discover/portal_protocol.go b/p2p/discover/portal_protocol.go index 8cb95ec77e..cceb256d11 100644 --- a/p2p/discover/portal_protocol.go +++ b/p2p/discover/portal_protocol.go @@ -666,7 +666,7 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) if err != nil { if metrics.Enabled { - p.portalMetrics.utpOutFailShutdown.Inc(1) + p.portalMetrics.utpOutFailDeadline.Inc(1) } p.Log.Error("failed to set write deadline", "err", err) return @@ -676,7 +676,7 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request * written, err = conn.Write(contentsPayload) if err != nil { if metrics.Enabled { - p.portalMetrics.utpOutFailTx.Inc(1) + p.portalMetrics.utpOutFailWrite.Inc(1) } p.Log.Error("failed to write to utp connection", "err", err) return @@ -768,7 +768,7 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout)) if err != nil { if metrics.Enabled { - p.portalMetrics.utpInFailShutdown.Inc(1) + p.portalMetrics.utpInFailDeadline.Inc(1) } return 0xff, nil, err } @@ -776,7 +776,7 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte, data, err := io.ReadAll(conn) if err != nil { if metrics.Enabled { - p.portalMetrics.utpInFailTx.Inc(1) + p.portalMetrics.utpInFailRead.Inc(1) } p.Log.Error("failed to read from utp connection", "err", err) return 0xff, nil, err @@ -1209,7 +1209,7 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) if err != nil { if metrics.Enabled { - p.portalMetrics.utpOutFailShutdown.Inc(1) + p.portalMetrics.utpOutFailDeadline.Inc(1) } p.Log.Error("failed to set write deadline", "err", err) return @@ -1219,7 +1219,7 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque n, err = conn.Write(content) if err != nil { if metrics.Enabled { - p.portalMetrics.utpOutFailTx.Inc(1) + p.portalMetrics.utpOutFailWrite.Inc(1) } p.Log.Error("failed to write content to utp connection", "err", err) return @@ -1343,7 +1343,7 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout)) if err != nil { if metrics.Enabled { - p.portalMetrics.utpInFailShutdown.Inc(1) + p.portalMetrics.utpInFailDeadline.Inc(1) } p.Log.Error("failed to set read deadline", "err", err) return @@ -1353,7 +1353,7 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po data, err = io.ReadAll(conn) if err != nil { if metrics.Enabled { - p.portalMetrics.utpInFailTx.Inc(1) + p.portalMetrics.utpInFailRead.Inc(1) } p.Log.Error("failed to read from utp connection", "err", err) return @@ -1365,9 +1365,6 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po err = p.handleOfferedContents(id, contentKeys, data) if err != nil { - if metrics.Enabled { - p.portalMetrics.utpInFailTx.Inc(1) - } p.Log.Error("failed to handle offered Contents", "err", err) return } @@ -1411,7 +1408,7 @@ func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, paylo contents, err := decodeContents(payload) if err != nil { if metrics.Enabled { - p.portalMetrics.contentInvalidated.Inc(1) + p.portalMetrics.contentDecodedFalse.Inc(1) } return err } @@ -1420,7 +1417,7 @@ func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, paylo contentLen := len(contents) if keyLen != contentLen { if metrics.Enabled { - p.portalMetrics.contentInvalidated.Inc(1) + p.portalMetrics.contentDecodedFalse.Inc(1) } return fmt.Errorf("content keys len %d doesn't match content values len %d", keyLen, contentLen) } @@ -1434,7 +1431,7 @@ func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, paylo p.contentQueue <- contentElement if metrics.Enabled { - p.portalMetrics.contentValidated.Inc(1) + p.portalMetrics.contentDecodedTrue.Inc(1) } return nil } diff --git a/p2p/discover/portal_protocol_metrics.go b/p2p/discover/portal_protocol_metrics.go index 86c02663a2..0bff030f5a 100644 --- a/p2p/discover/portal_protocol_metrics.go +++ b/p2p/discover/portal_protocol_metrics.go @@ -22,17 +22,17 @@ type portalMetrics struct { messagesSentPong metrics.Meter utpInFailConn metrics.Counter - utpInFailTx metrics.Counter - utpInFailShutdown metrics.Counter + utpInFailRead metrics.Counter + utpInFailDeadline metrics.Counter utpInSuccess metrics.Counter utpOutFailConn metrics.Counter - utpOutFailTx metrics.Counter - utpOutFailShutdown metrics.Counter + utpOutFailWrite metrics.Counter + utpOutFailDeadline metrics.Counter utpOutSuccess metrics.Counter - contentValidated metrics.Counter - contentInvalidated metrics.Counter + contentDecodedTrue metrics.Counter + contentDecodedFalse metrics.Counter } func newPortalMetrics(protocolName string) *portalMetrics { @@ -54,14 +54,14 @@ func newPortalMetrics(protocolName string) *portalMetrics { messagesSentPing: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/ping", nil), messagesSentPong: metrics.NewRegisteredMeter("portal/"+protocolName+"/sent/pong", nil), utpInFailConn: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_conn", nil), - utpInFailTx: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_tx", nil), - utpInFailShutdown: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_shutdown", nil), + utpInFailRead: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_read", nil), + utpInFailDeadline: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_deadline", nil), utpInSuccess: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/success", nil), utpOutFailConn: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_conn", nil), - utpOutFailTx: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_tx", nil), - utpOutFailShutdown: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_shutdown", nil), + utpOutFailWrite: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_write", nil), + utpOutFailDeadline: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/fail_deadline", nil), utpOutSuccess: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/outbound/success", nil), - contentValidated: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/validated", nil), - contentInvalidated: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/invalidated", nil), + contentDecodedTrue: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/decoded/true", nil), + contentDecodedFalse: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/decoded/false", nil), } }