From 811520d44215a537ac6aba60b85df88382b753bf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kurk=C3=B3=20Mih=C3=A1ly?= Date: Thu, 20 Sep 2018 20:12:54 +0300 Subject: [PATCH] p2p, dashboard: experiment 4 --- dashboard/message.go | 117 ++++++++++----------- dashboard/peers.go | 245 +++++++++++++++++++++++++++++++++++++++++-- p2p/metrics.go | 12 +-- 3 files changed, 302 insertions(+), 72 deletions(-) diff --git a/dashboard/message.go b/dashboard/message.go index af3cf2002b..d2affa20d4 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -17,7 +17,6 @@ package dashboard import ( - "container/list" "encoding/json" "time" ) @@ -63,57 +62,57 @@ type NetworkMessage struct { // getOrInitBundle returns the peer bundle belonging to the given IP, or // initializes the bundle if it doesn't exist. -func (m *NetworkMessage) getOrInitBundle(ip string) *PeerBundle { - if _, ok := m.PeerBundles[ip]; !ok { - m.PeerBundles[ip] = &PeerBundle{ - Peers: make(PeerMap), - FailedPeers: make(PeerMap), - } - } - return m.PeerBundles[ip] -} +//func (m *NetworkMessage) getOrInitBundle(ip string) *PeerBundle { +// if _, ok := m.PeerBundles[ip]; !ok { +// m.PeerBundles[ip] = &PeerBundle{ +// Peers: make(PeerMap), +// FailedPeers: make(PeerMap), +// } +// } +// return m.PeerBundles[ip] +//} +// +//// getOrInitPeer returns the peer belonging to the given IP and node id, or +//// initializes the peer if it doesn't exist. +//func (m *NetworkMessage) getOrInitPeer(ip, id string) *Peer { +// return m.getOrInitBundle(ip).getOrInitPeer(id) +//} -// getOrInitPeer returns the peer belonging to the given IP and node id, or -// initializes the peer if it doesn't exist. -func (m *NetworkMessage) getOrInitPeer(ip, id string) *Peer { - return m.getOrInitBundle(ip).getOrInitPeer(id) -} +//type PeerMap map[string]*Peer -type PeerMap map[string]*Peer +//func (pm PeerMap) getOrInit(id string) *Peer { +// if _, ok := pm[id]; !ok { +// pm[id] = new(Peer) +// } +// return pm[id] +//} +// +//func (pm PeerMap) remove(id string) { +// delete(pm, id) +//} -func (pm PeerMap) getOrInit(id string) *Peer { - if _, ok := pm[id]; !ok { - pm[id] = new(Peer) - } - return pm[id] -} - -func (pm PeerMap) remove(id string) { - delete(pm, id) -} - -// PeerBundle contains information about the peers pertaining to an IP address. -type PeerBundle struct { - Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP - Peers PeerMap `json:"peers,omitempty"` // the peers' node id is used as key - FailedPeers PeerMap `json:"failedPeers,omitempty"` -} - -func (b *PeerBundle) getOrInitPeer(id string) *Peer { - return b.Peers.getOrInit(id) -} - -func (b *PeerBundle) removePeer(id string) { - b.Peers.remove(id) -} - -func (b *PeerBundle) getOrInitFailedPeer(id string) *Peer { - return b.FailedPeers.getOrInit(id) -} - -func (b * PeerBundle) removeFailedPeer(id string) { - b.FailedPeers.remove(id) -} +//// PeerBundle contains information about the peers pertaining to an IP address. +//type PeerBundle struct { +// Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP +// Peers PeerMap `json:"peers,omitempty"` // the peers' node id is used as key +// FailedPeers PeerMap `json:"failedPeers,omitempty"` +//} +// +//func (b *PeerBundle) getOrInitPeer(id string) *Peer { +// return b.Peers.getOrInit(id) +//} +// +//func (b *PeerBundle) removePeer(id string) { +// b.Peers.remove(id) +//} +// +//func (b *PeerBundle) getOrInitFailedPeer(id string) *Peer { +// return b.FailedPeers.getOrInit(id) +//} +// +//func (b * PeerBundle) removeFailedPeer(id string) { +// b.FailedPeers.remove(id) +//} // GeoLocation contains geographical information. type GeoLocation struct { @@ -124,16 +123,16 @@ type GeoLocation struct { } // Peer contains lifecycle timestamps and traffic information of a given peer. -type Peer struct { - Connected []time.Time `json:"connected,omitempty"` - Disconnected []time.Time `json:"disconnected,omitempty"` - - Ingress ChartEntries `json:"ingress,omitempty"` - Egress ChartEntries `json:"egress,omitempty"` - - element *list.Element - ip, id string -} +//type Peer struct { +// Connected []time.Time `json:"connected,omitempty"` +// Disconnected []time.Time `json:"disconnected,omitempty"` +// +// Ingress ChartEntries `json:"ingress,omitempty"` +// Egress ChartEntries `json:"egress,omitempty"` +// +// element *list.Element +// ip, id string +//} // SystemMessage contains the metered system data samples. type SystemMessage struct { diff --git a/dashboard/peers.go b/dashboard/peers.go index 704b7dbd5a..6f9c3def9b 100644 --- a/dashboard/peers.go +++ b/dashboard/peers.go @@ -17,7 +17,11 @@ package dashboard import ( + "container/list" + "encoding/json" "fmt" + "github.com/ethereum/go-ethereum/metrics" + "strings" "time" "github.com/ethereum/go-ethereum/log" @@ -190,6 +194,214 @@ type removedPeer struct { ip, id string } + + + + + + + + +const ( + knownPeerLimit = p2p.MeteredPeerLimit + unknownPeerLimit = p2p.MeteredPeerLimit +) + +type PeerContainer struct { + Bundles map[string]*PeerBundle `json:"peerBundles,omitempty"` + + activeSeparator list.Element + knownPeers *list.List + unknownPeers *list.List + + refresh time.Duration +} + +func NewPeerContainer(refresh time.Duration) *PeerContainer { + return &PeerContainer{ + Bundles: make(map[string]*PeerBundle), + knownPeers: list.New(), + unknownPeers: list.New(), + refresh: refresh, + } +} + +func (pc *PeerContainer) getOrInit(ip string) *PeerBundle { + if _, ok := pc.Bundles[ip]; !ok { + pc.Bundles[ip] = &PeerBundle{ + KnownPeers: make(map[string]*KnownPeer), + parent: pc, + ip: ip, + } + } + return pc.Bundles[ip] +} + +func (pc *PeerContainer) updateKnown(ip, id string, cycle *PeerCycle) { + pc.getOrInit(ip).updateKnown(id, cycle) +} + +func (pc *PeerContainer) updateUnknown(ip string, peer *UnknownPeer) { + pc.getOrInit(ip).updateUnknown(peer) +} + +func (pc *PeerContainer) Update(event *peerEvent) { + switch event.t { + case peerConnected: + connected := time.Now().Add(-event.Elapsed) + pc.updateKnown(event.IP.String(), event.ID, &PeerCycle{ + Connected: &connected, + }) + case peerDisconnected: + now := time.Now() + pc.updateKnown(event.IP.String(), event.ID, &PeerCycle{ + Disconnected: &now, + Ingress: ChartEntries{ + &ChartEntry{ + Time: now, + Value: float64(event.Ingress), + }, + }, + Egress: ChartEntries{ + &ChartEntry{ + Time: now, + Value: float64(event.Egress), + }, + }, + }) + case peerHandshakeFailed: + now := time.Now() + pc.updateUnknown(event.IP.String(), &UnknownPeer{ + Connected: now.Add(-event.Elapsed), + Disconnected: now, + }) + case peerIngress: + fmt.Println("Ingress:", event) + pc.updateKnown(event.ip, event.id, &PeerCycle{ + Ingress: ChartEntries{ + &ChartEntry{ + Time: time.Now(), + Value: float64(event.traffic), + }, + }, + }) + case peerEgress: + fmt.Println("Egress:", event) + pc.updateKnown(event.ip, event.id, &PeerCycle{ + Egress: ChartEntries{ + &ChartEntry{ + Time: time.Now(), + Value: float64(event.traffic), + }, + }, + }) + default: + log.Error("Unknown peer event type", "type", event.Type) + } +} + +type PeerBundle struct { + Location *GeoLocation `json:"location,omitempty"` + KnownPeers map[string]*KnownPeer `json:"knownPeers,omitempty"` + UnknownPeers []*UnknownPeer `json:"unknownPeers,omitempty"` + + parent *PeerContainer + ip string +} + +func (b *PeerBundle) getOrInit(id string) *KnownPeer { + if _, ok := b.KnownPeers[id]; !ok { + b.KnownPeers[id] = &KnownPeer{ + parent: b, + id: id, + } + } + return b.KnownPeers[id] +} + +func (b *PeerBundle) updateKnown(id string, cycle *PeerCycle) { + b.getOrInit(id).update(cycle) +} + +func (b *PeerBundle) updateUnknown(peer *UnknownPeer) { + b.UnknownPeers = append(b.UnknownPeers, peer) +} + +type KnownPeer struct { + Cycles []*PeerCycle `json:"cycles,omitempty"` + + parent *PeerBundle + id string +} + +func (peer *KnownPeer) update(cycle *PeerCycle) { + if cycle.Connected != nil { + if peer.Cycles == nil { + peer.Cycles = append(peer.Cycles, &PeerCycle{ + Ingress: emptyChartEntries(time.Now(), 3, peer.parent.parent.refresh), + Egress: emptyChartEntries(time.Now(), 3, peer.parent.parent.refresh), + }) + } + peer.Cycles = append(peer.Cycles, cycle) + return + } + if cycle.Disconnected != nil { + if len(peer.Cycles) < 1 { + log.Error("Peer disconnect event appeared without connect") + return + } + last := peer.Cycles[len(peer.Cycles)-1] + last.Disconnected = cycle.Disconnected + last.Ingress = append(last.Ingress, cycle.Ingress...) + last.Egress = append(last.Egress, cycle.Egress...) + return + } +} + +type PeerCycle struct { + Connected *time.Time `json:"connected,omitempty"` + Disconnected *time.Time `json:"disconnected,omitempty"` + + Ingress ChartEntries `json:"ingress,omitempty"` + Egress ChartEntries `json:"egress,omitempty"` +} + +type UnknownPeer struct { + Connected time.Time `json:"connected,omitempty"` + Disconnected time.Time `json:"disconnected,omitempty"` +} + +type peerContainerUpdateType int + +const ( + peerConnected = peerContainerUpdateType(p2p.PeerConnected) + peerDisconnected = peerContainerUpdateType(p2p.PeerDisconnected) + peerHandshakeFailed = peerContainerUpdateType(p2p.PeerHandshakeFailed) + peerIngress = peerConnected + peerDisconnected + peerHandshakeFailed + iota + peerEgress +) + +type peerEvent struct { + *p2p.MeteredPeerEvent + ip string + id string + t peerContainerUpdateType + traffic uint64 +} + + + + + + + + + + + + + + // collectPeerData gathers data about the peers and sends it to the clients. func (db *Dashboard) collectPeerData() { defer db.wg.Done() @@ -210,16 +422,35 @@ func (db *Dashboard) collectPeerData() { ticker := time.NewTicker(db.config.Refresh) defer ticker.Stop() - db.peerLock.RLock() - //historyMaintainer := NewPeerMaintainer(p2p.MeteredPeerLimit) - //historyHandshakeFailedMaintainer := NewPeerMaintainer(p2p.MeteredPeerLimit) - //diffMaintainer := NewPeerMaintainer(p2p.MeteredPeerLimit) - //diffHandshakeFailedMaintainer := NewPeerMaintainer(p2p.MeteredPeerLimit) + pc := NewPeerContainer(db.config.Refresh) + + trafficUpdater := func(prefix string) (func (*map[string]int64) func(name string, i interface{})) { + return func (entryMap *map[string]int64) func(name string, i interface{}) { + return func(name string, i interface{}) { + if m, ok := i.(metrics.Meter); ok { + (*entryMap)[strings.TrimPrefix(name, prefix)] = m.Count() + } + } + } + } + updateIngress := trafficUpdater(p2p.MetricsInboundTraffic+"/") + updateEgress := trafficUpdater(p2p.MetricsOutboundTraffic+"/") + for { select { case event := <-peerCh: - fmt.Println(event) - //diffMaintainer.Update(event.IP.String(), event.ID) + pc.Update(&peerEvent{ + MeteredPeerEvent: &event, + t: peerContainerUpdateType(event.Type), + }) + case <-ticker.C: + ingress := make(map[string]int64) + egress := make(map[string]int64) + p2p.PeerIngressRegistry.Each(updateIngress(&ingress)) + p2p.PeerEgressRegistry.Each(updateEgress(&egress)) + //fmt.Println(ingress) + s, _ := json.MarshalIndent(pc, "", " ") + fmt.Println(string(s)) case err := <-subPeer.Err(): log.Warn("Peer subscription error", "err", err) return diff --git a/p2p/metrics.go b/p2p/metrics.go index 86401379df..8c0aa518ad 100644 --- a/p2p/metrics.go +++ b/p2p/metrics.go @@ -183,12 +183,6 @@ func (c *meteredConn) handshakeDone(nodeID discover.NodeID) { func (c *meteredConn) Close() error { err := c.Conn.Close() c.lock.RLock() - if !c.metered { - c.lock.RUnlock() - return err - } - // Decrement the metered peer count - atomic.AddInt32(&meteredPeerCount, -1) if c.id == "" { // If the peer disconnects before the handshake c.lock.RUnlock() @@ -199,6 +193,12 @@ func (c *meteredConn) Close() error { }) return err } + if !c.metered { + c.lock.RUnlock() + return err + } + // Decrement the metered peer count + atomic.AddInt32(&meteredPeerCount, -1) id, ingress, egress := c.id, uint64(c.ingressMeter.Count()), uint64(c.egressMeter.Count()) c.lock.RUnlock()