diff --git a/beacon/light/api/api_server.go b/beacon/light/api/api_server.go index c01f012957..8f20455365 100755 --- a/beacon/light/api/api_server.go +++ b/beacon/light/api/api_server.go @@ -59,8 +59,8 @@ func (s *ApiServer) Subscribe(eventCallback func(event request.Event)) { // SendRequest implements request.requestServer. func (s *ApiServer) SendRequest(req request.Request) request.ID { id := request.ID(atomic.AddUint64(&s.lastId, 1)) + s.eventCallback(request.Event{Type: request.EvRequest, Data: request.RequestResponse{ID: id, Request: req}}) go func() { - s.eventCallback(request.Event{Type: request.EvRequest, Data: request.RequestResponse{ID: id, Request: req}}) var resp request.Response switch data := req.(type) { case sync.ReqUpdates: diff --git a/beacon/light/request/server.go b/beacon/light/request/server.go index 683da27aa3..c10f744d33 100644 --- a/beacon/light/request/server.go +++ b/beacon/light/request/server.go @@ -55,9 +55,9 @@ const ( // requestServer can send requests in a non-blocking way and feed back events // through the event callback. When successfully sending a request it should -// send back an EvRequest event. When finished, it should send back either -// EvResponse or EvFail. Additionally, it may also send application-defined -// events that the Modules can interpret. +// send back an EvRequest event before returning from SendRequest. When finished, +// it should send back either EvResponse or EvFail. Additionally, it may also +// send application-defined events that the Modules can interpret. type requestServer interface { Subscribe(eventCallback func(event Event)) SendRequest(request Request) ID @@ -151,6 +151,11 @@ func (s *serverWithTimeout) subscribe(eventCallback func(event Event)) { s.parent.Subscribe(s.eventCallback) } +// sendRequest sends a request through the parent (requestServer). +func (s *serverWithTimeout) sendRequest(request Request) (reqId ID) { + return s.parent.SendRequest(request) +} + // eventCallback is called by parent (requestServer) event subscription. func (s *serverWithTimeout) eventCallback(event Event) { s.lock.Lock() @@ -159,6 +164,7 @@ func (s *serverWithTimeout) eventCallback(event Event) { switch event.Type { case EvRequest: s.startTimeout(event.Data.(RequestResponse)) + s.childEventCb(event) case EvResponse, EvFail: id := event.Data.(RequestResponse).ID if timer, ok := s.timeouts[id]; ok { @@ -205,11 +211,6 @@ func (s *serverWithTimeout) startTimeout(reqData RequestResponse) { }) } -// sendRequest sends a request through the parent (requestServer). -func (s *serverWithTimeout) sendRequest(request Request) (reqId ID) { - return s.parent.SendRequest(request) -} - // stop stops all goroutines associated with the server. func (s *serverWithTimeout) unsubscribe() { s.lock.Lock() @@ -315,9 +316,8 @@ func (s *serverWithLimits) eventCallback(event Event) { // sendRequest sends a request through the parent (serverWithTimeout). func (s *serverWithLimits) sendRequest(request Request) (reqId ID) { s.lock.Lock() - defer s.lock.Unlock() - s.pendingCount++ + s.lock.Unlock() return s.serverWithTimeout.sendRequest(request) } diff --git a/cmd/blsync/block_sync.go b/cmd/blsync/block_sync.go index e4627d8e8d..7305c2f422 100755 --- a/cmd/blsync/block_sync.go +++ b/cmd/blsync/block_sync.go @@ -33,8 +33,8 @@ type beaconBlockSync struct { serverHeads map[request.Server]common.Hash headTracker headTracker - lastHeadBlock *capella.BeaconBlock - headCh chan headData + lastHeadInfo types.HeadInfo + headCh chan headData } type headData struct { @@ -55,7 +55,7 @@ func newBeaconBlockSync(headTracker headTracker) *beaconBlockSync { recentBlocks: lru.NewCache[common.Hash, *capella.BeaconBlock](10), locked: make(map[common.Hash]struct{}), serverHeads: make(map[request.Server]common.Hash), - headCh: make(chan headData, 1), + headCh: make(chan headData, 1), } } @@ -86,18 +86,22 @@ func (s *beaconBlockSync) Process(events []request.Event) { if !ok { return } - finality, ok := s.headTracker.ValidatedFinality() //TODO fetch directly if subscription does not deliver + finality, ok := s.headTracker.ValidatedFinality() //TODO fetch directly if subscription does not deliver if !ok || head.Header.Epoch() != finality.Attested.Header.Epoch() { return } validatedHead := head.Header.Hash() headBlock, ok := s.recentBlocks.Get(validatedHead) - if !ok || headBlock == s.lastHeadBlock { + if !ok { + return + } + headInfo := blockHeadInfo(headBlock) + if headInfo == s.lastHeadInfo { return } select { case s.headCh <- headData{block: headBlock, update: finality}: - s.lastHeadBlock = headBlock + s.lastHeadInfo = headInfo default: } } @@ -122,3 +126,10 @@ func (s *beaconBlockSync) MakeRequest(server request.Server) (request.Request, f } return nil, 0 } + +func blockHeadInfo(block *capella.BeaconBlock) types.HeadInfo { + if block == nil { + return types.HeadInfo{} + } + return types.HeadInfo{Slot: uint64(block.Slot), BlockRoot: beaconBlockHash(block)} +} diff --git a/cmd/blsync/block_sync_test.go b/cmd/blsync/block_sync_test.go index 7547b061e9..93526f11d0 100644 --- a/cmd/blsync/block_sync_test.go +++ b/cmd/blsync/block_sync_test.go @@ -103,13 +103,6 @@ func TestBlockSync(t *testing.T) { expHeadBlock(5, testBlock2) } -func blockHeadInfo(block *capella.BeaconBlock) types.HeadInfo { - if block == nil { - return types.HeadInfo{} - } - return types.HeadInfo{Slot: uint64(block.Slot), BlockRoot: beaconBlockHash(block)} -} - func blockHeader(block *capella.BeaconBlock) types.Header { return types.Header{ Slot: uint64(block.Slot),