From 272444d11bf026aaf4731571e55edc783d4772e4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kurk=C3=B3=20Mih=C3=A1ly?= Date: Mon, 23 Jul 2018 16:50:57 +0300 Subject: [PATCH] dashboard, p2p: store traffic meters in maps --- dashboard/dashboard.go | 6 +- p2p/metrics.go | 165 ++++++++++++++++++++++++++++++++--------- p2p/rlpx.go | 2 +- 3 files changed, 136 insertions(+), 37 deletions(-) diff --git a/dashboard/dashboard.go b/dashboard/dashboard.go index a965110688..4c2b39be43 100644 --- a/dashboard/dashboard.go +++ b/dashboard/dashboard.go @@ -385,10 +385,10 @@ func (db *Dashboard) collectData() { DiskWrite: ChartEntries{diskWrite}, }, }) - - for k, v := range p2p.IngressTrafficMeters { - fmt.Println(k[:6], v.Count()) + for _, id := range p2p.TrafficMeter.GetIDs() { + fmt.Println(id) } + fmt.Println() } } } diff --git a/p2p/metrics.go b/p2p/metrics.go index 9ca5b86b64..00211513b6 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -22,25 +22,136 @@ import ( "net" "fmt" + "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/metrics" + "github.com/syndtr/goleveldb/leveldb/errors" + "sync" + "sync/atomic" ) var ( - ingressConnectMeter = metrics.NewRegisteredMeter("p2p/InboundConnects", nil) - ingressTrafficMeter = metrics.NewRegisteredMeter("p2p/InboundTraffic", nil) - egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil) - egressTrafficMeter = metrics.NewRegisteredMeter("p2p/OutboundTraffic", nil) - IngressTrafficMeters = make(map[string]metrics.Meter) - EgressTrafficMeters = make(map[string]metrics.Meter) + ingressConnectMeter = metrics.NewRegisteredMeter("p2p/InboundConnects", nil) + egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil) + TrafficMeter = newTrafficMeter() + nextDefaultID uint32 ) +const ( + IngressPrefix = "p2p/InboundTraffic" + EgressPrefix = "p2p/OutboundTraffic" +) + +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 +} + +// 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), + } +} + +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() + return nil + } + if ingressMeter == nil || egressMeter == nil { + return errors.New("Meter is not set") + } + if err := metrics.Register(fmt.Sprintf("%s/%s", IngressPrefix, id), ingressMeter); 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)) + return err + } + tm.lock.Lock() + tm.peerIngress[id] = ingressMeter + tm.peerEgress[id] = egressMeter + tm.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 (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 { + return errors.New(fmt.Sprintf("No meter with id %s", old)) + } + if err := tm.register(new, irm, erm); err != nil { + return err + } + tm.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() + if ok { + rm.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() + if ok { + rm.Mark(n) + } +} + +func (tm *trafficMeter) GetIDs() []string { + ids := make([]string, 0, len(tm.peerIngress)) + tm.lock.RLock() + for id := range tm.peerIngress { + ids = append(ids, id) + } + tm.lock.RUnlock() + return ids +} + // 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 - unmarkedIngress uint - unmarkedEgress uint + net.Conn // Network connection to wrap with metering + id string } // newMeteredConn creates a new metered connection, also bumping the ingress or @@ -57,19 +168,16 @@ func newMeteredConn(conn net.Conn, ingress bool) net.Conn { } else { egressConnectMeter.Mark(1) } - return &meteredConn{Conn: conn} + id := fmt.Sprintf("unidentified_%d", atomic.AddUint32(&nextDefaultID, 1)) + TrafficMeter.register(id, nil, nil) + return &meteredConn{Conn: conn, id: id} } // 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) - ingressTrafficMeter.Mark(int64(n)) - if rm, ok := IngressTrafficMeters[c.id]; ok { - rm.Mark(int64(n)) - return n, err - } - c.unmarkedIngress += uint(n) + TrafficMeter.markIngress(c.id, int64(n)) return n, err } @@ -77,29 +185,20 @@ 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) - egressTrafficMeter.Mark(int64(n)) - if rm, ok := EgressTrafficMeters[c.id]; ok { - rm.Mark(int64(n)) - return n, err - } - c.unmarkedEgress += uint(n) + TrafficMeter.markEgress(c.id, int64(n)) return n, err } -func (c *meteredConn) meterIndividually(id string) { +func (c *meteredConn) setPeerID(id string) { + if err := TrafficMeter.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 - IngressTrafficMeters[id] = metrics.NewRegisteredMeter(fmt.Sprintf("p2p/InboundTraffic/%s", id), nil) - EgressTrafficMeters[id] = metrics.NewRegisteredMeter(fmt.Sprintf("p2p/OutboundTraffic/%s", id), nil) - IngressTrafficMeters[id].Mark(int64(c.unmarkedIngress)) - EgressTrafficMeters[id].Mark(int64(c.unmarkedEgress)) - c.unmarkedIngress, c.unmarkedEgress = 0, 0 } func (c *meteredConn) Close() error { - delete(IngressTrafficMeters, c.id) - delete(EgressTrafficMeters, c.id) - metrics.Unregister(fmt.Sprintf("p2p/InboundTraffic/%s", c.id)) - metrics.Unregister(fmt.Sprintf("p2p/OutboundTraffic/%s", c.id)) + TrafficMeter.unregister(c.id, "") fmt.Println("Close", c.id) return c.Conn.Close() } diff --git a/p2p/rlpx.go b/p2p/rlpx.go index 05ab8bac9c..393676b815 100644 --- a/p2p/rlpx.go +++ b/p2p/rlpx.go @@ -587,7 +587,7 @@ func newRLPXFrameRW(conn io.ReadWriter, s secrets) *rlpxFrameRW { // for encryption is ephemeral. iv := make([]byte, encc.BlockSize()) if c, ok := conn.(*meteredConn); ok { - c.meterIndividually(s.RemoteID.String()) + c.setPeerID(s.RemoteID.String()) } return &rlpxFrameRW{ conn: conn,