beacon/light/api: pass event listener as struct

This commit is contained in:
Felix Lange 2024-03-04 16:37:28 +01:00
parent 96df24d81e
commit fa1bbe7a25
2 changed files with 41 additions and 26 deletions

View file

@ -40,18 +40,24 @@ func NewApiServer(api *BeaconLightApi) *ApiServer {
// Subscribe implements request.requestServer. // Subscribe implements request.requestServer.
func (s *ApiServer) Subscribe(eventCallback func(event request.Event)) { func (s *ApiServer) Subscribe(eventCallback func(event request.Event)) {
s.eventCallback = eventCallback s.eventCallback = eventCallback
s.unsubscribe = s.api.StartHeadListener(func(slot uint64, blockRoot common.Hash) { listener := HeadEventListener{
OnNewHead: func(slot uint64, blockRoot common.Hash) {
log.Debug("New head received", "slot", slot, "blockRoot", blockRoot) log.Debug("New head received", "slot", slot, "blockRoot", blockRoot)
eventCallback(request.Event{Type: sync.EvNewHead, Data: types.HeadInfo{Slot: slot, BlockRoot: blockRoot}}) eventCallback(request.Event{Type: sync.EvNewHead, Data: types.HeadInfo{Slot: slot, BlockRoot: blockRoot}})
}, func(head types.SignedHeader) { },
OnSignedHead: func(head types.SignedHeader) {
log.Debug("New signed head received", "slot", head.Header.Slot, "blockRoot", head.Header.Hash(), "signerCount", head.Signature.SignerCount()) log.Debug("New signed head received", "slot", head.Header.Slot, "blockRoot", head.Header.Hash(), "signerCount", head.Signature.SignerCount())
eventCallback(request.Event{Type: sync.EvNewSignedHead, Data: head}) eventCallback(request.Event{Type: sync.EvNewSignedHead, Data: head})
}, func(head types.FinalityUpdate) { },
OnFinality: func(head types.FinalityUpdate) {
log.Debug("New finality update received", "slot", head.Attested.Slot, "blockRoot", head.Attested.Hash(), "signerCount", head.Signature.SignerCount()) log.Debug("New finality update received", "slot", head.Attested.Slot, "blockRoot", head.Attested.Hash(), "signerCount", head.Signature.SignerCount())
eventCallback(request.Event{Type: sync.EvNewFinalityUpdate, Data: head}) eventCallback(request.Event{Type: sync.EvNewFinalityUpdate, Data: head})
}, func(err error) { },
OnError: func(err error) {
log.Warn("Head event stream error", "err", err) log.Warn("Head event stream error", "err", err)
}) },
}
s.unsubscribe = s.api.StartHeadListener(listener)
} }
// SendRequest implements request.requestServer. // SendRequest implements request.requestServer.

View file

@ -393,11 +393,18 @@ func decodeHeadEvent(enc []byte) (uint64, common.Hash, error) {
return uint64(data.Slot), data.Block, nil return uint64(data.Slot), data.Block, nil
} }
type HeadEventListener struct {
OnNewHead func(slot uint64, blockRoot common.Hash)
OnSignedHead func(head types.SignedHeader)
OnFinality func(head types.FinalityUpdate)
OnError func(err error)
}
// StartHeadListener creates an event subscription for heads and signed (optimistic) // StartHeadListener creates an event subscription for heads and signed (optimistic)
// head updates and calls the specified callback functions when they are received. // head updates and calls the specified callback functions when they are received.
// The callbacks are also called for the current head and optimistic head at startup. // The callbacks are also called for the current head and optimistic head at startup.
// They are never called concurrently. // They are never called concurrently.
func (api *BeaconLightApi) StartHeadListener(headFn func(slot uint64, blockRoot common.Hash), signedFn func(head types.SignedHeader), finalityFn func(head types.FinalityUpdate), errFn func(err error)) func() { func (api *BeaconLightApi) StartHeadListener(listener HeadEventListener) func() {
closeCh := make(chan struct{}) // initiate closing the stream closeCh := make(chan struct{}) // initiate closing the stream
closedCh := make(chan struct{}) // stream closed (or failed to create) closedCh := make(chan struct{}) // stream closed (or failed to create)
stoppedCh := make(chan struct{}) // sync loop stopped stoppedCh := make(chan struct{}) // sync loop stopped
@ -411,7 +418,7 @@ func (api *BeaconLightApi) StartHeadListener(headFn func(slot uint64, blockRoot
req, err := http.NewRequest("GET", api.url+ req, err := http.NewRequest("GET", api.url+
"/eth/v1/events?topics=head&topics=light_client_optimistic_update&topics=light_client_finality_update", nil) "/eth/v1/events?topics=head&topics=light_client_optimistic_update&topics=light_client_finality_update", nil)
if err != nil { if err != nil {
errFn(fmt.Errorf("error creating event subscription request: %v", err)) listener.OnError(fmt.Errorf("error creating event subscription request: %v", err))
return return
} }
for k, v := range api.customHeaders { for k, v := range api.customHeaders {
@ -419,7 +426,7 @@ func (api *BeaconLightApi) StartHeadListener(headFn func(slot uint64, blockRoot
} }
stream, err := eventsource.SubscribeWithRequest("", req) stream, err := eventsource.SubscribeWithRequest("", req)
if err != nil { if err != nil {
errFn(fmt.Errorf("error creating event subscription: %v", err)) listener.OnError(fmt.Errorf("error creating event subscription: %v", err))
close(streamCh) close(streamCh)
return return
} }
@ -427,22 +434,24 @@ func (api *BeaconLightApi) StartHeadListener(headFn func(slot uint64, blockRoot
<-closeCh <-closeCh
stream.Close() stream.Close()
}() }()
go func() { go func() {
defer close(stoppedCh) defer close(stoppedCh)
if head, err := api.GetHeader(common.Hash{}); err == nil { if head, err := api.GetHeader(common.Hash{}); err == nil {
headFn(head.Slot, head.Hash()) listener.OnNewHead(head.Slot, head.Hash())
} }
if signedHead, err := api.GetOptimisticHeadUpdate(); err == nil { if signedHead, err := api.GetOptimisticHeadUpdate(); err == nil {
signedFn(signedHead) listener.OnSignedHead(signedHead)
} }
if finalityUpdate, err := api.GetFinalityUpdate(); err == nil { if finalityUpdate, err := api.GetFinalityUpdate(); err == nil {
finalityFn(finalityUpdate) listener.OnFinality(finalityUpdate)
} }
stream := <-streamCh stream := <-streamCh
if stream == nil { if stream == nil {
return return
} }
for { for {
select { select {
case event, ok := <-stream.Events: case event, ok := <-stream.Events:
@ -452,30 +461,30 @@ func (api *BeaconLightApi) StartHeadListener(headFn func(slot uint64, blockRoot
switch event.Event() { switch event.Event() {
case "head": case "head":
if slot, blockRoot, err := decodeHeadEvent([]byte(event.Data())); err == nil { if slot, blockRoot, err := decodeHeadEvent([]byte(event.Data())); err == nil {
headFn(slot, blockRoot) listener.OnNewHead(slot, blockRoot)
} else { } else {
errFn(fmt.Errorf("error decoding head event: %v", err)) listener.OnError(fmt.Errorf("error decoding head event: %v", err))
} }
case "light_client_optimistic_update": case "light_client_optimistic_update":
if signedHead, err := decodeOptimisticHeadUpdate([]byte(event.Data())); err == nil { if signedHead, err := decodeOptimisticHeadUpdate([]byte(event.Data())); err == nil {
signedFn(signedHead) listener.OnSignedHead(signedHead)
} else { } else {
errFn(fmt.Errorf("error decoding optimistic update event: %v", err)) listener.OnError(fmt.Errorf("error decoding optimistic update event: %v", err))
} }
case "light_client_finality_update": case "light_client_finality_update":
if finalityUpdate, err := decodeFinalityUpdate([]byte(event.Data())); err == nil { if finalityUpdate, err := decodeFinalityUpdate([]byte(event.Data())); err == nil {
finalityFn(finalityUpdate) listener.OnFinality(finalityUpdate)
} else { } else {
errFn(fmt.Errorf("error decoding finality update event: %v", err)) listener.OnError(fmt.Errorf("error decoding finality update event: %v", err))
} }
default: default:
errFn(fmt.Errorf("unexpected event: %s", event.Event())) listener.OnError(fmt.Errorf("unexpected event: %s", event.Event()))
} }
case err, ok := <-stream.Errors: case err, ok := <-stream.Errors:
if !ok { if !ok {
break break
} }
errFn(err) listener.OnError(err)
} }
} }
}() }()