dashboard, p2p: initial version of peer traffic metering with event feeds

This commit is contained in:
Kurkó Mihály 2018-08-07 14:28:25 +03:00
parent 82ae16192a
commit 8fd0e4ad1e
7 changed files with 533 additions and 278 deletions

View file

@ -18,36 +18,59 @@
import React, {Component} from 'react';
import Table, {TableBody, TableHeader, TableHeaderColumn, TableRow, TableCell} from 'material-ui/Table';
import Table, {TableHead, TableBody, TableRow, TableCell} from 'material-ui/Table';
import type {Network as NetworkType, Peer} 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: {[number]: Peer}, prev: {[number]: Peer}) => {
Object.keys(update).forEach((k) => {
if (!prev[k]) {
prev[k] = update[k];
export const inserter = (update: {[string]: {[string]: Peer}}, prev: {[string]: {[string]: Peer}}) => {
Object.keys(update).forEach((ip) => {
if (!prev[ip]) {
prev[ip] = update[ip];
return;
}
const u: Peer = update[k];
const p: Peer = prev[k];
if (u.id) {
p.id = u.id;
if (!update[ip]) {
return;
}
if (u.ip) {
p.ip = u.ip;
Object.keys(update[ip]).forEach((id) => {
if (!prev[ip][id]) {
prev[ip][id] = update[ip][id];
return;
}
if (u.lifecycle) {
if (u.lifecycle.handshake) {
p.lifecycle.handshake = u.lifecycle.handshake;
const u: Peer = update[ip][id];
const p: Peer = prev[ip][id];
if (u.connected) {
if (!Array.isArray(p.connected)) {
p.connected = [];
}
if (u.lifecycle.disconnected) {
p.lifecycle.disconnected = u.lifecycle.disconnected;
p.connected = [...p.connected, ...u.connected];
}
if (u.handshake) {
if (!Array.isArray(p.handshake)) {
p.handshake = [];
}
p.handshake = [...p.handshake, ...u.handshake];
}
if (u.disconnected) {
if (!Array.isArray(p.disconnected)) {
p.disconnected = [];
}
p.disconnected = [...p.disconnected, ...u.disconnected];
}
if (Array.isArray(u.ingress)) {
if (!Array.isArray(p.ingress)) {
p.ingress = [];
}
p.ingress = [...p.ingress, ...u.ingress].slice(-200);
}
if (Array.isArray(u.egress)) {
if (!Array.isArray(p.egress)) {
p.egress = [];
}
p.egress = [...p.egress, ...u.egress].slice(-200);
prev[k] = p;
}
prev[ip][id] = p;
});
});
return prev;
};
@ -63,19 +86,65 @@ export type Props = {
// Network renders the network page.
class Network extends Component<Props, State> {
formatTime = (t) => {
const time = new Date(t);
if (isNaN(time)) {
return '';
}
const month = `0${time.getMonth() + 1}`.slice(-2);
const date = `0${time.getDate()}`.slice(-2);
const hours = `0${time.getHours()}`.slice(-2);
const minutes = `0${time.getMinutes()}`.slice(-2);
const seconds = `0${time.getSeconds()}`.slice(-2);
return `${month}/${date}/${hours}:${minutes}:${seconds}`;
};
render() {
return (
<Table>
<TableBody>
{Object.entries(this.props.content.peers).map(([k, v]) => (
<TableHead>
<TableRow>
<TableCell>{k}</TableCell>
<TableCell>{v.id ? v.id.substring(0, 6) : ''}</TableCell>
<TableCell>{v.ip}</TableCell>
<TableCell>{v.ingress.value}</TableCell>
<TableCell>{v.egress.value}</TableCell>
<TableCell>{JSON.stringify(v.location)}</TableCell>
<TableCell>{JSON.stringify(v.lifecycle)}</TableCell>
<TableCell>IP</TableCell>
<TableCell>Peer ID</TableCell>
<TableCell>Location</TableCell>
<TableCell>Ingress</TableCell>
<TableCell>Egress</TableCell>
<TableCell>Connected</TableCell>
<TableCell>Handshake</TableCell>
<TableCell>Disconnected</TableCell>
</TableRow>
</TableHead>
<TableBody>
{Object.entries(this.props.content.peers).map(([ip, peers]) => (
<TableRow key={ip}>
<TableCell>{ip}</TableCell>
<TableCell>
{Object.keys(peers).map(id => id.substring(0, 10)).join(' ')}
</TableCell>
<TableCell>
{(() => {
const k = Object.keys(peers)[0];
return k && peers[k].location ? (() => {
const l = peers[k].location;
return `${l.country}${l.city ? `/${l.city}` : ''} ${l.latitude} ${l.longitude}`;
})() : '';
})()}
</TableCell>
<TableCell>
{Object.keys(peers).map((id) => peers[id].ingress && peers[id].ingress.map(sample => sample.value).join(' ')).join(', ')}
</TableCell>
<TableCell>
{Object.keys(peers).map((id) => peers[id].egress && peers[id].egress.map(sample => sample.value).join(' ')).join(', ')}
</TableCell>
<TableCell>
{Object.keys(peers).map((id) => peers[id].connected && peers[id].connected.map(time => this.formatTime(time)).join(' ')).join(', ')}
</TableCell>
<TableCell>
{Object.keys(peers).map((id) => peers[id].handshake && peers[id].handshake.map(time => this.formatTime(time)).join(' ')).join(', ')}
</TableCell>
<TableCell>
{Object.keys(peers).map((id) => peers[id].disconnected && peers[id].disconnected.map(time => this.formatTime(time)).join(' ')).join(', ')}
</TableCell>
</TableRow>
))}
</TableBody>

View file

@ -51,14 +51,14 @@ export type TxPool = {
};
export type Network = {
peers: {[number]: Peer},
peers: {[string]: {[string]: Peer}},
};
export type Peer = {
id: string,
ip: string,
location: PeerLocation,
lifecycle: PeerLifecycle,
connected: Array<Date>,
handshake: Array<Date>,
disconnected: Array<Date>,
ingress: ChartEntries,
egress: ChartEntries,
};
@ -70,12 +70,6 @@ export type PeerLocation = {
longitude: number,
};
export type PeerLifecycle = {
connected: Date,
handshake: Date,
disconnected: Date,
};
export type System = {
activeMemory: ChartEntries,
virtualMemory: ChartEntries,

View file

@ -54,9 +54,9 @@ const (
diskReadSampleLimit = 200 // Maximum number of disk read data samples
diskWriteSampleLimit = 200 // Maximum number of disk write data samples
storedDisconnectedPeerLimit = p2p.MeteredPeerLimit - p2p.DefaultMaxPendingPeers
peerIngressSampleLimit = 200
peerEgressSampleLimit = 200
peerTrafficSampleLimit = 200
peerIngressSampleLimit = peerTrafficSampleLimit
peerEgressSampleLimit = peerTrafficSampleLimit
)
var nextID uint32 // Next connection id
@ -68,10 +68,13 @@ type Dashboard struct {
listener net.Listener
conns map[uint32]*client // Currently live websocket connections
history *Message
peerHistory *NetworkMessage
lock sync.RWMutex // Lock protecting the dashboard's internals
peerLock sync.RWMutex
geodb *GeoDB
disconnectedPeerIDs chan uint
peersByIP map[string]*Peer
peersByID map[string]*Peer
logdir string
quit chan chan error // Channel used for graceful exit
@ -111,11 +114,10 @@ func New(config *Config, commit string, logdir string) *Dashboard {
DiskRead: emptyChartEntries(now, diskReadSampleLimit, config.Refresh),
DiskWrite: emptyChartEntries(now, diskWriteSampleLimit, config.Refresh),
},
Network: &NetworkMessage{
Peers: make(map[uint]*Peer),
},
peerHistory: &NetworkMessage{
Peers: make(map[string]map[string]*Peer),
},
disconnectedPeerIDs: make(chan uint, storedDisconnectedPeerLimit),
logdir: logdir,
}
}
@ -142,16 +144,10 @@ func (db *Dashboard) APIs() []rpc.API { return nil }
func (db *Dashboard) Start(server *p2p.Server) error {
log.Info("Starting dashboard")
var err error
db.geodb, err = OpenGeoDB()
if err != nil {
log.Warn("Failed to open geodb", "err", err)
return err
}
db.wg.Add(2)
db.wg.Add(3)
go db.collectData()
go db.streamLogs()
go db.collectPeerData()
http.HandleFunc("/", db.webHandler)
http.Handle("/api", websocket.Handler(db.apiHandler))
@ -170,8 +166,6 @@ func (db *Dashboard) Start(server *p2p.Server) error {
// Stop stops the data collection thread and the connection listener of the dashboard.
// Implements the node.Service interface.
func (db *Dashboard) Stop() error {
db.geodb.Close()
// Close the connection listener.
var errs []error
if err := db.listener.Close(); err != nil {
@ -179,7 +173,7 @@ func (db *Dashboard) Stop() error {
}
// Close the collectors.
errc := make(chan error, 1)
for i := 0; i < 2; i++ {
for i := 0; i < 3; i++ {
db.quit <- errc
if err := <-errc; err != nil {
errs = append(errs, err)
@ -252,10 +246,17 @@ func (db *Dashboard) apiHandler(conn *websocket.Conn) {
}
}()
db.lock.Lock()
// Send the past data.
client.msg <- deepcopy.Copy(db.history).(*Message)
db.lock.RLock()
h := deepcopy.Copy(db.history).(*Message)
db.lock.RUnlock()
db.peerLock.RLock()
h.Network = deepcopy.Copy(db.peerHistory).(*NetworkMessage)
db.peerLock.RUnlock()
client.msg <- h
// Start tracking the connection and drop at connection loss.
db.lock.Lock()
db.conns[id] = client
db.lock.Unlock()
defer func() {
@ -392,84 +393,6 @@ func (db *Dashboard) collectData() {
sys.DiskWrite = append(sys.DiskWrite[1:], diskWrite)
db.lock.Unlock()
peers := p2p.PeerTrafficMeters.Peers()
network := new(NetworkMessage)
if len(peers) > 0 {
network.Peers = make(map[uint]*Peer)
}
for id, peer := range peers {
peerIngress := &ChartEntry{
Time: now,
Value: float64(peer.Ingress),
}
peerEgress := &ChartEntry{
Time: now,
Value: float64(peer.Egress),
}
p := &Peer{
Ingress: ChartEntries{peerIngress},
Egress: ChartEntries{peerEgress},
}
if prevMetrics, ok := db.history.Network.Peers[id]; !ok {
p.ID = peer.ID
p.IP = peer.IP
location := db.geodb.Lookup(peer.IP)
p.Location = &PeerLocation{
Country: location.Country.Names.English,
City: location.City.Names.English,
Latitude: location.Location.Latitude,
Longitude: location.Location.Longitude,
}
p.Lifecycle = &PeerLifecycle{
Connected: peer.Connected,
}
if peer.Handshake != nil {
p.Lifecycle.Handshake = peer.Handshake
}
if peer.Disconnected != nil {
p.Lifecycle.Disconnected = peer.Disconnected
}
db.history.Network.Peers[id] = p
} else {
if prevMetrics.ID != peer.ID {
prevMetrics.ID, p.ID = peer.ID, peer.ID
}
if prevMetrics.Lifecycle.Handshake == nil && peer.Handshake != nil {
p.Lifecycle = new(PeerLifecycle)
prevMetrics.Lifecycle.Handshake, p.Lifecycle.Handshake = peer.Handshake, peer.Handshake
}
if prevMetrics.Lifecycle.Disconnected == nil && peer.Disconnected != nil {
if p.Lifecycle == nil {
p.Lifecycle = new(PeerLifecycle)
}
prevMetrics.Lifecycle.Disconnected, p.Lifecycle.Disconnected = peer.Disconnected, peer.Disconnected
}
first := 0
if len(prevMetrics.Ingress) >= peerIngressSampleLimit {
first = len(prevMetrics.Ingress) - peerIngressSampleLimit + 1
}
prevMetrics.Ingress = append(prevMetrics.Ingress[first:], peerIngress)
first = 0
if len(prevMetrics.Egress) >= peerEgressSampleLimit {
first = len(prevMetrics.Egress) - peerEgressSampleLimit + 1
}
prevMetrics.Egress = append(prevMetrics.Egress[first:], peerEgress)
}
if p.Lifecycle != nil && p.Lifecycle.Disconnected != nil {
select {
case db.disconnectedPeerIDs <- id:
default:
// if the number of the stored disconnected peers exceeds the limit, remove the firstly disconnected peer
delete(db.history.Network.Peers, <-db.disconnectedPeerIDs)
db.disconnectedPeerIDs <- id
}
}
network.Peers[id] = p
}
db.lock.Unlock()
//s, _ := json.MarshalIndent(network, "", " ")
//fmt.Println(string(s))
db.sendToAll(&Message{
System: &SystemMessage{
ActiveMemory: ChartEntries{activeMemory},
@ -481,7 +404,6 @@ func (db *Dashboard) collectData() {
DiskRead: ChartEntries{diskRead},
DiskWrite: ChartEntries{diskWrite},
},
Network: network,
})
}
}

View file

@ -18,7 +18,6 @@ package dashboard
import (
"encoding/json"
"net"
"time"
)
@ -56,15 +55,16 @@ type TxPoolMessage struct {
/* TODO (kurkomisi) */
}
// k1: IP, k2: ID
type NetworkMessage struct {
Peers map[uint]*Peer `json:"peers,omitempty"`
Peers map[string]map[string]*Peer `json:"peers,omitempty"`
}
type Peer struct {
ID string `json:"id,omitempty"`
IP net.IP `json:"ip,omitempty"`
Location *PeerLocation `json:"location,omitempty"`
Lifecycle *PeerLifecycle `json:"lifecycle,omitempty"`
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"`
}
@ -76,12 +76,6 @@ type PeerLocation struct {
Longitude float64 `json:"longitude,omitempty"`
}
type PeerLifecycle struct {
Connected *time.Time `json:"connected,omitempty"`
Handshake *time.Time `json:"handshake,omitempty"`
Disconnected *time.Time `json:"disconnected,omitempty"`
}
type SystemMessage struct {
ActiveMemory ChartEntries `json:"activeMemory,omitempty"`
VirtualMemory ChartEntries `json:"virtualMemory,omitempty"`
@ -93,7 +87,7 @@ type SystemMessage struct {
DiskWrite ChartEntries `json:"diskWrite,omitempty"`
}
// LogsMessage wraps up a log chunk. If Source isn't present, the chunk is a stream chunk.
// LogsMessage wraps up a log chunk. If 'Source' isn't present, the chunk is a stream chunk.
type LogsMessage struct {
Source *LogFile `json:"source,omitempty"` // Attributes of the log file.
Chunk json.RawMessage `json:"chunk"` // Contains log records.

301
dashboard/peers.go Normal file
View file

@ -0,0 +1,301 @@
package dashboard
import (
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"time"
"github.com/mohae/deepcopy"
)
const eventBufferLimit = 128
func getOrInitPeer(m *NetworkMessage, ip, id string) *Peer {
if _, ok := m.Peers[ip]; !ok {
m.Peers[ip] = make(map[string]*Peer)
}
if _, ok := m.Peers[ip][id]; !ok {
m.Peers[ip][id] = new(Peer)
}
return m.Peers[ip][id]
}
func (db *Dashboard) collectPeerData() {
defer db.wg.Done()
var err error
db.geodb, err = OpenGeoDB()
if err != nil {
log.Warn("Failed to open geodb", "err", err)
return
}
defer db.geodb.Close()
var (
quit = make(chan struct{})
connectCh = make(chan *p2p.PeerConnectEvent, eventBufferLimit)
handshakeCh = make(chan *p2p.PeerHandshakeEvent, eventBufferLimit)
disconnectCh = make(chan *p2p.PeerDisconnectEvent, eventBufferLimit)
readCh = make(chan *p2p.PeerReadEvent, eventBufferLimit)
writeCh = make(chan *p2p.PeerWriteEvent, eventBufferLimit)
)
go func() {
var (
peerConnectEventCh = make(chan p2p.PeerConnectEvent, eventBufferLimit)
peerHandshakeEventCh = make(chan p2p.PeerHandshakeEvent, eventBufferLimit)
peerDisconnectEventCh = make(chan p2p.PeerDisconnectEvent, eventBufferLimit)
peerReadEventCh = make(chan p2p.PeerReadEvent, eventBufferLimit)
peerWriteEventCh = make(chan p2p.PeerWriteEvent, eventBufferLimit)
subConnect = p2p.SubscribePeerConnectEvent(peerConnectEventCh)
subHandshake = p2p.SubscribePeerHandshakeEvent(peerHandshakeEventCh)
subDisconnect = p2p.SubscribePeerDisconnectEvent(peerDisconnectEventCh)
subRead = p2p.SubscribePeerReadEvent(peerReadEventCh)
subWrite = p2p.SubscribePeerWriteEvent(peerWriteEventCh)
)
defer func() {
subConnect.Unsubscribe()
subHandshake.Unsubscribe()
subDisconnect.Unsubscribe()
subRead.Unsubscribe()
subWrite.Unsubscribe()
}()
for {
select {
case event := <-peerConnectEventCh:
select {
case connectCh <- &event:
default:
log.Warn("Failed to handle connect event", "event", event)
}
case event := <-peerHandshakeEventCh:
select {
case handshakeCh <- &event:
default:
log.Warn("Failed to handle handshake event", "event", event)
}
case event := <-peerDisconnectEventCh:
select {
case disconnectCh <- &event:
default:
log.Warn("Failed to handle disconnect event", "event", event)
}
case event := <-peerReadEventCh:
select {
case readCh <- &event:
default:
log.Warn("Failed to handle read event", "event", event)
}
case event := <-peerWriteEventCh:
select {
case writeCh <- &event:
default:
log.Warn("Failed to handle write event", "event", event)
}
case <-quit:
return
}
}
}()
go db.cleanPeerHistory(quit)
ticker := time.NewTicker(db.config.Refresh)
defer ticker.Stop()
network := &NetworkMessage{
Peers: make(map[string]map[string]*Peer),
}
for {
select {
case event := <-connectCh:
ip := event.IP.String()
p := getOrInitPeer(network, ip, event.ID)
if p.Location == nil {
db.peerLock.RLock()
peers := db.peerHistory.Peers
lookup := peers[ip] == nil || peers[ip][event.ID] == nil || peers[ip][event.ID].Location == nil
db.peerLock.RUnlock()
if lookup {
location := db.geodb.Lookup(event.IP)
p.Location = &PeerLocation{
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 := <-handshakeCh:
ip := event.IP.String()
p := getOrInitPeer(network, ip, event.DefaultID)
if p.Handshake == nil {
p.Handshake = []time.Time{event.Handshake}
} else {
p.Handshake = append(p.Handshake, event.Handshake)
}
delete(network.Peers[ip], event.DefaultID)
getOrInitPeer(network, ip, event.ID)
network.Peers[ip][event.ID] = p // interleave instead
// Remove the peer from history in case the metering was before the handshake.
db.peerLock.RLock()
stored := db.peerHistory.Peers[ip] != nil && db.peerHistory.Peers[ip][event.DefaultID] != nil
db.peerLock.RUnlock()
if stored {
db.peerLock.Lock()
hp := getOrInitPeer(db.peerHistory, ip, event.DefaultID)
delete(db.peerHistory.Peers[ip], event.DefaultID)
getOrInitPeer(db.peerHistory, ip, event.ID)
db.peerHistory.Peers[ip][event.ID] = hp // interleave instead
db.peerLock.Unlock()
}
case event := <-disconnectCh:
p := getOrInitPeer(network, event.IP.String(), event.ID)
if p.Disconnected == nil {
p.Disconnected = []time.Time{event.Disconnected}
} else {
p.Disconnected = append(p.Disconnected, event.Disconnected)
}
case event := <-readCh:
// Sum up the ingress between two updates.
p := getOrInitPeer(network, event.IP.String(), event.ID)
if len(p.Ingress) <= 0 {
p.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Ingress)}}
} else {
p.Ingress[0].Value += float64(event.Ingress)
}
case event := <-writeCh:
// Sum up the egress between two updates.
p := getOrInitPeer(network, event.IP.String(), event.ID)
if len(p.Egress) <= 0 {
p.Egress = ChartEntries{&ChartEntry{Value: float64(event.Egress)}}
} else {
p.Egress[0].Value += float64(event.Egress)
}
case <-ticker.C:
now := time.Now()
db.peerLock.Lock()
for ip, peers := range network.Peers {
for id, peer := range peers {
peerHistory := getOrInitPeer(db.peerHistory, ip, id)
if peer.Location != nil {
peerHistory.Location = peer.Location
}
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
//peerHistory.Ingress = ChartEntries{ingress}
} else {
peer.Ingress = ChartEntries{ingress}
peerHistory.Ingress = append(peerHistory.Ingress[1:], ingress)
//peerHistory.Ingress = append(peerHistory.Ingress, 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
//peerHistory.Egress = ChartEntries{egress}
} else {
peer.Egress = ChartEntries{egress}
peerHistory.Egress = append(peerHistory.Egress[1:], egress)
//peerHistory.Egress = append(peerHistory.Egress, egress)
}
}
}
db.peerLock.Unlock()
db.sendToAll(&Message{Network: deepcopy.Copy(network).(*NetworkMessage)})
//fmt.Println()
//s, _ := json.MarshalIndent(network, "", " ")
//fmt.Println(string(s))
for ip, peers := range network.Peers {
for id := range peers {
peers[id] = nil
delete(peers, id)
}
delete(network.Peers, ip)
}
case errc := <-db.quit:
close(quit)
errc <- nil
return
}
}
}
func (db *Dashboard) cleanPeerHistory(quit chan struct{}) {
cleanRate := db.config.Refresh * peerTrafficSampleLimit
for {
select {
case <-time.After(cleanRate):
// clear disconnected
validAfter := time.Now().Add(-cleanRate)
db.peerLock.Lock()
for ip, peers := range db.peerHistory.Peers {
for id, peer := range peers {
if len(peer.Disconnected) > 0 && peer.Disconnected[len(peer.Disconnected)-1].Before(validAfter) {
db.peerHistory.Peers[ip][id].Location = nil
db.peerHistory.Peers[ip][id] = nil
delete(db.peerHistory.Peers[ip], id)
}
}
if len(peers) <= 0 {
delete(db.peerHistory.Peers, ip)
}
}
var lenCount int
for _, peers := range db.peerHistory.Peers {
lenCount += len(peers)
}
if lenCount > p2p.MeteredPeerLimit {
outerLoop:
for ip, peers := range db.peerHistory.Peers {
for id, peer := range peers {
if peer.Disconnected != nil {
db.peerHistory.Peers[ip][id].Location = nil
db.peerHistory.Peers[ip][id] = nil
delete(db.peerHistory.Peers[ip], id)
lenCount--
if lenCount <= p2p.MeteredPeerLimit {
if len(peers) <= 0 {
delete(db.peerHistory.Peers, ip)
}
break outerLoop
}
}
}
if len(peers) <= 0 {
delete(db.peerHistory.Peers, ip)
}
}
}
db.peerLock.Unlock()
case <-quit:
return
}
}
}

View file

@ -21,12 +21,13 @@ package p2p
import (
"net"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/event"
"github.com/ethereum/go-ethereum/metrics"
"github.com/mohae/deepcopy"
"sync"
"sync/atomic"
"time"
"github.com/ethereum/go-ethereum/log"
"sync/atomic"
"fmt"
)
const (
@ -35,9 +36,6 @@ const (
MetricsOutboundTraffic = "p2p/OutboundTraffic"
MetricsOutboundConnects = "p2p/OutboundConnects"
MetricsRegistryIngressPrefix = MetricsInboundTraffic + "/"
MetricsRegistryEgressPrefix = MetricsOutboundTraffic + "/"
MeteredPeerLimit = 16384
)
@ -47,101 +45,79 @@ var (
egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil)
egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil)
PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix)
PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix)
PeerTrafficMeters = newPeerTrafficMeters()
nextDefaultID uint32
NME = &networkMeterEvents{}
)
type PeerMetrics struct {
ID string
type networkMeterEvents struct {
connectFeed event.Feed
handshakeFeed event.Feed
disconnectFeed event.Feed
readFeed event.Feed
writeFeed event.Feed
scope event.SubscriptionScope
defaultID uint64
}
type PeerConnectEvent struct {
IP net.IP
// TODO: -*
Connected *time.Time
Handshake *time.Time
Disconnected *time.Time
Ingress int64
Egress int64
traffic func() (ingress, egress int64)
ID string
Connected time.Time
}
type peerTrafficMeters struct {
peers map[uint]*PeerMetrics
lock sync.RWMutex
type PeerHandshakeEvent struct {
IP net.IP
DefaultID string
ID string
Handshake time.Time
}
func newPeerTrafficMeters() *peerTrafficMeters {
return &peerTrafficMeters{
peers: make(map[uint]*PeerMetrics),
}
type PeerDisconnectEvent struct {
IP net.IP
ID string
Disconnected time.Time
}
func (m *peerTrafficMeters) register(id uint, ip net.IP, traffic func() (ingress, egress int64)) error {
now := time.Now()
peer := &PeerMetrics{
IP: ip,
Connected: &now,
traffic: traffic,
}
m.lock.Lock()
m.peers[id] = peer
m.lock.Unlock()
return nil
type PeerReadEvent struct {
IP net.IP
ID string
Ingress int
}
func (m *peerTrafficMeters) handshakeDone(id uint, peerID string, traffic func() (ingress, egress int64)) {
now := time.Now()
m.lock.Lock()
if peer, ok := m.peers[id]; ok {
peer.Handshake = &now
peer.ID = peerID
peer.traffic = traffic
}
m.lock.Unlock()
type PeerWriteEvent struct {
IP net.IP
ID string
Egress int
}
func (m *peerTrafficMeters) close(id uint) {
now := time.Now()
m.lock.Lock()
m.peers[id].Disconnected = &now
m.lock.Unlock()
func SubscribePeerConnectEvent(ch chan<- PeerConnectEvent) event.Subscription {
return NME.scope.Track(NME.connectFeed.Subscribe(ch))
}
func SubscribePeerHandshakeEvent(ch chan<- PeerHandshakeEvent) event.Subscription {
return NME.scope.Track(NME.handshakeFeed.Subscribe(ch))
}
func SubscribePeerDisconnectEvent(ch chan<- PeerDisconnectEvent) event.Subscription {
return NME.scope.Track(NME.disconnectFeed.Subscribe(ch))
}
func SubscribePeerReadEvent(ch chan<- PeerReadEvent) event.Subscription {
return NME.scope.Track(NME.readFeed.Subscribe(ch))
}
func SubscribePeerWriteEvent(ch chan<- PeerWriteEvent) event.Subscription {
return NME.scope.Track(NME.writeFeed.Subscribe(ch))
}
func (m *peerTrafficMeters) Peers() map[uint]*PeerMetrics {
peers := make(map[uint]*PeerMetrics)
m.lock.Lock()
for id, peer := range m.peers {
peer.Ingress, peer.Egress = peer.traffic()
peers[id] = deepcopy.Copy(peer).(*PeerMetrics)
if peer.Disconnected != nil {
PeerIngressRegistry.Unregister(peer.ID)
PeerEgressRegistry.Unregister(peer.ID)
delete(m.peers, id)
}
}
m.lock.Unlock()
return peers
}
type networkMeter struct {
ingress metrics.Meter
egress metrics.Meter
func closeNME() {
NME.scope.Close()
}
// 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 uint
meter *networkMeter
ingressBeforeHandshake int64
egressBeforeHandshake int64
ip net.IP
id string
lock sync.RWMutex
}
@ -154,8 +130,8 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn {
if !metrics.Enabled {
return conn
}
if len(PeerTrafficMeters.peers) >= MeteredPeerLimit {
log.Warn("Metered peer limit exceeded")
if ip.IsUnspecified() {
log.Warn("peer IP is unspecified")
return conn
}
// Otherwise bump the connection counters and wrap the connection
@ -164,15 +140,17 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn {
} else {
egressConnectMeter.Mark(1)
}
id := uint(atomic.AddUint32(&nextDefaultID, 1))
c := &meteredConn{
id := fmt.Sprintf("peer_%d", atomic.AddUint64(&NME.defaultID, 1))
NME.connectFeed.Send(PeerConnectEvent{
IP: ip,
ID: id,
Connected: time.Now(),
})
return &meteredConn{
Conn: conn,
ip: ip,
id: id,
}
PeerTrafficMeters.register(id, ip, func() (ingress, egress int64) {
return atomic.LoadInt64(&c.ingressBeforeHandshake), atomic.LoadInt64(&c.egressBeforeHandshake)
})
return c
}
// Read delegates a network read to the underlying connection, bumping the ingress
@ -181,12 +159,13 @@ func (c *meteredConn) Read(b []byte) (n int, err error) {
n, err = c.Conn.Read(b)
ingressTrafficMeter.Mark(int64(n))
c.lock.RLock()
if c.meter == nil {
atomic.AddInt64(&c.ingressBeforeHandshake, int64(n))
} else {
c.meter.ingress.Mark(int64(n))
}
id := c.id
c.lock.RUnlock()
NME.readFeed.Send(PeerReadEvent{
IP: c.ip,
ID: id,
Ingress: n,
})
return n, err
}
@ -196,42 +175,37 @@ func (c *meteredConn) Write(b []byte) (n int, err error) {
n, err = c.Conn.Write(b)
egressTrafficMeter.Mark(int64(n))
c.lock.RLock()
if c.meter == nil {
atomic.AddInt64(&c.egressBeforeHandshake, int64(n))
} else {
c.meter.egress.Mark(int64(n))
}
id := c.id
c.lock.RUnlock()
NME.writeFeed.Send(PeerWriteEvent{
IP: c.ip,
ID: id,
Egress: n,
})
return n, err
}
func (c *meteredConn) Close() error {
PeerTrafficMeters.close(c.id)
c.lock.RLock()
id := c.id
c.lock.RUnlock()
NME.disconnectFeed.Send(PeerDisconnectEvent{
IP: c.ip,
ID: id,
Disconnected: time.Now(),
})
return c.Conn.Close()
}
func (c *meteredConn) handshakeDone(peerID string) {
m := &networkMeter{
ingress: metrics.NewRegisteredMeter(peerID, PeerIngressRegistry),
egress: metrics.NewRegisteredMeter(peerID, PeerEgressRegistry),
}
c.lock.Lock()
m.ingress.Mark(atomic.LoadInt64(&c.ingressBeforeHandshake))
m.egress.Mark(atomic.LoadInt64(&c.egressBeforeHandshake))
c.meter = m
defaultID := c.id
c.id = peerID
c.lock.Unlock()
ingressMeter, oki := PeerIngressRegistry.Get(peerID).(metrics.Meter)
egressMeter, oke := PeerEgressRegistry.Get(peerID).(metrics.Meter)
traffic := func() (ingress, egress int64) {
return 0, 0
}
if oki && oke {
traffic = func() (ingress, egress int64) {
return ingressMeter.Count(), egressMeter.Count()
}
} else {
log.Warn("Failed to get traffic meter", "peerID", peerID)
}
PeerTrafficMeters.handshakeDone(c.id, peerID, traffic)
NME.handshakeFeed.Send(PeerHandshakeEvent{
IP: c.ip,
DefaultID: defaultID,
ID: peerID,
Handshake: time.Now(),
})
}

View file

@ -388,6 +388,7 @@ func (srv *Server) Stop() {
close(srv.quit)
srv.lock.Unlock()
srv.loopWG.Wait()
closeNME()
}
// sharedUDPConn implements a shared connection. Write sends messages to the underlying connection while read returns