From 8fd0e4ad1eb38a8cc3e2cadf60c72788ee0d2a8e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kurk=C3=B3=20Mih=C3=A1ly?= Date: Tue, 7 Aug 2018 14:28:25 +0300 Subject: [PATCH] dashboard, p2p: initial version of peer traffic metering with event feeds --- dashboard/assets/components/Network.jsx | 129 +++++++--- dashboard/assets/types/content.jsx | 20 +- dashboard/dashboard.go | 128 ++-------- dashboard/message.go | 24 +- dashboard/peers.go | 301 ++++++++++++++++++++++++ p2p/metrics.go | 208 +++++++--------- p2p/server.go | 1 + 7 files changed, 533 insertions(+), 278 deletions(-) create mode 100644 dashboard/peers.go diff --git a/dashboard/assets/components/Network.jsx b/dashboard/assets/components/Network.jsx index 19773a8a34..aad95403f6 100644 --- a/dashboard/assets/components/Network.jsx +++ b/dashboard/assets/components/Network.jsx @@ -18,36 +18,59 @@ import React, {Component} from 'react'; -import Table, {TableBody, TableHeader, TableHeaderColumn, TableRow, TableCell} from 'material-ui/Table'; +import Table, {TableHead, TableBody, TableRow, TableCell} from 'material-ui/Table'; import type {Network as NetworkType, Peer} from '../types/content'; // inserter is a state updater function for the main component, which inserts the new log chunk into the chunk array. // limit is the maximum length of the chunk array, used in order to prevent the browser from OOM. -export const inserter = (update: {[number]: Peer}, prev: {[number]: Peer}) => { - Object.keys(update).forEach((k) => { - if (!prev[k]) { - prev[k] = update[k]; +export const inserter = (update: {[string]: {[string]: Peer}}, prev: {[string]: {[string]: Peer}}) => { + Object.keys(update).forEach((ip) => { + if (!prev[ip]) { + prev[ip] = update[ip]; return; } - const u: Peer = update[k]; - const p: Peer = prev[k]; - if (u.id) { - p.id = u.id; + if (!update[ip]) { + return; } - if (u.ip) { - p.ip = u.ip; - } - if (u.lifecycle) { - if (u.lifecycle.handshake) { - p.lifecycle.handshake = u.lifecycle.handshake; + Object.keys(update[ip]).forEach((id) => { + if (!prev[ip][id]) { + prev[ip][id] = update[ip][id]; + return; } - if (u.lifecycle.disconnected) { - p.lifecycle.disconnected = u.lifecycle.disconnected; + const u: Peer = update[ip][id]; + const p: Peer = prev[ip][id]; + if (u.connected) { + if (!Array.isArray(p.connected)) { + p.connected = []; + } + p.connected = [...p.connected, ...u.connected]; } - } - p.ingress = [...p.ingress, ...u.ingress].slice(-200); - p.egress = [...p.egress, ...u.egress].slice(-200); - prev[k] = p; + 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 = []; + } + p.disconnected = [...p.disconnected, ...u.disconnected]; + } + if (Array.isArray(u.ingress)) { + if (!Array.isArray(p.ingress)) { + p.ingress = []; + } + p.ingress = [...p.ingress, ...u.ingress].slice(-200); + } + if (Array.isArray(u.egress)) { + if (!Array.isArray(p.egress)) { + p.egress = []; + } + p.egress = [...p.egress, ...u.egress].slice(-200); + } + prev[ip][id] = p; + }); }); return prev; }; @@ -63,19 +86,65 @@ export type Props = { // Network renders the network page. class Network extends Component { + formatTime = (t) => { + const time = new Date(t); + if (isNaN(time)) { + return ''; + } + const month = `0${time.getMonth() + 1}`.slice(-2); + const date = `0${time.getDate()}`.slice(-2); + const hours = `0${time.getHours()}`.slice(-2); + const minutes = `0${time.getMinutes()}`.slice(-2); + const seconds = `0${time.getSeconds()}`.slice(-2); + return `${month}/${date}/${hours}:${minutes}:${seconds}`; + }; + render() { return ( + + + IP + Peer ID + Location + Ingress + Egress + Connected + Handshake + Disconnected + + - {Object.entries(this.props.content.peers).map(([k, v]) => ( - - {k} - {v.id ? v.id.substring(0, 6) : ''} - {v.ip} - {v.ingress.value} - {v.egress.value} - {JSON.stringify(v.location)} - {JSON.stringify(v.lifecycle)} + {Object.entries(this.props.content.peers).map(([ip, peers]) => ( + + {ip} + + {Object.keys(peers).map(id => id.substring(0, 10)).join(' ')} + + + {(() => { + const k = Object.keys(peers)[0]; + return k && peers[k].location ? (() => { + const l = peers[k].location; + return `${l.country}${l.city ? `/${l.city}` : ''} ${l.latitude} ${l.longitude}`; + })() : ''; + })()} + + + {Object.keys(peers).map((id) => peers[id].ingress && peers[id].ingress.map(sample => sample.value).join(' ')).join(', ')} + + + {Object.keys(peers).map((id) => peers[id].egress && peers[id].egress.map(sample => sample.value).join(' ')).join(', ')} + + + {Object.keys(peers).map((id) => peers[id].connected && peers[id].connected.map(time => this.formatTime(time)).join(' ')).join(', ')} + + + {Object.keys(peers).map((id) => peers[id].handshake && peers[id].handshake.map(time => this.formatTime(time)).join(' ')).join(', ')} + + + {Object.keys(peers).map((id) => peers[id].disconnected && peers[id].disconnected.map(time => this.formatTime(time)).join(' ')).join(', ')} + ))} diff --git a/dashboard/assets/types/content.jsx b/dashboard/assets/types/content.jsx index 691db4987a..94f70537a7 100644 --- a/dashboard/assets/types/content.jsx +++ b/dashboard/assets/types/content.jsx @@ -51,16 +51,16 @@ export type TxPool = { }; export type Network = { - peers: {[number]: Peer}, + peers: {[string]: {[string]: Peer}}, }; export type Peer = { - id: string, - ip: string, - location: PeerLocation, - lifecycle: PeerLifecycle, - ingress: ChartEntries, - egress: ChartEntries, + location: PeerLocation, + connected: Array, + handshake: Array, + disconnected: Array, + ingress: ChartEntries, + egress: ChartEntries, }; export type PeerLocation = { @@ -70,12 +70,6 @@ export type PeerLocation = { longitude: number, }; -export type PeerLifecycle = { - connected: Date, - handshake: Date, - disconnected: Date, -}; - export type System = { activeMemory: ChartEntries, virtualMemory: ChartEntries, diff --git a/dashboard/dashboard.go b/dashboard/dashboard.go index f92a4c6d64..dbd44b0620 100644 --- a/dashboard/dashboard.go +++ b/dashboard/dashboard.go @@ -54,9 +54,9 @@ const ( diskReadSampleLimit = 200 // Maximum number of disk read data samples diskWriteSampleLimit = 200 // Maximum number of disk write data samples - storedDisconnectedPeerLimit = p2p.MeteredPeerLimit - p2p.DefaultMaxPendingPeers - peerIngressSampleLimit = 200 - peerEgressSampleLimit = 200 + peerTrafficSampleLimit = 200 + peerIngressSampleLimit = peerTrafficSampleLimit + peerEgressSampleLimit = peerTrafficSampleLimit ) var nextID uint32 // Next connection id @@ -68,11 +68,14 @@ type Dashboard struct { listener net.Listener conns map[uint32]*client // Currently live websocket connections history *Message + peerHistory *NetworkMessage lock sync.RWMutex // Lock protecting the dashboard's internals + peerLock sync.RWMutex - geodb *GeoDB - disconnectedPeerIDs chan uint - logdir string + geodb *GeoDB + peersByIP map[string]*Peer + peersByID map[string]*Peer + logdir string quit chan chan error // Channel used for graceful exit wg sync.WaitGroup @@ -111,12 +114,11 @@ func New(config *Config, commit string, logdir string) *Dashboard { DiskRead: emptyChartEntries(now, diskReadSampleLimit, config.Refresh), DiskWrite: emptyChartEntries(now, diskWriteSampleLimit, config.Refresh), }, - Network: &NetworkMessage{ - Peers: make(map[uint]*Peer), - }, }, - disconnectedPeerIDs: make(chan uint, storedDisconnectedPeerLimit), - logdir: logdir, + peerHistory: &NetworkMessage{ + Peers: make(map[string]map[string]*Peer), + }, + logdir: logdir, } } @@ -142,16 +144,10 @@ func (db *Dashboard) APIs() []rpc.API { return nil } func (db *Dashboard) Start(server *p2p.Server) error { log.Info("Starting dashboard") - var err error - db.geodb, err = OpenGeoDB() - if err != nil { - log.Warn("Failed to open geodb", "err", err) - return err - } - - db.wg.Add(2) + db.wg.Add(3) go db.collectData() go db.streamLogs() + go db.collectPeerData() http.HandleFunc("/", db.webHandler) http.Handle("/api", websocket.Handler(db.apiHandler)) @@ -170,8 +166,6 @@ func (db *Dashboard) Start(server *p2p.Server) error { // Stop stops the data collection thread and the connection listener of the dashboard. // Implements the node.Service interface. func (db *Dashboard) Stop() error { - db.geodb.Close() - // Close the connection listener. var errs []error if err := db.listener.Close(); err != nil { @@ -179,7 +173,7 @@ func (db *Dashboard) Stop() error { } // Close the collectors. errc := make(chan error, 1) - for i := 0; i < 2; i++ { + for i := 0; i < 3; i++ { db.quit <- errc if err := <-errc; err != nil { errs = append(errs, err) @@ -252,10 +246,17 @@ func (db *Dashboard) apiHandler(conn *websocket.Conn) { } }() - db.lock.Lock() // Send the past data. - client.msg <- deepcopy.Copy(db.history).(*Message) + db.lock.RLock() + h := deepcopy.Copy(db.history).(*Message) + db.lock.RUnlock() + db.peerLock.RLock() + h.Network = deepcopy.Copy(db.peerHistory).(*NetworkMessage) + db.peerLock.RUnlock() + client.msg <- h + // Start tracking the connection and drop at connection loss. + db.lock.Lock() db.conns[id] = client db.lock.Unlock() defer func() { @@ -392,84 +393,6 @@ func (db *Dashboard) collectData() { sys.DiskWrite = append(sys.DiskWrite[1:], diskWrite) db.lock.Unlock() - peers := p2p.PeerTrafficMeters.Peers() - network := new(NetworkMessage) - if len(peers) > 0 { - network.Peers = make(map[uint]*Peer) - } - for id, peer := range peers { - peerIngress := &ChartEntry{ - Time: now, - Value: float64(peer.Ingress), - } - peerEgress := &ChartEntry{ - Time: now, - Value: float64(peer.Egress), - } - p := &Peer{ - Ingress: ChartEntries{peerIngress}, - Egress: ChartEntries{peerEgress}, - } - if prevMetrics, ok := db.history.Network.Peers[id]; !ok { - p.ID = peer.ID - p.IP = peer.IP - location := db.geodb.Lookup(peer.IP) - p.Location = &PeerLocation{ - Country: location.Country.Names.English, - City: location.City.Names.English, - Latitude: location.Location.Latitude, - Longitude: location.Location.Longitude, - } - p.Lifecycle = &PeerLifecycle{ - Connected: peer.Connected, - } - if peer.Handshake != nil { - p.Lifecycle.Handshake = peer.Handshake - } - if peer.Disconnected != nil { - p.Lifecycle.Disconnected = peer.Disconnected - } - db.history.Network.Peers[id] = p - } else { - if prevMetrics.ID != peer.ID { - prevMetrics.ID, p.ID = peer.ID, peer.ID - } - if prevMetrics.Lifecycle.Handshake == nil && peer.Handshake != nil { - p.Lifecycle = new(PeerLifecycle) - prevMetrics.Lifecycle.Handshake, p.Lifecycle.Handshake = peer.Handshake, peer.Handshake - } - if prevMetrics.Lifecycle.Disconnected == nil && peer.Disconnected != nil { - if p.Lifecycle == nil { - p.Lifecycle = new(PeerLifecycle) - } - prevMetrics.Lifecycle.Disconnected, p.Lifecycle.Disconnected = peer.Disconnected, peer.Disconnected - } - first := 0 - if len(prevMetrics.Ingress) >= peerIngressSampleLimit { - first = len(prevMetrics.Ingress) - peerIngressSampleLimit + 1 - } - prevMetrics.Ingress = append(prevMetrics.Ingress[first:], peerIngress) - first = 0 - if len(prevMetrics.Egress) >= peerEgressSampleLimit { - first = len(prevMetrics.Egress) - peerEgressSampleLimit + 1 - } - prevMetrics.Egress = append(prevMetrics.Egress[first:], peerEgress) - } - if p.Lifecycle != nil && p.Lifecycle.Disconnected != nil { - select { - case db.disconnectedPeerIDs <- id: - default: - // if the number of the stored disconnected peers exceeds the limit, remove the firstly disconnected peer - delete(db.history.Network.Peers, <-db.disconnectedPeerIDs) - db.disconnectedPeerIDs <- id - } - } - network.Peers[id] = p - } - db.lock.Unlock() - //s, _ := json.MarshalIndent(network, "", " ") - //fmt.Println(string(s)) - db.sendToAll(&Message{ System: &SystemMessage{ ActiveMemory: ChartEntries{activeMemory}, @@ -481,7 +404,6 @@ func (db *Dashboard) collectData() { DiskRead: ChartEntries{diskRead}, DiskWrite: ChartEntries{diskWrite}, }, - Network: network, }) } } diff --git a/dashboard/message.go b/dashboard/message.go index c1c6a32445..cb9bd0cc8b 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -18,7 +18,6 @@ package dashboard import ( "encoding/json" - "net" "time" ) @@ -56,17 +55,18 @@ type TxPoolMessage struct { /* TODO (kurkomisi) */ } +// k1: IP, k2: ID type NetworkMessage struct { - Peers map[uint]*Peer `json:"peers,omitempty"` + Peers map[string]map[string]*Peer `json:"peers,omitempty"` } type Peer struct { - ID string `json:"id,omitempty"` - IP net.IP `json:"ip,omitempty"` - Location *PeerLocation `json:"location,omitempty"` - Lifecycle *PeerLifecycle `json:"lifecycle,omitempty"` - Ingress ChartEntries `json:"ingress,omitempty"` - Egress ChartEntries `json:"egress,omitempty"` + Location *PeerLocation `json:"location,omitempty"` + Connected []time.Time `json:"connected,omitempty"` + Handshake []time.Time `json:"handshake,omitempty"` + Disconnected []time.Time `json:"disconnected,omitempty"` + Ingress ChartEntries `json:"ingress,omitempty"` + Egress ChartEntries `json:"egress,omitempty"` } type PeerLocation struct { @@ -76,12 +76,6 @@ type PeerLocation struct { Longitude float64 `json:"longitude,omitempty"` } -type PeerLifecycle struct { - Connected *time.Time `json:"connected,omitempty"` - Handshake *time.Time `json:"handshake,omitempty"` - Disconnected *time.Time `json:"disconnected,omitempty"` -} - type SystemMessage struct { ActiveMemory ChartEntries `json:"activeMemory,omitempty"` VirtualMemory ChartEntries `json:"virtualMemory,omitempty"` @@ -93,7 +87,7 @@ type SystemMessage struct { DiskWrite ChartEntries `json:"diskWrite,omitempty"` } -// LogsMessage wraps up a log chunk. If Source isn't present, the chunk is a stream chunk. +// LogsMessage wraps up a log chunk. If 'Source' isn't present, the chunk is a stream chunk. type LogsMessage struct { Source *LogFile `json:"source,omitempty"` // Attributes of the log file. Chunk json.RawMessage `json:"chunk"` // Contains log records. diff --git a/dashboard/peers.go b/dashboard/peers.go new file mode 100644 index 0000000000..06d80a0704 --- /dev/null +++ b/dashboard/peers.go @@ -0,0 +1,301 @@ +package dashboard + +import ( + "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/p2p" + "time" + "github.com/mohae/deepcopy" +) + +const eventBufferLimit = 128 + +func getOrInitPeer(m *NetworkMessage, ip, id string) *Peer { + if _, ok := m.Peers[ip]; !ok { + m.Peers[ip] = make(map[string]*Peer) + } + if _, ok := m.Peers[ip][id]; !ok { + m.Peers[ip][id] = new(Peer) + } + return m.Peers[ip][id] +} + +func (db *Dashboard) collectPeerData() { + defer db.wg.Done() + + var err error + db.geodb, err = OpenGeoDB() + if err != nil { + log.Warn("Failed to open geodb", "err", err) + return + } + defer db.geodb.Close() + + var ( + quit = make(chan struct{}) + connectCh = make(chan *p2p.PeerConnectEvent, eventBufferLimit) + handshakeCh = make(chan *p2p.PeerHandshakeEvent, eventBufferLimit) + disconnectCh = make(chan *p2p.PeerDisconnectEvent, eventBufferLimit) + readCh = make(chan *p2p.PeerReadEvent, eventBufferLimit) + writeCh = make(chan *p2p.PeerWriteEvent, eventBufferLimit) + ) + go func() { + var ( + peerConnectEventCh = make(chan p2p.PeerConnectEvent, eventBufferLimit) + peerHandshakeEventCh = make(chan p2p.PeerHandshakeEvent, eventBufferLimit) + peerDisconnectEventCh = make(chan p2p.PeerDisconnectEvent, eventBufferLimit) + peerReadEventCh = make(chan p2p.PeerReadEvent, eventBufferLimit) + peerWriteEventCh = make(chan p2p.PeerWriteEvent, eventBufferLimit) + + subConnect = p2p.SubscribePeerConnectEvent(peerConnectEventCh) + subHandshake = p2p.SubscribePeerHandshakeEvent(peerHandshakeEventCh) + subDisconnect = p2p.SubscribePeerDisconnectEvent(peerDisconnectEventCh) + subRead = p2p.SubscribePeerReadEvent(peerReadEventCh) + subWrite = p2p.SubscribePeerWriteEvent(peerWriteEventCh) + ) + defer func() { + subConnect.Unsubscribe() + subHandshake.Unsubscribe() + subDisconnect.Unsubscribe() + subRead.Unsubscribe() + subWrite.Unsubscribe() + }() + for { + select { + case event := <-peerConnectEventCh: + select { + case connectCh <- &event: + default: + log.Warn("Failed to handle connect event", "event", event) + } + case event := <-peerHandshakeEventCh: + select { + case handshakeCh <- &event: + default: + log.Warn("Failed to handle handshake event", "event", event) + } + case event := <-peerDisconnectEventCh: + select { + case disconnectCh <- &event: + default: + log.Warn("Failed to handle disconnect event", "event", event) + } + case event := <-peerReadEventCh: + select { + case readCh <- &event: + default: + log.Warn("Failed to handle read event", "event", event) + } + case event := <-peerWriteEventCh: + select { + case writeCh <- &event: + default: + log.Warn("Failed to handle write event", "event", event) + } + case <-quit: + return + } + } + }() + go db.cleanPeerHistory(quit) + + ticker := time.NewTicker(db.config.Refresh) + defer ticker.Stop() + + network := &NetworkMessage{ + Peers: make(map[string]map[string]*Peer), + } + for { + select { + case event := <-connectCh: + ip := event.IP.String() + p := getOrInitPeer(network, ip, event.ID) + if p.Location == nil { + db.peerLock.RLock() + peers := db.peerHistory.Peers + lookup := peers[ip] == nil || peers[ip][event.ID] == nil || peers[ip][event.ID].Location == nil + db.peerLock.RUnlock() + if lookup { + location := db.geodb.Lookup(event.IP) + p.Location = &PeerLocation{ + 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() + p := getOrInitPeer(network, ip, event.DefaultID) + if p.Handshake == nil { + p.Handshake = []time.Time{event.Handshake} + } else { + p.Handshake = append(p.Handshake, event.Handshake) + } + delete(network.Peers[ip], event.DefaultID) + getOrInitPeer(network, ip, event.ID) + network.Peers[ip][event.ID] = p // interleave instead + // Remove the peer from history in case the metering was before the handshake. + db.peerLock.RLock() + stored := db.peerHistory.Peers[ip] != nil && db.peerHistory.Peers[ip][event.DefaultID] != nil + db.peerLock.RUnlock() + if stored { + db.peerLock.Lock() + hp := getOrInitPeer(db.peerHistory, ip, event.DefaultID) + delete(db.peerHistory.Peers[ip], event.DefaultID) + getOrInitPeer(db.peerHistory, ip, event.ID) + db.peerHistory.Peers[ip][event.ID] = hp // interleave instead + db.peerLock.Unlock() + } + case event := <-disconnectCh: + p := getOrInitPeer(network, 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 := getOrInitPeer(network, 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 := getOrInitPeer(network, 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 <-ticker.C: + now := time.Now() + db.peerLock.Lock() + for ip, peers := range network.Peers { + for id, peer := range peers { + peerHistory := getOrInitPeer(db.peerHistory, ip, id) + if peer.Location != nil { + peerHistory.Location = peer.Location + } + if peer.Connected != nil { + peerHistory.Connected = append(peerHistory.Connected, peer.Connected...) + } + if peer.Handshake != nil { + peerHistory.Handshake = append(peerHistory.Handshake, peer.Handshake...) + } + if peer.Disconnected != nil { + peerHistory.Disconnected = append(peerHistory.Disconnected, peer.Disconnected...) + } + ingress := &ChartEntry{ + Time: now, + } + if len(peer.Ingress) > 0 { + ingress.Value = peer.Ingress[0].Value + } + if peerHistory.Ingress == nil { + peer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), peerIngressSampleLimit-1, db.config.Refresh), ingress) + peerHistory.Ingress = peer.Ingress + //peerHistory.Ingress = ChartEntries{ingress} + } else { + peer.Ingress = ChartEntries{ingress} + peerHistory.Ingress = append(peerHistory.Ingress[1:], ingress) + //peerHistory.Ingress = append(peerHistory.Ingress, ingress) + } + egress := &ChartEntry{ + Time: now, + } + if len(peer.Egress) > 0 { + egress.Value = peer.Egress[0].Value + } + if peerHistory.Egress == nil { + peer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), peerEgressSampleLimit-1, db.config.Refresh), egress) + peerHistory.Egress = peer.Egress + //peerHistory.Egress = ChartEntries{egress} + } else { + peer.Egress = ChartEntries{egress} + peerHistory.Egress = append(peerHistory.Egress[1:], egress) + //peerHistory.Egress = append(peerHistory.Egress, egress) + } + } + } + db.peerLock.Unlock() + db.sendToAll(&Message{Network: deepcopy.Copy(network).(*NetworkMessage)}) + + //fmt.Println() + //s, _ := json.MarshalIndent(network, "", " ") + //fmt.Println(string(s)) + + for ip, peers := range network.Peers { + for id := range peers { + peers[id] = nil + delete(peers, id) + } + delete(network.Peers, ip) + } + case errc := <-db.quit: + close(quit) + errc <- nil + return + } + } +} + +func (db *Dashboard) cleanPeerHistory(quit chan struct{}) { + cleanRate := db.config.Refresh * peerTrafficSampleLimit + for { + select { + case <-time.After(cleanRate): + // clear disconnected + validAfter := time.Now().Add(-cleanRate) + db.peerLock.Lock() + for ip, peers := range db.peerHistory.Peers { + for id, peer := range peers { + if len(peer.Disconnected) > 0 && peer.Disconnected[len(peer.Disconnected)-1].Before(validAfter) { + db.peerHistory.Peers[ip][id].Location = nil + db.peerHistory.Peers[ip][id] = nil + delete(db.peerHistory.Peers[ip], id) + } + } + if len(peers) <= 0 { + delete(db.peerHistory.Peers, ip) + } + } + var lenCount int + for _, peers := range db.peerHistory.Peers { + lenCount += len(peers) + } + if lenCount > p2p.MeteredPeerLimit { + outerLoop: + for ip, peers := range db.peerHistory.Peers { + for id, peer := range peers { + if peer.Disconnected != nil { + db.peerHistory.Peers[ip][id].Location = nil + db.peerHistory.Peers[ip][id] = nil + delete(db.peerHistory.Peers[ip], id) + lenCount-- + if lenCount <= p2p.MeteredPeerLimit { + if len(peers) <= 0 { + delete(db.peerHistory.Peers, ip) + } + break outerLoop + } + } + } + if len(peers) <= 0 { + delete(db.peerHistory.Peers, ip) + } + } + } + db.peerLock.Unlock() + case <-quit: + return + } + } +} diff --git a/p2p/metrics.go b/p2p/metrics.go index 2ab133698c..836ab560fe 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -21,12 +21,13 @@ package p2p import ( "net" + "github.com/ethereum/go-ethereum/event" + "github.com/ethereum/go-ethereum/metrics" + "sync" + "time" "github.com/ethereum/go-ethereum/log" - "github.com/ethereum/go-ethereum/metrics" - "github.com/mohae/deepcopy" - "sync" "sync/atomic" - "time" + "fmt" ) const ( @@ -35,9 +36,6 @@ const ( MetricsOutboundTraffic = "p2p/OutboundTraffic" MetricsOutboundConnects = "p2p/OutboundConnects" - MetricsRegistryIngressPrefix = MetricsInboundTraffic + "/" - MetricsRegistryEgressPrefix = MetricsOutboundTraffic + "/" - MeteredPeerLimit = 16384 ) @@ -47,101 +45,79 @@ var ( egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) - PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) - PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) - PeerTrafficMeters = newPeerTrafficMeters() - - nextDefaultID uint32 + NME = &networkMeterEvents{} ) -type PeerMetrics struct { - ID string - IP net.IP +type networkMeterEvents struct { + connectFeed event.Feed + handshakeFeed event.Feed + disconnectFeed event.Feed - // TODO: -* - Connected *time.Time - Handshake *time.Time - Disconnected *time.Time + readFeed event.Feed + writeFeed event.Feed - Ingress int64 - Egress int64 + scope event.SubscriptionScope - traffic func() (ingress, egress int64) + defaultID uint64 } -type peerTrafficMeters struct { - peers map[uint]*PeerMetrics - lock sync.RWMutex +type PeerConnectEvent struct { + IP net.IP + ID string + Connected time.Time } -func newPeerTrafficMeters() *peerTrafficMeters { - return &peerTrafficMeters{ - peers: make(map[uint]*PeerMetrics), - } +type PeerHandshakeEvent struct { + IP net.IP + DefaultID string + ID string + Handshake time.Time } -func (m *peerTrafficMeters) register(id uint, ip net.IP, traffic func() (ingress, egress int64)) error { - now := time.Now() - peer := &PeerMetrics{ - IP: ip, - Connected: &now, - traffic: traffic, - } - m.lock.Lock() - m.peers[id] = peer - m.lock.Unlock() - - return nil +type PeerDisconnectEvent struct { + IP net.IP + ID string + Disconnected time.Time } -func (m *peerTrafficMeters) handshakeDone(id uint, peerID string, traffic func() (ingress, egress int64)) { - now := time.Now() - m.lock.Lock() - if peer, ok := m.peers[id]; ok { - peer.Handshake = &now - peer.ID = peerID - peer.traffic = traffic - } - m.lock.Unlock() +type PeerReadEvent struct { + IP net.IP + ID string + Ingress int } -func (m *peerTrafficMeters) close(id uint) { - now := time.Now() - m.lock.Lock() - m.peers[id].Disconnected = &now - m.lock.Unlock() +type PeerWriteEvent struct { + IP net.IP + ID string + Egress int } -func (m *peerTrafficMeters) Peers() map[uint]*PeerMetrics { - peers := make(map[uint]*PeerMetrics) - m.lock.Lock() - for id, peer := range m.peers { - peer.Ingress, peer.Egress = peer.traffic() - peers[id] = deepcopy.Copy(peer).(*PeerMetrics) - if peer.Disconnected != nil { - PeerIngressRegistry.Unregister(peer.ID) - PeerEgressRegistry.Unregister(peer.ID) - delete(m.peers, id) - } - } - m.lock.Unlock() - return peers +func SubscribePeerConnectEvent(ch chan<- PeerConnectEvent) event.Subscription { + return NME.scope.Track(NME.connectFeed.Subscribe(ch)) +} +func SubscribePeerHandshakeEvent(ch chan<- PeerHandshakeEvent) event.Subscription { + return NME.scope.Track(NME.handshakeFeed.Subscribe(ch)) +} +func SubscribePeerDisconnectEvent(ch chan<- PeerDisconnectEvent) event.Subscription { + return NME.scope.Track(NME.disconnectFeed.Subscribe(ch)) +} +func SubscribePeerReadEvent(ch chan<- PeerReadEvent) event.Subscription { + return NME.scope.Track(NME.readFeed.Subscribe(ch)) +} +func SubscribePeerWriteEvent(ch chan<- PeerWriteEvent) event.Subscription { + return NME.scope.Track(NME.writeFeed.Subscribe(ch)) } -type networkMeter struct { - ingress metrics.Meter - egress metrics.Meter +func closeNME() { + NME.scope.Close() } // 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 - id uint - meter *networkMeter - - ingressBeforeHandshake int64 - egressBeforeHandshake int64 + ip net.IP + id string lock sync.RWMutex } @@ -154,8 +130,8 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { if !metrics.Enabled { return conn } - if len(PeerTrafficMeters.peers) >= MeteredPeerLimit { - log.Warn("Metered peer limit exceeded") + if ip.IsUnspecified() { + log.Warn("peer IP is unspecified") return conn } // Otherwise bump the connection counters and wrap the connection @@ -164,15 +140,17 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { } else { egressConnectMeter.Mark(1) } - id := uint(atomic.AddUint32(&nextDefaultID, 1)) - c := &meteredConn{ + id := fmt.Sprintf("peer_%d", atomic.AddUint64(&NME.defaultID, 1)) + NME.connectFeed.Send(PeerConnectEvent{ + IP: ip, + ID: id, + Connected: time.Now(), + }) + return &meteredConn{ Conn: conn, + ip: ip, id: id, } - PeerTrafficMeters.register(id, ip, func() (ingress, egress int64) { - return atomic.LoadInt64(&c.ingressBeforeHandshake), atomic.LoadInt64(&c.egressBeforeHandshake) - }) - return c } // Read delegates a network read to the underlying connection, bumping the ingress @@ -181,12 +159,13 @@ func (c *meteredConn) Read(b []byte) (n int, err error) { n, err = c.Conn.Read(b) ingressTrafficMeter.Mark(int64(n)) c.lock.RLock() - if c.meter == nil { - atomic.AddInt64(&c.ingressBeforeHandshake, int64(n)) - } else { - c.meter.ingress.Mark(int64(n)) - } + id := c.id c.lock.RUnlock() + NME.readFeed.Send(PeerReadEvent{ + IP: c.ip, + ID: id, + Ingress: n, + }) return n, err } @@ -196,42 +175,37 @@ func (c *meteredConn) Write(b []byte) (n int, err error) { n, err = c.Conn.Write(b) egressTrafficMeter.Mark(int64(n)) c.lock.RLock() - if c.meter == nil { - atomic.AddInt64(&c.egressBeforeHandshake, int64(n)) - } else { - c.meter.egress.Mark(int64(n)) - } + id := c.id c.lock.RUnlock() + NME.writeFeed.Send(PeerWriteEvent{ + IP: c.ip, + ID: id, + Egress: n, + }) return n, err } func (c *meteredConn) Close() error { - PeerTrafficMeters.close(c.id) + c.lock.RLock() + id := c.id + c.lock.RUnlock() + NME.disconnectFeed.Send(PeerDisconnectEvent{ + IP: c.ip, + ID: id, + Disconnected: time.Now(), + }) return c.Conn.Close() } func (c *meteredConn) handshakeDone(peerID string) { - m := &networkMeter{ - ingress: metrics.NewRegisteredMeter(peerID, PeerIngressRegistry), - egress: metrics.NewRegisteredMeter(peerID, PeerEgressRegistry), - } c.lock.Lock() - m.ingress.Mark(atomic.LoadInt64(&c.ingressBeforeHandshake)) - m.egress.Mark(atomic.LoadInt64(&c.egressBeforeHandshake)) - c.meter = m + defaultID := c.id + c.id = peerID c.lock.Unlock() - - ingressMeter, oki := PeerIngressRegistry.Get(peerID).(metrics.Meter) - egressMeter, oke := PeerEgressRegistry.Get(peerID).(metrics.Meter) - traffic := func() (ingress, egress int64) { - return 0, 0 - } - if oki && oke { - traffic = func() (ingress, egress int64) { - return ingressMeter.Count(), egressMeter.Count() - } - } else { - log.Warn("Failed to get traffic meter", "peerID", peerID) - } - PeerTrafficMeters.handshakeDone(c.id, peerID, traffic) + NME.handshakeFeed.Send(PeerHandshakeEvent{ + IP: c.ip, + DefaultID: defaultID, + ID: peerID, + Handshake: time.Now(), + }) } diff --git a/p2p/server.go b/p2p/server.go index c5087710f8..11fd9b6dfc 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -388,6 +388,7 @@ func (srv *Server) Stop() { close(srv.quit) srv.lock.Unlock() srv.loopWG.Wait() + closeNME() } // sharedUDPConn implements a shared connection. Write sends messages to the underlying connection while read returns