From dc16372b3906b667c7607e74d762485b66214393 Mon Sep 17 00:00:00 2001 From: Yanujz Date: Sat, 19 Sep 2026 19:38:12 +0200 Subject: [PATCH] Anchor probe ticks to the phase end so stagger actually spreads load The ticker was created at loop start, so ticks stayed anchored to creation and the phase wait (always shorter than one interval) never moved them: same-period services probed in lockstep despite #29. Start the ticker after the phase wait; steady period unchanged, shutdown during phasing still exits via context. Fixes #61 --- internal/scheduler/scheduler.go | 12 +-- internal/scheduler/scheduler_stagger_test.go | 81 ++++++++++++++++++++ 2 files changed, 88 insertions(+), 5 deletions(-) diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index f1a06ba..f92d2fb 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -101,17 +101,19 @@ func (s *Scheduler) loop(ctx context.Context, svc config.Service, sem chan struc } probe() // Stagger the first tick: an immediate probe per service plus aligned - // tickers would fire the whole fleet in lockstep every interval. - // The phase only offsets alignment; the steady period is unchanged. + // tickers would fire the whole fleet in lockstep every interval. The + // ticker starts after the phase wait so ticks anchor to the phase end, + // spreading same-period services across the interval; the steady + // period itself is unchanged. phaseTimer := time.NewTimer(s.phase(svc.Interval.Std())) - defer phaseTimer.Stop() - ticker := time.NewTicker(svc.Interval.Std()) - defer ticker.Stop() select { case <-ctx.Done(): + phaseTimer.Stop() return case <-phaseTimer.C: } + ticker := time.NewTicker(svc.Interval.Std()) + defer ticker.Stop() for { select { case <-ctx.Done(): diff --git a/internal/scheduler/scheduler_stagger_test.go b/internal/scheduler/scheduler_stagger_test.go index c3d85ab..38ce0b5 100644 --- a/internal/scheduler/scheduler_stagger_test.go +++ b/internal/scheduler/scheduler_stagger_test.go @@ -5,6 +5,7 @@ import ( "net/http" "net/http/httptest" "sync" + "sync/atomic" "testing" "time" @@ -95,6 +96,86 @@ func TestSchedulerNoStagger(t *testing.T) { } } +// spreadRecorder captures per-service probe timestamps. +type spreadRecorder struct { + mu sync.Mutex + ts map[string][]time.Duration + t0 time.Time +} + +func (r *spreadRecorder) RecordCheck(_ context.Context, c store.Check) error { + r.mu.Lock() + defer r.mu.Unlock() + r.ts[c.ServiceID] = append(r.ts[c.ServiceID], time.Since(r.t0)) + return nil +} + +func (r *spreadRecorder) Purge(context.Context, int, time.Time) (int64, error) { + return 0, nil +} + +func (r *spreadRecorder) stamps(id string) []time.Duration { + r.mu.Lock() + defer r.mu.Unlock() + return append([]time.Duration(nil), r.ts[id]...) +} + +// TestSchedulerStaggerSeparates runs two same-interval services with +// in-range phases 5ms vs 45ms and asserts their steady ticks stay ~40ms +// apart instead of aligned. (The pre-fix code anchored ticks to loop +// start, yielding a ~0 gap.) +func TestSchedulerStaggerSeparates(t *testing.T) { + target := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {})) + defer target.Close() + + spec := config.StatusSpec{} + if err := spec.UnmarshalJSON([]byte("[200]")); err != nil { + t.Fatal(err) + } + mk := func(id string) config.Service { + return config.Service{ + ID: id, Name: id, URL: target.URL, Method: "GET", + Interval: config.Duration(50 * time.Millisecond), + Timeout: config.Duration(2 * time.Second), + ExpectStatus: spec, + } + } + cfg := &config.Config{Services: []config.Service{mk("a"), mk("b")}} + rec := &spreadRecorder{ts: map[string][]time.Duration{}, t0: time.Now()} + s := New(cfg, rec, staggerObserver{}) + // Loop start order across goroutines is unspecified, so hand out the + // two phases atomically without assuming which service is first: the + // assertion below uses the absolute gap. + var calls atomic.Int64 + s.phase = func(time.Duration) time.Duration { + if calls.Add(1) == 1 { + return 5 * time.Millisecond + } + return 45 * time.Millisecond + } + + ctx, cancel := context.WithCancel(context.Background()) + s.Run(ctx) + time.Sleep(160 * time.Millisecond) + cancel() + s.Stop() + + a, b := rec.stamps("a"), rec.stamps("b") + if len(a) < 2 || len(b) < 2 { + t.Fatalf("too few probes: a=%d b=%d", len(a), len(b)) + } + if gap := absDuration(b[1] - a[1]); gap < 20*time.Millisecond { + t.Errorf("second-tick gap = %v, want ~40ms (ticks must not align)", gap) + } +} + +func absDuration(d time.Duration) time.Duration { + if d < 0 { + return -d + } + return d +} + // TestRandomPhaseBounds pins the production phase source to [0, interval). func TestRandomPhaseBounds(t *testing.T) { iv := 60 * time.Second