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