From e4cb1e2e430f378668151c77966622703af61814 Mon Sep 17 00:00:00 2001 From: ShawnKung Date: Tue, 8 Sep 2026 11:55:17 +0800 Subject: [PATCH] fix(align): spread hold time across windows Co-Authored-By: Claude Sonnet 4.6 --- internal/align/align_test.go | 56 +++++++++++++++++++++++---- internal/align/scheduler.go | 73 ++++++++++++++++++++++++++++-------- internal/cli/cli.go | 6 ++- 3 files changed, 110 insertions(+), 25 deletions(-) diff --git a/internal/align/align_test.go b/internal/align/align_test.go index 40b99de..3c00952 100644 --- a/internal/align/align_test.go +++ b/internal/align/align_test.go @@ -32,17 +32,22 @@ func TestDecideNoTargetsAlwaysPings(t *testing.T) { } } -func TestDecideFarFromTargetPingsImmediately(t *testing.T) { - // 目标每晚 00:00;现在 08:00,锚点应为 19:00,距今 11h,远超预算 → 立即 ping。 +func TestDecideSpreadsDelayAcrossReachableWindows(t *testing.T) { + // 目标每晚 00:00;现在 08:00,锚点为 19:00。11h 可以拆成两个 5h 窗口 + // 加三段 20min hold,因此不应把全部延迟堆到最后一跳。 planner := mustPlanner(t, []string{"0 0 * * *"}, DefaultMaxDelayMinutes) now := time.Date(2026, 9, 7, 8, 0, 0, 0, time.Local) decision := planner.Decide(now) - if decision.Action != ActionPing { - t.Fatalf("expected ping when far from anchor, got %+v", decision) + if decision.Action != ActionHold { + t.Fatalf("expected hold when delay can be spread, got %+v", decision) } if !decision.HasTarget { t.Fatalf("expected a resolved target") } + wantPing := time.Date(2026, 9, 7, 8, 20, 0, 0, time.Local) + if !decision.PlannedPing.Equal(wantPing) { + t.Fatalf("planned ping = %s, want %s", decision.PlannedPing, wantPing) + } } func TestDecideHoldsWithinBudget(t *testing.T) { @@ -106,8 +111,22 @@ func TestDecideUsesPerTargetMaxDelay(t *testing.T) { } } +func TestDecidePingsWhenSpreadWouldExceedBudget(t *testing.T) { + planner, err := NewPlanner(&Config{Targets: []Target{ + {ID: "aaaa", Cron: "0 0 * * *", MaxDelayMinutes: 15}, + }}) + if err != nil { + t.Fatalf("NewPlanner: %v", err) + } + now := time.Date(2026, 9, 7, 8, 0, 0, 0, time.Local) + decision := planner.Decide(now) + if decision.Action != ActionPing { + t.Fatalf("expected ping when spread delay exceeds budget, got %+v", decision) + } +} + func TestNextPingTimesConvergeThenAlign(t *testing.T) { - // 从 08:00 起链式:前几次每 5h,临近后把某次拖到 19:00,使刷新落在 00:00。 + // 从 08:00 起把额外等待均匀摊到后续窗口,最终使刷新落在 00:00。 planner := mustPlanner(t, []string{"0 0 * * *"}, DefaultMaxDelayMinutes) start := time.Date(2026, 9, 7, 8, 0, 0, 0, time.Local) plans := planner.NextPingTimes(start, 5) @@ -132,6 +151,23 @@ func TestNextPingTimesConvergeThenAlign(t *testing.T) { } } +func TestNextPingTimesSpreadDailyMidnightTarget(t *testing.T) { + planner := mustPlanner(t, []string{"0 0 * * *"}, 240) + start := time.Date(2026, 9, 9, 0, 0, 0, 0, time.Local) + plans := planner.NextPingTimes(start, 4) + wantPings := []time.Time{ + time.Date(2026, 9, 9, 1, 0, 0, 0, time.Local), + time.Date(2026, 9, 9, 7, 0, 0, 0, time.Local), + time.Date(2026, 9, 9, 13, 0, 0, 0, time.Local), + time.Date(2026, 9, 9, 19, 0, 0, 0, time.Local), + } + for i, want := range wantPings { + if !plans[i].PingAt.Equal(want) { + t.Fatalf("plan[%d].PingAt = %s, want %s; plans=%v", i, plans[i].PingAt, want, plans) + } + } +} + func TestUpcomingUsesRunningWindow(t *testing.T) { planner := mustPlanner(t, []string{"0 0 * * *"}, DefaultMaxDelayMinutes) now := time.Date(2026, 9, 7, 8, 0, 0, 0, time.Local) @@ -153,9 +189,9 @@ func TestUpcomingNormalizesPersistedUTCLastPing(t *testing.T) { lastPingLocal := time.Date(2026, 9, 8, 10, 41, 0, 0, time.Local) lastPingUTC := lastPingLocal.UTC() - plans := planner.Upcoming(now, &lastPingUTC, 3) - if len(plans) != 3 { - t.Fatalf("want 3 plans, got %d", len(plans)) + plans := planner.Upcoming(now, &lastPingUTC, 5) + if len(plans) != 5 { + t.Fatalf("want 5 plans, got %d", len(plans)) } wantPing := time.Date(2026, 9, 8, 19, 0, 0, 0, time.Local) wantRefresh := time.Date(2026, 9, 9, 0, 0, 0, 0, time.Local) @@ -165,6 +201,10 @@ func TestUpcomingNormalizesPersistedUTCLastPing(t *testing.T) { if plans[0].PingAt.Location() != time.Local || plans[0].RefreshAt.Location() != time.Local { t.Fatalf("plan should use local timezone, got ping=%s refresh=%s", plans[0].PingAt.Location(), plans[0].RefreshAt.Location()) } + wantNextPing := time.Date(2026, 9, 9, 1, 0, 0, 0, time.Local) + if !plans[1].PingAt.Equal(wantNextPing) { + t.Fatalf("plan[1].PingAt = %s, want smoothed next ping %s", plans[1].PingAt, wantNextPing) + } } func TestParseCronRejectsInvalid(t *testing.T) { diff --git a/internal/align/scheduler.go b/internal/align/scheduler.go index e7290f5..6e5b4c9 100644 --- a/internal/align/scheduler.go +++ b/internal/align/scheduler.go @@ -32,7 +32,10 @@ type Decision struct { Target time.Time // Anchor 是理想的 ping 锚点 A* = T - Window,在此锚定可使窗口于 T 重置。 Anchor time.Time - // Delay 是从现在到锚点还需等待的时长(Hold 时为正,Ping 时约为 0 或不适用)。 + // PlannedPing 是本轮满额窗口应执行 ping 的时间。它可能早于最终锚点, + // 用于把对齐目标前的总等待时间均匀摊到多个窗口里。 + PlannedPing time.Time + // Delay 是从现在到本轮计划 ping 还需等待的时长(Hold 时为正,Ping 时约为 0 或不适用)。 Delay time.Duration // MaxDelay 是本次瞄准目标自己的延时预算。 MaxDelay time.Duration @@ -78,36 +81,57 @@ func (p *Planner) HasTargets() bool { return len(p.targets) > 0 } -// Decide 在窗口满额的前提下,判断此刻应 ping 还是 hold 延后。判定是无状态的, -// 仅依赖当前时间与 cron 目标:延后中的每一分钟都会重新评估,随着时间逼近锚点, -// 待等时长单调减小,绝不会超过预算。每个目标使用它自己的延时预算。 +// Decide 在窗口满额的前提下,判断此刻应 ping 还是 hold 延后。 +// 调用方若知道当前窗口的满额起点,应优先使用 DecideAtFull 或 DecideFromLastPing, +// 避免均匀摊分计划在反复轮询时漂移。每个目标使用它自己的延时预算。 func (p *Planner) Decide(now time.Time) Decision { + return p.DecideAtFull(now, now) +} + +// DecideFromLastPing 基于最近一次成功 ping 推导当前窗口满额起点,再做对齐决策。 +// 这样 cron 每分钟重试时会围绕同一个满额起点稳定计算本轮计划 ping 时间。 +func (p *Planner) DecideFromLastPing(now time.Time, lastPing *time.Time) Decision { + now = now.In(time.Local) + full := now + if lastPing != nil { + nextFull := lastPing.In(time.Local).Add(p.window) + if !nextFull.After(now) { + full = nextFull + } + } + return p.DecideAtFull(full, now) +} + +// DecideAtFull 在给定窗口满额起点 full 的前提下,判断 now 是否应 ping。 +func (p *Planner) DecideAtFull(full, now time.Time) Decision { + full = full.In(time.Local) now = now.In(time.Local) if len(p.targets) == 0 { return Decision{Action: ActionPing, Now: now, Reason: "未配置对齐目标,满额即 ping"} } // 只考虑锚点仍可达(A* = T - window >= now - tol)的目标,取锚点最近的一个。 - threshold := now.Add(p.window).Add(-p.tolerance) + threshold := full.Add(p.window).Add(-p.tolerance) target, maxDelay, ok := p.nextOnOrAfter(threshold) if !ok { return Decision{Action: ActionPing, Now: now, Reason: "无可用 cron 目标,满额即 ping"} } anchor := target.Add(-p.window) - delay := anchor.Sub(now) - base := Decision{Now: now, HasTarget: true, Target: target, Anchor: anchor, Delay: delay, MaxDelay: maxDelay} + plannedPing, feasible := p.plannedPingAt(full, anchor, maxDelay) + delay := plannedPing.Sub(now) + base := Decision{Now: now, HasTarget: true, Target: target, Anchor: anchor, PlannedPing: plannedPing, Delay: delay, MaxDelay: maxDelay} switch { - case delay > maxDelay: + case !feasible: base.Action = ActionPing - base.Reason = fmt.Sprintf("距目标锚点还有 %s,超过延时预算 %s,正常链式 ping", - roundDuration(delay), roundDuration(maxDelay)) + base.Reason = fmt.Sprintf("距目标锚点还需累计等待 %s,超过延时预算 %s,正常链式 ping", + roundDuration(anchor.Sub(full)), roundDuration(maxDelay)) case delay <= p.tolerance: base.Action = ActionPing - base.Reason = fmt.Sprintf("已到目标锚点,ping 使窗口约在 %s 重置", - target.Format("01-02 15:04")) + base.Reason = fmt.Sprintf("已到本轮计划点,ping 使窗口约在 %s 重置", + plannedPing.Add(p.window).Format("01-02 15:04")) default: base.Action = ActionHold - base.Reason = fmt.Sprintf("在预算内延后:等到 %s 再 ping,使窗口约在 %s 重置", - anchor.Format("01-02 15:04"), target.Format("01-02 15:04")) + base.Reason = fmt.Sprintf("均匀延后:等到 %s 再 ping,逐步对齐 %s 刷新", + plannedPing.Format("01-02 15:04"), target.Format("01-02 15:04")) } return base } @@ -131,10 +155,10 @@ func (p *Planner) NextPingTimes(start time.Time, n int) []PingPlan { plans := make([]PingPlan, 0, n) full := start.In(time.Local) for i := 0; i < n; i++ { - decision := p.Decide(full) + decision := p.DecideAtFull(full, full) pingAt := full if decision.Action == ActionHold { - pingAt = decision.Anchor + pingAt = decision.PlannedPing } if pingAt.Before(full) { pingAt = full @@ -145,6 +169,23 @@ func (p *Planner) NextPingTimes(start time.Time, n int) []PingPlan { return plans } +func (p *Planner) plannedPingAt(full, anchor time.Time, maxDelay time.Duration) (time.Time, bool) { + untilAnchor := anchor.Sub(full) + if untilAnchor <= p.tolerance { + return full, true + } + windowsBeforeAnchor := int(untilAnchor / p.window) + plannedPings := windowsBeforeAnchor + 1 + totalHold := untilAnchor - time.Duration(windowsBeforeAnchor)*p.window + if totalHold < 0 { + totalHold = 0 + } + if totalHold > maxDelay*time.Duration(plannedPings) { + return full, false + } + return full.Add(totalHold / time.Duration(plannedPings)), true +} + // nextOnOrAfter 返回所有 cron 中不早于 t 的最近一次触发时刻,及该目标自己的延时预算。 func (p *Planner) nextOnOrAfter(t time.Time) (time.Time, time.Duration, bool) { // cron.Schedule.Next 返回严格晚于入参的时刻,减一纳秒以纳入 t 本身。 diff --git a/internal/cli/cli.go b/internal/cli/cli.go index 48c91bc..d7b39b2 100644 --- a/internal/cli/cli.go +++ b/internal/cli/cli.go @@ -338,7 +338,11 @@ func executePing(ctx context.Context, stdin io.Reader, stdout, stderr io.Writer, } fmt.Fprintln(stdout, "5h 可用额度为 100%,满足触发条件。") if planner != nil && planner.HasTargets() { - decision := planner.Decide(time.Now()) + lastSuccessfulPing, err := pingstate.LastSuccessfulPing() + if err != nil { + return fmt.Errorf("读取最近成功 ping: %w", err) + } + decision := planner.DecideFromLastPing(time.Now(), lastSuccessfulPing) fmt.Fprintf(stdout, "对齐决策:%s\n", decision.Reason) if decision.Action == align.ActionHold { fmt.Fprintf(stdout, "本次延后 ping(目标 %s)。\n", decision.Target.Format("01-02 15:04"))