diff --git a/dashboard/message.go b/dashboard/message.go index dfe949bcce..f756cc1f63 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -83,8 +83,9 @@ func (m *NetworkMessage) getOrInitPeer(ip, id string) *Peer { // PeerBundle contains information about the peers pertaining to an IP address. type PeerBundle struct { - Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP - Peers map[string]*Peer `json:"peers,omitempty"` // the peers' node id is used as key + Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP + Peers map[string]*Peer `json:"peers,omitempty"` // the peers' node id is used as key + FailedPeers []*Peer `json:"failedPeers,omitempty"` } // GeoLocation contains geographical information. diff --git a/dashboard/peers.go b/dashboard/peers.go index eb64bb4ffc..7afb1e27a7 100644 --- a/dashboard/peers.go +++ b/dashboard/peers.go @@ -17,6 +17,8 @@ package dashboard import ( + "fmt" + "net" "time" "github.com/ethereum/go-ethereum/log" @@ -24,7 +26,8 @@ import ( "github.com/mohae/deepcopy" ) -const eventBufferLimit = 128 // Maximum number of buffered peer events for each event type +const eventBufferLimit = 128 // Maximum number of buffered peer events +const trafficEventBufferLimit = p2p.MeteredPeerLimit // collectPeerData gathers data about the peers and sends it to the clients. func (db *Dashboard) collectPeerData() { @@ -42,29 +45,27 @@ func (db *Dashboard) collectPeerData() { var ( // Peer event channels. connectCh = make(chan p2p.PeerConnectEvent, eventBufferLimit) - handshakeCh = make(chan p2p.PeerHandshakeEvent, eventBufferLimit) + failedCh = make(chan p2p.PeerFailedEvent, eventBufferLimit) disconnectCh = make(chan p2p.PeerDisconnectEvent, eventBufferLimit) - //readCh = make(chan p2p.PeerReadEvent, eventBufferLimit) - //writeCh = make(chan p2p.PeerWriteEvent, eventBufferLimit) + ingressCh = make(chan p2p.PeerTrafficEvent, trafficEventBufferLimit) + egressCh = make(chan p2p.PeerTrafficEvent, trafficEventBufferLimit) // Subscribe to peer events. subConnect = p2p.SubscribePeerConnectEvent(connectCh) - subHandshake = p2p.SubscribePeerHandshakeEvent(handshakeCh) + subFailed = p2p.SubscribePeerFailedEvent(failedCh) subDisconnect = p2p.SubscribePeerDisconnectEvent(disconnectCh) - //subRead = p2p.SubscribePeerReadEvent(readCh) - //subWrite = p2p.SubscribePeerWriteEvent(writeCh) + subIngress = p2p.SubscribePeerIngressEvent(ingressCh) + subEgress = p2p.SubscribePeerEgressEvent(egressCh) ) defer func() { // Unsubscribe at the end. subConnect.Unsubscribe() - subHandshake.Unsubscribe() + subFailed.Unsubscribe() subDisconnect.Unsubscribe() - //subRead.Unsubscribe() - //subWrite.Unsubscribe() + subIngress.Unsubscribe() + subEgress.Unsubscribe() }() - //go db.keepPeerHistoryClean(quit) - ticker := time.NewTicker(db.config.Refresh) defer ticker.Stop() @@ -74,30 +75,30 @@ func (db *Dashboard) collectPeerData() { } for { select { - //case event := <-connectCh: - // ip := event.IP.String() - // p := diff.getOrInitPeer(ip, event.ID) - // if diff.PeerBundles[ip].Location == nil { - // db.peerLock.RLock() - // lookup := db.history.Network.PeerBundles[ip] == nil || db.history.Network.PeerBundles[ip].Location == nil - // db.peerLock.RUnlock() - // if lookup { - // location := db.geodb.Lookup(event.IP) - // diff.PeerBundles[ip].Location = &GeoLocation{ - // Country: location.Country.Names.English, - // City: location.City.Names.English, - // Latitude: location.Location.Latitude, - // Longitude: location.Location.Longitude, - // } - // } - // } - // if p.Connected == nil { - // p.Connected = []time.Time{event.Connected} - // } else { - // p.Connected = append(p.Connected, event.Connected) - // } - //case event := <-handshakeCh: - // ip := event.IP.String() + case event := <-connectCh: + ip := event.IP + p := diff.getOrInitPeer(ip, event.ID) + if diff.PeerBundles[ip].Location == nil { + db.peerLock.RLock() + lookup := db.history.Network.PeerBundles[ip] == nil || db.history.Network.PeerBundles[ip].Location == nil + db.peerLock.RUnlock() + if lookup { + location := db.geodb.Lookup(net.ParseIP(event.IP)) + diff.PeerBundles[ip].Location = &GeoLocation{ + Country: location.Country.Names.English, + City: location.City.Names.English, + Latitude: location.Location.Latitude, + Longitude: location.Location.Longitude, + } + } + } + if p.Connected == nil { + p.Connected = []time.Time{event.Connected} + } else { + p.Connected = append(p.Connected, event.Connected) + } + //case event := <-failedCh: + // ip := event.IP // p := diff.getOrInitPeer(ip, event.AutoID) // p.DefaultID = event.AutoID // if p.Handshake == nil { @@ -120,29 +121,31 @@ func (db *Dashboard) collectPeerData() { // db.history.Network.PeerBundles[ip].Peers[event.ID] = hp // TODO (kurkomisi): Merge. // db.peerLock.Unlock() // } - //case event := <-disconnectCh: - // p := diff.getOrInitPeer(event.IP.String(), event.ID) - // if p.Disconnected == nil { - // p.Disconnected = []time.Time{event.Disconnected} - // } else { - // p.Disconnected = append(p.Disconnected, event.Disconnected) - // } - //case event := <-readCh: - // // Sum up the ingress between two updates. - // p := diff.getOrInitPeer(event.IP.String(), event.ID) - // if len(p.Ingress) <= 0 { - // p.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Ingress)}} - // } else { - // p.Ingress[0].Value += float64(event.Ingress) - // } - //case event := <-writeCh: - // // Sum up the egress between two updates. - // p := diff.getOrInitPeer(event.IP.String(), event.ID) - // if len(p.Egress) <= 0 { - // p.Egress = ChartEntries{&ChartEntry{Value: float64(event.Egress)}} - // } else { - // p.Egress[0].Value += float64(event.Egress) - // } + case event := <-disconnectCh: + p := diff.getOrInitPeer(event.IP, event.ID) + if p.Disconnected == nil { + p.Disconnected = []time.Time{event.Disconnected} + } else { + p.Disconnected = append(p.Disconnected, event.Disconnected) + } + case event := <-ingressCh: + fmt.Println("ingress", event.IP, event.Amount) + // Sum up the ingress between two updates. + p := diff.getOrInitPeer(event.IP, event.ID) + if len(p.Ingress) <= 0 { + p.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} + } else { + p.Ingress[0].Value += float64(event.Amount) + } + case event := <-egressCh: + fmt.Println("egress ", event.IP, event.Amount) + // Sum up the egress between two updates. + p := diff.getOrInitPeer(event.IP, event.ID) + if len(p.Egress) <= 0 { + p.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} + } else { + p.Egress[0].Value += float64(event.Amount) + } case <-ticker.C: now := time.Now() // Merge the diff with the history. @@ -206,18 +209,18 @@ func (db *Dashboard) collectPeerData() { case err := <-subConnect.Err(): log.Warn("Peer connect subscription error", "err", err) return - case err := <-subHandshake.Err(): - log.Warn("Peer handshake subscription error", "err", err) + case err := <-subFailed.Err(): + log.Warn("Peer failed subscription error", "err", err) return case err := <-subDisconnect.Err(): log.Warn("Peer disconnect subscription error", "err", err) return - //case err := <-subRead.Err(): - // log.Warn("Peer read subscription error", "err", err) - // return - //case err := <-subWrite.Err(): - // log.Warn("Peer write subscription error", "err", err) - // return + case err := <-subIngress.Err(): + log.Warn("Peer ingress subscription error", "err", err) + return + case err := <-subEgress.Err(): + log.Warn("Peer egress subscription error", "err", err) + return case errc := <-db.quit: errc <- nil return diff --git a/p2p/metrics.go b/p2p/metrics.go index 9031201931..b5ad84ea49 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -56,17 +56,16 @@ var ( metricsFeed = new(peerMetricsFeed) // Peer event feed for metrics - meteredPeerAutoID uint64 // Used to create unique id for the metered connection before the handshake - meteredPeerCount uint64 + meteredPeerCount uint64 ) // peerMetricsFeed delivers the peer metrics to the subscribed channels. type peerMetricsFeed struct { - connect event.Feed // Event feed to notify the connection of a peer - handshake event.Feed // Event feed to notify the handshake with a peer + connect event.Feed // Event feed to notify the connection and the successful handshake of a peer + ingress event.Feed // Event feed to notify the amount of read bytes of a peer + egress event.Feed // Event feed to notify the amount of written bytes of a peer disconnect event.Feed // Event feed to notify the disconnection of a peer - read event.Feed // Event feed to notify the amount of read bytes of a peer - write event.Feed // Event feed to notify the amount of written bytes of a peer + failed event.Feed // Event feed to notify the connection of a peer and its disconnection before the handshake scope event.SubscriptionScope // Facility to unsubscribe all the subscriptions at once @@ -75,79 +74,86 @@ type peerMetricsFeed struct { // PeerConnectEvent contains information about the connection of a peer. type PeerConnectEvent struct { - Key string + IP string + ID string Connected time.Time -} - -// PeerHandshakeEvent contains information about the handshake with a peer. -type PeerHandshakeEvent struct { - AutoKey string - Key string - Ingress int64 - Egress int64 Handshake time.Time } // PeerDisconnectEvent contains information about the disconnection of a peer. type PeerDisconnectEvent struct { - Key string - Ingress int64 - Egress int64 + IP string + ID string Disconnected time.Time } -// PeerReadEvent contains information about the read operation of a peer. -type PeerTrafficEvent map[string]int64 +type PeerTrafficEvent struct { + IP string + ID string + Amount int64 +} + +type PeerFailedEvent struct { + IP string + Connected time.Time + Disconnected time.Time +} // SubscribePeerConnectEvent registers a subscription of PeerConnectEvent func SubscribePeerConnectEvent(ch chan<- PeerConnectEvent) event.Subscription { return metricsFeed.scope.Track(metricsFeed.connect.Subscribe(ch)) } -// SubscribePeerHandshakeEvent registers a subscription of PeerHandshakeEvent -func SubscribePeerHandshakeEvent(ch chan<- PeerHandshakeEvent) event.Subscription { - return metricsFeed.scope.Track(metricsFeed.handshake.Subscribe(ch)) -} - // SubscribePeerDisconnectEvent registers a subscription of PeerDisconnectEvent func SubscribePeerDisconnectEvent(ch chan<- PeerDisconnectEvent) event.Subscription { return metricsFeed.scope.Track(metricsFeed.disconnect.Subscribe(ch)) } -// SubscribePeerReadEvent registers a subscription of PeerReadEvent -func SubscribePeerReadEvent(ch chan<- PeerTrafficEvent) event.Subscription { - return metricsFeed.scope.Track(metricsFeed.read.Subscribe(ch)) +// SubscribePeerTrafficEvent registers a subscription of PeerTrafficEvent + +func SubscribePeerIngressEvent(ch chan<- PeerTrafficEvent) event.Subscription { + return metricsFeed.scope.Track(metricsFeed.ingress.Subscribe(ch)) } -// SubscribePeerWriteEvent registers a subscription of PeerWriteEvent -func SubscribePeerWriteEvent(ch chan<- PeerTrafficEvent) event.Subscription { - return metricsFeed.scope.Track(metricsFeed.write.Subscribe(ch)) +func SubscribePeerEgressEvent(ch chan<- PeerTrafficEvent) event.Subscription { + return metricsFeed.scope.Track(metricsFeed.egress.Subscribe(ch)) } -func startTrafficNotifier(refresh time.Duration) { +// SubscribePeerFailedEvent registers a subscription of PeerFailedEvent +func SubscribePeerFailedEvent(ch chan<- PeerFailedEvent) event.Subscription { + return metricsFeed.scope.Track(metricsFeed.failed.Subscribe(ch)) +} + +func runMetricsFeedHelper(refresh time.Duration) { metricsFeed.quit = make(chan chan error) ticker := time.NewTicker(refresh) defer ticker.Stop() + + // It is possible to send all of the traffic events together, but it is risky to use pointers in the events. + trafficEventSender := func(prefix string, feed *event.Feed) func(name string, i interface{}) { + return func(name string, i interface{}) { + if m, ok := i.(metrics.Meter); ok { + // Trim the common prefix and split the peer specific part in order to get the ip and the node id. + if key := strings.Split(strings.TrimPrefix(name, prefix), "/"); len(key) == 2 { + feed.Send(PeerTrafficEvent{ + IP: key[0], + ID: key[1], + Amount: m.Count(), + }) + } else { + log.Warn("Invalid peer metrics name", "name", name) + } + } + } + } + sendIngress := trafficEventSender(MetricsRegistryIngressPrefix, &metricsFeed.ingress) + sendEgress := trafficEventSender(MetricsRegistryEgressPrefix, &metricsFeed.egress) + for { select { case <-ticker.C: - // send read and write - ingressEvents, egressEvents := make(PeerTrafficEvent), make(PeerTrafficEvent) - PeerIngressRegistry.Each(func(name string, i interface{}) { - if m, ok := i.(metrics.Meter); ok { - ingressEvents[strings.TrimPrefix(name, MetricsRegistryIngressPrefix)] = m.Count() - } - }) - PeerEgressRegistry.Each(func(name string, i interface{}) { - if m, ok := i.(metrics.Meter); ok { - egressEvents[strings.TrimPrefix(name, MetricsRegistryEgressPrefix)] = m.Count() - } - }) - metricsFeed.read.Send(ingressEvents) - metricsFeed.write.Send(egressEvents) - //fmt.Println(ingressEvents) - //fmt.Println(egressEvents) - //fmt.Println() + PeerIngressRegistry.Each(sendIngress) + PeerEgressRegistry.Each(sendEgress) case errc := <-metricsFeed.quit: errc <- nil return @@ -170,10 +176,11 @@ func closeMetricsFeed() { // meteredConn is a wrapper around a net.Conn that meters both the // inbound and outbound network traffic. type meteredConn struct { - net.Conn // Network connection to wrap with metering - ip string // The IP address of the peer + net.Conn // Network connection to wrap with metering - key string + connected time.Time + ip string // The IP address of the peer + id string // The NodeID of the peer ingressMeter metrics.Meter egressMeter metrics.Meter @@ -196,24 +203,18 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { log.Warn("Metered peer count reached the limit") return conn } + // Increment the metered peer count atomic.AddUint64(&meteredPeerCount, 1) - // Otherwise bump the connection counters and wrap the connection + // Bump the connection counters and wrap the connection if ingress { ingressConnectMeter.Mark(1) } else { egressConnectMeter.Mark(1) } - key := fmt.Sprintf("%s/%s", ip.String(), fmt.Sprintf("peer_%d", atomic.AddUint64(&meteredPeerAutoID, 1))) - metricsFeed.connect.Send(PeerConnectEvent{ - Key: key, - Connected: time.Now(), - }) return &meteredConn{ - Conn: conn, - key: key, - ip: ip.String(), - ingressMeter: metrics.NewRegisteredMeter(key, PeerIngressRegistry), - egressMeter: metrics.NewRegisteredMeter(key, PeerEgressRegistry), + Conn: conn, + ip: ip.String(), + connected: time.Now(), } } @@ -222,7 +223,11 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { func (c *meteredConn) Read(b []byte) (n int, err error) { n, err = c.Conn.Read(b) ingressTrafficMeter.Mark(int64(n)) - c.ingressMeter.Mark(int64(n)) + c.lock.RLock() + if c.ingressMeter != nil { + c.ingressMeter.Mark(int64(n)) + } + c.lock.RUnlock() return n, err } @@ -231,55 +236,73 @@ func (c *meteredConn) Read(b []byte) (n int, err error) { func (c *meteredConn) Write(b []byte) (n int, err error) { n, err = c.Conn.Write(b) egressTrafficMeter.Mark(int64(n)) - c.egressMeter.Mark(int64(n)) + c.lock.RLock() + if c.egressMeter != nil { + c.egressMeter.Mark(int64(n)) + } + c.lock.RUnlock() return n, err } // Close closes the underlying connection. func (c *meteredConn) Close() error { - // Decrement the metered peer count. + // Decrement the metered peer count atomic.AddUint64(&meteredPeerCount, ^uint64(0)) - c.ingressMeter.Stop() - c.egressMeter.Stop() + err, now := c.Conn.Close(), time.Now() + c.lock.RLock() - key := c.key - metricsFeed.disconnect.Send(PeerDisconnectEvent{ - Key: key, - Ingress: c.ingressMeter.Count(), - Egress: c.egressMeter.Count(), - Disconnected: time.Now(), - }) + ip, id := c.ip, c.id c.lock.RUnlock() + + // If the peer disconnects before the handshake + if id == "" { + metricsFeed.failed.Send(PeerFailedEvent{ + IP: ip, + Connected: c.connected, + Disconnected: now, + }) + return err + } + c.lock.RLock() + //ingress, egress := c.ingressMeter.Count(), c.egressMeter.Count() + c.lock.RUnlock() + + // Unregister the peer from the metrics registry + key := fmt.Sprintf("%s/%s", ip, id) PeerIngressRegistry.Unregister(key) PeerEgressRegistry.Unregister(key) - return c.Conn.Close() + + //metricsFeed.ingress.Send(PeerTrafficEvent{ + // IP: ip, + // ID: id, + // Amount: ingress, + //}) + //metricsFeed.egress.Send(PeerTrafficEvent{ + // IP: ip, + // ID: id, + // Amount: egress, + //}) + metricsFeed.disconnect.Send(PeerDisconnectEvent{ + IP: ip, + ID: id, + Disconnected: now, + }) + return err } // handshakeDone changes the default id to the peer's node id. func (c *meteredConn) handshakeDone(id discover.NodeID) { - c.ingressMeter.Stop() - c.egressMeter.Stop() c.lock.Lock() - - autoKey := c.key - key := fmt.Sprintf("%s/%s", c.ip, id.String()) - ingressMeter := metrics.NewRegisteredMeter(key, PeerIngressRegistry) - egressMeter := metrics.NewRegisteredMeter(key, PeerEgressRegistry) - ingressMeter.Mark(c.ingressMeter.Count()) - egressMeter.Mark(c.egressMeter.Count()) - PeerIngressRegistry.Unregister(c.key) - PeerEgressRegistry.Unregister(c.key) - c.key = key - c.ingressMeter = ingressMeter - c.egressMeter = egressMeter - + c.id = id.String() + 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.handshake.Send(PeerHandshakeEvent{ - AutoKey: autoKey, - Key: key, - //Ingress: nil, - //Egress: nil, + metricsFeed.connect.Send(PeerConnectEvent{ + IP: c.ip, + ID: id.String(), + Connected: c.connected, Handshake: time.Now(), }) } diff --git a/p2p/server.go b/p2p/server.go index a9bd9c7c54..eaa6f1f933 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -542,7 +542,7 @@ func (srv *Server) Start() (err error) { srv.loopWG.Add(1) go srv.run(dialer) - go startTrafficNotifier(2 * time.Second) + go runMetricsFeedHelper(5 * time.Second) srv.running = true return nil }