dashboard, p2p: experiment 1

This commit is contained in:
Kurkó Mihály 2018-09-17 13:16:13 +03:00
parent f2e52e2fee
commit cac8611230
5 changed files with 118 additions and 146 deletions

View file

@ -42,19 +42,7 @@ import (
)
const (
activeMemorySampleLimit = 200 // Maximum number of active memory data samples
virtualMemorySampleLimit = 200 // Maximum number of virtual memory data samples
networkIngressSampleLimit = 200 // Maximum number of network ingress data samples
networkEgressSampleLimit = 200 // Maximum number of network egress data samples
processCPUSampleLimit = 200 // Maximum number of process cpu data samples
systemCPUSampleLimit = 200 // Maximum number of system cpu data samples
diskReadSampleLimit = 200 // Maximum number of disk read data samples
diskWriteSampleLimit = 200 // Maximum number of disk write data samples
peerLimit = 1000 // Maximum number of metered peers
peerTrafficSampleLimit = 200 // Maximum number of traffic data samples for a peer
peerIngressSampleLimit = peerTrafficSampleLimit // Maximum number of ingress data samples for a peer
peerEgressSampleLimit = peerTrafficSampleLimit // Maximum number of egress data samples for a peer
sampleLimit = 200 // Maximum number of data samples
)
// Dashboard contains the dashboard internals.
@ -103,14 +91,14 @@ func New(config *Config, commit string, logdir string) *Dashboard {
Version: fmt.Sprintf("v%d.%d.%d%s", params.VersionMajor, params.VersionMinor, params.VersionPatch, versionMeta),
},
System: &SystemMessage{
ActiveMemory: emptyChartEntries(now, activeMemorySampleLimit, config.Refresh),
VirtualMemory: emptyChartEntries(now, virtualMemorySampleLimit, config.Refresh),
NetworkIngress: emptyChartEntries(now, networkIngressSampleLimit, config.Refresh),
NetworkEgress: emptyChartEntries(now, networkEgressSampleLimit, config.Refresh),
ProcessCPU: emptyChartEntries(now, processCPUSampleLimit, config.Refresh),
SystemCPU: emptyChartEntries(now, systemCPUSampleLimit, config.Refresh),
DiskRead: emptyChartEntries(now, diskReadSampleLimit, config.Refresh),
DiskWrite: emptyChartEntries(now, diskWriteSampleLimit, config.Refresh),
ActiveMemory: emptyChartEntries(now, sampleLimit, config.Refresh),
VirtualMemory: emptyChartEntries(now, sampleLimit, config.Refresh),
NetworkIngress: emptyChartEntries(now, sampleLimit, config.Refresh),
NetworkEgress: emptyChartEntries(now, sampleLimit, config.Refresh),
ProcessCPU: emptyChartEntries(now, sampleLimit, config.Refresh),
SystemCPU: emptyChartEntries(now, sampleLimit, config.Refresh),
DiskRead: emptyChartEntries(now, sampleLimit, config.Refresh),
DiskWrite: emptyChartEntries(now, sampleLimit, config.Refresh),
},
Network: &NetworkMessage{
PeerBundles: make(map[string]*PeerBundle),

View file

@ -48,6 +48,9 @@ type GeoDB struct {
geodb *freegeoip.DB
}
// TODO (kurkomisi): freegeoip newReader - error opening the maxmindDB file (/tmp/freegeoip/db.gz)
// error message: "gzip: invalid header" - possibly a bad update of the file
// Open creates a new geoip database with an up-to-date database from the internet.
func OpenGeoDB() (*GeoDB, error) {
// Initiate a geoip database to cross reference locations
@ -56,11 +59,11 @@ func OpenGeoDB() (*GeoDB, error) {
return nil, err
}
// Wait until the database is updated to the latest data
select {
case <-db.NotifyOpen():
case err := <-db.NotifyError():
return nil, err
}
//select {
//case <-db.NotifyOpen():
//case err := <-db.NotifyError():
// return nil, err
//}
// Assemble and return our custom wrapper
return &GeoDB{geodb: db}, nil
}
@ -77,3 +80,14 @@ func (db *GeoDB) Lookup(ip net.IP) *GeoDBInfo {
db.geodb.Lookup(ip, result)
return result
}
func (db *GeoDB) Location (ip string) *GeoLocation {
//location := db.Lookup(net.ParseIP(ip))
location := &GeoDBInfo{}
return &GeoLocation{
Country: location.Country.Names.English,
City: location.City.Names.English,
Latitude: location.Location.Latitude,
Longitude: location.Location.Longitude,
}
}

View file

@ -17,6 +17,7 @@
package dashboard
import (
"container/list"
"encoding/json"
"time"
)
@ -74,11 +75,7 @@ func (m *NetworkMessage) getOrInitBundle(ip string) *PeerBundle {
// 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 {
b := m.getOrInitBundle(ip)
if _, ok := b.Peers[id]; !ok {
b.Peers[id] = new(Peer)
}
return b.Peers[id]
return m.getOrInitBundle(ip).getOrInitPeer(id)
}
// PeerBundle contains information about the peers pertaining to an IP address.
@ -88,6 +85,13 @@ type PeerBundle struct {
FailedPeers []*Peer `json:"failedPeers,omitempty"`
}
func (b *PeerBundle) getOrInitPeer(id string) *Peer {
if _, ok := b.Peers[id]; !ok {
b.Peers[id] = new(Peer)
}
return b.Peers[id]
}
// GeoLocation contains geographical information.
type GeoLocation struct {
Country string `json:"country,omitempty"`
@ -99,13 +103,14 @@ type GeoLocation struct {
// Peer contains lifecycle timestamps and traffic information of a given peer.
type Peer struct {
Connected []time.Time `json:"connected,omitempty"`
Handshake []time.Time `json:"handshake,omitempty"`
Disconnected []time.Time `json:"disconnected,omitempty"`
Ingress ChartEntries `json:"ingress,omitempty"`
Egress ChartEntries `json:"egress,omitempty"`
DefaultID string `json:"defaultID,omitempty"`
element *list.Element
}
// SystemMessage contains the metered system data samples.

View file

@ -17,13 +17,14 @@
package dashboard
import (
"container/list"
"encoding/json"
"fmt"
"net"
"github.com/mohae/deepcopy"
"time"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/mohae/deepcopy"
)
const eventBufferLimit = 128 // Maximum number of buffered peer events
@ -69,6 +70,15 @@ func (db *Dashboard) collectPeerData() {
ticker := time.NewTicker(db.config.Refresh)
defer ticker.Stop()
purgeOrder := list.New()
failedPurgeOrder := list.New()
update := func(peer *Peer, l *list.List) {
if peer.element == nil {
peer.element = purgeOrder.PushBack(peer)
} else {
purgeOrder.MoveToBack(peer.element)
}
}
// Listen for events, and prepare the difference between two metering.
diff := &NetworkMessage{
PeerBundles: make(map[string]*PeerBundle),
@ -76,135 +86,92 @@ func (db *Dashboard) collectPeerData() {
for {
select {
case event := <-connectCh:
ip := event.IP
p := diff.getOrInitPeer(ip, event.ID)
if diff.PeerBundles[ip].Location == nil {
db.peerLock.RLock()
lookup := db.history.Network.PeerBundles[ip] == nil || db.history.Network.PeerBundles[ip].Location == nil
db.peerLock.RUnlock()
if lookup {
location := db.geodb.Lookup(net.ParseIP(event.IP))
diff.PeerBundles[ip].Location = &GeoLocation{
Country: location.Country.Names.English,
City: location.City.Names.English,
Latitude: location.Location.Latitude,
Longitude: location.Location.Longitude,
}
}
}
if p.Connected == nil {
p.Connected = []time.Time{event.Connected}
} else {
p.Connected = append(p.Connected, event.Connected)
}
//case event := <-failedCh:
// ip := event.IP
// p := diff.getOrInitPeer(ip, event.AutoID)
// p.DefaultID = event.AutoID
// if p.Handshake == nil {
// p.Handshake = []time.Time{event.Handshake}
// } else {
// p.Handshake = append(p.Handshake, event.Handshake)
// }
// delete(diff.PeerBundles[ip].Peers, event.AutoID)
// diff.getOrInitPeer(ip, event.ID)
// diff.PeerBundles[ip].Peers[event.ID] = p // TODO (kurkomisi): Merge instead in order to keep the previous connection.
// // Remove the peer from history in case the metering was before the handshake.
// db.peerLock.RLock()
// stored := db.history.Network.PeerBundles[ip] != nil && db.history.Network.PeerBundles[ip].Peers[event.AutoID] != nil
// db.peerLock.RUnlock()
// if stored {
// db.peerLock.Lock()
// hp := db.history.Network.getOrInitPeer(ip, event.AutoID)
// delete(db.history.Network.PeerBundles[ip].Peers, event.AutoID)
// db.history.Network.getOrInitPeer(ip, event.ID)
// db.history.Network.PeerBundles[ip].Peers[event.ID] = hp // TODO (kurkomisi): Merge.
// db.peerLock.Unlock()
// }
diffBundle := diff.getOrInitBundle(event.IP)
diffBundle.Location = db.geodb.Location(event.IP)
diffPeer := diffBundle.getOrInitPeer(event.ID)
diffPeer.Connected = append(diffPeer.Connected, event.Connected)
case event := <-failedCh:
diffBundle := diff.getOrInitBundle(event.IP)
diffBundle.Location = db.geodb.Location(event.IP)
diffBundle.FailedPeers = append(diffBundle.FailedPeers, &Peer{
Connected: []time.Time{event.Connected},
Disconnected: []time.Time{event.Disconnected},
})
case event := <-disconnectCh:
p := diff.getOrInitPeer(event.IP, event.ID)
if p.Disconnected == nil {
p.Disconnected = []time.Time{event.Disconnected}
} else {
p.Disconnected = append(p.Disconnected, event.Disconnected)
}
diffPeer := diff.getOrInitPeer(event.IP, event.ID)
diffPeer.Disconnected = append(diffPeer.Disconnected, event.Disconnected)
case event := <-ingressCh:
fmt.Println("ingress", event.IP, event.Amount)
// Sum up the ingress between two updates.
p := diff.getOrInitPeer(event.IP, event.ID)
if len(p.Ingress) <= 0 {
p.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}}
diffPeer := diff.getOrInitPeer(event.IP, event.ID)
if len(diffPeer.Ingress) != 1 {
diffPeer.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}}
} else {
p.Ingress[0].Value += float64(event.Amount)
diffPeer.Ingress[0].Value = float64(event.Amount)
}
case event := <-egressCh:
fmt.Println("egress ", event.IP, event.Amount)
// Sum up the egress between two updates.
p := diff.getOrInitPeer(event.IP, event.ID)
if len(p.Egress) <= 0 {
p.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}}
diffPeer := diff.getOrInitPeer(event.IP, event.ID)
if len(diffPeer.Egress) != 1 {
diffPeer.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}}
} else {
p.Egress[0].Value += float64(event.Amount)
diffPeer.Egress[0].Value = float64(event.Amount)
}
case <-ticker.C:
now := time.Now()
// Merge the diff with the history.
db.peerLock.Lock()
for ip, bundle := range diff.PeerBundles {
if bundle.Location != nil {
b := db.history.Network.getOrInitBundle(ip)
b.Location = bundle.Location
for ip, diffBundle := range diff.PeerBundles {
historyBundle := db.history.Network.getOrInitBundle(ip)
historyBundle.Location = diffBundle.Location
for id, diffPeer := range diffBundle.Peers {
historyPeer := historyBundle.getOrInitPeer(id)
historyPeer.Connected = append(historyPeer.Connected, diffPeer.Connected...)
historyPeer.Disconnected = append(historyPeer.Disconnected, diffPeer.Disconnected...)
if len(diffPeer.Ingress) == 1 {
diffPeer.Ingress[0].Time = now
if historyPeer.Ingress == nil {
historyPeer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), sampleLimit-1, db.config.Refresh), diffPeer.Ingress[0])
// The first message about a diffPeer should contain the whole list
diffPeer.Ingress = historyPeer.Ingress
} else {
historyPeer.Ingress = append(historyPeer.Ingress, diffPeer.Ingress[0])[1:]
}
}
if len(diffPeer.Egress) == 1 {
diffPeer.Egress[0].Time = now
if historyPeer.Egress == nil {
historyPeer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), sampleLimit-1, db.config.Refresh), diffPeer.Egress[0])
// The first message about a diffPeer should contain the whole list
diffPeer.Egress = historyPeer.Egress
} else {
historyPeer.Egress = append(historyPeer.Egress, diffPeer.Egress[0])[1:]
}
}
update(historyPeer, purgeOrder)
}
for id, peer := range bundle.Peers {
peerHistory := db.history.Network.getOrInitPeer(ip, id)
if peer.Connected != nil {
peerHistory.Connected = append(peerHistory.Connected, peer.Connected...)
}
if peer.Handshake != nil {
peerHistory.Handshake = append(peerHistory.Handshake, peer.Handshake...)
}
if peer.Disconnected != nil {
peerHistory.Disconnected = append(peerHistory.Disconnected, peer.Disconnected...)
}
ingress := &ChartEntry{
Time: now,
}
if len(peer.Ingress) > 0 {
ingress.Value = peer.Ingress[0].Value
}
if peerHistory.Ingress == nil {
peer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), peerIngressSampleLimit-1, db.config.Refresh), ingress)
peerHistory.Ingress = peer.Ingress
} else {
peer.Ingress = ChartEntries{ingress}
peerHistory.Ingress = append(peerHistory.Ingress[1:], ingress)
}
egress := &ChartEntry{
Time: now,
}
if len(peer.Egress) > 0 {
egress.Value = peer.Egress[0].Value
}
if peerHistory.Egress == nil {
peer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), peerEgressSampleLimit-1, db.config.Refresh), egress)
peerHistory.Egress = peer.Egress
} else {
peer.Egress = ChartEntries{egress}
peerHistory.Egress = append(peerHistory.Egress[1:], egress)
}
historyBundle.FailedPeers = append(historyBundle.FailedPeers, diffBundle.FailedPeers...)
for _, fp := range diffBundle.FailedPeers {
update(fp, failedPurgeOrder)
}
}
for purgeOrder.Len() > p2p.MeteredPeerLimit {
purgeOrder.Remove(purgeOrder.Front())
}
for failedPurgeOrder.Len() > p2p.MeteredPeerLimit {
failedPurgeOrder.Remove(failedPurgeOrder.Front())
}
//ss, _ := json.MarshalIndent(db.history.Network, "", " ")
//fmt.Println(string(ss))
db.peerLock.Unlock()
// Send the diff to the clients.
db.sendToAll(&Message{Network: deepcopy.Copy(diff).(*NetworkMessage)})
//s, _ := json.MarshalIndent(diff, "", " ")
//fmt.Println(string(s))
s, _ := json.MarshalIndent(diff, "", " ")
fmt.Println(string(s))
// Prepare for the next metering, clear the diff variable.
for ip, bundle := range diff.PeerBundles {
for id := range bundle.Peers {
bundle.Peers[id] = nil
delete(bundle.Peers, id)
}
delete(diff.PeerBundles, ip)
diff = &NetworkMessage{
PeerBundles: make(map[string]*PeerBundle),
}
case err := <-subConnect.Err():
log.Warn("Peer connect subscription error", "err", err)

View file

@ -77,7 +77,6 @@ type PeerConnectEvent struct {
IP string
ID string
Connected time.Time
Handshake time.Time
}
// PeerDisconnectEvent contains information about the disconnection of a peer.
@ -303,6 +302,5 @@ func (c *meteredConn) handshakeDone(id discover.NodeID) {
IP: c.ip,
ID: id.String(),
Connected: c.connected,
Handshake: time.Now(),
})
}