From a7ae7156f40f2783ce9445d007065c1f166f90fb Mon Sep 17 00:00:00 2001 From: wangjia <809946525@qq.com> Date: Wed, 17 Jun 2026 07:06:48 +0800 Subject: [PATCH] =?UTF-8?q?fix(scheduler):=20=E4=BF=AE=E5=A4=8D=20#15H=20?= =?UTF-8?q?=E4=B8=A4=E4=B8=AA=20flaky=20leader-election=20=E6=B5=8B?= =?UTF-8?q?=E8=AF=95=E7=9A=84=E6=A0=B9=E5=9B=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit TestGracefulShutdown(真 bug):runLoop 在 SetNX 获取 leader 锁后、进入 runLeaderLoop(其 defer 负责释放)之前,若 ctx 恰好在此刻 cancel, `if ctx.Err() != nil { return }` 会丢弃已获取的 key 而不释放 → 关停后 leader 键残留(DEL 从未执行)。早 cancel 时三个 loop 各自命中此竞态, 故残留的键不固定。修复:将 ctx.Err() 早退限定在 SetNX 出错(未获取)的分支; 一旦获取成功就必定进入 runLeaderLoop,由其 defer 保证释放。 TestFollowerTakeover(测试设计竞态):测试同时启动 leader/follower 两实例却 假定名为 "leader" 的实例赢得选举——而选举是先到先得,"follower" 可能先抢到 detect 锁,导致 engLeader 永不 tick("leader never ticked")。修复:先单独 启动 leader 并等其 tick(确认占锁),再启动 follower,消除选举非确定性。 两修复均移除原 t.Skip,恢复测试。验证:各自隔离 12/12 通过、scheduler 全包 10/10、全量 go test ./... 3 连跑 0 失败、scheduler -race 无数据竞争。 Co-Authored-By: Claude Opus 4.8 --- server/internal/scheduler/scheduler.go | 17 +++++++---- server/internal/scheduler/scheduler_test.go | 31 +++++++-------------- 2 files changed, 22 insertions(+), 26 deletions(-) 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()