Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 6 additions & 6 deletions src/providers/registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ export interface ProviderRegistryEntry {
* reframes as Responses events. Use only for upstreams whose streaming response
* can omit or indefinitely delay the terminal event.
*/
modelWebsocketUpstreamStreaming?: Record<string, boolean>;
modelResponsesUpstreamStreaming?: Record<string, boolean>;
/**
* Responses-API resource path for providers whose route is not `/v1/responses`.
* Unlike `modelWireDefaults` above, this IS seeded into saved config: it describes
Expand Down Expand Up @@ -1143,7 +1143,7 @@ export const PROVIDER_REGISTRY: readonly ProviderRegistryEntry[] = [
// DeepSeek's Codex Responses stream can deliver output without closing on the
// terminal event. Keep Codex on WebSocket, but use the provider's bounded JSON
// response upstream so the bridge can synthesize a complete WS event sequence.
modelWebsocketUpstreamStreaming: { "deepseek-v4-flash": false },
modelResponsesUpstreamStreaming: { "deepseek-v4-flash": false },
// DeepSeek's Responses route is `POST /responses` with no `/v1` segment. Without
// this the passthrough adapter falls back to its legacy `/v1/responses`
// construction and the wire above can never route.
Expand Down Expand Up @@ -1815,15 +1815,15 @@ export function providerModelWireDefault(
return wire !== undefined && allowedWires.has(wire) ? wire : undefined;
}

/** Resolve a registry-only upstream-streaming compatibility hint for WS turns. */
export function providerModelWebsocketUpstreamStreaming(
/** Resolve a registry-only upstream-streaming compatibility hint for Responses turns. */
export function providerModelResponsesUpstreamStreaming(
id: string,
provider: Pick<OcxProviderConfig, "baseUrl" | "adapter"> & Partial<Pick<OcxProviderConfig, "authMode">>,
modelId: string,
): boolean | undefined {
const entry = getProviderRegistryEntry(id);
if (!entry?.modelWebsocketUpstreamStreaming || !providerMatchesRegistryTransport(id, provider)) return undefined;
return entry.modelWebsocketUpstreamStreaming[modelId.trim().toLowerCase()];
if (!entry?.modelResponsesUpstreamStreaming || !providerMatchesRegistryTransport(id, provider)) return undefined;
return entry.modelResponsesUpstreamStreaming[modelId.trim().toLowerCase()];
}

/**
Expand Down
22 changes: 22 additions & 0 deletions src/server/responses-item-id-repair.ts
Original file line number Diff line number Diff line change
Expand Up @@ -222,3 +222,25 @@ export function hasResponsesItemIdRepair(config: ResponsesItemIdRepairConfig | u
|| (config?.message?.length ?? 0) > 0
|| (config?.reasoning?.length ?? 0) > 0;
}

/**
* Client-facing id normalization for a WHOLE bounded-JSON Responses object.
*
* The bounded-JSON policy (#875) answers a streaming client by synthesizing SSE
* from a completed JSON body, and reframes the same body into events for WS
* turns. Neither path goes through the SSE relay, so neither picks up the SSE
* item-id rewrite — a provider that needs id repair would get it on a streaming
* response and silently lose it the moment the reliability policy switched the
* upstream to bounded JSON. This applies the same rewrite to the object so all
* three paths agree. Raw recorded state is untouched: recording happens before
* any normalization.
*/
export function repairResponsesJsonItemIds(
response: Record<string, unknown>,
config: ResponsesItemIdRepairConfig,
budget?: TranslatorBudget,
): Record<string, unknown> {
const state = createRepairState(config, budget);
const rewritten = rewriteResponseSnapshot(state, response);
return rewritten.changed ? rewritten.response : response;
}
52 changes: 52 additions & 0 deletions src/server/responses-json-events.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
/**
* Shared bounded-JSON → Responses event sequence (#875): the same pure event
* list used by the WebSocket bridge (sendResponsesJsonAsEvents) and by the
* HTTP SSE synthesis for models whose reliability policy forces a bounded
* JSON upstream. One algorithm, two serializations — no duplicated drift.
*/

export type ResponsesJsonEventFrame = Record<string, unknown>;

/**
* The canonical minimal sequence Codex commits: response.created (empty
* output, in_progress) → one response.output_item.done per output item → a
* status-preserving terminal (completed / failed / incomplete).
*/
export function responsesJsonEventSequence(
response: Record<string, unknown>,
rewritePayload?: (payload: Record<string, unknown>) => Record<string, unknown>,
): ResponsesJsonEventFrame[] {
const rewrite = rewritePayload ?? ((payload: Record<string, unknown>) => payload);
const output = Array.isArray(response.output) ? response.output : [];
const finalStatus = response.status === "failed" || response.status === "incomplete"
? response.status
: "completed";
return [
rewrite({
type: "response.created",
response: { ...response, status: "in_progress", output: [] },
}),
...output.map((item, outputIndex) => rewrite({
type: "response.output_item.done",
output_index: outputIndex,
item,
})),
rewrite({
type: `response.${finalStatus}`,
response: { ...response, status: finalStatus },
}),
];
}

/**
* Serialize the event sequence as one SSE body with exactly one
* `data: [DONE]\n\n` trailer.
*/
export function responsesJsonToSseBody(
response: Record<string, unknown>,
rewritePayload?: (payload: Record<string, unknown>) => Record<string, unknown>,
): string {
const frames = responsesJsonEventSequence(response, rewritePayload)
.map(frame => `data: ${JSON.stringify(frame)}\n\n`);
return `${frames.join("")}data: [DONE]\n\n`;
}
70 changes: 64 additions & 6 deletions src/server/responses/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@ import { applyOpenAiVirtualModel, resolveOpenAiCompactModel } from "../../provid
import { isUsageDebugEnabled } from "../../usage/debug";
import { readJsonRequestBody, DecompressedBodyTooLargeError, UnsupportedContentEncodingError } from "../request-decompress";
import { resolveAdapter, resolveWireProtocolOverride } from "../adapter-resolve";
import { providerModelWebsocketUpstreamStreaming, type InboundWire } from "../../providers/registry";
import { providerModelResponsesUpstreamStreaming, type InboundWire } from "../../providers/registry";
import type { AdapterRequest } from "../../adapters/base";
import {
hasKeyPoolFailover,
Expand Down Expand Up @@ -163,6 +163,7 @@ import { cancelBodyOnAbort } from "../../lib/abort";
import {
createResponsesItemIdPayloadRewrite,
hasResponsesItemIdRepair,
repairResponsesJsonItemIds,
} from "../responses-item-id-repair";
import {
createImageGenCallRestoreRewrite,
Expand All @@ -187,6 +188,7 @@ import {
payloadRewriteAsBlockRewrite,
relaySseWithBlockRewrite,
} from "../sse-payload-rewrite";
import { responsesJsonToSseBody } from "../responses-json-events";
import { guardTerminalEventStream } from "./terminal-guard";

/**
Expand Down Expand Up @@ -860,9 +862,13 @@ async function applyFinalRouteRequestNormalization(args: {
}
parsed.modelId = route.modelId;
}
const websocketUpstreamStreaming = inboundTransport === "websocket"
? providerModelWebsocketUpstreamStreaming(route.providerName, route.provider, route.modelId)
: undefined;
// 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,
);

// Settle the wire once so logging, fast-mode, auth, and sidecars read the adapter
// this request will actually use (#404).
Expand All @@ -872,7 +878,7 @@ async function applyFinalRouteRequestNormalization(args: {
logCtx.providerAdapter = route.provider.adapter;
logCtx.routeDecision = route.routeDecision;

if (websocketUpstreamStreaming === false) {
if (responsesUpstreamStreaming === false && route.provider.adapter === "openai-responses") {
parsed.stream = false;
if (parsed._rawBody && typeof parsed._rawBody === "object") {
(parsed._rawBody as Record<string, unknown>).stream = false;
Expand Down Expand Up @@ -1501,6 +1507,11 @@ async function handleResponsesInner(
);
}

// Captured before normalization: whether the CLIENT asked for SSE. The
// transport-neutral upstream-streaming policy below may force a bounded JSON
// upstream for reliability (#875); the answer must then be reframed to SSE
// for streaming clients.
const clientRequestedStream = parsed.stream;
await applyFinalRouteRequestNormalization({
parsed,
route,
Expand Down Expand Up @@ -2235,7 +2246,54 @@ async function handleResponsesInner(
}
return repairResponsesSnapshotJson(restored, outbound);
})();
return new Response(clientJson, {
// #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 →
// terminal → [DONE]) so Codex commits the turn instead of hanging on a
// 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") {
try {
let completed = JSON.parse(clientJson) as Record<string, unknown>;
// 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)) {
completed = repairResponsesJsonItemIds(completed, route.provider.responsesItemIdRepair!, translatorBudget);
}
const sseHeaders = sanitizePassthroughHeaders(headers);
sseHeaders.set("content-type", "text/event-stream");
sseHeaders.set("cache-control", "no-store");
return new Response(responsesJsonToSseBody(completed), {
status: upstreamResponse.status,
statusText: upstreamResponse.statusText,
headers: sseHeaders,
});
} catch {
// Non-JSON despite content-type: fall through to the plain relay.
}
}
// 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<string, unknown>,
route.provider.responsesItemIdRepair!,
translatorBudget,
));
} catch {
return clientJson;
}
})()
: clientJson;
return new Response(outboundJson, {
status: upstreamResponse.status,
statusText: upstreamResponse.statusText,
headers,
Expand Down
20 changes: 4 additions & 16 deletions src/server/ws-bridge.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { ServerWebSocket } from "bun";
import { responsesJsonEventSequence } from "./responses-json-events";
import { FORWARD_HEADERS } from "../adapters/openai-responses";
import type { CodexAuthContext } from "../codex/auth-context";
import { headersForCodexAuthContext } from "../codex/auth-context";
Expand Down Expand Up @@ -307,25 +308,12 @@ export function sendResponsesJsonAsEvents(
}
sendTextFrame(ws, text);
};
const output = Array.isArray(response.output) ? response.output : [];
sendObservedFrame({
type: "response.created",
response: { ...response, status: "in_progress", output: [] },
});
output.forEach((item, outputIndex) => {
sendObservedFrame({
type: "response.output_item.done",
output_index: outputIndex,
item,
});
});
const finalStatus = response.status === "failed" || response.status === "incomplete"
? response.status
: "completed";
sendObservedFrame({
type: `response.${finalStatus}` as "response.completed" | "response.failed" | "response.incomplete",
response: { ...response, status: finalStatus },
});
for (const frame of responsesJsonEventSequence(response)) {
sendObservedFrame(frame);
}
onTerminal?.(finalStatus);
}

Expand Down
11 changes: 7 additions & 4 deletions structure/04_transports-and-sidecars.md
Original file line number Diff line number Diff line change
Expand Up @@ -206,10 +206,13 @@ 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 may keep the client-facing WebSocket while asking
the upstream Responses endpoint for bounded JSON. The bridge reframes that JSON into the same
Responses event sequence. DeepSeek V4 Flash uses this path because its Codex streaming response can
deliver output without closing on a terminal event; ordinary HTTP/SSE calls remain streaming.
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
(`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]`). DeepSeek V4 Flash uses this path because its Codex
streaming response can deliver output without closing on a terminal event.

`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
Expand Down
Loading
Loading