beacon/light: refactored events

This commit is contained in:
Zsolt Felfoldi 2024-01-10 02:28:33 +01:00 committed by Felix Lange
parent ed5005adb5
commit 26813d2a6c
12 changed files with 299 additions and 300 deletions

View file

@ -78,9 +78,9 @@ func (s *ApiServer) SendRequest(req request.Request) request.ID {
default: default:
} }
if resp != nil { 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 { } 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 return id

View file

@ -48,30 +48,7 @@ type Module interface {
// Note: Process functions of different modules are never called concurrently; // Note: Process functions of different modules are never called concurrently;
// they are called by Scheduler in the same order of priority as they were // they are called by Scheduler in the same order of priority as they were
// registered in. // registered in.
Process(Tracker, []RequestEvent, []ServerEvent) bool Process(Tracker, []Event) 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)
} }
// Scheduler is a modular network data retrieval framework that coordinates multiple // Scheduler is a modular network data retrieval framework that coordinates multiple
@ -86,7 +63,7 @@ type Scheduler struct {
trackers map[Module]*tracker trackers map[Module]*tracker
servers map[server]struct{} servers map[server]struct{}
pending map[ServerAndID]pendingRequest pending map[ServerAndID]pendingRequest
serverEvents []ServerEvent serverEvents []Event
stopCh chan chan struct{} stopCh chan chan struct{}
triggerCh chan struct{} // restarts waiting sync loop 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 // 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 // pendingRequest keeps track of sent and not finalized requests and their sender
// modules and whether a soft timeout has already happened. // modules and whether a soft timeout has already happened.
type pendingRequest struct { type pendingRequest struct {
request Request request Request
module Module module Module
timeout bool
} }
// NewScheduler creates a new Scheduler. // NewScheduler creates a new Scheduler.
@ -165,11 +117,12 @@ func (s *Scheduler) RegisterServer(rs requestServer) {
defer s.lock.Unlock() defer s.lock.Unlock()
server := newServer(rs, s.clock) server := newServer(rs, s.clock)
s.handleEvent(server, Event{Type: EvRegistered}) s.handleEvent(Event{Type: EvRegistered, Server: server})
server.subscribe(func(event Event) { server.subscribe(func(event Event) {
s.lock.Lock() s.lock.Lock()
if _, ok := s.servers[server]; ok { if _, ok := s.servers[server]; ok {
s.handleEvent(server, event) event.Server = server
s.handleEvent(event)
} else { } else {
log.Error("Event received from unsubscribed server") 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 { if sl, ok := server.(*serverWithLimits); ok && sl.parent == rs {
server.unsubscribe() server.unsubscribe()
delete(s.servers, server) delete(s.servers, server)
s.handleEvent(server, Event{Type: EvUnregistered}) s.handleEvent(Event{Type: EvUnregistered, Server: server})
return return
} }
} }
@ -248,7 +201,7 @@ func (s *Scheduler) processModules() {
eventTypes := make([]string, len(serverEvents)) eventTypes := make([]string, len(serverEvents))
for i, ev := range 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) log.Debug("Processing modules", "servers", len(servers), "server events", eventTypes)
@ -262,17 +215,18 @@ func (s *Scheduler) processModules() {
var respCount, failCount, timeoutCount int var respCount, failCount, timeoutCount int
for _, ev := range requestEvents { for _, ev := range requestEvents {
if ev.Response != nil { switch ev.Type {
case EvResponse:
respCount++ respCount++
} else if ev.Finalized { case EvFail:
failCount++ failCount++
} else { case EvTimeout:
timeoutCount++ timeoutCount++
} }
} }
log.Debug("Processing module", "name", s.names[module], "responses", respCount, "fails", failCount, "timeouts", 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() s.Trigger()
} }
} }
@ -289,24 +243,12 @@ func (s *Scheduler) Trigger() {
// addRequestEvent adds a request event to the sender module's Tracker, ensuring // addRequestEvent adds a request event to the sender module's Tracker, ensuring
// that the module receives it in the next processing round. // that the module receives it in the next processing round.
func (s *Scheduler) addRequestEvent(server Server, id ID, response Response, timeout, finalized bool) { func (s *Scheduler) addRequestEvent(event Event) {
sid := ServerAndID{Server: server, ID: id} sid, _, _ := event.RequestInfo()
if pr, ok := s.pending[sid]; ok { if pr, ok := s.pending[sid]; ok {
tracker := s.trackers[pr.module] tracker := s.trackers[pr.module]
timeout = timeout || pr.timeout tracker.requestEvents = append(tracker.requestEvents, event)
tracker.requestEvents = append(tracker.requestEvents, RequestEvent{ if event.Type != EvTimeout {
RequestWithID: RequestWithID{
ServerAndID: sid,
Request: pr.request,
},
Response: response,
Timeout: timeout,
Finalized: finalized,
})
if timeout && !finalized {
pr.timeout = true
s.pending[sid] = pr
} else {
delete(s.pending, sid) 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 // addServerEvent adds a server event to the global server event list, ensuring
// that all modules receive it in the next processing round. // that all modules receive it in the next processing round.
func (s *Scheduler) addServerEvent(server Server, event Event) { func (s *Scheduler) addServerEvent(event Event) {
s.serverEvents = append(s.serverEvents, ServerEvent{Server: server, Type: event.Type, Data: event.Data}) s.serverEvents = append(s.serverEvents, event)
} }
// handleEvent processes an Event and adds it either as a request event or a // 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 // it also closes all pending requests to the given server by emitting a failed
// request event (Finalized without Response), ensuring that all requests get // request event (Finalized without Response), ensuring that all requests get
// finalized and thereby allowing the module logic to be safe and simple. // 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() s.Trigger()
switch event.Type { if event.IsRequestEvent() {
case EvResponse: s.addRequestEvent(event)
idr := event.Data.(IdAndResponse) return
s.addRequestEvent(server, idr.ID, idr.Response, false, true) }
case EvFail: if event.Type == EvUnregistered {
s.addRequestEvent(server, event.Data.(ID), nil, false, true) for id, pending := range s.pending {
case EvTimeout: if id.Server != event.Server {
s.addRequestEvent(server, event.Data.(ID), nil, true, false)
case EvUnregistered:
for id, _ := range s.pending {
if id.Server != server {
continue continue
} }
s.addRequestEvent(server, id.ID, nil, false, true) s.addRequestEvent(Event{
} Type: EvFail,
s.addServerEvent(server, event) Server: event.Server,
default: Data: RequestResponse{
s.addServerEvent(server, event) ID: id.ID,
Request: pending.request,
},
})
} }
} }
s.addServerEvent(event)
}

View file

@ -28,13 +28,13 @@ import (
var ( var (
// request events // request events
EvResponse = "response" // data: IdAndResponse; sent by requestServer EvResponse = &EventType{Name: "response", requestEvent: true} // data: RequestResponse
EvFail = "fail" // data: ID; sent by requestServer EvFail = &EventType{Name: "fail", requestEvent: true} // data: RequestResponse
EvTimeout = "timeout" // data: ID; sent by serverWithTimeout EvTimeout = &EventType{Name: "timeout", requestEvent: true} // data: RequestResponse
// server events // server events
EvRegistered = "registered" // data: nil; sent by Scheduler EvRegistered = &EventType{Name: "registered"} // data: nil; sent by Scheduler
EvUnregistered = "unregistered" // data: nil; sent by Scheduler EvUnregistered = &EventType{Name: "unregistered"} // data: nil; sent by Scheduler
EvCanRequestAgain = "canRequestAgain" // data: nil; sent by server EvCanRequestAgain = &EventType{Name: "canRequestAgain"} // data: nil; sent by serverWithLimits
) )
const ( const (
@ -81,13 +81,29 @@ func newServer(rs requestServer, clock mclock.Clock) server {
type serverSet map[server]struct{} type serverSet map[server]struct{}
type EventType struct {
Name string
requestEvent bool
}
type Event struct { type Event struct {
Type string Type *EventType
Server Server // filled by Scheduler
Data any Data any
} }
type IdAndResponse struct { 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 ID ID
Request Request
Response Response Response Response
} }
@ -123,12 +139,7 @@ func (s *serverWithTimeout) eventCallback(event Event) {
switch event.Type { switch event.Type {
case EvResponse, EvFail: case EvResponse, EvFail:
var id ID id := event.Data.(RequestResponse).ID
if event.Type == EvResponse {
id = event.Data.(IdAndResponse).ID
} else {
id = event.Data.(ID)
}
if timer, ok := s.timeouts[id]; ok { if timer, ok := s.timeouts[id]; ok {
// Note: if stopping the timer is unsuccessful then the resulting AfterFunc // Note: if stopping the timer is unsuccessful then the resulting AfterFunc
// call will just do nothing // call will just do nothing
@ -167,11 +178,11 @@ func (s *serverWithTimeout) sendRequest(request Request) (reqId ID) {
delete(s.timeouts, reqId) delete(s.timeouts, reqId)
childEventCb := s.childEventCb childEventCb := s.childEventCb
s.lock.Unlock() s.lock.Unlock()
childEventCb(Event{Type: EvFail, Data: reqId}) childEventCb(Event{Type: EvFail, Data: RequestResponse{ID: reqId, Request: request}})
}) })
childEventCb := s.childEventCb childEventCb := s.childEventCb
s.lock.Unlock() s.lock.Unlock()
childEventCb(Event{Type: EvTimeout, Data: reqId}) childEventCb(Event{Type: EvTimeout, Data: RequestResponse{ID: reqId, Request: request}})
}) })
return reqId return reqId
} }
@ -238,7 +249,8 @@ func (s *serverWithLimits) eventCallback(event Event) {
var sendCanRequestAgain bool var sendCanRequestAgain bool
switch event.Type { switch event.Type {
case EvTimeout: case EvTimeout:
s.softTimeouts[event.Data.(ID)] = struct{}{} id := event.Data.(RequestResponse).ID
s.softTimeouts[id] = struct{}{}
s.timeoutCount++ s.timeoutCount++
s.parallelLimit -= parallelAdjustDown s.parallelLimit -= parallelAdjustDown
if s.parallelLimit < minParallelLimit { if s.parallelLimit < minParallelLimit {
@ -246,12 +258,7 @@ func (s *serverWithLimits) eventCallback(event Event) {
} }
log.Debug("Server timeout", "count", s.timeoutCount, "parallelLimit", s.parallelLimit) log.Debug("Server timeout", "count", s.timeoutCount, "parallelLimit", s.parallelLimit)
case EvResponse, EvFail: case EvResponse, EvFail:
var id ID id := event.Data.(RequestResponse).ID
if event.Type == EvResponse {
id = event.Data.(IdAndResponse).ID
} else {
id = event.Data.(ID)
}
if _, ok := s.softTimeouts[id]; ok { if _, ok := s.softTimeouts[id]; ok {
delete(s.softTimeouts, id) delete(s.softTimeouts, id)
s.timeoutCount-- s.timeoutCount--

View file

@ -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 // one per sync process
type tracker struct { type tracker struct {
servers serverSet // one per trigger servers serverSet // one per trigger
scheduler *Scheduler scheduler *Scheduler
module Module module Module
requestEvents []RequestEvent requestEvents []Event
} }
func (p *tracker) TryRequest(requestFn func(server Server) (Request, float32)) (RequestWithID, bool) { func (p *tracker) TryRequest(requestFn func(server Server) (Request, float32)) (RequestWithID, bool) {

View file

@ -58,13 +58,13 @@ func NewHeadSync(headTracker headTracker, chain committeeChain) *HeadSync {
} }
// Process implements request.Module // 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() nextPeriod, chainInit := s.chain.NextSyncPeriod()
if nextPeriod != s.nextSyncPeriod || chainInit != s.chainInit { if nextPeriod != s.nextSyncPeriod || chainInit != s.chainInit {
s.nextSyncPeriod, s.chainInit = nextPeriod, chainInit s.nextSyncPeriod, s.chainInit = nextPeriod, chainInit
trigger = s.processUnvalidatedHeadsHeads() trigger = s.processUnvalidatedHeads()
} }
for _, event := range serverEvents { for _, event := range events {
switch event.Type { switch event.Type {
case EvNewHead: case EvNewHead:
if s.setServerHead(event.Server, event.Data.(types.HeadInfo)) { 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 return updated
} }
func (s *HeadSync) processUnvalidatedHeadsHeads() (trigger bool) { func (s *HeadSync) processUnvalidatedHeads() (trigger bool) {
if !s.chainInit { if !s.chainInit {
return false return false
} }

View file

@ -50,38 +50,38 @@ func TestValidatedHead(t *testing.T) {
headSync := NewHeadSync(ht, chain) headSync := NewHeadSync(ht, chain)
ht.ExpValidated(t, 1, nil) 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: request.EvRegistered},
{Server: testServer1, Type: EvNewSignedHead, Data: testSHead1}, {Server: testServer1, Type: EvNewSignedHead, Data: testSHead1},
})) }))
ht.ExpValidated(t, 2, nil) ht.ExpValidated(t, 2, nil)
chain.SetNextSyncPeriod(0) 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}) ht.ExpValidated(t, 3, []types.SignedHeader{testSHead1})
chain.SetNextSyncPeriod(1) 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: testServer1, Type: EvNewSignedHead, Data: testSHead2},
{Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: request.EvRegistered},
{Server: testServer2, Type: EvNewSignedHead, Data: testSHead2}, {Server: testServer2, Type: EvNewSignedHead, Data: testSHead2},
})) }))
ht.ExpValidated(t, 4, []types.SignedHeader{testSHead2, 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: testServer1, Type: EvNewSignedHead, Data: testSHead3},
{Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: request.EvRegistered},
{Server: testServer3, Type: EvNewSignedHead, Data: testSHead4}, {Server: testServer3, Type: EvNewSignedHead, Data: testSHead4},
})) }))
ht.ExpValidated(t, 5, nil) ht.ExpValidated(t, 5, nil)
chain.SetNextSyncPeriod(2) 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}) 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}, {Server: testServer3, Type: request.EvUnregistered},
})) }))
ht.ExpValidated(t, 7, nil) ht.ExpValidated(t, 7, nil)
chain.SetNextSyncPeriod(3) 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) 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}, {Server: testServer2, Type: EvNewSignedHead, Data: testSHead4},
})) }))
ht.ExpValidated(t, 9, []types.SignedHeader{testSHead4}) ht.ExpValidated(t, 9, []types.SignedHeader{testSHead4})
@ -94,47 +94,47 @@ func TestPrefetchHead(t *testing.T) {
headSync := NewHeadSync(ht, chain) headSync := NewHeadSync(ht, chain)
ht.ExpPrefetch(t, 1, testHead0) // no servers registered 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: request.EvRegistered},
{Server: testServer1, Type: EvNewHead, Data: testHead1}, {Server: testServer1, Type: EvNewHead, Data: testHead1},
})) }))
ht.ExpPrefetch(t, 2, testHead1) // s1: h1 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: request.EvRegistered},
{Server: testServer2, Type: EvNewHead, Data: testHead2}, {Server: testServer2, Type: EvNewHead, Data: testHead2},
})) }))
ht.ExpPrefetch(t, 3, testHead2) // s1: h1, s2: h2 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}, {Server: testServer1, Type: EvNewHead, Data: testHead2},
})) }))
ht.ExpPrefetch(t, 4, testHead2) // s1: h2, s2: h2 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: request.EvRegistered},
{Server: testServer3, Type: EvNewHead, Data: testHead3}, {Server: testServer3, Type: EvNewHead, Data: testHead3},
})) }))
ht.ExpPrefetch(t, 5, testHead2) // s1: h2, s2: h2, s3: h3 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: request.EvRegistered},
{Server: testServer4, Type: EvNewHead, Data: testHead4}, {Server: testServer4, Type: EvNewHead, Data: testHead4},
})) }))
ht.ExpPrefetch(t, 6, testHead2) // s1: h2, s2: h2, s3: h3, s4: h4 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}, {Server: testServer2, Type: EvNewHead, Data: testHead3},
})) }))
ht.ExpPrefetch(t, 7, testHead3) // s1: h2, s2: h3, s3: h3, s4: h4 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}, {Server: testServer3, Type: request.EvUnregistered},
})) }))
ht.ExpPrefetch(t, 8, testHead4) // s1: h2, s2: h3, s4: h4 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}, {Server: testServer1, Type: request.EvUnregistered},
})) }))
ht.ExpPrefetch(t, 9, testHead4) // s2: h3, s4: h4 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}, {Server: testServer4, Type: request.EvUnregistered},
})) }))
ht.ExpPrefetch(t, 10, testHead3) // s2: h3 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}, {Server: testServer2, Type: request.EvUnregistered},
})) }))
ht.ExpPrefetch(t, 11, testHead0) // no servers registered ht.ExpPrefetch(t, 11, testHead0) // no servers registered

View file

@ -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 { type TestCommitteeChain struct {
fsp, nsp uint64 fsp, nsp uint64
init bool init bool

View file

@ -17,13 +17,14 @@
package sync package sync
import ( import (
"github.com/ethereum/go-ethereum/beacon/light/request"
"github.com/ethereum/go-ethereum/beacon/types" "github.com/ethereum/go-ethereum/beacon/types"
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
) )
const ( var (
EvNewHead = "newHead" EvNewHead = &request.EventType{Name: "newHead"} // data: types.HeadInfo
EvNewSignedHead = "newSignedHead" EvNewSignedHead = &request.EventType{Name: "newSignedHead"} // data: types.SignedHeader
) )
type ( type (

View file

@ -37,7 +37,7 @@ type committeeChain interface {
type CheckpointInit struct { type CheckpointInit struct {
chain committeeChain chain committeeChain
checkpointHash common.Hash checkpointHash common.Hash
pending bool locked bool
initialized bool initialized bool
} }
@ -49,28 +49,30 @@ func NewCheckpointInit(chain committeeChain, checkpointHash common.Hash) *Checkp
} }
// Process implements request.Module // 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 { if s.initialized {
return false return false
} }
for _, event := range requestEvents { for _, event := range events {
if event.Timeout != event.Finalized { if !event.IsRequestEvent() {
s.pending = false continue
} }
if event.Response != nil { s.locked = false
if checkpoint, ok := event.Response.(*types.BootstrapData); ok && checkpoint.Header.Hash() == common.Hash(event.Request.(ReqCheckpointData)) { 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.chain.CheckpointInit(*checkpoint) //TODO
s.initialized = true s.initialized = true
return 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) { if _, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) {
return ReqCheckpointData(s.checkpointHash), 0 return ReqCheckpointData(s.checkpointHash), 0
}); ok { }); ok {
s.pending = true s.locked = true
} }
} }
return false return false
@ -79,7 +81,8 @@ func (s *CheckpointInit) Process(tracker request.Tracker, requestEvents []reques
type ForwardUpdateSync struct { type ForwardUpdateSync struct {
chain committeeChain chain committeeChain
rangeLock rangeLock rangeLock rangeLock
processQueue []request.RequestEvent lockedIDs map[request.ServerAndID]struct{}
processQueue []request.Event
nextSyncPeriod map[request.Server]uint64 nextSyncPeriod map[request.Server]uint64
} }
@ -87,6 +90,7 @@ func NewForwardUpdateSync(chain committeeChain) *ForwardUpdateSync {
return &ForwardUpdateSync{ return &ForwardUpdateSync{
chain: chain, chain: chain,
rangeLock: make(rangeLock), rangeLock: make(rangeLock),
lockedIDs: make(map[request.ServerAndID]struct{}),
nextSyncPeriod: make(map[request.Server]uint64), nextSyncPeriod: make(map[request.Server]uint64),
} }
} }
@ -123,12 +127,30 @@ func (r rangeLock) firstUnlocked(start, maxCount uint64) (first, count uint64) {
return return
} }
func (s *ForwardUpdateSync) verifyRange(event request.RequestEvent) bool { func (s *ForwardUpdateSync) lockRange(sid request.ServerAndID, req request.Request) {
request, ok := event.Request.(ReqUpdates) 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 { if !ok {
return false return false
} }
response, ok := event.Response.(RespUpdates) response, ok := resp.(RespUpdates)
if !ok { if !ok {
return false return false
} }
@ -144,8 +166,9 @@ func (s *ForwardUpdateSync) verifyRange(event request.RequestEvent) bool {
} }
// returns true for partial success // returns true for partial success
func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event request.RequestEvent) (success bool) { func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event request.Event) (success bool) {
response, ok := event.Response.(RespUpdates) sid, _, resp := event.RequestInfo()
response, ok := resp.(RespUpdates)
if !ok { if !ok {
return false return false
} }
@ -157,7 +180,7 @@ func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event reque
return return
} }
if err == light.ErrInvalidUpdate || err == light.ErrWrongCommitteeRoot || err == light.ErrCannotReorg { if err == light.ErrInvalidUpdate || err == light.ErrWrongCommitteeRoot || err == light.ErrCannotReorg {
tracker.InvalidResponse(event.ServerAndID, "invalid update received") tracker.InvalidResponse(sid, "invalid update received")
} else { } else {
log.Error("Unexpected InsertUpdate error", "error", err) log.Error("Unexpected InsertUpdate error", "error", err)
} }
@ -168,34 +191,38 @@ func (s *ForwardUpdateSync) processResponse(tracker request.Tracker, event reque
return return
} }
type updateResponseList []request.RequestEvent type updateResponseList []request.Event
func (u updateResponseList) Len() int { return len(u) } 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) Swap(i, j int) { u[i], u[j] = u[j], u[i] }
func (u updateResponseList) Less(i, j int) bool { 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 // 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 // iterate events and add responses to process queue
for _, event := range requestEvents { for _, event := range events {
if event.Response != nil && !s.verifyRange(event) { switch event.Type {
tracker.InvalidResponse(event.ServerAndID, "invalid update range") case request.EvResponse, request.EvFail, request.EvTimeout:
event.Response = nil sid, req, resp := event.RequestInfo()
if event.Type == request.EvResponse && !s.verifyRange(req, resp) {
tracker.InvalidResponse(sid, "invalid update range")
resp = nil
} }
req := event.Request.(ReqUpdates) if resp != nil {
if event.Response != nil {
// there is a response with a valid format; put it in the process queue // there is a response with a valid format; put it in the process queue
s.processQueue = append(s.processQueue, event) s.processQueue = append(s.processQueue, event)
if event.Timeout { s.lockRange(sid, req)
// it was already timed out and unlocked; lock again until processed } else {
s.rangeLock.lock(req.FirstPeriod, req.Count, 1) s.unlockRange(sid, req)
} }
} else if event.Timeout != event.Finalized { case EvNewSignedHead:
// unlock if timed out or returned with an invalid response without signedHead := event.Data.(types.SignedHeader)
// previously being unlocked by a timeout s.nextSyncPeriod[event.Server] = types.SyncPeriod(signedHead.SignatureSlot + 256)
s.rangeLock.lock(req.FirstPeriod, req.Count, -1) case request.EvUnregistered:
delete(s.nextSyncPeriod, event.Server)
} }
} }
@ -207,25 +234,14 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []req
break break
} }
trigger = true trigger = true
req := event.Request.(ReqUpdates) sid, req, _ := event.RequestInfo()
s.rangeLock.lock(req.FirstPeriod, req.Count, -1) s.unlockRange(sid, req)
s.processQueue = s.processQueue[1:] s.processQueue = s.processQueue[1:]
if len(s.processQueue) == 0 { if len(s.processQueue) == 0 {
s.processQueue = nil 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 // start new requests if necessary
startPeriod, chainInit := s.chain.NextSyncPeriod() startPeriod, chainInit := s.chain.NextSyncPeriod()
if !chainInit { if !chainInit {
@ -233,7 +249,7 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []req
} }
for { for {
firstPeriod, maxCount := s.rangeLock.firstUnlocked(startPeriod, maxUpdateRequest) 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] nextPeriod := s.nextSyncPeriod[server]
if nextPeriod <= firstPeriod { if nextPeriod <= firstPeriod {
return nil, 0 return nil, 0
@ -244,8 +260,7 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, requestEvents []req
} }
return ReqUpdates{FirstPeriod: firstPeriod, Count: count}, float32(count) return ReqUpdates{FirstPeriod: firstPeriod, Count: count}, float32(count)
}); ok { }); ok {
req := request.Request.(ReqUpdates) s.lockRange(reqWithID.ServerAndID, reqWithID.Request)
s.rangeLock.lock(req.FirstPeriod, req.Count, 1)
} else { } else {
break break
} }

View file

@ -32,37 +32,41 @@ func TestCheckpointInit(t *testing.T) {
checkpoint := &types.BootstrapData{Header: types.Header{Slot: 0x2000*4 + 0x1000}} // period 4 checkpoint := &types.BootstrapData{Header: types.Header{Slot: 0x2000*4 + 0x1000}} // period 4
checkpointHash := checkpoint.Header.Hash() checkpointHash := checkpoint.Header.Hash()
chkInit := NewCheckpointInit(chain, checkpointHash) 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: testServer1, Type: request.EvRegistered},
{Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: request.EvRegistered},
})) }))
// expect bootstrap request to server 1 // expect bootstrap request to server 1
req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 1}, Request: ReqCheckpointData(checkpointHash)} req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 1}, Request: ReqCheckpointData(checkpointHash)}
tracker.ExpRequests(t, 1, []request.RequestWithID{req1}) tracker.ExpRequests(t, 1, []request.RequestWithID{req1})
// request times out; expect request to server 2 // req1 times out; expect request to server 2
ExpTrigger(t, 2, false, chkInit.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 2, false, chkInit.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req1, Timeout: true}, TestReqEvent(request.EvTimeout, req1, nil),
}, nil)) }))
req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqCheckpointData(checkpointHash)} req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqCheckpointData(checkpointHash)}
tracker.ExpRequests(t, 2, []request.RequestWithID{req2}) tracker.ExpRequests(t, 2, []request.RequestWithID{req2})
// invalid response to req2; expect init state to still be false // invalid response to req2; expect init state to still be false
wrongCheckpoint := &types.BootstrapData{Header: types.Header{Slot: 123456}} wrongCheckpoint := &types.BootstrapData{Header: types.Header{Slot: 123456}}
ExpTrigger(t, 3, false, chkInit.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 3, false, chkInit.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2, Response: wrongCheckpoint, Finalized: true}, TestReqEvent(request.EvResponse, req2, wrongCheckpoint),
}, nil)) }))
// req1 fails (hard timeout)
ExpTrigger(t, 4, false, chkInit.Process(tracker, []request.Event{
TestReqEvent(request.EvFail, req1, nil),
}))
chain.ExpInit(t, false) chain.ExpInit(t, false)
// server 3 is registered // server 3 is registered
tracker.AddServer(testServer3, 1) 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}, {Server: testServer3, Type: request.EvRegistered},
})) }))
// expect bootstrap request to server 3 // expect bootstrap request to server 3
req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 3}, Request: ReqCheckpointData(checkpointHash)} req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 3}, Request: ReqCheckpointData(checkpointHash)}
tracker.ExpRequests(t, 3, []request.RequestWithID{req3}) tracker.ExpRequests(t, 3, []request.RequestWithID{req3})
// valid response to req3; expect chain to be initialized // valid response to req3; expect chain to be initialized
ExpTrigger(t, 5, true, chkInit.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 6, true, chkInit.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req3, Response: checkpoint, Finalized: true}, TestReqEvent(request.EvResponse, req3, checkpoint),
}, nil)) }))
chain.ExpInit(t, true) chain.ExpInit(t, true)
} }
@ -74,7 +78,7 @@ func TestUpdateSyncParallel(t *testing.T) {
chain := &TestCommitteeChain{} chain := &TestCommitteeChain{}
chain.SetNextSyncPeriod(0) chain.SetNextSyncPeriod(0)
updateSync := NewForwardUpdateSync(chain) 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: request.EvRegistered},
{Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*100 + 0x1000}}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*100 + 0x1000}},
{Server: testServer2, Type: request.EvRegistered}, {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}) tracker.ExpRequests(t, 1, []request.RequestWithID{req1, req2, req3, req4, req5, req6})
// valid response to request 1 // valid response to request 1
tracker.AddAllowance(testServer1, 1) tracker.AddAllowance(testServer1, 1)
ExpTrigger(t, 2, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 2, true, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req1, Response: testRespUpdate(req1), Finalized: true}, TestReqEvent(request.EvResponse, req1, testRespUpdate(req1)),
}, nil)) }))
// expect 8 periods synced and a new request started // expect 8 periods synced and a new request started
chain.ExpNextSyncPeriod(t, 8) chain.ExpNextSyncPeriod(t, 8)
req7 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 7}, Request: ReqUpdates{FirstPeriod: 48, Count: 8}} req7 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 7}, Request: ReqUpdates{FirstPeriod: 48, Count: 8}}
tracker.ExpRequests(t, 2, []request.RequestWithID{req7}) tracker.ExpRequests(t, 2, []request.RequestWithID{req7})
// valid response to requests 4 and 5 // valid response to requests 4 and 5
tracker.AddAllowance(testServer2, 2) tracker.AddAllowance(testServer2, 2)
ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req4, Response: testRespUpdate(req4), Finalized: true}, TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)),
request.RequestEvent{RequestWithID: req5, Response: testRespUpdate(req5), Finalized: true}, TestReqEvent(request.EvResponse, req5, testRespUpdate(req5)),
}, nil)) }))
// expect 2 more requests but no sync progress (responses 4 and 5 cannot be added before 2 and 3) // expect 2 more requests but no sync progress (responses 4 and 5 cannot be added before 2 and 3)
chain.ExpNextSyncPeriod(t, 8) chain.ExpNextSyncPeriod(t, 8)
req8 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 8}, Request: ReqUpdates{FirstPeriod: 56, Count: 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}} req9 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 9}, Request: ReqUpdates{FirstPeriod: 64, Count: 8}}
tracker.ExpRequests(t, 3, []request.RequestWithID{req8, req9}) tracker.ExpRequests(t, 3, []request.RequestWithID{req8, req9})
// soft timeout for requests 2 and 3 (server 1 is overloaded) // soft timeout for requests 2 and 3 (server 1 is overloaded)
ExpTrigger(t, 4, false, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 4, false, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2, Timeout: true}, TestReqEvent(request.EvTimeout, req2, nil),
request.RequestEvent{RequestWithID: req3, Timeout: true}, TestReqEvent(request.EvTimeout, req3, nil),
}, nil)) }))
// no allowance, no more requests // no allowance, no more requests
tracker.ExpRequests(t, 4, nil) tracker.ExpRequests(t, 4, nil)
// valid response to requests 6 and 8 and 9 // valid response to requests 6 and 8 and 9
tracker.AddAllowance(testServer2, 3) tracker.AddAllowance(testServer2, 3)
ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req6, Response: testRespUpdate(req6), Finalized: true}, TestReqEvent(request.EvResponse, req6, testRespUpdate(req6)),
request.RequestEvent{RequestWithID: req8, Response: testRespUpdate(req8), Finalized: true}, TestReqEvent(request.EvResponse, req8, testRespUpdate(req8)),
request.RequestEvent{RequestWithID: req9, Response: testRespUpdate(req9), Finalized: true}, TestReqEvent(request.EvResponse, req9, testRespUpdate(req9)),
}, nil)) }))
// server 2 can now resend requests 2 and 3 (timed out by server 1) and also send a new one // 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}} 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}} 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}) tracker.ExpRequests(t, 5, []request.RequestWithID{req2r, req3r, req10})
// server 1 finally answers timed out request 2 // server 1 finally answers timed out request 2
tracker.AddAllowance(testServer1, 1) tracker.AddAllowance(testServer1, 1)
ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.Event{
// note that Timeout flag has to be true once the request timed out, even if answered later TestReqEvent(request.EvResponse, req2, testRespUpdate(req2)),
request.RequestEvent{RequestWithID: req2, Response: testRespUpdate(req2), Timeout: true, Finalized: true}, }))
}, nil))
// expect sync progress and one new request // expect sync progress and one new request
chain.ExpNextSyncPeriod(t, 16) chain.ExpNextSyncPeriod(t, 16)
req11 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 13}, Request: ReqUpdates{FirstPeriod: 80, Count: 8}} req11 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 13}, Request: ReqUpdates{FirstPeriod: 80, Count: 8}}
tracker.ExpRequests(t, 6, []request.RequestWithID{req11}) tracker.ExpRequests(t, 6, []request.RequestWithID{req11})
// server 2 answers re-sent requests 2 and 3 // server 2 answers re-sent requests 2 and 3
tracker.AddAllowance(testServer2, 2) tracker.AddAllowance(testServer2, 2)
ExpTrigger(t, 7, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 7, true, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2r, Response: testRespUpdate(req2r), Finalized: true}, TestReqEvent(request.EvResponse, req2r, testRespUpdate(req2r)),
request.RequestEvent{RequestWithID: req3r, Response: testRespUpdate(req3r), Finalized: true}, TestReqEvent(request.EvResponse, req3r, testRespUpdate(req3r)),
}, nil)) }))
// finally the gap is filled, update can process responses up to req6 // finally the gap is filled, update can process responses up to req6
chain.ExpNextSyncPeriod(t, 48) chain.ExpNextSyncPeriod(t, 48)
// expect 2 new requests from server 2 (now the available range is covered) // 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}} req13 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 15}, Request: ReqUpdates{FirstPeriod: 96, Count: 4}}
tracker.ExpRequests(t, 7, []request.RequestWithID{req12, req13}) tracker.ExpRequests(t, 7, []request.RequestWithID{req12, req13})
// all remaining requests are answered // all remaining requests are answered
ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req3, Response: testRespUpdate(req3), Timeout: true, Finalized: true}, TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)),
request.RequestEvent{RequestWithID: req7, Response: testRespUpdate(req7), Finalized: true}, TestReqEvent(request.EvResponse, req7, testRespUpdate(req7)),
request.RequestEvent{RequestWithID: req10, Response: testRespUpdate(req10), Finalized: true}, TestReqEvent(request.EvResponse, req10, testRespUpdate(req10)),
request.RequestEvent{RequestWithID: req11, Response: testRespUpdate(req11), Finalized: true}, TestReqEvent(request.EvResponse, req11, testRespUpdate(req11)),
request.RequestEvent{RequestWithID: req12, Response: testRespUpdate(req12), Finalized: true}, TestReqEvent(request.EvResponse, req12, testRespUpdate(req12)),
request.RequestEvent{RequestWithID: req13, Response: testRespUpdate(req13), Finalized: true}, TestReqEvent(request.EvResponse, req13, testRespUpdate(req13)),
}, nil)) }))
// expect chain to be fully synced // expect chain to be fully synced
chain.ExpNextSyncPeriod(t, 100) chain.ExpNextSyncPeriod(t, 100)
} }
@ -171,7 +174,7 @@ func TestUpdateSyncDifferentHeads(t *testing.T) {
chain := &TestCommitteeChain{} chain := &TestCommitteeChain{}
chain.SetNextSyncPeriod(10) chain.SetNextSyncPeriod(10)
updateSync := NewForwardUpdateSync(chain) 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: request.EvRegistered},
{Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*15 + 0x1000}}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*15 + 0x1000}},
{Server: testServer2, Type: request.EvRegistered}, {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}} req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 1}, Request: ReqUpdates{FirstPeriod: 10, Count: 7}}
tracker.ExpRequests(t, 1, []request.RequestWithID{req1}) tracker.ExpRequests(t, 1, []request.RequestWithID{req1})
// request times out, expect request to the next best head // request times out, expect request to the next best head
ExpTrigger(t, 2, false, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 2, false, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req1, Timeout: true}, TestReqEvent(request.EvTimeout, req1, nil),
}, nil)) }))
req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqUpdates{FirstPeriod: 10, Count: 6}} req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqUpdates{FirstPeriod: 10, Count: 6}}
tracker.ExpRequests(t, 2, []request.RequestWithID{req2}) tracker.ExpRequests(t, 2, []request.RequestWithID{req2})
// request times out, expect request to the last available server // request times out, expect request to the last available server
ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2, Timeout: true}, TestReqEvent(request.EvTimeout, req2, nil),
}, nil)) }))
req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 3}, Request: ReqUpdates{FirstPeriod: 10, Count: 5}} req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 3}, Request: ReqUpdates{FirstPeriod: 10, Count: 5}}
tracker.ExpRequests(t, 3, []request.RequestWithID{req3}) tracker.ExpRequests(t, 3, []request.RequestWithID{req3})
// valid response to request 3, expect chain synced to period 15 // valid response to request 3, expect chain synced to period 15
tracker.AddAllowance(testServer1, 1) tracker.AddAllowance(testServer1, 1)
ExpTrigger(t, 4, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 4, true, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req3, Response: testRespUpdate(req3), Finalized: true}, TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)),
}, nil)) }))
chain.ExpNextSyncPeriod(t, 15) chain.ExpNextSyncPeriod(t, 15)
// invalid response to request 1, server can only deliver updates up to period 15 despite announced head // 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}} req1x := request.RequestWithID{ServerAndID: req1.ServerAndID, Request: ReqUpdates{FirstPeriod: 10, Count: 5}}
ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req1, Response: testRespUpdate(req1x), Timeout: true, Finalized: true}, TestReqEvent(request.EvResponse, req1, testRespUpdate(req1x)),
}, nil)) }))
// expect no progress of chain head // expect no progress of chain head
chain.ExpNextSyncPeriod(t, 15) chain.ExpNextSyncPeriod(t, 15)
// valid response to request 2, expect chain synced to period 16 // valid response to request 2, expect chain synced to period 16
tracker.AddAllowance(testServer2, 1) tracker.AddAllowance(testServer2, 1)
ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2, Response: testRespUpdate(req2), Timeout: true, Finalized: true}, TestReqEvent(request.EvResponse, req2, testRespUpdate(req2)),
}, nil)) }))
chain.ExpNextSyncPeriod(t, 16) chain.ExpNextSyncPeriod(t, 16)
// a new server is registered with announced head period 17 // a new server is registered with announced head period 17
tracker.AddServer(testServer4, 1) 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: request.EvRegistered},
{Server: testServer4, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*17 + 0x1000}}, {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}) tracker.ExpRequests(t, 4, []request.RequestWithID{req4})
// valid response, expect chain synced to period 17 // valid response, expect chain synced to period 17
tracker.AddAllowance(testServer1, 1) tracker.AddAllowance(testServer1, 1)
ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.RequestEvent{ ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req4, Response: testRespUpdate(req4), Finalized: true}, TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)),
}, nil)) }))
chain.ExpNextSyncPeriod(t, 17) chain.ExpNextSyncPeriod(t, 17)
} }

View file

@ -40,7 +40,7 @@ import (
type beaconBlockSync struct { type beaconBlockSync struct {
recentBlocks *lru.Cache[common.Hash, *capella.BeaconBlock] recentBlocks *lru.Cache[common.Hash, *capella.BeaconBlock]
validatedHead common.Hash validatedHead common.Hash
pending map[common.Hash]struct{} locked map[common.Hash]struct{}
serverHeads map[request.Server]common.Hash serverHeads map[request.Server]common.Hash
headTracker headTracker headTracker headTracker
} }
@ -54,36 +54,31 @@ func newBeaconBlockSyncer(headTracker headTracker) *beaconBlockSync {
return &beaconBlockSync{ return &beaconBlockSync{
headTracker: headTracker, headTracker: headTracker,
recentBlocks: lru.NewCache[common.Hash, *capella.BeaconBlock](10), 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), serverHeads: make(map[request.Server]common.Hash),
} }
} }
// Process implements request.Module // 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{}) { if header := s.headTracker.ValidatedHead().Header; header != (types.Header{}) {
s.validatedHead = header.Hash() s.validatedHead = header.Hash()
} }
// iterate events and add valid responses to recentBlocks // iterate events and add valid responses to recentBlocks
for _, event := range requestEvents { for _, event := range events {
blockRoot := common.Hash(event.Request.(sync.ReqBeaconBlock)) switch event.Type {
if event.Response != nil { case request.EvResponse, request.EvFail, request.EvTimeout:
block := event.Response.(*capella.BeaconBlock) _, req, resp := event.RequestInfo()
blockRoot := common.Hash(req.(sync.ReqBeaconBlock))
if resp != nil {
block := resp.(*capella.BeaconBlock)
s.recentBlocks.Add(blockRoot, block) s.recentBlocks.Add(blockRoot, block)
if blockRoot == s.validatedHead { if blockRoot == s.validatedHead {
trigger = true trigger = true
} }
} }
if event.Timeout || event.Finalized { delete(s.locked, blockRoot)
// unlock if timed out or returned with an invalid response
delete(s.pending, blockRoot)
}
}
// update server heads
for _, event := range serverEvents {
switch event.Type {
case sync.EvNewHead: case sync.EvNewHead:
s.serverHeads[event.Server] = event.Data.(types.HeadInfo).BlockRoot s.serverHeads[event.Server] = event.Data.(types.HeadInfo).BlockRoot
case request.EvUnregistered: case request.EvUnregistered:
@ -111,7 +106,7 @@ func (s *beaconBlockSync) tryRequestBlock(tracker request.Tracker, blockRoot com
if _, ok := s.recentBlocks.Get(blockRoot); ok { if _, ok := s.recentBlocks.Get(blockRoot); ok {
return return
} }
if _, ok := s.pending[blockRoot]; ok { if _, ok := s.locked[blockRoot]; ok {
return return
} }
if _, ok := tracker.TryRequest(func(server request.Server) (request.Request, float32) { 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 return sync.ReqBeaconBlock(blockRoot), 0
}); ok { }); ok {
s.pending[blockRoot] = struct{}{} s.locked[blockRoot] = struct{}{}
} }
} }
@ -185,7 +180,7 @@ type engineApiUpdater struct {
} }
// Process implements request.Module // 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 { if atomic.LoadUint32(&s.updating) == 1 {
return false return false
} }

View file

@ -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: testServer1, Type: request.EvRegistered},
{Server: testServer2, 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 // set block 1 as prefetch head, announced by server 2
head1 := blockHeadInfo(testBlock1) head1 := blockHeadInfo(testBlock1)
ht.prefetch = head1 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}, {Server: testServer2, Type: sync.EvNewHead, Data: head1},
})) }))
// expect request to server 2 which has announced the head // 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}) tracker.ExpRequests(t, 2, []request.RequestWithID{req1})
// valid response // valid response
tracker.AddAllowance(testServer2, 1) tracker.AddAllowance(testServer2, 1)
sync.ExpTrigger(t, 3, false, blockSync.Process(tracker, []request.RequestEvent{ sync.ExpTrigger(t, 3, false, blockSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req1, Response: testBlock1, Finalized: true}, sync.TestReqEvent(request.EvResponse, req1, testBlock1),
}, nil)) }))
// head block still not expected as the fetched block is not the validated head yet // head block still not expected as the fetched block is not the validated head yet
expHeadBlock(2, nil) expHeadBlock(2, nil)
// set as validated head, expect no further requests but block 1 set as head block // set as validated head, expect no further requests but block 1 set as head block
ht.validated.Header = blockHeader(testBlock1) 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) tracker.ExpRequests(t, 3, nil)
expHeadBlock(3, testBlock1) expHeadBlock(3, testBlock1)
// set block 2 as prefetch head, announced by server 1 // set block 2 as prefetch head, announced by server 1
head2 := blockHeadInfo(testBlock2) head2 := blockHeadInfo(testBlock2)
ht.prefetch = head2 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}, {Server: testServer1, Type: sync.EvNewHead, Data: head2},
})) }))
// expect request to server 1 // expect request to server 1
req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 2}, Request: sync.ReqBeaconBlock(head2.BlockRoot)} req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 2}, Request: sync.ReqBeaconBlock(head2.BlockRoot)}
tracker.ExpRequests(t, 4, []request.RequestWithID{req2}) tracker.ExpRequests(t, 4, []request.RequestWithID{req2})
// req2 times out but no further requests expected because server 2 has not announced it // req2 fails, no further requests expected because server 2 has not announced it
sync.ExpTrigger(t, 6, false, blockSync.Process(tracker, []request.RequestEvent{ sync.ExpTrigger(t, 6, false, blockSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2, Timeout: true}, sync.TestReqEvent(request.EvFail, req2, nil),
}, nil)) }))
tracker.ExpRequests(t, 5, nil) tracker.ExpRequests(t, 5, nil)
// set as validated head before retrieving block; now it's assumed to be available from server 2 too // set as validated head before retrieving block; now it's assumed to be available from server 2 too
ht.validated.Header = blockHeader(testBlock2) 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 // now head block is unavailable again
expHeadBlock(4, nil) expHeadBlock(4, nil)
// expect req2 retry to server 2 // expect req2 retry to server 2
req2r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 3}, Request: sync.ReqBeaconBlock(head2.BlockRoot)} req2r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 3}, Request: sync.ReqBeaconBlock(head2.BlockRoot)}
tracker.ExpRequests(t, 6, []request.RequestWithID{req2r}) tracker.ExpRequests(t, 6, []request.RequestWithID{req2r})
// valid response, now head block should be block 2 immediately as it is already validated // 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{ sync.ExpTrigger(t, 8, true, blockSync.Process(tracker, []request.Event{
request.RequestEvent{RequestWithID: req2r, Response: testBlock2, Finalized: true}, sync.TestReqEvent(request.EvResponse, req2r, testBlock2),
}, nil)) }))
expHeadBlock(5, testBlock2) expHeadBlock(5, testBlock2)
} }