fix(scheduler): 修复 #15H 两个 flaky leader-election 测试的根因
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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 <leaseTTL ms>).
|
||||
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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user