beacon/light/request: bugs fixed

This commit is contained in:
Zsolt Felfoldi 2024-01-14 04:26:50 +01:00 committed by Felix Lange
parent 1db975fd34
commit 30bdde91ae
2 changed files with 16 additions and 11 deletions

View file

@ -152,14 +152,8 @@ func (s *Scheduler) RegisterServer(rs requestServer) {
server := newServer(rs, s.clock) server := newServer(rs, s.clock)
s.addEvent(Event{Type: EvRegistered, Server: server}) s.addEvent(Event{Type: EvRegistered, Server: server})
server.subscribe(func(event Event) { server.subscribe(func(event Event) {
s.lock.Lock() event.Server = server
if _, ok := s.servers[server]; ok { s.addEvent(event)
event.Server = server
s.addEvent(event)
} else {
log.Error("Event received from unsubscribed server")
}
s.lock.Unlock()
}) })
s.servers[server] = struct{}{} s.servers[server] = struct{}{}
} }
@ -205,8 +199,11 @@ func (s *Scheduler) syncLoop() {
for { for {
s.lock.Lock() s.lock.Lock()
s.handleEvents() s.handleEvents()
for s.targetChanged() { for {
s.processModules() s.processModules()
if !s.targetChanged() {
break
}
} }
s.sendRequests() s.sendRequests()
s.lock.Unlock() s.lock.Unlock()
@ -330,7 +327,11 @@ func (s *Scheduler) handleEvents() {
s.events = nil s.events = nil
s.eventLock.Unlock() s.eventLock.Unlock()
for _, event := range events { for _, event := range events {
s.handleEvent(event) if _, ok := s.servers[event.Server.(server)]; ok {
s.handleEvent(event)
} else {
log.Error("Event received from unsubscribed server")
}
} }
} }

View file

@ -55,6 +55,10 @@ func newBeaconBlockSync(headTracker headTracker) *beaconBlockSync {
func (s *beaconBlockSync) HandleEvent(event request.Event) { func (s *beaconBlockSync) HandleEvent(event request.Event) {
switch event.Type { switch event.Type {
case request.EvRequest:
_, req, _ := event.RequestInfo()
blockRoot := common.Hash(req.(sync.ReqBeaconBlock))
s.locked[blockRoot] = struct{}{}
case request.EvResponse, request.EvFail, request.EvTimeout: case request.EvResponse, request.EvFail, request.EvTimeout:
_, req, resp := event.RequestInfo() _, req, resp := event.RequestInfo()
blockRoot := common.Hash(req.(sync.ReqBeaconBlock)) blockRoot := common.Hash(req.(sync.ReqBeaconBlock))
@ -95,7 +99,7 @@ func (s *beaconBlockSync) MakeRequest(server request.Server) (request.Request, f
} }
} }
// request prefetch head if the given server has announced it // request prefetch head if the given server has announced it
if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead != (common.Hash{}) && prefetchHead != s.serverHeads[server] { if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead == s.serverHeads[server] {
if _, ok := s.recentBlocks.Get(prefetchHead); !ok { if _, ok := s.recentBlocks.Get(prefetchHead); !ok {
if _, ok := s.locked[prefetchHead]; !ok { if _, ok := s.locked[prefetchHead]; !ok {
return sync.ReqBeaconBlock(prefetchHead), 0 return sync.ReqBeaconBlock(prefetchHead), 0