diff --git a/docs-site/src/content/docs/ja/reference/configuration/providers.md b/docs-site/src/content/docs/ja/reference/configuration/providers.md index 29e09c174..bc0b69a74 100644 --- a/docs-site/src/content/docs/ja/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ja/reference/configuration/providers.md @@ -86,6 +86,7 @@ namespace 付き combo alias はその namespace prefix に selector を再利 | `noPenaltyModels?` | `string[]` |存在/周波数ペナルティを拒否するモデル。 | | `parallelToolCalls?` | `boolean` |並列ツール呼び出しを切り替えます。 OpenAI Chat はデフォルトでオンになっています。非チャット アダプターは明示的な `true` でのみアドバタイズします。 | | `responsesItemIdRepair?` | `{ message?: string[]; reasoning?: string[]; repairMissingTerminalIds?: boolean }` |正確なプレースホルダー ID および欠落している端末 ID に対するダウンストリーム SSE 修復はデフォルトで無効になっています。関数呼び出し ID は決して書き換えられません。 | +| `responsesSnapshotRepair?` | `boolean` | デフォルトで無効のクライアント向け修復です。SSE と JSON の Responses ライフサイクルで欠落した status、output、ツールメタデータを補完し、raw 検査と永続化は変更しません。 | | `retryOn429?` | `{ enabled?: boolean; attempts?: number; intervalMs?: number; maxIntervalMs?: number; respectRetryAfter?: boolean }` | API-key プロバイダーのみ(`authMode: "key"`)。オプトインの同一ターゲット 429 リトライ: `retryOn429` が無ければ無効で、オブジェクトがあれば `enabled: false` でない限り有効になります。429 時に待機(上流の `Retry-After` または固定間隔)してから、キー フェイルオーバーの前に同一キーで同一リクエストを再送します — メインのテキストターン回復ループ、Responses passthrough、画像/動画ブリッジ、web-search サイドカー、ターミナル継続要求をすべてカバーします。再送の対象はプリストリームの HTTP 429 応答のみで、カスタム `runTurn` トランスポートは HTTP リトライループの対象外です。`attempts` は最初の 429 以降の同一キー再送回数(合計送信数 = `attempts` + 1)で、メインの回復ループ・ターミナルガード継続・ブリッジ再試行で共有されるリクエスト単位の予算です。`attempts` を使い切っても同一キーでの再送が止まるだけで、通常のキー フェイルオーバーまたは最終エラー処理が利用可能なターゲットに応じて続きます — キー認証の passthrough ワイヤにはフェイルオーバーがないため、使い切った 429 はそのまま返ります。Codex 自体は 429 をリトライしないため、単一キーのプロバイダーでは唯一の防御です。デフォルト: `enabled: true`、`attempts: 3`、`intervalMs: 5000`、`maxIntervalMs: 60000`(1回の待機は `maxIntervalMs` で上限、その上限は 600000)、`respectRetryAfter: true`。 | | `autoToolChoiceOnlyModels?` | `string[]` | `tool_choice` が `auto` または `none` のみを受け入れるモデル。強制的な選択は格下げされます。 | | `preserveReasoningContentModels?` | `string[]` |チャット履歴に以前のアシスタント `reasoning_content` が必要なモデル。 | diff --git a/docs-site/src/content/docs/ko/reference/configuration/providers.md b/docs-site/src/content/docs/ko/reference/configuration/providers.md index f321e7307..964ca558a 100644 --- a/docs-site/src/content/docs/ko/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ko/reference/configuration/providers.md @@ -86,6 +86,7 @@ target도 selector로 재사용할 수 없습니다. raw account id와 email은 | `noPenaltyModels?` | `string[]` | presence/frequency penalty를 허용하지 않는 모델입니다. | | `parallelToolCalls?` | `boolean` | 병렬 도구 호출을 켜거나 끕니다. OpenAI Chat은 기본으로 켜져 있고, 비-chat 어댑터는 명시적으로 `true`일 때만 이를 노출합니다. | | `responsesItemIdRepair?` | `{ message?: string[]; reasoning?: string[]; repairMissingTerminalIds?: boolean }` | 기본값이 꺼진 downstream SSE 복구입니다. 정확한 자리표시자 id와 누락된 종료 id를 복구합니다. function-call id는 다시 쓰지 않습니다. | +| `responsesSnapshotRepair?` | `boolean` | 기본값이 꺼진 클라이언트용 복구입니다. SSE와 JSON의 Responses 수명 주기에서 누락된 status, output, 도구 메타데이터를 채우며 raw 검사와 영속화는 변경하지 않습니다. | | `retryOn429?` | `{ enabled?: boolean; attempts?: number; intervalMs?: number; maxIntervalMs?: number; respectRetryAfter?: boolean }` | API-key 프로바이더 전용(`authMode: "key"`). 동일 대상 429 재시도: `retryOn429`가 없으면 기능이 꺼져 있고, 객체가 있으면 `enabled: false`가 아닌 한 활성화됩니다. 429 시 대기(업스트림 `Retry-After` 또는 고정 간격) 후 키 장애 조치 전에 동일 키로 동일 요청을 재전송합니다 — 일반 텍스트 턴 복구 루프, Responses passthrough, 이미지/비디오 브리지, web-search 사이드카, 터미널 연속 요청을 모두 포함합니다. 재전송 대상은 프리스트림 HTTP 429 응답뿐이며, 커스텀 `runTurn` 전송은 HTTP 재시도 루프에서 제외됩니다. `attempts`는 첫 429 이후의 동일 키 재전송 횟수(총 전송 = `attempts` + 1)이며, 메인 복구 루프·터미널 가드 연속 요청·브리지 재시도가 공유하는 요청 단위 예산입니다. `attempts`를 모두 소진해도 동일 키 재전송만 중단되며, 이후에는 일반 키 장애 조치 또는 최종 오류 처리가 사용 가능한 대상에 따라 진행됩니다 — 키 인증 passthrough 와이어에는 장애 조치가 없으므로 소진된 429가 그대로 반환됩니다. Codex 자체는 429를 재시도하지 않으므로 단일 키 프로바이더의 유일한 방어선입니다. 기본값: `enabled: true`, `attempts: 3`, `intervalMs: 5000`, `maxIntervalMs: 60000`(단일 대기는 `maxIntervalMs`로 상한, 그 자체는 600000으로 상한), `respectRetryAfter: true`. | | `autoToolChoiceOnlyModels?` | `string[]` | `tool_choice`가 `auto` 또는 `none`만 받는 모델입니다. 강제 선택은 낮은 수준으로 바뀝니다. | | `preserveReasoningContentModels?` | `string[]` | chat 기록에서 이전 assistant `reasoning_content`가 필요한 모델입니다. | diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index 4c7545254..98bb446c9 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -93,6 +93,7 @@ differing backup and rewrites known legacy namespaced selected ids to bare ids. | `noPenaltyModels?` | `string[]` | Models that reject presence/frequency penalties. | | `parallelToolCalls?` | `boolean` | Toggle parallel tool calls. OpenAI Chat defaults on; non-chat adapters advertise only on explicit `true`. | | `responsesItemIdRepair?` | `{ message?: string[]; reasoning?: string[]; repairMissingTerminalIds?: boolean }` | Disabled-by-default downstream SSE repair for exact placeholder ids and missing terminal ids. Function-call ids are never rewritten. | +| `responsesSnapshotRepair?` | `boolean` | Disabled-by-default client-facing repair for sparse Responses lifecycle snapshots in SSE and JSON. Fills missing canonical status, output, and tool metadata while raw inspection and persistence remain unchanged. | | `retryOn429?` | `{ enabled?: boolean; attempts?: number; intervalMs?: number; maxIntervalMs?: number; respectRetryAfter?: boolean }` | API-key providers only (`authMode: "key"`). Opt-in same-target 429 retry: when `retryOn429` is absent the feature is off; object presence enables it unless `enabled: false`. On 429 the proxy waits (upstream `Retry-After` or the fixed interval) and replays the identical request on the same key before any key failover — across the main text-turn recovery loop, the Responses passthrough wire, the image/video bridge, the web-search sidecar, and terminal continuations. Only pre-stream HTTP 429 responses are eligible for replay; custom `runTurn` transports are outside the HTTP retry loop. `attempts` counts same-key replays after the first 429 (total sends = `attempts` + 1) and is one request-wide budget shared by the main recovery loop, the terminal-guard continuation, and bridge retries. Exhausting `attempts` only stops further same-key replays: normal key failover or final-error handling then applies per the available targets — on the key-auth passthrough wire there is no failover, so the exhausted 429 surfaces as-is. Codex itself never retries 429, so this is the only defense for single-key providers. Defaults: `enabled: true`, `attempts: 3`, `intervalMs: 5000`, `maxIntervalMs: 60000` (any single wait is capped at `maxIntervalMs`, itself capped at 600000), `respectRetryAfter: true`. | | `autoToolChoiceOnlyModels?` | `string[]` | Models whose `tool_choice` accepts only `auto` or `none`; forced choices are downgraded. | | `preserveReasoningContentModels?` | `string[]` | Models requiring prior assistant `reasoning_content` in chat history. | diff --git a/docs-site/src/content/docs/ru/reference/configuration/providers.md b/docs-site/src/content/docs/ru/reference/configuration/providers.md index 6f11149e5..5984c48f8 100644 --- a/docs-site/src/content/docs/ru/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ru/reference/configuration/providers.md @@ -96,6 +96,7 @@ cross-route credential fallback не существует. Строки API GPT- | `noPenaltyModels?` | `string[]` | Модели, отвергающие penalty presence/frequency. | | `parallelToolCalls?` | `boolean` | Переключатель parallel tool call'ов. Для OpenAI Chat по умолчанию включено; не-chat adapter'ы рекламируют это только при явном `true`. | | `responsesItemIdRepair?` | `{ message?: string[]; reasoning?: string[]; repairMissingTerminalIds?: boolean }` | По умолчанию выключенная downstream SSE-repair для exact placeholder-id и отсутствующих terminal-id. Function-call id никогда не переписываются. | +| `responsesSnapshotRepair?` | `boolean` | По умолчанию выключенная клиентская repair для неполных lifecycle snapshot'ов Responses в SSE и JSON. Добавляет отсутствующие status, output и tool metadata, не меняя raw inspection и persistence. | | `retryOn429?` | `{ enabled?: boolean; attempts?: number; intervalMs?: number; maxIntervalMs?: number; respectRetryAfter?: boolean }` | Только для провайдеров с API-ключом (`authMode: "key"`). Опциональный повтор при 429 на том же таргете: если `retryOn429` отсутствует, функция выключена; наличие объекта включает её, если только `enabled: false`. При 429: ожидание (`Retry-After` апстрима или фиксированный интервал) и повтор идентичного запроса на том же ключе до любого фейловера ключей — покрывает основной цикл восстановления текстовых ходов, passthrough-канал Responses, мост изображений/видео, sidecar web-search и терминальные продолжения. Повтор допустим только для HTTP 429, полученных до начала потока; пользовательские транспорты `runTurn` не входят в цикл HTTP-повторов. `attempts` — это число повторов на том же ключе после первого 429 (всего отправок = `attempts` + 1) и единый бюджет на запрос, общий для основного цикла восстановления, терминального продолжения и повторов моста. Исчерпание `attempts` лишь останавливает дальнейшие повторы на том же ключе; далее применяется обычный фейловер ключей или финальная обработка ошибки в зависимости от доступных таргетов — на passthrough-канале с ключевой аутентификацией фейловера нет, поэтому исчерпанный 429 возвращается как есть. Codex сам никогда не повторяет 429, поэтому это единственная защита для провайдеров с одним ключом. По умолчанию: `enabled: true`, `attempts: 3`, `intervalMs: 5000`, `maxIntervalMs: 60000` (любое ожидание ограничено `maxIntervalMs`, который сам ограничен 600000), `respectRetryAfter: true`. | | `autoToolChoiceOnlyModels?` | `string[]` | Модели, у которых `tool_choice` принимает только `auto` или `none`; forced choice понижается. | | `preserveReasoningContentModels?` | `string[]` | Модели, которым нужен предыдущий assistant `reasoning_content` в chat history. | diff --git a/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md b/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md index 36dd4e288..56befac3d 100644 --- a/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md +++ b/docs-site/src/content/docs/zh-cn/reference/configuration/providers.md @@ -85,6 +85,7 @@ pool account id(不能是内部 `__main__`),或用 `"@main"` 表示 Codex | `noPenaltyModels?` | `string[]` | 会拒绝 presence/frequency penalty 的模型。 | | `parallelToolCalls?` | `boolean` | 切换并行工具调用。OpenAI Chat 默认开启;非 chat 适配器只有显式 `true` 时才会声明支持。 | | `responsesItemIdRepair?` | `{ message?: string[]; reasoning?: string[]; repairMissingTerminalIds?: boolean }` | 默认关闭的下游 SSE 修复,用于精确占位 id 和缺失的终止 id。function-call id 永远不会被重写。 | +| `responsesSnapshotRepair?` | `boolean` | 默认关闭的客户端修复,用于补全 SSE 与 JSON 中稀疏 Responses 生命周期快照缺失的 status、output 和工具元数据;原始检查与持久化保持不变。 | | `retryOn429?` | `{ enabled?: boolean; attempts?: number; intervalMs?: number; maxIntervalMs?: number; respectRetryAfter?: boolean }` | 仅限 API-key 提供商(`authMode: "key"`)。可选的同目标 429 重试:未配置 `retryOn429` 时功能关闭;对象存在即启用,除非 `enabled: false`。收到 429 时等待(上游 `Retry-After` 或固定间隔)后在相同 key 上重放完全相同请求,再进入任何 key 故障转移——覆盖主文本恢复循环、Responses passthrough、图像/视频桥、web-search 侧车与终结续接。重放仅适用于流开始前的 HTTP 429 响应;自定义 `runTurn` 传输不在 HTTP 重试循环范围内。`attempts` 是首个 429 之后的同 key 重放次数(总发送次数 = `attempts` + 1),是主恢复循环、终结守卫续接与桥接重试共享的按请求统一预算;`attempts` 耗尽只会停止进一步的同 key 重放:随后按可用目标进行正常的 key 故障转移或最终错误处理——key 认证的 passthrough 线路上没有故障转移,因此耗尽的 429 会原样透出。Codex 自身从不重试 429,因此这是单 key 提供商唯一的防线。默认值:`enabled: true`、`attempts: 3`、`intervalMs: 5000`、`maxIntervalMs: 60000`(单次等待以 `maxIntervalMs` 为上限,其本身上限 600000)、`respectRetryAfter: true`。 | | `autoToolChoiceOnlyModels?` | `string[]` | `tool_choice` 只接受 `auto` 或 `none` 的模型;强制选择会被降级。 | | `preserveReasoningContentModels?` | `string[]` | 需要在聊天历史中保留先前 assistant `reasoning_content` 的模型。 | diff --git a/src/config.ts b/src/config.ts index 72defeb9c..08f803108 100644 --- a/src/config.ts +++ b/src/config.ts @@ -598,6 +598,7 @@ const providerConfigSchema = z.object({ reasoning: z.array(z.string().min(1)).optional(), repairMissingTerminalIds: z.boolean().optional(), }).strict().optional(), + responsesSnapshotRepair: z.boolean().optional(), }).passthrough(); const RESERVED_PROVIDER_NAMES = new Set(["__proto__", "prototype", "constructor"]); diff --git a/src/server/auth-cors.ts b/src/server/auth-cors.ts index 2acea2848..6de271da2 100644 --- a/src/server/auth-cors.ts +++ b/src/server/auth-cors.ts @@ -413,7 +413,9 @@ export function providerManagementConfigError(name: unknown, provider: unknown): return "provider openai codexAccountMode must be pool or direct"; } if (seed) seed.codexAccountMode = raw.codexAccountMode; - const canonical = seed && sameCanonicalProviderSeed(raw, seed); + const canonicalCandidate = { ...raw }; + delete canonicalCandidate.responsesSnapshotRepair; + const canonical = seed && sameCanonicalProviderSeed(canonicalCandidate, seed); if (!canonical) { return `provider ${name} must equal the canonical built-in provider seed`; } @@ -453,6 +455,9 @@ export function providerManagementConfigError(name: unknown, provider: unknown): typed, ); if (preferHostedToolsError) return `provider ${name} ${preferHostedToolsError}`; + if (raw.responsesSnapshotRepair !== undefined && typeof raw.responsesSnapshotRepair !== "boolean") { + return `provider ${name} responsesSnapshotRepair must be a boolean`; + } const defaultMaxOutputError = positiveIntegerConfigError(raw.defaultMaxOutputTokens, "defaultMaxOutputTokens"); if (defaultMaxOutputError) return `provider ${name} ${defaultMaxOutputError}`; const maxOutputError = positiveIntegerRecordConfigError(raw.modelMaxOutputTokens, "modelMaxOutputTokens"); diff --git a/src/server/relay-eager.ts b/src/server/relay-eager.ts index 5c1241f22..8865a9923 100644 --- a/src/server/relay-eager.ts +++ b/src/server/relay-eager.ts @@ -27,8 +27,10 @@ import { buildFailedTailPayload, createSseTerminalOutputBoundary } from "./relay"; import { nextSseBlock, + payloadRewriteAsBlockRewrite, replaceSseDataPayload, sseDataPayload, + type SseBlockRewrite, type SsePayloadRewrite, } from "./sse-payload-rewrite"; import type { TranslatorBudget } from "../lib/translator-budget"; @@ -43,6 +45,12 @@ export type EagerRelayHooks = { * Bun#32111-unsafe tee()+JS-pull chain (#864). */ rewritePayload?: SsePayloadRewrite; + /** + * Optional block-level rewrite (zero or more blocks out per upstream + * block) for lifecycle event injection (#893). Takes precedence over + * rewritePayload when both are set. + */ + rewriteBlocks?: SseBlockRewrite; /** Flush inspection at upstream end (createSseInspector.finish). */ finishInspection: () => void; /** Drop inspector-owned frame/item state during producer teardown. */ @@ -93,9 +101,10 @@ export function relaySseEagerBounded( const reader = body.getReader(); const terminalBoundary = createSseTerminalOutputBoundary(); - const rewrite = hooks.rewritePayload; - const rewriteDecoder = rewrite ? new TextDecoder() : null; - const rewriteEncoder = rewrite ? new TextEncoder() : null; + const activeRewrite: SseBlockRewrite | undefined = hooks.rewriteBlocks + ?? (hooks.rewritePayload ? payloadRewriteAsBlockRewrite(hooks.rewritePayload) : undefined); + const rewriteDecoder = activeRewrite ? new TextDecoder() : null; + const rewriteEncoder = activeRewrite ? new TextEncoder() : null; const rewriteBudget = opts?.rewriteBudget; let frameBuffer = ""; let frameBufferBytes = 0; @@ -122,15 +131,9 @@ export function relaySseEagerBounded( for (;;) { const next = nextSseBlock(frameBuffer); if (!next) break; - const payload = sseDataPayload(next.block); - const rewrittenPayload = payload === null ? null : rewrite!(payload); - // Replace only on an actual change: replaceSseDataPayload collapses - // multi-data-line events and normalizes newline style even when the - // payload is identical, which corrupts valid streams. - const block = payload !== null && rewrittenPayload !== payload - ? replaceSseDataPayload(next.block, rewrittenPayload!) - : next.block; - out += block + next.delimiter; + for (const outBlock of activeRewrite!(next.block)) { + out += outBlock + next.delimiter; + } frameBuffer = next.rest; } if (rewriteBudget) { @@ -144,14 +147,13 @@ export function relaySseEagerBounded( }; /** Flush any trailing partial block at upstream end (rewrite applied, matching the pull relay). */ const flushRewriteTail = (): Uint8Array => { - if (!rewrite) return new Uint8Array(0); + if (!activeRewrite) return new Uint8Array(0); // Decoder-flushed bytes logically follow everything already decoded. let tail = frameBuffer + rewriteDecoder!.decode(); - const payload = sseDataPayload(tail); - if (payload !== null) { - const rewrittenPayload = rewrite(payload); - if (rewrittenPayload !== payload) tail = replaceSseDataPayload(tail, rewrittenPayload); - } + const rewritten = activeRewrite(tail); + // Multiple emitted blocks must stay separately framed (#893 review); + // join places the delimiter only between blocks, never after the last. + tail = rewritten.join(tail.includes("\r\n") ? "\r\n\r\n" : "\n\n"); frameBuffer = ""; if (rewriteBudget && frameBufferBytes > 0) { rewriteBudget.releaseRetained(frameBufferBytes, { kind: "live_transient" }); @@ -219,7 +221,7 @@ export function relaySseEagerBounded( if (upstreamDone) { hooks.finishInspection(); const boundedTail = terminalBoundary.finish(); - if (rewrite) { + if (activeRewrite) { const rewritten = rewriteOutbound(boundedTail); const tail = joinUint8Arrays(rewritten, flushRewriteTail()); if (tail.byteLength > 0 && !cancelled) { @@ -245,7 +247,7 @@ export function relaySseEagerBounded( continue; } const terminalBounded = terminalBoundary.feed(value); - const outbound = rewrite ? rewriteOutbound(terminalBounded) : terminalBounded; + const outbound = activeRewrite ? rewriteOutbound(terminalBounded) : terminalBounded; if (outbound.byteLength > 0) { queuedBytes += outbound.byteLength; try { @@ -312,6 +314,7 @@ export function relaySseEagerBounded( try { controllerRef?.close(); } catch { /* already closed/errored */ } } try { hooks.disposeInspection?.(); } catch { /* inspection teardown must not block lifecycle cleanup */ } + try { activeRewrite?.dispose?.(); } catch { /* rewrite teardown must not block lifecycle cleanup */ } fireDone(); } }; diff --git a/src/server/responses-snapshot-repair.ts b/src/server/responses-snapshot-repair.ts new file mode 100644 index 000000000..6818e6492 --- /dev/null +++ b/src/server/responses-snapshot-repair.ts @@ -0,0 +1,601 @@ +/** + * Provider-opt-in repair for Responses-compatible gateways that return sparse + * lifecycle snapshots (#893): field backfills for missing canonical fields, + * AND lifecycle event injection so Codex clients actually commit the turn. + * + * Field-repair semantics (snapshot fields, item/part backfills, retention + * bounds) are adopted from PR #928 (0xWinner98) with attribution. The + * lifecycle-completion layer is new: #928's exact-issue test expects no + * `output_item.done` and a terminal `output: []` — normalization without + * commitment, which the canonical bridge contract (src/bridge.ts) shows is + * not enough for Codex to commit the message. + * + * Design: + * - One stateful tracker per stream. Existing upstream values are always + * authoritative; only absent or structurally invalid fields are backfilled. + * - Event injection happens only for items PROVEN open at a terminal event + * (never for content the gateway closed itself). + * - Ambiguous, gapped, malformed, oversized, or contradictory shapes taint + * the tracker: no injections and no output reconstruction afterwards + * (fail closed; explicit `output: []` from the gateway stays authoritative). + */ + +import type { TranslatorBudget } from "../lib/translator-budget"; +import { + MAX_COMPLETED_OUTPUT_ITEMS, + MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES, +} from "./relay"; +import { sseDataPayload, type SseBlockRewrite } from "./sse-payload-rewrite"; + +const RESPONSE_EVENT_STATUSES: Readonly> = { + "response.created": "in_progress", + "response.in_progress": "in_progress", + "response.completed": "completed", + "response.failed": "failed", + "response.incomplete": "incomplete", + "response.queued": "queued", +}; + +type RequestDefaults = { + parallelToolCalls: boolean; + toolChoice: unknown; + tools: unknown[]; +}; + +function isPlainObject(value: unknown): value is Record { + return !!value && typeof value === "object" && !Array.isArray(value); +} + +function isStructurallyValidToolChoice(value: unknown): boolean { + return (typeof value === "string" && value.trim().length > 0) + || (isPlainObject(value) && typeof value.type === "string" && value.type.trim().length > 0); +} + +function requestDefaults(requestBody: unknown): RequestDefaults { + const request = isPlainObject(requestBody) ? requestBody : {}; + return { + parallelToolCalls: typeof request.parallel_tool_calls === "boolean" + ? request.parallel_tool_calls + : true, + toolChoice: isStructurallyValidToolChoice(request.tool_choice) ? request.tool_choice : "auto", + tools: Array.isArray(request.tools) ? request.tools : [], + }; +} + +function repairOutputTextPart(part: Record): Record { + if (part.type !== "output_text") return part; + const needsText = typeof part.text !== "string"; + const needsAnnotations = !Array.isArray(part.annotations); + if (!needsText && !needsAnnotations) return part; + return { + ...part, + ...(needsText ? { text: "" } : {}), + ...(needsAnnotations ? { annotations: [] } : {}), + }; +} + +function repairSummaryPart(part: Record): Record { + if (part.type !== "summary_text" || typeof part.text === "string") return part; + return { ...part, text: "" }; +} + +function repairOutputItem( + item: Record, + inferredStatus?: string, +): Record { + let repaired = item; + let changed = false; + + if (item.type === "reasoning") { + const rawSummary = item.summary; + const summary = Array.isArray(rawSummary) + ? rawSummary.map((part) => isPlainObject(part) ? repairSummaryPart(part) : part) + : []; + changed = !Array.isArray(rawSummary) + || summary.some((part, index) => part !== rawSummary[index]); + if (changed) repaired = { ...repaired, summary }; + } else if (item.type === "message") { + const rawContent = item.content; + const content = Array.isArray(rawContent) + ? rawContent.map((part) => isPlainObject(part) ? repairOutputTextPart(part) : part) + : []; + changed = !Array.isArray(rawContent) + || content.some((part, index) => part !== rawContent[index]); + // Responses output-message roles are the literal "assistant"; input roles are invalid here. + changed = changed || item.role !== "assistant"; + if (changed) repaired = { ...repaired, content, role: "assistant" }; + } + + if (inferredStatus && (typeof repaired.status !== "string" || repaired.status.trim().length === 0)) { + repaired = { ...repaired, status: inferredStatus }; + } + return repaired; +} + +function repairResponseSnapshot( + response: Record, + defaultStatus: string, + defaults: RequestDefaults, + reconstructedOutput?: Record[], +): Record { + const repaired = { ...response }; + let changed = false; + const effectiveResponseStatus = typeof response.status === "string" && response.status.trim().length > 0 + ? response.status + : defaultStatus; + const outputStatus = effectiveResponseStatus === "completed" || effectiveResponseStatus === "incomplete" + ? effectiveResponseStatus + : undefined; + + if (reconstructedOutput) { + repaired.output = reconstructedOutput; + changed = true; + } else if (Array.isArray(repaired.output)) { + const output = repaired.output.map((item) => { + if (!isPlainObject(item)) return item; + const next = repairOutputItem(item, outputStatus); + changed = changed || next !== item; + return next; + }); + if (changed) repaired.output = output; + } + if (typeof repaired.parallel_tool_calls !== "boolean") { + repaired.parallel_tool_calls = defaults.parallelToolCalls; + changed = true; + } + if (!isStructurallyValidToolChoice(repaired.tool_choice)) { + repaired.tool_choice = defaults.toolChoice; + changed = true; + } + if (!Array.isArray(repaired.tools)) { + repaired.tools = defaults.tools; + changed = true; + } + if (typeof repaired.status !== "string" || repaired.status.trim().length === 0) { + repaired.status = defaultStatus; + changed = true; + } + + return changed ? repaired : response; +} + +type OpenItem = { + itemId: string; + outputIndex: number; + type: string; + /** Message/reasoning items get lifecycle injections; other types never do. */ + injectable: boolean; + /** The gateway opened a content part explicitly (or we injected one). */ + contentPartOpen: boolean; + /** The gateway sent output_text.done for this item. */ + textDone: boolean; + /** The gateway sent content_part.done for this item. */ + partDone: boolean; + text: string; + item: Record; +}; + +/** Open-item retention bounds (#893 review): aggregate, not just per-item. */ +const MAX_OPEN_ITEMS = MAX_COMPLETED_OUTPUT_ITEMS; +const MAX_OPEN_ITEM_AGGREGATE_TEXT_BYTES = MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES; + +type RetainedOutputItem = { + item: Record; + sourceBytes: number; +}; + +function jsonBlock(event: Record): string { + return `data: ${JSON.stringify(event)}`; +} + +/** + * Stateful block rewrite: field backfills + lifecycle completion injection. + * `budget` bounds retained completed items (reconstruction only). + */ +export function createResponsesSnapshotBlockRewrite( + requestBody?: unknown, + budget?: TranslatorBudget, +): SseBlockRewrite { + const defaults = requestDefaults(requestBody); + const openItems = new Map(); + const completedItems = new Map(); + let aggregateItemBytes = 0; + let aggregateOpenTextBytes = 0; + let tainted = false; + + const releaseRetained = (): void => { + if (aggregateItemBytes > 0) { + budget?.releaseRetained(aggregateItemBytes, { kind: "retained_collectors" }); + } + completedItems.clear(); + openItems.clear(); + aggregateItemBytes = 0; + aggregateOpenTextBytes = 0; + tainted = false; + }; + + const taintAndRelease = (): void => { + // Fail closed means stop RETAINING too: open state and charged collectors + // are released immediately, and no further state accumulates (#893 review). + if (aggregateItemBytes > 0) { + budget?.releaseRetained(aggregateItemBytes, { kind: "retained_collectors" }); + } + completedItems.clear(); + openItems.clear(); + aggregateItemBytes = 0; + aggregateOpenTextBytes = 0; + tainted = true; + }; + + /** Drop one open item and refund its accumulated text from the aggregate. */ + const closeOpenItem = (index: number): void => { + const open = openItems.get(index); + if (!open) return; + aggregateOpenTextBytes -= Buffer.byteLength(open.text, "utf8"); + openItems.delete(index); + }; + + const retainCompletedItem = (index: number, item: Record): void => { + if (tainted) return; // fail closed means stop retaining (#893 review) + const sourceBytes = Buffer.byteLength(JSON.stringify(item), "utf8"); + const previous = completedItems.get(index); + if (sourceBytes > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES) { + if (previous) { + completedItems.delete(index); + aggregateItemBytes -= previous.sourceBytes; + budget?.releaseRetained(previous.sourceBytes, { kind: "retained_collectors" }); + } + taintAndRelease(); + return; + } + const retainedDelta = sourceBytes - (previous?.sourceBytes ?? 0); + if (retainedDelta > 0) { + budget?.chargeRetained(retainedDelta, { kind: "retained_collectors" }); + } else if (retainedDelta < 0) { + budget?.releaseRetained(-retainedDelta, { kind: "retained_collectors" }); + } + completedItems.set(index, { item, sourceBytes }); + aggregateItemBytes += retainedDelta; + while (completedItems.size > MAX_COMPLETED_OUTPUT_ITEMS + || aggregateItemBytes > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES) { + let highestIndex = -1; + for (const retainedIndex of completedItems.keys()) { + if (retainedIndex > highestIndex) highestIndex = retainedIndex; + } + const evicted = completedItems.get(highestIndex); + if (!evicted) break; + completedItems.delete(highestIndex); + aggregateItemBytes -= evicted.sourceBytes; + budget?.releaseRetained(evicted.sourceBytes, { kind: "retained_collectors" }); + taintAndRelease(); + return; + } + }; + + const completedItemSnapshot = (open: OpenItem): Record => { + if (open.type === "message") { + return repairOutputItem({ + ...open.item, + status: "completed", + role: "assistant", + content: [{ type: "output_text", text: open.text, annotations: [] }], + }, "completed"); + } + return repairOutputItem({ ...open.item, status: "completed" }, "completed"); + }; + + const rewrite: SseBlockRewrite = (block: string): readonly string[] => { + const payload = sseDataPayload(block); + if (payload === null) return [block]; + let event: unknown; + try { + event = JSON.parse(payload); + } catch { + return [block]; + } + if (!isPlainObject(event)) return [block]; + const type = typeof event.type === "string" ? event.type : ""; + let nextEvent: Record = event; + let changed = false; + const outputIndex = Number.isInteger(event.output_index) && (event.output_index as number) >= 0 + ? event.output_index as number + : undefined; + const itemId = typeof event.item_id === "string" ? event.item_id : undefined; + + // --- item lifecycle tracking + field backfills ------------------------- + if (type === "response.output_item.added" && isPlainObject(event.item)) { + const itemType = typeof event.item.type === "string" ? event.item.type : ""; + const id = typeof event.item.id === "string" ? event.item.id : undefined; + if (outputIndex === undefined || !id) { + taintAndRelease(); + } else { + if (openItems.has(outputIndex)) taintAndRelease(); // contradictory reuse of an open index + if (openItems.size >= MAX_OPEN_ITEMS) taintAndRelease(); // aggregate retention bound + // Unsupported types (function_call, …) reserve their index but are + // never injected — a sparse gateway that leaves one open blocks only + // the terminal reconstruction, not message repair (#893 review). + const injectable = itemType === "message" || itemType === "reasoning"; + if (!tainted) { + openItems.set(outputIndex, { + itemId: id, + outputIndex, + type: itemType, + injectable, + contentPartOpen: false, + textDone: false, + partDone: false, + text: "", + item: event.item, + }); + } + } + const item = repairOutputItem(event.item, "in_progress"); + if (item !== event.item) { + nextEvent = { ...nextEvent, item }; + changed = true; + } + } + + if (type === "response.output_item.done" && !isPlainObject(event.item)) { + taintAndRelease(); + } + + if (type === "response.output_item.done" && isPlainObject(event.item)) { + const item = repairOutputItem(event.item, "completed"); + if (item !== event.item) { + nextEvent = { ...nextEvent, item }; + changed = true; + } + // Identity correlation: a done event only closes the tracked item when + // id AND type agree — a mismatched done is a contradictory stream and + // must go fail-closed rather than close/reconstruct the wrong item. + const tracked = outputIndex !== undefined ? openItems.get(outputIndex) : undefined; + const doneId = typeof event.item.id === "string" ? event.item.id : undefined; + const doneType = typeof event.item.type === "string" ? event.item.type : undefined; + if (tracked && (doneId !== tracked.itemId || doneType !== tracked.type)) { + taintAndRelease(); + return [changed ? jsonBlock(nextEvent) : block]; + } + if (outputIndex !== undefined && typeof item.type === "string" && item.type.trim().length > 0) { + closeOpenItem(outputIndex); + retainCompletedItem(outputIndex, item); + } else { + taintAndRelease(); + } + } + + if ((type === "response.content_part.added" || type === "response.content_part.done") + && isPlainObject(event.part)) { + if (outputIndex !== undefined) { + const open = openItems.get(outputIndex); + // Correlate by item_id when present: a mismatched event must not + // mutate (or suppress injections for) the tracked item (#893 review). + if (open && (itemId === undefined || itemId === open.itemId)) { + if (type === "response.content_part.added") open.contentPartOpen = true; + else open.partDone = true; + } + } + const part = repairOutputTextPart(event.part); + if (part !== event.part) { + nextEvent = { ...nextEvent, part }; + changed = true; + } + } + + if ((type === "response.reasoning_summary_part.added" + || type === "response.reasoning_summary_part.done") && isPlainObject(event.part)) { + const part = repairSummaryPart(event.part); + if (part !== event.part) { + nextEvent = { ...nextEvent, part }; + changed = true; + } + } + + if ((type === "response.output_text.delta" || type === "response.output_text.done") + && !Array.isArray(event.logprobs)) { + nextEvent = { ...nextEvent, logprobs: [] }; + changed = true; + } + if (type === "response.output_text.delta" && typeof event.delta === "string" && outputIndex !== undefined) { + const open = openItems.get(outputIndex); + if (open && (itemId === undefined || itemId === open.itemId)) { + const deltaBytes = Buffer.byteLength(event.delta, "utf8"); + open.text += event.delta; + aggregateOpenTextBytes += deltaBytes; + // Unbounded text accumulation is a retention hole (#893 review): + // overshoot per-item or aggregate caps and the stream goes fail-closed. + if (Buffer.byteLength(open.text, "utf8") > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES + || aggregateOpenTextBytes > MAX_OPEN_ITEM_AGGREGATE_TEXT_BYTES) { + taintAndRelease(); + } + } + } + if (type === "response.output_text.done") { + if (outputIndex !== undefined) { + const open = openItems.get(outputIndex); + // Identity correlation, same contract as output_item.done above: a + // text-done carrying a DIFFERENT item_id for a tracked index is a + // contradictory stream. Silently ignoring it used to leave the item + // open, so the terminal then synthesized a second output_text.done and + // output_item.done — a double close built on a stream we do not + // understand (#1025 review blocker 2). Go fail-closed instead: forward + // the block untouched and stop injecting for the rest of the stream. + if (open && itemId !== undefined && itemId !== open.itemId) { + taintAndRelease(); + return [changed ? jsonBlock(nextEvent) : block]; + } + if (open && (itemId === undefined || itemId === open.itemId)) { + open.textDone = true; + if (typeof event.text === "string") { + // Replacement must keep the aggregate honest and bounded, the + // same as delta accumulation (#893 review). + const previousBytes = Buffer.byteLength(open.text, "utf8"); + const nextBytes = Buffer.byteLength(event.text, "utf8"); + open.text = event.text; + aggregateOpenTextBytes += nextBytes - previousBytes; + if (nextBytes > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES + || aggregateOpenTextBytes > MAX_OPEN_ITEM_AGGREGATE_TEXT_BYTES) { + taintAndRelease(); + } + } + } + } + if (typeof event.text !== "string") { + nextEvent = { ...nextEvent, text: "" }; + changed = true; + } + } + + // --- lifecycle response snapshots -------------------------------------- + const responseStatus = Object.prototype.hasOwnProperty.call(RESPONSE_EVENT_STATUSES, type) + ? RESPONSE_EVENT_STATUSES[type] + : undefined; + const isTerminal = type === "response.completed" + || type === "response.failed" + || type === "response.incomplete"; + + if (responseStatus && isPlainObject(event.response)) { + // Terminal output reconstruction stays in the injection layer below; + // here the snapshot only gets field backfills. + const response = repairResponseSnapshot(event.response, responseStatus, defaults); + if (response !== event.response) { + nextEvent = { ...nextEvent, response }; + changed = true; + } + } + + // --- lifecycle completion injection ------------------------------------ + // Only response.completed justifies synthesized closing events: a failed + // or incomplete terminal must never fabricate completed output (#893 review). + const isCompletedTerminal = type === "response.completed"; + const out: string[] = []; + if (!tainted) { + // A text delta for an item whose content part was never opened needs the + // opening event first, or Codex never binds the text to a committed part. + if (type === "response.output_text.delta" && outputIndex !== undefined) { + const open = openItems.get(outputIndex); + if (open && open.type === "message" && !open.contentPartOpen) { + out.push(jsonBlock({ + type: "response.content_part.added", + item_id: open.itemId, + output_index: open.outputIndex, + content_index: 0, + part: { type: "output_text", text: "", annotations: [] }, + })); + open.contentPartOpen = true; + } + } + + if (isCompletedTerminal && openItems.size > 0) { + const closing = [...openItems.values()].sort((a, b) => a.outputIndex - b.outputIndex); + for (const open of closing) { + if (open.injectable && open.type === "message") { + if (!open.contentPartOpen) { + out.push(jsonBlock({ + type: "response.content_part.added", + item_id: open.itemId, + output_index: open.outputIndex, + content_index: 0, + part: { type: "output_text", text: "", annotations: [] }, + })); + open.contentPartOpen = true; + } + if (!open.textDone) { + out.push(jsonBlock({ + type: "response.output_text.done", + item_id: open.itemId, + output_index: open.outputIndex, + content_index: 0, + logprobs: [], + text: open.text, + })); + open.textDone = true; + } + if (!open.partDone) { + out.push(jsonBlock({ + type: "response.content_part.done", + item_id: open.itemId, + output_index: open.outputIndex, + content_index: 0, + part: { type: "output_text", text: open.text, annotations: [] }, + })); + open.partDone = true; + } + } + if (!open.injectable) continue; // never fabricate completions for other types + const snapshot = completedItemSnapshot(open); + out.push(jsonBlock({ + type: "response.output_item.done", + output_index: open.outputIndex, + item: snapshot, + })); + closeOpenItem(open.outputIndex); + retainCompletedItem(open.outputIndex, snapshot); + } + } + } + + // Terminal output reconstruction/canonical backfill: an explicit + // `output: []` is authoritative; absent or malformed output on a + // completed terminal becomes the retained items or the canonical empty + // list (#893 review). Leftover non-injectable open items block + // reconstruction only, not the message closing above. + if (isCompletedTerminal && !tainted && isPlainObject(nextEvent.response)) { + const response = nextEvent.response; + const outputValue = (response as Record).output; + const outputExplicitArray = Array.isArray(outputValue); + if (!outputExplicitArray) { + const injectableOpen = [...openItems.values()].filter(open => open.injectable); + const nonInjectableOpen = [...openItems.values()].filter(open => !open.injectable); + if (completedItems.size > 0 && injectableOpen.length === 0 && nonInjectableOpen.length === 0) { + const ordered = [...completedItems.entries()].sort(([left], [right]) => left - right); + if (ordered.every(([index], position) => index === position)) { + nextEvent = { + ...nextEvent, + response: { ...response, output: ordered.map(([, retained]) => retained.item) }, + }; + changed = true; + } else { + // A gap means at least one completed item is missing. Never compact + // later indexes into a shorter array that only appears complete. + tainted = true; + } + } else if (completedItems.size === 0 && openItems.size === 0) { + nextEvent = { ...nextEvent, response: { ...response, output: [] } }; + changed = true; + } + } + } + + out.push(changed ? jsonBlock(nextEvent) : block); + if (isTerminal) releaseRetained(); + return out; + }; + // Relays call this on every teardown path (terminal, EOF, cancel, error): + // retained collectors and open-item state never outlive the stream. + rewrite.dispose = releaseRetained; + return rewrite; +} + +/** Repair a non-streaming Responses JSON object without changing raw inspection state. */ +export function repairResponsesSnapshotJson(payload: string, requestBody?: unknown): string { + let response: unknown; + try { + response = JSON.parse(payload); + } catch { + return payload; + } + if (!isPlainObject(response)) return payload; + // Canonical output for non-stream completed responses: absent or malformed + // output becomes the canonical empty list; explicit arrays stay authoritative. + if (!Array.isArray(response.output)) { + const repaired = repairResponseSnapshot({ ...response, output: [] }, "completed", requestDefaults(requestBody)); + return JSON.stringify(repaired); + } + const repaired = repairResponseSnapshot(response, "completed", requestDefaults(requestBody)); + return repaired === response ? payload : JSON.stringify(repaired); +} + +export function hasResponsesSnapshotRepair(enabled: boolean | undefined): enabled is true { + return enabled === true; +} diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index dd0f24e60..383c24b50 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -177,6 +177,16 @@ import { hasUnreadableEncryptedAgentTask, looksLikeBackendCiphertext, sanitizeEn import { fetchWithHeaderTimeout, providerFetch, safeHostLabel, safeOriginLabel } from "./fetch-helpers"; import { classifyTransportFailureKind, transportErrorCode } from "../../lib/upstream-reachability"; import { recordUpstreamHostFailure, resetUpstreamHostHealth, upstreamHostHealthKey } from "../../codex/upstream-host-health"; +import { + createResponsesSnapshotBlockRewrite, + hasResponsesSnapshotRepair, + repairResponsesSnapshotJson, +} from "../responses-snapshot-repair"; +import { + composeSseBlockRewrites, + payloadRewriteAsBlockRewrite, + relaySseWithBlockRewrite, +} from "../sse-payload-rewrite"; import { guardTerminalEventStream } from "./terminal-guard"; /** @@ -2013,7 +2023,8 @@ async function handleResponsesInner( // The bundled known-bad runtime remains on tee by default on both platforms. if (isEventStream && upstreamResponse.body) { const repairConfig = route.provider.responsesItemIdRepair; - const needsClientRewrite = imageGenCallAliases.size > 0 || hasResponsesItemIdRepair(repairConfig); + const snapshotRepairEnabled = hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair); + const needsClientRewrite = imageGenCallAliases.size > 0 || hasResponsesItemIdRepair(repairConfig) || snapshotRepairEnabled; // Compose opt-in payload rewrites into one parse/stringify pass (image-gen restore first). const payloadRewrites = [ createImageGenCallRestoreRewrite(imageGenCallAliases), @@ -2021,6 +2032,23 @@ async function handleResponsesInner( ? createResponsesItemIdPayloadRewrite(repairConfig!, translatorBudget) : undefined, ].filter((rewrite): rewrite is NonNullable => rewrite !== undefined); + // #893: sparse-snapshot gateways get field backfills AND lifecycle event + // injection at the block level, after payload rewrites. Defaults come + // from the finalized OUTBOUND body — the normalized internal tool shapes + // are not the Responses wire shapes the snapshot must mirror. + const snapshotDefaultsRequest = (() => { + try { + return JSON.parse(request.body) as unknown; + } catch { + return undefined; + } + })(); + const clientBlockRewrite = snapshotRepairEnabled + ? composeSseBlockRewrites( + payloadRewriteAsBlockRewrite(composeSsePayloadRewrites(...payloadRewrites)), + createResponsesSnapshotBlockRewrite(snapshotDefaultsRequest, translatorBudget), + ) + : undefined; // #864: win32 rewrite traffic must never enter the tee()+JS-pull chain // (Bun#32111 JS-sink segfault — text frames pass, the terminal block is // lost). The eager single reader applies the same rewrites inline. @@ -2072,6 +2100,9 @@ async function handleResponsesInner( ...(win32EagerRewrite ? { rewritePayload: composeSsePayloadRewrites(...payloadRewrites) } : {}), + ...(clientBlockRewrite + ? { rewriteBlocks: clientBlockRewrite } + : {}), onSynthetic: kind => { if (!reportNativeTerminal) return; if (kind === "incomplete") { @@ -2159,8 +2190,8 @@ async function handleResponsesInner( // Windows was handled by the eager terminal-aware branch above. Remaining // tee traffic can use the JS relay to close on a protocol terminal and to // convert a mid-stream reset into a clean response.failed event. - const rewrittenBody = payloadRewrites.length > 0 - ? relaySseWithPayloadRewrite(nativeBody, composeSsePayloadRewrites(...payloadRewrites), translatorBudget) + const rewrittenBody = clientBlockRewrite !== undefined || payloadRewrites.length > 0 + ? relaySseWithBlockRewrite(nativeBody, clientBlockRewrite ?? payloadRewriteAsBlockRewrite(composeSsePayloadRewrites(...payloadRewrites)), translatorBudget) : nativeBody; const clientBody = relaySseWithFailedTail(rewrittenBody, upstream, reason => clientGone.abort(reason)); return markNativePassthroughSseResponse(new Response(clientBody, { @@ -2193,7 +2224,18 @@ async function handleResponsesInner( rememberPassthroughResponse(JSON.parse(text) as { id?: unknown; output?: unknown; status?: unknown }); } catch { /* non-JSON despite content-type; recording is best-effort */ } } - return new Response(restoreImageGenCallsInJson(text, imageGenCallAliases), { + const clientJson = (() => { + const restored = restoreImageGenCallsInJson(text, imageGenCallAliases); + if (!hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair)) return restored; + let outbound: unknown; + try { + outbound = JSON.parse(request.body); + } catch { + outbound = undefined; + } + return repairResponsesSnapshotJson(restored, outbound); + })(); + return new Response(clientJson, { status: upstreamResponse.status, statusText: upstreamResponse.statusText, headers, diff --git a/src/server/sse-payload-rewrite.ts b/src/server/sse-payload-rewrite.ts index 712853cee..a79151532 100644 --- a/src/server/sse-payload-rewrite.ts +++ b/src/server/sse-payload-rewrite.ts @@ -9,6 +9,56 @@ import type { TranslatorBudget } from "../lib/translator-budget"; export type SsePayloadRewrite = (payload: string) => string; +/** + * Block-level SSE rewrite: maps one complete SSE event block (without its + * blank-line delimiter) to zero or more replacement blocks. This is the + * contract lifecycle repair needs: injecting missing canonical events + * (#893) is impossible in the one-payload-in/one-payload-out model. + * + * `dispose` releases any retained state (budget-charged collectors) when the + * relay tears down — terminal, EOF, cancel, or error. Relays call it exactly + * once per teardown path. + */ +export type SseBlockRewrite = ((block: string) => readonly string[]) & { + dispose?: () => void; +}; + +/** Adapt a payload rewrite to the block contract (replace only on change). */ +export function payloadRewriteAsBlockRewrite(rewrite: SsePayloadRewrite): SseBlockRewrite { + return (block) => { + const payload = sseDataPayload(block); + if (payload === null) return [block]; + const rewritten = rewrite(payload); + return rewritten !== payload ? [replaceSseDataPayload(block, rewritten)] : [block]; + }; +} + +/** Chain block rewrites: every block stage N emits feeds stage N+1. */ +export function composeSseBlockRewrites(...rewrites: SseBlockRewrite[]): SseBlockRewrite { + const active = rewrites.filter(Boolean); + if (active.length === 0) return Object.assign((block: string) => [block], {}); + const composed: SseBlockRewrite = (block: string) => { + let blocks: readonly string[] = [block]; + for (const rewrite of active) { + const next: string[] = []; + for (const current of blocks) next.push(...rewrite(current)); + blocks = next; + } + return blocks; + }; + // Child disposal is part of the contract: one idempotent disposer for the + // whole chain, so relay teardown never leaks a nested collector. + let disposed = false; + composed.dispose = () => { + if (disposed) return; + disposed = true; + for (const rewrite of active) { + try { rewrite.dispose?.(); } catch { /* teardown must not throw */ } + } + }; + return composed; +} + /** Split one complete SSE event block while retaining its original blank-line delimiter. */ export function nextSseBlock(buffer: string): { block: string; delimiter: string; rest: string } | null { const match = buffer.match(/\r?\n\r?\n/); @@ -69,12 +119,34 @@ export function relaySseWithPayloadRewrite( body: ReadableStream, rewrite: SsePayloadRewrite, translatorBudget: TranslatorBudget, +): ReadableStream { + return relaySseWithBlockRewrite(body, payloadRewriteAsBlockRewrite(rewrite), translatorBudget); +} + +/** + * Relay an SSE body through a single JS pull wrapper, applying a block-level + * rewrite that may emit zero or more blocks per upstream event (lifecycle + * event injection, #893). The original stream's delimiter style is preserved + * for every emitted block. + */ +export function relaySseWithBlockRewrite( + body: ReadableStream, + rewrite: SseBlockRewrite, + translatorBudget: TranslatorBudget, ): ReadableStream { const reader = body.getReader(); const decoder = new TextDecoder(); const encoder = new TextEncoder(); let buffer = ""; let bufferBytes = 0; + // Relays have several independent teardown paths; disposal is exactly once. + let disposed = false; + let cancelled = false; + const disposeRewrite = (): void => { + if (disposed) return; + disposed = true; + try { rewrite.dispose?.(); } catch { /* teardown must not throw */ } + }; const appendBuffer = (fragment: string): void => { if (!fragment) return; @@ -130,20 +202,18 @@ export function relaySseWithPayloadRewrite( let next: { block: string; delimiter: string; rest: string } | null; while ((next = nextSseBlock(buffer))) { replaceBuffer(next.rest); - const payload = sseDataPayload(next.block); - const rewrittenPayload = payload ? rewrite(payload) : undefined; - const block = payload && rewrittenPayload !== undefined && rewrittenPayload !== payload - ? replaceSseDataPayload(next.block, rewrittenPayload) - : next.block; - enqueueText(controller, block + next.delimiter); + for (const outBlock of rewrite(next.block)) { + enqueueText(controller, outBlock + next.delimiter); + } } if (flushFinal && buffer.length > 0) { - const payload = sseDataPayload(buffer); - const rewrittenPayload = payload ? rewrite(payload) : undefined; - const block = payload && rewrittenPayload !== undefined && rewrittenPayload !== payload - ? replaceSseDataPayload(buffer, rewrittenPayload) - : buffer; - enqueueText(controller, block); + const tailBlocks = rewrite(buffer); + // A trailing fragment has no delimiter of its own; multiple emitted + // blocks must still be framed as separate events (#893 review). + const tailDelimiter = buffer.includes("\r\n") ? "\r\n\r\n" : "\n\n"; + for (let i = 0; i < tailBlocks.length; i++) { + enqueueText(controller, tailBlocks[i]! + (i < tailBlocks.length - 1 ? tailDelimiter : "")); + } releaseBuffer(); } }; @@ -152,10 +222,14 @@ export function relaySseWithPayloadRewrite( async pull(controller) { try { const { done, value } = await reader.read(); + // A cancel raced this pending read: never feed the rewriter again + // after its disposal (#893 review). + if (cancelled) return; if (done) { appendBuffer(decoder.decode()); emitProcessedBlocks(controller, true); releaseBuffer(); + disposeRewrite(); controller.close(); return; } @@ -163,12 +237,15 @@ export function relaySseWithPayloadRewrite( emitProcessedBlocks(controller); } catch (error) { releaseBuffer(); + disposeRewrite(); try { await reader.cancel(error); } catch { /* already closed */ } controller.error(error); } }, cancel(reason) { + cancelled = true; releaseBuffer(); + disposeRewrite(); reader.cancel(reason).catch(() => {}); }, }); diff --git a/src/types.ts b/src/types.ts index ae84aa674..3c23b8e48 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1111,6 +1111,12 @@ export interface OcxProviderConfig { * Use for non-forward Responses gateways that reserve a hosted tool namespace server-side. */ modelPreferHostedTools?: Record; + /** + * Provider-local repair for Responses gateways whose lifecycle snapshots omit canonical + * fields or closing events (#893). Disabled by default and applied only to client-facing + * SSE/JSON; raw inspection state remains authoritative. + */ + responsesSnapshotRepair?: boolean; /** Provider-wide mapping from Codex effort labels to upstream `reasoning_effort` values. */ reasoningEffortMap?: Record; /** Model-specific mapping from Codex effort labels to upstream `reasoning_effort` values. */ diff --git a/tests/config.test.ts b/tests/config.test.ts index e5855aa49..579ca6a68 100644 --- a/tests/config.test.ts +++ b/tests/config.test.ts @@ -648,6 +648,36 @@ describe("opencodex config defaults", () => { expect(readConfigDiagnostics().error).toContain("responsesItemIdRepair"); }); + test("accepts only a boolean responsesSnapshotRepair opt-in", () => { + writeConfig({ + port: 12345, + providers: { + custom: { + adapter: "openai-responses", + baseUrl: "https://example.test/v1", + responsesSnapshotRepair: true, + }, + }, + defaultProvider: "custom", + }); + expect(readConfigDiagnostics().error).toBeNull(); + expect(readConfigDiagnostics().config.providers.custom.responsesSnapshotRepair).toBe(true); + + writeConfig({ + port: 12345, + providers: { + custom: { + adapter: "openai-responses", + baseUrl: "https://example.test/v1", + responsesSnapshotRepair: { enabled: true }, + }, + }, + defaultProvider: "custom", + }); + expect(readConfigDiagnostics().source).toBe("fallback"); + expect(readConfigDiagnostics().error).toContain("responsesSnapshotRepair"); + }); + test("accepts a relative responsesPath", () => { writeResponsesPathConfig("/responses"); diff --git a/tests/management-provider-validation.test.ts b/tests/management-provider-validation.test.ts index cf8195b0d..3cff04c38 100644 --- a/tests/management-provider-validation.test.ts +++ b/tests/management-provider-validation.test.ts @@ -192,6 +192,27 @@ describe("provider management validation", () => { })).toContain("not supported on forward-auth"); }); + test("provider management permits snapshot repair only on canonical OpenAI forward seeds", () => { + for (const mode of ["pool", "direct"] as const) { + expect(providerManagementConfigError("openai", { + ...canonicalDirect, + codexAccountMode: mode, + responsesSnapshotRepair: true, + })).toBeNull(); + } + + expect(providerManagementConfigError("openai", { + ...canonicalDirect, + responsesSnapshotRepair: { enabled: true }, + })).toBe("provider openai responsesSnapshotRepair must be a boolean"); + + expect(providerManagementConfigError("openai", { + ...canonicalDirect, + responsesSnapshotRepair: true, + noVisionModels: ["gpt-5.6"], + })).toContain("canonical built-in provider seed"); + }); + test("provider management validates retryOn429 bounds and unknown keys", () => { const base = { adapter: "openai-chat", baseUrl: "https://api.openai.com/v1" }; expect(providerManagementConfigError("custom", { diff --git a/tests/passthrough-abort.test.ts b/tests/passthrough-abort.test.ts index d702452cb..5c782b607 100644 --- a/tests/passthrough-abort.test.ts +++ b/tests/passthrough-abort.test.ts @@ -52,7 +52,7 @@ describe("passthrough relayWithAbort (RC2, passthrough path)", () => { expect(sseBranch).toContain("const repairConfig = route.provider.responsesItemIdRepair;"); expect(sseBranch).toContain("const needsClientRewrite = imageGenCallAliases.size > 0"); expect(sseBranch).toContain("new Response(eagerBody"); - expect(sseBranch).toContain("const rewrittenBody = payloadRewrites.length > 0"); + expect(sseBranch).toContain("const rewrittenBody = clientBlockRewrite !== undefined || payloadRewrites.length > 0"); expect(sseBranch).toContain("eagerPath?.useEagerRelay || win32EagerRewrite"); expect(sseBranch).not.toContain("win32TerminalRelay"); // #864: win32 traffic that DOES need a client rewrite takes the eager single diff --git a/tests/responses-snapshot-repair-server.test.ts b/tests/responses-snapshot-repair-server.test.ts new file mode 100644 index 000000000..ddfba352d --- /dev/null +++ b/tests/responses-snapshot-repair-server.test.ts @@ -0,0 +1,147 @@ +import { afterEach, beforeEach, describe, expect, setDefaultTimeout, test } from "bun:test"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { saveConfig } from "../src/config"; +import { startServer } from "../src/server"; +import type { OcxConfig } from "../src/types"; +import { installIsolatedCodexHome, type IsolatedCodexHome } from "./helpers/isolated-codex-home"; + +setDefaultTimeout(30_000); + +const originalFetch = globalThis.fetch; +let TEST_DIR = ""; +let isolated: IsolatedCodexHome; + +const SPARSE_EVENTS = [ + { type: "response.created", response: { id: "resp_sparse" } }, + { type: "response.output_item.added", output_index: 0, item: { type: "message", id: "msg_sparse" } }, + { type: "response.output_text.delta", item_id: "msg_sparse", output_index: 0, delta: "hello" }, + { type: "response.completed", response: { id: "resp_sparse" } }, +]; + +function sparseSseBody(): ReadableStream { + return new ReadableStream({ + start(controller) { + const encoder = new TextEncoder(); + for (const event of SPARSE_EVENTS) { + controller.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`)); + } + controller.enqueue(encoder.encode("data: [DONE]\n\n")); + controller.close(); + }, + }); +} + +function stubSparseGateway(origin: string): void { + globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => { + const requestUrl = typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url; + const url = new URL(requestUrl); + if (url.origin === origin && url.pathname.endsWith("/models")) { + return Response.json({ data: [] }); + } + if (url.origin === origin && url.pathname.endsWith("/responses")) { + return new Response(sparseSseBody(), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + } + return originalFetch(input, init); + }) as typeof fetch; +} + +beforeEach(() => { + TEST_DIR = mkdtempSync(join(tmpdir(), "ocx-snapshot-repair-server-")); + process.env.OPENCODEX_HOME = TEST_DIR; + isolated = installIsolatedCodexHome("ocx-snapshot-repair-codex-"); +}); + +afterEach(async () => { + globalThis.fetch = originalFetch; + await isolated.restore(); + rmSync(TEST_DIR, { recursive: true, force: true }); +}); + +describe("responsesSnapshotRepair through /v1/responses", () => { + test("an opt-in gateway's sparse stream reaches the client as the full canonical lifecycle", async () => { + const gateway = "https://sparse.example.test"; + stubSparseGateway(gateway); + saveConfig({ + port: 0, + defaultProvider: "sparse", + providers: { + sparse: { + adapter: "openai-responses", + baseUrl: `${gateway}/v1`, + authMode: "key", + apiKey: "test-key", + responsesSnapshotRepair: true, + }, + }, + } as OcxConfig); + + const server = startServer(0); + try { + const response = await originalFetch(new URL("/v1/responses", server.url), { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "sparse-model", input: "hi", stream: true }), + }); + expect(response.status).toBe(200); + const text = await response.text(); + const sequence = [...text.matchAll(/"type":"([^"]+)"/g)].map(match => match[1]); + for (const expected of [ + "response.created", + "response.output_item.added", + "response.content_part.added", + "response.output_text.delta", + "response.output_text.done", + "response.content_part.done", + "response.output_item.done", + "response.completed", + ]) { + expect(sequence).toContain(expected); + } + // The terminal snapshot carries the reconstructed committed message. + const completedLine = text.split("\n").find(line => line.includes('"response.completed"')); + expect(completedLine).toBeDefined(); + const completed = JSON.parse(completedLine!.replace(/^data: /, "")) as { response: { output: { id: string }[] } }; + expect(completed.response.output[0]?.id).toBe("msg_sparse"); + } finally { + await server.stop(true); + } + }); + + test("the same gateway without the opt-in relays the sparse stream unchanged", async () => { + const gateway = "https://sparse-off.example.test"; + stubSparseGateway(gateway); + saveConfig({ + port: 0, + defaultProvider: "sparse", + providers: { + sparse: { + adapter: "openai-responses", + baseUrl: `${gateway}/v1`, + authMode: "key", + apiKey: "test-key", + }, + }, + } as OcxConfig); + + const server = startServer(0); + try { + const response = await originalFetch(new URL("/v1/responses", server.url), { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "sparse-model", input: "hi", stream: true }), + }); + expect(response.status).toBe(200); + const text = await response.text(); + expect(text).not.toContain("response.content_part.added"); + expect(text).not.toContain("response.output_text.done"); + expect(text).not.toContain("response.output_item.done"); + } finally { + await server.stop(true); + } + }); +}); diff --git a/tests/responses-snapshot-repair.test.ts b/tests/responses-snapshot-repair.test.ts new file mode 100644 index 000000000..08093c4b6 --- /dev/null +++ b/tests/responses-snapshot-repair.test.ts @@ -0,0 +1,357 @@ +import { describe, expect, test } from "bun:test"; +import { + createResponsesSnapshotBlockRewrite, + hasResponsesSnapshotRepair, + repairResponsesSnapshotJson, +} from "../src/server/responses-snapshot-repair"; +import { + composeSseBlockRewrites, + payloadRewriteAsBlockRewrite, + relaySseWithBlockRewrite, +} from "../src/server/sse-payload-rewrite"; +import { createTestTranslatorBudget } from "./helpers/translator-budget"; + + +function dataBlock(payload: unknown): string { + return `data: ${typeof payload === "string" ? payload : JSON.stringify(payload)}`; +} + +function eventsOf(blocks: readonly string[]): Record[] { + return blocks.map((block) => JSON.parse(block.replace(/^data: /, "")) as Record); +} + +function typesOf(blocks: readonly string[]): string[] { + return eventsOf(blocks).map(event => String(event.type)); +} + +const ISSUE_FIXTURE = { + created: { type: "response.created", response: { id: "resp_1" } }, + itemAdded: { type: "response.output_item.added", output_index: 0, item: { type: "message", id: "msg_1" } }, + delta: { type: "response.output_text.delta", item_id: "msg_1", output_index: 0, delta: "hello" }, + completed: { type: "response.completed", response: { id: "resp_1" } }, +}; + +describe("createResponsesSnapshotBlockRewrite", () => { + test("the exact #893 issue fixture yields the full canonical lifecycle and a committed message", () => { + const rewrite = createResponsesSnapshotBlockRewrite(undefined, createTestTranslatorBudget()); + const out: string[] = []; + for (const event of [ISSUE_FIXTURE.created, ISSUE_FIXTURE.itemAdded, ISSUE_FIXTURE.delta, ISSUE_FIXTURE.completed]) { + out.push(...rewrite(dataBlock(event))); + } + expect(typesOf(out)).toEqual([ + "response.created", + "response.output_item.added", + "response.content_part.added", + "response.output_text.delta", + "response.output_text.done", + "response.content_part.done", + "response.output_item.done", + "response.completed", + ]); + // Snapshot field backfills. + const created = eventsOf(out)[0]!.response as Record; + expect(created.status).toBe("in_progress"); + expect(created.parallel_tool_calls).toBe(true); + expect(created.tool_choice).toBe("auto"); + expect(created.tools).toEqual([]); + // The repaired item-added has role/status/content. + const added = eventsOf(out)[1]!.item as Record; + expect(added.role).toBe("assistant"); + expect(added.status).toBe("in_progress"); + expect(Array.isArray(added.content)).toBe(true); + // The injected done carries the accumulated text. + const textDone = eventsOf(out)[4]!; + expect(textDone.text).toBe("hello"); + // The committed done item and terminal output contain the message. + const doneItem = eventsOf(out)[6]!.item as Record; + expect(doneItem.status).toBe("completed"); + const terminal = eventsOf(out)[7]!.response as Record; + expect(terminal.status).toBe("completed"); + const output = terminal.output as Record[]; + expect(output).toHaveLength(1); + expect(output[0]).toMatchObject({ type: "message", id: "msg_1", status: "completed" }); + expect((output[0]!.content as Record[])[0]).toMatchObject({ text: "hello" }); + }); + + test("an explicit output: [] in the terminal snapshot is authoritative (never reconstructed)", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock(ISSUE_FIXTURE.created)); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); + rewrite(dataBlock(ISSUE_FIXTURE.delta)); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", status: "completed", output: [] }, + })); + const terminal = eventsOf(out).find(event => event.type === "response.completed")!; + expect((terminal.response as Record).output).toEqual([]); + }); + + test("already-canonical streams pass through byte-identical", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + const canonical = [ + { type: "response.created", response: { id: "r", status: "in_progress", output: [], parallel_tool_calls: true, tool_choice: "auto", tools: [] } }, + { type: "response.output_item.added", output_index: 0, item: { type: "message", id: "m", role: "assistant", status: "in_progress", content: [] } }, + { type: "response.content_part.added", item_id: "m", output_index: 0, content_index: 0, part: { type: "output_text", text: "", annotations: [] } }, + { type: "response.output_text.delta", item_id: "m", output_index: 0, logprobs: [], delta: "hi" }, + { type: "response.output_text.done", item_id: "m", output_index: 0, logprobs: [], text: "hi" }, + { type: "response.content_part.done", item_id: "m", output_index: 0, content_index: 0, part: { type: "output_text", text: "hi", annotations: [] } }, + { type: "response.output_item.done", output_index: 0, item: { type: "message", id: "m", role: "assistant", status: "completed", content: [{ type: "output_text", text: "hi", annotations: [] }] } }, + { type: "response.completed", response: { id: "r", status: "completed", output: [{ type: "message", id: "m", role: "assistant", status: "completed", content: [{ type: "output_text", text: "hi", annotations: [] }] }], parallel_tool_calls: true, tool_choice: "auto", tools: [] } }, + ]; + const out: string[] = []; + for (const event of canonical) out.push(...rewrite(dataBlock(event))); + expect(out).toHaveLength(canonical.length); + // No injections: exactly the same event type sequence. + expect(typesOf(out)).toEqual(canonical.map(event => event.type)); + }); + + test("malformed JSON and non-data blocks pass through untouched", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + expect(rewrite("data: {not json")).toEqual(["data: {not json"]); + expect(rewrite(": comment-only")).toEqual([": comment-only"]); + }); + + test("a contradictory duplicate open index taints the stream: fail closed, no injection", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); // same open output_index reused + rewrite(dataBlock(ISSUE_FIXTURE.delta)); + const out = rewrite(dataBlock(ISSUE_FIXTURE.completed)); + // Fail closed: the terminal arrives with no injected closing events and no reconstruction. + expect(typesOf(out)).toEqual(["response.completed"]); + const terminal = eventsOf(out)[0]!.response as Record; + expect(Object.hasOwn(terminal, "output")).toBe(false); + }); + + test("an oversized completed item taints reconstruction but field repairs continue", () => { + const rewrite = createResponsesSnapshotBlockRewrite(undefined, createTestTranslatorBudget()); + rewrite(dataBlock({ type: "response.output_item.added", output_index: 0, item: { type: "message", id: "big" } })); + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", id: "big", role: "assistant", status: "completed", + content: [{ type: "output_text", text: "x".repeat(9 * 1024 * 1024), annotations: [] }], + }, + })); + const out = rewrite(dataBlock({ type: "response.completed", response: { id: "r" } })); + expect(typesOf(out)).toEqual(["response.completed"]); + const terminal = eventsOf(out)[0]!.response as Record; + // Field backfills still applied, but no output reconstruction. + expect(terminal.status).toBe("completed"); + expect(Object.hasOwn(terminal, "output")).toBe(false); + }); + + test("gateway-closed items get no injected duplicates", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); + rewrite(dataBlock({ type: "response.content_part.added", item_id: "msg_1", output_index: 0, content_index: 0, part: { type: "output_text", text: "", annotations: [] } })); + rewrite(dataBlock(ISSUE_FIXTURE.delta)); + rewrite(dataBlock({ type: "response.output_text.done", item_id: "msg_1", output_index: 0, text: "hello" })); + rewrite(dataBlock({ type: "response.content_part.done", item_id: "msg_1", output_index: 0, content_index: 0, part: { type: "output_text", text: "hello", annotations: [] } })); + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: { type: "message", id: "msg_1", status: "completed", content: [{ type: "output_text", text: "hello", annotations: [] }] } })); + const out = rewrite(dataBlock(ISSUE_FIXTURE.completed)); + expect(typesOf(out)).toEqual(["response.completed"]); + }); + + test("request defaults fill snapshot fields from the request", () => { + const rewrite = createResponsesSnapshotBlockRewrite({ + parallel_tool_calls: false, + tool_choice: { type: "function", name: "search" }, + tools: [{ type: "function", name: "search" }], + }); + const out = rewrite(dataBlock(ISSUE_FIXTURE.created)); + const response = eventsOf(out)[0]!.response as Record; + expect(response.parallel_tool_calls).toBe(false); + expect(response.tool_choice).toEqual({ type: "function", name: "search" }); + expect(response.tools).toEqual([{ type: "function", name: "search" }]); + }); + + test("failed and incomplete terminals never fabricate completed output", () => { + for (const terminalType of ["response.failed", "response.incomplete"]) { + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); + rewrite(dataBlock(ISSUE_FIXTURE.delta)); + const out = rewrite(dataBlock({ type: terminalType, response: { id: "r" } })); + expect(typesOf(out)).toEqual([terminalType]); + const response = eventsOf(out)[0]!.response as Record; + expect(Object.hasOwn(response, "output")).toBe(false); + expect(JSON.stringify(out)).not.toContain("output_item.done"); + } + }); + + test("dispose releases retained collectors on EOF without a terminal", () => { + const budget = createTestTranslatorBudget(); + const rewrite = createResponsesSnapshotBlockRewrite(undefined, budget); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 0, + item: { type: "message", id: "msg_1", status: "completed", content: [{ type: "output_text", text: "hello", annotations: [] }] }, + })); + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + rewrite.dispose?.(); + expect(budget.snapshot().currentBytes).toBe(0); + }); + + test("function_call items do not taint later sparse messages", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock({ type: "response.output_item.added", output_index: 0, item: { type: "function_call", id: "fc_1", call_id: "c1", name: "search" } })); + rewrite(dataBlock({ ...ISSUE_FIXTURE.itemAdded, output_index: 1 })); + const fromDelta = rewrite(dataBlock({ ...ISSUE_FIXTURE.delta, output_index: 1, item_id: "msg_1" })); + const out = [...fromDelta, ...rewrite(dataBlock(ISSUE_FIXTURE.completed))]; + // Message lifecycle completed; the open function_call blocks only reconstruction. + expect(typesOf(out)).toContain("response.content_part.added"); + expect(typesOf(out)).toContain("response.output_item.done"); + const terminal = eventsOf(out).find(event => event.type === "response.completed")!; + expect(Object.hasOwn(terminal.response as Record, "output")).toBe(false); + }); + + test("a mismatched item_id on output_text.done goes fail-closed instead of double-closing", () => { + // #1025 review blocker 2: this used to ignore the foreign done, leave the + // item open, and let the terminal synthesize a SECOND output_text.done + + // output_item.done — a double close on a stream whose identity model we + // just proved wrong. A contradictory item_id now taints, exactly like a + // mismatched output_item.done. + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); + const foreign = rewrite(dataBlock({ type: "response.output_text.done", item_id: "msg_OTHER", output_index: 0, text: "foreign" })); + // The contradictory block itself is forwarded untouched, not swallowed. + expect(eventsOf(foreign).some(event => event.item_id === "msg_OTHER")).toBe(true); + rewrite(dataBlock(ISSUE_FIXTURE.delta)); + const out = rewrite(dataBlock(ISSUE_FIXTURE.completed)); + // No synthetic closure and no reconstruction after the taint. + expect(typesOf(out)).not.toContain("response.output_text.done"); + expect(typesOf(out)).not.toContain("response.output_item.done"); + const terminal = eventsOf(out).find(event => event.type === "response.completed")!; + expect(Object.hasOwn(terminal.response as Record, "output")).toBe(false); + }); + + test("a completed terminal with absent output and zero items gets the canonical empty list", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + const out = rewrite(dataBlock({ type: "response.completed", response: { id: "r" } })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + }); + + test("a completed terminal with malformed output is replaced canonically", () => { + for (const malformed of [null, { bogus: true }, "not-an-array"]) { + const rewrite = createResponsesSnapshotBlockRewrite(); + const out = rewrite(dataBlock({ type: "response.completed", response: { id: "r", output: malformed } })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + } + }); + + test("an unterminated final block with injections stays separately framed", async () => { + // EOF immediately after a complete-but-undelimited block: the injected + // closing events must not fuse into one invalid data record. + const upstream = new ReadableStream({ + start(controller) { + const encoder = new TextEncoder(); + controller.enqueue(encoder.encode(`${dataBlock(ISSUE_FIXTURE.itemAdded)}\n\n${dataBlock(ISSUE_FIXTURE.delta)}`)); + controller.close(); + }, + }); + const rewrite = createResponsesSnapshotBlockRewrite(undefined, createTestTranslatorBudget()); + const body = relaySseWithBlockRewrite(upstream, rewrite, createTestTranslatorBudget()); + const text = await new Response(body).text(); + const frames = text.split(/\r?\n\r?\n/).filter(frame => frame.trim().length > 0); + for (const frame of frames) { + expect(frame.startsWith("data:")).toBe(true); + expect(frame.indexOf("\ndata:")).toBe(-1); + } + }); + + test("a done event with a mismatched identity closes nothing and fails closed", () => { + const rewrite = createResponsesSnapshotBlockRewrite(); + rewrite(dataBlock(ISSUE_FIXTURE.itemAdded)); // tracks msg_1 + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 0, + item: { type: "message", id: "msg_WRONG", status: "completed", content: [{ type: "output_text", text: "wrong", annotations: [] }] }, + })); + const out = rewrite(dataBlock(ISSUE_FIXTURE.completed)); + // Fail closed: no injection, and the wrong item is never reconstructed. + expect(typesOf(out)).toEqual(["response.completed"]); + expect(JSON.stringify(out)).not.toContain("msg_WRONG"); + const terminal = eventsOf(out)[0]!.response as Record; + expect(Object.hasOwn(terminal, "output")).toBe(false); + }); + + test("composed chains dispose every child exactly once", () => { + let disposed = 0; + const childA = Object.assign((block: string) => [block], { dispose: () => { disposed++; } }); + const childB = Object.assign((block: string) => [block], { dispose: () => { disposed++; } }); + const chain = composeSseBlockRewrites(childA, childB); + chain.dispose?.(); + chain.dispose?.(); + expect(disposed).toBe(2); + }); +}); + +describe("repairResponsesSnapshotJson", () => { + test("backfills missing fields on a non-streaming response", () => { + const repaired = repairResponsesSnapshotJson(JSON.stringify({ id: "r", output: [] })); + const parsed = JSON.parse(repaired) as Record; + expect(parsed.status).toBe("completed"); + expect(parsed.parallel_tool_calls).toBe(true); + expect(parsed.output).toEqual([]); + }); + + test("invalid JSON passes through", () => { + expect(repairResponsesSnapshotJson("{nope")).toBe("{nope"); + }); +}); + +describe("hasResponsesSnapshotRepair", () => { + test("only explicit true enables", () => { + expect(hasResponsesSnapshotRepair(true)).toBe(true); + expect(hasResponsesSnapshotRepair(false)).toBe(false); + expect(hasResponsesSnapshotRepair(undefined)).toBe(false); + }); +}); + +describe("relaySseWithBlockRewrite (stream level)", () => { + test("the issue fixture flows through the relay as the full canonical sequence", async () => { + const upstream = new ReadableStream({ + start(controller) { + const encoder = new TextEncoder(); + for (const event of [ISSUE_FIXTURE.created, ISSUE_FIXTURE.itemAdded, ISSUE_FIXTURE.delta, ISSUE_FIXTURE.completed]) { + controller.enqueue(encoder.encode(`${dataBlock(event)}\n\n`)); + } + controller.enqueue(encoder.encode("data: [DONE]\n\n")); + controller.close(); + }, + }); + const rewrite = createResponsesSnapshotBlockRewrite(undefined, createTestTranslatorBudget()); + const body = relaySseWithBlockRewrite(upstream, rewrite, createTestTranslatorBudget()); + const text = await new Response(body).text(); + const sequence = [...text.matchAll(/"type":"([^"]+)"/g)].map(match => match[1]); + for (const expected of [ + "response.created", + "response.output_item.added", + "response.content_part.added", + "response.output_text.delta", + "response.output_text.done", + "response.content_part.done", + "response.output_item.done", + "response.completed", + ]) { + expect(sequence).toContain(expected); + } + expect(text).toContain("data: [DONE]"); + }); + + test("composition: payload rewrite then lifecycle repair", async () => { + const chain = composeSseBlockRewrites( + payloadRewriteAsBlockRewrite(payload => payload.replace("hello", "hola")), + createResponsesSnapshotBlockRewrite(), + ); + chain(dataBlock(ISSUE_FIXTURE.itemAdded)); + const out = chain(dataBlock(ISSUE_FIXTURE.delta)); + // content_part.added injected first; the delta's payload got the payload rewrite. + expect(typesOf(out)).toEqual(["response.content_part.added", "response.output_text.delta"]); + expect(JSON.stringify(out[1])).toContain("hola"); + }); +});