mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
beacon/light: ensure EvRequest before SendRequest returns
This commit is contained in:
parent
d42b0e9b55
commit
42d6237eda
4 changed files with 28 additions and 24 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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)}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
Loading…
Reference in a new issue