From a87086b6084c0f623005a68ffc9883fbcc0d509d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kurk=C3=B3=20Mih=C3=A1ly?= Date: Tue, 18 Sep 2018 19:02:05 +0300 Subject: [PATCH] dashboard, p2p: experiment 2 --- dashboard/assets.go | 2 +- dashboard/assets/components/Network.jsx | 30 +-- dashboard/assets/webpack.config.dev.js | 2 +- dashboard/message.go | 42 +++- dashboard/peers.go | 322 ++++++++++++++---------- p2p/metrics.go | 274 +++++++------------- p2p/server.go | 2 - 7 files changed, 321 insertions(+), 353 deletions(-) diff --git a/dashboard/assets.go b/dashboard/assets.go index 649cd03a6b..e09aac76c0 100644 --- a/dashboard/assets.go +++ b/dashboard/assets.go @@ -11363,7 +11363,7 @@ var _bundleJs = []byte((((((((((`!function(modules) { key: "render", value: function() { var _this2 = this; - return _react2.default.createElement(_Table2.default, null, _react2.default.createElement(_Table.TableHead, null, _react2.default.createElement(_Table.TableRow, null, _react2.default.createElement(_Table.TableCell, null, "IP"), _react2.default.createElement(_Table.TableCell, null, "Location"), _react2.default.createElement(_Table.TableCell, null, "Peer ID"), _react2.default.createElement(_Table.TableCell, null, "Ingress"), _react2.default.createElement(_Table.TableCell, null, "Egress"), _react2.default.createElement(_Table.TableCell, null, "Connected"), _react2.default.createElement(_Table.TableCell, null, "Handshake"), _react2.default.createElement(_Table.TableCell, null, "Disconnected"))), _react2.default.createElement(_Table.TableBody, null, Object.entries(this.props.content.peerBundles).map(function(_ref4) { + return _react2.default.createElement(_Table2.default, null, _react2.default.createElement(_Table.TableHead, null, _react2.default.createElement(_Table.TableRow, null, _react2.default.createElement(_Table.TableCell, null, "IP"), _react2.default.createElement(_Table.TableCell, null, "Location"), _react2.default.createElement(_Table.TableCell, null, "Peer ID"), _react2.default.createElement(_Table.TableCell, null, "Ingress"), _react2.default.createElement(_Table.TableCell, null, "Egress"), _react2.default.createElement(_Table.TableCell, null, "Connected"), _react2.default.createElement(_Table.TableCell, null, "Handshake"), _react2.default.createElement(_Table.TableCell, null, "Time"))), _react2.default.createElement(_Table.TableBody, null, Object.entries(this.props.content.peerBundles).map(function(_ref4) { var _ref5 = _slicedToArray(_ref4, 2), ip = _ref5[0], bundle = _ref5[1]; return _react2.default.createElement(_Table.TableRow, { key: ip diff --git a/dashboard/assets/components/Network.jsx b/dashboard/assets/components/Network.jsx index da89d1a3d2..b54633ea01 100644 --- a/dashboard/assets/components/Network.jsx +++ b/dashboard/assets/components/Network.jsx @@ -46,12 +46,6 @@ export const inserter = (update: {[string]: PeerBundle}, prev: {[string]: PeerBu prev[ip].peers[id] = u; return; } - // If the handshake was between two metering - if (u.defaultID && prev[ip].peers[u.defaultID]) { - // TODO (kurkomisi): merge the two in order to keep the previous connection. - prev[ip].peers[id] = prev[ip].peers[u.defaultID]; - delete prev[ip].peers[u.defaultID]; - } const p: Peer = prev[ip].peers[id]; if (u.connected) { if (!Array.isArray(p.connected)) { @@ -59,12 +53,6 @@ export const inserter = (update: {[string]: PeerBundle}, prev: {[string]: PeerBu } p.connected = [...p.connected, ...u.connected]; } - if (u.handshake) { - if (!Array.isArray(p.handshake)) { - p.handshake = []; - } - p.handshake = [...p.handshake, ...u.handshake]; - } if (u.disconnected) { if (!Array.isArray(p.disconnected)) { p.disconnected = []; @@ -126,12 +114,11 @@ class Network extends Component { Ingress Egress Connected - Handshake Disconnected - {Object.entries(this.props.content.peerBundles).map(([ip, bundle]) => ( + {Object.entries(this.props.content.peerBundles).map(([ip, bundle]) => { console.log(ip, bundle); return ( {ip} @@ -141,25 +128,22 @@ class Network extends Component { })() : ''} - {Object.keys(bundle.peers).map(id => id.substring(0, 10)).join(' ')} + {bundle.peers && Object.keys(bundle.peers).map(id => id.substring(0, 10)).join(' ')} - {Object.values(bundle.peers).map(peer => peer.ingress && peer.ingress.map(sample => sample.value).join(' ')).join(', ')} + {bundle.peers && Object.values(bundle.peers).map(peer => peer.ingress && peer.ingress.map(sample => sample.value).join(' ')).join(', ')} - {Object.values(bundle.peers).map(peer => peer.egress && peer.egress.map(sample => sample.value).join(' ')).join(', ')} + {bundle.peers && Object.values(bundle.peers).map(peer => peer.egress && peer.egress.map(sample => sample.value).join(' ')).join(', ')} - {Object.values(bundle.peers).map(peer => peer.connected && peer.connected.map(time => this.formatTime(time)).join(' ')).join(', ')} + {bundle.peers && Object.values(bundle.peers).map(peer => peer.connected && peer.connected.map(time => this.formatTime(time)).join(' ')).join(', ')} - {Object.values(bundle.peers).map(peer => peer.handshake && peer.handshake.map(time => this.formatTime(time)).join(' ')).join(', ')} - - - {Object.values(bundle.peers).map(peer => peer.disconnected && peer.disconnected.map(time => this.formatTime(time)).join(' ')).join(', ')} + {bundle.peers && Object.values(bundle.peers).map(peer => peer.disconnected && peer.disconnected.map(time => this.formatTime(time)).join(' ')).join(', ')} - ))} + )})} ); diff --git a/dashboard/assets/webpack.config.dev.js b/dashboard/assets/webpack.config.dev.js index 3f9cc1d795..4223aa4897 100644 --- a/dashboard/assets/webpack.config.dev.js +++ b/dashboard/assets/webpack.config.dev.js @@ -24,7 +24,7 @@ module.exports = merge(common, { new webpack.HotModuleReplacementPlugin(), ], // devtool: 'eval', - devtool: 'inline-source-map', + devtool: 'source-map', devServer: { port: 8081, hot: true, diff --git a/dashboard/message.go b/dashboard/message.go index 69309d053e..af3cf2002b 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -66,7 +66,8 @@ type NetworkMessage struct { func (m *NetworkMessage) getOrInitBundle(ip string) *PeerBundle { if _, ok := m.PeerBundles[ip]; !ok { m.PeerBundles[ip] = &PeerBundle{ - Peers: make(map[string]*Peer), + Peers: make(PeerMap), + FailedPeers: make(PeerMap), } } return m.PeerBundles[ip] @@ -78,18 +79,40 @@ func (m *NetworkMessage) getOrInitPeer(ip, id string) *Peer { return m.getOrInitBundle(ip).getOrInitPeer(id) } +type PeerMap map[string]*Peer + +func (pm PeerMap) getOrInit(id string) *Peer { + if _, ok := pm[id]; !ok { + pm[id] = new(Peer) + } + return pm[id] +} + +func (pm PeerMap) remove(id string) { + delete(pm, id) +} + // 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 - FailedPeers []*Peer `json:"failedPeers,omitempty"` + Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP + Peers PeerMap `json:"peers,omitempty"` // the peers' node id is used as key + FailedPeers PeerMap `json:"failedPeers,omitempty"` } func (b *PeerBundle) getOrInitPeer(id string) *Peer { - if _, ok := b.Peers[id]; !ok { - b.Peers[id] = new(Peer) - } - return b.Peers[id] + return b.Peers.getOrInit(id) +} + +func (b *PeerBundle) removePeer(id string) { + b.Peers.remove(id) +} + +func (b *PeerBundle) getOrInitFailedPeer(id string) *Peer { + return b.FailedPeers.getOrInit(id) +} + +func (b * PeerBundle) removeFailedPeer(id string) { + b.FailedPeers.remove(id) } // GeoLocation contains geographical information. @@ -108,9 +131,8 @@ type Peer struct { Ingress ChartEntries `json:"ingress,omitempty"` Egress ChartEntries `json:"egress,omitempty"` - DefaultID string `json:"defaultID,omitempty"` - element *list.Element + ip, id string } // SystemMessage contains the metered system data samples. diff --git a/dashboard/peers.go b/dashboard/peers.go index d29f7ae926..2a2cb60608 100644 --- a/dashboard/peers.go +++ b/dashboard/peers.go @@ -18,9 +18,7 @@ package dashboard import ( "container/list" - "encoding/json" "fmt" - "github.com/mohae/deepcopy" "time" "github.com/ethereum/go-ethereum/log" @@ -29,6 +27,58 @@ import ( const eventBufferLimit = 128 // Maximum number of buffered peer events const trafficEventBufferLimit = p2p.MeteredPeerLimit +const connectionLimit = 100 + +var autoID int64 + +type peerLimiter struct { + underlying *NetworkMessage + failed bool + l *list.List +} + +func NewPeerLimiter(underlying *NetworkMessage, failed bool) *peerLimiter { + return &peerLimiter{l: list.New(), underlying: underlying, failed: failed} +} + +func (pl *peerLimiter) update(peer *Peer) { + return + if peer.element == nil { + peer.element = pl.l.PushBack(peer) + } else { + pl.l.MoveToBack(peer.element) + } + for pl.l.Len() > 2 {//p2p.MeteredPeerLimit { + pl.remove(pl.l.Front()) + } +} + +func (pl *peerLimiter) remove(e *list.Element) { + return + elem := pl.l.Remove(e) + if peer, ok := elem.(*Peer); ok { + if pl.failed { + fmt.Println(peer.ip, peer.id) + pl.underlying.PeerBundles[peer.ip].FailedPeers.remove(peer.id) + } else { + fmt.Println(peer.ip, peer.id[:10]) + pl.underlying.PeerBundles[peer.ip].Peers.remove(peer.id) + } + } +} + +func (pl *peerLimiter) clear() { + for pl.l.Front() != nil { + pl.remove(pl.l.Front()) + } +} +// +//func tail(arr []time.Time) []time.Time { +// if first := len(arr)-connectionLimit; first > 0 { +// return arr[first:] +// } +// return arr +//} // collectPeerData gathers data about the peers and sends it to the clients. func (db *Dashboard) collectPeerData() { @@ -45,148 +95,154 @@ func (db *Dashboard) collectPeerData() { var ( // Peer event channels. - connectCh = make(chan p2p.PeerConnectEvent, eventBufferLimit) - failedCh = make(chan p2p.PeerFailedEvent, eventBufferLimit) - disconnectCh = make(chan p2p.PeerDisconnectEvent, eventBufferLimit) - ingressCh = make(chan p2p.PeerTrafficEvent, trafficEventBufferLimit) - egressCh = make(chan p2p.PeerTrafficEvent, trafficEventBufferLimit) - + peerCh = make(chan p2p.MeteredPeerEvent, eventBufferLimit) // Subscribe to peer events. - subConnect = p2p.SubscribePeerConnectEvent(connectCh) - subFailed = p2p.SubscribePeerFailedEvent(failedCh) - subDisconnect = p2p.SubscribePeerDisconnectEvent(disconnectCh) - subIngress = p2p.SubscribePeerIngressEvent(ingressCh) - subEgress = p2p.SubscribePeerEgressEvent(egressCh) + subPeer = p2p.SubscribePeerEvent(peerCh) ) defer func() { // Unsubscribe at the end. - subConnect.Unsubscribe() - subFailed.Unsubscribe() - subDisconnect.Unsubscribe() - subIngress.Unsubscribe() - subEgress.Unsubscribe() + subPeer.Unsubscribe() }() ticker := time.NewTicker(db.config.Refresh) defer ticker.Stop() - purgeOrder := list.New() - failedPurgeOrder := list.New() - update := func(peer *Peer, l *list.List) { - if peer.element == nil { - peer.element = purgeOrder.PushBack(peer) - } else { - purgeOrder.MoveToBack(peer.element) - } - } - // Listen for events, and prepare the difference between two metering. - diff := &NetworkMessage{ - PeerBundles: make(map[string]*PeerBundle), - } + db.peerLock.RLock() + //historyPeerLimiter := NewPeerLimiter(db.history.Network, false) + //historyFailedPeerLimiter := NewPeerLimiter(db.history.Network, true) + //db.peerLock.RUnlock() + //// Listen for events, and prepare the difference between two metering. + //diff := &NetworkMessage{ + // PeerBundles: make(map[string]*PeerBundle), + //} + //// Needed in order to keep the limit in the diff + //diffPeerLimiter := NewPeerLimiter(diff, false) + //diffFailedPeerLimiter := NewPeerLimiter(diff, true) for { select { - case event := <-connectCh: - diffBundle := diff.getOrInitBundle(event.IP) - diffBundle.Location = db.geodb.Location(event.IP) - diffPeer := diffBundle.getOrInitPeer(event.ID) - diffPeer.Connected = append(diffPeer.Connected, event.Connected) - case event := <-failedCh: - diffBundle := diff.getOrInitBundle(event.IP) - diffBundle.Location = db.geodb.Location(event.IP) - diffBundle.FailedPeers = append(diffBundle.FailedPeers, &Peer{ - Connected: []time.Time{event.Connected}, - Disconnected: []time.Time{event.Disconnected}, - }) - case event := <-disconnectCh: - diffPeer := diff.getOrInitPeer(event.IP, event.ID) - diffPeer.Disconnected = append(diffPeer.Disconnected, event.Disconnected) - case event := <-ingressCh: - diffPeer := diff.getOrInitPeer(event.IP, event.ID) - if len(diffPeer.Ingress) != 1 { - diffPeer.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} - } else { - diffPeer.Ingress[0].Value = float64(event.Amount) - } - case event := <-egressCh: - diffPeer := diff.getOrInitPeer(event.IP, event.ID) - if len(diffPeer.Egress) != 1 { - diffPeer.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} - } else { - diffPeer.Egress[0].Value = float64(event.Amount) - } - case <-ticker.C: - now := time.Now() - // Merge the diff with the history. - db.peerLock.Lock() - for ip, diffBundle := range diff.PeerBundles { - historyBundle := db.history.Network.getOrInitBundle(ip) - historyBundle.Location = diffBundle.Location - for id, diffPeer := range diffBundle.Peers { - historyPeer := historyBundle.getOrInitPeer(id) - historyPeer.Connected = append(historyPeer.Connected, diffPeer.Connected...) - historyPeer.Disconnected = append(historyPeer.Disconnected, diffPeer.Disconnected...) - if len(diffPeer.Ingress) == 1 { - diffPeer.Ingress[0].Time = now - if historyPeer.Ingress == nil { - historyPeer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), sampleLimit-1, db.config.Refresh), diffPeer.Ingress[0]) - // The first message about a diffPeer should contain the whole list - diffPeer.Ingress = historyPeer.Ingress - } else { - historyPeer.Ingress = append(historyPeer.Ingress, diffPeer.Ingress[0])[1:] - } - } - if len(diffPeer.Egress) == 1 { - diffPeer.Egress[0].Time = now - if historyPeer.Egress == nil { - historyPeer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), sampleLimit-1, db.config.Refresh), diffPeer.Egress[0]) - // The first message about a diffPeer should contain the whole list - diffPeer.Egress = historyPeer.Egress - } else { - historyPeer.Egress = append(historyPeer.Egress, diffPeer.Egress[0])[1:] - } - } - update(historyPeer, purgeOrder) - } - historyBundle.FailedPeers = append(historyBundle.FailedPeers, diffBundle.FailedPeers...) - for _, fp := range diffBundle.FailedPeers { - update(fp, failedPurgeOrder) - } - } - for purgeOrder.Len() > p2p.MeteredPeerLimit { - purgeOrder.Remove(purgeOrder.Front()) - } - for failedPurgeOrder.Len() > p2p.MeteredPeerLimit { - failedPurgeOrder.Remove(failedPurgeOrder.Front()) - } - - //ss, _ := json.MarshalIndent(db.history.Network, "", " ") - //fmt.Println(string(ss)) - db.peerLock.Unlock() - // Send the diff to the clients. - db.sendToAll(&Message{Network: deepcopy.Copy(diff).(*NetworkMessage)}) - //s, _ := json.MarshalIndent(diff, "", " ") - //fmt.Println(string(s)) - - s, _ := json.MarshalIndent(diff, "", " ") - fmt.Println(string(s)) - // Prepare for the next metering, clear the diff variable. - diff = &NetworkMessage{ - PeerBundles: make(map[string]*PeerBundle), - } - case err := <-subConnect.Err(): - log.Warn("Peer connect subscription error", "err", err) - return - 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 := <-subIngress.Err(): - log.Warn("Peer ingress subscription error", "err", err) - return - case err := <-subEgress.Err(): - log.Warn("Peer egress subscription error", "err", err) + case event := <-peerCh: + fmt.Println(event) + //case event := <-connectCh: + // diffBundle := diff.getOrInitBundle(event.IP) + // diffBundle.Location = db.geodb.Location(event.IP) + // diffPeer := diffBundle.getOrInitPeer(event.ID) + // diffPeer.Connected = append(diffPeer.Connected, event.Connected) + // if first := len(diffPeer.Connected)-connectionLimit; first > 0 { + // diffPeer.Connected = diffPeer.Connected[first:] + // } + // diffPeer.ip = event.IP + // diffPeer.id = event.ID + // diffPeerLimiter.update(diffPeer) + //case event := <-disconnectCh: + // diffPeer := diff.getOrInitPeer(event.IP, event.ID) + // diffPeer.Disconnected = append(diffPeer.Disconnected, event.Time) + // if first := len(diffPeer.Connected)-connectionLimit; first > 0 { + // diffPeer.Connected = diffPeer.Connected[first:] + // } + // diffPeer.ip = event.IP + // diffPeer.id = event.ID + // diffPeerLimiter.update(diffPeer) + //case event := <-ingressCh: + // diffPeer := diff.getOrInitPeer(event.IP, event.ID) + // if len(diffPeer.Ingress) != 1 { + // diffPeer.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} + // } else { + // diffPeer.Ingress[0].Value = float64(event.Amount) + // } + // diffPeer.ip = event.IP + // diffPeer.id = event.ID + // diffPeerLimiter.update(diffPeer) + //case event := <-egressCh: + // diffPeer := diff.getOrInitPeer(event.IP, event.ID) + // if len(diffPeer.Egress) != 1 { + // diffPeer.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} + // } else { + // diffPeer.Egress[0].Value = float64(event.Amount) + // } + // diffPeer.ip = event.IP + // diffPeer.id = event.ID + // diffPeerLimiter.update(diffPeer) + //case event := <-failedCh: + // diffBundle := diff.getOrInitBundle(event.IP) + // diffBundle.Location = db.geodb.Location(event.IP) + // id := fmt.Sprintf("peer_%d", atomic.AddInt64(&autoID, 1)) + // failedPeer := diffBundle.FailedPeers.getOrInit(id) + // failedPeer.Connected = []time.Time{event.Connected} + // failedPeer.Disconnected = []time.Time{event.Disconnected} + // failedPeer.ip = event.IP + // failedPeer.id = id + // diffFailedPeerLimiter.update(failedPeer) + //case <-ticker.C: + // now := time.Now() + // // Merge the diff with the history. + // db.peerLock.Lock() + // for ip, diffBundle := range diff.PeerBundles { + // historyBundle := db.history.Network.getOrInitBundle(ip) + // historyBundle.Location = diffBundle.Location + // for id, diffPeer := range diffBundle.Peers { + // historyPeer := historyBundle.getOrInitPeer(id) + // historyPeer.Connected = append(historyPeer.Connected, diffPeer.Connected...) + // if first := len(historyPeer.Connected)-connectionLimit; first > 0 { + // historyPeer.Connected = historyPeer.Connected[first:] + // } + // historyPeer.Disconnected = append(historyPeer.Disconnected, diffPeer.Disconnected...) + // if first := len(historyPeer.Disconnected)-connectionLimit; first > 0 { + // historyPeer.Disconnected = historyPeer.Disconnected[first:] + // } + // if len(diffPeer.Ingress) == 1 { + // diffPeer.Ingress[0].Time = now + // if historyPeer.Ingress == nil { + // historyPeer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), 3/*sampleLimit-1*/, db.config.Refresh), diffPeer.Ingress[0]) + // // The first message about a diffPeer should contain the whole list + // diffPeer.Ingress = historyPeer.Ingress + // } else { + // historyPeer.Ingress = append(historyPeer.Ingress, diffPeer.Ingress[0])[1:] + // } + // } + // if len(diffPeer.Egress) == 1 { + // diffPeer.Egress[0].Time = now + // if historyPeer.Egress == nil { + // historyPeer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), 3/*sampleLimit-1*/, db.config.Refresh), diffPeer.Egress[0]) + // // The first message about a diffPeer should contain the whole list + // diffPeer.Egress = historyPeer.Egress + // } else { + // historyPeer.Egress = append(historyPeer.Egress, diffPeer.Egress[0])[1:] + // } + // } + // historyPeer.ip = diffPeer.ip + // historyPeer.id = diffPeer.id + // historyPeerLimiter.update(historyPeer) + // } + // for id, diffFailedPeer := range diffBundle.FailedPeers { + // historyFailedPeer := historyBundle.getOrInitFailedPeer(id) + // historyFailedPeer.Connected = diffFailedPeer.Connected + // historyFailedPeer.Disconnected = diffFailedPeer.Disconnected + // historyFailedPeer.ip = diffFailedPeer.ip + // historyFailedPeer.id = diffFailedPeer.id + // historyFailedPeerLimiter.update(historyFailedPeer) + // } + // } + // //for elem := historyPeerLimiter.l.Front(); elem != nil; elem = elem.Next() { + // // s, _ := json.MarshalIndent(elem.Value, "", " ") + // // fmt.Println(string(s)) + // //} + // //fmt.Println() + // db.peerLock.Unlock() + // + // //s, _ := json.MarshalIndent(deepcopy.Copy(diff), "", " ") + // //fmt.Println(string(s)) + // //fmt.Println() + // // Send the diff to the clients. + // db.sendToAll(&Message{Network: deepcopy.Copy(diff).(*NetworkMessage)}) + // + // // Prepare for the next metering, clear the diff variable. + // diffPeerLimiter.clear() + // diffFailedPeerLimiter.clear() + // diff = &NetworkMessage{ + // PeerBundles: make(map[string]*PeerBundle), + // } + case err := <-subPeer.Err(): + log.Warn("Peer subscription error", "err", err) return case errc := <-db.quit: errc <- nil diff --git a/p2p/metrics.go b/p2p/metrics.go index 0bbb55e5ef..1af13df5a9 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -19,10 +19,8 @@ package p2p import ( - "net" - "strings" - "fmt" + "net" "sync" "sync/atomic" "time" @@ -46,130 +44,47 @@ const ( ) var ( - ingressConnectMeter = metrics.NewRegisteredMeter(MetricsInboundConnects, nil) // meter counting the ingress connections - ingressTrafficMeter = metrics.NewRegisteredMeter(MetricsInboundTraffic, nil) // meter metering the cumulative ingress traffic - egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) // meter counting the egress connections - egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) // meter metering the cumulative egress traffic + ingressConnectMeter = metrics.NewRegisteredMeter(MetricsInboundConnects, nil) // Meter counting the ingress connections + ingressTrafficMeter = metrics.NewRegisteredMeter(MetricsInboundTraffic, nil) // Meter metering the cumulative ingress traffic + 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) - PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) + PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) // Registry containing the peer ingress + PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) // Registry containing the peer egress - metricsFeed = new(peerMetricsFeed) // Peer event feed for metrics - - meteredPeerCount uint64 + metricsFeed event.Feed // Event feed for peer metrics + meteredPeerCount uint64 // Actually stored peer connection count ) -// peerMetricsFeed delivers the peer metrics to the subscribed channels. -type peerMetricsFeed struct { - 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 - failed event.Feed // Event feed to notify the connection of a peer and its disconnection before the handshake +// MeteredPeerEventType is the type of peer events emitted by a metered connection. +type MeteredPeerEventType int - scope event.SubscriptionScope // Facility to unsubscribe all the subscriptions at once +const ( + // PeerConnected is the type of event emitted when a peer successfully + // made the handshake. + PeerConnected MeteredPeerEventType = iota - quit chan chan error + // PeerDisconnected is the type of event emitted when a peer disconnects. + PeerDisconnected + + // PeerHandshakeFailed is the type of event emitted when a peer fails to + // make the handshake or disconnects before the handshake. + PeerHandshakeFailed +) + +// 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 } -// PeerConnectEvent contains information about the connection of a peer. -type PeerConnectEvent struct { - IP string - ID string - Connected time.Time -} - -// PeerDisconnectEvent contains information about the disconnection of a peer. -type PeerDisconnectEvent struct { - IP string - ID string - Disconnected time.Time -} - -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)) -} - -// SubscribePeerDisconnectEvent registers a subscription of PeerDisconnectEvent -func SubscribePeerDisconnectEvent(ch chan<- PeerDisconnectEvent) event.Subscription { - return metricsFeed.scope.Track(metricsFeed.disconnect.Subscribe(ch)) -} - -// SubscribePeerTrafficEvent registers a subscription of PeerTrafficEvent - -func SubscribePeerIngressEvent(ch chan<- PeerTrafficEvent) event.Subscription { - return metricsFeed.scope.Track(metricsFeed.ingress.Subscribe(ch)) -} - -func SubscribePeerEgressEvent(ch chan<- PeerTrafficEvent) event.Subscription { - return metricsFeed.scope.Track(metricsFeed.egress.Subscribe(ch)) -} - -// 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: - PeerIngressRegistry.Each(sendIngress) - PeerEgressRegistry.Each(sendEgress) - case errc := <-metricsFeed.quit: - errc <- nil - return - } - } -} - -// closeMetricsFeed closes all the tracked subscriptions. -func closeMetricsFeed() { - if metricsFeed.quit != nil { - errc := make(chan error) - metricsFeed.quit <- errc - <-errc - } - metricsFeed.scope.Close() - PeerIngressRegistry.UnregisterAll() - PeerEgressRegistry.UnregisterAll() +// SubscribePeerEvent registers a subscription of PeerEvent +func SubscribePeerEvent(ch chan<- MeteredPeerEvent) event.Subscription { + return metricsFeed.Subscribe(ch) } // meteredConn is a wrapper around a net.Conn that meters both the @@ -177,18 +92,19 @@ func closeMetricsFeed() { type meteredConn struct { net.Conn // Network connection to wrap with metering - connected time.Time - ip string // The IP address of the peer - id string // The NodeID of the peer - ingressMeter metrics.Meter - egressMeter metrics.Meter + connected time.Time // Connection time of the peer + ip net.IP // IP address of the peer + id string // NodeID 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 lock sync.RWMutex // Lock protecting the metered connection's internals } -// newMeteredConn creates a new metered connection, also bumping the ingress or -// egress connection meter. If the metrics system is disabled, this function -// returns the original object. +// newMeteredConn creates a new metered connection, bumps the ingress or egress +// connection meter and also increases the metered peer count. If the metrics +// system is disabled, the IP address is unspecified or the metered peer count +// reached the limit, this function returns the original object. func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { // Short circuit if metrics are disabled if !metrics.Enabled { @@ -212,13 +128,13 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { } return &meteredConn{ Conn: conn, - ip: ip.String(), + ip: ip, connected: time.Now(), } } -// Read delegates a network read to the underlying connection, bumping the ingress -// traffic meter along the way. +// Read delegates a network read to the underlying connection, bumping the common +// and the peer ingress traffic meters along the way. func (c *meteredConn) Read(b []byte) (n int, err error) { n, err = c.Conn.Read(b) ingressTrafficMeter.Mark(int64(n)) @@ -230,8 +146,8 @@ func (c *meteredConn) Read(b []byte) (n int, err error) { return n, err } -// Write delegates a network write to the underlying connection, bumping the -// egress traffic meter along the way. +// Write delegates a network write to the underlying connection, bumping the common +// and the peer egress traffic meters along the way. func (c *meteredConn) Write(b []byte) (n int, err error) { n, err = c.Conn.Write(b) egressTrafficMeter.Mark(int64(n)) @@ -243,53 +159,9 @@ func (c *meteredConn) Write(b []byte) (n int, err error) { return n, err } -// Close closes the underlying connection. -func (c *meteredConn) Close() error { - // Decrement the metered peer count - atomic.AddUint64(&meteredPeerCount, ^uint64(0)) - err, now := c.Conn.Close(), time.Now() - - c.lock.RLock() - 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) - - //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. +// handshakeDone is called when a peer handshake is done. Registers the peer to +// the ingress and the egress traffic registries using the peer's IP and NodeID, +// also emits connect event. func (c *meteredConn) handshakeDone(id discover.NodeID) { c.lock.Lock() c.id = id.String() @@ -298,9 +170,45 @@ func (c *meteredConn) handshakeDone(id discover.NodeID) { c.egressMeter = metrics.NewRegisteredMeter(key, PeerEgressRegistry) c.lock.Unlock() - metricsFeed.connect.Send(PeerConnectEvent{ - IP: c.ip, - ID: id.String(), - Connected: c.connected, + metricsFeed.Send(MeteredPeerEvent{ + Type: PeerConnected, + IP: c.ip, + ID: id.String(), + Elapsed: time.Now().Sub(c.connected), }) } + +// 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.id == "" { + c.lock.RUnlock() + metricsFeed.Send(MeteredPeerEvent{ + Type: PeerHandshakeFailed, + IP: c.ip, + Elapsed: time.Now().Sub(c.connected), + }) + return c.Conn.Close() + } + id, ingress, egress := c.id, uint64(c.ingressMeter.Count()), uint64(c.egressMeter.Count()) + c.lock.RUnlock() + + // Unregister the peer from the traffic registries + key := fmt.Sprintf("%s/%s", c.ip, id) + PeerIngressRegistry.Unregister(key) + PeerEgressRegistry.Unregister(key) + + metricsFeed.Send(MeteredPeerEvent{ + Type: PeerDisconnected, + IP: c.ip, + ID: id, + Ingress: ingress, + Egress: egress, + }) + return c.Conn.Close() +} diff --git a/p2p/server.go b/p2p/server.go index eaa6f1f933..a432e27269 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -387,7 +387,6 @@ func (srv *Server) Stop() { } close(srv.quit) srv.lock.Unlock() - closeMetricsFeed() srv.loopWG.Wait() } @@ -542,7 +541,6 @@ func (srv *Server) Start() (err error) { srv.loopWG.Add(1) go srv.run(dialer) - go runMetricsFeedHelper(5 * time.Second) srv.running = true return nil }