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
56 changes: 48 additions & 8 deletions internal/align/align_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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)
Expand All @@ -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) {
Expand Down
73 changes: 57 additions & 16 deletions internal/align/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
}
Expand All @@ -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
Expand All @@ -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 本身。
Expand Down
6 changes: 5 additions & 1 deletion internal/cli/cli.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
Expand Down
Loading