diff --git a/devlog/_plan/260907_axis1_bugfixes/021_source_review.md b/devlog/_plan/260907_axis1_bugfixes/021_source_review.md new file mode 100644 index 0000000000..516a760813 --- /dev/null +++ b/devlog/_plan/260907_axis1_bugfixes/021_source_review.md @@ -0,0 +1,7 @@ +# wp1 source review + +Three bounded patches implemented with regression coverage. Hooke independently passed the physical-response quota observer wiring; Tesla independently passed quota/recovery security and source review with zero blockers. Version comparator and status/doctor projections inspected by main. All source workers report no local suite/typecheck/build execution. + +Quota source: #3809, Éverton Toffanetto; Co-authored-by included in f215f79b4. Version report: garysassano; Reported-by included in f91e3953a. Recovery report: Hu9956; Reported-by included in recovery commit. + +Source-only checks: git diff --check and documentation fence/whitespace inspection. These do not prove runtime correctness. wp2 final cumulative hosted CI is still mandatory. Final CI dispatch includes Windows because ordinary PR workflow omits it. No release/deploy workflow will be dispatched. diff --git a/docs-site/src/content/docs/ko/reference/cli/lifecycle.md b/docs-site/src/content/docs/ko/reference/cli/lifecycle.md index 4847614674..068807025b 100644 --- a/docs-site/src/content/docs/ko/reference/cli/lifecycle.md +++ b/docs-site/src/content/docs/ko/reference/cli/lifecycle.md @@ -82,6 +82,19 @@ dedicated-provider history도 포함됩니다. 상태를 백업하고 이 전체 ### `ocx status [--json]` +status와 `ocx doctor`는 현재 CLI와 실행 중인 프록시의 버전을 비교합니다. CLI가 더 새로우면 +원하는 최신 설치로 프록시를 재시작하십시오. 백그라운드 서비스라면 `ocx service repair`를 +실행합니다(`ocx service restart`는 별칭). 프록시가 더 새로우면 CLI를 업그레이드하거나 +`PATH`가 원하는 설치를 가리키도록 수정하십시오. 이 진단은 서비스를 복구하거나 요청 허용 +여부를 바꾸지 않습니다. + +버전 문자열이 같거나 어느 쪽이 `unknown` / `0.0.0`이면 경고하지 않으며, 프록시 버전이 없어도 +경고하지 않습니다. doctor는 placeholder를 버전 일치로 확정하지 않습니다. 엄격한 SemVer로 +해석할 수 없는 서로 다른 문자열이나 build metadata만 다른 버전은 어느 쪽이 오래됐다고 +단정하지 않는 중립 경고를 표시합니다. 공백을 제거하거나 앞의 `v`를 정규화하지 않습니다. +JSON의 `versionSkew`에도 같은 안내가 들어가며 필드는 `cliVersion`, `proxyVersion`, `skewed`, +`warning` 그대로입니다. + 읽기 전용 진단 요약을 출력합니다. 프록시 PID, `/healthz` 도달 가능 여부, 대시보드 URL, 설정 경로, 기본 공급자, Codex 자동 시작 설정, 서비스 상태, shim 상태, 그리고 마스킹된 실제로 적용되는 Codex 홈이 포함됩니다. 명시적이고 높은 신뢰도의 Windows Orca 런타임 홈 시그니처만 diff --git a/docs-site/src/content/docs/reference/architecture.md b/docs-site/src/content/docs/reference/architecture.md index 8e1e361a81..e0fbc8bcba 100644 --- a/docs-site/src/content/docs/reference/architecture.md +++ b/docs-site/src/content/docs/reference/architecture.md @@ -233,7 +233,17 @@ response is not cacheable. Post-commit and 5xx errors keep the no-resend path. When encrypted agent-task recovery refuses a routed task, its existing 400 error can include a bounded `recovery_reason`: `unsupported_envelope`, -`admission_denied`, `recovery_unavailable`, `caller_cancelled`, or `input_changed`. -The field is omitted when no classified recovery result exists. +`admission_denied`, `recovery_unavailable`, `caller_cancelled`, `input_changed`, +`recovery_http_rejected`, `recovery_timeout`, `recovery_aborted`, +`recovery_transport_error`, or `recovery_invalid_output`. +HTTP rejection requires an observed non-success response. Invalid output includes +invalid UTF-8, oversized bodies, malformed or incomplete recovery streams, and +invalid or conflicting assignments. A caller's cancellation takes precedence over +an owned deadline, which takes precedence over decode/transport failures. +`recovery_aborted` describes a shared recovery cancelled independently of that caller. +Shared-flight waiters receive the same underlying failure unless individually cancelled; +only successful plaintext is cached. Diagnostics contain no upstream error or payload text. +The field is omitted when no classified recovery result exists, and existing combo +branches that return the original target failure keep that response. `recovery_unavailable` includes cache/singleflight capacity and does not prove an upstream request was attempted. No retry or broader envelope acceptance is enabled. diff --git a/docs-site/src/content/docs/reference/cli/lifecycle.md b/docs-site/src/content/docs/reference/cli/lifecycle.md index e75a2b6241..0dda487b3a 100644 --- a/docs-site/src/content/docs/reference/cli/lifecycle.md +++ b/docs-site/src/content/docs/reference/cli/lifecycle.md @@ -88,6 +88,19 @@ are left in place. ### `ocx status [--json]` +Status and `ocx doctor` compare this CLI's version with the running proxy. If the CLI is newer, +restart the proxy using the intended current installation; for a background service, run +`ocx service repair` (`ocx service restart` is an alias). If the proxy is newer, upgrade the CLI +or resolve `PATH` to the intended installation. These diagnostics do not repair the service or +change whether requests are allowed. + +Identical version strings and the `unknown` / `0.0.0` placeholders suppress the warning, as does +an absent proxy version. Doctor does not report placeholders as a confirmed match. Different +strings still produce a neutral warning when they cannot be strictly parsed as SemVer or differ +only in build metadata; neither side is called older. Versions are not trimmed and a leading `v` +is not normalized. JSON exposes the same advice in `versionSkew`, whose fields remain +`cliVersion`, `proxyVersion`, `skewed`, and `warning`. + Print a read-only diagnostic summary: proxy PID, `/healthz` reachability, dashboard URL, config path, default provider, Codex autostart setting, service state, shim state, and the redacted effective Codex home. Only the explicit, high-confidence Windows Orca runtime-home signature adds an actionable App-home @@ -261,9 +274,10 @@ bundled Bun paths are deliberately rediscovered after upgrades instead of being Definitions installed before this change still carry the old versioned paths and cannot migrate themselves — once the old executable is deleted, no opencodex code runs to fix it. Run `ocx service repair` once after upgrading; after that, each service start follows the launcher. -An already-running proxy is not replaced by an external upgrade: restart the service (or run -`ocx service repair`) so the new build serves, and treat a CLI/proxy version mismatch warning as -exactly that signal. +An already-running proxy is not replaced by an external upgrade: when the installed CLI is newer +than the running proxy, restart the service (or run `ocx service repair`) so the new build serves. +If the proxy is newer instead, check the CLI installation and `PATH` as described under +[`ocx status`](#ocx-status---json). | Subcommand | Action | | --- | --- | diff --git a/docs-site/src/content/docs/ru/reference/cli/lifecycle.md b/docs-site/src/content/docs/ru/reference/cli/lifecycle.md index 1ace7cc10f..7be5d5ad77 100644 --- a/docs-site/src/content/docs/ru/reference/cli/lifecycle.md +++ b/docs-site/src/content/docs/ru/reference/cli/lifecycle.md @@ -89,6 +89,19 @@ ocx eject back ### `ocx status [--json]` +Status и `ocx doctor` сравнивают версии текущего CLI и работающего прокси. Если CLI новее, +перезапустите прокси из нужной актуальной установки. Для фоновой службы используйте +`ocx service repair` (`ocx service restart` — её псевдоним). Если новее прокси, обновите CLI +или исправьте `PATH`, чтобы он указывал на нужную установку. Диагностика не ремонтирует службу +и не меняет разрешение запросов. + +При одинаковых строках версий, значениях `unknown` / `0.0.0` или отсутствии версии прокси +предупреждение подавляется. Doctor не считает placeholder подтверждённым совпадением. +Разные строки, которые нельзя строго разобрать как SemVer, и версии, отличающиеся только +build metadata, вызывают нейтральное предупреждение без указания устаревшей стороны. +Пробелы не удаляются, префикс `v` не нормализуется. JSON содержит ту же рекомендацию в +`versionSkew` с прежними полями `cliVersion`, `proxyVersion`, `skewed` и `warning`. + Печатает read-only диагностическую сводку: PID прокси, достижимость `/healthz`, URL дашборда, путь к конфигу, провайдера по умолчанию, настройку автозапуска Codex, состояние службы, состояние shim'а и redacted effective Codex home. Только явная и высокоуверенная сигнатура mismatch diff --git a/src/cli/doctor.ts b/src/cli/doctor.ts index 1ab4fe9f1b..d7148530a0 100644 --- a/src/cli/doctor.ts +++ b/src/cli/doctor.ts @@ -1157,11 +1157,11 @@ export async function runDoctor(args: string[] = []): Promise { // No extra probe -- findLiveProxy already carried the version back. { const { packageVersion } = await import("./help"); - const { computeVersionSkew } = await import("./version-skew"); + const { computeVersionSkew, isConfirmedVersionMatch } = await import("./version-skew"); const skew = computeVersionSkew(packageVersion(), live?.version); if (skew.skewed && skew.warning) { console.log(`!! ${skew.warning}`); - } else if (skew.proxyVersion !== null) { + } else if (isConfirmedVersionMatch(skew)) { console.log(`ok ocx ${skew.cliVersion} matches the running proxy`); } } diff --git a/src/cli/version-skew.ts b/src/cli/version-skew.ts index 588d29a307..48b71a51ee 100644 --- a/src/cli/version-skew.ts +++ b/src/cli/version-skew.ts @@ -1,5 +1,5 @@ /** - * CLI-versus-proxy version skew (#2701). + * CLI-versus-proxy version skew (#2701, #3464). * * The reported failure: `ocx` on PATH is an older install than the running proxy, so its * help describes commands the proxy does not have and its output describes a different @@ -9,6 +9,7 @@ * comparison instead of reimplementing it -- two diagnostics disagreeing about whether an * install is stale would be worse than neither reporting it. */ +import { parseStrictSemver, type StrictSemver } from "../lib/strict-semver"; /** Placeholder versions that mean "unknown", not "different". */ const PLACEHOLDERS = new Set(["unknown", "0.0.0"]); @@ -22,6 +23,30 @@ export interface VersionSkew { readonly warning: string | null; } +/** Suppressed comparisons are not confirmed matches, even when both placeholders agree. */ +export function isConfirmedVersionMatch(skew: VersionSkew): boolean { + return skew.proxyVersion === skew.cliVersion && !PLACEHOLDERS.has(skew.cliVersion); +} + +/** SemVer precedence ignores build metadata; raw equality is handled separately. */ +function compareVersions(cli: StrictSemver, proxy: StrictSemver): number { + for (let i = 0; i < cli.core.length; i++) { + if (cli.core[i]! !== proxy.core[i]!) return cli.core[i]! > proxy.core[i]! ? 1 : -1; + } + if (cli.prerelease.length === 0) return proxy.prerelease.length === 0 ? 0 : 1; + if (proxy.prerelease.length === 0) return -1; + for (let i = 0; i < Math.max(cli.prerelease.length, proxy.prerelease.length); i++) { + const left = cli.prerelease[i]; + const right = proxy.prerelease[i]; + if (left === right) continue; + if (left === undefined) return -1; + if (right === undefined) return 1; + if (typeof left !== typeof right) return typeof left === "bigint" ? -1 : 1; + return left > right ? 1 : -1; + } + return 0; +} + /** * Compare the running CLI against the live proxy. * @@ -36,11 +61,19 @@ export function computeVersionSkew(cliVersion: string, proxyVersion: string | un if (proxy === null || PLACEHOLDERS.has(proxy) || PLACEHOLDERS.has(cliVersion) || proxy === cliVersion) { return { cliVersion, proxyVersion: proxy, skewed: false, warning: null }; } + const cliSemver = parseStrictSemver(cliVersion); + const proxySemver = parseStrictSemver(proxy); + const order = cliSemver && proxySemver ? compareVersions(cliSemver, proxySemver) : 0; + const advice = order > 0 + ? "the running proxy is older than this CLI. Restart the proxy using the intended current installation. " + + "For a background service, run ocx service repair (ocx service restart is an alias)." + : order < 0 + ? "this ocx on PATH is older than the running proxy. Upgrade the CLI or resolve PATH to the intended installation." + : "the versions differ, but neither can be identified as older. Check which installations the CLI and proxy use."; return { cliVersion, proxyVersion: proxy, skewed: true, - warning: `CLI ${cliVersion} does not match the running proxy ${proxy} — this ocx on PATH is stale. ` - + "Its help and features describe a different build. Reinstall, or run the proxy's own binary.", + warning: `CLI ${cliVersion} does not match the running proxy ${proxy} — ${advice}`, }; } diff --git a/src/lib/bounded-body.ts b/src/lib/bounded-body.ts index 4016a0a753..0975268560 100644 --- a/src/lib/bounded-body.ts +++ b/src/lib/bounded-body.ts @@ -212,13 +212,28 @@ export async function readBoundedResponseBytes( } } -function decodeUtf8(chunks: readonly Uint8Array[], fatal: boolean): string { +// Mark only exceptions thrown by our decoder, preserving their identity and TypeError contract. +// Timeout-path flushing may fail too; retain that origin so callers do not lose the deadline. +const decodeFailures = new WeakMap(); + +export function boundedBodyDecodeFailure(error: unknown): "invalid_utf8" | "timeout" | undefined { + return error !== null && typeof error === "object" ? decodeFailures.get(error) : undefined; +} + +function decodeUtf8(chunks: readonly Uint8Array[], fatal: boolean, timedOut = false): string { const decoder = new TextDecoder("utf-8", { fatal }); - let text = ""; - for (const chunk of chunks) text += decoder.decode(chunk, { stream: true }); - // Flush an incomplete trailing UTF-8 sequence deterministically. - text += decoder.decode(); - return text; + try { + let text = ""; + for (const chunk of chunks) text += decoder.decode(chunk, { stream: true }); + // Flush an incomplete trailing UTF-8 sequence deterministically. + text += decoder.decode(); + return text; + } catch (error) { + if (error !== null && typeof error === "object") { + decodeFailures.set(error, timedOut ? "timeout" : "invalid_utf8"); + } + throw error; + } } /** @@ -297,7 +312,7 @@ export async function readBoundedResponseBody( "TimeoutError", ); return { - text: decodeUtf8([retained.subarray(0, retainedBytes)], options.fatalUtf8 === true), + text: decodeUtf8([retained.subarray(0, retainedBytes)], options.fatalUtf8 === true, true), truncated: true, timedOut: true, totalTimedOut: outcome === TOTAL_TIMEOUT, diff --git a/src/server/responses/agent-task-recovery-cache.ts b/src/server/responses/agent-task-recovery-cache.ts index 93d0c1778b..398a0feba4 100644 --- a/src/server/responses/agent-task-recovery-cache.ts +++ b/src/server/responses/agent-task-recovery-cache.ts @@ -2,6 +2,20 @@ const MAX_CACHE_BYTES = 8 * 1024 * 1024; const MAX_CONCURRENT_RECOVERIES = 32; const CACHE_TTL_MS = 15 * 60 * 1000; +export type AgentTaskRecoveryResolutionFailureReason = + | "recovery_unavailable" + | "caller_cancelled" + | "recovery_http_rejected" + | "recovery_timeout" + | "recovery_aborted" + | "recovery_transport_error" + | "recovery_invalid_output"; + +/** Shared flights carry bounded failures; only successful plaintext enters the cache. */ +export type AgentTaskRecoveryResolution = + | { readonly recovered: true; readonly assignment: string } + | { readonly recovered: false; readonly reason: AgentTaskRecoveryResolutionFailureReason }; + interface RecoveryCacheEntry { assignment: string; bytes: number; @@ -11,7 +25,7 @@ interface RecoveryCacheEntry { interface RecoveryFlight { controller: AbortController; - promise: Promise; + promise: Promise; waiters: number; settled: boolean; } @@ -63,7 +77,7 @@ function insertRecoveryCacheEntry(key: string, assignment: string, maxEntries: n function startRecoveryFlight( key: string, maxEntries: number, - request: (signal: AbortSignal) => Promise, + request: (signal: AbortSignal) => Promise, ): RecoveryFlight | null { const active = RECOVERY_FLIGHTS.get(key); if (active) return active; @@ -72,15 +86,15 @@ function startRecoveryFlight( const controller = new AbortController(); const flight: RecoveryFlight = { controller, - promise: Promise.resolve(null), + promise: Promise.resolve({ recovered: false, reason: "recovery_unavailable" }), waiters: 0, settled: false, }; flight.promise = request(controller.signal) - .then((assignment) => { - if (!assignment || controller.signal.aborted) return null; - insertRecoveryCacheEntry(key, assignment, maxEntries); - return assignment; + .then((result): AgentTaskRecoveryResolution => { + if (controller.signal.aborted) return { recovered: false, reason: "recovery_aborted" }; + if (result.recovered) insertRecoveryCacheEntry(key, result.assignment, maxEntries); + return result; }) .finally(() => { flight.settled = true; @@ -93,14 +107,14 @@ function startRecoveryFlight( async function waitForRecoveryFlight( flight: RecoveryFlight, abortSignal?: AbortSignal, -): Promise { - if (abortSignal?.aborted) return null; +): Promise { + if (abortSignal?.aborted) return { recovered: false, reason: "caller_cancelled" }; flight.waiters += 1; let onAbort: (() => void) | undefined; try { if (!abortSignal) return await flight.promise; - const cancelled = new Promise((resolve) => { - onAbort = () => resolve(null); + const cancelled = new Promise((resolve) => { + onAbort = () => resolve({ recovered: false, reason: "caller_cancelled" }); abortSignal.addEventListener("abort", onAbort, { once: true }); if (abortSignal.aborted) onAbort(); }); @@ -120,12 +134,27 @@ export async function resolveCachedAgentTaskRecovery( request: (signal: AbortSignal) => Promise, abortSignal?: AbortSignal, ): Promise { - if (abortSignal?.aborted) return null; + const result = await resolveCachedAgentTaskRecoveryWithResult(key, maxEntries, async signal => { + const assignment = await request(signal); + return assignment + ? { recovered: true, assignment } + : { recovered: false, reason: "recovery_unavailable" }; + }, abortSignal); + return result.recovered ? result.assignment : null; +} + +export async function resolveCachedAgentTaskRecoveryWithResult( + key: string, + maxEntries: number, + request: (signal: AbortSignal) => Promise, + abortSignal?: AbortSignal, +): Promise { + if (abortSignal?.aborted) return { recovered: false, reason: "caller_cancelled" }; sweepRecoveryCache(Date.now(), maxEntries); const cached = RECOVERY_CACHE.get(key)?.assignment; - if (cached) return cached; + if (cached) return { recovered: true, assignment: cached }; const flight = startRecoveryFlight(key, maxEntries, request); - return flight ? waitForRecoveryFlight(flight, abortSignal) : null; + return flight ? waitForRecoveryFlight(flight, abortSignal) : { recovered: false, reason: "recovery_unavailable" }; } export function discardCachedAgentTaskRecovery(key: string): void { diff --git a/src/server/responses/agent-task-recovery.ts b/src/server/responses/agent-task-recovery.ts index 22b7a4e66b..a15a2563ca 100644 --- a/src/server/responses/agent-task-recovery.ts +++ b/src/server/responses/agent-task-recovery.ts @@ -1,14 +1,16 @@ import { createHash, createHmac, randomBytes } from "node:crypto"; import { decodeJwtPayload, extractAccountId } from "../../oauth/chatgpt"; import type { OcxConfig } from "../../types"; -import { readBoundedResponseBody } from "../../lib/bounded-body"; +import { boundedBodyDecodeFailure, readBoundedResponseBody } from "../../lib/bounded-body"; import { isApiAuthRequired, isProxyAdmissionSecret } from "../auth-cors"; import { structurallyValidFernetTokens } from "./encrypted-payload"; import { cachedAgentTaskRecovery, discardCachedAgentTaskRecovery, resetAgentTaskRecoveryCache, - resolveCachedAgentTaskRecovery, + resolveCachedAgentTaskRecoveryWithResult, + type AgentTaskRecoveryResolution, + type AgentTaskRecoveryResolutionFailureReason, } from "./agent-task-recovery-cache"; /** Experimental opt-in normalization through ChatGPT's fixed Codex endpoint. */ @@ -44,9 +46,8 @@ export interface AgentTaskRecoveryOptions { export type AgentTaskRecoveryFailureReason = | "unsupported_envelope" | "admission_denied" - // Includes cache capacity rejection; does not imply an upstream request was attempted. - | "recovery_unavailable" - | "caller_cancelled" + // recovery_unavailable includes capacity rejection, which does not imply an upstream attempt. + | AgentTaskRecoveryResolutionFailureReason | "input_changed"; export type AgentTaskRecoveryResult = @@ -436,7 +437,7 @@ async function requestRecovery( envelope: AgentEnvelope, options: AgentTaskRecoveryOptions, abortSignal?: AbortSignal, -): Promise { +): Promise { const controller = new AbortController(); const timeout = setTimeout( () => controller.abort(new DOMException("Agent task recovery timed out", "TimeoutError")), @@ -454,8 +455,11 @@ async function requestRecovery( redirect: "error", }); if (!response.ok) { - try { await response.body?.cancel(); } catch { /* already closed */ } - return null; + // A rejected or never-settling cancellation must not extend the recovery deadline. + try { void response.body?.cancel().catch(() => undefined); } catch { /* already closed */ } + if (abortSignal?.aborted) return { recovered: false, reason: "recovery_aborted" }; + if (controller.signal.aborted) return { recovered: false, reason: "recovery_timeout" }; + return { recovered: false, reason: "recovery_http_rejected" }; } const body = await readBoundedResponseBody(response, { signal, @@ -465,10 +469,18 @@ async function requestRecovery( inactivityTimeoutMs: options.timeoutMs ?? 45_000, firstByteTimeoutMs: options.timeoutMs ?? 45_000, }); - if (body.truncated || body.oversized || body.timedOut || !body.displaySafe) return null; - return assignmentFromRecoverySse(body.text, envelope); - } catch { - return null; + if (abortSignal?.aborted) return { recovered: false, reason: "recovery_aborted" }; + if (controller.signal.aborted || body.timedOut) return { recovered: false, reason: "recovery_timeout" }; + if (body.truncated || body.oversized || !body.displaySafe) return { recovered: false, reason: "recovery_invalid_output" }; + const assignment = assignmentFromRecoverySse(body.text, envelope); + return assignment === null + ? { recovered: false, reason: "recovery_invalid_output" } + : { recovered: true, assignment }; + } catch (error) { + if (abortSignal?.aborted) return { recovered: false, reason: "recovery_aborted" }; + const decodeFailure = boundedBodyDecodeFailure(error); + if (controller.signal.aborted || decodeFailure === "timeout") return { recovered: false, reason: "recovery_timeout" }; + return { recovered: false, reason: decodeFailure === "invalid_utf8" ? "recovery_invalid_output" : "recovery_transport_error" }; } finally { clearTimeout(timeout); } @@ -497,23 +509,23 @@ export async function recoverEncryptedAgentTaskWithResult( const admitted = admittedRecovery(req, input, config, context.parentThreadId); if (!admitted.admitted) return { recovered: false, reason: admitted.reason }; const { admission, cacheKey, envelope } = admitted.recovery; - const assignment = await resolveCachedAgentTaskRecovery( + const result = await resolveCachedAgentTaskRecoveryWithResult( cacheKey, options.cacheEntries ?? 200, signal => requestRecovery(admission, envelope, options, signal), context.abortSignal, ); - if (!assignment) { + if (!result.recovered) { return { recovered: false, - reason: context.abortSignal?.aborted ? "caller_cancelled" : "recovery_unavailable", + reason: context.abortSignal?.aborted ? "caller_cancelled" : result.reason, }; } if (context.abortSignal?.aborted) { discardCachedAgentTaskRecovery(cacheKey); return { recovered: false, reason: "caller_cancelled" }; } - if (!injectAssignment(input, envelope, assignment)) { + if (!injectAssignment(input, envelope, result.assignment)) { discardCachedAgentTaskRecovery(cacheKey); return { recovered: false, reason: "input_changed" }; } diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index 7e41dc2222..d27022c2ad 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -1719,7 +1719,17 @@ response is not cacheable. Post-commit and 5xx errors keep the no-resend path. When encrypted agent-task recovery refuses a routed task, its existing 400 error can include a bounded `recovery_reason`: `unsupported_envelope`, -`admission_denied`, `recovery_unavailable`, `caller_cancelled`, or `input_changed`. -The field is omitted when no classified recovery result exists. +`admission_denied`, `recovery_unavailable`, `caller_cancelled`, `input_changed`, +`recovery_http_rejected`, `recovery_timeout`, `recovery_aborted`, +`recovery_transport_error`, or `recovery_invalid_output`. +HTTP rejection requires an observed non-success response. Invalid output includes +invalid UTF-8, oversized bodies, malformed or incomplete recovery streams, and +invalid or conflicting assignments. A caller's cancellation takes precedence over +an owned deadline, which takes precedence over decode/transport failures. +`recovery_aborted` describes a shared recovery cancelled independently of that caller. +Shared-flight waiters receive the same underlying failure unless individually cancelled; +only successful plaintext is cached. Diagnostics contain no upstream error or payload text. +The field is omitted when no classified recovery result exists, and existing combo +branches that return the original target failure keep that response. `recovery_unavailable` includes cache/singleflight capacity and does not prove an upstream request was attempted. No retry or broader envelope acceptance is enabled. diff --git a/tests/cli/cli-status-json.test.ts b/tests/cli/cli-status-json.test.ts index 31371baa33..10ab4f110e 100644 --- a/tests/cli/cli-status-json.test.ts +++ b/tests/cli/cli-status-json.test.ts @@ -10,9 +10,11 @@ import { fileURLToPath } from "node:url"; import { isConnectionRefused, isUncleanExitEvidence, proxyHealthFailureReason, resolveStatusPid, selectListenTarget } from "../../src/cli/status"; import * as statusFacade from "../../src/cli/status"; import * as statusProbes from "../../src/cli/status-probes"; +import { packageVersion } from "../../src/cli/help"; +import { getDefaultConfig } from "../../src/config"; import { findDeadPid } from "../helpers/dead-pid"; import { removeTreeWithRetry } from "../helpers/remove-tree"; -import { STORE_BUDGET_MS } from "../helpers/test-budget"; +import { INTERNAL_DEADLINE_MS, SPAWN_BUDGET_MS, STORE_BUDGET_MS } from "../helpers/test-budget"; import { inspectClientRotationRecoveryGate, readClientConnectionState } from "../../src/client/state"; import * as lifecycleLock from "../../src/client/lifecycle-lock"; import { writeDesktopDisconnectReceipt } from "../../src/claude/desktop-remote-store"; @@ -28,6 +30,84 @@ function runStatusJson(opencodexHome: string) { }); } +describe("status version skew projection", () => { + test.each([ + ["0.0.1", "the running proxy is older"], + ["999999.0.0", "this ocx on PATH is older"], + [packageVersion(), null], + [`${packageVersion()}+skew-fixture`, "neither can be identified as older"], + ["not-a-version", "neither can be identified as older"], + ["unknown", null], + ["0.0.0", null], + [undefined, null], + ] as const)("projects proxy %s in JSON and human output", async (proxyVersion, expected) => { + const home = mkdtempSync(join(tmpdir(), "ocx-status-skew-")); + const codexHome = join(home, "codex"); + let server: ReturnType | undefined; + try { + // Explicit CODEX_HOME must exist before the CLI imports codex/paths.ts. + mkdirSync(codexHome, { recursive: true }); + server = Bun.serve({ + hostname: "127.0.0.1", port: 0, + fetch(request) { + return new URL(request.url).pathname === "/healthz" + ? Response.json({ service: "opencodex", status: "ok", version: proxyVersion, uptime: 1 }) + : new Response("not found", { status: 404 }); + }, + }); + writeFileSync(join(home, "config.json"), JSON.stringify({ + ...getDefaultConfig(), port: server.port, hostname: "127.0.0.1", codexAutoStart: false, + })); + for (const json of [true, false]) { + // Async child execution lets the fixture answer the real identity/health probes. + const child = Bun.spawn([process.execPath, cliPath, "status", ...(json ? ["--json"] : [])], { + cwd: repoRoot, + env: { ...process.env, OPENCODEX_HOME: home, CODEX_HOME: codexHome }, + stdout: "pipe", stderr: "pipe", + }); + let timedOut = false; + const timer = setTimeout(() => { + timedOut = true; + child.kill("SIGKILL"); + }, INTERNAL_DEADLINE_MS); + try { + const [stdout, stderr, exitCode] = await Promise.all([ + new Response(child.stdout).text(), new Response(child.stderr).text(), child.exited, + ]); + expect(timedOut).toBe(false); + // Preserve both gates while surfacing the child error when startup fails. + expect({ exitCode, stderr }).toEqual({ exitCode: 0, stderr: "" }); + if (json) { + const parsed = JSON.parse(stdout); + expect(parsed.schemaVersion).toBe(1); + expect(Object.keys(parsed.versionSkew).sort()).toEqual(["cliVersion", "proxyVersion", "skewed", "warning"]); + expect(parsed.versionSkew.cliVersion).toBe(packageVersion()); + expect(parsed.versionSkew.proxyVersion).toBe(proxyVersion ?? null); + expect(parsed.versionSkew.skewed).toBe(expected !== null); + if (expected === null) expect(parsed.versionSkew.warning).toBeNull(); + else expect(parsed.versionSkew.warning).toContain(expected); + } else if (expected === null) { + expect(stdout).not.toContain("does not match the running proxy"); + } else { + expect(stdout).toContain(expected); + } + } finally { + clearTimeout(timer); + if (child.exitCode === null) child.kill("SIGKILL"); + await child.exited; + } + } + expect(existsSync(join(home, "ocx.pid"))).toBe(false); + } finally { + try { + await server?.stop(true); + } finally { + removeTreeWithRetry(home); + } + } + }, SPAWN_BUDGET_MS); +}); + function withRecoveryStatusFixture(work: (fixture: { home: string; lockDeps: { lockPath: string }; diff --git a/tests/cli/cli-version-skew.test.ts b/tests/cli/cli-version-skew.test.ts index 6e45f83c28..36fb6845f9 100644 --- a/tests/cli/cli-version-skew.test.ts +++ b/tests/cli/cli-version-skew.test.ts @@ -1,5 +1,5 @@ import { describe, expect, test } from "bun:test"; -import { computeVersionSkew } from "../../src/cli/version-skew"; +import { computeVersionSkew, isConfirmedVersionMatch } from "../../src/cli/version-skew"; import { packageVersion } from "../../src/cli/help"; /** @@ -7,20 +7,83 @@ import { packageVersion } from "../../src/cli/help"; * build, and nothing surfaced it because the CLI never compared the two versions. */ describe("version skew detection", () => { - test("reports skew when the proxy reports a different version", () => { + test("directs an older CLI to upgrade or resolve PATH", () => { const skew = computeVersionSkew("2.35.0", "2.36.1"); expect(skew.skewed).toBe(true); expect(skew.cliVersion).toBe("2.35.0"); expect(skew.proxyVersion).toBe("2.36.1"); expect(skew.warning).toContain("2.35.0"); expect(skew.warning).toContain("2.36.1"); - expect(skew.warning).toContain("stale"); + expect(skew.warning).toContain("this ocx on PATH is older"); + expect(skew.warning).toContain("Upgrade the CLI or resolve PATH"); + expect(skew.warning).not.toContain("ocx service repair"); + }); + + test("#3464 directs a newer CLI to restart the older proxy", () => { + const skew = computeVersionSkew("2.42.0", "2.10.1-preview.20260805"); + expect(skew).toEqual({ + cliVersion: "2.42.0", + proxyVersion: "2.10.1-preview.20260805", + skewed: true, + warning: "CLI 2.42.0 does not match the running proxy 2.10.1-preview.20260805 — " + + "the running proxy is older than this CLI. Restart the proxy using the intended current installation. " + + "For a background service, run ocx service repair (ocx service restart is an alias).", + }); + expect(skew.warning).not.toContain("this ocx on PATH is older"); + }); + + test.each([ + ["2.43.0", "2.43.0-preview.1"], + ["2.43.0-preview.10", "2.43.0-preview.2"], + ["2.43.0-preview.beta", "2.43.0-preview.10"], + ["2.43.0-preview.1", "2.43.0-preview"], + ["2.43.0-beta", "2.43.0-alpha"], + ["2.44.0-preview.1", "2.43.0"], + ["10.0.0", "9.99.99"], + ["2.43.1", "2.43.0"], + ["2.43.0-preview.9007199254740993", "2.43.0-preview.9007199254740992"], + ])("orders %s above %s in both directions", (newer, older) => { + expect(computeVersionSkew(newer, older).warning).toContain("the running proxy is older"); + expect(computeVersionSkew(older, newer).warning).toContain("this ocx on PATH is older"); + }); + + test.each([ + ["2.43.0+build.1", "2.43.0+build.2"], + ["2.43.0", "2.43.0+build.1"], + ["2.43.0-preview.1+a", "2.43.0-preview.1+b"], + ["invalid", "2.43.0"], + ["2.43", "2.43.0"], + ["v2.43.0", "2.43.0"], + [" 2.43.0", "2.43.0"], + ["2.43.0 ", "2.43.0"], + ["2.43.0-preview.01", "2.43.0-preview.1"], + ["", "2.43.0"], + ])("keeps raw unequal %s / %s neutral in both directions", (left, right) => { + for (const [cli, proxy] of [[left, right], [right, left]]) { + const skew = computeVersionSkew(cli!, proxy!); + expect(skew.cliVersion).toBe(cli); + expect(skew.proxyVersion).toBe(proxy); + expect(skew.skewed).toBe(true); + expect(skew.warning).toContain("neither can be identified as older"); + expect(skew.warning).not.toContain("ocx service repair"); + expect(isConfirmedVersionMatch(skew)).toBe(false); + } + }); + + test.each(["unknown", "0.0.0"])("suppresses %s on either side without confirming a match", placeholder => { + for (const [cli, proxy] of [[placeholder, "2.43.0"], ["2.43.0", placeholder], [placeholder, placeholder]]) { + const skew = computeVersionSkew(cli!, proxy!); + expect(skew.skewed).toBe(false); + expect(skew.warning).toBeNull(); + expect(isConfirmedVersionMatch(skew)).toBe(false); + } }); test("stays quiet when the versions match", () => { const skew = computeVersionSkew("2.35.0", "2.35.0"); expect(skew.skewed).toBe(false); expect(skew.warning).toBeNull(); + expect(isConfirmedVersionMatch(skew)).toBe(true); }); test("stays quiet when nothing is live", () => { @@ -28,6 +91,7 @@ describe("version skew detection", () => { expect(skew.skewed).toBe(false); expect(skew.proxyVersion).toBeNull(); expect(skew.warning).toBeNull(); + expect(isConfirmedVersionMatch(skew)).toBe(false); }); test("suppresses the warning when the proxy reports the 0.0.0 placeholder", () => { diff --git a/tests/codex-integration/doctor.test.ts b/tests/codex-integration/doctor.test.ts index 9fdb7ee30d..acb0f1b87e 100644 --- a/tests/codex-integration/doctor.test.ts +++ b/tests/codex-integration/doctor.test.ts @@ -1,4 +1,7 @@ -import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import * as proxyLiveness from "../../src/server/proxy-liveness"; +import * as cliHelp from "../../src/cli/help"; +import { getDefaultConfig } from "../../src/config"; import { spawnSync } from "node:child_process"; import { existsSync, mkdirSync, mkdtempSync, utimesSync, writeFileSync } from "node:fs"; import { join } from "node:path"; @@ -32,6 +35,7 @@ import { } from "../../src/lib/local-management-capability"; import { findDeadPid } from "../helpers/dead-pid"; import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { STORE_BUDGET_MS } from "../helpers/test-budget"; const TEST_DIR = join(import.meta.dir, ".tmp-doctor-test"); const TEST_CODEX_HOME = join(TEST_DIR, "codex"); @@ -780,6 +784,63 @@ describe("doctor abandoned response-state temps", () => { }); }); +describe("doctor version skew projection", () => { + test.each([ + ["2.42.0", "2.10.1-preview.20260805", "the running proxy is older"], + ["2.35.0", "2.36.1", "this ocx on PATH is older"], + ["2.43.0", "2.43.0", "ok ocx 2.43.0 matches the running proxy"], + ["2.43.0+a", "2.43.0+b", "neither can be identified as older"], + ["v2.43.0", "2.43.0", "neither can be identified as older"], + ["2.43.0", "unknown", null], + ["unknown", "2.43.0", null], + ["2.43.0", "0.0.0", null], + ["0.0.0", "0.0.0", null], + ["unknown", "unknown", null], + ["2.43.0", undefined, null], + ] as const)("projects CLI %s / proxy %s without false matches", async (cli, proxy, expected) => { + const home = mkdtempSync(join(tmpdir(), "ocx-doctor-skew-")); + const codexHome = join(home, "codex"); + const previousHome = process.env.OPENCODEX_HOME; + const previousCodexHome = process.env.CODEX_HOME; + const previousExitCode = process.exitCode; + const restore: Array<() => void> = []; + try { + // Runtime history diagnostics resolve and stat an explicit CODEX_HOME. + mkdirSync(codexHome, { recursive: true }); + process.env.OPENCODEX_HOME = home; + process.env.CODEX_HOME = codexHome; + writeFileSync(join(home, "config.json"), JSON.stringify({ ...getDefaultConfig(), port: 9, codexAutoStart: false })); + const logged: string[] = []; + const log = spyOn(console, "log").mockImplementation((...args: unknown[]) => { logged.push(args.map(String).join(" ")); }); + restore.push(() => log.mockRestore()); + const version = spyOn(cliHelp, "packageVersion").mockReturnValue(cli); + restore.push(() => version.mockRestore()); + // Other doctor sections probe upstream health; this diagnostic fixture must stay offline. + const fetch = spyOn(globalThis, "fetch").mockImplementation(async () => new Response(null, { status: 503 })); + restore.push(() => fetch.mockRestore()); + const proxyInfo: proxyLiveness.LiveProxy = { + pid: null, port: 9, hostname: "127.0.0.1", source: "config", ...(proxy === undefined ? {} : { version: proxy }), + }; + const live = spyOn(proxyLiveness, "findLiveProxy").mockResolvedValue(proxyInfo); + restore.push(() => live.mockRestore()); + await runDoctor([]); + const output = logged.join("\n"); + if (expected !== null) expect(output).toContain(expected); + else expect(output).not.toContain("does not match the running proxy"); + if (cli !== "2.43.0" || proxy !== "2.43.0") expect(output).not.toContain("matches the running proxy"); + if (expected === "the running proxy is older") expect(output).toContain("ocx service repair"); + } finally { + for (const cleanup of restore.reverse()) cleanup(); + process.exitCode = previousExitCode; + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + if (previousCodexHome === undefined) delete process.env.CODEX_HOME; + else process.env.CODEX_HOME = previousCodexHome; + removeTreeWithRetry(home); + } + }, STORE_BUDGET_MS); +}); + describe("doctor reclaim wiring (end to end)", () => { // The formatter tests above cannot observe deletion. This covers the call site itself: // inverting the report/reclaim ternary in runDoctor must fail a test. diff --git a/tests/server/agent-task-recovery-cache.test.ts b/tests/server/agent-task-recovery-cache.test.ts index 2ee994f8ce..35b5ba7050 100644 --- a/tests/server/agent-task-recovery-cache.test.ts +++ b/tests/server/agent-task-recovery-cache.test.ts @@ -29,16 +29,23 @@ describe("agent task recovery cache", () => { resetAgentTaskRecoveryCache(); }); - test("shared failure gives each waiter its own result without contaminating another key", async () => { + test.each([ + { kind: "http", reason: "recovery_http_rejected" }, + { kind: "reader", reason: "recovery_transport_error" }, + { kind: "decode", reason: "recovery_invalid_output" }, + ] as const)("shared $kind failure gives each waiter its own result without contaminating another key", async ({ kind, reason }) => { let release: (() => void) | undefined; const gate = new Promise(resolve => { release = resolve; }); let fetches = 0; globalThis.fetch = (async () => { const requestNumber = ++fetches; await gate; - return requestNumber === 1 - ? new Response("raw-failure-sentinel", { status: 503 }) - : new Response(recoverySse("Independent assignment.")); + if (requestNumber !== 1) return new Response(recoverySse("Independent assignment.")); + if (kind === "decode") return new Response(new Uint8Array([0xff])); + if (kind === "reader") return new Response(new ReadableStream({ + pull(controller) { controller.error(new TypeError("private-reader-failure")); }, + })); + return new Response("raw-failure-sentinel", { status: 503 }); }) as typeof fetch; const req = new Request("http://localhost/v1/responses", { headers: codexHeaders() }); const config = routedConfig(); @@ -53,8 +60,8 @@ describe("agent task recovery cache", () => { expect(fetches).toBe(2); release?.(); const [firstResult, secondResult, otherResult] = await Promise.all([first, second, other]); - expect(firstResult).toEqual({ recovered: false, reason: "recovery_unavailable" }); - expect(secondResult).toEqual({ recovered: false, reason: "recovery_unavailable" }); + expect(firstResult).toEqual({ recovered: false, reason }); + expect(secondResult).toEqual({ recovered: false, reason }); expect(firstResult).not.toBe(secondResult); expect(otherResult).toEqual({ recovered: true }); expect(firstInput).toEqual(encryptedInput()); @@ -68,6 +75,39 @@ describe("agent task recovery cache", () => { } }); + test("shared flight reset reports abort to surviving callers and never caches late plaintext", async () => { + let release!: () => void; + const gate = new Promise(resolve => { release = resolve; }); + let fetches = 0; + globalThis.fetch = (async () => { + fetches++; + await gate; + return new Response(recoverySse("private-late-assignment")); + }) as typeof fetch; + const req = new Request("http://localhost/v1/responses", { headers: codexHeaders() }); + const firstInput = encryptedInput(); + const secondInput = encryptedInput(); + const first = recoverEncryptedAgentTaskWithResult(req, firstInput, {}, routedConfig()); + const second = recoverEncryptedAgentTaskWithResult(req, secondInput, {}, routedConfig()); + try { + expect(fetches).toBe(1); + resetAgentTaskRecoveryCache(); + release(); + const results = await Promise.all([first, second]); + expect(results).toEqual([ + { recovered: false, reason: "recovery_aborted" }, + { recovered: false, reason: "recovery_aborted" }, + ]); + expect(results[0]).not.toBe(results[1]); + expect(firstInput).toEqual(encryptedInput()); + expect(secondInput).toEqual(encryptedInput()); + expect(agentTaskRecoveryCacheSnapshotForTests()).toEqual({ entries: 0, bytes: 0 }); + } finally { + release(); + await Promise.all([first, second]); + } + }); + for (const succeeds of [true, false]) { test(`caller cancellation stays local when the remaining waiter ${succeeds ? "succeeds" : "fails"}`, async () => { let release: (() => void) | undefined; @@ -95,7 +135,7 @@ describe("agent task recovery cache", () => { release?.(); expect(await second).toEqual(succeeds ? { recovered: true } - : { recovered: false, reason: "recovery_unavailable" }); + : { recovered: false, reason: "recovery_http_rejected" }); expect(fetches).toBe(1); expect(restoreCachedEncryptedAgentTasks(req, encryptedInput(), config)).toBe(succeeds ? 1 : 0); } finally { diff --git a/tests/server/agent-task-recovery.test.ts b/tests/server/agent-task-recovery.test.ts index ceb1c5b6b5..a168f2c364 100644 --- a/tests/server/agent-task-recovery.test.ts +++ b/tests/server/agent-task-recovery.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; import { createTranslatorBudget } from "../../src/lib/translator-budget"; import { warnAgentTaskRecoveryStartup } from "../../src/server"; import { @@ -7,6 +7,7 @@ import { recoverEncryptedAgentTaskWithResult, resetAgentTaskRecoveryState, restoreCachedEncryptedAgentTasks, + type AgentTaskRecoveryFailureReason, } from "../../src/server/responses/agent-task-recovery"; import { agentTaskRecoveryWaiterCountForTests } from "../../src/server/responses/agent-task-recovery-cache"; import { @@ -78,24 +79,36 @@ describe("agent task recovery (opt-in, default off)", () => { }); } - const failedRecoveries: Array<[string, () => Response]> = [ - ["HTTP 503", () => new Response("raw-error-sentinel", { status: 503 })], - ["network exception", () => { throw new Error("raw-error-sentinel"); }], - ["malformed SSE", () => new Response("data: {not-json}\n\n")], - ["missing completion", () => new Response(recoverySse("payload-sentinel").split("data: {\"type\":\"response.completed\"")[0])], - ["conflicting assignment", () => new Response(recoverySse("payload-sentinel") + recoveryCompletedSse("other-payload-sentinel"))], - ["failed terminal", () => new Response(recoverySse("payload-sentinel") + 'data: {"type":"response.failed","response":{"error":{"message":"raw-error-sentinel"}}}\n\n')], - ["incomplete terminal", () => new Response(recoverySse("payload-sentinel") + 'data: {"type":"response.incomplete"}\n\n')], - ["bare error", () => new Response(recoverySse("payload-sentinel") + 'data: {"type":"error","error":{"message":"raw-error-sentinel"}}\n\n')], + const failedRecoveries: Array<[string, () => Response, AgentTaskRecoveryFailureReason]> = [ + ["HTTP 401", () => new Response("private-error", { status: 401 }), "recovery_http_rejected"], + ["HTTP 403", () => new Response("private-error", { status: 403 }), "recovery_http_rejected"], + ["HTTP 429", () => new Response("private-error", { status: 429 }), "recovery_http_rejected"], + ["fetch TypeError", () => { throw new TypeError("private-error"); }, "recovery_transport_error"], + ["unowned TimeoutError", () => { throw new DOMException("private-error", "TimeoutError"); }, "recovery_transport_error"], + ["reader TypeError", () => new Response(new ReadableStream({ + pull(controller) { controller.error(new TypeError("private-reader-error")); }, + })), "recovery_transport_error"], + ["invalid UTF-8", () => new Response(new Uint8Array([0xff])), "recovery_invalid_output"], + ["trailing UTF-8", () => new Response(new Uint8Array([0xe2, 0x82])), "recovery_invalid_output"], + ["oversized body", () => new Response(new Uint8Array(4 * 1024 * 1024 + 1)), "recovery_invalid_output"], + ["invalid arguments", () => new Response(recoverySse("task").replace('{\\"assignment\\":\\"task\\"}', '{broken')), "recovery_invalid_output"], + ["HTTP 503", () => new Response("raw-error-sentinel", { status: 503 }), "recovery_http_rejected"], + ["network exception", () => { throw new Error("raw-error-sentinel"); }, "recovery_transport_error"], + ["malformed SSE", () => new Response("data: {not-json}\n\n"), "recovery_invalid_output"], + ["missing completion", () => new Response(recoverySse("payload-sentinel").split("data: {\"type\":\"response.completed\"")[0]), "recovery_invalid_output"], + ["conflicting assignment", () => new Response(recoverySse("payload-sentinel") + recoveryCompletedSse("other-payload-sentinel")), "recovery_invalid_output"], + ["failed terminal", () => new Response(recoverySse("payload-sentinel") + 'data: {"type":"response.failed","response":{"error":{"message":"raw-error-sentinel"}}}\n\n'), "recovery_invalid_output"], + ["incomplete terminal", () => new Response(recoverySse("payload-sentinel") + 'data: {"type":"response.incomplete"}\n\n'), "recovery_invalid_output"], + ["bare error", () => new Response(recoverySse("payload-sentinel") + 'data: {"type":"error","error":{"message":"raw-error-sentinel"}}\n\n'), "recovery_invalid_output"], // Exact-case events are also used by the pinned official Codex source. Recovery's // additional completed-status requirement remains deliberately stricter. - ["mixed-case completion", () => new Response(recoverySse("payload-sentinel").replace("response.completed", "Response.Completed"))], - ["mixed-case status", () => new Response(recoverySse("payload-sentinel").replace('"status":"completed"', '"status":"Completed"'))], - ["missing status", () => new Response(recoverySse("payload-sentinel").replace('"status":"completed",', ""))], - ["ciphertext assignment", () => new Response(recoverySse(FERNET_TASK))], + ["mixed-case completion", () => new Response(recoverySse("payload-sentinel").replace("response.completed", "Response.Completed")), "recovery_invalid_output"], + ["mixed-case status", () => new Response(recoverySse("payload-sentinel").replace('"status":"completed"', '"status":"Completed"')), "recovery_invalid_output"], + ["missing status", () => new Response(recoverySse("payload-sentinel").replace('"status":"completed",', "")), "recovery_invalid_output"], + ["ciphertext assignment", () => new Response(recoverySse(FERNET_TASK)), "recovery_invalid_output"], ]; - for (const [name, response] of failedRecoveries) { - test(`typed recovery keeps ${name} coarse and preserves false without retrying`, async () => { + for (const [name, response, reason] of failedRecoveries) { + test(`typed recovery classifies ${name} and preserves false without retrying`, async () => { const req = new Request("http://localhost/v1/responses", { headers: codexHeaders() }); const config = routedConfig(); let fetches = 0; @@ -103,7 +116,7 @@ describe("agent task recovery (opt-in, default off)", () => { const input = encryptedInput(); const original = structuredClone(input); expect(await recoverEncryptedAgentTaskWithResult(req, input, {}, config)) - .toEqual({ recovered: false, reason: "recovery_unavailable" }); + .toEqual({ recovered: false, reason }); expect(input).toEqual(original); expect(fetches).toBe(1); expect(restoreCachedEncryptedAgentTasks(req, encryptedInput(), config)).toBe(0); @@ -113,6 +126,69 @@ describe("agent task recovery (opt-in, default off)", () => { }); } + test.each(["pending", "rejecting"] as const)("HTTP refusal does not await %s body cancellation", async mode => { + let cancels = 0; + let reads = 0; + let releaseCancel: (() => void) | undefined; + const cancellation = new Promise(resolve => { releaseCancel = resolve; }); + globalThis.fetch = (async () => new Response(new ReadableStream({ + pull() { reads++; }, + cancel() { + cancels++; + return mode === "pending" ? cancellation : Promise.reject(new Error("private-cancel-error")); + }, + }, { highWaterMark: 0 }), { status: 503 })) as typeof fetch; + try { + const result = await recoverEncryptedAgentTaskWithResult( + new Request("http://localhost/v1/responses", { headers: codexHeaders() }), encryptedInput(), {}, routedConfig(), + ); + expect(result).toEqual({ recovered: false, reason: "recovery_http_rejected" }); + expect(cancels).toBe(1); + expect(reads).toBe(0); + } finally { + releaseCancel?.(); + } + }); + + test.each(["headers", "body", "caller"] as const)("owned deadline classification at %s preserves cancellation precedence", async site => { + const callbacks: Array<() => void> = []; + const timers = spyOn(globalThis, "setTimeout").mockImplementation(((callback: () => void) => { + callbacks.push(callback); + return 0 as unknown as ReturnType; + }) as typeof setTimeout); + const caller = new AbortController(); + let started!: () => void; + const ready = new Promise(resolve => { started = resolve; }); + let fetches = 0; + globalThis.fetch = ((_, init) => { + fetches++; + if (site === "body") return Promise.resolve(new Response(new ReadableStream({ + pull(controller) { + controller.enqueue(new Uint8Array([0xe2, 0x82])); + started(); + return new Promise(() => {}); + }, + }, { highWaterMark: 0 }))); + return new Promise((_resolve, reject) => { + init?.signal?.addEventListener("abort", () => reject(init.signal?.reason), { once: true }); + started(); + }); + }) as typeof fetch; + try { + const pending = recoverEncryptedAgentTaskWithResult( + new Request("http://localhost/v1/responses", { headers: codexHeaders() }), encryptedInput(), {}, routedConfig(), + { abortSignal: caller.signal }, + ); + await ready; + callbacks[0]!(); // Fire the owned deadline without wall-clock sleeps. + if (site === "caller") caller.abort(new TypeError("private-caller-error")); + expect(await pending).toEqual({ recovered: false, reason: site === "caller" ? "caller_cancelled" : "recovery_timeout" }); + expect(fetches).toBe(1); + } finally { + timers.mockRestore(); + } + }); + test("keeps the disabled fail-fast response byte-identical to the absent feature", async () => { const snapshot = async (config: ReturnType) => { let fetchCalls = 0; @@ -226,7 +302,7 @@ describe("agent task recovery (opt-in, default off)", () => { expect(response.status).toBe(400); expect(json.error?.code).toBe("unreadable_encrypted_agent_task"); - expect(json.error?.recovery_reason).toBe("recovery_unavailable"); + expect(json.error?.recovery_reason).toBe("recovery_invalid_output"); expect(fetchedUrls.length).toBeGreaterThan(0); expect(fetchedUrls[0]).toContain("chatgpt.com/backend-api/codex"); }); @@ -778,7 +854,7 @@ describe("agent task recovery (opt-in, default off)", () => { expect(fetchedUrls).toHaveLength(1); expect(fetchedUrls[0]).toContain("chatgpt.com/backend-api/codex/responses"); expect(await response.json()).toMatchObject({ - error: { code: "unreadable_encrypted_agent_task", recovery_reason: "recovery_unavailable" }, + error: { code: "unreadable_encrypted_agent_task", recovery_reason: "recovery_transport_error" }, }); }); }); diff --git a/tests/server/bounded-body.test.ts b/tests/server/bounded-body.test.ts index f5223d34a4..0bf5e0ae1b 100644 --- a/tests/server/bounded-body.test.ts +++ b/tests/server/bounded-body.test.ts @@ -1,7 +1,8 @@ -import { describe, expect, test } from "bun:test"; +import { describe, expect, spyOn, test } from "bun:test"; import { BOUNDED_BODY_MAX_BYTES, boundedBodyBufferGrowthsForTests, + boundedBodyDecodeFailure, readBoundedResponseBytes, readBoundedResponseBody, } from "../../src/lib/bounded-body"; @@ -21,6 +22,65 @@ function responseFromChunks(...chunks: Uint8Array[]): Response { } describe("readBoundedResponseBody", () => { + test("only actual decoder exceptions carry the decode discriminator", async () => { + for (const bytes of [new Uint8Array([0xff]), new Uint8Array([0xe2, 0x82])]) { + let caught: unknown; + try { await readBoundedResponseBody(responseFromChunks(bytes), { fatalUtf8: true }); } + catch (error) { caught = error; } + expect(caught).toBeInstanceOf(TypeError); + expect(boundedBodyDecodeFailure(caught)).toBe("invalid_utf8"); + } + const readerError = new TypeError("private-reader-error"); + const response = new Response(new ReadableStream({ pull(controller) { controller.error(readerError); } })); + let caught: unknown; + try { await readBoundedResponseBody(response, { fatalUtf8: true }); } + catch (error) { caught = error; } + expect(caught).toBe(readerError); + expect(boundedBodyDecodeFailure(caught)).toBeUndefined(); + }); + + test("fatal UTF-8 abort retains the exact caller reason without a decode mark", async () => { + const caller = new AbortController(); + const reason = new TypeError("private-caller-error"); + const pending = readBoundedResponseBody(new Response(new ReadableStream({})), { signal: caller.signal, fatalUtf8: true }); + caller.abort(reason); + let caught: unknown; + try { await pending; } catch (error) { caught = error; } + expect(caught).toBe(reason); + expect(boundedBodyDecodeFailure(caught)).toBeUndefined(); + }); + + test.each([0, 1])("fatal timeout flush retains deadline origin %s and cancels without waiting", async deadline => { + const callbacks: Array<() => void> = []; + const timers = spyOn(globalThis, "setTimeout").mockImplementation(((callback: () => void) => { + callbacks.push(callback); + return 0 as unknown as ReturnType; + }) as typeof setTimeout); + let stalled!: () => void; + const ready = new Promise(resolve => { stalled = resolve; }); + let pulls = 0; + let cancelled = false; + const response = new Response(new ReadableStream({ + pull(controller) { + if (pulls++ === 0) controller.enqueue(new Uint8Array([0xe2, 0x82])); + else { stalled(); return new Promise(() => {}); } + }, + cancel() { cancelled = true; return new Promise(() => {}); }, + }, { highWaterMark: 0 })); + try { + const pending = readBoundedResponseBody(response, { fatalUtf8: true }); + await ready; + callbacks[deadline === 0 ? 0 : callbacks.length - 1]!(); + let caught: unknown; + try { await pending; } catch (error) { caught = error; } + expect(caught).toBeInstanceOf(TypeError); + expect(boundedBodyDecodeFailure(caught)).toBe("timeout"); + expect(cancelled).toBe(true); + } finally { + timers.mockRestore(); + } + }); + test("the bounded JSON caller allows a full total deadline for its first byte", () => { expect(UPSTREAM_JSON_BODY_READ_OPTIONS.firstByteTimeoutMs) .toBe(UPSTREAM_JSON_BODY_READ_OPTIONS.totalTimeoutMs);