diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index cef81f3..4331dd9 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -6,6 +6,19 @@ on: pull_request: jobs: + + quota-state-platforms: + strategy: + matrix: + os: [ubuntu-latest, macos-latest, windows-latest] + runs-on: ${{ matrix.os }} + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-go@v5 + with: + go-version-file: go.mod + - run: go test ./internal/codexstate + build: runs-on: ubuntu-latest steps: diff --git a/README.md b/README.md index f9b6bbe..5ac3f12 100644 --- a/README.md +++ b/README.md @@ -27,8 +27,8 @@ after the terminal closes. ``` claude ✓ pinged (6.6s) -codex ✓ pinged (13.6s) -spark ✓ pinged (12.4s) +codex CLI trigger returned without error (13.6s); turn completion is not verified +spark CLI trigger returned without error (12.4s); turn completion is not verified ``` ## Highlights @@ -112,6 +112,27 @@ and pings as soon as the window resets. Claude/Codex tokens are reused from the official tools (no separate login) and refreshed on 401. Spark reuses the Codex token. +### Codex/Spark window verification + +At 0% usage, one quota API response cannot always tell whether a window has started. +For example, compare these five-hour reset times: + +| Pattern (both show 0% used) | Read at 10:00 | Read at 10:01 | +| --- | --- | --- | +| Started: reset stays fixed | 15:00 | 15:00 | +| Not started: reset slides forward | 15:00 | 15:01 | + +To distinguish these patterns, limitping compares reads at least one minute apart. +Inconclusive results stay unconfirmed. `ping` returns without waiting that minute +and suggests a later `status` check; `watch`/`bg` rechecks automatically. + +If the pre-ping quota check fails, both manual and automatic pings report the +reason and stop without sending. Authentication reload/refresh on HTTP 401 is +still attempted. Only the watcher waits before retrying: authentication/permission +failures back off from 30 seconds to one hour; other read failures cap at ten +minutes, with `Retry-After` respected. Restart the watcher to retry immediately +after fixing access. Manual pings do not wait through this backoff. + ## Install `limitping` ships as a single self-contained binary — **no Go required**. @@ -232,9 +253,9 @@ elapsed time only: claude → claude --model haiku . claude ✓ pinged (6.6s) codex → codex -c model_reasoning_effort=low -m gpt-5.4-mini ok -codex ✓ pinged (13.6s) +codex CLI trigger returned without error (13.6s); turn completion is not verified spark → codex -c model_reasoning_effort=low -m gpt-5.3-codex-spark ok -spark ✓ pinged (12.4s) +spark CLI trigger returned without error (12.4s); turn completion is not verified ``` Use `status` or `bg status` for the authoritative 5h/weekly window view after a diff --git a/README.zh-CN.md b/README.zh-CN.md index c64cf41..fe7e48c 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -23,8 +23,8 @@ Claude Code、Codex 和 Spark 的订阅限额按 **5 小时滚动窗口**(外加 ``` claude ✓ pinged (6.6s) -codex ✓ pinged (13.6s) -spark ✓ pinged (12.4s) +codex CLI trigger returned without error (13.6s); turn completion is not verified +spark CLI trigger returned without error (12.4s); turn completion is not verified ``` ## 亮点 @@ -97,6 +97,25 @@ limitping bg logs -f Claude/Codex 的 token 直接复用官方工具(无需另外登录),遇到 401 会自动刷新。Spark 复用 Codex token。 +### Codex/Spark 窗口验证 + +用量为 0% 时,单次限额 API 响应不一定能判断窗口是否已启动。 +以下以五小时窗口的重置时间为例: + +| 情况(用量均为 0%) | 10:00 查询 | 10:01 查询 | +| --- | --- | --- | +| 已启动:重置时间固定 | 15:00 | 15:00 | +| 未启动:重置时间向后滑动 | 15:00 | 15:01 | + +因此,limitping 会比较至少间隔一分钟的查询结果;证据不足时仍显示未确认。 +`ping` 不会等待这一分钟,而是提示稍后运行 `status`;`watch`/`bg` 会自动复查。 + +发送前的限额检查失败时,手动和自动 ping 都会说明原因并停止,不发送请求。 +HTTP 401 仍会尝试重新加载或刷新凭据。只有 watcher 负责重试等待: +认证或权限错误从 30 秒退避至最多一小时,其他读取错误最多十分钟, +并遵守 `Retry-After`。修复访问权限后可重启 watcher 立即重试; +手动 ping 不等待这些退避间隔。 + ## 安装 `limitping` 是一个自包含的单文件二进制——**普通用户无需安装 Go**。 @@ -214,9 +233,9 @@ limitping uninstall # 删除 limitping 以及配置/缓存(简称: rm claude → claude --model haiku . claude ✓ pinged (6.6s) codex → codex -c model_reasoning_effort=low -m gpt-5.4-mini ok -codex ✓ pinged (13.6s) +codex CLI trigger returned without error (13.6s); turn completion is not verified spark → codex -c model_reasoning_effort=low -m gpt-5.3-codex-spark ok -spark ✓ pinged (12.4s) +spark CLI trigger returned without error (12.4s); turn completion is not verified ``` ping 后请用 `status` 或 `bg status` 查看权威的 5h/周窗口状态。 diff --git a/internal/cli/bg.go b/internal/cli/bg.go index 21a0497..4d0812d 100644 --- a/internal/cli/bg.go +++ b/internal/cli/bg.go @@ -386,6 +386,8 @@ func parseBgPingAttempt(line string) (bgPingAttempt, bool) { status = bgPingFailed case strings.Contains(msg, "ping sent, new window started"): status = bgPingSucceeded + case strings.Contains(msg, "ping trigger returned; checking window"): + status = bgPingSucceeded default: return bgPingAttempt{}, false } diff --git a/internal/cli/i18n.go b/internal/cli/i18n.go index fd44eeb..3ab6496 100644 --- a/internal/cli/i18n.go +++ b/internal/cli/i18n.go @@ -6,10 +6,23 @@ import ( ) type cliText struct { - rootShort string - rootLong string - helpFlag string - usageTemplate string + verifyAlreadyActive string + verifyQuotaStateFmt string + verifyStarted string + verifyNotStarted string + verifyUnknown string + verifyStatusCheck string + verifyRetryFmt string + verifyRunning string + verifyReady string + verifyUnavailable string + verifyDisabled string + verifyNoBaseline string + verifyCheckFmt string + rootShort string + rootLong string + helpFlag string + usageTemplate string helpCommandShort string helpCommandLong string @@ -56,13 +69,14 @@ type cliText struct { statusNowWord string statusWeekdays [7]string // Sunday first; zero value = Go's "Mon" names - pingShort string - pingLong string - pingDryRunFlag string - pingWouldRunFmt string // provider, command - pingSendingFmt string // provider, spinner frame, elapsed - pingFailedFmt string // provider, elapsed, error - pingSuccessFmt string // provider, elapsed, usage suffix + pingShort string + pingLong string + pingDryRunFlag string + pingWouldRunFmt string // provider, command + pingSendingFmt string // provider, spinner frame, elapsed + pingFailedFmt string // provider, elapsed, error + pingSuccessFmt string // provider, elapsed, usage suffix + pingTriggerReturnedFmt string watchShort string watchLong string @@ -186,9 +200,22 @@ func isChineseLocale() bool { } var enText = cliText{ - rootShort: "Keep Claude Code / Codex / Spark rate-limit windows back-to-back", - rootLong: "limitping pings your AI coding provider the moment its 5h rate-limit window resets, so the next window starts immediately and stays aligned. Usage is read via zero-quota endpoints; pings go through the official CLIs.", - helpFlag: "help for this command", + verifyAlreadyActive: "already active before ping", + verifyQuotaStateFmt: " %s quota state: %s\n", + verifyStarted: "window started", + verifyNotStarted: "window not started (reset is an estimate)", + verifyUnknown: "window start unconfirmed (reset is an estimate)", + verifyStatusCheck: "Window start unconfirmed; run `limitping status` again in about a minute.", + verifyRetryFmt: "Automatic ping retry eligible at %s.", + verifyRunning: "Ping in progress.", + verifyReady: "Window not started; eligible for automatic ping, subject to watcher checks.", + verifyUnavailable: "Quota-window state unavailable; run `limitping status` again in about a minute.", + verifyDisabled: " This provider is disabled in config and will not appear in `limitping status`.", + verifyNoBaseline: " Run `limitping status` to collect a baseline, then check again after 60s.", + verifyCheckFmt: " Run `limitping status` after %s to recheck (no background check was scheduled by this command).\n", + rootShort: "Keep Claude Code / Codex / Spark rate-limit windows back-to-back", + rootLong: "limitping pings your AI coding provider the moment its 5h rate-limit window resets, so the next window starts immediately and stays aligned. Usage is read via zero-quota endpoints; pings go through the official CLIs.", + helpFlag: "help for this command", usageTemplate: `Usage:{{if .Runnable}} {{.UseLine}}{{end}}{{if .HasAvailableSubCommands}} {{.CommandPath}} [command]{{end}}{{if gt (len .Aliases) 0}} @@ -272,11 +299,12 @@ Examples: limitping ping limitping p claude limitping ping codex --dry-run`, - pingDryRunFlag: "print the command without sending", - pingWouldRunFmt: "%-7s would run: %s\n", - pingSendingFmt: "\r%-7s %c sending… %s", - pingFailedFmt: "%-7s ✗ failed after %s: %v\n", - pingSuccessFmt: "%-7s ✓ pinged (%s%s)\n", + pingDryRunFlag: "print the command without sending", + pingWouldRunFmt: "%-7s would run: %s\n", + pingSendingFmt: "\r%-7s %c sending… %s", + pingFailedFmt: "%-7s ✗ failed after %s: %v\n", + pingSuccessFmt: "%-7s ✓ pinged (%s%s)\n", + pingTriggerReturnedFmt: "%-7s CLI trigger returned without error (%s%s); turn completion is not verified\n", watchShort: "Run the foreground daemon and ping each provider when its 5h window resets", watchLong: `Run the foreground daemon. When a provider's 5h window resets, limitping sends the minimal message to start the next window. @@ -444,9 +472,22 @@ Examples: } var zhText = cliText{ - rootShort: "让 Claude Code / Codex / Spark 的限额窗口自动接龙", - rootLong: "limitping 会在 AI 编程 Provider 的 5h 限额窗口重置时立即发送 ping,让下一个窗口马上开始并保持对齐。用量读取走零消耗接口;ping 通过官方 CLI 发送。", - helpFlag: "显示此命令的帮助", + verifyAlreadyActive: "ping 前已启动", + verifyQuotaStateFmt: " %s 限额状态:%s\n", + verifyStarted: "窗口已启动", + verifyNotStarted: "窗口未启动(重置时间为估计值)", + verifyUnknown: "窗口启动尚未确认(重置时间为估计值)", + verifyStatusCheck: "窗口启动尚未确认;请约一分钟后再次运行 `limitping status`。", + verifyRetryFmt: "自动 ping 最早可于 %s 重试。", + verifyRunning: "Ping 正在进行。", + verifyReady: "窗口未启动;可尝试自动 ping,仍需通过监视器检查。", + verifyUnavailable: "窗口状态不可用;请约一分钟后再次运行 `limitping status`。", + verifyDisabled: " 此服务在配置中已禁用,不会出现在 `limitping status` 中。", + verifyNoBaseline: " 运行 `limitping status` 采集基准,60 秒后再次检查。", + verifyCheckFmt: " %s 后运行 `limitping status` 再次检查(本命令未安排后台检查)。\n", + rootShort: "让 Claude Code / Codex / Spark 的限额窗口自动接龙", + rootLong: "limitping 会在 AI 编程 Provider 的 5h 限额窗口重置时立即发送 ping,让下一个窗口马上开始并保持对齐。用量读取走零消耗接口;ping 通过官方 CLI 发送。", + helpFlag: "显示此命令的帮助", usageTemplate: `用法:{{if .Runnable}} {{.UseLine}}{{end}}{{if .HasAvailableSubCommands}} {{.CommandPath}} [command]{{end}}{{if gt (len .Aliases) 0}} @@ -532,11 +573,12 @@ var zhText = cliText{ limitping ping limitping p claude limitping ping codex --dry-run`, - pingDryRunFlag: "只打印将执行的命令,不真正发送", - pingWouldRunFmt: "%-7s 将执行: %s\n", - pingSendingFmt: "\r%-7s %c 发送中… %s", - pingFailedFmt: "%-7s ✗ 失败 (耗时 %s): %v\n", - pingSuccessFmt: "%-7s ✓ 已 ping (%s%s)\n", + pingDryRunFlag: "只打印将执行的命令,不真正发送", + pingWouldRunFmt: "%-7s 将执行: %s\n", + pingSendingFmt: "\r%-7s %c 发送中… %s", + pingFailedFmt: "%-7s ✗ 失败 (耗时 %s): %v\n", + pingSuccessFmt: "%-7s ✓ 已 ping (%s%s)\n", + pingTriggerReturnedFmt: "%-7s CLI 触发已返回且未报错(%s%s);尚未验证轮次完成\n", watchShort: "以前台守护方式运行,并在每个 Provider 的 5h 窗口重置时自动 ping", watchLong: `以前台守护方式运行。某个 Provider 的 5h 窗口重置后,limitping 会发送最小消息来开启下一个窗口。 diff --git a/internal/cli/ping.go b/internal/cli/ping.go index 31a8717..08c944b 100644 --- a/internal/cli/ping.go +++ b/internal/cli/ping.go @@ -5,6 +5,7 @@ import ( "fmt" "io" "os" + "sync/atomic" "time" "github.com/spf13/cobra" @@ -72,6 +73,9 @@ func runPing(parent context.Context, out io.Writer, text cliText, p provider.Pro defer cancel() start := time.Now() + var stage atomic.Value + stage.Store("") + ctx = provider.WithPingStage(ctx, func(s string) { stage.Store(s) }) type outcome struct { res *provider.TriggerResult err error @@ -100,18 +104,27 @@ func runPing(parent context.Context, out io.Writer, text cliText, p provider.Pro report(out, text, name, start, o.res, o.err) return o.err case <-ticker.C: - fmt.Fprintf(out, text.pingSendingFmt, name, frames[i%len(frames)], elapsed(start)) + label := name + if s := stage.Load().(string); s != "" { + label += " [" + s + "]" + } + fmt.Fprintf(out, text.pingSendingFmt, label, frames[i%len(frames)], elapsed(start)) i++ } } } func report(out io.Writer, text cliText, name string, start time.Time, res *provider.TriggerResult, err error) { + defer reportVerification(out, text, res) if err != nil { fmt.Fprintf(out, text.pingFailedFmt, name, elapsed(start), localizedProviderError(text, err)) return } - fmt.Fprintf(out, text.pingSuccessFmt, name, elapsed(start), usageSuffix(res)) + format := text.pingSuccessFmt + if res != nil && res.Verification != nil { + format = text.pingTriggerReturnedFmt + } + fmt.Fprintf(out, format, name, elapsed(start), usageSuffix(res)) } // usageSuffix renders the token/cost tail, e.g. ", 32,934 tok, $0.0110". diff --git a/internal/cli/ping_quota_state_test.go b/internal/cli/ping_quota_state_test.go new file mode 100644 index 0000000..f3e6bef --- /dev/null +++ b/internal/cli/ping_quota_state_test.go @@ -0,0 +1,46 @@ +package cli + +import ( + "bytes" + "fmt" + "strings" + "testing" + "time" + + "github.com/wavever/CCLimitPing/internal/provider" + "github.com/wavever/CCLimitPing/internal/usage" +) + +func TestPingAlreadyActiveBeforePing(t *testing.T) { + for _, target := range []string{"weekly", "five_hour"} { + for _, before := range []string{"", "unknown", "not_started", "started"} { + for _, after := range []string{"unknown", "not_started", "started"} { + for _, failed := range []bool{false, true} { + res := &provider.TriggerResult{StatusEnabled: true, + Verification: &usage.Verification{Target: target, + FiveHour: usage.StartStatus{State: after}, Weekly: usage.StartStatus{State: after}}, + } + if before != "" { + res.PreVerification = &usage.Verification{Target: target, + FiveHour: usage.StartStatus{State: before}, Weekly: usage.StartStatus{State: before}} + } + var err error + if failed { + err = fmt.Errorf("CLI failed") + } + for _, text := range []cliText{enText, zhText} { + var out bytes.Buffer + report(&out, text, "codex", time.Now(), res, err) + want := before == "started" && after == "started" + if strings.Contains(out.String(), text.verifyAlreadyActive) != want { + t.Fatalf("%s/%s/%s/failed=%v: %s", target, before, after, failed, out.String()) + } + if want && !strings.Contains(out.String(), fmt.Sprintf(text.verifyQuotaStateFmt, target, text.verifyAlreadyActive)) { + t.Fatal(out.String()) + } + } + } + } + } + } +} diff --git a/internal/cli/status.go b/internal/cli/status.go index ff84f73..9e921f7 100644 --- a/internal/cli/status.go +++ b/internal/cli/status.go @@ -48,6 +48,7 @@ func runStatus(ctx context.Context, out, progress io.Writer, text cliText, provi progress = io.Discard } display = normalizeUsageDisplay(display) + diagnostics := progress // In JSON mode keep stdout a single valid document: suppress the // "Fetching..." progress chatter that would otherwise interleave. if jsonOut { @@ -72,6 +73,9 @@ func runStatus(ctx context.Context, out, progress io.Writer, text cliText, provi continue } if jsonOut { + if u.Verification != nil && u.Verification.Warning != "" { + fmt.Fprintln(diagnostics, u.Verification.Warning) + } entries = append(entries, newStatusJSON(u, verbose)) continue } @@ -115,12 +119,14 @@ type statusJSON struct { } type windowJSON struct { - UsedPercent float64 `json:"used_percent"` - RemainingPercent float64 `json:"remaining_percent"` - Active bool `json:"active"` - ResetsAt string `json:"resets_at,omitempty"` - RemainingSeconds int `json:"remaining_seconds"` - WindowSeconds int `json:"window_seconds,omitempty"` + UsedPercent float64 `json:"used_percent"` + RemainingPercent float64 `json:"remaining_percent"` + Active bool `json:"active"` + ResetsAt string `json:"resets_at,omitempty"` + RemainingSeconds int `json:"remaining_seconds"` + WindowSeconds int `json:"window_seconds,omitempty"` + StartState string `json:"start_state,omitempty"` + VerificationDueAt string `json:"verification_due_at,omitempty"` } type creditsJSON struct { @@ -155,6 +161,10 @@ func newStatusJSON(u *usage.Usage, verbose bool) statusJSON { if !u.Weekly.Missing() { s.Weekly = newWindowJSON(u.Weekly) } + if u.Verification != nil { + addStartJSON(s.FiveHour, u.Verification.FiveHour) + addStartJSON(s.Weekly, u.Verification.Weekly) + } if !u.FetchedAt.IsZero() { s.FetchedAt = u.FetchedAt.Format(time.RFC3339) } @@ -174,6 +184,16 @@ func newStatusJSON(u *usage.Usage, verbose bool) statusJSON { return s } +func addStartJSON(w *windowJSON, s usage.StartStatus) { + if w == nil { + return + } + w.StartState = s.State + if !s.DueAt.IsZero() { + w.VerificationDueAt = s.DueAt.Format(time.RFC3339) + } +} + func newWindowJSON(w usage.Window) *windowJSON { j := &windowJSON{ UsedPercent: w.UsedPercent, @@ -218,8 +238,34 @@ func printUsage(out io.Writer, text cliText, u *usage.Usage, verbose bool, displ plan = " (" + plan + ")" } fmt.Fprintf(out, "%s%s\n", u.Provider, plan) - fmt.Fprintf(out, text.statusFiveHourLineFmt, fmtWindow(text, u.FiveHour, display)) - fmt.Fprintf(out, text.statusWeeklyLineFmt, fmtWindow(text, u.Weekly, display)) + five, week := fmtWindow(text, u.FiveHour, display), fmtWindow(text, u.Weekly, display) + if v := u.Verification; v != nil { + if !u.FiveHour.Missing() { + five += " — " + startDescription(text, v.FiveHour) + } + if !u.Weekly.Missing() { + week += " — " + startDescription(text, v.Weekly) + } + } + fmt.Fprintf(out, text.statusFiveHourLineFmt, five) + fmt.Fprintf(out, text.statusWeeklyLineFmt, week) + if v := u.Verification; v != nil { + switch v.Recovery { + case "verifying": + fmt.Fprintln(out, " "+text.verifyStatusCheck) + case "backoff", "cooldown": + fmt.Fprintf(out, " "+text.verifyRetryFmt+"\n", fmtClock(text, v.NextEligible)) + case "ping_running": + fmt.Fprintln(out, " "+text.verifyRunning) + case "ready": + fmt.Fprintln(out, " "+text.verifyReady) + case "unavailable": + fmt.Fprintln(out, " "+text.verifyUnavailable) + } + if v.Warning != "" { + fmt.Fprintln(out, " "+v.Warning) + } + } if u.Credits != nil && (u.Credits.HasCredits || u.Credits.Unlimited) { if u.Credits.Unlimited { fmt.Fprint(out, text.statusCreditsUnlimited) diff --git a/internal/cli/status_recovery_test.go b/internal/cli/status_recovery_test.go new file mode 100644 index 0000000..7c4c312 --- /dev/null +++ b/internal/cli/status_recovery_test.go @@ -0,0 +1,52 @@ +package cli + +import ( + "bytes" + "strings" + "testing" + "time" + + "github.com/wavever/CCLimitPing/internal/usage" +) + +func TestStatusRecoveryGuidance(t *testing.T) { + for _, text := range []cliText{enText, zhText} { + for _, tc := range []struct { + state string + want string + }{ + {"window_started", ""}, + {"verifying", text.verifyStatusCheck}, + {"backoff", fmtClock(text, time.Date(2026, 9, 7, 12, 4, 0, 0, time.UTC))}, + {"cooldown", fmtClock(text, time.Date(2026, 9, 7, 12, 4, 0, 0, time.UTC))}, + {"ping_running", text.verifyRunning}, + {"ready", text.verifyReady}, + {"unavailable", text.verifyUnavailable}, + } { + t.Run(tc.state+"/"+text.verifyStarted, func(t *testing.T) { + u := &usage.Usage{Provider: "codex"} + var plain bytes.Buffer + printUsage(&plain, text, u, false, "") + u.Verification = &usage.Verification{ + Recovery: tc.state, + NextEligible: time.Date(2026, 9, 7, 12, 4, 0, 0, time.UTC), + } + var out bytes.Buffer + printUsage(&out, text, u, false, "") + if tc.want == "" { + if out.String() != plain.String() { + t.Fatal(out.String()) + } + } else if !strings.Contains(out.String(), tc.want) { + t.Fatal(out.String()) + } + u.Verification.Warning = "observation warning" + out.Reset() + printUsage(&out, text, u, false, "") + if !strings.Contains(out.String(), "observation warning") { + t.Fatal(out.String()) + } + }) + } + } +} diff --git a/internal/cli/verification.go b/internal/cli/verification.go new file mode 100644 index 0000000..64f938e --- /dev/null +++ b/internal/cli/verification.go @@ -0,0 +1,70 @@ +package cli + +import ( + "fmt" + "io" + "time" + + "github.com/wavever/CCLimitPing/internal/codexstate" + "github.com/wavever/CCLimitPing/internal/provider" + "github.com/wavever/CCLimitPing/internal/usage" +) + +func startDescription(text cliText, s usage.StartStatus) string { + switch s.State { + case codexstate.Started: + return text.verifyStarted + case codexstate.NotStarted: + return text.verifyNotStarted + default: + if !s.DueAt.IsZero() { + d := time.Until(s.DueAt).Round(time.Second) + if d < 0 { + d = 0 + } + return fmt.Sprintf("%s; recheck after %s", text.verifyUnknown, d) + } + return text.verifyUnknown + } +} + +func reportVerification(out io.Writer, text cliText, res *provider.TriggerResult) { + if res == nil || res.Verification == nil { + return + } + v := res.Verification + s := v.FiveHour + if v.Target == "weekly" { + s = v.Weekly + } + description := startDescription(text, s) + if pre := res.PreVerification; pre != nil && pre.Target == v.Target && s.State == codexstate.Started { + before := pre.FiveHour + if v.Target == "weekly" { + before = pre.Weekly + } + if before.State == codexstate.Started { + description = text.verifyAlreadyActive + } + } + fmt.Fprintf(out, text.verifyQuotaStateFmt, v.Target, description) + if v.Warning != "" { + fmt.Fprintln(out, " "+v.Warning) + } + if s.State == codexstate.Started { + return + } + if !res.StatusEnabled { + fmt.Fprintln(out, text.verifyDisabled) + return + } + if s.DueAt.IsZero() { + fmt.Fprintln(out, text.verifyNoBaseline) + return + } + d := time.Until(s.DueAt).Round(time.Second) + if d < 0 { + d = 0 + } + fmt.Fprintf(out, text.verifyCheckFmt, d) +} diff --git a/internal/cli/verification_test.go b/internal/cli/verification_test.go new file mode 100644 index 0000000..8efe849 --- /dev/null +++ b/internal/cli/verification_test.go @@ -0,0 +1,104 @@ +package cli + +import ( + "bytes" + "fmt" + "strings" + "testing" + "time" + + "github.com/wavever/CCLimitPing/internal/codexstate" + "github.com/wavever/CCLimitPing/internal/provider" + "github.com/wavever/CCLimitPing/internal/usage" +) + +func TestBusyPingOutputSuggestsRetry(t *testing.T) { + var out bytes.Buffer + report(&out, enText, "codex", time.Now(), nil, fmt.Errorf("ping not sent: %w", codexstate.ErrBusy)) + want := "ping not sent: another ping or quota-state update is in progress.\nPlease try again in about a minute." + if got := out.String(); !strings.Contains(got, want) || strings.Contains(got, "✓") { + t.Fatal(got) + } +} + +func TestPrecheckFailureOutputSaysNotSent(t *testing.T) { + var out bytes.Buffer + err := fmt.Errorf("ping not sent: quota precheck failed: %w", &provider.UsageHTTPError{StatusCode: 403}) + report(&out, enText, "codex", time.Now(), nil, err) + if got := out.String(); !strings.Contains(got, "ping not sent: quota precheck failed") || !strings.Contains(got, "403") || strings.Contains(got, "✓") { + t.Fatal(got) + } +} + +func TestVerificationGuidance(t *testing.T) { + for _, tc := range []struct { + name, state string + baseline, enabled bool + want string + }{ + {"pending", "unknown", true, true, "after"}, + {"no-baseline", "unknown", false, true, "collect a baseline"}, + {"disabled", "unknown", true, false, "disabled"}, + {"confirmed", "started", false, true, "window started"}, + } { + t.Run(tc.name, func(t *testing.T) { + v := &usage.Verification{Target: "weekly", Weekly: usage.StartStatus{State: tc.state}} + if tc.baseline { + v.Weekly.DueAt = time.Now().Add(time.Minute) + } + var out bytes.Buffer + reportVerification(&out, enText, &provider.TriggerResult{Verification: v, StatusEnabled: tc.enabled}) + if !strings.Contains(out.String(), tc.want) { + t.Fatal(out.String()) + } + if tc.state == "started" && strings.Contains(out.String(), "Run `") { + t.Fatal(out.String()) + } + }) + } +} + +func TestStartJSONDoesNotChangeLegacyActiveOrClaude(t *testing.T) { + u := &usage.Usage{Provider: "codex", Weekly: usage.Window{ResetsAt: time.Now().Add(time.Hour), WindowSeconds: 604800}, + Verification: &usage.Verification{Weekly: usage.StartStatus{State: "started"}}} + j := newStatusJSON(u, false) + if j.Weekly.Active || j.Weekly.StartState != "started" { + t.Fatal(j.Weekly) + } + u.Verification = nil + u.Provider = "claude" + j = newStatusJSON(u, false) + if j.Weekly.StartState != "" { + t.Fatal(j.Weekly) + } +} + +func TestBackgroundVerificationDoesNotCountAsPing(t *testing.T) { + for _, tc := range []struct { + msg string + count bool + }{ + {"ping request completed; checking window", false}, + {"ping trigger returned; checking window", true}, + {"window started after verification", false}, + {"quota read failed: timeout", false}, + {"ping failed: notification timeout; verifying quota before retry", true}, + {"ping sent, new window started", true}, + } { + _, ok := parseBgPingAttempt("2026/09/01 12:00:00 [codex] " + tc.msg) + if ok != tc.count { + t.Fatalf("%s: %v", tc.msg, ok) + } + } +} + +func TestReturnedTriggerDoesNotClaimTurnCompletion(t *testing.T) { + var out bytes.Buffer + report(&out, enText, "codex", time.Now(), &provider.TriggerResult{ + Verification: &usage.Verification{Target: "weekly", Weekly: usage.StartStatus{State: "unknown"}}, + }, nil) + if !strings.Contains(out.String(), "CLI trigger returned without error") || + !strings.Contains(out.String(), "turn completion is not verified") { + t.Fatal(out.String()) + } +} diff --git a/internal/codexstate/file_unix.go b/internal/codexstate/file_unix.go new file mode 100644 index 0000000..930ad21 --- /dev/null +++ b/internal/codexstate/file_unix.go @@ -0,0 +1,26 @@ +//go:build !windows + +package codexstate + +import ( + "errors" + "golang.org/x/sys/unix" + "os" +) + +func lock(path string) (func(), error) { + f, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0600) + if err != nil { + return nil, err + } + if err = unix.Flock(int(f.Fd()), unix.LOCK_EX|unix.LOCK_NB); err != nil { + f.Close() + if errors.Is(err, unix.EWOULDBLOCK) { + return nil, ErrBusy + } + return nil, err + } + return func() { _ = unix.Flock(int(f.Fd()), unix.LOCK_UN); _ = f.Close() }, nil +} + +func replace(from, to string) error { return os.Rename(from, to) } diff --git a/internal/codexstate/file_windows.go b/internal/codexstate/file_windows.go new file mode 100644 index 0000000..d6e3af4 --- /dev/null +++ b/internal/codexstate/file_windows.go @@ -0,0 +1,36 @@ +package codexstate + +import ( + "errors" + "golang.org/x/sys/windows" + "os" +) + +func lock(path string) (func(), error) { + f, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0600) + if err != nil { + return nil, err + } + var o windows.Overlapped + err = windows.LockFileEx(windows.Handle(f.Fd()), windows.LOCKFILE_EXCLUSIVE_LOCK|windows.LOCKFILE_FAIL_IMMEDIATELY, 0, 1, 0, &o) + if err != nil { + f.Close() + if errors.Is(err, windows.ERROR_LOCK_VIOLATION) { + return nil, ErrBusy + } + return nil, err + } + return func() { _ = windows.UnlockFileEx(windows.Handle(f.Fd()), 0, 1, 0, &o); _ = f.Close() }, nil +} + +func replace(from, to string) error { + a, err := windows.UTF16PtrFromString(from) + if err != nil { + return err + } + b, err := windows.UTF16PtrFromString(to) + if err != nil { + return err + } + return windows.MoveFileEx(a, b, windows.MOVEFILE_REPLACE_EXISTING|windows.MOVEFILE_WRITE_THROUGH) +} diff --git a/internal/codexstate/state.go b/internal/codexstate/state.go new file mode 100644 index 0000000..ceab2a1 --- /dev/null +++ b/internal/codexstate/state.go @@ -0,0 +1,356 @@ +// Package codexstate interprets Codex quota observations. It never sends requests. +package codexstate + +import ( + "crypto/rand" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "time" + + "github.com/wavever/CCLimitPing/internal/usage" +) + +const ( + Started = "started" + NotStarted = "not_started" + Unknown = "unknown" + Interval = time.Minute + claimDuration = 3*time.Minute + 15*time.Second + tolerance = 5 * time.Second +) + +var ErrBusy = errors.New("another ping or quota-state update is in progress.\nPlease try again in about a minute.") +var ErrDeferred = errors.New("automatic ping deferred; quota verification or cooldown is pending") + +type Sample struct { + At time.Time + Reset time.Time + Seconds int + Used float64 +} + +type window struct { + Baseline Sample + Latest Sample + State string + // ConfirmedReset is independent of the post-attempt failure baseline. + ConfirmedReset time.Time + PreviousReset time.Time +} + +type bucket struct { + Plan string + FiveHour window + Weekly window + LastRead time.Time + AttemptEnd time.Time + ClaimID string + ClaimUntil time.Time + Attempts []time.Time + Retry int + LastAuto time.Time +} + +type diskState struct { + Version int + Buckets map[string]*bucket +} + +type Store struct { + Dir string + ResetBuffer time.Duration // automatic sends only; never shifts an observation +} + +func sample(w usage.Window, now time.Time) Sample { + return Sample{At: now, Reset: w.ResetsAt, Seconds: w.WindowSeconds, Used: w.UsedPercent} +} + +func near(a, b time.Time) bool { return a.Sub(b).Abs() <= tolerance } +func valid(s Sample) bool { + return s.Seconds > 0 && s.Used >= 0 && s.Used <= 100 && s.Reset.After(s.At) && + s.Reset.Sub(s.At) <= time.Duration(s.Seconds)*time.Second+tolerance +} + +func observe(w *window, s Sample, duringAttempt bool) usage.StartStatus { + previousReset := w.PreviousReset + if w.Latest.Seconds != s.Seconds || s.At.Before(w.Latest.At) { + previousReset = time.Time{} + } else if !w.ConfirmedReset.IsZero() && !w.ConfirmedReset.After(s.At) { + previousReset = w.ConfirmedReset + } + if !valid(s) { + *w = window{PreviousReset: previousReset, Latest: s} + return usage.StartStatus{State: Unknown} + } + old := w.Baseline + if w.Latest.Seconds != s.Seconds || s.At.Before(w.Latest.At) || + (!w.ConfirmedReset.IsZero() && (!w.ConfirmedReset.After(s.At) || !near(w.ConfirmedReset, s.Reset))) { + *w = window{PreviousReset: previousReset} + old = Sample{} + } + state := Unknown + if s.Used > 0 || (!w.ConfirmedReset.IsZero() && near(w.ConfirmedReset, s.Reset)) { + state = Started + } else if !duringAttempt && w.State == NotStarted && s.At.Sub(w.Latest.At) >= 0 && s.At.Sub(w.Latest.At) <= Interval && + near(s.Reset, s.At.Add(time.Duration(s.Seconds)*time.Second)) && + near(s.Reset, w.Latest.Reset.Add(s.At.Sub(w.Latest.At))) { + // A fresh pre-send read may extend a recently established sliding pair. + // It must not erase that evidence merely because it arrives within 60s. + state = NotStarted + } else if valid(old) && old.Seconds == s.Seconds && old.Reset.After(s.At) && s.At.Sub(old.At) >= Interval { + if near(old.Reset, s.Reset) { + state = Started + } else if !duringAttempt && s.At.Sub(old.At) <= 10*time.Minute && + near(old.Reset, old.At.Add(time.Duration(old.Seconds)*time.Second)) && + near(s.Reset, s.At.Add(time.Duration(s.Seconds)*time.Second)) && + (s.Reset.Sub(old.Reset)-s.At.Sub(old.At)).Abs() <= tolerance { + state = NotStarted + } + } + if state == Started && w.ConfirmedReset.IsZero() { + w.ConfirmedReset = s.Reset + } + // During a live claim do not collect evidence that could authorize its retry. + if !duringAttempt && (old.At.IsZero() || state != Unknown || !old.Reset.After(s.At) || + s.At.Sub(old.At) > 10*time.Minute) { + w.Baseline = s + } + w.Latest, w.State = s, state + result := usage.StartStatus{State: state} + if state == Unknown && !w.Baseline.At.IsZero() { + result.DueAt = w.Baseline.At.Add(Interval) + if !result.DueAt.After(s.At) { + result.DueAt = s.At.Add(Interval) + } + } + return result +} + +func target(b *bucket) (string, *window) { + if b.FiveHour.Latest.Seconds != 0 { + return "five_hour", &b.FiveHour + } + if b.Weekly.Latest.Seconds != 0 { + return "weekly", &b.Weekly + } + return "", nil +} + +func prune(b *bucket, now time.Time) { + kept := b.Attempts[:0] + for _, at := range b.Attempts { + if at.Add(time.Hour).After(now) { + kept = append(kept, at) + } + } + b.Attempts = kept +} + +func view(b *bucket, now time.Time) *usage.Verification { + v := &usage.Verification{FiveHour: status(b.FiveHour, now), Weekly: status(b.Weekly, now)} + name, w := target(b) + v.Target = name + if w != nil { + v.PreviousReset = w.PreviousReset + } + switch { + case b.ClaimID != "" && b.ClaimUntil.After(now): + v.Recovery, v.NextEligible = "ping_running", b.ClaimUntil + case w == nil: + v.Recovery = "unavailable" + case w.State == Started: + v.Recovery, v.NextEligible = "window_started", w.Latest.Reset + default: + v.Recovery = "verifying" + v.NextEligible = now.Add(Interval) + if w.State == NotStarted { + v.Recovery, v.NextEligible = "ready", now + } + if d := retryDue(b); d.After(v.NextEligible) { + v.Recovery, v.NextEligible = "backoff", d + } + if len(b.Attempts) >= 4 && b.Attempts[0].Add(time.Hour).After(v.NextEligible) { + v.Recovery, v.NextEligible = "cooldown", b.Attempts[0].Add(time.Hour) + } + } + return v +} + +func status(w window, now time.Time) usage.StartStatus { + s := usage.StartStatus{State: w.State} + if s.State == "" { + s.State = Unknown + } + if s.State == Unknown && !w.Baseline.At.IsZero() { + s.DueAt = w.Baseline.At.Add(Interval) + if !s.DueAt.After(now) { + s.DueAt = now.Add(Interval) + } + } + return s +} + +func retryDue(b *bucket) time.Time { + delay := time.Minute + if b.Retry >= 2 { + delay = 5 * time.Minute + } + if b.Retry >= 3 { + delay = 15 * time.Minute + } + due := b.LastAuto.Add(delay) + if b.AttemptEnd.Add(Interval).After(due) { + due = b.AttemptEnd.Add(Interval) + } + return due +} + +func (s Store) Observe(account, key string, u *usage.Usage) (*usage.Verification, error) { + var v *usage.Verification + err := s.update(account, key, func(b *bucket) error { + now := u.FetchedAt + prune(b, now) + if b.Plan != u.Plan || now.Before(b.LastRead) { + b.FiveHour, b.Weekly = window{}, window{} + } + b.Plan, b.LastRead = u.Plan, now + if b.ClaimID != "" && !b.ClaimUntil.After(now) { + // Crash recovery: first observation after an expired claim is a new baseline. + clearEvidence(b) + b.AttemptEnd, b.ClaimID = now, "" + } + live := b.ClaimID != "" + observe(&b.FiveHour, sample(u.FiveHour, now), live) + observe(&b.Weekly, sample(u.Weekly, now), live) + _, w := target(b) + if w != nil && w.State == Started { + b.Retry, b.LastAuto = 0, time.Time{} + } + v = view(b, now) + return nil + }) + return v, err +} + +func clearEvidence(b *bucket) { + for _, w := range []*window{&b.FiveHour, &b.Weekly} { + w.Baseline = Sample{} + if w.ConfirmedReset.IsZero() { + w.State = Unknown + } + } +} + +// Begin is a short transaction, never a network-length lock. observedAt prevents +// an automatic sender from acting on a snapshot superseded by another process. +func (s Store) Begin(account, key string, automatic bool, observedAt, now time.Time) (string, error) { + idBytes := make([]byte, 16) + if _, err := rand.Read(idBytes); err != nil { + return "", err + } + id := hex.EncodeToString(idBytes) + err := s.update(account, key, func(b *bucket) error { + prune(b, now) + if b.ClaimID != "" && b.ClaimUntil.After(now) { + return ErrBusy + } + if automatic { + v := view(b, now) + _, w := target(b) + if b.ClaimID != "" || !b.LastRead.Equal(observedAt) || now.Sub(observedAt) > 3*time.Second || + w == nil || w.State != NotStarted || v.NextEligible.After(now) || + (!v.PreviousReset.IsZero() && v.PreviousReset.Add(s.ResetBuffer).After(now)) { + return ErrDeferred + } + b.Attempts = append(b.Attempts, now) + b.LastAuto = now + if b.Retry < 3 { + b.Retry++ + } + } + clearEvidence(b) + b.ClaimID, b.ClaimUntil = id, now.Add(claimDuration) + return nil + }) + return id, err +} + +func (s Store) Finish(account, key, id string, now time.Time, invalidate bool) error { + return s.update(account, key, func(b *bucket) error { + if id != b.ClaimID { + return ErrBusy + } + clearEvidence(b) + if invalidate { + b.FiveHour, b.Weekly = window{}, window{} + } + b.AttemptEnd, b.ClaimID, b.ClaimUntil = now, "", time.Time{} + return nil + }) +} + +func (s Store) Invalidate(account, key string) error { + return s.update(account, key, func(b *bucket) error { b.FiveHour, b.Weekly = window{}, window{}; return nil }) +} + +func (s Store) update(account, key string, fn func(*bucket) error) error { + if account == "" || key == "" { + return errors.New("quota account or bucket identity unavailable") + } + if err := os.MkdirAll(s.Dir, 0700); err != nil { + return err + } + hash := sha256.Sum256([]byte(account)) + path := filepath.Join(s.Dir, hex.EncodeToString(hash[:])+".json") + unlock, err := lock(path + ".lock") + if err != nil { + return err + } + defer unlock() + d := diskState{Version: 1, Buckets: map[string]*bucket{}} + data, err := os.ReadFile(path) + if err == nil { + if err := json.Unmarshal(data, &d); err != nil { + return fmt.Errorf("invalid quota state: %w", err) + } + if d.Version != 1 || d.Buckets == nil { + return errors.New("unsupported quota state") + } + } else if !os.IsNotExist(err) { + return err + } + b := d.Buckets[key] + if b == nil { + b = &bucket{} + d.Buckets[key] = b + } + if err := fn(b); err != nil { + return err + } + data, err = json.Marshal(d) + if err != nil { + return err + } + f, err := os.CreateTemp(s.Dir, ".quota-*") + if err != nil { + return err + } + defer os.Remove(f.Name()) + _, err = f.Write(data) + if err == nil { + err = f.Sync() + } + closeErr := f.Close() + if err != nil { + return err + } + if closeErr != nil { + return closeErr + } + return replace(f.Name(), path) +} diff --git a/internal/codexstate/state_test.go b/internal/codexstate/state_test.go new file mode 100644 index 0000000..f2ec954 --- /dev/null +++ b/internal/codexstate/state_test.go @@ -0,0 +1,387 @@ +package codexstate + +import ( + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + "github.com/wavever/CCLimitPing/internal/usage" +) + +var epoch = time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC) + +func quota(at, reset time.Time, used float64, weekly bool) *usage.Usage { + u := &usage.Usage{Plan: "pro", FetchedAt: at} + w := usage.Window{UsedPercent: used, ResetsAt: reset, WindowSeconds: 18000} + if weekly { + w.WindowSeconds = 604800 + u.Weekly = w + } else { + u.FiveHour = w + } + return u +} + +func TestClassify(t *testing.T) { + for _, tc := range []struct { + name string + gap, move time.Duration + used float64 + want string + }{ + {"first-minute", 59 * time.Second, 0, 0, Unknown}, + {"fixed", time.Minute, 0, 0, Started}, + {"fixed-six-minutes", 6 * time.Minute, 0, 0, Started}, + {"fixed-hours", 2 * time.Hour, 0, 0, Started}, + {"sliding", time.Minute, time.Minute, 0, NotStarted}, + {"jitter", time.Minute, time.Second, 0, Started}, + {"ambiguous", time.Minute, 30 * time.Second, 0, Unknown}, + {"stale-sliding", 11 * time.Minute, 11 * time.Minute, 0, Unknown}, + {"positive", time.Second, 0, 1, Started}, + } { + t.Run(tc.name, func(t *testing.T) { + w := &window{} + reset := epoch.Add(5 * time.Hour) + observe(w, Sample{At: epoch, Reset: reset, Seconds: 18000}, false) + got := observe(w, Sample{At: epoch.Add(tc.gap), Reset: reset.Add(tc.move), Seconds: 18000, Used: tc.used}, false) + if got.State != tc.want { + t.Fatalf("got %s want %s", got.State, tc.want) + } + }) + } +} + +func TestPollingAndInvalidation(t *testing.T) { + w := &window{} + reset := epoch.Add(5 * time.Hour) + for i := 0; i <= 60; i += 10 { + v := observe(w, Sample{At: epoch.Add(time.Duration(i) * time.Second), Reset: reset, Seconds: 18000}, false) + if i == 60 && v.State != Started { + t.Fatal(v) + } + } + for _, s := range []Sample{ + {At: epoch.Add(61 * time.Second), Reset: reset, Seconds: 604800}, + {At: epoch.Add(-time.Second), Reset: reset, Seconds: 18000}, + {At: reset.Add(time.Second), Reset: reset, Seconds: 18000}, + } { + copy := *w + if v := observe(©, s, false); v.State != Unknown { + t.Fatal(v) + } + } +} + +func read(t *testing.T, s Store, key string, u *usage.Usage) *usage.Verification { + t.Helper() + v, err := s.Observe("account", key, u) + if err != nil { + t.Fatal(err) + } + return v +} + +func TestPostAttemptBaselineAndExistingStarted(t *testing.T) { + for _, started := range []bool{false, true} { + t.Run(map[bool]string{true: "already-started", false: "new-start"}[started], func(t *testing.T) { + s := Store{Dir: t.TempDir()} + reset := epoch.Add(5 * time.Hour) + read(t, s, "codex", quota(epoch, reset, 0, false)) + if !started { + reset = reset.Add(time.Minute) + } + pre := quota(epoch.Add(time.Minute), reset, 0, false) + read(t, s, "codex", pre) + id, err := s.Begin("account", "codex", false, pre.FetchedAt, pre.FetchedAt) + if err != nil { + t.Fatal(err) + } + end := pre.FetchedAt.Add(5 * time.Second) + if err := s.Finish("account", "codex", id, end, false); err != nil { + t.Fatal(err) + } + if !started { + reset = end.Add(5 * time.Hour) + } + v := read(t, s, "codex", quota(end, reset, 0, false)) + want := Unknown + if started { + want = Started + } + if v.FiveHour.State != want { + t.Fatal(v) + } + v = read(t, s, "codex", quota(end.Add(time.Minute), reset, 0, false)) + if v.FiveHour.State != Started { + t.Fatal(v) + } + }) + } +} + +func TestBucketIdentityAndTarget(t *testing.T) { + s := Store{Dir: t.TempDir()} + for _, key := range []string{"codex", "spark:model"} { + u := quota(epoch, epoch.Add(5*time.Hour), 0, false) + u.Weekly = usage.Window{UsedPercent: 5, ResetsAt: epoch.Add(7 * 24 * time.Hour), WindowSeconds: 604800} + v := read(t, s, key, u) + if v.Target != "five_hour" || v.Recovery == "window_started" { + t.Fatal(v) + } + } + v := read(t, s, "codex", quota(epoch.Add(time.Minute), epoch.Add(5*time.Hour), 0, false)) + if v.FiveHour.State != Started { + t.Fatal(v) + } + v = read(t, s, "spark:model", quota(epoch.Add(time.Minute), epoch.Add(5*time.Hour+time.Minute), 0, false)) + if v.FiveHour.State != NotStarted { + t.Fatal(v) + } + u := quota(epoch, epoch.Add(7*24*time.Hour), 0, true) + v = read(t, s, "weekly-only", u) + if v.Target != "weekly" { + t.Fatal(v) + } + u.FetchedAt = epoch.Add(10 * time.Minute) + v = read(t, s, "weekly-only", u) + if v.Weekly.State != Started { + t.Fatal(v) + } + v, err := s.Observe("other-account", "codex", quota(epoch.Add(time.Minute), epoch.Add(5*time.Hour), 0, false)) + if err != nil || v.FiveHour.State != Unknown { + t.Fatal(v, err) + } + if err := s.Invalidate("account", "codex"); err != nil { + t.Fatal(err) + } + v = read(t, s, "codex", quota(epoch.Add(2*time.Minute), epoch.Add(5*time.Hour), 0, false)) + if v.FiveHour.State != Unknown { + t.Fatal(v) + } +} + +func TestImmediatePreSendReadRetainsSlidingEvidence(t *testing.T) { + s := Store{Dir: t.TempDir()} + for _, d := range []time.Duration{0, time.Minute, time.Minute + time.Second} { + now := epoch.Add(d) + v := read(t, s, "codex", quota(now, now.Add(5*time.Hour), 0, false)) + if d >= time.Minute && v.FiveHour.State != NotStarted { + t.Fatal(v) + } + } + now := epoch.Add(time.Minute + time.Second) + if _, err := s.Begin("account", "codex", true, now, now); err != nil { + t.Fatal(err) + } +} + +func TestJitterCannotAccumulateIntoSlidingReset(t *testing.T) { + w := &window{} + reset := epoch.Add(5 * time.Hour) + observe(w, Sample{At: epoch, Reset: reset, Seconds: 18000}, false) + observe(w, Sample{At: epoch.Add(time.Minute), Reset: reset, Seconds: 18000}, false) + for i := 1; i <= 10; i++ { + v := observe(w, Sample{At: epoch.Add(time.Minute + time.Duration(i)*time.Second), Reset: reset.Add(time.Duration(i) * time.Second), Seconds: 18000}, false) + if i > 5 && v.State == Started { + t.Fatal("jitter accumulated beyond original anchor", v) + } + } +} + +func TestBudgetAndCrashRecovery(t *testing.T) { + s := Store{Dir: t.TempDir()} + now := epoch + // Persistent transitions use fake time, including the backoff schedule. + for i := 0; i < 4; i++ { + read(t, s, "codex", quota(now, now.Add(5*time.Hour), 0, false)) + now = now.Add(time.Minute) + u := quota(now, now.Add(5*time.Hour), 0, false) + v := read(t, s, "codex", u) + if v.NextEligible.After(now) { + now = v.NextEligible + // Fresh short-interval evidence after a long backoff. + read(t, s, "codex", quota(now, now.Add(5*time.Hour), 0, false)) + now = now.Add(time.Minute) + u = quota(now, now.Add(5*time.Hour), 0, false) + read(t, s, "codex", u) + } + id, err := s.Begin("account", "codex", true, u.FetchedAt, now) + if err != nil { + t.Fatalf("attempt %d: %v", i, err) + } + if _, err := s.Begin("account", "codex", false, time.Time{}, now); !errors.Is(err, ErrBusy) { + t.Fatal(err) + } + if err := s.Finish("account", "codex", id, now, false); err != nil { + t.Fatal(err) + } + read(t, s, "codex", quota(now, now.Add(5*time.Hour), 0, false)) + } + now = now.Add(time.Minute) + u := quota(now, now.Add(5*time.Hour), 0, false) + v := read(t, Store{Dir: s.Dir}, "codex", u) + if v.Recovery != "cooldown" { + t.Fatal(v) + } + if _, err := s.Begin("account", "codex", true, u.FetchedAt, now); !errors.Is(err, ErrDeferred) { + t.Fatal(err) + } + now = epoch.Add(2 * time.Hour) + read(t, s, "codex", quota(now, now.Add(5*time.Hour), 0, false)) + now = now.Add(time.Minute) + u = quota(now, now.Add(5*time.Hour), 0, false) + read(t, s, "codex", u) + _, err := s.Begin("account", "codex", true, now, now) + if err != nil { + t.Fatal(err) + } + now = now.Add(claimDuration + time.Second) + v = read(t, s, "codex", quota(now, now.Add(5*time.Hour), 0, false)) + if v.FiveHour.State != Unknown { + t.Fatal("expired claim reused old evidence", v) + } +} + +func TestLockProcess(t *testing.T) { + if path := os.Getenv("LIMITPING_TEST_LOCK"); path != "" { + _, err := lock(path) + if !errors.Is(err, ErrBusy) { + os.Exit(2) + } + os.Exit(0) + } + path := filepath.Join(t.TempDir(), "test.lock") + unlock, err := lock(path) + if err != nil { + t.Fatal(err) + } + cmd := exec.Command(os.Args[0], "-test.run=^TestLockProcess$") + cmd.Env = append(os.Environ(), "LIMITPING_TEST_LOCK="+path) + err = cmd.Run() + unlock() + if err != nil { + t.Fatal(err) + } + unlock, err = lock(path) + if err != nil { + t.Fatal(err) + } + unlock() +} + +func TestLockReleasedAfterProcessCrash(t *testing.T) { + if path := os.Getenv("LIMITPING_TEST_CRASH_LOCK"); path != "" { + _, err := lock(path) + if err != nil { + os.Exit(2) + } + if err := os.WriteFile(path+".ready", []byte("ready"), 0600); err != nil { + os.Exit(3) + } + for { + time.Sleep(time.Hour) + } + } + path := filepath.Join(t.TempDir(), "crash.lock") + cmd := exec.Command(os.Args[0], "-test.run=^TestLockReleasedAfterProcessCrash$") + cmd.Env = append(os.Environ(), "LIMITPING_TEST_CRASH_LOCK="+path) + if err := cmd.Start(); err != nil { + t.Fatal(err) + } + defer func() { _ = cmd.Process.Kill(); _ = cmd.Wait() }() + deadline := time.Now().Add(5 * time.Second) + for { + if _, err := os.Stat(path + ".ready"); err == nil { + break + } + if time.Now().After(deadline) { + t.Fatal("child did not acquire lock") + } + time.Sleep(10 * time.Millisecond) + } + _ = cmd.Process.Kill() + _ = cmd.Wait() + unlock, err := lock(path) + if err != nil { + t.Fatal(err) + } + unlock() +} + +func TestAutomaticBudgetsAreBucketSpecific(t *testing.T) { + s := Store{Dir: t.TempDir()} + if err := s.update("account", "codex", func(b *bucket) error { b.Attempts = []time.Time{epoch, epoch, epoch, epoch}; return nil }); err != nil { + t.Fatal(err) + } + read(t, s, "spark:model", quota(epoch, epoch.Add(5*time.Hour), 0, false)) + now := epoch.Add(time.Minute) + read(t, s, "spark:model", quota(now, now.Add(5*time.Hour), 0, false)) + if _, err := s.Begin("account", "spark:model", true, now, now); err != nil { + t.Fatal("Codex budget blocked Spark", err) + } +} + +func TestResetBufferUsesPersistedPreviousBoundary(t *testing.T) { + for _, weekly := range []bool{false, true} { + for _, tc := range []struct { + name string + buffer time.Duration + known, manual bool + blocked bool + }{ + {"ten-minute-buffer", 10 * time.Minute, true, false, true}, + {"verification-covers-buffer", 35 * time.Second, true, false, false}, + {"unknown-reset", 10 * time.Minute, false, false, false}, + {"manual-bypass", 10 * time.Minute, true, true, false}, + } { + t.Run(fmt.Sprintf("%s/weekly=%t", tc.name, weekly), func(t *testing.T) { + dir := t.TempDir() + s := Store{Dir: dir, ResetBuffer: tc.buffer} + length := 5 * time.Hour + if weekly { + length = 7 * 24 * time.Hour + } + reset := epoch.Add(time.Minute) + if tc.known { + read(t, s, "codex", quota(epoch, reset, 1, weekly)) + // The expired snapshot itself must not erase the known boundary. + read(t, s, "codex", quota(reset, reset, 0, weekly)) + } + for _, d := range []time.Duration{time.Second, 61 * time.Second} { + now := reset.Add(d) + v := read(t, Store{Dir: dir}, "codex", quota(now, now.Add(length), 0, weekly)) + if tc.known && !v.PreviousReset.Equal(reset) { + t.Fatalf("boundary lost on restart: %v", v.PreviousReset) + } + if !tc.known && !v.PreviousReset.IsZero() { + t.Fatal("invented boundary") + } + } + now := reset.Add(61 * time.Second) + _, err := s.Begin("account", "codex", !tc.manual, now, now) + if tc.blocked { + if !errors.Is(err, ErrDeferred) { + t.Fatalf("early send: %v", err) + } + // Repeated moving resets do not move the buffer deadline. + for d := 2 * time.Minute; d <= 10*time.Minute; d += time.Minute { + now = reset.Add(d) + v := read(t, Store{Dir: dir}, "codex", quota(now, now.Add(length), 0, weekly)) + if !v.PreviousReset.Equal(reset) { + t.Fatal(v.PreviousReset) + } + } + if _, err := s.Begin("account", "codex", true, now, now); err != nil { + t.Fatalf("buffer counted twice: %v", err) + } + } else if err != nil { + t.Fatal(err) + } + }) + } + } +} diff --git a/internal/provider/codex.go b/internal/provider/codex.go index 7fb0c45..9f11637 100644 --- a/internal/provider/codex.go +++ b/internal/provider/codex.go @@ -52,8 +52,7 @@ const ( // via the interactive, TTY-backed Codex CLI. Headless `codex exec` can consume // tokens without anchoring the subscription-backed Codex window. type Codex struct { - cfg config.ProviderConfig - auth *auth.CodexAuth + cfg config.ProviderConfig redeemMu sync.Mutex lastRedeem time.Time // last automatic redemption attempt, for the cooldown @@ -61,8 +60,7 @@ type Codex struct { func NewCodex(cfg config.ProviderConfig) *Codex { return &Codex{ - cfg: cfg, - auth: auth.NewCodexAuth(), + cfg: cfg, } } @@ -73,23 +71,16 @@ func (c *Codex) ActiveTask(ctx context.Context) (string, bool, error) { } func (c *Codex) ReadUsage(ctx context.Context) (*usage.Usage, error) { - body, r, err := readCodexUsage(ctx, c.auth) - if err != nil { - return nil, err - } - u := codexUsageToUsage(c.Name(), body, r, r.RateLimit) - if credits, err := readCodexResetCredits(ctx, c.auth); err == nil { - u.ResetCredits = credits - } else if r.ResetCredits != nil { - // The detail endpoint is private and may go away; the usage response - // itself now embeds the available count, so keep at least that. - u.ResetCredits = &usage.ResetCredits{AvailableCount: r.ResetCredits.AvailableCount} - } - return u, nil + u, _, err := readVerifiedUsage(ctx, c.Name(), c.cfg, true) + return u, err } func (c *Codex) Trigger(ctx context.Context, dryRun bool) (*TriggerResult, error) { - return triggerCodex(ctx, c.cfg, dryRun) + return pingVerified(ctx, c.Name(), c.cfg, dryRun, nil) +} + +func (c *Codex) TriggerWithReservation(ctx context.Context, reserve PingReservation) (*TriggerResult, error) { + return pingVerified(ctx, c.Name(), c.cfg, false, reserve) } // RedeemResetCredit spends the next available reset credit right now. Each call @@ -107,6 +98,9 @@ func (c *Codex) AutoRedeemResetCredit(ctx context.Context, u *usage.Usage) (stri if !ok { return "", nil } + if u.QuotaAccount == "" { + return "", fmt.Errorf("automatic reset credit requires an identified quota observation") + } c.redeemMu.Lock() if time.Since(c.lastRedeem) < codexRedeemCooldown { c.redeemMu.Unlock() @@ -114,7 +108,7 @@ func (c *Codex) AutoRedeemResetCredit(ctx context.Context, u *usage.Usage) (stri } c.lastRedeem = time.Now() c.redeemMu.Unlock() - return c.consumeResetCredit(ctx, creditIdempotencyKey(credit)) + return c.consumeResetCreditFor(ctx, creditIdempotencyKey(credit), u.QuotaAccount) } // consumeResetCredit redeems one banked reset credit. The credit id is @@ -122,12 +116,21 @@ func (c *Codex) AutoRedeemResetCredit(ctx context.Context, u *usage.Usage) (stri // same one the policy targets — so we don't depend on an id field this private // endpoint doesn't document. func (c *Codex) consumeResetCredit(ctx context.Context, idempotencyKey string) (string, error) { + return c.consumeResetCreditFor(ctx, idempotencyKey, "") +} + +func (c *Codex) consumeResetCreditFor(ctx context.Context, idempotencyKey, expectedAccount string) (string, error) { payload, err := json.Marshal(map[string]string{"idempotency_key": idempotencyKey}) if err != nil { return "", err } - accountID, _ := c.auth.AccountID(ctx) - body, err := fetchWithAuth(ctx, c.auth, func(token string) (*http.Request, error) { + a := auth.NewCodexAuth() + accountID := "" + body, err := fetchWithAuth(ctx, a, func(token string) (*http.Request, error) { + accountID, _ = a.AccountID(ctx) + if expectedAccount != "" && accountID != expectedAccount { + return nil, fmt.Errorf("Codex account changed; automatic reset credit cancelled") + } req, err := http.NewRequestWithContext(ctx, http.MethodPost, codexConsumeURL(), bytes.NewReader(payload)) if err != nil { return nil, err @@ -155,7 +158,13 @@ func (c *Codex) consumeResetCredit(ctx context.Context, idempotencyKey string) ( if r.Code == "" { return "", fmt.Errorf("codex reset credit consume: no outcome in response: %s", truncate(body, 200)) } - return normalizeRedeemOutcome(r.Code), nil + outcome := normalizeRedeemOutcome(r.Code) + if outcome == RedeemReset { + if store, err := quotaStore(); err == nil { + _ = store.Invalidate(accountID, "codex") + } + } + return outcome, nil } // normalizeRedeemOutcome folds the two spellings of the same outcomes into the @@ -193,8 +202,7 @@ func creditIdempotencyKey(c usage.ResetCredit) string { // Spark is a separate provider backed by Codex auth and CLI transport. // Its usage window is the Spark-specific entry inside the Codex usage payload. type Spark struct { - cfg config.ProviderConfig - auth *auth.CodexAuth + cfg config.ProviderConfig } // NewSpark returns the Spark provider. It shares Codex credentials and the @@ -204,8 +212,7 @@ func NewSpark(cfg config.ProviderConfig) *Spark { cfg.Model = sparkDefaultModel } return &Spark{ - cfg: cfg, - auth: auth.NewCodexAuth(), + cfg: cfg, } } @@ -216,19 +223,16 @@ func (s *Spark) ActiveTask(ctx context.Context) (string, bool, error) { } func (s *Spark) ReadUsage(ctx context.Context) (*usage.Usage, error) { - body, r, err := readCodexUsage(ctx, s.auth) - if err != nil { - return nil, err - } - rateLimit, err := sparkRateLimitFromResponse(r, s.cfg.Model) - if err != nil { - return nil, err - } - return codexUsageToUsage(s.Name(), body, r, rateLimit), nil + u, _, err := readVerifiedUsage(ctx, s.Name(), s.cfg, false) + return u, err } func (s *Spark) Trigger(ctx context.Context, dryRun bool) (*TriggerResult, error) { - return triggerCodex(ctx, s.cfg, dryRun) + return pingVerified(ctx, s.Name(), s.cfg, dryRun, nil) +} + +func (s *Spark) TriggerWithReservation(ctx context.Context, reserve PingReservation) (*TriggerResult, error) { + return pingVerified(ctx, s.Name(), s.cfg, false, reserve) } func codexActiveTask(_ context.Context) (string, bool, error) { @@ -298,9 +302,14 @@ type codexResetCredit struct { } func readCodexUsage(ctx context.Context, auth *auth.CodexAuth) ([]byte, codexUsageResp, error) { + // Codex quota retries belong to the watcher. Immediate 401 recovery remains. + ctx = context.WithValue(ctx, noUsageRetryKey{}, true) var r codexUsageResp - accountID, _ := auth.AccountID(ctx) + if _, err := auth.Token(ctx); err != nil { + return nil, r, &AuthenticationError{Err: err} + } body, err := fetchWithAuth(ctx, auth, func(token string) (*http.Request, error) { + accountID, _ := auth.AccountID(ctx) req, err := http.NewRequestWithContext(ctx, http.MethodGet, codexUsageURL(), nil) if err != nil { return nil, err @@ -325,8 +334,8 @@ func readCodexUsage(ctx context.Context, auth *auth.CodexAuth) ([]byte, codexUsa func readCodexResetCredits(ctx context.Context, auth *auth.CodexAuth) (*usage.ResetCredits, error) { var r codexResetCreditsResp - accountID, _ := auth.AccountID(ctx) body, err := fetchWithAuth(ctx, auth, func(token string) (*http.Request, error) { + accountID, _ := auth.AccountID(ctx) req, err := http.NewRequestWithContext(ctx, http.MethodGet, codexResetCreditsURL(), nil) if err != nil { return nil, err diff --git a/internal/provider/codex_test.go b/internal/provider/codex_test.go index e36a1d6..5ab928b 100644 --- a/internal/provider/codex_test.go +++ b/internal/provider/codex_test.go @@ -22,6 +22,7 @@ func fakeCodexHome(t *testing.T) { t.Helper() home := t.TempDir() t.Setenv("CODEX_HOME", home) + t.Setenv("XDG_CONFIG_HOME", t.TempDir()) authJSON := `{"tokens":{"access_token":"access-token","refresh_token":"refresh-token","account_id":"account-123"}}` if err := os.WriteFile(filepath.Join(home, "auth.json"), []byte(authJSON), 0o600); err != nil { t.Fatal(err) @@ -476,7 +477,7 @@ func TestCodexAutoRedeemSkipsUntilExpiryAndThenThrottles(t *testing.T) { t.Fatalf("requests = %d, want 0", requests) } - expiring := &usage.Usage{ResetCredits: &usage.ResetCredits{Credits: []usage.ResetCredit{ + expiring := &usage.Usage{QuotaAccount: "account-123", ResetCredits: &usage.ResetCredits{Credits: []usage.ResetCredit{ {Status: "available", ExpiresAt: time.Now().Add(30 * time.Minute)}, }}} if outcome, err := c.AutoRedeemResetCredit(context.Background(), expiring); outcome != RedeemNothingToReset || err != nil { diff --git a/internal/provider/codex_verification.go b/internal/provider/codex_verification.go new file mode 100644 index 0000000..1889a8d --- /dev/null +++ b/internal/provider/codex_verification.go @@ -0,0 +1,186 @@ +package provider + +import ( + "context" + "errors" + "fmt" + "path/filepath" + "time" + + "github.com/wavever/CCLimitPing/internal/auth" + "github.com/wavever/CCLimitPing/internal/codexstate" + "github.com/wavever/CCLimitPing/internal/config" + "github.com/wavever/CCLimitPing/internal/usage" +) + +const quotaReadBudget = 3 * time.Second + +// PingReservation lets the watcher gate and reserve a send using fresh quota data. +// It returns a state claim, which the common ping path finishes after the CLI exits. +type PingReservation func(codexstate.Store, string, string, *usage.Usage) (string, error) + +type noQuotaStateKey struct{} +type pingStageKey struct{} + +// WithPingStage supplies optional live CLI progress without coupling providers +// to terminal output. The callback must be safe to call from a worker goroutine. +func WithPingStage(ctx context.Context, f func(string)) context.Context { + return context.WithValue(ctx, pingStageKey{}, f) +} +func pingStage(ctx context.Context, stage string) { + if f, ok := ctx.Value(pingStageKey{}).(func(string)); ok { + f(stage) + } +} + +// WithoutQuotaState preserves watch --dry-run's read-only storage behavior. +func WithoutQuotaState(ctx context.Context) context.Context { + return context.WithValue(ctx, noQuotaStateKey{}, true) +} + +func quotaStore() (codexstate.Store, error) { + dir, err := config.Dir() + return codexstate.Store{Dir: filepath.Join(dir, "state", "codex")}, err +} + +func quotaBucket(name string, cfg config.ProviderConfig) string { + if name == "spark" { + return "spark:" + normalizeCodexLimitName(cfg.Model) + } + return "codex" +} + +func currentCodexAccount(ctx context.Context) (string, error) { + a := auth.NewCodexAuth() + id, err := a.AccountID(ctx) + if err == nil && id == "" { + err = errors.New("Codex account identity unavailable") + } + return id, err +} + +func unknownVerification(warning string) *usage.Verification { + return &usage.Verification{FiveHour: usage.StartStatus{State: codexstate.Unknown}, + Weekly: usage.StartStatus{State: codexstate.Unknown}, Recovery: "verifying", Warning: warning} +} + +// Each observation owns its credential snapshot; a watcher cannot retain a +// previous login indefinitely. Authentication retries update that same snapshot. +func readVerifiedUsage(ctx context.Context, name string, cfg config.ProviderConfig, details bool) (*usage.Usage, string, error) { + a := auth.NewCodexAuth() + body, r, err := readCodexUsage(ctx, a) + if err != nil { + return nil, "", err + } + account, err := a.AccountID(ctx) + if err != nil { + return nil, "", err + } + rl := r.RateLimit + if name == "spark" { + rl, err = sparkRateLimitFromResponse(r, cfg.Model) + if err != nil { + return nil, account, err + } + } + u := codexUsageToUsage(name, body, r, rl) + u.QuotaAccount = account + if r.ResetCredits != nil { + u.ResetCredits = &usage.ResetCredits{AvailableCount: r.ResetCredits.AvailableCount} + } + store, err := quotaStore() + if ctx.Value(noQuotaStateKey{}) == true { + u.Verification = unknownVerification("") + } else if err == nil { + u.Verification, err = store.Observe(account, quotaBucket(name, cfg), u) + } + if err != nil { + u.Verification = unknownVerification("quota state unavailable: " + err.Error()) + } + if details { + if credits, err := readCodexResetCredits(ctx, a); err == nil { + if id, _ := a.AccountID(ctx); id == account { + u.ResetCredits = credits + } + } + } + return u, account, nil +} + +func boundedQuota(ctx context.Context, name string, cfg config.ProviderConfig) (*usage.Usage, string, error) { + readCtx, cancel := context.WithTimeout(ctx, quotaReadBudget) + defer cancel() + return readVerifiedUsage(readCtx, name, cfg, false) +} + +func pingVerified(ctx context.Context, name string, cfg config.ProviderConfig, dry bool, reserve PingReservation) (*TriggerResult, error) { + if dry { + return triggerCodex(ctx, cfg, true) + } + pingStage(ctx, "checking quota before ping") + pre, account, readErr := boundedQuota(ctx, name, cfg) + if ctx.Err() != nil { + return nil, ctx.Err() + } + if readErr != nil { + return nil, fmt.Errorf("ping not sent: quota precheck failed: %w", readErr) + } + warning := pre.Verification.Warning + identity, identityErr := currentCodexAccount(ctx) + if identityErr != nil { + return nil, fmt.Errorf("ping not sent: account check failed: %w", identityErr) + } + if identity != account { + return nil, fmt.Errorf("ping not sent: Codex account changed during quota precheck; retry") + } + store, storeErr := quotaStore() + key := quotaBucket(name, cfg) + claim := "" + if storeErr == nil { + if reserve != nil { + claim, storeErr = reserve(store, account, key, pre) + } else { + claim, storeErr = store.Begin(account, key, false, pre.FetchedAt, time.Now()) + } + } + if storeErr != nil { + if reserve != nil || errors.Is(storeErr, codexstate.ErrBusy) || errors.Is(storeErr, codexstate.ErrDeferred) { + return nil, fmt.Errorf("ping not sent: %w", storeErr) + } + warning = "quota coordination unavailable: " + storeErr.Error() + } + // The PTY deadline and claim deadline must describe the same bounded operation. + pingStage(ctx, "sending ping") + triggerCtx, cancel := context.WithTimeout(ctx, 3*time.Minute) + res, triggerErr := triggerCodex(triggerCtx, cfg, false) + cancel() + if res == nil { + res = &TriggerResult{} + } + res.StatusEnabled = cfg.Enabled + res.PreVerification = pre.Verification + after, afterErr := currentCodexAccount(ctx) + changed := account == "" || afterErr != nil || after != account + if claim != "" { + if err := store.Finish(account, key, claim, time.Now(), changed); err != nil { + warning = "quota state update failed: " + err.Error() + } + } + res.Verification = unknownVerification(warning) + if ctx.Err() == nil { + pingStage(ctx, "checking quota after ping") + post, postAccount, err := boundedQuota(ctx, name, cfg) + if err != nil { + res.PostcheckErr = err + res.Verification.Warning = "post-ping quota read failed: " + err.Error() + } else if !changed && postAccount == account { + res.Verification = post.Verification + if warning != "" { + res.Verification.Warning = warning + } + } else { + res.Verification.Warning = "account changed; ping result cannot be attributed to this quota" + } + } + return res, triggerErr +} diff --git a/internal/provider/codex_verification_test.go b/internal/provider/codex_verification_test.go new file mode 100644 index 0000000..6a05682 --- /dev/null +++ b/internal/provider/codex_verification_test.go @@ -0,0 +1,276 @@ +package provider + +import ( + "context" + "errors" + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "runtime" + "strings" + "testing" + "time" + + "github.com/wavever/CCLimitPing/internal/codexstate" + "github.com/wavever/CCLimitPing/internal/config" + "github.com/wavever/CCLimitPing/internal/usage" +) + +func fakeCodexCLI(t *testing.T, script string) { + t.Helper() + if runtime.GOOS == "windows" { + t.Skip("native Windows PTY is unsupported") + } + dir := t.TempDir() + if err := os.WriteFile(filepath.Join(dir, "codex"), []byte("#!/bin/sh\n"+script+"\n"), 0700); err != nil { + t.Fatal(err) + } + t.Setenv("PATH", dir+string(os.PathListSeparator)+os.Getenv("PATH")) +} + +func quotaResponse(reset int64) string { + return fmt.Sprintf(`{"plan_type":"pro","rate_limit":{"primary_window":{"used_percent":0,"limit_window_seconds":604800,"reset_at":%d}}}`, reset) +} + +func TestVerifiedPingReadsBeforeAndAfterWithoutWaitingMinute(t *testing.T) { + fakeCodexHome(t) + fakeCodexCLI(t, "exit 0") + old := usageHTTPClient + defer func() { usageHTTPClient = old }() + reads := 0 + usageHTTPClient = &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) { + if !strings.HasSuffix(req.URL.Path, "/usage") { + t.Fatalf("unexpected detail read %s", req.URL.Path) + } + reads++ + return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(quotaResponse(time.Now().Add(7 * 24 * time.Hour).Unix()))), Header: make(http.Header)}, nil + })} + start := time.Now() + res, err := NewCodex(config.ProviderConfig{Enabled: true}).Trigger(context.Background(), false) + if err != nil { + t.Fatal(err) + } + if reads != 2 { + t.Fatalf("reads=%d", reads) + } + if time.Since(start) > 5*time.Second { + t.Fatal("manual ping blocked") + } + if res.Verification == nil || res.Verification.Weekly.State != codexstate.Unknown || res.Verification.Weekly.DueAt.IsZero() { + t.Fatalf("%+v", res.Verification) + } +} + +func TestFailedTriggerStillReadsQuotaAfterwards(t *testing.T) { + fakeCodexHome(t) + fakeCodexCLI(t, "exit 1") + old := usageHTTPClient + defer func() { usageHTTPClient = old }() + reads := 0 + usageHTTPClient = &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) { + reads++ + return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(quotaResponse(time.Now().Add(7 * 24 * time.Hour).Unix()))), Header: make(http.Header)}, nil + })} + res, err := NewCodex(config.ProviderConfig{}).Trigger(context.Background(), false) + if err == nil || reads != 2 || res.Verification == nil { + t.Fatal(err, reads, res) + } +} + +func TestAccountSwitchReloadsIdentity(t *testing.T) { + fakeCodexHome(t) + old := usageHTTPClient + defer func() { usageHTTPClient = old }() + expected := "account-123" + usageHTTPClient = &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) { + if req.Header.Get("ChatGPT-Account-Id") != expected { + t.Fatalf("stale account header: %s", req.Header.Get("ChatGPT-Account-Id")) + } + body := quotaResponse(time.Now().Add(7 * 24 * time.Hour).Unix()) + if strings.HasSuffix(req.URL.Path, "reset-credits") { + body = `{"available_count":0}` + } + return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(body)), Header: make(http.Header)}, nil + })} + c := NewCodex(config.ProviderConfig{}) + if _, err := c.ReadUsage(context.Background()); err != nil { + t.Fatal(err) + } + expected = "account-new" + if err := os.WriteFile(filepath.Join(os.Getenv("CODEX_HOME"), "auth.json"), []byte(`{"tokens":{"access_token":"different-valid-token","account_id":"account-new"}}`), 0600); err != nil { + t.Fatal(err) + } + u, err := c.ReadUsage(context.Background()) + if err != nil { + t.Fatal(err) + } + if u.Verification.Weekly.State != codexstate.Unknown { + t.Fatal(u.Verification) + } +} + +func TestDryRunNeverReadsOrWritesState(t *testing.T) { + t.Setenv("CODEX_HOME", t.TempDir()) + dir := t.TempDir() + t.Setenv("XDG_CONFIG_HOME", dir) + old := usageHTTPClient + defer func() { usageHTTPClient = old }() + usageHTTPClient = &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) { t.Fatal("dry run read usage"); return nil, nil })} + if _, err := NewCodex(config.ProviderConfig{}).Trigger(context.Background(), true); err != nil { + t.Fatal(err) + } + entries, _ := os.ReadDir(dir) + if len(entries) != 0 { + t.Fatal(entries) + } +} + +func TestQuotaPrecheckFailureNeverSends(t *testing.T) { + for _, status := range []int{403, 503} { + for _, automatic := range []bool{false, true} { + t.Run(fmt.Sprint(status, automatic), func(t *testing.T) { + fakeCodexHome(t) + marker := filepath.Join(t.TempDir(), "sent") + t.Setenv("TEST_SENT", marker) + fakeCodexCLI(t, `touch "$TEST_SENT"`) + reads := 0 + useTransport(t, func(*http.Request) (*http.Response, error) { + reads++ + return &http.Response{StatusCode: status, Body: io.NopCloser(strings.NewReader("{}")), Header: http.Header{"Retry-After": []string{"120"}}}, nil + }) + var reserve PingReservation + if automatic { + reserve = func(codexstate.Store, string, string, *usage.Usage) (string, error) { + t.Fatal("failed precheck must not reserve a ping") + return "", nil + } + } + start := time.Now() + res, err := pingVerified(context.Background(), "codex", config.ProviderConfig{}, false, reserve) + var httpErr *UsageHTTPError + if res != nil || !errors.As(err, &httpErr) || httpErr.StatusCode != status || !strings.Contains(err.Error(), "ping not sent: quota precheck failed") { + t.Fatal(res, err) + } + if reads != 1 || time.Since(start) > 5*time.Second { + t.Fatal("precheck retried or waited", reads) + } + if !httpErr.RetryAfter.After(start.Add(time.Minute)) { + t.Fatal("Retry-After was lost") + } + if _, err := os.Stat(marker); !os.IsNotExist(err) { + t.Fatal("CLI was executed", err) + } + }) + } + } +} + +func TestQuotaNetworkFailureDoesNotRetry(t *testing.T) { + fakeCodexHome(t) + reads := 0 + useTransport(t, func(*http.Request) (*http.Response, error) { + reads++ + return nil, io.ErrUnexpectedEOF + }) + res, err := NewCodex(config.ProviderConfig{}).Trigger(context.Background(), false) + if res != nil || !errors.Is(err, io.ErrUnexpectedEOF) || reads != 1 { + t.Fatal(res, err, reads) + } +} + +func TestPostcheckFailureReportsSentWithoutRetry(t *testing.T) { + fakeCodexHome(t) + fakeCodexCLI(t, `printf '\033]9;done\007'`) + reads := 0 + useTransport(t, func(*http.Request) (*http.Response, error) { + reads++ + status, body := 200, quotaResponse(time.Now().Add(7*24*time.Hour).Unix()) + if reads > 1 { + status, body = 503, "{}" + } + return &http.Response{StatusCode: status, Body: io.NopCloser(strings.NewReader(body)), Header: http.Header{"Retry-After": []string{"1200"}}}, nil + }) + res, err := NewCodex(config.ProviderConfig{}).Trigger(context.Background(), false) + if err != nil || reads != 2 || res == nil || !strings.Contains(res.Verification.Warning, "post-ping quota read failed") { + t.Fatal(res, err, reads) + } + var postErr *UsageHTTPError + if res.PreVerification == nil || res.PreVerification == res.Verification { + t.Fatal("pre-ping verification was not retained separately") + } + if !errors.As(res.PostcheckErr, &postErr) || postErr.StatusCode != 503 || !postErr.RetryAfter.After(time.Now().Add(19*time.Minute)) { + t.Fatal("postcheck retry metadata lost", res.PostcheckErr) + } +} + +func TestMissingCredentialsIsAuthenticationFailure(t *testing.T) { + t.Setenv("CODEX_HOME", t.TempDir()) + res, err := NewCodex(config.ProviderConfig{}).Trigger(context.Background(), false) + var authErr *AuthenticationError + if res != nil || !errors.As(err, &authErr) || !strings.Contains(err.Error(), "ping not sent") { + t.Fatal(res, err) + } +} + +func TestPrecheckReloadsCredentialsOn401(t *testing.T) { + fakeCodexHome(t) + fakeCodexCLI(t, `printf '\033]9;done\007'`) + reads := 0 + useTransport(t, func(req *http.Request) (*http.Response, error) { + reads++ + if reads == 1 { + if err := os.WriteFile(filepath.Join(os.Getenv("CODEX_HOME"), "auth.json"), []byte(`{"tokens":{"access_token":"new-token","account_id":"account-123"}}`), 0600); err != nil { + t.Fatal(err) + } + return &http.Response{StatusCode: 401, Body: io.NopCloser(strings.NewReader("{}")), Header: make(http.Header)}, nil + } + if req.Header.Get("Authorization") != "Bearer new-token" { + t.Fatal("stale token") + } + return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(quotaResponse(time.Now().Add(7 * 24 * time.Hour).Unix()))), Header: make(http.Header)}, nil + }) + if _, err := NewCodex(config.ProviderConfig{}).Trigger(context.Background(), false); err != nil || reads != 3 { + t.Fatal(err, reads) + } +} + +func TestAutoRedeemRejectsChangedObservationAccount(t *testing.T) { + fakeCodexHome(t) + requests := 0 + useTransport(t, func(req *http.Request) (*http.Response, error) { + requests++ + t.Fatal("must not spend another account's credit") + return nil, nil + }) + u := &usage.Usage{QuotaAccount: "previous-account", ResetCredits: &usage.ResetCredits{Credits: []usage.ResetCredit{ + {Status: "available", ExpiresAt: time.Now().Add(30 * time.Minute)}, + }}} + _, err := NewCodex(config.ProviderConfig{}).AutoRedeemResetCredit(context.Background(), u) + if err == nil || requests != 0 { + t.Fatal(err, requests) + } +} + +func TestAutoRedeemRechecksIdentityOnAuthenticationRetry(t *testing.T) { + fakeCodexHome(t) + requests := 0 + useTransport(t, func(req *http.Request) (*http.Response, error) { + requests++ + if requests > 1 { + t.Fatal("retried redemption against changed account") + } + if err := os.WriteFile(filepath.Join(os.Getenv("CODEX_HOME"), "auth.json"), []byte(`{"tokens":{"access_token":"new-token","account_id":"new-account"}}`), 0600); err != nil { + t.Fatal(err) + } + return &http.Response{StatusCode: 401, Body: io.NopCloser(strings.NewReader(`{}`)), Header: make(http.Header)}, nil + }) + u := &usage.Usage{QuotaAccount: "account-123", ResetCredits: &usage.ResetCredits{Credits: []usage.ResetCredit{ + {Status: "available", ExpiresAt: time.Now().Add(30 * time.Minute)}, + }}} + _, err := NewCodex(config.ProviderConfig{}).AutoRedeemResetCredit(context.Background(), u) + if err == nil || requests != 1 { + t.Fatal(err, requests) + } +} diff --git a/internal/provider/provider.go b/internal/provider/provider.go index 24f6047..09100b9 100644 --- a/internal/provider/provider.go +++ b/internal/provider/provider.go @@ -87,14 +87,29 @@ type ResetCreditRedeemer interface { // consumed (parsed from the CLI's machine-readable output). CostUSD is 0 when // the provider doesn't report a cost (e.g. Codex). type TriggerResult struct { - Command string - HasUsage bool - InputTokens int - OutputTokens int - TotalTokens int - CostUSD float64 + Command string + HasUsage bool + InputTokens int + OutputTokens int + TotalTokens int + CostUSD float64 + Verification *usage.Verification + PreVerification *usage.Verification + PostcheckErr error // quota-read failure, independent of the CLI outcome + StatusEnabled bool } +// VerifiedTrigger uses the same ping path with a watcher-owned pre-send reservation. +type VerifiedTrigger interface { + TriggerWithReservation(context.Context, PingReservation) (*TriggerResult, error) +} + +// AuthenticationError preserves credential failures for watcher retry policy. +type AuthenticationError struct{ Err error } + +func (e *AuthenticationError) Error() string { return e.Err.Error() } +func (e *AuthenticationError) Unwrap() error { return e.Err } + // UsageHTTPError preserves usage endpoint HTTP failures so callers can make // status-aware scheduling decisions instead of treating every failure alike. type UsageHTTPError struct { @@ -142,7 +157,7 @@ func fetchWithAuth(ctx context.Context, src tokenSource, buildReq func(token str if status == http.StatusUnauthorized { t, rerr := src.Refresh(ctx) if rerr != nil { - return nil, fmt.Errorf("unauthorized and refresh failed: %w", rerr) + return nil, &AuthenticationError{Err: fmt.Errorf("unauthorized and refresh failed: %w", rerr)} } token = t if body, status, header, err = doGet(ctx, token, buildReq); err != nil { @@ -207,8 +222,10 @@ func doGet(ctx context.Context, token string, buildReq func(token string) (*http return lastBody, lastStatus, lastHeader, lastErr } +type noUsageRetryKey struct{} + func shouldRetryUsageGET(ctx context.Context, attempt, status int, err error) bool { - if attempt >= usageGETAttempts || ctx.Err() != nil { + if ctx.Value(noUsageRetryKey{}) == true || attempt >= usageGETAttempts || ctx.Err() != nil { return false } if err != nil { diff --git a/internal/scheduler/codex.go b/internal/scheduler/codex.go new file mode 100644 index 0000000..bc17e96 --- /dev/null +++ b/internal/scheduler/codex.go @@ -0,0 +1,197 @@ +package scheduler + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/wavever/CCLimitPing/internal/codexstate" + "github.com/wavever/CCLimitPing/internal/provider" + "github.com/wavever/CCLimitPing/internal/usage" +) + +// runVerifiedTarget is shared by Codex and Spark, never by Claude. A successful +// transport does not advance a quota schedule without observation evidence. +func (s *Scheduler) runVerifiedTarget(ctx context.Context, t Target, p provider.VerifiedTrigger) { + name := t.Provider.Name() + var reads, prechecks quotaRetry + aligned := t.AlignStart.IsZero() + wait := func(reason string, d time.Duration) bool { + if d <= 0 { + d = time.Second + } + s.live.set(name, reason, time.Now().Add(d)) + return sleepCtx(ctx, d) + } + for ctx.Err() == nil { + s.live.set(name, "checking quota…", time.Time{}) + rctx, cancel := context.WithTimeout(ctx, readTimeout) + u, err := t.Provider.ReadUsage(rctx) + cancel() + if err != nil { + d := reads.next(err, time.Now()) + s.log.Printf("[%s] quota read failed: %v (retry in %s)", name, err, d) + if !wait("quota read failed", d) { + return + } + continue + } + reads = quotaRetry{} + if s.redeemExpiringCredit(ctx, t, u) { + continue + } + if s.weeklyExhausted(u) { + d := u.Weekly.Remaining() + if d <= 0 { + // A stale or missing reset must not turn quota reads into a tight loop. + d = 5 * time.Minute + } + if d > 5*time.Minute { + d = 5 * time.Minute + } + if !wait("weekly limit reached", d) { + return + } + continue + } + v := u.Verification + if v == nil || v.Warning != "" { + if v != nil { + s.log.Printf("[%s] %s; automatic ping deferred", name, v.Warning) + } + if !wait("quota state unavailable", codexstate.Interval) { + return + } + continue + } + if !aligned { + aligned = true + if d := time.Until(t.AlignStart); d > 0 { + if !wait("waiting for align_start", d) { + return + } + continue // re-read after alignment; user activity may have started it + } + } + if v.Recovery == "window_started" { + prechecks = quotaRetry{} + // Observe at rollover; the verification interval counts toward the buffer. + d := time.Until(v.NextEligible) + if d > 5*time.Minute { + d = 5 * time.Minute + } + if !wait(v.Target+" started", d) { + return + } + continue + } + if v.Recovery != "ready" || v.NextEligible.After(time.Now()) { + d := time.Until(v.NextEligible) + if d <= 0 || d > codexstate.Interval { + d = codexstate.Interval + } + if !wait(v.Target+" "+v.Recovery, d) { + return + } + continue + } + if !v.PreviousReset.IsZero() { + if d := time.Until(v.PreviousReset.Add(s.cfg.ResetBuffer.Duration)); d > 0 { + if d > 5*time.Minute { + d = 5 * time.Minute + } + if !wait("waiting for reset_buffer", d) { + return + } + continue // polling and verification count toward the same fixed deadline + } + } + if desc, active, err := activeProviderTask(ctx, t.Provider); err != nil || active { + if !wait(desc+" active or activity unavailable", activeTaskPoll) { + return + } + continue + } + s.live.set(name, "checking and sending ping…", time.Time{}) + res, err := p.TriggerWithReservation(ctx, func(store codexstate.Store, account, key string, pre *usage.Usage) (string, error) { + if pre.Verification == nil || pre.Verification.Warning != "" { + return "", fmt.Errorf("quota state unavailable") + } + if pre.WeeklyExhausted(s.cfg.WeeklyThreshold) { + return "", codexstate.ErrDeferred + } + if _, active, err := activeProviderTask(ctx, t.Provider); err != nil || active { + return "", codexstate.ErrDeferred + } + store.ResetBuffer = s.cfg.ResetBuffer.Duration + return store.Begin(account, key, true, pre.FetchedAt, time.Now()) + }) + if errors.Is(err, codexstate.ErrBusy) || errors.Is(err, codexstate.ErrDeferred) { + prechecks = quotaRetry{} + if !wait("ping deferred", codexstate.Interval) { + return + } + continue + } + if err != nil && res == nil { + // No trigger occurred: do not manufacture a failed ping-history entry. + d := prechecks.next(err, time.Now()) + s.log.Printf("[%s] pre-ping check failed: %v; observing again in %s", name, err, d) + if !wait("pre-ping check unavailable", d) { + return + } + continue + } + prechecks = quotaRetry{} + if err != nil { + s.log.Printf("[%s] ping failed: %v; verifying quota before retry", name, err) + s.notify(name+": ping failed", "Verifying quota before another attempt") + } else { + s.log.Printf("[%s] ping trigger returned; checking window%s", name, triggerCost(res)) + s.notify(name+": CLI trigger returned", "Turn completion is unverified; checking quota separately") + } + if res != nil && res.Verification != nil && res.Verification.Warning != "" { + s.log.Printf("[%s] %s", name, res.Verification.Warning) + } + delay := codexstate.Interval + if res != nil && res.PostcheckErr != nil { + delay = reads.next(res.PostcheckErr, time.Now()) + } + if !wait("verifying quota", delay) { + return + } + } +} + +type quotaRetry struct { + delay time.Duration + auth bool +} + +func (r *quotaRetry) next(err error, now time.Time) time.Duration { + var httpErr *provider.UsageHTTPError + var authErr *provider.AuthenticationError + auth := errors.As(err, &authErr) + if errors.As(err, &httpErr) { + auth = auth || httpErr.StatusCode == 401 || httpErr.StatusCode == 403 + if !httpErr.RetryAfter.IsZero() || httpErr.StatusCode == 429 { + *r = quotaRetry{} + return usageRateLimitWait(httpErr.RetryAfter, now) + } + } + if r.delay == 0 || r.auth != auth { + r.delay = minBackoff + } else { + r.delay *= 2 + } + r.auth = auth + cap := maxBackoff + if auth { + cap = time.Hour + } + if r.delay > cap { + r.delay = cap + } + return r.delay +} diff --git a/internal/scheduler/codex_test.go b/internal/scheduler/codex_test.go new file mode 100644 index 0000000..36b9bb0 --- /dev/null +++ b/internal/scheduler/codex_test.go @@ -0,0 +1,193 @@ +package scheduler + +import ( + "context" + "errors" + "io" + "testing" + "time" + + "github.com/wavever/CCLimitPing/internal/provider" + "github.com/wavever/CCLimitPing/internal/usage" +) + +func TestQuotaRetryPolicy(t *testing.T) { + for _, tc := range []struct { + name string + err error + cap time.Duration + }{ + {"unauthorized", &provider.UsageHTTPError{StatusCode: 401}, time.Hour}, + {"forbidden", &provider.UsageHTTPError{StatusCode: 403}, time.Hour}, + {"credentials", &provider.AuthenticationError{Err: errors.New("missing credentials")}, time.Hour}, + {"server", &provider.UsageHTTPError{StatusCode: 503}, 10 * time.Minute}, + {"network", errors.New("connection failed"), 10 * time.Minute}, + } { + t.Run(tc.name, func(t *testing.T) { + var retry quotaRetry + want := 30 * time.Second + for i := 0; i < 12; i++ { + if got := retry.next(tc.err, time.Now()); got != want { + t.Fatalf("step %d: %s, want %s", i, got, want) + } + want *= 2 + if want > tc.cap { + want = tc.cap + } + } + }) + } + now := time.Now() + var retry quotaRetry + if got := retry.next(&provider.UsageHTTPError{StatusCode: 429, RetryAfter: now.Add(20 * time.Minute)}, now); got != 20*time.Minute { + t.Fatal(got) + } + if got := retry.next(&provider.UsageHTTPError{StatusCode: 429}, now); got != 5*time.Minute { + t.Fatal(got) + } + if got := retry.next(&provider.UsageHTTPError{StatusCode: 403}, now); got != 30*time.Second { + t.Fatal(got) + } +} + +type verifiedStub struct { + stubProvider + postcheckErr error +} + +func (p *verifiedStub) TriggerWithReservation(ctx context.Context, _ provider.PingReservation) (*provider.TriggerResult, error) { + res, err := p.Trigger(ctx, false) + if res != nil { + res.PostcheckErr = p.postcheckErr + } + return res, err +} + +func TestPostcheckFailureSchedulesQuotaRetry(t *testing.T) { + for _, tc := range []struct { + name string + status int + retryAfter, want time.Duration + }{ + {"rate-limit-deadline", 429, 20 * time.Minute, 20 * time.Minute}, + {"server-deadline", 503, time.Hour, time.Hour}, + {"rate-limit-fallback", 429, 0, 5 * time.Minute}, + {"permission-backoff", 403, 0, 30 * time.Second}, + } { + t.Run(tc.name, func(t *testing.T) { + start := time.Now() + err := &provider.UsageHTTPError{StatusCode: tc.status} + if tc.retryAfter > 0 { + err.RetryAfter = start.Add(tc.retryAfter) + } + p := &verifiedStub{stubProvider: stubProvider{usage: &usage.Usage{ + Verification: &usage.Verification{Target: "weekly", Recovery: "ready"}, + }}, postcheckErr: err} + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + s := New(testConfig(), []Target{{Provider: p}}, false, false, io.Discard) + s.live.enabled = true + s.runVerifiedTarget(ctx, Target{Provider: p}, p) + item := s.live.items[p.Name()] + if item.state != "verifying quota" || item.deadline.Sub(start.Add(tc.want)).Abs() > time.Second { + t.Fatalf("next read: %+v, want delay %s", item, tc.want) + } + if reads, sends := p.counts(); reads != 1 || sends != 1 { + t.Fatal(reads, sends) + } + }) + } +} + +func TestVerifiedSchedulerGates(t *testing.T) { + for _, tc := range []struct { + name, recovery string + warn, active bool + weekly float64 + want int + }{ + {"unknown", "verifying", false, false, 0, 0}, + {"ready", "ready", false, false, 0, 1}, + {"cooldown", "cooldown", false, false, 0, 0}, + {"storage-error", "ready", true, false, 0, 0}, + {"active-user", "ready", false, true, 0, 0}, + {"weekly-guard", "ready", false, false, 100, 0}, + } { + t.Run(tc.name, func(t *testing.T) { + v := &usage.Verification{Target: "five_hour", Recovery: tc.recovery} + if tc.warn { + v.Warning = "cannot write state" + } + p := &verifiedStub{stubProvider: stubProvider{active: tc.active, usage: &usage.Usage{ + Weekly: usage.Window{UsedPercent: tc.weekly, ResetsAt: time.Now().Add(time.Hour)}, Verification: v}}} + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + s := New(testConfig(), []Target{{Provider: p}}, false, false, io.Discard) + s.Run(ctx) + _, n := p.counts() + if n != tc.want { + t.Fatalf("triggers=%d", n) + } + }) + } +} + +func TestVerifiedWeeklyExhaustionPolling(t *testing.T) { + for _, tc := range []struct { + name string + reset time.Time + want time.Duration + }{ + {"missing-reset", time.Time{}, 5 * time.Minute}, + {"stale-reset", time.Now().Add(-time.Minute), 5 * time.Minute}, + {"near-reset", time.Now().Add(30 * time.Second), 30 * time.Second}, + {"distant-reset", time.Now().Add(time.Hour), 5 * time.Minute}, + } { + t.Run(tc.name, func(t *testing.T) { + p := &verifiedStub{stubProvider: stubProvider{usage: &usage.Usage{ + Weekly: usage.Window{UsedPercent: 100, ResetsAt: tc.reset}, + }}} + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + s := New(testConfig(), []Target{{Provider: p}}, false, false, io.Discard) + s.live.enabled = true + start := time.Now() + s.runVerifiedTarget(ctx, Target{Provider: p}, p) + item := s.live.items[p.Name()] + if item.state != "weekly limit reached" || item.deadline.Sub(start.Add(tc.want)).Abs() > time.Second { + t.Fatalf("next read: %+v, want delay %s", item, tc.want) + } + if reads, sends := p.counts(); reads != 1 || sends != 0 { + t.Fatalf("reads=%d sends=%d", reads, sends) + } + }) + } +} + +func TestVerifiedSchedulerHonorsResetBuffer(t *testing.T) { + for _, tc := range []struct { + name string + previousReset time.Time + buffer time.Duration + want int + }{ + {"known-boundary", time.Now().Add(-time.Minute), 10 * time.Minute, 0}, + {"short-buffer", time.Now().Add(-time.Minute), 35 * time.Second, 1}, + {"unknown-boundary", time.Time{}, 10 * time.Minute, 1}, + } { + t.Run(tc.name, func(t *testing.T) { + p := &verifiedStub{stubProvider: stubProvider{usage: &usage.Usage{Verification: &usage.Verification{ + Target: "weekly", Recovery: "ready", PreviousReset: tc.previousReset, + }}}} + cfg := testConfig() + cfg.ResetBuffer.Duration = tc.buffer + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + New(cfg, []Target{{Provider: p}}, false, false, io.Discard).Run(ctx) + _, n := p.counts() + if n != tc.want { + t.Fatalf("triggers=%d want %d", n, tc.want) + } + }) + } +} diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index 6b6de14..63f7258 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -100,6 +100,10 @@ func (s *Scheduler) Run(ctx context.Context) { } func (s *Scheduler) runTarget(ctx context.Context, t Target) { + if p, ok := t.Provider.(provider.VerifiedTrigger); ok && !s.dryRun { + s.runVerifiedTarget(ctx, t, p) + return + } name := t.Provider.Name() backoff := minBackoff aligned := t.AlignStart.IsZero() // whether the align gate has been passed @@ -112,6 +116,9 @@ func (s *Scheduler) runTarget(ctx context.Context, t Target) { s.live.set(name, "checking usage…", time.Time{}) rctx, cancel := context.WithTimeout(ctx, readTimeout) + if s.dryRun { + rctx = provider.WithoutQuotaState(rctx) + } u, err := t.Provider.ReadUsage(rctx) cancel() if err != nil { diff --git a/internal/usage/usage.go b/internal/usage/usage.go index aafd8f2..0150859 100644 --- a/internal/usage/usage.go +++ b/internal/usage/usage.go @@ -80,7 +80,25 @@ type Usage struct { ResetCredits *ResetCredits LimitReached bool FetchedAt time.Time - Raw []byte // raw JSON body, for `status -v` + Raw []byte // raw JSON body, for `status -v` + Verification *Verification // Codex-backed providers only; not a shared window predicate + QuotaAccount string // identity of the quota request; never part of status JSON +} + +// Verification separates a completed CLI request from an observed quota window. +type Verification struct { + FiveHour StartStatus + Weekly StartStatus + Target string + Recovery string + NextEligible time.Time + PreviousReset time.Time // known reset boundary of the target's previous window + Warning string +} + +type StartStatus struct { + State string + DueAt time.Time } // Reset-credit auto-redeem policy. A credit that lapses unused is worth