diff --git a/docs-site/src/content/docs/guides/codex-integration.md b/docs-site/src/content/docs/guides/codex-integration.md index 712d5d9036c..eb236fff474 100644 --- a/docs-site/src/content/docs/guides/codex-integration.md +++ b/docs-site/src/content/docs/guides/codex-integration.md @@ -477,6 +477,42 @@ While the mode is active, the realtime voice sideband override (`experimental_realtime_ws_base_url`) is not injected — the dedicated provider-table form cannot carry it — so Codex Desktop voice uses its native endpoint rather than the proxy. +### Emergency compaction model (opt-in) + +`compactionRecovery` leaves the initial compaction on the conversation's selected route. It +permits one emergency attempt only after a supported, pre-output compaction failure. It is +separate from `compactionRouting`, which chooses another model before compaction starts, and +from `codexClientCompaction`, which changes Codex's provider form. + +```json +{ + "compactionRecovery": { + "enabled": true, + "model": "provider/emergency-model", + "allowDevinInvalidArgument": false + } +} +``` + +Use an independently configured, authorized model with enough context for the failed input. +Enabling recovery permits that model's provider to receive the compaction history and charge +for the extra attempt when recovery runs; ordinary successful compactions incur no extra call. +The option is off when absent or disabled. The existing authenticated management API accepts +this block through `PUT /api/settings`; send `compactionRecovery: null` to remove it. A direct +file edit should follow the normal stopped-proxy configuration workflow. This setting does not +change sign-in, the conversation's ordinary model, Codex's provider ID, or the desktop composer. + +Recovery does not replay after cancellation, semantic output, tool side effects, an exhausted +send budget, or an authentication, admission or policy refusal. Generic `400` errors do not +enable fallback. The separately opted-in Devin `invalid_argument` case applies only to an +identified compaction failure from that adapter. The emergency attempt shares the original +request's send budget and never starts a second recovery attempt. + +Native encrypted compaction is outside this recovery path: its original error is retained. +There is no automatic local truncation mode. A response being accepted is not proof that a +long conversation retained its goals; verify the next turn on the original model before +treating an emergency summary as a recovered task. + ### Authless Codex Desktop (opt-in) In **Dashboard → Overview**, **Open Codex without signing in** controls this existing diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 26aad251318..7f2519ca7cf 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -168,6 +168,9 @@ } }, "explicit": { + "responses-compaction-recovery.test.ts": "responses", + "compaction-recovery-settings.test.ts": "config", + "responses-compaction-recovery-policy.test.ts": "responses", "deepseek-artifact-tool-schema.test.ts": "providers", "client-config-export-output-limit.test.ts": "config", "openai-chat-serialized-tool-call-scaling.test.ts": "adapters/openai", diff --git a/src/config/diagnostics.ts b/src/config/diagnostics.ts index deb761dc2e8..ca7b899e3d6 100644 --- a/src/config/diagnostics.ts +++ b/src/config/diagnostics.ts @@ -2,6 +2,7 @@ import { createHash } from "node:crypto"; import { lstatSync, readFileSync } from "node:fs"; import { join } from "node:path"; import * as z from "zod/v4"; +import { compactionRecoveryConfigError } from "./schema/compaction-recovery"; import type { OcxConfig } from "../types"; import { configReasoningPinsConfigError } from "./provider-validation"; import { loopbackCompanionAllowed } from "../codex/loopback-target"; @@ -598,7 +599,7 @@ export function validateConfigCandidate(value: unknown): { ok: true; config: Ocx if (compactionRouting !== undefined && !compactionRoutingSchema.safeParse(compactionRouting).success) { return { ok: false, error: "schema_invalid: compactionRouting: requires a nonblank model, an optional valid reasoningEffort, and optional non-repeating triggers drawn from \"manual\" and \"auto\"" }; } - const boundaryError = configReasoningPinsConfigError(value) + const boundaryError = compactionRecoveryConfigError(value) ?? configReasoningPinsConfigError(value) ?? blankHostnameError(value) ?? claudeSubagentEffortError(value) ?? appOwnedMemoryBudgetError(value) diff --git a/src/config/load-degrade.ts b/src/config/load-degrade.ts index 412d011b2f3..aa8cc5e94c4 100644 --- a/src/config/load-degrade.ts +++ b/src/config/load-degrade.ts @@ -1,5 +1,6 @@ import { chmodSync, existsSync } from "node:fs"; import { join } from "node:path"; +import { compactionRecoveryConfigError } from "./schema/compaction-recovery"; import { modelPinnedEffortsConfigError, pinnedReasoningEffortConfigError, @@ -118,6 +119,7 @@ export function warnDegradedCompactionRouting(rawParsed: unknown, validated: Ocx * the ratchet only ever moves down: a per-block call there costs a line the file does not have. */ export function warnDegradedTopLevelOptIns(rawParsed: unknown, validated: OcxConfig): void { + if (compactionRecoveryConfigError(rawParsed)) console.warn("⚠️ invalid compactionRecovery disabled; the original compaction failure is preserved"); warnDegradedStreamMode(rawParsed, validated); warnDegradedCompactionRouting(rawParsed, validated); } diff --git a/src/config/schema/compaction-recovery.ts b/src/config/schema/compaction-recovery.ts new file mode 100644 index 00000000000..0282074b722 --- /dev/null +++ b/src/config/schema/compaction-recovery.ts @@ -0,0 +1,15 @@ +import * as z from "zod/v4"; + +/** Failure-only, opt-in recovery; it never chooses the initial compaction model. */ +export const compactionRecoverySchema = z.object({ + enabled: z.boolean(), + model: z.string().trim().min(1).max(512).regex(/^[^\s\u0000-\u0020\u007f-\u009f]+$/), + allowDevinInvalidArgument: z.boolean().optional(), +}).strict(); + +export function compactionRecoveryConfigError(value: unknown): string | null { + if (!value || typeof value !== "object" || Array.isArray(value)) return null; + const recovery = (value as Record).compactionRecovery; + return recovery === undefined || compactionRecoverySchema.safeParse(recovery).success + ? null : "schema_invalid: compactionRecovery: requires enabled, a nonblank model, and optional boolean allowDevinInvalidArgument"; +} diff --git a/src/config/schema/config-schema.ts b/src/config/schema/config-schema.ts index c881a269fc1..48460d3667e 100644 --- a/src/config/schema/config-schema.ts +++ b/src/config/schema/config-schema.ts @@ -1,4 +1,5 @@ import * as z from "zod/v4"; +import { compactionRecoverySchema } from "./compaction-recovery"; import { agentTaskRecoverySchema, catalogAutoRefreshSchema, @@ -156,6 +157,7 @@ export const configSchema = z.object({ providers: z.record(z.string(), providerConfigSchema), modelPinnedEfforts: modelPinnedEffortsSchema.optional(), compactionRouting: compactionRoutingSchema.optional().catch(undefined), + compactionRecovery: compactionRecoverySchema.optional().catch(undefined), defaultProvider: z.string().min(1).default("openai"), defaultModelAliases: z.boolean().optional(), // Malformed hand edits disable this opt-in projection without rejecting providers. diff --git a/src/server/management/config-routes.ts b/src/server/management/config-routes.ts index d439abb137c..1efba0c4b88 100644 --- a/src/server/management/config-routes.ts +++ b/src/server/management/config-routes.ts @@ -1,4 +1,5 @@ import { compactionRoutingSchema } from "../../config/schema/leaf-validators"; +import { compactionRecoverySchema } from "../../config/schema/compaction-recovery"; import { captureConfigTopLevelRollback } from "../../config/rebase-provenance"; import type { IntegrationClientId } from "../../integrations/registry"; import { randomUUID } from "node:crypto"; @@ -369,6 +370,7 @@ export async function handleConfigRoutes(ctx: ManagementContext): Promise { try { - yield* readResponseStreamWithInactivity( + for await (const event of readResponseStreamWithInactivity( upstreamResponse, upstream.signal, bodyInactivityMs, response => transportState.activeAdapter.parseStream(response, translatorBudget, logCtx.activeTierMetadata), - ); + )) { + options.onCompactionRecoveryAdapterEvent?.(event); + yield event; + } } catch (error) { if (error instanceof ResponseBodyInactivityError) { yield { @@ -197,6 +200,7 @@ export async function deliverAdapterResponse( bodyInactivityMs, response => transportState.activeAdapter.parseResponse!(response, translatorBudget, logCtx.activeTierMetadata), ); + for (const event of initialEvents) options.onCompactionRecoveryAdapterEvent?.(event); let guardedEvents: AdapterEvent[]; if (terminalGuardEnabled) { guardedEvents = []; diff --git a/src/server/responses/adapter-dispatch.ts b/src/server/responses/adapter-dispatch.ts index e6f0c742e85..2370fe705d4 100644 --- a/src/server/responses/adapter-dispatch.ts +++ b/src/server/responses/adapter-dispatch.ts @@ -320,11 +320,18 @@ export async function prepareAdapterExchange( // legacy direct-Google exception is preserved exactly; every other adapter still keeps // reset-only semantics so combo failover hops on the first 5xx. const transientPolicy = transientRetryPolicyFor(route.provider); + const compactPrepaid = options.compactionRecoveryAttempted ? sendBudgetState.pendingHopPermit : undefined; + if (compactPrepaid) sendBudgetState.pendingHopPermit = undefined; + let compactPrepaidUsed = false; const fetchWithRetryPolicy = (route.provider.adapter === "google" || transientPolicy) ? fetchWithTransientRetry : fetchWithResetRetry; upstreamResponse = await fetchWithRetryPolicy( recovery => { + if (compactPrepaid && !compactPrepaidUsed) { + if (!compactPrepaid.use()) throw new SendBudgetExhaustedError(safeHostLabel(builtInitialRequest.url)); + compactPrepaidUsed = true; + } transportState.noteRoutedAttemptSend(inputTokenEstimate, recovery); return fetchWithHeaderTimeout(builtInitialRequest.url, applyUpstreamRecoveryInit({ method: builtInitialRequest.method, @@ -340,12 +347,15 @@ export async function prepareAdapterExchange( { abortSignal: upstream.signal, label: safeHostLabel(builtInitialRequest.url), - ...(transientPolicy + ...(transientPolicy || compactPrepaid // Draws the remainder, not the raw policy. A combo child inherits the parent's // holder but used to take a fresh full allowance on its own first send, so the // shared counter was inherited without ever being read as a limit. ? { - attempts: remainingTransientSendBudget(transientPolicy.attempts), + // The first emergency send is already paid for. Only retries consume the + // remaining allowance; treating the booking as unavailable blocks a cap of two. + attempts: Math.min(transientPolicy?.attempts ?? 1, + remainingTransientSendBudget(transientPolicy?.attempts ?? 1) + (compactPrepaid ? 1 : 0)), onSendsConsumed: noteTransientSends, } : {}), @@ -1102,6 +1112,11 @@ export async function prepareAdapterExchange( // material before it reaches the client-facing error surface. const upstreamRetryAfter = upstreamResponse.headers.get("retry-after"); const normalized = normalizeUpstreamErrorText(errorText, "unknown error"); + options.onCompactionRecoveryAdapterEvent?.({ + type: "error", status: upstreamResponse.status, + errorType: normalized.type, code: normalized.code, + message: "Structured upstream failure observed before client formatting", + }); const message = normalized.cyberPolicy ? normalized.message ?? (isCyberPolicyCode(normalized.code) ? CYBER_POLICY_FALLBACK_MESSAGE : normalized.safeText) diff --git a/src/server/responses/compact.ts b/src/server/responses/compact.ts index c9a3cc5e495..33e57939c9e 100644 --- a/src/server/responses/compact.ts +++ b/src/server/responses/compact.ts @@ -1433,7 +1433,7 @@ export async function handleResponsesCompact( // The routed compaction turn is a handoff inside the same logical request, so it draws the // REMAINDER. Minting here is what let a native attempt spend three sends and the routed // fallback spend four more. - const response = await handleResponses(internalReq, config, logCtx, { abortSignal: req.signal, turnAdmissionLease, sendBudget, compactionRoutingOverride: options.compactionRoutingOverride, ...(admission ? { admission } : {}) }); + const response = await handleResponses(internalReq, config, logCtx, { abortSignal: req.signal, turnAdmissionLease, sendBudget, compactionRecoveryKind: "compaction-v1", compactionRoutingOverride: options.compactionRoutingOverride, ...(admission ? { admission } : {}) }); if (!response.ok) return response; let json: { output?: unknown[]; status?: unknown; error?: unknown }; if (response.headers.get("content-type")?.includes("text/event-stream")) { diff --git a/src/server/responses/compaction-recovery-policy.ts b/src/server/responses/compaction-recovery-policy.ts new file mode 100644 index 00000000000..7d74e176496 --- /dev/null +++ b/src/server/responses/compaction-recovery-policy.ts @@ -0,0 +1,116 @@ +/** Pure eligibility policy. The caller owns admission, buffering, dispatch and the shared budget. */ +export interface CompactionRecoveryConfig { + enabled: true; + model: string; + allowDevinInvalidArgument?: boolean; +} + +/** Evidence must describe the failed attempt, never request text or a guessed error message. */ +export interface CompactionRecoveryEvidence { + requestKind: "ordinary" | "compaction-v1" | "compaction-v2"; + recoveryAttempts: number; + cancelled: boolean; + nonReplayable: boolean; + /** Any semantic output observed, including output retained privately before delivery. */ + partialOutput: boolean; + toolEffects: boolean; + remainingSends: number; + /** Canonical serving identities resolved by the caller, including provider/account identity. */ + originalModel: string; + fallbackModel: string; + /** Actual serving adapter/provider, not a prefix inferred from the requested selector. */ + provider: string; + httpStatus?: number; + /** A completed result is never retried. Output validation belongs to the caller. */ + responseStatus: "completed" | "failed" | "incomplete" | "unknown"; + errorCode?: string; + errorType?: string; + authenticationDenied: boolean; + policyDenied: boolean; + budgetDenied: boolean; + refusal: boolean; + upstreamFailure: boolean; +} + +export type CompactionRecoveryDecision = + | { recover: true; model: string; reason: "context-overflow" | "compaction-output" | "upstream-unavailable" | "devin-invalid-argument" } + | { recover: false; reason: "disabled" | "invalid-evidence" | "ordinary-request" | "already-attempted" | "cancelled" | "unsafe-replay" | "protected-failure" | "budget-exhausted" | "same-model" | "succeeded" | "unclassified-failure" }; + +const MODEL_LIMIT = 512; +const safeName = (value: unknown, limit: number): value is string => + typeof value === "string" && value.length > 0 && value.length <= limit + && value === value.trim() && !/[\u0000-\u0020\u007f-\u009f]/.test(value); + +/** Invalid persisted values disable recovery rather than widening it. No configurable retry count. */ +export function readCompactionRecoveryConfig(value: unknown): CompactionRecoveryConfig | null { + if (!value || typeof value !== "object" || Array.isArray(value)) return null; + const raw = value as Record; + if (Object.keys(raw).some(key => !["enabled", "model", "allowDevinInvalidArgument"].includes(key)) + || raw.enabled !== true || !safeName(raw.model, MODEL_LIMIT) + || (raw.allowDevinInvalidArgument !== undefined && typeof raw.allowDevinInvalidArgument !== "boolean")) return null; + return { + enabled: true, + model: raw.model, + ...(raw.allowDevinInvalidArgument !== undefined ? { allowDevinInvalidArgument: raw.allowDevinInvalidArgument } : {}), + }; +} + +// Structured tokens only. The normalizer must also supply the explicit denial flags above. +const PROTECTED_CODE = /(?:^|_)(?:auth|authentication|authorization|unauthenticated|unauthorized|permission|forbidden|policy|refusal|refused|budget|quota|billing|safety|content_filter|origin_rejected|admission|scope)(?:_|$)/; +const CONTEXT_CODES = new Set(["context_length_exceeded", "context_window_exceeded", "input_too_long"]); +const COMPACTION_CODES = new Set(["compaction_failed", "invalid_compaction_output", "empty_compaction_output"]); +const SERVER_CODES = new Set(["upstream_error", "upstream_server_error", "server_is_overloaded", "internal", "internal_error", "server_error", "unavailable", "service_unavailable", "gateway_timeout"]); + +/** Decides eligibility only; returning true neither spends nor grants another send. */ +export function decideCompactionRecovery( + value: unknown, + evidence: CompactionRecoveryEvidence, +): CompactionRecoveryDecision { + const config = readCompactionRecoveryConfig(value); + if (!config) return { recover: false, reason: "disabled" }; + const flags = [evidence.cancelled, evidence.nonReplayable, evidence.partialOutput, evidence.toolEffects, + evidence.authenticationDenied, evidence.policyDenied, evidence.budgetDenied, evidence.refusal, evidence.upstreamFailure]; + if (flags.some(flag => typeof flag !== "boolean") + || !Number.isSafeInteger(evidence.recoveryAttempts) || evidence.recoveryAttempts < 0 + || !Number.isSafeInteger(evidence.remainingSends) || evidence.remainingSends < 0 + || !safeName(evidence.originalModel, MODEL_LIMIT) || !safeName(evidence.fallbackModel, MODEL_LIMIT) + || !safeName(evidence.provider, 128) + || !["ordinary", "compaction-v1", "compaction-v2"].includes(evidence.requestKind) + || !["completed", "failed", "incomplete", "unknown"].includes(evidence.responseStatus) + || (evidence.httpStatus !== undefined && (!Number.isInteger(evidence.httpStatus) || evidence.httpStatus < 100 || evidence.httpStatus > 599)) + || [evidence.errorCode, evidence.errorType].some(code => code !== undefined && !safeName(code, 128))) { + return { recover: false, reason: "invalid-evidence" }; + } + if (evidence.requestKind === "ordinary") return { recover: false, reason: "ordinary-request" }; + if (evidence.recoveryAttempts !== 0) return { recover: false, reason: "already-attempted" }; + const status = evidence.httpStatus; + const codes = [evidence.errorCode, evidence.errorType].filter((code): code is string => code !== undefined); + if (evidence.cancelled || status === 499 || codes.some(code => code === "cancelled" || code === "canceled" || code === "client_cancelled")) { + return { recover: false, reason: "cancelled" }; + } + if (evidence.nonReplayable || evidence.partialOutput || evidence.toolEffects) return { recover: false, reason: "unsafe-replay" }; + if (evidence.authenticationDenied || evidence.policyDenied || evidence.budgetDenied || evidence.refusal + || status === 401 || status === 403 || status === 402 || status === 429 + || codes.some(code => PROTECTED_CODE.test(code.toLowerCase()) || code === "failed_precondition")) { + return { recover: false, reason: "protected-failure" }; + } + if (evidence.remainingSends === 0) return { recover: false, reason: "budget-exhausted" }; + if (evidence.originalModel === evidence.fallbackModel) return { recover: false, reason: "same-model" }; + if (evidence.responseStatus === "completed") return { recover: false, reason: "succeeded" }; + const failedHttp = status !== undefined && status >= 400; + if (!evidence.upstreamFailure || (!failedHttp && evidence.responseStatus !== "failed")) { + return { recover: false, reason: "unclassified-failure" }; + } + const code = evidence.errorCode; + const requestFailure = status === undefined || status === 200 || status === 400 || status === 413 || status === 422 || status >= 500; + if (requestFailure && code && CONTEXT_CODES.has(code)) return { recover: true, model: config.model, reason: "context-overflow" }; + if (requestFailure && code && COMPACTION_CODES.has(code)) return { recover: true, model: config.model, reason: "compaction-output" }; + if (config.allowDevinInvalidArgument === true && evidence.provider === "devin" && code === "invalid_argument" + && (evidence.responseStatus === "failed" || (status !== undefined && status >= 400 && status < 500))) { + return { recover: true, model: config.model, reason: "devin-invalid-argument" }; + } + if (status !== undefined && status >= 500 && (code === undefined || SERVER_CODES.has(code))) { + return { recover: true, model: config.model, reason: "upstream-unavailable" }; + } + return { recover: false, reason: "unclassified-failure" }; +} diff --git a/src/server/responses/compaction-recovery.ts b/src/server/responses/compaction-recovery.ts new file mode 100644 index 00000000000..d04cd08d8c3 --- /dev/null +++ b/src/server/responses/compaction-recovery.ts @@ -0,0 +1,323 @@ +import type { AdapterEvent, OcxConfig } from "../../types"; +import { routeConcreteModel, type RouteResult } from "../../router"; +import { copyPlainData } from "../../lib/plain-data"; +import { jsonUtf8Bytes } from "../../lib/json-byte-size"; +import type { TranslatorBudget } from "../../lib/translator-budget"; +import { readBoundedResponseBytes } from "../../lib/bounded-body"; +import { isRequestExecutionBudget, type RequestExecutionBudget, type SingleUseDispatchPermit } from "../../lib/request-execution-budget"; +import { isNonReplayableResponse, isNonReplayableUpstreamCode, markResponseNonReplayable, TRANSIENT_RETRY_MAX_ATTEMPTS } from "../../lib/upstream-retry"; +import { isCyberPolicyCode, isTerminalRefusalCode } from "../../lib/errors"; +import { isCanonicalOpenAiForwardProvider, supportsNativeResponsesCompactEndpoint } from "../../providers/openai-tiers"; +import { bridgeToResponsesSSE, formatErrorResponse } from "../../bridge"; +import { buildCompactV1Output, decodeCompactionSummary, encodeCompactionSummary, extractCompactUserMessages } from "../../responses/compaction"; +import { finishRequestAttempt, usageFromResponsesPayload, type RequestLogContext } from "../request-log"; +import { linkRequestSessionLane } from "../request-log-conversation"; +import { isNativePassthroughSseResponse, markNativePassthroughSseResponse, isEagerRelaySseResponse, markEagerRelaySseResponse } from "../relay"; +import type { HandleResponsesOptions } from "./core-options"; +import { consumeComboFailure, createChildPassthroughCallbackGate } from "./core-combo-failure"; +import { preflightComboStreamResponse } from "./combo-stream-preflight"; +import { conversationCarriesUploadedFiles } from "./account-change-state"; +import { selfContainedResponsesBody } from "./reset-replay"; +import { decideCompactionRecovery, readCompactionRecoveryConfig } from "./compaction-recovery-policy"; + +type Options = HandleResponsesOptions & { translatorBudget: TranslatorBudget }; +type Dispatch = (req: Request, config: OcxConfig, log: RequestLogContext, options: Options) => Promise; +const MAX_BYTES = 32 * 1024 * 1024; +const RETAINED_USER_CHARS = 80_000; +const record = (value: unknown): value is Record => !!value && typeof value === "object" && !Array.isArray(value); +const token = (value: unknown): string | undefined => typeof value === "string" && /^[a-zA-Z0-9_-]{1,128}$/.test(value) ? value : undefined; + +function identity(route: RouteResult): string { + return JSON.stringify([route.providerName, route.modelId, route.codexAccountMode ?? "", route.codexAccountNamespace ?? ""]); +} + +function physicalSends(log: RequestLogContext): number { + // activeAttempt normally also belongs to attempts: count each receipt exactly once. + const attempts = new Set([...(log.attempts ?? []), ...(log.activeAttempt ? [log.activeAttempt] : [])]); + return [...attempts].reduce((sum, attempt) => sum + Math.max(0, attempt.sendCount), 0); +} + +/** Reconcile only this leg's physical receipts; never charge already-booked adapter sends again. */ +function settlePhysicalSends(log: RequestLogContext, budget: RequestExecutionBudget, beforeSends: number, beforeUsed: number, reported: number, permit?: SingleUseDispatchPermit): number { + const sent = Math.max(0, physicalSends(log) - beforeSends); + // A legacy fetch leg does not claim the hop through adapterDispatchBudget. Settle its + // prepaid booking explicitly; an adapter-owned leg already claimed it, making this a no-op. + if (sent > 0) permit?.assumeCharge(); + else permit?.release(); + // An external report may settle a prepaid booking without changing used. Its explicit + // receipt outranks the numeric delta, or the same source send would be charged twice. + const unreported = Math.max(0, sent - Math.max(reported, Math.max(0, budget.used - beforeUsed))); + if (unreported > 0) budget.used += unreported; + return sent; +} + +function portableBody(body: Record): boolean { + if (!Array.isArray(body.input) || body.store === true || conversationCarriesUploadedFiles(body)) return false; + // Native ciphertext cannot be summarized by another provider. Never silently replace it with a note. + if (body.input.some(item => record(item) && ["compaction", "compaction_summary", "context_compaction"].includes(String(item.type)) + && typeof item.encrypted_content === "string" && !item.encrypted_content.startsWith("ocx1:"))) return false; + const input = body.input.filter(item => !record(item) || item.type !== "compaction_trigger"); + return selfContainedResponsesBody({ ...body, store: false, input }); +} + +function routed(route: RouteResult): boolean { + return !route.combo && route.routeKind !== "policy" && route.routeReason !== "default-provider" + && !isCanonicalOpenAiForwardProvider(route.provider) + && !supportsNativeResponsesCompactEndpoint(route.providerName, route.provider); +} + +/** One reader, bounded bytes, and an exact replacement body; never clone a live stream. */ +async function bufferedJson(response: Response, signal: AbortSignal): Promise<{ response: Response; json?: Record }> { + const bytes = await readBoundedResponseBytes(response, { signal, maxBytes: MAX_BYTES, inactivityTimeoutMs: 300_000 }); + if (bytes.oversized) return { response: formatErrorResponse(502, "translation_buffer_limit", "Compaction recovery response exceeded its byte limit") }; + const headers = new Headers(response.headers); + headers.delete("content-length"); + headers.delete("content-encoding"); + const replacement = new Response(bytes.bytes, { status: response.status, statusText: response.statusText, headers }); + if (isNonReplayableResponse(response)) markResponseNonReplayable(replacement); + try { + const json: unknown = JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(bytes.bytes)); + return { response: replacement, ...(record(json) ? { json } : {}) }; + } catch { return { response: replacement }; } +} + +/** + * Opt-in recovery for routed compaction only. Native compact/ciphertext, stored continuations, + * policy/combo routes and hosted tools retain their original failure. This owns no credentials. + */ +export async function runWithCompactionRecovery( + req: Request, config: OcxConfig, logCtx: RequestLogContext, options: Options, dispatch: Dispatch, +): Promise { + const recovery = readCompactionRecoveryConfig(config.compactionRecovery); + if (!recovery || options.compactionRecoveryAttempted || options.comboAttempt || (options.inboundWire && options.inboundWire !== "responses")) { + return dispatch(req, config, logCtx, options); + } + let snapshot: Record | undefined; + let snapshotBytes = 0; + let sourceRoute: RouteResult | undefined; + let partialOutput = false; + let replayUnsafe = false; + let adapterError: Extract | undefined; + let sourceFailure: Response | undefined; + let recoveryPermit: SingleUseDispatchPermit | undefined; + let restoreSourceLog: ((attemptStatus?: number) => void) | undefined; + let sourceReportedSends = 0; + const gate = createChildPassthroughCallbackGate(options); + const signal = options.abortSignal ?? req.signal; + const spentBefore = options.sendBudget?.used ?? 0; + const sendsBefore = physicalSends(logCtx); + const firstOptions: Options = { + ...options, + onCompactionRecoverySendsReported(count) { + sourceReportedSends += count; + options.onCompactionRecoverySendsReported?.(count); + }, + onRequestBodyParsed(body) { + options.onRequestBodyParsed?.(body); + if (!record(body) || !Array.isArray(body.input) || typeof body.model !== "string" + || !body.input.some(item => record(item) && item.type === "compaction_trigger") || !portableBody(body)) return; + const users = extractCompactUserMessages(body.input); + if ((users.at(-1)?.length ?? 0) > RETAINED_USER_CHARS) return; + try { + const bytes = jsonUtf8Bytes(body, MAX_BYTES); + const reservation = options.translatorBudget.reserveTransient(bytes, { kind: "request_copies" }); + try { + const copy = copyPlainData(body); + if (!copy.ok) return; + snapshot = copy.value; + snapshotBytes = bytes; + reservation.commitRetained(); + } finally { reservation.release(); } + } catch { /* Optional recovery cannot reject an otherwise valid original request. */ } + }, + onCompactionRecoveryRoute(route) { + options.onCompactionRecoveryRoute?.(route); + if (snapshot && routed(route)) sourceRoute = { ...route }; + }, + onCompactionRecoveryAdapterEvent(event) { + options.onCompactionRecoveryAdapterEvent?.(event); + if (!snapshot) return; + if (event.type === "heartbeat") replayUnsafe ||= event.replayUnsafe === true; + else if (event.type === "error") adapterError = event; + else if (event.type !== "done") partialOutput = true; + }, + onResponseComplete: model => snapshot ? gate.onResponseComplete(model) : options.onResponseComplete?.(model), + onNativePassthroughTerminal: status => snapshot ? gate.onTerminal(status) : options.onNativePassthroughTerminal?.(status), + onNativePassthroughCancel: () => snapshot ? gate.onCancel() : options.onNativePassthroughCancel?.(), + }; + try { + let response = await dispatch(req, config, logCtx, firstOptions); + const keep = (value: Response) => { gate.commit(); return value; }; + if (!snapshot || !sourceRoute || signal.aborted || req.signal.aborted || isNonReplayableResponse(response)) return keep(response); + let target: RouteResult; + try { target = routeConcreteModel(config, recovery.model); } catch { return keep(response); } + if (!routed(target) || identity(sourceRoute) === identity(target)) return keep(response); + const originalModel = firstOptions.compactionRoutingOverride?.sourceModel ?? String(snapshot.model); + // Use the established protocol commit boundary. For runTurn streams the direct event + // observer additionally preserves side-effect heartbeats that the bridge does not publish. + if (response.ok && response.headers.get("content-type")?.includes("text/event-stream")) { + const native = isNativePassthroughSseResponse(response); + const eager = isEagerRelaySseResponse(response); + const preflight = await preflightComboStreamResponse(response, logCtx); + response = preflight.response; + if (preflight.kind !== "failed") { + if (native) markNativePassthroughSseResponse(response); + if (eager) markEagerRelaySseResponse(response); + return keep(response); + } + } else if (response.ok) { + const buffered = await bufferedJson(response, signal); + response = buffered.response; + const json = buffered.json; + if (!json || json.status !== "failed" || (Array.isArray(json.output) && json.output.length > 0)) return keep(response); + // HTTP 200 can carry a failed terminal. The original structured error remains intact. + response = Response.json({ error: json.error, response: json }, { status: adapterError?.status ?? 502 }); + } + if (response.ok || replayUnsafe || partialOutput || signal.aborted || req.signal.aborted) { + if (replayUnsafe) markResponseNonReplayable(response); + return keep(response); + } + const failure = await consumeComboFailure(response, signal); + response = failure.response; + const code = token(adapterError?.code) ?? token(failure.upstreamCode); + const errorType = token(adapterError?.errorType) ?? token(failure.upstreamType); + const budget = options.sendBudget; + // Reset-only fetch legs report their physical receipt but historically leave used alone. + // Reconcile only a failed, eligible compaction; successful/disabled requests stay unchanged. + const sourceSends = budget && isRequestExecutionBudget(budget) + ? settlePhysicalSends(logCtx, budget, sendsBefore, spentBefore, sourceReportedSends) : 0; + const decision = decideCompactionRecovery(recovery, { + requestKind: options.compactionRecoveryKind ?? "compaction-v2", recoveryAttempts: 0, + cancelled: signal.aborted || req.signal.aborted, nonReplayable: !!failure.nonReplayable || isNonReplayableUpstreamCode(code), + partialOutput, toolEffects: replayUnsafe, + remainingSends: budget && isRequestExecutionBudget(budget) ? budget.remainingBaseSends(TRANSIENT_RETRY_MAX_ATTEMPTS) : 0, + originalModel: identity(sourceRoute), fallbackModel: identity(target), provider: sourceRoute.provider.adapter, + httpStatus: adapterError?.status ?? response.status, responseStatus: "failed", errorCode: code, errorType, + authenticationDenied: response.status === 401 || response.status === 403, + policyDenied: isCyberPolicyCode(code), budgetDenied: code === "translation_buffer_limit", + refusal: isTerminalRefusalCode(code), upstreamFailure: sourceSends > 0, + }); + if (!decision.recover || !budget || !isRequestExecutionBudget(budget)) return keep(response); + const fallbackBeforeSends = physicalSends(logCtx); + const fallbackBeforeUsed = budget.used; + const reservation = budget.reserveDispatch({ sendClass: "combo-failover", targetKey: `compaction:${identity(target)}`, countedExternally: true, replaySafe: true }); + if (!reservation.allowed) return keep(response); + recoveryPermit = reservation.permit; + sourceFailure = response; + gate.discard(); + // Finish the first physical attempt while retaining its receipt in attempts[]. + if (logCtx.activeAttempt) finishRequestAttempt(logCtx.activeAttempt, response.status, + Math.max(0, Date.now() - (logCtx.activeAttemptStartedAt ?? Date.now())), logCtx.activeAttempt.usage ?? logCtx.usage); + // Snapshot the original log fields before the fallback rewrites them: a failed fallback + // returns the original failure, so the log must keep describing that failure, not the + // fallback's model, provider, route decision or terminal error. + const SOURCE_LOG_FIELDS = [ + "model", "provider", "providerAdapter", "requestedAlias", "servedModel", "wireModel", + "resolvedModel", "routeDecision", "tierOutcome", "activeTierMetadata", "usage", + "usageFromBridge", "upstreamError", "terminalHttpStatus", "terminalErrorCode", + "terminalIncompleteReason", "errorCode", + ] as const; + const sourceLog: Partial> = {}; + const logFields = logCtx as unknown as Record; + for (const field of SOURCE_LOG_FIELDS) sourceLog[field] = logCtx[field]; + restoreSourceLog = (attemptStatus?: number) => { + if (logCtx.activeAttempt) finishRequestAttempt(logCtx.activeAttempt, attemptStatus ?? sourceFailure!.status, + Math.max(0, Date.now() - (logCtx.activeAttemptStartedAt ?? Date.now())), logCtx.activeAttempt.usage ?? logCtx.usage); + delete logCtx.activeAttempt; + delete logCtx.activeAttemptStartedAt; + for (const field of SOURCE_LOG_FIELDS) { + if (sourceLog[field] === undefined) delete logFields[field]; + else logFields[field] = sourceLog[field]; + } + }; + delete logCtx.activeAttempt; + delete logCtx.activeAttemptStartedAt; + delete logCtx.usage; + delete logCtx.usageFromBridge; + delete logCtx.upstreamError; + delete logCtx.terminalHttpStatus; + delete logCtx.terminalErrorCode; + delete logCtx.terminalIncompleteReason; + const headers = new Headers(req.headers); + headers.delete("authorization"); + headers.delete("chatgpt-account-id"); + headers.delete("content-length"); + headers.delete("content-encoding"); + headers.set("content-type", "application/json"); + const nextBody = { ...snapshot, model: decision.model, stream: false, store: false }; + const bytes = jsonUtf8Bytes(nextBody, MAX_BYTES); + const serialization = options.translatorBudget.reserveTransient(bytes, { kind: "request_copies" }); + let fallback: Response; + let fallbackReportedSends = 0; + try { + const child = new Request(req.url, { method: "POST", headers, body: JSON.stringify(nextBody), signal: req.signal }); + linkRequestSessionLane(req, child); + fallback = await dispatch(child, config, logCtx, { + ...options, compactionRecoveryAttempted: true, compactionRecoveryPermit: recoveryPermit, + compactionRoutingOverride: { sourceModel: originalModel }, + onRequestBodyRead: undefined, onRequestBodyParsed: undefined, + onCompactionRecoveryRoute: undefined, onCompactionRecoveryAdapterEvent: undefined, + onCompactionRecoverySendsReported: count => { + fallbackReportedSends += count; + options.onCompactionRecoverySendsReported?.(count); + }, + onResponseComplete: undefined, onNativePassthroughTerminal: undefined, onNativePassthroughCancel: undefined, + }); + } finally { + settlePhysicalSends(logCtx, budget, fallbackBeforeSends, fallbackBeforeUsed, fallbackReportedSends, recoveryPermit); + serialization.release(); + } + if (signal.aborted || req.signal.aborted) { + void fallback.body?.cancel().catch(() => undefined); + return formatErrorResponse(499, "client_cancelled", "Client cancelled compact request"); + } + if (!fallback.ok) { + void fallback.body?.cancel().catch(() => undefined); + restoreSourceLog?.(fallback.status); + return response; + } + const completed = await bufferedJson(fallback, signal); + const json = completed.json; + const items = json && Array.isArray(json.output) ? json.output : []; + const compactions = items.filter(value => record(value) && value.type === "compaction"); + const permitted = items.every(value => record(value) && (value.type === "compaction" || value.type === "reasoning")); + const item = permitted && compactions.length === 1 ? compactions[0] as Record : undefined; + const summary = item && typeof item.encrypted_content === "string" ? decodeCompactionSummary(item.encrypted_content) : null; + if (json?.status !== "completed" || !summary?.trim()) { + void completed.response.body?.cancel().catch(() => undefined); + restoreSourceLog?.(completed.response.status); + return response; + } + // v1 unpacks this text through buildCompactV1Output, which already re-adds retained user + // messages as items; embedding them here too would duplicate the same text in the output. + const preserved = options.compactionRecoveryKind === "compaction-v1" ? summary : (() => { + const retained = buildCompactV1Output(extractCompactUserMessages(snapshot.input), summary).slice(0, -1); + const userText = extractCompactUserMessages(retained).map((text, index) => `User message ${index + 1}:\n${text}`).join("\n\n"); + return `${summary}\n\nRetained original user messages (verbatim; preserve their goals and constraints):\n${userText}`; + })(); + item!.encrypted_content = encodeCompactionSummary(preserved); + json.model = originalModel; + void completed.response.body?.cancel().catch(() => undefined); + void response.body?.cancel().catch(() => undefined); + if (snapshot.stream === true) { + const usage = usageFromResponsesPayload(json.usage); + async function* events(): AsyncGenerator { + yield { type: "text_delta", text: preserved }; + yield { type: "done", ...(usage ? { usage } : {}) }; + } + return new Response(bridgeToResponsesSSE(events(), originalModel, undefined, undefined, undefined, undefined, 2_000, + { compaction: true, translatorBudget: options.translatorBudget, onCompletedResponse: () => options.onResponseComplete?.(originalModel) }), { headers: { "content-type": "text/event-stream" } }); + } + options.onResponseComplete?.(originalModel); + return Response.json(json); + } catch (error) { + gate.discard(); + if (signal.aborted || req.signal.aborted) return formatErrorResponse(499, "client_cancelled", "Client cancelled compact request"); + if (sourceFailure) { restoreSourceLog?.(); return sourceFailure; } + throw error; + } finally { + if (snapshotBytes) options.translatorBudget.releaseRetained(snapshotBytes, { kind: "request_copies" }); + recoveryPermit?.release(); + snapshot = undefined; + } +} diff --git a/src/server/responses/core-combo-failure.ts b/src/server/responses/core-combo-failure.ts index e222c2ccaea..18e9f6e5c6a 100644 --- a/src/server/responses/core-combo-failure.ts +++ b/src/server/responses/core-combo-failure.ts @@ -117,6 +117,7 @@ export async function consumeComboFailure( ...(nonReplayable ? { nonReplayable: true } : {}), classificationText, ...(normalizedUpstreamCode !== undefined ? { upstreamCode: normalizedUpstreamCode } : {}), + ...(upstreamType !== undefined ? { upstreamType } : {}), ...(!cyberFailure && cooldownRetryAfter !== undefined ? { retryAfter: cooldownRetryAfter } : {}), // The EFFECTIVE classification decides, not the raw status. An upstream that wraps a quota // refusal in a 5xx still carries `x-codex-*-reset-at`, and gating on 402/429 alone threw diff --git a/src/server/responses/core-options.ts b/src/server/responses/core-options.ts index 47c2c982c40..a8cc1e06835 100644 --- a/src/server/responses/core-options.ts +++ b/src/server/responses/core-options.ts @@ -1,5 +1,7 @@ import type { NativeResponseControl } from "./native-response-control"; -import type { OcxUsage, OcxProviderContinuationState, OcxConfig } from "../../types"; +import type { AdapterEvent, OcxUsage, OcxProviderContinuationState, OcxConfig } from "../../types"; +import type { RouteResult } from "../../router"; +import type { SingleUseDispatchPermit } from "../../lib/request-execution-budget"; import type { CodexAuthPolicyConfig, CodexAuthContext } from "../../codex/auth-context"; import type { AdmissionLease } from "../../lib/admission"; import type { DataPlaneAdmission } from "../auth-cors"; @@ -22,6 +24,8 @@ export interface ConsumedComboFailure { classificationText: string; /** Structured upstream `error.code` when present in the failure body. */ upstreamCode?: string; + /** Complete structured provider type, retained for conservative recovery classification. */ + upstreamType?: string; /** Valid numeric/date value used only for cooldown calculation. */ retryAfter?: string; /** Upstream Codex quota-window reset timestamps used for combo cooldowns. */ @@ -52,6 +56,14 @@ export interface ClientEncoderOption { } export interface HandleResponsesOptions { + /** Internal routed-compaction recovery: one logical request, one emergency target. */ + compactionRecoveryAttempted?: boolean; + compactionRecoveryPermit?: SingleUseDispatchPermit; + compactionRecoveryKind?: "compaction-v1" | "compaction-v2"; + onCompactionRecoveryRoute?: (route: RouteResult) => void; + onCompactionRecoveryAdapterEvent?: (event: AdapterEvent) => void; + /** Physical-send reports already delivered to the shared used setter, including booking settlement. */ + onCompactionRecoverySendsReported?: (count: number) => void; /** Internal Claude replay identity; consumed only by the final canonical Go transport. */ claudeGoAffinity?: { sessionLane?: string }; /** Validated Claude metadata identity; projected only into final canonical attempt headers. */ diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 635d5e37c17..168c28c6a30 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -27,6 +27,7 @@ import { createAdapterContinuations } from "./adapter-continuation"; import { deliverAdapterResponse } from "./adapter-delivery"; import { releaseUpstreamHostAdmission } from "../../codex/upstream-host-health"; import { releaseCodexAuthContextProbeLease } from "../../codex/auth-context"; +import { runWithCompactionRecovery } from "./compaction-recovery"; /** Public Responses entry and compatibility exports. Implementations live with their owners. */ @@ -43,7 +44,7 @@ export async function handleResponses( const ownsBudget = options.translatorBudget === undefined; const translatorBudget = options.translatorBudget ?? createTranslatorBudget(); try { - const response = await handleResponsesInner(req, config, logCtx, { + const response = await runWithCompactionRecovery(req, config, logCtx, { ...options, openAiSidecarAuth: options.openAiSidecarAuth === undefined ? captureExplicitOpenAiCallerAuth(req.headers, config) : options.openAiSidecarAuth, @@ -57,7 +58,7 @@ export async function handleResponses( translatorBudget, // Once at ingress, spend observer included: a combo child inherits the parent's holder. sendBudget: options.sendBudget ?? createInferenceSendBudget(req, logCtx), - }); + }, handleResponsesInner); return ownsBudget ? finalizeOwnedTranslatorBudget(response, translatorBudget) : response; } catch (error) { if (ownsBudget) translatorBudget.dispose(); @@ -101,6 +102,7 @@ async function handleResponsesInner( if (requestState instanceof Response) return requestState; const transportState = await prepareResponsesTransport(requestContext, admissionState, requestState); if (transportState instanceof Response) return transportState; + options.onCompactionRecoveryRoute?.(requestState.route); const sidecarState = await prepareResponsesSidecarAuth(requestContext, requestState, transportState); if (sidecarState instanceof Response) return sidecarState; const responseEffects = createResponsesEffects( diff --git a/src/server/responses/request-send-budget.ts b/src/server/responses/request-send-budget.ts index 33f9b3da212..4f2d600bd62 100644 --- a/src/server/responses/request-send-budget.ts +++ b/src/server/responses/request-send-budget.ts @@ -60,6 +60,7 @@ export function createResponsesSendBudget( const noteTransientSends = (used: number): void => { const charged = Math.max(0, used); sendBudget.used += charged; + options.onCompactionRecoverySendsReported?.(charged); chargeWorkflowSends(workflowRootId, charged); }; // Refused before any dispatch, and deliberately not by evicting the root's ledger entry: @@ -156,7 +157,7 @@ export function createResponsesSendBudget( * refused and the request would answer with a synthetic 502 in place of the real 429 the hop * was recovering from. */ - let pendingHopPermit: SingleUseDispatchPermit | undefined; + let pendingHopPermit: SingleUseDispatchPermit | undefined = options.compactionRecoveryPermit; /** * The budget an adapter's OWN dispatch ladder reserves against. * diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 762dbe3eec9..a1d6bf5f256 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -216,6 +216,10 @@ export async function executeResponsesRunTurn( turnParsed: PreparedResponsesRequest["parsed"] = parsed, ): Promise => { const attemptSeq = ++runTurnAttemptSeq; + const emit = (event: AdapterEvent) => { + options.onCompactionRecoveryAdapterEvent?.(event); + targetQueue.push(event); + }; try { if (!pacingSlotAcquired) { await waitForProviderRequestSlot(route.providerName, route.provider, route.modelId, runTurnAbort.signal); @@ -265,7 +269,7 @@ export async function executeResponsesRunTurn( ), onRecoveryWithheld: noteAdapterRecoveryWithheld, }, - targetQueue.push, + emit, ); // LOCAL PATCH (runturn-websearch): adapters may write conversation/ // continuation state onto the object they received; merge it back so @@ -279,7 +283,7 @@ export async function executeResponsesRunTurn( Object.assign(parsed, routeState); } } catch (err) { - targetQueue.push(err instanceof RequestPacingQueueOverloadError + emit(err instanceof RequestPacingQueueOverloadError ? { type: "error", status: 429, diff --git a/src/types/config.ts b/src/types/config.ts index ad59d6f9f65..8c89c4dd824 100644 --- a/src/types/config.ts +++ b/src/types/config.ts @@ -722,6 +722,8 @@ export interface OcxConfig { /** Compaction triggers this override covers; omission means `["manual"]`. */ triggers?: ("manual" | "auto")[]; }; + /** Opt-in failure-only recovery; never replaces the initial compaction model. */ + compactionRecovery?: { enabled: boolean; model: string; allowDevinInvalidArgument?: boolean }; /** * Models hidden from Codex discovery without blocking direct proxy calls. Routed provider ids * are excluded from the catalog + /v1/models entirely. Account-qualified native ids hide only diff --git a/structure/config.md b/structure/config.md index 32dc63be3ad..fa36a5f3cd9 100644 --- a/structure/config.md +++ b/structure/config.md @@ -21,13 +21,13 @@ Native main reauthentication follows the [CLI JSON output contract](runtime.md#n The Codex restart command follows the [CLI restart scope contract](runtime.md#cli-codex-restart-scope). -`src/cli/account-orca-import.ts` exposes an explicit-source, preview-first local import command. -Apply adds pool configuration under the shared mutation lock; the -[source-owned credential contract](codex-home.md#orca-source-owned-account-import) governs -deduplication and credential storage separately from Codex config injection. +`src/cli/account-orca-import.ts` exposes an explicit-source, preview-first local import command. Apply adds pool configuration under the shared mutation lock; +the [source-owned credential contract](codex-home.md#orca-source-owned-account-import) governs deduplication and credential storage separately from Codex config injection. ## Config surface +`src/config/schema/compaction-recovery.ts` strictly validates opt-in `compactionRecovery`; invalid disk values disable it with a warning, while candidate writes reject them. The [failure-only contract](transports/responses-failover.md) leaves provider identity, accounts and client compaction unchanged. + Google providers may persist `googleToolSchemaPolicy` as `compatible` or `reject-lossy`. `ocx provider add --google-tool-schema-policy` is one authoring path and is accepted only when the effective adapter is `google`. Omission remains absent in `config.json`; the adapter resolves it to diff --git a/structure/gui-and-management-api.md b/structure/gui-and-management-api.md index 9f23cb400b3..b388ed4c615 100644 --- a/structure/gui-and-management-api.md +++ b/structure/gui-and-management-api.md @@ -22,6 +22,11 @@ linked sources after the quota await to reject revoked or rotated captures. ## Compact desktop usage +`GET /api/settings` and successful `PUT /api/settings` report `compactionRecovery` as a block or +null. The write accepts the strict failure-recovery schema or null to remove it, and restores the +prior field if persistence fails. Changing this field alone does not converge catalogs or inject +Codex configuration. `tests/config/compaction-recovery-settings.test.ts` covers that boundary. + The standalone `/#/tray` GUI route presents local usage and account limits without the full dashboard navigation. It reuses the existing API session and fetch wrapper; it has no Tauri IPC capability. Companion settings control its sections and chart. Account diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 76c41755933..e6bb05158c3 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -1,5 +1,22 @@ # Responses Failover And Replay +`src/server/responses/compaction-recovery-policy.ts` is a pure eligibility policy, not a dispatcher. +It requires explicit configuration and normalized attempt evidence, preserves ordinary requests, +and refuses cancellation, committed semantic output, tool effects, protected failures, exhausted +send budgets and repeated recovery. Its Devin `invalid_argument` exception is separately opted in; +an opaque HTTP 400 alone never grants replay. The caller owns canonical target resolution, output +validation and shared-budget reservation. `compaction-recovery.ts` connects this policy to +self-contained routed v1/v2 compaction in `core.ts` and `compact.ts`; normal and successful +requests keep their original route. One configured emergency target shares the original send +and translation budgets. Physical-send receipts and explicit retry-helper reports reconcile legacy +fetch sends without double charging external reservations; one prepaid emergency permit is shared +with adapter dispatch, and only additional retries draw from the remainder. Adapter observers retain partial-output and structured denial evidence +before response projection. Native encrypted compaction, uploaded files, stored continuations, +and policy/combo routes are excluded. Emergency output must contain one readable portable +compaction item; recent original user messages are retained verbatim, and recovery failure keeps +the original failure. `tests/responses/responses-compaction-recovery-policy.test.ts` and +`tests/responses/responses-compaction-recovery.test.ts` pin these boundaries. + Retry, replay, and combo failover on the Responses data plane: upstream reset retry, the ambiguous-resend gate and replay boundary, combo quota fallback and commit boundaries, compaction routing overrides, and output headroom. The endpoint and dispatch rules they build on are in diff --git a/tests/config/compaction-recovery-settings.test.ts b/tests/config/compaction-recovery-settings.test.ts new file mode 100644 index 00000000000..b5cff096e47 --- /dev/null +++ b/tests/config/compaction-recovery-settings.test.ts @@ -0,0 +1,67 @@ +import { describe, expect, test } from "bun:test"; +import { getDefaultConfig, validateConfigCandidate } from "../../src/config"; +import { configSchema } from "../../src/config/schema/config-schema"; +import { compactionRecoverySchema } from "../../src/config/schema/compaction-recovery"; +import { handleConfigRoutes } from "../../src/server/management/config-routes"; +import type { ManagementContext } from "../../src/server/management/context"; +import type { OcxConfig } from "../../src/types"; + +const enabled = { enabled: true, model: "emergency/model", allowDevinInvalidArgument: true }; +function harness(failSave = false) { + const config = getDefaultConfig(); + let saves = 0; + const call = async (body: unknown) => { + const url = new URL("http://localhost/api/settings"); + return handleConfigRoutes({ + url, req: new Request(url, { method: "PUT", headers: { "content-type": "application/json" }, body: JSON.stringify(body) }), + config, version: "test", deps: { saveConfigPreservingClaudeCode: (_c: OcxConfig) => { + saves++; if (failSave) throw new Error("fixture-save-failed"); + } }, + } as unknown as ManagementContext); + }; + return { config, call, saves: () => saves }; +} + +describe("compaction failure recovery settings", () => { + test("absent is off; an explicit valid configuration survives parsing", () => { + expect(getDefaultConfig().compactionRecovery).toBeUndefined(); + const result = validateConfigCandidate({ ...getDefaultConfig(), compactionRecovery: enabled }); + expect(result.ok).toBe(true); + if (result.ok) expect(result.config.compactionRecovery).toEqual(enabled); + expect(compactionRecoverySchema.parse({ enabled: false, model: "emergency/model" }).enabled).toBe(false); + }); + test("invalid hand edits degrade to off but mutation validation rejects them", () => { + for (const value of [{ ...enabled, retryCount: 99 }, { ...enabled, enabled: "true" }, { ...enabled, model: "" }, { ...enabled, model: "x\ny" }, { ...enabled, allowDevinInvalidArgument: 1 }]) { + const input = { ...getDefaultConfig(), compactionRecovery: value }; + expect(configSchema.parse(input).compactionRecovery).toBeUndefined(); + expect(validateConfigCandidate(input).ok).toBe(false); + } + }); + test("management writes and clears recovery without toggling routing or login mode", async () => { + const h = harness(); + h.config.compactionRouting = { model: "original/manual" }; + const res = await h.call({ compactionRecovery: enabled }); + expect(res?.status).toBe(200); + expect(h.config.compactionRecovery).toEqual(enabled); + expect(h.config.compactionRouting).toEqual({ model: "original/manual" }); + expect(h.config.codexDesktopAuthless).toBeUndefined(); + expect(h.config.codexClientCompaction).toBeUndefined(); + expect((await res!.json() as any).compactionRecovery).toEqual(enabled); + expect((await h.call({ compactionRecovery: null }))?.status).toBe(200); + expect(h.config.compactionRecovery).toBeUndefined(); + expect(h.saves()).toBe(2); + }); + test("rejected management configuration never reaches persistence", async () => { + const h = harness(); + expect((await h.call({ compactionRecovery: { ...enabled, enabled: 1 } }))?.status).toBe(400); + expect(h.saves()).toBe(0); + expect(h.config.compactionRecovery).toBeUndefined(); + }); + test("a save failure restores the prior recovery field", async () => { + const h = harness(true); + const original = { enabled: false, model: "kept/model" }; + h.config.compactionRecovery = original; + await expect(h.call({ compactionRecovery: enabled })).rejects.toThrow("fixture-save-failed"); + expect(h.config.compactionRecovery).toEqual(original); + }); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 2eed0eac3af..2393c577c23 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1,4 +1,7 @@ { + "responses-compaction-recovery.test.ts": "responses", + "compaction-recovery-settings.test.ts": "config", + "responses-compaction-recovery-policy.test.ts": "responses", "deepseek-artifact-tool-schema.test.ts": "providers", "client-config-export-output-limit.test.ts": "config", "openai-chat-serialized-tool-call-scaling.test.ts": "adapters/openai", diff --git a/tests/helpers/responses-core-source.ts b/tests/helpers/responses-core-source.ts index 2d3805ea0b7..7460b64ada9 100644 --- a/tests/helpers/responses-core-source.ts +++ b/tests/helpers/responses-core-source.ts @@ -32,6 +32,8 @@ export const RESPONSES_CORE_MODULES = [ "request-prepare.ts", "shadow-target-availability.ts", "compaction-routing.ts", + "compaction-recovery.ts", + "compaction-recovery-policy.ts", "request-transport.ts", "request-sidecar-auth.ts", "response-effects.ts", diff --git a/tests/responses/responses-compaction-recovery-policy.test.ts b/tests/responses/responses-compaction-recovery-policy.test.ts new file mode 100644 index 00000000000..87c7e2734a6 --- /dev/null +++ b/tests/responses/responses-compaction-recovery-policy.test.ts @@ -0,0 +1,102 @@ +import { describe, expect, test } from "bun:test"; +import { decideCompactionRecovery, readCompactionRecoveryConfig, type CompactionRecoveryEvidence } from "../../src/server/responses/compaction-recovery-policy"; + +const config = { enabled: true, model: "emergency-alias", allowDevinInvalidArgument: true }; +const failure = (patch: Partial = {}): CompactionRecoveryEvidence => ({ + requestKind: "compaction-v2", recoveryAttempts: 0, cancelled: false, nonReplayable: false, + partialOutput: false, toolEffects: false, remainingSends: 1, + originalModel: "devin/main/swe-2", fallbackModel: "google/main/gemini", provider: "devin", + httpStatus: 400, responseStatus: "failed", errorCode: "invalid_argument", + authenticationDenied: false, policyDenied: false, budgetDenied: false, refusal: false, + upstreamFailure: true, ...patch, +}); + +describe("compaction emergency recovery policy", () => { + test("accepts the explicit Devin exception on v1 and HTTP-200 failed v2", () => { + expect(decideCompactionRecovery(config, failure({ requestKind: "compaction-v1" }))).toEqual({ + recover: true, model: "emergency-alias", reason: "devin-invalid-argument", + }); + expect(decideCompactionRecovery(config, failure({ httpStatus: 200 }))).toMatchObject({ recover: true }); + }); + + test("successful compaction and ordinary generation never request fallback", () => { + expect(decideCompactionRecovery(config, failure({ httpStatus: 200, responseStatus: "completed" }))).toEqual({ recover: false, reason: "succeeded" }); + expect(decideCompactionRecovery(config, failure({ requestKind: "ordinary" }))).toEqual({ recover: false, reason: "ordinary-request" }); + }); + + test("the exception needs both opt-ins, exact serving provider and structured code", () => { + for (const candidate of [undefined, null, {}, { ...config, enabled: false }]) { + expect(decideCompactionRecovery(candidate, failure()).recover).toBe(false); + } + for (const patch of [ + { provider: "google" }, { errorCode: "invalid_request_error" }, { errorCode: undefined }, + { errorCode: "message_contains_invalid_argument" }, { upstreamFailure: false }, + { httpStatus: 200, responseStatus: "unknown" as const }, + { httpStatus: 502, responseStatus: "unknown" as const }, + ]) expect(decideCompactionRecovery(config, failure(patch)).recover).toBe(false); + expect(decideCompactionRecovery({ enabled: true, model: "google/gemini" }, failure()).recover).toBe(false); + }); + + test("caller cancellation and semantic output or tool effects stop replay", () => { + for (const patch of [{ cancelled: true }, { httpStatus: 499 }, { errorCode: "client_cancelled" }]) { + expect(decideCompactionRecovery(config, failure(patch))).toEqual({ recover: false, reason: "cancelled" }); + } + for (const patch of [{ nonReplayable: true }, { partialOutput: true }, { toolEffects: true }]) { + expect(decideCompactionRecovery(config, failure(patch))).toEqual({ recover: false, reason: "unsafe-replay" }); + } + }); + + test("auth, policy, refusal and budget evidence outrank the provider exception", () => { + for (const patch of [ + { authenticationDenied: true }, { policyDenied: true }, { budgetDenied: true }, { refusal: true }, + { httpStatus: 401 }, { httpStatus: 403 }, { httpStatus: 402 }, { httpStatus: 429 }, + ...["authentication_error", "permission_denied", "origin_rejected", "cyber_policy_violation", + "content_filter", "request_send_budget_exhausted", "failed_precondition", "admission_model_denied"] + .map(errorType => ({ errorType })), + ]) expect(decideCompactionRecovery(config, failure(patch))).toEqual({ recover: false, reason: "protected-failure" }); + }); + + test("same canonical target, recursive recovery and exhausted allowance stop dispatch", () => { + expect(decideCompactionRecovery(config, failure({ fallbackModel: "devin/main/swe-2" }))).toEqual({ recover: false, reason: "same-model" }); + for (const recoveryAttempts of [1, 2, 100]) { + expect(decideCompactionRecovery(config, failure({ recoveryAttempts }))).toEqual({ recover: false, reason: "already-attempted" }); + } + expect(decideCompactionRecovery(config, failure({ remainingSends: 0 }))).toEqual({ recover: false, reason: "budget-exhausted" }); + }); + + test("recognizes typed overflow and compact-output failures without a Devin exception", () => { + const configured = { enabled: true, model: "google/gemini" }; + expect(decideCompactionRecovery(configured, failure({ provider: "google", errorCode: "context_length_exceeded" }))).toMatchObject({ recover: true, reason: "context-overflow" }); + expect(decideCompactionRecovery(configured, failure({ httpStatus: 200, errorCode: "invalid_compaction_output" }))).toMatchObject({ recover: true, reason: "compaction-output" }); + expect(decideCompactionRecovery(configured, failure({ httpStatus: 503, errorCode: "unavailable" }))).toMatchObject({ recover: true, reason: "upstream-unavailable" }); + // classifyError normalizes codeless 5xx bodies and overloads before the evidence reaches us. + expect(decideCompactionRecovery(configured, failure({ httpStatus: 500, errorCode: "upstream_server_error" }))).toMatchObject({ recover: true, reason: "upstream-unavailable" }); + expect(decideCompactionRecovery(configured, failure({ httpStatus: 503, errorCode: "server_is_overloaded" }))).toMatchObject({ recover: true, reason: "upstream-unavailable" }); + expect(decideCompactionRecovery(configured, failure({ httpStatus: 503, errorCode: "unrecognized_failure" })).recover).toBe(false); + }); + + test("missing or malformed evidence cannot silently grant a send", () => { + for (const patch of [ + { recoveryAttempts: NaN }, { remainingSends: Infinity }, { remainingSends: -1 }, + { partialOutput: undefined }, { toolEffects: undefined }, { originalModel: "" }, + { fallbackModel: "bad\nmodel" }, { httpStatus: 99 }, { errorCode: "x".repeat(129) }, + ]) expect(decideCompactionRecovery(config, failure(patch as Partial))).toEqual({ recover: false, reason: "invalid-evidence" }); + }); + + test("configuration is bounded, detached and requires explicit booleans", () => { + expect(readCompactionRecoveryConfig(config)).toEqual(config); + expect(readCompactionRecoveryConfig(config)).not.toBe(config); + for (const model of ["", " x", "x ", "x\ny", "x".repeat(513)]) { + expect(readCompactionRecoveryConfig({ ...config, model })).toBeNull(); + } + expect(readCompactionRecoveryConfig({ ...config, allowDevinInvalidArgument: "true" })).toBeNull(); + expect(readCompactionRecoveryConfig({ ...config, retries: 10 })).toBeNull(); + }); + + test("decision is pure and does not consume the caller's shared send allowance", () => { + const evidence = Object.freeze(failure()); + expect(decideCompactionRecovery(Object.freeze(config), evidence).recover).toBe(true); + expect(evidence.remainingSends).toBe(1); + expect(evidence.recoveryAttempts).toBe(0); + }); +}); diff --git a/tests/responses/responses-compaction-recovery.test.ts b/tests/responses/responses-compaction-recovery.test.ts new file mode 100644 index 00000000000..a60f3af0105 --- /dev/null +++ b/tests/responses/responses-compaction-recovery.test.ts @@ -0,0 +1,348 @@ +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import { ADAPTER_REGISTRY } from "../../src/adapters/registry"; +import { getDefaultConfig } from "../../src/config"; +import { handleResponses, handleResponsesCompact } from "../../src/server/responses"; +import { decodeCompactionSummary } from "../../src/responses/compaction"; +import { createRequestExecutionBudget } from "../../src/lib/request-execution-budget"; +import { createTranslatorBudget } from "../../src/lib/translator-budget"; +import { jsonUtf8Bytes } from "../../src/lib/json-byte-size"; +import type { AdapterEvent, OcxConfig, OcxParsedRequest } from "../../src/types"; +import type { RequestLogContext } from "../../src/server/request-log"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; + +const originalFetch = globalThis.fetch; +const sourceError: AdapterEvent = { type: "error", status: 400, code: "invalid_argument", message: "Source rejected compact fixture" }; +let sourceEvents: AdapterEvent[]; +let fallbackEvents: AdapterEvent[]; +let calls: Array<{ model: string; parsed: OcxParsedRequest }>; +let releaseSpend: (() => void) | undefined; +let restoreFactory: (() => void) | undefined; +let restoreChatFactory: (() => void) | undefined; +let abortOnSource: AbortController | undefined; + +function settings(): OcxConfig { + return { + ...getDefaultConfig(), defaultProvider: "source", + providers: { + source: { adapter: "devin", authMode: "key", apiKey: "fixture-only", baseUrl: "https://source.example" }, + emergency: { adapter: "devin", authMode: "key", apiKey: "fixture-only", baseUrl: "https://emergency.example" }, + }, + compactionRecovery: { enabled: true, model: "emergency/rescue", allowDevinInvalidArgument: true }, + }; +} + +function body(stream = false, compact = true): Record { + return { + model: "source/swe-2", stream, store: false, max_output_tokens: 512, + input: [ + { type: "message", role: "user", content: "Remember marker ALPHA-729." }, + { type: "message", role: "assistant", content: "Recorded." }, + { type: "message", role: "user", content: "Latest goal: finish the report, preserve the marker." }, + ...(compact ? [{ type: "compaction_trigger" }] : []), + ], + }; +} + +function request(payload = body(), path = "responses", signal?: AbortSignal): Request { + return new Request(`http://localhost/v1/${path}`, { + method: "POST", headers: { "content-type": "application/json", session_id: "recovery-fixture" }, + body: JSON.stringify(payload), signal, + }); +} + +beforeEach(() => { + releaseSpend = acquireOwnedSpendHome(); + calls = []; + sourceEvents = [sourceError]; + fallbackEvents = [ + { type: "thinking_delta", thinking: "Prepare the handoff." }, + { type: "text_delta", text: "Work is pending; resume the report." }, + { type: "done", usage: { inputTokens: 4, outputTokens: 6, totalTokens: 10 } }, + ]; + abortOnSource = undefined; + globalThis.fetch = (async () => { throw new Error("Unexpected network request in recovery fixture"); }) as typeof fetch; + const factory = spyOn(ADAPTER_REGISTRY.devin, "create").mockImplementation((_provider, context) => ({ + name: "devin", + reportsPhysicalSends: true, + buildRequest() { throw new Error("runTurn fixture must not build HTTP requests"); }, + async *parseStream() { throw new Error("runTurn fixture must not parse HTTP responses"); }, + async runTurn(parsed, incoming, emit) { + const send = incoming.sendBudget?.reserveDispatch({ sendClass: "initial", targetKey: `${context.providerId}/${parsed.modelId}` }); + if (send && (!send.allowed || !send.permit.use())) { + emit({ type: "error", status: 429, code: "request_send_budget_exhausted", message: "Fixture shared send allowance exhausted" }); + return; + } + incoming.onPhysicalSend?.({ ordinal: 1 }); + calls.push({ model: parsed.modelId, parsed: structuredClone(parsed) }); + const isSource = context.providerId === "source"; + for (const event of isSource ? sourceEvents : fallbackEvents) emit(event); + if (isSource) abortOnSource?.abort(); + }, + })); + restoreFactory = () => factory.mockRestore(); +}); + +afterEach(() => { + releaseSpend?.(); + releaseSpend = undefined; + restoreFactory?.(); + restoreFactory = undefined; + restoreChatFactory?.(); + restoreChatFactory = undefined; + globalThis.fetch = originalFetch; +}); + +describe("routed compaction emergency integration", () => { + test("source success and ordinary requests never use the emergency model", async () => { + sourceEvents = [{ type: "text_delta", text: "Source summary" }, { type: "done" }]; + for (const compact of [true, false]) { + const response = await handleResponses(request(body(false, compact)), settings(), { model: "", provider: "" }); + expect((await response.json()).status).toBe("completed"); + } + expect(calls.map(call => call.model)).toEqual(["swe-2", "swe-2"]); + }); + + test.each([false, true])("v2 failed terminal recovers once and retains user goals (stream=%s)", async stream => { + const config = settings(); + const before = structuredClone(config); + const completed: string[] = []; + const log: RequestLogContext = { model: "", provider: "" }; + const response = await handleResponses(request(body(stream)), config, log, { onResponseComplete: model => completed.push(model) }); + let summary: string | null; + if (stream) { + const text = await response.text(); + expect(text).not.toContain("Source rejected compact fixture"); + const terminal = text.split("\n").filter(line => line.startsWith("data: ") && line !== "data: [DONE]") + .map(line => JSON.parse(line.slice(6))).find(event => event.type === "response.completed"); + summary = decodeCompactionSummary(terminal.response.output.find((item: { type: string }) => item.type === "compaction").encrypted_content); + expect(terminal.response.usage.total_tokens).toBe(10); + } else { + const json = await response.json(); + expect(json.status).toBe("completed"); + summary = decodeCompactionSummary(json.output.find((item: { type: string }) => item.type === "compaction").encrypted_content); + } + expect(summary).toContain("ALPHA-729"); + expect(summary).toContain("Latest goal: finish the report"); + expect(calls.map(call => call.model)).toEqual(["swe-2", "rescue"]); + expect(calls[1]!.parsed.options.maxOutputTokens).toBe(512); + expect(calls[1]!.parsed.context.tools).toBeUndefined(); + expect(completed).toEqual(["source/swe-2"]); + expect(log.provider).toBe("emergency"); + expect(config).toEqual(before); + }); + + test("routed v1 returns replacement history retaining original user text once", async () => { + const response = await handleResponsesCompact(request(body(false, false), "responses/compact"), settings(), { model: "", provider: "" }); + expect(response.status).toBe(200); + const json = await response.json(); + const output = JSON.stringify(json.output); + expect(output).toContain("ALPHA-729"); + expect(output).toContain("Latest goal: finish the report"); + // Retained messages belong to v1 output items; embedding them in the summary duplicates the text. + expect(output).not.toContain("Retained original user messages"); + expect(output.split("ALPHA-729").length - 1).toBe(1); + expect(calls.map(call => call.model)).toEqual(["swe-2", "rescue"]); + }); + + test("a failed fallback returns and logs the original failure", async () => { + fallbackEvents = [{ type: "error", status: 503, code: "server_is_overloaded", message: "Emergency overloaded fixture" }]; + const log: RequestLogContext = { model: "", provider: "" }; + const response = await handleResponses(request(body(false)), settings(), log); + expect(response.status).toBe(400); + expect(calls.map(call => call.model)).toEqual(["swe-2", "rescue"]); + // The returned failure is the source's, so the log must describe it too — not the fallback's. + expect(log.provider).toBe("source"); + expect(log.model).toBe("swe-2"); + expect(log.requestedAlias).toBe("source/swe-2"); + expect(log.activeAttempt).toBeUndefined(); + // The fallback's own failed attempt stays recorded; its wire status was a failed 200 terminal. + expect(log.attempts?.map(attempt => attempt.status)).toEqual([400, 200]); + }); + + test("an existing unconditional override keeps its original logical model on recovery", async () => { + const config = settings(); + config.compactionRouting = { model: "source/swe-2" }; + const payload = { ...body(), model: "source/normal", client_metadata: { + "x-codex-turn-metadata": JSON.stringify({ request_kind: "compaction", compaction: { trigger: "manual" } }), + } }; + const completed: string[] = []; + const response = await handleResponses(request(payload), config, { model: "", provider: "" }, { onResponseComplete: model => completed.push(model) }); + expect((await response.json()).model).toBe("source/normal"); + expect(calls.map(call => call.model)).toEqual(["swe-2", "rescue"]); + expect(completed).toEqual(["source/normal"]); + }); + + test.each([false, true])("fetch adapter hidden text then failure cannot replay (stream=%s)", async stream => { + const config = settings(); + config.providers.source = { adapter: "openai-chat", authMode: "key", apiKey: "fixture-only", baseUrl: "https://source.example/v1" }; + const events: AdapterEvent[] = [{ type: "text_delta", text: "Private partial compact text" }, { type: "error", status: 500, errorType: "upstream_error", message: "Source failed after partial text" }]; + const factory = spyOn(ADAPTER_REGISTRY["openai-chat"], "create").mockImplementation(() => ({ + name: "openai-chat", buildRequest() { return { url: "https://source.example/v1/chat/completions", method: "POST", headers: {}, body: "{}" }; }, + async *parseStream() { yield* events; }, async parseResponse() { return events; }, + })); + restoreChatFactory = () => factory.mockRestore(); + let fetches = 0; + globalThis.fetch = (async () => { fetches++; return Response.json({ fixture: true }); }) as typeof fetch; + const response = await handleResponses(request(body(stream)), config, { model: "", provider: "" }); + expect(await response.text()).toContain("Source failed after partial text"); + expect(fetches).toBe(1); + expect(calls).toHaveLength(0); + }); + + test("fetch HTTP 500 authentication type survives client formatting and forbids recovery", async () => { + const config = settings(); + config.providers.source = { adapter: "openai-chat", authMode: "key", apiKey: "fixture-only", baseUrl: "https://source.example/v1" }; + let fetches = 0; + globalThis.fetch = (async () => { + fetches++; + return Response.json({ error: { type: "authentication_error", message: "denied" } }, { status: 500 }); + }) as typeof fetch; + const response = await handleResponses(request(), config, { model: "", provider: "" }); + expect(response.status).toBe(500); + await response.text(); + expect(fetches).toBe(1); + expect(calls).toHaveLength(0); + }); + + test.each([[1, false], [2, false], [2, true]] as const)("fetch source and fetch emergency share cap=%s transient=%s", async (cap, transient) => { + const config = settings(); + config.providers.source = { adapter: "openai-chat", authMode: "key", apiKey: "fixture-only", baseUrl: "https://source.example/v1" }; + config.providers.emergency = { adapter: "openai-chat", authMode: "key", apiKey: "fixture-only", baseUrl: "https://emergency.example/v1", ...(transient ? { transientRetryOn5xx: { attempts: 3 } } : {}) }; + const requests: string[] = []; + globalThis.fetch = (async (input: unknown) => { + const url = String(input); + requests.push(url); + return url.includes("source.example") + ? Response.json({ error: { type: "invalid_request_error", code: "context_length_exceeded", message: "Source input context is full" } }, { status: 400 }) + : Response.json({ choices: [{ message: { role: "assistant", content: "Resume the report." }, finish_reason: "stop" }], usage: { prompt_tokens: 4, completion_tokens: 6, total_tokens: 10 } }); + }) as typeof fetch; + const budget = createRequestExecutionBudget({ maxTotalModelSends: cap, baseSendAllowance: cap, finalRecoveryAllowance: 0, maxAlternateTargetSends: 1, maxTargetTransitions: 1 }); + const response = await handleResponses(request(), config, { model: "", provider: "" }, { sendBudget: budget }); + const json = await response.json(); + expect(requests).toHaveLength(cap); + expect(budget.used).toBe(cap); + if (cap === 2) { + expect(json.status).toBe("completed"); + expect(decodeCompactionSummary(json.output.find((item: { type: string }) => item.type === "compaction").encrypted_content)).toContain("ALPHA-729"); + } else expect(json.error.code).toBe("context_length_exceeded"); + }); + + test.each([false, true])("externally booked source settles once (transient=%s)", async transient => { + const config = settings(); + config.providers.source = { adapter: "openai-chat", authMode: "key", apiKey: "fixture-only", baseUrl: "https://source.example/v1", ...(transient ? { transientRetryOn5xx: { attempts: 3 } } : {}) }; + let sourceRequests = 0; + globalThis.fetch = (async () => { + sourceRequests++; + return Response.json({ error: { code: "context_length_exceeded", message: "Source context is full" } }, { status: 400 }); + }) as typeof fetch; + const budget = createRequestExecutionBudget({ maxTotalModelSends: 2, baseSendAllowance: 2, finalRecoveryAllowance: 0, maxAlternateTargetSends: 1, maxTargetTransitions: 1 }); + const reservation = budget.reserveDispatch({ sendClass: "initial", targetKey: "source/swe-2", countedExternally: true }); + expect(reservation.allowed).toBe(true); + const response = await handleResponses(request(), config, { model: "", provider: "" }, { sendBudget: budget }); + expect((await response.json()).status).toBe("completed"); + expect(sourceRequests).toBe(1); + expect(calls.map(call => call.model)).toEqual(["rescue"]); + expect(budget.used).toBe(2); + }); + + test("runTurn emergency consumes its prepaid permit once rather than taking another send", async () => { + const budget = createRequestExecutionBudget({ maxTotalModelSends: 2, baseSendAllowance: 2, finalRecoveryAllowance: 0, maxAlternateTargetSends: 1, maxTargetTransitions: 1 }); + const response = await handleResponses(request(), settings(), { model: "", provider: "" }, { sendBudget: budget }); + expect((await response.json()).status).toBe("completed"); + expect(calls.map(call => call.model)).toEqual(["swe-2", "rescue"]); + expect(budget.used).toBe(2); + }); + + test("emergency transient 5xx retry cannot exceed the shared cap", async () => { + const config = settings(); + config.providers.emergency = { adapter: "openai-chat", authMode: "key", apiKey: "fixture-only", baseUrl: "https://emergency.example/v1", transientRetryOn5xx: { attempts: 3 } }; + let emergencyRequests = 0; + globalThis.fetch = (async () => { + emergencyRequests++; + return Response.json({ error: { code: "server_error", message: "Emergency unavailable" } }, { status: 500 }); + }) as typeof fetch; + const budget = createRequestExecutionBudget({ maxTotalModelSends: 4, baseSendAllowance: 4, finalRecoveryAllowance: 0, maxAlternateTargetSends: 1, maxTargetTransitions: 1 }); + const response = await handleResponses(request(), config, { model: "", provider: "" }, { sendBudget: budget }); + expect(await response.text()).toContain("Source rejected compact fixture"); + expect(emergencyRequests).toBe(3); + expect(1 + emergencyRequests).toBeLessThanOrEqual(4); + expect(budget.used).toBe(1 + emergencyRequests); + }); + + test("an emergency rejected before sending refunds the unused reservation", async () => { + const config = settings(); + config.providers.emergency = { adapter: "openai-chat", authMode: "key", baseUrl: "https://emergency.example/v1" }; + const budget = createRequestExecutionBudget(); + const response = await handleResponses(request(), config, { model: "", provider: "" }, { sendBudget: budget }); + expect(await response.text()).toContain("Source rejected compact fixture"); + expect(calls.map(call => call.model)).toEqual(["swe-2"]); + expect(budget.used).toBe(1); + expect(budget.alternateTargetSends).toBe(0); + }); + + test.each(["disabled", "generic-400", "policy", "auth", "partial", "side-effect", "same-model", "opaque", "continuation"])("keeps source failure: %s", async variant => { + const config = settings(); + const payload = body(true); + if (variant === "disabled") config.compactionRecovery!.allowDevinInvalidArgument = false; + if (variant === "generic-400") sourceEvents = [{ ...sourceError, code: "invalid_request_error" } as AdapterEvent]; + if (variant === "policy") sourceEvents = [{ ...sourceError, code: "cyber_policy" } as AdapterEvent]; + if (variant === "auth") sourceEvents = [{ ...sourceError, status: 403, code: "permission_denied" } as AdapterEvent]; + if (variant === "partial") sourceEvents = [{ type: "text_delta", text: "Partial source text" }, sourceError]; + if (variant === "side-effect") sourceEvents = [{ type: "heartbeat", replayUnsafe: true }, sourceError]; + if (variant === "same-model") config.compactionRecovery!.model = "source/swe-2"; + if (variant === "opaque") (payload.input as unknown[]).unshift({ type: "compaction", encrypted_content: "native-opaque-fixture" }); + if (variant === "continuation") payload.previous_response_id = "missing-fixture"; + const response = await handleResponses(request(payload), config, { model: "", provider: "" }); + await response.text(); + expect(calls.filter(call => call.model === "rescue")).toHaveLength(0); + if (variant !== "continuation") expect(calls.map(call => call.model)).toEqual(["swe-2"]); + }); + + test.each(["error", "empty", "truncated"])("failed emergency %s preserves the source error without recursive recovery", async outcome => { + fallbackEvents = outcome === "error" ? [{ ...sourceError, message: "Different emergency failure" } as AdapterEvent] + : outcome === "empty" ? [{ type: "done" }] + : [{ type: "text_delta", text: "Truncated summary" }, { type: "done", stopReason: "max_tokens" }]; + const response = await handleResponses(request(), settings(), { model: "", provider: "" }); + const text = await response.text(); + expect(text).toContain("Source rejected compact fixture"); + expect(text).not.toContain("Different emergency failure"); + expect(calls.map(call => call.model)).toEqual(["swe-2", "rescue"]); + }); + + test("cancellation after source error does not dispatch emergency", async () => { + abortOnSource = new AbortController(); + const response = await handleResponses(request(body(), "responses", abortOnSource.signal), settings(), { model: "", provider: "" }); + await response.text(); + expect(calls.map(call => call.model)).toEqual(["swe-2"]); + }); + + test("one shared send budget blocks emergency when the source consumes the allowance", async () => { + const sendBudget = createRequestExecutionBudget({ maxTotalModelSends: 1, baseSendAllowance: 1, finalRecoveryAllowance: 0, maxAlternateTargetSends: 0, maxTargetTransitions: 0 }); + const translatorBudget = createTranslatorBudget(); + try { + const response = await handleResponses(request(), settings(), { model: "", provider: "" }, { sendBudget, translatorBudget }); + await response.text(); + expect(calls.map(call => call.model)).toEqual(["swe-2"]); + expect(sendBudget.used).toBe(1); + // The ingress-owned body observation remains until its caller disposes the shared budget; + // the additional recovery snapshot has already released its separate retained charge. + expect(translatorBudget.snapshot().currentBytes).toBe(jsonUtf8Bytes(body())); + } finally { translatorBudget.dispose(); } + expect(translatorBudget.snapshot().currentBytes).toBe(0); + }); + + test("canonical native v1 stays on its existing compact path and never invokes routed recovery", async () => { + const config = settings(); + config.providers["openai-apikey"] = { adapter: "openai-responses", authMode: "key", baseUrl: "https://api.openai.com/v1", apiKey: "fixture-only" }; + const urls: string[] = []; + globalThis.fetch = (async (input: unknown) => { + urls.push(String(input)); + return Response.json({ error: { code: "invalid_argument", message: "Native compact fixture failure" } }, { status: 400 }); + }) as typeof fetch; + const response = await handleResponsesCompact(request({ ...body(false, false), model: "openai-apikey/gpt-4.1" }, "responses/compact"), config, { model: "", provider: "" }); + expect(response.status).toBe(400); + await response.text(); + expect(urls).toEqual(["https://api.openai.com/v1/responses/compact"]); + expect(calls).toHaveLength(0); + }); +});