p2p: requested changes

This commit is contained in:
Kurkó Mihály 2018-09-19 11:42:37 +03:00
parent 93e6fdf58a
commit 263304dfa0

View file

@ -37,10 +37,7 @@ const (
MetricsOutboundConnects = "p2p/OutboundConnects" // Name for the registered outbound connects meter MetricsOutboundConnects = "p2p/OutboundConnects" // Name for the registered outbound connects meter
MetricsOutboundTraffic = "p2p/OutboundTraffic" // Name for the registered outbound traffic meter MetricsOutboundTraffic = "p2p/OutboundTraffic" // Name for the registered outbound traffic meter
MetricsRegistryIngressPrefix = MetricsInboundTraffic + "/" MeteredPeerLimit = 1024 // This amount of peers are individually metered
MetricsRegistryEgressPrefix = MetricsOutboundTraffic + "/"
MeteredPeerLimit = 1024
) )
var ( var (
@ -49,11 +46,11 @@ var (
egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) // Meter counting the egress connections egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) // Meter counting the egress connections
egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) // Meter metering the cumulative egress traffic egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) // Meter metering the cumulative egress traffic
PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) // Registry containing the peer ingress PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsInboundTraffic+"/") // Registry containing the peer ingress
PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) // Registry containing the peer egress PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsOutboundTraffic+"/") // Registry containing the peer egress
metricsFeed event.Feed // Event feed for peer metrics meteredPeerFeed event.Feed // Event feed for peer metrics
meteredPeerCount uint64 // Actually stored peer connection count meteredPeerCount int32 // Actually stored peer connection count
) )
// MeteredPeerEventType is the type of peer events emitted by a metered connection. // MeteredPeerEventType is the type of peer events emitted by a metered connection.
@ -72,19 +69,20 @@ const (
PeerHandshakeFailed 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 MeteredPeerEvent struct {
Type MeteredPeerEventType // Type of peer event Type MeteredPeerEventType // Type of peer event
IP net.IP // IP address of the peer IP net.IP // IP address of the peer
ID string // NodeID of the peer ID string // NodeID of the peer
Elapsed time.Duration // Time elapsed between the connection and the handshake/disconnection Elapsed time.Duration // Time elapsed between the connection and the handshake/disconnection
Ingress uint64 // Ingress count in the moment of disconnection Ingress uint64 // Ingress count at the moment of the event
Egress uint64 // Egress count in the moment of disconnection 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 { 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 // 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 connected time.Time // Connection time of the peer
ip net.IP // IP address of the peer ip net.IP // IP address of the peer
id string // NodeID 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 ingressMeter metrics.Meter // Meter for the read bytes of the peer
egressMeter metrics.Meter // Meter for the written 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") log.Warn("Peer IP is unspecified")
return conn 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 // Bump the connection counters and wrap the connection
if ingress { if ingress {
ingressConnectMeter.Mark(1) 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, // the ingress and the egress traffic registries using the peer's IP and NodeID,
// also emits connect event. // also emits connect event.
func (c *meteredConn) handshakeDone(id enode.ID) { 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.lock.Lock()
c.id = id.String() c.id, c.metered = id.String(), true
key := fmt.Sprintf("%s/%s", c.ip, c.id) key := fmt.Sprintf("%s/%s", c.ip, c.id)
c.ingressMeter = metrics.NewRegisteredMeter(key, PeerIngressRegistry) c.ingressMeter = metrics.NewRegisteredMeter(key, PeerIngressRegistry)
c.egressMeter = metrics.NewRegisteredMeter(key, PeerEgressRegistry) c.egressMeter = metrics.NewRegisteredMeter(key, PeerEgressRegistry)
c.lock.Unlock() c.lock.Unlock()
metricsFeed.Send(MeteredPeerEvent{ meteredPeerFeed.Send(MeteredPeerEvent{
Type: PeerConnected, Type: PeerConnected,
IP: c.ip, IP: c.ip,
ID: id.String(), ID: id.String(),
@ -181,14 +180,15 @@ func (c *meteredConn) handshakeDone(id enode.ID) {
// Close delegates a close operation to the underlying connection, unregisters // Close delegates a close operation to the underlying connection, unregisters
// the peer from the traffic registries and emits close event. // the peer from the traffic registries and emits close event.
func (c *meteredConn) Close() error { func (c *meteredConn) Close() error {
// Decrement the metered peer count
atomic.AddUint64(&meteredPeerCount, ^uint64(0))
c.lock.RLock() 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 c.id == "" {
// If the peer disconnects before the handshake
c.lock.RUnlock() c.lock.RUnlock()
metricsFeed.Send(MeteredPeerEvent{ meteredPeerFeed.Send(MeteredPeerEvent{
Type: PeerHandshakeFailed, Type: PeerHandshakeFailed,
IP: c.ip, IP: c.ip,
Elapsed: time.Since(c.connected), Elapsed: time.Since(c.connected),
@ -203,7 +203,7 @@ func (c *meteredConn) Close() error {
PeerIngressRegistry.Unregister(key) PeerIngressRegistry.Unregister(key)
PeerEgressRegistry.Unregister(key) PeerEgressRegistry.Unregister(key)
metricsFeed.Send(MeteredPeerEvent{ meteredPeerFeed.Send(MeteredPeerEvent{
Type: PeerDisconnected, Type: PeerDisconnected,
IP: c.ip, IP: c.ip,
ID: id, ID: id,