diff --git a/p2p/metrics.go b/p2p/metrics.go index 27109ab1a9..ffc280db70 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -37,10 +37,7 @@ const ( MetricsOutboundConnects = "p2p/OutboundConnects" // Name for the registered outbound connects meter MetricsOutboundTraffic = "p2p/OutboundTraffic" // Name for the registered outbound traffic meter - MetricsRegistryIngressPrefix = MetricsInboundTraffic + "/" - MetricsRegistryEgressPrefix = MetricsOutboundTraffic + "/" - - MeteredPeerLimit = 1024 + MeteredPeerLimit = 1024 // This amount of peers are individually metered ) var ( @@ -49,11 +46,11 @@ var ( egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) // Meter counting the egress connections egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) // Meter metering the cumulative egress traffic - PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) // Registry containing the peer ingress - PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) // Registry containing the peer egress + PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsInboundTraffic+"/") // Registry containing the peer ingress + PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsOutboundTraffic+"/") // Registry containing the peer egress - metricsFeed event.Feed // Event feed for peer metrics - meteredPeerCount uint64 // Actually stored peer connection count + meteredPeerFeed event.Feed // Event feed for peer metrics + meteredPeerCount int32 // Actually stored peer connection count ) // MeteredPeerEventType is the type of peer events emitted by a metered connection. @@ -72,19 +69,20 @@ const ( PeerHandshakeFailed ) -// MeteredPeerEvent is an event emitted when peers connect or disconnect +// MeteredPeerEvent is an event emitted when peers connect or disconnect. type MeteredPeerEvent struct { Type MeteredPeerEventType // Type of peer event IP net.IP // IP address of the peer ID string // NodeID of the peer Elapsed time.Duration // Time elapsed between the connection and the handshake/disconnection - Ingress uint64 // Ingress count in the moment of disconnection - Egress uint64 // Egress count in the moment of disconnection + Ingress uint64 // Ingress count at the moment of the event + Egress uint64 // Egress count at the moment of the event } -// SubscribeMeteredPeerEvent registers a subscription of MeteredPeerEvent +// SubscribeMeteredPeerEvent registers a subscription for peer life-cycle events +// if metrics collection is enabled. func SubscribeMeteredPeerEvent(ch chan<- MeteredPeerEvent) event.Subscription { - return metricsFeed.Subscribe(ch) + return meteredPeerFeed.Subscribe(ch) } // meteredConn is a wrapper around a net.Conn that meters both the @@ -95,6 +93,7 @@ type meteredConn struct { connected time.Time // Connection time of the peer ip net.IP // IP address of the peer id string // NodeID of the peer + metered bool // Used when the connection is closed to check if the peer was metered ingressMeter metrics.Meter // Meter for the read bytes of the peer egressMeter metrics.Meter // Meter for the written bytes of the peer @@ -114,12 +113,6 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { log.Warn("Peer IP is unspecified") return conn } - if atomic.LoadUint64(&meteredPeerCount) >= MeteredPeerLimit { - log.Warn("Metered peer count reached the limit") - return conn - } - // Increment the metered peer count - atomic.AddUint64(&meteredPeerCount, 1) // Bump the connection counters and wrap the connection if ingress { ingressConnectMeter.Mark(1) @@ -163,14 +156,20 @@ func (c *meteredConn) Write(b []byte) (n int, err error) { // the ingress and the egress traffic registries using the peer's IP and NodeID, // also emits connect event. func (c *meteredConn) handshakeDone(id enode.ID) { + if atomic.LoadInt32(&meteredPeerCount) >= MeteredPeerLimit { + log.Warn("Metered peer count reached the limit") + return + } + // Increment the metered peer count + atomic.AddInt32(&meteredPeerCount, 1) c.lock.Lock() - c.id = id.String() + c.id, c.metered = id.String(), true key := fmt.Sprintf("%s/%s", c.ip, c.id) c.ingressMeter = metrics.NewRegisteredMeter(key, PeerIngressRegistry) c.egressMeter = metrics.NewRegisteredMeter(key, PeerEgressRegistry) c.lock.Unlock() - metricsFeed.Send(MeteredPeerEvent{ + meteredPeerFeed.Send(MeteredPeerEvent{ Type: PeerConnected, IP: c.ip, ID: id.String(), @@ -181,14 +180,15 @@ func (c *meteredConn) handshakeDone(id enode.ID) { // Close delegates a close operation to the underlying connection, unregisters // the peer from the traffic registries and emits close event. func (c *meteredConn) Close() error { - // Decrement the metered peer count - atomic.AddUint64(&meteredPeerCount, ^uint64(0)) - c.lock.RLock() - // If the peer disconnects before the handshake + if c.metered { + // Decrement the metered peer count + atomic.AddInt32(&meteredPeerCount, -1) + } if c.id == "" { + // If the peer disconnects before the handshake c.lock.RUnlock() - metricsFeed.Send(MeteredPeerEvent{ + meteredPeerFeed.Send(MeteredPeerEvent{ Type: PeerHandshakeFailed, IP: c.ip, Elapsed: time.Since(c.connected), @@ -203,7 +203,7 @@ func (c *meteredConn) Close() error { PeerIngressRegistry.Unregister(key) PeerEgressRegistry.Unregister(key) - metricsFeed.Send(MeteredPeerEvent{ + meteredPeerFeed.Send(MeteredPeerEvent{ Type: PeerDisconnected, IP: c.ip, ID: id,