dashboard, p2p: experiment 2

This commit is contained in:
Kurkó Mihály 2018-09-18 19:02:05 +03:00
parent cac8611230
commit a87086b608
7 changed files with 321 additions and 353 deletions

View file

@ -11363,7 +11363,7 @@ var _bundleJs = []byte((((((((((`!function(modules) {
key: "render", key: "render",
value: function() { value: function() {
var _this2 = this; var _this2 = this;
return _react2.default.createElement(_Table2.default, null, _react2.default.createElement(_Table.TableHead, null, _react2.default.createElement(_Table.TableRow, null, _react2.default.createElement(_Table.TableCell, null, "IP"), _react2.default.createElement(_Table.TableCell, null, "Location"), _react2.default.createElement(_Table.TableCell, null, "Peer ID"), _react2.default.createElement(_Table.TableCell, null, "Ingress"), _react2.default.createElement(_Table.TableCell, null, "Egress"), _react2.default.createElement(_Table.TableCell, null, "Connected"), _react2.default.createElement(_Table.TableCell, null, "Handshake"), _react2.default.createElement(_Table.TableCell, null, "Disconnected"))), _react2.default.createElement(_Table.TableBody, null, Object.entries(this.props.content.peerBundles).map(function(_ref4) { return _react2.default.createElement(_Table2.default, null, _react2.default.createElement(_Table.TableHead, null, _react2.default.createElement(_Table.TableRow, null, _react2.default.createElement(_Table.TableCell, null, "IP"), _react2.default.createElement(_Table.TableCell, null, "Location"), _react2.default.createElement(_Table.TableCell, null, "Peer ID"), _react2.default.createElement(_Table.TableCell, null, "Ingress"), _react2.default.createElement(_Table.TableCell, null, "Egress"), _react2.default.createElement(_Table.TableCell, null, "Connected"), _react2.default.createElement(_Table.TableCell, null, "Handshake"), _react2.default.createElement(_Table.TableCell, null, "Time"))), _react2.default.createElement(_Table.TableBody, null, Object.entries(this.props.content.peerBundles).map(function(_ref4) {
var _ref5 = _slicedToArray(_ref4, 2), ip = _ref5[0], bundle = _ref5[1]; var _ref5 = _slicedToArray(_ref4, 2), ip = _ref5[0], bundle = _ref5[1];
return _react2.default.createElement(_Table.TableRow, { return _react2.default.createElement(_Table.TableRow, {
key: ip key: ip

View file

@ -46,12 +46,6 @@ export const inserter = (update: {[string]: PeerBundle}, prev: {[string]: PeerBu
prev[ip].peers[id] = u; prev[ip].peers[id] = u;
return; return;
} }
// If the handshake was between two metering
if (u.defaultID && prev[ip].peers[u.defaultID]) {
// TODO (kurkomisi): merge the two in order to keep the previous connection.
prev[ip].peers[id] = prev[ip].peers[u.defaultID];
delete prev[ip].peers[u.defaultID];
}
const p: Peer = prev[ip].peers[id]; const p: Peer = prev[ip].peers[id];
if (u.connected) { if (u.connected) {
if (!Array.isArray(p.connected)) { if (!Array.isArray(p.connected)) {
@ -59,12 +53,6 @@ export const inserter = (update: {[string]: PeerBundle}, prev: {[string]: PeerBu
} }
p.connected = [...p.connected, ...u.connected]; 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 (u.disconnected) {
if (!Array.isArray(p.disconnected)) { if (!Array.isArray(p.disconnected)) {
p.disconnected = []; p.disconnected = [];
@ -126,12 +114,11 @@ class Network extends Component<Props, State> {
<TableCell>Ingress</TableCell> <TableCell>Ingress</TableCell>
<TableCell>Egress</TableCell> <TableCell>Egress</TableCell>
<TableCell>Connected</TableCell> <TableCell>Connected</TableCell>
<TableCell>Handshake</TableCell>
<TableCell>Disconnected</TableCell> <TableCell>Disconnected</TableCell>
</TableRow> </TableRow>
</TableHead> </TableHead>
<TableBody> <TableBody>
{Object.entries(this.props.content.peerBundles).map(([ip, bundle]) => ( {Object.entries(this.props.content.peerBundles).map(([ip, bundle]) => { console.log(ip, bundle); return (
<TableRow key={ip}> <TableRow key={ip}>
<TableCell>{ip}</TableCell> <TableCell>{ip}</TableCell>
<TableCell> <TableCell>
@ -141,25 +128,22 @@ class Network extends Component<Props, State> {
})() : ''} })() : ''}
</TableCell> </TableCell>
<TableCell> <TableCell>
{Object.keys(bundle.peers).map(id => id.substring(0, 10)).join(' ')} {bundle.peers && Object.keys(bundle.peers).map(id => id.substring(0, 10)).join(' ')}
</TableCell> </TableCell>
<TableCell> <TableCell>
{Object.values(bundle.peers).map(peer => peer.ingress && peer.ingress.map(sample => sample.value).join(' ')).join(', ')} {bundle.peers && Object.values(bundle.peers).map(peer => peer.ingress && peer.ingress.map(sample => sample.value).join(' ')).join(', ')}
</TableCell> </TableCell>
<TableCell> <TableCell>
{Object.values(bundle.peers).map(peer => peer.egress && peer.egress.map(sample => sample.value).join(' ')).join(', ')} {bundle.peers && Object.values(bundle.peers).map(peer => peer.egress && peer.egress.map(sample => sample.value).join(' ')).join(', ')}
</TableCell> </TableCell>
<TableCell> <TableCell>
{Object.values(bundle.peers).map(peer => peer.connected && peer.connected.map(time => this.formatTime(time)).join(' ')).join(', ')} {bundle.peers && Object.values(bundle.peers).map(peer => peer.connected && peer.connected.map(time => this.formatTime(time)).join(' ')).join(', ')}
</TableCell> </TableCell>
<TableCell> <TableCell>
{Object.values(bundle.peers).map(peer => peer.handshake && peer.handshake.map(time => this.formatTime(time)).join(' ')).join(', ')} {bundle.peers && Object.values(bundle.peers).map(peer => peer.disconnected && peer.disconnected.map(time => this.formatTime(time)).join(' ')).join(', ')}
</TableCell>
<TableCell>
{Object.values(bundle.peers).map(peer => peer.disconnected && peer.disconnected.map(time => this.formatTime(time)).join(' ')).join(', ')}
</TableCell> </TableCell>
</TableRow> </TableRow>
))} )})}
</TableBody> </TableBody>
</Table> </Table>
); );

View file

@ -24,7 +24,7 @@ module.exports = merge(common, {
new webpack.HotModuleReplacementPlugin(), new webpack.HotModuleReplacementPlugin(),
], ],
// devtool: 'eval', // devtool: 'eval',
devtool: 'inline-source-map', devtool: 'source-map',
devServer: { devServer: {
port: 8081, port: 8081,
hot: true, hot: true,

View file

@ -66,7 +66,8 @@ type NetworkMessage struct {
func (m *NetworkMessage) getOrInitBundle(ip string) *PeerBundle { func (m *NetworkMessage) getOrInitBundle(ip string) *PeerBundle {
if _, ok := m.PeerBundles[ip]; !ok { if _, ok := m.PeerBundles[ip]; !ok {
m.PeerBundles[ip] = &PeerBundle{ m.PeerBundles[ip] = &PeerBundle{
Peers: make(map[string]*Peer), Peers: make(PeerMap),
FailedPeers: make(PeerMap),
} }
} }
return m.PeerBundles[ip] return m.PeerBundles[ip]
@ -78,18 +79,40 @@ func (m *NetworkMessage) getOrInitPeer(ip, id string) *Peer {
return m.getOrInitBundle(ip).getOrInitPeer(id) return m.getOrInitBundle(ip).getOrInitPeer(id)
} }
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)
}
// PeerBundle contains information about the peers pertaining to an IP address. // PeerBundle contains information about the peers pertaining to an IP address.
type PeerBundle struct { type PeerBundle struct {
Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP Location *GeoLocation `json:"location,omitempty"` // geographical information based on IP
Peers map[string]*Peer `json:"peers,omitempty"` // the peers' node id is used as key Peers PeerMap `json:"peers,omitempty"` // the peers' node id is used as key
FailedPeers []*Peer `json:"failedPeers,omitempty"` FailedPeers PeerMap `json:"failedPeers,omitempty"`
} }
func (b *PeerBundle) getOrInitPeer(id string) *Peer { func (b *PeerBundle) getOrInitPeer(id string) *Peer {
if _, ok := b.Peers[id]; !ok { return b.Peers.getOrInit(id)
b.Peers[id] = new(Peer)
} }
return b.Peers[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. // GeoLocation contains geographical information.
@ -108,9 +131,8 @@ type Peer struct {
Ingress ChartEntries `json:"ingress,omitempty"` Ingress ChartEntries `json:"ingress,omitempty"`
Egress ChartEntries `json:"egress,omitempty"` Egress ChartEntries `json:"egress,omitempty"`
DefaultID string `json:"defaultID,omitempty"`
element *list.Element element *list.Element
ip, id string
} }
// SystemMessage contains the metered system data samples. // SystemMessage contains the metered system data samples.

View file

@ -18,9 +18,7 @@ package dashboard
import ( import (
"container/list" "container/list"
"encoding/json"
"fmt" "fmt"
"github.com/mohae/deepcopy"
"time" "time"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
@ -29,6 +27,58 @@ import (
const eventBufferLimit = 128 // Maximum number of buffered peer events const eventBufferLimit = 128 // Maximum number of buffered peer events
const trafficEventBufferLimit = p2p.MeteredPeerLimit const trafficEventBufferLimit = p2p.MeteredPeerLimit
const connectionLimit = 100
var autoID int64
type peerLimiter struct {
underlying *NetworkMessage
failed bool
l *list.List
}
func NewPeerLimiter(underlying *NetworkMessage, failed bool) *peerLimiter {
return &peerLimiter{l: list.New(), underlying: underlying, failed: failed}
}
func (pl *peerLimiter) update(peer *Peer) {
return
if peer.element == nil {
peer.element = pl.l.PushBack(peer)
} else {
pl.l.MoveToBack(peer.element)
}
for pl.l.Len() > 2 {//p2p.MeteredPeerLimit {
pl.remove(pl.l.Front())
}
}
func (pl *peerLimiter) remove(e *list.Element) {
return
elem := pl.l.Remove(e)
if peer, ok := elem.(*Peer); ok {
if pl.failed {
fmt.Println(peer.ip, peer.id)
pl.underlying.PeerBundles[peer.ip].FailedPeers.remove(peer.id)
} else {
fmt.Println(peer.ip, peer.id[:10])
pl.underlying.PeerBundles[peer.ip].Peers.remove(peer.id)
}
}
}
func (pl *peerLimiter) clear() {
for pl.l.Front() != nil {
pl.remove(pl.l.Front())
}
}
//
//func tail(arr []time.Time) []time.Time {
// if first := len(arr)-connectionLimit; first > 0 {
// return arr[first:]
// }
// return arr
//}
// collectPeerData gathers data about the peers and sends it to the clients. // collectPeerData gathers data about the peers and sends it to the clients.
func (db *Dashboard) collectPeerData() { func (db *Dashboard) collectPeerData() {
@ -45,148 +95,154 @@ func (db *Dashboard) collectPeerData() {
var ( var (
// Peer event channels. // Peer event channels.
connectCh = make(chan p2p.PeerConnectEvent, eventBufferLimit) peerCh = make(chan p2p.MeteredPeerEvent, eventBufferLimit)
failedCh = make(chan p2p.PeerFailedEvent, eventBufferLimit)
disconnectCh = make(chan p2p.PeerDisconnectEvent, eventBufferLimit)
ingressCh = make(chan p2p.PeerTrafficEvent, trafficEventBufferLimit)
egressCh = make(chan p2p.PeerTrafficEvent, trafficEventBufferLimit)
// Subscribe to peer events. // Subscribe to peer events.
subConnect = p2p.SubscribePeerConnectEvent(connectCh) subPeer = p2p.SubscribePeerEvent(peerCh)
subFailed = p2p.SubscribePeerFailedEvent(failedCh)
subDisconnect = p2p.SubscribePeerDisconnectEvent(disconnectCh)
subIngress = p2p.SubscribePeerIngressEvent(ingressCh)
subEgress = p2p.SubscribePeerEgressEvent(egressCh)
) )
defer func() { defer func() {
// Unsubscribe at the end. // Unsubscribe at the end.
subConnect.Unsubscribe() subPeer.Unsubscribe()
subFailed.Unsubscribe()
subDisconnect.Unsubscribe()
subIngress.Unsubscribe()
subEgress.Unsubscribe()
}() }()
ticker := time.NewTicker(db.config.Refresh) ticker := time.NewTicker(db.config.Refresh)
defer ticker.Stop() defer ticker.Stop()
purgeOrder := list.New() db.peerLock.RLock()
failedPurgeOrder := list.New() //historyPeerLimiter := NewPeerLimiter(db.history.Network, false)
update := func(peer *Peer, l *list.List) { //historyFailedPeerLimiter := NewPeerLimiter(db.history.Network, true)
if peer.element == nil { //db.peerLock.RUnlock()
peer.element = purgeOrder.PushBack(peer) //// Listen for events, and prepare the difference between two metering.
} else { //diff := &NetworkMessage{
purgeOrder.MoveToBack(peer.element) // PeerBundles: make(map[string]*PeerBundle),
} //}
} //// Needed in order to keep the limit in the diff
// Listen for events, and prepare the difference between two metering. //diffPeerLimiter := NewPeerLimiter(diff, false)
diff := &NetworkMessage{ //diffFailedPeerLimiter := NewPeerLimiter(diff, true)
PeerBundles: make(map[string]*PeerBundle),
}
for { for {
select { select {
case event := <-connectCh: case event := <-peerCh:
diffBundle := diff.getOrInitBundle(event.IP) fmt.Println(event)
diffBundle.Location = db.geodb.Location(event.IP) //case event := <-connectCh:
diffPeer := diffBundle.getOrInitPeer(event.ID) // diffBundle := diff.getOrInitBundle(event.IP)
diffPeer.Connected = append(diffPeer.Connected, event.Connected) // diffBundle.Location = db.geodb.Location(event.IP)
case event := <-failedCh: // diffPeer := diffBundle.getOrInitPeer(event.ID)
diffBundle := diff.getOrInitBundle(event.IP) // diffPeer.Connected = append(diffPeer.Connected, event.Connected)
diffBundle.Location = db.geodb.Location(event.IP) // if first := len(diffPeer.Connected)-connectionLimit; first > 0 {
diffBundle.FailedPeers = append(diffBundle.FailedPeers, &Peer{ // diffPeer.Connected = diffPeer.Connected[first:]
Connected: []time.Time{event.Connected}, // }
Disconnected: []time.Time{event.Disconnected}, // diffPeer.ip = event.IP
}) // diffPeer.id = event.ID
case event := <-disconnectCh: // diffPeerLimiter.update(diffPeer)
diffPeer := diff.getOrInitPeer(event.IP, event.ID) //case event := <-disconnectCh:
diffPeer.Disconnected = append(diffPeer.Disconnected, event.Disconnected) // diffPeer := diff.getOrInitPeer(event.IP, event.ID)
case event := <-ingressCh: // diffPeer.Disconnected = append(diffPeer.Disconnected, event.Time)
diffPeer := diff.getOrInitPeer(event.IP, event.ID) // if first := len(diffPeer.Connected)-connectionLimit; first > 0 {
if len(diffPeer.Ingress) != 1 { // diffPeer.Connected = diffPeer.Connected[first:]
diffPeer.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} // }
} else { // diffPeer.ip = event.IP
diffPeer.Ingress[0].Value = float64(event.Amount) // diffPeer.id = event.ID
} // diffPeerLimiter.update(diffPeer)
case event := <-egressCh: //case event := <-ingressCh:
diffPeer := diff.getOrInitPeer(event.IP, event.ID) // diffPeer := diff.getOrInitPeer(event.IP, event.ID)
if len(diffPeer.Egress) != 1 { // if len(diffPeer.Ingress) != 1 {
diffPeer.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}} // diffPeer.Ingress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}}
} else { // } else {
diffPeer.Egress[0].Value = float64(event.Amount) // diffPeer.Ingress[0].Value = float64(event.Amount)
} // }
case <-ticker.C: // diffPeer.ip = event.IP
now := time.Now() // diffPeer.id = event.ID
// Merge the diff with the history. // diffPeerLimiter.update(diffPeer)
db.peerLock.Lock() //case event := <-egressCh:
for ip, diffBundle := range diff.PeerBundles { // diffPeer := diff.getOrInitPeer(event.IP, event.ID)
historyBundle := db.history.Network.getOrInitBundle(ip) // if len(diffPeer.Egress) != 1 {
historyBundle.Location = diffBundle.Location // diffPeer.Egress = ChartEntries{&ChartEntry{Value: float64(event.Amount)}}
for id, diffPeer := range diffBundle.Peers { // } else {
historyPeer := historyBundle.getOrInitPeer(id) // diffPeer.Egress[0].Value = float64(event.Amount)
historyPeer.Connected = append(historyPeer.Connected, diffPeer.Connected...) // }
historyPeer.Disconnected = append(historyPeer.Disconnected, diffPeer.Disconnected...) // diffPeer.ip = event.IP
if len(diffPeer.Ingress) == 1 { // diffPeer.id = event.ID
diffPeer.Ingress[0].Time = now // diffPeerLimiter.update(diffPeer)
if historyPeer.Ingress == nil { //case event := <-failedCh:
historyPeer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), sampleLimit-1, db.config.Refresh), diffPeer.Ingress[0]) // diffBundle := diff.getOrInitBundle(event.IP)
// The first message about a diffPeer should contain the whole list // diffBundle.Location = db.geodb.Location(event.IP)
diffPeer.Ingress = historyPeer.Ingress // id := fmt.Sprintf("peer_%d", atomic.AddInt64(&autoID, 1))
} else { // failedPeer := diffBundle.FailedPeers.getOrInit(id)
historyPeer.Ingress = append(historyPeer.Ingress, diffPeer.Ingress[0])[1:] // failedPeer.Connected = []time.Time{event.Connected}
} // failedPeer.Disconnected = []time.Time{event.Disconnected}
} // failedPeer.ip = event.IP
if len(diffPeer.Egress) == 1 { // failedPeer.id = id
diffPeer.Egress[0].Time = now // diffFailedPeerLimiter.update(failedPeer)
if historyPeer.Egress == nil { //case <-ticker.C:
historyPeer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), sampleLimit-1, db.config.Refresh), diffPeer.Egress[0]) // now := time.Now()
// The first message about a diffPeer should contain the whole list // // Merge the diff with the history.
diffPeer.Egress = historyPeer.Egress // db.peerLock.Lock()
} else { // for ip, diffBundle := range diff.PeerBundles {
historyPeer.Egress = append(historyPeer.Egress, diffPeer.Egress[0])[1:] // historyBundle := db.history.Network.getOrInitBundle(ip)
} // historyBundle.Location = diffBundle.Location
} // for id, diffPeer := range diffBundle.Peers {
update(historyPeer, purgeOrder) // historyPeer := historyBundle.getOrInitPeer(id)
} // historyPeer.Connected = append(historyPeer.Connected, diffPeer.Connected...)
historyBundle.FailedPeers = append(historyBundle.FailedPeers, diffBundle.FailedPeers...) // if first := len(historyPeer.Connected)-connectionLimit; first > 0 {
for _, fp := range diffBundle.FailedPeers { // historyPeer.Connected = historyPeer.Connected[first:]
update(fp, failedPurgeOrder) // }
} // historyPeer.Disconnected = append(historyPeer.Disconnected, diffPeer.Disconnected...)
} // if first := len(historyPeer.Disconnected)-connectionLimit; first > 0 {
for purgeOrder.Len() > p2p.MeteredPeerLimit { // historyPeer.Disconnected = historyPeer.Disconnected[first:]
purgeOrder.Remove(purgeOrder.Front()) // }
} // if len(diffPeer.Ingress) == 1 {
for failedPurgeOrder.Len() > p2p.MeteredPeerLimit { // diffPeer.Ingress[0].Time = now
failedPurgeOrder.Remove(failedPurgeOrder.Front()) // if historyPeer.Ingress == nil {
} // historyPeer.Ingress = append(emptyChartEntries(now.Add(-db.config.Refresh), 3/*sampleLimit-1*/, db.config.Refresh), diffPeer.Ingress[0])
// // The first message about a diffPeer should contain the whole list
//ss, _ := json.MarshalIndent(db.history.Network, "", " ") // diffPeer.Ingress = historyPeer.Ingress
//fmt.Println(string(ss)) // } else {
db.peerLock.Unlock() // historyPeer.Ingress = append(historyPeer.Ingress, diffPeer.Ingress[0])[1:]
// Send the diff to the clients. // }
db.sendToAll(&Message{Network: deepcopy.Copy(diff).(*NetworkMessage)}) // }
//s, _ := json.MarshalIndent(diff, "", " ") // if len(diffPeer.Egress) == 1 {
//fmt.Println(string(s)) // diffPeer.Egress[0].Time = now
// if historyPeer.Egress == nil {
s, _ := json.MarshalIndent(diff, "", " ") // historyPeer.Egress = append(emptyChartEntries(now.Add(-db.config.Refresh), 3/*sampleLimit-1*/, db.config.Refresh), diffPeer.Egress[0])
fmt.Println(string(s)) // // The first message about a diffPeer should contain the whole list
// Prepare for the next metering, clear the diff variable. // diffPeer.Egress = historyPeer.Egress
diff = &NetworkMessage{ // } else {
PeerBundles: make(map[string]*PeerBundle), // historyPeer.Egress = append(historyPeer.Egress, diffPeer.Egress[0])[1:]
} // }
case err := <-subConnect.Err(): // }
log.Warn("Peer connect subscription error", "err", err) // historyPeer.ip = diffPeer.ip
return // historyPeer.id = diffPeer.id
case err := <-subFailed.Err(): // historyPeerLimiter.update(historyPeer)
log.Warn("Peer failed subscription error", "err", err) // }
return // for id, diffFailedPeer := range diffBundle.FailedPeers {
case err := <-subDisconnect.Err(): // historyFailedPeer := historyBundle.getOrInitFailedPeer(id)
log.Warn("Peer disconnect subscription error", "err", err) // historyFailedPeer.Connected = diffFailedPeer.Connected
return // historyFailedPeer.Disconnected = diffFailedPeer.Disconnected
case err := <-subIngress.Err(): // historyFailedPeer.ip = diffFailedPeer.ip
log.Warn("Peer ingress subscription error", "err", err) // historyFailedPeer.id = diffFailedPeer.id
return // historyFailedPeerLimiter.update(historyFailedPeer)
case err := <-subEgress.Err(): // }
log.Warn("Peer egress subscription error", "err", err) // }
// //for elem := historyPeerLimiter.l.Front(); elem != nil; elem = elem.Next() {
// // s, _ := json.MarshalIndent(elem.Value, "", " ")
// // fmt.Println(string(s))
// //}
// //fmt.Println()
// db.peerLock.Unlock()
//
// //s, _ := json.MarshalIndent(deepcopy.Copy(diff), "", " ")
// //fmt.Println(string(s))
// //fmt.Println()
// // Send the diff to the clients.
// db.sendToAll(&Message{Network: deepcopy.Copy(diff).(*NetworkMessage)})
//
// // Prepare for the next metering, clear the diff variable.
// diffPeerLimiter.clear()
// diffFailedPeerLimiter.clear()
// diff = &NetworkMessage{
// PeerBundles: make(map[string]*PeerBundle),
// }
case err := <-subPeer.Err():
log.Warn("Peer subscription error", "err", err)
return return
case errc := <-db.quit: case errc := <-db.quit:
errc <- nil errc <- nil

View file

@ -19,10 +19,8 @@
package p2p package p2p
import ( import (
"net"
"strings"
"fmt" "fmt"
"net"
"sync" "sync"
"sync/atomic" "sync/atomic"
"time" "time"
@ -46,130 +44,47 @@ const (
) )
var ( var (
ingressConnectMeter = metrics.NewRegisteredMeter(MetricsInboundConnects, nil) // meter counting the ingress connections ingressConnectMeter = metrics.NewRegisteredMeter(MetricsInboundConnects, nil) // Meter counting the ingress connections
ingressTrafficMeter = metrics.NewRegisteredMeter(MetricsInboundTraffic, nil) // meter metering the cumulative ingress traffic ingressTrafficMeter = metrics.NewRegisteredMeter(MetricsInboundTraffic, nil) // Meter metering the cumulative ingress traffic
egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) // meter counting the egress connections egressConnectMeter = metrics.NewRegisteredMeter(MetricsOutboundConnects, nil) // Meter counting the egress connections
egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) // meter metering the cumulative egress traffic egressTrafficMeter = metrics.NewRegisteredMeter(MetricsOutboundTraffic, nil) // Meter metering the cumulative egress traffic
PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) PeerIngressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryIngressPrefix) // Registry containing the peer ingress
PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) PeerEgressRegistry = metrics.NewPrefixedChildRegistry(metrics.DefaultRegistry, MetricsRegistryEgressPrefix) // Registry containing the peer egress
metricsFeed = new(peerMetricsFeed) // Peer event feed for metrics metricsFeed event.Feed // Event feed for peer metrics
meteredPeerCount uint64 // Actually stored peer connection count
meteredPeerCount uint64
) )
// peerMetricsFeed delivers the peer metrics to the subscribed channels. // MeteredPeerEventType is the type of peer events emitted by a metered connection.
type peerMetricsFeed struct { type MeteredPeerEventType int
connect event.Feed // Event feed to notify the connection and the successful handshake of a peer
ingress event.Feed // Event feed to notify the amount of read bytes of a peer
egress event.Feed // Event feed to notify the amount of written bytes of a peer
disconnect event.Feed // Event feed to notify the disconnection of a peer
failed event.Feed // Event feed to notify the connection of a peer and its disconnection before the handshake
scope event.SubscriptionScope // Facility to unsubscribe all the subscriptions at once const (
// PeerConnected is the type of event emitted when a peer successfully
// made the handshake.
PeerConnected MeteredPeerEventType = iota
quit chan chan error // PeerDisconnected is the type of event emitted when a peer disconnects.
PeerDisconnected
// PeerHandshakeFailed is the type of event emitted when a peer fails to
// make the handshake or disconnects before the handshake.
PeerHandshakeFailed
)
// MeteredPeerEvent is an event emitted when peers connect or disconnect
type MeteredPeerEvent struct {
Type MeteredPeerEventType // Type of peer event
IP net.IP // IP address of the peer
ID string // NodeID of the peer
Elapsed time.Duration // Time elapsed between the connection and the handshake/disconnection
Ingress uint64 // Ingress count in the moment of disconnection
Egress uint64 // Egress count in the moment of disconnection
} }
// PeerConnectEvent contains information about the connection of a peer. // SubscribePeerEvent registers a subscription of PeerEvent
type PeerConnectEvent struct { func SubscribePeerEvent(ch chan<- MeteredPeerEvent) event.Subscription {
IP string return metricsFeed.Subscribe(ch)
ID string
Connected time.Time
}
// PeerDisconnectEvent contains information about the disconnection of a peer.
type PeerDisconnectEvent struct {
IP string
ID string
Disconnected time.Time
}
type PeerTrafficEvent struct {
IP string
ID string
Amount int64
}
type PeerFailedEvent struct {
IP string
Connected time.Time
Disconnected time.Time
}
// SubscribePeerConnectEvent registers a subscription of PeerConnectEvent
func SubscribePeerConnectEvent(ch chan<- PeerConnectEvent) event.Subscription {
return metricsFeed.scope.Track(metricsFeed.connect.Subscribe(ch))
}
// SubscribePeerDisconnectEvent registers a subscription of PeerDisconnectEvent
func SubscribePeerDisconnectEvent(ch chan<- PeerDisconnectEvent) event.Subscription {
return metricsFeed.scope.Track(metricsFeed.disconnect.Subscribe(ch))
}
// SubscribePeerTrafficEvent registers a subscription of PeerTrafficEvent
func SubscribePeerIngressEvent(ch chan<- PeerTrafficEvent) event.Subscription {
return metricsFeed.scope.Track(metricsFeed.ingress.Subscribe(ch))
}
func SubscribePeerEgressEvent(ch chan<- PeerTrafficEvent) event.Subscription {
return metricsFeed.scope.Track(metricsFeed.egress.Subscribe(ch))
}
// SubscribePeerFailedEvent registers a subscription of PeerFailedEvent
func SubscribePeerFailedEvent(ch chan<- PeerFailedEvent) event.Subscription {
return metricsFeed.scope.Track(metricsFeed.failed.Subscribe(ch))
}
func runMetricsFeedHelper(refresh time.Duration) {
metricsFeed.quit = make(chan chan error)
ticker := time.NewTicker(refresh)
defer ticker.Stop()
// It is possible to send all of the traffic events together, but it is risky to use pointers in the events.
trafficEventSender := func(prefix string, feed *event.Feed) func(name string, i interface{}) {
return func(name string, i interface{}) {
if m, ok := i.(metrics.Meter); ok {
// Trim the common prefix and split the peer specific part in order to get the ip and the node id.
if key := strings.Split(strings.TrimPrefix(name, prefix), "/"); len(key) == 2 {
feed.Send(PeerTrafficEvent{
IP: key[0],
ID: key[1],
Amount: m.Count(),
})
} else {
log.Warn("Invalid peer metrics name", "name", name)
}
}
}
}
sendIngress := trafficEventSender(MetricsRegistryIngressPrefix, &metricsFeed.ingress)
sendEgress := trafficEventSender(MetricsRegistryEgressPrefix, &metricsFeed.egress)
for {
select {
case <-ticker.C:
PeerIngressRegistry.Each(sendIngress)
PeerEgressRegistry.Each(sendEgress)
case errc := <-metricsFeed.quit:
errc <- nil
return
}
}
}
// closeMetricsFeed closes all the tracked subscriptions.
func closeMetricsFeed() {
if metricsFeed.quit != nil {
errc := make(chan error)
metricsFeed.quit <- errc
<-errc
}
metricsFeed.scope.Close()
PeerIngressRegistry.UnregisterAll()
PeerEgressRegistry.UnregisterAll()
} }
// meteredConn is a wrapper around a net.Conn that meters both the // meteredConn is a wrapper around a net.Conn that meters both the
@ -177,18 +92,19 @@ func closeMetricsFeed() {
type meteredConn struct { type meteredConn struct {
net.Conn // Network connection to wrap with metering net.Conn // Network connection to wrap with metering
connected time.Time connected time.Time // Connection time of the peer
ip string // The IP address of the peer ip net.IP // IP address of the peer
id string // The NodeID of the peer id string // NodeID of the peer
ingressMeter metrics.Meter ingressMeter metrics.Meter // Meter for the read bytes of the peer
egressMeter metrics.Meter egressMeter metrics.Meter // Meter for the written bytes of the peer
lock sync.RWMutex // Lock protecting the metered connection's internals lock sync.RWMutex // Lock protecting the metered connection's internals
} }
// newMeteredConn creates a new metered connection, also bumping the ingress or // newMeteredConn creates a new metered connection, bumps the ingress or egress
// egress connection meter. If the metrics system is disabled, this function // connection meter and also increases the metered peer count. If the metrics
// returns the original object. // system is disabled, the IP address is unspecified or the metered peer count
// reached the limit, this function returns the original object.
func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn { func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn {
// Short circuit if metrics are disabled // Short circuit if metrics are disabled
if !metrics.Enabled { if !metrics.Enabled {
@ -212,13 +128,13 @@ func newMeteredConn(conn net.Conn, ingress bool, ip net.IP) net.Conn {
} }
return &meteredConn{ return &meteredConn{
Conn: conn, Conn: conn,
ip: ip.String(), ip: ip,
connected: time.Now(), connected: time.Now(),
} }
} }
// Read delegates a network read to the underlying connection, bumping the ingress // Read delegates a network read to the underlying connection, bumping the common
// traffic meter along the way. // and the peer ingress traffic meters along the way.
func (c *meteredConn) Read(b []byte) (n int, err error) { func (c *meteredConn) Read(b []byte) (n int, err error) {
n, err = c.Conn.Read(b) n, err = c.Conn.Read(b)
ingressTrafficMeter.Mark(int64(n)) ingressTrafficMeter.Mark(int64(n))
@ -230,8 +146,8 @@ func (c *meteredConn) Read(b []byte) (n int, err error) {
return n, err return n, err
} }
// Write delegates a network write to the underlying connection, bumping the // Write delegates a network write to the underlying connection, bumping the common
// egress traffic meter along the way. // and the peer egress traffic meters along the way.
func (c *meteredConn) Write(b []byte) (n int, err error) { func (c *meteredConn) Write(b []byte) (n int, err error) {
n, err = c.Conn.Write(b) n, err = c.Conn.Write(b)
egressTrafficMeter.Mark(int64(n)) egressTrafficMeter.Mark(int64(n))
@ -243,53 +159,9 @@ func (c *meteredConn) Write(b []byte) (n int, err error) {
return n, err return n, err
} }
// Close closes the underlying connection. // handshakeDone is called when a peer handshake is done. Registers the peer to
func (c *meteredConn) Close() error { // the ingress and the egress traffic registries using the peer's IP and NodeID,
// Decrement the metered peer count // also emits connect event.
atomic.AddUint64(&meteredPeerCount, ^uint64(0))
err, now := c.Conn.Close(), time.Now()
c.lock.RLock()
ip, id := c.ip, c.id
c.lock.RUnlock()
// If the peer disconnects before the handshake
if id == "" {
metricsFeed.failed.Send(PeerFailedEvent{
IP: ip,
Connected: c.connected,
Disconnected: now,
})
return err
}
c.lock.RLock()
//ingress, egress := c.ingressMeter.Count(), c.egressMeter.Count()
c.lock.RUnlock()
// Unregister the peer from the metrics registry
key := fmt.Sprintf("%s/%s", ip, id)
PeerIngressRegistry.Unregister(key)
PeerEgressRegistry.Unregister(key)
//metricsFeed.ingress.Send(PeerTrafficEvent{
// IP: ip,
// ID: id,
// Amount: ingress,
//})
//metricsFeed.egress.Send(PeerTrafficEvent{
// IP: ip,
// ID: id,
// Amount: egress,
//})
metricsFeed.disconnect.Send(PeerDisconnectEvent{
IP: ip,
ID: id,
Disconnected: now,
})
return err
}
// handshakeDone changes the default id to the peer's node id.
func (c *meteredConn) handshakeDone(id discover.NodeID) { func (c *meteredConn) handshakeDone(id discover.NodeID) {
c.lock.Lock() c.lock.Lock()
c.id = id.String() c.id = id.String()
@ -298,9 +170,45 @@ func (c *meteredConn) handshakeDone(id discover.NodeID) {
c.egressMeter = metrics.NewRegisteredMeter(key, PeerEgressRegistry) c.egressMeter = metrics.NewRegisteredMeter(key, PeerEgressRegistry)
c.lock.Unlock() c.lock.Unlock()
metricsFeed.connect.Send(PeerConnectEvent{ metricsFeed.Send(MeteredPeerEvent{
Type: PeerConnected,
IP: c.ip, IP: c.ip,
ID: id.String(), ID: id.String(),
Connected: c.connected, Elapsed: time.Now().Sub(c.connected),
}) })
} }
// Close delegates a close operation to the underlying connection, unregisters
// the peer from the traffic registries and emits close event.
func (c *meteredConn) Close() error {
// Decrement the metered peer count
atomic.AddUint64(&meteredPeerCount, ^uint64(0))
c.lock.RLock()
// If the peer disconnects before the handshake
if c.id == "" {
c.lock.RUnlock()
metricsFeed.Send(MeteredPeerEvent{
Type: PeerHandshakeFailed,
IP: c.ip,
Elapsed: time.Now().Sub(c.connected),
})
return c.Conn.Close()
}
id, ingress, egress := c.id, uint64(c.ingressMeter.Count()), uint64(c.egressMeter.Count())
c.lock.RUnlock()
// Unregister the peer from the traffic registries
key := fmt.Sprintf("%s/%s", c.ip, id)
PeerIngressRegistry.Unregister(key)
PeerEgressRegistry.Unregister(key)
metricsFeed.Send(MeteredPeerEvent{
Type: PeerDisconnected,
IP: c.ip,
ID: id,
Ingress: ingress,
Egress: egress,
})
return c.Conn.Close()
}

View file

@ -387,7 +387,6 @@ func (srv *Server) Stop() {
} }
close(srv.quit) close(srv.quit)
srv.lock.Unlock() srv.lock.Unlock()
closeMetricsFeed()
srv.loopWG.Wait() srv.loopWG.Wait()
} }
@ -542,7 +541,6 @@ func (srv *Server) Start() (err error) {
srv.loopWG.Add(1) srv.loopWG.Add(1)
go srv.run(dialer) go srv.run(dialer)
go runMetricsFeedHelper(5 * time.Second)
srv.running = true srv.running = true
return nil return nil
} }