diff --git a/docs-site/src/content/docs/guides/codex-integration.md b/docs-site/src/content/docs/guides/codex-integration.md index 28c40eddcb0..2a487676341 100644 --- a/docs-site/src/content/docs/guides/codex-integration.md +++ b/docs-site/src/content/docs/guides/codex-integration.md @@ -911,6 +911,15 @@ it is not; `ocx doctor` reports restart safety (service/shim coverage). ## Routed models during Codex reserve mode +Codex Pool can optionally protect stored pool accounts at a selected 5-hour or weekly usage +threshold. The Desktop/main account keeps its separate 98% hard lock. +Set `codexPool.lowQuotaProtection` in configuration to pause accounts, record a log-and-API +alert, or both; see [routing configuration](/reference/configuration/routing/#codex-pool-low-quota-protection). +A pause takes effect for the next selection immediately, while saving it to disk is deferred. +Check this server’s authenticated `GET /api/codex-auth/low-quota-events` history for `logged` +alerts or save failures. Manual resume remains in force for the current quota episode. The +default alert reaches only the log and API; it does not produce a desktop or OS notification. + When the ChatGPT 5-hour quota is exhausted, Codex may offer a reserve fallback model (`gpt-reserve` / Luna Reserve). While that state is active, the Codex model picker can make **every other entry unselectable — including opencodex routed models**, even though those diff --git a/docs-site/src/content/docs/reference/configuration/routing.md b/docs-site/src/content/docs/reference/configuration/routing.md index ecc0ea666c5..016eb74884f 100644 --- a/docs-site/src/content/docs/reference/configuration/routing.md +++ b/docs-site/src/content/docs/reference/configuration/routing.md @@ -13,6 +13,39 @@ Routing turns the model id sent by a client into one concrete provider and upstr | `combos?` | `Record` | `{}` | Virtual `combo/` models built from ordered provider/model targets. | | `routingProfiles?` | `Record` | `{}` | Virtual `policy/` models that select among an explicit candidate allowlist using hard capability requirements and deterministic scoring. | +### Codex Pool low-quota protection + +`codexPool.lowQuotaProtection` applies only to stored Codex Pool accounts. The Desktop/main +account keeps its separate 98% hard lock. The optional shape is: + +```json +{"codexPool":{"lowQuotaProtection":{"enabled":true,"threshold":80,"actions":{"pause":true,"notify":true},"windows":{"short":true,"weekly":true}}}} +``` + +Absence or `enabled: false` disables it. `threshold` is a finite inclusive percentage from 1 +through 100. An enabled policy needs at least one true action and one true window. A selected +5-hour or weekly window triggers at `usage >= threshold`; either selected window is enough. +Only fresh accepted observations qualify: there is no extra polling, and credits-only, +cached, expired, monthly, custom-window-only, and raw out-of-range usage observations are ignored +for this policy. Display bars may still clamp invalid upstream percentages. The policy is independent +of proactive account switching. + +With `pause`, the account leaves Pool selection in memory before the next request. The server +coalesces a config save after the observation turn and retries failed saves for a bounded time. +Normal shutdown waits briefly for pending saves. A timed-out save that is already running remains +`pending` until it actually succeeds or fails; a queued save is cancelled. A failed save is visible +in the event history. A restart before a successful save cannot preserve the pause. In-flight requests keep their +captured account. Manual resume suppresses repause across all currently high windows for that +account until a below-threshold reading or a new reset boundary re-arms a window. A reset never resumes an account automatically. + +With `notify`, the server writes a local log line with window and percentage but no account id +and records a bounded event. Authenticated `GET /api/codex-auth/low-quota-events` exposes only +that server’s events, including account id, window, usage percentage, known reset time, timestamp, +and status. The default status is `logged`: the alert reached the log and event history only. +There is no desktop or OS notification. Each account/window logs once per server-local episode. +If an injected notification sink fails, its failure is recorded and a later eligible observation +can retry it; only a successful sink is marked `delivered`. + ## Model resolution order opencodex resolves the requested model in this order: diff --git a/docs-site/src/content/docs/reference/management-api.md b/docs-site/src/content/docs/reference/management-api.md index 4665a23e79a..95cce7f26c3 100644 --- a/docs-site/src/content/docs/reference/management-api.md +++ b/docs-site/src/content/docs/reference/management-api.md @@ -683,6 +683,7 @@ manager. Its routes are: | `PUT, PATCH /api/codex-auth/pool-strategy` | Update Codex account-pool selection strategy | 400 invalid strategy/config | | `PUT /api/codex-auth/failover` | Set the account failover threshold | 400 invalid threshold | | `GET /api/codex-auth/quota` | Read cached quota state by account | — | +| `GET /api/codex-auth/low-quota-events?limit=20` | Read only this server’s last 0–100 low-quota log/notice and pause-save events (default 20); includes account id and status (`logged` for the default log-only alert; `delivered` for a successful injected notice sink; `succeeded` for a completed pause save) | 400 invalid limit; management authentication required | | `GET /api/codex-auth/reset-credits` | Inspect reset-credit eligibility for an account | 400 missing account id; upstream status passthrough; 500 lookup failure | | `POST /api/codex-auth/reset-credits/consume` | Consume an eligible reset credit. Optional `operationId` (UUIDv4) makes the redemption idempotent: the same id replays one durable outcome instead of spending a second credit. | 400 missing account id or invalid `operationId`; 409 `identity_mismatch` when the id belongs to another account; upstream status passthrough; 503 `server_busy`, `capacity`, or `unavailable`; 500 consume failure | | `POST /api/codex-auth/login` | Start Codex login or reauthentication | 400 invalid request; conflict/busy login states | diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 819f07151f5..3b8a4affc13 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -3,6 +3,7 @@ "root": "tests", "explicit": { "pnpm-command-isolation.test.ts": "update", "provider-antigravity-quota-retry.test.ts": "providers", "project-config-warning-snapshot.test.ts": "codex-integration", "codex-quota-auto-refresh-generation.test.ts": "codex-integration", "codex-account-clear-paused.test.ts": "codex-integration", + "low-quota-protection.test.ts": "codex-integration", "responses-compaction-recovery.test.ts": "responses", "compaction-recovery-settings.test.ts": "config", "responses-compaction-recovery-policy.test.ts": "responses", "plugin-loader.test.ts": "lib", "plugin-upstream-hooks.test.ts": "lib", "cli-kiro-auto-selection.test.ts": "cli", "codebuddy-live-models.test.ts": "providers", "kiro-auto-selection.test.ts": "providers/kiro", "kiro-quota-metrics.test.ts": "providers/kiro", "management-provider-request-pacing.test.ts": "server", "desktop-supervised-restart.test.ts": "clients", "cli-restart-handoff.test.ts": "cli", diff --git a/src/codex/auth-api/login-flow.ts b/src/codex/auth-api/login-flow.ts index 74dce7b47ab..89adc254cee 100644 --- a/src/codex/auth-api/login-flow.ts +++ b/src/codex/auth-api/login-flow.ts @@ -1,6 +1,6 @@ import { withCodexAccountLogLabel } from "../account-label"; import { getCodexAccountCredential, markCodexAccountValidated, readCodexAccountRecord, saveCodexAccountCredential, CodexCredentialRefreshLockTimeoutError, CodexCredentialRefreshBusyError, CodexCredentialRefreshStaleError } from "../account-store"; -import { clearAccountQuota, isCodexQuotaExhausted, parseUsageQuota, setAccountQuotaFromParsed } from "../quota"; +import { clearAccountQuota, isCodexQuotaExhausted, isValidWhamHistoryObservation, parseUsageQuota, setAccountQuotaFromParsed } from "../quota"; import type { StoredAccountQuota, WhamUsageResponse } from "../quota"; import { ConfigMutationLockError, withConfigMutationLockSync } from "../../config"; import { appendDefaultCodexAccountNamespace, codexAccountPickerEnabled } from "../account-namespaces"; @@ -265,6 +265,7 @@ export async function handleCodexAuthLoginStart(req: Request, config: OcxConfig, let email = cred.email || accountId; let plan: string | undefined; let quota: Omit | null = null; + let policyQuota: Omit | null = null; try { const tokens = { access_token: cred.access, account_id: oauthAccountId }; const resp = await fetch("https://chatgpt.com/backend-api/wham/usage", { @@ -276,6 +277,7 @@ export async function handleCodexAuthLoginStart(req: Request, config: OcxConfig, email = data.email ?? email; plan = nonEmptyPlan(data.plan_type) ?? undefined; quota = parseUsageQuota(data); + policyQuota = isValidWhamHistoryObservation(data) ? quota : null; } } catch { /* wham fetch is non-blocking */ } // Reauth must refresh the same ChatGPT identity already bound to this pool slot. @@ -390,7 +392,7 @@ export async function handleCodexAuthLoginStart(req: Request, config: OcxConfig, clearCodexPoolRefreshFailure(accountId); if (warmup.validatedAt !== undefined) markCodexAccountValidated(accountId, warmup.validatedAt, generation); clearAccountNeedsReauth(accountId); - if (quota) setAccountQuotaFromParsed(accountId, quota); + if (quota) setAccountQuotaFromParsed(accountId, quota, undefined, undefined, policyQuota); // Keep the pool id stable; refresh display metadata after a successful login/reauth. accounts[existingIdx] = withCodexAccountLogLabel({ ...accounts[existingIdx], @@ -420,7 +422,7 @@ export async function handleCodexAuthLoginStart(req: Request, config: OcxConfig, // A new quota row is generation-gated by live account ownership. Reconcile the // durable config owner first so a partial prior sweep cannot reject this write. if (newAccountPersistence?.status === "committed" && quota) { - setAccountQuotaFromParsed(accountId, quota); + setAccountQuotaFromParsed(accountId, quota, undefined, undefined, policyQuota); } const { catalogRefreshPending } = await convergeAccountNamespaceCatalog( latestConfig, diff --git a/src/codex/auth-api/pool-quota-probe.ts b/src/codex/auth-api/pool-quota-probe.ts index 538f6023240..d2510d9f109 100644 --- a/src/codex/auth-api/pool-quota-probe.ts +++ b/src/codex/auth-api/pool-quota-probe.ts @@ -338,8 +338,9 @@ export async function commitPoolQuotaResponse( if (!isCodexAccountGenerationLive(accountId, generation)) { return { quota: null, needsReauth: false, credentialGeneration: generation }; } - setAccountQuotaFromParsed(accountId, quota, writerGeneration, undefined, quota, - ctx.poolWriter && isValidWhamHistoryObservation(data) ? { writer: ctx.poolWriter, observedAt, source: "wham", raw: quota } : undefined); + const validPolicyObservation = isValidWhamHistoryObservation(data); + setAccountQuotaFromParsed(accountId, quota, writerGeneration, undefined, validPolicyObservation ? quota : null, + ctx.poolWriter && validPolicyObservation ? { writer: ctx.poolWriter, observedAt, source: "wham", raw: quota } : undefined); return { quota: getAccountQuota(accountId), needsReauth: false, diff --git a/src/codex/low-quota-events.ts b/src/codex/low-quota-events.ts new file mode 100644 index 00000000000..0a5bd0c04e8 --- /dev/null +++ b/src/codex/low-quota-events.ts @@ -0,0 +1,28 @@ +/** Bounded event projection owned by one live low-quota registration. */ +export type LowQuotaEvent = { + accountId: string; + window: "short" | "weekly"; + percentUsed: number; + resetAt: number | null; + timestamp: number; + status: "pending" | "logged" | "delivered" | "succeeded" | "failed" | "cancelled"; + delivery: "notice" | "pause-save"; +}; + +const CAPACITY = 100; + +export function createLowQuotaEventLedger(): { + publish(event: LowQuotaEvent): void; + list(limit?: number): LowQuotaEvent[]; +} { + const events: LowQuotaEvent[] = []; + return { + publish(event) { + events.unshift({ ...event }); + if (events.length > CAPACITY) events.length = CAPACITY; + }, + list(limit = 20) { + return events.slice(0, Math.max(0, Math.min(CAPACITY, limit))).map(event => ({ ...event })); + }, + }; +} diff --git a/src/codex/low-quota-observer.ts b/src/codex/low-quota-observer.ts new file mode 100644 index 00000000000..4db45a6958d --- /dev/null +++ b/src/codex/low-quota-observer.ts @@ -0,0 +1,23 @@ +import type { StoredAccountQuota } from "./quota-types"; + +type Observer = (accountId: string, quota: Omit) => void; +const observers = new Map(); + +/** The composition root owns the live config; quota writers know only this synchronous slot. */ +export function registerLowQuotaObserver(observer: Observer): () => void { + const owner = Symbol(); + observers.set(owner, observer); + return () => { observers.delete(owner); }; +} + +/** Only newly accepted evidence belongs here, never carried or disk-hydrated windows. */ +export function observeCodexLowQuota(accountId: string, quota: Omit): void { + for (const observer of observers.values()) { + try { + observer(accountId, quota); + } catch { + // One server's optional policy cannot suppress another server's observation. + console.warn("[codex-low-quota] protection action failed"); + } + } +} diff --git a/src/codex/low-quota-protection.ts b/src/codex/low-quota-protection.ts new file mode 100644 index 00000000000..356efe10d33 --- /dev/null +++ b/src/codex/low-quota-protection.ts @@ -0,0 +1,272 @@ +import type { OcxConfig } from "../types"; +import { isSelectableCodexPoolAccount } from "./account-id"; +import { isCodexAccountPaused, setCodexAccountPaused } from "./account-pause"; +import { createLowQuotaEventLedger, type LowQuotaEvent } from "./low-quota-events"; +import { registerLowQuotaObserver } from "./low-quota-observer"; +import { resetAtToMs } from "./quota-types"; + +type Window = "short" | "weekly"; +type Notice = { window: Window; percentUsed: number; threshold: number }; +type Dependencies = { + persist?: (config: OcxConfig) => void | Promise; + notify?: (notice: Notice) => void | Promise; +}; +type EventBase = Omit; +type Episode = { reset: string; notice: "in-flight" | "logged" | "delivered" | "failed" | undefined; pausedByUs: boolean; noticeBase?: EventBase }; +export type LowQuotaRegistration = (() => void) & { + hasPendingSave(): boolean; + flush(): Promise; + listEvents(limit?: number): LowQuotaEvent[]; +}; + +const RETRY_DELAYS_MS = [100, 250]; +const FLUSH_DEADLINE_MS = 500; + +/** Each server owns its own policy, episode state and deferred writer. */ +export function registerCodexLowQuotaProtection(config: OcxConfig, deps: Dependencies = {}): LowQuotaRegistration { + const episodes = new Map(); + const autoPausedAccounts = new Set(); + const ledger = createLowQuotaEventLedger(); + let policyKey: string | undefined; + let closed = false; + let generation = 0; + let dirty = false; + let attempts = 0; + let retryDelay: number | undefined; + let saveFlight: Promise | null = null; + let activeSave: { events: Map; started: boolean; settled: boolean } | null = null; + let saveTimer: ReturnType | null = null; + const pendingPauseEvents = new Map(); + const pauseStatus = new WeakMap(); + + function event(base: EventBase, + delivery: LowQuotaEvent["delivery"], status: LowQuotaEvent["status"]): void { + if (delivery === "pause-save") pauseStatus.set(base, status); + ledger.publish({ ...base, delivery, status, timestamp: Date.now() }); + } + function cancelTimer(): void { + if (saveTimer) clearTimeout(saveTimer); + saveTimer = null; + } + function scheduleSave(delay = 0): void { + if (closed || saveTimer || saveFlight) return; + const ownerGeneration = generation; + saveTimer = setTimeout(() => { + saveTimer = null; + if (closed || ownerGeneration !== generation) { + for (const base of pendingPauseEvents.values()) { + if (pauseStatus.get(base) === "pending") event(base, "pause-save", "cancelled"); + } + return; + } + void runSave(); + }, delay); + saveTimer.unref?.(); + } + function runSave(): Promise { + if (closed || !dirty) return saveFlight ?? Promise.resolve(); + if (saveFlight) return saveFlight; + cancelTimer(); + const ownerGeneration = generation; + const saving = new Map(pendingPauseEvents); + const flightState = { events: saving, started: false, settled: false }; + activeSave = flightState; + dirty = false; + saveFlight = Promise.resolve().then(async () => { + if (closed || ownerGeneration !== generation) { + return; + } + if (deps.persist) { + flightState.started = true; + return deps.persist(config); + } + // Import only for an actual deferred save, then recheck ownership: an import + // that resolves after a timed-out flush must not start a late config write. + const { saveConfigPreservingClaudeCode } = await import("../config/live-reconcile"); + if (closed || ownerGeneration !== generation) { + return; + } + flightState.started = true; + saveConfigPreservingClaudeCode(config); + }).then(() => { + flightState.settled = true; + if (flightState.started) { + attempts = 0; + for (const [key, base] of saving) { + event(base, "pause-save", isCodexAccountPaused(config, base.accountId) ? "succeeded" : "cancelled"); + if (pendingPauseEvents.get(key) === base) pendingPauseEvents.delete(key); + } + } + }, () => { + flightState.settled = true; + if (flightState.started || !closed) { + for (const base of saving.values()) event(base, "pause-save", "failed"); + } + if (flightState.started) { + for (const [key, base] of saving) { + if (pendingPauseEvents.get(key) === base && closed) pendingPauseEvents.delete(key); + } + } + if (closed || ownerGeneration !== generation) return; + console.warn("[codex-low-quota] pause persistence failed"); + dirty = true; + retryDelay = RETRY_DELAYS_MS[attempts++]; + }).finally(() => { + if (activeSave === flightState) activeSave = null; + saveFlight = null; + if (dirty && !closed) { + if (retryDelay !== undefined) scheduleSave(retryDelay); + else if (attempts === 0) scheduleSave(); + retryDelay = undefined; + } + }); + return saveFlight; + } + function close(): void { + if (closed) return; + closed = true; + generation++; + cancelTimer(); + for (const [key, base] of pendingPauseEvents) { + if (activeSave?.started && activeSave.events.get(key) === base) continue; + if (pauseStatus.get(base) === "pending") event(base, "pause-save", "cancelled"); + pendingPauseEvents.delete(key); + } + for (const episode of episodes.values()) { + if (episode.notice === "in-flight" && episode.noticeBase) event(episode.noticeBase, "notice", "cancelled"); + } + dirty = false; + autoPausedAccounts.clear(); + unregister(); + } + async function flush(): Promise { + if (closed) return; + const deadline = Date.now() + FLUSH_DEADLINE_MS; + // Drain the in-flight write and any coalesced save before closing the owner. + while (saveFlight || dirty) { + cancelTimer(); + const flight = saveFlight ?? runSave(); + let timeout: ReturnType | undefined; + const completed = await Promise.race([ + flight.then(() => true), + new Promise(resolve => { + timeout = setTimeout(() => resolve(false), Math.max(0, deadline - Date.now())); + }), + ]); + if (timeout) clearTimeout(timeout); + if (!completed || Date.now() >= deadline) { + const inFlight = activeSave; + const stillRunning = inFlight !== null && inFlight.started && !inFlight.settled; + if (inFlight && stillRunning) { + for (const base of inFlight.events.values()) event(base, "pause-save", "pending"); + } + console.warn(stillRunning + ? "[codex-low-quota] pause persistence flush timed out; in-flight save remains pending" + : "[codex-low-quota] pause persistence flush reached its deadline"); + break; + } + if (attempts > RETRY_DELAYS_MS.length) break; + } + close(); + } + + const unregister = registerLowQuotaObserver((accountId, quota) => { + if (closed) return; + const policy = config.codexPool?.lowQuotaProtection; + const nextKey = JSON.stringify(policy); + if (policyKey !== nextKey) { + for (const episode of episodes.values()) { + if (episode.notice === "in-flight" && episode.noticeBase) event(episode.noticeBase, "notice", "cancelled"); + } + episodes.clear(); + autoPausedAccounts.clear(); + policyKey = nextKey; + } + if (!policy?.enabled) return; + const liveIds = new Set((config.codexAccounts ?? []).filter(isSelectableCodexPoolAccount).map(account => account.id)); + for (const id of autoPausedAccounts) if (!liveIds.has(id)) autoPausedAccounts.delete(id); + for (const [key, episode] of episodes) { + if (liveIds.has(key.split("\u0000")[0]!)) continue; + if (episode.notice === "in-flight" && episode.noticeBase) event(episode.noticeBase, "notice", "cancelled"); + episodes.delete(key); + } + if (!liveIds.has(accountId)) return; + // A manual resume ends this policy's active pause, but each already-high + // window keeps its episode marker until recovery or a new reset. + if (!isCodexAccountPaused(config, accountId)) autoPausedAccounts.delete(accountId); + for (const window of ["short", "weekly"] as const) { + if (!policy.windows[window]) continue; + const percentUsed = quota[`${window}Percent`]; + const rawReset = quota[`${window}ResetAt`]; + if (typeof percentUsed !== "number" || !Number.isFinite(percentUsed) || percentUsed < 0 || percentUsed > 100) continue; + if (rawReset !== undefined && (!Number.isFinite(rawReset) || resetAtToMs(rawReset) <= Date.now())) continue; + const key = `${accountId}\u0000${window}`; + if (percentUsed < policy.threshold) { + const prior = episodes.get(key); + if (prior?.notice === "in-flight" && prior.noticeBase) event(prior.noticeBase, "notice", "cancelled"); + episodes.delete(key); + continue; + } + const resetAt = rawReset === undefined ? null : resetAtToMs(rawReset); + const reset = resetAt === null ? "unknown" : String(resetAt); + let episode = episodes.get(key); + if (!episode || episode.reset !== reset) { + episode = { reset, notice: undefined, pausedByUs: false }; + episodes.set(key, episode); + } + const base = { accountId, window, percentUsed, resetAt }; + if (autoPausedAccounts.has(accountId) && isCodexAccountPaused(config, accountId)) { + episode.pausedByUs = true; + } + if (dirty && attempts > RETRY_DELAYS_MS.length) { + attempts = 0; + scheduleSave(); + } + if (policy.actions.pause && !isCodexAccountPaused(config, accountId) && !episode.pausedByUs) { + setCodexAccountPaused(config, accountId, true); + autoPausedAccounts.add(accountId); + episode.pausedByUs = true; + pendingPauseEvents.set(key, base); + dirty = true; + attempts = 0; + retryDelay = undefined; + event(base, "pause-save", "pending"); + scheduleSave(); + } + // A manual resume leaves pausedByUs set until recovery or a new reset episode. + if (policy.actions.notify && episode.notice !== "in-flight" && episode.notice !== "logged" && episode.notice !== "delivered") { + console.warn(`[codex-low-quota] ${window === "short" ? "5-hour" : "weekly"} quota reached ${percentUsed}% used (threshold ${policy.threshold}%)`); + if (!deps.notify) { + event(base, "notice", "logged"); + episode.notice = "logged"; + continue; + } + episode.notice = "in-flight"; + episode.noticeBase = base; + event(base, "notice", "pending"); + try { + void Promise.resolve(deps.notify({ window, percentUsed, threshold: policy.threshold })).then(() => { + if (closed || episodes.get(key) !== episode) return; + event(base, "notice", "delivered"); + episode.notice = "delivered"; + episode.noticeBase = undefined; + }, () => { + if (closed || episodes.get(key) !== episode) return; + event(base, "notice", "failed"); + episode.notice = "failed"; + episode.noticeBase = undefined; + }); + } catch { + event(base, "notice", "failed"); + episode.notice = "failed"; + episode.noticeBase = undefined; + } + } + } + }); + return Object.assign(close, { + hasPendingSave: () => dirty || saveFlight !== null || saveTimer !== null, + flush, + listEvents: ledger.list, + }); +} diff --git a/src/codex/quota.ts b/src/codex/quota.ts index 6f27fd32a2b..57444410a85 100644 --- a/src/codex/quota.ts +++ b/src/codex/quota.ts @@ -4,6 +4,7 @@ import { atomicWriteFile, getConfigDir } from "../config"; import { captureConfigGeneration, type GenerationContext } from "../lib/state-store-sweeper"; import { isThirtyDayOnlyCodexPlan } from "./plan"; import { stampCodexQuotaUsageObservation } from "./quota-observation-freshness"; +import { observeCodexLowQuota } from "./low-quota-observer"; import { MAIN_CODEX_ACCOUNT_ID } from "./account-id"; import { getObservedMainQuotaIdentityKey, isMainQuotaWriterLive, type MainQuotaWriter } from "./main-account-cache"; @@ -276,7 +277,7 @@ function snapshotHasUsage(quota: Omit): boolean return snapshotHasWeekly(quota) || snapshotHasMonthly(quota) || snapshotHasShort(quota) || snapshotHasCustom(quota); } /** - * Publish parsed display quota and separately validated main-policy evidence after writer checks. + * Publish parsed display quota and separately validated policy evidence after writer checks. * A null policy observation retains only the matching main identity's previous evidence; * transient replacement markers are consumed during merging and never enter stored snapshots. */ @@ -285,7 +286,7 @@ export function setAccountQuotaFromParsed( quota: Omit | null, writerGeneration = captureConfigGeneration(), mainWriter?: MainQuotaWriter, - policyQuota: MainPolicyQuotaObservation | null = quota, + policyQuota: MainPolicyQuotaObservation | null = accountId === MAIN_CODEX_ACCOUNT_ID ? quota : null, historyEvidence?: QuotaObservationEvidence, ): void { quota = withoutRetiredCodexQuota(quota); @@ -322,6 +323,7 @@ export function setAccountQuotaFromParsed( schedulePersistAccountQuotas(); // Credits carry the previous usage tuple; they must not refresh its observation clock. if (!(quota.resetCredits !== undefined && !snapshotHasUsage(quota))) { + if (!isMain && policyQuota) observeCodexLowQuota(accountId, policyQuota); notifyCodexQuotaSnapshot(accountId, next); } } @@ -650,6 +652,9 @@ export function updateAccountQuota( // otherwise bypass detection AND leave a stale baseline that corrupts the next real diff. // The credits-only path at setAccountQuotaFromParsed deliberately does not notify; this one // writes window percentages and deadlines, so it must. + if (accountId !== MAIN_CODEX_ACCOUNT_ID && nextWeekly !== undefined && !isInvalidPolicyUsagePercent(weekly)) observeCodexLowQuota(accountId, { + weeklyPercent: nextWeekly, weeklyResetAt: nextWeeklyResetAt, + }); notifyCodexQuotaSnapshot(accountId, quota); } diff --git a/src/config/schema/leaf-validators.ts b/src/config/schema/leaf-validators.ts index 4b6c6805e38..92ddb4ae396 100644 --- a/src/config/schema/leaf-validators.ts +++ b/src/config/schema/leaf-validators.ts @@ -960,6 +960,26 @@ export const clientConnectionSchema = z.object({ export const codexPoolSchema = z.object({ excludedPlans: z.array(z.string().trim().min(1)).optional(), startIdleWindows: z.boolean().optional(), + lowQuotaProtection: z.object({ + enabled: z.boolean(), + threshold: z.number().finite().min(1).max(100), + actions: z.object({ + pause: z.boolean(), + notify: z.boolean(), + }).strict(), + windows: z.object({ + short: z.boolean(), + weekly: z.boolean(), + }).strict(), + }).strict().superRefine((policy, ctx) => { + if (!policy.enabled) return; + if (!policy.actions.pause && !policy.actions.notify) { + ctx.addIssue({ code: "custom", path: ["actions"], message: "enabled low-quota protection needs an action" }); + } + if (!policy.windows.short && !policy.windows.weekly) { + ctx.addIssue({ code: "custom", path: ["windows"], message: "enabled low-quota protection needs a window" }); + } + }).optional(), }).strict(); /** diff --git a/src/server/background-lifecycle.ts b/src/server/background-lifecycle.ts index a2cf4da136e..d12a366ef31 100644 --- a/src/server/background-lifecycle.ts +++ b/src/server/background-lifecycle.ts @@ -1,4 +1,6 @@ -import type { StorageCleanupPolicy } from "../types"; +import type { OcxConfig, StorageCleanupPolicy } from "../types"; +import { registerCodexLowQuotaProtection, type LowQuotaRegistration } from "../codex/low-quota-protection"; +import type { LowQuotaEvent } from "../codex/low-quota-events"; import { startStateStoreSweeper } from "../lib/state-store-sweeper"; import { abortStorageCleanupPolicyJobAsync, @@ -40,10 +42,12 @@ type LeaseOwner = { token: symbol; applyPolicy: PolicyApply; resources: ServerResourceOwnerLease; + lowQuota: LowQuotaRegistration; }; export type ServerBackgroundLifecycleLease = { scheduleStartupRun(): void; + listLowQuotaEvents(limit?: number): LowQuotaEvent[]; release(): Promise; releaseAfterFailedStart(): void; }; @@ -148,6 +152,7 @@ function removeOwner(owner: LeaseOwner): boolean { function releaseOwnerSynchronously(owner: LeaseOwner): "inactive" | "shared" | "last" { if (!removeOwner(owner)) return "inactive"; + owner.lowQuota(); owner.resources.release(); const nextOwner = owners.at(-1); if (nextOwner) { @@ -168,6 +173,7 @@ function releaseOwnerSynchronously(owner: LeaseOwner): "inactive" | "shared" | " */ export function acquireServerBackgroundLifecycle( applyPolicy: PolicyApply, + config: OcxConfig, ): ServerBackgroundLifecycleLease { if (cleanupInProgress) { throw new Error("server background lifecycle cleanup is still in progress"); @@ -177,6 +183,7 @@ export function acquireServerBackgroundLifecycle( token: Symbol("server-background-lifecycle"), applyPolicy, resources: acquireServerResourceOwner(), + lowQuota: registerCodexLowQuotaProtection(config), }; try { if (!processLoops) { @@ -186,12 +193,20 @@ export function acquireServerBackgroundLifecycle( } owners.push(owner); } catch (error) { + owner.lowQuota(); owner.resources.release(); throw error; } let releaseFlight: Promise | null = null; + function finishRelease(): Promise { + const outcome = releaseOwnerSynchronously(owner); + if (outcome !== "last") return Promise.resolve(); + cleanupInProgress = true; + return stopStoragePolicyWorker().finally(() => { cleanupInProgress = false; }); + } return { + listLowQuotaEvents(limit) { return owner.lowQuota.listEvents(limit); }, scheduleStartupRun() { if (owners.some(candidate => candidate.token === owner.token)) { scheduleStorageCleanupStartupRun(); @@ -199,19 +214,15 @@ export function acquireServerBackgroundLifecycle( }, release() { if (releaseFlight) return releaseFlight; - const outcome = releaseOwnerSynchronously(owner); - if (outcome !== "last") { - releaseFlight = Promise.resolve(); - return releaseFlight; - } - cleanupInProgress = true; - releaseFlight = stopStoragePolicyWorker().finally(() => { - cleanupInProgress = false; - }); + // Default installs have no low-quota write. Preserve synchronous owner removal; + // an unconditional await here extends process-global ownership into later starts. + if (!owner.lowQuota.hasPendingSave()) releaseFlight = finishRelease(); + else releaseFlight = owner.lowQuota.flush().then(finishRelease, finishRelease); return releaseFlight; }, releaseAfterFailedStart() { if (releaseFlight) return; + owner.lowQuota(); // Startup evaluation is scheduled only after both listeners bind, so a // failed start cannot have spawned a Worker for this lease. Keep rollback // synchronous so a caller may immediately retry a different port. diff --git a/src/server/index.ts b/src/server/index.ts index eef1f39a802..bc827af06fe 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -640,6 +640,7 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe let remoteWorkspaceShutdown: (() => Promise) | undefined; const managementApiDeps: ManagementApiDeps = { ...deps.managementApi, + listLowQuotaEvents: limit => backgroundLifecycle?.listLowQuotaEvents(limit) ?? [], remoteWorkspaceStopping: () => remoteWorkspaceStopping, onRemoteWorkspaceShutdown: shutdown => { remoteWorkspaceShutdown = shutdown; }, linkSupervisor: () => optionalListeners.linkSupervisor(), linkListener: () => optionalListeners, }; @@ -655,13 +656,10 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe return workspaceRuntimeFlight; }; try { - backgroundLifecycle = acquireServerBackgroundLifecycle(applyPolicy); + backgroundLifecycle = acquireServerBackgroundLifecycle(applyPolicy, config); unregisterQuotaAutoRefresh = (deps.registerCodexQuotaAutoRefreshWorker ?? registerCodexQuotaAutoRefreshWorker)(config); - // External `ocx config set` / direct config.json edits run in other - // processes; poll the file so Logs/Usage display prices follow them live. - // Started inside the guarded startup transaction so the catch below can - // release the owner-scoped lease on any listener failure. + // Poll external pricing edits; the startup catch releases this owner-scoped lease. userCostOverlayReconciler = startUserCostOverlayReconciler({ liveConfig: config }); const serveOptions = createServeOptions({ drainingResponse, @@ -754,9 +752,7 @@ function startServerWithSpendLedgerOwner(port: number | undefined, deps: StartSe value: async (closeActiveConnections?: boolean): Promise => { remoteWorkspaceStopping = true; liveCallBindings.clear(); - // Disarm the package-tree restart timer before listener teardown: a queued - // replacement callback must not call acceptSystemRestart() after stop() has - // begun, or it would schedule a drain-and-restart on a stopped server. + // Disarm the package-tree restart timer before teardown can schedule another restart. if (!packageRefreshStopped) { packageRefreshStopped = true; stopPackageRefresh(); diff --git a/src/server/management-api.ts b/src/server/management-api.ts index dc8390c84ca..da36c718634 100644 --- a/src/server/management-api.ts +++ b/src/server/management-api.ts @@ -149,6 +149,12 @@ async function handleQuotaResetRoutesOnDemand(ctx: ManagementContext): Promise { + if (!pathInManagementNamespace(ctx.url.pathname, "/api/codex-auth/low-quota-events", false)) return null; + const { handleLowQuotaRoutes } = await import("./management/low-quota-routes"); + return handleLowQuotaRoutes(ctx); +} + /** * Lazy like the Lab and routing-profile handlers, and for the same recorded reason: this file is * mounted for every dashboard request, so a static import would put the workflow-budget ledger @@ -328,6 +334,7 @@ export async function handleManagementAPI( ?? (await handleLogsUsageRoutes(ctx)) ?? (await handleRequestHistoryRoutes(ctx)) ?? (await handleQuotaResetRoutesOnDemand(ctx)) + ?? (await handleLowQuotaRoutesOnDemand(ctx)) ?? (await handleWorkflowBudgetRoutesOnDemand(ctx)) ?? (await handleProtocolRoutesOnDemand(ctx)) ?? (await handleGrokCouponRoutesOnDemand(ctx)) diff --git a/src/server/management/context.ts b/src/server/management/context.ts index 0a515b78104..b9a0bba0644 100644 --- a/src/server/management/context.ts +++ b/src/server/management/context.ts @@ -1,4 +1,5 @@ import type { OcxConfig } from "../../types"; +import type { LowQuotaEvent } from "../../codex/low-quota-events"; import type { Channel } from "../../update/index"; import type { UpdateCheckResult } from "../../update/job"; import type { NativeProfileApiDeps } from "../../codex/native-profile-api"; @@ -44,6 +45,8 @@ export interface ManagementRequestIngress { } export interface ManagementApiDeps { + /** Bound to this server's lifecycle owner; absent in direct route tests. */ + listLowQuotaEvents?: (limit?: number) => LowQuotaEvent[]; /** Bound Claude intercept state, injectable for isolated management-route tests. */ getClaudeInterceptState?: typeof import("../../claude/intercept/runtime").getClaudeInterceptState; /** Reconciliation seam for field-scoped rollback tests. */ diff --git a/src/server/management/low-quota-routes.ts b/src/server/management/low-quota-routes.ts new file mode 100644 index 00000000000..998c6478ae6 --- /dev/null +++ b/src/server/management/low-quota-routes.ts @@ -0,0 +1,14 @@ +/** Auth is enforced by the management dispatcher before this lazy route runs. */ +import { jsonResponse } from "../auth-cors"; +import type { ManagementContext } from "./context"; + +export async function handleLowQuotaRoutes(ctx: ManagementContext): Promise { + const { url, req, config } = ctx; + if (url.pathname !== "/api/codex-auth/low-quota-events" || req.method !== "GET") return null; + const raw = url.searchParams.get("limit"); + if (raw !== null && (!/^[0-9]+$/.test(raw) || !Number.isSafeInteger(Number(raw)))) { + return jsonResponse({ error: { code: "invalid_limit", message: "limit must be a non-negative integer" } }, 400, req, config); + } + const limit = raw === null ? 20 : Math.min(100, Number(raw)); + return jsonResponse({ events: ctx.deps.listLowQuotaEvents?.(limit) ?? [] }, 200, req, config); +} diff --git a/src/server/management/route-registry.ts b/src/server/management/route-registry.ts index 2bf5db58488..59fb43adf26 100644 --- a/src/server/management/route-registry.ts +++ b/src/server/management/route-registry.ts @@ -103,6 +103,7 @@ export const MANAGEMENT_ROUTES: readonly ManagementRoute[] = [ { method: "GET", path: "/api/codex-auth/login-status", module: "codex/auth-api/routes", mutates: false }, { method: "GET", path: "/api/codex-auth/quota", module: "codex/auth-api/routes", mutates: false }, { method: "GET", path: "/api/codex-auth/quota/history", module: "codex/auth-api/routes", mutates: false }, + { method: "GET", path: "/api/codex-auth/low-quota-events", module: "server/management/low-quota-routes", mutates: false, exempt: { reason: "deferred-verb", why: "The authenticated event history is an operator diagnostic with no CLI verb yet.", owner: "rt5 low-quota", ownerDoc: "structure/gui-and-management-api.md" } }, { method: "GET", path: "/api/codex-auth/reset-credits", module: "codex/auth-api/routes", mutates: false }, { method: "PATCH", path: "/api/codex-auth/pool-strategy", module: "codex/auth-api/routes", mutates: true }, { method: "POST", path: "/api/codex-auth/accounts", module: "codex/auth-api/routes", mutates: true }, diff --git a/src/types.ts b/src/types.ts index ccca3817fb5..b44a67b0c69 100644 --- a/src/types.ts +++ b/src/types.ts @@ -80,6 +80,7 @@ export type { OcxConfig, SkillsCatalogRefresh, OcxSkillsConfig, + CodexLowQuotaProtectionConfig, OcxAccountPoolRotationStrategy, OcxAccountPoolQuotaWindow, OcxComboCooldownWaitPolicy, diff --git a/src/types/config.ts b/src/types/config.ts index d8bc49ee311..0b8f05856ba 100644 --- a/src/types/config.ts +++ b/src/types/config.ts @@ -1569,6 +1569,23 @@ export interface OcxCodexPoolConfig { * operator never meant to exclude. */ excludedPlans?: string[]; + /** Optional per-account response to fresh quota observations at or above a usage percentage. */ + lowQuotaProtection?: CodexLowQuotaProtectionConfig; +} + +/** Optional policy for pausing accounts and notifying when selected quota windows are low. */ +export interface CodexLowQuotaProtectionConfig { + enabled: boolean; + /** Inclusive usage percentage from 1 to 100. */ + threshold: number; + actions: { + pause: boolean; + notify: boolean; + }; + windows: { + short: boolean; + weekly: boolean; + }; } /** diff --git a/structure/config.md b/structure/config.md index cf0f065355d..b12aff43145 100644 --- a/structure/config.md +++ b/structure/config.md @@ -492,16 +492,13 @@ the residual directory for manual review; there is no recursive-delete fallback. ## Remote client key files -The connection's `tokenFingerprint` participates in -[`ocx status` credential binding](runtime.md#remote-hub-status-credential-binding). +The connection's `tokenFingerprint` participates in [`ocx status` credential binding](runtime.md#remote-hub-status-credential-binding). -Client catalog readiness observes the selected Codex runtime without creating or rewriting -`codex-runtime.json`; general status reuses its already-resolved command under the [runtime contract](runtime.md#remote-hub-hardening-ownership). +Client catalog readiness observes the selected Codex runtime without creating or rewriting `codex-runtime.json`; general status reuses its already-resolved command under the [runtime contract](runtime.md#remote-hub-hardening-ownership). Client connection metadata stores a stable `apiKeyId` and a non-secret rotation `pendingOperation`. The current data secret remains only in `service-api-token`; a bounded rotation temporarily keeps the old secret in owner-only `service-api-token.prev`. Commit or recovery clears the marker before orphan cleanup. `ocx disconnect` is local-only and leaves remote revocation to the hub's **Integrations → API Keys** page. Hub and local usage stores are not mirrored. -Codex display-cache expiry, retained blocking main-policy evidence, and reset history follow the -[quota cache contract](providers/openai-tiers.md#quota-cache-and-short-window-history). +Codex display-cache expiry, retained blocking main-policy evidence, and reset history follow the [quota cache contract](providers/openai-tiers.md#quota-cache-and-short-window-history). `codexPool.excludedPlans` is interpreted only by automatic selection; its all-excluded and explicit-route behavior follows the [plan exclusion contract](providers/openai-accounts.md#automatic-pool-plan-exclusions). Optional `codexPool.startIdleWindows` defaults off and follows the [idle-window steering contract](providers/openai-accounts.md#idle-window-steering), using real new requests to start observed idle 5-hour windows. @@ -517,6 +514,10 @@ Usage consumers preserve positive incomplete-history metadata as specified in [u malformed persisted values stay disabled. It controls only the allowlisted client-output hints described in [Responses transport](transports/responses.md), not upstream policy or model selection. +## Codex Pool low-quota protection + +`codexPool.lowQuotaProtection` is opt-in, requires a 1–100 threshold and a selected action/window when enabled, covers pool accounts only and is independent of proactive switching and the main account’s 98% hard lock. `src/codex/low-quota-protection.ts` pauses in live `pausedCodexAccountIds` before selection, then coalesces a deferred config save with bounded retry and shutdown flush. Fresh accepted observations reach `src/codex/low-quota-observer.ts`; credits-only and expired windows do not act. Manual resume suppresses repause across currently qualifying window episodes; a new reset boundary or below-threshold reading re-arms the policy, but never automatically resumes an account. A timed-out in-flight save remains pending until its eventual success or failure; queued work is cancelled at owner close. An unsuccessful save does not survive restart. The default alert is log-and-API only and records `logged`, not notification delivery. + ## Management-backed CLI commands need a management plane `src/cli/runtime-api.ts` is the single client every headless management subcommand calls through, so @@ -596,5 +597,4 @@ so wrong types and unknown nested fields are rejected rather than silently saved `apiSurfaces` and `protocols` on `src/types/config.ts` are parsed by `src/protocols/settings.ts` only; [Protocol Paths](data-planes/protocol-paths.md#settings) owns their schema handling, meaning and the one writer (`PATCH /api/protocols/settings`), including why closing Messages also writes `claudeCode.enabled` through `commitClaudeCodeBlock` (`src/claude/claude-code-block.ts`, the sentinel-stamping block writer every management route uses). Stored Direct substitution follows the [credential identity contract](providers/openai-accounts.md#sidecars-management-and-ui): both synchronous and asynchronous materializers discard the caller account header before applying the stored credential; ordinary native Direct passthrough is unchanged. - Proxy activation and credential-safe CLI output follow [Proxy Configuration](config-proxy.md). diff --git a/structure/gui-and-management-api.md b/structure/gui-and-management-api.md index efd02380539..8f97307c1e0 100644 --- a/structure/gui-and-management-api.md +++ b/structure/gui-and-management-api.md @@ -266,6 +266,7 @@ after the save, and a non-retryable skip (no managed catalog) is a clean save. C | Protocol paths | `src/server/management/protocol-routes.ts` (lazy-loaded) — `GET /api/protocols` returns the contract version, the resolved API surfaces and protocol settings, the policy revision and the feature vocabulary; with `?provider=` (one non-empty name of at most 200 characters without control characters, else 400 `invalid_provider`; an unconfigured name is 404 `unknown_provider`, neither echoing the name) it adds a `provider` block from `src/protocols/provider-summary.ts` ([Protocol Paths](data-planes/protocol-paths.md#provider-wire-summary)); `POST /api/protocols/plan` takes `{ model, inbound, features? }` (model at most 200 characters, at most 24 features, any other key refused with 400) and returns a `ProtocolPlanV1` with `basis: "preview"` from `src/protocols/plan-snapshot.ts`. Both are read-only, never log their input, and send nothing upstream. `PATCH /api/protocols/settings` takes `{ messagesEnabled?, unrepresentable?, rollout? }` (strict: unknown keys and wrong types are 400), validates and applies through `src/server/management/protocol-settings-patch.ts`, persists with `saveConfigPreservingClaudeCode`, restores the live config if the save fails (409 on lock contention, 500 otherwise), and answers with the fresh `GET /api/protocols` body; closing Messages also writes `claudeCode.enabled = false` through `commitClaudeCodeBlock` ([Protocol Paths](data-planes/protocol-paths.md#settings)). The CLI drives all three through `ocx api protocols`, `ocx api explain` and `ocx api policy` ([Protocol Paths](data-planes/protocol-paths.md#cli)); none carries a route-registry exemption. | | Workflow budget | `src/server/management/workflow-budget-routes.ts` — `GET /api/workflow-budget` reads the tracked roots or one root, and `POST /api/workflow-budget/clear` clears exactly one. The clear moves the windowed send ring and the child map and nothing else: `active` belongs to turns still in flight, the spend ledger is a token budget an operator did not ask to forgive, and the lifetime send total survives so a clear cannot launder the record. A refusal event carries `spendScope` and `spendLimit` when a token ceiling fired, so the reason is readable without the config open beside it; no scope id is ever attached, because root ids are client thread headers and identity ids are credentials. Both are `deferred-verb` in the route registry — they are owed CLI verbs, and because the ledger is process memory there is no local projection the CLI could read instead. See [`../devlog/_plan/260915_workflow_budget_window/030_wfc_diff_plan.md`](../devlog/_plan/260915_workflow_budget_window/030_wfc_diff_plan.md). | | Codex accounts | `src/codex/auth-api/routes.ts` — `GET/POST/DELETE /api/codex-auth/accounts`, `PUT /api/codex-auth/accounts/alias`, `PUT /api/codex-auth/accounts/pause`, `PUT /api/codex-auth/accounts/pause-exhausted`, `POST /api/codex-auth/accounts/clear-cooldown`, `GET/PUT /api/codex-auth/active`, `PUT /api/codex-auth/auto-switch`, `PUT /api/codex-auth/pool-strategy`, `PUT /api/codex-auth/failover`, `GET /api/codex-auth/quota`, `GET /api/codex-auth/reset-credits` with `POST /api/codex-auth/reset-credits/consume`, and the login flow `POST /api/codex-auth/login`, `POST /api/codex-auth/login/code`, `POST /api/codex-auth/login/cancel`, `GET /api/codex-auth/login-status`. Per-account quota activation uses the existing `GET/PUT /api/settings` surface and `src/codex/quota-auto-refresh.ts`, keeping scheduled spending separate from credential/authentication mutation. Account ids are opaque handles and are serialized so the GUI can address an account; emails are masked and tokens are never serialized. New-account config commits add UI-managed selector bindings in the same config save; deletion deliberately retains existing bindings for fail-closed exact routing and re-add stability. Account mutations request catalog convergence only after config durability and expose only the boolean `catalogRefreshPending` completion projection. | +| Low-quota history | `src/server/management/low-quota-routes.ts` lazily serves authenticated `GET /api/codex-auth/low-quota-events` from the bounded ledger owned by the server handling the request. The response includes only that server’s account ids and notice/pause-save status; logs omit identifiers. The default notice records `logged`; a notice records `delivered` only after an injected sink succeeds, while a pause-save records `succeeded` after persistence succeeds. A timed-out in-flight save remains `pending` until it settles. Invalid limits return 400. No OS notification is emitted. | | Sidebar | `src/server/management/sidebar-routes.ts` — `GET/POST /api/github/star`, `GET /api/update/badge`, and `POST /api/update/desktop-snapshot`. The snapshot POST accepts the raw admin-token principal or the dedicated `local-desktop-snapshot-capability`; GUI sessions and requests carrying `Origin` cannot publish desktop state. A capability's bounded raw body is verified against its authenticated digest before JSON parsing or storage. Badge state is cosmetic and a failed poll degrades silently. | | Logs | `src/server/management/logs-usage-routes.ts` — `GET /api/logs`, `GET /api/claude/inbound-debug`, and `GET /api/debug/injection-logs` join the debug streams described above. | diff --git a/structure/providers/openai-accounts.md b/structure/providers/openai-accounts.md index 5010979c458..fe3b3763b20 100644 --- a/structure/providers/openai-accounts.md +++ b/structure/providers/openai-accounts.md @@ -282,6 +282,10 @@ materialized headers pass the proxy-credential exclusion check before owner matc `capturePoolQuotaWriter` captures the exact dispatched access/account pair and generation. Legacy identity initialization rechecks under the credential mutation lock, persists metadata without advancing credential generation or mutation epoch, and fails to no optional evidence on read/lock/write errors. Append admission uses the captured generation and tag; history retention compares the tag across ordinary refresh. Native main is excluded from this pool proof. These interfaces supply the bounded observation layer; the identity alone is neither a quota sample nor proof of capacity. +## Low-quota protection + +`src/codex/quota.ts` sends accepted usage observations through `src/codex/low-quota-observer.ts`; credits-only updates never replay carried usage into protection, and pool WHAM/header observations reach the policy only after raw percentages pass validation; clamped display bars cannot authorize a pause. The pool-account-only [configuration policy](../config.md#codex-pool-low-quota-protection) pauses live selection immediately, defers a bounded config save, and deduplicates account/window notices. Manual resume is respected across every qualifying window already active for that account until recovery or a new reset episode. Native main keeps its separate 98% hard lock. Each server registration owns a bounded status ledger; its authenticated management route exposes only its own account ids. The default alert is a log line plus a `logged` event, with no OS notification. + ## Bounded pool quota observations `src/codex/quota-history.ts` retains at most 200 raw observations per stored pool account for 30 days, bounded globally to 64 identities, 4096 observations and 2 MiB. `src/codex/quota.ts` persists these alongside the latest quota cache; the file reader caps allocation at 4 MiB and rejects nonregular/oversized input. Invalid history envelopes are discarded without blocking inference. Atomic cache replacement is best-effort single-writer persistence, not cross-process merging. diff --git a/structure/runtime.md b/structure/runtime.md index 1a197584da0..c7daa2d5471 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -182,7 +182,7 @@ their own files. ## Lifecycle Startup catalog sync and native restore apply the [retired-native policy](catalog.md#shared-catalog). -Codex quota processing has shared and Reserve scopes; retired model evidence is suppressed as +`src/server/background-lifecycle.ts` registers low-quota protection for each live server, flushes pending saves before releasing its owner, binds its own event ledger to management requests, and unregisters it on failed startup or stop; with no pending save, owner release remains synchronous. The protection writer loads config persistence only when a threshold schedules a save; see [configuration](config.md#codex-pool-low-quota-protection). Codex quota processing has shared and Reserve scopes; retired model evidence is suppressed as described in [OpenAI quota ownership](providers/openai-tiers.md#public-provider-contract). `ocx start` refuses a duplicate PID, starts the proxy, writes `~/.opencodex/ocx.pid` and diff --git a/tests/codex-integration/low-quota-protection.test.ts b/tests/codex-integration/low-quota-protection.test.ts new file mode 100644 index 00000000000..883bc6ac8da --- /dev/null +++ b/tests/codex-integration/low-quota-protection.test.ts @@ -0,0 +1,456 @@ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { validateConfigCandidate } from "../../src/config"; +import { flushConfigDirHardeningForTests } from "../../src/config/paths"; +import { observeCodexLowQuota, registerLowQuotaObserver } from "../../src/codex/low-quota-observer"; +import { createLowQuotaEventLedger } from "../../src/codex/low-quota-events"; +import { registerCodexLowQuotaProtection, type LowQuotaRegistration } from "../../src/codex/low-quota-protection"; +import { setCodexAccountPaused } from "../../src/codex/account-pause"; +import { MAIN_CODEX_ACCOUNT_ID } from "../../src/codex/account-id"; +import { saveCodexAccountCredential } from "../../src/codex/account-store"; +import { commitPoolQuotaResponse } from "../../src/codex/auth-api/pool-quota-probe"; +import { captureConfigGeneration } from "../../src/lib/state-store-sweeper"; +import { applyAccountQuotaFromUpstreamHeaders, clearAccountQuota, getAccountQuota, setAccountQuotaFromParsed, updateAccountQuota } from "../../src/codex/quota"; +import type { CodexAccount } from "../../src/types/accounts"; +import type { CodexLowQuotaProtectionConfig, OcxConfig } from "../../src/types/config"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; + +const ACCOUNT_A = "low-quota-a"; +const ACCOUNT_B = "low-quota-b"; +type Notice = { window: "short" | "weekly"; percentUsed: number; threshold: number }; + +let home = ""; +let previousHome: string | undefined; +let cleanups: Array<() => void> = []; + +function account(id: string): CodexAccount { + return { id, email: `${id}@example.test`, isMain: false, plan: "team" }; +} + +function protection(overrides: Partial = {}): CodexLowQuotaProtectionConfig { + return { + enabled: true, + threshold: 80, + actions: { pause: true, notify: false }, + windows: { short: true, weekly: true }, + ...overrides, + }; +} + +function configWith(policy?: CodexLowQuotaProtectionConfig): OcxConfig { + return { + port: 10100, + defaultProvider: "fixture", + providers: { fixture: { adapter: "openai-chat", baseUrl: "https://fixture.example/v1", apiKey: "fixture-key" } }, + codexAccounts: [account(ACCOUNT_A), account(ACCOUNT_B)], + ...(policy === undefined ? {} : { codexPool: { lowQuotaProtection: policy } }), + }; +} + +function register(config: OcxConfig, deps: { + persist?: (next: OcxConfig) => void; + notify?: (notice: Notice) => void | Promise; +} = {}): LowQuotaRegistration { + const cleanup = registerCodexLowQuotaProtection(config, deps); + cleanups.push(cleanup); + return cleanup; +} + +beforeEach(() => { + previousHome = process.env.OPENCODEX_HOME; + home = mkdtempSync(join(tmpdir(), "ocx-low-quota-protection-")); + process.env.OPENCODEX_HOME = home; + cleanups = []; + clearAccountQuota(); +}); + +afterEach(async () => { + for (const cleanup of cleanups.reverse()) cleanup(); + clearAccountQuota(); + await flushConfigDirHardeningForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + removeTreeWithRetry(home); +}); + +describe("low quota protection", () => { + test("persists a paused-account snapshot only after the threshold is reached", async () => { + const config = configWith(protection()); + const persisted: OcxConfig[] = []; + const registration = register(config, { persist: next => { persisted.push(structuredClone(next)); } }); + expect(registration.hasPendingSave()).toBe(false); + + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 79 }); + expect(config.pausedCodexAccountIds).toBeUndefined(); + expect(persisted).toEqual([]); + + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 80 }); + expect(registration.hasPendingSave()).toBe(true); + expect(config.pausedCodexAccountIds).toEqual([ACCOUNT_A]); + expect(persisted).toHaveLength(0); + await registration.flush(); + expect(registration.hasPendingSave()).toBe(false); + expect(persisted).toHaveLength(1); + expect(persisted[0]?.pausedCodexAccountIds).toContain(ACCOUNT_A); + observeCodexLowQuota(ACCOUNT_B, { weeklyPercent: 79 }); + expect(config.pausedCodexAccountIds).toEqual([ACCOUNT_A]); + expect(persisted).toHaveLength(1); + }); + + test("notifies once per account and window, then rearms for a reset or a recovery", () => { + const config = configWith(protection({ actions: { pause: false, notify: true } })); + const notices: Notice[] = []; + register(config, { notify: async notice => { notices.push(notice); } }); + const firstReset = Date.now() + 60_000; + const secondReset = firstReset + 60_000; + + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 80, weeklyResetAt: firstReset }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 95, weeklyResetAt: firstReset }); + expect(notices).toHaveLength(1); + + observeCodexLowQuota(ACCOUNT_B, { weeklyPercent: 80, weeklyResetAt: firstReset }); + observeCodexLowQuota(ACCOUNT_A, { shortPercent: 80, shortResetAt: firstReset }); + expect(notices.map(notice => notice.window)).toEqual(["weekly", "weekly", "short"]); + + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90, weeklyResetAt: secondReset }); + expect(notices).toHaveLength(4); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 95, weeklyResetAt: secondReset }); + expect(notices).toHaveLength(4); + + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 10, weeklyResetAt: secondReset }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90, weeklyResetAt: secondReset }); + expect(notices).toHaveLength(5); + expect(notices.at(-1)).toEqual({ window: "weekly", percentUsed: 90, threshold: 80 }); + observeCodexLowQuota(ACCOUNT_A, { shortPercent: 95, shortResetAt: firstReset }); + observeCodexLowQuota(ACCOUNT_B, { weeklyPercent: 95, weeklyResetAt: firstReset }); + expect(notices).toHaveLength(5); + expect(config.pausedCodexAccountIds).toBeUndefined(); + }); + + test("ignores credits-only and expired observations", () => { + const config = configWith(protection({ actions: { pause: true, notify: true } })); + const persisted: OcxConfig[] = []; + const notices: Notice[] = []; + register(config, { + persist: next => persisted.push(structuredClone(next)), + notify: async notice => { notices.push(notice); }, + }); + + observeCodexLowQuota(ACCOUNT_A, { resetCredits: 1 }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 100, weeklyResetAt: Date.now() - 1 }); + observeCodexLowQuota(ACCOUNT_A, { shortPercent: 100, shortResetAt: Math.floor(Date.now() / 1000) - 1 }); + + expect(config.pausedCodexAccountIds).toBeUndefined(); + expect(persisted).toEqual([]); + expect(notices).toEqual([]); + }); + + test("does not treat carried usage in a credits-only quota write as a fresh observation", async () => { + setAccountQuotaFromParsed(ACCOUNT_A, { weeklyPercent: 99 }); + const config = configWith(protection({ actions: { pause: true, notify: true } })); + const persisted: OcxConfig[] = []; + const notices: Notice[] = []; + const registration = register(config, { + persist: next => persisted.push(structuredClone(next)), + notify: async notice => { notices.push(notice); }, + }); + + setAccountQuotaFromParsed(ACCOUNT_A, { resetCredits: 1 }); + + expect(config.pausedCodexAccountIds).toBeUndefined(); + expect(persisted).toEqual([]); + expect(notices).toEqual([]); + const accepted = { weeklyPercent: 99 }; + setAccountQuotaFromParsed(ACCOUNT_A, accepted, undefined, undefined, accepted); + await registration.flush(); + expect(persisted[0]?.pausedCodexAccountIds).toEqual([ACCOUNT_A]); + expect(notices).toEqual([{ window: "weekly", percentUsed: 99, threshold: 80 }]); + }); + + test("stops acting after unregister", () => { + const config = configWith(protection({ actions: { pause: true, notify: true } })); + const persisted: OcxConfig[] = []; + const notices: Notice[] = []; + const unregister = register(config, { + persist: next => persisted.push(structuredClone(next)), + notify: async notice => { notices.push(notice); }, + }); + + unregister(); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 100 }); + + expect(config.pausedCodexAccountIds).toBeUndefined(); + expect(persisted).toEqual([]); + expect(notices).toEqual([]); + }); + + test("leaves protection inactive when disabled or absent", () => { + for (const config of [ + configWith(protection({ enabled: false, actions: { pause: true, notify: true } })), + configWith(), + ]) { + const persisted: OcxConfig[] = []; + const notices: Notice[] = []; + register(config, { + persist: next => persisted.push(structuredClone(next)), + notify: async notice => { notices.push(notice); }, + }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 100 }); + expect(config.pausedCodexAccountIds).toBeUndefined(); + expect(persisted).toEqual([]); + expect(notices).toEqual([]); + cleanups.pop()?.(); + } + }); + + test("main-account observations never enter pool low-quota protection", () => { + const config = configWith(protection({ actions: { pause: true, notify: true } })); + let writes = 0; + const registration = register(config, { persist: () => { writes++; } }); + observeCodexLowQuota(MAIN_CODEX_ACCOUNT_ID, { shortPercent: 95, weeklyPercent: 95 }); + expect(config.pausedCodexAccountIds).toBeUndefined(); + expect(registration.listEvents()).toEqual([]); + expect(writes).toBe(0); + }); + + test("rejects unknown low-quota policy, action, and window keys", () => { + const policy = protection(); + expect(validateConfigCandidate(configWith(policy)).ok).toBe(true); + const candidates = [ + { name: "policy", value: { ...policy, unexpected: true } }, + { name: "actions", value: { ...policy, actions: { ...policy.actions, unexpected: true } } }, + { name: "windows", value: { ...policy, windows: { ...policy.windows, unexpected: true } } }, + ]; + + for (const candidate of candidates) { + const result = validateConfigCandidate({ + ...configWith(), + codexPool: { lowQuotaProtection: candidate.value }, + }); + expect(result.ok, candidate.name).toBe(false); + } + for (const candidate of [ + { ...policy, threshold: 0 }, + { ...policy, actions: { pause: false, notify: false } }, + { ...policy, windows: { short: false, weekly: false } }, + ]) { + expect(validateConfigCandidate({ ...configWith(), codexPool: { lowQuotaProtection: candidate } }).ok).toBe(false); + } + }); + + test("a manual resume suppresses repause until recovery or a new reset", () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + register(config, { persist: () => {} }); + const firstReset = Date.now() + 60_000; + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 85, weeklyResetAt: firstReset }); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + setCodexAccountPaused(config, ACCOUNT_A, false); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90, weeklyResetAt: firstReset }); + expect(config.pausedCodexAccountIds).toBeUndefined(); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 50, weeklyResetAt: firstReset }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90, weeklyResetAt: firstReset }); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + setCodexAccountPaused(config, ACCOUNT_A, false); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90, weeklyResetAt: firstReset + 60_000 }); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + }); + + test("both active windows honor one manual resume until a new reset episode", () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + register(config, { persist: () => {} }); + const reset = Date.now() + 60_000; + const high = { shortPercent: 90, shortResetAt: reset, weeklyPercent: 90, weeklyResetAt: reset }; + observeCodexLowQuota(ACCOUNT_A, high); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + setCodexAccountPaused(config, ACCOUNT_A, false); + observeCodexLowQuota(ACCOUNT_A, high); + observeCodexLowQuota(ACCOUNT_A, { ...high, shortPercent: 95, weeklyPercent: 95 }); + expect(config.pausedCodexAccountIds).toBeUndefined(); + observeCodexLowQuota(ACCOUNT_A, { ...high, shortResetAt: reset + 60_000 }); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + }); + + test("invalid raw WHAM usage stays display-only while valid usage pauses", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + const registration = register(config, { persist: () => {} }); + const generation = saveCodexAccountCredential(ACCOUNT_A, { + accessToken: "fixture-access", refreshToken: "fixture-refresh", + expiresAt: Date.now() + 60_000, chatgptAccountId: "fixture-chatgpt-account", + }); + const publish = (usedPercent: number) => commitPoolQuotaResponse( + new Response(JSON.stringify({ rate_limit: { primary_window: { + used_percent: usedPercent, limit_window_seconds: 604_800, + } } }), { status: 200 }), + { accountId: ACCOUNT_A, existing: null, configuredPlan: "plus", generation, + writerGeneration: captureConfigGeneration() }, + ); + await publish(150); + expect(getAccountQuota(ACCOUNT_A)?.weeklyPercent).toBe(100); + expect(config.pausedCodexAccountIds).toBeUndefined(); + await publish(90); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + await registration.flush(); + }); + + test("invalid raw response-header usage stays display-only while valid usage pauses", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + const registration = register(config, { persist: () => {} }); + const headers = (usedPercent: number) => new Headers({ + "x-codex-primary-used-percent": String(usedPercent), + "x-codex-primary-window-minutes": "10080", + }); + applyAccountQuotaFromUpstreamHeaders(ACCOUNT_A, headers(150)); + expect(getAccountQuota(ACCOUNT_A)?.weeklyPercent).toBe(100); + expect(config.pausedCodexAccountIds).toBeUndefined(); + applyAccountQuotaFromUpstreamHeaders(ACCOUNT_A, headers(90)); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + await registration.flush(); + }); + + test("invalid raw legacy weekly usage stays display-only while valid usage pauses", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + const registration = register(config, { persist: () => {} }); + updateAccountQuota(ACCOUNT_A, 150); + expect(getAccountQuota(ACCOUNT_A)?.weeklyPercent).toBe(100); + expect(config.pausedCodexAccountIds).toBeUndefined(); + updateAccountQuota(ACCOUNT_A, 90); + expect(config.pausedCodexAccountIds).toContain(ACCOUNT_A); + await registration.flush(); + }); + + test("notice failure retries on a later observation and only then reports delivery", async () => { + const config = configWith(protection({ actions: { pause: false, notify: true } })); + let calls = 0; + const registration = register(config, { notify: async () => { if (++calls === 1) throw new Error("sink failed"); } }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 80 }); + await Promise.resolve(); + expect(registration.listEvents(1)[0]?.status).toBe("failed"); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 82 }); + await Promise.resolve(); + expect(registration.listEvents(1)[0]?.status).toBe("delivered"); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90 }); + expect(calls).toBe(2); + expect(registration.listEvents(100).every(event => Object.keys(event).sort().join(",") === + "accountId,delivery,percentUsed,resetAt,status,timestamp,window")).toBe(true); + }); + + test("the default headless alert is logged, never reported as delivered", () => { + const config = configWith(protection({ actions: { pause: false, notify: true } })); + const registration = register(config); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 85 }); + expect(registration.listEvents(1)[0]).toMatchObject({ + accountId: ACCOUNT_A, delivery: "notice", status: "logged", + }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90 }); + expect(registration.listEvents(100).filter(event => event.delivery === "notice")).toHaveLength(1); + }); + + test("a blocked save times out flush and later work is fenced", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + let resolveSave: (() => void) | undefined; + let writes = 0; + const registration = registerCodexLowQuotaProtection(config, { + persist: async () => { writes++; await new Promise(resolve => { resolveSave = resolve; }); }, + }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90 }); + await registration.flush(); + expect(writes).toBe(1); + expect(registration.listEvents(1)[0]?.status).toBe("pending"); + resolveSave?.(); + const deadline = Date.now() + 1_000; + while (registration.listEvents(1)[0]?.status !== "succeeded" && Date.now() < deadline) { + await Bun.sleep(1); + } + expect(registration.listEvents(1)[0]?.status).toBe("succeeded"); + expect(registration.listEvents(100).filter(event => event.status === "cancelled")).toHaveLength(0); + observeCodexLowQuota(ACCOUNT_B, { weeklyPercent: 90 }); + await new Promise(resolve => setTimeout(resolve, 300)); + expect(writes).toBe(1); + }); + + test("closing before the queued save starts cancels it without a config write", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + let writes = 0; + const registration = register(config, { persist: () => { writes++; } }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90 }); + registration(); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(writes).toBe(0); + expect(registration.listEvents(1)[0]?.status).toBe("cancelled"); + }); + + test("a failed deferred save retries and publishes durable status", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + let writes = 0; + let saved: (() => void) | undefined; + const succeeded = new Promise(resolve => { saved = resolve; }); + const registration = register(config, { persist: () => { + if (++writes === 1) throw new Error("transient save failure"); + saved?.(); + } }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 88 }); + await Promise.race([succeeded, new Promise((_, reject) => setTimeout(() => reject(new Error("save retry timeout")), 1_000))]); + await registration.flush(); + expect(writes).toBe(2); + expect(registration.listEvents(100).map(event => event.status)).toContain("failed"); + expect(registration.listEvents(1)[0]?.status).toBe("succeeded"); + }); + + test("fan-out isolates throwing observers and independent server configs", async () => { + const older = configWith(protection({ actions: { pause: false, notify: true } })); + older.codexAccounts = [account(ACCOUNT_A)]; + const newer = configWith(protection({ actions: { pause: true, notify: false } })); + newer.codexAccounts = [account(ACCOUNT_B)]; + const notices: Notice[] = []; + const first = register(older, { notify: notice => { notices.push(notice); } }); + const throwing = registerLowQuotaObserver(() => { throw new Error("isolated"); }); + cleanups.push(throwing); + const second = register(newer, { persist: () => {} }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 85 }); + expect(newer.pausedCodexAccountIds).toBeUndefined(); + expect(notices).toHaveLength(1); + observeCodexLowQuota(ACCOUNT_B, { weeklyPercent: 85 }); + expect(newer.pausedCodexAccountIds).toEqual([ACCOUNT_B]); + await second.flush(); + first(); + throwing(); + observeCodexLowQuota(ACCOUNT_A, { shortPercent: 85 }); + expect(notices).toHaveLength(1); + }); + + test("the event ledger retains only its newest hundred sanitized entries", () => { + const ledger = createLowQuotaEventLedger(); + for (let i = 0; i < 120; i++) { + ledger.publish({ accountId: `account-${i}`, window: "weekly", percentUsed: 80, + resetAt: null, timestamp: i, status: "succeeded", delivery: "pause-save" }); + } + const events = ledger.list(999); + expect(events).toHaveLength(100); + expect(events[0]?.accountId).toBe("account-119"); + expect(events.at(-1)?.accountId).toBe("account-20"); + }); + + test("one coalesced save reports durability for both paused accounts", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + let writes = 0; + const registration = register(config, { persist: () => { writes++; } }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90 }); + observeCodexLowQuota(ACCOUNT_B, { weeklyPercent: 90 }); + expect(writes).toBe(0); + await registration.flush(); + expect(writes).toBe(1); + expect(new Set(registration.listEvents(100).filter(event => event.status === "succeeded") + .map(event => event.accountId))).toEqual(new Set([ACCOUNT_A, ACCOUNT_B])); + }); + + test("a manual resume before the deferred save is not reported as a durable pause", async () => { + const config = configWith(protection({ actions: { pause: true, notify: false } })); + const registration = register(config, { persist: () => {} }); + observeCodexLowQuota(ACCOUNT_A, { weeklyPercent: 90 }); + setCodexAccountPaused(config, ACCOUNT_A, false); + await registration.flush(); + expect(registration.listEvents(1)[0]).toMatchObject({ accountId: ACCOUNT_A, delivery: "pause-save", status: "cancelled" }); + }); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index aad4d16599e..c0b2cae048f 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -3,6 +3,7 @@ "pnpm-command-isolation.test.ts": "update", "project-config-warning-snapshot.test.ts": "codex-integration", "codex-quota-auto-refresh-generation.test.ts": "codex-integration", + "low-quota-protection.test.ts": "codex-integration", "provider-antigravity-quota-retry.test.ts": "providers", "responses-compaction-recovery.test.ts": "responses", "compaction-recovery-settings.test.ts": "config", diff --git a/tests/server/server-background-lifecycle.test.ts b/tests/server/server-background-lifecycle.test.ts index 8d402e39f09..58e350ab362 100644 --- a/tests/server/server-background-lifecycle.test.ts +++ b/tests/server/server-background-lifecycle.test.ts @@ -14,10 +14,13 @@ import { } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { saveConfig } from "../../src/config"; +import { loadConfig, saveConfig } from "../../src/config"; +import { observeCodexLowQuota } from "../../src/codex/low-quota-observer"; +import { MAIN_CODEX_ACCOUNT_ID } from "../../src/codex/account-id"; import { startServer, type StartServerDeps } from "../../src/server"; import { registerStateSweepAfterTick } from "../../src/lib/state-store-sweeper"; import { getActiveMemoryWatchdog } from "../../src/server/memory-watchdog"; +import { acquireServerBackgroundLifecycle } from "../../src/server/background-lifecycle"; import { getStorageCleanupPolicyJobState, requestStorageCleanupPolicyRun, @@ -291,6 +294,134 @@ afterEach(async () => { }); describe("server background lifecycle", () => { + test("a server without low-quota work releases its process owner synchronously", async () => { + const lease = acquireServerBackgroundLifecycle(() => {}, baseConfig()); + const completion = lease.release(); + try { + expect(getActiveMemoryWatchdog()).toBeNull(); + } finally { + await completion; + } + }); + + test("authenticated low-quota history is bounded and an unauthenticated reader is refused", async () => { + const config = baseConfig(); + config.codexAccounts = [{ id: "low-quota-pool", email: "pool@example.com", isMain: false }]; + config.codexPool = { lowQuotaProtection: { + enabled: true, threshold: 80, + windows: { short: false, weekly: true }, actions: { pause: false, notify: true }, + } }; + saveConfig(config); + const server = trackedStart(); + observeCodexLowQuota("low-quota-pool", { weeklyPercent: 85, weeklyResetAt: Date.now() + 60_000 }); + await Promise.resolve(); + const url = new URL("/api/codex-auth/low-quota-events?limit=1", server.url); + const refused = await fetch(url); + expect(refused.status).toBe(401); + const allowed = await managementFetch(url); + expect(allowed.status).toBe(200); + const body = await allowed.json() as { events: Array<{ accountId: string; window: string; percentUsed: number; status: string }> }; + expect(body.events).toHaveLength(1); + expect(body.events[0]).toMatchObject({ accountId: "low-quota-pool", window: "weekly", percentUsed: 85, status: "logged" }); + expect((await managementFetch(new URL("/api/codex-auth/low-quota-events?limit=oops", server.url))).status).toBe(400); + await stopTracked(server); + }); + + test("two live owners observe independently and either stop order preserves the survivor", async () => { + const config = baseConfig(); + config.codexAccounts = [{ id: "low-quota-pool", email: "pool@example.com", isMain: false }]; + config.codexPool = { lowQuotaProtection: { + enabled: true, threshold: 80, + windows: { short: false, weekly: true }, actions: { pause: false, notify: true }, + } }; + saveConfig(config); + const older = trackedStart(); + const newer = trackedStart(); + const reset = Date.now() + 60_000; + observeCodexLowQuota("low-quota-pool", { weeklyPercent: 85, weeklyResetAt: reset }); + await Promise.resolve(); + const eventsOf = async (server: StartedServer) => { + const response = await managementFetch(new URL("/api/codex-auth/low-quota-events", server.url)); + return (await response.json() as { events: Array<{ status: string }> }).events; + }; + expect((await eventsOf(older)).filter(event => event.status === "logged")).toHaveLength(1); + expect((await eventsOf(newer)).filter(event => event.status === "logged")).toHaveLength(1); + await stopTracked(older); + observeCodexLowQuota("low-quota-pool", { weeklyPercent: 85, weeklyResetAt: reset + 60_000 }); + await Promise.resolve(); + expect((await eventsOf(newer)).filter(event => event.status === "logged")).toHaveLength(2); + await stopTracked(newer); + }); + + test("each server GET excludes another live owner's account events", async () => { + const firstConfig = baseConfig(); + firstConfig.codexAccounts = [{ id: "low-quota-a", email: "a@example.com", isMain: false }]; + firstConfig.codexPool = { lowQuotaProtection: { + enabled: true, threshold: 80, + windows: { short: false, weekly: true }, actions: { pause: false, notify: true }, + } }; + saveConfig(firstConfig); + const first = trackedStart(); + const secondConfig = baseConfig(); + secondConfig.codexAccounts = [{ id: "low-quota-b", email: "b@example.com", isMain: false }]; + secondConfig.codexPool = firstConfig.codexPool; + saveConfig(secondConfig); + const second = trackedStart(); + observeCodexLowQuota("low-quota-a", { weeklyPercent: 85 }); + observeCodexLowQuota("low-quota-b", { weeklyPercent: 85 }); + const accountIds = async (server: StartedServer) => { + const response = await managementFetch(new URL("/api/codex-auth/low-quota-events", server.url)); + expect(response.status).toBe(200); + const body = await response.json() as { events: Array<{ accountId: string }> }; + return body.events.map(event => event.accountId); + }; + await expect(accountIds(first)).resolves.toEqual(["low-quota-a"]); + await expect(accountIds(second)).resolves.toEqual(["low-quota-b"]); + await stopTracked(second); + await stopTracked(first); + }); + + test("a newer disabled server cannot suppress the older low-quota owner", async () => { + const enabled = baseConfig(); + enabled.codexAccounts = [{ id: "low-quota-pool", email: "pool@example.com", isMain: false }]; + enabled.codexPool = { lowQuotaProtection: { + enabled: true, threshold: 80, + windows: { short: false, weekly: true }, actions: { pause: false, notify: true }, + } }; + saveConfig(enabled); + const older = trackedStart(); + saveConfig(baseConfig()); + const newer = trackedStart(); + observeCodexLowQuota("low-quota-pool", { weeklyPercent: 85 }); + await Promise.resolve(); + const olderUrl = new URL("/api/codex-auth/low-quota-events", older.url); + const newerUrl = new URL("/api/codex-auth/low-quota-events", newer.url); + expect((await (await managementFetch(olderUrl)).json() as { events: unknown[] }).events).toHaveLength(1); + expect((await (await managementFetch(newerUrl)).json() as { events: unknown[] }).events).toHaveLength(0); + await stopTracked(newer); + observeCodexLowQuota("low-quota-pool", { weeklyPercent: 90, weeklyResetAt: Date.now() + 60_000 }); + await Promise.resolve(); + expect((await (await managementFetch(olderUrl)).json() as { events: unknown[] }).events).toHaveLength(2); + await stopTracked(older); + }); + test("low-quota pause survives reload and server stop unregisters protection", async () => { + const config = baseConfig(); + config.codexAccounts = [{ id: "low-quota-pool", email: "pool@example.com", isMain: false }]; + config.codexPool = { lowQuotaProtection: { + enabled: true, threshold: 80, + windows: { short: true, weekly: true }, actions: { pause: true, notify: false }, + } }; + saveConfig(config); + const server = trackedStart(); + observeCodexLowQuota("low-quota-pool", { shortPercent: 85 }); + expect(loadConfig().pausedCodexAccountIds).toBeUndefined(); + await stopTracked(server); + expect(loadConfig().pausedCodexAccountIds).toContain("low-quota-pool"); + // An account not previously paused would act if stop leaked the registration. + observeCodexLowQuota(MAIN_CODEX_ACCOUNT_ID, { weeklyPercent: 95 }); + expect(loadConfig().pausedCodexAccountIds).toEqual(["low-quota-pool"]); + }); + test("stopping the newer server preserves the older server's process-wide work", async () => { saveConfig(baseConfig()); const probe = installBackgroundTimerProbe();