beacon/light: automatic target data trigger

This commit is contained in:
Zsolt Felfoldi 2024-01-11 17:34:41 +01:00 committed by Felix Lange
parent 7dd8190f63
commit c7d884c770
14 changed files with 326 additions and 312 deletions

View file

@ -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 {

View file

@ -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
}

View file

@ -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...))
}
}

View file

@ -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

View file

@ -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)

View file

@ -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 {

View file

@ -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
}

View file

@ -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,

View file

@ -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
}

View file

@ -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)
}

View file

@ -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
}

View file

@ -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)
}

134
cmd/blsync/engine_api.go Normal file
View file

@ -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 <http://www.gnu.org/licenses/>.
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
}

View file

@ -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
}