From 0b1d6f7ee78b94361b34f2316aa31904063a5b37 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Thu, 17 Sep 2026 03:00:00 +0000 Subject: [PATCH 1/4] feat(responses): relay native multi-agent function-result injection Follow up on #4782 with a separate default-off injection owner, bounded serial acknowledgements, same-account caller continuations, accepted-result replay and regression coverage. --- .../content/docs/guides/codex-integration.md | 74 ++++ .../docs/reference/configuration/server.md | 1 + scripts/test-layout/layout.json | 1 + src/config/schema/config-schema.ts | 1 + src/server/index/websocket-handler.ts | 16 +- src/server/responses/codex-ws-exchange.ts | 5 +- src/server/responses/core-options.ts | 4 +- src/server/responses/fetch-helpers.ts | 4 +- .../responses/native-injection-protocol.ts | 52 +++ .../responses/native-injection-replay.ts | 100 +++++ src/server/responses/native-injection.ts | 233 ++++++++++ .../responses/native-response-control.ts | 27 ++ src/server/responses/native-steering-log.ts | 2 +- src/server/responses/passthrough-dispatch.ts | 17 +- src/server/responses/ws-upstream.ts | 21 +- src/server/ws-bridge.ts | 4 +- src/types/config.ts | 2 + structure/adapters/registry.md | 2 + structure/catalog.md | 2 + structure/clients/claude-desktop.md | 2 + structure/config.md | 2 + structure/data-planes/images.md | 2 + structure/data-planes/inbound-compat.md | 2 + structure/gui-and-management-api.md | 2 + structure/ops/service-and-sidecars.md | 2 + structure/providers/xai-grok.md | 2 + structure/runtime.md | 2 + structure/subagents.md | 2 + structure/transports/byte-accounting.md | 2 + structure/transports/inventory.md | 2 + structure/transports/responses.md | 2 + structure/transports/streaming-health.md | 47 +++ tests/fixtures/test-layout-expected.json | 1 + tests/helpers/native-injection-fixture.ts | 98 +++++ tests/helpers/responses-core-source.ts | 3 + tests/responses/ws-native-injection.test.ts | 398 ++++++++++++++++++ 36 files changed, 1116 insertions(+), 23 deletions(-) create mode 100644 src/server/responses/native-injection-protocol.ts create mode 100644 src/server/responses/native-injection-replay.ts create mode 100644 src/server/responses/native-injection.ts create mode 100644 src/server/responses/native-response-control.ts create mode 100644 tests/helpers/native-injection-fixture.ts create mode 100644 tests/responses/ws-native-injection.test.ts diff --git a/docs-site/src/content/docs/guides/codex-integration.md b/docs-site/src/content/docs/guides/codex-integration.md index 774a4dc5e4f..165253b8704 100644 --- a/docs-site/src/content/docs/guides/codex-integration.md +++ b/docs-site/src/content/docs/guides/codex-integration.md @@ -920,3 +920,77 @@ The implementation has synthetic protocol and regression coverage, not live Astr certification. Keep the option disabled for production work until your client/model path has been verified. Set `codexNativeSteering` to `false` and restart to restore the existing single-response relay; no account or conversation files need to be deleted. + + +## Experimental native function-result injection + +For a compatible client that sends OpenAI multi-agent `response.inject` messages, +merge these keys into the existing OpenCodex configuration and restart before a +fresh turn. Do not replace your provider or account settings: + +```json +{ + "websockets": true, + "codexNativeInjection": true +} +``` + +The initial `response.create` must explicitly include `"multi_agent": { "enabled": true }`. +OpenCodex does not enable it based on a model name. A public OpenAI API provider +must use `adapter: "openai-responses"`, `baseUrl: "https://api.openai.com/v1"`, +its normal API-key authentication and `upstreamWebsocket: true`. Route the initial +model through that provider's configured prefix. The relay adds the required +`responses_multi_agent=v1` beta token on that public API connection only, preserving +other configured beta tokens. It does not substitute a subscription credential, +create an API account or automatically switch to a separately billed API. + +Canonical ChatGPT forward connections can opt into the same transport experimentally, +but the public API contract does **not** establish ChatGPT subscription or Codex +App/CLI support. A compatible upstream model and execution mode are still required. +See the [OpenAI multi-agent protocol](https://developers.openai.com/api/docs/guides/responses-multi-agent). + +Return a saved tool result after the matching developer function call has completed: + +```json +{ + "type": "response.inject", + "response_id": "resp_example", + "input": [ + { "type": "function_call_output", "call_id": "call_example", "output": "saved result" } + ] +} +``` + +Use the response/call IDs from the **same connection**, not these example IDs. +The first version accepts string-valued `function_call_output` only. User/system +messages, rich output arrays, hosted tools and simultaneous `response.steer` are +not accepted in an injection turn. Multiple saved function results can share a +single injection. Each call can be submitted only once, including while queued. + +Parallel tool results are queued and sent one frame at a time, since the success +event identifies the response rather than an individual injection. The relay +preserves `response.inject.created` and `response.inject.failed`. It keeps the +connection alive after a response terminal while submitted results await confirmation +or advertised calls await results, so late asynchronous results are not discarded. + +When the server rejects an injection with `response_already_completed`, use its +returned saved outputs in **one client-sent** `response.create` with the completed +`previous_response_id`, unchanged model/settings and the same lane. Include each +outstanding result exactly once; do not include already accepted outputs. The +relay keeps that continuation on the original account/socket and preserves normal +request pacing. It never runs the tool again or creates a recovery request itself. +Other failures remain visible for the client to handle. + +A missing acknowledgement or a disconnect means delivery can be **unknown**. Do +not automatically resend a result, restart a tool or change accounts to retry it. +The pending queue is limited to 32 frames and 8 MiB, with 1,024 advertised function +calls, a 32 MiB replay journal and at most 128 responses per owned connection. +Each sent injection has a 90-second acknowledgement deadline that unrelated output +cannot extend; a saved-result wait is limited to 30 minutes. Existing frame limits +and stall timeouts still apply. + +Translated providers, custom gateways, Combo/sidecar paths and HTTP fallback do +not gain injection support. Unsupported attempts return an explicit error instead +of disappearing. The option stays off by default; synthetic transport tests are +not live compatibility certification. Set `codexNativeInjection` to `false` and +restart to roll back. No account or conversation files need to be removed. diff --git a/docs-site/src/content/docs/reference/configuration/server.md b/docs-site/src/content/docs/reference/configuration/server.md index 7b93c8f3f72..d67c9995693 100644 --- a/docs-site/src/content/docs/reference/configuration/server.md +++ b/docs-site/src/content/docs/reference/configuration/server.md @@ -22,6 +22,7 @@ runs helper features around provider requests. | `shutdownTimeoutMs?` | `number` | `5000` | Graceful drain deadline before active turns are aborted. | | `websockets?` | `boolean` | `false` | Advertise and admit the client-facing Responses WebSocket path. False keeps clients on HTTP/SSE; it does not disable an eligible canonical ChatGPT upstream WS optimization. Complete-input requests may reuse an upstream connection within the same selected credential, account, thread and turn; changed handshake policy or missing identity keeps requests on separate connections. This does not trim HTTP input or create previous-response IDs. | | `codexNativeSteering?` | `boolean` | `false` | Experimental, native-only mid-turn steering on the Responses WebSocket endpoint. Requires `websockets: true`, a compatible upstream/client, and unchanged model/settings for saved-tool-result continuations. Does not enable translated models or HTTP fallback. See [native steering](/guides/codex-integration/#experimental-native-mid-turn-steering). | +| `codexNativeInjection?` | `boolean` | `false` | Experimental saved function-result injection on compatible native multi-agent WebSocket turns. Requires `websockets: true`, explicit `multi_agent.enabled`, and an eligible provider. Separate from steering; no automatic tool rerun or recovery create. See [native injection](/guides/codex-integration/#experimental-native-function-result-injection). | | `corsAllowOrigins?` | `string[]` | `[]` | Additional exact origins allowed by CORS. Loopback origins are always allowed. Authority-based browser extension origins such as `chrome-extension://` are supported; `*` is not a wildcard. Firefox and Safari regenerate the extension UUID (per install / per browser launch), so update the entry when the origin changes. | | `apiKeys?` | `OcxApiKey[]` | `[]` | Generated `ocx_…` credentials accepted by management and data-plane auth on non-loopback binds. Dashboard-managed. | | `storageCleanupPolicy?` | `StorageCleanupPolicy` | disabled | Opt-in archived-session cleanup policy. Never enabled implicitly. | diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index e9e550717dd..b6b17bafda6 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -1488,6 +1488,7 @@ "responses-account-change-scrub.test.ts": "responses", "response-log-inspection.test.ts": "server", "request-log-nonstream.test.ts": "usage", + "ws-native-injection.test.ts": "responses", "ws-native-steering.test.ts": "responses" }, "migrated": [ diff --git a/src/config/schema/config-schema.ts b/src/config/schema/config-schema.ts index 8d66ebe7321..ea8f7f72e84 100644 --- a/src/config/schema/config-schema.ts +++ b/src/config/schema/config-schema.ts @@ -60,6 +60,7 @@ import { DEFAULT_APP_OWNED_MEMORY_BUDGET_BYTES, MAX_APP_OWNED_MEMORY_BUDGET_MB, export const configSchema = z.object({ codexNativeSteering: z.boolean().optional().catch(false), + codexNativeInjection: z.boolean().optional().catch(false), port: z.number().int().min(0).max(65535).default(10100), // A malformed hand edit must disable only remote-role behavior, not discard // providers or data-plane keys. Live writes are rejected explicitly below. diff --git a/src/server/index/websocket-handler.ts b/src/server/index/websocket-handler.ts index f1c564c65a1..692a48c0a75 100644 --- a/src/server/index/websocket-handler.ts +++ b/src/server/index/websocket-handler.ts @@ -1,3 +1,6 @@ +import type { NativeResponseControl } from "../responses/native-response-control"; +import { NativeInjectionChannel } from "../responses/native-injection"; +import { isInjectionRequest } from "../responses/native-injection-protocol"; import { NativeSteeringChannel, NativeSteeringError } from "../responses/native-steering"; import { createNativeSteeringLogObserver } from "../responses/native-steering-log"; import type { Server, ServerWebSocket } from "bun"; @@ -190,8 +193,13 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { } catch { return; // text-only contract; ignore unparseable frames } - if (frame.type === "response.steer" || (frame.type === "response.create" && ws.data.nativeSteering)) { + if (frame.type === "response.inject" || frame.type === "response.steer" || (frame.type === "response.create" && ws.data.nativeSteering)) { try { + if (frame.type === "response.inject") { + if (!ws.data.nativeSteering?.inject) throw new NativeSteeringError("injection_not_supported", "Native injection is disabled or unavailable on this route."); + ws.data.nativeSteering.inject(frame); + return; + } if (frame.type === "response.steer") { if (!ws.data.nativeSteering) throw new NativeSteeringError("steering_not_supported", "Native steering is disabled or unavailable on this route."); ws.data.nativeSteering.steer(frame); @@ -214,11 +222,13 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { ws.data.cancel?.(); // A superseded turn must not keep ownership during warmup or refusal. ws.data.nativeSteering = undefined; - let nativeSteering: NativeSteeringChannel | undefined; + let nativeSteering: NativeResponseControl | undefined; try { const idleMs = typeof config.stallTimeoutSec === "number" && Number.isFinite(config.stallTimeoutSec) ? Math.max(1, config.stallTimeoutSec) * 1000 : 300_000; - nativeSteering = config.codexNativeSteering === true ? new NativeSteeringChannel(frame, idleMs) : undefined; + nativeSteering = config.codexNativeInjection === true && isInjectionRequest(frame) + ? new NativeInjectionChannel(frame, idleMs) + : config.codexNativeSteering === true ? new NativeSteeringChannel(frame, idleMs) : undefined; } catch { sendJsonFrame(ws, buildWsErrorFrame(400, { type: "invalid_request_error", message: "Invalid native steering request settings" })); return; diff --git a/src/server/responses/codex-ws-exchange.ts b/src/server/responses/codex-ws-exchange.ts index 213a91fd79f..0c686de150a 100644 --- a/src/server/responses/codex-ws-exchange.ts +++ b/src/server/responses/codex-ws-exchange.ts @@ -1,4 +1,5 @@ -import { markNativeSteeringResponse, type NativeSteeringChannel } from "./native-steering"; +import { markNativeSteeringResponse } from "./native-steering"; +import type { NativeResponseControl } from "./native-response-control"; import { MAX_CLIENT_SSE_FRAME_BYTES } from "../sse-frame-buffer"; import { isSafeResponseHeader } from "../safe-response-headers"; import { CodexWsMetadata, type CodexWsQuotaObserver } from "./codex-ws-metadata"; @@ -11,7 +12,7 @@ import { UPGRADE_DEADLINE_MS, CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPO type CodexWsFailureStage, type CodexWsStageRecord } from "./codex-ws-wire"; interface ExchangeOptions { - nativeSteering?: NativeSteeringChannel; + nativeSteering?: NativeResponseControl; beforeContinuation?: () => Promise; session: CodexWsSession; url: string; diff --git a/src/server/responses/core-options.ts b/src/server/responses/core-options.ts index 0c84494a284..fdebef3c4b7 100644 --- a/src/server/responses/core-options.ts +++ b/src/server/responses/core-options.ts @@ -1,4 +1,4 @@ -import type { NativeSteeringChannel } from "./native-steering"; +import type { NativeResponseControl } from "./native-response-control"; import type { OcxUsage, OcxProviderContinuationState, OcxConfig } from "../../types"; import type { CodexAuthPolicyConfig, CodexAuthContext } from "../../codex/auth-context"; import type { AdmissionLease } from "../../lib/admission"; @@ -53,7 +53,7 @@ export interface HandleResponsesOptions { onRequestBodyRead?: () => void; forceEmptyResponseId?: boolean; /** Internal, connection-owned control channel; never reconstructed from headers. */ - nativeSteering?: NativeSteeringChannel; + nativeSteering?: NativeResponseControl; abortSignal?: AbortSignal; /** One-shot TTFT callback: first non-empty model output observed (WP4). */ onFirstOutput?: () => void; diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index 2c867ba0d93..4a7ebb85d16 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -1,4 +1,4 @@ -import type { NativeSteeringChannel } from "./native-steering"; +import type { NativeResponseControl } from "./native-response-control"; import type { Server } from "bun"; import { codexWsUpstreamFetch, @@ -54,7 +54,7 @@ export interface PaceAwareFetch { export type ProviderFetch = typeof globalThis.fetch & PaceAwareFetch; export interface ProviderFetchOptions { - nativeSteering?: NativeSteeringChannel; + nativeSteering?: NativeResponseControl; providerName?: string; modelId?: string; /** One pacing slot was acquired immediately before this fetch wrapper was created. */ diff --git a/src/server/responses/native-injection-protocol.ts b/src/server/responses/native-injection-protocol.ts new file mode 100644 index 00000000000..6c88e37646c --- /dev/null +++ b/src/server/responses/native-injection-protocol.ts @@ -0,0 +1,52 @@ +import { createHash } from "node:crypto"; +import { CODEX_WS_ID_MAX_BYTES } from "./codex-ws-correlation"; +import { NativeSteeringError } from "./native-steering"; + +export type InjectionFrame = Record; +export type FunctionResult = { type: "function_call_output"; call_id: string; output: string }; +export const MAX_NATIVE_INJECTIONS = 32; +export const MAX_NATIVE_INJECTION_BYTES = 8 * 1024 * 1024; +export const MAX_NATIVE_INJECTION_CALLS = 1024; +export const NATIVE_INJECTION_ACK_MS = 90_000; +export const NATIVE_INJECTION_TOOL_MS = 30 * 60_000; + +/** Narrow a JSON object without accepting arrays or null. */ +export function injectionRecord(value: unknown): value is InjectionFrame { + return value !== null && typeof value === "object" && !Array.isArray(value); +} +/** Bound identities and exclude control characters, without changing their spelling. */ +export function injectionId(value: unknown): value is string { + return typeof value === "string" && value.length > 0 && Buffer.byteLength(value) <= CODEX_WS_ID_MAX_BYTES + && !/[\u0000-\u001f\u007f]/.test(value); +} +/** Stable setting comparison; only the digest is retained by the connection owner. */ +export function injectionFingerprint(value: unknown): string { + const canonical = (item: unknown): string => Array.isArray(item) ? `[${item.map(canonical).join(",")}]` + : injectionRecord(item) ? `{${Object.keys(item).sort().map(key => `${JSON.stringify(key)}:${canonical(item[key])}`).join(",")}}` + : JSON.stringify(item) ?? "null"; + return createHash("sha256").update(canonical(value)).digest("hex"); +} +/** Throw only fixed, content-free errors, never tool output or caller identifiers. */ +export function injectionError(code: string, message: string): never { + throw new NativeSteeringError(code, message); +} +/** First-version contract: nonempty arrays of string-valued, client-owned function results. */ +export function injectionResults(value: unknown): FunctionResult[] { + if (!Array.isArray(value) || !value.length || value.length > MAX_NATIVE_INJECTION_CALLS) { + injectionError("invalid_injection", "Supply a bounded, nonempty array of saved function results."); + } + const seen = new Set(); + for (const item of value) { + if (!injectionRecord(item) || item.type !== "function_call_output" || !injectionId(item.call_id) + || typeof item.output !== "string" || Object.keys(item).some(key => !["type", "call_id", "output"].includes(key))) { + injectionError("invalid_injection", "Only function_call_output with call_id and string output is supported; no messages or hosted-tool results."); + } + if (seen.has(item.call_id)) injectionError("duplicate_injection", "A function result must occur exactly once."); + seen.add(item.call_id); + } + return value as FunctionResult[]; +} +/** Detect explicit multi-agent opt-in without inferring it from the model name. */ +export function isInjectionRequest(frame: InjectionFrame): boolean { + return injectionRecord(frame.multi_agent) && frame.multi_agent.enabled === true; +} diff --git a/src/server/responses/native-injection-replay.ts b/src/server/responses/native-injection-replay.ts new file mode 100644 index 00000000000..3a71566c5ca --- /dev/null +++ b/src/server/responses/native-injection-replay.ts @@ -0,0 +1,100 @@ +import { MAX_NATIVE_STEERING_REPLAY_BYTES, type NativeSteeringReplayObserver } from "./native-steering-replay"; +import { injectionRecord as record, type InjectionFrame as Frame, type FunctionResult } from "./native-injection-protocol"; + +/** A bounded journal of accepted tool results, independent of user steering history. */ +export class NativeInjectionReplay implements NativeSteeringReplayObserver { + private prefix: unknown[]; + private bytes = 0; + private current?: string; + private output = new Map(); + private accepted = new Map(); + private pending?: FunctionResult[]; + private explicit: unknown[] = []; + private previous: unknown[] = []; + + /** Capture a private initial prefix; existing persistence eligibility is checked by the caller. */ + constructor(input: unknown, private readonly remember: (input: unknown[], response: Frame) => void) { + this.prefix = typeof input === "string" + ? [{ type: "message", role: "user", content: [{ type: "input_text", text: input }] }] + : Array.isArray(input) ? [...input] : []; + this.reserve(this.prefix); + } + /** Charge serialized bytes, refusing rather than truncating an over-budget transcript. */ + private reserve(value: unknown): number { + const bytes = Buffer.byteLength(JSON.stringify(value)); + if (this.bytes + bytes > MAX_NATIVE_STEERING_REPLAY_BYTES) throw new Error("Native injection replay exceeded its history budget."); + this.bytes += bytes; + return bytes; + } + /** Journal before physical send, with rollback usable only for a known unsent frame. */ + submitted(frame: Frame): () => void { + const input = Array.isArray(frame.input) ? frame.input : []; + const bytes = this.reserve(input); + if (frame.type === "response.inject") this.pending = input as FunctionResult[]; + else this.explicit = input; + return () => { + if (frame.type === "response.inject") this.pending = undefined; + else this.explicit = []; + this.bytes -= bytes; + }; + } + /** Keep wire output order and insert each accepted result after its owning function call. */ + private completedOutput(response: Frame): unknown[] { + const output = Array.isArray(response.output) && response.output.length + ? response.output : [...this.output.entries()].sort((a, b) => a[0] - b[0]).map(([, item]) => item); + const echoed = new Set(); + for (const item of output) { + if (!record(item) || item.type !== "function_call_output" || typeof item.call_id !== "string") continue; + const accepted = this.accepted.get(item.call_id); + if (accepted) { + if (echoed.has(item.call_id) || item.output !== accepted.output) throw new Error("Native injection replay result mismatch."); + echoed.add(item.call_id); + } + } + const merged: unknown[] = []; + const found = new Set(echoed); + for (const item of output) { + merged.push(item); + if (!record(item) || item.type !== "function_call" || typeof item.call_id !== "string") continue; + const accepted = this.accepted.get(item.call_id); + if (accepted && !found.has(item.call_id)) { merged.push(accepted); found.add(item.call_id); } + } + if (found.size !== this.accepted.size) throw new Error("Native injection replay is missing an accepted result's call."); + return merged; + } + /** Terminals are supplied by the owner only after all injection acknowledgements settle. */ + observe(frame: Frame): void { + if (frame.type === "response.created") { + if (this.current) { + for (const item of this.previous) this.prefix.push(item); + for (const item of this.explicit) this.prefix.push(item); + } + this.current = String(record(frame.response) ? frame.response.id : ""); + this.output.clear(); this.accepted.clear(); this.explicit = []; this.previous = []; + } else if (frame.type === "response.inject.created" || frame.type === "response.inject.failed") { + if (!this.pending) throw new Error("Native injection replay acknowledgement has no pending input."); + if (frame.type === "response.inject.created") { + for (const item of this.pending) this.accepted.set(item.call_id, item); + } else this.bytes -= Buffer.byteLength(JSON.stringify(this.pending)); + this.pending = undefined; + } else if (frame.type === "response.output_item.done") { + if (!Number.isSafeInteger(frame.output_index) || (frame.output_index as number) < 0 + || (frame.output_index as number) > 10_000 || !record(frame.item)) throw new Error("Native injection replay output identity is invalid."); + const old = this.output.get(frame.output_index as number); + if (old) this.bytes -= Buffer.byteLength(JSON.stringify(old)); + this.reserve(frame.item); this.output.set(frame.output_index as number, frame.item); + } else if (record(frame.response) && ["response.completed", "response.failed", "response.incomplete"].includes(String(frame.type))) { + if (this.pending) throw new Error("Native injection replay cannot commit an unacknowledged result."); + const output = this.completedOutput(frame.response); + for (const item of this.output.values()) this.bytes -= Buffer.byteLength(JSON.stringify(item)); + for (const item of this.accepted.values()) this.bytes -= Buffer.byteLength(JSON.stringify(item)); + this.reserve(output); this.output.clear(); this.accepted.clear(); this.previous = output; + if (frame.type === "response.completed") this.remember(this.prefix, { ...frame.response, output }); + } + } + /** Drop all retained bodies at cancellation, connection teardown or unknown delivery. */ + dispose(): void { + this.prefix = []; this.output.clear(); this.accepted.clear(); this.pending = undefined; + this.explicit = []; this.previous = []; this.bytes = 0; + } +} diff --git a/src/server/responses/native-injection.ts b/src/server/responses/native-injection.ts new file mode 100644 index 00000000000..09f762f8969 --- /dev/null +++ b/src/server/responses/native-injection.ts @@ -0,0 +1,233 @@ +import { CodexWsCorrelation } from "./codex-ws-correlation"; +import type { NativeResponseControl } from "./native-response-control"; +import type { NativeSteeringReplayObserver } from "./native-steering-replay"; +import { + injectionError, injectionFingerprint, injectionId, injectionRecord as record, injectionResults, + isInjectionRequest, MAX_NATIVE_INJECTIONS, MAX_NATIVE_INJECTION_BYTES, MAX_NATIVE_INJECTION_CALLS, + NATIVE_INJECTION_ACK_MS, NATIVE_INJECTION_TOOL_MS, + type FunctionResult, type InjectionFrame as Frame, +} from "./native-injection-protocol"; + +type Call = { itemId: unknown; state: "available" | "queued" | "accepted" | "failed"; result?: FunctionResult; recoverable?: boolean }; +type Submission = { frame: Frame; results: FunctionResult[]; bytes: number }; +const ENVELOPE = new Set(["type", "input", "previous_response_id", "stream", "stream_id"]); + +/** + * Injection-only multi-agent owner. Exactly one injection is on the physical wire + * at a time because created acknowledgements contain a response ID, not a request + * ID. Queued submissions, terminal delivery and caller-owned recovery are distinct. + */ +export class NativeInjectionChannel implements NativeResponseControl { + readonly kind = "injection" as const; + relayActive = false; + replayFactory?: () => NativeSteeringReplayObserver; + private replay?: NativeSteeringReplayObserver; + private send?: (frame: Frame) => void; + private onFailure?: (error: Error) => void; + private currentId?: string; + private correlation?: CodexWsCorrelation; + private terminal?: Frame; + private terminalRecorded = false; + private continuationSent = false; + private finished = false; + private everAttached = false; + private readonly seen = new Set(); + private readonly calls = new Map(); + private callBytes = 0; + private readonly queue: Submission[] = []; + private queueBytes = 0; + private inFlight?: Submission; + private lastAckSequence = -1; + private ackTimer?: ReturnType; + private idleTimer?: ReturnType; + private readonly settings = new Map(); + private readonly lane: unknown; + + /** Pin the original settings and lane; construction never opens a connection. */ + constructor(initial: Frame, private readonly idleMs = 300_000, + private readonly deadlines = { ackMs: NATIVE_INJECTION_ACK_MS, toolMs: NATIVE_INJECTION_TOOL_MS }) { + if (!isInjectionRequest(initial)) injectionError("injection_not_supported", "Native injection requires explicit multi_agent.enabled."); + this.lane = initial.stream_id ?? undefined; + for (const [key, value] of Object.entries(initial)) if (!ENVELOPE.has(key)) this.settings.set(key, injectionFingerprint(value)); + } + /** Report that the real dispatch boundary has selected this owner. */ + get attached(): boolean { return this.everAttached; } + /** A terminal is not final until submitted results have acknowledgements. */ + get ended(): boolean { return this.finished; } + + /** Attach once, after routing/auth/admission, retaining no global response-ID lookup. */ + attach(send: (frame: Frame) => void, fail: (error: Error) => void): () => void { + if (this.everAttached) throw new Error("Native injection transport is already owned."); + this.replay = this.replayFactory?.(); + this.send = send; this.onFailure = fail; this.everAttached = true; + return () => { + if (this.send !== send) return; + this.send = undefined; this.onFailure = undefined; this.finished = true; + clearTimeout(this.ackTimer); clearTimeout(this.idleTimer); + this.correlation?.finish(); this.calls.clear(); this.seen.clear(); this.queue.length = 0; + this.inFlight = undefined; this.queueBytes = 0; this.callBytes = 0; this.terminal = undefined; + this.replay?.dispose(); this.replay = undefined; + }; + } + /** Never reinterpret user steering as a function result, or silently mix beta modes. */ + steer(_frame: Frame): never { + return injectionError("native_control_mode_mismatch", "This multi-agent turn owns an injection-only channel; start a separate turn for steering."); + } + /** Abort unknown-delivery state without HTTP fallback, resends or invented acceptance. */ + private fail(): void { + this.finished = true; + clearTimeout(this.ackTimer); clearTimeout(this.idleTimer); + this.onFailure?.(new Error("Native injection transport failed or timed out; delivery is unknown. Do not automatically resend or rerun tools.")); + } + /** Require the same live owner; an unbound or detached channel cannot authorize a send. */ + private live(): void { + if (!this.send || this.finished) injectionError("injection_not_supported", "No live native injection transport is available on this route."); + } + /** Record only completed developer function calls, excluding hosted agent/tool actions. */ + private advertise(item: unknown): void { + if (!record(item) || item.type !== "function_call") return; + if (!injectionId(item.call_id)) throw new Error("Native injection function identity is invalid."); + const old = this.calls.get(item.call_id); + if (old) { + if (old.itemId !== item.id) throw new Error("Native injection function identity was reused."); + return; + } + const bytes = Buffer.byteLength(item.call_id) + (typeof item.id === "string" ? Buffer.byteLength(item.id) : 0); + if (this.calls.size >= MAX_NATIVE_INJECTION_CALLS || this.callBytes + bytes > 256 * 1024) throw new Error("Native injection call budget exceeded."); + this.calls.set(item.call_id, { itemId: item.id, state: "available" }); this.callBytes += bytes; + } + /** Queue validated saved results; reserve each call before any possibly synchronous send. */ + inject(frame: Frame): void { + this.live(); + if (frame.type !== "response.inject" || !injectionId(frame.response_id) + || Object.keys(frame).some(key => !["type", "response_id", "input", "stream_id"].includes(key))) { + injectionError("invalid_injection", "Invalid native injection envelope."); + } + if (frame.response_id !== this.currentId || (frame.stream_id ?? undefined) !== this.lane || this.continuationSent) { + injectionError("injection_response_mismatch", "Injection must target the current response and lane on this connection."); + } + const results = injectionResults(frame.input); + for (const item of results) { + const call = this.calls.get(item.call_id); + if (!call) injectionError("injection_call_not_found", "The result does not match a completed function call on this connection."); + if (call.state !== "available") injectionError("duplicate_injection", "This function result was already submitted; do not replay it."); + } + const text = JSON.stringify(frame); + const bytes = Buffer.byteLength(text); + if (this.queue.length >= MAX_NATIVE_INJECTIONS || bytes + this.queueBytes > MAX_NATIVE_INJECTION_BYTES) { + injectionError("injection_queue_full", "Native injection queue count or byte limit reached; no result was sent."); + } + // Detach from caller-owned objects before keeping data across asynchronous callbacks. + const copy = JSON.parse(text) as Frame; + const submission = { frame: copy, results: copy.input as FunctionResult[], bytes }; + for (const item of results) this.calls.get(item.call_id)!.state = "queued"; + this.queue.push(submission); this.queueBytes += bytes; + this.pump(); + } + /** Dispatch one queued frame; unrelated output cannot reset its acknowledgement deadline. */ + private pump(): void { + if (this.inFlight || !this.queue.length || this.finished) return; + this.live(); + const submission = this.queue[0]; + this.inFlight = submission; + this.ackTimer = setTimeout(() => this.fail(), this.deadlines.ackMs); + this.ackTimer.unref?.(); + try { + this.replay?.submitted(submission.frame); + this.send!(submission.frame); + } catch { + this.fail(); + injectionError("injection_delivery_unknown", "Injection dispatch failed; do not automatically resend or rerun tools."); + } + } + /** Match an acknowledgement to the sole in-flight frame, before releasing the next send. */ + private acknowledge(event: Frame): void { + const pending = this.inFlight; + if (!pending || event.response_id !== this.currentId || !Number.isSafeInteger(event.sequence_number) + || (event.sequence_number as number) <= this.lastAckSequence) throw new Error("Native injection acknowledgement identity mismatch."); + const failed = event.type === "response.inject.failed"; + if (failed && (!record(event.error) || typeof event.error.code !== "string" + || injectionFingerprint(event.input) !== injectionFingerprint(pending.results))) throw new Error("Native injection rejection does not match submitted results."); + this.replay?.observe(event); + this.lastAckSequence = event.sequence_number as number; + for (const item of pending.results) { + const call = this.calls.get(item.call_id)!; + call.state = failed ? "failed" : "accepted"; + // Retain a digest, not another result body, for an explicitly rejected continuation. + call.recoverable = failed && record(event.error) && event.error.code === "response_already_completed"; + if (call.recoverable) call.result = { ...item, output: injectionFingerprint(item.output) }; + } + clearTimeout(this.ackTimer); this.ackTimer = undefined; + this.inFlight = undefined; this.queue.shift(); this.queueBytes -= pending.bytes; + // Do not let a synchronous fake peer publish the next ack before this event is relayed. + if (this.queue.length) queueMicrotask(() => { try { this.pump(); } catch { this.fail(); } }); + } + /** Commit terminal replay only when no submitted injection can change its accepted inputs. */ + private recordTerminal(): void { + if (this.terminal && !this.terminalRecorded) { this.replay?.observe(this.terminal); this.terminalRecorded = true; } + } + /** Recover saved, explicitly rejected results on the same socket only when the client asks. */ + continue(frame: Frame): boolean { + if (!this.send || this.finished) return false; + if (this.queue.length || this.continuationSent) injectionError("injection_pending", "Wait for every injection acknowledgement before creating another response."); + if (!this.terminal || frame.previous_response_id !== this.currentId) return false; + if (this.terminal.type !== "response.completed") injectionError("injection_response_failed", "The parent response did not complete successfully."); + if ((frame.stream_id ?? undefined) !== this.lane || frame.generate === false) injectionError("invalid_injection", "Use the same lane for an injection continuation."); + for (const [key, value] of Object.entries(frame)) { + if (!ENVELOPE.has(key) && this.settings.get(key) !== injectionFingerprint(value)) injectionError("injection_settings_changed", "A native injection continuation cannot change the pinned model or settings."); + } + const results = injectionResults(frame.input); + const required = [...this.calls.entries()].filter(([, call]) => call.state !== "accepted"); + if (!required.length || results.length !== required.length) injectionError("invalid_injection", "Supply every outstanding saved function result exactly once."); + for (const item of results) { + const call = this.calls.get(item.call_id); + if (!call || call.state === "accepted" || call.state === "queued" + || (call.state === "failed" && (!call.recoverable || call.result?.output !== injectionFingerprint(item.output)))) { + injectionError("invalid_injection", "Continuation input must match unsent or explicitly completion-rejected function results."); + } + } + if (Buffer.byteLength(JSON.stringify(frame)) > MAX_NATIVE_INJECTION_BYTES) injectionError("invalid_injection", "Native injection continuation exceeds its byte limit."); + this.continuationSent = true; + try { this.recordTerminal(); this.replay?.submitted(frame); this.send(frame); } + catch { this.fail(); injectionError("injection_delivery_unknown", "Continuation delivery is unknown; do not automatically resend results."); } + if (!this.finished) this.armIdle(this.deadlines.ackMs); + return true; + } + /** Manage response/tool liveness independently of the non-resettable acknowledgement timer. */ + private armIdle(ms: number): void { + clearTimeout(this.idleTimer); + this.idleTimer = setTimeout(() => this.fail(), ms); this.idleTimer.unref?.(); + } + /** Validate ordered upstream events while allowing acknowledgements after a response terminal. */ + observe(event: Frame): boolean { + if ((event.stream_id ?? undefined) !== this.lane && !(event.type === "error" && event.stream_id == null)) throw new Error("Native injection lane mismatch."); + const type = event.type; + if (type === "error") { this.finished = true; return true; } + if (type === "response.inject.created" || type === "response.inject.failed") this.acknowledge(event); + else { + if (typeof type === "string" && (type.startsWith("response.inject.") || type.startsWith("response.steer."))) throw new Error("Unsupported native injection control event."); + const response = record(event.response) ? event.response : undefined; + if (type === "response.created") { + if (!injectionId(response?.id) || this.seen.has(response.id) || this.seen.size >= 128) throw new Error("Native injection response identity or chain limit violated."); + if (this.currentId && (!this.continuationSent || !this.terminal || this.queue.length + || (response.previous_response_id != null && response.previous_response_id !== this.currentId))) throw new Error("Unexpected native injection successor."); + this.currentId = response.id; this.seen.add(response.id); this.calls.clear(); this.callBytes = 0; + this.terminal = undefined; this.terminalRecorded = false; this.continuationSent = false; this.lastAckSequence = -1; + this.correlation?.finish(); this.correlation = new CodexWsCorrelation(true, () => false); + } else if (!this.currentId || this.terminal) throw new Error("Unexpected native injection event outside an active response."); + this.correlation?.accept({ ...event, stream_id: undefined }); + if (type === "response.output_item.done") this.advertise(event.item); + if (["response.completed", "response.failed", "response.incomplete"].includes(String(type))) { + if (!this.currentId || response?.id !== this.currentId) throw new Error("Native injection terminal identity mismatch."); + this.terminal = event; + if (Array.isArray(response.output)) for (const item of response.output) this.advertise(item); + } else this.replay?.observe(event); + } + const unresolved = [...this.calls.values()].some(call => call.state !== "accepted"); + this.finished = Boolean(this.terminal && !this.queue.length && !this.continuationSent + && (this.terminal.type !== "response.completed" || !unresolved)); + if (this.finished) { this.recordTerminal(); clearTimeout(this.idleTimer); } + else this.armIdle(this.continuationSent ? this.deadlines.ackMs : unresolved ? this.deadlines.toolMs : this.idleMs); + return this.finished; + } +} diff --git a/src/server/responses/native-response-control.ts b/src/server/responses/native-response-control.ts new file mode 100644 index 00000000000..516933ebfbb --- /dev/null +++ b/src/server/responses/native-response-control.ts @@ -0,0 +1,27 @@ +import type { OcxProviderConfig } from "../../types"; +import { isCanonicalOpenAiForwardProvider } from "../../providers/openai-tiers"; +import type { NativeSteeringReplayObserver } from "./native-steering-replay"; + +/** Shared transport ownership, not a shared steer/inject protocol state machine. */ +export interface NativeResponseControl { + readonly kind?: "steering" | "injection"; + relayActive: boolean; + replayFactory?: () => NativeSteeringReplayObserver; + readonly attached: boolean; + readonly ended: boolean; + attach(send: (frame: Record) => void, fail: (error: Error) => void): () => void; + observe(frame: Record): boolean; + steer(frame: Record): void; + inject?(frame: Record): void; + continue(frame: Record): boolean; +} + +export const OPENAI_API_RESPONSES_URL = "https://api.openai.com/v1/responses"; + +/** Preserve canonical ChatGPT eligibility; public API injection is separately opted in. */ +export function nativeResponseControlEligible(provider: OcxProviderConfig, control?: NativeResponseControl): boolean { + if (isCanonicalOpenAiForwardProvider(provider)) return true; + return control?.kind === "injection" && provider.adapter === "openai-responses" + && provider.upstreamWebsocket === true && provider.authMode !== "forward" + && provider.baseUrl?.replace(/\/+$/, "") === "https://api.openai.com/v1"; +} diff --git a/src/server/responses/native-steering-log.ts b/src/server/responses/native-steering-log.ts index d36911d127a..17f5a67dfae 100644 --- a/src/server/responses/native-steering-log.ts +++ b/src/server/responses/native-steering-log.ts @@ -9,7 +9,7 @@ export function createNativeSteeringLogObserver(logCtx: RequestLogContext, onFir return payload => { let event: { type?: string; delta?: unknown; response?: { id?: string; usage?: unknown; incomplete_details?: { reason?: string } } }; try { event = JSON.parse(payload); } catch { return; } - if (event.type?.startsWith("response.steer.")) return; + if (event.type?.startsWith("response.steer.") || event.type?.startsWith("response.inject.")) return; if (!outputSeen && event.type?.endsWith(".delta") && typeof event.delta === "string" && event.delta.length) { outputSeen = true; onFirstOutput?.(); } diff --git a/src/server/responses/passthrough-dispatch.ts b/src/server/responses/passthrough-dispatch.ts index 91633ce2d58..cd1dc0a9831 100644 --- a/src/server/responses/passthrough-dispatch.ts +++ b/src/server/responses/passthrough-dispatch.ts @@ -1,3 +1,5 @@ +import { nativeResponseControlEligible } from "./native-response-control"; +import { NativeInjectionReplay } from "./native-injection-replay"; import { NativeSteeringReplay } from "./native-steering-replay"; import type { ResponsesRequestContext, @@ -252,10 +254,11 @@ export async function preparePassthroughExchange( ? (response: { id?: unknown; output?: unknown; status?: unknown }) => rememberResponseState(parsed._rawBody, response, undefined, responseStateOptions(true)) : undefined; - if (options.nativeSteering && isCanonicalOpenAiForwardProvider(route.provider) + if (options.nativeSteering && nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt) { const body = parsed._rawBody as Record; - options.nativeSteering.replayFactory = () => new NativeSteeringReplay(body.input, (input, response) => { + const Replay = options.nativeSteering.kind === "injection" ? NativeInjectionReplay : NativeSteeringReplay; + options.nativeSteering.replayFactory = () => new Replay(body.input, (input, response) => { if (passthroughRecordEligible && !isBodyNonPersistable(body)) { rememberResponseState({ ...body, input }, response, undefined, responseStateOptions(true)); } @@ -782,7 +785,7 @@ export async function preparePassthroughExchange( body: request.body, }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: isCanonicalOpenAiForwardProvider(route.provider) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 ? options.nativeSteering : undefined, dispatchOverride: oauthDispatch(request), @@ -881,7 +884,7 @@ export async function preparePassthroughExchange( body: request.body, }, innerRecovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: isCanonicalOpenAiForwardProvider(route.provider) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 ? options.nativeSteering : undefined, dispatchOverride: oauthDispatch(request), @@ -990,7 +993,7 @@ export async function preparePassthroughExchange( // here on is a genuine transport attempt. storedPoolReplayDispatchNotifier( providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: isCanonicalOpenAiForwardProvider(route.provider) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 ? options.nativeSteering : undefined, dispatchOverride: oauthDispatch(request), @@ -1114,7 +1117,7 @@ export async function preparePassthroughExchange( body: request.body, }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: isCanonicalOpenAiForwardProvider(route.provider) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 ? options.nativeSteering : undefined, dispatchOverride: oauthDispatch(request), @@ -1239,7 +1242,7 @@ export async function preparePassthroughExchange( body: request.body, }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: isCanonicalOpenAiForwardProvider(route.provider) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 ? options.nativeSteering : undefined, dispatchOverride: oauthDispatch(request), diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index 31857e38d5b..a0fb1c9bd0c 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -1,4 +1,6 @@ -import type { NativeSteeringChannel } from "./native-steering"; +import { OPENAI_API_RESPONSES_URL } from "./native-response-control"; +import { isInjectionRequest } from "./native-injection-protocol"; +import type { NativeResponseControl } from "./native-response-control"; // Upstream WebSocket transport for the ChatGPT Codex backend. // // Why this exists: the Codex backend serves the responses_websockets path from @@ -130,7 +132,7 @@ export function codexWsUpstreamFetch( runtime: BunRuntimeGateInput = currentBunRuntimeIdentity(), onQuota?: CodexWsQuotaObserver, beforeDispatch?: (headers: Headers) => void, - nativeSteering?: NativeSteeringChannel, + nativeSteering?: NativeResponseControl, beforeContinuation?: () => Promise, ): Promise { const prepared = prepareCodexWsRequest(url, init); @@ -145,6 +147,17 @@ export function codexWsUpstreamFetch( } const { frameText, headers } = prepared; + // Never infer backend support from a model name or enable controls on a gateway. + const control = nativeSteering?.kind === "injection" + ? ((prepared.canonical || url === OPENAI_API_RESPONSES_URL) && isInjectionRequest(JSON.parse(frameText)) ? nativeSteering : undefined) + : prepared.canonical ? nativeSteering : undefined; + if (control?.kind === "injection" && url === OPENAI_API_RESPONSES_URL) { + const beta = headers["openai-beta"]; + if (!beta?.split(",").some(value => value.trim() === "responses_multi_agent=v1")) { + headers["openai-beta"] = beta ? `${beta}, responses_multi_agent=v1` : "responses_multi_agent=v1"; + } + } + // Decide before dialing. Once the socket is open the caller already holds a // streaming Response, so the oversized close can only be surfaced as a stream @@ -174,7 +187,7 @@ export function codexWsUpstreamFetch( try { // Steering keeps a private physical connection across successor responses; it // must never enter the idle-socket pool or move to a different credential. - const identity = nativeSteering && prepared.canonical ? null : codexWsReuseIdentity(url, headers, frameText, proxy); + const identity = control ? null : codexWsReuseIdentity(url, headers, frameText, proxy); session = (identity ? codexWsPool.acquire(identity, wsUrl, headers, proxy) : null) ?? new CodexWsSession(wsUrl, headers, false, undefined, proxy); if (!session.busy && !session.reserve()) { @@ -186,7 +199,7 @@ export function codexWsUpstreamFetch( } return codexWsExchange({ session, url, init, prepared, sseFallback, onQuota, beforeDispatch, - nativeSteering: prepared.canonical ? nativeSteering : undefined, + nativeSteering: control, beforeContinuation, bunVersion: typeof runtime === "string" ? runtime : runtime.version, }); diff --git a/src/server/ws-bridge.ts b/src/server/ws-bridge.ts index 451272e7073..fa7e0a4b652 100644 --- a/src/server/ws-bridge.ts +++ b/src/server/ws-bridge.ts @@ -1,4 +1,4 @@ -import type { NativeSteeringChannel } from "./responses/native-steering"; +import type { NativeResponseControl } from "./responses/native-response-control"; import type { ServerWebSocket } from "bun"; import { responsesJsonEventSequence } from "./responses-json-events"; import { FORWARD_HEADERS } from "../adapters/openai-responses"; @@ -18,7 +18,7 @@ type ResponsesTerminalReporter = (status: ResponsesTerminalStatus) => void; type ResponsesPayloadObserver = (payload: string) => void; export interface WsData { - nativeSteering?: NativeSteeringChannel; + nativeSteering?: NativeResponseControl; headers?: Headers; // base inbound forward headers only; per-turn auth refresh injects current pool tokens /** * Resolved once at the handshake. Auth is handshake-time only on this path, so diff --git a/src/types/config.ts b/src/types/config.ts index 3c427e808f4..59cdc64d29d 100644 --- a/src/types/config.ts +++ b/src/types/config.ts @@ -751,6 +751,8 @@ export interface OcxConfig { websockets?: boolean; /** Experimental single-lane native OpenAI WebSocket steering; default off. */ codexNativeSteering?: boolean; + /** Experimental, default-off saved function-result injection on native multi-agent WebSockets. */ + codexNativeInjection?: boolean; /** * Opt-in auto-cleanup policy for archived Codex sessions (issue #42 Phase 3). * Default OFF (`enabled` false / unset). Never enabled implicitly. diff --git a/structure/adapters/registry.md b/structure/adapters/registry.md index be6ebdaeec1..5bc604082cf 100644 --- a/structure/adapters/registry.md +++ b/structure/adapters/registry.md @@ -1,5 +1,7 @@ # Adapter Registry Authority +Native function-result injection follows [the separate opt-in control contract](../transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](../transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Request-local adapter bindings are separate from registry authority in the Responses diff --git a/structure/catalog.md b/structure/catalog.md index b42a190e690..d602e50a961 100644 --- a/structure/catalog.md +++ b/structure/catalog.md @@ -1,5 +1,7 @@ # Model Catalog +Native function-result injection follows [the separate opt-in control contract](transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Catalog discovery remains separate from the Responses final-route diff --git a/structure/clients/claude-desktop.md b/structure/clients/claude-desktop.md index 94caaf10684..b9e253082ec 100644 --- a/structure/clients/claude-desktop.md +++ b/structure/clients/claude-desktop.md @@ -1,5 +1,7 @@ # Claude Desktop Integration +Native function-result injection follows [the separate opt-in control contract](../transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](../transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Desktop callers retain their existing ingress through the Responses diff --git a/structure/config.md b/structure/config.md index 53a6112a9d1..7103bfe73bf 100644 --- a/structure/config.md +++ b/structure/config.md @@ -1,5 +1,7 @@ # Config Surface +Native function-result injection follows [the separate opt-in control contract](transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Catalog HTTP acquisition follows the [proxy-routing contract](catalog.md#remote-catalog-http-proxy-routing). diff --git a/structure/data-planes/images.md b/structure/data-planes/images.md index a621e175466..8281557d9c9 100644 --- a/structure/data-planes/images.md +++ b/structure/data-planes/images.md @@ -1,5 +1,7 @@ # Images Data Plane +Native function-result injection follows [the separate opt-in control contract](../transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](../transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Vision preprocessing and image/video/search execution use the Responses diff --git a/structure/data-planes/inbound-compat.md b/structure/data-planes/inbound-compat.md index 07eba5f5b97..cdef291e987 100644 --- a/structure/data-planes/inbound-compat.md +++ b/structure/data-planes/inbound-compat.md @@ -1,5 +1,7 @@ # Inbound Compatibility Surfaces +Native function-result injection follows [the separate opt-in control contract](../transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](../transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Compatibility callers retain the public Responses ingress described by the diff --git a/structure/gui-and-management-api.md b/structure/gui-and-management-api.md index cb3b7b23f36..eec2337dd5b 100644 --- a/structure/gui-and-management-api.md +++ b/structure/gui-and-management-api.md @@ -1,5 +1,7 @@ # GUI And Management API +Native function-result injection follows [the separate opt-in control contract](transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. The shared server request path follows the Responses diff --git a/structure/ops/service-and-sidecars.md b/structure/ops/service-and-sidecars.md index c9aced60cd4..2321a1db2e9 100644 --- a/structure/ops/service-and-sidecars.md +++ b/structure/ops/service-and-sidecars.md @@ -1,5 +1,7 @@ # Background Service And Sidecars +Native function-result injection follows [the separate opt-in control contract](../transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](../transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Service endpoints are unchanged by the Responses diff --git a/structure/providers/xai-grok.md b/structure/providers/xai-grok.md index 950942980b1..e35f40c637a 100644 --- a/structure/providers/xai-grok.md +++ b/structure/providers/xai-grok.md @@ -1,5 +1,7 @@ # xAI Grok Provider +Native function-result injection follows [the separate opt-in control contract](../transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](../transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. xAI uses the same shared credential and delivery policies through the Responses diff --git a/structure/runtime.md b/structure/runtime.md index a5b1fe7df03..7c29d16d926 100644 --- a/structure/runtime.md +++ b/structure/runtime.md @@ -1,5 +1,7 @@ # Runtime +Native function-result injection follows [the separate opt-in control contract](transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Responses admission and finalization are composed through the diff --git a/structure/subagents.md b/structure/subagents.md index fdb108d7437..76d7e551cd5 100644 --- a/structure/subagents.md +++ b/structure/subagents.md @@ -1,5 +1,7 @@ # Subagents And Multi-Agent Surface +Native function-result injection follows [the separate opt-in control contract](transports/streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](transports/streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Encrypted-task and fallback request handling follow the Responses diff --git a/structure/transports/byte-accounting.md b/structure/transports/byte-accounting.md index 0c2e59e2ba2..d28fc7d663a 100644 --- a/structure/transports/byte-accounting.md +++ b/structure/transports/byte-accounting.md @@ -1,5 +1,7 @@ # Byte Accounting +Native function-result injection follows [the separate opt-in control contract](streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. Responses body-reader limits and lifetime handling follow the diff --git a/structure/transports/inventory.md b/structure/transports/inventory.md index e820658bb91..29750a4d4bb 100644 --- a/structure/transports/inventory.md +++ b/structure/transports/inventory.md @@ -1,5 +1,7 @@ # Transport Inventory +Native function-result injection follows [the separate opt-in control contract](streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. The existing Responses transport is divided by responsibility in the diff --git a/structure/transports/responses.md b/structure/transports/responses.md index 60524b0b221..7fe0a98e174 100644 --- a/structure/transports/responses.md +++ b/structure/transports/responses.md @@ -1,5 +1,7 @@ # Responses Transport +Native function-result injection follows [the separate opt-in control contract](streaming-health.md#experimental-native-function-result-injection); this surface does not infer upstream support or alter its defaults. + Native steering follows [the shared WebSocket contract](streaming-health.md#experimental-native-mid-turn-steering); this surface's defaults remain unchanged. The configuration-only [plaintext V2 contract](../subagents.md#plaintext-v2-agent-messages) diff --git a/structure/transports/streaming-health.md b/structure/transports/streaming-health.md index c06436dd014..cd59e55da13 100644 --- a/structure/transports/streaming-health.md +++ b/structure/transports/streaming-health.md @@ -308,3 +308,50 @@ compatibility certification. End-to-end live client/backend verification remains before promoting this experimental option to a default. Shared response-log retention and native SSE inspection pacing follow the [bounded inspection contract](byte-accounting.md#response-log-inspection); other subsystem behavior remains unchanged. + + +## Experimental native function-result injection + +`codexNativeInjection` is a separate, default-off opt-in on the Responses WebSocket +ingress. An initial request must explicitly set `multi_agent.enabled: true`. The +canonical ChatGPT forward route remains experimental; public API injection requires +the exact `https://api.openai.com/v1` provider, non-forward authentication and +`upstreamWebsocket: true`. Only that public route adds `responses_multi_agent=v1` +to the outgoing beta header. No client/model capability or subscription entitlement +is inferred. Translated, Combo, sidecar, plaintext-restoration and HTTP-fallback +paths cannot receive controls. The common interface lives in +`src/server/responses/native-response-control.ts`; it shares transport ownership, +not protocol semantics, with steering. Mixed steer/inject turns are rejected. + +`src/server/responses/native-injection.ts` retains the normally selected credential +and private socket. `src/server/responses/native-injection-protocol.ts` validates +only string-valued developer `function_call_output` items for completed calls +advertised by that response and lane. IDs are never global lookup keys. One physical +injection awaits acknowledgement at a time because success carries a response ID, +not an injection ID; further submissions remain in a bounded FIFO. Repeated call +results, mismatched/repeated acknowledgements and unsupported shapes fail closed. + +A response terminal is relayed immediately, but pending acknowledgements and +unreturned advertised calls retain the socket. Late tool results still reach that +socket. A `response_already_completed` failure is relayed unchanged, including its +returned input; only an explicit same-parent/lane/settings client create can supply +those saved outputs once. The existing continuation pacing and captured dispatch +guard run again. The proxy never reruns tools, invents acceptance, switches accounts +or automatically creates a recovery response. Unknown delivery terminates without +HTTP fallback or replay. The client decides how to recover other failures. + +`src/server/responses/native-injection-replay.ts` commits accepted outputs only, +after all acknowledgements settle, inserting results after their owning calls and +preserving the original non-persistable-body policy. Failed inputs do not enter +continuation history. The existing numeric usage observer excludes inject events, +including echoed failed tool outputs, from log samples. All private bodies and +timers are disposed on teardown. + +Limits: 32 pending submissions including the on-wire frame, 8 MiB queued frame +bytes, 1,024 function identities / 256 KiB identity bytes, 128 response IDs and a +32 MiB replay journal. An on-wire injection has an absolute 90-second acknowledgement +deadline independent of incoming output; saved-tool-result waits use 30 minutes. +Existing socket/SSE frame limits and the active-response stall deadline also apply. +`tests/responses/ws-native-injection.test.ts` exercises the real handler, captured +auth, dispatch, relay, replay and synthetic failure paths. It is not live backend +or Codex App/CLI compatibility certification. diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index d0a98dd7627..0cbf4554fc4 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1320,5 +1320,6 @@ "responses-account-change-scrub.test.ts": "responses", "response-log-inspection.test.ts": "server", "request-log-nonstream.test.ts": "usage", + "ws-native-injection.test.ts": "responses", "ws-native-steering.test.ts": "responses" } diff --git a/tests/helpers/native-injection-fixture.ts b/tests/helpers/native-injection-fixture.ts new file mode 100644 index 00000000000..571c2bfe1d7 --- /dev/null +++ b/tests/helpers/native-injection-fixture.ts @@ -0,0 +1,98 @@ +import { afterEach, beforeEach, expect } from "bun:test"; +import type { ServerWebSocket } from "bun"; +import type { OcxConfig } from "../../src/types"; +import { createWebsocketHandler } from "../../src/server/index/websocket-handler"; +import type { ServeOptionsContext } from "../../src/server/index/serve-options"; +import type { WsData } from "../../src/server/ws-bridge"; +import { clearRequestLogsForTests } from "../../src/server/request-log"; +import { runOptionalShutdownHooks } from "../../src/lib/optional-shutdown-hooks"; + +export type Frame = Record; +const realSocket = globalThis.WebSocket; +const realFetch = globalThis.fetch; +const proxyKeys = ["HTTP_PROXY", "HTTPS_PROXY", "ALL_PROXY", "NO_PROXY", "http_proxy", "https_proxy", "all_proxy", "no_proxy"]; +let savedProxy: Record; +export let fallbackCalls = 0; +let nextId = 0; + +/** In-process upstream: all model traffic remains synthetic and network attempts fail. */ +export class InjectionSocket extends EventTarget { + static OPEN = 1; + static all: InjectionSocket[] = []; + readyState = 0; + frames: Frame[] = []; + readonly root = `inject-${++nextId}`; + throwOnInject = false; + constructor(readonly url: string, readonly options: { headers: Record }) { + super(); InjectionSocket.all.push(this); + queueMicrotask(() => { this.readyState = 1; this.dispatchEvent(new Event("open")); }); + } + send(text: string) { + const frame = JSON.parse(text); + if (frame.type === "response.inject" && this.throwOnInject) throw new Error("fixture send failure"); + this.frames.push(frame); + if (this.frames.length === 1) queueMicrotask(() => this.emit({ type: "response.created", response: { id: this.root, status: "in_progress", output: [] } })); + } + emit(frame: Frame) { + const lane = this.frames[0]?.stream_id; + this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify({ ...(lane !== undefined ? { stream_id: lane } : {}), ...frame }) })); + } + close() { if (this.readyState === 3) return; this.readyState = 3; this.dispatchEvent(new Event("close")); } +} +/** Public API and subscription fixtures have distinct, never-live credentials. */ +export const injectionConfig = (api = false): OcxConfig => ({ port: 0, defaultProvider: api ? "api" : "openai", websockets: true, codexNativeInjection: true, + providers: api + ? { api: { adapter: "openai-responses", baseUrl: "https://api.openai.com/v1", apiKey: "fixture-public-key", upstreamWebsocket: true, headers: { "openai-beta": "fixture_beta=v1" } } } + : { openai: { adapter: "openai-responses", baseUrl: "https://chatgpt.com/backend-api/codex", authMode: "forward", codexAccountMode: "direct" } }, +} as OcxConfig); +export const waitForInjection = async (condition: () => boolean) => { + for (let i = 0; i < 1000; i++) { if (condition()) return; await Bun.sleep(1); } + throw new Error("injection fixture condition timed out"); +}; +export function injectionClient(fields: Frame = {}, settings = injectionConfig(), credential = "test") { + const handler = createWebsocketHandler({ config: settings, deps: {} } as ServeOptionsContext); + const sent: Frame[] = []; + const ws = { readyState: 1, data: { headers: new Headers({ authorization: `Bearer ${credential}`, "thread-id": `injection-fixture-${++nextId}`, "openai-beta": "fixture_beta=v1" }) } as WsData, + send: (text: string) => { sent.push(JSON.parse(text)); return 1; }, close() { handler.close(ws, 1000, "fixture close"); }, + } as unknown as ServerWebSocket; + const send = (frame: Frame) => handler.message(ws, JSON.stringify(frame)); + send({ type: "response.create", model: settings.defaultProvider === "api" ? "api/gpt-5.6-sol" : "gpt-5.6-sol", input: "initial", multi_agent: { enabled: true }, + tools: [{ type: "function", name: "get_value", parameters: { type: "object", properties: {} } }], ...fields }); + return { ws, sent, send, handler }; +} +export async function beginInjection(fields: Frame = {}, settings = injectionConfig(), credential = "test") { + const client = injectionClient(fields, settings, credential); + await waitForInjection(() => client.sent.some(frame => frame.type === "response.created")); + const socket = InjectionSocket.all.at(-1)!; + expect(socket).toBeDefined(); + return { ...client, socket, id: socket.root }; +} +export function advertiseInjection(socket: InjectionSocket, call = "call-1", index = 0) { + const item = { id: `item-${call}`, type: "function_call", call_id: call, name: "get_value", arguments: "{}" }; + socket.emit({ type: "response.output_item.added", output_index: index, item }); + socket.emit({ type: "response.output_item.done", output_index: index, item }); + return item; +} +export const savedResult = (call = "call-1", output = "saved result") => ({ type: "function_call_output", call_id: call, output }); +export function completeInjection(socket: InjectionSocket, extra: Frame = {}, id = socket.root) { + socket.emit({ type: "response.completed", response: { id, status: "completed", output: [], ...extra } }); +} +export function acknowledgeInjection(socket: InjectionSocket, sequence = 100, id = socket.root) { + socket.emit({ type: "response.inject.created", response_id: id, sequence_number: sequence }); +} +export function installInjectionFixture() { + beforeEach(() => { + nextId = 0; fallbackCalls = 0; + savedProxy = Object.fromEntries(proxyKeys.map(key => [key, process.env[key]])); + for (const key of proxyKeys) delete process.env[key]; + globalThis.WebSocket = InjectionSocket as unknown as typeof WebSocket; + globalThis.fetch = (async () => { fallbackCalls++; throw new Error("network disabled in injection fixture"); }) as typeof fetch; + clearRequestLogsForTests(); + }); + afterEach(() => { + for (const socket of InjectionSocket.all) socket.close(); + InjectionSocket.all = []; runOptionalShutdownHooks(); + globalThis.WebSocket = realSocket; globalThis.fetch = realFetch; + for (const key of proxyKeys) { delete process.env[key]; if (savedProxy[key] !== undefined) process.env[key] = savedProxy[key]; } + }); +} diff --git a/tests/helpers/responses-core-source.ts b/tests/helpers/responses-core-source.ts index 5131188bd37..6a0223fc163 100644 --- a/tests/helpers/responses-core-source.ts +++ b/tests/helpers/responses-core-source.ts @@ -8,6 +8,9 @@ import { repoPath } from "./repo-root"; export const RESPONSES_CORE_MODULES = [ "core.ts", "core-options.ts", + "native-response-control.ts", + "native-injection-protocol.ts", + "native-injection-replay.ts", "native-steering.ts", "native-steering-replay.ts", "codex-ws-correlation.ts", diff --git a/tests/responses/ws-native-injection.test.ts b/tests/responses/ws-native-injection.test.ts new file mode 100644 index 00000000000..fe444ec993a --- /dev/null +++ b/tests/responses/ws-native-injection.test.ts @@ -0,0 +1,398 @@ +import { expect, test } from "bun:test"; +import { + installInjectionFixture, beginInjection, injectionClient, injectionConfig, InjectionSocket, + advertiseInjection, completeInjection, acknowledgeInjection, savedResult, waitForInjection, fallbackCalls, +} from "../helpers/native-injection-fixture"; +import { configSchema } from "../../src/config/schema/config-schema"; +import { getRequestLogEntries } from "../../src/server/request-log"; +import { createNativeSteeringLogObserver } from "../../src/server/responses/native-steering-log"; +import { NativeInjectionChannel } from "../../src/server/responses/native-injection"; +import { MAX_NATIVE_INJECTIONS, MAX_NATIVE_INJECTION_BYTES, injectionResults } from "../../src/server/responses/native-injection-protocol"; +import { nativeResponseControlEligible } from "../../src/server/responses/native-response-control"; +import { NativeInjectionReplay } from "../../src/server/responses/native-injection-replay"; +import type { RequestLogContext } from "../../src/server/request-log"; + +installInjectionFixture(); + +test("injection configuration is default-off, invalid values fail closed, and explicit opt-in survives parsing", () => { + const config = injectionConfig(); + expect(configSchema.parse(config).codexNativeInjection).toBe(true); + expect(configSchema.parse({ ...config, codexNativeInjection: "true" }).codexNativeInjection).toBe(false); + delete config.codexNativeInjection; + expect(configSchema.parse(config).codexNativeInjection).not.toBe(true); +}); + +test.each([false, true])("real handler sends saved results over the same connection (public API = %s)", async api => { + const { socket, send, sent, ws, id } = await beginInjection({}, injectionConfig(api)); + const call = advertiseInjection(socket); + const frame = { type: "response.inject", response_id: id, input: [savedResult()] }; + send(frame); + expect(socket.frames[1]).toEqual(frame); + acknowledgeInjection(socket); + completeInjection(socket, { output: [call], usage: { input_tokens: 10, output_tokens: 5 } }); + await waitForInjection(() => !ws.data.nativeSteering); + expect(sent.some(event => event.type === "response.inject.created")).toBe(true); + expect(sent.at(-1)?.type).toBe("response.completed"); + expect(InjectionSocket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); + expect(socket.options.headers.authorization).toBe(api ? "Bearer fixture-public-key" : "Bearer test"); + if (api) { + expect(socket.url).toBe("wss://api.openai.com/v1/responses"); + expect(socket.options.headers["openai-beta"]).toContain("responses_multi_agent=v1"); + expect(socket.options.headers["openai-beta"]).toContain("fixture_beta=v1"); + } else expect(socket.options.headers["openai-beta"]).not.toContain("responses_multi_agent=v1"); + expect(socket.frames[0].multi_agent).toEqual({ enabled: true }); + expect(getRequestLogEntries().at(-1)?.usage).toMatchObject({ inputTokens: 10, outputTokens: 5 }); +}); + +test("terminal before acknowledgement is relayed without dropping the late successful acknowledgement", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + const call = advertiseInjection(socket); + send({ type: "response.inject", response_id: id, input: [savedResult()] }); + completeInjection(socket, { output: [call] }); + await waitForInjection(() => sent.some(event => event.type === "response.completed")); + expect(socket.readyState).toBe(1); expect(ws.data.nativeSteering).toBeDefined(); + acknowledgeInjection(socket); + await waitForInjection(() => !ws.data.nativeSteering); + expect(sent.at(-1)?.type).toBe("response.inject.created"); + expect(socket.readyState).toBe(3); expect(fallbackCalls).toBe(0); +}); + +test("asynchronous tool completion after the response terminal still reaches the original socket", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + const call = advertiseInjection(socket); + completeInjection(socket, { output: [call] }); + await waitForInjection(() => sent.some(event => event.type === "response.completed")); + expect(socket.readyState).toBe(1); + send({ type: "response.inject", response_id: id, input: [savedResult()] }); + expect(socket.frames[1]?.type).toBe("response.inject"); + acknowledgeInjection(socket); + await waitForInjection(() => !ws.data.nativeSteering); + expect(sent.at(-1)?.type).toBe("response.inject.created"); +}); + +test("completion rejection is preserved; only an explicit caller continuation resubmits the saved result", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + const call = advertiseInjection(socket); + const input = [savedResult()]; + send({ type: "response.inject", response_id: id, input }); + completeInjection(socket, { output: [call] }); + const failed = { type: "response.inject.failed", response_id: id, sequence_number: 100, input, + error: { code: "response_already_completed", message: "upstream rejected completed response" } }; + socket.emit(failed); + await waitForInjection(() => sent.some(event => event.type === failed.type)); + expect(sent.at(-1)).toEqual(failed); expect(socket.frames).toHaveLength(2); + send({ type: "response.create", previous_response_id: id, input: [savedResult("call-1", "changed output")] }); + expect(sent.at(-1)?.error.code).toBe("invalid_injection"); expect(socket.frames).toHaveLength(2); + const continuation = { type: "response.create", previous_response_id: id, input }; + send(continuation); send(continuation); + await waitForInjection(() => socket.frames.length === 3); + expect(socket.frames[2].input).toEqual(input); expect(socket.frames[2].previous_response_id).toBe(id); + expect(socket.frames[2].multi_agent.enabled).toBe(true); + expect(sent.at(-1)?.error.code).toBe("injection_pending"); + socket.emit({ type: "response.created", response: { id: "successor", previous_response_id: id } }); + completeInjection(socket, {}, "successor"); + await waitForInjection(() => !ws.data.nativeSteering); + expect(InjectionSocket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); +}); + +test("parallel tool results serialize by acknowledgement without losing caller order", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + const calls = [advertiseInjection(socket), advertiseInjection(socket, "call-2", 1)]; + send({ type: "response.inject", response_id: id, input: [savedResult()] }); + send({ type: "response.inject", response_id: id, input: [savedResult("call-2")] }); + expect(socket.frames).toHaveLength(2); + completeInjection(socket, { output: calls }); + acknowledgeInjection(socket); + await waitForInjection(() => socket.frames.length === 3); + expect(socket.frames[2].input).toEqual([savedResult("call-2")]); + expect(ws.data.nativeSteering).toBeDefined(); + acknowledgeInjection(socket, 101); + await waitForInjection(() => !ws.data.nativeSteering); + expect(sent.filter(event => event.type === "response.inject.created")).toHaveLength(2); +}); + +test("duplicate results are refused both while pending and after successful acceptance", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + const call = advertiseInjection(socket); + const frame = { type: "response.inject", response_id: id, input: [savedResult()] }; + send(frame); send(frame); + expect(sent.at(-1)?.error.code).toBe("duplicate_injection"); expect(socket.frames).toHaveLength(2); + acknowledgeInjection(socket); send(frame); + expect(sent.at(-1)?.error.code).toBe("duplicate_injection"); expect(socket.frames).toHaveLength(2); + completeInjection(socket, { output: [call] }); + await waitForInjection(() => !ws.data.nativeSteering); +}); + +test("different response, lane, unadvertised and hosted-tool results never reach upstream", async () => { + const { socket, send, sent, ws, id } = await beginInjection({ stream_id: "lane-A" }); + advertiseInjection(socket); + const base = { type: "response.inject", response_id: id, stream_id: "lane-A", input: [savedResult()] }; + for (const frame of [ + { ...base, response_id: "foreign" }, { ...base, stream_id: "lane-B" }, + { ...base, input: [savedResult("foreign-call")] }, + { ...base, input: [{ type: "multi_agent_call_output", call_id: "hosted", output: "no" }] }, + { ...base, input: [{ type: "message", role: "system", content: "no" }] }, + ]) { send(frame); expect(sent.at(-1)?.type).toBe("error"); } + expect(socket.frames).toHaveLength(1); + ws.close(); await waitForInjection(() => !ws.data.nativeSteering); +}); + +test("two connections cannot inject results into each other's response or credentials", async () => { + const a = await beginInjection({}, injectionConfig(), "fixture-account-A"); + const b = await beginInjection({}, injectionConfig(), "fixture-account-B"); + advertiseInjection(a.socket); advertiseInjection(b.socket); + a.send({ type: "response.inject", response_id: b.id, input: [savedResult()] }); + expect(a.sent.at(-1)?.error.code).toBe("injection_response_mismatch"); + expect(a.socket.frames).toHaveLength(1); expect(b.socket.frames).toHaveLength(1); + a.ws.close(); b.ws.close(); + await waitForInjection(() => !a.ws.data.nativeSteering && !b.ws.data.nativeSteering); +}); + +test.each(["disabled", "no-multi-agent", "warmup", "steering-only"])("unsupported %s does not silently discard an injection", async mode => { + const cfg = injectionConfig(); + if (mode === "disabled" || mode === "steering-only") cfg.codexNativeInjection = false; + if (mode === "steering-only") cfg.codexNativeSteering = true; + const client = injectionClient(mode === "warmup" ? { generate: false } : mode === "no-multi-agent" ? { multi_agent: { enabled: false } } : {}, cfg); + await waitForInjection(() => client.sent.some(event => event.type === "response.created")); + client.send({ type: "response.inject", response_id: "unused", input: [savedResult()] }); + expect(client.sent.at(-1)?.error.code).toBe("injection_not_supported"); + expect(InjectionSocket.all.every(socket => socket.frames.length === 1)).toBe(true); + client.ws.close(); +}); + +test("a pending injection prevents a new create from cancelling the owned socket", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + advertiseInjection(socket); + send({ type: "response.inject", response_id: id, input: [savedResult()] }); + send({ type: "response.create", model: "different-model", input: "new work" }); + expect(sent.at(-1)?.error.code).toBe("injection_pending"); + expect(socket.readyState).toBe(1); expect(socket.frames).toHaveLength(2); + ws.close(); await waitForInjection(() => !ws.data.nativeSteering); +}); + +test("unknown delivery closes without HTTP fallback or resending the control", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + advertiseInjection(socket); socket.throwOnInject = true; + send({ type: "response.inject", response_id: id, input: [savedResult()] }); + await waitForInjection(() => !ws.data.nativeSteering); + expect(socket.readyState).toBe(3); expect(fallbackCalls).toBe(0); + expect(InjectionSocket.all).toHaveLength(1); + expect(sent.some(event => event.error?.message?.includes("unknown"))).toBe(true); +}); + +test("control failures are never sampled into the request log or counted as response usage", () => { + const log = { model: "fixture", provider: "fixture" } as RequestLogContext; + const inspect = createNativeSteeringLogObserver(log); + inspect(JSON.stringify({ type: "response.inject.failed", input: [savedResult("call-1", "PRIVATE_TOOL_RESULT")], error: { code: "x", message: "PRIVATE_TOOL_RESULT" } })); + expect(JSON.stringify(log)).not.toContain("PRIVATE_TOOL_RESULT"); + expect(log.usage).toBeUndefined(); expect(log.upstreamError).toBeUndefined(); +}); + +test("accepted function results survive ordinary subsequent delta turns; no user message is invented", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + const call = advertiseInjection(socket); + send({ type: "response.inject", response_id: id, input: [savedResult("call-1", "accepted-result")] }); + acknowledgeInjection(socket); completeInjection(socket, { output: [call] }); + await waitForInjection(() => !ws.data.nativeSteering); + send({ type: "response.create", model: "gpt-5.6-sol", multi_agent: { enabled: true }, previous_response_id: id, input: "followup" }); + await waitForInjection(() => InjectionSocket.all.length === 2 && InjectionSocket.all[1].frames.length > 0); + const next = InjectionSocket.all[1]; + const history = next.frames[0].input as Array>; + expect(history.filter(item => item.type === "function_call_output")).toEqual([savedResult("call-1", "accepted-result")]); + expect(history.findIndex(item => item.type === "function_call_output")).toBe(history.findIndex(item => item.type === "function_call") + 1); + expect(JSON.stringify(history)).toContain("followup"); + completeInjection(next); await waitForInjection(() => !ws.data.nativeSteering); + expect(sent.at(-1)?.type).toBe("response.completed"); +}); + +/** Unit owner fixture uses the same event contract without a network or application home. */ +function unitChannel(deadlines = { ackMs: 90_000, toolMs: 1_800_000 }) { + const sent: Array> = []; + const failures: Error[] = []; + const channel = new NativeInjectionChannel({ multi_agent: { enabled: true }, model: "fixture" }, 1000, deadlines); + const detach = channel.attach(frame => sent.push(frame), error => failures.push(error)); + channel.observe({ type: "response.created", response: { id: "root" } }); + const advertise = (call: string, index = 0) => { + const item = { id: `item-${call}`, type: "function_call", call_id: call, name: "fixture", arguments: "{}" }; + channel.observe({ type: "response.output_item.added", output_index: index, item }); + channel.observe({ type: "response.output_item.done", output_index: index, item }); + }; + return { channel, sent, failures, detach, advertise }; +} + +test("injection queue counts include the in-flight frame and refuse the next frame without sending it", () => { + const { channel, advertise, sent, detach } = unitChannel(); + try { + for (let i = 0; i <= MAX_NATIVE_INJECTIONS; i++) advertise(`c${i}`, i); + for (let i = 0; i < MAX_NATIVE_INJECTIONS; i++) channel.inject({ type: "response.inject", response_id: "root", input: [savedResult(`c${i}`)] }); + expect(sent).toHaveLength(1); + expect(() => channel.inject({ type: "response.inject", response_id: "root", input: [savedResult(`c${MAX_NATIVE_INJECTIONS}`)] })).toThrow("limit reached"); + expect(sent).toHaveLength(1); + } finally { detach(); } +}); + +test("serialized-byte cap rejects oversized output before a physical send", () => { + const { channel, advertise, sent, detach } = unitChannel(); + try { + advertise("c"); + expect(() => channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c", "x".repeat(MAX_NATIVE_INJECTION_BYTES))] })).toThrow("byte limit"); + expect(sent).toHaveLength(0); + channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c", "small")] }); + expect(sent).toHaveLength(1); // refusal did not consume the call's reservation + } finally { detach(); } +}); + +test("ack deadline remains absolute even while unrelated valid output keeps arriving", async () => { + const { channel, advertise, sent, failures, detach } = unitChannel({ ackMs: 20, toolMs: 1000 }); + try { + advertise("c"); channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c")] }); + for (let i = 0; i < 10 && !failures.length; i++) { + await Bun.sleep(5); + if (!channel.ended) channel.observe({ type: "response.in_progress", response: { id: "root" } }); + } + expect(failures).toHaveLength(1); expect(failures[0].message).toContain("delivery is unknown"); + expect(sent).toHaveLength(1); expect(channel.ended).toBe(true); + } finally { detach(); } +}); + +test("waiting for an asynchronously computed result is bounded and does not execute the tool", async () => { + const { channel, advertise, sent, failures, detach } = unitChannel({ ackMs: 1000, toolMs: 10 }); + try { + advertise("c"); + expect(channel.observe({ type: "response.completed", response: { id: "root", status: "completed", output: [] } })).toBe(false); + await waitForInjection(() => failures.length > 0); + expect(sent).toHaveLength(0); expect(failures).toHaveLength(1); + } finally { detach(); } +}); + +test.each([ + { type: "response.inject.created", response_id: "foreign", sequence_number: 10 }, + { type: "response.inject.created", response_id: "root", sequence_number: -1 }, + { type: "response.inject.created", response_id: "root" }, + { type: "response.inject.failed", response_id: "root", sequence_number: 10, input: [savedResult("c", "different")], error: { code: "response_already_completed" } }, +])("unknown or mismatched acknowledgements are rejected without committing results", event => { + const { channel, advertise, sent, detach } = unitChannel(); + try { + advertise("c"); channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c")] }); + expect(() => channel.observe(event)).toThrow(); expect(sent).toHaveLength(1); + } finally { detach(); } +}); + +test("a repeated acknowledgement cannot consume the next queued injection", async () => { + const { channel, advertise, sent, detach } = unitChannel(); + try { + advertise("c1"); advertise("c2", 1); + for (const call of ["c1", "c2"]) channel.inject({ type: "response.inject", response_id: "root", input: [savedResult(call)] }); + const ack = { type: "response.inject.created", response_id: "root", sequence_number: 10 }; + channel.observe(ack); await Bun.sleep(0); + expect(sent).toHaveLength(2); + expect(() => channel.observe(ack)).toThrow("identity mismatch"); + } finally { detach(); } +}); + +test("submitted objects are copied so later caller mutation cannot alter a queued send", async () => { + const { channel, advertise, sent, detach } = unitChannel(); + try { + advertise("c1"); advertise("c2", 1); + channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c1")] }); + const result = savedResult("c2", "original"); + channel.inject({ type: "response.inject", response_id: "root", input: [result] }); + result.output = "mutated"; + channel.observe({ type: "response.inject.created", response_id: "root", sequence_number: 10 }); + await Bun.sleep(0); + expect(sent[1].input).toEqual([savedResult("c2", "original")]); + } finally { detach(); } +}); + +test("batch results receive one acknowledgement and cannot reserve a call twice", () => { + const { channel, advertise, sent, detach } = unitChannel(); + try { + advertise("c1"); advertise("c2", 1); + expect(() => channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c1"), savedResult("c1")] })).toThrow("exactly once"); + channel.inject({ type: "response.inject", response_id: "root", input: [savedResult("c1"), savedResult("c2")] }); + expect(sent).toHaveLength(1); + channel.observe({ type: "response.inject.created", response_id: "root", sequence_number: 10 }); + expect(channel.observe({ type: "response.completed", response: { id: "root", status: "completed", output: [] } })).toBe(true); + } finally { detach(); } +}); + +test("tool-result validation rejects empty arrays, privileged roles, extra fields and unsupported rich outputs", () => { + for (const input of [[], "text", null, [savedResult("bad\n")], [{ ...savedResult(), role: "system" }], [{ ...savedResult(), output: [] }]]) { + expect(() => injectionResults(input)).toThrow(); + } +}); + +test("public injection excludes custom gateways, forwarded auth and an unopted API provider", () => { + const channel = new NativeInjectionChannel({ multi_agent: { enabled: true } }); + const provider = injectionConfig(true).providers.api; + expect(nativeResponseControlEligible(provider, channel)).toBe(true); + expect(nativeResponseControlEligible({ ...provider, baseUrl: "https://api.openai.com.attacker.invalid/v1" }, channel)).toBe(false); + expect(nativeResponseControlEligible({ ...provider, baseUrl: "http://api.openai.com/v1" }, channel)).toBe(false); + expect(nativeResponseControlEligible({ ...provider, upstreamWebsocket: false }, channel)).toBe(false); + expect(nativeResponseControlEligible({ ...provider, authMode: "forward" }, channel)).toBe(false); + expect(nativeResponseControlEligible(provider)).toBe(false); +}); + +test("injection mode refuses simultaneous steering instead of fabricating protocol equivalence", () => { + const { channel, sent, detach } = unitChannel(); + try { + expect(() => channel.steer({ type: "response.steer", previous_response_id: "root", input: "change plan" })).toThrow("injection-only"); + expect(sent).toHaveLength(0); + } finally { detach(); } +}); + +test("replay stores only accepted outputs and does not duplicate results echoed by the backend", () => { + const remembered: Array<{ input: unknown[]; response: Record }> = []; + const replay = new NativeInjectionReplay("original", (input, response) => remembered.push({ input, response })); + const call = { type: "function_call", call_id: "c1" }; + const result = savedResult("c1", "accepted"); + replay.observe({ type: "response.created", response: { id: "r1" } }); + replay.submitted({ type: "response.inject", input: [result] }); + replay.observe({ type: "response.inject.created" }); + replay.observe({ type: "response.completed", response: { id: "r1", output: [call, result] } }); + expect(remembered[0].response.output).toEqual([call, result]); + replay.dispose(); + const failed = new NativeInjectionReplay([], (input, response) => remembered.push({ input, response })); + failed.observe({ type: "response.created", response: { id: "r2" } }); + failed.submitted({ type: "response.inject", input: [savedResult("c1", "REJECTED_CONTENT")] }); + failed.observe({ type: "response.inject.failed" }); + failed.observe({ type: "response.completed", response: { id: "r2", output: [call] } }); + expect(JSON.stringify(remembered)).not.toContain("REJECTED_CONTENT"); + failed.dispose(); +}); + +test("a foreign acknowledgement closes the real exchange without exposing input or retrying", async () => { + const { socket, send, sent, ws, id } = await beginInjection(); + advertiseInjection(socket); + send({ type: "response.inject", response_id: id, input: [savedResult("call-1", "PRIVATE_FIXTURE_RESULT")] }); + acknowledgeInjection(socket, 100, "another-response"); + await waitForInjection(() => !ws.data.nativeSteering); + expect(socket.readyState).toBe(3); + expect(InjectionSocket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); + expect(sent.some(event => event.type === "response.inject.created")).toBe(false); + expect(JSON.stringify(sent)).not.toContain("PRIVATE_FIXTURE_RESULT"); +}); + +test("downstream disconnect discards queued injection without a second physical send", async () => { + const { socket, send, ws, handler, id } = await beginInjection(); + advertiseInjection(socket); advertiseInjection(socket, "call-2", 1); + for (const call of ["call-1", "call-2"]) send({ type: "response.inject", response_id: id, input: [savedResult(call)] }); + handler.close(ws, 1000, "fixture disconnect"); + await waitForInjection(() => !ws.data.nativeSteering); + acknowledgeInjection(socket); + await Bun.sleep(0); + expect(socket.frames.filter(frame => frame.type === "response.inject")).toHaveLength(1); + expect(socket.readyState).toBe(3); expect(fallbackCalls).toBe(0); +}); + +test("HTTP fallback never acquires injection ownership or replays a control frame", async () => { + const config = injectionConfig(true); + config.providers.api.upstreamWebsocket = false; + const { send, sent, ws } = injectionClient({}, config); + await waitForInjection(() => sent.some(event => event.type === "error")); + const requests = fallbackCalls; + send({ type: "response.inject", response_id: "unknown", input: [savedResult()] }); + expect(sent.at(-1)?.error.code).toBe("injection_not_supported"); + expect(fallbackCalls).toBe(requests); expect(InjectionSocket.all).toHaveLength(0); + expect(ws.data.nativeSteering).toBeUndefined(); +}); From 7c426e840b580a4ed8887e87bd7ca7372bd89dd1 Mon Sep 17 00:00:00 2001 From: luvs01 Date: Thu, 17 Sep 2026 13:59:06 +0900 Subject: [PATCH 2/4] refactor(responses): name the shared native control field nativeControl The field now holds either a steering or an injection owner, so the generic name matches the NativeResponseControl contract. Move the response-ownership markers next to the shared interface. Behavior is unchanged. --- src/server/index/websocket-handler.ts | 28 +++++++------- src/server/responses/codex-ws-exchange.ts | 16 ++++---- src/server/responses/core-options.ts | 2 +- src/server/responses/fetch-helpers.ts | 4 +- .../responses/native-response-control.ts | 6 +++ src/server/responses/native-steering.ts | 5 --- src/server/responses/passthrough-delivery.ts | 6 +-- src/server/responses/passthrough-dispatch.ts | 26 ++++++------- src/server/responses/ws-upstream.ts | 10 ++--- src/server/ws-bridge.ts | 2 +- tests/responses/ws-native-injection.test.ts | 34 ++++++++--------- tests/responses/ws-native-steering.test.ts | 38 +++++++++---------- 12 files changed, 89 insertions(+), 88 deletions(-) diff --git a/src/server/index/websocket-handler.ts b/src/server/index/websocket-handler.ts index 692a48c0a75..d117f960fad 100644 --- a/src/server/index/websocket-handler.ts +++ b/src/server/index/websocket-handler.ts @@ -193,19 +193,19 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { } catch { return; // text-only contract; ignore unparseable frames } - if (frame.type === "response.inject" || frame.type === "response.steer" || (frame.type === "response.create" && ws.data.nativeSteering)) { + if (frame.type === "response.inject" || frame.type === "response.steer" || (frame.type === "response.create" && ws.data.nativeControl)) { try { if (frame.type === "response.inject") { - if (!ws.data.nativeSteering?.inject) throw new NativeSteeringError("injection_not_supported", "Native injection is disabled or unavailable on this route."); - ws.data.nativeSteering.inject(frame); + if (!ws.data.nativeControl?.inject) throw new NativeSteeringError("injection_not_supported", "Native injection is disabled or unavailable on this route."); + ws.data.nativeControl.inject(frame); return; } if (frame.type === "response.steer") { - if (!ws.data.nativeSteering) throw new NativeSteeringError("steering_not_supported", "Native steering is disabled or unavailable on this route."); - ws.data.nativeSteering.steer(frame); + if (!ws.data.nativeControl) throw new NativeSteeringError("steering_not_supported", "Native steering is disabled or unavailable on this route."); + ws.data.nativeControl.steer(frame); return; } - if (ws.data.nativeSteering?.continue(frame)) return; + if (ws.data.nativeControl?.continue(frame)) return; } catch (error) { sendJsonFrame(ws, buildWsErrorFrame(400, { type: "invalid_request_error", @@ -221,12 +221,12 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { ws.data.cancel?.(); // A superseded turn must not keep ownership during warmup or refusal. - ws.data.nativeSteering = undefined; - let nativeSteering: NativeResponseControl | undefined; + ws.data.nativeControl = undefined; + let nativeControl: NativeResponseControl | undefined; try { const idleMs = typeof config.stallTimeoutSec === "number" && Number.isFinite(config.stallTimeoutSec) ? Math.max(1, config.stallTimeoutSec) * 1000 : 300_000; - nativeSteering = config.codexNativeInjection === true && isInjectionRequest(frame) + nativeControl = config.codexNativeInjection === true && isInjectionRequest(frame) ? new NativeInjectionChannel(frame, idleMs) : config.codexNativeSteering === true ? new NativeSteeringChannel(frame, idleMs) : undefined; } catch { @@ -268,7 +268,7 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { } // Only a genuinely admitted turn may receive steering or continuations. - ws.data.nativeSteering = nativeSteering; + ws.data.nativeControl = nativeControl; const payload: Record = { ...frame }; delete payload.type; turnAdmissionLease.bindAbortController(turnAbort); @@ -309,7 +309,7 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { ...(wsAdmission ? { admission: wsAdmission } : {}), forceEmptyResponseId: true, inboundTransport: "websocket", - nativeSteering, + nativeControl, abortSignal: turnAbort.signal, turnAdmissionLease, onFirstOutput: () => recordFirstOutput(logCtx, start), @@ -320,8 +320,8 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { }, }); await sendResponseToWebSocket(ws, response, isCurrent, { - untilEof: nativeSteering?.relayActive === true, - onSsePayload: nativeSteering?.relayActive + untilEof: nativeControl?.relayActive === true, + onSsePayload: nativeControl?.relayActive ? createNativeSteeringLogObserver(logCtx, () => recordFirstOutput(logCtx, start)) : payload => inspectResponseLogSsePayload(logCtx, payload), onTerminal: status => { @@ -359,7 +359,7 @@ export function createWebsocketHandler(ctx: ServeOptionsContext) { } } finally { turnAdmissionLease.release(); - if (ws.data.nativeSteering === nativeSteering) ws.data.nativeSteering = undefined; + if (ws.data.nativeControl === nativeControl) ws.data.nativeControl = undefined; if (!logged && turnAbort.signal.aborted) finalizeLog(499); if (ws.data.cancel === cancelTurn) ws.data.cancel = undefined; } diff --git a/src/server/responses/codex-ws-exchange.ts b/src/server/responses/codex-ws-exchange.ts index 0c686de150a..6074c445082 100644 --- a/src/server/responses/codex-ws-exchange.ts +++ b/src/server/responses/codex-ws-exchange.ts @@ -1,4 +1,4 @@ -import { markNativeSteeringResponse } from "./native-steering"; +import { markNativeControlResponse } from "./native-response-control"; import type { NativeResponseControl } from "./native-response-control"; import { MAX_CLIENT_SSE_FRAME_BYTES } from "../sse-frame-buffer"; import { isSafeResponseHeader } from "../safe-response-headers"; @@ -12,7 +12,7 @@ import { UPGRADE_DEADLINE_MS, CODEX_WS_LIVENESS_PING_INTERVAL_MS, CODEX_WS_RESPO type CodexWsFailureStage, type CodexWsStageRecord } from "./codex-ws-wire"; interface ExchangeOptions { - nativeSteering?: NativeResponseControl; + nativeControl?: NativeResponseControl; beforeContinuation?: () => Promise; session: CodexWsSession; url: string; @@ -90,7 +90,7 @@ function wrappedRejectionResponse(payload: Record, prelude: Hea /** The sole SSE exchange state machine for both one-shot and retained sockets. */ export function codexWsExchange(options: ExchangeOptions): Promise { - const { session, url, init, prepared, sseFallback, onQuota, beforeDispatch, bunVersion, nativeSteering, beforeContinuation } = options; + const { session, url, init, prepared, sseFallback, onQuota, beforeDispatch, bunVersion, nativeControl, beforeContinuation } = options; const { frameText, headers } = prepared; const signal = init.signal ?? undefined; return new Promise((resolve, reject) => { @@ -204,7 +204,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise { const response = new Response(stream, { status: 200, headers: responseHeaders }); metadata?.commit(); markCodexWsResponse(response, Boolean(metadata && onQuota)); - if (nativeSteering) markNativeSteeringResponse(response); + if (nativeControl) markNativeControlResponse(response); markCodexWsStage(response, stageRecord(null)); committedResponse = response; resolve(response); @@ -337,8 +337,8 @@ export function codexWsExchange(options: ExchangeOptions): Promise { if (terminal || settledPreOpen || signal?.aborted) return; sent = true; try { - if (nativeSteering) { - detachSteering = nativeSteering.attach(frame => { + if (nativeControl) { + detachSteering = nativeControl.attach(frame => { const sendControl = () => { if (terminal || signal?.aborted || session.closed || ws.readyState !== WebSocket.OPEN) { throw new Error("Native steering connection is no longer available"); @@ -442,7 +442,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise { let steeringEnded = false; if (!controlFrame) { try { - if (nativeSteering) steeringEnded = nativeSteering.observe(normalized.payload); + if (nativeControl) steeringEnded = nativeControl.observe(normalized.payload); else correlation?.accept(normalized.payload); } catch (error) { failStream(error); return; } // Correlation must run first: a reused socket's foreign-stream error settles as a @@ -486,7 +486,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise { return; } if (!controlFrame) relayedEvents += 1; - if (nativeSteering ? steeringEnded : (type === "response.completed" || type === "response.failed" || type === "response.incomplete" || type === "error")) { + if (nativeControl ? steeringEnded : (type === "response.completed" || type === "response.failed" || type === "response.incomplete" || type === "error")) { const completedId = correlation?.completed(normalized.payload) ?? null; terminal = true; cleanup(); diff --git a/src/server/responses/core-options.ts b/src/server/responses/core-options.ts index fdebef3c4b7..4c790fd4984 100644 --- a/src/server/responses/core-options.ts +++ b/src/server/responses/core-options.ts @@ -53,7 +53,7 @@ export interface HandleResponsesOptions { onRequestBodyRead?: () => void; forceEmptyResponseId?: boolean; /** Internal, connection-owned control channel; never reconstructed from headers. */ - nativeSteering?: NativeResponseControl; + nativeControl?: NativeResponseControl; abortSignal?: AbortSignal; /** One-shot TTFT callback: first non-empty model output observed (WP4). */ onFirstOutput?: () => void; diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index 4a7ebb85d16..b776bfbe6fb 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -54,7 +54,7 @@ export interface PaceAwareFetch { export type ProviderFetch = typeof globalThis.fetch & PaceAwareFetch; export interface ProviderFetchOptions { - nativeSteering?: NativeResponseControl; + nativeControl?: NativeResponseControl; providerName?: string; modelId?: string; /** One pacing slot was acquired immediately before this fetch wrapper was created. */ @@ -103,7 +103,7 @@ export function providerFetch( // used, protocol pin included: a WS turn that falls back is serving the // request over HTTP, and dropping the provider's `upstreamHttpVersion` // there would silently negotiate a transport the operator ruled out. - return codexWsUpstreamFetch(input, init, httpFetch, runtime, options.onCodexWsQuota, options.beforeDispatch, options.nativeSteering, + return codexWsUpstreamFetch(input, init, httpFetch, runtime, options.onCodexWsQuota, options.beforeDispatch, options.nativeControl, () => waitForPacing(init.signal ?? undefined)); } return httpFetch(input, init); diff --git a/src/server/responses/native-response-control.ts b/src/server/responses/native-response-control.ts index 516933ebfbb..0b779da033b 100644 --- a/src/server/responses/native-response-control.ts +++ b/src/server/responses/native-response-control.ts @@ -18,6 +18,12 @@ export interface NativeResponseControl { export const OPENAI_API_RESPONSES_URL = "https://api.openai.com/v1/responses"; +const nativeControlResponses = new WeakSet(); +/** Mark the exact response for multi-response delivery without serializing a wire field. */ +export function markNativeControlResponse(response: Response): Response { nativeControlResponses.add(response); return response; } +/** Recognize a marked native response by identity, not by caller-controlled content. */ +export function isNativeControlResponse(response: Response): boolean { return nativeControlResponses.has(response); } + /** Preserve canonical ChatGPT eligibility; public API injection is separately opted in. */ export function nativeResponseControlEligible(provider: OcxProviderConfig, control?: NativeResponseControl): boolean { if (isCanonicalOpenAiForwardProvider(provider)) return true; diff --git a/src/server/responses/native-steering.ts b/src/server/responses/native-steering.ts index 741606ad433..05ccf3e482f 100644 --- a/src/server/responses/native-steering.ts +++ b/src/server/responses/native-steering.ts @@ -10,11 +10,6 @@ export const NATIVE_STEERING_TOOL_WAIT_MS = 30 * 60_000; type Frame = Record; type Send = (frame: Frame) => void; -const responses = new WeakSet(); -/** Mark the exact response for multi-response delivery without serializing a wire field. */ -export function markNativeSteeringResponse(response: Response): Response { responses.add(response); return response; } -/** Recognize a marked native response by identity, not by caller-controlled content. */ -export function isNativeSteeringResponse(response: Response): boolean { return responses.has(response); } /** Narrow JSON object envelopes while excluding arrays and null. */ function record(value: unknown): value is Frame { diff --git a/src/server/responses/passthrough-delivery.ts b/src/server/responses/passthrough-delivery.ts index b45701dbcfc..dbc3b15133e 100644 --- a/src/server/responses/passthrough-delivery.ts +++ b/src/server/responses/passthrough-delivery.ts @@ -1,4 +1,4 @@ -import { isNativeSteeringResponse } from "./native-steering"; +import { isNativeControlResponse } from "./native-response-control"; import type { ResponsesRequestContext, ResponsesAdmissionState } from "./core-options"; import type { PreparedResponsesRequest } from "./request-prepare"; import type { ResponsesTransport } from "./request-transport"; @@ -313,11 +313,11 @@ export async function deliverPassthroughResponse( }); } - if (options.nativeSteering && isNativeSteeringResponse(upstreamResponse) && upstreamResponse.body) { + if (options.nativeControl && isNativeControlResponse(upstreamResponse) && upstreamResponse.body) { // A native chain carries several response terminals. Ordinary SSE repair, // cancellation-on-terminal and local previous-response replay are single-response // contracts and would truncate it. Keep the bounded upstream as the sole reader. - options.nativeSteering.relayActive = true; + options.nativeControl.relayActive = true; commitReasoningReplayServingRoute(nativeExchange.request.headers); const body = trackStreamLifetime(upstreamResponse.body, upstream, undefined, options.turnAdmissionLease); return new Response(body, { status: upstreamResponse.status, headers }); diff --git a/src/server/responses/passthrough-dispatch.ts b/src/server/responses/passthrough-dispatch.ts index cd1dc0a9831..f38f3a02af3 100644 --- a/src/server/responses/passthrough-dispatch.ts +++ b/src/server/responses/passthrough-dispatch.ts @@ -254,11 +254,11 @@ export async function preparePassthroughExchange( ? (response: { id?: unknown; output?: unknown; status?: unknown }) => rememberResponseState(parsed._rawBody, response, undefined, responseStateOptions(true)) : undefined; - if (options.nativeSteering && nativeResponseControlEligible(route.provider, options.nativeSteering) + if (options.nativeControl && nativeResponseControlEligible(route.provider, options.nativeControl) && options.inboundTransport === "websocket" && !options.comboAttempt) { const body = parsed._rawBody as Record; - const Replay = options.nativeSteering.kind === "injection" ? NativeInjectionReplay : NativeSteeringReplay; - options.nativeSteering.replayFactory = () => new Replay(body.input, (input, response) => { + const Replay = options.nativeControl.kind === "injection" ? NativeInjectionReplay : NativeSteeringReplay; + options.nativeControl.replayFactory = () => new Replay(body.input, (input, response) => { if (passthroughRecordEligible && !isBodyNonPersistable(body)) { rememberResponseState({ ...body, input }, response, undefined, responseStateOptions(true)); } @@ -785,9 +785,9 @@ export async function preparePassthroughExchange( body: request.body, }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeControl: nativeResponseControlEligible(route.provider, options.nativeControl) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 - ? options.nativeSteering : undefined, + ? options.nativeControl : undefined, dispatchOverride: oauthDispatch(request), providerName: route.providerName, modelId: route.modelId, @@ -884,9 +884,9 @@ export async function preparePassthroughExchange( body: request.body, }, innerRecovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeControl: nativeResponseControlEligible(route.provider, options.nativeControl) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 - ? options.nativeSteering : undefined, + ? options.nativeControl : undefined, dispatchOverride: oauthDispatch(request), providerName: route.providerName, modelId: route.modelId, @@ -993,9 +993,9 @@ export async function preparePassthroughExchange( // here on is a genuine transport attempt. storedPoolReplayDispatchNotifier( providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeControl: nativeResponseControlEligible(route.provider, options.nativeControl) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 - ? options.nativeSteering : undefined, + ? options.nativeControl : undefined, dispatchOverride: oauthDispatch(request), providerName: route.providerName, modelId: route.modelId, @@ -1117,9 +1117,9 @@ export async function preparePassthroughExchange( body: request.body, }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeControl: nativeResponseControlEligible(route.provider, options.nativeControl) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 - ? options.nativeSteering : undefined, + ? options.nativeControl : undefined, dispatchOverride: oauthDispatch(request), providerName: route.providerName, modelId: route.modelId, @@ -1242,9 +1242,9 @@ export async function preparePassthroughExchange( body: request.body, }, recovery), upstream.signal, connectMs, parsed.stream, providerFetch(route.provider, options.codexWsRuntimeIdentity, { - nativeSteering: nativeResponseControlEligible(route.provider, options.nativeSteering) && options.inboundTransport === "websocket" && !options.comboAttempt + nativeControl: nativeResponseControlEligible(route.provider, options.nativeControl) && options.inboundTransport === "websocket" && !options.comboAttempt && responseEffects.plaintextV2AgentMessageToolNames.size === 0 - ? options.nativeSteering : undefined, + ? options.nativeControl : undefined, dispatchOverride: oauthDispatch(request), providerName: route.providerName, modelId: route.modelId, diff --git a/src/server/responses/ws-upstream.ts b/src/server/responses/ws-upstream.ts index a0fb1c9bd0c..69b5539b16e 100644 --- a/src/server/responses/ws-upstream.ts +++ b/src/server/responses/ws-upstream.ts @@ -132,7 +132,7 @@ export function codexWsUpstreamFetch( runtime: BunRuntimeGateInput = currentBunRuntimeIdentity(), onQuota?: CodexWsQuotaObserver, beforeDispatch?: (headers: Headers) => void, - nativeSteering?: NativeResponseControl, + nativeControl?: NativeResponseControl, beforeContinuation?: () => Promise, ): Promise { const prepared = prepareCodexWsRequest(url, init); @@ -148,9 +148,9 @@ export function codexWsUpstreamFetch( const { frameText, headers } = prepared; // Never infer backend support from a model name or enable controls on a gateway. - const control = nativeSteering?.kind === "injection" - ? ((prepared.canonical || url === OPENAI_API_RESPONSES_URL) && isInjectionRequest(JSON.parse(frameText)) ? nativeSteering : undefined) - : prepared.canonical ? nativeSteering : undefined; + const control = nativeControl?.kind === "injection" + ? ((prepared.canonical || url === OPENAI_API_RESPONSES_URL) && isInjectionRequest(JSON.parse(frameText)) ? nativeControl : undefined) + : prepared.canonical ? nativeControl : undefined; if (control?.kind === "injection" && url === OPENAI_API_RESPONSES_URL) { const beta = headers["openai-beta"]; if (!beta?.split(",").some(value => value.trim() === "responses_multi_agent=v1")) { @@ -199,7 +199,7 @@ export function codexWsUpstreamFetch( } return codexWsExchange({ session, url, init, prepared, sseFallback, onQuota, beforeDispatch, - nativeSteering: control, + nativeControl: control, beforeContinuation, bunVersion: typeof runtime === "string" ? runtime : runtime.version, }); diff --git a/src/server/ws-bridge.ts b/src/server/ws-bridge.ts index fa7e0a4b652..eb562b43ae0 100644 --- a/src/server/ws-bridge.ts +++ b/src/server/ws-bridge.ts @@ -18,7 +18,7 @@ type ResponsesTerminalReporter = (status: ResponsesTerminalStatus) => void; type ResponsesPayloadObserver = (payload: string) => void; export interface WsData { - nativeSteering?: NativeResponseControl; + nativeControl?: NativeResponseControl; headers?: Headers; // base inbound forward headers only; per-turn auth refresh injects current pool tokens /** * Resolved once at the handshake. Auth is handshake-time only on this path, so diff --git a/tests/responses/ws-native-injection.test.ts b/tests/responses/ws-native-injection.test.ts index fe444ec993a..7adf7c3e185 100644 --- a/tests/responses/ws-native-injection.test.ts +++ b/tests/responses/ws-native-injection.test.ts @@ -30,7 +30,7 @@ test.each([false, true])("real handler sends saved results over the same connect expect(socket.frames[1]).toEqual(frame); acknowledgeInjection(socket); completeInjection(socket, { output: [call], usage: { input_tokens: 10, output_tokens: 5 } }); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(sent.some(event => event.type === "response.inject.created")).toBe(true); expect(sent.at(-1)?.type).toBe("response.completed"); expect(InjectionSocket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); @@ -50,9 +50,9 @@ test("terminal before acknowledgement is relayed without dropping the late succe send({ type: "response.inject", response_id: id, input: [savedResult()] }); completeInjection(socket, { output: [call] }); await waitForInjection(() => sent.some(event => event.type === "response.completed")); - expect(socket.readyState).toBe(1); expect(ws.data.nativeSteering).toBeDefined(); + expect(socket.readyState).toBe(1); expect(ws.data.nativeControl).toBeDefined(); acknowledgeInjection(socket); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("response.inject.created"); expect(socket.readyState).toBe(3); expect(fallbackCalls).toBe(0); }); @@ -66,7 +66,7 @@ test("asynchronous tool completion after the response terminal still reaches the send({ type: "response.inject", response_id: id, input: [savedResult()] }); expect(socket.frames[1]?.type).toBe("response.inject"); acknowledgeInjection(socket); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("response.inject.created"); }); @@ -91,7 +91,7 @@ test("completion rejection is preserved; only an explicit caller continuation re expect(sent.at(-1)?.error.code).toBe("injection_pending"); socket.emit({ type: "response.created", response: { id: "successor", previous_response_id: id } }); completeInjection(socket, {}, "successor"); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(InjectionSocket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); }); @@ -105,9 +105,9 @@ test("parallel tool results serialize by acknowledgement without losing caller o acknowledgeInjection(socket); await waitForInjection(() => socket.frames.length === 3); expect(socket.frames[2].input).toEqual([savedResult("call-2")]); - expect(ws.data.nativeSteering).toBeDefined(); + expect(ws.data.nativeControl).toBeDefined(); acknowledgeInjection(socket, 101); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(sent.filter(event => event.type === "response.inject.created")).toHaveLength(2); }); @@ -120,7 +120,7 @@ test("duplicate results are refused both while pending and after successful acce acknowledgeInjection(socket); send(frame); expect(sent.at(-1)?.error.code).toBe("duplicate_injection"); expect(socket.frames).toHaveLength(2); completeInjection(socket, { output: [call] }); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); }); test("different response, lane, unadvertised and hosted-tool results never reach upstream", async () => { @@ -134,7 +134,7 @@ test("different response, lane, unadvertised and hosted-tool results never reach { ...base, input: [{ type: "message", role: "system", content: "no" }] }, ]) { send(frame); expect(sent.at(-1)?.type).toBe("error"); } expect(socket.frames).toHaveLength(1); - ws.close(); await waitForInjection(() => !ws.data.nativeSteering); + ws.close(); await waitForInjection(() => !ws.data.nativeControl); }); test("two connections cannot inject results into each other's response or credentials", async () => { @@ -145,7 +145,7 @@ test("two connections cannot inject results into each other's response or creden expect(a.sent.at(-1)?.error.code).toBe("injection_response_mismatch"); expect(a.socket.frames).toHaveLength(1); expect(b.socket.frames).toHaveLength(1); a.ws.close(); b.ws.close(); - await waitForInjection(() => !a.ws.data.nativeSteering && !b.ws.data.nativeSteering); + await waitForInjection(() => !a.ws.data.nativeControl && !b.ws.data.nativeControl); }); test.each(["disabled", "no-multi-agent", "warmup", "steering-only"])("unsupported %s does not silently discard an injection", async mode => { @@ -167,14 +167,14 @@ test("a pending injection prevents a new create from cancelling the owned socket send({ type: "response.create", model: "different-model", input: "new work" }); expect(sent.at(-1)?.error.code).toBe("injection_pending"); expect(socket.readyState).toBe(1); expect(socket.frames).toHaveLength(2); - ws.close(); await waitForInjection(() => !ws.data.nativeSteering); + ws.close(); await waitForInjection(() => !ws.data.nativeControl); }); test("unknown delivery closes without HTTP fallback or resending the control", async () => { const { socket, send, sent, ws, id } = await beginInjection(); advertiseInjection(socket); socket.throwOnInject = true; send({ type: "response.inject", response_id: id, input: [savedResult()] }); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(socket.readyState).toBe(3); expect(fallbackCalls).toBe(0); expect(InjectionSocket.all).toHaveLength(1); expect(sent.some(event => event.error?.message?.includes("unknown"))).toBe(true); @@ -193,7 +193,7 @@ test("accepted function results survive ordinary subsequent delta turns; no user const call = advertiseInjection(socket); send({ type: "response.inject", response_id: id, input: [savedResult("call-1", "accepted-result")] }); acknowledgeInjection(socket); completeInjection(socket, { output: [call] }); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); send({ type: "response.create", model: "gpt-5.6-sol", multi_agent: { enabled: true }, previous_response_id: id, input: "followup" }); await waitForInjection(() => InjectionSocket.all.length === 2 && InjectionSocket.all[1].frames.length > 0); const next = InjectionSocket.all[1]; @@ -201,7 +201,7 @@ test("accepted function results survive ordinary subsequent delta turns; no user expect(history.filter(item => item.type === "function_call_output")).toEqual([savedResult("call-1", "accepted-result")]); expect(history.findIndex(item => item.type === "function_call_output")).toBe(history.findIndex(item => item.type === "function_call") + 1); expect(JSON.stringify(history)).toContain("followup"); - completeInjection(next); await waitForInjection(() => !ws.data.nativeSteering); + completeInjection(next); await waitForInjection(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("response.completed"); }); @@ -366,7 +366,7 @@ test("a foreign acknowledgement closes the real exchange without exposing input advertiseInjection(socket); send({ type: "response.inject", response_id: id, input: [savedResult("call-1", "PRIVATE_FIXTURE_RESULT")] }); acknowledgeInjection(socket, 100, "another-response"); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); expect(socket.readyState).toBe(3); expect(InjectionSocket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); expect(sent.some(event => event.type === "response.inject.created")).toBe(false); @@ -378,7 +378,7 @@ test("downstream disconnect discards queued injection without a second physical advertiseInjection(socket); advertiseInjection(socket, "call-2", 1); for (const call of ["call-1", "call-2"]) send({ type: "response.inject", response_id: id, input: [savedResult(call)] }); handler.close(ws, 1000, "fixture disconnect"); - await waitForInjection(() => !ws.data.nativeSteering); + await waitForInjection(() => !ws.data.nativeControl); acknowledgeInjection(socket); await Bun.sleep(0); expect(socket.frames.filter(frame => frame.type === "response.inject")).toHaveLength(1); @@ -394,5 +394,5 @@ test("HTTP fallback never acquires injection ownership or replays a control fram send({ type: "response.inject", response_id: "unknown", input: [savedResult()] }); expect(sent.at(-1)?.error.code).toBe("injection_not_supported"); expect(fallbackCalls).toBe(requests); expect(InjectionSocket.all).toHaveLength(0); - expect(ws.data.nativeSteering).toBeUndefined(); + expect(ws.data.nativeControl).toBeUndefined(); }); diff --git a/tests/responses/ws-native-steering.test.ts b/tests/responses/ws-native-steering.test.ts index 67708cea6c8..e2e99831093 100644 --- a/tests/responses/ws-native-steering.test.ts +++ b/tests/responses/ws-native-steering.test.ts @@ -103,7 +103,7 @@ test("real handler -> auth/dispatch -> native exchange -> downstream preserves a socket.emit({ type: "response.incomplete", response: { id, status: "incomplete", output: [], incomplete_details: { reason: "steered" }, usage: { input_tokens: 10, output_tokens: 2 } } }); socket.emit({ type: "response.created", response: { id: "successor", previous_response_id: id, output: [] } }); complete(socket, "successor", { usage: { input_tokens: 20, output_tokens: 3 } }); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.map(frame => frame.type)).toEqual(["response.created", "response.steer.accepted", "response.incomplete", "response.created", "response.completed"]); expect(sent.at(-1)?.response.id).toBe("successor"); expect(Socket.all).toHaveLength(1); @@ -123,7 +123,7 @@ test("normal completion before acceptance still retains the socket and successor accept(socket, id); socket.emit({ type: "response.created", response: { id: "r2", previous_response_id: id } }); complete(socket, "r2"); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.filter(frame => frame.type === "response.completed").map(frame => frame.response.id)).toEqual([id, "r2"]); }); @@ -146,7 +146,7 @@ test("pending results use one same-account/lane create and never replay accepted expect(sent.at(-1)?.error.code).toBe("duplicate_continuation"); socket.emit({ type: "response.created", response: { id: "r2", previous_response_id: id } }); complete(socket, "r2"); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.at(-1)?.response.id).toBe("r2"); expect(Socket.all).toHaveLength(1); }); @@ -157,7 +157,7 @@ test("subsequent ordinary turns retain committed steering through the scoped rep accept(socket, id); complete(socket, id); socket.emit({ type: "response.created", response: { id: "cached-successor", previous_response_id: id } }); complete(socket, "cached-successor"); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); send({ type: "response.create", model: "gpt-5.5", previous_response_id: "cached-successor", input: "ordinary next turn" }); await waitFor(() => Socket.all.length === 2 && Socket.all[1].frames.length > 0); const next = Socket.all[1]; @@ -165,7 +165,7 @@ test("subsequent ordinary turns retain committed steering through the scoped rep expect(JSON.stringify(next.frames[0].input)).toContain("initial"); expect(JSON.stringify(next.frames[0].input)).toContain("ordinary next turn"); complete(next, next.root); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("response.completed"); }); @@ -174,7 +174,7 @@ test("rejected steering after terminal settles without an invented successor", a send({ type: "response.steer", previous_response_id: id, input: "not supported" }); complete(socket, id); socket.emit({ type: "response.steer.failed", steer: { previous_response_id: id, input: "not supported" }, error: { code: "steering_not_supported", message: "model does not support steering" } }); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("response.steer.failed"); expect(sent.filter(frame => frame.type === "response.created")).toHaveLength(1); expect(socket.frames).toHaveLength(2); @@ -263,7 +263,7 @@ test("HTTP upgrade fallback keeps ordinary streaming and rejects steering explic send({ type: "response.steer", previous_response_id: "http-response", input: "not delivered" }); expect(sent.at(-1)?.error.code).toBe("steering_not_supported"); finish(); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("response.completed"); expect(fallbackCalls).toBe(1); expect(Socket.all).toHaveLength(0); @@ -274,7 +274,7 @@ test("post-send disconnect never replays accepted steering through HTTP or anoth send({ type: "response.steer", previous_response_id: id, input: "delivery unknown" }); accept(socket, id); socket.close(); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.at(-1)?.type).toBe("error"); expect(fallbackCalls).toBe(0); expect(Socket.all).toHaveLength(1); @@ -287,7 +287,7 @@ test("downstream disconnect closes the dedicated upstream while steering is pend accept(socket, id); complete(socket, id); socket.emit({ type: "response.steer.pending", steer: { id: "s1", previous_response_id: id }, reason: "waiting_for_required_input", required_input: [{ type: "function_call_output", call_id: "saved-call" }] }); handler.close(ws); - await waitFor(() => socket.readyState === 3 && !ws.data.nativeSteering); + await waitFor(() => socket.readyState === 3 && !ws.data.nativeControl); expect(fallbackCalls).toBe(0); expect(socket.frames).toHaveLength(2); }); @@ -319,7 +319,7 @@ test("saved tool results may arrive before pending and retain extra user input w expect(sent.some(frame => frame.type === "error")).toBe(false); socket.emit({ type: "response.created", response: { id: "early-successor", previous_response_id: id } }); complete(socket, "early-successor"); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(Socket.all).toHaveLength(1); expect(fallbackCalls).toBe(0); }); @@ -354,7 +354,7 @@ test("a steering failure cannot close an already submitted explicit continuation expect(socket.readyState).toBe(1); socket.emit({ type: "response.created", response: { id: "explicit-successor", previous_response_id: id } }); complete(socket, "explicit-successor"); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); expect(sent.at(-1)?.response.id).toBe("explicit-successor"); expect(fallbackCalls).toBe(0); }); @@ -386,18 +386,18 @@ test("early continuation validates advertised call and approval identities and r test("warmup leaves no steering owner and the next ordinary turn gets a fresh channel", async () => { const { ws, sent, send } = downstream({ generate: false }); expect(sent.map(frame => frame.type)).toEqual(["response.created", "response.completed"]); - expect(ws.data.nativeSteering).toBeUndefined(); + expect(ws.data.nativeControl).toBeUndefined(); expect(ws.data.cancel).toBeUndefined(); expect(Socket.all).toHaveLength(0); send({ type: "response.steer", previous_response_id: sent[0].response.id, input: "not a running turn" }); expect(sent.at(-1)?.error.code).toBe("steering_not_supported"); send({ type: "response.create", model: "gpt-5.5", input: "real turn" }); await waitFor(() => Socket.all.length === 1 && sent.filter(frame => frame.type === "response.created").length === 2); - expect(ws.data.nativeSteering?.attached).toBe(true); + expect(ws.data.nativeControl?.attached).toBe(true); const socket = Socket.all[0]; expect(socket.frames[0].input).toBe("real turn"); complete(socket, socket.root); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); }); test("admission refusal leaves no steering owner and a later admitted turn is independent", async () => { @@ -410,7 +410,7 @@ test("admission refusal leaves no steering owner and a later admitted turn is in expect(leases.length).toBeGreaterThan(0); const { ws, sent, send } = downstream(); expect(sent.at(-1)?.error.code).toBe("server_busy"); - expect(ws.data.nativeSteering).toBeUndefined(); + expect(ws.data.nativeControl).toBeUndefined(); expect(ws.data.cancel).toBeUndefined(); expect(Socket.all).toHaveLength(0); for (const lease of leases) lease.release(); @@ -419,9 +419,9 @@ test("admission refusal leaves no steering owner and a later admitted turn is in expect(Socket.all).toHaveLength(1); const socket = Socket.all[0]; expect(socket.frames[0].input).toBe("after admission"); - expect(ws.data.nativeSteering?.attached).toBe(true); + expect(ws.data.nativeControl?.attached).toBe(true); complete(socket, socket.root); - await waitFor(() => !ws.data.nativeSteering); + await waitFor(() => !ws.data.nativeControl); } finally { for (const lease of leases) lease.release(); } @@ -429,9 +429,9 @@ test("admission refusal leaves no steering owner and a later admitted turn is in test("superseding an active turn with warmup clears its steering owner immediately", async () => { const { ws, socket, send } = await begin(); - expect(ws.data.nativeSteering?.attached).toBe(true); + expect(ws.data.nativeControl?.attached).toBe(true); send({ type: "response.create", model: "gpt-5.5", input: "warmup", generate: false }); - expect(ws.data.nativeSteering).toBeUndefined(); + expect(ws.data.nativeControl).toBeUndefined(); expect(ws.data.cancel).toBeUndefined(); await waitFor(() => socket.readyState === 3); }); From 0e1b5f964bdcf24c42600b96aab3e9b2a278343d Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Thu, 17 Sep 2026 19:51:13 +0900 Subject: [PATCH 3/4] fix(responses): fail the turn on a native-control attach conflict instead of falling back to HTTP --- src/server/responses/codex-ws-exchange.ts | 9 +++++++ tests/responses/ws-upstream.test.ts | 31 +++++++++++++++++++++++ 2 files changed, 40 insertions(+) diff --git a/src/server/responses/codex-ws-exchange.ts b/src/server/responses/codex-ws-exchange.ts index 6074c445082..89f3096f0cf 100644 --- a/src/server/responses/codex-ws-exchange.ts +++ b/src/server/responses/codex-ws-exchange.ts @@ -369,6 +369,15 @@ export function codexWsExchange(options: ExchangeOptions): Promise { } else sendControl(); }, error => failStream(error)); } + } catch (error) { + // An attach failure is an ownership conflict, not a failed send: no frame + // left the process, but the channel can never bind, so resolving the HTTP + // fallback here would silently degrade a multi-agent turn into an ordinary + // one. Fail the turn visibly instead. + failStream(error); + return; + } + try { ws.send(frameText); sentAt = Date.now(); } catch { diff --git a/tests/responses/ws-upstream.test.ts b/tests/responses/ws-upstream.test.ts index c9afb0cd7d9..5257ceb0f50 100644 --- a/tests/responses/ws-upstream.test.ts +++ b/tests/responses/ws-upstream.test.ts @@ -1148,6 +1148,37 @@ describe("codexWsUpstreamFetch", () => { expect([...ws.listeners.values()].every(listeners => listeners.length === 0)).toBe(true); } finally { session.dispose(); } }); + + test("a native-control attach conflict fails the turn instead of falling back to HTTP", async () => { + installFake(ws => { ws.emit("open", {}); }); + const init = streamingInit(); + const prepared = prepareCodexWsRequest(CODEX_URL, init)!; + const session = new CodexWsSession("wss://chatgpt.com/backend-api/codex/responses", prepared.headers, true); + let fallbacks = 0; + const nativeControl = { + kind: "injection" as const, + relayActive: false, + attached: true, + ended: false, + attach() { throw new Error("Native injection transport is already owned."); }, + observe() { return false; }, + steer() { throw new Error("unreachable"); }, + continue() { return false; }, + }; + const options = { session, url: CODEX_URL, init, prepared, nativeControl, + sseFallback: (async () => { fallbacks++; throw new Error("attach conflict must not fall back"); }) as typeof fetch }; + try { + expect(session.reserve()).toBe(true); + const response = await codexWsExchange(options); + const ws = FakeWebSocket.instances.at(-1)!; + expect(fallbacks).toBe(0); + expect(ws.sent).toHaveLength(0); + expect(response.status).toBe(502); + expect(((await response.json()) as { error: { message: string } }).error.message).toContain("already owned"); + expect(ws.closed).toBe(true); + expect([...ws.listeners.values()].every(listeners => listeners.length === 0)).toBe(true); + } finally { session.dispose(); } + }); }); test.each(["error", "response.completed"])("multiline upstream %s JSON remains one valid SSE data value", async type => { From 92ec3630dc6b1d053010bd79d09ca268791d55d8 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Thu, 17 Sep 2026 22:01:04 +0900 Subject: [PATCH 4/4] test(responses): move the attach-conflict regression into ws-failure-stage --- tests/responses/ws-failure-stage.test.ts | 39 +++++++++++++++++++++++- tests/responses/ws-upstream.test.ts | 31 ------------------- 2 files changed, 38 insertions(+), 32 deletions(-) diff --git a/tests/responses/ws-failure-stage.test.ts b/tests/responses/ws-failure-stage.test.ts index d4af512669e..e84145b04dc 100644 --- a/tests/responses/ws-failure-stage.test.ts +++ b/tests/responses/ws-failure-stage.test.ts @@ -17,6 +17,9 @@ import { codexWsUpstreamFetch, CODEX_WS_RESPONSE_PRELUDE_TIMEOUT_MS, } from "../../src/server/responses/ws-upstream"; +import { codexWsExchange } from "../../src/server/responses/codex-ws-exchange"; +import { CodexWsSession } from "../../src/server/responses/codex-ws-session"; +import { prepareCodexWsRequest } from "../../src/server/responses/codex-ws-request"; /** * #4191: a long Codex thread died only through the proxy, and every variant of @@ -37,6 +40,7 @@ class FakeWebSocket { static script: (ws: FakeWebSocket) => void = () => {}; url: string; sent: string[] = []; + closed = false; listeners = new Map(); constructor(url: string) { @@ -63,7 +67,7 @@ class FakeWebSocket { this.sent.push(data); } - close() {} + close() { this.closed = true; } } const RealWebSocket = globalThis.WebSocket; @@ -471,3 +475,36 @@ describe("codex ws stage record marker (#4191)", () => { } }); }); + +describe("native-control attach conflict", () => { + test("an already-owned channel fails the turn instead of falling back to HTTP", async () => { + installFake(ws => { ws.emit("open", {}); }); + const init = streamingInit(); + const prepared = prepareCodexWsRequest(CODEX_URL, init)!; + const session = new CodexWsSession("wss://chatgpt.com/backend-api/codex/responses", prepared.headers, true); + let fallbacks = 0; + const nativeControl = { + kind: "injection" as const, + relayActive: false, + attached: true, + ended: false, + attach() { throw new Error("Native injection transport is already owned."); }, + observe() { return false; }, + steer() { throw new Error("unreachable"); }, + continue() { return false; }, + }; + const options = { session, url: CODEX_URL, init, prepared, nativeControl, + sseFallback: (async () => { fallbacks++; throw new Error("attach conflict must not fall back"); }) as typeof fetch }; + try { + expect(session.reserve()).toBe(true); + const response = await codexWsExchange(options); + const ws = FakeWebSocket.instances.at(-1)!; + expect(fallbacks).toBe(0); + expect(ws.sent).toHaveLength(0); + expect(response.status).toBe(502); + expect(((await response.json()) as { error: { message: string } }).error.message).toContain("already owned"); + expect(ws.closed).toBe(true); + expect([...ws.listeners.values()].every(listeners => listeners.length === 0)).toBe(true); + } finally { session.dispose(); } + }); +}); diff --git a/tests/responses/ws-upstream.test.ts b/tests/responses/ws-upstream.test.ts index 5257ceb0f50..c9afb0cd7d9 100644 --- a/tests/responses/ws-upstream.test.ts +++ b/tests/responses/ws-upstream.test.ts @@ -1148,37 +1148,6 @@ describe("codexWsUpstreamFetch", () => { expect([...ws.listeners.values()].every(listeners => listeners.length === 0)).toBe(true); } finally { session.dispose(); } }); - - test("a native-control attach conflict fails the turn instead of falling back to HTTP", async () => { - installFake(ws => { ws.emit("open", {}); }); - const init = streamingInit(); - const prepared = prepareCodexWsRequest(CODEX_URL, init)!; - const session = new CodexWsSession("wss://chatgpt.com/backend-api/codex/responses", prepared.headers, true); - let fallbacks = 0; - const nativeControl = { - kind: "injection" as const, - relayActive: false, - attached: true, - ended: false, - attach() { throw new Error("Native injection transport is already owned."); }, - observe() { return false; }, - steer() { throw new Error("unreachable"); }, - continue() { return false; }, - }; - const options = { session, url: CODEX_URL, init, prepared, nativeControl, - sseFallback: (async () => { fallbacks++; throw new Error("attach conflict must not fall back"); }) as typeof fetch }; - try { - expect(session.reserve()).toBe(true); - const response = await codexWsExchange(options); - const ws = FakeWebSocket.instances.at(-1)!; - expect(fallbacks).toBe(0); - expect(ws.sent).toHaveLength(0); - expect(response.status).toBe(502); - expect(((await response.json()) as { error: { message: string } }).error.message).toContain("already owned"); - expect(ws.closed).toBe(true); - expect([...ws.listeners.values()].every(listeners => listeners.length === 0)).toBe(true); - } finally { session.dispose(); } - }); }); test.each(["error", "response.completed"])("multiline upstream %s JSON remains one valid SSE data value", async type => {