diff --git a/beacon/light/sync/head_sync_test.go b/beacon/light/sync/head_sync_test.go index ee0d8563b5..48d6d85c72 100644 --- a/beacon/light/sync/head_sync_test.go +++ b/beacon/light/sync/head_sync_test.go @@ -44,97 +44,95 @@ var ( ) func TestValidatedHead(t *testing.T) { - tracker := &TestTracker{} chain := &TestCommitteeChain{} ht := &TestHeadTracker{} headSync := NewHeadSync(ht, chain) ht.ExpValidated(t, 1, nil) - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: testSHead1}, }) ht.ExpValidated(t, 2, nil) chain.SetNextSyncPeriod(0) - headSync.Process(tracker, nil) + headSync.Process(nil) ht.ExpValidated(t, 3, []types.SignedHeader{testSHead1}) chain.SetNextSyncPeriod(1) - headSync.Process(tracker, []request.Event{ + headSync.Process([]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}) - headSync.Process(tracker, []request.Event{ + headSync.Process([]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) - headSync.Process(tracker, nil) + headSync.Process(nil) ht.ExpValidated(t, 6, []types.SignedHeader{testSHead3}) - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer3, Type: request.EvUnregistered}, }) ht.ExpValidated(t, 7, nil) chain.SetNextSyncPeriod(3) - headSync.Process(tracker, nil) + headSync.Process(nil) ht.ExpValidated(t, 8, nil) - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer2, Type: EvNewSignedHead, Data: testSHead4}, }) ht.ExpValidated(t, 9, []types.SignedHeader{testSHead4}) } func TestPrefetchHead(t *testing.T) { - tracker := &TestTracker{} chain := &TestCommitteeChain{} ht := &TestHeadTracker{} headSync := NewHeadSync(ht, chain) ht.ExpPrefetch(t, 1, testHead0) // no servers registered - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewHead, Data: testHead1}, }) ht.ExpPrefetch(t, 2, testHead1) // s1: h1 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer2, Type: request.EvRegistered}, {Server: testServer2, Type: EvNewHead, Data: testHead2}, }) ht.ExpPrefetch(t, 3, testHead2) // s1: h1, s2: h2 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer1, Type: EvNewHead, Data: testHead2}, }) ht.ExpPrefetch(t, 4, testHead2) // s1: h2, s2: h2 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer3, Type: request.EvRegistered}, {Server: testServer3, Type: EvNewHead, Data: testHead3}, }) ht.ExpPrefetch(t, 5, testHead2) // s1: h2, s2: h2, s3: h3 - headSync.Process(tracker, []request.Event{ + headSync.Process([]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 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer2, Type: EvNewHead, Data: testHead3}, }) ht.ExpPrefetch(t, 7, testHead3) // s1: h2, s2: h3, s3: h3, s4: h4 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer3, Type: request.EvUnregistered}, }) ht.ExpPrefetch(t, 8, testHead4) // s1: h2, s2: h3, s4: h4 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer1, Type: request.EvUnregistered}, }) ht.ExpPrefetch(t, 9, testHead4) // s2: h3, s4: h4 - headSync.Process(tracker, []request.Event{ + headSync.Process([]request.Event{ {Server: testServer4, Type: request.EvUnregistered}, }) ht.ExpPrefetch(t, 10, testHead3) // s2: h3 - headSync.Process(tracker, []request.Event{ + headSync.Process([]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 4503552a08..c7899f8a9f 100644 --- a/beacon/light/sync/test_helpers.go +++ b/beacon/light/sync/test_helpers.go @@ -24,33 +24,94 @@ import ( "github.com/ethereum/go-ethereum/beacon/types" ) -type TestTracker struct { +type TestServer struct{ ID int } + +func (ts *TestServer) Fail(desc string) {} + +type TestScheduler struct { + t *testing.T + module request.Module + events []request.Event servers []request.Server allowance map[request.Server]int - sent []request.RequestWithID + sent map[int]request.RequestWithID lastId request.ID } -func (tt *TestTracker) AddServer(server request.Server, allowance int) { - tt.servers = append(tt.servers, server) - if tt.allowance == nil { - tt.allowance = make(map[request.Server]int) +func NewTestScheduler(t *testing.T, module request.Module) *TestScheduler { + return &TestScheduler{ + t: t, + module: module, + allowance: make(map[request.Server]int), + sent: make(map[int]request.RequestWithID), } - tt.allowance[server] = allowance } -func (tt *TestTracker) AddAllowance(server request.Server, allowance int) { - tt.allowance[server] += allowance +func (ts *TestScheduler) Run(testIndex int, expServer request.Server, expReq request.Request) { + ts.module.Process(ts.events) + ts.events = nil + expReqWithID := request.RequestWithID{ + ServerAndID: request.ServerAndID{Server: expServer, ID: ts.lastId + 1}, + Request: expReq, + } + req, ok := ts.tryRequest(testIndex, ts.module.MakeRequest) + if expReq == nil { + if ok { + ts.t.Errorf("Unexpected request in test case #%d (expected none, got %v)", testIndex, req) + } + return + } + if !ok { + ts.t.Errorf("Missing request in test case #%d (expected none, got %v)", testIndex, expReqWithID) + return + } + if req != expReqWithID { + ts.t.Errorf("Wrong request in test case #%d (expected %v, got %v)", testIndex, req, expReqWithID) + } } -func (tt *TestTracker) TryRequest(requestFn func(server request.Server) (request.Request, float32)) (request.RequestWithID, bool) { +func (ts *TestScheduler) ServerEvent(evType *request.EventType, server request.Server, data any) { + ts.events = append(ts.events, request.Event{ + Type: evType, + Server: server, + Data: data, + }) +} + +func (ts *TestScheduler) RequestEvent(evType *request.EventType, testIndex int, resp request.Response) { + req, ok := ts.sent[testIndex] + if !ok { + ts.t.Errorf("Missing request from test case %v", testIndex) + return + } + ts.events = append(ts.events, request.Event{ + Type: evType, + Server: req.ServerAndID.Server, + Data: request.RequestResponse{ + ID: req.ServerAndID.ID, + Request: req.Request, + Response: resp, + }, + }) +} + +func (ts *TestScheduler) AddServer(server request.Server, allowance int) { + ts.servers = append(ts.servers, server) + ts.allowance[server] = allowance +} + +func (ts *TestScheduler) AddAllowance(server request.Server, allowance int) { + ts.allowance[server] += allowance +} + +func (ts *TestScheduler) tryRequest(testIndex int, requestFn func(server request.Server) (request.Request, float32)) (request.RequestWithID, bool) { var ( bestServer request.Server bestReq request.Request bestPri float32 ) - for _, server := range tt.servers { - if tt.allowance[server] == 0 { + for _, server := range ts.servers { + if ts.allowance[server] == 0 { continue } req, pri := requestFn(server) @@ -61,48 +122,16 @@ func (tt *TestTracker) TryRequest(requestFn func(server request.Server) (request if bestServer == nil { return request.RequestWithID{}, false } - tt.allowance[bestServer]-- - tt.lastId++ + ts.allowance[bestServer]-- + ts.lastId++ req := request.RequestWithID{ - ServerAndID: request.ServerAndID{Server: bestServer, ID: tt.lastId}, + ServerAndID: request.ServerAndID{Server: bestServer, ID: ts.lastId}, Request: bestReq, } - tt.sent = append(tt.sent, req) + ts.sent[testIndex] = req return req, true } -func (tt *TestTracker) ExpRequests(t *testing.T, tci int, expSent []request.RequestWithID) { - for i, expReq := range expSent { - if i >= len(tt.sent) { - t.Errorf("Missing sent request in test case #%d index #%d (expected %v, got none)", tci, i, expReq) - continue - } - if tt.sent[i] != expReq { - t.Errorf("Wrong sent request in test case #%d index #%d (expected %v, got %v)", tci, i, expReq, tt.sent[i]) - } - } - for i := len(expSent); i < len(tt.sent); i++ { - t.Errorf("Unexpected sent request in test case #%d index #%d (expected none, got %v)", tci, i, tt.sent[i]) - } - tt.sent = nil -} - -func (tt *TestTracker) InvalidResponse(id request.ServerAndID, desc string) { - return -} - -func TestReqEvent(evType *request.EventType, req request.RequestWithID, response request.Response) request.Event { - return request.Event{ - Type: evType, - Server: req.ServerAndID.Server, - Data: request.RequestResponse{ - ID: req.ServerAndID.ID, - Request: req.Request, - Response: response, - }, - } -} - type TestCommitteeChain struct { fsp, nsp uint64 init bool diff --git a/beacon/light/sync/update_sync_test.go b/beacon/light/sync/update_sync_test.go index 0f1b9356e6..f17aebf712 100644 --- a/beacon/light/sync/update_sync_test.go +++ b/beacon/light/sync/update_sync_test.go @@ -24,61 +24,63 @@ import ( ) func TestCheckpointInit(t *testing.T) { - tracker := &TestTracker{} - // add 2 servers - tracker.AddServer(testServer1, 1) - tracker.AddServer(testServer2, 1) chain := &TestCommitteeChain{} checkpoint := &types.BootstrapData{Header: types.Header{Slot: 0x2000*4 + 0x1000}} // period 4 checkpointHash := checkpoint.Header.Hash() chkInit := NewCheckpointInit(chain, checkpointHash) - chkInit.Process(tracker, []request.Event{ + ts := NewTestScheduler(t, chkInit) + // add 2 servers + ts.AddServer(testServer1, 1) + ts.AddServer(testServer2, 1) + + chkInit.Process([]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}) + ts.ExpRequests(t, 1, []request.RequestWithID{req1}) // req1 times out; expect request to server 2 - chkInit.Process(tracker, []request.Event{ + chkInit.Process([]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}) + ts.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}} - chkInit.Process(tracker, []request.Event{ + chkInit.Process([]request.Event{ TestReqEvent(request.EvResponse, req2, wrongCheckpoint), }) // req1 fails (hard timeout) - chkInit.Process(tracker, []request.Event{ + chkInit.Process([]request.Event{ TestReqEvent(request.EvFail, req1, nil), }) chain.ExpInit(t, false) // server 3 is registered - tracker.AddServer(testServer3, 1) - chkInit.Process(tracker, []request.Event{ + ts.AddServer(testServer3, 1) + chkInit.Process([]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}) + ts.ExpRequests(t, 3, []request.RequestWithID{req3}) // valid response to req3; expect chain to be initialized - chkInit.Process(tracker, []request.Event{ + chkInit.Process([]request.Event{ TestReqEvent(request.EvResponse, req3, checkpoint), }) chain.ExpInit(t, true) } func TestUpdateSyncParallel(t *testing.T) { - tracker := &TestTracker{} - // add 2 servers, head at period 100; allow 3-3 parallel requests for each - tracker.AddServer(testServer1, 3) - tracker.AddServer(testServer2, 3) chain := &TestCommitteeChain{} chain.SetNextSyncPeriod(0) updateSync := NewForwardUpdateSync(chain) - updateSync.Process(tracker, []request.Event{ + ts := NewTestScheduler(t, updateSync) + // add 2 servers, head at period 100; allow 3-3 parallel requests for each + ts.AddServer(testServer1, 3) + ts.AddServer(testServer2, 3) + + updateSync.Process([]request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*100 + 0x1000}}, {Server: testServer2, Type: request.EvRegistered}, @@ -91,19 +93,19 @@ func TestUpdateSyncParallel(t *testing.T) { req4 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 4}, Request: ReqUpdates{FirstPeriod: 24, Count: 8}} req5 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 5}, Request: ReqUpdates{FirstPeriod: 32, Count: 8}} req6 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 6}, Request: ReqUpdates{FirstPeriod: 40, Count: 8}} - tracker.ExpRequests(t, 1, []request.RequestWithID{req1, req2, req3, req4, req5, req6}) + ts.ExpRequests(t, 1, []request.RequestWithID{req1, req2, req3, req4, req5, req6}) // valid response to request 1 - tracker.AddAllowance(testServer1, 1) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer1, 1) + updateSync.Process([]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}) + ts.ExpRequests(t, 2, []request.RequestWithID{req7}) // valid response to requests 4 and 5 - tracker.AddAllowance(testServer2, 2) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer2, 2) + updateSync.Process([]request.Event{ TestReqEvent(request.EvResponse, req4, testRespUpdate(req4)), TestReqEvent(request.EvResponse, req5, testRespUpdate(req5)), }) @@ -111,17 +113,17 @@ func TestUpdateSyncParallel(t *testing.T) { 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}) + ts.ExpRequests(t, 3, []request.RequestWithID{req8, req9}) // soft timeout for requests 2 and 3 (server 1 is overloaded) - updateSync.Process(tracker, []request.Event{ + updateSync.Process([]request.Event{ TestReqEvent(request.EvTimeout, req2, nil), TestReqEvent(request.EvTimeout, req3, nil), }) // no allowance, no more requests - tracker.ExpRequests(t, 4, nil) + ts.ExpRequests(t, 4, nil) // valid response to requests 6 and 8 and 9 - tracker.AddAllowance(testServer2, 3) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer2, 3) + updateSync.Process([]request.Event{ TestReqEvent(request.EvResponse, req6, testRespUpdate(req6)), TestReqEvent(request.EvResponse, req8, testRespUpdate(req8)), TestReqEvent(request.EvResponse, req9, testRespUpdate(req9)), @@ -130,19 +132,19 @@ func TestUpdateSyncParallel(t *testing.T) { 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}} req10 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 12}, Request: ReqUpdates{FirstPeriod: 72, Count: 8}} - tracker.ExpRequests(t, 5, []request.RequestWithID{req2r, req3r, req10}) + ts.ExpRequests(t, 5, []request.RequestWithID{req2r, req3r, req10}) // server 1 finally answers timed out request 2 - tracker.AddAllowance(testServer1, 1) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer1, 1) + updateSync.Process([]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}) + ts.ExpRequests(t, 6, []request.RequestWithID{req11}) // server 2 answers re-sent requests 2 and 3 - tracker.AddAllowance(testServer2, 2) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer2, 2) + updateSync.Process([]request.Event{ TestReqEvent(request.EvResponse, req2r, testRespUpdate(req2r)), TestReqEvent(request.EvResponse, req3r, testRespUpdate(req3r)), }) @@ -151,9 +153,9 @@ func TestUpdateSyncParallel(t *testing.T) { // expect 2 new requests from server 2 (now the available range is covered) req12 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 14}, Request: ReqUpdates{FirstPeriod: 88, Count: 8}} req13 := request.RequestWithID{ServerAndID: request.ServerAndID{Server: testServer2, ID: 15}, Request: ReqUpdates{FirstPeriod: 96, Count: 4}} - tracker.ExpRequests(t, 7, []request.RequestWithID{req12, req13}) + ts.ExpRequests(t, 7, []request.RequestWithID{req12, req13}) // all remaining requests are answered - updateSync.Process(tracker, []request.Event{ + updateSync.Process([]request.Event{ TestReqEvent(request.EvResponse, req3, testRespUpdate(req3)), TestReqEvent(request.EvResponse, req7, testRespUpdate(req7)), TestReqEvent(request.EvResponse, req10, testRespUpdate(req10)), @@ -166,15 +168,16 @@ func TestUpdateSyncParallel(t *testing.T) { } func TestUpdateSyncDifferentHeads(t *testing.T) { - tracker := &TestTracker{} - // add 3 servers with different announced head periods - tracker.AddServer(testServer1, 1) - tracker.AddServer(testServer2, 1) - tracker.AddServer(testServer3, 1) chain := &TestCommitteeChain{} chain.SetNextSyncPeriod(10) updateSync := NewForwardUpdateSync(chain) - updateSync.Process(tracker, []request.Event{ + ts := NewTestScheduler(t, updateSync) + // add 3 servers with different announced head periods + ts.AddServer(testServer1, 1) + ts.AddServer(testServer2, 1) + ts.AddServer(testServer3, 1) + + updateSync.Process([]request.Event{ {Server: testServer1, Type: request.EvRegistered}, {Server: testServer1, Type: EvNewSignedHead, Data: types.SignedHeader{SignatureSlot: 0x2000*15 + 0x1000}}, {Server: testServer2, Type: request.EvRegistered}, @@ -184,50 +187,50 @@ func TestUpdateSyncDifferentHeads(t *testing.T) { }) // 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}) + ts.ExpRequests(t, 1, []request.RequestWithID{req1}) // request times out, expect request to the next best head - updateSync.Process(tracker, []request.Event{ + updateSync.Process([]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}) + ts.ExpRequests(t, 2, []request.RequestWithID{req2}) // request times out, expect request to the last available server - updateSync.Process(tracker, []request.Event{ + updateSync.Process([]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}) + ts.ExpRequests(t, 3, []request.RequestWithID{req3}) // valid response to request 3, expect chain synced to period 15 - tracker.AddAllowance(testServer1, 1) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer1, 1) + updateSync.Process([]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}} - updateSync.Process(tracker, []request.Event{ + updateSync.Process([]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) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer2, 1) + updateSync.Process([]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) - updateSync.Process(tracker, []request.Event{ + ts.AddServer(testServer4, 1) + updateSync.Process([]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}) + ts.ExpRequests(t, 4, []request.RequestWithID{req4}) // valid response, expect chain synced to period 17 - tracker.AddAllowance(testServer1, 1) - updateSync.Process(tracker, []request.Event{ + ts.AddAllowance(testServer1, 1) + updateSync.Process([]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 b70649715e..924fbc5962 100755 --- a/cmd/blsync/block_sync.go +++ b/cmd/blsync/block_sync.go @@ -99,7 +99,7 @@ func (s *beaconBlockSync) MakeRequest(server request.Server) (request.Request, f } } // request prefetch head if the given server has announced it - if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead == s.serverHeads[server] { + if prefetchHead := s.headTracker.PrefetchHead().BlockRoot; prefetchHead != (common.Hash{}) && prefetchHead == s.serverHeads[server] { if _, ok := s.recentBlocks.Get(prefetchHead); !ok { if _, ok := s.locked[prefetchHead]; !ok { return sync.ReqBeaconBlock(prefetchHead), 0 diff --git a/cmd/blsync/block_sync_test.go b/cmd/blsync/block_sync_test.go index f950ae5921..e82c379c37 100644 --- a/cmd/blsync/block_sync_test.go +++ b/cmd/blsync/block_sync_test.go @@ -29,83 +29,79 @@ import ( ) var ( - testServer1 = 1 - testServer2 = 2 + testServer1 = &sync.TestServer{ID: 1} + testServer2 = &sync.TestServer{ID: 2} testBlock1 = &capella.BeaconBlock{Slot: 123} testBlock2 = &capella.BeaconBlock{Slot: 124} ) func TestBlockSync(t *testing.T) { - tracker := &sync.TestTracker{} - tracker.AddServer(testServer1, 1) - tracker.AddServer(testServer2, 1) ht := &testHeadTracker{} blockSync := newBeaconBlockSync(ht) + ts := sync.NewTestScheduler(t, blockSync) + ts.AddServer(testServer1, 1) + ts.AddServer(testServer2, 1) expHeadBlock := func(tci int, expHead *capella.BeaconBlock) { expInfo := blockHeadInfo(expHead) - headInfo := blockHeadInfo(blockSync.getHeadBlock()) + var head *capella.BeaconBlock + select { + case head = <-blockSync.headBlockCh: + default: + } + headInfo := blockHeadInfo(head) if headInfo != expInfo { t.Errorf("Wrong head block in test case #%d (expected {slot %d blockRoot %x}, got {slot %d blockRoot %x})", tci, expInfo.Slot, expInfo.BlockRoot, headInfo.Slot, headInfo.BlockRoot) } } - blockSync.Process(tracker, []request.Event{ - {Server: testServer1, Type: request.EvRegistered}, - {Server: testServer2, Type: request.EvRegistered}, - }) + ts.ServerEvent(request.EvRegistered, testServer1, nil) + ts.ServerEvent(request.EvRegistered, testServer2, nil) // no block requests expected until head tracker knows about a head - tracker.ExpRequests(t, 1, nil) + ts.Run(1, nil, nil) expHeadBlock(1, nil) + // set block 1 as prefetch head, announced by server 2 head1 := blockHeadInfo(testBlock1) ht.prefetch = head1 - blockSync.Process(tracker, []request.Event{ - {Server: testServer2, Type: sync.EvNewHead, Data: head1}, - }) + ts.ServerEvent(sync.EvNewHead, testServer2, 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}) + ts.Run(2, testServer2, sync.ReqBeaconBlock(head1.BlockRoot)) + // valid response - tracker.AddAllowance(testServer2, 1) - blockSync.Process(tracker, []request.Event{ - sync.TestReqEvent(request.EvResponse, req1, testBlock1), - }) + ts.RequestEvent(request.EvResponse, 2, testBlock1) + ts.AddAllowance(testServer2, 1) + ts.Run(3, nil, nil) // head block still not expected as the fetched block is not the validated head yet - expHeadBlock(2, nil) + expHeadBlock(3, nil) + // set as validated head, expect no further requests but block 1 set as head block ht.validated.Header = blockHeader(testBlock1) - blockSync.Process(tracker, nil) - tracker.ExpRequests(t, 3, nil) - expHeadBlock(3, testBlock1) + ts.Run(4, nil, nil) + expHeadBlock(4, testBlock1) // set block 2 as prefetch head, announced by server 1 head2 := blockHeadInfo(testBlock2) ht.prefetch = head2 - blockSync.Process(tracker, []request.Event{ - {Server: testServer1, Type: sync.EvNewHead, Data: head2}, - }) + ts.ServerEvent(sync.EvNewHead, testServer1, 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}) + ts.Run(5, testServer1, sync.ReqBeaconBlock(head2.BlockRoot)) + // req2 fails, no further requests expected because server 2 has not announced it - blockSync.Process(tracker, []request.Event{ - sync.TestReqEvent(request.EvFail, req2, nil), - }) - tracker.ExpRequests(t, 5, nil) + ts.RequestEvent(request.EvFail, 5, nil) + ts.Run(6, nil, nil) + // set as validated head before retrieving block; now it's assumed to be available from server 2 too ht.validated.Header = blockHeader(testBlock2) - 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}) + ts.Run(7, testServer2, sync.ReqBeaconBlock(head2.BlockRoot)) + // now head block should be unavailable again + expHeadBlock(4, nil) + // valid response, now head block should be block 2 immediately as it is already validated - blockSync.Process(tracker, []request.Event{ - sync.TestReqEvent(request.EvResponse, req2r, testBlock2), - }) + ts.RequestEvent(request.EvResponse, 7, testBlock2) + ts.Run(8, nil, nil) expHeadBlock(5, testBlock2) }