Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 7 additions & 5 deletions internal/scheduler/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand Down
81 changes: 81 additions & 0 deletions internal/scheduler/scheduler_stagger_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"net/http"
"net/http/httptest"
"sync"
"sync/atomic"
"testing"
"time"

Expand Down Expand Up @@ -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
Expand Down
Loading