diff --git a/server/internal/scheduler/scheduler.go b/server/internal/scheduler/scheduler.go index 993e952..158d21c 100644 --- a/server/internal/scheduler/scheduler.go +++ b/server/internal/scheduler/scheduler.go @@ -249,12 +249,19 @@ func (s *Scheduler) runLoop(ctx context.Context, name string, interval time.Dura for { // Try to acquire the leader lease (SET NX PX ). ok, err := s.rdb.SetNX(ctx, key, s.cfg.InstanceID, s.cfg.LeaseTTL).Result() - if ctx.Err() != nil { - return // context cancelled — don't log spurious errors - } - if err != nil { + switch { + case err != nil: + // On a cancelled context SetNX errors without acquiring anything, + // so it is safe to exit silently. Any other error is logged and retried. + if ctx.Err() != nil { + return + } slog.Error("scheduler: leader setnx error", "loop", name, "error", err) - } else if ok { + case ok: + // Lease acquired. Always enter runLeaderLoop — even if ctx was + // cancelled in the race immediately after SetNX — because its defer + // is what releases the key. Returning here on ctx.Err() would leak + // the just-acquired leader key (no clean DEL on shutdown). slog.Info("scheduler: acquired leader lease", "loop", name, "instance", s.cfg.InstanceID) s.runLeaderLoop(ctx, name, interval, tick, key) slog.Info("scheduler: relinquished leader", "loop", name, "instance", s.cfg.InstanceID) diff --git a/server/internal/scheduler/scheduler_test.go b/server/internal/scheduler/scheduler_test.go index 2bd5f37..fd97eba 100644 --- a/server/internal/scheduler/scheduler_test.go +++ b/server/internal/scheduler/scheduler_test.go @@ -185,14 +185,6 @@ func TestLeaderElection(t *testing.T) { // ───────────────────────────────────────────────────────────────────────────── func TestFollowerTakeover(t *testing.T) { - // FLAKY under load (passes ~5/6 in isolation, flakes more under full-suite - // CPU contention even with a 3s deadline): follower takeover is correct but - // the leader/renew/retry loops can be starved long enough that takeover is - // not observed in time. Same root area as TestGracefulShutdown. Skipped to - // keep CI green; TODO(#15H): make the leader loops less scheduling-sensitive - // (or inject a clock) and re-enable. - t.Skip("flaky: leader-election loops starve under load — tracked for #15H follow-up") - rdb, mr := newTestRedis(t) leaseTTL, renew, retry, tickTout, interval := testIntervals() @@ -235,9 +227,10 @@ func TestFollowerTakeover(t *testing.T) { wgLeader.Add(1) go func() { defer wgLeader.Done(); _ = schedLeader.Start(ctxLeader) }() - go func() { _ = schedFollower.Start(ctxFollower) }() - - // Wait for leader to tick at least once (generous under load). + // Start ONLY the leader first and wait until it actually ticks, so it is the + // confirmed owner of the detect lease before the follower joins. Starting both + // instances at once races the leader election — either could win the lease — + // which previously flaked this test as "leader never ticked". for i := 0; i < 150; i++ { if atomic.LoadInt64(&ticksLeader) > 0 { break @@ -248,6 +241,9 @@ func TestFollowerTakeover(t *testing.T) { t.Fatal("leader never ticked") } + // Now start the follower; it blocks retrying to acquire the held leases. + go func() { _ = schedFollower.Start(ctxFollower) }() + prevFollower := atomic.LoadInt64(&ticksFollower) // Stop leader — its defer releases the leader key immediately. @@ -259,9 +255,9 @@ func TestFollowerTakeover(t *testing.T) { mr.FastForward(leaseTTL + 10*time.Millisecond) // Follower should take over shortly after the lease frees (≈ RetryPeriod + - // one loop interval ≈ 130ms). Deadline widened to 3s so goroutine-scheduling - // jitter under full-suite CPU load doesn't flake this (the assertion still - // fails if takeover genuinely never happens). + // one loop interval ≈ 130ms). 3s deadline leaves generous headroom for + // goroutine-scheduling jitter under full-suite CPU load; the assertion still + // fails if takeover genuinely never happens. deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { if atomic.LoadInt64(&ticksFollower) > prevFollower { @@ -287,13 +283,6 @@ func (e *counterEngine) Tick(_ context.Context) error { atomic.AddInt64(e.n, 1); // ───────────────────────────────────────────────────────────────────────────── func TestGracefulShutdown(t *testing.T) { - // FLAKY (~2/3 fail under full-package load, passes in isolation): the leader - // renew goroutine can SET a leader key immediately after the shutdown path - // DELs it, resurrecting it — a release/renew race in scheduler shutdown that - // predates this integration merge. Skipped to keep CI green; needs a real fix - // (serialize renew vs release on ctx cancel). TODO(#15H): fix the race, re-enable. - t.Skip("flaky: leader release/renew shutdown race — tracked for #15H follow-up") - rdb, _ := newTestRedis(t) leaseTTL, renew, retry, tickTout, interval := testIntervals()