From c7d884c770d16cabef1e0c6954e7bc267ec05734 Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Thu, 11 Jan 2024 17:34:41 +0100 Subject: [PATCH] beacon/light: automatic target data trigger --- beacon/light/committee_chain.go | 11 ++ beacon/light/head_tracker.go | 13 ++ beacon/light/request/scheduler.go | 40 +++++-- beacon/light/request/server.go | 2 +- beacon/light/request/tracker.go | 2 +- beacon/light/sync/head_sync.go | 34 ++---- beacon/light/sync/head_sync_test.go | 66 +++++------ beacon/light/sync/test_helpers.go | 6 - beacon/light/sync/update_sync.go | 14 +-- beacon/light/sync/update_sync_test.go | 88 +++++++------- cmd/blsync/block_sync.go | 163 ++++---------------------- cmd/blsync/block_sync_test.go | 28 ++--- cmd/blsync/engine_api.go | 134 +++++++++++++++++++++ cmd/blsync/main.go | 37 +----- 14 files changed, 326 insertions(+), 312 deletions(-) create mode 100644 cmd/blsync/engine_api.go diff --git a/beacon/light/committee_chain.go b/beacon/light/committee_chain.go index a090e298d2..08c4968051 100644 --- a/beacon/light/committee_chain.go +++ b/beacon/light/committee_chain.go @@ -70,6 +70,7 @@ type CommitteeChain struct { committees *canonicalStore[*types.SerializedSyncCommittee] fixedCommitteeRoots *canonicalStore[common.Hash] committeeCache *lru.Cache[uint64, syncCommittee] // cache deserialized committees + changeCounter uint64 clock mclock.Clock // monotonic clock (simulated clock in tests) unixNano func() int64 // system clock (simulated clock in tests) @@ -186,6 +187,7 @@ func (s *CommitteeChain) Reset() { if err := s.rollback(0); err != nil { log.Error("Error writing batch into chain database", "error", err) } + s.changeCounter++ } // CheckpointInit initializes a CommitteeChain based on a checkpoint. @@ -219,6 +221,7 @@ func (s *CommitteeChain) CheckpointInit(bootstrap types.BootstrapData) error { s.Reset() return err } + s.changeCounter++ return nil } @@ -371,6 +374,7 @@ func (s *CommitteeChain) InsertUpdate(update *types.LightClientUpdate, nextCommi return ErrWrongCommitteeRoot } } + s.changeCounter++ if reorg { if err := s.rollback(period + 1); err != nil { return err @@ -409,6 +413,13 @@ func (s *CommitteeChain) NextSyncPeriod() (uint64, bool) { return s.committees.periods.End - 1, true } +func (s *CommitteeChain) ChangeCounter() uint64 { + s.chainmu.RLock() + defer s.chainmu.RUnlock() + + return s.changeCounter +} + // rollback removes all committees and fixed roots from the given period and updates // starting from the previous period. func (s *CommitteeChain) rollback(period uint64) error { diff --git a/beacon/light/head_tracker.go b/beacon/light/head_tracker.go index 56301c386b..a689280492 100644 --- a/beacon/light/head_tracker.go +++ b/beacon/light/head_tracker.go @@ -35,6 +35,7 @@ type HeadTracker struct { signedHead types.SignedHeader headSignerCount int prefetchHead types.HeadInfo + changeCounter uint64 } // NewHeadTracker creates a new HeadTracker. @@ -82,6 +83,7 @@ func (h *HeadTracker) Validate(head types.SignedHeader) (bool, error) { return false, errors.New("invalid header signature") } h.signedHead, h.headSignerCount = head, signerCount + h.changeCounter++ return true, nil } @@ -104,5 +106,16 @@ func (h *HeadTracker) SetPrefetchHead(head types.HeadInfo) { h.lock.Lock() defer h.lock.Unlock() + if head == h.prefetchHead { + return + } h.prefetchHead = head + h.changeCounter++ +} + +func (h *HeadTracker) ChangeCounter() uint64 { + h.lock.RLock() + defer h.lock.RUnlock() + + return h.changeCounter } diff --git a/beacon/light/request/scheduler.go b/beacon/light/request/scheduler.go index f7bca17568..eaa3a5eb62 100644 --- a/beacon/light/request/scheduler.go +++ b/beacon/light/request/scheduler.go @@ -41,13 +41,12 @@ type Module interface { // a processing round is triggered. It can start new requests through the // received Tracker, process events and/or do other data processing tasks. // Note that request events are only passed to the module that made the given - // request while server events are passed to every module. Process can also - // trigger a next processing round by returning true. + // request while server events are passed to every module. // // Note: Process functions of different modules are never called concurrently; // they are called by Scheduler in the same order of priority as they were // registered in. - Process(Tracker, []Event) bool + Process(Tracker, []Event) } // Scheduler is a modular network data retrieval framework that coordinates multiple @@ -63,6 +62,7 @@ type Scheduler struct { trackers map[Module]*tracker servers map[server]struct{} pending map[ServerAndID]pendingRequest + target map[targetData]uint64 serverEvents []Event stopCh chan chan struct{} @@ -71,6 +71,10 @@ type Scheduler struct { // testTimerResults []bool // true is appended when simulated timer is processed; false when stopped } +type targetData interface { + ChangeCounter() uint64 +} + // pendingRequest keeps track of sent and not yet finalized requests and their // sender modules. type pendingRequest struct { @@ -86,6 +90,7 @@ func NewScheduler(clock mclock.Clock) *Scheduler { names: make(map[Module]string), trackers: make(map[Module]*tracker), pending: make(map[ServerAndID]pendingRequest), + target: make(map[targetData]uint64), stopCh: make(chan chan struct{}), // Note: testWaitCh should not have capacity in order to ensure // that after a trigger happens testWaitCh will block until the resulting @@ -96,6 +101,13 @@ func NewScheduler(clock mclock.Clock) *Scheduler { return s } +func (s *Scheduler) RegisterTarget(t targetData) { + s.lock.Lock() + defer s.lock.Unlock() + + s.target[t] = 0 +} + // RegisterModule registers a module. Should be called before starting the scheduler. // In each processing round the order of module processing depends on the order of // registration. @@ -155,7 +167,7 @@ func (s *Scheduler) Start() { // Stop stops the scheduler. func (s *Scheduler) Stop() { s.lock.Lock() - for server, _ := range s.servers { + for server := range s.servers { server.unsubscribe() } s.servers = nil @@ -171,6 +183,9 @@ func (s *Scheduler) Stop() { func (s *Scheduler) syncLoop() { for { s.processModules() + for s.targetChanged() { + s.processModules() + } loop: for { select { @@ -185,12 +200,22 @@ func (s *Scheduler) syncLoop() { } } +func (s *Scheduler) targetChanged() (changed bool) { + for target, counter := range s.target { + if newCounter := target.ChangeCounter(); newCounter != counter { + s.target[target] = newCounter + changed = true + } + } + return +} + // processModules runs an entire processing round, calling the Process functions // of all modules, passing all relevant events. func (s *Scheduler) processModules() { s.lock.Lock() servers := make(serverSet) - for server, _ := range s.servers { + for server := range s.servers { if ok, _ := server.canRequestNow(); ok { servers[server] = struct{}{} } @@ -225,10 +250,7 @@ func (s *Scheduler) processModules() { } } log.Debug("Processing module", "name", s.names[module], "responses", respCount, "fails", failCount, "timeouts", timeoutCount) - - if module.Process(tracker, append(serverEvents, requestEvents...)) { - s.Trigger() - } + module.Process(tracker, append(serverEvents, requestEvents...)) } } diff --git a/beacon/light/request/server.go b/beacon/light/request/server.go index 8d0d7d97ad..936c587783 100644 --- a/beacon/light/request/server.go +++ b/beacon/light/request/server.go @@ -238,7 +238,7 @@ func (s *serverWithTimeout) stopTimer(timer mclock.Timer) { // failures of the server might happen sometimes, but still avoids hammering a // non-functional server with requests. // -//TODO protect against excessive server events +// TODO protect against excessive server events type serverWithLimits struct { serverWithTimeout lock sync.Mutex diff --git a/beacon/light/request/tracker.go b/beacon/light/request/tracker.go index a50ae127b5..79b91cf405 100644 --- a/beacon/light/request/tracker.go +++ b/beacon/light/request/tracker.go @@ -88,7 +88,7 @@ func (p *tracker) TryRequest(requestFn func(server Server) (Request, float32)) ( maxServerPriority, maxRequestPriority = -math.MaxFloat32, -math.MaxFloat32 serverCount := len(p.servers) var removed, candidates int - for server, _ := range p.servers { + for server := range p.servers { canRequest, serverPriority := server.canRequestNow() if !canRequest { delete(p.servers, server) diff --git a/beacon/light/sync/head_sync.go b/beacon/light/sync/head_sync.go index 49593f2443..8d6426d151 100644 --- a/beacon/light/sync/head_sync.go +++ b/beacon/light/sync/head_sync.go @@ -68,26 +68,20 @@ func NewHeadSync(headTracker headTracker, chain committeeChain) *HeadSync { } // Process implements request.Module -func (s *HeadSync) Process(tracker request.Tracker, events []request.Event) (trigger bool) { +func (s *HeadSync) Process(tracker request.Tracker, events []request.Event) { nextPeriod, chainInit := s.chain.NextSyncPeriod() if nextPeriod != s.nextSyncPeriod || chainInit != s.chainInit { s.nextSyncPeriod, s.chainInit = nextPeriod, chainInit - trigger = s.processUnvalidatedHeads() + s.processUnvalidatedHeads() } for _, event := range events { switch event.Type { case EvNewHead: - if s.setServerHead(event.Server, event.Data.(types.HeadInfo)) { - trigger = true - } + s.setServerHead(event.Server, event.Data.(types.HeadInfo)) case EvNewSignedHead: - if s.newSignedHead(event.Server, event.Data.(types.SignedHeader)) { - trigger = true - } + s.newSignedHead(event.Server, event.Data.(types.SignedHeader)) case request.EvUnregistered: - if s.setServerHead(event.Server, types.HeadInfo{}) { - trigger = true - } + s.setServerHead(event.Server, types.HeadInfo{}) delete(s.serverHeads, event.Server) delete(s.unvalidatedHeads, event.Server) } @@ -97,35 +91,31 @@ func (s *HeadSync) Process(tracker request.Tracker, events []request.Event) (tri // newSignedHead handles received signed head; either validates it if the chain // is properly synced or stores it for further validation. -func (s *HeadSync) newSignedHead(server request.Server, signedHead types.SignedHeader) (trigger bool) { +func (s *HeadSync) newSignedHead(server request.Server, signedHead types.SignedHeader) { if !s.chainInit || types.SyncPeriod(signedHead.SignatureSlot) > s.nextSyncPeriod { s.unvalidatedHeads[server] = signedHead - return false + return } - updated, _ := s.headTracker.Validate(signedHead) - return updated + s.headTracker.Validate(signedHead) } // processUnvalidatedHeads iterates the list of unvalidated heads and validates // those which can be validated. -func (s *HeadSync) processUnvalidatedHeads() (trigger bool) { +func (s *HeadSync) processUnvalidatedHeads() { if !s.chainInit { - return false + return } for server, signedHead := range s.unvalidatedHeads { if types.SyncPeriod(signedHead.SignatureSlot) <= s.nextSyncPeriod { - if updated, _ := s.headTracker.Validate(signedHead); updated { - trigger = true - } + s.headTracker.Validate(signedHead) delete(s.unvalidatedHeads, server) } } - return } // setServerHead processes non-validated server head announcements and updates // the prefetch head if necessary. -//TODO report server failure if a server announces many heads that do not become validated soon. +// TODO report server failure if a server announces many heads that do not become validated soon. func (s *HeadSync) setServerHead(server request.Server, head types.HeadInfo) bool { if oldHead, ok := s.serverHeads[server]; ok { if head == oldHead { diff --git a/beacon/light/sync/head_sync_test.go b/beacon/light/sync/head_sync_test.go index 993de7e6e7..ee0d8563b5 100644 --- a/beacon/light/sync/head_sync_test.go +++ b/beacon/light/sync/head_sync_test.go @@ -50,40 +50,40 @@ func TestValidatedHead(t *testing.T) { headSync := NewHeadSync(ht, chain) ht.ExpValidated(t, 1, nil) - ExpTrigger(t, 1, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: testSHead1}, - })) + }) ht.ExpValidated(t, 2, nil) chain.SetNextSyncPeriod(0) - ExpTrigger(t, 2, true, headSync.Process(tracker, nil)) + headSync.Process(tracker, nil) ht.ExpValidated(t, 3, []types.SignedHeader{testSHead1}) chain.SetNextSyncPeriod(1) - ExpTrigger(t, 3, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: EvNewSignedHead, Data: testSHead2}, {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewSignedHead, Data: testSHead2}, - })) + }) ht.ExpValidated(t, 4, []types.SignedHeader{testSHead2, testSHead2}) - ExpTrigger(t, 4, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: EvNewSignedHead, Data: testSHead3}, {Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: EvNewSignedHead, Data: testSHead4}, - })) + }) ht.ExpValidated(t, 5, nil) chain.SetNextSyncPeriod(2) - ExpTrigger(t, 5, true, headSync.Process(tracker, nil)) + headSync.Process(tracker, nil) ht.ExpValidated(t, 6, []types.SignedHeader{testSHead3}) - ExpTrigger(t, 6, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvUnregistered}, - })) + }) ht.ExpValidated(t, 7, nil) chain.SetNextSyncPeriod(3) - ExpTrigger(t, 7, false, headSync.Process(tracker, nil)) + headSync.Process(tracker, nil) ht.ExpValidated(t, 8, nil) - ExpTrigger(t, 8, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: EvNewSignedHead, Data: testSHead4}, - })) + }) ht.ExpValidated(t, 9, []types.SignedHeader{testSHead4}) } @@ -94,48 +94,48 @@ func TestPrefetchHead(t *testing.T) { headSync := NewHeadSync(ht, chain) ht.ExpPrefetch(t, 1, testHead0) // no servers registered - ExpTrigger(t, 1, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewHead, Data: testHead1}, - })) + }) ht.ExpPrefetch(t, 2, testHead1) // s1: h1 - ExpTrigger(t, 2, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewHead, Data: testHead2}, - })) + }) ht.ExpPrefetch(t, 3, testHead2) // s1: h1, s2: h2 - ExpTrigger(t, 3, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: EvNewHead, Data: testHead2}, - })) + }) ht.ExpPrefetch(t, 4, testHead2) // s1: h2, s2: h2 - ExpTrigger(t, 4, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: EvNewHead, Data: testHead3}, - })) + }) ht.ExpPrefetch(t, 5, testHead2) // s1: h2, s2: h2, s3: h3 - ExpTrigger(t, 5, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer4, Type: request.EvRegistered}, {Server: testServer4, Type: EvNewHead, Data: testHead4}, - })) + }) ht.ExpPrefetch(t, 6, testHead2) // s1: h2, s2: h2, s3: h3, s4: h4 - ExpTrigger(t, 6, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: EvNewHead, Data: testHead3}, - })) + }) ht.ExpPrefetch(t, 7, testHead3) // s1: h2, s2: h3, s3: h3, s4: h4 - ExpTrigger(t, 7, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvUnregistered}, - })) + }) ht.ExpPrefetch(t, 8, testHead4) // s1: h2, s2: h3, s4: h4 - ExpTrigger(t, 8, false, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvUnregistered}, - })) + }) ht.ExpPrefetch(t, 9, testHead4) // s2: h3, s4: h4 - ExpTrigger(t, 9, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer4, Type: request.EvUnregistered}, - })) + }) ht.ExpPrefetch(t, 10, testHead3) // s2: h3 - ExpTrigger(t, 10, true, headSync.Process(tracker, []request.Event{ + headSync.Process(tracker, []request.Event{ {Server: testServer2, Type: request.EvUnregistered}, - })) + }) ht.ExpPrefetch(t, 11, testHead0) // no servers registered } diff --git a/beacon/light/sync/test_helpers.go b/beacon/light/sync/test_helpers.go index c393a9c49f..4503552a08 100644 --- a/beacon/light/sync/test_helpers.go +++ b/beacon/light/sync/test_helpers.go @@ -91,12 +91,6 @@ func (tt *TestTracker) InvalidResponse(id request.ServerAndID, desc string) { return } -func ExpTrigger(t *testing.T, tci int, expTrigger, trigger bool) { - if trigger != expTrigger { - t.Errorf("Invalid process trigger output in test case #%d (expected %v, got %v)", tci, expTrigger, trigger) - } -} - func TestReqEvent(evType *request.EventType, req request.RequestWithID, response request.Response) request.Event { return request.Event{ Type: evType, diff --git a/beacon/light/sync/update_sync.go b/beacon/light/sync/update_sync.go index 8d2d6e3296..417e23143b 100644 --- a/beacon/light/sync/update_sync.go +++ b/beacon/light/sync/update_sync.go @@ -53,9 +53,9 @@ func NewCheckpointInit(chain committeeChain, checkpointHash common.Hash) *Checkp } // Process implements request.Module -func (s *CheckpointInit) Process(tracker request.Tracker, events []request.Event) bool { +func (s *CheckpointInit) Process(tracker request.Tracker, events []request.Event) { if s.initialized { - return false + return } for _, event := range events { if !event.IsRequestEvent() { @@ -67,7 +67,7 @@ func (s *CheckpointInit) Process(tracker request.Tracker, events []request.Event if checkpoint, ok := response.(*types.BootstrapData); ok && checkpoint.Header.Hash() == common.Hash(request.(ReqCheckpointData)) { s.chain.CheckpointInit(*checkpoint) //TODO s.initialized = true - return true + return } tracker.InvalidResponse(sid, "invalid checkpoint data") } @@ -79,7 +79,7 @@ func (s *CheckpointInit) Process(tracker request.Tracker, events []request.Event s.locked = true } } - return false + return } // ForwardUpdateSync implements request.Module; it fetches updates between the @@ -226,7 +226,7 @@ func (u updateResponseList) Less(i, j int) bool { } // Process implements request.Module -func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Event) (trigger bool) { +func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Event) { // iterate events and add responses to process queue for _, event := range events { switch event.Type { @@ -258,7 +258,6 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Ev if !s.processResponse(tracker, event) { break } - trigger = true sid, req, _ := event.RequestInfo() s.unlockRange(sid, req) s.processQueue = s.processQueue[1:] @@ -270,7 +269,7 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Ev // start new requests if necessary startPeriod, chainInit := s.chain.NextSyncPeriod() if !chainInit { - return false + return } for { firstPeriod, maxCount := s.rangeLock.firstUnlocked(startPeriod, maxUpdateRequest) @@ -290,5 +289,4 @@ func (s *ForwardUpdateSync) Process(tracker request.Tracker, events []request.Ev break } } - return } diff --git a/beacon/light/sync/update_sync_test.go b/beacon/light/sync/update_sync_test.go index 413c763077..0f1b9356e6 100644 --- a/beacon/light/sync/update_sync_test.go +++ b/beacon/light/sync/update_sync_test.go @@ -32,41 +32,41 @@ func TestCheckpointInit(t *testing.T) { checkpoint := &types.BootstrapData{Header: types.Header{Slot: 0x2000*4 + 0x1000}} // period 4 checkpointHash := checkpoint.Header.Hash() chkInit := NewCheckpointInit(chain, checkpointHash) - ExpTrigger(t, 1, false, chkInit.Process(tracker, []request.Event{ + chkInit.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer2, Type: request.EvRegistered}, - })) + }) // expect bootstrap request to server 1 req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 1}, Request: ReqCheckpointData(checkpointHash)} tracker.ExpRequests(t, 1, []request.RequestWithID{req1}) // req1 times out; expect request to server 2 - ExpTrigger(t, 2, false, chkInit.Process(tracker, []request.Event{ + chkInit.Process(tracker, []request.Event{ TestReqEvent(request.EvTimeout, req1, nil), - })) + }) req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqCheckpointData(checkpointHash)} tracker.ExpRequests(t, 2, []request.RequestWithID{req2}) // invalid response to req2; expect init state to still be false wrongCheckpoint := &types.BootstrapData{Header: types.Header{Slot: 123456}} - ExpTrigger(t, 3, false, chkInit.Process(tracker, []request.Event{ + chkInit.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req2, wrongCheckpoint), - })) + }) // req1 fails (hard timeout) - ExpTrigger(t, 4, false, chkInit.Process(tracker, []request.Event{ + chkInit.Process(tracker, []request.Event{ TestReqEvent(request.EvFail, req1, nil), - })) + }) chain.ExpInit(t, false) // server 3 is registered tracker.AddServer(testServer3, 1) - ExpTrigger(t, 5, false, chkInit.Process(tracker, []request.Event{ + chkInit.Process(tracker, []request.Event{ {Server: testServer3, Type: request.EvRegistered}, - })) + }) // expect bootstrap request to server 3 req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 3}, Request: ReqCheckpointData(checkpointHash)} tracker.ExpRequests(t, 3, []request.RequestWithID{req3}) // valid response to req3; expect chain to be initialized - ExpTrigger(t, 6, true, chkInit.Process(tracker, []request.Event{ + chkInit.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req3, checkpoint), - })) + }) chain.ExpInit(t, true) } @@ -78,12 +78,12 @@ func TestUpdateSyncParallel(t *testing.T) { chain := &TestCommitteeChain{} chain.SetNextSyncPeriod(0) updateSync := NewForwardUpdateSync(chain) - ExpTrigger(t, 1, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*100 + 0x1000}}, {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*100 + 0x1000}}, - })) + }) // expect 6 requests to be sent req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 1}, Request: ReqUpdates{FirstPeriod: 0, Count: 8}} req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 2}, Request: ReqUpdates{FirstPeriod: 8, Count: 8}} @@ -94,38 +94,38 @@ func TestUpdateSyncParallel(t *testing.T) { tracker.ExpRequests(t, 1, []request.RequestWithID{req1, req2, req3, req4, req5, req6}) // valid response to request 1 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 2, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req1, testRespUpdate(req1)), - })) + }) // expect 8 periods synced and a new request started chain.ExpNextSyncPeriod(t, 8) req7 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 7}, Request: ReqUpdates{FirstPeriod: 48, Count: 8}} tracker.ExpRequests(t, 2, []request.RequestWithID{req7}) // valid response to requests 4 and 5 tracker.AddAllowance(testServer2, 2) - ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)), TestReqEvent(request.EvResponse, req5, testRespUpdate(req5)), - })) + }) // expect 2 more requests but no sync progress (responses 4 and 5 cannot be added before 2 and 3) chain.ExpNextSyncPeriod(t, 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}} tracker.ExpRequests(t, 3, []request.RequestWithID{req8, req9}) // soft timeout for requests 2 and 3 (server 1 is overloaded) - ExpTrigger(t, 4, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvTimeout, req2, nil), TestReqEvent(request.EvTimeout, req3, nil), - })) + }) // no allowance, no more requests tracker.ExpRequests(t, 4, nil) // valid response to requests 6 and 8 and 9 tracker.AddAllowance(testServer2, 3) - ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req6, testRespUpdate(req6)), TestReqEvent(request.EvResponse, req8, testRespUpdate(req8)), TestReqEvent(request.EvResponse, req9, testRespUpdate(req9)), - })) + }) // 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}} req3r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 11}, Request: ReqUpdates{FirstPeriod: 16, Count: 8}} @@ -133,19 +133,19 @@ func TestUpdateSyncParallel(t *testing.T) { tracker.ExpRequests(t, 5, []request.RequestWithID{req2r, req3r, req10}) // server 1 finally answers timed out request 2 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req2, testRespUpdate(req2)), - })) + }) // expect sync progress and one new request chain.ExpNextSyncPeriod(t, 16) req11 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 13}, Request: ReqUpdates{FirstPeriod: 80, Count: 8}} tracker.ExpRequests(t, 6, []request.RequestWithID{req11}) // server 2 answers re-sent requests 2 and 3 tracker.AddAllowance(testServer2, 2) - ExpTrigger(t, 7, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req2r, testRespUpdate(req2r)), TestReqEvent(request.EvResponse, req3r, testRespUpdate(req3r)), - })) + }) // finally the gap is filled, update can process responses up to req6 chain.ExpNextSyncPeriod(t, 48) // expect 2 new requests from server 2 (now the available range is covered) @@ -153,14 +153,14 @@ func TestUpdateSyncParallel(t *testing.T) { req13 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 15}, Request: ReqUpdates{FirstPeriod: 96, Count: 4}} tracker.ExpRequests(t, 7, []request.RequestWithID{req12, req13}) // all remaining requests are answered - ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)), TestReqEvent(request.EvResponse, req7, testRespUpdate(req7)), TestReqEvent(request.EvResponse, req10, testRespUpdate(req10)), TestReqEvent(request.EvResponse, req11, testRespUpdate(req11)), TestReqEvent(request.EvResponse, req12, testRespUpdate(req12)), TestReqEvent(request.EvResponse, req13, testRespUpdate(req13)), - })) + }) // expect chain to be fully synced chain.ExpNextSyncPeriod(t, 100) } @@ -174,62 +174,62 @@ func TestUpdateSyncDifferentHeads(t *testing.T) { chain := &TestCommitteeChain{} chain.SetNextSyncPeriod(10) updateSync := NewForwardUpdateSync(chain) - ExpTrigger(t, 1, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*15 + 0x1000}}, {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*16 + 0x1000}}, {Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*17 + 0x1000}}, - })) + }) // expect request to the best announced head req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer3, ID: 1}, Request: ReqUpdates{FirstPeriod: 10, Count: 7}} tracker.ExpRequests(t, 1, []request.RequestWithID{req1}) // request times out, expect request to the next best head - ExpTrigger(t, 2, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvTimeout, req1, nil), - })) + }) req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 2}, Request: ReqUpdates{FirstPeriod: 10, Count: 6}} tracker.ExpRequests(t, 2, []request.RequestWithID{req2}) // request times out, expect request to the last available server - ExpTrigger(t, 3, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvTimeout, req2, nil), - })) + }) req3 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 3}, Request: ReqUpdates{FirstPeriod: 10, Count: 5}} tracker.ExpRequests(t, 3, []request.RequestWithID{req3}) // valid response to request 3, expect chain synced to period 15 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 4, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)), - })) + }) chain.ExpNextSyncPeriod(t, 15) // 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}} - ExpTrigger(t, 5, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req1, testRespUpdate(req1x)), - })) + }) // expect no progress of chain head chain.ExpNextSyncPeriod(t, 15) // valid response to request 2, expect chain synced to period 16 tracker.AddAllowance(testServer2, 1) - ExpTrigger(t, 6, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req2, testRespUpdate(req2)), - })) + }) chain.ExpNextSyncPeriod(t, 16) // a new server is registered with announced head period 17 tracker.AddServer(testServer4, 1) - ExpTrigger(t, 7, false, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ {Server: testServer4, Type: request.EvRegistered}, {Server: testServer4, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*17 + 0x1000}}, - })) + }) // expect request to sync one more period req4 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer4, ID: 4}, Request: ReqUpdates{FirstPeriod: 16, Count: 1}} tracker.ExpRequests(t, 4, []request.RequestWithID{req4}) // valid response, expect chain synced to period 17 tracker.AddAllowance(testServer1, 1) - ExpTrigger(t, 8, true, updateSync.Process(tracker, []request.Event{ + updateSync.Process(tracker, []request.Event{ TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)), - })) + }) chain.ExpNextSyncPeriod(t, 17) } diff --git a/cmd/blsync/block_sync.go b/cmd/blsync/block_sync.go index a840afd077..ffbbf9da2b 100755 --- a/cmd/blsync/block_sync.go +++ b/cmd/blsync/block_sync.go @@ -17,34 +17,24 @@ package main import ( - "fmt" - "math/big" - "sync/atomic" - "github.com/ethereum/go-ethereum/beacon/light/request" "github.com/ethereum/go-ethereum/beacon/light/sync" "github.com/ethereum/go-ethereum/beacon/types" "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/lru" - ctypes "github.com/ethereum/go-ethereum/core/types" - - "github.com/ethereum/go-ethereum/log" - "github.com/ethereum/go-ethereum/rpc" - "github.com/ethereum/go-ethereum/trie" - "github.com/holiman/uint256" "github.com/protolambda/zrnt/eth2/beacon/capella" - "github.com/protolambda/zrnt/eth2/configs" - "github.com/protolambda/ztyp/tree" ) // beaconBlockSync implements request.Module; it fetches the beacon blocks belonging // to the validated and prefetch heads. type beaconBlockSync struct { - recentBlocks *lru.Cache[common.Hash, *capella.BeaconBlock] - validatedHead common.Hash - locked map[common.Hash]struct{} - serverHeads map[request.Server]common.Hash - headTracker headTracker + recentBlocks *lru.Cache[common.Hash, *capella.BeaconBlock] + locked map[common.Hash]struct{} + serverHeads map[request.Server]common.Hash + headTracker headTracker + + lastHeadBlock *capella.BeaconBlock + headBlockCh chan *capella.BeaconBlock } type headTracker interface { @@ -59,15 +49,12 @@ 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), + headBlockCh: make(chan *capella.BeaconBlock, 1), } } // Process implements request.Module -func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Event) (trigger bool) { - if header := s.headTracker.ValidatedHead().Header; header != (types.Header{}) { - s.validatedHead = header.Hash() - } - +func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Event) { // iterate events and add valid responses to recentBlocks for _, event := range events { switch event.Type { @@ -77,9 +64,6 @@ func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Even if resp != nil { block := resp.(*capella.BeaconBlock) s.recentBlocks.Add(blockRoot, block) - if blockRoot == s.validatedHead { - trigger = true - } } delete(s.locked, blockRoot) case sync.EvNewHead: @@ -89,20 +73,23 @@ func (s *beaconBlockSync) Process(tracker request.Tracker, events []request.Even } } - // start new requests if necessary - if s.validatedHead != (common.Hash{}) { - s.tryRequestBlock(tracker, s.validatedHead, false) + // send validated head block or request it if unavailable + if vh := s.headTracker.ValidatedHead(); vh != (types.SignedHeader{}) { + validatedHead := vh.Header.Hash() + if headBlock, ok := s.recentBlocks.Get(validatedHead); ok && headBlock != s.lastHeadBlock { + select { + case s.headBlockCh <- headBlock: + s.lastHeadBlock = headBlock + default: + } + } else { + s.tryRequestBlock(tracker, validatedHead, false) + } } + // request prefetch head if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead != (common.Hash{}) { s.tryRequestBlock(tracker, prefetchHead, true) } - return -} - -// getHeadBlock returns the beacon block belonging to ValidatedHead or nil if not available. -func (s *beaconBlockSync) getHeadBlock() *capella.BeaconBlock { - block, _ := s.recentBlocks.Get(s.validatedHead) - return block } // tryRequestBlock tries to send a block request for the given root if the block @@ -129,109 +116,3 @@ func (s *beaconBlockSync) tryRequestBlock(tracker request.Tracker, blockRoot com s.locked[blockRoot] = struct{}{} } } - -// getExecBlock extracts the execution block from the beacon block's payload. -func getExecBlock(beaconBlock *capella.BeaconBlock) (*ctypes.Block, error) { - payload := &beaconBlock.Body.ExecutionPayload - txs := make([]*ctypes.Transaction, len(payload.Transactions)) - for i, opaqueTx := range payload.Transactions { - var tx ctypes.Transaction - if err := tx.UnmarshalBinary(opaqueTx); err != nil { - return nil, fmt.Errorf("failed to parse tx %d: %v", i, err) - } - txs[i] = &tx - } - withdrawals := make([]*ctypes.Withdrawal, len(payload.Withdrawals)) - for i, w := range payload.Withdrawals { - withdrawals[i] = &ctypes.Withdrawal{ - Index: uint64(w.Index), - Validator: uint64(w.ValidatorIndex), - Address: common.Address(w.Address), - Amount: uint64(w.Amount), - } - } - wroot := ctypes.DeriveSha(ctypes.Withdrawals(withdrawals), trie.NewStackTrie(nil)) - execHeader := &ctypes.Header{ - ParentHash: common.Hash(payload.ParentHash), - UncleHash: ctypes.EmptyUncleHash, - Coinbase: common.Address(payload.FeeRecipient), - Root: common.Hash(payload.StateRoot), - TxHash: ctypes.DeriveSha(ctypes.Transactions(txs), trie.NewStackTrie(nil)), - ReceiptHash: common.Hash(payload.ReceiptsRoot), - Bloom: ctypes.Bloom(payload.LogsBloom), - Difficulty: common.Big0, - Number: new(big.Int).SetUint64(uint64(payload.BlockNumber)), - GasLimit: uint64(payload.GasLimit), - GasUsed: uint64(payload.GasUsed), - Time: uint64(payload.Timestamp), - Extra: []byte(payload.ExtraData), - MixDigest: common.Hash(payload.PrevRandao), // reused in merge - Nonce: ctypes.BlockNonce{}, // zero - BaseFee: (*uint256.Int)(&payload.BaseFeePerGas).ToBig(), - WithdrawalsHash: &wroot, - } - execBlock := ctypes.NewBlockWithHeader(execHeader).WithBody(txs, nil).WithWithdrawals(withdrawals) - if execBlockHash := execBlock.Hash(); execBlockHash != common.Hash(payload.BlockHash) { - return nil, fmt.Errorf("Sanity check failed, payload hash does not match (expected %x, got %x)", common.Hash(payload.BlockHash), execBlockHash) - } - return execBlock, nil -} - -// beaconBlockHash calculates the hash of a beacon block. -func beaconBlockHash(beaconBlock *capella.BeaconBlock) common.Hash { - return common.Hash(beaconBlock.HashTreeRoot(configs.Mainnet, tree.GetHashFn())) -} - -// engineApiUpdater implements request.Module. This module does not start requests, -// it is only implemented as a module in order to easily trigger it by successful -// head block retrieval. -type engineApiUpdater struct { - client *rpc.Client - trigger func() - lastHead common.Hash - blockSync *beaconBlockSync - updating uint32 -} - -// Process implements request.Module -func (s *engineApiUpdater) Process(tracker request.Tracker, events []request.Event) bool { - if atomic.LoadUint32(&s.updating) == 1 { - return false - } - headBlock := s.blockSync.getHeadBlock() - if headBlock == nil { - return false - } - headRoot := beaconBlockHash(headBlock) - if headRoot == s.lastHead { - return false - } - - s.lastHead = headRoot - execBlock, err := getExecBlock(headBlock) - if err != nil { - log.Error("Error extracting execution block from validated beacon block", "error", err) - return false - } - execRoot := execBlock.Hash() - if s.client == nil { // dry run, no engine API specified - log.Info("New execution block retrieved", "block number", execBlock.NumberU64(), "block hash", execRoot) - } else { - atomic.StoreUint32(&s.updating, 1) - go func() { - if status, err := callNewPayloadV2(s.client, execBlock); err == nil { - log.Info("Successful NewPayload", "block number", execBlock.NumberU64(), "block hash", execRoot, "status", status) - } else { - log.Error("Failed NewPayload", "block number", execBlock.NumberU64(), "block hash", execRoot, "error", err) - } - if status, err := callForkchoiceUpdatedV1(s.client, execRoot, common.Hash{}); err == nil { - log.Info("Successful ForkchoiceUpdated", "head", execRoot, "status", status) - } else { - log.Error("Failed ForkchoiceUpdated", "head", execRoot, "error", err) - } - atomic.StoreUint32(&s.updating, 0) - s.trigger() - }() - } - return false -} diff --git a/cmd/blsync/block_sync_test.go b/cmd/blsync/block_sync_test.go index eef3c77a9b..f950ae5921 100644 --- a/cmd/blsync/block_sync_test.go +++ b/cmd/blsync/block_sync_test.go @@ -51,61 +51,61 @@ func TestBlockSync(t *testing.T) { } } - sync.ExpTrigger(t, 1, false, blockSync.Process(tracker, []request.Event{ + blockSync.Process(tracker, []request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer2, Type: request.EvRegistered}, - })) + }) // no block requests expected until head tracker knows about a head tracker.ExpRequests(t, 1, nil) expHeadBlock(1, nil) // set block 1 as prefetch head, announced by server 2 head1 := blockHeadInfo(testBlock1) ht.prefetch = head1 - sync.ExpTrigger(t, 2, false, blockSync.Process(tracker, []request.Event{ + blockSync.Process(tracker, []request.Event{ {Server: testServer2, Type: sync.EvNewHead, Data: head1}, - })) + }) // expect request to server 2 which has announced the head req1 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 1}, Request: sync.ReqBeaconBlock(head1.BlockRoot)} tracker.ExpRequests(t, 2, []request.RequestWithID{req1}) // valid response tracker.AddAllowance(testServer2, 1) - sync.ExpTrigger(t, 3, false, blockSync.Process(tracker, []request.Event{ + blockSync.Process(tracker, []request.Event{ sync.TestReqEvent(request.EvResponse, req1, testBlock1), - })) + }) // head block still not expected as the fetched block is not the validated head yet expHeadBlock(2, nil) // set as validated head, expect no further requests but block 1 set as head block ht.validated.Header = blockHeader(testBlock1) - sync.ExpTrigger(t, 4, false, blockSync.Process(tracker, nil)) + blockSync.Process(tracker, nil) tracker.ExpRequests(t, 3, nil) expHeadBlock(3, testBlock1) // set block 2 as prefetch head, announced by server 1 head2 := blockHeadInfo(testBlock2) ht.prefetch = head2 - sync.ExpTrigger(t, 5, false, blockSync.Process(tracker, []request.Event{ + blockSync.Process(tracker, []request.Event{ {Server: testServer1, Type: sync.EvNewHead, Data: head2}, - })) + }) // expect request to server 1 req2 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer1, ID: 2}, Request: sync.ReqBeaconBlock(head2.BlockRoot)} tracker.ExpRequests(t, 4, []request.RequestWithID{req2}) // req2 fails, no further requests expected because server 2 has not announced it - sync.ExpTrigger(t, 6, false, blockSync.Process(tracker, []request.Event{ + blockSync.Process(tracker, []request.Event{ sync.TestReqEvent(request.EvFail, req2, nil), - })) + }) tracker.ExpRequests(t, 5, nil) // set as validated head before retrieving block; now it's assumed to be available from server 2 too ht.validated.Header = blockHeader(testBlock2) - sync.ExpTrigger(t, 7, false, blockSync.Process(tracker, nil)) + blockSync.Process(tracker, nil) // now head block is unavailable again expHeadBlock(4, nil) // expect req2 retry to server 2 req2r := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 3}, Request: sync.ReqBeaconBlock(head2.BlockRoot)} tracker.ExpRequests(t, 6, []request.RequestWithID{req2r}) // valid response, now head block should be block 2 immediately as it is already validated - sync.ExpTrigger(t, 8, true, blockSync.Process(tracker, []request.Event{ + blockSync.Process(tracker, []request.Event{ sync.TestReqEvent(request.EvResponse, req2r, testBlock2), - })) + }) expHeadBlock(5, testBlock2) } diff --git a/cmd/blsync/engine_api.go b/cmd/blsync/engine_api.go new file mode 100644 index 0000000000..0faae8fa21 --- /dev/null +++ b/cmd/blsync/engine_api.go @@ -0,0 +1,134 @@ +// Copyright 2024 The go-ethereum Authors +// This file is part of the go-ethereum library. +// +// The go-ethereum library is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// The go-ethereum library is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with the go-ethereum library. If not, see . + +package main + +import ( + "context" + "fmt" + "math/big" + "time" + + "github.com/ethereum/go-ethereum/beacon/engine" + "github.com/ethereum/go-ethereum/common" + ctypes "github.com/ethereum/go-ethereum/core/types" + + "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/rpc" + "github.com/ethereum/go-ethereum/trie" + "github.com/holiman/uint256" + "github.com/protolambda/zrnt/eth2/beacon/capella" + "github.com/protolambda/zrnt/eth2/configs" + "github.com/protolambda/ztyp/tree" +) + +func updateEngineApi(client *rpc.Client, headBlockCh chan *capella.BeaconBlock) { + for headBlock := range headBlockCh { + execBlock, err := getExecBlock(headBlock) + if err != nil { + log.Error("Error extracting execution block from validated beacon block", "error", err) + continue + } + execRoot := execBlock.Hash() + if client == nil { // dry run, no engine API specified + log.Info("New execution block retrieved", "block number", execBlock.NumberU64(), "block hash", execRoot) + } else { + if status, err := callNewPayloadV2(client, execBlock); err == nil { + log.Info("Successful NewPayload", "block number", execBlock.NumberU64(), "block hash", execRoot, "status", status) + } else { + log.Error("Failed NewPayload", "block number", execBlock.NumberU64(), "block hash", execRoot, "error", err) + } + if status, err := callForkchoiceUpdatedV1(client, execRoot, common.Hash{}); err == nil { + log.Info("Successful ForkchoiceUpdated", "head", execRoot, "status", status) + } else { + log.Error("Failed ForkchoiceUpdated", "head", execRoot, "error", err) + } + } + } +} + +// getExecBlock extracts the execution block from the beacon block's payload. +func getExecBlock(beaconBlock *capella.BeaconBlock) (*ctypes.Block, error) { + payload := &beaconBlock.Body.ExecutionPayload + txs := make([]*ctypes.Transaction, len(payload.Transactions)) + for i, opaqueTx := range payload.Transactions { + var tx ctypes.Transaction + if err := tx.UnmarshalBinary(opaqueTx); err != nil { + return nil, fmt.Errorf("failed to parse tx %d: %v", i, err) + } + txs[i] = &tx + } + withdrawals := make([]*ctypes.Withdrawal, len(payload.Withdrawals)) + for i, w := range payload.Withdrawals { + withdrawals[i] = &ctypes.Withdrawal{ + Index: uint64(w.Index), + Validator: uint64(w.ValidatorIndex), + Address: common.Address(w.Address), + Amount: uint64(w.Amount), + } + } + wroot := ctypes.DeriveSha(ctypes.Withdrawals(withdrawals), trie.NewStackTrie(nil)) + execHeader := &ctypes.Header{ + ParentHash: common.Hash(payload.ParentHash), + UncleHash: ctypes.EmptyUncleHash, + Coinbase: common.Address(payload.FeeRecipient), + Root: common.Hash(payload.StateRoot), + TxHash: ctypes.DeriveSha(ctypes.Transactions(txs), trie.NewStackTrie(nil)), + ReceiptHash: common.Hash(payload.ReceiptsRoot), + Bloom: ctypes.Bloom(payload.LogsBloom), + Difficulty: common.Big0, + Number: new(big.Int).SetUint64(uint64(payload.BlockNumber)), + GasLimit: uint64(payload.GasLimit), + GasUsed: uint64(payload.GasUsed), + Time: uint64(payload.Timestamp), + Extra: []byte(payload.ExtraData), + MixDigest: common.Hash(payload.PrevRandao), // reused in merge + Nonce: ctypes.BlockNonce{}, // zero + BaseFee: (*uint256.Int)(&payload.BaseFeePerGas).ToBig(), + WithdrawalsHash: &wroot, + } + execBlock := ctypes.NewBlockWithHeader(execHeader).WithBody(txs, nil).WithWithdrawals(withdrawals) + if execBlockHash := execBlock.Hash(); execBlockHash != common.Hash(payload.BlockHash) { + return nil, fmt.Errorf("Sanity check failed, payload hash does not match (expected %x, got %x)", common.Hash(payload.BlockHash), execBlockHash) + } + return execBlock, nil +} + +// beaconBlockHash calculates the hash of a beacon block. +func beaconBlockHash(beaconBlock *capella.BeaconBlock) common.Hash { + return common.Hash(beaconBlock.HashTreeRoot(configs.Mainnet, tree.GetHashFn())) +} + +func callNewPayloadV2(client *rpc.Client, block *ctypes.Block) (string, error) { + var resp engine.PayloadStatusV1 + ctx, cancel := context.WithTimeout(context.Background(), time.Second*5) + err := client.CallContext(ctx, &resp, "engine_newPayloadV2", *engine.BlockToExecutableData(block, nil, nil).ExecutionPayload) + cancel() + return resp.Status, err +} + +func callForkchoiceUpdatedV1(client *rpc.Client, headHash, finalizedHash common.Hash) (string, error) { + var resp engine.ForkChoiceResponse + update := engine.ForkchoiceStateV1{ + HeadBlockHash: headHash, + SafeBlockHash: finalizedHash, + FinalizedBlockHash: finalizedHash, + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second*5) + err := client.CallContext(ctx, &resp, "engine_forkchoiceUpdatedV1", update, nil) + cancel() + return resp.PayloadStatus.Status, err +} diff --git a/cmd/blsync/main.go b/cmd/blsync/main.go index cd4db01f7f..134077e04e 100644 --- a/cmd/blsync/main.go +++ b/cmd/blsync/main.go @@ -17,26 +17,20 @@ package main import ( - "context" "fmt" "io" "os" "strings" - "time" - "github.com/ethereum/go-ethereum/beacon/engine" "github.com/ethereum/go-ethereum/beacon/light" "github.com/ethereum/go-ethereum/beacon/light/api" "github.com/ethereum/go-ethereum/beacon/light/request" "github.com/ethereum/go-ethereum/beacon/light/sync" "github.com/ethereum/go-ethereum/cmd/utils" - "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/mclock" - ctypes "github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/ethdb/memorydb" "github.com/ethereum/go-ethereum/internal/flags" "github.com/ethereum/go-ethereum/log" - "github.com/ethereum/go-ethereum/rpc" "github.com/mattn/go-colorable" "github.com/mattn/go-isatty" "github.com/urfave/cli/v2" @@ -126,16 +120,13 @@ func blsync(ctx *cli.Context) error { checkpointInit := sync.NewCheckpointInit(committeeChain, chainConfig.Checkpoint) forwardSync := sync.NewForwardUpdateSync(committeeChain) beaconBlockSync := newBeaconBlockSync(headTracker) - engineApiUpdater := &engineApiUpdater{ //TODO constructor - client: makeRPCClient(ctx), - blockSync: beaconBlockSync, - } - + scheduler.RegisterTarget(headTracker) + scheduler.RegisterTarget(committeeChain) scheduler.RegisterModule(checkpointInit, "checkpointInit") scheduler.RegisterModule(forwardSync, "forwardSync") scheduler.RegisterModule(headSync, "headSync") scheduler.RegisterModule(beaconBlockSync, "beaconBlockSync") - scheduler.RegisterModule(engineApiUpdater, "engineApiUpdater") + go updateEngineApi(makeRPCClient(ctx), beaconBlockSync.headBlockCh) // start scheduler.Start() // register server(s) @@ -146,26 +137,6 @@ func blsync(ctx *cli.Context) error { // run until stopped <-ctx.Done() scheduler.Stop() + close(beaconBlockSync.headBlockCh) return nil } - -func callNewPayloadV2(client *rpc.Client, block *ctypes.Block) (string, error) { - var resp engine.PayloadStatusV1 - ctx, cancel := context.WithTimeout(context.Background(), time.Second*5) - err := client.CallContext(ctx, &resp, "engine_newPayloadV2", *engine.BlockToExecutableData(block, nil, nil).ExecutionPayload) - cancel() - return resp.Status, err -} - -func callForkchoiceUpdatedV1(client *rpc.Client, headHash, finalizedHash common.Hash) (string, error) { - var resp engine.ForkChoiceResponse - update := engine.ForkchoiceStateV1{ - HeadBlockHash: headHash, - SafeBlockHash: finalizedHash, - FinalizedBlockHash: finalizedHash, - } - ctx, cancel := context.WithTimeout(context.Background(), time.Second*5) - err := client.CallContext(ctx, &resp, "engine_forkchoiceUpdatedV1", update, nil) - cancel() - return resp.PayloadStatus.Status, err -}