diff --git a/dashboard/assets/components/Dashboard.jsx b/dashboard/assets/components/Dashboard.jsx index 4b949478e3..c11c40c195 100644 --- a/dashboard/assets/components/Dashboard.jsx +++ b/dashboard/assets/components/Dashboard.jsx @@ -24,7 +24,8 @@ import Header from './Header'; import Body from './Body'; import {MENU} from '../common'; import type {Content} from '../types/content'; -import {inserter as logInserter} from './Logs'; +import {inserter as logInserter, SAME} from './Logs'; +import {inserter as peerInserter} from './Network'; // deepUpdate updates an object corresponding to the given update data, which has // the shape of the same structure as the original object. updater also has the same @@ -89,8 +90,7 @@ const defaultContent: () => Content = () => ({ chain: {}, txpool: {}, network: { - peers: [], - changed: [], + peers: {}, }, system: { activeMemory: [], @@ -106,8 +106,8 @@ const defaultContent: () => Content = () => ({ chunks: [], endTop: false, endBottom: true, - topChanged: 0, - bottomChanged: 0, + topChanged: SAME, + bottomChanged: SAME, }, }); @@ -123,8 +123,7 @@ const updaters = { chain: null, txpool: null, network: { - peers: appender(200), - changed: appender(200), + peers: peerInserter, }, system: { activeMemory: appender(200), @@ -204,9 +203,6 @@ class Dashboard extends Component { return; } this.update(msg); - if (msg.network) { - console.log(msg.network); - } }; server.onclose = () => { this.setState({server: null}); diff --git a/dashboard/assets/components/Logs.jsx b/dashboard/assets/components/Logs.jsx index 57d6328d35..ce0e4396bc 100644 --- a/dashboard/assets/components/Logs.jsx +++ b/dashboard/assets/components/Logs.jsx @@ -104,9 +104,9 @@ const createChunk = (records: Array) => { // ADDED, SAME and REMOVED are used to track the change of the log chunk array. // The scroll position is set using these values. -const ADDED = 1; -const SAME = 0; -const REMOVED = -1; +export const ADDED = 1; +export const SAME = 0; +export const REMOVED = -1; // 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. diff --git a/dashboard/assets/components/Main.jsx b/dashboard/assets/components/Main.jsx index 3102274a73..77c569fa69 100644 --- a/dashboard/assets/components/Main.jsx +++ b/dashboard/assets/components/Main.jsx @@ -21,6 +21,7 @@ import React, {Component} from 'react'; import withStyles from 'material-ui/styles/withStyles'; import {MENU} from '../common'; +import Network from './Network'; import Logs from './Logs'; import Footer from './Footer'; import type {Content} from '../types/content'; @@ -98,7 +99,7 @@ class Main extends Component { children =
Work in progress.
; break; case MENU.get('network').id: - children =
Work in progress.
; + children = ; break; case MENU.get('system').id: children =
Work in progress.
; diff --git a/dashboard/assets/components/Network.jsx b/dashboard/assets/components/Network.jsx index 41a73c7c03..19773a8a34 100644 --- a/dashboard/assets/components/Network.jsx +++ b/dashboard/assets/components/Network.jsx @@ -18,41 +18,68 @@ import React, {Component} from 'react'; -import type {Content, Network as NetworkType} from '../types/content'; - +import Table, {TableBody, TableHeader, TableHeaderColumn, 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: NetworkType, prev: LogsType) => prev; +export const inserter = (update: {[number]: Peer}, prev: {[number]: Peer}) => { + Object.keys(update).forEach((k) => { + if (!prev[k]) { + prev[k] = update[k]; + return; + } + const u: Peer = update[k]; + const p: Peer = prev[k]; + if (u.id) { + p.id = u.id; + } + if (u.ip) { + p.ip = u.ip; + } + if (u.lifecycle) { + if (u.lifecycle.handshake) { + p.lifecycle.handshake = u.lifecycle.handshake; + } + if (u.lifecycle.disconnected) { + p.lifecycle.disconnected = u.lifecycle.disconnected; + } + } + p.ingress = [...p.ingress, ...u.ingress].slice(-200); + p.egress = [...p.egress, ...u.egress].slice(-200); + prev[k] = p; + }); + return prev; +}; // styles contains the constant styles of the component. const styles = {}; export type Props = { container: Object, - content: Content, + content: NetworkType, shouldUpdate: Object, }; -type State = { - peers: List, -}; - // Network renders the network page. class Network extends Component { - constructor(props: Props) { - super(props); - this.content = React.createRef(); - this.state = { - peers: [], - }; - } - render() { return ( -
- Alma -
+ + + {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)} + + ))} + +
); } } diff --git a/dashboard/assets/types/content.jsx b/dashboard/assets/types/content.jsx index ff82c3572d..691db4987a 100644 --- a/dashboard/assets/types/content.jsx +++ b/dashboard/assets/types/content.jsx @@ -51,7 +51,29 @@ export type TxPool = { }; export type Network = { - /* TODO (kurkomisi) */ + peers: {[number]: Peer}, +}; + +export type Peer = { + id: string, + ip: string, + location: PeerLocation, + lifecycle: PeerLifecycle, + ingress: ChartEntries, + egress: ChartEntries, +}; + +export type PeerLocation = { + country: string, + city: string, + latitude: number, + longitude: number, +}; + +export type PeerLifecycle = { + connected: Date, + handshake: Date, + disconnected: Date, }; export type System = { diff --git a/dashboard/dashboard.go b/dashboard/dashboard.go index c7b6c9433b..f92a4c6d64 100644 --- a/dashboard/dashboard.go +++ b/dashboard/dashboard.go @@ -42,7 +42,6 @@ import ( "github.com/ethereum/go-ethereum/rpc" "github.com/mohae/deepcopy" "golang.org/x/net/websocket" - "encoding/json" ) const ( @@ -54,6 +53,10 @@ const ( systemCPUSampleLimit = 200 // Maximum number of system cpu data samples 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 ) var nextID uint32 // Next connection id @@ -67,7 +70,9 @@ type Dashboard struct { history *Message lock sync.RWMutex // Lock protecting the dashboard's internals - logdir string + geodb *GeoDB + disconnectedPeerIDs chan uint + logdir string quit chan chan error // Channel used for graceful exit wg sync.WaitGroup @@ -106,8 +111,12 @@ 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), + }, }, - logdir: logdir, + disconnectedPeerIDs: make(chan uint, storedDisconnectedPeerLimit), + logdir: logdir, } } @@ -133,6 +142,13 @@ 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) go db.collectData() go db.streamLogs() @@ -154,6 +170,8 @@ 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 { @@ -262,7 +280,7 @@ func (db *Dashboard) apiHandler(conn *websocket.Conn) { // meterCollector returns a function, which retrieves a specific meter. func meterCollector(name string) func() int64 { - if meter := metrics.DefaultRegistry.Get(name); meter != nil { + if meter := metrics.Get(name); meter != nil { m := meter.(metrics.Meter) return func() int64 { return m.Count() @@ -282,8 +300,8 @@ func (db *Dashboard) collectData() { var ( mem runtime.MemStats - collectNetworkIngress = meterCollector("p2p/InboundTraffic") - collectNetworkEgress = meterCollector("p2p/OutboundTraffic") + collectNetworkIngress = meterCollector(p2p.MetricsInboundTraffic) + collectNetworkEgress = meterCollector(p2p.MetricsOutboundTraffic) collectDiskRead = meterCollector("eth/db/chaindata/disk/read") collectDiskWrite = meterCollector("eth/db/chaindata/disk/write") @@ -374,21 +392,83 @@ func (db *Dashboard) collectData() { sys.DiskWrite = append(sys.DiskWrite[1:], diskWrite) db.lock.Unlock() - peerIDs := p2p.TrafficMeterCollector.GetIDs() - nm := &NetworkMessage{ - Changed: p2p.TrafficMeterCollector.GetAndClearChanged(), + peers := p2p.PeerTrafficMeters.Peers() + network := new(NetworkMessage) + if len(peers) > 0 { + network.Peers = make(map[uint]*Peer) } - for _, id := range peerIDs { - nm.Peers = append(nm.Peers, &Peer{ - ID: id[len(id)-6:], - Ingress: p2p.PeerIngressRegistry.Get(id).(metrics.Meter).Count(), - Egress: p2p.PeerEgressRegistry.Get(id).(metrics.Meter).Count(), - }) - fmt.Println(metrics.DefaultRegistry.Get(fmt.Sprintf("%s/%s", p2p.IngressPrefix, id)).(metrics.Meter).Count()) - fmt.Println(metrics.DefaultRegistry.Get(fmt.Sprintf("%s/%s", p2p.EgressPrefix, id)).(metrics.Meter).Count()) + 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 } - s, _ := json.Marshal(nm) - fmt.Println(string(s)) + db.lock.Unlock() + //s, _ := json.MarshalIndent(network, "", " ") + //fmt.Println(string(s)) db.sendToAll(&Message{ System: &SystemMessage{ @@ -401,7 +481,7 @@ func (db *Dashboard) collectData() { DiskRead: ChartEntries{diskRead}, DiskWrite: ChartEntries{diskWrite}, }, - Network: nm, + Network: network, }) } } diff --git a/dashboard/geoip.go b/dashboard/geoip.go new file mode 100644 index 0000000000..cc702dac90 --- /dev/null +++ b/dashboard/geoip.go @@ -0,0 +1,79 @@ +// Copyright 2018 The go-ethereum Authors +// This file is part of the go-ethereum library. +// +// The go-ethereum library is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// The go-ethereum library is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with the go-ethereum library. If not, see . + +package dashboard + +import ( + "github.com/apilayer/freegeoip" + "time" + "net" +) +// Package geoip contains utility methods for converting IPs to geographical data. + +// GeoDBInfo contains all the geographical information we could extract based on an IP +// address. +type GeoDBInfo struct { + Country struct { + Names struct { + English string `maxminddb:"en" json:"en,omitempty"` + } `maxminddb:"names" json:"names,omitempty"` + } `maxminddb:"country" json:"country,omitempty"` + City struct { + Names struct { + English string `maxminddb:"en" json:"en,omitempty"` + } `maxminddb:"names" json:"names,omitempty"` + } `maxminddb:"city" json:"city,omitempty"` + Location struct { + Latitude float64 `maxminddb:"latitude" json:"latitude,omitempty"` + Longitude float64 `maxminddb:"longitude" json:"longitude,omitempty"` + } `maxminddb:"location" json:"location,omitempty"` +} + +// GeoDB represents a geoip database that can be queried for IP to geographical +// information conversions. +type GeoDB struct { + geodb *freegeoip.DB +} + +// Open creats a new geoip database with an up-to-date database from the internet. +func OpenGeoDB() (*GeoDB, error) { + // Initiate a geoip database to cross reference locations + db, err := freegeoip.OpenURL(freegeoip.MaxMindDB, 24*time.Hour, time.Hour) + if err != nil { + return nil, err + } + // Wait until the database is updated to the latest data + select { + case <-db.NotifyOpen(): + case err := <-db.NotifyError(): + return nil, err + } + // Assemble and return our custom wrapper + return &GeoDB{geodb: db}, nil +} + +// Close terminates the database background updater. +func (db *GeoDB) Close() error { + db.geodb.Close() + return nil +} + +// Lookup converts an IP address to a geographical location. +func (db *GeoDB) Lookup(ip net.IP) *GeoDBInfo { + result := new(GeoDBInfo) + db.geodb.Lookup(ip, result) + return result +} diff --git a/dashboard/message.go b/dashboard/message.go index 99050d854d..c1c6a32445 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -18,6 +18,7 @@ package dashboard import ( "encoding/json" + "net" "time" ) @@ -34,8 +35,8 @@ type Message struct { type ChartEntries []*ChartEntry type ChartEntry struct { - Time time.Time `json:"time,omitempty"` - Value float64 `json:"value,omitempty"` + Time time.Time `json:"time"` + Value float64 `json:"value"` } type GeneralMessage struct { @@ -56,14 +57,29 @@ type TxPoolMessage struct { } type NetworkMessage struct { - Peers []*Peer `json:"peers,omitempty"` - Changed []string `json:"changed,omitempty"` + Peers map[uint]*Peer `json:"peers,omitempty"` } type Peer struct { - ID string `json:"id,omitempty"` - Ingress int64 `json:"ingress,omitempty"` - Egress int64 `json:"egress,omitempty"` + 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"` +} + +type PeerLocation struct { + Country string `json:"country,omitempty"` + City string `json:"city,omitempty"` + Latitude float64 `json:"latitude,omitempty"` + 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 { diff --git a/p2p/dial.go b/p2p/dial.go index d8feceb9f3..e2aba5587b 100644 --- a/p2p/dial.go +++ b/p2p/dial.go @@ -350,7 +350,7 @@ func (t *dialTask) dial(srv *Server, dest *discover.Node) error { if err != nil { return &dialError{err} } - mfd := newMeteredConn(fd, false) + mfd := newMeteredConn(fd, false, dest.IP) return srv.SetupConn(mfd, t.flags, dest) } diff --git a/p2p/metrics.go b/p2p/metrics.go index dfabad0873..2ab133698c 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -21,181 +21,172 @@ package p2p import ( "net" - "fmt" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/metrics" - "github.com/syndtr/goleveldb/leveldb/errors" + "github.com/mohae/deepcopy" "sync" "sync/atomic" + "time" ) const ( - IngressPrefix = "p2p/InboundTraffic" - EgressPrefix = "p2p/OutboundTraffic" + MetricsInboundTraffic = "p2p/InboundTraffic" + MetricsInboundConnects = "p2p/InboundConnects" + MetricsOutboundTraffic = "p2p/OutboundTraffic" + MetricsOutboundConnects = "p2p/OutboundConnects" + + MetricsRegistryIngressPrefix = MetricsInboundTraffic + "/" + MetricsRegistryEgressPrefix = MetricsOutboundTraffic + "/" + + MeteredPeerLimit = 16384 ) var ( - ingressConnectMeter = metrics.NewRegisteredMeter("p2p/InboundConnects", nil) - egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil) + ingressConnectMeter = metrics.NewRegisteredMeter(MetricsInboundConnects, nil) + ingressTrafficMeter = metrics.NewRegisteredMeter(MetricsInboundTraffic, nil) + egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) + egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) - PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, IngressPrefix+"/") - PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, EgressPrefix+"/") - TrafficMeterCollector = newTrafficMeterCollector() + PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) + PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) + PeerTrafficMeters = newPeerTrafficMeters() nextDefaultID uint32 ) -type trafficMeter struct { +type PeerMetrics struct { + ID string + IP net.IP + + // TODO: -* + Connected *time.Time + Handshake *time.Time + Disconnected *time.Time + + Ingress int64 + Egress int64 + + traffic func() (ingress, egress int64) +} + +type peerTrafficMeters struct { + peers map[uint]*PeerMetrics + lock sync.RWMutex +} + +func newPeerTrafficMeters() *peerTrafficMeters { + return &peerTrafficMeters{ + peers: make(map[uint]*PeerMetrics), + } +} + +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 +} + +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() +} + +func (m *peerTrafficMeters) close(id uint) { + now := time.Now() + m.lock.Lock() + m.peers[id].Disconnected = &now + m.lock.Unlock() +} + +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 +} + +type networkMeter struct { ingress metrics.Meter egress metrics.Meter } -type trafficMeterCollector struct { - common *trafficMeter - peers map[string]*trafficMeter - changed []string - - lock sync.RWMutex - cLock sync.Mutex -} - -func newTrafficMeterCollector() *trafficMeterCollector { - return &trafficMeterCollector{ - common: &trafficMeter{ - ingress: metrics.NewRegisteredMeter(IngressPrefix, nil), - egress: metrics.NewRegisteredMeter(EgressPrefix, nil), - }, - peers: make(map[string]*trafficMeter), - changed: make([]string, 0, 128), - } -} - -func (tmc *trafficMeterCollector) register(id string, tm *trafficMeter) error { - if tm == nil { - peer := &trafficMeter{ - ingress: metrics.NewRegisteredMeter(id, PeerIngressRegistry), - egress: metrics.NewRegisteredMeter(id, PeerEgressRegistry), - } - tmc.lock.Lock() - tmc.peers[id] = peer - tmc.lock.Unlock() - return nil - } - if tm.ingress == nil || tm.egress == nil { - return errors.New("Meter is not set correctly") - } - if err := PeerIngressRegistry.Register(id, tm.ingress); err != nil { - return err - } - if err := PeerEgressRegistry.Register(id, tm.egress); err != nil { - PeerIngressRegistry.Unregister(id) - return err - } - tmc.lock.Lock() - tmc.peers[id] = tm - tmc.lock.Unlock() - return nil -} - -func (tmc *trafficMeterCollector) unregister(old, new string) { - PeerIngressRegistry.Unregister(old) - PeerEgressRegistry.Unregister(old) - - tmc.lock.Lock() - delete(tmc.peers, old) - tmc.lock.Unlock() - - tmc.cLock.Lock() - tmc.changed = append(tmc.changed, old, new) - tmc.cLock.Unlock() -} - -func (tmc *trafficMeterCollector) changeID(old, new string) error { - tmc.lock.RLock() - peer, ok := tmc.peers[old] - tmc.lock.RUnlock() - if !ok { - return errors.New(fmt.Sprintf("No meter with id %s", old)) - } - if err := tmc.register(new, peer); err != nil { - return err - } - tmc.unregister(old, new) - return nil -} - -func (tmc *trafficMeterCollector) markIngress(id string, n int64) { - tmc.common.ingress.Mark(n) - tmc.lock.RLock() - peer, ok := tmc.peers[id] - tmc.lock.RUnlock() - if ok { - peer.ingress.Mark(n) - } -} - -func (tmc *trafficMeterCollector) markEgress(id string, n int64) { - tmc.common.egress.Mark(n) - tmc.lock.RLock() - peer, ok := tmc.peers[id] - tmc.lock.RUnlock() - if ok { - peer.egress.Mark(n) - } -} - -func (tmc *trafficMeterCollector) GetIDs() []string { - tmc.lock.RLock() - ids := make([]string, 0, len(tmc.peers)) - for id := range tmc.peers { - ids = append(ids, id) - } - tmc.lock.RUnlock() - return ids -} - -func (tmc *trafficMeterCollector) GetAndClearChanged() []string { - tmc.cLock.Lock() - defer tmc.cLock.Unlock() - - changed := make([]string, len(tmc.changed)) - copy(changed, tmc.changed) - tmc.changed = tmc.changed[:0] - - return changed -} - // 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 string + id uint + meter *networkMeter + + ingressBeforeHandshake int64 + egressBeforeHandshake int64 + + lock sync.RWMutex } // 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. -func newMeteredConn(conn net.Conn, ingress bool) net.Conn { +func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { // Short circuit if metrics are disabled if !metrics.Enabled { return conn } + if len(PeerTrafficMeters.peers) >= MeteredPeerLimit { + log.Warn("Metered peer limit exceeded") + return conn + } // Otherwise bump the connection counters and wrap the connection if ingress { ingressConnectMeter.Mark(1) } else { egressConnectMeter.Mark(1) } - id := fmt.Sprintf("unidentified_%d", atomic.AddUint32(&nextDefaultID, 1)) - TrafficMeterCollector.register(id, nil) - return &meteredConn{Conn: conn, id: id} + id := uint(atomic.AddUint32(&nextDefaultID, 1)) + c := &meteredConn{ + Conn: conn, + 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 // traffic meter along the way. func (c *meteredConn) Read(b []byte) (n int, err error) { n, err = c.Conn.Read(b) - TrafficMeterCollector.markIngress(c.id, int64(n)) + ingressTrafficMeter.Mark(int64(n)) + c.lock.RLock() + if c.meter == nil { + atomic.AddInt64(&c.ingressBeforeHandshake, int64(n)) + } else { + c.meter.ingress.Mark(int64(n)) + } + c.lock.RUnlock() return n, err } @@ -203,19 +194,44 @@ func (c *meteredConn) Read(b []byte) (n int, err error) { // egress traffic meter along the way. func (c *meteredConn) Write(b []byte) (n int, err error) { n, err = c.Conn.Write(b) - TrafficMeterCollector.markEgress(c.id, int64(n)) + egressTrafficMeter.Mark(int64(n)) + c.lock.RLock() + if c.meter == nil { + atomic.AddInt64(&c.egressBeforeHandshake, int64(n)) + } else { + c.meter.egress.Mark(int64(n)) + } + c.lock.RUnlock() return n, err } func (c *meteredConn) Close() error { - TrafficMeterCollector.unregister(c.id, "") + PeerTrafficMeters.close(c.id) return c.Conn.Close() } -func (c *meteredConn) setPeerID(id string) { - if err := TrafficMeterCollector.changeID(c.id, id); err != nil { - log.Warn("Failed to set peer id", "id", fmt.Sprintf("%s...%s", id[:6], id[len(id)-6:]), "err", err) - return +func (c *meteredConn) handshakeDone(peerID string) { + m := &networkMeter{ + ingress: metrics.NewRegisteredMeter(peerID, PeerIngressRegistry), + egress: metrics.NewRegisteredMeter(peerID, PeerEgressRegistry), } - c.id = id + c.lock.Lock() + m.ingress.Mark(atomic.LoadInt64(&c.ingressBeforeHandshake)) + m.egress.Mark(atomic.LoadInt64(&c.egressBeforeHandshake)) + c.meter = m + 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) } diff --git a/p2p/rlpx.go b/p2p/rlpx.go index 393676b815..46b6668695 100644 --- a/p2p/rlpx.go +++ b/p2p/rlpx.go @@ -586,9 +586,6 @@ func newRLPXFrameRW(conn io.ReadWriter, s secrets) *rlpxFrameRW { // we use an all-zeroes IV for AES because the key used // for encryption is ephemeral. iv := make([]byte, encc.BlockSize()) - if c, ok := conn.(*meteredConn); ok { - c.setPeerID(s.RemoteID.String()) - } return &rlpxFrameRW{ conn: conn, enc: cipher.NewCTR(encc, iv), diff --git a/p2p/server.go b/p2p/server.go index 8f3a511f31..c5087710f8 100644 --- a/p2p/server.go +++ b/p2p/server.go @@ -41,7 +41,7 @@ const ( // Connectivity defaults. maxActiveDialTasks = 16 - defaultMaxPendingPeers = 50 + DefaultMaxPendingPeers = 50 defaultDialRatio = 3 // Maximum time allowed for reading a complete message. @@ -799,7 +799,7 @@ func (srv *Server) listenLoop() { defer srv.loopWG.Done() srv.log.Info("RLPx listener up", "self", srv.makeSelf(srv.listener, srv.ntab)) - tokens := defaultMaxPendingPeers + tokens := DefaultMaxPendingPeers if srv.MaxPendingPeers > 0 { tokens = srv.MaxPendingPeers } @@ -838,7 +838,11 @@ func (srv *Server) listenLoop() { } } - fd = newMeteredConn(fd, true) + var ip net.IP + if tcp, ok := fd.RemoteAddr().(*net.TCPAddr); ok { + ip = tcp.IP + } + fd = newMeteredConn(fd, true, ip) srv.log.Trace("Accepted connection", "addr", fd.RemoteAddr()) go func() { srv.SetupConn(fd, inboundConn, nil) @@ -878,6 +882,9 @@ func (srv *Server) setupConn(c *conn, flags connFlag, dialDest *discover.Node) e srv.log.Trace("Failed RLPx handshake", "addr", c.fd.RemoteAddr(), "conn", c.flags, "err", err) return err } + if conn, ok := c.fd.(*meteredConn); ok { + conn.handshakeDone(c.id.String()) + } clog := srv.log.New("id", c.id, "addr", c.fd.RemoteAddr(), "conn", c.flags) // For dialed connections, check that the remote public key matches. if dialDest != nil && c.id != dialDest.ID {