dashboard, metrics, p2p: log peers on client side

This commit is contained in:
Kurkó Mihály 2018-07-24 15:13:56 +03:00
parent 272444d11b
commit 0125c61e78
7 changed files with 210 additions and 95 deletions

View file

@ -88,7 +88,10 @@ const defaultContent: () => Content = () => ({
home: {}, home: {},
chain: {}, chain: {},
txpool: {}, txpool: {},
network: {}, network: {
peers: [],
changed: [],
},
system: { system: {
activeMemory: [], activeMemory: [],
virtualMemory: [], virtualMemory: [],
@ -119,7 +122,10 @@ const updaters = {
home: null, home: null,
chain: null, chain: null,
txpool: null, txpool: null,
network: null, network: {
peers: appender(200),
changed: appender(200),
},
system: { system: {
activeMemory: appender(200), activeMemory: appender(200),
virtualMemory: appender(200), virtualMemory: appender(200),
@ -198,6 +204,9 @@ class Dashboard extends Component<Props, State> {
return; return;
} }
this.update(msg); this.update(msg);
if (msg.network) {
console.log(msg.network);
}
}; };
server.onclose = () => { server.onclose = () => {
this.setState({server: null}); this.setState({server: null});

View file

@ -89,9 +89,17 @@ class Main extends Component<Props> {
let children = null; let children = null;
switch (active) { switch (active) {
case MENU.get('home').id: case MENU.get('home').id:
children = <div>Work in progress.</div>;
break;
case MENU.get('chain').id: case MENU.get('chain').id:
children = <div>Work in progress.</div>;
break;
case MENU.get('txpool').id: case MENU.get('txpool').id:
children = <div>Work in progress.</div>;
break;
case MENU.get('network').id: case MENU.get('network').id:
children = <div>Work in progress.</div>;
break;
case MENU.get('system').id: case MENU.get('system').id:
children = <div>Work in progress.</div>; children = <div>Work in progress.</div>;
break; break;

View file

@ -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 <http://www.gnu.org/licenses/>.
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<string>,
};
// Network renders the network page.
class Network extends Component<Props, State> {
constructor(props: Props) {
super(props);
this.content = React.createRef();
this.state = {
peers: [],
};
}
render() {
return (
<div>
Alma
</div>
);
}
}
export default Network;

View file

@ -42,6 +42,7 @@ import (
"github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/rpc"
"github.com/mohae/deepcopy" "github.com/mohae/deepcopy"
"golang.org/x/net/websocket" "golang.org/x/net/websocket"
"encoding/json"
) )
const ( const (
@ -261,8 +262,8 @@ func (db *Dashboard) apiHandler(conn *websocket.Conn) {
// meterCollector returns a function, which retrieves a specific meter. // meterCollector returns a function, which retrieves a specific meter.
func meterCollector(name string) func() int64 { func meterCollector(name string) func() int64 {
if metric := metrics.DefaultRegistry.Get(name); metric != nil { if meter := metrics.DefaultRegistry.Get(name); meter != nil {
m := metric.(metrics.Meter) m := meter.(metrics.Meter)
return func() int64 { return func() int64 {
return m.Count() return m.Count()
} }
@ -373,6 +374,22 @@ func (db *Dashboard) collectData() {
sys.DiskWrite = append(sys.DiskWrite[1:], diskWrite) sys.DiskWrite = append(sys.DiskWrite[1:], diskWrite)
db.lock.Unlock() 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{ db.sendToAll(&Message{
System: &SystemMessage{ System: &SystemMessage{
ActiveMemory: ChartEntries{activeMemory}, ActiveMemory: ChartEntries{activeMemory},
@ -384,11 +401,8 @@ func (db *Dashboard) collectData() {
DiskRead: ChartEntries{diskRead}, DiskRead: ChartEntries{diskRead},
DiskWrite: ChartEntries{diskWrite}, DiskWrite: ChartEntries{diskWrite},
}, },
Network: nm,
}) })
for _, id := range p2p.TrafficMeter.GetIDs() {
fmt.Println(id)
}
fmt.Println()
} }
} }
} }

View file

@ -56,7 +56,14 @@ type TxPoolMessage struct {
} }
type NetworkMessage 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 { type SystemMessage struct {

View file

@ -311,7 +311,7 @@ func (r *PrefixedRegistry) UnregisterAll() {
r.underlying.UnregisterAll() r.underlying.UnregisterAll()
} }
var DefaultRegistry Registry = NewRegistry() var DefaultRegistry = NewRegistry()
// Call the given function for each registered metric. // Call the given function for each registered metric.
func Each(f func(string, interface{})) { func Each(f func(string, interface{})) {

View file

@ -29,124 +29,142 @@ import (
"sync/atomic" "sync/atomic"
) )
var (
ingressConnectMeter = metrics.NewRegisteredMeter("p2p/InboundConnects", nil)
egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil)
TrafficMeter = newTrafficMeter()
nextDefaultID uint32
)
const ( const (
IngressPrefix = "p2p/InboundTraffic" IngressPrefix = "p2p/InboundTraffic"
EgressPrefix = "p2p/OutboundTraffic" 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 { type trafficMeter struct {
peerIngress map[string]metrics.Meter ingress metrics.Meter
peerEgress map[string]metrics.Meter egress metrics.Meter
commonIngress metrics.Meter }
commonEgress metrics.Meter
changed map[string]string type trafficMeterCollector struct {
common *trafficMeter
peers map[string]*trafficMeter
changed []string
lock sync.RWMutex lock sync.RWMutex
cLock sync.RWMutex cLock sync.Mutex
} }
// Create a new registry. func newTrafficMeterCollector() *trafficMeterCollector {
func newTrafficMeter() *trafficMeter { return &trafficMeterCollector{
return &trafficMeter{ common: &trafficMeter{
peerIngress: make(map[string]metrics.Meter), ingress: metrics.NewRegisteredMeter(IngressPrefix, nil),
peerEgress: make(map[string]metrics.Meter), egress: metrics.NewRegisteredMeter(EgressPrefix, nil),
commonIngress: metrics.NewRegisteredMeter("p2p/InboundTraffic", nil), },
commonEgress: metrics.NewRegisteredMeter("p2p/OutboundTraffic", nil), peers: make(map[string]*trafficMeter),
changed: make(map[string]string), changed: make([]string, 0, 128),
} }
} }
func (tm *trafficMeter) register(id string, ingressMeter, egressMeter metrics.Meter) error { func (tmc *trafficMeterCollector) register(id string, tm *trafficMeter) error {
if ingressMeter == nil && egressMeter == nil { if tm == nil {
tm.lock.Lock() peer := &trafficMeter{
tm.peerIngress[id] = metrics.NewRegisteredMeter(fmt.Sprintf("%s/%s", IngressPrefix, id), nil) ingress: metrics.NewRegisteredMeter(id, PeerIngressRegistry),
tm.peerEgress[id] = metrics.NewRegisteredMeter(fmt.Sprintf("%s/%s", EgressPrefix, id), nil) egress: metrics.NewRegisteredMeter(id, PeerEgressRegistry),
tm.lock.Unlock() }
tmc.lock.Lock()
tmc.peers[id] = peer
tmc.lock.Unlock()
return nil return nil
} }
if ingressMeter == nil || egressMeter == nil { if tm.ingress == nil || tm.egress == nil {
return errors.New("Meter is not set") 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 return err
} }
if err := metrics.Register(fmt.Sprintf("%s/%s", EgressPrefix, id), egressMeter); err != nil { if err := PeerEgressRegistry.Register(id, tm.egress); err != nil {
metrics.Unregister(fmt.Sprintf("%s/%s", IngressPrefix, id)) PeerIngressRegistry.Unregister(id)
return err return err
} }
tm.lock.Lock() tmc.lock.Lock()
tm.peerIngress[id] = ingressMeter tmc.peers[id] = tm
tm.peerEgress[id] = egressMeter tmc.lock.Unlock()
tm.lock.Unlock()
return nil return nil
} }
func (tm *trafficMeter) delete(id string) { func (tmc *trafficMeterCollector) unregister(old, new string) {
tm.lock.Lock() PeerIngressRegistry.Unregister(old)
delete(tm.peerIngress, id) PeerEgressRegistry.Unregister(old)
delete(tm.peerEgress, id)
tm.lock.Unlock() tmc.lock.Lock()
} delete(tmc.peers, old)
func (tm *trafficMeter) unregister(old, new string) { tmc.lock.Unlock()
metrics.Unregister(fmt.Sprintf("%s/%s", IngressPrefix, old))
metrics.Unregister(fmt.Sprintf("%s/%s", EgressPrefix, old)) tmc.cLock.Lock()
tm.delete(old) tmc.changed = append(tmc.changed, old, new)
tm.cLock.Lock() tmc.cLock.Unlock()
tm.changed[old] = new
tm.cLock.Unlock()
} }
func (tm *trafficMeter) changeID(old, new string) error { func (tmc *trafficMeterCollector) changeID(old, new string) error {
tm.lock.RLock() tmc.lock.RLock()
irm, oki := tm.peerIngress[old] peer, ok := tmc.peers[old]
erm, oke := tm.peerEgress[old] tmc.lock.RUnlock()
tm.lock.RUnlock() if !ok {
if !oki || !oke {
return errors.New(fmt.Sprintf("No meter with id %s", old)) 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 return err
} }
tm.unregister(old, new) tmc.unregister(old, new)
return nil return nil
} }
func (tm *trafficMeter) markIngress(id string, n int64) { func (tmc *trafficMeterCollector) markIngress(id string, n int64) {
tm.commonIngress.Mark(n) tmc.common.ingress.Mark(n)
tm.lock.RLock() tmc.lock.RLock()
rm, ok := tm.peerIngress[id] peer, ok := tmc.peers[id]
tm.lock.RUnlock() tmc.lock.RUnlock()
if ok { if ok {
rm.Mark(n) peer.ingress.Mark(n)
} }
} }
func (tm *trafficMeter) markEgress(id string, n int64) { func (tmc *trafficMeterCollector) markEgress(id string, n int64) {
tm.commonEgress.Mark(n) tmc.common.egress.Mark(n)
tm.lock.RLock() tmc.lock.RLock()
rm, ok := tm.peerEgress[id] peer, ok := tmc.peers[id]
tm.lock.RUnlock() tmc.lock.RUnlock()
if ok { if ok {
rm.Mark(n) peer.egress.Mark(n)
} }
} }
func (tm *trafficMeter) GetIDs() []string { func (tmc *trafficMeterCollector) GetIDs() []string {
ids := make([]string, 0, len(tm.peerIngress)) tmc.lock.RLock()
tm.lock.RLock() ids := make([]string, 0, len(tmc.peers))
for id := range tm.peerIngress { for id := range tmc.peers {
ids = append(ids, id) ids = append(ids, id)
} }
tm.lock.RUnlock() tmc.lock.RUnlock()
return ids 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 // meteredConn is a wrapper around a net.Conn that meters both the
// inbound and outbound network traffic. // inbound and outbound network traffic.
type meteredConn struct { type meteredConn struct {
@ -169,7 +187,7 @@ func newMeteredConn(conn net.Conn, ingress bool) net.Conn {
egressConnectMeter.Mark(1) egressConnectMeter.Mark(1)
} }
id := fmt.Sprintf("unidentified_%d", atomic.AddUint32(&nextDefaultID, 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} return &meteredConn{Conn: conn, id: id}
} }
@ -177,7 +195,7 @@ func newMeteredConn(conn net.Conn, ingress bool) net.Conn {
// traffic meter along the way. // traffic meter along the way.
func (c *meteredConn) Read(b []byte) (n int, err error) { func (c *meteredConn) Read(b []byte) (n int, err error) {
n, err = c.Conn.Read(b) n, err = c.Conn.Read(b)
TrafficMeter.markIngress(c.id, int64(n)) TrafficMeterCollector.markIngress(c.id, int64(n))
return n, err return n, err
} }
@ -185,20 +203,19 @@ func (c *meteredConn) Read(b []byte) (n int, err error) {
// egress traffic meter along the way. // egress traffic meter along the way.
func (c *meteredConn) Write(b []byte) (n int, err error) { func (c *meteredConn) Write(b []byte) (n int, err error) {
n, err = c.Conn.Write(b) n, err = c.Conn.Write(b)
TrafficMeter.markEgress(c.id, int64(n)) TrafficMeterCollector.markEgress(c.id, int64(n))
return n, err return n, err
} }
func (c *meteredConn) Close() error {
TrafficMeterCollector.unregister(c.id, "")
return c.Conn.Close()
}
func (c *meteredConn) setPeerID(id string) { 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) log.Warn("Failed to set peer id", "id", fmt.Sprintf("%s...%s", id[:6], id[len(id)-6:]), "err", err)
return return
} }
c.id = id c.id = id
} }
func (c *meteredConn) Close() error {
TrafficMeter.unregister(c.id, "")
fmt.Println("Close", c.id)
return c.Conn.Close()
}