From 2aac0b6ae2f40a8c6919ce28ced4202fe36e29bd Mon Sep 17 00:00:00 2001 From: t Date: Mon, 7 Sep 2026 01:55:05 +0900 Subject: [PATCH 1/5] fix(catalog): bound proven native custom reasoning efforts [skip ci] Normalize canonical Codex-forward custom ladders and defaults after metadata merging. Preserve explicit empty lists and prevent fresh custom projections from gaining max again during catalog merge. Add existing-fixture coverage for provenance negatives, repeat sync, observed convergence, and direct model discovery. Keep arbitrary gateways, shared effort maps, and request-time policy unchanged. Verification: source review and git diff --check only. Tests, typecheck, build, install, prepush, live proxy, and CI NOT RUN by explicit user instruction. Issue #3775 remains partial pending gateway provenance and exact Desktop runtime evidence. --- .../content/docs/guides/codex-app-models.md | 17 +++ .../docs/reference/configuration/providers.md | 8 ++ src/codex/catalog/provider-fetch.ts | 47 ++++++-- src/codex/catalog/sync.ts | 6 +- structure/03_catalog-and-subagents.md | 18 +++ .../claude-models-discovery.test.ts | 30 +++++ .../codex-catalog-sync-hardening.test.ts | 56 ++++++++- tests/codex-integration/codex-catalog.test.ts | 110 ++++++++++++++++++ ...odex-convergence-account-selectors.test.ts | 30 +++++ 9 files changed, 310 insertions(+), 12 deletions(-) diff --git a/docs-site/src/content/docs/guides/codex-app-models.md b/docs-site/src/content/docs/guides/codex-app-models.md index 4f3a2218f58..bcda44080b8 100644 --- a/docs-site/src/content/docs/guides/codex-app-models.md +++ b/docs-site/src/content/docs/guides/codex-app-models.md @@ -64,6 +64,23 @@ or grant account entitlement. The separately billed `openai-apikey/daybreak-blue-latest` API row is a different route and its 1,050,000 / 922,000 limits are never copied into the Codex-login row. +For custom Astra and Daybreak rows on that canonical `openai` Codex-forward destination, +explicit `reasoningEfforts` are bounded by the model's pinned Codex capabilities. A custom +`["none", "minimal", "low"]` becomes `["low"]` in the catalog; a nonempty list with no +supported values also falls back to the native default as a single choice. An explicit `[]` +stays empty and has no advertised default. A declared default is retained only if it belongs to +the resulting list; otherwise the native default is used when present, then the first surviving +choice. Stored custom configuration is unchanged, and repeated syncs do not add `max` back to a +narrow custom list. + +This requires the exact provider, destination, and capability-backed model identity. An arbitrary +gateway such as `YYLJ/gpt-6-astra` does not inherit native capabilities from its name. Its explicit +custom ladder continues to override discovered provider metadata under the normal routed rules. +Codex's native Astra `ultra` choice is retained: it is a client delegation mode converted to a +supported wire effort, distinct from the [API model's effort list](https://developers.openai.com/api/docs/models/gpt-6-astra). +Catalog normalization does not rewrite existing thread settings or establish support for a +particular installed Desktop version. + When the `codexAccountNamespaces` map is empty, account-qualified picker rows are off. If `codexAccountPickerEnabled` is omitted with a non-empty map, they are treated as enabled for backward compatibility. Set it to `false` to hide generated qualified rows and restore bare native diff --git a/docs-site/src/content/docs/reference/configuration/providers.md b/docs-site/src/content/docs/reference/configuration/providers.md index cd5c9657d69..02d5ba43237 100644 --- a/docs-site/src/content/docs/reference/configuration/providers.md +++ b/docs-site/src/content/docs/reference/configuration/providers.md @@ -201,6 +201,14 @@ predictions. Explicit provider/model price overrides still take precedence. | `unsafeAllowNativeLocalExec?` | `boolean` | Cursor legacy boolean, equivalent to `nativeLocalExec: "on"` only when the newer field is unset. | | `nativeLocalExec?` | `"off" \| "codex-sandbox" \| "on"` | Cursor local-exec policy. `off` is default; `codex-sandbox` currently fails closed like `off`. | +Custom-model `reasoningEfforts` normally override discovered provider metadata. The bounded +exception is an explicit Astra or Daybreak custom row on the canonical `openai` Codex-forward +destination: its advertised list is intersected with that model's pinned native capabilities. +An explicit empty list remains empty with no default; a nonempty incompatible list falls back +to the native default as a single choice. Defaults must belong to the final list. This changes +the catalog projection, not stored configuration or arbitrary gateway models sharing a GPT name. +See [custom native catalog examples](/guides/codex-app-models/). + ### Discovered model display names Use `modelDisplayNames` when a provider returns machine friendly ids but the Codex model picker diff --git a/src/codex/catalog/provider-fetch.ts b/src/codex/catalog/provider-fetch.ts index 55b90cfea45..58e21cf82cf 100644 --- a/src/codex/catalog/provider-fetch.ts +++ b/src/codex/catalog/provider-fetch.ts @@ -2130,6 +2130,31 @@ async function gatherRoutedModelsWithAuth( return models; } +/** Bound a proven Codex-forward custom row without changing its stored configuration. */ +function boundCustomNativeReasoning( + model: CatalogModel, + allowed: readonly string[], + nativeDefault: string | undefined, +): CatalogModel { + if (allowed.length === 0 || model.reasoningEfforts === undefined) return model; + const bounded = { ...model }; + if (model.reasoningEfforts.length === 0) { + bounded.reasoningEfforts = []; + delete bounded.defaultReasoningEffort; + return bounded; + } + const declared = new Set(model.reasoningEfforts); + const surviving = [...new Set(allowed)].filter(effort => declared.has(effort)); + const fallback = nativeDefault && allowed.includes(nativeDefault) ? nativeDefault : allowed[0]!; + // A nonempty but incompatible declaration is not an explicit no-reasoning setting. + bounded.reasoningEfforts = surviving.length > 0 ? surviving : [fallback]; + bounded.defaultReasoningEffort = model.defaultReasoningEffort + && bounded.reasoningEfforts.includes(model.defaultReasoningEffort) + ? model.defaultReasoningEffort + : bounded.reasoningEfforts.includes(fallback) ? fallback : bounded.reasoningEfforts[0]!; + return bounded; +} + async function gatherRoutedModelsUncached( config: OcxConfig, capture: GatherFlightCapture, @@ -2394,7 +2419,7 @@ async function gatherRoutedModelsUncached( ...(typeof supportsReasoningSummaries === "boolean" ? { supportsReasoningSummaries } : {}), // Native-alias defaults apply only where the custom row declares nothing: the explicit // spreads below must win (later in object order), so a stored `[]` stays empty and a - // declared ladder is never replaced by the alias's native ladder. + // declared ladder is narrowed to proven native capabilities after the merge below. ...(codexForwardNativeCapabilityAlias ? { codexForwardNativeCapabilityAlias: true, @@ -2409,7 +2434,8 @@ async function gatherRoutedModelsUncached( : {}), // Explicit custom-row ladder wins over the inherited provider row below: the merge only // gap-fills, so a stored `[]` (explicit "no reasoning") or a declared ladder is kept - // verbatim instead of being replaced by the replaced row's metadata. + // instead of being replaced by that row's metadata. Only proven native aliases are + // bounded against their own capability source after the merge. ...(Array.isArray(cm.reasoningEfforts) ? { reasoningEfforts: [...cm.reasoningEfforts] } : {}), ...(cm.defaultReasoningEffort ? { defaultReasoningEffort: cm.defaultReasoningEffort } : {}), ...(typeof supportsServiceTier === "boolean" ? { supportsServiceTier } : {}), @@ -2462,22 +2488,25 @@ async function gatherRoutedModelsUncached( ...(base.codexToolMode === undefined && replaced.codexToolMode !== undefined ? { codexToolMode: replaced.codexToolMode } : {}), ...(base.capabilities === undefined && replaced.capabilities !== undefined ? { capabilities: replaced.capabilities } : {}), } : base; + const reasoningBounded = codexForwardNativeCapabilityAlias + ? boundCustomNativeReasoning(merged, nativeReasoningEfforts(cm.modelId), nativeAliasDefaultEffort) + : merged; // Vision-sidecar coverage only: when the enriched provider's shared predicate matches // noVisionModels or text-without-image modelInputModalities, advertise image input so the // Codex app lets images reach the sidecar (#349/#344). Deliberately NOT the full // applyProviderConfigHints pass — custom rows are a // user override, so their explicit contextWindow / inputModalities / reasoning fields must be // preserved verbatim (the hint pass would cap context and overwrite modalities from registry). - const mergedContext = typeof merged.contextWindow === "number" && merged.contextWindow > 0 - ? merged.contextWindow + const mergedContext = typeof reasoningBounded.contextWindow === "number" && reasoningBounded.contextWindow > 0 + ? reasoningBounded.contextWindow : undefined; - const boundedMergedMaxInput = typeof merged.maxInputTokens === "number" && merged.maxInputTokens > 0 - ? (mergedContext !== undefined ? Math.min(merged.maxInputTokens, mergedContext) : merged.maxInputTokens) + const boundedMergedMaxInput = typeof reasoningBounded.maxInputTokens === "number" && reasoningBounded.maxInputTokens > 0 + ? (mergedContext !== undefined ? Math.min(reasoningBounded.maxInputTokens, mergedContext) : reasoningBounded.maxInputTokens) : undefined; const mergedWithHardBounds = boundedMergedMaxInput !== undefined - && boundedMergedMaxInput !== merged.maxInputTokens - ? { ...merged, maxInputTokens: boundedMergedMaxInput } - : merged; + && boundedMergedMaxInput !== reasoningBounded.maxInputTokens + ? { ...reasoningBounded, maxInputTokens: boundedMergedMaxInput } + : reasoningBounded; const mergedSoftCandidates = [mergedWithHardBounds.autoCompactTokenLimit, configuredAutoCompact] .filter((value): value is number => typeof value === "number" && value > 0); const mergedWithAutoCompact: CatalogModel = mergedContext !== undefined && mergedSoftCandidates.length > 0 diff --git a/src/codex/catalog/sync.ts b/src/codex/catalog/sync.ts index 81a27eb01e7..22d235dbccf 100644 --- a/src/codex/catalog/sync.ts +++ b/src/codex/catalog/sync.ts @@ -903,6 +903,10 @@ export function mergeCatalogEntriesFromObservedState({ const detachedBaselineCatalogModels = baselineCatalogModels .map(entry => structuredClone(entry) as RawEntry); const detachedRoutedEntries = routedEntries.map(entry => structuredClone(entry) as RawEntry); + // Track this invocation's generated custom rows, not ownership markers read from disk. + // Their builder already finalized exact native ladders and ordinary routed mock tiers. + const freshCustomEntries = new Set(detachedRoutedEntries.filter(entry => + entry.opencodex_catalog_kind === CODEX_CUSTOM_MODEL_CATALOG_KIND)); const detachedAccountBoundEntries = accountBoundEntries .map(entry => structuredClone(entry) as RawEntry); const disabledModelKeys = new Set([...disabledModels].map(slugEquivalenceKey)); @@ -1223,7 +1227,7 @@ export function mergeCatalogEntriesFromObservedState({ // Mock-max universality (260709): preserved routed entries from disk may predate // the max rung — ensure it here so subagent max spawns validate on every // reasoning-capable entry. max only: 5.6 exact ladders (luna: no ultra) stay intact. - if (!exactCombo && !reserveProjection && !String(e.slug ?? "").startsWith("opencode-go/")) { + if (!freshCustomEntries.has(m) && !exactCombo && !reserveProjection && !String(e.slug ?? "").startsWith("opencode-go/")) { const levels = Array.isArray(e.supported_reasoning_levels) ? e.supported_reasoning_levels as Array<{ effort?: string }> : []; diff --git a/structure/03_catalog-and-subagents.md b/structure/03_catalog-and-subagents.md index 84dfcab25af..f9f46a06533 100644 --- a/structure/03_catalog-and-subagents.md +++ b/structure/03_catalog-and-subagents.md @@ -37,6 +37,24 @@ custom catalog remains the native metadata/template authority even when a bundle warm. Both paths may use an admitted matching bundled memo only as installed-runtime capability evidence to remove unsupported reasoning efforts; convergence never probes Codex itself. +Custom Astra and Daybreak rows acquire native reasoning capability only through the existing +canonical `openai` forward destination and explicit capability-source predicate. The shared +custom-row producer bounds their merged effort lists against pinned per-model Codex metadata, +preserves an explicit empty list without a default, and recovers an incompatible nonempty list +to the native default singleton. A default must belong to the projected list. Other custom rows +keep their declaration precedence; a GPT model name, display alias, or arbitrary gateway is not +native provenance. Stored configuration and native capability maps are unchanged. + +The observed-state merge tracks the current invocation's freshly generated custom row objects +after detaching its inputs. Those rows already own their complete reasoning projection, so the +merge does not append `max` again. This also keeps a generic none-only custom row none-only; +ordinary retained provider rows still receive the existing mock-tier policy. A persisted custom +marker alone never grants this exemption. Both gather entry points, retained sync, management +convergence and direct Codex model discovery use the same producer. The legacy runtime effort +union clamp remains separate; it is not a per-model or per-client-version grammar oracle. +Existing thread settings and the reported Desktop 0.153.4 gateway rejection require separate +runtime evidence. Codex's native `ultra` mode is preserved and is not a literal API wire promise. + When account selectors are enabled, the sync path may also observe exact, visible, API-supported OpenAI-family ids from Codex's user-owned catalog/cache. Only rows with native catalog provenance are trusted; unknown ids are carried through startup cache invalidation as hidden observations and diff --git a/tests/claude-integration/claude-models-discovery.test.ts b/tests/claude-integration/claude-models-discovery.test.ts index 5fd77c96eca..608a1ffa800 100644 --- a/tests/claude-integration/claude-models-discovery.test.ts +++ b/tests/claude-integration/claude-models-discovery.test.ts @@ -141,6 +141,36 @@ test("per-surface id style: ?ids= wins, claude-code UA gets readable, unknown UA } }); +test("Codex discovery bounds proven custom Astra before any disk sync and preserves a gateway namesake", async () => { + const config = configWithStaticModels(); + config.providers.openai = { + adapter: "openai-responses", + baseUrl: "https://chatgpt.com/backend-api/codex", + authMode: "forward", + liveModels: false, + }; + config.providers.YYLJ = { adapter: "openai-chat", baseUrl: "https://gateway.example.test/v1", liveModels: false, models: ["gpt-6-astra"] }; + config.customModels = ["openai", "YYLJ"].map(provider => ({ + id: `${provider}-astra`, provider, modelId: "gpt-6-astra", + reasoningEfforts: ["none", "minimal", "low"], defaultReasoningEffort: "minimal", + })); + saveConfig(config); + const server = startServer(0); + try { + const response = await fetch(new URL("/v1/models?client_version=0.153.4", server.url)); + expect(response.status).toBe(200); + const catalog = await response.json() as { models: Array<{ slug: string; supported_reasoning_levels: Array<{ effort: string }>; default_reasoning_level?: string }> }; + const canonical = catalog.models.find(row => row.slug === "openai/gpt-6-astra"); + expect(canonical?.supported_reasoning_levels.map(level => level.effort)).toEqual(["low"]); + expect(canonical?.default_reasoning_level).toBe("low"); + const gateway = catalog.models.find(row => row.slug === "YYLJ/gpt-6-astra"); + expect(gateway?.supported_reasoning_levels.map(level => level.effort)).toEqual(["none", "minimal", "low", "max", "ultra"]); + expect(gateway?.default_reasoning_level).toBe("minimal"); + } finally { + await server.stop(true); + } +}); + test("OpenAI list shape and Codex catalog shape stay unchanged", async () => { saveConfig(configWithStaticModels()); const server = startServer(0); diff --git a/tests/codex-integration/codex-catalog-sync-hardening.test.ts b/tests/codex-integration/codex-catalog-sync-hardening.test.ts index 6f7e7e38fd8..21c60d58353 100644 --- a/tests/codex-integration/codex-catalog-sync-hardening.test.ts +++ b/tests/codex-integration/codex-catalog-sync-hardening.test.ts @@ -28,9 +28,9 @@ function runScript( return { stdout: result.stdout?.trim() ?? "", stderr: result.stderr ?? "", status: result.status ?? 1 }; } -function createCodexCatalogFixture(dir: string): string { +function createCodexCatalogFixture(dir: string, models = [nativeEntry("gpt-5.5", 0)]): string { const scriptPath = join(dir, "codex-catalog-fixture.js"); - const bundled = JSON.stringify({ models: [nativeEntry("gpt-5.5", 0)] }); + const bundled = JSON.stringify({ models }); writeFileSync(scriptPath, [ 'if (process.argv.includes("--version")) {', ' console.log("codex-cli 0.999.0");', @@ -447,6 +447,58 @@ describe("Codex catalog sync hardening", () => { expect(rows.filter(row => row.slug === "gpt-daybreak-blue-latest")).toHaveLength(1); }); + test("canonical custom Astra repairs stale efforts and keeps a narrow ladder across syncs", () => { + const catalogPath = join(codexHome, "catalog.json"); + const runtime = createCodexCatalogFixture(codexHome, [{ + ...nativeEntry("gpt-5.5", 0), + // Another model permits sentinels, so the global union cannot perform this repair. + supported_reasoning_levels: ["none", "minimal", "low", "medium", "high", "xhigh", "max", "ultra"].map(effort => ({ effort, description: effort })), + }]); + writeFileSync(join(codexHome, "config.toml"), 'model_catalog_json = "catalog.json"\n'); + writeFileSync(catalogPath, JSON.stringify({ models: [{ + ...ocxAuthoredEntry("openai/gpt-6-astra", 5), + opencodex_catalog_kind: "custom-model-v1", + supported_reasoning_levels: [{ effort: "minimal", description: "stale" }], + default_reasoning_level: "minimal", + }] })); + const result = runScript(codexHome, opencodexHome, ` + const { readFileSync } = require("node:fs"); + const { saveConfig } = require("./src/config"); + const { syncCatalogModels } = require("./src/codex/catalog"); + const config = { + port: 10100, + defaultProvider: "openai", + providers: { openai: { adapter: "openai-responses", baseUrl: "https://chatgpt.com/backend-api/codex", authMode: "forward", codexAccountMode: "pool" } }, + codexAccountPickerEnabled: false, + customModels: [{ id: "astra", provider: "openai", modelId: "gpt-6-astra", reasoningEfforts: ["none", "minimal", "low"], defaultReasoningEffort: "minimal" }] + }; + saveConfig(config); + (async () => { + const first = await syncCatalogModels(config, { allowWhenDesiredDisabled: true }); + const firstBytes = readFileSync(first.path, "utf8"); + const second = await syncCatalogModels(config, { allowWhenDesiredDisabled: true }); + const secondBytes = readFileSync(second.path, "utf8"); + config.customModels = []; + saveConfig(config); + await syncCatalogModels(config, { allowWhenDesiredDisabled: true }); + console.log(JSON.stringify({ + first: JSON.parse(firstBytes), second: JSON.parse(secondBytes), + unchanged: firstBytes === secondBytes, + deleted: JSON.parse(readFileSync(first.path, "utf8")) + })); + })(); + `, { CODEX_CLI_PATH: runtime }); + expect(result.status).toBe(0); + const output = JSON.parse(result.stdout); + for (const catalog of [output.first, output.second]) { + const astra = catalog.models.find((row: { slug: string }) => row.slug === "openai/gpt-6-astra"); + expect(astra.supported_reasoning_levels.map((level: { effort: string }) => level.effort)).toEqual(["low"]); + expect(astra.default_reasoning_level).toBe("low"); + } + expect(output.unchanged).toBe(true); + expect(output.deleted.models.some((row: { slug: string }) => row.slug === "openai/gpt-6-astra")).toBe(false); + }); + test("explicit Codex-forward Daybreak survives sync with Sol metadata while account picker is off", () => { const catalogPath = join(codexHome, "catalog.json"); writeFileSync(join(codexHome, "config.toml"), 'model_catalog_json = "catalog.json"\n', "utf8"); diff --git a/tests/codex-integration/codex-catalog.test.ts b/tests/codex-integration/codex-catalog.test.ts index 916b6a20979..5820edc6dd9 100644 --- a/tests/codex-integration/codex-catalog.test.ts +++ b/tests/codex-integration/codex-catalog.test.ts @@ -46,6 +46,7 @@ import { enrichProviderFromCatalog } from "../../src/oauth/key-providers"; import { handleManagementAPI } from "../../src/server/management-api"; import { OAUTH_PROVIDERS } from "../../src/oauth"; import { + catalogEntryEfforts, clampCatalogModelsToObservedCodexSupport, supportedCodexReasoningEffortsFromObservedCatalog, } from "../../src/codex/catalog/effort"; @@ -3704,6 +3705,114 @@ describe("Codex catalog routed normalization", () => { expect(astra?.base_instructions).not.toContain("daybreak"); }); + const nativeCustomEffortCases: Array<{ + name: string; + efforts?: string[]; + defaultEffort?: string; + expected: string[]; + expectedDefault?: string; + }> = [ + { name: "legacy sentinels", efforts: ["none", "minimal", "low", "medium", "high", "xhigh", "max"], defaultEffort: "minimal", expected: ["low", "medium", "high", "xhigh", "max"], expectedDefault: "low" }, + { name: "native default", expected: ["low", "medium", "high", "xhigh", "max", "ultra"], expectedDefault: "low" }, + { name: "empty declaration", efforts: [], defaultEffort: "minimal", expected: [] }, + { name: "no compatible rung", efforts: ["none", "minimal"], defaultEffort: "minimal", expected: ["low"], expectedDefault: "low" }, + { name: "narrow subset", efforts: ["low"], defaultEffort: "high", expected: ["low"], expectedDefault: "low" }, + { name: "valid explicit default", efforts: ["high", "medium", "high"], defaultEffort: "high", expected: ["medium", "high"], expectedDefault: "high" }, + { name: "first survivor default", efforts: ["high", "medium"], defaultEffort: "minimal", expected: ["medium", "high"], expectedDefault: "medium" }, + { name: "Ultra mode", efforts: ["minimal", "ultra"], defaultEffort: "minimal", expected: ["ultra"], expectedDefault: "ultra" }, + ]; + + test.each(nativeCustomEffortCases)("canonical custom Astra bounds $name through gather/build/merge", async fixture => { + globalThis.fetch = (() => { throw new Error("canonical forward discovery must not fetch"); }) as typeof fetch; + const config = withStubbedProviderFetch({ + port: 10100, + defaultProvider: "openai", + providers: { openai: { adapter: "openai-responses", baseUrl: "https://chatgpt.com/backend-api/codex", codexAccountMode: "pool" } }, + customModels: [{ + id: "astra-effort", + provider: "openai", + modelId: "gpt-6-astra", + ...(fixture.efforts !== undefined ? { reasoningEfforts: fixture.efforts } : {}), + ...(fixture.defaultEffort !== undefined ? { defaultReasoningEffort: fixture.defaultEffort } : {}), + }], + }); + const beforeConfig = JSON.stringify(config); + const beforeNative = JSON.stringify(upstreamNativeEntry("gpt-6-astra")); + const models = await gatherRoutedModelsDirect(config); + const custom = models.find(row => row.provider === "openai" && row.id === "gpt-6-astra"); + expect(custom?.codexForwardNativeCapabilityAlias).toBe(true); + expect(custom?.reasoningEfforts).toEqual(fixture.expected); + expect(custom?.defaultReasoningEffort).toBe(fixture.expectedDefault); + + const entries = buildCatalogEntries(nativeTemplate(), [], models); + const first = mergeCatalogEntriesForSync([], entries, new Map(), [], false); + const second = mergeCatalogEntriesForSync(first, buildCatalogEntries(nativeTemplate(), [], models), new Map(), [], false); + // Another model's sentinels make the legacy union permissive: it cannot mask this bug. + const observed = { models: [{ + slug: "other-model", + supported_reasoning_levels: ["none", "minimal", "low", "medium", "high", "xhigh", "max", "ultra"].map(effort => ({ effort })), + }] }; + const beforeObserved = JSON.stringify(observed); + for (const projection of [entries, first, second]) { + clampCatalogModelsToObservedCodexSupport(projection, supportedCodexReasoningEffortsFromObservedCatalog(observed)); + const row = projection.find(entry => entry.slug === "openai/gpt-6-astra"); + expect(row ? catalogEntryEfforts(row) : undefined).toEqual(fixture.expected); + expect(row?.default_reasoning_level).toBe(fixture.expectedDefault); + expect(row?.use_responses_lite).toBe(true); + expect(row?.multi_agent_reasoning_effort).toBe("xhigh"); + if (fixture.expected.length === 0) expect(row).not.toHaveProperty("default_reasoning_level"); + } + expect(JSON.stringify(config)).toBe(beforeConfig); + expect(JSON.stringify(upstreamNativeEntry("gpt-6-astra"))).toBe(beforeNative); + expect(JSON.stringify(observed)).toBe(beforeObserved); + }); + + test.each([ + { name: "YYLJ", adapter: "openai-responses", baseUrl: "https://gateway.example.test/v1", authMode: "key", modelId: "gpt-6-astra" }, + { name: "openai", adapter: "openai-responses", baseUrl: "https://gateway.example.test/v1", authMode: "forward", modelId: "gpt-6-astra" }, + { name: "openai", adapter: "openai-responses", baseUrl: "https://chatgpt.com/backend-api/codex", authMode: "key", modelId: "gpt-6-astra" }, + { name: "openai", adapter: "openai-chat", baseUrl: "https://chatgpt.com/backend-api/codex", authMode: "key", modelId: "gpt-6-astra" }, + { name: "openai-apikey", adapter: "openai-responses", baseUrl: "https://api.openai.com/v1", authMode: "key", modelId: "gpt-6-astra" }, + { name: "openai", adapter: "openai-responses", baseUrl: "https://chatgpt.com/backend-api/codex", authMode: "forward", modelId: "gpt-unproven" }, + ] satisfies Array<{ name: string; adapter: OcxProviderConfig["adapter"]; baseUrl: string; authMode: OcxProviderConfig["authMode"]; modelId: string }>)( + "custom $name/$modelId does not infer native effort capability from $baseUrl / $authMode / $adapter", + async fixture => { + const models = await gatherRoutedModels({ + port: 10100, + defaultProvider: fixture.name, + providers: { [fixture.name]: { adapter: fixture.adapter, baseUrl: fixture.baseUrl, authMode: fixture.authMode, liveModels: false, models: [fixture.modelId] } }, + customModels: [{ id: "unproven", provider: fixture.name, modelId: fixture.modelId, displayName: "Astra", reasoningEfforts: ["none", "minimal", "low"], defaultReasoningEffort: "minimal" }], + }); + const custom = models.find(row => row.provider === fixture.name && row.id === fixture.modelId); + expect(custom?.codexForwardNativeCapabilityAlias).toBeUndefined(); + expect(custom?.reasoningEfforts).toEqual(["none", "minimal", "low"]); + expect(custom?.defaultReasoningEffort).toBe("minimal"); + const entries = buildCatalogEntries(nativeTemplate(), [], models); + const row = entries.find(entry => entry.slug === `${fixture.name}/${fixture.modelId}`); + expect(row ? catalogEntryEfforts(row) : undefined) + .toEqual(["none", "minimal", "low", "max", "ultra"]); + }, + ); + + test("fresh none-only custom rows keep their ladder while retained provider rows still gain max", async () => { + const models = await gatherRoutedModels({ + port: 10100, + defaultProvider: "custom-provider", + providers: { "custom-provider": { adapter: "openai-chat", baseUrl: "https://example.invalid/v1", liveModels: false } }, + customModels: [{ id: "none-only", provider: "custom-provider", modelId: "none-only", reasoningEfforts: ["none"] }], + }); + const entries = buildCatalogEntries(nativeTemplate(), [], models); + const retained = { ...nativeTemplate(), slug: "foreign/model", supported_reasoning_levels: [{ effort: "low", description: "Low" }] }; + const stale = { ...entries[0]!, slug: "custom-provider/deleted" }; + const merged = mergeCatalogEntriesForSync([retained, stale], entries, new Map(), [], false); + expect(merged.find(row => row.slug === "custom-provider/none-only")?.supported_reasoning_levels) + .toEqual(entries[0]!.supported_reasoning_levels); + const foreign = merged.find(row => row.slug === "foreign/model"); + expect(foreign ? catalogEntryEfforts(foreign) : undefined) + .toEqual(["low", "max"]); + expect(merged.some(row => row.slug === "custom-provider/deleted")).toBe(false); + }); + test("Astra refresh repairs only built-in speed text and does not leak native effort", () => { const pinned = upstreamNativeEntry(NATIVE_GPT6_ASTRA_MODEL)!; expect(pinned.service_tiers).toEqual([{ id: "priority", name: "Fast", description: "2x speed, increased usage" }]); @@ -3771,6 +3880,7 @@ describe("Codex catalog routed normalization", () => { // not overwrite it — otherwise the catalog would advertise reasoning the user // explicitly disabled for this row. reasoningEfforts: [], + defaultReasoningEffort: "minimal", }], }); const model = models.find(row => row.provider === "openai" && row.id === NATIVE_DAYBREAK_BLUE_MODEL); diff --git a/tests/codex-integration/codex-convergence-account-selectors.test.ts b/tests/codex-integration/codex-convergence-account-selectors.test.ts index ae83c1f2bf4..9602e913095 100644 --- a/tests/codex-integration/codex-convergence-account-selectors.test.ts +++ b/tests/codex-integration/codex-convergence-account-selectors.test.ts @@ -38,6 +38,7 @@ import { getConfigPath, saveConfig } from "../../src/config"; import { CODEX_FORWARD_BASE_URL } from "../../src/providers/openai-tiers"; import type { OcxConfig } from "../../src/types"; import { setBundledCatalogCacheForTests } from "../../src/codex/catalog/bundled"; +import { catalogEntryEfforts } from "../../src/codex/catalog/effort"; import { resetCodexRuntimeResolveCacheForTests, setCodexRuntimeResolveCacheForTests, @@ -577,6 +578,35 @@ test("convergence preserves only provider-local degraded rows", async () => { expect(models.some(entry => entry.slug === "external/vendor-model")).toBe(true); }); +test.each([ + { efforts: ["none", "minimal", "low"], expected: ["low"], defaultEffort: "low" }, + { efforts: ["none", "minimal"], expected: ["low"], defaultEffort: "low" }, + { efforts: [], expected: [], defaultEffort: undefined }, +])("observed convergence bounds canonical custom efforts $efforts without reviving stale max", async fixture => { + seedObservedRuntimeSupport(["none", "minimal", "low", "medium", "high", "xhigh", "max", "ultra"]); + const nextConfig = config(false); + nextConfig.customModels = [{ + id: "astra-custom", + provider: "openai", + modelId: "gpt-6-astra", + reasoningEfforts: fixture.efforts, + defaultReasoningEffort: "minimal", + }]; + writeCatalog([nativeEntry(), { + ...generatedRoutedEntry("openai/gpt-6-astra"), + opencodex_catalog_kind: "custom-model-v1", + supported_reasoning_levels: [{ effort: "minimal", description: "Stale" }, { effort: "max", description: "Stale max" }], + default_reasoning_level: "minimal", + }]); + for (let pass = 0; pass < 2; pass++) { + const catalog = await convergeCatalog(nextConfig); + const row = catalog.models?.find(entry => entry.slug === "openai/gpt-6-astra"); + expect(row ? catalogEntryEfforts(row) : undefined).toEqual(fixture.expected); + expect(row?.default_reasoning_level).toBe(fixture.defaultEffort); + if (fixture.expected.length === 0) expect(row).not.toHaveProperty("default_reasoning_level"); + } +}); + function legacyCustomDeletionConfig(): OcxConfig { const nextConfig = config(false); nextConfig.providers.offline = { From 73f1d2c1357a71e2d8952431941dc31d5e20eb01 Mon Sep 17 00:00:00 2001 From: t Date: Mon, 7 Sep 2026 01:59:54 +0900 Subject: [PATCH 2/5] fix(chat): preserve refusal across completion projections --- .../content/docs/reference/proxy-formats.md | 8 + src/chat/outbound.ts | 259 ++++++++-- src/server/chat-native-sse.ts | 1 + structure/04_transports-and-sidecars.md | 19 + tests/responses/chat-refusal.test.ts | 451 ++++++++++++++++++ 5 files changed, 711 insertions(+), 27 deletions(-) create mode 100644 tests/responses/chat-refusal.test.ts diff --git a/docs-site/src/content/docs/reference/proxy-formats.md b/docs-site/src/content/docs/reference/proxy-formats.md index d2167bbfa45..7008194e194 100644 --- a/docs-site/src/content/docs/reference/proxy-formats.md +++ b/docs-site/src/content/docs/reference/proxy-formats.md @@ -260,6 +260,14 @@ It does not issue an additional inference request. An incomplete response caused token limit or content filtering retains `length` or `content_filter`, even if it includes tool output. Other incomplete boundaries return an upstream error instead of claiming a normal finish. +Refusal text stays separate from answer text: JSON completions use nullable `message.refusal`, +and streaming chunks use `delta.refusal`. Native Chat JSON-to-SSE and SSE-to-JSON conversions +preserve that field; native streaming relay preserves the provider's refusal deltas. On translated +Responses streams, refusal parts are buffered until the terminal event and emitted once in their +original output/content order. Compatible repeated or sparse snapshots do not duplicate or erase +text. Contradictory refusal snapshots and buffer overflow produce a typed error without a successful +finish or `[DONE]`. This preserves the upstream refusal; it does not introduce a proxy policy decision. + Because the internal execution path is Responses-based, a provider adapter can impose a narrower feature set. For example, a request feature that cannot be represented by the selected adapter is returned as an error instead of silently changing its meaning. diff --git a/src/chat/outbound.ts b/src/chat/outbound.ts index 812f1b553f8..7646b63efba 100644 --- a/src/chat/outbound.ts +++ b/src/chat/outbound.ts @@ -8,7 +8,11 @@ type Rec = Record; import { decodeServerSentEvents, sseFieldValue } from "../lib/sse-decoder"; -import { isTranslatorBudgetExceededError, type TranslatorBudget } from "../lib/translator-budget"; +import { + isTranslatorBudgetExceededError, + type TranslatorBudget, + type TranslatorTransientReservation, +} from "../lib/translator-budget"; import { classifyError, cyberPolicyErrorType, @@ -156,6 +160,14 @@ function appendedUtf8Bytes(previous: string, previousBytes: number, fragment: st return nextBytes; } +function refusalTranslationError(): ChatCompletionsStreamError { + // Never include provider-controlled refusal text or correlation IDs in diagnostics. + return new ChatCompletionsStreamError("upstream refusal representations are inconsistent", { + type: "upstream_error", + code: "invalid_refusal", + }); +} + /** * Streaming: Responses SSE bytes -> Chat Completions SSE bytes. */ @@ -188,6 +200,115 @@ export function responsesSseToChatCompletionsSse( let emittedFrames = 0; let stepping = false; let decoderStarted = false; + // Raw output/content positions are the ordering authority; IDs only constrain identity. + // Charge a fixed entry allowance as well as keys/IDs so empty parts remain bounded. + const refusalEntryBytes = 64; + const refusalItems = new Map; + }>(); + let refusalMetadataBytes = 0; + let refusalTextBytes = 0; + const releaseRefusals = () => { + refusalItems.clear(); + translatorBudget.releaseRetained(refusalMetadataBytes, { kind: "item_ids" }); + translatorBudget.releaseRetained(refusalTextBytes, { kind: "retained_collectors" }); + refusalMetadataBytes = 0; + refusalTextBytes = 0; + }; + const position = (value: unknown): number => { + if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) { + throw refusalTranslationError(); + } + return value; + }; + const chargeRefusalMetadata = (bytes: number) => { + translatorBudget.chargeRetained(bytes, { kind: "item_ids" }); + refusalMetadataBytes += bytes; + }; + const refusalItem = (outputIndex: unknown, source: Rec, idField: string) => { + const index = position(outputIndex); + const hasId = Object.hasOwn(source, idField); + const candidate = source[idField]; + if (hasId && typeof candidate !== "string") throw refusalTranslationError(); + let item = refusalItems.get(index); + if (!item) { + chargeRefusalMetadata(refusalEntryBytes + Buffer.byteLength(String(index))); + item = { parts: new Map() }; + refusalItems.set(index, item); + } + if (hasId && typeof candidate === "string") { + if (item.id !== undefined && item.id !== candidate) throw refusalTranslationError(); + if (item.id === undefined) { + chargeRefusalMetadata(Buffer.byteLength(candidate)); + item.id = candidate; + } + } + return item; + }; + const retainRefusal = (outputIndex: unknown, contentIndex: unknown, source: Rec, + idField: string, evidence: Rec, field: string, delta = false) => { + const item = refusalItem(outputIndex, source, idField); + const index = position(contentIndex); + let part = item.parts.get(index); + if (!part) { + chargeRefusalMetadata(refusalEntryBytes + Buffer.byteLength(String(index))); + part = { text: "", bytes: 0, present: false }; + item.parts.set(index, part); + } + if (!Object.hasOwn(evidence, field)) { + if (delta) throw refusalTranslationError(); + return; + } + const candidate = evidence[field]; + if (typeof candidate !== "string") throw refusalTranslationError(); + part.present = true; + if (!delta) { + // Equal, empty, and stale-prefix snapshots add no evidence; never erase deltas. + if (part.text.startsWith(candidate)) return; + if (!candidate.startsWith(part.text)) throw refusalTranslationError(); + } + const nextBytes = delta ? appendedUtf8Bytes(part.text, part.bytes, candidate) : Buffer.byteLength(candidate); + const reservation = translatorBudget.reserveTransient(nextBytes, { kind: "retained_collectors" }); + try { + const next = delta ? part.text + candidate : candidate; + reservation.commitRetained(); + translatorBudget.releaseRetained(part.bytes, { kind: "retained_collectors" }); + refusalTextBytes += nextBytes - part.bytes; + part.text = next; + part.bytes = nextBytes; + } catch (error) { + reservation.release(); + throw error; + } + }; + const snapshotRefusalItem = (outputIndex: unknown, item: Rec) => { + const existing = typeof outputIndex === "number" ? refusalItems.get(outputIndex) : undefined; + if (item.type !== "message") { + if (existing && existing.parts.size > 0 && item.type !== undefined) throw refusalTranslationError(); + return; + } + // Unrelated sparse text messages historically need no position metadata. + if (outputIndex === undefined && (!Array.isArray(item.content) + || !item.content.some(part => isRec(part) && part.type === "refusal"))) return; + const known = refusalItem(outputIndex, item, "id"); + if (!Array.isArray(item.content)) return; + item.content.forEach((part: unknown, contentIndex: number) => { + if (!isRec(part)) return; + if (part.type === "refusal") { + retainRefusal(outputIndex, contentIndex, item, "id", part, "refusal"); + } else if (part.type !== undefined && known.parts.has(contentIndex)) { + throw refusalTranslationError(); + } + }); + }; + const snapshotRefusals = (response: Rec) => { + if (!Array.isArray(response.output)) return; + response.output.forEach((item: unknown, outputIndex: number) => { + if (isRec(item)) snapshotRefusalItem(outputIndex, item); + }); + }; + let terminalBatch: Array<{ frame: Uint8Array; reservation: TranslatorTransientReservation }> | undefined; const queuedLiveFrameBytes: number[] = []; const enqueueLiveFrame = (frame: Uint8Array) => { const reservation = translatorBudget.reserveTransient(frame.byteLength, { kind: "live_transient" }); @@ -235,8 +356,20 @@ export function responsesSseToChatCompletionsSse( }; const emit = (payload: Rec | "[DONE]") => { if (failed) return; - enqueueLiveFrame(encoder.encode(dataFrame(payload))); - emittedFrames++; + if (terminalBatch) { + const serialized = dataFrame(payload); + const stringReservation = translatorBudget.reserveTransient(Buffer.byteLength(serialized), { kind: "live_transient" }); + try { + const frame = encoder.encode(serialized); + const reservation = translatorBudget.reserveTransient(frame.byteLength, { kind: "live_transient" }); + terminalBatch.push({ frame, reservation }); + } finally { + stringReservation.release(); + } + } else { + enqueueLiveFrame(encoder.encode(dataFrame(payload))); + emittedFrames++; + } }; const ensureRole = () => { if (started) return; @@ -298,38 +431,68 @@ export function responsesSseToChatCompletionsSse( }; const finish = (finishReason: string, usage: unknown) => { if (terminated) return; - // A valid completed/incomplete terminal frame may arrive without output_item.done. - // Preserve any known tool call before emitting its finish reason. - flushPendingToolCalls(); + // Admit every pending role/tool/refusal/finish/DONE frame before exposing any + // of this terminal batch. Serialization and encoded bytes coexist and both count. + const batch: NonNullable = []; + terminalBatch = batch; + try { + flushPendingToolCalls(); + ensureRole(); + for (const [, item] of [...refusalItems.entries()].sort(([a], [b]) => a - b)) { + for (const [, part] of [...item.parts.entries()].sort(([a], [b]) => a - b)) { + if (!part.present) continue; + const refusal = chunkBase(id, model, created); + refusal.choices = [{ index: 0, delta: { refusal: part.text }, finish_reason: null }]; + emit(refusal); + } + } + const frame = chunkBase(id, model, created); + frame.choices = [{ index: 0, delta: {}, finish_reason: finishReason }]; + if (usage) frame.usage = chatCompletionsUsage(usage); + emit(frame); + emit("[DONE]"); + } catch (error) { + for (const staged of batch) staged.reservation.release(); + throw error; + } finally { + terminalBatch = undefined; + } + for (const staged of batch) { + controller.enqueue(staged.frame); + staged.reservation.commitRetained(); + queuedLiveFrameBytes.push(staged.frame.byteLength); + emittedFrames++; + } terminated = true; - ensureRole(); - const frame = chunkBase(id, model, created); - frame.choices = [{ index: 0, delta: {}, finish_reason: finishReason }]; - if (usage) frame.usage = chatCompletionsUsage(usage); - emit(frame); - emit("[DONE]"); + releaseRefusals(); }; const fail = (message: string, details?: { code?: string | null; type?: string; status?: number }) => { if (terminated) return; terminated = true; failed = true; + releaseRefusals(); + closeToolCalls(); + upstreamAbort.abort(new Error("upstream chat translation failed")); + try { void sseIterator?.return(undefined).catch(() => {}); } catch { /* already closed */ } // OpenAI-compatible clients need a real error event, not a success completion // that embeds `[error] ...` text followed by a clean [DONE]. // Deliver the error frame then close the stream abnormally (no [DONE]). // Do not controller.error() — that can drop already-enqueued bytes from consumers // like response.text(). - const safeMessage = redactSecretString(message); + const translatorOverflow = details?.code === "translation_buffer_limit"; + const safeMessage = translatorOverflow ? "upstream translation buffer exceeded the safe limit" + : details?.code === "invalid_refusal" ? "upstream refusal representations are inconsistent" + : redactSecretString(message); const statusHint = details?.status ?? streamErrorStatus(safeMessage); const classified = classifyError(statusHint, details?.type ?? "upstream_error", safeMessage); - const translatorOverflow = details?.code === "translation_buffer_limit"; if (translatorOverflow) { - upstreamAbort.abort(new Error("upstream translation buffer exceeded the safe limit")); - closeToolCalls(); - try { void sseIterator?.return(undefined).catch(() => {}); } catch { /* already closed */ } classified.code = "translation_buffer_limit"; // Provider-controlled overflow is an upstream failure on every path: // streaming frame, collector, and defensive JSON agree on 502. classified.type = "upstream_error"; + } else if (details?.code === "invalid_refusal") { + classified.code = details.code; + classified.type = "upstream_error"; } else if (isCyberPolicyCode(details?.code) || classified.code === CYBER_POLICY_ERROR_CODE) { classified.code = CYBER_POLICY_ERROR_CODE; classified.type = cyberPolicyErrorType(details?.type); @@ -345,9 +508,9 @@ export function responsesSseToChatCompletionsSse( code: classified.code, }, })); - // The budget is already exhausted. This bounded emergency frame is the sole - // typed overflow closure and therefore cannot reserve from that budget again. - if (translatorOverflow) controller.enqueue(frame); + // These fixed, bounded failures must survive even when decoder-owned input + // still fills the budget. They contain no provider text or IDs. + if (translatorOverflow || details?.code === "invalid_refusal") controller.enqueue(frame); else enqueueLiveFrame(frame); emittedFrames++; } catch { @@ -371,8 +534,29 @@ export function responsesSseToChatCompletionsSse( if (typeof data.delta === "string") emitReasoning(data.delta); break; } + case "response.refusal.delta": + case "response.refusal.done": { + const delta = eventName === "response.refusal.delta"; + retainRefusal(data.output_index, data.content_index, data, "item_id", data, delta ? "delta" : "refusal", delta); + break; + } + case "response.content_part.added": + case "response.content_part.done": { + const part = isRec(data.part) ? data.part : null; + if (part?.type === "refusal") { + retainRefusal(data.output_index, data.content_index, data, "item_id", part, "refusal"); + } else if (typeof data.output_index === "number" && refusalItems.has(data.output_index)) { + const item = refusalItem(data.output_index, data, "item_id"); + if (part?.type !== undefined && item.parts.has(position(data.content_index))) throw refusalTranslationError(); + } + break; + } case "response.output_item.added": { const item = isRec(data.item) ? data.item : null; + if (item?.type === "message") { + snapshotRefusalItem(data.output_index, item); + if (Object.hasOwn(data, "item_id")) refusalItem(data.output_index, data, "item_id"); + } if (!item || item.type !== "function_call") break; ensureRole(); sawToolUse = true; @@ -416,6 +600,8 @@ export function responsesSseToChatCompletionsSse( case "response.output_item.done": { const item = isRec(data.item) ? data.item : null; if (!item) break; + snapshotRefusalItem(data.output_index, item); + if (item.type === "message" && Object.hasOwn(data, "item_id")) refusalItem(data.output_index, data, "item_id"); if (item.type === "function_call") { sawToolUse = true; const callId = typeof item.call_id === "string" ? item.call_id : ""; @@ -445,6 +631,7 @@ export function responsesSseToChatCompletionsSse( } case "response.completed": { const response = isRec(data.response) ? data.response : {}; + snapshotRefusals(response); finish(sawToolUse ? "tool_calls" : "stop", response.usage); break; } @@ -456,6 +643,7 @@ export function responsesSseToChatCompletionsSse( : undefined; if (reason !== undefined) { // Truthful OpenAI-compatible finish reasons: the turn ended, just early. + snapshotRefusals(response); finish(reason, response.usage); } else { // upstream_stall_timeout / adapter_eof / proxy-synthesized incompletes are @@ -500,6 +688,7 @@ export function responsesSseToChatCompletionsSse( while (!cancelled && emittedFrames === emittedAtStart) { decoderStarted = true; const next = await sseIterator!.next(); + if (cancelled) break; if (next.done) { if (!cancelled && !terminated) { fail("upstream stream ended before a terminal frame (truncated response)"); @@ -525,6 +714,8 @@ export function responsesSseToChatCompletionsSse( upstreamAbort.abort(err); closeToolCalls(); fail(err.message, { status: 502, type: "upstream_error", code: err.code }); + } else if (isChatCompletionsStreamError(err)) { + fail(err.message, { status: err.status, type: err.type, code: err.code }); } else { fail(err instanceof Error ? err.message : String(err)); } @@ -545,6 +736,7 @@ export function responsesSseToChatCompletionsSse( }, cancel(reason) { cancelled = true; + releaseRefusals(); while (queuedLiveFrameBytes.length > 0) releaseDeliveredFrame(); closeToolCalls(); // Abort first: it cancels the decoder's underlying reader, settling any in-flight @@ -573,6 +765,8 @@ export function responsesJsonToChatCompletion(json: unknown, model: string, tran } const output = Array.isArray(body.output) ? body.output : []; let content = ""; + let refusal: string | null = null; + let refusalBytes = 0; let reasoning = ""; let contentBytes = 0; let reasoningBytes = 0; @@ -599,6 +793,11 @@ export function responsesJsonToChatCompletion(json: unknown, model: string, tran for (const part of raw.content) { if (isRec(part) && part.type === "output_text" && typeof part.text === "string") { ({ text: content, bytes: contentBytes } = append(content, contentBytes, part.text)); + } else if (isRec(part) && part.type === "refusal" && Object.hasOwn(part, "refusal")) { + if (typeof part.refusal !== "string") throw refusalTranslationError(); + const next = append(refusal ?? "", refusalBytes, part.refusal); + refusal = next.text; + refusalBytes = next.bytes; } } } else if (raw.type === "reasoning") { @@ -645,6 +844,7 @@ export function responsesJsonToChatCompletion(json: unknown, model: string, tran const message: Rec = { role: "assistant", content: content || null, + refusal, }; if (reasoning) message.reasoning_content = reasoning; if (toolCalls.length > 0) message.tool_calls = toolCalls; @@ -673,6 +873,7 @@ export async function collectChatCompletion( const decoder = new TextDecoder(); let buffer = ""; let content = ""; + let refusal: string | null = null; let reasoning = ""; const toolCalls = new Map(); // Per-call budget scopes (2 MiB/call enforced by the budget): the map key is the @@ -680,7 +881,6 @@ export async function collectChatCompletion( const callScope = (index: number) => `chat_collect_${index}`; let finishReason = "stop"; let usage: unknown; - let streamError: ChatCompletionsStreamError | null = null; const replaceRetained = (previous: string, next: string, kind: "live_transient" | "retained_collectors") => { const reservation = translatorBudget.reserveTransient(Buffer.byteLength(next), { kind }); reservation.commitRetained(); @@ -733,12 +933,12 @@ export async function collectChatCompletion( : code === CYBER_POLICY_ERROR_CODE || isCyberPolicyMessage(message) ? 400 : streamErrorStatus(message); - streamError = new ChatCompletionsStreamError(message, { + const streamError = new ChatCompletionsStreamError(message, { status, type: code === "translation_buffer_limit" ? "upstream_error" : type, code, }); - continue; + throw streamError; } if (parsed.usage) usage = parsed.usage; const choices = Array.isArray(parsed.choices) ? parsed.choices : []; @@ -748,6 +948,10 @@ export async function collectChatCompletion( const delta = isRec(choice.delta) ? choice.delta : null; if (!delta) continue; if (typeof delta.content === "string") content = replaceRetained(content, content + delta.content, "retained_collectors"); + if (delta.refusal !== undefined && delta.refusal !== null) { + if (typeof delta.refusal !== "string") throw refusalTranslationError(); + refusal = replaceRetained(refusal ?? "", (refusal ?? "") + delta.refusal, "retained_collectors"); + } if (typeof delta.reasoning_content === "string") reasoning = replaceRetained(reasoning, reasoning + delta.reasoning_content, "retained_collectors"); if (Array.isArray(delta.tool_calls)) { for (const tc of delta.tool_calls) { @@ -785,6 +989,10 @@ export async function collectChatCompletion( } } } catch (error) { + // Processing may fail between reads; cancel while we still own the lock so the + // upstream translator releases its maps and stops any pending provider read. + try { await reader.cancel(error); } catch { /* preserve the original failure */ } + translatorBudget.releaseRetained(Buffer.byteLength(refusal ?? ""), { kind: "retained_collectors" }); // Never leak an open call scope on the error path; the turn budget's // dispose is a backstop, not the owner of this transfer. for (const index of toolCalls.keys()) translatorBudget.closeCall(callScope(index)); @@ -801,14 +1009,11 @@ export async function collectChatCompletion( } finally { reader.releaseLock(); } - if (streamError) { - for (const index of toolCalls.keys()) translatorBudget.closeCall(callScope(index)); - throw streamError; - } const message: Rec = { role: "assistant", content: content || null, + refusal, }; if (reasoning) message.reasoning_content = reasoning; if (toolCalls.size > 0) { diff --git a/src/server/chat-native-sse.ts b/src/server/chat-native-sse.ts index 6570e62a4b6..0d807330574 100644 --- a/src/server/chat-native-sse.ts +++ b/src/server/chat-native-sse.ts @@ -77,6 +77,7 @@ export function jsonCompletionSse(value: Rec, requestedModel: string, budget?: T }]; const delta: Rec = {}; if (typeof message.content === "string" && message.content.length > 0) delta.content = message.content; + if (typeof message.refusal === "string") delta.refusal = message.refusal; if (typeof message.reasoning_content === "string" && message.reasoning_content.length > 0) { delta.reasoning_content = message.reasoning_content; } diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index 0c87f3f50d0..f0c855bea1b 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -1323,6 +1323,25 @@ release on consumption or cancellation. Actual upstream SSE and native Chat bypa - 다른 대안 대신 이 방식을 선택한 이유: One conversion authority prevents the streaming fallback from drifting from non-streaming semantics without changing routing or retry behavior. - 장점, 단점 및 영향: No additional upstream request or dependency; this remains buffered delivery, not token-by-token upstream streaming. Handler regressions cover tools, reasoning, length, ordinary and empty completions, and budget release. +### Chat refusal projection + +`src/chat/outbound.ts` keeps Responses refusal parts separate from ordinary content. JSON output +and the stream collector expose nullable `message.refusal`; `jsonCompletionSse` preserves it as +`delta.refusal`, while the native SSE relay remains opaque. The translated live stream keys refusal +state by raw `output_index` / `content_index`, validates present item IDs as correlation constraints, +and emits buffered parts in that order only at a valid completed/incomplete terminal. Deltas append; +equal, empty, absent, and shorter-prefix snapshots preserve existing text; extending snapshots add +only new text. Non-string or contradictory snapshots fail with a content-free typed error. + +The existing turn budget accounts for refusal text and map metadata, including empty entries, and +releases that state on terminal, failure, or cancellation. Pending role/tool/refusal/finish/`[DONE]` +frames form one terminal batch: all serialized strings and encoded frames must be admitted before +any batch frame is enqueued. Admission failure releases the batch and refusal state, cancels upstream, +and emits only the bounded overflow error. Collector processing failures cancel their reader before +releasing its lock, so upstream translation cannot continue after failed JSON collection. The outer +response finalizer continues to own retained response bytes. These are projection rules, not new +refusal policy or changes to ordinary content/tool semantics. + ## Parallel tool calls (default-on for chat providers) The openai-chat adapter buffers ALL streamed `tool_calls` deltas (keyed by `index`, falling back to diff --git a/tests/responses/chat-refusal.test.ts b/tests/responses/chat-refusal.test.ts new file mode 100644 index 00000000000..4b22f690d13 --- /dev/null +++ b/tests/responses/chat-refusal.test.ts @@ -0,0 +1,451 @@ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { + ChatCompletionsStreamError, + collectChatCompletion, + responsesJsonToChatCompletion, + responsesSseToChatCompletionsSse, +} from "../../src/chat/outbound"; +import { jsonCompletionSse, nativeChatSse } from "../../src/server/chat-native-sse"; +import type { OcxConfig } from "../../src/types"; +import { createTestTranslatorBudget } from "../helpers/translator-budget"; +import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; +import { resetProviderRequestPacingForTest } from "../../src/providers/request-pacing"; + +type Rec = Record; +type Frame = { + error?: { code: string; type: string; message: string }; + choices?: Array<{ delta?: Rec; finish_reason?: string | null }>; +}; +const encoder = new TextEncoder(); +const model = "refusal-fixture/model"; +const event = (type: string, fields: Rec = {}): Rec => ({ type, ...fields }); +const part = (refusal: unknown): Rec => ({ type: "refusal", refusal }); +const message = (content: unknown[], id?: unknown): Rec => ({ + type: "message", role: "assistant", content, ...(id === undefined ? {} : { id }), +}); +const terminal = (output?: unknown[], reason?: string): Rec => event( + reason ? "response.incomplete" : "response.completed", + { response: { + status: reason ? "incomplete" : "completed", + ...(output ? { output } : {}), + ...(reason ? { incomplete_details: { reason } } : {}), + } }, +); +const refusalDelta = (delta: unknown, output_index = 0, content_index = 0, ids: Rec = {}): Rec => + event("response.refusal.delta", { output_index, content_index, delta, ...ids }); +const refusalDone = (refusal: unknown, output_index = 0, content_index = 0, ids: Rec = {}): Rec => + event("response.refusal.done", { output_index, content_index, refusal, ...ids }); +function wireEvent(value: Rec): string { + return `event: ${value.type}\ndata: ${JSON.stringify(value)}\n\n`; +} +function bytesSource(chunks: string[], onCancel = () => {}, close = true): ReadableStream { + let next = 0; + return new ReadableStream({ + pull(controller) { + if (next < chunks.length) controller.enqueue(encoder.encode(chunks[next++]!)); + else if (close) controller.close(); + }, + cancel: onCancel, + }, { highWaterMark: 0 }); +} +function translated(events: Rec[], budget = createTestTranslatorBudget(), onCancel = () => {}, close = true) { + return responsesSseToChatCompletionsSse(bytesSource(events.map(wireEvent), onCancel, close), model, { + translatorBudget: budget, + }); +} +function frames(wire: string): Frame[] { + return wire.split("\n\n").filter(block => block.startsWith("data: ") && block !== "data: [DONE]") + .map(block => JSON.parse(block.slice(6)) as Frame); +} +function refusalText(wire: string): string { + return frames(wire).map(frame => frame.choices?.[0]?.delta?.refusal ?? "").join(""); +} +function firstChoice(completion: Rec) { + return (completion.choices as Array<{ message: Rec; finish_reason: string }>)[0]!; +} +function expectSuccess(wire: string, refusal: string, reason = "stop") { + expect(refusalText(wire)).toBe(refusal); + expect(frames(wire).filter(frame => frame.error)).toHaveLength(0); + expect(frames(wire).filter(frame => frame.choices?.[0]?.finish_reason)).toEqual([ + expect.objectContaining({ choices: [{ index: 0, delta: {}, finish_reason: reason }] }), + ]); + expect(wire.match(/data: \[DONE\]/g)).toHaveLength(1); +} +function expectFailure(wire: string, code: string) { + expect(frames(wire).filter(frame => frame.error)).toEqual([ + { error: expect.objectContaining({ type: "upstream_error", code }) }, + ]); + expect(frames(wire).some(frame => frame.choices?.[0]?.finish_reason)).toBe(false); + expect(wire).not.toContain("data: [DONE]"); +} + +// Independent oracle: OpenAI SDK ChatCompletionMessage.refusal: string | null, +// Choice.Delta.refusal?: string | null; Responses refusal.delta.delta and +// refusal.done.refusal, keyed by raw output_index/content_index. All text is inert. +describe("Chat refusal projection", () => { + test("JSON keeps ordered refusal separate from answer, reasoning and tools", () => { + const completion = responsesJsonToChatCompletion({ output: [ + { type: "reasoning", summary: [{ type: "summary_text", text: "reason" }] }, + message([{ type: "output_text", text: "answer" }, part("fixture A"), part(" + B")]), + { type: "function_call", call_id: "call_fixture", name: "fixture", arguments: "{}" }, + message([part(" + C")]), + ] }, model); + expect(firstChoice(completion)).toMatchObject({ + message: { content: "answer", refusal: "fixture A + B + C", reasoning_content: "reason", + tool_calls: [{ id: "call_fixture", function: { name: "fixture", arguments: "{}" } }] }, + finish_reason: "tool_calls", + }); + expect(firstChoice(responsesJsonToChatCompletion({ output: [] }, model)).message.refusal).toBeNull(); + expect(firstChoice(responsesJsonToChatCompletion({ output: [message([part("")])] }, model)).message.refusal).toBe(""); + expect(() => responsesJsonToChatCompletion({ output: [message([part(null)])] }, model)) + .toThrow(ChatCompletionsStreamError); + }); + + test("split deltas and all repeated final representations contribute each suffix once", async () => { + const wire = await new Response(translated([ + event("response.output_item.added", { output_index: 2, item: message([], "item_fixture") }), + refusalDelta("fixture ", 2, 1, { item_id: "item_fixture" }), + refusalDelta("A", 2, 1), + refusalDone("fixture A", 2, 1), + event("response.content_part.done", { output_index: 2, content_index: 1, part: part("fixture AB") }), + event("response.output_item.done", { output_index: 2, + item: message([{ type: "output_text", text: "" }, part("fixture AB")], "item_fixture") }), + terminal([{}, { type: "reasoning" }, message([{}, part("fixture ABC")], "item_fixture")]), + ])).text(); + expectSuccess(wire, "fixture ABC"); + expect(frames(wire).filter(frame => frame.choices?.[0]?.delta?.refusal !== undefined)).toHaveLength(1); + }); + + for (const representation of ["done", "part", "item", "terminal"] as const) { + test(`${representation}-only refusal survives without deltas`, async () => { + const item = message([part("fixture")]); + const events = representation === "done" ? [refusalDone("fixture")] + : representation === "part" ? [event("response.content_part.done", { output_index: 0, content_index: 0, part: part("fixture") })] + : representation === "item" ? [event("response.output_item.done", { output_index: 0, item })] : []; + events.push(terminal(representation === "terminal" ? [item] : undefined)); + expectSuccess(await new Response(translated(events)).text(), "fixture"); + }); + } + + test("interleaved parts emit in raw output/content order and leave text live", async () => { + const wire = await new Response(translated([ + refusalDelta("C", 3, 0), refusalDelta("B", 1, 2), refusalDelta("A", 1, 0), + event("response.output_text.delta", { delta: "answer" }), + refusalDelta("2", 1, 2), refusalDelta("1", 1, 0), + terminal([{}, message([part("A1"), { type: "output_text", text: "answer" }, part("B2")]), {}, message([part("C")])]), + ])).text(); + expectSuccess(wire, "A1B2C"); + const deltas = frames(wire).flatMap(frame => frame.choices?.map(choice => choice.delta) ?? []); + expect(deltas.filter(delta => delta?.refusal !== undefined).map(delta => delta?.refusal)).toEqual(["A1", "B2", "C"]); + expect(deltas.filter(delta => delta?.content).map(delta => delta?.content)).toEqual(["answer"]); + expect(deltas.findIndex(delta => delta?.content === "answer")).toBeLessThan(deltas.findIndex(delta => delta?.refusal === "A1")); + }); + + test("missing, empty and stale-prefix snapshots preserve text and split Unicode", async () => { + const wire = await new Response(translated([ + refusalDelta("fixture \ud83d"), refusalDelta("\ude00"), + event("response.refusal.done", { output_index: 0, content_index: 0 }), + refusalDone(""), refusalDone("fixture"), + event("response.content_part.done", { output_index: 0, content_index: 0, part: { type: "refusal" } }), + event("response.output_item.done", { output_index: 0, item: message([]) }), + terminal([message([part("fixture ")])]), + ])).text(); + expectSuccess(wire, "fixture 😀"); + }); + + for (const reason of ["max_output_tokens", "content_filter"]) { + test(`valid incomplete ${reason} flushes refusal with truthful live finish`, async () => { + expectSuccess(await new Response(translated([ + refusalDelta("fixture"), terminal([message([part("fixture suffix")])], reason), + ])).text(), "fixture suffix", reason === "max_output_tokens" ? "length" : "content_filter"); + }); + } + + const invalidEvents: Array<[string, Rec]> = [ + ["contradictory done", refusalDone("other")], + ["nonstring delta", refusalDelta(42)], + ["nonstring done", refusalDone(null)], + ["nonstring content part", event("response.content_part.done", { output_index: 0, content_index: 0, part: part([]) })], + ["contradictory item", event("response.output_item.done", { output_index: 0, item: message([part("other")]) })], + ["contradictory terminal", terminal([message([part("other")])])], + ["nonstring terminal", terminal([message([part({})])])], + ["delta ID mismatch", refusalDelta("suffix", 0, 0, { item_id: "other" })], + ["nonstring event ID", refusalDone("fixture", 0, 0, { item_id: null })], + ["snapshot ID mismatch", terminal([message([part("fixture")], "other")])], + ["nonstring snapshot ID", terminal([message([part("fixture")], 5)])], + ["different part type", terminal([message([{ type: "output_text", text: "fixture" }])])], + ["negative position", refusalDelta("fixture", -1)], + ["fractional position", refusalDelta("fixture", 0, 0.5)], + ]; + for (const [label, invalid] of invalidEvents) { + test(`${label} fails without refusal or success terminal and cancels upstream`, async () => { + let cancelled = 0; + const wire = await new Response(translated([ + refusalDelta("fixture", 0, 0, { item_id: "item_fixture" }), invalid, terminal(), + ], createTestTranslatorBudget(), () => { cancelled++; }, false)).text(); + expectFailure(wire, "invalid_refusal"); + expect(refusalText(wire)).toBe(""); + expect(cancelled).toBe(1); + }); + } + + test("failure and unknown incomplete terminals discard buffered refusal", async () => { + for (const end of [terminal(undefined, "adapter_eof"), event("response.failed", { + response: { error: { message: "fixture failure" } }, + })]) { + const wire = await new Response(translated([refusalDelta("fixture"), end])).text(); + expect(refusalText(wire)).toBe(""); + expect(frames(wire).filter(frame => frame.error)).toHaveLength(1); + expect(wire).not.toContain("data: [DONE]"); + expect(frames(wire).some(frame => frame.choices?.[0]?.finish_reason)).toBe(false); + } + }); + + test("one item's optional ID constrains all of its content parts", async () => { + const wire = await new Response(translated([ + refusalDelta("A", 0, 0, { item_id: "first" }), + refusalDelta("B", 0, 1, { item_id: "second" }), terminal(), + ])).text(); + expectFailure(wire, "invalid_refusal"); + }); + + test("refusal text overflow is bounded and cancels the source", async () => { + let cancelled = 0; + const budget = createTestTranslatorBudget({ maxTurnBytes: 4096 }); + const wire = await new Response(translated([ + ...Array.from({ length: 50 }, () => refusalDelta("x".repeat(100))), terminal(), + ], budget, () => { cancelled++; }, false)).text(); + expectFailure(wire, "translation_buffer_limit"); + expect(budget.snapshot().highWaterBytes).toBeLessThanOrEqual(4096); + expect(cancelled).toBe(1); + }); + + test("zero-length parts consume metadata budget", async () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 2048 }); + let cancelled = 0; + const wire = await new Response(translated([ + ...Array.from({ length: 100 }, (_, index) => refusalDone("", 0, index)), terminal(), + ], budget, () => { cancelled++; }, false)).text(); + expectFailure(wire, "translation_buffer_limit"); + expect(budget.snapshot().overflows).toBe(1); + expect(cancelled).toBe(1); + }); + + test("small turn budget rejects the whole final batch, including pending role", async () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 900 }); + let cancelled = 0; + const wire = await new Response(translated([refusalDone("fixture"), terminal()], budget, + () => { cancelled++; }, false)).text(); + expectFailure(wire, "translation_buffer_limit"); + expect(frames(wire).filter(frame => frame.choices)).toHaveLength(0); + expect(cancelled).toBe(1); + }); + + test("a reservation failure at DONE cannot leak pending tool/refusal/finish frames", async () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 4096 }); + const reserve = budget.reserveTransient.bind(budget); + let rejectedDone = false; + budget.reserveTransient = (bytes, scope) => { + if (bytes === encoder.encode("data: [DONE]\n\n").byteLength) { + rejectedDone = true; + // Exhaust the real configured budget at this precise admission boundary. + return reserve(4097, scope); + } + return reserve(bytes, scope); + }; + let cancelled = 0; + const wire = await new Response(translated([ + event("response.output_item.added", { output_index: 0, + item: { type: "function_call", id: "tool_fixture", call_id: "call_fixture", name: "f", arguments: "{}" } }), + refusalDone("fixture", 1), terminal(), + ], budget, () => { cancelled++; }, false)).text(); + expect(rejectedDone).toBe(true); + expectFailure(wire, "translation_buffer_limit"); + expect(frames(wire).flatMap(frame => frame.choices ?? []).every(choice => + !choice.delta?.tool_calls && choice.delta?.refusal === undefined)).toBe(true); + expect(budget.snapshot().activeCalls).toBe(0); + expect(cancelled).toBe(1); + }); + + test("successful terminal flush preserves pending tools and releases charged metadata", async () => { + const budget = createTestTranslatorBudget(); + const charge = budget.chargeRetained.bind(budget); + const release = budget.releaseRetained.bind(budget); + let metadataCharged = 0; + let metadataReleased = 0; + budget.chargeRetained = (bytes, scope) => { + charge(bytes, scope); + if (scope.kind === "item_ids") metadataCharged += bytes; + }; + budget.releaseRetained = (bytes, scope) => { + release(bytes, scope); + if (scope.kind === "item_ids") metadataReleased += bytes; + }; + const wire = await new Response(translated([ + event("response.output_item.added", { output_index: 0, + item: { type: "function_call", id: "tool_fixture", call_id: "call_fixture", name: "fixture" } }), + event("response.function_call_arguments.delta", { item_id: "tool_fixture", delta: "{}" }), + refusalDelta("fixture", 2, 0, { item_id: "message_fixture" }), terminal(), + ], budget)).text(); + expectSuccess(wire, "fixture", "tool_calls"); + expect(frames(wire).flatMap(frame => frame.choices?.[0]?.delta?.tool_calls ?? [])).toEqual([ + { index: 0, id: "call_fixture", type: "function", function: { name: "fixture", arguments: "{}" } }, + ]); + expect(metadataCharged).toBeGreaterThan(0); + expect(metadataReleased).toBe(metadataCharged); + }); + + test("cancellation releases buffered refusal text and map metadata", async () => { + const budget = createTestTranslatorBudget(); + let cancelled = 0; + const reader = translated([ + refusalDelta("fixture"), event("response.heartbeat"), + ], budget, () => { cancelled++; }, false).getReader(); + await reader.read(); // Heartbeat role proves the preceding refusal was retained. + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + await reader.cancel(); + reader.releaseLock(); + expect(cancelled).toBe(1); + expect(budget.snapshot().currentBytes).toBe(0); + }); +}); + +describe("Chat refusal collection and native serialization", () => { + test("collector preserves nullable refusal and native JSON-to-SSE round trip", async () => { + const completion = responsesJsonToChatCompletion({ output: [message([part("fixture")])] }, model); + const wire = jsonCompletionSse(completion, model); + expectSuccess(wire, "fixture"); + const collected = await collectChatCompletion(bytesSource([wire]), model, createTestTranslatorBudget()); + expect(firstChoice(collected).message).toMatchObject({ content: null, refusal: "fixture" }); + const empty = await collectChatCompletion(bytesSource([jsonCompletionSse({ choices: [{ message: { content: "answer", refusal: null } }] }, model)]), model, createTestTranslatorBudget()); + expect(firstChoice(empty).message).toMatchObject({ content: "answer", refusal: null }); + }); + + test("absent-only refusal evidence stays null while explicit empty refusal stays empty", async () => { + for (const evidence of [{ type: "refusal" }, part("")]) { + const budget = createTestTranslatorBudget(); + const completion = await collectChatCompletion(translated([terminal([message([evidence])])], budget), model, budget); + expect(firstChoice(completion).message.refusal).toBe(Object.hasOwn(evidence, "refusal") ? "" : null); + } + }); + + test("translated SSE collection uses the same ordered refusal contract", async () => { + const budget = createTestTranslatorBudget(); + const completion = await collectChatCompletion(translated([ + refusalDelta("B", 1), refusalDelta("A", 0), terminal(), + ], budget), model, budget); + expect(firstChoice(completion).message).toMatchObject({ content: null, refusal: "AB" }); + }); + + test("native SSE relay leaves refusal deltas intact", async () => { + const budget = createTestTranslatorBudget(); + const wire = [ + 'data: {"choices":[{"index":0,"delta":{"content":"answer","refusal":"fixture"},"finish_reason":null}]}\n\n', + 'data: {"choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}\n\n', + "data: [DONE]\n\n", + ].join(""); + const relayed = nativeChatSse(bytesSource([wire]), { + requestedModel: model, translatorBudget: budget, signal: new AbortController().signal, onUsage() {}, + }); + expectSuccess(await new Response(relayed).text(), "fixture"); + }); + + test("collector processing overflow cancels its reader and never returns partial JSON", async () => { + let cancelled = 0; + const budget = createTestTranslatorBudget({ maxTurnBytes: 1024 }); + const reserve = budget.reserveTransient.bind(budget); + let failedKind = ""; + budget.reserveTransient = (bytes, scope) => { + try { return reserve(bytes, scope); } catch (error) { failedKind = scope.kind; throw error; } + }; + const chunk = `data: ${JSON.stringify({ choices: [{ delta: { refusal: "x".repeat(100) } }] })}\n\n`; + await expect(collectChatCompletion(bytesSource(Array(20).fill(chunk), () => { cancelled++; }, false), model, budget)) + .rejects.toMatchObject({ status: 502, type: "upstream_error", code: "translation_buffer_limit" }); + expect(failedKind).toBe("retained_collectors"); + expect(cancelled).toBe(1); + }); + + test("collector error cancellation reaches an upstream translator", async () => { + const translatorBudget = createTestTranslatorBudget(); + const collectorBudget = createTestTranslatorBudget({ maxTurnBytes: 100 }); + let cancelled = 0; + const stream = translated([refusalDelta("fixture"), event("response.heartbeat")], translatorBudget, + () => { cancelled++; }, false); + await expect(collectChatCompletion(stream, model, collectorBudget)) + .rejects.toMatchObject({ code: "translation_buffer_limit" }); + expect(cancelled).toBe(1); + expect(translatorBudget.snapshot().currentBytes).toBe(0); + }); + + test("malformed native refusal and typed error frames cancel without partial JSON", async () => { + for (const payload of [ + { choices: [{ delta: { refusal: 17 } }] }, + { error: { message: "fixture error", type: "upstream_error", code: "fixture_error" } }, + ]) { + let cancelled = 0; + await expect(collectChatCompletion(bytesSource([`data: ${JSON.stringify(payload)}\n\n`], () => { cancelled++; }, false), + model, createTestTranslatorBudget())).rejects.toBeInstanceOf(ChatCompletionsStreamError); + expect(cancelled).toBe(1); + } + }); +}); + +// Handler coverage uses only an external fetch stub, never a mocked converter/handler. +// It also exercises #3770's shared JSON fallback after the parent commit is applied. +describe("refusal handler delivery matrix", () => { + const originalFetch = globalThis.fetch; + let isolatedHome: IsolatedCodexHome | undefined; + let previousOcxHome: string | undefined; + beforeEach(() => { + previousOcxHome = process.env.OPENCODEX_HOME; + isolatedHome = installIsolatedCodexHome("ocx-refusal-fixture-"); + process.env.OPENCODEX_HOME = isolatedHome.path; + globalThis.fetch = (async () => { throw new Error("unstubbed external transport"); }) as typeof fetch; + }); + afterEach(() => { + globalThis.fetch = originalFetch; + resetProviderRequestPacingForTest(); + if (previousOcxHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousOcxHome; + isolatedHome?.restore(); + }); + for (const native of [true, false]) { + for (const upstreamSse of [true, false]) { + for (const clientSse of [true, false]) { + test(`${native ? "native" : "translated"} upstream ${upstreamSse ? "SSE" : "JSON"} -> client ${clientSse ? "SSE" : "JSON"}`, async () => { + const { handleChatCompletions } = await import("../../src/server/chat-completions"); + const responseJson = { id: "resp_fixture", status: "completed", output: [message([part("fixture")], "item_fixture")] }; + const chatJson = { id: "chatcmpl_fixture", object: "chat.completion", created: 1, model, + choices: [{ index: 0, message: { role: "assistant", content: null, refusal: "fixture" }, finish_reason: "stop" }] }; + const seen: string[] = []; + globalThis.fetch = (async (input: RequestInfo | URL) => { + const url = new URL(input instanceof Request ? input.url : String(input)); + expect(url.origin).toBe("https://refusal.example.test"); + seen.push(url.pathname); + if (!upstreamSse) return Response.json(native ? chatJson : responseJson); + const wire = native ? [ + 'data: {"choices":[{"index":0,"delta":{"refusal":"fixture"},"finish_reason":null}]}\n\n', + 'data: {"choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}\n\n', + "data: [DONE]\n\n", + ].join("") : [refusalDelta("fixture", 0, 0, { item_id: "item_fixture" }), terminal(responseJson.output)].map(wireEvent).join(""); + return new Response(wire, { headers: { "content-type": "text/event-stream" } }); + }) as typeof fetch; + const config = { + port: 0, defaultProvider: "refusal-fixture", providers: { + "refusal-fixture": { adapter: native ? "openai-chat" : "openai-responses", + baseUrl: "https://refusal.example.test/v1", apiKey: "fixture-key", authMode: "key" }, + }, + } as OcxConfig; + const response = await handleChatCompletions(new Request("http://localhost/v1/chat/completions", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model, stream: clientSse, messages: [{ role: "user", content: "inert fixture" }] }), + }), config, { model: "", provider: "" }); + expect(response.status).toBe(200); + if (clientSse) expectSuccess(await response.text(), "fixture"); + else expect(firstChoice(await response.json() as Rec).message).toMatchObject({ content: null, refusal: "fixture" }); + expect(seen).toEqual([native ? "/v1/chat/completions" : "/v1/responses"]); + }); + } + } + } +}); From fd10390985c0000929b219cfebb793e55fcf4160 Mon Sep 17 00:00:00 2001 From: t Date: Mon, 7 Sep 2026 02:01:57 +0900 Subject: [PATCH 3/5] test(chat): register refusal regression coverage --- scripts/test-layout/layout.json | 1 + tests/fixtures/test-layout-expected.json | 1 + 2 files changed, 2 insertions(+) diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 2e2b806deea..3d84e77bd27 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -271,6 +271,7 @@ "catalog-vision-sidecar-modalities.test.ts": "codex-integration", "chat-completions-endpoint.test.ts": "responses", "chat-json-sse-fallback.test.ts": "responses", + "chat-refusal.test.ts": "responses", "chatgpt-device-auth.test.ts": "oauth", "chatgpt-oauth.test.ts": "oauth", "chatgpt-token-expiry.test.ts": "oauth", diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index a1f0198e44a..dc1bb60a246 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -106,6 +106,7 @@ "catalog-vision-sidecar-modalities.test.ts": "codex-integration", "chat-completions-endpoint.test.ts": "responses", "chat-json-sse-fallback.test.ts": "responses", + "chat-refusal.test.ts": "responses", "chatgpt-device-auth.test.ts": "oauth", "chatgpt-oauth.test.ts": "oauth", "chatgpt-token-expiry.test.ts": "oauth", From 660967be4ec3f055fcdaa6225075b13faf87cd3f Mon Sep 17 00:00:00 2001 From: t Date: Mon, 7 Sep 2026 02:09:13 +0900 Subject: [PATCH 4/5] fix(chat): constrain sparse refusal identities across output positions --- src/chat/outbound.ts | 9 ++++++++- tests/responses/chat-refusal.test.ts | 22 ++++++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/src/chat/outbound.ts b/src/chat/outbound.ts index 7646b63efba..be9cbc7c11a 100644 --- a/src/chat/outbound.ts +++ b/src/chat/outbound.ts @@ -207,10 +207,12 @@ export function responsesSseToChatCompletionsSse( id?: string; parts: Map; }>(); + const refusalIndexById = new Map(); let refusalMetadataBytes = 0; let refusalTextBytes = 0; const releaseRefusals = () => { refusalItems.clear(); + refusalIndexById.clear(); translatorBudget.releaseRetained(refusalMetadataBytes, { kind: "item_ids" }); translatorBudget.releaseRetained(refusalTextBytes, { kind: "retained_collectors" }); refusalMetadataBytes = 0; @@ -238,10 +240,13 @@ export function responsesSseToChatCompletionsSse( refusalItems.set(index, item); } if (hasId && typeof candidate === "string") { + const knownIndex = refusalIndexById.get(candidate); + if (knownIndex !== undefined && knownIndex !== index) throw refusalTranslationError(); if (item.id !== undefined && item.id !== candidate) throw refusalTranslationError(); if (item.id === undefined) { - chargeRefusalMetadata(Buffer.byteLength(candidate)); + chargeRefusalMetadata(refusalEntryBytes + Buffer.byteLength(candidate)); item.id = candidate; + refusalIndexById.set(candidate, index); } } return item; @@ -284,6 +289,8 @@ export function responsesSseToChatCompletionsSse( }; const snapshotRefusalItem = (outputIndex: unknown, item: Rec) => { const existing = typeof outputIndex === "number" ? refusalItems.get(outputIndex) : undefined; + // Sparse final snapshots may omit type/content but cannot change a known ID. + if (existing && Object.hasOwn(item, "id")) refusalItem(outputIndex, item, "id"); if (item.type !== "message") { if (existing && existing.parts.size > 0 && item.type !== undefined) throw refusalTranslationError(); return; diff --git a/tests/responses/chat-refusal.test.ts b/tests/responses/chat-refusal.test.ts index 4b22f690d13..fc178da01ae 100644 --- a/tests/responses/chat-refusal.test.ts +++ b/tests/responses/chat-refusal.test.ts @@ -173,6 +173,9 @@ describe("Chat refusal projection", () => { ["nonstring event ID", refusalDone("fixture", 0, 0, { item_id: null })], ["snapshot ID mismatch", terminal([message([part("fixture")], "other")])], ["nonstring snapshot ID", terminal([message([part("fixture")], 5)])], + ["sparse snapshot ID mismatch", terminal([{ id: "other" }])], + ["sparse nonstring snapshot ID", terminal([{ id: null }])], + ["same ID at another position", terminal([{}, {}, message([part("fixture")], "item_fixture")])], ["different part type", terminal([message([{ type: "output_text", text: "fixture" }])])], ["negative position", refusalDelta("fixture", -1)], ["fractional position", refusalDelta("fixture", 0, 0.5)], @@ -449,3 +452,22 @@ describe("refusal handler delivery matrix", () => { } } }); + +test("sparse terminal preserves a matching refusal ID", async () => { + const wire = await new Response(translated([ + refusalDelta("fixture", 0, 0, { item_id: "item_fixture" }), + terminal([{ id: "item_fixture" }]), + ])).text(); + expectSuccess(wire, "fixture"); +}); + +test("JSON refusal charges joined surrogate bytes and rejects retained overflow", () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 16 }); + const completion = responsesJsonToChatCompletion({ output: [message([part("\ud83d"), part("\ude00")])] }, model, budget); + expect(firstChoice(completion).message.refusal).toBe("😀"); + expect(budget.snapshot().currentBytes).toBe(4); + const small = createTestTranslatorBudget({ maxTurnBytes: 4 }); + expect(() => responsesJsonToChatCompletion({ output: [message([part("fixture")])] }, model, small)).toThrow(); + expect(small.snapshot().overflows).toBe(1); + expect(small.snapshot().currentBytes).toBe(0); +}); From 14ee44bc2cb61d7c904ea6b164fe68b06b623aea Mon Sep 17 00:00:00 2001 From: t Date: Mon, 7 Sep 2026 02:12:54 +0900 Subject: [PATCH 5/5] fix(chat): validate known refusal IDs in sparse moved snapshots [skip ci] --- src/chat/outbound.ts | 5 ++++- tests/responses/chat-refusal.test.ts | 1 + 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/src/chat/outbound.ts b/src/chat/outbound.ts index be9cbc7c11a..e69d03f499d 100644 --- a/src/chat/outbound.ts +++ b/src/chat/outbound.ts @@ -290,7 +290,10 @@ export function responsesSseToChatCompletionsSse( const snapshotRefusalItem = (outputIndex: unknown, item: Rec) => { const existing = typeof outputIndex === "number" ? refusalItems.get(outputIndex) : undefined; // Sparse final snapshots may omit type/content but cannot change a known ID. - if (existing && Object.hasOwn(item, "id")) refusalItem(outputIndex, item, "id"); + if (Object.hasOwn(item, "id") + && (existing || (typeof item.id === "string" && refusalIndexById.has(item.id)))) { + refusalItem(outputIndex, item, "id"); + } if (item.type !== "message") { if (existing && existing.parts.size > 0 && item.type !== undefined) throw refusalTranslationError(); return; diff --git a/tests/responses/chat-refusal.test.ts b/tests/responses/chat-refusal.test.ts index fc178da01ae..27f6ff86efd 100644 --- a/tests/responses/chat-refusal.test.ts +++ b/tests/responses/chat-refusal.test.ts @@ -176,6 +176,7 @@ describe("Chat refusal projection", () => { ["sparse snapshot ID mismatch", terminal([{ id: "other" }])], ["sparse nonstring snapshot ID", terminal([{ id: null }])], ["same ID at another position", terminal([{}, {}, message([part("fixture")], "item_fixture")])], + ["sparse same ID at another position", terminal([{}, {}, { id: "item_fixture" }])], ["different part type", terminal([message([{ type: "output_text", text: "fixture" }])])], ["negative position", refusalDelta("fixture", -1)], ["fractional position", refusalDelta("fixture", 0, 0.5)],