From 84530c6b0b171892bd4eb180e505204edaf74874 Mon Sep 17 00:00:00 2001 From: zsfelfoldi Date: Sat, 30 Dec 2023 00:34:27 +0100 Subject: [PATCH] beacon/light: unexport types, use more interfaces --- beacon/light/request/request.go | 16 ++++--- beacon/light/request/scheduler.go | 47 +++++++++++------- beacon/light/request/server.go | 80 +++++++++++++++++-------------- beacon/light/sync/head_sync.go | 39 ++++++++------- beacon/light/sync/update_sync.go | 32 ++++++++----- cmd/blsync/block_sync.go | 12 ++--- cmd/blsync/main.go | 2 +- 7 files changed, 127 insertions(+), 101 deletions(-) diff --git a/beacon/light/request/request.go b/beacon/light/request/request.go index 2a04b59212..70a7e4ba04 100644 --- a/beacon/light/request/request.go +++ b/beacon/light/request/request.go @@ -27,30 +27,30 @@ type ( Response any ID uint64 ServerAndId struct { - Server Server + Server any Id ID } ) // one per sync process -type RequestTracker struct { +type tracker struct { servers serverSet // one per trigger scheduler *Scheduler module Module requestEvents []RequestEvent } -func (p *RequestTracker) TryRequest(requestFn func(server Server) (Request, float32)) (ServerAndId, Request) { +func (p *tracker) TryRequest(requestFn func(server any) (Request, float32)) (ServerAndId, Request) { var ( maxServerPriority, maxRequestPriority float32 - bestServer Server + bestServer server bestRequest Request ) maxServerPriority, maxRequestPriority = -math.MaxFloat32, -math.MaxFloat32 serverCount := len(p.servers) var removed, candidates int for server, _ := range p.servers { - canRequest, serverPriority := server.CanRequestNow() + canRequest, serverPriority := server.canRequestNow() if !canRequest { delete(p.servers, server) removed++ @@ -71,7 +71,11 @@ func (p *RequestTracker) TryRequest(requestFn func(server Server) (Request, floa if bestServer == nil { return ServerAndId{}, nil } - id := ServerAndId{Server: bestServer, Id: bestServer.SendRequest(bestRequest)} + id := ServerAndId{Server: bestServer, Id: bestServer.sendRequest(bestRequest)} p.scheduler.pending[id] = pendingRequest{request: bestRequest, module: p.module} return id, bestRequest } + +func (p *tracker) InvalidResponse(id ServerAndId, desc string) { + id.Server.(server).fail(desc) +} diff --git a/beacon/light/request/scheduler.go b/beacon/light/request/scheduler.go index c61000c731..00cecf6f13 100644 --- a/beacon/light/request/scheduler.go +++ b/beacon/light/request/scheduler.go @@ -37,7 +37,12 @@ type Module interface { // Note: Process functions of different modules are never called concurrently; // they are called by Scheduler in the same order of priority as they were // registered in. - Process(*RequestTracker, []RequestEvent, []ServerEvent) bool + Process(Tracker, []RequestEvent, []ServerEvent) bool +} + +type Tracker interface { + TryRequest(requestFn func(server any) (Request, float32)) (ServerAndId, Request) + InvalidResponse(id ServerAndId, desc string) } // Scheduler is a modular network data retrieval framework that coordinates multiple @@ -49,8 +54,8 @@ type Scheduler struct { clock mclock.Clock modules []Module // first has highest priority names map[Module]string - trackers map[Module]*RequestTracker - servers map[Server]struct{} + trackers map[Module]*tracker + servers map[server]struct{} pending map[ServerAndId]pendingRequest serverEvents []ServerEvent stopCh chan chan struct{} @@ -61,7 +66,7 @@ type Scheduler struct { } type ServerEvent struct { - Server Server + Server any Type string Data any } @@ -83,9 +88,9 @@ type pendingRequest struct { func NewScheduler(clock mclock.Clock) *Scheduler { s := &Scheduler{ clock: clock, - servers: make(map[Server]struct{}), + servers: make(map[server]struct{}), names: make(map[Module]string), - trackers: make(map[Module]*RequestTracker), + trackers: make(map[Module]*tracker), pending: make(map[ServerAndId]pendingRequest), stopCh: make(chan chan struct{}), // Note: testWaitCh should not have capacity in order to ensure @@ -105,7 +110,7 @@ func (s *Scheduler) RegisterModule(m Module, name string) { defer s.lock.Unlock() s.modules = append(s.modules, m) - s.trackers[m] = &RequestTracker{ + s.trackers[m] = &tracker{ scheduler: s, module: m, } @@ -113,12 +118,13 @@ func (s *Scheduler) RegisterModule(m Module, name string) { } // RegisterServer registers a new server. -func (s *Scheduler) RegisterServer(server Server) { +func (s *Scheduler) RegisterServer(rs requestServer) { s.lock.Lock() defer s.lock.Unlock() + server := newServer(rs, s.clock) s.handleEvent(server, Event{Type: EvRegistered}) - server.Subscribe(func(event Event) { + server.subscribe(func(event Event) { s.lock.Lock() if _, ok := s.servers[server]; ok { s.handleEvent(server, event) @@ -131,13 +137,18 @@ func (s *Scheduler) RegisterServer(server Server) { } // UnregisterServer removes a registered server. -func (s *Scheduler) UnregisterServer(server Server) { +func (s *Scheduler) UnregisterServer(rs requestServer) { s.lock.Lock() defer s.lock.Unlock() - server.Unsubscribe() - delete(s.servers, server) - s.handleEvent(server, Event{Type: EvUnregistered}) + for server := range s.servers { + if sl, ok := server.(*serverWithLimits); ok && sl.parent == rs { + server.unsubscribe() + delete(s.servers, server) + s.handleEvent(server, Event{Type: EvUnregistered}) + return + } + } } // Start starts the scheduler. It should be called after registering all modules @@ -150,7 +161,7 @@ func (s *Scheduler) Start() { func (s *Scheduler) Stop() { s.lock.Lock() for server, _ := range s.servers { - server.Unsubscribe() + server.unsubscribe() } s.servers = nil s.lock.Unlock() @@ -186,7 +197,7 @@ func (s *Scheduler) processModules() { s.lock.Lock() servers := make(serverSet) for server, _ := range s.servers { - if ok, _ := server.CanRequestNow(); ok { + if ok, _ := server.canRequestNow(); ok { servers[server] = struct{}{} } } @@ -233,7 +244,7 @@ func (s *Scheduler) Trigger() { } } -func (s *Scheduler) addRequestEvent(server Server, id ID, response Response, timeout, finalized bool) { +func (s *Scheduler) addRequestEvent(server any, id ID, response Response, timeout, finalized bool) { sid := ServerAndId{Server: server, Id: id} if pr, ok := s.pending[sid]; ok { tracker := s.trackers[pr.module] @@ -254,11 +265,11 @@ func (s *Scheduler) addRequestEvent(server Server, id ID, response Response, tim } } -func (s *Scheduler) addServerEvent(server Server, event Event) { +func (s *Scheduler) addServerEvent(server any, event Event) { s.serverEvents = append(s.serverEvents, ServerEvent{Server: server, Type: event.Type, Data: event.Data}) } -func (s *Scheduler) handleEvent(server Server, event Event) { +func (s *Scheduler) handleEvent(server any, event Event) { s.Trigger() switch event.Type { case EvResponse: diff --git a/beacon/light/request/server.go b/beacon/light/request/server.go index a72cb7816a..5ffe706686 100644 --- a/beacon/light/request/server.go +++ b/beacon/light/request/server.go @@ -28,13 +28,13 @@ import ( var ( // request events - EvResponse = "response" // data: IdAndResponse; sent by RequestServer - EvFail = "fail" // data: ID; sent by RequestServer + EvResponse = "response" // data: IdAndResponse; sent by requestServer + EvFail = "fail" // data: ID; sent by requestServer EvTimeout = "timeout" // data: ID; sent by serverWithTimeout // server events EvRegistered = "registered" // data: nil; sent by Scheduler EvUnregistered = "unregistered" // data: nil; sent by Scheduler - EvCanRequestAgain = "canRequestAgain" // data: nil; sent by serverWithLimits + EvCanRequestAgain = "canRequestAgain" // data: nil; sent by server ) const ( @@ -51,31 +51,35 @@ const ( maxFailureDelay = time.Minute ) -// RequestServer can send a set of requests pre-defined by the application and +// requestServer can send a set of requests pre-defined by the application and // signal events through the event callback. After each request, it should send // back either EvResponse or EvFail. Additionally, it may also send application- // defined events that the Modules can interpret. -type RequestServer interface { +// +//TODO ?separate request and server events here? +type requestServer interface { Subscribe(eventCallback func(event Event)) SendRequest(request Request) ID Unsubscribe() } -type Server interface { - RequestServer - CanRequestNow() (bool, float32) - Fail(desc string) +type server interface { + subscribe(eventCallback func(event Event)) + canRequestNow() (bool, float32) + sendRequest(request Request) ID + fail(desc string) + unsubscribe() } -func NewServer(rs RequestServer, clock mclock.Clock) Server { +func newServer(rs requestServer, clock mclock.Clock) server { s := &serverWithLimits{} - s.serverWithTimeout.RequestServer = rs + s.parent = rs s.serverWithTimeout.init(clock) s.init() return s } -type serverSet map[Server]struct{} +type serverSet map[server]struct{} type Event struct { Type string @@ -87,13 +91,13 @@ type IdAndResponse struct { Response Response } -// serverWithTimeout wraps a RequestServer and implements timeouts. After +// serverWithTimeout wraps a requestServer and implements timeouts. After // softRequestTimeout it sends an EvTimeout after which and EvResponse or an // EvFail will still follow (EvTimeout cannot follow the latter two). // After hardRequestTimeout it sends an EvFail and blocks any further events -// related to the given request coming from the parent RequestServer. +// related to the given request coming from the parent requestServer. type serverWithTimeout struct { - RequestServer + parent requestServer lock sync.Mutex clock mclock.Clock childEventCb func(event Event) @@ -105,12 +109,12 @@ func (s *serverWithTimeout) init(clock mclock.Clock) { s.timeouts = make(map[ID]mclock.Timer) } -func (s *serverWithTimeout) Subscribe(eventCallback func(event Event)) { +func (s *serverWithTimeout) subscribe(eventCallback func(event Event)) { s.lock.Lock() defer s.lock.Unlock() s.childEventCb = eventCallback - s.RequestServer.Subscribe(s.eventCallback) + s.parent.Subscribe(s.eventCallback) } func (s *serverWithTimeout) eventCallback(event Event) { @@ -137,11 +141,11 @@ func (s *serverWithTimeout) eventCallback(event Event) { } } -func (s *serverWithTimeout) SendRequest(request Request) (reqId ID) { +func (s *serverWithTimeout) sendRequest(request Request) (reqId ID) { s.lock.Lock() defer s.lock.Unlock() - reqId = s.RequestServer.SendRequest(request) + reqId = s.parent.SendRequest(request) s.timeouts[reqId] = s.clock.AfterFunc(softRequestTimeout, func() { /*if s.testTimerResults != nil { s.testTimerResults = append(s.testTimerResults, true) // simulated timer finished @@ -173,7 +177,7 @@ func (s *serverWithTimeout) SendRequest(request Request) (reqId ID) { } // stop stops all goroutines associated with the server. -func (s *serverWithTimeout) Unsubscribe() { +func (s *serverWithTimeout) unsubscribe() { s.lock.Lock() defer s.lock.Unlock() @@ -183,7 +187,7 @@ func (s *serverWithTimeout) Unsubscribe() { } } s.childEventCb = nil - s.RequestServer.Unsubscribe() + s.parent.Unsubscribe() } func (s *serverWithTimeout) stopTimer(timer mclock.Timer) { @@ -193,13 +197,15 @@ func (s *serverWithTimeout) stopTimer(timer mclock.Timer) { }*/ } -// serverWithLimits wraps serverWithTimeout and implements Server. It limits the +// serverWithLimits wraps serverWithTimeout and implements server. It limits the // number of parallel in-flight requests and prevents sending new requests when a // pending one has already timed out. It also implements a failure delay mechanism // that adds an exponentially growing delay each time a request fails (wrong answer // or hard timeout). This makes the syncing mechanism less brittle as temporary // failures of the server might happen sometimes, but still avoids hammering a // non-functional server with requests. +// +//TODO protect against excessive server events type serverWithLimits struct { serverWithTimeout lock sync.Mutex @@ -219,12 +225,12 @@ func (s *serverWithLimits) init() { s.parallelLimit = defaultParallelLimit } -func (s *serverWithLimits) Subscribe(eventCallback func(event Event)) { +func (s *serverWithLimits) subscribe(eventCallback func(event Event)) { s.lock.Lock() defer s.lock.Unlock() s.childEventCb = eventCallback - s.serverWithTimeout.Subscribe(s.eventCallback) + s.serverWithTimeout.subscribe(s.eventCallback) } func (s *serverWithLimits) eventCallback(event Event) { @@ -255,12 +261,12 @@ func (s *serverWithLimits) eventCallback(event Event) { s.parallelLimit += parallelAdjustUp } s.pendingCount-- - if canRequest, _ := s.canRequestNow(); canRequest { + if canRequest, _ := s.canRequest(); canRequest { sendCanRequestAgain = s.sendEvent s.sendEvent = false } if event.Type == EvFail { - s.fail("failed request") + s.failLocked("failed request") } } childEventCb := s.childEventCb @@ -271,17 +277,17 @@ func (s *serverWithLimits) eventCallback(event Event) { } } -func (s *serverWithLimits) SendRequest(request Request) (reqId ID) { +func (s *serverWithLimits) sendRequest(request Request) (reqId ID) { s.lock.Lock() defer s.lock.Unlock() s.pendingCount++ - id := s.serverWithTimeout.SendRequest(request) + id := s.serverWithTimeout.sendRequest(request) return id } // stop stops all goroutines associated with the server. -func (s *serverWithLimits) Unsubscribe() { +func (s *serverWithLimits) unsubscribe() { s.lock.Lock() defer s.lock.Unlock() @@ -290,10 +296,10 @@ func (s *serverWithLimits) Unsubscribe() { s.delayTimer = nil } s.childEventCb = nil - s.serverWithTimeout.Unsubscribe() + s.serverWithTimeout.unsubscribe() } -func (s *serverWithLimits) canRequestNow() (bool, float32) { +func (s *serverWithLimits) canRequest() (bool, float32) { if s.delayTimer != nil || s.pendingCount >= int(s.parallelLimit) { return false, 0 } @@ -304,10 +310,10 @@ func (s *serverWithLimits) canRequestNow() (bool, float32) { } // EvCanRequestAgain guaranteed if it returns false -func (s *serverWithLimits) CanRequestNow() (bool, float32) { +func (s *serverWithLimits) canRequestNow() (bool, float32) { var sendCanRequestAgain bool s.lock.Lock() - canRequest, priority := s.canRequestNow() + canRequest, priority := s.canRequest() if canRequest { sendCanRequestAgain = s.sendEvent s.sendEvent = false @@ -340,7 +346,7 @@ func (s *serverWithLimits) delay(delay time.Duration) { s.lock.Lock() if s.delayTimer != nil && s.delayCounter == delayCounter { // do nothing if there is a new timer now s.delayTimer = nil - if canRequest, _ := s.canRequestNow(); canRequest { + if canRequest, _ := s.canRequest(); canRequest { sendCanRequestAgain = s.sendEvent s.sendEvent = false } @@ -353,14 +359,14 @@ func (s *serverWithLimits) delay(delay time.Duration) { }) } -func (s *serverWithLimits) Fail(desc string) { +func (s *serverWithLimits) fail(desc string) { s.lock.Lock() defer s.lock.Unlock() - s.fail(desc) + s.failLocked(desc) } -func (s *serverWithLimits) fail(desc string) { +func (s *serverWithLimits) failLocked(desc string) { log.Debug("Server error", "description", desc) s.failureDelay *= 2 now := s.clock.Now() diff --git a/beacon/light/sync/head_sync.go b/beacon/light/sync/head_sync.go index 94a5e4e0ff..d0c19ea5d6 100644 --- a/beacon/light/sync/head_sync.go +++ b/beacon/light/sync/head_sync.go @@ -17,21 +17,24 @@ package sync import ( - "fmt" "math" - "github.com/ethereum/go-ethereum/beacon/light" "github.com/ethereum/go-ethereum/beacon/light/request" "github.com/ethereum/go-ethereum/beacon/types" ) +type headTracker interface { + Validate(head types.SignedHeader) (bool, error) + SetPrefetchHead(head types.HeadInfo) +} + type HeadSync struct { - headTracker *light.HeadTracker - chain *light.CommitteeChain + headTracker headTracker + chain committeeChain nextSyncPeriod uint64 chainInit bool - queuedHeads map[request.Server][]types.SignedHeader - serverHeads map[request.Server]types.HeadInfo + queuedHeads map[any][]types.SignedHeader + serverHeads map[any]types.HeadInfo headServerCount map[types.HeadInfo]headServerCount headCounter uint64 prefetchHead types.HeadInfo @@ -42,20 +45,20 @@ type headServerCount struct { headCounter uint64 } -func NewHeadSync(headTracker *light.HeadTracker, chain *light.CommitteeChain) *HeadSync { +func NewHeadSync(headTracker headTracker, chain committeeChain) *HeadSync { s := &HeadSync{ headTracker: headTracker, chain: chain, nextSyncPeriod: math.MaxUint64, - queuedHeads: make(map[request.Server][]types.SignedHeader), - serverHeads: make(map[request.Server]types.HeadInfo), + queuedHeads: make(map[any][]types.SignedHeader), + serverHeads: make(map[any]types.HeadInfo), headServerCount: make(map[types.HeadInfo]headServerCount), } return s } // Process implements request.Module -func (s *HeadSync) Process(tracker *request.RequestTracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { +func (s *HeadSync) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { nextPeriod, chainInit := s.chain.NextSyncPeriod() if nextPeriod != s.nextSyncPeriod || chainInit != s.chainInit { s.nextSyncPeriod, s.chainInit = nextPeriod, chainInit @@ -76,24 +79,20 @@ func (s *HeadSync) Process(tracker *request.RequestTracker, requestEvents []requ return } -func (s *HeadSync) newSignedHead(server request.Server, signedHead types.SignedHeader) { +func (s *HeadSync) newSignedHead(server any, signedHead types.SignedHeader) { if signedHead.Header.SyncPeriod() > s.nextSyncPeriod { - s.queuedHeads[server] = append(s.queuedHeads[server], signedHead) //TODO protect against future period spam + s.queuedHeads[server] = append(s.queuedHeads[server], signedHead) return } - if _, err := s.headTracker.Validate(signedHead); err != nil { - server.Fail(fmt.Sprintf("Invalid signed head: %v", err)) - } + s.headTracker.Validate(signedHead) } func (s *HeadSync) processQueuedHeads() { for server, queued := range s.queuedHeads { j := len(queued) for i := len(queued) - 1; i >= 0; i-- { - if signedHead := queued[i]; signedHead.Header.SyncPeriod() <= s.nextSyncPeriod { - if _, err := s.headTracker.Validate(signedHead); err != nil { - server.Fail(fmt.Sprintf("Invalid queued head: %v", err)) - } + if signedHead := queued[i]; types.SyncPeriod(signedHead.SignatureSlot) <= s.nextSyncPeriod { + s.headTracker.Validate(signedHead) } else { j-- if j != i { @@ -110,7 +109,7 @@ func (s *HeadSync) processQueuedHeads() { // setServerHead processes non-validated server head announcements and updates // the prefetch head if necessary. //TODO report server failure if a server announces many heads that do not become validated soon. -func (s *HeadSync) setServerHead(server request.Server, head types.HeadInfo) bool { +func (s *HeadSync) setServerHead(server any, head types.HeadInfo) bool { if oldHead, ok := s.serverHeads[server]; ok { if head == oldHead { return false diff --git a/beacon/light/sync/update_sync.go b/beacon/light/sync/update_sync.go index e9a1ef84ee..50378b2bad 100644 --- a/beacon/light/sync/update_sync.go +++ b/beacon/light/sync/update_sync.go @@ -28,14 +28,20 @@ import ( const maxUpdateRequest = 8 +type committeeChain interface { + CheckpointInit(bootstrap types.BootstrapData) error + InsertUpdate(update *types.LightClientUpdate, nextCommittee *types.SerializedSyncCommittee) error + NextSyncPeriod() (uint64, bool) +} + type CheckpointInit struct { - chain *light.CommitteeChain + chain committeeChain checkpointHash common.Hash pending bool initialized bool } -func NewCheckpointInit(chain *light.CommitteeChain, checkpointHash common.Hash) *CheckpointInit { +func NewCheckpointInit(chain committeeChain, checkpointHash common.Hash) *CheckpointInit { return &CheckpointInit{ chain: chain, checkpointHash: checkpointHash, @@ -43,7 +49,7 @@ func NewCheckpointInit(chain *light.CommitteeChain, checkpointHash common.Hash) } // Process implements request.Module -func (s *CheckpointInit) Process(tracker *request.RequestTracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) bool { +func (s *CheckpointInit) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) bool { if s.initialized { return false } @@ -57,11 +63,11 @@ func (s *CheckpointInit) Process(tracker *request.RequestTracker, requestEvents s.initialized = true return true } - event.Server.Fail("invalid checkpoint data") + tracker.InvalidResponse(event.ServerAndId, "invalid checkpoint data") } } if !s.pending { - if _, request := tracker.TryRequest(func(server request.Server) (request.Request, float32) { + if _, request := tracker.TryRequest(func(server any) (request.Request, float32) { return ReqCheckpointData(s.checkpointHash), 0 }); request != nil { s.pending = true @@ -74,14 +80,14 @@ type ForwardUpdateSync struct { chain *light.CommitteeChain rangeLock rangeLock processQueue []request.RequestEvent - nextSyncPeriod map[request.Server]uint64 + nextSyncPeriod map[any]uint64 } func NewForwardUpdateSync(chain *light.CommitteeChain) *ForwardUpdateSync { return &ForwardUpdateSync{ chain: chain, rangeLock: make(rangeLock), - nextSyncPeriod: make(map[request.Server]uint64), + nextSyncPeriod: make(map[any]uint64), } } @@ -138,7 +144,7 @@ func (s *ForwardUpdateSync) verifyRange(event request.RequestEvent) bool { } // returns true for partial success -func (s *ForwardUpdateSync) processResponse(event request.RequestEvent) (success bool) { +func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event request.RequestEvent) (success bool) { response, ok := event.Response.(RespUpdates) if !ok { return false @@ -151,7 +157,7 @@ func (s *ForwardUpdateSync) processResponse(event request.RequestEvent) (success return } if err == light.ErrInvalidUpdate || err == light.ErrWrongCommitteeRoot || err == light.ErrCannotReorg { - event.Server.Fail("invalid update received") + tracker.InvalidResponse(event.ServerAndId, "invalid update received") } else { log.Error("Unexpected InsertUpdate error", "error", err) } @@ -171,11 +177,11 @@ func (u updateResponseList) Less(i, j int) bool { } // Process implements request.Module -func (s *ForwardUpdateSync) Process(tracker *request.RequestTracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { +func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { // iterate events and add responses to process queue for _, event := range requestEvents { if event.Response != nil && !s.verifyRange(event) { - event.Server.Fail("invalid update range") + tracker.InvalidResponse(event.ServerAndId, "invalid update range") event.Response = nil } req := event.Request.(ReqUpdates) @@ -197,7 +203,7 @@ func (s *ForwardUpdateSync) Process(tracker *request.RequestTracker, requestEven sort.Sort(updateResponseList(s.processQueue)) //TODO for s.processQueue != nil { event := s.processQueue[0] - if !s.processResponse(event) { + if !s.processResponse(tracker, event) { break } trigger = true @@ -227,7 +233,7 @@ func (s *ForwardUpdateSync) Process(tracker *request.RequestTracker, requestEven } for { firstPeriod, maxCount := s.rangeLock.firstUnlocked(startPeriod, maxUpdateRequest) - if _, request := tracker.TryRequest(func(server request.Server) (request.Request, float32) { + if _, request := tracker.TryRequest(func(server any) (request.Request, float32) { nextPeriod := s.nextSyncPeriod[server] if nextPeriod <= firstPeriod { return nil, 0 diff --git a/cmd/blsync/block_sync.go b/cmd/blsync/block_sync.go index 5410de02db..19d585e95e 100755 --- a/cmd/blsync/block_sync.go +++ b/cmd/blsync/block_sync.go @@ -44,7 +44,7 @@ type beaconBlockSync struct { recentBlocks *lru.Cache[common.Hash, *capella.BeaconBlock] validatedHead types.Header pending map[common.Hash]struct{} - serverHeads map[request.Server]common.Hash + serverHeads map[any]common.Hash headTracker *light.HeadTracker } @@ -53,12 +53,12 @@ func newBeaconBlockSyncer(headTracker *light.HeadTracker) *beaconBlockSync { headTracker: headTracker, recentBlocks: lru.NewCache[common.Hash, *capella.BeaconBlock](10), pending: make(map[common.Hash]struct{}), - serverHeads: make(map[request.Server]common.Hash), + serverHeads: make(map[any]common.Hash), } } // Process implements request.Module -func (s *beaconBlockSync) Process(tracker *request.RequestTracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { +func (s *beaconBlockSync) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { s.validatedHead = s.headTracker.ValidatedHead().Header if s.validatedHead == (types.Header{}) { return false @@ -105,14 +105,14 @@ func (s *beaconBlockSync) getHeadBlock() *capella.BeaconBlock { return block } -func (s *beaconBlockSync) tryRequestBlock(tracker *request.RequestTracker, blockRoot common.Hash, prefetch bool) { +func (s *beaconBlockSync) tryRequestBlock(tracker request.Tracker, blockRoot common.Hash, prefetch bool) { if _, ok := s.recentBlocks.Get(blockRoot); ok { return } if _, ok := s.pending[blockRoot]; ok { return } - if _, request := tracker.TryRequest(func(server request.Server) (request.Request, float32) { + if _, request := tracker.TryRequest(func(server any) (request.Request, float32) { if prefetch && s.serverHeads[server] != blockRoot { // when requesting a not yet validated head, request it from someone // who has announced it already @@ -179,7 +179,7 @@ type engineApiUpdater struct { } // Process implements request.Module -func (s *engineApiUpdater) Process(tracker *request.RequestTracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) bool { +func (s *engineApiUpdater) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) bool { if atomic.LoadUint32(&s.updating) == 1 { return false } diff --git a/cmd/blsync/main.go b/cmd/blsync/main.go index 5bf0aa13bb..68b3d0a218 100644 --- a/cmd/blsync/main.go +++ b/cmd/blsync/main.go @@ -141,7 +141,7 @@ func blsync(ctx *cli.Context) error { // register server(s) for _, url := range ctx.StringSlice(utils.BeaconApiFlag.Name) { beaconApi := api.NewBeaconLightApi(url, customHeader) - scheduler.RegisterServer(request.NewServer(api.NewApiServer(beaconApi), &mclock.System{})) + scheduler.RegisterServer(api.NewApiServer(beaconApi)) } // run until stopped <-ctx.Done()