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 42b5e5997..03851a3d0 100644 --- a/docs-site/src/content/docs/ja/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ja/reference/configuration/providers.md @@ -78,6 +78,7 @@ description: プロバイダー エントリ、認証、エンドポイント、 | `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 検査と永続化は変更しません。 | | `autoToolChoiceOnlyModels?` | `string[]` | `tool_choice` が `auto` または `none` のみを受け入れるモデル。強制的な選択は格下げされます。 | | `preserveReasoningContentModels?` | `string[]` |チャット履歴に以前のアシスタント `reasoning_content` が必要なモデル。 | | `thinkingToggleModels?` | `string[]` |エフォート ラダーではなく `thinking.enabled` を使用してモデルをチャットします。 | @@ -191,7 +192,8 @@ Anthropic アカウント ポリシーのリスクを理解していない限り "reasoning": ["rs_0"], "message": ["msg_0"], "repairMissingTerminalIds": true - } + }, + "responsesSnapshotRepair": true } } } 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 25b40955d..5f61522de 100644 --- a/docs-site/src/content/docs/ko/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ko/reference/configuration/providers.md @@ -78,6 +78,7 @@ description: 공급자 항목, 인증, 엔드포인트, 모델 카탈로그, 할 | `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 검사와 영속화는 변경하지 않습니다. | | `autoToolChoiceOnlyModels?` | `string[]` | `tool_choice`가 `auto` 또는 `none`만 받는 모델입니다. 강제 선택은 낮은 수준으로 바뀝니다. | | `preserveReasoningContentModels?` | `string[]` | chat 기록에서 이전 assistant `reasoning_content`가 필요한 모델입니다. | | `thinkingToggleModels?` | `string[]` | effort 계층 대신 `thinking.enabled`를 쓰는 chat 모델입니다. | @@ -194,7 +195,8 @@ Anthropic 계정 정책 위험을 이해하지 못한다면 이 기능은 꺼두 "reasoning": ["rs_0"], "message": ["msg_0"], "repairMissingTerminalIds": true - } + }, + "responsesSnapshotRepair": true } } } diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index 80c7ed520..2f5101d65 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -89,6 +89,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. | | `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. | | `thinkingToggleModels?` | `string[]` | Chat models using `thinking.enabled` rather than an effort ladder. | @@ -235,7 +236,8 @@ For a broken `openai-responses` gateway, repair belongs on the provider object: "reasoning": ["rs_0"], "message": ["msg_0"], "repairMissingTerminalIds": true - } + }, + "responsesSnapshotRepair": true } } } 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 1ece0bb05..8bbc02423 100644 --- a/docs-site/src/content/docs/ru/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ru/reference/configuration/providers.md @@ -94,6 +94,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. | | `autoToolChoiceOnlyModels?` | `string[]` | Модели, у которых `tool_choice` принимает только `auto` или `none`; forced choice понижается. | | `preserveReasoningContentModels?` | `string[]` | Модели, которым нужен предыдущий assistant `reasoning_content` в chat history. | | `thinkingToggleModels?` | `string[]` | Chat-модели, использующие `thinking.enabled` вместо effort-ladder. | @@ -245,7 +246,8 @@ Beijing, а `alibaba-token-plan-intl` обслуживает междунаро "reasoning": ["rs_0"], "message": ["msg_0"], "repairMissingTerminalIds": true - } + }, + "responsesSnapshotRepair": true } } } 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 9dc3f50d8..eb35571fe 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 @@ -78,6 +78,7 @@ description: 提供者条目、身份验证、端点、模型目录、配额、 | `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 和工具元数据;原始检查与持久化保持不变。 | | `autoToolChoiceOnlyModels?` | `string[]` | `tool_choice` 只接受 `auto` 或 `none` 的模型;强制选择会被降级。 | | `preserveReasoningContentModels?` | `string[]` | 需要在聊天历史中保留先前 assistant `reasoning_content` 的模型。 | | `thinkingToggleModels?` | `string[]` | 使用 `thinking.enabled` 而不是 effort 阶梯的 chat 模型。 | @@ -188,7 +189,8 @@ affinity。这些策略不能规避 provider enforcement。 "reasoning": ["rs_0"], "message": ["msg_0"], "repairMissingTerminalIds": true - } + }, + "responsesSnapshotRepair": true } } } diff --git a/src/config.ts b/src/config.ts index fe523323a..24da92803 100644 --- a/src/config.ts +++ b/src/config.ts @@ -568,6 +568,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 79b02c329..1fe4e3546 100644 --- a/src/server/auth-cors.ts +++ b/src/server/auth-cors.ts @@ -406,7 +406,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`; } @@ -440,6 +442,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/responses-snapshot-repair.ts b/src/server/responses-snapshot-repair.ts new file mode 100644 index 000000000..093141056 --- /dev/null +++ b/src/server/responses-snapshot-repair.ts @@ -0,0 +1,350 @@ +import type { TranslatorBudget } from "../lib/translator-budget"; +import { + MAX_COMPLETED_OUTPUT_ITEMS, + MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES, +} from "./relay"; +import type { SsePayloadRewrite } 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[]; +}; + +type RetainedOutputItem = { + item: Record; + sourceBytes: number; +}; + +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; + } else { + repaired.output = []; + changed = true; + } + 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; +} + +/** + * Repair required fields omitted by a few Responses-compatible gateways across lifecycle events. + * Existing upstream values remain authoritative; only absent or structurally invalid fields are + * backfilled. + */ +export function createResponsesSnapshotPayloadRewrite( + requestBody?: unknown, + budget?: TranslatorBudget, +): SsePayloadRewrite { + const defaults = requestDefaults(requestBody); + const completedItems = new Map(); + const unfinishedItemIndexes = new Set(); + let aggregateItemBytes = 0; + let reconstructionTainted = false; + + const clearReconstructionState = (): void => { + if (aggregateItemBytes > 0) { + budget?.releaseRetained(aggregateItemBytes, { kind: "retained_collectors" }); + } + completedItems.clear(); + unfinishedItemIndexes.clear(); + aggregateItemBytes = 0; + reconstructionTainted = false; + }; + + const retainCompletedItem = ( + index: number, + item: Record, + sourceBytes: number, + ): void => { + 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" }); + } + reconstructionTainted = true; + 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" }); + reconstructionTainted = true; + } + }; + + return (payload: string): string => { + let event: unknown; + try { + event = JSON.parse(payload); + } catch { + return payload; + } + if (!isPlainObject(event)) return payload; + const type = typeof event.type === "string" ? event.type : ""; + let nextEvent = event; + let changed = false; + const outputIndex = Number.isInteger(event.output_index) && (event.output_index as number) >= 0 + ? event.output_index as number + : undefined; + + if (type === "response.output_item.added") { + if (outputIndex === undefined) reconstructionTainted = true; + else if (!unfinishedItemIndexes.has(outputIndex) + && unfinishedItemIndexes.size >= MAX_COMPLETED_OUTPUT_ITEMS) { + reconstructionTainted = true; + } else { + unfinishedItemIndexes.add(outputIndex); + } + } + + if (type === "response.output_item.done" && !isPlainObject(event.item)) { + reconstructionTainted = true; + } + + if ((type === "response.output_item.added" || type === "response.output_item.done") + && isPlainObject(event.item)) { + const itemStatus = type === "response.output_item.done" ? "completed" : "in_progress"; + const item = repairOutputItem(event.item, itemStatus); + if (item !== event.item) { + nextEvent = { ...nextEvent, item }; + changed = true; + } + if (type === "response.output_item.done") { + if (outputIndex !== undefined + && typeof item.type === "string" + && item.type.trim().length > 0) { + unfinishedItemIndexes.delete(outputIndex); + retainCompletedItem( + outputIndex, + item, + Buffer.byteLength(JSON.stringify(item), "utf8"), + ); + } else { + reconstructionTainted = true; + } + } + } + + const responseStatus = Object.prototype.hasOwnProperty.call(RESPONSE_EVENT_STATUSES, type) + ? RESPONSE_EVENT_STATUSES[type] + : undefined; + if (responseStatus && isPlainObject(event.response)) { + const outputIsAbsent = !Object.hasOwn(event.response, "output"); + let reconstructedOutput: Record[] | undefined; + if ((type === "response.completed" || type === "response.incomplete") + && outputIsAbsent + && !reconstructionTainted + && unfinishedItemIndexes.size === 0 + && completedItems.size > 0) { + const orderedItems = [...completedItems.entries()] + .sort(([left], [right]) => left - right); + if (orderedItems.every(([index], position) => index === position)) { + reconstructedOutput = orderedItems.map(([, retained]) => retained.item); + } else { + // A gap means at least one completed item is missing. Never compact later indexes into a + // shorter array that appears complete to persistence or the client. + reconstructionTainted = true; + } + } + const response = repairResponseSnapshot( + event.response, + responseStatus, + defaults, + reconstructedOutput, + ); + if (response !== event.response) { + nextEvent = { ...nextEvent, response }; + changed = true; + } + } + + const shouldClearReconstructionState = type === "response.completed" + || type === "response.failed" + || type === "response.incomplete"; + + if ((type === "response.content_part.added" || type === "response.content_part.done") + && isPlainObject(event.part)) { + 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.done" && typeof event.text !== "string") { + nextEvent = { ...nextEvent, text: "" }; + changed = true; + } + + const result = changed ? JSON.stringify(nextEvent) : payload; + if (shouldClearReconstructionState) clearReconstructionState(); + return result; + }; +} + +/** 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; + 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 3577d85fd..5403d701c 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -150,6 +150,11 @@ import { createResponsesItemIdPayloadRewrite, hasResponsesItemIdRepair, } from "../responses-item-id-repair"; +import { + createResponsesSnapshotPayloadRewrite, + hasResponsesSnapshotRepair, + repairResponsesSnapshotJson, +} from "../responses-snapshot-repair"; import { createImageGenCallRestoreRewrite, imageGenToolCallAliases, @@ -1817,13 +1822,18 @@ 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 needsClientRewrite = imageGenCallAliases.size > 0 + || hasResponsesItemIdRepair(repairConfig) + || hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair); // Compose opt-in payload rewrites into one parse/stringify pass (image-gen restore first). const payloadRewrites = [ createImageGenCallRestoreRewrite(imageGenCallAliases), hasResponsesItemIdRepair(repairConfig) ? createResponsesItemIdPayloadRewrite(repairConfig!, translatorBudget) : undefined, + hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair) + ? createResponsesSnapshotPayloadRewrite(parsed._rawBody, translatorBudget) + : undefined, ].filter((rewrite): rewrite is NonNullable => rewrite !== 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 @@ -1996,7 +2006,11 @@ 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 restoredText = restoreImageGenCallsInJson(text, imageGenCallAliases); + const clientText = hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair) + ? repairResponsesSnapshotJson(restoredText, parsed._rawBody) + : restoredText; + return new Response(clientText, { status: upstreamResponse.status, statusText: upstreamResponse.statusText, headers, diff --git a/src/types.ts b/src/types.ts index 0bbdc7e41..560cf88ba 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1109,6 +1109,12 @@ export interface OcxProviderConfig { * Disabled by default; function_call ids and call_id pairing are never rewritten. */ responsesItemIdRepair?: ResponsesItemIdRepairConfig; + /** + * Provider-local repair for Responses gateways whose lifecycle snapshots omit canonical fields. + * Disabled by default and applied only to client-facing SSE/JSON; raw inspection state remains + * authoritative. + */ + responsesSnapshotRepair?: boolean; /** Model ids whose tool_choice only accepts `auto` or `none`; forced/named choices are downgraded. */ autoToolChoiceOnlyModels?: string[]; /** Model ids that expect prior assistant `reasoning_content` to be preserved in chat history. */ 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 7f4d9e67e..e08dcbf05 100644 --- a/tests/management-provider-validation.test.ts +++ b/tests/management-provider-validation.test.ts @@ -192,6 +192,71 @@ 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"); + + expect(providerManagementConfigError("custom-forward", { + adapter: "openai-responses", + baseUrl: canonicalDirect.baseUrl, + authMode: "forward", + responsesSnapshotRepair: true, + })).toContain('authMode "forward"'); + }); + + test("provider management persists snapshot repair for canonical OpenAI pool and direct modes", async () => { + if (existsSync(TEST_DIR)) rmSync(TEST_DIR, { recursive: true }); + mkdirSync(TEST_DIR, { recursive: true }); + process.env.OPENCODEX_HOME = TEST_DIR; + const liveConfig: OcxConfig = { + port: 0, + defaultProvider: "openai", + openaiProviderTierVersion: 2, + providers: { openai: canonicalDirect }, + }; + saveConfig(liveConfig); + const resolvedError = spyOn(destinationPolicy, "providerDestinationResolvedError").mockResolvedValue(null); + + try { + for (const mode of ["pool", "direct"] as const) { + const request = new Request("http://127.0.0.1/api/providers", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + name: "openai", + provider: { ...canonicalDirect, codexAccountMode: mode, responsesSnapshotRepair: true }, + }), + }); + const response = await handleManagementAPI(request, new URL(request.url), liveConfig, { + refreshCodexCatalog: async () => undefined, + }); + expect(response?.status).toBe(200); + expect(loadConfig().providers.openai).toMatchObject({ + codexAccountMode: mode, + responsesSnapshotRepair: true, + }); + } + } finally { + resolvedError.mockRestore(); + } + }); + test("provider discovery status is additive and omitted before an attempt", async () => { markProviderDiscoveryFailed("auth-broken", { reason: "http", httpStatus: 401 }); try { @@ -247,6 +312,7 @@ describe("provider management validation", () => { adapter: "openai-responses", baseUrl: "https://attacker.example/backend-api/codex", authMode: "forward", + responsesSnapshotRepair: true, }, }), }); @@ -918,6 +984,54 @@ describe("provider management validation", () => { } }); + test("provider management accepts only a boolean responsesSnapshotRepair value", async () => { + if (existsSync(TEST_DIR)) rmSync(TEST_DIR, { recursive: true }); + mkdirSync(TEST_DIR, { recursive: true }); + process.env.OPENCODEX_HOME = TEST_DIR; + saveConfig(config("127.0.0.1")); + stubModelDiscoveryFor("http://127.0.0.1:8080"); + + const server = startServer(0); + try { + const valid = await fetch(new URL("/api/providers", server.url), { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + name: "snapshot-valid", + provider: { + adapter: "openai-responses", + baseUrl: "http://127.0.0.1:8080/v1", + allowPrivateNetwork: true, + responsesSnapshotRepair: true, + }, + }), + }); + expect(valid.status).toBe(200); + expect(loadConfig().providers["snapshot-valid"]?.responsesSnapshotRepair).toBe(true); + + const invalid = await fetch(new URL("/api/providers", server.url), { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + name: "snapshot-invalid", + provider: { + adapter: "openai-responses", + baseUrl: "http://127.0.0.1:8080/v1", + allowPrivateNetwork: true, + responsesSnapshotRepair: { enabled: true }, + }, + }), + }); + expect(invalid.status).toBe(400); + expect(await invalid.json()).toMatchObject({ + error: "provider snapshot-invalid responsesSnapshotRepair must be a boolean", + }); + expect(loadConfig().providers["snapshot-invalid"]).toBeUndefined(); + } finally { + await server.stop(true); + } + }); + test("provider management rejects sensitive or injectable provider headers", async () => { if (existsSync(TEST_DIR)) rmSync(TEST_DIR, { recursive: true }); mkdirSync(TEST_DIR, { recursive: true }); diff --git a/tests/responses-snapshot-repair.test.ts b/tests/responses-snapshot-repair.test.ts new file mode 100644 index 000000000..9b2f2d932 --- /dev/null +++ b/tests/responses-snapshot-repair.test.ts @@ -0,0 +1,636 @@ +import { describe, expect, test } from "bun:test"; +import { createTranslatorBudget } from "../src/lib/translator-budget"; +import { isEagerRelaySseResponse, MAX_COMPLETED_OUTPUT_ITEMS } from "../src/server/relay"; +import { handleResponses } from "../src/server/responses"; +import { + createResponsesSnapshotPayloadRewrite, + hasResponsesSnapshotRepair, + repairResponsesSnapshotJson, +} from "../src/server/responses-snapshot-repair"; +import type { OcxConfig } from "../src/types"; + +function parseSseEvents(text: string): Array> { + return text.split(/\r?\n\r?\n/) + .map(block => block.split(/\r?\n/).find(line => line.startsWith("data:"))) + .map(line => line?.slice(5).trim()) + .filter((payload): payload is string => !!payload && payload !== "[DONE]") + .map(payload => JSON.parse(payload) as Record); +} + +describe("Responses sparse-snapshot repair", () => { + test("backfills canonical lifecycle fields from request metadata", () => { + const request = { + parallel_tool_calls: false, + tool_choice: { type: "function", name: "lookup" }, + tools: [{ type: "function", name: "lookup", parameters: {} }], + }; + const rewrite = createResponsesSnapshotPayloadRewrite(request); + + for (const type of ["response.created", "response.in_progress"] as const) { + const event = JSON.parse(rewrite(JSON.stringify({ + type, + response: { id: "resp_sparse", object: "response" }, + }))) as { response: Record }; + + expect(event.response).toMatchObject({ + id: "resp_sparse", + status: "in_progress", + output: [], + parallel_tool_calls: false, + tool_choice: request.tool_choice, + tools: request.tools, + }); + } + }); + + test("repairs sparse output items, parts, text events, and terminal snapshots", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + const cases = [ + { + input: { + type: "response.output_item.added", + item: { id: "rs_1", type: "reasoning", summary: [{ type: "summary_text" }] }, + }, + expected: { + item: { status: "in_progress", summary: [{ type: "summary_text", text: "" }] }, + }, + }, + { + input: { type: "response.output_item.done", item: { id: "msg_1", type: "message" } }, + expected: { item: { status: "completed", role: "assistant", content: [] } }, + }, + { + input: { type: "response.reasoning_summary_part.added", part: { type: "summary_text" } }, + expected: { part: { type: "summary_text", text: "" } }, + }, + { + input: { type: "response.content_part.done", part: { type: "output_text" } }, + expected: { part: { type: "output_text", text: "", annotations: [] } }, + }, + { + input: { type: "response.output_text.delta", delta: "ok" }, + expected: { logprobs: [] }, + }, + { + input: { type: "response.output_text.done" }, + expected: { text: "", logprobs: [] }, + }, + { + input: { + type: "response.completed", + response: { + id: "resp_terminal", + object: "response", + output: [{ + id: "msg_1", + type: "message", + content: [{ type: "output_text", text: "ok" }], + }], + }, + }, + expected: { + response: { + status: "completed", + parallel_tool_calls: true, + tool_choice: "auto", + tools: [], + output: [{ + status: "completed", + role: "assistant", + content: [{ type: "output_text", text: "ok", annotations: [] }], + }], + }, + }, + }, + ]; + + for (const { input, expected } of cases) { + expect(JSON.parse(rewrite(JSON.stringify(input)))).toMatchObject(expected); + } + }); + + test("replaces structurally invalid statuses, roles, and tool choices", () => { + for (const invalidStatus of [undefined, null, 42, ""]) { + const rewrite = createResponsesSnapshotPayloadRewrite(); + const event = JSON.parse(rewrite(JSON.stringify({ + type: "response.completed", + response: { + id: "resp_invalid", + object: "response", + ...(invalidStatus === undefined ? {} : { status: invalidStatus }), + output: [{ id: "msg_invalid", type: "message", role: "user" }], + }, + }))) as { response: { status?: unknown; output: Array> } }; + expect(event.response.status).toBe("completed"); + expect(event.response.output[0]?.role).toBe("assistant"); + } + + const requestChoice = { type: "function", name: "lookup" }; + for (const invalidChoice of [undefined, null, 42, "", " ", [], {}, { type: "" }, { type: 42 }]) { + const rewrite = createResponsesSnapshotPayloadRewrite({ tool_choice: requestChoice }); + const event = JSON.parse(rewrite(JSON.stringify({ + type: "response.created", + response: { + id: "resp_choice", + object: "response", + ...(invalidChoice === undefined ? {} : { tool_choice: invalidChoice }), + }, + }))) as { response: Record }; + expect(event.response.tool_choice).toEqual(requestChoice); + } + }); + + test("preserves valid upstream and future-compatible values", () => { + for (const status of ["queued", "in_progress", "completed", "failed", "future_status"]) { + const rewrite = createResponsesSnapshotPayloadRewrite(); + const event = JSON.parse(rewrite(JSON.stringify({ + type: "response.created", + response: { + id: "resp_status", + object: "response", + status, + output: [], + parallel_tool_calls: false, + tool_choice: { type: "future_choice" }, + tools: [], + }, + }))) as { response: Record }; + expect(event.response.status).toBe(status); + expect(event.response.tool_choice).toEqual({ type: "future_choice" }); + } + + const valid = JSON.stringify({ + type: "response.created", + response: { + id: "resp_complete", + object: "response", + status: "queued", + output: [], + parallel_tool_calls: false, + tool_choice: "none", + tools: [{ type: "function", name: "lookup" }], + }, + }); + const unrelated = JSON.stringify({ type: "response.reasoning_summary_text.delta", delta: "thinking" }); + const rewrite = createResponsesSnapshotPayloadRewrite(); + expect(rewrite(valid)).toBe(valid); + expect(rewrite(unrelated)).toBe(unrelated); + expect(rewrite("{not-json}")).toBe("{not-json}"); + }); + + test("does not invent nested statuses for in-progress snapshots", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + const event = JSON.parse(rewrite(JSON.stringify({ + type: "response.in_progress", + response: { + id: "resp_mixed", + object: "response", + output: [ + { id: "msg_done", type: "message", role: "assistant", status: "completed" }, + { id: "msg_active", type: "message", role: "assistant" }, + ], + }, + }))) as { response: { output: Array> } }; + + expect(event.response.output[0]?.status).toBe("completed"); + expect(event.response.output[1]?.status).toBeUndefined(); + }); + + test("reconstructs missing terminal output only from contiguous done indexes", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + rewrite(JSON.stringify({ + type: "response.output_item.added", + output_index: 0, + item: { id: "msg_0", type: "message" }, + })); + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: { + id: "msg_0", + type: "message", + content: [{ type: "output_text", text: "hello" }], + }, + })); + const reconstructed = JSON.parse(rewrite(JSON.stringify({ + type: "response.completed", + response: { id: "resp_reconstructed", object: "response" }, + }))) as { response: { output?: unknown } }; + + expect(reconstructed.response.output).toEqual([{ + id: "msg_0", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "hello", annotations: [] }], + }]); + + const gapped = createResponsesSnapshotPayloadRewrite(); + for (const outputIndex of [0, 2]) { + gapped(JSON.stringify({ + type: "response.output_item.done", + output_index: outputIndex, + item: { id: `msg_${outputIndex}`, type: "message" }, + })); + } + const terminal = JSON.parse(gapped(JSON.stringify({ + type: "response.completed", + response: { id: "resp_gapped", object: "response" }, + }))) as { response: { output?: unknown } }; + expect(terminal.response.output).toEqual([]); + }); + + test("preserves empty terminal output and repairs malformed values without reconstruction", () => { + for (const output of [[], null, "invalid"]) { + const rewrite = createResponsesSnapshotPayloadRewrite(); + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: { + id: "msg_explicit_output", + type: "message", + content: [{ type: "output_text", text: "must-not-be-reconstructed" }], + }, + })); + + const terminal = JSON.parse(rewrite(JSON.stringify({ + type: "response.completed", + response: { id: "resp_explicit_output", object: "response", output }, + }))) as { response: { output?: unknown } }; + + expect(terminal.response.output).toEqual([]); + } + }); + + test("suppresses reconstruction while an observed added item remains unfinished", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: { id: "msg_done", type: "message", content: [{ type: "output_text", text: "partial" }] }, + })); + rewrite(JSON.stringify({ + type: "response.output_item.added", + output_index: 1, + item: { id: "msg_pending", type: "message" }, + })); + + const terminal = JSON.parse(rewrite(JSON.stringify({ + type: "response.completed", + response: { id: "resp_pending_item", object: "response" }, + }))) as { response: { output?: unknown } }; + + expect(terminal.response.output).toEqual([]); + }); + + test("suppresses reconstruction after malformed done events", () => { + const malformedDoneEvents = [ + { type: "response.output_item.done", output_index: 1 }, + { type: "response.output_item.done", output_index: 1, item: "invalid" }, + { type: "response.output_item.done", output_index: "1", item: { type: "message" } }, + { type: "response.output_item.done", output_index: 1, item: { id: "msg_1" } }, + { type: "response.output_item.done", output_index: 1, item: { id: "msg_1", type: "" } }, + { type: "response.output_item.done", output_index: 1, item: { id: "msg_1", type: " " } }, + ]; + + for (const malformedDoneEvent of malformedDoneEvents) { + const rewrite = createResponsesSnapshotPayloadRewrite(); + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: { id: "msg_0", type: "message" }, + })); + rewrite(JSON.stringify(malformedDoneEvent)); + + const terminal = JSON.parse(rewrite(JSON.stringify({ + type: "response.completed", + response: { id: "resp_malformed_done", object: "response" }, + }))) as { response: { output?: unknown } }; + + expect(terminal.response.output).toEqual([]); + } + }); + + test("reconstructs contiguous done items for incomplete terminal responses", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: { id: "msg_partial", type: "message", content: [{ type: "output_text", text: "partial" }] }, + })); + + const terminal = JSON.parse(rewrite(JSON.stringify({ + type: "response.incomplete", + response: { id: "resp_incomplete", object: "response" }, + }))) as { response: { output?: unknown } }; + + expect(terminal.response.output).toEqual([{ + id: "msg_partial", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "partial", annotations: [] }], + }]); + }); + + test("suppresses partial reconstruction after retained-item bounds are exceeded", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + for (let outputIndex = 0; outputIndex <= MAX_COMPLETED_OUTPUT_ITEMS; outputIndex += 1) { + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: outputIndex, + item: { id: `msg_${outputIndex}`, type: "message" }, + })); + } + const terminal = JSON.parse(rewrite(JSON.stringify({ + type: "response.completed", + response: { id: "resp_capped", object: "response" }, + }))) as { response: { output?: unknown } }; + expect(terminal.response.output).toEqual([]); + }); + + test("charges reconstructed items to the translator budget and releases them at terminal", () => { + for (const terminalType of ["response.completed", "response.failed", "response.incomplete"] as const) { + const budget = createTranslatorBudget(); + const rewrite = createResponsesSnapshotPayloadRewrite(undefined, budget); + rewrite(JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: { + id: "msg_budget", + type: "message", + content: [{ type: "output_text", text: "x".repeat(4096) }], + }, + })); + expect(budget.snapshot().currentBytes).toBeGreaterThan(4096); + rewrite(JSON.stringify({ + type: terminalType, + response: { id: "resp_budget", object: "response" }, + })); + expect(budget.snapshot().currentBytes).toBe(0); + budget.dispose(); + } + }); + + test("ignores inherited lifecycle names and requires explicit provider opt-in", () => { + const rewrite = createResponsesSnapshotPayloadRewrite(); + const inherited = JSON.stringify({ type: "__proto__", response: { id: "resp_proto" } }); + expect(rewrite(inherited)).toBe(inherited); + expect(hasResponsesSnapshotRepair(undefined)).toBe(false); + expect(hasResponsesSnapshotRepair(false)).toBe(false); + expect(hasResponsesSnapshotRepair(true)).toBe(true); + }); + + test("repairs non-streaming JSON without changing malformed payloads", () => { + const repaired = JSON.parse(repairResponsesSnapshotJson(JSON.stringify({ + id: "resp_json", + object: "response", + output: [{ id: "msg_json", type: "message", content: [{ type: "output_text" }] }], + }))) as Record; + expect(repaired).toMatchObject({ + status: "completed", + parallel_tool_calls: true, + tool_choice: "auto", + tools: [], + output: [{ + id: "msg_json", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "", annotations: [] }], + }], + }); + expect(repairResponsesSnapshotJson("{not-json}")).toBe("{not-json}"); + }); + + test("does not complete sparse output items in an incomplete non-streaming response", () => { + const repaired = JSON.parse(repairResponsesSnapshotJson(JSON.stringify({ + id: "resp_incomplete_json", + object: "response", + status: "incomplete", + output: [{ id: "msg_partial", type: "message" }], + }))) as { status?: unknown; output: Array> }; + + expect(repaired.status).toBe("incomplete"); + expect(repaired.output[0]).toEqual({ + id: "msg_partial", + type: "message", + role: "assistant", + status: "incomplete", + content: [], + }); + }); + + test("handleResponses keeps SSE repair default-off and preserves the existing relay gate", async () => { + const savedFetch = globalThis.fetch; + const doneItem = { + id: "msg_1", + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "hello" }], + }; + const upstream = [ + `event: response.output_item.done\ndata: ${JSON.stringify({ + type: "response.output_item.done", + output_index: 0, + item: doneItem, + })}\n\n`, + `event: response.completed\ndata: ${JSON.stringify({ + type: "response.completed", + response: { id: "resp_1", object: "response", status: "completed" }, + })}\n\n`, + "data: [DONE]\n\n", + ].join(""); + globalThis.fetch = (async () => new Response(upstream, { + headers: { "content-type": "text/event-stream" }, + })) as typeof fetch; + + try { + for (const enabled of [false, true]) { + const config = { + port: 0, + streamMode: "eager-relay", + defaultProvider: "fixture", + providers: { + fixture: { + adapter: "openai-responses", + baseUrl: "https://fixture.test/v1", + authMode: "key", + apiKey: "fixture-key", + ...(enabled ? { responsesSnapshotRepair: true } : {}), + }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + model: "fixture/example-model", + stream: true, + input: "hello", + parallel_tool_calls: false, + tool_choice: { type: "function", name: "lookup" }, + tools: [{ type: "function", name: "lookup", parameters: {} }], + }), + }), config, { model: "", provider: "" }); + + const expectedEager = process.platform === "win32" + || (process.platform === "darwin" && !enabled); + expect(isEagerRelaySseResponse(response)).toBe(expectedEager); + const events = parseSseEvents(await response.text()); + const terminal = events.find(event => event.type === "response.completed") as { + response: Record; + } | undefined; + if (!enabled) { + expect(terminal?.response).toEqual({ id: "resp_1", object: "response", status: "completed" }); + continue; + } + expect(terminal?.response).toMatchObject({ + parallel_tool_calls: false, + tool_choice: { type: "function", name: "lookup" }, + tools: [{ type: "function", name: "lookup", parameters: {} }], + output: [{ + ...doneItem, + status: "completed", + content: [{ type: "output_text", text: "hello", annotations: [] }], + }], + }); + } + } finally { + globalThis.fetch = savedFetch; + } + }); + + test("handleResponses canonicalizes the exact sparse added-delta-completed stream from issue #893", async () => { + const savedFetch = globalThis.fetch; + const upstream = [ + 'event: response.created\ndata: {"type":"response.created","response":{"id":"resp_example","object":"response"}}\n\n', + 'event: response.output_item.added\ndata: {"type":"response.output_item.added","item":{"id":"msg_example","type":"message"}}\n\n', + 'event: response.output_text.delta\ndata: {"type":"response.output_text.delta","delta":"hello"}\n\n', + 'event: response.completed\ndata: {"type":"response.completed","response":{"id":"resp_example","object":"response"}}\n\n', + "data: [DONE]\n\n", + ].join(""); + globalThis.fetch = (async () => new Response(upstream, { + headers: { "content-type": "text/event-stream" }, + })) as typeof fetch; + + try { + const config = { + port: 0, + defaultProvider: "fixture", + providers: { + fixture: { + adapter: "openai-responses", + baseUrl: "https://fixture.test/v1", + authMode: "key", + apiKey: "fixture-key", + responsesSnapshotRepair: true, + }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "fixture/example-model", stream: true, input: "hello" }), + }), config, { model: "", provider: "" }); + const events = parseSseEvents(await response.text()); + + // The issue fixture has no output_item.done, so the safe contract is to preserve its text + // delta and emit a canonical empty terminal output, not invent a completed message item. + expect(events).toEqual([ + { + type: "response.created", + response: { + id: "resp_example", + object: "response", + status: "in_progress", + output: [], + parallel_tool_calls: true, + tool_choice: "auto", + tools: [], + }, + }, + { + type: "response.output_item.added", + item: { + id: "msg_example", + type: "message", + status: "in_progress", + role: "assistant", + content: [], + }, + }, + { type: "response.output_text.delta", delta: "hello", logprobs: [] }, + { + type: "response.completed", + response: { + id: "resp_example", + object: "response", + status: "completed", + output: [], + parallel_tool_calls: true, + tool_choice: "auto", + tools: [], + }, + }, + ]); + } finally { + globalThis.fetch = savedFetch; + } + }); + + test("handleResponses repairs non-streaming JSON only for an opted-in provider", async () => { + const savedFetch = globalThis.fetch; + globalThis.fetch = (async () => new Response(JSON.stringify({ + id: "resp_json", + object: "response", + output: [{ id: "msg_json", type: "message", content: [{ type: "output_text" }] }], + }), { headers: { "content-type": "application/json" } })) as typeof fetch; + + try { + for (const enabled of [false, true]) { + const config = { + port: 0, + defaultProvider: "fixture", + providers: { + fixture: { + adapter: "openai-responses", + baseUrl: "https://fixture.test/v1", + authMode: "key", + apiKey: "fixture-key", + ...(enabled ? { responsesSnapshotRepair: true } : {}), + }, + }, + } as OcxConfig; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "fixture/example-model", input: "hello" }), + }), config, { model: "", provider: "" }); + const body = await response.json() as { + status?: unknown; + output: Array>; + }; + if (!enabled) { + expect(body.status).toBeUndefined(); + expect(body.output[0]).toEqual({ + id: "msg_json", + type: "message", + content: [{ type: "output_text" }], + }); + continue; + } + expect(body.status).toBe("completed"); + expect(body.output[0]).toEqual({ + id: "msg_json", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "", annotations: [] }], + }); + } + } finally { + globalThis.fetch = savedFetch; + } + }); +});