beacon/light: colored logs, log levels, scheduler debug logs

This commit is contained in:
Zsolt Felfoldi 2023-12-23 03:32:41 +01:00 committed by Felix Lange
parent 56423e602e
commit f3e7b9fa7d
6 changed files with 67 additions and 20 deletions

View file

@ -40,8 +40,10 @@ func NewApiServer(api *BeaconLightApi) *ApiServer {
func (s *ApiServer) Subscribe(eventCallback func(event request.Event)) {
s.eventCallback = eventCallback
s.unsubscribe = s.api.StartHeadListener(func(slot uint64, blockRoot common.Hash) {
log.Debug("New head received", "slot", slot, "blockRoot", blockRoot)
eventCallback(request.Event{Type: sync.EvNewHead, Data: types.HeadInfo{Slot: slot, BlockRoot: blockRoot}})
}, func(head types.SignedHeader) {
log.Debug("New signed head received", "slot", head.Header.Slot, "blockRoot", head.Header.Hash(), "signerCount", head.Signature.SignerCount())
eventCallback(request.Event{Type: sync.EvNewSignedHead, Data: head})
}, func(err error) {
log.Warn("Head event stream error", "err", err)

View file

@ -16,6 +16,12 @@
package request
import (
"math"
"github.com/ethereum/go-ethereum/log"
)
type (
Request any
Response any
@ -40,14 +46,20 @@ func (p *RequestTracker) TryRequest(requestFn func(server Server) (Request, floa
bestServer Server
bestRequest Request
)
maxServerPriority, maxRequestPriority = -1000, -1000
maxServerPriority, maxRequestPriority = -math.MaxFloat32, -math.MaxFloat32
serverCount := len(p.servers)
var removed, candidates int
for server, _ := range p.servers {
canRequest, serverPriority := server.CanRequestNow()
if !canRequest {
delete(p.servers, server)
removed++
continue
}
request, requestPriority := requestFn(server)
if request != nil {
candidates++
}
if request == nil || requestPriority < maxRequestPriority ||
(requestPriority == maxRequestPriority && serverPriority <= maxServerPriority) {
continue
@ -55,6 +67,7 @@ func (p *RequestTracker) TryRequest(requestFn func(server Server) (Request, floa
maxServerPriority, maxRequestPriority = serverPriority, requestPriority
bestServer, bestRequest = server, request
}
log.Debug("Request attempt", "serverCount", serverCount, "removedServers", removed, "requestCandidates", candidates)
if bestServer == nil {
return ServerAndId{}, nil
}

View file

@ -48,6 +48,7 @@ type Scheduler struct {
lock sync.Mutex
clock mclock.Clock
modules []Module // first has highest priority
names map[Module]string
trackers map[Module]*RequestTracker
servers map[Server]struct{}
pending map[ServerAndId]pendingRequest
@ -61,7 +62,7 @@ type Scheduler struct {
type ServerEvent struct {
Server Server
Type int
Type string
Data any
}
@ -83,6 +84,7 @@ func NewScheduler(clock mclock.Clock) *Scheduler {
s := &Scheduler{
clock: clock,
servers: make(map[Server]struct{}),
names: make(map[Module]string),
trackers: make(map[Module]*RequestTracker),
pending: make(map[ServerAndId]pendingRequest),
stopCh: make(chan chan struct{}),
@ -98,7 +100,7 @@ func NewScheduler(clock mclock.Clock) *Scheduler {
// RegisterModule registers a module. Should be called before starting the scheduler.
// In each processing round the order of module processing depends on the order of
// registration.
func (s *Scheduler) RegisterModule(m Module) {
func (s *Scheduler) RegisterModule(m Module, name string) {
s.lock.Lock()
defer s.lock.Unlock()
@ -107,6 +109,7 @@ func (s *Scheduler) RegisterModule(m Module) {
scheduler: s,
module: m,
}
s.names[m] = name
}
// RegisterServer registers a new server.
@ -191,6 +194,12 @@ func (s *Scheduler) processModules() {
s.serverEvents = nil
s.lock.Unlock()
eventTypes := make([]string, len(serverEvents))
for i, ev := range serverEvents {
eventTypes[i] = ev.Type
}
log.Debug("Processing modules", "servers", len(servers), "server events", eventTypes)
for _, module := range s.modules {
s.lock.Lock()
tracker := s.trackers[module]
@ -198,6 +207,19 @@ func (s *Scheduler) processModules() {
requestEvents := tracker.requestEvents
tracker.requestEvents = nil
s.lock.Unlock()
var respCount, failCount, timeoutCount int
for _, ev := range requestEvents {
if ev.Response != nil {
respCount++
} else if ev.Finalized {
failCount++
} else {
timeoutCount++
}
}
log.Debug("Processing module", "name", s.names[module], "responses", respCount, "fails", failCount, "timeouts", timeoutCount)
if module.Process(tracker, requestEvents, serverEvents) {
s.Trigger()
}

View file

@ -26,16 +26,15 @@ import (
"github.com/ethereum/go-ethereum/log"
)
const (
var (
// request events
EvResponse = iota // data: IdAndResponse; sent by RequestServer
EvFail // data: ID; sent by RequestServer
EvTimeout // data: ID; sent by serverWithTimeout
EvResponse = "response" // data: IdAndResponse; sent by RequestServer
EvFail = "fail" // data: ID; sent by RequestServer
EvTimeout = "timeout" // data: ID; sent by serverWithTimeout
// server events
EvRegistered // data: nil; sent by Scheduler
EvUnregistered // data: nil; sent by Scheduler
EvCanRequestAgain // data: nil; sent by serverWithLimits
EvAppSpecific // application specific events (sent by RequestServer) start at this index
EvRegistered = "registered" // data: nil; sent by Scheduler
EvUnregistered = "unregistered" // data: nil; sent by Scheduler
EvCanRequestAgain = "canRequestAgain" // data: nil; sent by serverWithLimits
)
const (
@ -79,7 +78,7 @@ func NewServer(rs RequestServer, clock mclock.Clock) Server {
type serverSet map[Server]struct{}
type Event struct {
Type int
Type string
Data any
}

View file

@ -17,14 +17,13 @@
package sync
import (
"github.com/ethereum/go-ethereum/beacon/light/request"
"github.com/ethereum/go-ethereum/beacon/types"
"github.com/ethereum/go-ethereum/common"
)
const (
EvNewHead = iota + request.EvAppSpecific
EvNewSignedHead
EvNewHead = "newHead"
EvNewSignedHead = "newSignedHead"
)
type (

View file

@ -19,6 +19,7 @@ package main
import (
"context"
"fmt"
"io"
"os"
"strings"
"time"
@ -34,7 +35,10 @@ import (
ctypes "github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/ethdb/memorydb"
"github.com/ethereum/go-ethereum/internal/flags"
"github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/rpc"
"github.com/mattn/go-colorable"
"github.com/mattn/go-isatty"
"github.com/urfave/cli/v2"
)
@ -83,6 +87,14 @@ func main() {
}
func blsync(ctx *cli.Context) error {
usecolor := (isatty.IsTerminal(os.Stderr.Fd()) || isatty.IsCygwinTerminal(os.Stderr.Fd())) && os.Getenv("TERM") != "dumb"
output := io.Writer(os.Stderr)
if usecolor {
output = colorable.NewColorable(os.Stderr)
}
verbosity := log.FromLegacyLevel(ctx.Int(verbosityFlag.Name))
log.SetDefault(log.NewLogger(log.NewTerminalHandlerWithLevel(output, verbosity, usecolor)))
if !ctx.IsSet(utils.BeaconApiFlag.Name) {
utils.Fatalf("Beacon node light client API URL not specified")
}
@ -119,11 +131,11 @@ func blsync(ctx *cli.Context) error {
blockSync: beaconBlockSync,
}
scheduler.RegisterModule(checkpointInit)
scheduler.RegisterModule(forwardSync)
scheduler.RegisterModule(headSync)
scheduler.RegisterModule(beaconBlockSync)
scheduler.RegisterModule(engineApiUpdater)
scheduler.RegisterModule(checkpointInit, "checkpointInit")
scheduler.RegisterModule(forwardSync, "forwardSync")
scheduler.RegisterModule(headSync, "headSync")
scheduler.RegisterModule(beaconBlockSync, "beaconBlockSync")
scheduler.RegisterModule(engineApiUpdater, "engineApiUpdater")
// start
scheduler.Start()
// register server(s)