From 26813d2a6c5758afb0d2490c97bc8be8a48620bb Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Wed, 10 Jan 2024 02:28:33 +0100 Subject: [PATCH] beacon/light: refactored events --- beacon/light/api/api_server.go | 4 +- beacon/light/request/scheduler.go | 127 +++++------------ beacon/light/request/server.go | 57 ++++---- .../light/request/{request.go => tracker.go} | 25 +++- beacon/light/sync/head_sync.go | 8 +- beacon/light/sync/head_sync_test.go | 36 ++--- beacon/light/sync/test_helpers.go | 12 ++ beacon/light/sync/types.go | 7 +- beacon/light/sync/update_sync.go | 119 +++++++++------- beacon/light/sync/update_sync_test.go | 133 +++++++++--------- cmd/blsync/block_sync.go | 41 +++--- cmd/blsync/block_sync_test.go | 30 ++-- 12 files changed, 299 insertions(+), 300 deletions(-) rename beacon/light/request/{request.go => tracker.go} (63%) diff --git a/beacon/light/api/api_server.go b/beacon/light/api/api_server.go index 9b4ead07f2..92fe494278 100755 --- a/beacon/light/api/api_server.go +++ b/beacon/light/api/api_server.go @@ -78,9 +78,9 @@ func (s *ApiServer) SendRequest(req request.Request) request.ID { default: } if resp != nil { - s.eventCallback(request.Event{Type: request.EvResponse, Data: request.IdAndResponse{ID: id, Response: resp}}) + s.eventCallback(request.Event{Type: request.EvResponse, Data: request.RequestResponse{ID: id, Request: req, Response: resp}}) } else { - s.eventCallback(request.Event{Type: request.EvFail, Data: id}) + s.eventCallback(request.Event{Type: request.EvFail, Data: request.RequestResponse{ID: id, Request: req}}) } }() return id diff --git a/beacon/light/request/scheduler.go b/beacon/light/request/scheduler.go index 3cc7552a96..1de4c8aa41 100644 --- a/beacon/light/request/scheduler.go +++ b/beacon/light/request/scheduler.go @@ -48,30 +48,7 @@ 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, []RequestEvent, []ServerEvent) bool -} - -// 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) + Process(Tracker, []Event) bool } // Scheduler is a modular network data retrieval framework that coordinates multiple @@ -86,7 +63,7 @@ type Scheduler struct { trackers map[Module]*tracker servers map[server]struct{} pending map[ServerAndID]pendingRequest - serverEvents []ServerEvent + serverEvents []Event stopCh chan chan struct{} triggerCh chan struct{} // restarts waiting sync loop @@ -94,36 +71,11 @@ type Scheduler struct { // testTimerResults []bool // true is appended when simulated timer is processed; false when stopped } -// ServerEvent represents a server-related event. These events are passed to all -// modules. Scheduler generates an EvRegister and EvUnregister event for each -// server when added or removed. Other, application-specific server events may -// be emitted by the servers themselves and are also passed to the modules. -type ServerEvent struct { - Server Server - Type string - Data any // data type defined by application-specific events -} - -// RequestEvent represents a request-related event. These events are passed to -// the module that sent the given request. A Finalized event means either a -// response or a hard timeout (depending on whether there is also a Response), -// after which no further event related to the same request is emitted. -// Once a module has successfully sent a request, Scheduler guarantees that it -// receives exactly one Finalized event related to it. If a request reaches a -// soft timeout, a Timeout event is emitted that is not Finalized yet. In this -// case later the Finalized event will also have its Timeout flag set. -type RequestEvent struct { - RequestWithID - Response Response - Timeout, Finalized bool -} - // pendingRequest keeps track of sent and not finalized requests and their sender // modules and whether a soft timeout has already happened. type pendingRequest struct { request Request module Module - timeout bool } // NewScheduler creates a new Scheduler. @@ -165,11 +117,12 @@ func (s *Scheduler) RegisterServer(rs requestServer) { defer s.lock.Unlock() server := newServer(rs, s.clock) - s.handleEvent(server, Event{Type: EvRegistered}) + s.handleEvent(Event{Type: EvRegistered, Server: server}) server.subscribe(func(event Event) { s.lock.Lock() if _, ok := s.servers[server]; ok { - s.handleEvent(server, event) + event.Server = server + s.handleEvent(event) } else { log.Error("Event received from unsubscribed server") } @@ -187,7 +140,7 @@ func (s *Scheduler) UnregisterServer(rs requestServer) { if sl, ok := server.(*serverWithLimits); ok && sl.parent == rs { server.unsubscribe() delete(s.servers, server) - s.handleEvent(server, Event{Type: EvUnregistered}) + s.handleEvent(Event{Type: EvUnregistered, Server: server}) return } } @@ -248,7 +201,7 @@ func (s *Scheduler) processModules() { eventTypes := make([]string, len(serverEvents)) for i, ev := range serverEvents { - eventTypes[i] = ev.Type + eventTypes[i] = ev.Type.Name } log.Debug("Processing modules", "servers", len(servers), "server events", eventTypes) @@ -262,17 +215,18 @@ func (s *Scheduler) processModules() { var respCount, failCount, timeoutCount int for _, ev := range requestEvents { - if ev.Response != nil { + switch ev.Type { + case EvResponse: respCount++ - } else if ev.Finalized { + case EvFail: failCount++ - } else { + case EvTimeout: timeoutCount++ } } log.Debug("Processing module", "name", s.names[module], "responses", respCount, "fails", failCount, "timeouts", timeoutCount) - if module.Process(tracker, requestEvents, serverEvents) { + if module.Process(tracker, append(serverEvents, requestEvents...)) { s.Trigger() } } @@ -289,24 +243,12 @@ 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(server Server, id ID, response Response, timeout, finalized bool) { - sid := ServerAndID{Server: server, ID: id} +func (s *Scheduler) addRequestEvent(event Event) { + sid, _, _ := event.RequestInfo() if pr, ok := s.pending[sid]; ok { tracker := s.trackers[pr.module] - timeout = timeout || pr.timeout - tracker.requestEvents = append(tracker.requestEvents, RequestEvent{ - RequestWithID: RequestWithID{ - ServerAndID: sid, - Request: pr.request, - }, - Response: response, - Timeout: timeout, - Finalized: finalized, - }) - if timeout && !finalized { - pr.timeout = true - s.pending[sid] = pr - } else { + tracker.requestEvents = append(tracker.requestEvents, event) + if event.Type != EvTimeout { delete(s.pending, sid) } } @@ -314,8 +256,8 @@ func (s *Scheduler) addRequestEvent(server Server, id ID, response Response, tim // 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(server Server, event Event) { - s.serverEvents = append(s.serverEvents, ServerEvent{Server: server, Type: event.Type, Data: event.Data}) +func (s *Scheduler) addServerEvent(event Event) { + s.serverEvents = append(s.serverEvents, event) } // handleEvent processes an Event and adds it either as a request event or a @@ -323,25 +265,26 @@ func (s *Scheduler) addServerEvent(server Server, event Event) { // it also closes all pending requests to the given server by emitting a failed // request event (Finalized without Response), ensuring that all requests get // finalized and thereby allowing the module logic to be safe and simple. -func (s *Scheduler) handleEvent(server Server, event Event) { +func (s *Scheduler) handleEvent(event Event) { s.Trigger() - switch event.Type { - case EvResponse: - idr := event.Data.(IdAndResponse) - s.addRequestEvent(server, idr.ID, idr.Response, false, true) - case EvFail: - s.addRequestEvent(server, event.Data.(ID), nil, false, true) - case EvTimeout: - s.addRequestEvent(server, event.Data.(ID), nil, true, false) - case EvUnregistered: - for id, _ := range s.pending { - if id.Server != server { + if event.IsRequestEvent() { + s.addRequestEvent(event) + return + } + if event.Type == EvUnregistered { + for id, pending := range s.pending { + if id.Server != event.Server { continue } - s.addRequestEvent(server, id.ID, nil, false, true) + s.addRequestEvent(Event{ + Type: EvFail, + Server: event.Server, + Data: RequestResponse{ + ID: id.ID, + Request: pending.request, + }, + }) } - s.addServerEvent(server, event) - default: - s.addServerEvent(server, event) } + s.addServerEvent(event) } diff --git a/beacon/light/request/server.go b/beacon/light/request/server.go index 5ffe706686..3718db6f83 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 - EvTimeout = "timeout" // data: ID; sent by serverWithTimeout + EvResponse = &EventType{Name: "response", requestEvent: true} // data: RequestResponse + EvFail = &EventType{Name: "fail", requestEvent: true} // data: RequestResponse + EvTimeout = &EventType{Name: "timeout", requestEvent: true} // data: RequestResponse // server events - EvRegistered = "registered" // data: nil; sent by Scheduler - EvUnregistered = "unregistered" // data: nil; sent by Scheduler - EvCanRequestAgain = "canRequestAgain" // data: nil; sent by server + EvRegistered = &EventType{Name: "registered"} // data: nil; sent by Scheduler + EvUnregistered = &EventType{Name: "unregistered"} // data: nil; sent by Scheduler + EvCanRequestAgain = &EventType{Name: "canRequestAgain"} // data: nil; sent by serverWithLimits ) const ( @@ -81,13 +81,29 @@ func newServer(rs requestServer, clock mclock.Clock) server { type serverSet map[server]struct{} -type Event struct { - Type string - Data any +type EventType struct { + Name string + requestEvent bool } -type IdAndResponse struct { +type Event struct { + Type *EventType + Server Server // filled by Scheduler + Data any +} + +func (e *Event) IsRequestEvent() bool { + return e.Type.requestEvent +} + +func (e *Event) RequestInfo() (ServerAndID, Request, Response) { + data := e.Data.(RequestResponse) + return ServerAndID{Server: e.Server, ID: data.ID}, data.Request, data.Response +} + +type RequestResponse struct { ID ID + Request Request Response Response } @@ -123,12 +139,7 @@ func (s *serverWithTimeout) eventCallback(event Event) { switch event.Type { case EvResponse, EvFail: - var id ID - if event.Type == EvResponse { - id = event.Data.(IdAndResponse).ID - } else { - id = event.Data.(ID) - } + id := event.Data.(RequestResponse).ID if timer, ok := s.timeouts[id]; ok { // Note: if stopping the timer is unsuccessful then the resulting AfterFunc // call will just do nothing @@ -167,11 +178,11 @@ func (s *serverWithTimeout) sendRequest(request Request) (reqId ID) { delete(s.timeouts, reqId) childEventCb := s.childEventCb s.lock.Unlock() - childEventCb(Event{Type: EvFail, Data: reqId}) + childEventCb(Event{Type: EvFail, Data: RequestResponse{ID: reqId, Request: request}}) }) childEventCb := s.childEventCb s.lock.Unlock() - childEventCb(Event{Type: EvTimeout, Data: reqId}) + childEventCb(Event{Type: EvTimeout, Data: RequestResponse{ID: reqId, Request: request}}) }) return reqId } @@ -238,7 +249,8 @@ func (s *serverWithLimits) eventCallback(event Event) { var sendCanRequestAgain bool switch event.Type { case EvTimeout: - s.softTimeouts[event.Data.(ID)] = struct{}{} + id := event.Data.(RequestResponse).ID + s.softTimeouts[id] = struct{}{} s.timeoutCount++ s.parallelLimit -= parallelAdjustDown if s.parallelLimit < minParallelLimit { @@ -246,12 +258,7 @@ func (s *serverWithLimits) eventCallback(event Event) { } log.Debug("Server timeout", "count", s.timeoutCount, "parallelLimit", s.parallelLimit) case EvResponse, EvFail: - var id ID - if event.Type == EvResponse { - id = event.Data.(IdAndResponse).ID - } else { - id = event.Data.(ID) - } + id := event.Data.(RequestResponse).ID if _, ok := s.softTimeouts[id]; ok { delete(s.softTimeouts, id) s.timeoutCount-- diff --git a/beacon/light/request/request.go b/beacon/light/request/tracker.go similarity index 63% rename from beacon/light/request/request.go rename to beacon/light/request/tracker.go index 36128b468f..12099e951d 100644 --- a/beacon/light/request/request.go +++ b/beacon/light/request/tracker.go @@ -37,12 +37,35 @@ type ( } ) +// 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) +} + // one per sync process type tracker struct { servers serverSet // one per trigger scheduler *Scheduler module Module - requestEvents []RequestEvent + requestEvents []Event } func (p *tracker) TryRequest(requestFn func(server Server) (Request, float32)) (RequestWithID, bool) { diff --git a/beacon/light/sync/head_sync.go b/beacon/light/sync/head_sync.go index b0988f9006..d69d382ece 100644 --- a/beacon/light/sync/head_sync.go +++ b/beacon/light/sync/head_sync.go @@ -58,13 +58,13 @@ func NewHeadSync(headTracker headTracker, chain committeeChain) *HeadSync { } // Process implements request.Module -func (s *HeadSync) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { +func (s *HeadSync) Process(tracker request.Tracker, events []request.Event) (trigger bool) { nextPeriod, chainInit := s.chain.NextSyncPeriod() if nextPeriod != s.nextSyncPeriod || chainInit != s.chainInit { s.nextSyncPeriod, s.chainInit = nextPeriod, chainInit - trigger = s.processUnvalidatedHeadsHeads() + trigger = s.processUnvalidatedHeads() } - for _, event := range serverEvents { + for _, event := range events { switch event.Type { case EvNewHead: if s.setServerHead(event.Server, event.Data.(types.HeadInfo)) { @@ -94,7 +94,7 @@ func (s *HeadSync) newSignedHead(server request.Server, signedHead types.SignedH return updated } -func (s *HeadSync) processUnvalidatedHeadsHeads() (trigger bool) { +func (s *HeadSync) processUnvalidatedHeads() (trigger bool) { if !s.chainInit { return false } diff --git a/beacon/light/sync/head_sync_test.go b/beacon/light/sync/head_sync_test.go index ec2e553514..993de7e6e7 100644 --- a/beacon/light/sync/head_sync_test.go +++ b/beacon/light/sync/head_sync_test.go @@ -50,38 +50,38 @@ func TestValidatedHead(t *testing.T) { headSync := NewHeadSync(ht, chain) ht.ExpValidated(t, 1, nil) - ExpTrigger(t, 1, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 1, false, headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: testSHead1}, })) ht.ExpValidated(t, 2, nil) chain.SetNextSyncPeriod(0) - ExpTrigger(t, 2, true, headSync.Process(tracker, nil, nil)) + ExpTrigger(t, 2, true, headSync.Process(tracker, nil)) ht.ExpValidated(t, 3, []types.SignedHeader{testSHead1}) chain.SetNextSyncPeriod(1) - ExpTrigger(t, 3, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 3, true, headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: EvNewSignedHead, Data: testSHead2}, {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewSignedHead, Data: testSHead2}, })) ht.ExpValidated(t, 4, []types.SignedHeader{testSHead2, testSHead2}) - ExpTrigger(t, 4, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 4, false, headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: EvNewSignedHead, Data: testSHead3}, {Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: EvNewSignedHead, Data: testSHead4}, })) ht.ExpValidated(t, 5, nil) chain.SetNextSyncPeriod(2) - ExpTrigger(t, 5, true, headSync.Process(tracker, nil, nil)) + ExpTrigger(t, 5, true, headSync.Process(tracker, nil)) ht.ExpValidated(t, 6, []types.SignedHeader{testSHead3}) - ExpTrigger(t, 6, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 6, false, headSync.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvUnregistered}, })) ht.ExpValidated(t, 7, nil) chain.SetNextSyncPeriod(3) - ExpTrigger(t, 7, false, headSync.Process(tracker, nil, nil)) + ExpTrigger(t, 7, false, headSync.Process(tracker, nil)) ht.ExpValidated(t, 8, nil) - ExpTrigger(t, 8, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 8, true, headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: EvNewSignedHead, Data: testSHead4}, })) ht.ExpValidated(t, 9, []types.SignedHeader{testSHead4}) @@ -94,47 +94,47 @@ func TestPrefetchHead(t *testing.T) { headSync := NewHeadSync(ht, chain) ht.ExpPrefetch(t, 1, testHead0) // no servers registered - ExpTrigger(t, 1, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 1, true, headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewHead, Data: testHead1}, })) ht.ExpPrefetch(t, 2, testHead1) // s1: h1 - ExpTrigger(t, 2, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 2, true, headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewHead, Data: testHead2}, })) ht.ExpPrefetch(t, 3, testHead2) // s1: h1, s2: h2 - ExpTrigger(t, 3, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 3, false, headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: EvNewHead, Data: testHead2}, })) ht.ExpPrefetch(t, 4, testHead2) // s1: h2, s2: h2 - ExpTrigger(t, 4, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 4, false, headSync.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: EvNewHead, Data: testHead3}, })) ht.ExpPrefetch(t, 5, testHead2) // s1: h2, s2: h2, s3: h3 - ExpTrigger(t, 5, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 5, false, headSync.Process(tracker, []request.Event{ {Server: testServer4, Type: request.EvRegistered}, {Server: testServer4, Type: EvNewHead, Data: testHead4}, })) ht.ExpPrefetch(t, 6, testHead2) // s1: h2, s2: h2, s3: h3, s4: h4 - ExpTrigger(t, 6, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 6, true, headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: EvNewHead, Data: testHead3}, })) ht.ExpPrefetch(t, 7, testHead3) // s1: h2, s2: h3, s3: h3, s4: h4 - ExpTrigger(t, 7, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 7, true, headSync.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvUnregistered}, })) ht.ExpPrefetch(t, 8, testHead4) // s1: h2, s2: h3, s4: h4 - ExpTrigger(t, 8, false, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 8, false, headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvUnregistered}, })) ht.ExpPrefetch(t, 9, testHead4) // s2: h3, s4: h4 - ExpTrigger(t, 9, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 9, true, headSync.Process(tracker, []request.Event{ {Server: testServer4, Type: request.EvUnregistered}, })) ht.ExpPrefetch(t, 10, testHead3) // s2: h3 - ExpTrigger(t, 10, true, headSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 10, true, headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: request.EvUnregistered}, })) ht.ExpPrefetch(t, 11, testHead0) // no servers registered diff --git a/beacon/light/sync/test_helpers.go b/beacon/light/sync/test_helpers.go index efd2f866be..c393a9c49f 100644 --- a/beacon/light/sync/test_helpers.go +++ b/beacon/light/sync/test_helpers.go @@ -97,6 +97,18 @@ func ExpTrigger(t *testing.T, tci int, expTrigger, trigger bool) { } } +func TestReqEvent(evType *request.EventType, req request.RequestWithID, response request.Response) request.Event { + return request.Event{ + Type: evType, + Server: req.ServerAndID.Server, + Data: request.RequestResponse{ + ID: req.ServerAndID.ID, + Request: req.Request, + Response: response, + }, + } +} + type TestCommitteeChain struct { fsp, nsp uint64 init bool diff --git a/beacon/light/sync/types.go b/beacon/light/sync/types.go index e3c6ea76cd..5e482aa268 100644 --- a/beacon/light/sync/types.go +++ b/beacon/light/sync/types.go @@ -17,13 +17,14 @@ 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 = "newHead" - EvNewSignedHead = "newSignedHead" +var ( + EvNewHead = &request.EventType{Name: "newHead"} // data: types.HeadInfo + EvNewSignedHead = &request.EventType{Name: "newSignedHead"} // data: types.SignedHeader ) type ( diff --git a/beacon/light/sync/update_sync.go b/beacon/light/sync/update_sync.go index 5f626f2f8b..92b4e7f840 100644 --- a/beacon/light/sync/update_sync.go +++ b/beacon/light/sync/update_sync.go @@ -37,7 +37,7 @@ type committeeChain interface { type CheckpointInit struct { chain committeeChain checkpointHash common.Hash - pending bool + locked bool initialized bool } @@ -49,28 +49,30 @@ func NewCheckpointInit(chain committeeChain, checkpointHash common.Hash) *Checkp } // Process implements request.Module -func (s *CheckpointInit) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) bool { +func (s *CheckpointInit) Process(tracker request.Tracker, events []request.Event) bool { if s.initialized { return false } - for _, event := range requestEvents { - if event.Timeout != event.Finalized { - s.pending = false + for _, event := range events { + if !event.IsRequestEvent() { + continue } - if event.Response != nil { - if checkpoint, ok := event.Response.(*types.BootstrapData); ok && checkpoint.Header.Hash() == common.Hash(event.Request.(ReqCheckpointData)) { + 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 true } - tracker.InvalidResponse(event.ServerAndID, "invalid checkpoint data") + tracker.InvalidResponse(sid, "invalid checkpoint data") } } - if !s.pending { + if !s.locked { if _, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { return ReqCheckpointData(s.checkpointHash), 0 }); ok { - s.pending = true + s.locked = true } } return false @@ -79,7 +81,8 @@ func (s *CheckpointInit) Process(tracker request.Tracker, requestEvents []reques type ForwardUpdateSync struct { chain committeeChain rangeLock rangeLock - processQueue []request.RequestEvent + lockedIDs map[request.ServerAndID]struct{} + processQueue []request.Event nextSyncPeriod map[request.Server]uint64 } @@ -87,6 +90,7 @@ func NewForwardUpdateSync(chain committeeChain) *ForwardUpdateSync { return &ForwardUpdateSync{ chain: chain, rangeLock: make(rangeLock), + lockedIDs: make(map[request.ServerAndID]struct{}), nextSyncPeriod: make(map[request.Server]uint64), } } @@ -123,12 +127,30 @@ func (r rangeLock) firstUnlocked(start, maxCount uint64) (first, count uint64) { return } -func (s *ForwardUpdateSync) verifyRange(event request.RequestEvent) bool { - request, ok := event.Request.(ReqUpdates) +func (s *ForwardUpdateSync) lockRange(sid request.ServerAndID, req request.Request) { + if _, ok := s.lockedIDs[sid]; ok { + return + } + s.lockedIDs[sid] = struct{}{} + r := req.(ReqUpdates) + s.rangeLock.lock(r.FirstPeriod, r.Count, 1) +} + +func (s *ForwardUpdateSync) unlockRange(sid request.ServerAndID, req request.Request) { + if _, ok := s.lockedIDs[sid]; !ok { + return + } + delete(s.lockedIDs, sid) + r := req.(ReqUpdates) + s.rangeLock.lock(r.FirstPeriod, r.Count, -1) +} + +func (s *ForwardUpdateSync) verifyRange(req request.Request, resp request.Response) bool { + request, ok := req.(ReqUpdates) if !ok { return false } - response, ok := event.Response.(RespUpdates) + response, ok := resp.(RespUpdates) if !ok { return false } @@ -144,8 +166,9 @@ func (s *ForwardUpdateSync) verifyRange(event request.RequestEvent) bool { } // returns true for partial success -func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event request.RequestEvent) (success bool) { - response, ok := event.Response.(RespUpdates) +func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event request.Event) (success bool) { + sid, _, resp := event.RequestInfo() + response, ok := resp.(RespUpdates) if !ok { return false } @@ -157,7 +180,7 @@ func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event reque return } if err == light.ErrInvalidUpdate || err == light.ErrWrongCommitteeRoot || err == light.ErrCannotReorg { - tracker.InvalidResponse(event.ServerAndID, "invalid update received") + tracker.InvalidResponse(sid, "invalid update received") } else { log.Error("Unexpected InsertUpdate error", "error", err) } @@ -168,34 +191,38 @@ func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event reque return } -type updateResponseList []request.RequestEvent +type updateResponseList []request.Event func (u updateResponseList) Len() int { return len(u) } func (u updateResponseList) Swap(i, j int) { u[i], u[j] = u[j], u[i] } func (u updateResponseList) Less(i, j int) bool { - return u[i].Request.(ReqUpdates).FirstPeriod < u[j].Request.(ReqUpdates).FirstPeriod + return u[i].Data.(request.RequestResponse).Request.(ReqUpdates).FirstPeriod < + u[j].Data.(request.RequestResponse).Request.(ReqUpdates).FirstPeriod } // Process implements request.Module -func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { +func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Event) (trigger bool) { // iterate events and add responses to process queue - for _, event := range requestEvents { - if event.Response != nil && !s.verifyRange(event) { - tracker.InvalidResponse(event.ServerAndID, "invalid update range") - event.Response = nil - } - req := event.Request.(ReqUpdates) - if event.Response != nil { - // there is a response with a valid format; put it in the process queue - s.processQueue = append(s.processQueue, event) - if event.Timeout { - // it was already timed out and unlocked; lock again until processed - s.rangeLock.lock(req.FirstPeriod, req.Count, 1) + 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 } - } else if event.Timeout != event.Finalized { - // unlock if timed out or returned with an invalid response without - // previously being unlocked by a timeout - s.rangeLock.lock(req.FirstPeriod, req.Count, -1) + 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) } } @@ -207,25 +234,14 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []req break } trigger = true - req := event.Request.(ReqUpdates) - s.rangeLock.lock(req.FirstPeriod, req.Count, -1) + sid, req, _ := event.RequestInfo() + s.unlockRange(sid, req) s.processQueue = s.processQueue[1:] if len(s.processQueue) == 0 { s.processQueue = nil } } - // update nextSyncPeriod of servers based on server events - for _, event := range serverEvents { - switch event.Type { - case EvNewSignedHead: - signedHead := event.Data.(types.SignedHeader) - s.nextSyncPeriod[event.Server] = types.SyncPeriod(signedHead.SignatureSlot + 256) - case request.EvUnregistered: - delete(s.nextSyncPeriod, event.Server) - } - } - // start new requests if necessary startPeriod, chainInit := s.chain.NextSyncPeriod() if !chainInit { @@ -233,7 +249,7 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []req } for { firstPeriod, maxCount := s.rangeLock.firstUnlocked(startPeriod, maxUpdateRequest) - if request, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { + if reqWithID, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { nextPeriod := s.nextSyncPeriod[server] if nextPeriod <= firstPeriod { return nil, 0 @@ -244,8 +260,7 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []req } return ReqUpdates{FirstPeriod: firstPeriod, Count: count}, float32(count) }); ok { - req := request.Request.(ReqUpdates) - s.rangeLock.lock(req.FirstPeriod, req.Count, 1) + s.lockRange(reqWithID.ServerAndID, reqWithID.Request) } else { break } diff --git a/beacon/light/sync/update_sync_test.go b/beacon/light/sync/update_sync_test.go index 9b5b4be018..413c763077 100644 --- a/beacon/light/sync/update_sync_test.go +++ b/beacon/light/sync/update_sync_test.go @@ -32,37 +32,41 @@ func TestCheckpointInit(t *testing.T) { checkpoint := &types.BootstrapData{Header: types.Header{Slot: 0x2000*4 + 0x1000}} // period 4 checkpointHash := checkpoint.Header.Hash() chkInit := NewCheckpointInit(chain, checkpointHash) - ExpTrigger(t, 1, false, chkInit.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 1, false, chkInit.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer2, Type: request.EvRegistered}, })) // expect bootstrap request to server 1 req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 1}, Request: ReqCheckpointData(checkpointHash)} tracker.ExpRequests(t, 1, []request.RequestWithID{req1}) - // request times out; expect request to server 2 - ExpTrigger(t, 2, false, chkInit.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req1, Timeout: true}, - }, nil)) + // req1 times out; expect request to server 2 + ExpTrigger(t, 2, false, chkInit.Process(tracker, []request.Event{ + TestReqEvent(request.EvTimeout, req1, nil), + })) req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqCheckpointData(checkpointHash)} tracker.ExpRequests(t, 2, []request.RequestWithID{req2}) // invalid response to req2; expect init state to still be false wrongCheckpoint := &types.BootstrapData{Header: types.Header{Slot: 123456}} - ExpTrigger(t, 3, false, chkInit.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2, Response: wrongCheckpoint, Finalized: true}, - }, nil)) + ExpTrigger(t, 3, false, chkInit.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req2, wrongCheckpoint), + })) + // req1 fails (hard timeout) + ExpTrigger(t, 4, false, chkInit.Process(tracker, []request.Event{ + TestReqEvent(request.EvFail, req1, nil), + })) chain.ExpInit(t, false) // server 3 is registered tracker.AddServer(testServer3, 1) - ExpTrigger(t, 4, false, chkInit.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 5, false, chkInit.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvRegistered}, })) // expect bootstrap request to server 3 req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 3}, Request: ReqCheckpointData(checkpointHash)} tracker.ExpRequests(t, 3, []request.RequestWithID{req3}) // valid response to req3; expect chain to be initialized - ExpTrigger(t, 5, true, chkInit.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req3, Response: checkpoint, Finalized: true}, - }, nil)) + ExpTrigger(t, 6, true, chkInit.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req3, checkpoint), + })) chain.ExpInit(t, true) } @@ -74,7 +78,7 @@ func TestUpdateSyncParallel(t *testing.T) { chain := &TestCommitteeChain{} chain.SetNextSyncPeriod(0) updateSync := NewForwardUpdateSync(chain) - ExpTrigger(t, 1, false, updateSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 1, false, updateSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*100 + 0x1000}}, {Server: testServer2, Type: request.EvRegistered}, @@ -90,38 +94,38 @@ func TestUpdateSyncParallel(t *testing.T) { tracker.ExpRequests(t, 1, []request.RequestWithID{req1, req2, req3, req4, req5, req6}) // valid response to request 1 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 2, true, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req1, Response: testRespUpdate(req1), Finalized: true}, - }, nil)) + ExpTrigger(t, 2, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req1, testRespUpdate(req1)), + })) // expect 8 periods synced and a new request started chain.ExpNextSyncPeriod(t, 8) req7 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 7}, Request: ReqUpdates{FirstPeriod: 48, Count: 8}} tracker.ExpRequests(t, 2, []request.RequestWithID{req7}) // valid response to requests 4 and 5 tracker.AddAllowance(testServer2, 2) - ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req4, Response: testRespUpdate(req4), Finalized: true}, - request.RequestEvent{RequestWithID: req5, Response: testRespUpdate(req5), Finalized: true}, - }, nil)) + ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)), + TestReqEvent(request.EvResponse, req5, testRespUpdate(req5)), + })) // expect 2 more requests but no sync progress (responses 4 and 5 cannot be added before 2 and 3) chain.ExpNextSyncPeriod(t, 8) req8 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 8}, Request: ReqUpdates{FirstPeriod: 56, Count: 8}} req9 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 9}, Request: ReqUpdates{FirstPeriod: 64, Count: 8}} tracker.ExpRequests(t, 3, []request.RequestWithID{req8, req9}) // soft timeout for requests 2 and 3 (server 1 is overloaded) - ExpTrigger(t, 4, false, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2, Timeout: true}, - request.RequestEvent{RequestWithID: req3, Timeout: true}, - }, nil)) + ExpTrigger(t, 4, false, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvTimeout, req2, nil), + TestReqEvent(request.EvTimeout, req3, nil), + })) // no allowance, no more requests tracker.ExpRequests(t, 4, nil) // valid response to requests 6 and 8 and 9 tracker.AddAllowance(testServer2, 3) - ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req6, Response: testRespUpdate(req6), Finalized: true}, - request.RequestEvent{RequestWithID: req8, Response: testRespUpdate(req8), Finalized: true}, - request.RequestEvent{RequestWithID: req9, Response: testRespUpdate(req9), Finalized: true}, - }, nil)) + ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req6, testRespUpdate(req6)), + TestReqEvent(request.EvResponse, req8, testRespUpdate(req8)), + TestReqEvent(request.EvResponse, req9, testRespUpdate(req9)), + })) // server 2 can now resend requests 2 and 3 (timed out by server 1) and also send a new one req2r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 10}, Request: ReqUpdates{FirstPeriod: 8, Count: 8}} req3r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 11}, Request: ReqUpdates{FirstPeriod: 16, Count: 8}} @@ -129,20 +133,19 @@ func TestUpdateSyncParallel(t *testing.T) { tracker.ExpRequests(t, 5, []request.RequestWithID{req2r, req3r, req10}) // server 1 finally answers timed out request 2 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.RequestEvent{ - // note that Timeout flag has to be true once the request timed out, even if answered later - request.RequestEvent{RequestWithID: req2, Response: testRespUpdate(req2), Timeout: true, Finalized: true}, - }, nil)) + ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req2, testRespUpdate(req2)), + })) // expect sync progress and one new request chain.ExpNextSyncPeriod(t, 16) req11 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 13}, Request: ReqUpdates{FirstPeriod: 80, Count: 8}} tracker.ExpRequests(t, 6, []request.RequestWithID{req11}) // server 2 answers re-sent requests 2 and 3 tracker.AddAllowance(testServer2, 2) - ExpTrigger(t, 7, true, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2r, Response: testRespUpdate(req2r), Finalized: true}, - request.RequestEvent{RequestWithID: req3r, Response: testRespUpdate(req3r), Finalized: true}, - }, nil)) + ExpTrigger(t, 7, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req2r, testRespUpdate(req2r)), + TestReqEvent(request.EvResponse, req3r, testRespUpdate(req3r)), + })) // finally the gap is filled, update can process responses up to req6 chain.ExpNextSyncPeriod(t, 48) // expect 2 new requests from server 2 (now the available range is covered) @@ -150,14 +153,14 @@ func TestUpdateSyncParallel(t *testing.T) { req13 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 15}, Request: ReqUpdates{FirstPeriod: 96, Count: 4}} tracker.ExpRequests(t, 7, []request.RequestWithID{req12, req13}) // all remaining requests are answered - ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req3, Response: testRespUpdate(req3), Timeout: true, Finalized: true}, - request.RequestEvent{RequestWithID: req7, Response: testRespUpdate(req7), Finalized: true}, - request.RequestEvent{RequestWithID: req10, Response: testRespUpdate(req10), Finalized: true}, - request.RequestEvent{RequestWithID: req11, Response: testRespUpdate(req11), Finalized: true}, - request.RequestEvent{RequestWithID: req12, Response: testRespUpdate(req12), Finalized: true}, - request.RequestEvent{RequestWithID: req13, Response: testRespUpdate(req13), Finalized: true}, - }, nil)) + ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)), + TestReqEvent(request.EvResponse, req7, testRespUpdate(req7)), + TestReqEvent(request.EvResponse, req10, testRespUpdate(req10)), + TestReqEvent(request.EvResponse, req11, testRespUpdate(req11)), + TestReqEvent(request.EvResponse, req12, testRespUpdate(req12)), + TestReqEvent(request.EvResponse, req13, testRespUpdate(req13)), + })) // expect chain to be fully synced chain.ExpNextSyncPeriod(t, 100) } @@ -171,7 +174,7 @@ func TestUpdateSyncDifferentHeads(t *testing.T) { chain := &TestCommitteeChain{} chain.SetNextSyncPeriod(10) updateSync := NewForwardUpdateSync(chain) - ExpTrigger(t, 1, false, updateSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 1, false, updateSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*15 + 0x1000}}, {Server: testServer2, Type: request.EvRegistered}, @@ -183,39 +186,39 @@ func TestUpdateSyncDifferentHeads(t *testing.T) { req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 1}, Request: ReqUpdates{FirstPeriod: 10, Count: 7}} tracker.ExpRequests(t, 1, []request.RequestWithID{req1}) // request times out, expect request to the next best head - ExpTrigger(t, 2, false, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req1, Timeout: true}, - }, nil)) + ExpTrigger(t, 2, false, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvTimeout, req1, nil), + })) req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqUpdates{FirstPeriod: 10, Count: 6}} tracker.ExpRequests(t, 2, []request.RequestWithID{req2}) // request times out, expect request to the last available server - ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2, Timeout: true}, - }, nil)) + ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvTimeout, req2, nil), + })) req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 3}, Request: ReqUpdates{FirstPeriod: 10, Count: 5}} tracker.ExpRequests(t, 3, []request.RequestWithID{req3}) // valid response to request 3, expect chain synced to period 15 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 4, true, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req3, Response: testRespUpdate(req3), Finalized: true}, - }, nil)) + ExpTrigger(t, 4, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)), + })) chain.ExpNextSyncPeriod(t, 15) // invalid response to request 1, server can only deliver updates up to period 15 despite announced head req1x := request.RequestWithID{ServerAndID: req1.ServerAndID, Request: ReqUpdates{FirstPeriod: 10, Count: 5}} - ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req1, Response: testRespUpdate(req1x), Timeout: true, Finalized: true}, - }, nil)) + ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req1, testRespUpdate(req1x)), + })) // expect no progress of chain head chain.ExpNextSyncPeriod(t, 15) // valid response to request 2, expect chain synced to period 16 tracker.AddAllowance(testServer2, 1) - ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2, Response: testRespUpdate(req2), Timeout: true, Finalized: true}, - }, nil)) + ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req2, testRespUpdate(req2)), + })) chain.ExpNextSyncPeriod(t, 16) // a new server is registered with announced head period 17 tracker.AddServer(testServer4, 1) - ExpTrigger(t, 7, false, updateSync.Process(tracker, nil, []request.ServerEvent{ + ExpTrigger(t, 7, false, updateSync.Process(tracker, []request.Event{ {Server: testServer4, Type: request.EvRegistered}, {Server: testServer4, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*17 + 0x1000}}, })) @@ -224,9 +227,9 @@ func TestUpdateSyncDifferentHeads(t *testing.T) { tracker.ExpRequests(t, 4, []request.RequestWithID{req4}) // valid response, expect chain synced to period 17 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req4, Response: testRespUpdate(req4), Finalized: true}, - }, nil)) + ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.Event{ + TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)), + })) chain.ExpNextSyncPeriod(t, 17) } diff --git a/cmd/blsync/block_sync.go b/cmd/blsync/block_sync.go index 8ee03afd1d..c28094f285 100755 --- a/cmd/blsync/block_sync.go +++ b/cmd/blsync/block_sync.go @@ -40,7 +40,7 @@ import ( type beaconBlockSync struct { recentBlocks *lru.Cache[common.Hash, *capella.BeaconBlock] validatedHead common.Hash - pending map[common.Hash]struct{} + locked map[common.Hash]struct{} serverHeads map[request.Server]common.Hash headTracker headTracker } @@ -54,36 +54,31 @@ func newBeaconBlockSyncer(headTracker headTracker) *beaconBlockSync { return &beaconBlockSync{ headTracker: headTracker, recentBlocks: lru.NewCache[common.Hash, *capella.BeaconBlock](10), - pending: make(map[common.Hash]struct{}), + locked: make(map[common.Hash]struct{}), serverHeads: make(map[request.Server]common.Hash), } } // Process implements request.Module -func (s *beaconBlockSync) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) (trigger bool) { +func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Event) (trigger bool) { if header := s.headTracker.ValidatedHead().Header; header != (types.Header{}) { s.validatedHead = header.Hash() } // iterate events and add valid responses to recentBlocks - for _, event := range requestEvents { - blockRoot := common.Hash(event.Request.(sync.ReqBeaconBlock)) - if event.Response != nil { - block := event.Response.(*capella.BeaconBlock) - s.recentBlocks.Add(blockRoot, block) - if blockRoot == s.validatedHead { - trigger = true - } - } - if event.Timeout || event.Finalized { - // unlock if timed out or returned with an invalid response - delete(s.pending, blockRoot) - } - } - - // update server heads - for _, event := range serverEvents { + 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) + if blockRoot == s.validatedHead { + trigger = true + } + } + delete(s.locked, blockRoot) case sync.EvNewHead: s.serverHeads[event.Server] = event.Data.(types.HeadInfo).BlockRoot case request.EvUnregistered: @@ -111,7 +106,7 @@ func (s *beaconBlockSync) tryRequestBlock(tracker request.Tracker, blockRoot com if _, ok := s.recentBlocks.Get(blockRoot); ok { return } - if _, ok := s.pending[blockRoot]; ok { + if _, ok := s.locked[blockRoot]; ok { return } if _, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { @@ -122,7 +117,7 @@ func (s *beaconBlockSync) tryRequestBlock(tracker request.Tracker, blockRoot com } return sync.ReqBeaconBlock(blockRoot), 0 }); ok { - s.pending[blockRoot] = struct{}{} + s.locked[blockRoot] = struct{}{} } } @@ -185,7 +180,7 @@ type engineApiUpdater struct { } // Process implements request.Module -func (s *engineApiUpdater) Process(tracker request.Tracker, requestEvents []request.RequestEvent, serverEvents []request.ServerEvent) bool { +func (s *engineApiUpdater) Process(tracker request.Tracker, events []request.Event) bool { if atomic.LoadUint32(&s.updating) == 1 { return false } diff --git a/cmd/blsync/block_sync_test.go b/cmd/blsync/block_sync_test.go index ab82fd4076..1d5fdb0dea 100644 --- a/cmd/blsync/block_sync_test.go +++ b/cmd/blsync/block_sync_test.go @@ -51,7 +51,7 @@ func TestBlockSync(t *testing.T) { } } - sync.ExpTrigger(t, 1, false, blockSync.Process(tracker, nil, []request.ServerEvent{ + sync.ExpTrigger(t, 1, false, blockSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer2, Type: request.EvRegistered}, })) @@ -61,7 +61,7 @@ func TestBlockSync(t *testing.T) { // set block 1 as prefetch head, announced by server 2 head1 := blockHeadInfo(testBlock1) ht.prefetch = head1 - sync.ExpTrigger(t, 2, false, blockSync.Process(tracker, nil, []request.ServerEvent{ + sync.ExpTrigger(t, 2, false, blockSync.Process(tracker, []request.Event{ {Server: testServer2, Type: sync.EvNewHead, Data: head1}, })) // expect request to server 2 which has announced the head @@ -69,43 +69,43 @@ func TestBlockSync(t *testing.T) { tracker.ExpRequests(t, 2, []request.RequestWithID{req1}) // valid response tracker.AddAllowance(testServer2, 1) - sync.ExpTrigger(t, 3, false, blockSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req1, Response: testBlock1, Finalized: true}, - }, nil)) + sync.ExpTrigger(t, 3, false, blockSync.Process(tracker, []request.Event{ + sync.TestReqEvent(request.EvResponse, req1, testBlock1), + })) // head block still not expected as the fetched block is not the validated head yet expHeadBlock(2, nil) // set as validated head, expect no further requests but block 1 set as head block ht.validated.Header = blockHeader(testBlock1) - sync.ExpTrigger(t, 4, false, blockSync.Process(tracker, nil, nil)) + sync.ExpTrigger(t, 4, false, blockSync.Process(tracker, nil)) tracker.ExpRequests(t, 3, nil) expHeadBlock(3, testBlock1) // set block 2 as prefetch head, announced by server 1 head2 := blockHeadInfo(testBlock2) ht.prefetch = head2 - sync.ExpTrigger(t, 5, false, blockSync.Process(tracker, nil, []request.ServerEvent{ + sync.ExpTrigger(t, 5, false, blockSync.Process(tracker, []request.Event{ {Server: testServer1, Type: sync.EvNewHead, Data: head2}, })) // expect request to server 1 req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 2}, Request: sync.ReqBeaconBlock(head2.BlockRoot)} tracker.ExpRequests(t, 4, []request.RequestWithID{req2}) - // req2 times out but no further requests expected because server 2 has not announced it - sync.ExpTrigger(t, 6, false, blockSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2, Timeout: true}, - }, nil)) + // req2 fails, no further requests expected because server 2 has not announced it + sync.ExpTrigger(t, 6, false, blockSync.Process(tracker, []request.Event{ + sync.TestReqEvent(request.EvFail, req2, nil), + })) tracker.ExpRequests(t, 5, nil) // set as validated head before retrieving block; now it's assumed to be available from server 2 too ht.validated.Header = blockHeader(testBlock2) - sync.ExpTrigger(t, 7, false, blockSync.Process(tracker, nil, nil)) + sync.ExpTrigger(t, 7, false, blockSync.Process(tracker, nil)) // now head block is unavailable again expHeadBlock(4, nil) // expect req2 retry to server 2 req2r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 3}, Request: sync.ReqBeaconBlock(head2.BlockRoot)} tracker.ExpRequests(t, 6, []request.RequestWithID{req2r}) // valid response, now head block should be block 2 immediately as it is already validated - sync.ExpTrigger(t, 8, true, blockSync.Process(tracker, []request.RequestEvent{ - request.RequestEvent{RequestWithID: req2r, Response: testBlock2, Finalized: true}, - }, nil)) + sync.ExpTrigger(t, 8, true, blockSync.Process(tracker, []request.Event{ + sync.TestReqEvent(request.EvResponse, req2r, testBlock2), + })) expHeadBlock(5, testBlock2) }