From 58df98a6ab7d70d385c63761167a3d1243ee763c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kurk=C3=B3=20Mih=C3=A1ly?= Date: Tue, 9 Oct 2018 15:57:40 +0300 Subject: [PATCH] dashboard: experiment 6 --- dashboard/assets/components/Dashboard.jsx | 9 +- dashboard/assets/components/Network.jsx | 10 +- dashboard/assets/types/content.jsx | 32 +- dashboard/dashboard.go | 2 +- dashboard/geoip.go | 1 + dashboard/message.go | 205 +++++++- dashboard/peers.go | 582 +++++++--------------- 7 files changed, 412 insertions(+), 429 deletions(-) diff --git a/dashboard/assets/components/Dashboard.jsx b/dashboard/assets/components/Dashboard.jsx index 3dd8d5a11a..3d9fb78575 100644 --- a/dashboard/assets/components/Dashboard.jsx +++ b/dashboard/assets/components/Dashboard.jsx @@ -90,7 +90,12 @@ const defaultContent: () => Content = () => ({ chain: {}, txpool: {}, network: { - peerBundles: {}, + peers: { + bundles: {}, + removedKnownIP: [], + removedKnownID: [], + removedUnknownIP: [], + }, }, system: { activeMemory: [], @@ -123,7 +128,7 @@ const updaters = { chain: null, txpool: null, network: { - peerBundles: peerInserter, + peers: peerInserter, }, system: { activeMemory: appender(200), diff --git a/dashboard/assets/components/Network.jsx b/dashboard/assets/components/Network.jsx index b54633ea01..988a3cbd3c 100644 --- a/dashboard/assets/components/Network.jsx +++ b/dashboard/assets/components/Network.jsx @@ -23,11 +23,13 @@ import TableHead from '@material-ui/core/TableHead'; import TableBody from '@material-ui/core/TableBody'; import TableRow from '@material-ui/core/TableRow'; import TableCell from '@material-ui/core/TableCell'; -import type {Network as NetworkType, PeerBundle, Peer} from '../types/content'; +import type {Network as NetworkType, Peers} from '../types/content'; // inserter is a state updater function for the main component, which handles the peers. -export const inserter = (update: {[string]: PeerBundle}, prev: {[string]: PeerBundle}) => { - Object.keys(update).forEach((ip) => { +export const inserter = (update: Peers, prev: Peers) => { + console.log(update); + return prev; + Object.keys(update.bundles).forEach((ip) => { if (!update[ip]) { return; } @@ -118,7 +120,7 @@ class Network extends Component { - {Object.entries(this.props.content.peerBundles).map(([ip, bundle]) => { console.log(ip, bundle); return ( + {Object.entries(this.props.content.peers).map(([ip, bundle]) => { console.log(ip, bundle); return ( {ip} diff --git a/dashboard/assets/types/content.jsx b/dashboard/assets/types/content.jsx index a8324dbca5..96265a7a6d 100644 --- a/dashboard/assets/types/content.jsx +++ b/dashboard/assets/types/content.jsx @@ -51,21 +51,37 @@ export type TxPool = { }; export type Network = { - peerBundles: {[string]: PeerBundle}, + peers: Peers, +}; + +export type Peers = { + bundles: {[string]: PeerBundle}, + removedKnownIP: Array, + removedKnownID: Array, + removedUnknownIP: Array, }; export type PeerBundle = { - location: GeoLocation, - peers: {[string]: Peer}, + location: GeoLocation, + knownPeers: {[string]: KnownPeer}, + unknownPeers: Array, }; -export type Peer = { - connected: Array, - handshake: Array, - disconnected: Array, +export type KnownPeer = { + active: boolean, + sessions: Array +}; + +export type PeerSession = { + connected: Date, + disconnected: Date, ingress: ChartEntries, egress: ChartEntries, - defaultID: string, +}; + +export type UnknownPeer = { + connected: Date, + disconnected: Date, }; export type GeoLocation = { diff --git a/dashboard/dashboard.go b/dashboard/dashboard.go index ff3c602129..db87092beb 100644 --- a/dashboard/dashboard.go +++ b/dashboard/dashboard.go @@ -42,7 +42,7 @@ import ( ) const ( - sampleLimit = 5 // Maximum number of data samples + sampleLimit = 3 // Maximum number of data samples ) // Dashboard contains the dashboard internals. diff --git a/dashboard/geoip.go b/dashboard/geoip.go index 60d6b2eab7..edc0742988 100644 --- a/dashboard/geoip.go +++ b/dashboard/geoip.go @@ -88,6 +88,7 @@ func (db *GeoDB) Lookup(ip net.IP) *GeoDBInfo { func (db *GeoDB) Location (ip string) *GeoLocation { location := db.Lookup(net.ParseIP(ip)) + //location := new(GeoDBInfo) return &GeoLocation{ Country: location.Country.Names.English, City: location.City.Names.English, diff --git a/dashboard/message.go b/dashboard/message.go index 1f4f6790db..ce96bb2b3b 100644 --- a/dashboard/message.go +++ b/dashboard/message.go @@ -18,6 +18,7 @@ package dashboard import ( "encoding/json" + "github.com/ethereum/go-ethereum/log" "time" ) @@ -27,7 +28,6 @@ type Message struct { Chain *ChainMessage `json:"chain,omitempty"` TxPool *TxPoolMessage `json:"txpool,omitempty"` Network *NetworkMessage `json:"network,omitempty"` - Peers *PeersMessage `json:"peers,omitempty"` System *SystemMessage `json:"system,omitempty"` Logs *LogsMessage `json:"logs,omitempty"` } @@ -58,7 +58,208 @@ type TxPoolMessage struct { // NetworkMessage contains information about the peers organized based on the IP address. type NetworkMessage struct { - /* TODO (kurkomisi) */ + Peers *PeersMessage `json:"peers,omitempty"` +} + +type PeersMessage struct { + Bundles map[string]*PeerBundle `json:"bundles,omitempty"` + RemovedKnownIP []string `json:"removedKnownIP,omitempty"` + RemovedKnownID []string `json:"removedKnownID,omitempty"` + RemovedUnknownIP []string `json:"removedUnknownIP,omitempty"` +} + +func NewPeersMessage() *PeersMessage { + return &PeersMessage{ + Bundles: make(map[string]*PeerBundle), + } +} + +func (m *PeersMessage) hasBundle(ip string) bool { + _, ok := m.Bundles[ip] + return ok +} + +func (m *PeersMessage) hasKnownPeer(ip, id string) bool { + if m.hasBundle(ip) { + return m.Bundles[ip].has(id) + } + return false +} + +func (m *PeersMessage) initBundle(ip string) bool { + if !m.hasBundle(ip) { + m.Bundles[ip] = &PeerBundle{ + KnownPeers: make(map[string]*KnownPeer), + } + return true + } + return false +} + +func (m *PeersMessage) initKnownPeer(ip, id string) (bundle, peer bool) { + return m.initBundle(ip), m.Bundles[ip].initKnownPeer(id) +} + +func (m *PeersMessage) getOrInitBundle(ip string) *PeerBundle { + m.initBundle(ip) + return m.Bundles[ip] +} + +func (m *PeersMessage) getOrInitKnownPeer(ip, id string) *KnownPeer { + return m.getOrInitBundle(ip).getOrInitKnownPeer(id) +} + +func (m *PeersMessage) removeKnownPeer(ip, id string) { + if b, ok := m.Bundles[ip]; ok { + b.removeKnownPeer(id) + if len(b.KnownPeers) < 1 && len(b.UnknownPeers) < 1 { + delete(m.Bundles, ip) + } + } +} + +func (m *PeersMessage) removeUnknownPeer(ip string) { + if b, ok := m.Bundles[ip]; ok { + if len(b.UnknownPeers) > 0 { + b.UnknownPeers = b.UnknownPeers[1:] + } + if len(b.KnownPeers) < 1 && len(b.UnknownPeers) < 1 { + delete(m.Bundles, ip) + } + } +} + +func (m *PeersMessage) clear() { + for _, bundle := range m.Bundles { + bundle.Location = nil + for _, peer := range bundle.KnownPeers { + peer.clear() + } + bundle.UnknownPeers = bundle.UnknownPeers[:0] + } + m.RemovedKnownIP = m.RemovedKnownIP[:0] + m.RemovedKnownID = m.RemovedKnownID[:0] + m.RemovedUnknownIP = m.RemovedUnknownIP[:0] +} + +type PeerBundle struct { + Location *GeoLocation `json:"location,omitempty"` // Geographical location based on IP + KnownPeers map[string]*KnownPeer `json:"knownPeers,omitempty"` + UnknownPeers []*UnknownPeer `json:"unknownPeers,omitempty"` +} + +func (b *PeerBundle) has(id string) bool { + _, ok := b.KnownPeers[id] + return ok +} + +func (b *PeerBundle) initKnownPeer(id string) bool { + if !b.has(id) { + b.KnownPeers[id] = new(KnownPeer) + return true + } + return false +} + +func (b *PeerBundle) getOrInitKnownPeer(id string) *KnownPeer { + b.initKnownPeer(id) + return b.KnownPeers[id] +} + +func (b *PeerBundle) removeKnownPeer(id string) bool { + if b.has(id) { + b.KnownPeers[id].clear() + delete(b.KnownPeers, id) + return true + } + return false +} + +type KnownPeer struct { + Active bool `json:"active"` + Sessions []*PeerSession `json:"sessions,omitempty"` + sampleCount int +} + +func (peer *KnownPeer) append(session *PeerSession) { + if session == nil { + return + } + ingress, egress := session.Ingress, session.Egress + // Truncate the traffic arrays if they have more samples than the limit. + if first := len(ingress) - sampleLimit; first > 0 { + ingress = ingress[first:] + } + // If the length of the ingress and the egress arrays are different, + // cut the first part of the longer one. i.e. make sure they have the + // same length. + if first := len(ingress) - len(egress); first > 0 { + ingress = ingress[first:] + } else if first < 0 { + egress = egress[-first:] + } + if len(peer.Sessions) < 1 { + // If this is the first session. + peer.Sessions = append(peer.Sessions, session) + peer.sampleCount = len(ingress) + return + } + // Cut the old samples from the beginning if the + // count with the new samples exceeds the limit. + for l := sampleLimit + len(ingress) - peer.sampleCount; l > 0; l-- { + for len(peer.Sessions) > 0 && len(peer.Sessions[0].Ingress) < 1 { + peer.Sessions = peer.Sessions[1:] + } + if len(peer.Sessions) < 1 { + // This can only happen, when the sample count is greater than the + // sample limit. Theoretically impossible. + log.Warn("Empty session array with sample count greater than 0") + return + } + first := peer.Sessions[0] + first.Ingress = first.Ingress[1:] + first.Egress = first.Egress[1:] + peer.sampleCount-- + } + peer.sampleCount += len(ingress) + if session.Connected != nil { + peer.Sessions = append(peer.Sessions, session) + return + } + last := peer.Sessions[len(peer.Sessions)-1] + last.Disconnected = session.Disconnected + last.Ingress = append(last.Ingress, ingress...) + last.Egress = append(last.Egress, egress...) +} + +func (peer *KnownPeer) upgrade(p *KnownPeer) { + peer.Active = p.Active + for _, session := range p.Sessions { + peer.append(session) + } +} + +func (peer *KnownPeer) clear() { + for _, s := range peer.Sessions { + s.Connected = nil + s.Disconnected = nil + s.Ingress = nil + s.Egress = nil + } + peer.Sessions = peer.Sessions[:0] + peer.sampleCount = 0 +} + +type PeerSession 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"` + Disconnected time.Time `json:"disconnected"` } // SystemMessage contains the metered system data samples. diff --git a/dashboard/peers.go b/dashboard/peers.go index 490bf7d6f8..9c27907d41 100644 --- a/dashboard/peers.go +++ b/dashboard/peers.go @@ -21,7 +21,9 @@ import ( "encoding/json" "fmt" "github.com/ethereum/go-ethereum/metrics" + "github.com/mohae/deepcopy" "strings" + "sync" "time" "github.com/ethereum/go-ethereum/log" @@ -30,376 +32,145 @@ import ( const ( eventBufferLimit = 128 - knownPeerLimit = p2p.MeteredPeerLimit - unknownPeerLimit = p2p.MeteredPeerLimit + knownPeerLimit = 3 //p2p.MeteredPeerLimit + unknownPeerLimit = 3 //p2p.MeteredPeerLimit ) -// PeersMessage contains information about the node's peers. This data structure -// tries to maintain the metered peer data based on the different behaviours of -// the peers. -// -// Every peer has an IP address, and the peers that manage to make the handshake -// have node IDs. There can appear more peers with the same IP, therefore the -// peer maintainer data structure is a tree consisting of a map of maps, where -// the first key groups the peers by IP, while the second one groups them by the -// node ID. The peers failing before the handshake only have IP addresses, so -// they are stored as part of the value of the outer map. A peer can connect -// multiple times, so an array is needed to store the consecutive sessions. -// -// Another criteria is to limit the number of metered peers so that they don't -// fill the memory. The metered peers are selected based on their activity: the -// peers that are inactive for the longest time are thrown first. For the selection a fifo -// list is used which is linked to the bottom of the peer tree, and when a peer -// is removed from the list, it is also removed from the tree. The active peers -// have priority over the disconnected ones, therefore this list is extended by -// a separator, which is a pointer to a list element. The separator separates the -// active peers from the inactive ones, and it is the entry for the list. If the -// peer that is to be inserted is active, it goes before the separator, otherwise -// it goes after. This way the peers that are active for the longest time are at -// the beginning of the list, and the inactive ones move to the end. When a peer -// has some activity, it is removed from and reinserted into the list. -// -// The peers that don't manage to make handshake are not inserted into the list, -// their sessions are not stored, only their connection attempts are appended -// to the array belonging to their IP. In order to keep the fifo principle, a super -// array contains the order of the attempts, and when the peer count reaches the -// limit, the earliest attempt is removed from the beginning of its array. -type PeersMessage struct { - // Bundles is the outer map using the peer's IP address as key. - Bundles map[string]*PeerBundle `json:"bundles,omitempty"` - - RemovedKnown []string `json:"removedKnown,omitempty"` - RemovedUnknown []string `json:"removedUnknown,omitempty"` - - // activeSeparator is a pointer to the last active peer element, splitting - // the list into an active and inactive part, and forming the entry for the - // peer list. - activeSeparator *list.Element - - // knownPeers contains the peers that managed to make handshake. - knownPeers *list.List - - // unknownPeers contains pointers to the peer bundles that belong to the - // IP addresses, from which the peers attempted to connect then failed. - // Its values are appended in the moment of the attempt, so in chronological - // order, which means that the oldest attempt is at the beginning of the array. - // When the first element is removed, the first element of the linked bundle's - // attempt array is also removed, ensuring that always the latest attempts are stored. - unknownPeers []*PeerBundle - - // geodb is used to look up the geographical data based on the IP. - geodb *GeoDB +type knownPeerDiff struct { + *KnownPeer + activeListElement *list.Element + listElement *list.Element // Pointer to the peer element in the list. + ip, id string } -// NewPeersMessage returns a new instance of the metered peer maintainer data structure. -func NewPeersMessage(geodb *GeoDB) *PeersMessage { - return &PeersMessage{ - Bundles: make(map[string]*PeerBundle), - knownPeers: list.New(), - unknownPeers: make([]*PeerBundle, 0, unknownPeerLimit), - geodb: geodb, +type peerDiff struct { + *PeersMessage + root *PeersMessage + rootLock *sync.RWMutex + + knownActivePeerList *list.List + knownInactivePeerList *list.List + unknownPeers []string + + geodb *GeoDB + refresh time.Duration +} + +func newPeerDiff(root *PeersMessage, rootLock *sync.RWMutex, geodb *GeoDB, refresh time.Duration) *peerDiff { + return &peerDiff{ + PeersMessage: NewPeersMessage(), + root: root, + rootLock: rootLock, + knownActivePeerList: list.New(), + knownInactivePeerList: list.New(), + unknownPeers: make([]string, 0, unknownPeerLimit), + geodb: geodb, + refresh: refresh, } } -func (m *PeersMessage) hasBundle(ip string) bool { - _, ok := m.Bundles[ip] - return ok -} - -func (m *PeersMessage) hasPeer(ip, id string) bool { - if !m.hasBundle(ip) { - return false +func (diff *peerDiff) insert(ip, id string, session *PeerSession) { + newIP, newID := diff.initKnownPeer(ip, id) + bundle := diff.Bundles[ip] + if newIP { + bundle.Location = diff.geodb.Location(ip) } - return m.Bundles[ip].has(id) -} - -// getOrInitBundle returns the bundle belonging to the given IP. -// Inserts a new bundle into the map if it doesn't already exist. -func (m *PeersMessage) getOrInitBundle(ip string) *PeerBundle { - if _, ok := m.Bundles[ip]; !ok { - m.Bundles[ip] = &PeerBundle{ - Location: m.geodb.Location(ip), - KnownPeers: make(map[string]*KnownPeer), - root: m, - ip: ip, - } + peer := &knownPeerDiff{ + KnownPeer: bundle.KnownPeers[id], + ip: ip, + id: id, } - return m.Bundles[ip] -} - -func (m *PeersMessage) getOrInitKnownPeer(ip, id string) *KnownPeer { - return m.getOrInitBundle(ip).getOrInitKnownPeer(id) -} - -// updateKnownPeer updates the last session of the peer belonging to the given IP -// and ID and returns the IP and the ID of the removed peer if there is any, -// (nil, nil) otherwise. -func (m *PeersMessage) updateKnownPeer(ip, id string, session *PeerSession) (removedKey string) { - peer := m.getOrInitKnownPeer(ip, id) - if peer.listElement != nil { - // If the peer is already part of the list, remove it first. - if peer.listElement == m.activeSeparator { - m.activeSeparator = m.activeSeparator.Prev() - } - m.knownPeers.Remove(peer.listElement) - } - if m.knownPeers.Len() >= knownPeerLimit { - // If the peer count reached the limit, remove the last element of the - // list, which is the oldest inactive one. - removed := m.knownPeers.Remove(m.knownPeers.Back()) - if p, ok := removed.(*KnownPeer); ok { - removedKey = fmt.Sprintf("%s/%s", p.bundle.ip, p.id) - p.delete() // Remove the peer from the tree. - } else { - log.Warn("Bad value used as peer", "value", removed) - } - } - if m.activeSeparator == nil { - // If there isn't any active peer in the list, push the new peer to the - // front of the list. - peer.listElement = m.knownPeers.PushFront(peer) - } else { - // If there are active peers in the list, push the new peer after them. - peer.listElement = m.knownPeers.InsertAfter(peer, m.activeSeparator) - } - peer.update(session) - if peer.Sessions[len(peer.Sessions)-1].Disconnected == nil { - // If the new peer is active, step to it with the separator. - m.activeSeparator = peer.listElement - } - return removedKey -} - -// updateUnknownPeer inserts a peer connection attempt into the peer tree. -func (m *PeersMessage) updateUnknownPeer(ip string, peer *UnknownPeer) (removedKey string) { - if len(m.unknownPeers) >= unknownPeerLimit { - // If the count of the metered unknown peers reached the limit, - // remove the oldest attempt, which is the first element of the - // array belonging to the peer bundle pointed by the first element - // of the super unknown peer array. - removed := m.unknownPeers[0] - removed.UnknownPeers = removed.UnknownPeers[1:] - m.unknownPeers = m.unknownPeers[1:] - removedKey = removed.ip - } - bundle := m.getOrInitBundle(ip) - bundle.UnknownPeers = append(bundle.UnknownPeers, peer) - m.unknownPeers = append(m.unknownPeers, bundle) - - return removedKey -} - -// clear removes the elements of the metered peer list, and removes the linked -// peers from the peer tree along the way. -func (m *PeersMessage) clear() { - for m.knownPeers.Front() != nil { - first := m.knownPeers.Remove(m.knownPeers.Front()) - if p, ok := first.(*KnownPeer); ok { - p.delete() - } else { - log.Warn("Bad value used as peer", "value", first) - } - } - m.activeSeparator = nil -} - -func (m *PeersMessage) updateTraffic(ip, id string, ingress, egress float64) { - if !m.hasPeer(ip, id) { - m.updateKnownPeer(ip, id, &PeerSession{ - Ingress: ChartEntries{&ChartEntry{ - Value: ingress, - }}, - Egress: ChartEntries{&ChartEntry{ - Value: egress, - }}, + if newID { + now := time.Now() + peer.append(&PeerSession{ + Ingress: emptyChartEntries(now, sampleLimit, diff.refresh), + Egress: emptyChartEntries(now, sampleLimit, diff.refresh), }) - return } - m.getOrInitKnownPeer(ip, id).updateTraffic( - &ChartEntry{Value: ingress}, - &ChartEntry{Value: egress}, - ) -} - -//func (m *PeersMessage) updateTraffic(ingress, egress *map[string]float64) { -// now := time.Now() -// for k, v := range *ingress { -// -// } -// for e := m.knownPeers.Front(); e != nil; e = e.Next() { -// if p, ok := e.Value.(*KnownPeer); ok { -// key := fmt.Sprintf("%s/%s", p.bundle.ip, p.id) -// p.updateTraffic( -// &ChartEntry{Time: now, Value: (*ingress)[key]}, -// &ChartEntry{Time: now, Value: (*egress)[key]}, -// ) -// } else { -// log.Warn("Bad value used as peer", "value", e.Value) -// } -// } -//} - -// append -func (m *PeersMessage) append(n *PeersMessage) (removedKnown, removedUnknown []string) { - for _, bundle := range n.Bundles { - for _, peer := range bundle.KnownPeers { - for _, session := range peer.Sessions { - removedKey := m.updateKnownPeer(bundle.ip, peer.id, session) - if removedKey != "" { - removedKnown = append(removedKnown, removedKey) - } - } - } - for _, peer := range bundle.UnknownPeers { - removedKey := m.updateUnknownPeer(bundle.ip, peer) - if removedKey != "" { - removedUnknown = append(removedUnknown, removedKey) - } - } + peer.append(session) + if peer.activeListElement != nil { + diff.knownActivePeerList.Remove(peer.activeListElement) } - return removedKnown, removedUnknown -} - -// PeerBundle contains the peers belonging to a given IP address -type PeerBundle struct { - Location *GeoLocation `json:"location,omitempty"` // Geographical location based on IP - - // KnownPeers is the inner map of the metered peer maintainer data structure - // using the node ID as key. - KnownPeers map[string]*KnownPeer `json:"knownPeers,omitempty"` - - // UnknownPeers contains the failed connection attempts of the peers - // belonging to a given IP address in chronological order. - UnknownPeers []*UnknownPeer `json:"unknownPeers,omitempty"` - - root *PeersMessage // Pointer to the outer map. - ip string // Key of the bundle in the outer map. -} - -func (b *PeerBundle) has(id string) bool { - _, ok := b.KnownPeers[id] - return ok -} - -// getOrInitKnownPeer returns the peer belonging to the given ID. -// Initializes it if it doesn't already exist. -func (b *PeerBundle) getOrInitKnownPeer(id string) *KnownPeer { - if _, ok := b.KnownPeers[id]; !ok { - b.KnownPeers[id] = &KnownPeer{ - sampleCount: sampleLimit, - bundle: b, - id: id, - } + if peer.listElement != nil { + diff.knownInactivePeerList.Remove(peer.listElement) } - return b.KnownPeers[id] -} - -// KnownPeer contains the metered data of a particular peer. -type KnownPeer struct { - // Sessions contains the metered data in a session of a peer. - Sessions []*PeerSession `json:"sessions,omitempty"` - - sampleCount int - - bundle *PeerBundle // Pointer to the inner map. - id string // Key of the peer in the inner map. - - listElement *list.Element // Pointer to the peer element in the list. -} - -func (peer *KnownPeer) delete() { - delete(peer.bundle.KnownPeers, peer.id) - if len(peer.bundle.KnownPeers) < 1 && len(peer.bundle.UnknownPeers) < 1 { - delete(peer.bundle.root.Bundles, peer.bundle.ip) - } - peer.listElement = nil - peer.bundle = nil - for i := range peer.Sessions { - peer.Sessions[i] = nil - } - peer.Sessions = nil -} - -func (peer *KnownPeer) update(session *PeerSession) { - if peer.Sessions == nil { - peer.Sessions = []*PeerSession{session} - if session.Connected != nil && session.Disconnected != nil && session.Connected.After(*session.Disconnected) { - peer.Sessions = append(peer.Sessions, &PeerSession{Connected: session.Connected}) - session.Connected = nil - } - return - } - if session.Connected != nil && session.Disconnected != nil && session.Connected.After(*session.Disconnected) { - last := peer.Sessions[len(peer.Sessions)-1] - if last.Connected == nil { - // This may happen at the beginning, if the connection was established - // before the initialization of the peer event feed. Otherwise it should - // not happen to handle a disconnect event before the connect, because - // both of them are coming on the same channel consecutively. - log.Warn("Disconnect event appeared without connect") - } - last.Disconnected = session.Disconnected - last.Ingress = append(last.Ingress, session.Ingress...) - last.Egress = append(last.Egress, session.Egress...) - session.Disconnected = nil - session.Ingress = nil - session.Egress = nil - } - if session.Connected != nil { - peer.Sessions = append(peer.Sessions, session) - } - if session.Disconnected != nil { - peer.Sessions[len(peer.Sessions)-1].Disconnected = session.Disconnected - } - for i := 0; i < len(session.Ingress) && i < len(session.Egress); i++ { - peer.updateTraffic(session.Ingress[i], session.Egress[i]) - } -} - -func (peer *KnownPeer) updateTraffic(ingress, egress *ChartEntry) { - if ingress == nil { - ingress = new(ChartEntry) - } - if egress == nil { - egress = new(ChartEntry) - } - first := peer.Sessions[0] - if len(first.Ingress) < 1 || len(first.Egress) < 1 { - first.Ingress = append(first.Ingress, ingress) - first.Egress = append(first.Egress, egress) - peer.sampleCount = 1 - return - } - if peer.sampleCount >= sampleLimit { - if len(first.Ingress) < 2 || len(first.Egress) < 2 { - peer.Sessions = peer.Sessions[1:] + // Set peer activity + if len(peer.Sessions) > 0 { + peer.Active = peer.Sessions[len(peer.Sessions)-1].Disconnected == nil + } else { + diff.rootLock.RLock() + if diff.root.hasKnownPeer(ip, id) { + rootSessions := diff.root.Bundles[ip].KnownPeers[id].Sessions + peer.Active = len(rootSessions) > 0 && rootSessions[len(rootSessions)-1].Disconnected == nil } else { - first.Ingress = first.Ingress[1:] - first.Egress = first.Egress[1:] + peer.Active = false + } + diff.rootLock.RUnlock() + } + if peer.Active { + peer.activeListElement = diff.knownActivePeerList.PushBack(peer) + } else { + peer.listElement = diff.knownInactivePeerList.PushBack(peer) + } + for diff.knownActivePeerList.Len()+diff.knownInactivePeerList.Len() > knownPeerLimit { + var removed interface{} + if diff.knownInactivePeerList.Len() > 0 { + removed = diff.knownInactivePeerList.Remove(diff.knownInactivePeerList.Front()) + } else { + removed = diff.knownActivePeerList.Remove(diff.knownActivePeerList.Front()) + } + if p, ok := removed.(*knownPeerDiff); ok { + diff.removeKnownPeer(p.ip, p.id) + diff.rootLock.RLock() + if diff.root.hasKnownPeer(p.ip, p.id) { + diff.RemovedKnownIP = append(diff.RemovedKnownIP, p.ip) + diff.RemovedKnownID = append(diff.RemovedKnownID, p.id) + } + diff.rootLock.RUnlock() } - peer.sampleCount-- } - last := peer.Sessions[len(peer.Sessions)-1] - last.Ingress = append(last.Ingress, ingress) - last.Egress = append(last.Egress, egress) - peer.sampleCount++ } -func (peer *KnownPeer) len() int { - return peer.sampleCount +func (diff *peerDiff) insertUnknown(ip string, peer *UnknownPeer) { + newBundle := diff.initBundle(ip) + bundle := diff.Bundles[ip] + if newBundle { + bundle.Location = diff.geodb.Location(ip) + } + diff.unknownPeers = append(diff.unknownPeers, ip) + bundle.UnknownPeers = append(bundle.UnknownPeers, peer) + for len(diff.unknownPeers) > unknownPeerLimit { + rip := diff.unknownPeers[0] + diff.RemovedUnknownIP = append(diff.RemovedUnknownIP, rip) + diff.removeUnknownPeer(rip) + diff.unknownPeers = diff.unknownPeers[1:] + } } -type PeerSession 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"` - Disconnected time.Time `json:"disconnected"` +func (diff *peerDiff) dump() { + diff.rootLock.Lock() + for i := 0; i < len(diff.RemovedKnownIP); i++ { + diff.root.removeKnownPeer(diff.RemovedKnownIP[i], diff.RemovedKnownID[i]) + } + for _, rip := range diff.RemovedUnknownIP { + diff.root.removeUnknownPeer(rip) + } + for e := diff.knownActivePeerList.Front(); e != nil; e = e.Next() { + if peer, ok := e.Value.(*knownPeerDiff); ok { + diff.root.getOrInitKnownPeer(peer.ip, peer.id).upgrade(peer.KnownPeer) + } else { + log.Warn("Invalid value in the active peer metrics list") + } + } + for e := diff.knownInactivePeerList.Front(); e != nil; e = e.Next() { + if peer, ok := e.Value.(*knownPeerDiff); ok { + diff.root.getOrInitKnownPeer(peer.ip, peer.id).upgrade(peer.KnownPeer) + } else { + log.Warn("Invalid value in the inactive peer metrics list") + } + } + diff.rootLock.Unlock() + diff.clear() } // collectPeerData gathers data about the peers and sends it to the clients. @@ -422,13 +193,11 @@ func (db *Dashboard) collectPeerData() { ticker := time.NewTicker(db.config.Refresh) defer ticker.Stop() - db.peerLock.Lock() - db.history.Peers = NewPeersMessage(db.geodb) - db.peerLock.Unlock() - diff := NewPeersMessage(db.geodb) + type registryFunc func(name string, i interface{}) + type collectorFunc func(traffic *map[string]float64) registryFunc - trafficCollector := func(prefix string) func(*map[string]float64) func(name string, i interface{}) { - return func(traffic *map[string]float64) func(name string, i interface{}) { + trafficCollector := func(prefix string) collectorFunc { + return func(traffic *map[string]float64) registryFunc { return func(name string, i interface{}) { if m, ok := i.(metrics.Meter); ok { (*traffic)[strings.TrimPrefix(name, prefix)] = float64(m.Count()) @@ -441,6 +210,11 @@ func (db *Dashboard) collectPeerData() { collectIngress := trafficCollector(p2p.MetricsInboundTraffic + "/") collectEgress := trafficCollector(p2p.MetricsOutboundTraffic + "/") + db.peerLock.Lock() + db.history.Network = &NetworkMessage{Peers: NewPeersMessage()} + diff := newPeerDiff(db.history.Network.Peers, &db.peerLock, db.geodb, db.config.Refresh) + db.peerLock.Unlock() + for { select { case event := <-peerCh: @@ -448,11 +222,11 @@ func (db *Dashboard) collectPeerData() { switch event.Type { case p2p.PeerConnected: connected := now.Add(-event.Elapsed) - diff.updateKnownPeer(event.IP.String(), event.ID, &PeerSession{ + diff.insert(event.IP.String(), event.ID, &PeerSession{ Connected: &connected, }) case p2p.PeerDisconnected: - diff.updateKnownPeer(event.IP.String(), event.ID, &PeerSession{ + diff.insert(event.IP.String(), event.ID, &PeerSession{ Disconnected: &now, Ingress: ChartEntries{ &ChartEntry{ @@ -468,7 +242,7 @@ func (db *Dashboard) collectPeerData() { }, }) case p2p.PeerHandshakeFailed: - diff.updateUnknownPeer(event.IP.String(), &UnknownPeer{ + diff.insertUnknown(event.IP.String(), &UnknownPeer{ Connected: now.Add(-event.Elapsed), Disconnected: now, }) @@ -477,67 +251,51 @@ func (db *Dashboard) collectPeerData() { } case <-ticker.C: ingress, egress := make(map[string]float64), make(map[string]float64) - db.peerLock.Lock() - for e := db.history.Peers.knownPeers.Front(); e != nil; e = e.Next() { - if p, ok := e.Value.(*KnownPeer); ok { - key := fmt.Sprintf("%s/%s", p.bundle.ip, p.id) - ingress[key] = 0 - egress[key] = 0 - } - } p2p.PeerIngressRegistry.Each(collectIngress(&ingress)) p2p.PeerEgressRegistry.Each(collectEgress(&egress)) - //diff.updateTraffic(&ingress, &egress) - for key := range ingress { + + now := time.Now() + appendSample := func(key string, ingress, egress float64) { if k := strings.Split(key, "/"); len(k) == 2 { - diff.updateTraffic(k[0], k[1], ingress[key], egress[key]) + diff.insert(k[0], k[1], &PeerSession{ + Ingress: ChartEntries{&ChartEntry{ + Time: now, + Value: ingress, + }}, + Egress: ChartEntries{&ChartEntry{ + Time: now, + Value: egress, + }}, + }) } else { - log.Warn("Bad key", "key", key) + log.Warn("Invalid traffic key", "key", key) } } - for key := range egress { - if _, ok := ingress[key]; !ok { - if k := strings.Split(key, "/"); len(k) == 2 { - diff.updateTraffic(k[0], k[1], ingress[key], egress[key]) - } else { - log.Warn("Bad key", "key", key) - } + for key, val := range ingress { + appendSample(key, val, egress[key]) + } + for key, val := range egress { + if _, ok := ingress[key]; ok { + continue + } + appendSample(key, ingress[key], val) + } + for e := diff.knownInactivePeerList.Front(); e != nil; e = e.Next() { + if peer, ok := e.Value.(*knownPeerDiff); ok { + diff.insert(peer.ip, peer.id, &PeerSession{ + Ingress: ChartEntries{&ChartEntry{ + Time: now, + }}, + Egress: ChartEntries{&ChartEntry{ + Time: now, + }}, + }) } } - for _, bundle := range diff.Bundles { - for _, peer := range bundle.KnownPeers { - var ipExists, idExists bool - var b *PeerBundle - if b, ipExists = db.history.Peers.Bundles[bundle.ip]; ipExists { - _, idExists = b.KnownPeers[peer.id] - } - if !idExists || !ipExists { - fmt.Println("doesn't exist ", bundle.ip, peer.id) - //t, n := time.Now(), sampleLimit-peer.len() - //if len(peer.Sessions) > 0 && peer.Sessions[0].Connected != nil { - // t = *peer.Sessions[0].Connected - //} - //peer.Sessions = append([]*PeerSession{{ - // Ingress: emptyChartEntries(t, n, db.config.Refresh), - // Egress: emptyChartEntries(t, n, db.config.Refresh), - //}}, peer.Sessions...) - } else { - fmt.Println("does exist ", bundle.ip, peer.id) - s, _ := json.MarshalIndent(bundle, "", " ") - fmt.Println(string(s)) - fmt.Println(ingress[fmt.Sprintf("%s/%s", bundle.ip, peer.id)]) - fmt.Println(egress[fmt.Sprintf("%s/%s", bundle.ip, peer.id)]) - } - } - } - diff.RemovedKnown, diff.RemovedUnknown = db.history.Peers.append(diff) - //sh, _ := json.MarshalIndent(db.history.Network, "", " ") - //fmt.Println(string(sh)) - db.peerLock.Unlock() - //s, _ := json.MarshalIndent(diff, "", " ") - //fmt.Println(string(s)) - //db.sendToAll(&Message{Peers: diff}) - diff.clear() + db.sendToAll(&Message{Network: &NetworkMessage{Peers: deepcopy.Copy(diff.PeersMessage).(*PeersMessage)}}) + s, _ := json.MarshalIndent(deepcopy.Copy(diff), "", " ") + fmt.Println(string(s)) + diff.dump() case err := <-subPeer.Err(): log.Warn("Peer subscription error", "err", err) return