feat(ai): added parentTurnId support and capture turn-state refreshes in Codex

- Added first-class `parentTurnId` stream option and reserved metadata key for nested Codex requests.
- Captured `x-codex-turn-state` refreshes from `response.metadata` event headers in WebSocket and streaming sessions.
This commit is contained in:
can1357
2026-07-29 23:33:29 +02:00
parent b7cd8c6582
commit eab5b50834
3 changed files with 126 additions and 9 deletions
+2
View File
@@ -14,6 +14,8 @@
### Fixed
- Fixed aborted usage-limit recovery waiting on a local usage fetch or marking the resolved credential blocked after the caller had already moved to another session ([#6883](https://github.com/can1357/oh-my-pi/pull/6883) by [@paolomazzitti](https://github.com/paolomazzitti)).
- Codex sessions now capture `x-codex-turn-state` refreshes from `response.metadata` event headers, matching codex-rs, which mirrors per-response HTTP headers into that event on the WebSocket transport. Previously WebSocket sessions only picked up turn state from the connection handshake, so same-turn follow-ups could echo a stale or missing turn state.
- Added first-class `parentTurnId` support for nested Codex requests (codex-rs `parent_turn_id`, upstream #35835): stream options and `createOpenAICodexCompatibilityMetadata` accept the initiating turn's id and emit it in both upstream projections — the flat `client_metadata.parent_turn_id` key and the `x-codex-turn-metadata` JSON blob — with blank values ignored. The key is also reserved, so caller-supplied `clientMetadata` extras cannot forge nested-request provenance.
- OpenAI Responses requests for Harmony-dialect models (gpt-5.x / openai-codex) now escape reserved Harmony control-token spellings (e.g. `<|channel|>analysis`) in untrusted user and tool-result text before serialization, so ordinary data — documentation, code, logs, or `omp://` grep results — can no longer trip the provider's `invalid_prompt` / "Request blocked" validator and permanently poison the session. Only the transport copy is escaped; the persisted transcript is byte-for-byte unchanged, and non-Harmony models are untouched. Extends the existing compaction-payload escaping to the live agent-loop request boundary ([#6913](https://github.com/can1357/oh-my-pi/issues/6913)).
- Fixed named forced `tool_choice` no longer being enforced on string-only OpenAI-compatible hosts (llama.cpp, LM Studio): the named object degrades to `"required"`, which alone lets the host call any advertised tool, so the request now narrows the advertised tools to the forced one (mirroring the Ollama chat transport). When the forced tool is absent the full tool list is kept and the choice is dropped for an unforced turn ([#6925](https://github.com/can1357/oh-my-pi/issues/6925)).
- Fixed direct Anthropic Claude Opus 5 requests failing with HTTP 400 when the endpoint rejected `strict` tool fields; the existing non-strict retry now recognizes `tools.N.custom.strict: Extra inputs are not permitted` responses ([#6976](https://github.com/can1357/oh-my-pi/issues/6976)).
@@ -147,6 +147,14 @@ export interface OpenAICodexResponsesOptions extends StreamOptions {
* keys are ignored; extras are never emitted as top-level metadata fields.
*/
clientMetadata?: Record<string, string>;
/**
* Turn id of the initiating (parent) Codex turn for nested requests such as
* subagent spawns (codex-rs `parent_turn_id`, #35835). Emitted both as the
* flat `client_metadata.parent_turn_id` key and inside the
* `x-codex-turn-metadata` JSON blob; blank values are ignored. The key is
* reserved: `clientMetadata` extras cannot supply it.
*/
parentTurnId?: string;
/**
* Invoked when the server streams a `response.metadata` event carrying
* ChatGPT moderation metadata (`metadata.openai_chatgpt_moderation_metadata`)
@@ -165,6 +173,8 @@ export interface OpenAICodexCompatibilityMetadataOptions {
startNewTurn?: boolean;
turnStartedAtUnixMs?: number;
clientMetadata?: Readonly<Record<string, string>>;
/** Parent Codex turn id for nested requests; see {@link OpenAICodexResponsesOptions.parentTurnId}. */
parentTurnId?: string;
/** Add the direct installation header required by `/responses/compact`. */
includeInstallationHeader?: boolean;
}
@@ -495,6 +505,7 @@ const CODEX_RESERVED_METADATA_KEYS: Record<string, true> = {
turn_started_at_unix_ms: true,
forked_from_thread_id: true,
parent_thread_id: true,
parent_turn_id: true,
subagent_kind: true,
thread_source: true,
sandbox: true,
@@ -566,6 +577,7 @@ function createCodexRequestMetadata(
startNewTurn: boolean;
turnStartedAtUnixMs?: number;
clientMetadata?: Readonly<Record<string, string>>;
parentTurnId?: string;
compaction?: CodexCompactionRequestContext;
},
): CodexRequestMetadata {
@@ -574,6 +586,9 @@ function createCodexRequestMetadata(
session.turnStartedAtUnixMs = options.turnStartedAtUnixMs;
}
const identity = createCodexCompatibilityIdentity(session);
// codex-rs `set_parent_turn_id` ignores blank values; keep the original
// spelling when non-blank.
const parentTurnId = options.parentTurnId?.trim() ? options.parentTurnId : undefined;
const extra: Record<string, string> = {};
const callerMetadata = options.clientMetadata;
if (callerMetadata) {
@@ -589,6 +604,7 @@ function createCodexRequestMetadata(
window_id: identity.windowId,
request_kind: requestKind,
};
if (parentTurnId) turnMetadata.parent_turn_id = parentTurnId;
if (options.compaction) {
turnMetadata.compaction = {
trigger: options.compaction.trigger,
@@ -603,18 +619,22 @@ function createCodexRequestMetadata(
}
for (const key in extra) turnMetadata[key] = extra[key];
const turnMetadataJson = toAsciiJsonString(turnMetadata);
const clientMetadata: Record<string, string> = {
[OPENAI_HEADERS.INSTALLATION_ID]: identity.installationId,
session_id: identity.sessionId,
thread_id: identity.threadId,
[OPENAI_HEADERS.WINDOW_ID]: identity.windowId,
turn_id: session.turnId,
};
// Both projections, mirroring codex-rs `CodexResponsesMetadata::to_client_metadata`:
// the flat key above/below AND the field inside the turn-metadata JSON.
if (parentTurnId) clientMetadata.parent_turn_id = parentTurnId;
clientMetadata[OPENAI_HEADERS.TURN_METADATA] = turnMetadataJson;
return {
...identity,
turnId: session.turnId,
turnMetadataJson,
clientMetadata: {
[OPENAI_HEADERS.INSTALLATION_ID]: identity.installationId,
session_id: identity.sessionId,
thread_id: identity.threadId,
[OPENAI_HEADERS.WINDOW_ID]: identity.windowId,
turn_id: session.turnId,
[OPENAI_HEADERS.TURN_METADATA]: turnMetadataJson,
},
clientMetadata,
};
}
@@ -649,6 +669,7 @@ export function createOpenAICodexCompatibilityMetadata(
startNewTurn,
turnStartedAtUnixMs: options.turnStartedAtUnixMs ?? (startNewTurn || !session.turnId ? Date.now() : undefined),
clientMetadata: options.clientMetadata,
parentTurnId: options.parentTurnId,
compaction: options.compaction,
});
const headers = new Headers();
@@ -1411,6 +1432,7 @@ async function buildCodexRequestContext(
: undefined
: getCodexTurnStartedAtUnixMs(context),
clientMetadata: transformedBody.client_metadata,
parentTurnId: options?.parentTurnId,
compaction,
});
transformedBody.client_metadata = requestMetadata.clientMetadata;
@@ -2075,6 +2097,12 @@ class CodexStreamProcessor {
}
if (eventType === "response.metadata") {
// The WebSocket transport has no per-response HTTP headers; codex-rs
// mirrors them into this event's `headers` and reads
// `x-codex-turn-state` from there (ResponsesStreamEvent::turn_state).
// Pick up the refresh so same-turn follow-ups echo the latest turn
// state on either transport.
updateCodexSessionMetadataFromHeaders(this.requestContext.websocketState, toCodexHeaders(rawEvent.headers));
const moderation = asRecord(rawEvent.metadata)?.[CODEX_MODERATION_METADATA_KEY];
if (moderation !== undefined) {
try {
+88 -1
View File
@@ -1629,7 +1629,12 @@ describe("openai-codex streaming", () => {
sessionId: "ws-lite-session",
providerSessionState: new Map<string, ProviderSessionState>(),
responsesLite: true,
clientMetadata: { workspace_kind: "repo", "x-codex-turn-metadata": '{"thread_id":"caller"}' },
clientMetadata: {
workspace_kind: "repo",
parent_turn_id: "forged-parent-turn",
"x-codex-turn-metadata": '{"thread_id":"caller"}',
},
parentTurnId: "turn_parent-1",
},
).result();
@@ -1654,6 +1659,12 @@ describe("openai-codex streaming", () => {
request_kind: "turn",
workspace_kind: "repo",
});
// `parent_turn_id` is reserved (codex-rs PARENT_TURN_ID_KEY): only the
// first-class option feeds it — caller extras cannot forge provenance —
// and it lands in both projections: the flat client_metadata key and the
// x-codex-turn-metadata JSON blob.
expect(metadata.parent_turn_id).toBe("turn_parent-1");
expect(turnMetadata.parent_turn_id).toBe("turn_parent-1");
expect(capturedHeaders?.["x-codex-installation-id"]).toBeUndefined();
expect(metadata.session_id).toBe(capturedHeaders?.["session-id"]);
expect(metadata.thread_id).toBe(capturedHeaders?.["thread-id"]);
@@ -4946,6 +4957,82 @@ describe("openai-codex streaming", () => {
expect(requestTurnStates).toEqual([null, "turn-state-1", null]);
});
it("captures x-codex-turn-state from response.metadata event headers", async () => {
const tempDir = TempDir.createSync("@pi-codex-stream-");
setAgentDir(tempDir.path());
const requestTurnStates: Array<string | null> = [];
let callCount = 0;
const fetchMock = vi.fn(async (_input: string | URL, init?: RequestInit) => {
const headers = init?.headers instanceof Headers ? init.headers : new Headers(init?.headers);
requestTurnStates.push(headers.get("x-codex-turn-state"));
const index = callCount;
callCount += 1;
// No x-codex-turn-state HTTP response header: turn state arrives only
// via the response.metadata event's mirrored headers, the way the
// WebSocket transport delivers it.
const sse =
index === 0
? `${[
`data: ${JSON.stringify({ type: "response.metadata", headers: { "x-codex-turn-state": "meta-turn-state-1" } })}`,
`data: ${JSON.stringify({ type: "response.output_item.added", item: { type: "function_call", id: "fc_1", call_id: "call_1", name: "read_file", arguments: "" } })}`,
`data: ${JSON.stringify({ type: "response.output_item.done", item: { type: "function_call", id: "fc_1", call_id: "call_1", name: "read_file", arguments: '{"path":"README.md"}' } })}`,
`data: ${JSON.stringify({ type: "response.completed", response: { status: "completed", usage: { input_tokens: 5, output_tokens: 3, total_tokens: 8, input_tokens_details: { cached_tokens: 0 } } } })}`,
].join("\n\n")}\n\n`
: `${[
`data: ${JSON.stringify({ type: "response.output_item.added", item: { type: "message", id: "msg_1", role: "assistant", status: "in_progress", content: [] } })}`,
`data: ${JSON.stringify({ type: "response.output_item.done", item: { type: "message", id: "msg_1", role: "assistant", status: "completed", content: [{ type: "output_text", text: "Done" }] } })}`,
`data: ${JSON.stringify({ type: "response.completed", response: { status: "completed", usage: { input_tokens: 5, output_tokens: 3, total_tokens: 8, input_tokens_details: { cached_tokens: 0 } } } })}`,
].join("\n\n")}\n\n`;
return new Response(sse, { status: 200, headers: new Headers({ "content-type": "text/event-stream" }) });
});
const model: Model<"openai-codex-responses"> = buildModel({
id: "gpt-5.1-codex",
name: "GPT-5.1 Codex",
api: "openai-codex-responses",
provider: "openai-codex",
baseUrl: "https://chatgpt.com/backend-api",
reasoning: true,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 400000,
maxTokens: 128000,
});
const systemPrompt = ["You are a helpful assistant."];
const firstUser = { role: "user" as const, content: "Read the file", timestamp: Date.now() };
const providerSessionState = new Map<string, ProviderSessionState>();
const options = {
fetch: fetchMock as FetchImpl,
apiKey: createCodexTestToken(),
sessionId: "metadata-turn-state-session",
providerSessionState,
};
const first = await streamOpenAICodexResponses(model, { systemPrompt, messages: [firstUser] }, options).result();
const toolCall = first.content.find(
(c): c is Extract<(typeof first.content)[number], { type: "toolCall" }> => c.type === "toolCall",
);
expect(toolCall).toBeDefined();
const toolResult = {
role: "toolResult" as const,
toolCallId: toolCall!.id,
toolName: toolCall!.name,
content: [{ type: "text" as const, text: "file contents" }],
isError: false,
timestamp: Date.now(),
};
// The within-turn follow-up replays the turn state captured from the event.
await streamOpenAICodexResponses(
model,
{ systemPrompt, messages: [firstUser, first, toolResult] },
options,
).result();
expect(requestTurnStates).toEqual([null, "meta-turn-state-1"]);
});
it("drops stale frames from a prior response before sending the next websocket request", async () => {
const tempDir = TempDir.createSync("@pi-codex-stream-");
setAgentDir(tempDir.path());