From b098963ed1f147e0071f0e1af5b11ed6bc19c526 Mon Sep 17 00:00:00 2001 From: Victor Lyuboslavsky <2685025+getvictor@users.noreply.github.com> Date: Tue, 10 Feb 2026 20:26:43 -0600 Subject: [PATCH] Reworked how we handle server/worker delays to fix flaky tests (#39609) --- changes/39608-server-worker | 1 + server/worker/worker.go | 24 ++++++++++++++++-------- server/worker/worker_test.go | 9 +++------ 3 files changed, 20 insertions(+), 14 deletions(-) create mode 100644 changes/39608-server-worker diff --git a/changes/39608-server-worker b/changes/39608-server-worker new file mode 100644 index 0000000000..01b213670d --- /dev/null +++ b/changes/39608-server-worker @@ -0,0 +1 @@ +Reworked how we handle server/worker delays to fix flaky tests. diff --git a/server/worker/worker.go b/server/worker/worker.go index 67f7ec0169..e32961834d 100644 --- a/server/worker/worker.go +++ b/server/worker/worker.go @@ -67,6 +67,10 @@ type Worker struct { // For tests only, allows ignoring unknown jobs instead of failing them. TestIgnoreUnknownJobs bool + // delayPerRetry defines the delays between retries. If nil, the default + // delays are used. + delayPerRetry []time.Duration + registry map[string]Job } @@ -117,12 +121,12 @@ func QueueJobWithDelay(ctx context.Context, ds fleet.Datastore, name string, arg return ds.NewJob(ctx, job) } -// this defines the delays to add between retries (i.e. how the "not_before" -// timestamp of a job will be set for the next run). Keep in mind that at a -// minimum, the job will not be retried before the next cron run of the worker, -// but we want to ensure a minimum delay before retries to give a chance to -// e.g. transient network issues to resolve themselves. -var delayPerRetry = []time.Duration{ +// defaultDelayPerRetry defines the delays to add between retries (i.e. how +// the "not_before" timestamp of a job will be set for the next run). Keep in +// mind that at a minimum, the job will not be retried before the next cron run +// of the worker, but we want to ensure a minimum delay before retries to give +// a chance to e.g. transient network issues to resolve themselves. +var defaultDelayPerRetry = []time.Duration{ 1: 0, // i.e. for the first retry, do it ASAP (on the next worker run) 2: 5 * time.Minute, 3: 10 * time.Minute, @@ -184,8 +188,12 @@ func (w *Worker) ProcessJobs(ctx context.Context) error { if job.Retries < maxRetries { level.Debug(log).Log("msg", "will retry job") job.Retries += 1 - if job.Retries < len(delayPerRetry) { - job.NotBefore = time.Now().Add(delayPerRetry[job.Retries]) + delays := w.delayPerRetry + if delays == nil { + delays = defaultDelayPerRetry + } + if job.Retries < len(delays) { + job.NotBefore = time.Now().UTC().Add(delays[job.Retries]) } } else { job.State = fleet.JobStateFailure diff --git a/server/worker/worker_test.go b/server/worker/worker_test.go index 64d60b5344..81cf59e151 100644 --- a/server/worker/worker_test.go +++ b/server/worker/worker_test.go @@ -245,16 +245,13 @@ func TestWorkerWithRealDatastore(t *testing.T) { // call TruncateTables immediately, because a DB migration may create jobs mysql.TruncateTables(t, ds) - oldDelayPerRetry := delayPerRetry - delayPerRetry = []time.Duration{ + logger := kitlog.NewNopLogger() + w := NewWorker(ds, logger) + w.delayPerRetry = []time.Duration{ 1: 0, 2: 0, 3: time.Hour, } // retry twice on the next cron, then not before an hour - t.Cleanup(func() { delayPerRetry = oldDelayPerRetry }) - - logger := kitlog.NewNopLogger() - w := NewWorker(ds, logger) // register a test job var jobCallCount int