diff --git a/dashboard/assets/components/Dashboard.jsx b/dashboard/assets/components/Dashboard.jsx index 63c2186ada..4b949478e3 100644 --- a/dashboard/assets/components/Dashboard.jsx +++ b/dashboard/assets/components/Dashboard.jsx @@ -88,7 +88,10 @@ const defaultContent: () => Content = () => ({ home: {}, chain: {}, txpool: {}, - network: {}, + network: { + peers: [], + changed: [], + }, system: { activeMemory: [], virtualMemory: [], @@ -119,7 +122,10 @@ const updaters = { home: null, chain: null, txpool: null, - network: null, + network: { + peers: appender(200), + changed: appender(200), + }, system: { activeMemory: appender(200), virtualMemory: appender(200), @@ -198,6 +204,9 @@ 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/Main.jsx b/dashboard/assets/components/Main.jsx index 0018aaf75f..3102274a73 100644 --- a/dashboard/assets/components/Main.jsx +++ b/dashboard/assets/components/Main.jsx @@ -89,9 +89,17 @@ class Main extends Component { let children = null; switch (active) { case MENU.get('home').id: + children =
Work in progress.
; + break; case MENU.get('chain').id: + children =
Work in progress.
; + break; case MENU.get('txpool').id: + children =
Work in progress.
; + break; case MENU.get('network').id: + children =
Work in progress.
; + break; case MENU.get('system').id: children =
Work in progress.
; break; diff --git a/dashboard/assets/components/Network.jsx b/dashboard/assets/components/Network.jsx new file mode 100644 index 0000000000..41a73c7c03 --- /dev/null +++ b/dashboard/assets/components/Network.jsx @@ -0,0 +1,60 @@ +// @flow + +// 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 . + +import React, {Component} from 'react'; + +import type {Content, Network as NetworkType} 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; + +// styles contains the constant styles of the component. +const styles = {}; + +export type Props = { + container: Object, + content: Content, + 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 +
+ ); + } +} + +export default Network; diff --git a/dashboard/dashboard.go b/dashboard/dashboard.go index 4c2b39be43..c7b6c9433b 100644 --- a/dashboard/dashboard.go +++ b/dashboard/dashboard.go @@ -42,6 +42,7 @@ import ( "github.com/ethereum/go-ethereum/rpc" "github.com/mohae/deepcopy" "golang.org/x/net/websocket" + "encoding/json" ) const ( @@ -261,8 +262,8 @@ func (db *Dashboard) apiHandler(conn *websocket.Conn) { // meterCollector returns a function, which retrieves a specific meter. func meterCollector(name string) func() int64 { - if metric := metrics.DefaultRegistry.Get(name); metric != nil { - m := metric.(metrics.Meter) + if meter := metrics.DefaultRegistry.Get(name); meter != nil { + m := meter.(metrics.Meter) return func() int64 { return m.Count() } @@ -373,6 +374,22 @@ func (db *Dashboard) collectData() { sys.DiskWrite = append(sys.DiskWrite[1:], diskWrite) db.lock.Unlock() + peerIDs := p2p.TrafficMeterCollector.GetIDs() + nm := &NetworkMessage{ + Changed: p2p.TrafficMeterCollector.GetAndClearChanged(), + } + 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()) + } + s, _ := json.Marshal(nm) + fmt.Println(string(s)) + db.sendToAll(&Message{ System: &SystemMessage{ ActiveMemory: ChartEntries{activeMemory}, @@ -384,11 +401,8 @@ func (db *Dashboard) collectData() { DiskRead: ChartEntries{diskRead}, DiskWrite: ChartEntries{diskWrite}, }, + Network: nm, }) - for _, id := range p2p.TrafficMeter.GetIDs() { - fmt.Println(id) - } - fmt.Println() } } } diff --git a/dashboard/message.go b/dashboard/message.go index 46856b9e6e..99050d854d 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -56,7 +56,14 @@ type TxPoolMessage struct { } type NetworkMessage struct { - /* TODO (kurkomisi) */ + Peers []*Peer `json:"peers,omitempty"` + Changed []string `json:"changed,omitempty"` +} + +type Peer struct { + ID string `json:"id,omitempty"` + Ingress int64 `json:"ingress,omitempty"` + Egress int64 `json:"egress,omitempty"` } type SystemMessage struct { diff --git a/metrics/registry.go b/metrics/registry.go index cc34c9dfd2..baf01c554f 100644 --- a/metrics/registry.go +++ b/metrics/registry.go @@ -311,7 +311,7 @@ func (r *PrefixedRegistry) UnregisterAll() { r.underlying.UnregisterAll() } -var DefaultRegistry Registry = NewRegistry() +var DefaultRegistry = NewRegistry() // Call the given function for each registered metric. func Each(f func(string, interface{})) { diff --git a/p2p/metrics.go b/p2p/metrics.go index 00211513b6..dfabad0873 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -29,124 +29,142 @@ import ( "sync/atomic" ) -var ( - ingressConnectMeter = metrics.NewRegisteredMeter("p2p/InboundConnects", nil) - egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil) - TrafficMeter = newTrafficMeter() - nextDefaultID uint32 -) - const ( IngressPrefix = "p2p/InboundTraffic" EgressPrefix = "p2p/OutboundTraffic" ) +var ( + ingressConnectMeter = metrics.NewRegisteredMeter("p2p/InboundConnects", nil) + egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil) + + PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, IngressPrefix+"/") + PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, EgressPrefix+"/") + TrafficMeterCollector = newTrafficMeterCollector() + + nextDefaultID uint32 +) + type trafficMeter struct { - peerIngress map[string]metrics.Meter - peerEgress map[string]metrics.Meter - commonIngress metrics.Meter - commonEgress metrics.Meter - changed map[string]string - lock sync.RWMutex - cLock sync.RWMutex + ingress metrics.Meter + egress metrics.Meter } -// Create a new registry. -func newTrafficMeter() *trafficMeter { - return &trafficMeter{ - peerIngress: make(map[string]metrics.Meter), - peerEgress: make(map[string]metrics.Meter), - commonIngress: metrics.NewRegisteredMeter("p2p/InboundTraffic", nil), - commonEgress: metrics.NewRegisteredMeter("p2p/OutboundTraffic", nil), - changed: make(map[string]string), +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 (tm *trafficMeter) register(id string, ingressMeter, egressMeter metrics.Meter) error { - if ingressMeter == nil && egressMeter == nil { - tm.lock.Lock() - tm.peerIngress[id] = metrics.NewRegisteredMeter(fmt.Sprintf("%s/%s", IngressPrefix, id), nil) - tm.peerEgress[id] = metrics.NewRegisteredMeter(fmt.Sprintf("%s/%s", EgressPrefix, id), nil) - tm.lock.Unlock() +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 ingressMeter == nil || egressMeter == nil { - return errors.New("Meter is not set") + if tm.ingress == nil || tm.egress == nil { + return errors.New("Meter is not set correctly") } - if err := metrics.Register(fmt.Sprintf("%s/%s", IngressPrefix, id), ingressMeter); err != nil { + if err := PeerIngressRegistry.Register(id, tm.ingress); err != nil { return err } - if err := metrics.Register(fmt.Sprintf("%s/%s", EgressPrefix, id), egressMeter); err != nil { - metrics.Unregister(fmt.Sprintf("%s/%s", IngressPrefix, id)) + if err := PeerEgressRegistry.Register(id, tm.egress); err != nil { + PeerIngressRegistry.Unregister(id) return err } - tm.lock.Lock() - tm.peerIngress[id] = ingressMeter - tm.peerEgress[id] = egressMeter - tm.lock.Unlock() + tmc.lock.Lock() + tmc.peers[id] = tm + tmc.lock.Unlock() return nil } -func (tm *trafficMeter) delete(id string) { - tm.lock.Lock() - delete(tm.peerIngress, id) - delete(tm.peerEgress, id) - tm.lock.Unlock() -} -func (tm *trafficMeter) unregister(old, new string) { - metrics.Unregister(fmt.Sprintf("%s/%s", IngressPrefix, old)) - metrics.Unregister(fmt.Sprintf("%s/%s", EgressPrefix, old)) - tm.delete(old) - tm.cLock.Lock() - tm.changed[old] = new - tm.cLock.Unlock() +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 (tm *trafficMeter) changeID(old, new string) error { - tm.lock.RLock() - irm, oki := tm.peerIngress[old] - erm, oke := tm.peerEgress[old] - tm.lock.RUnlock() - if !oki || !oke { +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 := tm.register(new, irm, erm); err != nil { + if err := tmc.register(new, peer); err != nil { return err } - tm.unregister(old, new) + tmc.unregister(old, new) return nil } -func (tm *trafficMeter) markIngress(id string, n int64) { - tm.commonIngress.Mark(n) - tm.lock.RLock() - rm, ok := tm.peerIngress[id] - tm.lock.RUnlock() +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 { - rm.Mark(n) + peer.ingress.Mark(n) } } -func (tm *trafficMeter) markEgress(id string, n int64) { - tm.commonEgress.Mark(n) - tm.lock.RLock() - rm, ok := tm.peerEgress[id] - tm.lock.RUnlock() +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 { - rm.Mark(n) + peer.egress.Mark(n) } } -func (tm *trafficMeter) GetIDs() []string { - ids := make([]string, 0, len(tm.peerIngress)) - tm.lock.RLock() - for id := range tm.peerIngress { +func (tmc *trafficMeterCollector) GetIDs() []string { + tmc.lock.RLock() + ids := make([]string, 0, len(tmc.peers)) + for id := range tmc.peers { ids = append(ids, id) } - tm.lock.RUnlock() + 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 { @@ -169,7 +187,7 @@ func newMeteredConn(conn net.Conn, ingress bool) net.Conn { egressConnectMeter.Mark(1) } id := fmt.Sprintf("unidentified_%d", atomic.AddUint32(&nextDefaultID, 1)) - TrafficMeter.register(id, nil, nil) + TrafficMeterCollector.register(id, nil) return &meteredConn{Conn: conn, id: id} } @@ -177,7 +195,7 @@ func newMeteredConn(conn net.Conn, ingress bool) net.Conn { // traffic meter along the way. func (c *meteredConn) Read(b []byte) (n int, err error) { n, err = c.Conn.Read(b) - TrafficMeter.markIngress(c.id, int64(n)) + TrafficMeterCollector.markIngress(c.id, int64(n)) return n, err } @@ -185,20 +203,19 @@ 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) - TrafficMeter.markEgress(c.id, int64(n)) + TrafficMeterCollector.markEgress(c.id, int64(n)) return n, err } +func (c *meteredConn) Close() error { + TrafficMeterCollector.unregister(c.id, "") + return c.Conn.Close() +} + func (c *meteredConn) setPeerID(id string) { - if err := TrafficMeter.changeID(c.id, id); err != nil { + 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 } c.id = id } - -func (c *meteredConn) Close() error { - TrafficMeter.unregister(c.id, "") - fmt.Println("Close", c.id) - return c.Conn.Close() -}