Merge pull request #203 from r4f4ss/metrics_ethstats

[WIP] Metrics implementation
This commit is contained in:
rafaelss 2024-10-21 21:05:23 -04:00 committed by GitHub
commit 95c2637065
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 235 additions and 1 deletions

View file

@ -42,6 +42,10 @@ import (
"github.com/urfave/cli/v2" "github.com/urfave/cli/v2"
) )
var (
storageCapacity metrics.Gauge
)
const ( const (
privateKeyFileName = "clientKey" privateKeyFileName = "clientKey"
) )
@ -124,6 +128,11 @@ func shisui(ctx *cli.Context) error {
// Start system runtime metrics collection // Start system runtime metrics collection
go metrics.CollectProcessMetrics(3 * time.Second) 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) config, err := getPortalConfig(ctx)
if err != nil { if err != nil {
return nil return nil

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"
@ -199,6 +200,8 @@ type PortalProtocol struct {
portMappingRegister chan *portMapping portMappingRegister chan *portMapping
clock mclock.Clock clock mclock.Clock
NAT nat.Interface NAT nat.Interface
portalMetrics *portalMetrics
} }
func defaultContentIdFunc(contentKey []byte) []byte { func defaultContentIdFunc(contentKey []byte) []byte {
@ -237,6 +240,10 @@ func NewPortalProtocol(config *PortalProtocolConfig, protocolId portalwire.Proto
opt(protocol) opt(protocol)
} }
if metrics.Enabled {
protocol.portalMetrics = newPortalMetrics(protocolId.Name())
}
return protocol, nil return protocol, nil
} }
@ -428,6 +435,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 {
p.portalMetrics.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 +454,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 {
p.portalMetrics.messagesReceivedPong.Mark(1)
}
return p.processPong(node, talkResp) return p.processPong(node, talkResp)
} }
@ -463,6 +476,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 {
p.portalMetrics.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 +504,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 {
p.portalMetrics.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 +534,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 {
p.portalMetrics.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 +575,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 {
p.portalMetrics.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())
@ -631,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))) conn, err = utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId)))
conncancel() conncancel()
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpOutFailConn.Inc(1)
}
p.Log.Error("failed to dial utp connection", "err", err) p.Log.Error("failed to dial utp connection", "err", err)
return return
} }
err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout))
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpOutFailDeadline.Inc(1)
}
p.Log.Error("failed to set write deadline", "err", err) p.Log.Error("failed to set write deadline", "err", err)
return return
} }
@ -644,10 +675,17 @@ func (p *PortalProtocol) processOffer(target *enode.Node, resp []byte, request *
var written int var written int
written, err = conn.Write(contentsPayload) written, err = conn.Write(contentsPayload)
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpOutFailWrite.Inc(1)
}
p.Log.Error("failed to write to utp connection", "err", err) p.Log.Error("failed to write to utp connection", "err", err)
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 {
p.portalMetrics.messagesSentContent.Mark(1)
p.portalMetrics.utpOutSuccess.Inc(1)
}
return return
} }
} }
@ -677,6 +715,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 {
p.portalMetrics.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 +733,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 {
p.portalMetrics.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())
@ -706,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))) conn, err := utp.DialUTPOptions("utp", laddr, raddr, utp.WithContext(connctx), utp.WithSocketManager(p.utpSm), utp.WithConnId(uint32(connId)))
defer func() { defer func() {
if conn == nil { if conn == nil {
if metrics.Enabled {
p.portalMetrics.utpInFailConn.Inc(1)
}
return return
} }
err := conn.Close() err := conn.Close()
@ -720,15 +767,25 @@ func (p *PortalProtocol) processContent(target *enode.Node, resp []byte) (byte,
err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout)) err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout))
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpInFailDeadline.Inc(1)
}
return 0xff, nil, err return 0xff, nil, err
} }
// Read ALL the data from the connection until EOF and return it // Read ALL the data from the connection until EOF and return it
data, err := io.ReadAll(conn) data, err := io.ReadAll(conn)
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpInFailRead.Inc(1)
}
p.Log.Error("failed to read from utp connection", "err", err) p.Log.Error("failed to read from utp connection", "err", err)
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 {
p.portalMetrics.messagesReceivedContent.Mark(1)
p.portalMetrics.utpInSuccess.Inc(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 +796,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 {
p.portalMetrics.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 +864,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 {
p.portalMetrics.messagesReceivedNodes.Mark(1)
}
return nodes return nodes
} }
@ -821,6 +884,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 {
p.portalMetrics.messagesReceivedPong.Mark(1)
}
customPayload := &portalwire.PingPongCustomData{} customPayload := &portalwire.PingPongCustomData{}
err = customPayload.UnmarshalSSZ(pong.CustomPayload) err = customPayload.UnmarshalSSZ(pong.CustomPayload)
@ -829,6 +895,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 {
p.portalMetrics.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 +938,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 {
p.portalMetrics.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 +957,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 {
p.portalMetrics.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 +975,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 {
p.portalMetrics.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 +995,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 {
p.portalMetrics.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 +1039,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 {
p.portalMetrics.messagesSentPong.Mark(1)
}
pongBytes, err := pong.MarshalSSZ() pongBytes, err := pong.MarshalSSZ()
if err != nil { if err != nil {
@ -991,6 +1075,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 {
p.portalMetrics.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 +1129,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 {
p.portalMetrics.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 +1153,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 {
p.portalMetrics.messagesSentContent.Mark(1)
}
var rawContentMsgBytes []byte var rawContentMsgBytes []byte
rawContentMsgBytes, err = rawContentMsg.MarshalSSZ() rawContentMsgBytes, err = rawContentMsg.MarshalSSZ()
@ -1106,12 +1199,18 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend) conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
cancel() cancel()
if err != nil { 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) p.Log.Error("failed to accept utp connection for handle find content", "connId", connIdSend, "err", err)
return return
} }
err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout)) err = conn.SetWriteDeadline(time.Now().Add(defaultUTPWriteTimeout))
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpOutFailDeadline.Inc(1)
}
p.Log.Error("failed to set write deadline", "err", err) p.Log.Error("failed to set write deadline", "err", err)
return return
} }
@ -1119,10 +1218,16 @@ func (p *PortalProtocol) handleFindContent(id enode.ID, addr *net.UDPAddr, reque
var n int var n int
n, err = conn.Write(content) n, err = conn.Write(content)
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpOutFailWrite.Inc(1)
}
p.Log.Error("failed to write content to utp connection", "err", err) p.Log.Error("failed to write content to utp connection", "err", err)
return return
} }
if metrics.Enabled {
p.portalMetrics.utpOutSuccess.Inc(1)
}
p.Log.Trace("wrote content size to utp connection", "n", n) p.Log.Trace("wrote content size to utp connection", "n", n)
return return
} }
@ -1136,6 +1241,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 {
p.portalMetrics.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 +1272,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 {
p.portalMetrics.messagesSentAccept.Mark(1)
}
var acceptMsgBytes []byte var acceptMsgBytes []byte
acceptMsgBytes, err = acceptMsg.MarshalSSZ() acceptMsgBytes, err = acceptMsg.MarshalSSZ()
if err != nil { if err != nil {
@ -1222,12 +1333,18 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
conn, err = p.utp.AcceptUTPContext(ctx, connIdSend) conn, err = p.utp.AcceptUTPContext(ctx, connIdSend)
cancel() cancel()
if err != nil { 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) p.Log.Error("failed to accept utp connection for handle offer", "connId", connIdSend, "err", err)
return return
} }
err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout)) err = conn.SetReadDeadline(time.Now().Add(defaultUTPReadTimeout))
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpInFailDeadline.Inc(1)
}
p.Log.Error("failed to set read deadline", "err", err) p.Log.Error("failed to set read deadline", "err", err)
return return
} }
@ -1235,10 +1352,16 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
var data []byte var data []byte
data, err = io.ReadAll(conn) data, err = io.ReadAll(conn)
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.utpInFailRead.Inc(1)
}
p.Log.Error("failed to read from utp connection", "err", err) p.Log.Error("failed to read from utp connection", "err", err)
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 {
p.portalMetrics.messagesReceivedContent.Mark(1)
}
err = p.handleOfferedContents(id, contentKeys, data) err = p.handleOfferedContents(id, contentKeys, data)
if err != nil { if err != nil {
@ -1246,6 +1369,9 @@ func (p *PortalProtocol) handleOffer(id enode.ID, addr *net.UDPAddr, request *po
return return
} }
if metrics.Enabled {
p.portalMetrics.utpInSuccess.Inc(1)
}
return return
} }
} }
@ -1262,6 +1388,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 {
p.portalMetrics.messagesSentAccept.Mark(1)
}
var acceptMsgBytes []byte var acceptMsgBytes []byte
acceptMsgBytes, err = acceptMsg.MarshalSSZ() acceptMsgBytes, err = acceptMsg.MarshalSSZ()
if err != nil { if err != nil {
@ -1278,12 +1407,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 { func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, payload []byte) error {
contents, err := decodeContents(payload) contents, err := decodeContents(payload)
if err != nil { if err != nil {
if metrics.Enabled {
p.portalMetrics.contentDecodedFalse.Inc(1)
}
return err return err
} }
keyLen := len(keys) keyLen := len(keys)
contentLen := len(contents) contentLen := len(contents)
if keyLen != contentLen { if keyLen != contentLen {
if metrics.Enabled {
p.portalMetrics.contentDecodedFalse.Inc(1)
}
return fmt.Errorf("content keys len %d doesn't match content values len %d", keyLen, contentLen) return fmt.Errorf("content keys len %d doesn't match content values len %d", keyLen, contentLen)
} }
@ -1295,6 +1430,9 @@ func (p *PortalProtocol) handleOfferedContents(id enode.ID, keys [][]byte, paylo
p.contentQueue <- contentElement p.contentQueue <- contentElement
if metrics.Enabled {
p.portalMetrics.contentDecodedTrue.Inc(1)
}
return nil return nil
} }

View file

@ -0,0 +1,67 @@
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
utpInFailConn metrics.Counter
utpInFailRead metrics.Counter
utpInFailDeadline metrics.Counter
utpInSuccess metrics.Counter
utpOutFailConn metrics.Counter
utpOutFailWrite metrics.Counter
utpOutFailDeadline metrics.Counter
utpOutSuccess metrics.Counter
contentDecodedTrue metrics.Counter
contentDecodedFalse metrics.Counter
}
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),
utpInFailConn: metrics.NewRegisteredCounter("portal/"+protocolName+"/utp/inbound/fail_conn", 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),
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),
contentDecodedTrue: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/decoded/true", nil),
contentDecodedFalse: metrics.NewRegisteredCounter("portal/"+protocolName+"/content/decoded/false", nil),
}
}

View file

@ -14,12 +14,17 @@ import (
"github.com/ethereum/go-ethereum/common/hexutil" "github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/metrics"
"github.com/ethereum/go-ethereum/p2p/enode" "github.com/ethereum/go-ethereum/p2p/enode"
"github.com/ethereum/go-ethereum/portalnetwork/storage" "github.com/ethereum/go-ethereum/portalnetwork/storage"
"github.com/holiman/uint256" "github.com/holiman/uint256"
"github.com/mattn/go-sqlite3" "github.com/mattn/go-sqlite3"
) )
var (
radiusRatio metrics.GaugeFloat64
)
const ( const (
sqliteName = "history.sqlite" sqliteName = "history.sqlite"
contentDeletionFraction = 0.05 // 5% of the content will be deleted when the storage capacity is hit and radius gets adjusted. 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), log: log.New("storage", config.NetworkName),
} }
hs.radius.Store(storage.MaxDistance) hs.radius.Store(storage.MaxDistance)
if metrics.Enabled {
radiusRatio = metrics.NewRegisteredGaugeFloat64("portal/radius_ratio", nil)
radiusRatio.Update(1)
}
err := hs.createTable() err := hs.createTable()
if err != nil { if err != nil {
return nil, err return nil, err
@ -332,6 +341,12 @@ func (p *ContentStorage) EstimateNewRadius(currentRadius *uint256.Int) (*uint256
sizeRatio := currrentSize / p.storageCapacityInBytes sizeRatio := currrentSize / p.storageCapacityInBytes
if sizeRatio > 0 { if sizeRatio > 0 {
bigFormat := new(big.Int).SetUint64(sizeRatio) 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 new(uint256.Int).Div(currentRadius, uint256.MustFromBig(bigFormat)), nil
} }
return currentRadius, nil return currentRadius, nil
@ -391,6 +406,11 @@ func (p *ContentStorage) deleteContentFraction(fraction float64) (deleteCount in
return 0, err return 0, err
} }
p.radius.Store(dis) 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 // row must close first, or database is locked
// rows.Close() can call multi times // rows.Close() can call multi times