mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 10:22:23 +00:00
dashboard, p2p: extend meteredConn with peer metering
This commit is contained in:
parent
8051a0768a
commit
64a97809a6
3 changed files with 48 additions and 7 deletions
|
|
@ -385,6 +385,10 @@ func (db *Dashboard) collectData() {
|
||||||
DiskWrite: ChartEntries{diskWrite},
|
DiskWrite: ChartEntries{diskWrite},
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
|
for k, v := range p2p.IngressTrafficMeters {
|
||||||
|
fmt.Println(k[:6], v.Count())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -21,6 +21,7 @@ package p2p
|
||||||
import (
|
import (
|
||||||
"net"
|
"net"
|
||||||
|
|
||||||
|
"fmt"
|
||||||
"github.com/ethereum/go-ethereum/metrics"
|
"github.com/ethereum/go-ethereum/metrics"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -29,12 +30,17 @@ var (
|
||||||
ingressTrafficMeter = metrics.NewRegisteredMeter("p2p/InboundTraffic", nil)
|
ingressTrafficMeter = metrics.NewRegisteredMeter("p2p/InboundTraffic", nil)
|
||||||
egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil)
|
egressConnectMeter = metrics.NewRegisteredMeter("p2p/OutboundConnects", nil)
|
||||||
egressTrafficMeter = metrics.NewRegisteredMeter("p2p/OutboundTraffic", nil)
|
egressTrafficMeter = metrics.NewRegisteredMeter("p2p/OutboundTraffic", nil)
|
||||||
|
IngressTrafficMeters = make(map[string]metrics.Meter)
|
||||||
|
EgressTrafficMeters = make(map[string]metrics.Meter)
|
||||||
)
|
)
|
||||||
|
|
||||||
// meteredConn is a wrapper around a net.Conn that meters both the
|
// meteredConn is a wrapper around a net.Conn that meters both the
|
||||||
// inbound and outbound network traffic.
|
// inbound and outbound network traffic.
|
||||||
type meteredConn struct {
|
type meteredConn struct {
|
||||||
net.Conn // Network connection to wrap with metering
|
net.Conn // Network connection to wrap with metering
|
||||||
|
id string
|
||||||
|
unmarkedIngress uint
|
||||||
|
unmarkedEgress uint
|
||||||
}
|
}
|
||||||
|
|
||||||
// newMeteredConn creates a new metered connection, also bumping the ingress or
|
// newMeteredConn creates a new metered connection, also bumping the ingress or
|
||||||
|
|
@ -59,7 +65,12 @@ func newMeteredConn(conn net.Conn, ingress bool) net.Conn {
|
||||||
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))
|
||||||
return
|
if rm, ok := IngressTrafficMeters[c.id]; ok {
|
||||||
|
rm.Mark(int64(n))
|
||||||
|
return n, err
|
||||||
|
}
|
||||||
|
c.unmarkedIngress += uint(n)
|
||||||
|
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
|
||||||
|
|
@ -67,5 +78,28 @@ func (c *meteredConn) Read(b []byte) (n int, err error) {
|
||||||
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))
|
||||||
return
|
if rm, ok := EgressTrafficMeters[c.id]; ok {
|
||||||
|
rm.Mark(int64(n))
|
||||||
|
return n, err
|
||||||
|
}
|
||||||
|
c.unmarkedEgress += uint(n)
|
||||||
|
return n, err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *meteredConn) meterIndividually(id string) {
|
||||||
|
c.id = id
|
||||||
|
IngressTrafficMeters[id] = metrics.NewRegisteredMeter(fmt.Sprintf("p2p/InboundTraffic/%s", id), nil)
|
||||||
|
EgressTrafficMeters[id] = metrics.NewRegisteredMeter(fmt.Sprintf("p2p/OutboundTraffic/%s", id), nil)
|
||||||
|
IngressTrafficMeters[id].Mark(int64(c.unmarkedIngress))
|
||||||
|
EgressTrafficMeters[id].Mark(int64(c.unmarkedEgress))
|
||||||
|
c.unmarkedIngress, c.unmarkedEgress = 0, 0
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *meteredConn) Close() error {
|
||||||
|
delete(IngressTrafficMeters, c.id)
|
||||||
|
delete(EgressTrafficMeters, c.id)
|
||||||
|
metrics.Unregister(fmt.Sprintf("p2p/InboundTraffic/%s", c.id))
|
||||||
|
metrics.Unregister(fmt.Sprintf("p2p/OutboundTraffic/%s", c.id))
|
||||||
|
fmt.Println("Close", c.id)
|
||||||
|
return c.Conn.Close()
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -586,6 +586,9 @@ func newRLPXFrameRW(conn io.ReadWriter, s secrets) *rlpxFrameRW {
|
||||||
// we use an all-zeroes IV for AES because the key used
|
// we use an all-zeroes IV for AES because the key used
|
||||||
// for encryption is ephemeral.
|
// for encryption is ephemeral.
|
||||||
iv := make([]byte, encc.BlockSize())
|
iv := make([]byte, encc.BlockSize())
|
||||||
|
if c, ok := conn.(*meteredConn); ok {
|
||||||
|
c.meterIndividually(s.RemoteID.String())
|
||||||
|
}
|
||||||
return &rlpxFrameRW{
|
return &rlpxFrameRW{
|
||||||
conn: conn,
|
conn: conn,
|
||||||
enc: cipher.NewCTR(encc, iv),
|
enc: cipher.NewCTR(encc, iv),
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue