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 a99df988c..21a93f218 100644 --- a/docs-site/src/content/docs/ja/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ja/reference/configuration/providers.md @@ -84,6 +84,7 @@ account を削除しても mapping は保持され、同じ id を再追加す | `modelSupportsReasoningSummaries?` | `Record` |モデルを `false` に設定して、概要の広告を停止し、概要配信フィールドを削除します。 | | `modelReasoningSummaryDelivery?` | `Record` |モデルごとの応答配信列挙型。既存の配信フィールドを書き換えます。 | | `modelAdapters?` | `Record` | 混合配線ゲートウェイのモデルごとの `openai-chat` または `openai-responses` 配線オーバーライド。明示的なエントリはレジストリのデフォルトを破ります。DeepSeek のプリセットは `deepseek-v4-flash` のネイティブ Responses を選択でき、GitHub Copilot は GPT-5 ファミリー (`gpt-5.3-codex`, `gpt-5.4`, `gpt-5.4-mini`, `gpt-5.5`, `gpt-5.6-luna`, `gpt-5.6-sol`, `gpt-5.6-terra`) を Responses 専用デフォルトとして宣言します。これらのモデルはエージェント トラフィックで `/chat/completions` を拒否するためです。`gpt-5.4-nano` のようなビルトイン デフォルトのないモデルはここでオプトインできます。単線アップストリーム ピンと正規の ChatGPT 転送はオーバーライドを拒否します。 | +| `modelResponsesUpstreamStreaming?` | `Record` | forward 以外の `openai-responses` プロバイダー向けモデル別 upstream Responses ポリシーです。`false` は upstream に bounded JSON を要求し、検証済み terminal オブジェクトを streaming client 用 Responses イベントへ再構成します。`true` は registry の `false` 既定値を明示的に上書きします。照合は大文字小文字を区別せず、public virtual id を優先し、最終 wire-model id をフォールバックに使います。この correctness-first fallback では incremental delta がなくなり、bounded JSON のサイズと timeout 制限が適用されます。 | | `modelPreferHostedTools?` | `Record` | hosted tool namespace を予約する非 forward Responses gateway 向けの完全一致モデル opt-in。現在は `["image_generation"]` のみを受け付けます。一致したモデルは `openai-responses` wire を使い、その hosted tool をサポートする必要があります。競合するクライアント `image_gen` 宣言を除去し、呼び出し元の tool choice を維持するため selector も書き換えます。OpenAI API の仮想 `-pro` モデルでは、まず選択した公開 ID に一致させ、解決後のベース wire-model ID をフォールバックとして使用します。`modelAdapters` は公開 ID、次にベース ID の順に解決し、後者の結果が最終 wire を決めます。未設定のモデルは通常の alias 動作を維持します。 | | `reasoningEffortMap?` | `Record` |ラベルを推論するためのプロバイダー全体のワイヤ エイリアス。 | | `modelReasoningEffortMap?` | `Record>` |推論ラベルのモデルごとのワイヤ エイリアス。 | 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 b3a7cdb62..ae45c3afc 100644 --- a/docs-site/src/content/docs/ko/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ko/reference/configuration/providers.md @@ -84,6 +84,7 @@ managed map을 활성화하면 privacy-safe selector를 만들고, 이후 계정 | `modelSupportsReasoningSummaries?` | `Record` | 모델을 `false`로 두면 summary 광고를 멈추고 summary 전달 필드를 제거합니다. | | `modelReasoningSummaryDelivery?` | `Record` | 모델별 Responses 전달 enum입니다. 기존 delivery 필드를 다시 씁니다. | | `modelAdapters?` | `Record` | 혼합 와이어 게이트웨이를 위한 모델별 `openai-chat` 또는 `openai-responses` 와이어 재정의입니다. 명시적 항목이 레지스트리 기본값보다 우선합니다. DeepSeek 프리셋은 `deepseek-v4-flash`에 네이티브 Responses를 선택할 수 있고, GitHub Copilot은 GPT-5 계열(`gpt-5.3-codex`, `gpt-5.4`, `gpt-5.4-mini`, `gpt-5.5`, `gpt-5.6-luna`, `gpt-5.6-sol`, `gpt-5.6-terra`)을 Responses 전용 기본값으로 선언합니다. 이 모델들은 에이전트 트래픽에서 `/chat/completions`를 거부하기 때문입니다. `gpt-5.4-nano`처럼 기본값이 없는 모델은 여기서 직접 옵트인할 수 있습니다. 단일 와이어 상위 항목과 정식 ChatGPT forward는 재정의를 거부합니다. | +| `modelResponsesUpstreamStreaming?` | `Record` | Forward가 아닌 `openai-responses` provider의 모델별 upstream Responses 정책입니다. `false`는 upstream에 bounded JSON을 요청하고 검증된 terminal 객체를 streaming client용 Responses event로 다시 구성합니다. `true`는 registry의 `false` 기본값을 명시적으로 해제합니다. 대소문자를 구분하지 않으며 public virtual id가 먼저, 최종 wire-model id가 fallback으로 일치합니다. 이 correctness-first fallback에서는 incremental delta가 사라지고 bounded JSON 크기·시간 제한이 적용됩니다. | | `modelPreferHostedTools?` | `Record` | hosted tool namespace를 예약하는 non-forward Responses gateway용 정확한 모델 ID opt-in입니다. 현재 `["image_generation"]`만 허용하며, 일치하는 모델은 `openai-responses` wire를 사용하고 해당 hosted tool을 지원해야 합니다. 충돌하는 클라이언트 `image_gen` 선언을 제거하고 호출자의 tool choice를 유지하도록 selector도 다시 씁니다. OpenAI API 가상 `-pro` 모델은 선택한 공개 ID를 먼저 일치시키고, 해석된 기본 wire-model ID를 대체값으로 사용합니다. `modelAdapters`는 공개 ID를 먼저, 그 다음 기본 ID를 해석하며, 두 번째 결과가 최종 wire를 결정합니다. 설정하지 않은 모델은 일반 alias 동작을 유지합니다. | | `reasoningEffortMap?` | `Record` | reasoning 레이블의 공급자 전반 와이어 별칭입니다. | | `modelReasoningEffortMap?` | `Record>` | reasoning 레이블의 모델별 와이어 별칭입니다. | diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index dd5ca0b52..173efe0e2 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -94,6 +94,7 @@ differing backup and rewrites known legacy namespaced selected ids to bare ids. | `modelSupportsReasoningSummaries?` | `Record` | Set a model to `false` to stop advertising summaries and strip summary-delivery fields. | | `modelReasoningSummaryDelivery?` | `Record` | Per-model Responses delivery enum; rewrites an existing delivery field. | | `modelAdapters?` | `Record` | Per-model `openai-chat` or `openai-responses` wire override for mixed-wire gateways. Explicit entries beat registry defaults; DeepSeek's preset can select native Responses for `deepseek-v4-flash`, and GitHub Copilot declares Responses-only defaults for its GPT-5 family (`gpt-5.3-codex`, `gpt-5.4`, `gpt-5.4-mini`, `gpt-5.5`, `gpt-5.6-luna`, `gpt-5.6-sol`, `gpt-5.6-terra`) because those models reject `/chat/completions` for agent traffic. Models without a built-in default (for example `gpt-5.4-nano`) can be opted in here. Single-wire upstream pins and canonical ChatGPT forward reject overrides. | +| `modelResponsesUpstreamStreaming?` | `Record` | Per-model Responses upstream policy for non-forward `openai-responses` providers. `false` asks upstream for bounded JSON and reframes the validated terminal object as Responses events for a streaming client; `true` explicitly overrides a registry `false` default. Matching is case-insensitive and supports public virtual ids with their final wire-model id as fallback. This correctness-first fallback removes incremental deltas and remains subject to the bounded JSON size and timeout limits. | | `modelPreferHostedTools?` | `Record` | Exact-model opt-in for non-forward Responses gateways that reserve a hosted-tool namespace. Currently accepts only `["image_generation"]`; a matching model must use the `openai-responses` wire and support that hosted tool. It removes colliding client `image_gen` declarations and rewrites their selectors to preserve caller tool choice. For OpenAI API virtual `-pro` models, the selected public ID is matched first and the resolved base wire-model ID is a fallback. `modelAdapters` resolves the public ID first, then the base ID; the second resolution determines the final wire. Other models retain normal alias behavior. | | `reasoningEffortMap?` | `Record` | Provider-wide wire aliases for reasoning labels. | | `modelReasoningEffortMap?` | `Record>` | Per-model wire aliases for reasoning labels. | 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 7e55625e1..e2a179f1a 100644 --- a/docs-site/src/content/docs/ru/reference/configuration/providers.md +++ b/docs-site/src/content/docs/ru/reference/configuration/providers.md @@ -97,6 +97,7 @@ cross-route credential fallback не существует. Строки API GPT- | `modelSupportsReasoningSummaries?` | `Record` | Установите `false` для модели, чтобы перестать рекламировать summary и вырезать поля доставки summary. | | `modelReasoningSummaryDelivery?` | `Record` | Responses delivery enum по моделям; переписывает уже существующее поле delivery. | | `modelAdapters?` | `Record` | Wire-override по модели для `openai-chat` или `openai-responses` в gateway с несколькими wire-форматами. Явные записи имеют приоритет над default'ами registry; preset DeepSeek может выбирать native Responses для `deepseek-v4-flash`, а GitHub Copilot объявляет Responses-only default'ы для семейства GPT-5 (`gpt-5.3-codex`, `gpt-5.4`, `gpt-5.4-mini`, `gpt-5.5`, `gpt-5.6-luna`, `gpt-5.6-sol`, `gpt-5.6-terra`), потому что эти модели отклоняют `/chat/completions` для агентного трафика. Модели без встроенного default'а (например, `gpt-5.4-nano`) можно включить здесь. Single-wire upstream pin'ы и canonical ChatGPT forward override не принимают. | +| `modelResponsesUpstreamStreaming?` | `Record` | Модельная политика upstream Responses для non-forward провайдеров `openai-responses`. `false` запрашивает bounded JSON и преобразует проверенный terminal-объект в события Responses для streaming client; `true` явно перекрывает registry default `false`. Сопоставление не зависит от регистра: сначала используется public virtual id, затем итоговый wire-model id. Этот correctness-first fallback убирает incremental delta и подчиняется лимитам размера и timeout bounded JSON. | | `modelPreferHostedTools?` | `Record` | Opt-in для точного model ID в non-forward Responses gateway, который резервирует namespace hosted tool. Сейчас допускается только `["image_generation"]`; совпавшая модель должна использовать wire `openai-responses` и поддерживать этот hosted tool. Прокси удаляет конфликтующие клиентские объявления `image_gen` и переписывает их selectors, сохраняя caller tool choice. Для виртуальных моделей OpenAI API `-pro` сначала сопоставляется выбранный публичный ID, а затем в качестве fallback используется ID базовой wire-модели. `modelAdapters` сначала разрешается по публичному ID, затем по базовому ID; второй результат определяет итоговый wire. Остальные модели сохраняют обычное alias-поведение. | | `reasoningEffortMap?` | `Record` | Provider-wide wire-alias'ы для reasoning-label'ов. | | `modelReasoningEffortMap?` | `Record>` | Wire-alias'ы для reasoning-label'ов по отдельным моделям. | 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 a8855ebdb..31f18f493 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 @@ -84,6 +84,7 @@ selector,而不是分配一个新名称。 | `modelSupportsReasoningSummaries?` | `Record` | 将某个模型设为 `false`,即可停止暴露摘要并移除摘要交付字段。 | | `modelReasoningSummaryDelivery?` | `Record` | 按模型设置的 Responses 交付枚举;会重写现有的 delivery 字段。 | | `modelAdapters?` | `Record` | 按模型设置的 `openai-chat` 或 `openai-responses` 线协议覆盖项,用于混合线协议网关。显式条目优先于注册表默认值;DeepSeek 预设可以为 `deepseek-v4-flash` 选择原生 Responses,GitHub Copilot 则为 GPT-5 系列(`gpt-5.3-codex`、`gpt-5.4`、`gpt-5.4-mini`、`gpt-5.5`、`gpt-5.6-luna`、`gpt-5.6-sol`、`gpt-5.6-terra`)声明了 Responses 专用默认值,因为这些模型在代理流量下会拒绝 `/chat/completions`。没有内置默认值的模型(例如 `gpt-5.4-nano`)可以在此手动启用。单一线协议上游固定项和规范 ChatGPT forward 会拒绝覆盖。 | +| `modelResponsesUpstreamStreaming?` | `Record` | 非 forward `openai-responses` provider 的逐模型 upstream Responses 策略。`false` 请求 bounded JSON,并将验证后的 terminal 对象重构为 streaming client 所需的 Responses 事件;`true` 显式覆盖 registry 的 `false` 默认值。匹配不区分大小写,优先使用 public virtual id,并以最终 wire-model id 作为回退。此 correctness-first fallback 不提供 incremental delta,并受 bounded JSON 大小与 timeout 限制。 | | `modelPreferHostedTools?` | `Record` | 非 forward Responses gateway 的精确模型 ID opt-in,用于上游预留 hosted tool namespace 的情况。目前只支持 `["image_generation"]`;匹配模型必须使用 `openai-responses` wire 且支持该 hosted 工具。它会移除冲突的客户端 `image_gen` 声明,并改写其 selector 以保持调用方的 tool choice。对于 OpenAI API 的虚拟 `-pro` 模型,先匹配所选公开 ID,未命中时才使用解析出的基础 wire-model ID 作为回退。`modelAdapters` 会先按公开 ID、再按基础 ID 解析;后一次结果决定最终 wire。未配置模型保持普通 alias 行为。 | | `reasoningEffortMap?` | `Record` | 提供者级、用于推理标签的线协议别名。 | | `modelReasoningEffortMap?` | `Record>` | 按模型设置的推理标签线协议别名。 | diff --git a/src/config.ts b/src/config.ts index 32831fd39..d1e53340b 100644 --- a/src/config.ts +++ b/src/config.ts @@ -616,6 +616,7 @@ const providerConfigSchema = z.object({ requiresAdjacentResponsesToolResults: z.boolean().optional(), supportsServiceTier: z.boolean().optional(), preserveResponsesReasoningContent: z.boolean().optional(), + modelResponsesUpstreamStreaming: z.record(z.string(), z.boolean()).optional(), allowPrivateNetwork: z.boolean().optional(), retryOn429: retryOn429PolicySchema.optional(), codexAccountMode: z.enum(["pool", "direct"]).optional(), @@ -908,6 +909,83 @@ export function modelAdapterRecordConfigError( return null; } +/** Validate the opt-in bounded-JSON policy against the model's effective Responses wire. */ +export function modelResponsesUpstreamStreamingConfigError( + value: unknown, + field: string, + providerName: string, + provider: { adapter?: unknown; authMode?: unknown; baseUrl?: unknown; modelAdapters?: unknown }, +): string | null { + const shapeError = booleanRecordConfigError(value, field); + if (shapeError) return shapeError; + const entries = Object.entries((value ?? {}) as Record); + if (entries.length === 0) return null; + + const registry = getProviderRegistryEntry(providerName); + const registryTransportMatches = typeof provider.baseUrl === "string" + && providerMatchesRegistryTransport(providerName, { + baseUrl: provider.baseUrl, + adapter: provider.adapter as OcxProviderConfig["adapter"], + ...(typeof provider.authMode === "string" + ? { authMode: provider.authMode as OcxProviderConfig["authMode"] } + : {}), + }); + const effectiveForwardAuth = registryTransportMatches + ? registry?.authKind === "forward" + : provider.authMode === "forward"; + if (effectiveForwardAuth) { + return `${field} is not supported on forward-auth Responses providers`; + } + + const resolveEffectiveWire = (modelId: string, currentWire: unknown): unknown => { + const pinned = pinnedWireAdapter(providerName, modelId); + if (pinned) return pinned; + const configured = provider.modelAdapters && typeof provider.modelAdapters === "object" + && !Array.isArray(provider.modelAdapters) + ? (provider.modelAdapters as Record)[modelId] + : undefined; + if (typeof configured === "string" && MODEL_ADAPTER_OVERRIDE_ALLOWED.has(configured)) { + return configured; + } + const registryDefault = typeof currentWire === "string" && typeof provider.baseUrl === "string" + ? providerModelWireDefault( + providerName, + { + baseUrl: provider.baseUrl, + adapter: currentWire, + ...(typeof provider.authMode === "string" + ? { authMode: provider.authMode as OcxProviderConfig["authMode"] } + : {}), + }, + modelId, + MODEL_ADAPTER_OVERRIDE_ALLOWED, + "responses", + ) + : undefined; + return registryDefault ?? currentWire; + }; + + for (const [modelId] of entries) { + const baseWire = registryTransportMatches ? registry?.adapter ?? provider.adapter : provider.adapter; + const virtualSelectedModelId = Object.keys(registry?.virtualModels ?? {}).find( + candidate => candidate.toLowerCase() === modelId.trim().toLowerCase(), + ); + const effectiveSelectedModelId = virtualSelectedModelId ?? modelId; + let effectiveWire = resolveEffectiveWire(effectiveSelectedModelId, baseWire); + const virtualWireModel = resolveOpenAiVirtualModel( + providerName, + effectiveSelectedModelId, + )?.wireModelId; + if (virtualWireModel && virtualWireModel !== effectiveSelectedModelId) { + effectiveWire = resolveEffectiveWire(virtualWireModel, effectiveWire); + } + if (effectiveWire !== "openai-responses") { + return `${field}.${modelId} requires the openai-responses wire`; + } + } + return null; +} + const CODEX_ACCOUNT_NAMESPACES_RECORD_ERROR = "codexAccountNamespaces must be a plain object mapping account selectors to Codex account ids"; const CODEX_ACCOUNT_NAMESPACE_KEY_ERROR = @@ -1249,6 +1327,19 @@ const configSchema = z.object({ message: modelAdaptersError, }); } + const responsesUpstreamStreamingError = modelResponsesUpstreamStreamingConfigError( + (provider as { modelResponsesUpstreamStreaming?: unknown }).modelResponsesUpstreamStreaming, + "modelResponsesUpstreamStreaming", + name, + provider, + ); + if (responsesUpstreamStreamingError) { + ctx.addIssue({ + code: "custom", + path: ["providers", name, "modelResponsesUpstreamStreaming"], + message: responsesUpstreamStreamingError, + }); + } const preferHostedToolsError = modelPreferHostedToolsConfigError( (provider as { modelPreferHostedTools?: unknown }).modelPreferHostedTools, "modelPreferHostedTools", diff --git a/src/providers/registry.ts b/src/providers/registry.ts index 39353726e..55ac18ea2 100644 --- a/src/providers/registry.ts +++ b/src/providers/registry.ts @@ -155,11 +155,9 @@ export interface ProviderRegistryEntry { */ modelWireDefaults?: Record; /** - * Registry-only per-model override for the upstream request shape used behind a - * Codex Responses WebSocket turn. `false` keeps the client-facing WebSocket but - * asks the upstream Responses endpoint for bounded JSON, which the bridge then - * reframes as Responses events. Use only for upstreams whose streaming response - * can omit or indefinitely delay the terminal event. + * Registry default for the per-model upstream request shape used behind a Codex + * Responses transport. `false` asks the upstream endpoint for bounded JSON, which + * OpenCodex reframes as Responses events for streaming clients. */ modelResponsesUpstreamStreaming?: Record; /** @@ -2288,15 +2286,44 @@ export function providerModelWireDefault( return wire !== undefined && allowedWires.has(wire) ? wire : undefined; } -/** Resolve a registry-only upstream-streaming compatibility hint for Responses turns. */ +function responsesStreamingPolicyValue( + record: Record | undefined, + modelId: string, +): boolean | undefined { + if (!record) return undefined; + if (Object.prototype.hasOwnProperty.call(record, modelId)) return record[modelId]; + const folded = modelId.toLowerCase(); + for (const [key, value] of Object.entries(record)) { + if (key.toLowerCase() === folded) return value; + } + return undefined; +} + +/** Resolve an explicit upstream-streaming policy before the matching registry default. */ export function providerModelResponsesUpstreamStreaming( id: string, - provider: Pick & Partial>, + provider: Pick + & Partial>, modelId: string, + fallbackModelId?: string, ): boolean | undefined { const entry = getProviderRegistryEntry(id); - if (!entry?.modelResponsesUpstreamStreaming || !providerMatchesRegistryTransport(id, provider)) return undefined; - return entry.modelResponsesUpstreamStreaming[modelId.trim().toLowerCase()]; + const registryTransportMatches = !!entry && providerMatchesRegistryTransport(id, provider); + const effectiveForwardAuth = registryTransportMatches + ? entry.authKind === "forward" + : provider.authMode === "forward"; + if (!effectiveForwardAuth) { + const configured = responsesStreamingPolicyValue(provider.modelResponsesUpstreamStreaming, modelId) + ?? (fallbackModelId + ? responsesStreamingPolicyValue(provider.modelResponsesUpstreamStreaming, fallbackModelId) + : undefined); + if (configured !== undefined) return configured; + } + if (!entry?.modelResponsesUpstreamStreaming || !registryTransportMatches) return undefined; + return responsesStreamingPolicyValue(entry.modelResponsesUpstreamStreaming, modelId) + ?? (fallbackModelId + ? responsesStreamingPolicyValue(entry.modelResponsesUpstreamStreaming, fallbackModelId) + : undefined); } /** diff --git a/src/server/auth-cors.ts b/src/server/auth-cors.ts index a59bfba5f..7e96f3f9e 100644 --- a/src/server/auth-cors.ts +++ b/src/server/auth-cors.ts @@ -5,6 +5,7 @@ import { booleanRecordConfigError, modelAdapterRecordConfigError, modelPreferHostedToolsConfigError, + modelResponsesUpstreamStreamingConfigError, codexAutoStartEnabled, positiveIntegerConfigError, positiveIntegerRecordConfigError, @@ -447,6 +448,7 @@ export function providerManagementConfigError(name: unknown, provider: unknown): if (seed) seed.codexAccountMode = raw.codexAccountMode; const canonicalCandidate = { ...raw }; delete canonicalCandidate.responsesSnapshotRepair; + delete canonicalCandidate.modelResponsesUpstreamStreaming; const canonical = seed && sameCanonicalProviderSeed(canonicalCandidate, seed); if (!canonical) { return `provider ${name} must equal the canonical built-in provider seed`; @@ -484,6 +486,13 @@ export function providerManagementConfigError(name: unknown, provider: unknown): if (reasoningSummaryDeliveryError) return `provider ${name} ${reasoningSummaryDeliveryError}`; const modelAdaptersError = modelAdapterRecordConfigError(raw.modelAdapters, "modelAdapters", name, typed); if (modelAdaptersError) return `provider ${name} ${modelAdaptersError}`; + const responsesUpstreamStreamingError = modelResponsesUpstreamStreamingConfigError( + raw.modelResponsesUpstreamStreaming, + "modelResponsesUpstreamStreaming", + name, + typed, + ); + if (responsesUpstreamStreamingError) return `provider ${name} ${responsesUpstreamStreamingError}`; const preferHostedToolsError = modelPreferHostedToolsConfigError( raw.modelPreferHostedTools, "modelPreferHostedTools", diff --git a/src/server/responses-json-events.ts b/src/server/responses-json-events.ts index 6606efb72..7efb018ba 100644 --- a/src/server/responses-json-events.ts +++ b/src/server/responses-json-events.ts @@ -7,8 +7,87 @@ export type ResponsesJsonEventFrame = Record; +export type ResponsesJsonValidationResult = + | { ok: true; response: Record } + | { ok: false; message: string }; + export const MAX_SYNTHESIZED_OUTPUT_ITEMS = 10_000; +function isTokenCount(value: unknown): value is number { + return typeof value === "number" && Number.isSafeInteger(value) && value >= 0; +} + +function usageValidationError(value: unknown): string | null { + if (value === undefined || value === null) return null; + if (typeof value !== "object" || Array.isArray(value)) { + return "upstream Responses JSON usage must be an object or null"; + } + const usage = value as Record; + for (const field of ["input_tokens", "output_tokens"] as const) { + if (!isTokenCount(usage[field])) { + return "upstream Responses JSON usage token counts must be non-negative integers"; + } + } + if (usage.total_tokens !== undefined && !isTokenCount(usage.total_tokens)) { + return "upstream Responses JSON usage token counts must be non-negative integers"; + } + for (const field of ["input_tokens_details", "output_tokens_details"] as const) { + const details = usage[field]; + if (details === undefined || details === null) continue; + if (typeof details !== "object" || Array.isArray(details)) { + return "upstream Responses JSON usage details must be objects when present"; + } + } + const inputDetails = usage.input_tokens_details; + if (inputDetails && typeof inputDetails === "object" && !Array.isArray(inputDetails)) { + for (const field of ["cached_tokens", "cache_write_tokens"] as const) { + const count = (inputDetails as Record)[field]; + if (count !== undefined && !isTokenCount(count)) { + return "upstream Responses JSON usage details must contain non-negative integer token counts"; + } + } + } + const outputDetails = usage.output_tokens_details; + if (outputDetails && typeof outputDetails === "object" && !Array.isArray(outputDetails)) { + const reasoningTokens = (outputDetails as Record).reasoning_tokens; + if (reasoningTokens !== undefined && !isTokenCount(reasoningTokens)) { + return "upstream Responses JSON usage details must contain non-negative integer token counts"; + } + } + return null; +} + +/** Validate the terminal Responses object before turning one snapshot into lifecycle events. */ +export function validateResponsesJsonEventResponse(value: unknown): ResponsesJsonValidationResult { + if (!value || typeof value !== "object" || Array.isArray(value)) { + return { ok: false, message: "upstream Responses JSON must be an object" }; + } + const response = value as Record; + if (typeof response.id !== "string" || !response.id.trim()) { + return { ok: false, message: "upstream Responses JSON must include a non-empty id" }; + } + if (response.object !== undefined && response.object !== "response") { + return { ok: false, message: 'upstream Responses JSON object must be "response" when present' }; + } + if (response.status !== "completed" && response.status !== "failed" && response.status !== "incomplete") { + return { ok: false, message: "upstream Responses JSON must include a terminal status" }; + } + if (!Array.isArray(response.output)) { + return { ok: false, message: "upstream Responses JSON output must be an array" }; + } + for (const item of response.output) { + if (!item || typeof item !== "object" || Array.isArray(item)) { + return { ok: false, message: "upstream Responses JSON output items must be objects" }; + } + if (typeof (item as { type?: unknown }).type !== "string" || !(item as { type: string }).type.trim()) { + return { ok: false, message: "upstream Responses JSON output items must include a non-empty type" }; + } + } + const usageError = usageValidationError(response.usage); + if (usageError) return { ok: false, message: usageError }; + return { ok: true, response }; +} + /** * The canonical minimal sequence Codex commits: response.created (empty * output, in_progress) → one response.output_item.done per output item → a @@ -18,7 +97,9 @@ export function responsesJsonEventSequence( response: Record, rewritePayload?: (payload: Record) => Record, ): ResponsesJsonEventFrame[] { - return [...iterateResponsesJsonEvents(response, rewritePayload)]; + const validation = validateResponsesJsonEventResponse(response); + if (!validation.ok) throw new TypeError(validation.message); + return [...iterateResponsesJsonEvents(validation.response, rewritePayload)]; } function* iterateResponsesJsonEvents( @@ -32,9 +113,7 @@ function* iterateResponsesJsonEvents( `Responses JSON output contains ${output.length} items; maximum is ${MAX_SYNTHESIZED_OUTPUT_ITEMS}`, ); } - const finalStatus = response.status === "failed" || response.status === "incomplete" - ? response.status - : "completed"; + const finalStatus = response.status as "completed" | "failed" | "incomplete"; yield rewrite({ type: "response.created", response: { ...response, status: "in_progress", output: [] }, @@ -70,7 +149,10 @@ export function responsesJsonToSseStream( response: Record, rewritePayload?: (payload: Record) => Record, ): ReadableStream { - const output = Array.isArray(response.output) ? response.output : []; + const validation = validateResponsesJsonEventResponse(response); + if (!validation.ok) throw new TypeError(validation.message); + response = validation.response; + const output = response.output as unknown[]; if (output.length > MAX_SYNTHESIZED_OUTPUT_ITEMS) { throw new RangeError( `Responses JSON output contains ${output.length} items; maximum is ${MAX_SYNTHESIZED_OUTPUT_ITEMS}`, diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index 327376fc5..94d0ec3cd 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -200,7 +200,7 @@ import { relaySseWithBlockRewrite, } from "../sse-payload-rewrite"; import { createGithubCopilotResponsesBlockRewrite } from "../github-copilot-responses-repair"; -import { responsesJsonToSseStream } from "../responses-json-events"; +import { responsesJsonToSseStream, validateResponsesJsonEventResponse } from "../responses-json-events"; import { guardTerminalEventStream } from "./terminal-guard"; /** @@ -901,7 +901,7 @@ async function applyFinalRouteRequestNormalization(args: { logCtx: RequestLogContext; inboundWire: InboundWire; inboundTransport?: "websocket"; -}): Promise { +}): Promise { const { parsed, route, config, req, logCtx, inboundWire, inboundTransport } = args; // Only Anthropic message routes retain the Codex-facing selector. Other providers must keep @@ -917,40 +917,43 @@ async function applyFinalRouteRequestNormalization(args: { } parsed.modelId = route.modelId; } - // Transport-neutral reliability policy (#875): applies to any Responses - // upstream whose final adapter is openai-responses, not only WS turns. - const responsesUpstreamStreaming = providerModelResponsesUpstreamStreaming( - route.providerName, - route.provider, - route.modelId, - ); + // Final selected model before virtual wire-model rewriting (Pro aliases). + const finalSelectedModelId = route.modelId; + const policyProvider = route.provider; - // Settle the wire once so logging, fast-mode, auth, and sidecars read the adapter - // this request will actually use (#404). + // Resolve the public selector first; virtual aliases below may then resolve the base wire + // model once more. The streaming policy is captured only after both decisions settle. route.provider = resolveWireProtocolOverride(route.providerName, route.modelId, route.provider, inboundWire); if (preserveAnthropicResponseModel) parsed._responseModelId = responseModelId; logCtx.model = route.modelId; logCtx.provider = route.providerName; - logCtx.providerAdapter = route.provider.adapter; logCtx.routeDecision = route.routeDecision; - if (responsesUpstreamStreaming === false && route.provider.adapter === "openai-responses") { - parsed.stream = false; - if (parsed._rawBody && typeof parsed._rawBody === "object") { - (parsed._rawBody as Record).stream = false; - } - } - - // Final selected model before virtual wire-model rewriting (Pro aliases). - const finalSelectedModelId = route.modelId; - // Virtual model rewriting: Pro aliases → base model + reasoning.mode="pro". applyOpenAiVirtualModel(parsed, route, logCtx); + route.provider = resolveWireProtocolOverride(route.providerName, route.modelId, route.provider, inboundWire); + logCtx.providerAdapter = route.provider.adapter; if (parsed._responseModelId !== undefined && parsed._responseModelId !== parsed.modelId) { logCtx.resolvedModel = route.modelId; logCtx.preserveResolvedModelFromRoute = true; } + // Transport-neutral reliability policy (#875): applies to any Responses + // upstream whose final adapter is openai-responses, not only WS turns. A public + // virtual selector wins; its final wire-model id is the fallback lookup. + const responsesUpstreamStreaming = providerModelResponsesUpstreamStreaming( + route.providerName, + policyProvider, + finalSelectedModelId, + route.modelId, + ); + if (responsesUpstreamStreaming === false && route.provider.adapter === "openai-responses") { + parsed.stream = false; + if (parsed._rawBody && typeof parsed._rawBody === "object") { + (parsed._rawBody as Record).stream = false; + } + } + // Fast mode override for OpenAI-routed models, only where the provider's Responses // route documents `service_tier` support (capability gate below strips everywhere else). if (config.fastMode !== undefined && route.provider.adapter === "openai-responses" && route.provider.supportsServiceTier === true) { @@ -1016,6 +1019,7 @@ async function applyFinalRouteRequestNormalization(args: { route.modelId, logCtx.requestedServiceTier ?? logCtx.configuredServiceTier, ); + return responsesUpstreamStreaming; } @@ -1582,7 +1586,7 @@ async function handleResponsesInner( // upstream for reliability (#875); the answer must then be reframed to SSE // for streaming clients. const clientRequestedStream = parsed.stream; - await applyFinalRouteRequestNormalization({ + const responsesUpstreamStreaming = await applyFinalRouteRequestNormalization({ parsed, route, config, @@ -2103,6 +2107,8 @@ async function handleResponsesInner( const passthroughCt = headers.get("content-type")?.toLowerCase(); const isEventStream = passthroughCt?.includes("text/event-stream") || (upstreamResponse.ok && !!upstreamResponse.body && !passthroughCt && parsed.stream); + const forceBoundedResponsesJson = responsesUpstreamStreaming === false + && route.provider.adapter === "openai-responses"; const terminalRecorder = codexForwardTerminalOutcomeRecorder( config, authCtx, @@ -2199,6 +2205,15 @@ async function handleResponsesInner( // (src/server/relay-eager.ts; policy: // devlog/_fin/260731_macos_rss_retention/100_darwin_eager_optin.md). // The bundled known-bad runtime remains on tee by default on both platforms. + if (forceBoundedResponsesJson && isEventStream) { + upstream.abort(new DOMException("Unexpected event stream for bounded Responses request", "AbortError")); + try { void upstreamResponse.body?.cancel(upstream.signal.reason).catch(() => undefined); } catch { /* already closed */ } + return formatErrorResponse( + 502, + "upstream_error", + "upstream ignored the bounded Responses policy and returned an event stream", + ); + } if (isEventStream && upstreamResponse.body) { const repairConfig = route.provider.responsesItemIdRepair; const snapshotRepairEnabled = hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair); @@ -2399,7 +2414,16 @@ async function handleResponsesInner( // without limit. This path is no longer rare — WebSocket turns for models whose // streaming terminal event is unreliable are deliberately answered with bounded JSON. // Oversize and stall deadlines both fail closed; a partial body is never parsed. - const bounded = await readBoundedResponseBody(upstreamResponse, UPSTREAM_JSON_BODY_READ_OPTIONS); + let bounded; + try { + bounded = await readBoundedResponseBody(upstreamResponse, { + ...UPSTREAM_JSON_BODY_READ_OPTIONS, + signal: upstream.signal, + }); + } catch (error) { + if (upstream.signal.aborted || options.abortSignal?.aborted) return clientCancelledResponse(); + throw error; + } if (bounded.oversized) { return formatErrorResponse(502, "upstream_error", "upstream JSON response exceeded the safe body limit"); } @@ -2407,11 +2431,29 @@ async function handleResponsesInner( return formatErrorResponse(502, "upstream_error", "upstream JSON response stalled before completing"); } const text = bounded.text; + let validatedUpstreamResponse: Record | undefined; + if (forceBoundedResponsesJson) { + let upstreamCandidate: unknown; + try { + upstreamCandidate = JSON.parse(text) as unknown; + } catch { + return formatErrorResponse(502, "upstream_error", "upstream returned malformed Responses JSON"); + } + const upstreamValidation = validateResponsesJsonEventResponse(upstreamCandidate); + if (!upstreamValidation.ok) { + return formatErrorResponse(502, "upstream_error", upstreamValidation.message); + } + validatedUpstreamResponse = upstreamValidation.response; + } inspectResponseLogJson(logCtx, text); if (rememberPassthroughResponse) { - try { - rememberPassthroughResponse(JSON.parse(text) as { id?: unknown; output?: unknown; status?: unknown }); - } catch { /* non-JSON despite content-type; recording is best-effort */ } + if (validatedUpstreamResponse) { + rememberPassthroughResponse(validatedUpstreamResponse); + } else { + try { + rememberPassthroughResponse(JSON.parse(text) as { id?: unknown; output?: unknown; status?: unknown }); + } catch { /* non-JSON despite content-type; recording is best-effort */ } + } } const clientJson = (() => { const restored = restoreImageGenCallsInJson(text, imageGenCallAliases); @@ -2429,6 +2471,31 @@ async function handleResponsesInner( ? rewriteResponsesModelJson(repaired, parsed._responseModelId) : repaired; })(); + let boundedResponse: Record | undefined; + let outboundJson = clientJson; + if (forceBoundedResponsesJson) { + let candidate: unknown; + try { + candidate = JSON.parse(clientJson) as unknown; + if ((clientRequestedStream === true || options.inboundTransport === "websocket") + && hasResponsesItemIdRepair(route.provider.responsesItemIdRepair)) { + candidate = repairResponsesJsonItemIds( + candidate as Record, + route.provider.responsesItemIdRepair!, + translatorBudget, + ); + outboundJson = JSON.stringify(candidate); + } + } catch { + return formatErrorResponse(502, "upstream_error", "upstream returned malformed Responses JSON"); + } + const validation = validateResponsesJsonEventResponse(candidate); + if (!validation.ok) { + return formatErrorResponse(502, "upstream_error", validation.message); + } + boundedResponse = validation.response; + } + // #875: the transport-neutral reliability policy forced a bounded JSON // upstream for a client that asked for SSE. Reframe the completed JSON // as the canonical terminal SSE sequence (created → output_item.done → @@ -2436,30 +2503,11 @@ async function handleResponsesInner( // stream that never closes. Non-streaming clients keep the plain JSON. if (clientRequestedStream === true && options.inboundTransport !== "websocket" - && providerModelResponsesUpstreamStreaming(route.providerName, route.provider, route.modelId) === false - && route.provider.adapter === "openai-responses") { - let completed: Record | undefined; - try { - const parsedCompleted = JSON.parse(clientJson) as unknown; - if (!parsedCompleted || typeof parsedCompleted !== "object" || Array.isArray(parsedCompleted)) { - throw new TypeError("bounded Responses JSON is not an object"); - } - let candidate = parsedCompleted as Record; - // The bounded-JSON answer bypasses the SSE relay, so it also bypasses - // the SSE item-id rewrite. Apply the same client-facing normalization - // here or this policy would silently disable id repair for the very - // providers that need it (raw record already happened above). - if (hasResponsesItemIdRepair(route.provider.responsesItemIdRepair)) { - candidate = repairResponsesJsonItemIds(candidate, route.provider.responsesItemIdRepair!, translatorBudget); - } - completed = candidate; - } catch { - // Non-JSON despite content-type: fall through to the plain relay. - } - if (completed) { + && forceBoundedResponsesJson) { + if (boundedResponse) { let stream: ReadableStream; try { - stream = responsesJsonToSseStream(completed); + stream = responsesJsonToSseStream(boundedResponse); } catch (error) { if (error instanceof RangeError) { return formatErrorResponse( @@ -2480,29 +2528,21 @@ async function handleResponsesInner( }); } } - // WS turns reframe this JSON into events in the bridge, which is the - // other relay-free path — normalize ids so both bounded-JSON paths agree. - const outboundJson = options.inboundTransport === "websocket" - && providerModelResponsesUpstreamStreaming(route.providerName, route.provider, route.modelId) === false - && hasResponsesItemIdRepair(route.provider.responsesItemIdRepair) - ? (() => { - try { - return JSON.stringify(repairResponsesJsonItemIds( - JSON.parse(clientJson) as Record, - route.provider.responsesItemIdRepair!, - translatorBudget, - )); - } catch { - return clientJson; - } - })() - : clientJson; return new Response(outboundJson, { status: upstreamResponse.status, statusText: upstreamResponse.statusText, headers, }); } + if (forceBoundedResponsesJson) { + upstream.abort(new DOMException("Unexpected content type for bounded Responses request", "AbortError")); + try { void upstreamResponse.body?.cancel(upstream.signal.reason).catch(() => undefined); } catch { /* already closed */ } + return formatErrorResponse( + 502, + "upstream_error", + "upstream bounded Responses response must use application/json", + ); + } const body = relayWithAbort(upstreamResponse.body, upstream); const turnAc = new AbortController(); const tracked = body ? trackStreamLifetime(body, turnAc, undefined, options.turnAdmissionLease) : null; diff --git a/src/server/ws-bridge.ts b/src/server/ws-bridge.ts index da052c1b8..71ce40bb5 100644 --- a/src/server/ws-bridge.ts +++ b/src/server/ws-bridge.ts @@ -1,5 +1,5 @@ import type { ServerWebSocket } from "bun"; -import { responsesJsonEventSequence } from "./responses-json-events"; +import { responsesJsonEventSequence, validateResponsesJsonEventResponse } from "./responses-json-events"; import { FORWARD_HEADERS } from "../adapters/openai-responses"; import type { CodexAuthContext } from "../codex/auth-context"; import { headersForCodexAuthContext } from "../codex/auth-context"; @@ -304,6 +304,9 @@ export function sendResponsesJsonAsEvents( onTerminal?: ResponsesTerminalReporter, onPayload?: ResponsesPayloadObserver, ): void { + const validation = validateResponsesJsonEventResponse(response); + if (!validation.ok) throw new TypeError(validation.message); + response = validation.response; const sendObservedFrame = (payload: Record) => { const text = JSON.stringify(payload); try { @@ -313,15 +316,27 @@ export function sendResponsesJsonAsEvents( } sendTextFrame(ws, text); }; - const finalStatus = response.status === "failed" || response.status === "incomplete" - ? response.status - : "completed"; + const finalStatus = response.status as ResponsesTerminalStatus; for (const frame of responsesJsonEventSequence(response)) { sendObservedFrame(frame); } onTerminal?.(finalStatus); } +function sendInvalidResponsesJson( + ws: ServerWebSocket, + response: Response, + message: string, + onTerminal?: ResponsesTerminalReporter, +): void { + onTerminal?.("incomplete"); + sendJsonFrame(ws, buildWsErrorFrame(502, { + type: "protocol_error", + code: "websocket_protocol_error", + message, + }, response.headers)); +} + function errorPayloadFromText(text: string): Record { try { const json = JSON.parse(text) as { error?: unknown }; @@ -381,8 +396,19 @@ export async function sendResponseToWebSocket( if (contentType.includes("application/json")) { const text = await response.text(); if (!isCurrent()) return; - const json = JSON.parse(text) as Record; - sendResponsesJsonAsEvents(ws, json, options.onTerminal, options.onSsePayload); + let json: unknown; + try { + json = JSON.parse(text) as unknown; + } catch { + sendInvalidResponsesJson(ws, response, "Upstream returned malformed Responses JSON", options.onTerminal); + return; + } + const validation = validateResponsesJsonEventResponse(json); + if (!validation.ok) { + sendInvalidResponsesJson(ws, response, validation.message, options.onTerminal); + return; + } + sendResponsesJsonAsEvents(ws, validation.response, options.onTerminal, options.onSsePayload); return; } @@ -404,8 +430,19 @@ export async function sendResponseToWebSocket( if (!isCurrent()) return; const trimmed = text.trim(); if (trimmed.startsWith("{")) { - const json = JSON.parse(trimmed) as Record; - sendResponsesJsonAsEvents(ws, json, options.onTerminal, options.onSsePayload); + let json: unknown; + try { + json = JSON.parse(trimmed) as unknown; + } catch { + sendInvalidResponsesJson(ws, response, "Upstream returned malformed Responses JSON", options.onTerminal); + return; + } + const validation = validateResponsesJsonEventResponse(json); + if (!validation.ok) { + sendInvalidResponsesJson(ws, response, validation.message, options.onTerminal); + return; + } + sendResponsesJsonAsEvents(ws, validation.response, options.onTerminal, options.onSsePayload); return; } @@ -467,4 +504,4 @@ export async function readBoundedPrefix( export function looksLikeSse(prefix: Uint8Array): boolean { const text = new TextDecoder().decode(prefix); return /^\s*(event:|data:)/.test(text); -} \ No newline at end of file +} diff --git a/src/types.ts b/src/types.ts index b2b0c6d9e..ab64ac0ac 100644 --- a/src/types.ts +++ b/src/types.ts @@ -1149,6 +1149,13 @@ export interface OcxProviderConfig { * as before. */ modelAdapters?: Record; + /** + * Per-model Responses upstream streaming policy. `false` keeps the client-facing Responses + * transport but requests bounded JSON upstream, then reframes that terminal object for a + * streaming client. `true` explicitly overrides a matching registry `false` default. Use only + * with non-forward providers on the `openai-responses` wire. + */ + modelResponsesUpstreamStreaming?: Record; baseUrl: string; /** * Optional relative resource path for key-auth openai-responses requests. Must start with `/` diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index 23cb53570..8fe2abfd5 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -211,12 +211,16 @@ the upgrade with 426 so Codex falls back to HTTP cleanly. The endpoint handles `response.create`, ignores `response.processed`, supports warmup `generate: false`, and feeds the same request pipeline as HTTP/SSE. -Registry-declared per-model compatibility hints (`modelResponsesUpstreamStreaming`) may ask the -upstream Responses endpoint for bounded JSON on ANY client transport — WebSocket or ordinary -HTTP/SSE. The bridge reframes that JSON into the same Responses event sequence +Per-model compatibility hints (`modelResponsesUpstreamStreaming`) may be declared by the registry +or explicitly configured on a non-forward `openai-responses` provider. An explicit model value +wins over the matching registry default. `false` asks the upstream Responses endpoint for bounded +JSON on ANY client transport — WebSocket or ordinary HTTP/SSE. The validated terminal object is +reframed into the same Responses event sequence (`src/server/responses-json-events.ts`): WS turns send the frames as WebSocket messages, while HTTP clients that requested streaming receive a synthesized terminal SSE body (created → -output_item.done → terminal → `[DONE]`). No production registry entry currently opts in: +output_item.done → terminal → `[DONE]`). If the upstream ignores the policy and returns SSE or a +different successful content type, the request fails once with 502 instead of entering the tee +relay or replaying the model request. No production registry entry currently opts in: DeepSeek V4 Flash used this path while its public-beta Responses stream was suspected of not closing on the terminal event, but the official guide documents a `response.completed`/`response.incomplete`/`response.failed` terminal with no `data: [DONE]` @@ -228,9 +232,9 @@ rollback for upstreams that regress, kept suite-reachable by a synthetic-registr Synthesized output is capped at 10,000 items across HTTP and WebSocket reframing. HTTP frames are encoded incrementally, so bounded upstream JSON cannot expand into an unbounded event array or SSE string. -`ws-bridge.ts` preserves upstream `failed` and `incomplete` status values in the final WebSocket -frame rather than always emitting `response.completed`. If the response status is `failed`, a -`response.failed` frame is sent; otherwise `response.completed` carries through the original status. +`ws-bridge.ts` preserves the upstream terminal status in the final WebSocket frame rather than +always emitting `response.completed`: `completed`, `failed`, and `incomplete` become +`response.completed`, `response.failed`, and `response.incomplete`, respectively. ## Heartbeat and stall deadline diff --git a/tests/config.test.ts b/tests/config.test.ts index 0c1313b81..17b83f160 100644 --- a/tests/config.test.ts +++ b/tests/config.test.ts @@ -1030,6 +1030,90 @@ describe("opencodex config defaults", () => { } }); + test("modelResponsesUpstreamStreaming accepts only effective Responses model policies", () => { + writeConfig({ + port: 12345, + providers: { + custom: { + adapter: "openai-chat", + baseUrl: "https://example.test/v1", + modelAdapters: { BOUNDED: "openai-responses" }, + modelResponsesUpstreamStreaming: { BOUNDED: false }, + }, + }, + defaultProvider: "custom", + }); + expect(readConfigDiagnostics().error).toBeNull(); + expect(readConfigDiagnostics().config.providers.custom.modelResponsesUpstreamStreaming) + .toEqual({ BOUNDED: false }); + + for (const provider of [ + { + adapter: "openai-responses", + baseUrl: "https://example.test/v1", + modelResponsesUpstreamStreaming: { "": false }, + }, + { + adapter: "openai-responses", + baseUrl: "https://example.test/v1", + modelResponsesUpstreamStreaming: { bounded: "false" }, + }, + { + adapter: "openai-chat", + baseUrl: "https://example.test/v1", + modelResponsesUpstreamStreaming: { bounded: false }, + }, + { + adapter: "openai-responses", + baseUrl: "https://example.test/v1", + modelAdapters: { bounded: "openai-chat" }, + modelResponsesUpstreamStreaming: { bounded: false }, + }, + { + adapter: "openai-chat", + baseUrl: "https://example.test/v1", + modelAdapters: { bounded: "openai-responses" }, + modelResponsesUpstreamStreaming: { BOUNDED: false }, + }, + ]) { + writeConfig({ port: 12345, providers: { custom: provider }, defaultProvider: "custom" }); + expect(readConfigDiagnostics().source).toBe("fallback"); + expect(readConfigDiagnostics().error).toContain("modelResponsesUpstreamStreaming"); + } + + writeConfig({ + port: 12345, + providers: { + openai: { + adapter: "openai-responses", + baseUrl: "https://chatgpt.com/backend-api/codex", + authMode: "forward", + codexAccountMode: "pool", + modelResponsesUpstreamStreaming: { "gpt-test": false }, + }, + }, + defaultProvider: "openai", + openaiProviderTierVersion: 2, + }); + expect(readConfigDiagnostics().source).toBe("fallback"); + expect(readConfigDiagnostics().error).toContain("not supported on forward-auth"); + + writeConfig({ + port: 12345, + providers: { + "openai-apikey": { + adapter: "openai-responses", + baseUrl: "https://api.openai.com/v1", + modelAdapters: { "gpt-5.6-sol-pro": "openai-chat" }, + modelResponsesUpstreamStreaming: { "GPT-5.6-SOL-PRO": false }, + }, + }, + defaultProvider: "openai-apikey", + }); + expect(readConfigDiagnostics().source).toBe("fallback"); + expect(readConfigDiagnostics().error).toContain("requires the openai-responses wire"); + }); + test("modelReasoningSummaryDelivery validates known values and rejects summary opt-out conflicts (#538)", () => { writeConfig({ port: 12345, diff --git a/tests/deepseek-inbound-wire.test.ts b/tests/deepseek-inbound-wire.test.ts index cf9c4df11..dd45326f4 100644 --- a/tests/deepseek-inbound-wire.test.ts +++ b/tests/deepseek-inbound-wire.test.ts @@ -13,7 +13,7 @@ */ import { afterEach, beforeEach, describe, expect, test } from "bun:test"; import { enrichProviderFromRegistry, providerConfigSeed } from "../src/providers/derive"; -import { getProviderRegistryEntry, PROVIDER_REGISTRY } from "../src/providers/registry"; +import { getProviderRegistryEntry, providerModelResponsesUpstreamStreaming, PROVIDER_REGISTRY } from "../src/providers/registry"; import { createResponsesPassthroughAdapter as createResponsesPassthroughAdapterProduction } from "../src/adapters/openai-responses"; import { resolveWireProtocolOverride } from "../src/server/adapter-resolve"; import { handleResponses } from "../src/server/responses/core"; @@ -345,6 +345,27 @@ describe("the bounded-JSON mechanism stays alive behind a synthetic registry ent expect(text).toContain('"type":"response.completed"'); }); + test("an explicit true config value overrides a false registry default", async () => { + const provider = fixtureProvider({ modelResponsesUpstreamStreaming: { [FIXTURE_MODEL]: true } }); + expect(providerModelResponsesUpstreamStreaming(FIXTURE_ID, provider, FIXTURE_MODEL)).toBe(true); + const captured: Array<{ stream?: boolean }> = []; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + captured.push(JSON.parse(String(init?.body ?? "{}")) as { stream?: boolean }); + return completedWithPlaceholderIds(); + }) as typeof fetch; + const response = await driveFixture(provider); + expect(captured[0]?.stream).toBe(true); + expect(response.headers.get("content-type")).toContain("application/json"); + }); + + test("the policy lookup is case-insensitive exact and does not inherit colon families", () => { + const provider = fixtureProvider({ + modelResponsesUpstreamStreaming: { "FIXTURE-MODEL": false }, + }); + expect(providerModelResponsesUpstreamStreaming(FIXTURE_ID, provider, "fixture-model")).toBe(false); + expect(providerModelResponsesUpstreamStreaming(FIXTURE_ID, provider, "fixture-model:variant")).toBeUndefined(); + }); + test("the synthesized terminal SSE carries repaired item ids, not the upstream placeholders", async () => { globalThis.fetch = (async () => completedWithPlaceholderIds()) as typeof fetch; const response = await driveFixture(repairingFixtureProvider()); @@ -361,7 +382,10 @@ describe("the bounded-JSON mechanism stays alive behind a synthetic registry ent id: "resp_fixture", object: "response", status: "completed", - output: Array.from({ length: MAX_SYNTHESIZED_OUTPUT_ITEMS + 1 }, () => null), + output: Array.from( + { length: MAX_SYNTHESIZED_OUTPUT_ITEMS + 1 }, + (_, index) => ({ type: "message", id: `msg_${index}` }), + ), })) as typeof fetch; const response = await driveFixture(fixtureProvider()); expect(response.status).toBe(502); @@ -405,6 +429,178 @@ describe("the bounded-JSON mechanism stays alive behind a synthetic registry ent }); }); +describe("custom providers can opt into the bounded-JSON Responses path", () => { + const originalFetch = globalThis.fetch; + const PROVIDER = "custom-bounded-responses"; + const MODEL_ID = "fixture-model"; + + afterEach(() => { globalThis.fetch = originalFetch; }); + + function customProvider(overrides: Partial = {}): OcxProviderConfig { + return { + adapter: "openai-responses", + baseUrl: "https://custom-bounded.example.test/v1", + authMode: "key", + apiKey: "sk-test", + modelResponsesUpstreamStreaming: { "FIXTURE-MODEL": false }, + ...overrides, + }; + } + + async function drive( + provider: OcxProviderConfig, + websocket = false, + abortSignal: AbortSignal = AbortSignal.timeout(5_000), + ): Promise { + return handleResponses( + new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: `${PROVIDER}/${MODEL_ID}`, input: "ping", stream: true }), + }), + { providers: { [PROVIDER]: provider } } as unknown as OcxConfig, + { model: "", provider: "" }, + websocket + ? { inboundWire: "responses", inboundTransport: "websocket" } + : { abortSignal }, + ); + } + + test("a case-insensitive config match changes only upstream stream and synthesizes HTTP SSE", async () => { + const captured: Array> = []; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + captured.push(JSON.parse(String(init?.body ?? "{}")) as Record); + return Response.json({ + id: "resp_custom", + object: "response", + status: "completed", + output: [{ + id: "msg_custom", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "ok" }], + }], + usage: { input_tokens: 1, output_tokens: 1, total_tokens: 2 }, + }); + }) as typeof fetch; + + const response = await drive(customProvider()); + expect(captured).toHaveLength(1); + expect(captured[0]?.stream).toBe(false); + expect(captured[0]?.model).toBe(MODEL_ID); + expect(response.status).toBe(200); + expect(response.headers.get("content-type")).toContain("text/event-stream"); + const text = await response.text(); + expect(text.match(/"type":"response\.created"/g)).toHaveLength(1); + expect(text.match(/"type":"response\.output_item\.done"/g)).toHaveLength(1); + expect(text.match(/"type":"response\.completed"/g)).toHaveLength(1); + expect(text.match(/data: \[DONE\]/g)).toHaveLength(1); + }); + + test("failed and incomplete terminal statuses survive HTTP synthesis", async () => { + for (const status of ["failed", "incomplete"] as const) { + globalThis.fetch = (async () => Response.json({ + id: `resp_${status}`, + object: "response", + status, + output: [], + ...(status === "failed" + ? { error: { code: "upstream_failure", message: "fixture" } } + : { incomplete_details: { reason: "max_output_tokens" } }), + })) as typeof fetch; + const response = await drive(customProvider()); + const text = await response.text(); + expect(text).toContain(`"type":"response.${status}"`); + expect(text).not.toContain('"type":"response.completed"'); + } + }); + + test("an upstream SSE that ignores stream:false fails once without entering tee", async () => { + let calls = 0; + let cancelled = false; + let aborted = false; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + calls += 1; + init?.signal?.addEventListener("abort", () => { aborted = true; }, { once: true }); + return new Response(new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode('data: {"type":"response.completed"}\n\n')); + }, + cancel() { + cancelled = true; + return new Promise(() => {}); + }, + }), { status: 200, headers: { "content-type": "text/event-stream" } }); + }) as typeof fetch; + + const response = await Promise.race([ + drive(customProvider()), + Bun.sleep(250).then(() => { throw new Error("bounded policy waited for body cancellation"); }), + ]); + expect(response.status).toBe(502); + expect(calls).toBe(1); + expect(aborted).toBe(true); + expect(cancelled).toBe(true); + expect(await response.text()).not.toContain("response.completed"); + }); + + test("malformed terminal JSON fails closed for HTTP and WebSocket handoff", async () => { + for (const websocket of [false, true]) { + globalThis.fetch = (async () => Response.json({ + id: "resp_invalid", + object: "response", + status: "still_running", + output: [], + })) as typeof fetch; + const response = await drive(customProvider(), websocket); + expect(response.status).toBe(502); + expect(response.headers.get("content-type")).toContain("application/json"); + expect(await response.text()).toContain("terminal status"); + } + }); + + test("snapshot repair cannot promote a malformed bounded response to completed", async () => { + globalThis.fetch = (async () => Response.json({ + id: "resp_invalid", + object: "response", + })) as typeof fetch; + const response = await drive(customProvider({ responsesSnapshotRepair: true })); + expect(response.status).toBe(502); + expect(await response.text()).toContain("terminal status"); + }); + + test("a per-model Responses wire override is settled before the policy", async () => { + const captured: Array<{ stream?: boolean }> = []; + globalThis.fetch = (async (_input: RequestInfo | URL, init?: RequestInit) => { + captured.push(JSON.parse(String(init?.body ?? "{}")) as { stream?: boolean }); + return Response.json({ id: "resp_mixed", status: "completed", output: [] }); + }) as typeof fetch; + const response = await drive(customProvider({ + adapter: "openai-chat", + modelAdapters: { [MODEL_ID]: "openai-responses" }, + })); + expect(response.status).toBe(200); + expect(captured[0]?.stream).toBe(false); + expect(await response.text()).toContain('"type":"response.completed"'); + }); + + test("client cancellation aborts a pending bounded JSON read", async () => { + let cancelled = false; + globalThis.fetch = (async () => new Response(new ReadableStream({ + pull() { return new Promise(() => {}); }, + cancel() { cancelled = true; }, + }), { status: 200, headers: { "content-type": "application/json" } })) as typeof fetch; + const controller = new AbortController(); + const pending = drive(customProvider(), false, controller.signal); + await Bun.sleep(5); + controller.abort(new DOMException("fixture cancel", "AbortError")); + const response = await pending; + expect(response.status).toBe(499); + expect(cancelled).toBe(true); + }); +}); + /** * DeepSeek documents "the API is stateless: responses and conversations are not * stored on the server", so parameters that reference server-held state can never be diff --git a/tests/management-provider-validation.test.ts b/tests/management-provider-validation.test.ts index 81ab8133b..b9cb9f61d 100644 --- a/tests/management-provider-validation.test.ts +++ b/tests/management-provider-validation.test.ts @@ -214,6 +214,40 @@ describe("provider management validation", () => { })).toContain("canonical built-in provider seed"); }); + test("provider management validates bounded Responses model policies", () => { + expect(providerManagementConfigError("custom", { + adapter: "openai-responses", + baseUrl: "https://api.example.test/v1", + modelResponsesUpstreamStreaming: { bounded: false, streamed: true }, + })).toBeNull(); + expect(providerManagementConfigError("custom-mixed", { + adapter: "openai-chat", + baseUrl: "https://api.example.test/v1", + modelAdapters: { BOUNDED: "openai-responses" }, + modelResponsesUpstreamStreaming: { BOUNDED: false }, + })).toBeNull(); + expect(providerManagementConfigError("custom-mixed", { + adapter: "openai-chat", + baseUrl: "https://api.example.test/v1", + modelAdapters: { bounded: "openai-responses" }, + modelResponsesUpstreamStreaming: { BOUNDED: false }, + })).toContain("requires the openai-responses wire"); + expect(providerManagementConfigError("custom-chat", { + adapter: "openai-chat", + baseUrl: "https://api.example.test/v1", + modelResponsesUpstreamStreaming: { bounded: false }, + })).toContain("requires the openai-responses wire"); + expect(providerManagementConfigError("custom", { + adapter: "openai-responses", + baseUrl: "https://api.example.test/v1", + modelResponsesUpstreamStreaming: { bounded: "false" }, + })).toContain("must be a boolean"); + expect(providerManagementConfigError("openai", { + ...canonicalDirect, + modelResponsesUpstreamStreaming: { bounded: false }, + })).toContain("not supported on forward-auth"); + }); + 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/openai-api-virtual-models.test.ts b/tests/openai-api-virtual-models.test.ts index b14846387..f3a3b6f22 100644 --- a/tests/openai-api-virtual-models.test.ts +++ b/tests/openai-api-virtual-models.test.ts @@ -11,7 +11,8 @@ import { resolveOpenAiVirtualModel, validateOpenAiVirtualModelDefinition, } from "../src/providers/openai-virtual-models"; -import { PROVIDER_REGISTRY } from "../src/providers/registry"; +import { getProviderRegistryEntry, providerModelResponsesUpstreamStreaming, PROVIDER_REGISTRY } from "../src/providers/registry"; +import { providerConfigSeed } from "../src/providers/derive"; import { resolveWireProtocolOverride } from "../src/server/adapter-resolve"; import { saveConfig } from "../src/config"; import { startServer } from "../src/server"; @@ -64,6 +65,19 @@ describe("OpenAI API virtual model resolution", () => { expect(resolveOpenAiVirtualModel("openai-apikey", model)).toBeUndefined(); expect(resolveOpenAiCompactModel("openai-apikey", model)).toBeUndefined(); }); + + test("a base-model streaming policy applies after a Pro alias resolves", () => { + const provider = { + ...providerConfigSeed(getProviderRegistryEntry("openai-apikey")!), + modelResponsesUpstreamStreaming: { "gpt-5.6-luna": false }, + }; + expect(providerModelResponsesUpstreamStreaming( + "openai-apikey", + provider, + "gpt-5.6-luna-pro", + "gpt-5.6-luna", + )).toBe(false); + }); }); describe("applyOpenAiVirtualModel", () => { diff --git a/tests/responses-json-events.test.ts b/tests/responses-json-events.test.ts index 9e27ffc5c..135e7dd1a 100644 --- a/tests/responses-json-events.test.ts +++ b/tests/responses-json-events.test.ts @@ -4,6 +4,7 @@ import { responsesJsonEventSequence, responsesJsonToSseBody, responsesJsonToSseStream, + validateResponsesJsonEventResponse, } from "../src/server/responses-json-events"; describe("responsesJsonEventSequence", () => { @@ -35,11 +36,23 @@ describe("responsesJsonEventSequence", () => { } }); - test("empty and non-array outputs yield the minimal sequence", () => { - expect(responsesJsonEventSequence({ id: "r" }).map(frame => frame.type)) - .toEqual(["response.created", "response.completed"]); - expect(responsesJsonEventSequence({ id: "r", output: null }).map(frame => frame.type)) - .toEqual(["response.created", "response.completed"]); + test("invalid terminal snapshots are rejected instead of upgraded to completed", () => { + for (const response of [ + { id: "r", output: [] }, + { id: "r", status: "running", output: [] }, + { id: "r", status: "completed", output: null }, + ]) { + expect(() => responsesJsonEventSequence(response)).toThrow(TypeError); + } + }); + + test("parallel function calls remain byte-semantically intact", () => { + const output = [ + { id: "fc_1", type: "function_call", call_id: "call_1", name: "first", arguments: "{\"a\":1}" }, + { id: "fc_2", type: "function_call", call_id: "call_2", name: "second", arguments: "{\"b\":2}" }, + ]; + const frames = responsesJsonEventSequence({ id: "r", status: "completed", output }); + expect(frames.slice(1, 3).map(frame => frame.item)).toEqual(output); }); test("the payload rewrite hook runs on every frame (060 seam)", () => { @@ -51,6 +64,57 @@ describe("responsesJsonEventSequence", () => { }); }); +describe("validateResponsesJsonEventResponse", () => { + test("accepts all terminal statuses and absent, null, or valid usage", () => { + for (const status of ["completed", "failed", "incomplete"]) { + for (const usage of [ + undefined, + null, + { input_tokens: 1, output_tokens: 2 }, + { input_tokens: 1, output_tokens: 2, total_tokens: 3 }, + { + input_tokens: 1, + output_tokens: 2, + total_tokens: 3, + input_tokens_details: { cached_tokens: 0, cache_write_tokens: 0 }, + output_tokens_details: { reasoning_tokens: 1 }, + }, + ]) { + const response = { id: "r", object: "response", status, output: [], ...(usage !== undefined ? { usage } : {}) }; + expect(validateResponsesJsonEventResponse(response).ok).toBe(true); + } + } + }); + + test("rejects malformed snapshots without echoing their contents", () => { + const invalid = [ + null, + [], + { id: "", status: "completed", output: [] }, + { id: "r", object: "list", status: "completed", output: [] }, + { id: "r", status: "unknown", output: [] }, + { id: "r", status: "completed", output: {} }, + { id: "r", status: "completed", output: [null] }, + { id: "r", status: "completed", output: [{}] }, + { id: "r", status: "completed", output: [], usage: [] }, + { id: "r", status: "completed", output: [], usage: {} }, + { id: "r", status: "completed", output: [], usage: { input_tokens: "1", output_tokens: 1, total_tokens: 2 } }, + { + id: "r", + status: "completed", + output: [], + usage: { + input_tokens: 1, + output_tokens: 1, + total_tokens: 2, + input_tokens_details: { cached_tokens: "0" }, + }, + }, + ]; + for (const value of invalid) expect(validateResponsesJsonEventResponse(value).ok).toBe(false); + }); +}); + describe("responsesJsonToSseBody", () => { test("serializes the sequence with exactly one [DONE] trailer", () => { const body = responsesJsonToSseBody({ id: "r", status: "completed", output: [{ type: "message", id: "m" }] }); @@ -64,9 +128,13 @@ describe("responsesJsonToSseBody", () => { }); test("rejects output arrays that could amplify synthesized frames", () => { - const output = Array.from({ length: MAX_SYNTHESIZED_OUTPUT_ITEMS + 1 }, () => null); - expect(() => responsesJsonToSseBody({ id: "r", output })).toThrow(RangeError); - expect(() => responsesJsonToSseStream({ id: "r", output })).toThrow(RangeError); + const output = Array.from( + { length: MAX_SYNTHESIZED_OUTPUT_ITEMS + 1 }, + (_, index) => ({ type: "message", id: `m_${index}` }), + ); + const response = { id: "r", status: "completed", output }; + expect(() => responsesJsonToSseBody(response)).toThrow(RangeError); + expect(() => responsesJsonToSseStream(response)).toThrow(RangeError); }); test("streams one SSE frame per pull and ends with [DONE]", async () => { diff --git a/tests/ws-endpoint.test.ts b/tests/ws-endpoint.test.ts index 9554ad7a4..4607f069b 100644 --- a/tests/ws-endpoint.test.ts +++ b/tests/ws-endpoint.test.ts @@ -376,6 +376,26 @@ describe("WS endpoint re-framer (120/132)", () => { expect(JSON.parse(payloads.at(-1)!).response.usage).toEqual({ input_tokens: 11, output_tokens: 7 }); }); + test("invalid successful Responses JSON becomes one protocol error", async () => { + for (const body of [ + { id: "json", status: "running", output: [] }, + { id: "json", status: "completed", output: {}, usage: [] }, + ]) { + const { ws, sent } = mockWs(); + const terminals: string[] = []; + await sendResponseToWebSocket(ws, Response.json(body), () => true, { + onTerminal: status => terminals.push(status), + }); + expect(sent).toHaveLength(1); + expect(JSON.parse(sent[0])).toMatchObject({ + type: "error", + status: 502, + error: { code: "websocket_protocol_error" }, + }); + expect(terminals).toEqual(["incomplete"]); + } + }); + test("unexpected successful HTML and empty 204 become standalone protocol errors", async () => { const html = mockWs(); await sendResponseToWebSocket(html.ws, new Response("", {