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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 74 additions & 0 deletions docs-site/src/content/docs/guides/codex-integration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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://<extension-id>` 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. |
Expand Down
1 change: 1 addition & 0 deletions scripts/test-layout/layout.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": [
Expand Down
1 change: 1 addition & 0 deletions src/config/schema/config-schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
34 changes: 22 additions & 12 deletions src/server/index/websocket-handler.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -190,14 +193,19 @@ 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.nativeControl)) {
try {
if (frame.type === "response.inject") {
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",
Expand All @@ -213,12 +221,14 @@ 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;
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.codexNativeSteering === true ? new NativeSteeringChannel(frame, idleMs) : undefined;
nativeControl = 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;
Expand Down Expand Up @@ -258,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<string, unknown> = { ...frame };
delete payload.type;
turnAdmissionLease.bindAbortController(turnAbort);
Expand Down Expand Up @@ -299,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),
Expand All @@ -310,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 => {
Expand Down Expand Up @@ -349,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;
}
Expand Down
26 changes: 18 additions & 8 deletions src/server/responses/codex-ws-exchange.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { markNativeSteeringResponse, type NativeSteeringChannel } 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";
import { CodexWsMetadata, type CodexWsQuotaObserver } from "./codex-ws-metadata";
Expand All @@ -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;
nativeControl?: NativeResponseControl;
beforeContinuation?: () => Promise<void>;
session: CodexWsSession;
url: string;
Expand Down Expand Up @@ -89,7 +90,7 @@ function wrappedRejectionResponse(payload: Record<string, unknown>, prelude: Hea

/** The sole SSE exchange state machine for both one-shot and retained sockets. */
export function codexWsExchange(options: ExchangeOptions): Promise<Response> {
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<Response>((resolve, reject) => {
Expand Down Expand Up @@ -203,7 +204,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise<Response> {
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);
Expand Down Expand Up @@ -336,8 +337,8 @@ export function codexWsExchange(options: ExchangeOptions): Promise<Response> {
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");
Expand Down Expand Up @@ -368,6 +369,15 @@ export function codexWsExchange(options: ExchangeOptions): Promise<Response> {
} 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 {
Expand Down Expand Up @@ -441,7 +451,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise<Response> {
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
Expand Down Expand Up @@ -485,7 +495,7 @@ export function codexWsExchange(options: ExchangeOptions): Promise<Response> {
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();
Expand Down
4 changes: 2 additions & 2 deletions src/server/responses/core-options.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -53,7 +53,7 @@ export interface HandleResponsesOptions {
onRequestBodyRead?: () => void;
forceEmptyResponseId?: boolean;
/** Internal, connection-owned control channel; never reconstructed from headers. */
nativeSteering?: NativeSteeringChannel;
nativeControl?: NativeResponseControl;
abortSignal?: AbortSignal;
/** One-shot TTFT callback: first non-empty model output observed (WP4). */
onFirstOutput?: () => void;
Expand Down
6 changes: 3 additions & 3 deletions src/server/responses/fetch-helpers.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import type { NativeSteeringChannel } from "./native-steering";
import type { NativeResponseControl } from "./native-response-control";
import type { Server } from "bun";
import {
codexWsUpstreamFetch,
Expand Down Expand Up @@ -54,7 +54,7 @@ export interface PaceAwareFetch {
export type ProviderFetch = typeof globalThis.fetch & PaceAwareFetch;

export interface ProviderFetchOptions {
nativeSteering?: NativeSteeringChannel;
nativeControl?: NativeResponseControl;
providerName?: string;
modelId?: string;
/** One pacing slot was acquired immediately before this fetch wrapper was created. */
Expand Down Expand Up @@ -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);
Expand Down
Loading
Loading