mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
workersync: make the list retry backoff per-syncer [flake fix] (#1731)
The startup registered-worker scan read its backoff from a package-level var that tests overwrote. Syncers started by earlier tests outlive them and keep reading that var, so the write races the read and `go test -race` fails, in whichever test happens to do the write. Hold the schedule on the syncer instead, so a test only ever shrinks the backoff of the syncer it owns. Fixes a flake I encountered on one of my other PRs. I didn't see an open issue or fix but this is pretty straightforward.
This commit is contained in:
@@ -98,6 +98,12 @@ type WorkerPoolSyncer struct {
|
||||
workerInformer cache.SharedIndexInformer
|
||||
workerPoolInformer cache.SharedIndexInformer
|
||||
queue workqueue.TypedRateLimitingInterface[workerKey]
|
||||
|
||||
// Exponential backoff schedule for retrying a failed page of the startup
|
||||
// registered-worker scan. Per-syncer rather than package-level so a test
|
||||
// can shrink it without writing state another test's syncer is reading.
|
||||
listBackoff time.Duration
|
||||
listCap time.Duration
|
||||
}
|
||||
|
||||
// NewWorkerPoolSyncer creates a new WorkerPoolSyncer.
|
||||
@@ -107,6 +113,8 @@ func NewWorkerPoolSyncer(client ateapipb.ControlClient, workerInformer, workerPo
|
||||
workerInformer: workerInformer,
|
||||
workerPoolInformer: workerPoolInformer,
|
||||
queue: workqueue.NewTypedRateLimitingQueue(workqueue.DefaultTypedControllerRateLimiter[workerKey]()),
|
||||
listBackoff: defaultListBackoff,
|
||||
listCap: defaultListCap,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -398,12 +406,11 @@ func (s *WorkerPoolSyncer) reconcileDeadWorker(ctx context.Context, key workerKe
|
||||
return err
|
||||
}
|
||||
|
||||
// storedWorkerListBackoff and storedWorkerListCap are the exponential backoff
|
||||
// schedule for retrying a failed page of the startup registered-worker scan.
|
||||
// They are vars so tests can shrink them.
|
||||
var (
|
||||
storedWorkerListBackoff = 500 * time.Millisecond
|
||||
storedWorkerListCap = 30 * time.Second
|
||||
// The default retry backoff schedule for a failed page of the startup
|
||||
// registered-worker scan.
|
||||
const (
|
||||
defaultListBackoff = 500 * time.Millisecond
|
||||
defaultListCap = 30 * time.Second
|
||||
)
|
||||
|
||||
// enqueueRegisteredWorkers enqueues a key for every worker record in the
|
||||
@@ -452,7 +459,7 @@ func (s *WorkerPoolSyncer) enqueueRegisteredWorkers(ctx context.Context) {
|
||||
// resets it.
|
||||
func (s *WorkerPoolSyncer) listWorkersPageWithRetry(ctx context.Context, pageToken string) (*ateapipb.ListWorkersResponse, error) {
|
||||
backoff := wait.Backoff{
|
||||
Duration: storedWorkerListBackoff,
|
||||
Duration: s.listBackoff,
|
||||
Factor: 2.0,
|
||||
Jitter: 0.1,
|
||||
// Steps must be large enough for the ramp (Duration*Factor^n) to reach
|
||||
@@ -460,7 +467,7 @@ func (s *WorkerPoolSyncer) listWorkersPageWithRetry(ctx context.Context, pageTok
|
||||
// With Duration=500ms, Factor=2, the ramp hits Cap=30s at step 6
|
||||
// (0.5,1,2,4,8,16,30,30...).
|
||||
Steps: 6,
|
||||
Cap: storedWorkerListCap,
|
||||
Cap: s.listCap,
|
||||
}
|
||||
for {
|
||||
page, err := s.client.ListWorkers(ctx, &ateapipb.ListWorkersRequest{PageSize: 1000, PageToken: pageToken})
|
||||
|
||||
@@ -450,11 +450,6 @@ func TestSyncer_EnqueueRegisteredWorkers_RetriesTransientListError(t *testing.T)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// Shrink the retry backoff so the test's single retry is fast.
|
||||
prev := storedWorkerListBackoff
|
||||
storedWorkerListBackoff = time.Millisecond
|
||||
defer func() { storedWorkerListBackoff = prev }()
|
||||
|
||||
api := newFakeControl()
|
||||
api.put(registeredWorker("ns-enq-retry", "pool1", "worker-1", testPodUID, "10.0.0.10"))
|
||||
|
||||
@@ -468,6 +463,8 @@ func TestSyncer_EnqueueRegisteredWorkers_RetriesTransientListError(t *testing.T)
|
||||
})
|
||||
|
||||
s, _, _ := setupReconcileTest(t, api)
|
||||
// Shrink the retry backoff so the test's single retry is fast.
|
||||
s.listBackoff = time.Millisecond
|
||||
s.enqueueRegisteredWorkers(ctx)
|
||||
|
||||
if failsLeft != 0 {
|
||||
@@ -481,15 +478,12 @@ func TestSyncer_EnqueueRegisteredWorkers_RetriesTransientListError(t *testing.T)
|
||||
// listWorkersPageWithRetry must stop retrying and return once the context is
|
||||
// cancelled, rather than spinning forever, when the API stays unavailable.
|
||||
func TestSyncer_ListWorkersPageWithRetry_StopsOnContextCancel(t *testing.T) {
|
||||
prev := storedWorkerListBackoff
|
||||
storedWorkerListBackoff = time.Millisecond
|
||||
defer func() { storedWorkerListBackoff = prev }()
|
||||
|
||||
api := newFakeControl()
|
||||
api.setListHook(func(*ateapipb.ListWorkersRequest) error {
|
||||
return status.Error(codes.Unavailable, "still down")
|
||||
})
|
||||
s, _, _ := setupReconcileTest(t, api)
|
||||
s.listBackoff = time.Millisecond
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||
defer cancel()
|
||||
@@ -506,10 +500,6 @@ func TestSyncer_EnqueueRegisteredWorkers_StreamsPagesAndRetriesLatePage(t *testi
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
prev := storedWorkerListBackoff
|
||||
storedWorkerListBackoff = time.Millisecond
|
||||
defer func() { storedWorkerListBackoff = prev }()
|
||||
|
||||
api := newFakeControl()
|
||||
api.listPageSize = 1
|
||||
for i, uid := range []string{
|
||||
@@ -531,6 +521,7 @@ func TestSyncer_EnqueueRegisteredWorkers_StreamsPagesAndRetriesLatePage(t *testi
|
||||
})
|
||||
|
||||
s, _, _ := setupReconcileTest(t, api)
|
||||
s.listBackoff = time.Millisecond
|
||||
s.enqueueRegisteredWorkers(ctx)
|
||||
|
||||
if !failed {
|
||||
|
||||
Reference in New Issue
Block a user