mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les/csvlogger: add csv logger
This commit is contained in:
parent
2388e425f2
commit
99711a2edf
8 changed files with 311 additions and 28 deletions
28
les/api.go
28
les/api.go
|
|
@ -19,11 +19,13 @@ package les
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common/hexutil"
|
"github.com/ethereum/go-ethereum/common/hexutil"
|
||||||
"github.com/ethereum/go-ethereum/common/mclock"
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
|
"github.com/ethereum/go-ethereum/les/csvlogger"
|
||||||
"github.com/ethereum/go-ethereum/p2p/enode"
|
"github.com/ethereum/go-ethereum/p2p/enode"
|
||||||
"github.com/ethereum/go-ethereum/rpc"
|
"github.com/ethereum/go-ethereum/rpc"
|
||||||
)
|
)
|
||||||
|
|
@ -144,6 +146,8 @@ type priorityClientPool struct {
|
||||||
totalCap, totalCapAnnounced uint64
|
totalCap, totalCapAnnounced uint64
|
||||||
totalConnectedCap, freeClientCap uint64
|
totalConnectedCap, freeClientCap uint64
|
||||||
maxPeers, priorityCount int
|
maxPeers, priorityCount int
|
||||||
|
logger *csvlogger.Logger
|
||||||
|
logTotalPriConn *csvlogger.Channel
|
||||||
|
|
||||||
subs tcSubs
|
subs tcSubs
|
||||||
updateSchedule []scheduledUpdate
|
updateSchedule []scheduledUpdate
|
||||||
|
|
@ -164,12 +168,14 @@ type priorityClientInfo struct {
|
||||||
}
|
}
|
||||||
|
|
||||||
// newPriorityClientPool creates a new priority client pool
|
// newPriorityClientPool creates a new priority client pool
|
||||||
func newPriorityClientPool(freeClientCap uint64, ps *peerSet, child clientPool) *priorityClientPool {
|
func newPriorityClientPool(freeClientCap uint64, ps *peerSet, child clientPool, logger *csvlogger.Logger) *priorityClientPool {
|
||||||
return &priorityClientPool{
|
return &priorityClientPool{
|
||||||
clients: make(map[enode.ID]priorityClientInfo),
|
clients: make(map[enode.ID]priorityClientInfo),
|
||||||
freeClientCap: freeClientCap,
|
freeClientCap: freeClientCap,
|
||||||
ps: ps,
|
ps: ps,
|
||||||
child: child,
|
child: child,
|
||||||
|
logger: logger,
|
||||||
|
logTotalPriConn: logger.NewChannel("totalPriConn", 0),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -185,6 +191,7 @@ func (v *priorityClientPool) registerPeer(p *peer) {
|
||||||
|
|
||||||
id := p.ID()
|
id := p.ID()
|
||||||
c := v.clients[id]
|
c := v.clients[id]
|
||||||
|
v.logger.Event(fmt.Sprintf("priorityClientPool: registerPeer cap=%d connected=%v, %x", c.cap, c.connected, id.Bytes()))
|
||||||
if c.connected {
|
if c.connected {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -192,6 +199,7 @@ func (v *priorityClientPool) registerPeer(p *peer) {
|
||||||
v.child.registerPeer(p)
|
v.child.registerPeer(p)
|
||||||
}
|
}
|
||||||
if c.cap != 0 && v.totalConnectedCap+c.cap > v.totalCap {
|
if c.cap != 0 && v.totalConnectedCap+c.cap > v.totalCap {
|
||||||
|
v.logger.Event(fmt.Sprintf("priorityClientPool: rejected, %x", id.Bytes()))
|
||||||
go v.ps.Unregister(p.id)
|
go v.ps.Unregister(p.id)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -202,6 +210,8 @@ func (v *priorityClientPool) registerPeer(p *peer) {
|
||||||
if c.cap != 0 {
|
if c.cap != 0 {
|
||||||
v.priorityCount++
|
v.priorityCount++
|
||||||
v.totalConnectedCap += c.cap
|
v.totalConnectedCap += c.cap
|
||||||
|
v.logger.Event(fmt.Sprintf("priorityClientPool: accepted with %d capacity, %x", c.cap, id.Bytes()))
|
||||||
|
v.logTotalPriConn.Update(float64(v.totalConnectedCap))
|
||||||
if v.child != nil {
|
if v.child != nil {
|
||||||
v.child.setLimits(v.maxPeers-v.priorityCount, v.totalCap-v.totalConnectedCap)
|
v.child.setLimits(v.maxPeers-v.priorityCount, v.totalCap-v.totalConnectedCap)
|
||||||
}
|
}
|
||||||
|
|
@ -217,6 +227,7 @@ func (v *priorityClientPool) unregisterPeer(p *peer) {
|
||||||
|
|
||||||
id := p.ID()
|
id := p.ID()
|
||||||
c := v.clients[id]
|
c := v.clients[id]
|
||||||
|
v.logger.Event(fmt.Sprintf("priorityClientPool: unregisterPeer cap=%d connected=%v, %x", c.cap, c.connected, id.Bytes()))
|
||||||
if !c.connected {
|
if !c.connected {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -225,6 +236,7 @@ func (v *priorityClientPool) unregisterPeer(p *peer) {
|
||||||
v.clients[id] = c
|
v.clients[id] = c
|
||||||
v.priorityCount--
|
v.priorityCount--
|
||||||
v.totalConnectedCap -= c.cap
|
v.totalConnectedCap -= c.cap
|
||||||
|
v.logTotalPriConn.Update(float64(v.totalConnectedCap))
|
||||||
if v.child != nil {
|
if v.child != nil {
|
||||||
v.child.setLimits(v.maxPeers-v.priorityCount, v.totalCap-v.totalConnectedCap)
|
v.child.setLimits(v.maxPeers-v.priorityCount, v.totalCap-v.totalConnectedCap)
|
||||||
}
|
}
|
||||||
|
|
@ -299,8 +311,10 @@ func (v *priorityClientPool) setLimitsNow(count int, totalCap uint64) {
|
||||||
if v.priorityCount > count || v.totalConnectedCap > totalCap {
|
if v.priorityCount > count || v.totalConnectedCap > totalCap {
|
||||||
for id, c := range v.clients {
|
for id, c := range v.clients {
|
||||||
if c.connected {
|
if c.connected {
|
||||||
|
v.logger.Event(fmt.Sprintf("priorityClientPool: setLimitsNow kicked out, %x", id.Bytes()))
|
||||||
c.connected = false
|
c.connected = false
|
||||||
v.totalConnectedCap -= c.cap
|
v.totalConnectedCap -= c.cap
|
||||||
|
v.logTotalPriConn.Update(float64(v.totalConnectedCap))
|
||||||
v.priorityCount--
|
v.priorityCount--
|
||||||
v.clients[id] = c
|
v.clients[id] = c
|
||||||
go v.ps.Unregister(c.peer.id)
|
go v.ps.Unregister(c.peer.id)
|
||||||
|
|
@ -356,6 +370,7 @@ func (v *priorityClientPool) setClientCapacity(id enode.ID, cap uint64) error {
|
||||||
v.priorityCount--
|
v.priorityCount--
|
||||||
}
|
}
|
||||||
v.totalConnectedCap += cap - c.cap
|
v.totalConnectedCap += cap - c.cap
|
||||||
|
v.logTotalPriConn.Update(float64(v.totalConnectedCap))
|
||||||
if v.child != nil {
|
if v.child != nil {
|
||||||
v.child.setLimits(v.maxPeers-v.priorityCount, v.totalCap-v.totalConnectedCap)
|
v.child.setLimits(v.maxPeers-v.priorityCount, v.totalCap-v.totalConnectedCap)
|
||||||
}
|
}
|
||||||
|
|
@ -374,6 +389,9 @@ func (v *priorityClientPool) setClientCapacity(id enode.ID, cap uint64) error {
|
||||||
} else {
|
} else {
|
||||||
delete(v.clients, id)
|
delete(v.clients, id)
|
||||||
}
|
}
|
||||||
|
if c.connected {
|
||||||
|
v.logger.Event(fmt.Sprintf("priorityClientPool: changed capacity to %d, %x", cap, id.Bytes()))
|
||||||
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ package les
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
|
"fmt"
|
||||||
"math"
|
"math"
|
||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
|
|
@ -26,6 +27,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/common/mclock"
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
"github.com/ethereum/go-ethereum/eth"
|
"github.com/ethereum/go-ethereum/eth"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
|
"github.com/ethereum/go-ethereum/les/csvlogger"
|
||||||
"github.com/ethereum/go-ethereum/les/flowcontrol"
|
"github.com/ethereum/go-ethereum/les/flowcontrol"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
)
|
)
|
||||||
|
|
@ -99,16 +101,22 @@ type costTracker struct {
|
||||||
gfLock sync.RWMutex
|
gfLock sync.RWMutex
|
||||||
totalRechargeCh chan uint64
|
totalRechargeCh chan uint64
|
||||||
|
|
||||||
stats map[uint64][]uint64
|
stats map[uint64][]uint64
|
||||||
|
logger *csvlogger.Logger
|
||||||
|
logRecentUsage, logTotalRecharge, logRelCost *csvlogger.Channel
|
||||||
}
|
}
|
||||||
|
|
||||||
// newCostTracker creates a cost tracker and loads the cost factor statistics from the database
|
// newCostTracker creates a cost tracker and loads the cost factor statistics from the database
|
||||||
func newCostTracker(db ethdb.Database, config *eth.Config) *costTracker {
|
func newCostTracker(db ethdb.Database, config *eth.Config, logger *csvlogger.Logger) *costTracker {
|
||||||
utilTarget := float64(config.LightServ) * flowcontrol.FixedPointMultiplier / 100
|
utilTarget := float64(config.LightServ) * flowcontrol.FixedPointMultiplier / 100
|
||||||
ct := &costTracker{
|
ct := &costTracker{
|
||||||
db: db,
|
db: db,
|
||||||
stopCh: make(chan chan struct{}),
|
stopCh: make(chan chan struct{}),
|
||||||
utilTarget: utilTarget,
|
utilTarget: utilTarget,
|
||||||
|
logger: logger,
|
||||||
|
logRelCost: logger.NewMinMaxChannel("relativeCost", true),
|
||||||
|
logRecentUsage: logger.NewMinMaxChannel("recentUsage", true),
|
||||||
|
logTotalRecharge: logger.NewChannel("totalRecharge", 0.01),
|
||||||
}
|
}
|
||||||
if config.LightBandwidthIn > 0 {
|
if config.LightBandwidthIn > 0 {
|
||||||
ct.inSizeFactor = utilTarget / float64(config.LightBandwidthIn)
|
ct.inSizeFactor = utilTarget / float64(config.LightBandwidthIn)
|
||||||
|
|
@ -201,14 +209,23 @@ func (ct *costTracker) gfLoop() {
|
||||||
case r := <-ct.gfUpdateCh:
|
case r := <-ct.gfUpdateCh:
|
||||||
now := mclock.Now()
|
now := mclock.Now()
|
||||||
max := r.servingTime * gf
|
max := r.servingTime * gf
|
||||||
|
if ct.logRelCost != nil && r.avgTime > 1e-20 {
|
||||||
|
ct.logRelCost.Update(max / r.avgTime)
|
||||||
|
}
|
||||||
if r.avgTime > max {
|
if r.avgTime > max {
|
||||||
max = r.avgTime
|
max = r.avgTime
|
||||||
}
|
}
|
||||||
dt := float64(now - expUpdate)
|
dt := float64(now - expUpdate)
|
||||||
expUpdate = now
|
expUpdate = now
|
||||||
gfUsage = gfUsage*math.Exp(-dt/float64(gfUsageTC)) + max*1000000/float64(gfUsageTC)
|
gfUsage = gfUsage*math.Exp(-dt/float64(gfUsageTC)) + max*1000000/float64(gfUsageTC)
|
||||||
|
totalRecharge := ct.utilTarget * gf
|
||||||
|
ct.logRecentUsage.Update(gfUsage)
|
||||||
|
ct.logTotalRecharge.Update(totalRecharge)
|
||||||
|
if r.servingTime > 1000000000 {
|
||||||
|
ct.logger.Event(fmt.Sprintf("Very long servingTime = %f avgTime = %f costFactor = %f", r.servingTime, r.avgTime, gf))
|
||||||
|
}
|
||||||
|
|
||||||
if gfUsage >= gfUsageThreshold*ct.utilTarget*gf {
|
if gfUsage >= gfUsageThreshold*totalRecharge {
|
||||||
gfSum += r.avgTime
|
gfSum += r.avgTime
|
||||||
gfWeight += r.servingTime
|
gfWeight += r.servingTime
|
||||||
if time.Duration(now-lastUpdate) > time.Second {
|
if time.Duration(now-lastUpdate) > time.Second {
|
||||||
|
|
@ -224,7 +241,7 @@ func (ct *costTracker) gfLoop() {
|
||||||
ct.gfLock.Unlock()
|
ct.gfLock.Unlock()
|
||||||
if ch != nil {
|
if ch != nil {
|
||||||
select {
|
select {
|
||||||
case ct.totalRechargeCh <- uint64(ct.utilTarget * gf):
|
case ct.totalRechargeCh <- uint64(totalRecharge):
|
||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
210
les/csvlogger/csvlogger.go
Normal file
210
les/csvlogger/csvlogger.go
Normal file
|
|
@ -0,0 +1,210 @@
|
||||||
|
// Copyright 2019 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package csvlogger
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"os"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Logger struct {
|
||||||
|
file *os.File
|
||||||
|
started mclock.AbsTime
|
||||||
|
channels []*Channel
|
||||||
|
period time.Duration
|
||||||
|
stopCh, stopped chan struct{}
|
||||||
|
storeCh chan string
|
||||||
|
eventHeader string
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewLogger(fileName string, period time.Duration, eventHeader string) *Logger {
|
||||||
|
f, err := os.Create(fileName)
|
||||||
|
if err != nil {
|
||||||
|
log.Error("Error creating log file", "name", fileName, "error", err)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return &Logger{
|
||||||
|
file: f,
|
||||||
|
period: period,
|
||||||
|
stopCh: make(chan struct{}),
|
||||||
|
storeCh: make(chan string, 1),
|
||||||
|
eventHeader: eventHeader,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *Logger) NewChannel(name string, threshold float64) *Channel {
|
||||||
|
if l == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
c := &Channel{
|
||||||
|
logger: l,
|
||||||
|
name: name,
|
||||||
|
threshold: threshold,
|
||||||
|
}
|
||||||
|
l.channels = append(l.channels, c)
|
||||||
|
return c
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *Logger) NewMinMaxChannel(name string, zeroDefault bool) *Channel {
|
||||||
|
if l == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
c := &Channel{
|
||||||
|
logger: l,
|
||||||
|
name: name,
|
||||||
|
minmax: true,
|
||||||
|
mmZeroDefault: zeroDefault,
|
||||||
|
}
|
||||||
|
l.channels = append(l.channels, c)
|
||||||
|
return c
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *Logger) store(event string) {
|
||||||
|
s := fmt.Sprintf("%g", float64(mclock.Now()-l.started)/1000000000)
|
||||||
|
for _, ch := range l.channels {
|
||||||
|
s += ", " + ch.store()
|
||||||
|
}
|
||||||
|
if event != "" {
|
||||||
|
s += ", " + event
|
||||||
|
}
|
||||||
|
l.file.WriteString(s + "\n")
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *Logger) Start() {
|
||||||
|
if l == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
l.started = mclock.Now()
|
||||||
|
s := "Time"
|
||||||
|
for _, ch := range l.channels {
|
||||||
|
s += ", " + ch.header()
|
||||||
|
}
|
||||||
|
if l.eventHeader != "" {
|
||||||
|
s += ", " + l.eventHeader
|
||||||
|
}
|
||||||
|
l.file.WriteString(s + "\n")
|
||||||
|
fmt.Println(s)
|
||||||
|
go func() {
|
||||||
|
timer := time.NewTimer(l.period)
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-timer.C:
|
||||||
|
l.store("")
|
||||||
|
timer.Reset(l.period)
|
||||||
|
case event := <-l.storeCh:
|
||||||
|
l.store(event)
|
||||||
|
if !timer.Stop() {
|
||||||
|
<-timer.C
|
||||||
|
}
|
||||||
|
timer.Reset(l.period)
|
||||||
|
case <-l.stopCh:
|
||||||
|
close(l.stopped)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *Logger) Stop() {
|
||||||
|
if l == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
l.stopped = make(chan struct{})
|
||||||
|
close(l.stopCh)
|
||||||
|
<-l.stopped
|
||||||
|
l.file.Close()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *Logger) Event(event string) {
|
||||||
|
if l == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case l.storeCh <- event:
|
||||||
|
case <-l.stopCh:
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type Channel struct {
|
||||||
|
logger *Logger
|
||||||
|
lock sync.Mutex
|
||||||
|
name string
|
||||||
|
threshold, storeMin, storeMax, lastValue, min, max float64
|
||||||
|
minmax, mmSet, mmZeroDefault bool
|
||||||
|
}
|
||||||
|
|
||||||
|
func (lc *Channel) Update(value float64) {
|
||||||
|
if lc == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
lc.lock.Lock()
|
||||||
|
defer lc.lock.Unlock()
|
||||||
|
|
||||||
|
lc.lastValue = value
|
||||||
|
if lc.minmax {
|
||||||
|
if value > lc.max || !lc.mmSet {
|
||||||
|
lc.max = value
|
||||||
|
}
|
||||||
|
if value < lc.min || !lc.mmSet {
|
||||||
|
lc.min = value
|
||||||
|
}
|
||||||
|
lc.mmSet = true
|
||||||
|
} else {
|
||||||
|
if value < lc.storeMin || value > lc.storeMax {
|
||||||
|
select {
|
||||||
|
case lc.logger.storeCh <- "":
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (lc *Channel) store() (s string) {
|
||||||
|
lc.lock.Lock()
|
||||||
|
defer lc.lock.Unlock()
|
||||||
|
|
||||||
|
if lc.minmax {
|
||||||
|
s = fmt.Sprintf("%g, %g", lc.min, lc.max)
|
||||||
|
lc.mmSet = false
|
||||||
|
if lc.mmZeroDefault {
|
||||||
|
lc.min = 0
|
||||||
|
} else {
|
||||||
|
lc.min = lc.lastValue
|
||||||
|
}
|
||||||
|
lc.max = lc.min
|
||||||
|
} else {
|
||||||
|
s = fmt.Sprintf("%g", lc.lastValue)
|
||||||
|
lc.storeMin = lc.lastValue * (1 - lc.threshold)
|
||||||
|
lc.storeMax = lc.lastValue * (1 + lc.threshold)
|
||||||
|
if lc.lastValue < 0 {
|
||||||
|
lc.storeMin, lc.storeMax = lc.storeMax, lc.storeMin
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
func (lc *Channel) header() string {
|
||||||
|
if lc.minmax {
|
||||||
|
return lc.name + " (min), " + lc.name + " (max)"
|
||||||
|
}
|
||||||
|
return lc.name
|
||||||
|
}
|
||||||
|
|
@ -24,6 +24,7 @@ import (
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common/mclock"
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
"github.com/ethereum/go-ethereum/common/prque"
|
"github.com/ethereum/go-ethereum/common/prque"
|
||||||
|
"github.com/ethereum/go-ethereum/les/csvlogger"
|
||||||
)
|
)
|
||||||
|
|
||||||
// cmNodeFields are ClientNode fields used by the client manager
|
// cmNodeFields are ClientNode fields used by the client manager
|
||||||
|
|
@ -68,6 +69,8 @@ type ClientManager struct {
|
||||||
capLastUpdate mclock.AbsTime
|
capLastUpdate mclock.AbsTime
|
||||||
totalCapacityCh chan uint64
|
totalCapacityCh chan uint64
|
||||||
|
|
||||||
|
logTotalCap *csvlogger.Channel
|
||||||
|
|
||||||
// recharge integrator is increasing in each moment with a rate of
|
// recharge integrator is increasing in each moment with a rate of
|
||||||
// (totalRecharge / sumRecharge)*FixedPointMultiplier or 0 if sumRecharge==0
|
// (totalRecharge / sumRecharge)*FixedPointMultiplier or 0 if sumRecharge==0
|
||||||
rcLastUpdate mclock.AbsTime // last time the recharge integrator was updated
|
rcLastUpdate mclock.AbsTime // last time the recharge integrator was updated
|
||||||
|
|
@ -101,11 +104,12 @@ type ClientManager struct {
|
||||||
// starting from zero in order to not let a single low-priority client use up
|
// starting from zero in order to not let a single low-priority client use up
|
||||||
// the entire server capacity and thus ensure quick availability for others at
|
// the entire server capacity and thus ensure quick availability for others at
|
||||||
// any moment.
|
// any moment.
|
||||||
func NewClientManager(curve PieceWiseLinear, clock mclock.Clock) *ClientManager {
|
func NewClientManager(curve PieceWiseLinear, clock mclock.Clock, logger *csvlogger.Logger) *ClientManager {
|
||||||
cm := &ClientManager{
|
cm := &ClientManager{
|
||||||
clock: clock,
|
clock: clock,
|
||||||
rcQueue: prque.New(func(a interface{}, i int) { a.(*ClientNode).queueIndex = i }),
|
rcQueue: prque.New(func(a interface{}, i int) { a.(*ClientNode).queueIndex = i }),
|
||||||
capLastUpdate: clock.Now(),
|
capLastUpdate: clock.Now(),
|
||||||
|
logTotalCap: logger.NewChannel("totalCapacity", 0.01),
|
||||||
}
|
}
|
||||||
if curve != nil {
|
if curve != nil {
|
||||||
cm.SetRechargeCurve(curve)
|
cm.SetRechargeCurve(curve)
|
||||||
|
|
@ -337,6 +341,7 @@ func (cm *ClientManager) refreshCapacity() {
|
||||||
if totalCapacity >= cm.totalCapacity*0.999 && totalCapacity <= cm.totalCapacity*1.001 {
|
if totalCapacity >= cm.totalCapacity*0.999 && totalCapacity <= cm.totalCapacity*1.001 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
cm.logTotalCap.Update(totalCapacity)
|
||||||
cm.totalCapacity = totalCapacity
|
cm.totalCapacity = totalCapacity
|
||||||
if cm.totalCapacityCh != nil {
|
if cm.totalCapacityCh != nil {
|
||||||
select {
|
select {
|
||||||
|
|
|
||||||
|
|
@ -26,6 +26,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/common/mclock"
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
"github.com/ethereum/go-ethereum/common/prque"
|
"github.com/ethereum/go-ethereum/common/prque"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
|
"github.com/ethereum/go-ethereum/les/csvlogger"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
)
|
)
|
||||||
|
|
@ -52,6 +53,8 @@ type freeClientPool struct {
|
||||||
|
|
||||||
connectedLimit, totalLimit int
|
connectedLimit, totalLimit int
|
||||||
freeClientCap uint64
|
freeClientCap uint64
|
||||||
|
logger *csvlogger.Logger
|
||||||
|
logTotalFreeConn *csvlogger.Channel
|
||||||
|
|
||||||
addressMap map[string]*freeClientPoolEntry
|
addressMap map[string]*freeClientPoolEntry
|
||||||
connPool, disconnPool *prque.Prque
|
connPool, disconnPool *prque.Prque
|
||||||
|
|
@ -66,16 +69,18 @@ const (
|
||||||
)
|
)
|
||||||
|
|
||||||
// newFreeClientPool creates a new free client pool
|
// newFreeClientPool creates a new free client pool
|
||||||
func newFreeClientPool(db ethdb.Database, freeClientCap uint64, totalLimit int, clock mclock.Clock, removePeer func(string)) *freeClientPool {
|
func newFreeClientPool(db ethdb.Database, freeClientCap uint64, totalLimit int, clock mclock.Clock, removePeer func(string), logger *csvlogger.Logger) *freeClientPool {
|
||||||
pool := &freeClientPool{
|
pool := &freeClientPool{
|
||||||
db: db,
|
db: db,
|
||||||
clock: clock,
|
clock: clock,
|
||||||
addressMap: make(map[string]*freeClientPoolEntry),
|
addressMap: make(map[string]*freeClientPoolEntry),
|
||||||
connPool: prque.New(poolSetIndex),
|
connPool: prque.New(poolSetIndex),
|
||||||
disconnPool: prque.New(poolSetIndex),
|
disconnPool: prque.New(poolSetIndex),
|
||||||
freeClientCap: freeClientCap,
|
freeClientCap: freeClientCap,
|
||||||
totalLimit: totalLimit,
|
totalLimit: totalLimit,
|
||||||
removePeer: removePeer,
|
logger: logger,
|
||||||
|
logTotalFreeConn: logger.NewChannel("totalFreeConn", 0),
|
||||||
|
removePeer: removePeer,
|
||||||
}
|
}
|
||||||
pool.loadFromDb()
|
pool.loadFromDb()
|
||||||
return pool
|
return pool
|
||||||
|
|
@ -107,7 +112,9 @@ func (f *freeClientPool) connect(address, id string) bool {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
f.logger.Event("freeClientPool: connecting from " + address + ", " + id)
|
||||||
if f.connectedLimit == 0 {
|
if f.connectedLimit == 0 {
|
||||||
|
f.logger.Event("freeClientPool: rejected, " + id)
|
||||||
log.Debug("Client rejected", "address", address)
|
log.Debug("Client rejected", "address", address)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
@ -119,6 +126,7 @@ func (f *freeClientPool) connect(address, id string) bool {
|
||||||
f.addressMap[address] = e
|
f.addressMap[address] = e
|
||||||
} else {
|
} else {
|
||||||
if e.connected {
|
if e.connected {
|
||||||
|
f.logger.Event("freeClientPool: already connected, " + id)
|
||||||
log.Debug("Client already connected", "address", address)
|
log.Debug("Client already connected", "address", address)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
@ -131,9 +139,11 @@ func (f *freeClientPool) connect(address, id string) bool {
|
||||||
if e.linUsage+int64(connectedBias)-i.linUsage < 0 {
|
if e.linUsage+int64(connectedBias)-i.linUsage < 0 {
|
||||||
// kick it out and accept the new client
|
// kick it out and accept the new client
|
||||||
f.dropClient(i, now)
|
f.dropClient(i, now)
|
||||||
|
f.logger.Event("freeClientPool: kicked out, " + i.id)
|
||||||
} else {
|
} else {
|
||||||
// keep the old client and reject the new one
|
// keep the old client and reject the new one
|
||||||
f.connPool.Push(i, i.linUsage)
|
f.connPool.Push(i, i.linUsage)
|
||||||
|
f.logger.Event("freeClientPool: rejected, " + id)
|
||||||
log.Debug("Client rejected", "address", address)
|
log.Debug("Client rejected", "address", address)
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
@ -142,9 +152,11 @@ func (f *freeClientPool) connect(address, id string) bool {
|
||||||
e.connected = true
|
e.connected = true
|
||||||
e.id = id
|
e.id = id
|
||||||
f.connPool.Push(e, e.linUsage)
|
f.connPool.Push(e, e.linUsage)
|
||||||
|
f.logTotalFreeConn.Update(float64(uint64(f.connPool.Size()) * f.freeClientCap))
|
||||||
if f.connPool.Size()+f.disconnPool.Size() > f.totalLimit {
|
if f.connPool.Size()+f.disconnPool.Size() > f.totalLimit {
|
||||||
f.disconnPool.Pop()
|
f.disconnPool.Pop()
|
||||||
}
|
}
|
||||||
|
f.logger.Event("freeClientPool: accepted, " + id)
|
||||||
log.Debug("Client accepted", "address", address)
|
log.Debug("Client accepted", "address", address)
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
@ -174,9 +186,11 @@ func (f *freeClientPool) disconnect(address string) {
|
||||||
}
|
}
|
||||||
|
|
||||||
f.connPool.Remove(e.index)
|
f.connPool.Remove(e.index)
|
||||||
|
f.logTotalFreeConn.Update(float64(uint64(f.connPool.Size()) * f.freeClientCap))
|
||||||
f.calcLogUsage(e, now)
|
f.calcLogUsage(e, now)
|
||||||
e.connected = false
|
e.connected = false
|
||||||
f.disconnPool.Push(e, -e.logUsage)
|
f.disconnPool.Push(e, -e.logUsage)
|
||||||
|
f.logger.Event("freeClientPool: disconnected, " + e.id)
|
||||||
log.Debug("Client disconnected", "address", address)
|
log.Debug("Client disconnected", "address", address)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -194,6 +208,7 @@ func (f *freeClientPool) setLimits(count int, totalCap uint64) {
|
||||||
for f.connPool.Size() > f.connectedLimit {
|
for f.connPool.Size() > f.connectedLimit {
|
||||||
i := f.connPool.PopItem().(*freeClientPoolEntry)
|
i := f.connPool.PopItem().(*freeClientPoolEntry)
|
||||||
f.dropClient(i, now)
|
f.dropClient(i, now)
|
||||||
|
f.logger.Event("freeClientPool: setLimits kicked out, " + i.id)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -201,6 +216,7 @@ func (f *freeClientPool) setLimits(count int, totalCap uint64) {
|
||||||
// disconnected pool
|
// disconnected pool
|
||||||
func (f *freeClientPool) dropClient(i *freeClientPoolEntry, now mclock.AbsTime) {
|
func (f *freeClientPool) dropClient(i *freeClientPoolEntry, now mclock.AbsTime) {
|
||||||
f.connPool.Remove(i.index)
|
f.connPool.Remove(i.index)
|
||||||
|
f.logTotalFreeConn.Update(float64(uint64(f.connPool.Size()) * f.freeClientCap))
|
||||||
f.calcLogUsage(i, now)
|
f.calcLogUsage(i, now)
|
||||||
i.connected = false
|
i.connected = false
|
||||||
f.disconnPool.Push(i, -i.logUsage)
|
f.disconnPool.Push(i, -i.logUsage)
|
||||||
|
|
|
||||||
|
|
@ -61,7 +61,7 @@ func testFreeClientPool(t *testing.T, connLimit, clientCount int) {
|
||||||
}
|
}
|
||||||
disconnCh <- i
|
disconnCh <- i
|
||||||
}
|
}
|
||||||
pool = newFreeClientPool(db, 1, 10000, &clock, disconnFn)
|
pool = newFreeClientPool(db, 1, 10000, &clock, disconnFn, nil)
|
||||||
)
|
)
|
||||||
pool.setLimits(connLimit, uint64(connLimit))
|
pool.setLimits(connLimit, uint64(connLimit))
|
||||||
|
|
||||||
|
|
@ -130,7 +130,7 @@ func testFreeClientPool(t *testing.T, connLimit, clientCount int) {
|
||||||
|
|
||||||
// close and restart pool
|
// close and restart pool
|
||||||
pool.stop()
|
pool.stop()
|
||||||
pool = newFreeClientPool(db, 1, 10000, &clock, disconnFn)
|
pool = newFreeClientPool(db, 1, 10000, &clock, disconnFn, nil)
|
||||||
pool.setLimits(connLimit, uint64(connLimit))
|
pool.setLimits(connLimit, uint64(connLimit))
|
||||||
|
|
||||||
// try connecting all known peers (connLimit should be filled up)
|
// try connecting all known peers (connLimit should be filled up)
|
||||||
|
|
|
||||||
|
|
@ -34,6 +34,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/eth/downloader"
|
"github.com/ethereum/go-ethereum/eth/downloader"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
|
"github.com/ethereum/go-ethereum/les/csvlogger"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/p2p"
|
"github.com/ethereum/go-ethereum/p2p"
|
||||||
|
|
@ -118,6 +119,7 @@ type ProtocolManager struct {
|
||||||
|
|
||||||
wg *sync.WaitGroup
|
wg *sync.WaitGroup
|
||||||
eventMux *event.TypeMux
|
eventMux *event.TypeMux
|
||||||
|
logger *csvlogger.Logger
|
||||||
|
|
||||||
// Callbacks
|
// Callbacks
|
||||||
synced func() bool
|
synced func() bool
|
||||||
|
|
@ -272,6 +274,7 @@ func (pm *ProtocolManager) handle(p *peer) error {
|
||||||
// Ignore maxPeers if this is a trusted peer
|
// Ignore maxPeers if this is a trusted peer
|
||||||
// In server mode we try to check into the client pool after handshake
|
// In server mode we try to check into the client pool after handshake
|
||||||
if pm.client && pm.peers.Len() >= pm.maxPeers && !p.Peer.Info().Network.Trusted {
|
if pm.client && pm.peers.Len() >= pm.maxPeers && !p.Peer.Info().Network.Trusted {
|
||||||
|
pm.logger.Event("Rejected (too many peers), " + p.id)
|
||||||
return p2p.DiscTooManyPeers
|
return p2p.DiscTooManyPeers
|
||||||
}
|
}
|
||||||
// Reject light clients if server is not synced.
|
// Reject light clients if server is not synced.
|
||||||
|
|
@ -290,6 +293,7 @@ func (pm *ProtocolManager) handle(p *peer) error {
|
||||||
)
|
)
|
||||||
if err := p.Handshake(td, hash, number, genesis.Hash(), pm.server); err != nil {
|
if err := p.Handshake(td, hash, number, genesis.Hash(), pm.server); err != nil {
|
||||||
p.Log().Debug("Light Ethereum handshake failed", "err", err)
|
p.Log().Debug("Light Ethereum handshake failed", "err", err)
|
||||||
|
pm.logger.Event("Handshake error: " + err.Error() + ", " + p.id)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if p.fcClient != nil {
|
if p.fcClient != nil {
|
||||||
|
|
@ -303,9 +307,12 @@ func (pm *ProtocolManager) handle(p *peer) error {
|
||||||
// Register the peer locally
|
// Register the peer locally
|
||||||
if err := pm.peers.Register(p); err != nil {
|
if err := pm.peers.Register(p); err != nil {
|
||||||
p.Log().Error("Light Ethereum peer registration failed", "err", err)
|
p.Log().Error("Light Ethereum peer registration failed", "err", err)
|
||||||
|
pm.logger.Event("Peer registration error: " + err.Error() + ", " + p.id)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
pm.logger.Event("Connection established, " + p.id)
|
||||||
defer func() {
|
defer func() {
|
||||||
|
pm.logger.Event("Closed connection, " + p.id)
|
||||||
pm.removePeer(p.id)
|
pm.removePeer(p.id)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
|
@ -326,6 +333,7 @@ func (pm *ProtocolManager) handle(p *peer) error {
|
||||||
// main loop. handle incoming messages.
|
// main loop. handle incoming messages.
|
||||||
for {
|
for {
|
||||||
if err := pm.handleMsg(p); err != nil {
|
if err := pm.handleMsg(p); err != nil {
|
||||||
|
pm.logger.Event("Message handling error: " + err.Error() + ", " + p.id)
|
||||||
p.Log().Debug("Light Ethereum message handling failed", "err", err)
|
p.Log().Debug("Light Ethereum message handling failed", "err", err)
|
||||||
if p.fcServer != nil {
|
if p.fcServer != nil {
|
||||||
p.fcServer.DumpLogs()
|
p.fcServer.DumpLogs()
|
||||||
|
|
|
||||||
|
|
@ -19,6 +19,7 @@ package les
|
||||||
import (
|
import (
|
||||||
"crypto/ecdsa"
|
"crypto/ecdsa"
|
||||||
"sync"
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/common/mclock"
|
"github.com/ethereum/go-ethereum/common/mclock"
|
||||||
|
|
@ -26,6 +27,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/core/rawdb"
|
"github.com/ethereum/go-ethereum/core/rawdb"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/eth"
|
"github.com/ethereum/go-ethereum/eth"
|
||||||
|
"github.com/ethereum/go-ethereum/les/csvlogger"
|
||||||
"github.com/ethereum/go-ethereum/les/flowcontrol"
|
"github.com/ethereum/go-ethereum/les/flowcontrol"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
|
@ -47,6 +49,7 @@ type LesServer struct {
|
||||||
privateKey *ecdsa.PrivateKey
|
privateKey *ecdsa.PrivateKey
|
||||||
quitSync chan struct{}
|
quitSync chan struct{}
|
||||||
onlyAnnounce bool
|
onlyAnnounce bool
|
||||||
|
csvLogger *csvlogger.Logger
|
||||||
|
|
||||||
thcNormal, thcBlockProcessing int // serving thread count for normal operation and block processing mode
|
thcNormal, thcBlockProcessing int // serving thread count for normal operation and block processing mode
|
||||||
|
|
||||||
|
|
@ -83,6 +86,8 @@ func NewLesServer(eth *eth.Ethereum, config *eth.Config) (*LesServer, error) {
|
||||||
for i, pv := range AdvertiseProtocolVersions {
|
for i, pv := range AdvertiseProtocolVersions {
|
||||||
lesTopics[i] = lesTopic(eth.BlockChain().Genesis().Hash(), pv)
|
lesTopics[i] = lesTopic(eth.BlockChain().Genesis().Hash(), pv)
|
||||||
}
|
}
|
||||||
|
var csvLogger *csvlogger.Logger
|
||||||
|
csvLogger = csvlogger.NewLogger("/tmp/server.csv", time.Second*10, "event, peerId")
|
||||||
|
|
||||||
srv := &LesServer{
|
srv := &LesServer{
|
||||||
lesCommons: lesCommons{
|
lesCommons: lesCommons{
|
||||||
|
|
@ -93,20 +98,22 @@ func NewLesServer(eth *eth.Ethereum, config *eth.Config) (*LesServer, error) {
|
||||||
bloomTrieIndexer: light.NewBloomTrieIndexer(eth.ChainDb(), nil, params.BloomBitsBlocks, params.BloomTrieFrequency),
|
bloomTrieIndexer: light.NewBloomTrieIndexer(eth.ChainDb(), nil, params.BloomBitsBlocks, params.BloomTrieFrequency),
|
||||||
protocolManager: pm,
|
protocolManager: pm,
|
||||||
},
|
},
|
||||||
costTracker: newCostTracker(eth.ChainDb(), config),
|
costTracker: newCostTracker(eth.ChainDb(), config, csvLogger),
|
||||||
quitSync: quitSync,
|
quitSync: quitSync,
|
||||||
lesTopics: lesTopics,
|
lesTopics: lesTopics,
|
||||||
onlyAnnounce: config.OnlyAnnounce,
|
onlyAnnounce: config.OnlyAnnounce,
|
||||||
|
csvLogger: csvLogger,
|
||||||
}
|
}
|
||||||
|
|
||||||
logger := log.New()
|
logger := log.New()
|
||||||
pm.server = srv
|
pm.server = srv
|
||||||
|
pm.logger = csvLogger
|
||||||
srv.thcNormal = config.LightServ * 4 / 100
|
srv.thcNormal = config.LightServ * 4 / 100
|
||||||
if srv.thcNormal < 4 {
|
if srv.thcNormal < 4 {
|
||||||
srv.thcNormal = 4
|
srv.thcNormal = 4
|
||||||
}
|
}
|
||||||
srv.thcBlockProcessing = config.LightServ/100 + 1
|
srv.thcBlockProcessing = config.LightServ/100 + 1
|
||||||
srv.fcManager = flowcontrol.NewClientManager(nil, &mclock.System{})
|
srv.fcManager = flowcontrol.NewClientManager(nil, &mclock.System{}, csvLogger)
|
||||||
|
|
||||||
chtSectionCount, _, _ := srv.chtIndexer.Sections()
|
chtSectionCount, _, _ := srv.chtIndexer.Sections()
|
||||||
if chtSectionCount != 0 {
|
if chtSectionCount != 0 {
|
||||||
|
|
@ -205,10 +212,11 @@ func (s *LesServer) Start(srvr *p2p.Server) {
|
||||||
log.Warn("Light peer count limited", "specified", s.maxPeers, "allowed", freePeers)
|
log.Warn("Light peer count limited", "specified", s.maxPeers, "allowed", freePeers)
|
||||||
}
|
}
|
||||||
|
|
||||||
s.freeClientPool = newFreeClientPool(s.chainDb, s.freeClientCap, 10000, mclock.System{}, func(id string) { go s.protocolManager.removePeer(id) })
|
s.freeClientPool = newFreeClientPool(s.chainDb, s.freeClientCap, 10000, mclock.System{}, func(id string) { go s.protocolManager.removePeer(id) }, s.csvLogger)
|
||||||
s.priorityClientPool = newPriorityClientPool(s.freeClientCap, s.protocolManager.peers, s.freeClientPool)
|
s.priorityClientPool = newPriorityClientPool(s.freeClientCap, s.protocolManager.peers, s.freeClientPool, s.csvLogger)
|
||||||
|
|
||||||
s.protocolManager.peers.notify(s.priorityClientPool)
|
s.protocolManager.peers.notify(s.priorityClientPool)
|
||||||
|
s.csvLogger.Start()
|
||||||
s.startEventLoop()
|
s.startEventLoop()
|
||||||
s.protocolManager.Start(s.config.LightPeers)
|
s.protocolManager.Start(s.config.LightPeers)
|
||||||
if srvr.DiscV5 != nil {
|
if srvr.DiscV5 != nil {
|
||||||
|
|
@ -241,6 +249,7 @@ func (s *LesServer) Stop() {
|
||||||
s.freeClientPool.stop()
|
s.freeClientPool.stop()
|
||||||
s.costTracker.stop()
|
s.costTracker.stop()
|
||||||
s.protocolManager.Stop()
|
s.protocolManager.Stop()
|
||||||
|
s.csvLogger.Stop()
|
||||||
}
|
}
|
||||||
|
|
||||||
// todo(rjl493456442) separate client and server implementation.
|
// todo(rjl493456442) separate client and server implementation.
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue