diff --git a/beacon/light/request/scheduler.go b/beacon/light/request/scheduler.go index eaa3a5eb62..ab0cd30ed7 100644 --- a/beacon/light/request/scheduler.go +++ b/beacon/light/request/scheduler.go @@ -17,6 +17,7 @@ package request import ( + "math" "sync" "github.com/ethereum/go-ethereum/common/mclock" @@ -46,7 +47,9 @@ 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(Tracker, []Event) + HandleEvent(Event) + Process() + MakeRequest(Server) (Request, float32) } // Scheduler is a modular network data retrieval framework that coordinates multiple @@ -55,22 +58,45 @@ type Module interface { // of existing data structures or events coming from registered servers could // allow new operations. type Scheduler struct { - lock sync.Mutex - clock mclock.Clock - modules []Module // first has highest priority - names map[Module]string - trackers map[Module]*tracker - servers map[server]struct{} - pending map[ServerAndID]pendingRequest - target map[targetData]uint64 - serverEvents []Event - stopCh chan chan struct{} + lock sync.Mutex + clock mclock.Clock + modules []Module // first has highest priority + names map[Module]string + servers map[server]struct{} + targets map[targetData]uint64 + + pending map[ServerAndID]pendingRequest + eventLock sync.Mutex + events []Event + stopCh chan chan struct{} triggerCh chan struct{} // restarts waiting sync loop // testWaitCh chan struct{} // accepts sends when sync loop is waiting // testTimerResults []bool // true is appended when simulated timer is processed; false when stopped } +type ( + // Server identifies a server without allowing any direct interaction. + // Note: server interface is used by Scheduler and Tracker but not used by + // the modules that do not interact with them directly. + // In order to make module testing easier, Server interface is used in + // events and modules. + Server interface { + Fail(desc string) + } + Request any + Response any + ID uint64 + ServerAndID struct { + Server Server + ID ID + } + RequestWithID struct { + ServerAndID + Request Request + } +) + type targetData interface { ChangeCounter() uint64 } @@ -85,13 +111,12 @@ type pendingRequest struct { // NewScheduler creates a new Scheduler. func NewScheduler(clock mclock.Clock) *Scheduler { s := &Scheduler{ - clock: clock, - servers: make(map[server]struct{}), - names: make(map[Module]string), - trackers: make(map[Module]*tracker), - pending: make(map[ServerAndID]pendingRequest), - target: make(map[targetData]uint64), - stopCh: make(chan chan struct{}), + clock: clock, + servers: make(map[server]struct{}), + names: make(map[Module]string), + pending: make(map[ServerAndID]pendingRequest), + targets: make(map[targetData]uint64), + stopCh: make(chan chan struct{}), // Note: testWaitCh should not have capacity in order to ensure // that after a trigger happens testWaitCh will block until the resulting // processing round has been finished @@ -105,7 +130,7 @@ func (s *Scheduler) RegisterTarget(t targetData) { s.lock.Lock() defer s.lock.Unlock() - s.target[t] = 0 + s.targets[t] = 0 } // RegisterModule registers a module. Should be called before starting the scheduler. @@ -116,10 +141,6 @@ func (s *Scheduler) RegisterModule(m Module, name string) { defer s.lock.Unlock() s.modules = append(s.modules, m) - s.trackers[m] = &tracker{ - scheduler: s, - module: m, - } s.names[m] = name } @@ -129,12 +150,12 @@ func (s *Scheduler) RegisterServer(rs requestServer) { defer s.lock.Unlock() server := newServer(rs, s.clock) - s.handleEvent(Event{Type: EvRegistered, Server: server}) + s.addEvent(Event{Type: EvRegistered, Server: server}) server.subscribe(func(event Event) { s.lock.Lock() if _, ok := s.servers[server]; ok { event.Server = server - s.handleEvent(event) + s.addEvent(event) } else { log.Error("Event received from unsubscribed server") } @@ -152,7 +173,7 @@ func (s *Scheduler) UnregisterServer(rs requestServer) { if sl, ok := server.(*serverWithLimits); ok && sl.parent == rs { server.unsubscribe() delete(s.servers, server) - s.handleEvent(Event{Type: EvUnregistered, Server: server}) + s.addEvent(Event{Type: EvUnregistered, Server: server}) return } } @@ -182,10 +203,13 @@ func (s *Scheduler) Stop() { // fired during a processing round ensure that there is going to be a next round. func (s *Scheduler) syncLoop() { for { - s.processModules() + s.lock.Lock() + s.handleEvents() for s.targetChanged() { s.processModules() } + s.sendRequests() + s.lock.Unlock() loop: for { select { @@ -201,9 +225,9 @@ func (s *Scheduler) syncLoop() { } func (s *Scheduler) targetChanged() (changed bool) { - for target, counter := range s.target { + for target, counter := range s.targets { if newCounter := target.ChangeCounter(); newCounter != counter { - s.target[target] = newCounter + s.targets[target] = newCounter changed = true } } @@ -213,47 +237,77 @@ func (s *Scheduler) targetChanged() (changed bool) { // processModules runs an entire processing round, calling the Process functions // of all modules, passing all relevant events. func (s *Scheduler) processModules() { - s.lock.Lock() + for _, module := range s.modules { + module.Process() + } +} + +func (s *Scheduler) sendRequests() { servers := make(serverSet) for server := range s.servers { if ok, _ := server.canRequestNow(); ok { servers[server] = struct{}{} } } - serverEvents := s.serverEvents - s.serverEvents = nil - s.lock.Unlock() - - eventTypes := make([]string, len(serverEvents)) - for i, ev := range serverEvents { - eventTypes[i] = ev.Type.Name - } - log.Debug("Processing modules", "servers", len(servers), "server events", eventTypes) + log.Debug("Processing modules", "servers", len(servers)) for _, module := range s.modules { - s.lock.Lock() - tracker := s.trackers[module] - tracker.servers = servers - requestEvents := tracker.requestEvents - tracker.requestEvents = nil - s.lock.Unlock() - - var respCount, failCount, timeoutCount int - for _, ev := range requestEvents { - switch ev.Type { - case EvResponse: - respCount++ - case EvFail: - failCount++ - case EvTimeout: - timeoutCount++ + for { + if len(servers) == 0 { + return + } + if req, sent := s.tryRequest(module, servers); sent { + module.HandleEvent(Event{ + Type: EvRequest, + Server: req.Server, + Data: RequestResponse{ + ID: req.ID, + Request: req.Request, + }, + }) + } else { + break } } - log.Debug("Processing module", "name", s.names[module], "responses", respCount, "fails", failCount, "timeouts", timeoutCount) - module.Process(tracker, append(serverEvents, requestEvents...)) } } +func (s *Scheduler) tryRequest(module Module, servers serverSet) (RequestWithID, bool) { + var ( + maxServerPriority, maxRequestPriority float32 + bestServer server + bestRequest Request + ) + maxServerPriority, maxRequestPriority = -math.MaxFloat32, -math.MaxFloat32 + serverCount := len(servers) + var removed, candidates int + for server := range servers { + canRequest, serverPriority := server.canRequestNow() + if !canRequest { + delete(servers, server) + removed++ + continue + } + request, requestPriority := module.MakeRequest(server) + if request != nil { + candidates++ + } + if request == nil || requestPriority < maxRequestPriority || + (requestPriority == maxRequestPriority && serverPriority <= maxServerPriority) { + continue + } + maxServerPriority, maxRequestPriority = serverPriority, requestPriority + bestServer, bestRequest = server, request + } + log.Debug("Request attempt", "serverCount", serverCount, "removedServers", removed, "requestCandidates", candidates) + if bestServer == nil { + return RequestWithID{}, false + } + id := ServerAndID{Server: bestServer, ID: bestServer.sendRequest(bestRequest)} + s.pending[id] = pendingRequest{request: bestRequest, module: module} + return RequestWithID{ServerAndID: id, Request: bestRequest}, true +} + // Trigger starts a new processing round. If fired during processing, it ensures // another full round of processing all modules. func (s *Scheduler) Trigger() { @@ -263,23 +317,21 @@ func (s *Scheduler) Trigger() { } } -// addRequestEvent adds a request event to the sender module's tracker, ensuring -// that the module receives it in the next processing round. -func (s *Scheduler) addRequestEvent(event Event) { - sid, _, _ := event.RequestInfo() - if pr, ok := s.pending[sid]; ok { - tracker := s.trackers[pr.module] - tracker.requestEvents = append(tracker.requestEvents, event) - if event.Type != EvTimeout { - delete(s.pending, sid) - } - } +func (s *Scheduler) addEvent(event Event) { + s.eventLock.Lock() + s.events = append(s.events, event) + s.Trigger() + s.eventLock.Unlock() } -// addServerEvent adds a server event to the global server event list, ensuring -// that all modules receive it in the next processing round. -func (s *Scheduler) addServerEvent(event Event) { - s.serverEvents = append(s.serverEvents, event) +func (s *Scheduler) handleEvents() { + s.eventLock.Lock() + events := s.events + s.events = nil + s.eventLock.Unlock() + for _, event := range events { + s.handleEvent(event) + } } // handleEvent processes an Event and adds it either as a request event or a @@ -288,9 +340,14 @@ func (s *Scheduler) addServerEvent(event Event) { // request event (EvFail), ensuring that all requests get finalized and thereby // allowing the module logic to be safe and simple. func (s *Scheduler) handleEvent(event Event) { - s.Trigger() if event.IsRequestEvent() { - s.addRequestEvent(event) + sid, _, _ := event.RequestInfo() + if pr, ok := s.pending[sid]; ok { + pr.module.HandleEvent(event) + if event.Type != EvTimeout { + delete(s.pending, sid) + } + } return } if event.Type == EvUnregistered { @@ -298,7 +355,7 @@ func (s *Scheduler) handleEvent(event Event) { if id.Server != event.Server { continue } - s.addRequestEvent(Event{ + pending.module.HandleEvent(Event{ Type: EvFail, Server: event.Server, Data: RequestResponse{ @@ -308,5 +365,7 @@ func (s *Scheduler) handleEvent(event Event) { }) } } - s.addServerEvent(event) + for _, module := range s.modules { + module.HandleEvent(event) + } } diff --git a/beacon/light/request/server.go b/beacon/light/request/server.go index 936c587783..cce747faf3 100644 --- a/beacon/light/request/server.go +++ b/beacon/light/request/server.go @@ -28,6 +28,7 @@ import ( var ( // request events + EvRequest = &EventType{Name: "request", requestEvent: true} // data: RequestResponse; sent by Scheduler EvResponse = &EventType{Name: "response", requestEvent: true} // data: RequestResponse; sent by requestServer EvFail = &EventType{Name: "fail", requestEvent: true} // data: RequestResponse; sent by requestServer EvTimeout = &EventType{Name: "timeout", requestEvent: true} // data: RequestResponse; sent by serverWithTimeout @@ -71,7 +72,7 @@ type server interface { subscribe(eventCallback func(event Event)) canRequestNow() (bool, float32) sendRequest(request Request) ID - fail(desc string) + Fail(desc string) unsubscribe() } @@ -299,7 +300,7 @@ func (s *serverWithLimits) eventCallback(event Event) { s.sendEvent = false } if event.Type == EvFail { - s.failLocked("failed request") + s.fail("failed request") } } childEventCb := s.childEventCb @@ -404,15 +405,15 @@ func (s *serverWithLimits) delay(delay time.Duration) { // fail reports that a response from the server was found invalid by the processing // Module, disabling new requests for a dynamically adjused time period. -func (s *serverWithLimits) fail(desc string) { +func (s *serverWithLimits) Fail(desc string) { s.lock.Lock() defer s.lock.Unlock() - s.failLocked(desc) + s.fail(desc) } -// failLocked calculates the dynamic failure delay and applies it. -func (s *serverWithLimits) failLocked(desc string) { +// fail calculates the dynamic failure delay and applies it. +func (s *serverWithLimits) fail(desc string) { log.Debug("Server error", "description", desc) s.failureDelay *= 2 now := s.clock.Now() diff --git a/beacon/light/request/tracker.go b/beacon/light/request/tracker.go deleted file mode 100644 index 79b91cf405..0000000000 --- a/beacon/light/request/tracker.go +++ /dev/null @@ -1,121 +0,0 @@ -// Copyright 2023 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 . - -package request - -import ( - "math" - - "github.com/ethereum/go-ethereum/log" -) - -type ( - // Server identifies a server without allowing any direct interaction. - // Note: server interface is used by Scheduler and Tracker but not used by - // the modules that do not interact with them directly. - // In order to make module testing easier, Server interface is used in - // events and modules. - Server any - Request any - Response any - ID uint64 - ServerAndID struct { - Server Server - ID ID - } - RequestWithID struct { - ServerAndID - Request Request - } -) - -// Tracker allows Modules to start requests and provide feedback about responses -// that were found to be invalid during processing. -type Tracker interface { - // TryRequest iterates through currently available servers and selects the - // best server and request to send. The caller provides a callback function - // that generates a request candidate for each available server. Note that - // the module may keep track of relevant server specific info, such as assumed - // available range of data to request, and therefore it may generate different - // request candidates for different servers. The callback also returns a - // priority value. TryRequest selects the request candidate with the highest - // priority value. If multiple candidates belonging to multiple servers have - // the same highest priority then it selects based on server priority. - // If a request candidate and a server has been selected, the request is sent - // and also returned along with the target server and request ID. - TryRequest(requestFn func(server Server) (Request, float32)) (RequestWithID, bool) - // InvalidResponse signals that the given response was invalid. Note that - // certain responses can only be judged by modules, in the context of existing, - // partially synced data structures. Giving this signal results in blocking - // the given server for a certain amount of time, ensuring that the same - // request will not be instantly sent again to the same server. - InvalidResponse(id ServerAndID, desc string) -} - -// tracker implements Tracker. A separate instance is created for each Module. -type tracker struct { - // servers is a set of currently available servers; it is recreated at every - // processModule round. - servers serverSet - scheduler *Scheduler - module Module - // requestEvents is a list of events related to requests sent by the given - // module. It is reset before module processing and the previous contents are - // passed to Module.Process along with the globally collected server events. - requestEvents []Event -} - -// TryRequest implements Tracker. -func (p *tracker) TryRequest(requestFn func(server Server) (Request, float32)) (RequestWithID, bool) { - var ( - maxServerPriority, maxRequestPriority float32 - 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() - 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 - } - maxServerPriority, maxRequestPriority = serverPriority, requestPriority - bestServer, bestRequest = server, request - } - log.Debug("Request attempt", "serverCount", serverCount, "removedServers", removed, "requestCandidates", candidates) - if bestServer == nil { - return RequestWithID{}, false - } - id := ServerAndID{Server: bestServer, ID: bestServer.sendRequest(bestRequest)} - p.scheduler.pending[id] = pendingRequest{request: bestRequest, module: p.module} - return RequestWithID{ServerAndID: id, Request: bestRequest}, true -} - -// InvalidResponse implements Tracker. -func (p *tracker) InvalidResponse(id ServerAndID, desc string) { - id.Server.(server).fail(desc) -} diff --git a/beacon/light/sync/head_sync.go b/beacon/light/sync/head_sync.go index 8d6426d151..61a4291e86 100644 --- a/beacon/light/sync/head_sync.go +++ b/beacon/light/sync/head_sync.go @@ -67,26 +67,29 @@ func NewHeadSync(headTracker headTracker, chain committeeChain) *HeadSync { return s } -// Process implements request.Module -func (s *HeadSync) Process(tracker request.Tracker, events []request.Event) { +func (s *HeadSync) HandleEvent(event request.Event) { + switch event.Type { + case EvNewHead: + s.setServerHead(event.Server, event.Data.(types.HeadInfo)) + case EvNewSignedHead: + s.newSignedHead(event.Server, event.Data.(types.SignedHeader)) + case request.EvUnregistered: + s.setServerHead(event.Server, types.HeadInfo{}) + delete(s.serverHeads, event.Server) + delete(s.unvalidatedHeads, event.Server) + } +} + +func (s *HeadSync) Process() { nextPeriod, chainInit := s.chain.NextSyncPeriod() if nextPeriod != s.nextSyncPeriod || chainInit != s.chainInit { s.nextSyncPeriod, s.chainInit = nextPeriod, chainInit s.processUnvalidatedHeads() } - for _, event := range events { - switch event.Type { - case EvNewHead: - s.setServerHead(event.Server, event.Data.(types.HeadInfo)) - case EvNewSignedHead: - s.newSignedHead(event.Server, event.Data.(types.SignedHeader)) - case request.EvUnregistered: - s.setServerHead(event.Server, types.HeadInfo{}) - delete(s.serverHeads, event.Server) - delete(s.unvalidatedHeads, event.Server) - } - } - return +} + +func (s *HeadSync) MakeRequest(server request.Server) (request.Request, float32) { + return nil, 0 } // newSignedHead handles received signed head; either validates it if the chain @@ -102,9 +105,6 @@ func (s *HeadSync) newSignedHead(server request.Server, signedHead types.SignedH // processUnvalidatedHeads iterates the list of unvalidated heads and validates // those which can be validated. func (s *HeadSync) processUnvalidatedHeads() { - if !s.chainInit { - return - } for server, signedHead := range s.unvalidatedHeads { if types.SyncPeriod(signedHead.SignatureSlot) <= s.nextSyncPeriod { s.headTracker.Validate(signedHead) diff --git a/beacon/light/sync/update_sync.go b/beacon/light/sync/update_sync.go index 417e23143b..724178a1cb 100644 --- a/beacon/light/sync/update_sync.go +++ b/beacon/light/sync/update_sync.go @@ -40,7 +40,7 @@ type committeeChain interface { type CheckpointInit struct { chain committeeChain checkpointHash common.Hash - locked bool + locked request.ServerAndID initialized bool } @@ -52,34 +52,35 @@ func NewCheckpointInit(chain committeeChain, checkpointHash common.Hash) *Checkp } } -// Process implements request.Module -func (s *CheckpointInit) Process(tracker request.Tracker, events []request.Event) { - if s.initialized { +func (s *CheckpointInit) HandleEvent(event request.Event) { + if !event.IsRequestEvent() { return } - for _, event := range events { - if !event.IsRequestEvent() { - continue - } - s.locked = false - sid, request, response := event.RequestInfo() - if response != nil { - if checkpoint, ok := response.(*types.BootstrapData); ok && checkpoint.Header.Hash() == common.Hash(request.(ReqCheckpointData)) { - s.chain.CheckpointInit(*checkpoint) //TODO - s.initialized = true - return - } - tracker.InvalidResponse(sid, "invalid checkpoint data") - } + sid, req, resp := event.RequestInfo() + if event.Type == request.EvRequest { + s.locked = sid + return } - if !s.locked { - if _, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { - return ReqCheckpointData(s.checkpointHash), 0 - }); ok { - s.locked = true - } + if s.locked == sid { + s.locked = request.ServerAndID{} } - return + if resp != nil { + if checkpoint, ok := resp.(*types.BootstrapData); ok && checkpoint.Header.Hash() == common.Hash(req.(ReqCheckpointData)) { + s.chain.CheckpointInit(*checkpoint) + s.initialized = true + return + } + event.Server.Fail("invalid checkpoint data") + } +} + +func (s *CheckpointInit) Process() {} + +func (s *CheckpointInit) MakeRequest(server request.Server) (request.Request, float32) { + if s.initialized || s.locked != (request.ServerAndID{}) { + return nil, 0 + } + return ReqCheckpointData(s.checkpointHash), 0 } // ForwardUpdateSync implements request.Module; it fetches updates between the @@ -190,8 +191,8 @@ func (s *ForwardUpdateSync) verifyRange(req request.Request, resp request.Respon // processResponse adds the fetched updates and committees to the committee chain. // Returns true in case of full or partial success. -func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event request.Event) (success bool) { - sid, _, resp := event.RequestInfo() +func (s *ForwardUpdateSync) processResponse(event request.Event) (success bool) { + _, _, resp := event.RequestInfo() response, ok := resp.(RespUpdates) if !ok { return false @@ -204,7 +205,7 @@ func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event reque return } if err == light.ErrInvalidUpdate || err == light.ErrWrongCommitteeRoot || err == light.ErrCannotReorg { - tracker.InvalidResponse(sid, "invalid update received") + event.Server.Fail("invalid update received") } else { log.Error("Unexpected InsertUpdate error", "error", err) } @@ -225,38 +226,39 @@ func (u updateResponseList) Less(i, j int) bool { u[j].Data.(request.RequestResponse).Request.(ReqUpdates).FirstPeriod } -// Process implements request.Module -func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Event) { - // iterate events and add responses to process queue - for _, event := range events { - switch event.Type { - case request.EvResponse, request.EvFail, request.EvTimeout: - sid, req, resp := event.RequestInfo() - if event.Type == request.EvResponse && !s.verifyRange(req, resp) { - tracker.InvalidResponse(sid, "invalid update range") - resp = nil - } - if resp != nil { - // there is a response with a valid format; put it in the process queue - s.processQueue = append(s.processQueue, event) - s.lockRange(sid, req) - } else { - s.unlockRange(sid, req) - } - case EvNewSignedHead: - signedHead := event.Data.(types.SignedHeader) - s.nextSyncPeriod[event.Server] = types.SyncPeriod(signedHead.SignatureSlot + 256) - case request.EvUnregistered: - delete(s.nextSyncPeriod, event.Server) +func (s *ForwardUpdateSync) HandleEvent(event request.Event) { + switch event.Type { + case request.EvRequest: + sid, req, _ := event.RequestInfo() + s.lockRange(sid, req) + case request.EvResponse, request.EvFail, request.EvTimeout: + sid, req, resp := event.RequestInfo() + if event.Type == request.EvResponse && !s.verifyRange(req, resp) { + event.Server.Fail("invalid update range") + resp = nil } + if resp != nil { + // there is a response with a valid format; put it in the process queue + s.processQueue = append(s.processQueue, event) + s.lockRange(sid, req) + } else { + s.unlockRange(sid, req) + } + case EvNewSignedHead: + signedHead := event.Data.(types.SignedHeader) + s.nextSyncPeriod[event.Server] = types.SyncPeriod(signedHead.SignatureSlot + 256) + case request.EvUnregistered: + delete(s.nextSyncPeriod, event.Server) } +} +func (s *ForwardUpdateSync) Process() { // try processing ordered list of available responses sort.Sort(updateResponseList(s.processQueue)) //TODO for s.processQueue != nil { event := s.processQueue[0] - if !s.processResponse(tracker, event) { - break + if !s.processResponse(event) { + return } sid, req, _ := event.RequestInfo() s.unlockRange(sid, req) @@ -265,28 +267,21 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Ev s.processQueue = nil } } +} - // start new requests if necessary +func (s *ForwardUpdateSync) MakeRequest(server request.Server) (request.Request, float32) { startPeriod, chainInit := s.chain.NextSyncPeriod() if !chainInit { - return + return nil, 0 } - for { - firstPeriod, maxCount := s.rangeLock.firstUnlocked(startPeriod, maxUpdateRequest) - if reqWithID, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { - nextPeriod := s.nextSyncPeriod[server] - if nextPeriod <= firstPeriod { - return nil, 0 - } - count := maxCount - if nextPeriod < firstPeriod+maxCount { - count = nextPeriod - firstPeriod - } - return ReqUpdates{FirstPeriod: firstPeriod, Count: count}, float32(count) - }); ok { - s.lockRange(reqWithID.ServerAndID, reqWithID.Request) - } else { - break - } + firstPeriod, maxCount := s.rangeLock.firstUnlocked(startPeriod, maxUpdateRequest) + nextPeriod := s.nextSyncPeriod[server] + if nextPeriod <= firstPeriod { + return nil, 0 } + count := maxCount + if nextPeriod < firstPeriod+maxCount { + count = nextPeriod - firstPeriod + } + return ReqUpdates{FirstPeriod: firstPeriod, Count: count}, float32(count) } diff --git a/cmd/blsync/block_sync.go b/cmd/blsync/block_sync.go index ffbbf9da2b..8e3d99b0db 100755 --- a/cmd/blsync/block_sync.go +++ b/cmd/blsync/block_sync.go @@ -53,27 +53,25 @@ func newBeaconBlockSync(headTracker headTracker) *beaconBlockSync { } } -// Process implements request.Module -func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Event) { - // iterate events and add valid responses to recentBlocks - for _, event := range events { - switch event.Type { - case request.EvResponse, request.EvFail, request.EvTimeout: - _, req, resp := event.RequestInfo() - blockRoot := common.Hash(req.(sync.ReqBeaconBlock)) - if resp != nil { - block := resp.(*capella.BeaconBlock) - s.recentBlocks.Add(blockRoot, block) - } - delete(s.locked, blockRoot) - case sync.EvNewHead: - s.serverHeads[event.Server] = event.Data.(types.HeadInfo).BlockRoot - case request.EvUnregistered: - delete(s.serverHeads, event.Server) +func (s *beaconBlockSync) HandleEvent(event request.Event) { + switch event.Type { + case request.EvResponse, request.EvFail, request.EvTimeout: + _, req, resp := event.RequestInfo() + blockRoot := common.Hash(req.(sync.ReqBeaconBlock)) + if resp != nil { + block := resp.(*capella.BeaconBlock) + s.recentBlocks.Add(blockRoot, block) } + delete(s.locked, blockRoot) + case sync.EvNewHead: + s.serverHeads[event.Server] = event.Data.(types.HeadInfo).BlockRoot + case request.EvUnregistered: + delete(s.serverHeads, event.Server) } +} - // send validated head block or request it if unavailable +func (s *beaconBlockSync) Process() { + // send validated head block if vh := s.headTracker.ValidatedHead(); vh != (types.SignedHeader{}) { validatedHead := vh.Header.Hash() if headBlock, ok := s.recentBlocks.Get(validatedHead); ok && headBlock != s.lastHeadBlock { @@ -82,37 +80,27 @@ func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Even s.lastHeadBlock = headBlock default: } - } else { - s.tryRequestBlock(tracker, validatedHead, false) } } - // request prefetch head - if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead != (common.Hash{}) { - s.tryRequestBlock(tracker, prefetchHead, true) - } } -// tryRequestBlock tries to send a block request for the given root if the block -// is not available and the root is not locked by another pending request. -// If prefetch is true then the request is only sent to a server whose latest -// announced head has the same block root. If prefetch is false then a validated -// block is requested which is expected to be available at every properly synced -// server, therefore no such restriction is applied. -func (s *beaconBlockSync) tryRequestBlock(tracker request.Tracker, blockRoot common.Hash, prefetch bool) { - if _, ok := s.recentBlocks.Get(blockRoot); ok { - return - } - if _, ok := s.locked[blockRoot]; ok { - return - } - if _, ok := tracker.TryRequest(func(server request.Server) (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 - return nil, 0 +func (s *beaconBlockSync) MakeRequest(server request.Server) (request.Request, float32) { + // request validated head block if unavailable and not yet requested + if vh := s.headTracker.ValidatedHead(); vh != (types.SignedHeader{}) { + validatedHead := vh.Header.Hash() + if _, ok := s.recentBlocks.Get(validatedHead); !ok { + if _, ok := s.locked[validatedHead]; !ok { + return sync.ReqBeaconBlock(validatedHead), 1 + } } - return sync.ReqBeaconBlock(blockRoot), 0 - }); ok { - s.locked[blockRoot] = struct{}{} } + // request prefetch head if the given server has announced it + if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead != (common.Hash{}) && prefetchHead != s.serverHeads[server] { + if _, ok := s.recentBlocks.Get(prefetchHead); !ok { + if _, ok := s.locked[prefetchHead]; !ok { + return sync.ReqBeaconBlock(prefetchHead), 0 + } + } + } + return nil, 0 }