diff --git a/beacon/light/api/api_server.go b/beacon/light/api/api_server.go index a925bd3536..9b4ead07f2 100755 --- a/beacon/light/api/api_server.go +++ b/beacon/light/api/api_server.go @@ -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) diff --git a/beacon/light/request/request.go b/beacon/light/request/request.go index 15038ff20f..2a04b59212 100644 --- a/beacon/light/request/request.go +++ b/beacon/light/request/request.go @@ -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 } diff --git a/beacon/light/request/scheduler.go b/beacon/light/request/scheduler.go index 00c5aabaa7..87264235ed 100644 --- a/beacon/light/request/scheduler.go +++ b/beacon/light/request/scheduler.go @@ -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() } diff --git a/beacon/light/request/server.go b/beacon/light/request/server.go index 359a095554..8d13bfaf44 100644 --- a/beacon/light/request/server.go +++ b/beacon/light/request/server.go @@ -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 } diff --git a/beacon/light/sync/types.go b/beacon/light/sync/types.go index 5b2df1359a..e3c6ea76cd 100644 --- a/beacon/light/sync/types.go +++ b/beacon/light/sync/types.go @@ -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 ( diff --git a/cmd/blsync/main.go b/cmd/blsync/main.go index 555a4317b2..5bf0aa13bb 100644 --- a/cmd/blsync/main.go +++ b/cmd/blsync/main.go @@ -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)