diff --git a/packages/agent/src/compaction/compaction-v2-streaming.ts b/packages/agent/src/compaction/compaction-v2-streaming.ts index 251573775..14a2318f8 100644 --- a/packages/agent/src/compaction/compaction-v2-streaming.ts +++ b/packages/agent/src/compaction/compaction-v2-streaming.ts @@ -7,9 +7,14 @@ * compaction item as replacement history. */ -import type { Api, FetchImpl, Model } from "@oh-my-pi/pi-ai"; +import type { Api, CodexCompactionContext, FetchImpl, Model, ProviderSessionState } from "@oh-my-pi/pi-ai"; import { isTransientStatus, ProviderHttpError } from "@oh-my-pi/pi-ai/error"; import { applyCodexResponsesLiteShape } from "@oh-my-pi/pi-ai/providers/openai-codex/request-transformer"; +import { + createOpenAICodexCompactionRequestContext, + createOpenAICodexCompatibilityMetadata, + type OpenAICodexCompatibilityMetadata, +} from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { getOpenAIPromptCacheKey, getOpenAIResponsesRoutingSessionId, @@ -220,6 +225,8 @@ export async function requestCompactionV2Streaming( fetch?: FetchImpl; timeoutMs?: number; retryWait?: (delayMs: number, signal?: AbortSignal) => Promise; + providerSessionState?: Map; + codexCompaction?: CodexCompactionContext; }, ): Promise { const endpoint = getCompactionV2Endpoint(model); @@ -229,12 +236,32 @@ export async function requestCompactionV2Streaming( const fetchImpl = options?.fetch ?? globalThis.fetch; const retryWait = options?.retryWait ?? ((delayMs: number) => Bun.sleep(delayMs)); + const isCodexResponses = compactionV2Api(model) === "openai-codex-responses" || model.provider === "openai-codex"; + const codexMetadata = isCodexResponses + ? createOpenAICodexCompatibilityMetadata({ + sessionId: request.sessionId, + providerSessionState: options?.providerSessionState, + requestKind: "compaction", + compaction: createOpenAICodexCompactionRequestContext({ + context: options?.codexCompaction, + implementation: "responses_compaction_v2", + }), + }) + : undefined; let lastError: Error | undefined; for (let attempt = 0; attempt <= V2_COMPACTION_MAX_RETRIES; attempt++) { const timeoutSignal = withRequestTimeout(signal, options?.timeoutMs ?? V2_COMPACTION_TIMEOUT_MS); try { - return await attemptCompactionV2Streaming(endpoint, apiKey, model, request, fetchImpl, timeoutSignal); + return await attemptCompactionV2Streaming( + endpoint, + apiKey, + model, + request, + fetchImpl, + timeoutSignal, + codexMetadata, + ); } catch (err) { const error = err instanceof Error ? err : new Error(String(err)); if (signal?.aborted) throw error; @@ -265,6 +292,7 @@ async function attemptCompactionV2Streaming( request: CompactionV2Request, fetchImpl: FetchImpl, signal?: AbortSignal, + codexMetadata?: OpenAICodexCompatibilityMetadata, ): Promise { // Faithful to Codex: append the compaction trigger as the final input item // of an otherwise-normal Responses request, then stream the result. `store` @@ -287,6 +315,9 @@ async function attemptCompactionV2Streaming( ...(promptCacheKey ? { prompt_cache_key: promptCacheKey } : {}), ...(request.tools && request.tools.length > 0 ? { tools: request.tools, tool_choice: "auto" } : {}), }; + if (codexMetadata) { + body.client_metadata = codexMetadata.clientMetadata; + } // Responses Lite models take the same rewrite on the compaction stream: // instructions/tools ride as input items (codex-rs `compact_remote_v2` // builds through `build_responses_request`). @@ -295,7 +326,7 @@ async function attemptCompactionV2Streaming( } const response = await fetchImpl(endpoint, { method: "POST", - headers: buildCompactionV2Headers(model, apiKey, request), + headers: buildCompactionV2Headers(model, apiKey, request, codexMetadata), body: JSON.stringify(body), signal, }); @@ -320,7 +351,12 @@ async function attemptCompactionV2Streaming( return collectCompactionV2Output(response, request); } -function buildCompactionV2Headers(model: Model, apiKey: string, request: CompactionV2Request): Record { +function buildCompactionV2Headers( + model: Model, + apiKey: string, + request: CompactionV2Request, + codexMetadata?: OpenAICodexCompatibilityMetadata, +): Record { const api = compactionV2Api(model); const cacheOptions = { sessionId: request.sessionId, promptCacheKey: request.promptCacheKey }; const routingSessionId = getOpenAIResponsesRoutingSessionId(cacheOptions); @@ -355,6 +391,7 @@ function buildCompactionV2Headers(model: Model, apiKey: string, request: Compact headers[OPENAI_HEADERS.RESPONSES_LITE] = "true"; } } + if (codexMetadata) Object.assign(headers, codexMetadata.headers); return headers; } diff --git a/packages/agent/src/compaction/compaction.ts b/packages/agent/src/compaction/compaction.ts index 36b8feadd..772f0b232 100644 --- a/packages/agent/src/compaction/compaction.ts +++ b/packages/agent/src/compaction/compaction.ts @@ -9,18 +9,21 @@ import { type Api, type ApiKey, type AssistantMessage, + type CodexCompactionContext, type Context, Effort, type FetchImpl, type Message, type MessageAttribution, type Model, + type ProviderSessionState, type SimpleStreamOptions, type Tool, type Usage, withAuth, } from "@oh-my-pi/pi-ai"; import { ProviderHttpError } from "@oh-my-pi/pi-ai/error"; +import { createOpenAICodexCompactionRequestContext } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { convertTools } from "@oh-my-pi/pi-ai/providers/openai-responses"; import { buildResponsesInput, resolveOpenAICompatPolicy } from "@oh-my-pi/pi-ai/providers/openai-shared"; import { preferredDialect } from "@oh-my-pi/pi-catalog/identity"; @@ -736,6 +739,10 @@ export interface SummaryOptions { sessionId?: string; /** Prompt-cache key for remote compaction transports that support provider prefix caching. */ promptCacheKey?: string; + /** Mutable provider state used to keep Codex compaction on the live session identity. */ + providerSessionState?: Map; + /** Classification shared by every provider request in this logical compaction. */ + codexCompaction?: CodexCompactionContext; /** Provider-visible tools for remote compaction transports that replay native tool history. */ tools?: Tool[]; /** Optional fetch implementation threaded into remote compaction calls. */ @@ -755,6 +762,13 @@ export interface SummaryOptions { ) => Promise; } +function localCodexCompaction(options: SummaryOptions | undefined) { + return createOpenAICodexCompactionRequestContext({ + context: options?.codexCompaction, + implementation: "responses", + }); +} + function formatPreviousSnapcompactArchive(archiveText: string): string { return prompt.render(snapcompactArchiveContextPrompt, { archiveText }); } @@ -844,6 +858,11 @@ export async function generateSummary( reasoning: resolveCompactionEffort(model, options?.thinkingLevel), initiatorOverride: options?.initiatorOverride, metadata: options?.metadata, + fetch: options?.fetch, + sessionId: options?.sessionId, + promptCacheKey: options?.promptCacheKey, + providerSessionState: options?.providerSessionState, + codexCompaction: localCodexCompaction(options), }, { telemetry: options?.telemetry, oneshotKind: "compaction_summary", completeImpl: options?.completeImpl }, ); @@ -1047,6 +1066,11 @@ async function generateShortSummary( reasoning: resolveCompactionEffort(model, options?.thinkingLevel), initiatorOverride: options?.initiatorOverride, metadata: options?.metadata, + fetch: options?.fetch, + sessionId: options?.sessionId, + promptCacheKey: options?.promptCacheKey, + providerSessionState: options?.providerSessionState, + codexCompaction: localCodexCompaction(options), }, { telemetry: options?.telemetry, oneshotKind: "compaction_short_summary", completeImpl: options?.completeImpl }, ); @@ -1317,6 +1341,8 @@ export async function compact( thinkingLevel: options?.thinkingLevel, sessionId: options?.sessionId, promptCacheKey: options?.promptCacheKey, + providerSessionState: options?.providerSessionState, + codexCompaction: options?.codexCompaction, tools: options?.tools, fetch: options?.fetch, completeImpl: options?.completeImpl, @@ -1375,7 +1401,12 @@ export async function compact( ); const remote = await withAuth( apiKey, - key => requestCompactionV2Streaming(model, key, request, signal, { fetch: summaryOptions.fetch }), + key => + requestCompactionV2Streaming(model, key, request, signal, { + fetch: summaryOptions.fetch, + providerSessionState: summaryOptions.providerSessionState, + codexCompaction: summaryOptions.codexCompaction, + }), { signal }, ); preserveData = { ...(preserveData ?? {}), ...storeCompactionV2PreserveData(remote, model) }; @@ -1419,7 +1450,12 @@ export async function compact( remoteHistory, summaryOptions.remoteInstructions ?? SUMMARIZATION_SYSTEM_PROMPT, signal, - { fetch: summaryOptions.fetch }, + { + fetch: summaryOptions.fetch, + sessionId: summaryOptions.sessionId, + providerSessionState: summaryOptions.providerSessionState, + codexCompaction: summaryOptions.codexCompaction, + }, ), { signal }, ); @@ -1495,16 +1531,9 @@ export async function compact( const shortSummary = usedRemoteCompaction ? "Remote compaction" : await generateShortSummary(recentMessages, summary, model, reserveTokens, apiKey, signal, { + ...summaryOptions, extraContext: options?.extraContext, - remoteEndpoint: summaryOptions.remoteEndpoint, - initiatorOverride: summaryOptions.initiatorOverride, - metadata: summaryOptions.metadata, - telemetry: summaryOptions.telemetry, - // Same propagation as summaryOptions above β€” generateShortSummary - // resolves its own reasoning via resolveCompactionEffort. thinkingLevel: options?.thinkingLevel, - fetch: summaryOptions.fetch, - completeImpl: summaryOptions.completeImpl, }); // Compute file lists and append to summary @@ -1567,6 +1596,11 @@ async function generateTurnPrefixSummary( reasoning: resolveCompactionEffort(model, options?.thinkingLevel), initiatorOverride: options?.initiatorOverride, metadata: options?.metadata, + fetch: options?.fetch, + sessionId: options?.sessionId, + promptCacheKey: options?.promptCacheKey, + providerSessionState: options?.providerSessionState, + codexCompaction: localCodexCompaction(options), }, { telemetry: options?.telemetry, oneshotKind: "compaction_turn_prefix", completeImpl: options?.completeImpl }, ); diff --git a/packages/agent/src/compaction/openai.ts b/packages/agent/src/compaction/openai.ts index aa4f37bb4..fbac0a8b7 100644 --- a/packages/agent/src/compaction/openai.ts +++ b/packages/agent/src/compaction/openai.ts @@ -17,9 +17,21 @@ import { ProviderHttpError } from "@oh-my-pi/pi-ai/error"; import { applyCodexResponsesLiteShape } from "@oh-my-pi/pi-ai/providers/openai-codex/request-transformer"; +import { + createOpenAICodexCompactionRequestContext, + createOpenAICodexCompatibilityMetadata, +} from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { parseAzureDeploymentNameMap, parseTextSignature } from "@oh-my-pi/pi-ai/providers/openai-shared"; import { transformMessages } from "@oh-my-pi/pi-ai/providers/transform-messages"; -import type { Api, AssistantMessage, FetchImpl, Message, Model } from "@oh-my-pi/pi-ai/types"; +import type { + Api, + AssistantMessage, + CodexCompactionContext, + FetchImpl, + Message, + Model, + ProviderSessionState, +} from "@oh-my-pi/pi-ai/types"; import { getOpenAIResponsesHistoryItems, getOpenAIResponsesHistoryPayload, @@ -461,7 +473,13 @@ export async function requestOpenAiRemoteCompaction( compactInput: Array>, instructions: string, signal?: AbortSignal, - opts?: { fetch?: FetchImpl; timeoutMs?: number }, + opts?: { + fetch?: FetchImpl; + timeoutMs?: number; + sessionId?: string; + providerSessionState?: Map; + codexCompaction?: CodexCompactionContext; + }, ): Promise { const endpoint = resolveOpenAiCompactEndpoint(model); const requestModel = resolveOpenAiCompactModel(model); @@ -474,6 +492,8 @@ export async function requestOpenAiRemoteCompaction( instructions, }; const isAzureOpenAiResponses = (model.remoteCompaction?.api ?? model.api) === "azure-openai-responses"; + const isCodexResponses = + model.provider === "openai-codex" || (model.remoteCompaction?.api ?? model.api) === "openai-codex-responses"; const headers: Record = isAzureOpenAiResponses ? { "content-type": "application/json", @@ -487,13 +507,26 @@ export async function requestOpenAiRemoteCompaction( }; // Codex endpoints require additional auth headers - if (model.provider === "openai-codex") { + if (isCodexResponses) { const accountId = getCodexAccountId(apiKey); if (accountId) { headers[OPENAI_HEADERS.ACCOUNT_ID] = accountId; } headers[OPENAI_HEADERS.BETA] = OPENAI_HEADER_VALUES.BETA_RESPONSES; headers[OPENAI_HEADERS.ORIGINATOR] = OPENAI_HEADER_VALUES.ORIGINATOR_CODEX; + Object.assign( + headers, + createOpenAICodexCompatibilityMetadata({ + sessionId: opts?.sessionId, + providerSessionState: opts?.providerSessionState, + requestKind: "compaction", + compaction: createOpenAICodexCompactionRequestContext({ + context: opts?.codexCompaction, + implementation: "responses_compact", + }), + includeInstallationHeader: true, + }).headers, + ); // Responses Lite models take the same rewrite on `/responses/compact`: // instructions ride as an input item and the lite marker header is set // (codex-rs routes compaction through `build_responses_request`). diff --git a/packages/agent/test/remote-compaction.test.ts b/packages/agent/test/remote-compaction.test.ts index 237c81557..78d4cb97b 100644 --- a/packages/agent/test/remote-compaction.test.ts +++ b/packages/agent/test/remote-compaction.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, test, vi } from "bun:test"; +import { afterEach, beforeEach, describe, expect, test, vi } from "bun:test"; import { type CompactionPreparation, compact, @@ -18,10 +18,36 @@ import { shouldUseOpenAiRemoteCompaction, } from "@oh-my-pi/pi-agent-core/compaction/openai"; import * as ai from "@oh-my-pi/pi-ai"; -import type { AssistantMessage, FetchImpl, Model, ToolResultMessage } from "@oh-my-pi/pi-ai/types"; +import { getOpenAICodexTransportDetails } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; +import type { + AssistantMessage, + CodexCompactionContext, + FetchImpl, + Model, + ProviderSessionState, + ToolResultMessage, +} from "@oh-my-pi/pi-ai/types"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; import type { ModelSpec } from "@oh-my-pi/pi-catalog/types"; -import { isRecord } from "@oh-my-pi/pi-utils"; +import * as piUtils from "@oh-my-pi/pi-utils"; + +const { isRecord } = piUtils; +const TEST_INSTALLATION_ID = "00000000-0000-4000-8000-000000000001"; +const TEST_CODEX_COMPACTION: CodexCompactionContext = { + operationId: "compaction-operation-1", + trigger: "auto", + reason: "context_limit", + phase: "pre_turn", + strategy: "memento", +}; + +beforeEach(() => { + vi.spyOn(piUtils, "getInstallId").mockReturnValue(TEST_INSTALLATION_ID); +}); + +afterEach(() => { + vi.restoreAllMocks(); +}); function makeOpenAiModel(overrides: Partial> = {}): Model<"openai-responses"> { return buildModel({ @@ -394,7 +420,9 @@ describe("requestCompactionV2Streaming", () => { }); describe("Responses Lite remote compaction", () => { - function makeCodexLiteModel(): Model<"openai-codex-responses"> { + function makeCodexLiteModel( + overrides: Partial> = {}, + ): Model<"openai-codex-responses"> { return buildModel({ id: "gpt-5.6-terra", name: "GPT-5.6 Terra", @@ -402,12 +430,14 @@ describe("Responses Lite remote compaction", () => { provider: "openai-codex", baseUrl: "https://chatgpt.example/backend-api", reasoning: true, + preferWebsockets: false, input: ["text", "image"], cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, contextWindow: 372000, maxTokens: 128000, useResponsesLite: true, remoteCompaction: { enabled: true, api: "openai-codex-responses", v2StreamingEnabled: true }, + ...overrides, }); } @@ -415,22 +445,42 @@ describe("Responses Lite remote compaction", () => { instructions?: unknown; tools?: unknown; input?: Array>; + client_metadata?: unknown; } - function captureLite(init: RequestInit | undefined): { body: CapturedLiteRequest; liteHeader?: string } { + interface CapturedLiteExchange { + body: CapturedLiteRequest; + headers: Headers; + } + + function parseCodexTurnMetadata(value: unknown): Record { + if (typeof value !== "string") throw new Error("expected x-codex-turn-metadata"); + const parsed: unknown = JSON.parse(value); + if (!isRecord(parsed)) throw new Error("expected Codex turn metadata object"); + return parsed; + } + + function captureLite(init: RequestInit | undefined): CapturedLiteExchange { if (!init?.headers || init.headers instanceof Headers || Array.isArray(init.headers)) { throw new Error("Expected remote compaction to send headers as a plain object"); } - const rawLite = init.headers["x-openai-internal-codex-responses-lite"]; return { body: JSON.parse(String(init.body)) as CapturedLiteRequest, - liteHeader: typeof rawLite === "string" ? rawLite : undefined, + headers: new Headers(init.headers), + }; + } + + function captureStreamLite(init: RequestInit | undefined): CapturedLiteExchange { + if (!init?.headers) throw new Error("Expected local compaction request headers"); + return { + body: JSON.parse(String(init.body)) as CapturedLiteRequest, + headers: new Headers(init.headers), }; } test("V1 compaction sends the lite header and input-item instructions", async () => { const model = makeCodexLiteModel(); - let captured: { body: CapturedLiteRequest; liteHeader?: string } | undefined; + let captured: CapturedLiteExchange | undefined; const fetchMock: FetchImpl = async (_input, init) => { captured = captureLite(init); return Response.json({ output: [{ type: "compaction", encrypted_content: "enc" }] }); @@ -442,12 +492,29 @@ describe("Responses Lite remote compaction", () => { [{ type: "message", role: "user", content: [{ type: "input_text", text: "hi" }] }], "compact instructions", undefined, - { fetch: fetchMock }, + { + fetch: fetchMock, + sessionId: "codex-compaction-session", + providerSessionState: new Map(), + codexCompaction: TEST_CODEX_COMPACTION, + }, ); - expect(captured?.liteHeader).toBe("true"); + expect(captured?.headers.get("x-openai-internal-codex-responses-lite")).toBe("true"); expect(captured?.body.instructions).toBeUndefined(); expect(captured?.body.tools).toBeUndefined(); + expect(captured?.body.client_metadata).toBeUndefined(); + expect(captured?.headers.get("x-codex-installation-id")).toBe(TEST_INSTALLATION_ID); + expect(captured?.headers.get("session-id")).toBe("codex-compaction-session"); + const v1TurnMetadata = parseCodexTurnMetadata(captured?.headers.get("x-codex-turn-metadata")); + expect(v1TurnMetadata.request_kind).toBe("compaction"); + expect(v1TurnMetadata.compaction).toEqual({ + trigger: "auto", + reason: "context_limit", + implementation: "responses_compact", + phase: "pre_turn", + strategy: "memento", + }); expect(captured?.body.input?.[0]).toEqual({ type: "additional_tools", role: "developer", tools: [] }); expect(captured?.body.input?.[1]).toEqual({ type: "message", @@ -462,8 +529,9 @@ describe("Responses Lite remote compaction", () => { model, [{ type: "message", role: "user", content: [{ type: "input_text", text: "real user" }] }], "compact instructions", + { sessionId: "codex-compaction-session" }, ); - let captured: { body: CapturedLiteRequest; liteHeader?: string } | undefined; + let captured: CapturedLiteExchange | undefined; const fetchMock: FetchImpl = async (_input, init) => { captured = captureLite(init); return sseResponse([ @@ -477,11 +545,30 @@ describe("Responses Lite remote compaction", () => { }; expect(shouldUseCompactionV2Streaming(model)).toBe(true); - await requestCompactionV2Streaming(model, "test-key", request, undefined, { fetch: fetchMock }); + await requestCompactionV2Streaming(model, "test-key", request, undefined, { + fetch: fetchMock, + providerSessionState: new Map(), + codexCompaction: TEST_CODEX_COMPACTION, + }); - expect(captured?.liteHeader).toBe("true"); + expect(captured?.headers.get("x-openai-internal-codex-responses-lite")).toBe("true"); expect(captured?.body.instructions).toBeUndefined(); expect(captured?.body.tools).toBeUndefined(); + if (!isRecord(captured?.body.client_metadata)) throw new Error("expected V2 client_metadata"); + const v2ClientMetadata = captured.body.client_metadata; + const v2TurnMetadata = parseCodexTurnMetadata(v2ClientMetadata["x-codex-turn-metadata"]); + expect(captured.headers.get("x-codex-installation-id")).toBeNull(); + expect(v2ClientMetadata["x-codex-installation-id"]).toBe(TEST_INSTALLATION_ID); + expect(v2ClientMetadata.session_id).toBe(captured.headers.get("session-id")); + expect(v2ClientMetadata.thread_id).toBe(captured.headers.get("thread-id")); + expect(v2TurnMetadata.request_kind).toBe("compaction"); + expect(v2TurnMetadata.compaction).toEqual({ + trigger: "auto", + reason: "context_limit", + implementation: "responses_compaction_v2", + phase: "pre_turn", + strategy: "memento", + }); expect(captured?.body.input?.[0]).toEqual({ type: "additional_tools", role: "developer", tools: [] }); expect(captured?.body.input?.[1]).toEqual({ type: "message", @@ -490,6 +577,236 @@ describe("Responses Lite remote compaction", () => { }); expect(captured?.body.input?.at(-1)).toEqual({ type: "compaction_trigger" }); }); + + test("compact fan-out keeps local Codex summaries on one classified turn", async () => { + const model = makeCodexLiteModel(); + const captured: CapturedLiteExchange[] = []; + const fetchMock: FetchImpl = async (_input, init) => { + captured.push(captureStreamLite(init)); + return sseResponse([ + { + type: "response.output_item.added", + output_index: 0, + item: { type: "message", id: "msg_summary", role: "assistant", status: "in_progress", content: [] }, + }, + { + type: "response.content_part.added", + output_index: 0, + content_index: 0, + part: { type: "output_text", text: "" }, + }, + { type: "response.output_text.delta", output_index: 0, content_index: 0, delta: "local summary" }, + { + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", + id: "msg_summary", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "local summary" }], + }, + }, + { + type: "response.completed", + response: { + status: "completed", + usage: { + input_tokens: 8, + output_tokens: 2, + total_tokens: 10, + input_tokens_details: { cached_tokens: 0 }, + }, + }, + }, + ]); + }; + const preparation: CompactionPreparation = { + firstKeptEntryId: "kept-1", + messagesToSummarize: [{ role: "user", content: "long history", timestamp: 1 }], + turnPrefixMessages: [], + recentMessages: [{ role: "user", content: "recent", timestamp: 2 }], + isSplitTurn: false, + tokensBefore: 100_000, + fileOps: createFileOps(), + settings: { + ...DEFAULT_COMPACTION_SETTINGS, + remoteEnabled: false, + remoteStreamingV2Enabled: false, + }, + }; + + const result = await compact(preparation, model, "test-key", undefined, undefined, { + fetch: fetchMock, + sessionId: "codex-compaction-session", + providerSessionState: new Map(), + codexCompaction: TEST_CODEX_COMPACTION, + }); + + expect(result.summary).toContain("local summary"); + expect(captured).toHaveLength(2); + const turnIds: string[] = []; + for (const exchange of captured) { + if (!isRecord(exchange.body.client_metadata)) throw new Error("expected local client_metadata"); + const clientMetadata = exchange.body.client_metadata; + const turnMetadata = parseCodexTurnMetadata(clientMetadata["x-codex-turn-metadata"]); + expect(exchange.headers.get("x-codex-installation-id")).toBeNull(); + expect(clientMetadata["x-codex-installation-id"]).toBe(TEST_INSTALLATION_ID); + expect(turnMetadata.request_kind).toBe("compaction"); + expect(turnMetadata.compaction).toEqual({ + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "pre_turn", + strategy: "memento", + }); + if (typeof turnMetadata.turn_id !== "string") throw new Error("expected Codex turn id"); + turnIds.push(turnMetadata.turn_id); + } + expect(new Set(turnIds).size).toBe(1); + }); + + test("local Codex compaction isolates and closes transient websocket sessions", async () => { + const originalWebSocket = global.WebSocket; + const sockets: AgentCompactionWebSocket[] = []; + let responseCount = 0; + + class AgentCompactionWebSocket { + static readonly CONNECTING = 0; + static readonly OPEN = 1; + static readonly CLOSING = 2; + static readonly CLOSED = 3; + + readyState = AgentCompactionWebSocket.CONNECTING; + binaryType: "blob" | "arraybuffer" | "nodebuffer" = "blob"; + onopen: ((event: Event) => void) | null = null; + onmessage: ((event: MessageEvent) => void) | null = null; + onerror: ((event: Event) => void) | null = null; + onclose: ((event: Event) => void) | null = null; + readonly handshakeHeaders = { + "x-codex-turn-state": `agent-compaction-state-${sockets.length}`, + }; + + constructor( + readonly url: string, + readonly options?: { headers?: Record }, + ) { + sockets.push(this); + queueMicrotask(() => { + this.readyState = AgentCompactionWebSocket.OPEN; + this.onopen?.(new Event("open")); + }); + } + + send(_data: string): void { + responseCount += 1; + const responseId = `response-${responseCount}`; + const messageId = `message-${responseCount}`; + const text = sockets[0] === this ? "main response" : "local summary"; + const events: Record[] = [ + { + type: "response.output_item.added", + item: { type: "message", id: messageId, role: "assistant", status: "in_progress", content: [] }, + }, + { type: "response.content_part.added", part: { type: "output_text", text: "" } }, + { type: "response.output_text.delta", delta: text }, + { + type: "response.output_item.done", + item: { + type: "message", + id: messageId, + role: "assistant", + status: "completed", + content: [{ type: "output_text", text }], + }, + }, + { + type: "response.done", + response: { + id: responseId, + status: "completed", + usage: { + input_tokens: 8, + output_tokens: 2, + total_tokens: 10, + input_tokens_details: { cached_tokens: 0 }, + }, + }, + }, + ]; + for (const event of events) { + this.onmessage?.({ data: JSON.stringify(event) } as MessageEvent); + } + } + + close(): void { + this.readyState = AgentCompactionWebSocket.CLOSED; + } + } + + const providerSessionState = new Map(); + try { + global.WebSocket = AgentCompactionWebSocket as unknown as typeof WebSocket; + const model = makeCodexLiteModel({ preferWebsockets: true }); + const sessionId = "agent-compaction-isolation"; + const fetchMock: FetchImpl = async () => { + throw new Error("Codex websocket compaction unexpectedly used SSE"); + }; + const main = await ai + .streamSimple( + model, + { + systemPrompt: ["You are a helpful assistant."], + messages: [{ role: "user", content: "Start the turn", timestamp: Date.now() }], + }, + { apiKey: "test-key", fetch: fetchMock, sessionId, providerSessionState }, + ) + .result(); + expect(main.stopReason).toBe("stop"); + expect(sockets).toHaveLength(1); + expect(sockets[0]?.readyState).toBe(AgentCompactionWebSocket.OPEN); + + const preparation: CompactionPreparation = { + firstKeptEntryId: "kept-1", + messagesToSummarize: [{ role: "user", content: "long history", timestamp: 1 }], + turnPrefixMessages: [], + recentMessages: [{ role: "user", content: "recent", timestamp: 2 }], + isSplitTurn: false, + tokensBefore: 100_000, + fileOps: createFileOps(), + settings: { + ...DEFAULT_COMPACTION_SETTINGS, + remoteEnabled: false, + remoteStreamingV2Enabled: false, + }, + }; + const result = await compact(preparation, model, "test-key", undefined, undefined, { + fetch: fetchMock, + sessionId, + providerSessionState, + codexCompaction: TEST_CODEX_COMPACTION, + }); + + expect(result.summary).toContain("local summary"); + expect(sockets).toHaveLength(3); + expect(sockets[0]?.readyState).toBe(AgentCompactionWebSocket.OPEN); + expect(sockets[1]?.readyState).toBe(AgentCompactionWebSocket.CLOSED); + expect(sockets[2]?.readyState).toBe(AgentCompactionWebSocket.CLOSED); + expect( + getOpenAICodexTransportDetails(model, { + sessionId, + providerSessionState, + }), + ).toMatchObject({ + websocketConnected: true, + hasTurnState: true, + }); + } finally { + for (const state of providerSessionState.values()) state.close(); + providerSessionState.clear(); + global.WebSocket = originalWebSocket; + } + }); }); test("uses configured OpenAI-compatible compaction for custom providers", async () => { diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index cb7865ac7..05e484348 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -20,6 +20,7 @@ ### Fixed - Fixed concurrent reasoning summaries to ignore legacy streaming events under cutoff contract +- Fixed sequential-cutoff Codex reasoning summaries repeating earlier content when atomic summary snapshots are replayed or extended. ## [16.3.15] - 2026-07-09 diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 86ca0085c..609a570fb 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -12,6 +12,7 @@ import { $flag, asRecord, fetchWithRetry, + getInstallId, logger, parseStreamingJson, readSseJson, @@ -24,6 +25,8 @@ import { getEnvApiKey } from "../stream"; import type { Api, AssistantMessage, + CodexCompactionContext, + CodexCompactionRequestContext, Context, FetchImpl, Model, @@ -128,9 +131,9 @@ export interface OpenAICodexResponsesOptions extends StreamOptions { */ responsesLite?: boolean; /** - * Extra `client_metadata` to include in the request body on both transports. - * The canonical Codex envelope is `client_metadata["x-codex-turn-metadata"]` - * (JSON string of thread/turn identifiers); flat keys are also accepted. + * Additional fields embedded in the canonical + * `client_metadata["x-codex-turn-metadata"]` JSON blob. Reserved identity + * keys are ignored; extras are never emitted as top-level metadata fields. */ clientMetadata?: Record; /** @@ -142,6 +145,49 @@ export interface OpenAICodexResponsesOptions extends StreamOptions { onModerationMetadata?: (metadata: unknown) => void; } +/** Inputs for synthesizing Codex request identity outside the normal stream path. */ +export interface OpenAICodexCompatibilityMetadataOptions { + sessionId?: string; + providerSessionState?: Map; + requestKind: OpenAICodexRequestKind; + compaction?: CodexCompactionRequestContext; + startNewTurn?: boolean; + turnStartedAtUnixMs?: number; + clientMetadata?: Readonly>; + /** Add the direct installation header required by `/responses/compact`. */ + includeInstallationHeader?: boolean; +} + +/** Canonical Codex body metadata and compatibility headers for one request. */ +export interface OpenAICodexCompatibilityMetadata { + clientMetadata: Record; + headers: Record; +} + +/** Live Codex session state to preserve after a successful history rewrite. */ +export interface OpenAICodexCompactionResetOptions { + providerSessionState?: Map; + sessionId?: string; + compaction: CodexCompactionContext; +} + +/** Add the selected wire implementation to one logical compaction context. */ +export function createOpenAICodexCompactionRequestContext(options: { + context: CodexCompactionContext | undefined; + implementation: "responses" | "responses_compaction_v2" | "responses_compact"; +}): CodexCompactionRequestContext | undefined { + const context = options.context; + if (!context) return undefined; + return { + operationId: context.operationId, + trigger: context.trigger, + reason: context.reason, + implementation: options.implementation, + phase: context.phase, + strategy: context.strategy, + }; +} + const CODEX_DEBUG = $flag("PI_CODEX_DEBUG"); const CODEX_MAX_RETRIES = 5; const CODEX_RETRY_DELAY_MS = 500; @@ -345,6 +391,237 @@ type CodexWebSocketSessionState = { interface CodexProviderSessionState extends ProviderSessionState { webSocketSessions: Map; webSocketPublicToPrivate: Map; + metadataSessions: Map; +} + +/** Request classification encoded in Codex turn metadata. */ +export type OpenAICodexRequestKind = "turn" | "prewarm" | "compaction"; + +interface CodexMetadataSessionState { + sessionId: string; + threadId: string; + windowId: string; + turnId?: string; + turnStartedAtUnixMs?: number; + compactionOperationId?: string; + reuseTurnForNextRequest?: boolean; +} + +interface CodexCompatibilityIdentity { + installationId: string; + sessionId: string; + threadId: string; + windowId: string; + turnMetadataJson?: string; +} + +interface CodexRequestMetadata extends CodexCompatibilityIdentity { + turnId: string; + turnMetadataJson: string; + clientMetadata: Record; +} + +const CODEX_RESERVED_METADATA_KEYS: Record = { + installation_id: true, + [OPENAI_HEADERS.INSTALLATION_ID]: true, + session_id: true, + thread_id: true, + turn_id: true, + window_id: true, + [OPENAI_HEADERS.WINDOW_ID]: true, + [OPENAI_HEADERS.TURN_METADATA]: true, + [OPENAI_HEADERS.PARENT_THREAD_ID]: true, + [OPENAI_HEADERS.SUBAGENT]: true, + request_kind: true, + compaction: true, + turn_started_at_unix_ms: true, + forked_from_thread_id: true, + parent_thread_id: true, + subagent_kind: true, + thread_source: true, + sandbox: true, + workspaces: true, +}; + +function createCodexMetadataSessionState(sessionId: string): CodexMetadataSessionState { + return { + sessionId, + threadId: crypto.randomUUID(), + windowId: crypto.randomUUID(), + }; +} + +function getOrCreateCodexMetadataSessionState( + sessionId: string, + providerState: CodexProviderSessionState | undefined, +): CodexMetadataSessionState { + if (!providerState) return createCodexMetadataSessionState(sessionId); + const existing = providerState.metadataSessions.get(sessionId); + if (existing) return existing; + const created = createCodexMetadataSessionState(sessionId); + providerState.metadataSessions.set(sessionId, created); + return created; +} + +function createCodexCompatibilityIdentity(session: CodexMetadataSessionState): CodexCompatibilityIdentity { + return { + installationId: getInstallId(), + sessionId: session.sessionId, + threadId: session.threadId, + windowId: session.windowId, + }; +} + +function resolveCodexStartNewTurn( + session: CodexMetadataSessionState, + requestKind: OpenAICodexRequestKind, + compaction: CodexCompactionRequestContext | undefined, + override: boolean | undefined, +): boolean { + if (requestKind !== "compaction") { + if (requestKind === "turn") { + const reuseCompactionTurn = session.reuseTurnForNextRequest === true; + session.reuseTurnForNextRequest = false; + session.compactionOperationId = undefined; + if (reuseCompactionTurn) return false; + } + return override ?? requestKind === "turn"; + } + if (!compaction) return override ?? false; + const startsNewOperation = session.compactionOperationId !== compaction.operationId; + if (startsNewOperation) session.reuseTurnForNextRequest = false; + session.compactionOperationId = compaction.operationId; + return override ?? (compaction.phase !== "mid_turn" && startsNewOperation); +} + +function toAsciiJsonString(value: Record): string { + return JSON.stringify(value).replace( + /[\x7f-\uffff]/g, + char => `\\u${char.charCodeAt(0).toString(16).padStart(4, "0")}`, + ); +} + +function createCodexRequestMetadata( + session: CodexMetadataSessionState, + requestKind: OpenAICodexRequestKind, + options: { + startNewTurn: boolean; + turnStartedAtUnixMs?: number; + clientMetadata?: Readonly>; + compaction?: CodexCompactionRequestContext; + }, +): CodexRequestMetadata { + if (options.startNewTurn || !session.turnId) { + session.turnId = crypto.randomUUID(); + session.turnStartedAtUnixMs = options.turnStartedAtUnixMs; + } + const identity = createCodexCompatibilityIdentity(session); + const extra: Record = {}; + const callerMetadata = options.clientMetadata; + if (callerMetadata) { + for (const key in callerMetadata) { + if (!CODEX_RESERVED_METADATA_KEYS[key]) extra[key] = callerMetadata[key]; + } + } + const turnMetadata: Record = { + installation_id: identity.installationId, + session_id: identity.sessionId, + thread_id: identity.threadId, + turn_id: session.turnId, + window_id: identity.windowId, + request_kind: requestKind, + }; + if (options.compaction) { + turnMetadata.compaction = { + trigger: options.compaction.trigger, + reason: options.compaction.reason, + implementation: options.compaction.implementation, + phase: options.compaction.phase, + strategy: options.compaction.strategy, + }; + } + if (session.turnStartedAtUnixMs !== undefined) { + turnMetadata.turn_started_at_unix_ms = session.turnStartedAtUnixMs; + } + for (const key in extra) turnMetadata[key] = extra[key]; + const turnMetadataJson = toAsciiJsonString(turnMetadata); + 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, + }, + }; +} + +function applyCodexCompatibilityHeaders(headers: Headers, metadata: CodexCompatibilityIdentity): void { + headers.set(OPENAI_HEADERS.SCOPED_SESSION_ID, metadata.sessionId); + headers.set(OPENAI_HEADERS.THREAD_ID, metadata.threadId); + headers.set(OPENAI_HEADERS.WINDOW_ID, metadata.windowId); + if (metadata.turnMetadataJson) { + headers.set(OPENAI_HEADERS.TURN_METADATA, metadata.turnMetadataJson); + } else { + headers.delete(OPENAI_HEADERS.TURN_METADATA); + } +} + +/** + * Synthesize Codex request identity for raw provider routes such as remote + * compaction while reusing the live session's thread, window, and turn. + */ +export function createOpenAICodexCompatibilityMetadata( + options: OpenAICodexCompatibilityMetadataOptions, +): OpenAICodexCompatibilityMetadata { + const providerState = getCodexProviderSessionState(options.providerSessionState); + const sessionId = normalizeOpenAIPromptCacheKey(options.sessionId) ?? crypto.randomUUID(); + const session = getOrCreateCodexMetadataSessionState(sessionId, providerState); + const startNewTurn = resolveCodexStartNewTurn( + session, + options.requestKind, + options.compaction, + options.startNewTurn, + ); + const metadata = createCodexRequestMetadata(session, options.requestKind, { + startNewTurn, + turnStartedAtUnixMs: options.turnStartedAtUnixMs ?? (startNewTurn || !session.turnId ? Date.now() : undefined), + clientMetadata: options.clientMetadata, + compaction: options.compaction, + }); + const headers = new Headers(); + applyCodexCompatibilityHeaders(headers, metadata); + if (options.includeInstallationHeader) { + headers.set(OPENAI_HEADERS.INSTALLATION_ID, metadata.installationId); + } + return { + clientMetadata: { ...metadata.clientMetadata }, + headers: Object.fromEntries(headers.entries()), + }; +} + +/** + * Invalidate Codex history-dependent transport state after compaction while + * retaining the session identity and live connection. + */ +export function resetOpenAICodexHistoryAfterCompaction(options: OpenAICodexCompactionResetOptions): void { + const providerState = options.providerSessionState?.get(CODEX_PROVIDER_SESSION_STATE_KEY); + if (!isCodexProviderSessionState(providerState)) return; + for (const websocketState of providerState.webSocketSessions.values()) { + resetCodexWebSocketAppendState(websocketState); + if (options.compaction.phase !== "mid_turn") websocketState.turnState = undefined; + } + const sessionId = normalizeOpenAIPromptCacheKey(options.sessionId); + if (!sessionId) return; + const metadataSession = providerState.metadataSessions.get(sessionId); + if (!metadataSession) return; + metadataSession.windowId = crypto.randomUUID(); + metadataSession.compactionOperationId = undefined; + metadataSession.reuseTurnForNextRequest = options.compaction.phase !== "standalone_turn"; } interface CodexRequestContext { @@ -355,8 +632,10 @@ interface CodexRequestContext { requestHeaders: Record; transportSessionId?: string; providerSessionState?: CodexProviderSessionState; + isolatedTransportState?: CodexProviderSessionState; websocketState?: CodexWebSocketSessionState; responsesLite: boolean; + requestMetadata?: CodexRequestMetadata; transformedBody: RequestBody; rawRequestDump: RawHttpRequestDump; } @@ -613,23 +892,37 @@ function createCodexProviderSessionState(): CodexProviderSessionState { const state: CodexProviderSessionState = { webSocketSessions: new Map(), webSocketPublicToPrivate: new Map(), + metadataSessions: new Map(), close: () => { for (const session of state.webSocketSessions.values()) { session.connection?.close("session_disposed"); } state.webSocketSessions.clear(); state.webSocketPublicToPrivate.clear(); + state.metadataSessions.clear(); }, }; return state; } +function isCodexProviderSessionState(state: ProviderSessionState | undefined): state is CodexProviderSessionState { + return ( + state !== undefined && + "webSocketSessions" in state && + state.webSocketSessions instanceof Map && + "webSocketPublicToPrivate" in state && + state.webSocketPublicToPrivate instanceof Map && + "metadataSessions" in state && + state.metadataSessions instanceof Map + ); +} + function getCodexProviderSessionState( providerSessionState: Map | undefined, ): CodexProviderSessionState | undefined { if (!providerSessionState) return undefined; - const existing = providerSessionState.get(CODEX_PROVIDER_SESSION_STATE_KEY) as CodexProviderSessionState | undefined; - if (existing) return existing; + const existing = providerSessionState.get(CODEX_PROVIDER_SESSION_STATE_KEY); + if (isCodexProviderSessionState(existing)) return existing; const created = createCodexProviderSessionState(); providerSessionState.set(CODEX_PROVIDER_SESSION_STATE_KEY, created); return created; @@ -915,19 +1208,56 @@ async function buildCodexRequestContext( }; const providerSessionState = getCodexProviderSessionState(options?.providerSessionState); + const isolatedTransportState = options?.codexCompaction ? createCodexProviderSessionState() : undefined; + const transportProviderSessionState = isolatedTransportState ?? providerSessionState; const responsesLite = resolveCodexResponsesLite(model, options?.responsesLite); const sessionKey = getCodexWebSocketSessionKey(transportSessionId, model, accountId, apiKey, baseUrl, responsesLite); const publicSessionKey = transportSessionId ? `${baseUrl}:${model.id}:${transportSessionId}` : undefined; if (sessionKey && publicSessionKey) { - providerSessionState?.webSocketPublicToPrivate.set(publicSessionKey, sessionKey); + transportProviderSessionState?.webSocketPublicToPrivate.set(publicSessionKey, sessionKey); } + const sharedWebsocketState = + sessionKey && providerSessionState + ? isolatedTransportState + ? providerSessionState.webSocketSessions.get(sessionKey) + : getCodexWebSocketSessionState(sessionKey, providerSessionState) + : undefined; const websocketState = - sessionKey && providerSessionState ? getCodexWebSocketSessionState(sessionKey, providerSessionState) : undefined; - if (websocketState && !isCodexWithinTurnContinuation(context)) { - // codex-rs scopes `x-codex-turn-state` to a single user turn: tool-loop - // follow-ups echo it, a new user turn starts without it. + sessionKey && isolatedTransportState + ? getCodexWebSocketSessionState(sessionKey, isolatedTransportState) + : sharedWebsocketState; + if (isolatedTransportState && websocketState && sharedWebsocketState) { + websocketState.disableWebsocket = sharedWebsocketState.disableWebsocket; + websocketState.turnState = sharedWebsocketState.turnState; + websocketState.modelsEtag = sharedWebsocketState.modelsEtag; + } + const withinTurnContinuation = isCodexWithinTurnContinuation(context); + const metadataSessionId = transportSessionId ?? crypto.randomUUID(); + const metadataSession = getOrCreateCodexMetadataSessionState(metadataSessionId, providerSessionState); + const compaction = options?.codexCompaction; + const requestKind: OpenAICodexRequestKind = compaction ? "compaction" : "turn"; + const startNewTurn = resolveCodexStartNewTurn( + metadataSession, + requestKind, + compaction, + compaction ? undefined : !withinTurnContinuation, + ); + if (websocketState && startNewTurn) { + // Codex scopes turn-state to one turn. Mid-turn compaction and tool-loop + // follow-ups preserve it; new user or compaction turns start without it. websocketState.turnState = undefined; } + const requestMetadata = createCodexRequestMetadata(metadataSession, requestKind, { + startNewTurn, + turnStartedAtUnixMs: compaction + ? startNewTurn || !metadataSession.turnId + ? Date.now() + : undefined + : getCodexTurnStartedAtUnixMs(context), + clientMetadata: transformedBody.client_metadata, + compaction, + }); + transformedBody.client_metadata = requestMetadata.clientMetadata; return { apiKey, accountId, @@ -936,8 +1266,10 @@ async function buildCodexRequestContext( requestHeaders, transportSessionId, providerSessionState, + isolatedTransportState, websocketState, responsesLite, + requestMetadata, transformedBody, rawRequestDump, }; @@ -1065,19 +1397,21 @@ async function openCodexWebSocketTransport( }> { const canAppendBeforeRequest = websocketState.canAppend === true; const chainedBody = buildCodexChainedRequestBody(requestContext.transformedBody, websocketState); - // WebSocket frames cannot carry per-request HTTP headers, so the Responses - // Lite marker rides in `client_metadata` on every `response.create`. + // WebSocket frames cannot carry per-request HTTP headers. Canonical Codex + // request identity is already in `client_metadata`; connection-scoped + // compatibility values that can change after the upgrade ride alongside it + // on every `response.create`. + const websocketClientMetadata = { ...(chainedBody.client_metadata ?? {}) }; + if (requestContext.responsesLite) { + websocketClientMetadata[CODEX_WS_RESPONSES_LITE_CLIENT_METADATA_KEY] = "true"; + } + if (websocketState.turnState) { + websocketClientMetadata[X_CODEX_TURN_STATE_HEADER] = websocketState.turnState; + } let websocketRequest = { type: "response.create", ...chainedBody, - ...(requestContext.responsesLite - ? { - client_metadata: { - ...(chainedBody.client_metadata ?? {}), - [CODEX_WS_RESPONSES_LITE_CLIENT_METADATA_KEY]: "true", - }, - } - : {}), + client_metadata: websocketClientMetadata, }; const replacementWebsocketRequest = await options?.onPayload?.(websocketRequest, model); if (replacementWebsocketRequest !== undefined) { @@ -1092,6 +1426,7 @@ async function openCodexWebSocketTransport( "websocket", websocketState, requestContext.responsesLite, + requestContext.requestMetadata, ); const requestBodyForState = structuredCloneJSON(requestContext.transformedBody); // `onPayload` may rewrite the outgoing frame (e.g. drop `stream_options`); @@ -1137,6 +1472,16 @@ async function openCodexWebSocketTransport( }; } +function getCodexTurnStartedAtUnixMs(context: Context): number { + for (let i = context.messages.length - 1; i >= 0; i--) { + const message = context.messages[i]; + if (message?.role === "user" && Number.isFinite(message.timestamp)) { + return Math.trunc(message.timestamp); + } + } + return Date.now(); +} + /** * True when the request continues the current turn (everything after the * last assistant message is tool results), false when a new user turn starts. @@ -1177,6 +1522,7 @@ async function openCodexSseTransport( wireBody, state, requestContext.responsesLite, + requestContext.requestMetadata, requestSetup.requestSignal, requestSetup.firstEventTimeoutMs, event => options?.onSseEvent?.(event, model), @@ -1591,7 +1937,9 @@ class CodexStreamProcessor { const contentIndex = entry?.contentIndex ?? output.content.length - 1; if (item.type === "reasoning" && block?.type === "thinking") { - block.thinking = finalizeReasoningThinking(item, block.thinking); + block.thinking = finalizeReasoningThinking(item, block.thinking, { + cumulativeSummarySnapshots: this.#sequentialCutoffSummaries, + }); block.thinkingSignature = JSON.stringify(item); stream.push({ type: "thinking_end", @@ -2168,6 +2516,8 @@ export const streamOpenAICodexResponses: StreamFunction<"openai-codex-responses" stream.push({ type: "error", reason: "error", error: output }); } stream.end(); + } finally { + requestContext?.isolatedTransportState?.close(); } })(); @@ -2198,6 +2548,11 @@ export async function prewarmOpenAICodexResponses( if (!sessionKey || !providerSessionState) return; const state = getCodexWebSocketSessionState(sessionKey, providerSessionState); if (!shouldUseCodexWebSocket(model, state, options?.preferWebsockets)) return; + const metadataSession = getOrCreateCodexMetadataSessionState( + transportSessionId ?? crypto.randomUUID(), + providerSessionState, + ); + const requestIdentity = createCodexCompatibilityIdentity(metadataSession); const headers = logger.time( "prewarmCodex:createHeaders", createCodexHeaders, @@ -2208,6 +2563,7 @@ export async function prewarmOpenAICodexResponses( "websocket", state, responsesLite, + requestIdentity, ); await logger.time( "prewarmCodex:establishWs", @@ -2315,6 +2671,7 @@ export interface OpenAICodexTransportDetails { canAppend: boolean; prewarmed: boolean; hasSessionState: boolean; + hasTurnState: boolean; lastFallbackAt?: number; } @@ -2377,6 +2734,7 @@ export function getOpenAICodexTransportDetails( canAppend: state?.canAppend ?? false, prewarmed: state?.prewarmed ?? false, hasSessionState: state !== undefined, + hasTurnState: state?.turnState !== undefined, lastFallbackAt: state?.lastFallbackAt, }; } @@ -3320,12 +3678,22 @@ async function openCodexSseEventStream( body: RequestBody, state: CodexWebSocketSessionState | undefined, responsesLite: boolean, + requestMetadata: CodexRequestMetadata | undefined, signal: AbortSignal | undefined, firstEventTimeoutMs: number | undefined, onSseEvent?: OpenAICodexResponsesOptions["onSseEvent"], fetchOverride?: FetchImpl, ): Promise>> { - const headers = createCodexHeaders(requestHeaders, accountId, apiKey, sessionId, "sse", state, responsesLite); + const headers = createCodexHeaders( + requestHeaders, + accountId, + apiKey, + sessionId, + "sse", + state, + responsesLite, + requestMetadata, + ); CODEX_DEBUG && logger.debug("[codex] codex request", { url, @@ -3385,6 +3753,7 @@ function createCodexHeaders( transport: CodexTransport = "sse", state?: CodexWebSocketSessionState, responsesLite = false, + requestMetadata?: CodexCompatibilityIdentity, ): Headers { const headers = new Headers(initHeaders ?? {}); headers.delete("x-api-key"); @@ -3408,6 +3777,15 @@ function createCodexHeaders( headers.delete(OPENAI_HEADERS.SESSION_ID); headers.delete("x-client-request-id"); } + headers.delete(OPENAI_HEADERS.INSTALLATION_ID); + if (requestMetadata) { + applyCodexCompatibilityHeaders(headers, requestMetadata); + } else { + headers.delete(OPENAI_HEADERS.SCOPED_SESSION_ID); + headers.delete(OPENAI_HEADERS.THREAD_ID); + headers.delete(OPENAI_HEADERS.WINDOW_ID); + headers.delete(OPENAI_HEADERS.TURN_METADATA); + } if (state?.turnState) { headers.set(X_CODEX_TURN_STATE_HEADER, state.turnState); } else { @@ -3445,6 +3823,10 @@ function redactHeaders(headers: Headers): Record { lower.includes("account") || lower.includes("session") || lower.includes("conversation") || + lower.includes("thread") || + lower.includes("window") || + lower.includes("installation") || + lower.startsWith("x-codex-turn") || lower === "x-client-request-id" || lower === "cookie" ) { diff --git a/packages/ai/src/providers/openai-shared.ts b/packages/ai/src/providers/openai-shared.ts index 1820c24f5..789f31213 100644 --- a/packages/ai/src/providers/openai-shared.ts +++ b/packages/ai/src/providers/openai-shared.ts @@ -1702,9 +1702,36 @@ export function appendReasoningSummaryPart( item.summary.push(part); } -/** Chooses the final reasoning text without discarding content already streamed into the block. */ -export function finalizeReasoningThinking(item: ResponseReasoningItem, streamedThinking: string): string { - const summaryThinking = item.summary?.map(part => part.text).join("\n\n") ?? ""; +// Sequential-cutoff streams may repeat the full canonical summary as later parts. +function foldReasoningSummary(parts: ResponseReasoningItem["summary"] | undefined): string { + if (!parts) return ""; + let canonical = ""; + for (const part of parts) { + const text = part.text; + if (!text || text === canonical) continue; + const extendsCanonical = text.startsWith(canonical) && text[canonical.length] === "\n"; + canonical = !canonical || extendsCanonical ? text : `${canonical}\n\n${text}`; + } + return canonical; +} + +/** Chooses final reasoning text without making sequential-cutoff results disagree with emitted deltas. */ +export function finalizeReasoningThinking( + item: ResponseReasoningItem, + streamedThinking: string, + options: { cumulativeSummarySnapshots?: boolean } = {}, +): string { + const summaryThinking = options.cumulativeSummarySnapshots + ? foldReasoningSummary(item.summary) + : (item.summary?.map(part => part.text).join("\n\n") ?? ""); + if ( + options.cumulativeSummarySnapshots && + streamedThinking && + summaryThinking && + summaryThinking !== streamedThinking + ) { + return streamedThinking; + } if (summaryThinking) return summaryThinking; const contentThinking = item.content?.[0]?.type === "reasoning_text" ? (item.content[0].text ?? "") : ""; return contentThinking || streamedThinking || ""; @@ -1742,12 +1769,12 @@ export function appendReasoningSummaryPartDone( } /** - * Applies an atomic `response.reasoning_summary_text.done` event (concurrent - * reasoning summaries, `stream_options.reasoning_summary_delivery: - * "sequential_cutoff"`). The event carries the FULL text for `summaryIndex`; - * incremental `.delta`/`.part.*` events are ignored under this contract, so - * the whole part is stored and streamed here. Parts after the first are - * separated by a section break, mirroring codex-rs. + * Applies an atomic `response.reasoning_summary_text.done` snapshot. + * + * Sequential-cutoff streams can replay an index or send the accumulated + * summary as a later part. Rebuild the canonical summary and emit only its + * append-only suffix. Divergent corrections stay buffered until finalization + * so delta consumers never receive suffixes based on unseen replacement text. */ export function applyReasoningSummaryDone( item: ResponseReasoningItem, @@ -1763,8 +1790,11 @@ export function applyReasoningSummaryDone( item.summary.push({ type: "summary_text", text: "" }); } item.summary[summaryIndex].text = text; - const delta = summaryIndex > 0 ? `\n\n${text}` : text; - block.thinking += delta; + const after = foldReasoningSummary(item.summary); + if (!after.startsWith(block.thinking)) return; + const delta = after.slice(block.thinking.length); + if (!delta) return; + block.thinking = after; stream.push({ type: "thinking_delta", contentIndex, delta, partial: output }); } diff --git a/packages/ai/src/stream.ts b/packages/ai/src/stream.ts index 02d17ad2f..de922d926 100644 --- a/packages/ai/src/stream.ts +++ b/packages/ai/src/stream.ts @@ -1632,6 +1632,7 @@ function mapOptionsForApi( toolChoice: mapOpenAiToolChoice(options?.toolChoice), serviceTier: options?.serviceTier, preferWebsockets: options?.preferWebsockets, + codexCompaction: options?.codexCompaction, reasoningSummary: options?.hideThinkingSummary ? null : "detailed", textVerbosity: options?.textVerbosity, }); diff --git a/packages/ai/src/types.ts b/packages/ai/src/types.ts index e94884553..5b568de29 100644 --- a/packages/ai/src/types.ts +++ b/packages/ai/src/types.ts @@ -323,6 +323,30 @@ export interface RawSseEvent { raw: string[]; } +/** Lifecycle fields shared by every Codex compaction implementation. */ +export interface CodexCompactionContext { + /** Stable only for one logical compaction, including parallel summary calls. */ + operationId: string; + trigger: "manual" | "auto"; + reason: "user_requested" | "context_limit" | "model_downshift" | "comp_hash_changed"; + phase: "standalone_turn" | "pre_turn" | "mid_turn"; + strategy: "memento" | "prefix_compaction"; +} + +/** Canonical nested metadata serialized into the Codex turn envelope. */ +export interface CodexCompactionMetadata { + trigger: "manual" | "auto"; + reason: "user_requested" | "context_limit" | "model_downshift" | "comp_hash_changed"; + implementation: "responses" | "responses_compaction_v2" | "responses_compact"; + phase: "standalone_turn" | "pre_turn" | "mid_turn"; + strategy: "memento" | "prefix_compaction"; +} + +/** Dispatch context combining canonical metadata with its local operation identity. */ +export interface CodexCompactionRequestContext extends CodexCompactionMetadata { + operationId: string; +} + export interface StreamOptions { temperature?: number; topP?: number; @@ -398,6 +422,8 @@ export interface StreamOptions { * Providers can use this to persist transport/session state between turns. */ providerSessionState?: Map; + /** Canonical Codex compaction classification; ignored by other providers. */ + codexCompaction?: CodexCompactionRequestContext; /** * Force Gemini model-mode Interactions API transport for providers that support it. * When unset, those providers may still use Interactions to continue known diff --git a/packages/ai/test/issue-1701-repro.test.ts b/packages/ai/test/issue-1701-repro.test.ts index 4016d573d..79318532c 100644 --- a/packages/ai/test/issue-1701-repro.test.ts +++ b/packages/ai/test/issue-1701-repro.test.ts @@ -1,12 +1,23 @@ -import { describe, expect, it } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import { streamAzureOpenAIResponses } from "@oh-my-pi/pi-ai/providers/azure-openai-responses"; import { streamOpenAICodexResponses } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { streamOpenAICompletions } from "@oh-my-pi/pi-ai/providers/openai-completions"; import { streamOpenAIResponses } from "@oh-my-pi/pi-ai/providers/openai-responses"; import type { Context, Model, Tool, ToolChoice } from "@oh-my-pi/pi-ai/types"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; +import * as piUtils from "@oh-my-pi/pi-utils"; import { z } from "zod/v4"; +const TEST_INSTALLATION_ID = "00000000-0000-4000-8000-000000000001"; + +beforeEach(() => { + vi.spyOn(piUtils, "getInstallId").mockReturnValue(TEST_INSTALLATION_ID); +}); + +afterEach(() => { + vi.restoreAllMocks(); +}); + const completionsModel: Model<"openai-completions"> = buildModel({ id: "gpt-4o-mini-test", name: "GPT-4o Mini Test", diff --git a/packages/ai/test/openai-codex-responses-lite.test.ts b/packages/ai/test/openai-codex-responses-lite.test.ts index f70e57d43..169f578a9 100644 --- a/packages/ai/test/openai-codex-responses-lite.test.ts +++ b/packages/ai/test/openai-codex-responses-lite.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import { type InputItem, type RequestBody, @@ -7,13 +7,25 @@ import { import { buildTransformedCodexRequestBody, convertCodexResponsesMessages, + resetOpenAICodexHistoryAfterCompaction, streamOpenAICodexResponses, } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { isOpenAIResponsesProgressEvent } from "@oh-my-pi/pi-ai/providers/openai-shared"; -import type { Context, FetchImpl } from "@oh-my-pi/pi-ai/types"; +import type { CodexCompactionRequestContext, Context, FetchImpl, ProviderSessionState } from "@oh-my-pi/pi-ai/types"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; +import * as piUtils from "@oh-my-pi/pi-utils"; import { createCodexModel } from "./helpers"; +const TEST_INSTALLATION_ID = "00000000-0000-4000-8000-000000000001"; + +beforeEach(() => { + vi.spyOn(piUtils, "getInstallId").mockReturnValue(TEST_INSTALLATION_ID); +}); + +afterEach(() => { + vi.restoreAllMocks(); +}); + function createCodexTestToken(accountId = "acc_test"): string { const payload = Buffer.from( JSON.stringify({ "https://api.openai.com/auth": { chatgpt_account_id: accountId } }), @@ -64,6 +76,24 @@ interface CapturedCodexRequest { body: Record; } +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +function requireRecord(value: unknown, label: string): Record { + if (!isRecord(value)) { + throw new Error(`expected ${label} to be an object`); + } + return value; +} + +function parseTurnMetadata(clientMetadata: Record): Record { + const encoded = clientMetadata["x-codex-turn-metadata"]; + if (typeof encoded !== "string") throw new Error("expected x-codex-turn-metadata"); + const decoded: unknown = JSON.parse(encoded); + return requireRecord(decoded, "x-codex-turn-metadata"); +} + function createCodexFetchMock(sse: string, onRequest: (captured: CapturedCodexRequest) => void): FetchImpl { return (async (input: string | URL, init?: RequestInit) => { const url = typeof input === "string" ? input : input.toString(); @@ -370,15 +400,21 @@ describe("openai-codex fresh execution input shaping", () => { }); describe("openai-codex Responses Lite and client metadata wire format", () => { - it("sends the lite header and client_metadata body field over SSE", async () => { + it("sends canonical Codex metadata and protects reserved fields over SSE", async () => { const model = createCodexModel("gpt-5.1-codex"); - const clientMetadata = { "x-codex-turn-metadata": '{"thread_id":"thread_1","turn_id":"turn_1"}' }; + const context = createCodexTestContext(); + const clientMetadata = { + workspace_kind: "repo", + workspace_path: "東京/πŸš€", + session_id: "caller-session", + "x-codex-turn-metadata": '{"turn_id":"caller-turn"}', + }; let captured: CapturedCodexRequest | undefined; const fetchMock = createCodexFetchMock(createCodexSse(COMPLETED_CODEX_EVENTS), request => { captured = request; }); - const result = await streamOpenAICodexResponses(model, createCodexTestContext(), { + const result = await streamOpenAICodexResponses(model, context, { apiKey: createCodexTestToken(), fetch: fetchMock, responsesLite: true, @@ -386,8 +422,141 @@ describe("openai-codex Responses Lite and client metadata wire format", () => { }).result(); expect(result.stopReason).toBe("stop"); - expect(captured?.headers.get("x-openai-internal-codex-responses-lite")).toBe("true"); - expect(captured?.body.client_metadata).toEqual(clientMetadata); + if (!captured) throw new Error("expected a captured Codex request"); + expect(captured.headers.get("x-openai-internal-codex-responses-lite")).toBe("true"); + expect(captured.headers.get("x-codex-installation-id")).toBeNull(); + + const metadata = requireRecord(captured.body.client_metadata, "client_metadata"); + const turnMetadata = parseTurnMetadata(metadata); + expect(metadata.workspace_kind).toBeUndefined(); + expect(metadata.workspace_path).toBeUndefined(); + expect(metadata.session_id).not.toBe("caller-session"); + expect(turnMetadata.request_kind).toBe("turn"); + expect(turnMetadata.turn_started_at_unix_ms).toBe(context.messages[0]?.timestamp); + expect(turnMetadata.workspace_kind).toBe("repo"); + expect(turnMetadata.workspace_path).toBe("東京/πŸš€"); + expect(metadata["x-codex-installation-id"]).toMatch( + /^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i, + ); + expect(metadata.session_id).toBe(turnMetadata.session_id); + expect(metadata.thread_id).toBe(turnMetadata.thread_id); + expect(metadata.turn_id).toBe(turnMetadata.turn_id); + expect(metadata["x-codex-window-id"]).toBe(turnMetadata.window_id); + expect(metadata.session_id).toBe(captured.headers.get("session-id")); + expect(metadata.thread_id).toBe(captured.headers.get("thread-id")); + expect(metadata["x-codex-window-id"]).toBe(captured.headers.get("x-codex-window-id")); + expect(metadata["x-codex-turn-metadata"]).toBe(captured.headers.get("x-codex-turn-metadata")); + const turnMetadataHeader = captured.headers.get("x-codex-turn-metadata"); + expect(turnMetadataHeader).toMatch(/^[\x20-\x7e]+$/); + const reparsedTurnMetadata: unknown = turnMetadataHeader ? JSON.parse(turnMetadataHeader) : undefined; + expect(requireRecord(reparsedTurnMetadata, "round-tripped turn metadata").workspace_path).toBe("東京/πŸš€"); + }); + + it("keeps the installation identity stable across provider sessions", async () => { + const model = createCodexModel("gpt-5.1-codex"); + const captured: CapturedCodexRequest[] = []; + const fetchMock = createCodexFetchMock(createCodexSse(COMPLETED_CODEX_EVENTS), request => { + captured.push(request); + }); + + await streamOpenAICodexResponses(model, createCodexTestContext(), { + apiKey: createCodexTestToken(), + fetch: fetchMock, + sessionId: "metadata-session-one", + }).result(); + await streamOpenAICodexResponses(model, createCodexTestContext(), { + apiKey: createCodexTestToken(), + fetch: fetchMock, + sessionId: "metadata-session-two", + }).result(); + + const firstMetadata = requireRecord(captured[0]?.body.client_metadata, "first client_metadata"); + const secondMetadata = requireRecord(captured[1]?.body.client_metadata, "second client_metadata"); + expect(firstMetadata["x-codex-installation-id"]).toBe(secondMetadata["x-codex-installation-id"]); + expect(firstMetadata.session_id).toBe("metadata-session-one"); + expect(secondMetadata.session_id).toBe("metadata-session-two"); + expect(firstMetadata.thread_id).not.toBe(secondMetadata.thread_id); + }); + + it("rotates compaction turns by phase and reuses one operation across fan-out calls", async () => { + const model = createCodexModel("gpt-5.1-codex"); + const providerSessionState = new Map(); + const captured: CapturedCodexRequest[] = []; + const fetchMock = createCodexFetchMock(createCodexSse(COMPLETED_CODEX_EVENTS), request => { + captured.push(request); + }); + const send = async (codexCompaction?: CodexCompactionRequestContext): Promise => { + await streamOpenAICodexResponses(model, createCodexTestContext(), { + apiKey: createCodexTestToken(), + fetch: fetchMock, + sessionId: "compaction-lifecycle-session", + providerSessionState, + codexCompaction, + }).result(); + }; + const preTurn: CodexCompactionRequestContext = { + operationId: "pre-turn-operation", + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "pre_turn", + strategy: "memento", + }; + const midTurn: CodexCompactionRequestContext = { + ...preTurn, + operationId: "mid-turn-operation", + phase: "mid_turn", + }; + const standalone: CodexCompactionRequestContext = { + ...preTurn, + operationId: "standalone-operation", + trigger: "manual", + reason: "user_requested", + phase: "standalone_turn", + }; + + await send(); + await send(preTurn); + await send(preTurn); + resetOpenAICodexHistoryAfterCompaction({ + providerSessionState, + sessionId: "compaction-lifecycle-session", + compaction: preTurn, + }); + await send(); + await send(midTurn); + await send(standalone); + + const turns = captured.map((request, index) => + parseTurnMetadata(requireRecord(request.body.client_metadata, `client_metadata ${index}`)), + ); + expect(turns[0]?.request_kind).toBe("turn"); + expect(turns[1]?.turn_id).not.toBe(turns[0]?.turn_id); + expect(turns[2]?.turn_id).toBe(turns[1]?.turn_id); + expect(turns[2]?.turn_started_at_unix_ms).toBe(turns[1]?.turn_started_at_unix_ms); + expect(turns[3]?.request_kind).toBe("turn"); + expect(turns[3]?.turn_id).toBe(turns[1]?.turn_id); + expect(turns[3]?.window_id).not.toBe(turns[2]?.window_id); + expect(turns[4]?.turn_id).toBe(turns[1]?.turn_id); + expect(turns[5]?.turn_id).not.toBe(turns[4]?.turn_id); + expect(turns[1]?.thread_id).toBe(turns[5]?.thread_id); + expect(turns[1]?.compaction).toEqual({ + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "pre_turn", + strategy: "memento", + }); + const nestedCompaction = requireRecord(turns[1]?.compaction, "nested compaction metadata"); + expect(nestedCompaction.operationId).toBeUndefined(); + expect(nestedCompaction.operation_id).toBeUndefined(); + expect(turns[5]?.compaction).toEqual({ + trigger: "manual", + reason: "user_requested", + implementation: "responses", + phase: "standalone_turn", + strategy: "memento", + }); }); it("keeps lite and strips image detail when a lite request contains images", async () => { const model = buildModel({ @@ -461,7 +630,7 @@ describe("openai-codex Responses Lite and client metadata wire format", () => { expect((captured?.body.input as Array>)[0]?.type).toBe("additional_tools"); }); - it("omits the lite header and client_metadata when not requested", async () => { + it("omits the lite marker while retaining canonical client_metadata", async () => { const model = createCodexModel("gpt-5.1-codex"); let captured: CapturedCodexRequest | undefined; const fetchMock = createCodexFetchMock(createCodexSse(COMPLETED_CODEX_EVENTS), request => { @@ -475,7 +644,7 @@ describe("openai-codex Responses Lite and client metadata wire format", () => { expect(result.stopReason).toBe("stop"); expect(captured?.headers.get("x-openai-internal-codex-responses-lite")).toBeNull(); - expect(captured?.body.client_metadata).toBeUndefined(); + expect(captured?.body.client_metadata).toBeDefined(); }); }); @@ -559,7 +728,7 @@ describe("openai-codex concurrent reasoning summaries", () => { expect(unsupported.stream_options).toBeUndefined(); }); - it("decodes atomic summary dones and ignores legacy deltas under sequential cutoff", async () => { + it("deduplicates cumulative atomic summaries and ignores legacy deltas under sequential cutoff", async () => { const model = createCodexModel("gpt-5.6-terra"); const events: Array> = [ { @@ -586,14 +755,70 @@ describe("openai-codex concurrent reasoning summaries", () => { item_id: "reason_1", output_index: 0, summary_index: 0, - text: "First part", + text: "Plan", }, { type: "response.reasoning_summary_text.done", item_id: "reason_1", output_index: 0, summary_index: 1, - text: "Second part", + text: "Planning details", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 1, + text: "Planning details", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 2, + text: "Plan\n\nPlanning details\n\nInspect", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 2, + text: "Plan\n\nPlanning details\n\nInspect details", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 2, + text: "Plan\n\nPlanning details\n\nInspect details", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 3, + text: "Plan\n\nPlanning details\n\nInspect details", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 2, + text: "Plan\n\nPlanning details\n\nReview", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 2, + text: "Plan\n\nPlanning details\n\nReview output", + }, + { + type: "response.reasoning_summary_text.done", + item_id: "reason_1", + output_index: 0, + summary_index: 3, + text: "Plan\n\nPlanning details\n\nReview output", }, { type: "response.output_item.done", @@ -602,8 +827,10 @@ describe("openai-codex concurrent reasoning summaries", () => { type: "reasoning", id: "reason_1", summary: [ - { type: "summary_text", text: "First part" }, - { type: "summary_text", text: "Second part" }, + { type: "summary_text", text: "Plan" }, + { type: "summary_text", text: "Planning details" }, + { type: "summary_text", text: "Plan\n\nPlanning details\n\nInspect details\n\nUnseen final" }, + { type: "summary_text", text: "Plan\n\nPlanning details\n\nInspect details\n\nUnseen final" }, ], }, }, @@ -618,7 +845,7 @@ describe("openai-codex concurrent reasoning summaries", () => { type: "response.reasoning_summary_text.done", item_id: "reason_1", output_index: 0, - summary_index: 2, + summary_index: 4, text: "STALE", }, { @@ -662,10 +889,11 @@ describe("openai-codex concurrent reasoning summaries", () => { const result = await stream.result(); expect(captured?.body.stream_options).toEqual({ reasoning_summary_delivery: "sequential_cutoff" }); - expect(thinkingDeltas).toEqual(["First part", "\n\nSecond part"]); + expect(thinkingDeltas).toEqual(["Plan", "\n\nPlanning details", "\n\nInspect", " details"]); expect(result.stopReason).toBe("stop"); const thinking = result.content.find(block => block.type === "thinking"); - expect(thinking?.thinking).toBe("First part\n\nSecond part"); + expect(thinking?.thinking).toBe("Plan\n\nPlanning details\n\nInspect details"); + expect(thinking?.thinking).toBe(thinkingDeltas.join("")); const text = result.content.find(block => block.type === "text"); expect(text?.text).toBe("Hello"); }); diff --git a/packages/ai/test/openai-codex-stream.test.ts b/packages/ai/test/openai-codex-stream.test.ts index 505039c5c..32b0b52ec 100644 --- a/packages/ai/test/openai-codex-stream.test.ts +++ b/packages/ai/test/openai-codex-stream.test.ts @@ -1,18 +1,29 @@ -import { afterEach, describe, expect, it, vi } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import { streamSimple } from "@oh-my-pi/pi-ai"; import { getOpenAICodexTransportDetails, getOpenAICodexWebSocketDebugStats, prewarmOpenAICodexResponses, + resetOpenAICodexHistoryAfterCompaction, streamOpenAICodexResponses, } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; -import type { Context, FetchImpl, Model, ModelSpec, ProviderSessionState } from "@oh-my-pi/pi-ai/types"; +import type { + CodexCompactionRequestContext, + Context, + FetchImpl, + Model, + ModelSpec, + ProviderSessionState, +} from "@oh-my-pi/pi-ai/types"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; -import { getAgentDir, setAgentDir, TempDir } from "@oh-my-pi/pi-utils"; +import * as piUtils from "@oh-my-pi/pi-utils"; + +const { getAgentDir, setAgentDir, TempDir } = piUtils; const originalAgentDir = getAgentDir(); const originalWebSocket = global.WebSocket; const originalCodexWebSocketV2 = Bun.env.PI_CODEX_WEBSOCKET_V2; +const TEST_INSTALLATION_ID = "00000000-0000-4000-8000-000000000001"; function restoreEnv(name: string, value: string | undefined): void { if (value === undefined) { @@ -22,6 +33,10 @@ function restoreEnv(name: string, value: string | undefined): void { Bun.env[name] = value; } +beforeEach(() => { + vi.spyOn(piUtils, "getInstallId").mockReturnValue(TEST_INSTALLATION_ID); +}); + afterEach(() => { global.WebSocket = originalWebSocket; setAgentDir(originalAgentDir); @@ -60,6 +75,22 @@ function createCodexTestContext(): Context { }; } +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} + +function requireRecord(value: unknown, label: string): Record { + if (!isRecord(value)) throw new Error(`expected ${label} to be an object`); + return value; +} + +function parseTurnMetadata(clientMetadata: Record): Record { + const encoded = clientMetadata["x-codex-turn-metadata"]; + if (typeof encoded !== "string") throw new Error("expected x-codex-turn-metadata"); + const decoded: unknown = JSON.parse(encoded); + return requireRecord(decoded, "x-codex-turn-metadata"); +} + function createCompletedCodexSse(text: string): string { return `${[ `data: ${JSON.stringify({ type: "response.content_part.added", part: { type: "output_text", text: "" } })}`, @@ -1222,7 +1253,7 @@ describe("openai-codex streaming", () => { sessionId: "ws-lite-session", providerSessionState: new Map(), responsesLite: true, - clientMetadata: { "x-codex-turn-metadata": '{"thread_id":"t_1"}' }, + clientMetadata: { workspace_kind: "repo", "x-codex-turn-metadata": '{"thread_id":"caller"}' }, }, ).result(); @@ -1230,10 +1261,28 @@ describe("openai-codex streaming", () => { expect(capturedHeaders?.["x-openai-internal-codex-responses-lite"]).toBe("true"); expect(sentRequests).toHaveLength(1); expect(sentRequests[0]?.type).toBe("response.create"); - expect(sentRequests[0]?.client_metadata).toEqual({ - "x-codex-turn-metadata": '{"thread_id":"t_1"}', + const metadata = requireRecord(sentRequests[0]?.client_metadata, "client_metadata"); + const turnMetadata = parseTurnMetadata(metadata); + expect(metadata).toMatchObject({ + session_id: "ws-lite-session", ws_request_header_x_openai_internal_codex_responses_lite: "true", + "x-codex-installation-id": TEST_INSTALLATION_ID, }); + expect(metadata.workspace_kind).toBeUndefined(); + expect(turnMetadata).toMatchObject({ + installation_id: TEST_INSTALLATION_ID, + session_id: "ws-lite-session", + thread_id: metadata.thread_id, + turn_id: metadata.turn_id, + window_id: metadata["x-codex-window-id"], + request_kind: "turn", + workspace_kind: "repo", + }); + expect(capturedHeaders?.["x-codex-installation-id"]).toBeUndefined(); + expect(metadata.session_id).toBe(capturedHeaders?.["session-id"]); + expect(metadata.thread_id).toBe(capturedHeaders?.["thread-id"]); + expect(metadata["x-codex-window-id"]).toBe(capturedHeaders?.["x-codex-window-id"]); + expect(metadata["x-codex-turn-metadata"]).toBe(capturedHeaders?.["x-codex-turn-metadata"]); }); it("streams SSE responses into AssistantMessageEventStream", async () => { @@ -2110,7 +2159,7 @@ describe("openai-codex streaming", () => { expect(fallbackDetails.fallbackCount).toBe(1); }); - it("immediately falls back to SSE on fatal websocket connection errors", async () => { + it("carries fatal websocket fallback into isolated compaction transport", async () => { const tempDir = TempDir.createSync("@pi-codex-stream-"); setAgentDir(tempDir.path()); @@ -2173,6 +2222,23 @@ describe("openai-codex streaming", () => { expect(result.role).toBe("assistant"); expect(constructorCount).toBe(1); expect(fetchMock).toHaveBeenCalledTimes(1); + const compacted = await streamOpenAICodexResponses(model, context, { + fetch: fetchMock as FetchImpl, + apiKey: token, + sessionId: "ws-fatal-fallback-session", + providerSessionState, + codexCompaction: { + operationId: "fallback-compaction", + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "pre_turn", + strategy: "memento", + }, + }).result(); + expect(compacted.stopReason).toBe("stop"); + expect(constructorCount).toBe(1); + expect(fetchMock).toHaveBeenCalledTimes(2); const transportDetails = getOpenAICodexTransportDetails(model, { sessionId: "ws-fatal-fallback-session", providerSessionState, @@ -2182,7 +2248,7 @@ describe("openai-codex streaming", () => { expect(transportDetails.fallbackCount).toBe(1); }); - it("captures websocket handshake metadata and replays it on later SSE requests", async () => { + it("isolates compaction transport and preserves main mid-turn state", async () => { const tempDir = TempDir.createSync("@pi-codex-stream-"); setAgentDir(tempDir.path()); @@ -2199,12 +2265,21 @@ describe("openai-codex streaming", () => { `data: ${JSON.stringify({ type: "response.output_item.done", item: { type: "message", id: "msg_sse", role: "assistant", status: "completed", content: [{ type: "output_text", text: "Hello SSE" }] } })}`, `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`; + let firstRequest: Record | undefined; + let continuationRequest: Record | undefined; + let continuationHeaders: Headers | undefined; const fetchMock = vi.fn(async (_input: string | URL, init?: RequestInit) => { - const headers = init?.headers instanceof Headers ? init.headers : new Headers(init?.headers); - expect(headers.get("x-codex-turn-state")).toBe("ws-turn-state-1"); - expect(headers.get("x-models-etag")).toBe("models-etag-1"); + continuationHeaders = init?.headers instanceof Headers ? init.headers : new Headers(init?.headers); + expect(continuationHeaders.get("x-codex-turn-state")).toBe("ws-turn-state-1"); + expect(continuationHeaders.get("x-models-etag")).toBe("models-etag-1"); + if (typeof init?.body !== "string") throw new Error("expected an SSE request body"); + const body: unknown = JSON.parse(init.body); + continuationRequest = requireRecord(body, "SSE continuation request"); return new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" } }); }); + let websocketRequestCount = 0; + let websocketConstructorCount = 0; + const websocketInstances: MockWebSocket[] = []; class HandshakeWebSocket extends MockWebSocket { handshakeHeaders = { @@ -2215,11 +2290,29 @@ describe("openai-codex streaming", () => { constructor(url: string, options?: { headers?: WsHeaders }) { super(url, options); + websocketConstructorCount += 1; + websocketInstances.push(this); this.scheduleOpen(); } - send(): void { - this.emitCodexResponse({ messageId: "msg_ws", responseId: "resp_ws", text: "Hello WS" }); + send(data: string): void { + websocketRequestCount += 1; + const body: unknown = JSON.parse(data); + if (websocketRequestCount === 1) { + firstRequest = requireRecord(body, "websocket request"); + } + if (websocketRequestCount === 3) { + this.sendJson({ + type: "response.failed", + response: { error: { code: "invalid_request_error", message: "isolated compaction failed" } }, + }); + return; + } + this.emitCodexResponse({ + messageId: `msg_ws_${websocketRequestCount}`, + responseId: `resp_ws_${websocketRequestCount}`, + text: "Hello WS", + }); } } @@ -2248,12 +2341,65 @@ describe("openai-codex streaming", () => { messages: [{ role: "user", content: "Say hello", timestamp: Date.now() }], }; const providerSessionState = new Map(); + const midTurnCompaction: CodexCompactionRequestContext = { + operationId: "isolated-success", + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "mid_turn", + strategy: "memento", + }; const first = await streamOpenAICodexResponses(websocketModel, context, { fetch: fetchMock as FetchImpl, apiKey: token, sessionId: "ws-handshake-session", providerSessionState, }).result(); + expect(websocketInstances[0]?.readyState).toBe(MockWebSocket.OPEN); + const isolatedSuccess = await streamOpenAICodexResponses(websocketModel, createCodexTestContext(), { + fetch: fetchMock as FetchImpl, + apiKey: token, + sessionId: "ws-handshake-session", + providerSessionState, + codexCompaction: midTurnCompaction, + }).result(); + expect(isolatedSuccess.stopReason).toBe("stop"); + expect(websocketInstances[0]?.readyState).toBe(MockWebSocket.OPEN); + expect(websocketInstances[1]?.readyState).toBe(MockWebSocket.CLOSED); + expect(websocketInstances[1]?.options?.headers?.["x-codex-turn-state"]).toBe("ws-turn-state-1"); + expect(websocketInstances[1]?.options?.headers?.["x-models-etag"]).toBe("models-etag-1"); + resetOpenAICodexHistoryAfterCompaction({ + providerSessionState, + sessionId: "ws-handshake-session", + compaction: midTurnCompaction, + }); + const isolatedFailure = await streamOpenAICodexResponses(websocketModel, createCodexTestContext(), { + fetch: fetchMock as FetchImpl, + apiKey: token, + sessionId: "ws-handshake-session", + providerSessionState, + codexCompaction: { + operationId: "isolated-failure", + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "mid_turn", + strategy: "memento", + }, + }).result(); + expect(isolatedFailure.stopReason).toBe("error"); + expect(websocketInstances[0]?.readyState).toBe(MockWebSocket.OPEN); + expect(websocketInstances[2]?.readyState).toBe(MockWebSocket.CLOSED); + expect(websocketConstructorCount).toBe(3); + expect( + getOpenAICodexTransportDetails(websocketModel, { + sessionId: "ws-handshake-session", + providerSessionState, + }), + ).toMatchObject({ + websocketConnected: true, + hasTurnState: true, + }); // Turn-state is scoped to the current turn, so the SSE replay must be a // within-turn continuation (trailing tool result) to carry the header. const followUp: Context = { @@ -2285,6 +2431,155 @@ describe("openai-codex streaming", () => { providerSessionState, }).result(); expect(fetchMock).toHaveBeenCalledTimes(1); + if (!firstRequest || !continuationRequest || !continuationHeaders) { + throw new Error("expected both Codex transport requests"); + } + const firstMetadata = requireRecord(firstRequest.client_metadata, "first client_metadata"); + const continuationMetadata = requireRecord(continuationRequest.client_metadata, "continuation client_metadata"); + const firstTurnMetadata = parseTurnMetadata(firstMetadata); + const continuationTurnMetadata = parseTurnMetadata(continuationMetadata); + expect(continuationMetadata).toMatchObject({ + "x-codex-installation-id": TEST_INSTALLATION_ID, + session_id: firstMetadata.session_id, + thread_id: firstMetadata.thread_id, + turn_id: firstMetadata.turn_id, + }); + expect(continuationTurnMetadata).toMatchObject({ + installation_id: TEST_INSTALLATION_ID, + session_id: firstTurnMetadata.session_id, + thread_id: firstTurnMetadata.thread_id, + turn_id: firstTurnMetadata.turn_id, + window_id: continuationMetadata["x-codex-window-id"], + request_kind: "turn", + turn_started_at_unix_ms: context.messages[0]?.timestamp, + }); + expect(typeof continuationMetadata["x-codex-window-id"]).toBe("string"); + expect(continuationMetadata["x-codex-window-id"]).not.toBe(firstMetadata["x-codex-window-id"]); + expect(firstMetadata.session_id).toBe(continuationHeaders.get("session-id")); + expect(firstMetadata.thread_id).toBe(continuationHeaders.get("thread-id")); + expect(continuationMetadata["x-codex-window-id"]).toBe(continuationHeaders.get("x-codex-window-id")); + expect(continuationMetadata["x-codex-turn-metadata"]).toBe(continuationHeaders.get("x-codex-turn-metadata")); + }); + + it("clears stale main turn-state after pre-turn compaction", async () => { + const tempDir = TempDir.createSync("@pi-codex-stream-"); + setAgentDir(tempDir.path()); + const token = createCodexTestToken(); + const websocketInstances: MockWebSocket[] = []; + let websocketRequestCount = 0; + + class PreTurnCompactionWebSocket extends MockWebSocket { + handshakeHeaders = { + "x-codex-turn-state": "stale-main-turn-state", + "x-models-etag": "models-etag-1", + }; + + constructor(url: string, options?: { headers?: WsHeaders }) { + super(url, options); + websocketInstances.push(this); + queueMicrotask(() => { + this.readyState = MockWebSocket.OPEN; + this.emit("open", new Event("open")); + }); + } + + send(_data: string): void { + websocketRequestCount += 1; + this.emitCodexResponse({ + messageId: `msg_pre_turn_${websocketRequestCount}`, + responseId: `resp_pre_turn_${websocketRequestCount}`, + text: "Hello WS", + }); + } + } + + global.WebSocket = PreTurnCompactionWebSocket as unknown as typeof WebSocket; + const websocketModel = createCodexTestModel("https://chatgpt.com/backend-api"); + const sseModel: Model<"openai-codex-responses"> = buildModel({ + id: websocketModel.id, + name: websocketModel.name, + api: "openai-codex-responses", + provider: websocketModel.provider, + baseUrl: websocketModel.baseUrl, + reasoning: true, + preferWebsockets: false, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 128000, + maxTokens: 128000, + }); + const providerSessionState = new Map(); + const sessionId = "pre-turn-reset-session"; + let sseHeaders: Headers | undefined; + const fetchMock = vi.fn(async (_input: string | URL, init?: RequestInit) => { + sseHeaders = init?.headers instanceof Headers ? init.headers : new Headers(init?.headers); + return new Response(createCompletedCodexSse("Hello SSE"), { + headers: { "content-type": "text/event-stream" }, + }); + }); + const compaction: CodexCompactionRequestContext = { + operationId: "pre-turn-reset-operation", + trigger: "auto", + reason: "context_limit", + implementation: "responses", + phase: "pre_turn", + strategy: "memento", + }; + + try { + await streamOpenAICodexResponses(websocketModel, createCodexTestContext(), { + apiKey: token, + fetch: fetchMock as FetchImpl, + sessionId, + providerSessionState, + }).result(); + await streamOpenAICodexResponses(websocketModel, createCodexTestContext(), { + apiKey: token, + fetch: fetchMock as FetchImpl, + sessionId, + providerSessionState, + codexCompaction: compaction, + }).result(); + expect(fetchMock).not.toHaveBeenCalled(); + expect(websocketInstances).toHaveLength(2); + expect(websocketInstances[0]?.readyState).toBe(MockWebSocket.OPEN); + expect(websocketInstances[1]?.readyState).toBe(MockWebSocket.CLOSED); + expect(websocketInstances[1]?.options?.headers?.["x-codex-turn-state"]).toBeUndefined(); + expect(websocketInstances[1]?.options?.headers?.["x-models-etag"]).toBe("models-etag-1"); + + resetOpenAICodexHistoryAfterCompaction({ + providerSessionState, + sessionId, + compaction, + }); + expect( + getOpenAICodexTransportDetails(websocketModel, { + sessionId, + providerSessionState, + }), + ).toMatchObject({ + websocketConnected: true, + hasTurnState: false, + }); + await streamOpenAICodexResponses( + sseModel, + { + systemPrompt: ["You are a helpful assistant."], + messages: [{ role: "user", content: "Continue after compaction", timestamp: Date.now() }], + }, + { + apiKey: token, + fetch: fetchMock as FetchImpl, + sessionId, + providerSessionState, + }, + ).result(); + expect(fetchMock).toHaveBeenCalledTimes(1); + expect(sseHeaders?.get("x-codex-turn-state")).toBeNull(); + } finally { + for (const state of providerSessionState.values()) state.close(); + providerSessionState.clear(); + } }); it("includes service_tier in websocket payloads when requested", async () => { @@ -2472,6 +2767,17 @@ describe("openai-codex streaming", () => { expect(deltaItems[0]?.role).toBe("user"); expect(JSON.stringify(deltaItems)).toContain("Second question"); expect(JSON.stringify(deltaItems)).not.toContain("First answer"); + const firstMetadata = requireRecord(sentRequests[0]?.client_metadata, "first client_metadata"); + const secondMetadata = requireRecord(sentRequests[1]?.client_metadata, "second client_metadata"); + expect(secondMetadata).toMatchObject({ + "x-codex-installation-id": firstMetadata["x-codex-installation-id"], + session_id: firstMetadata.session_id, + thread_id: firstMetadata.thread_id, + "x-codex-window-id": firstMetadata["x-codex-window-id"], + }); + expect(secondMetadata.turn_id).not.toBe(firstMetadata.turn_id); + expect(parseTurnMetadata(firstMetadata).turn_started_at_unix_ms).toBe(firstContext.messages[0]?.timestamp); + expect(parseTurnMetadata(secondMetadata).turn_started_at_unix_ms).toBe(secondContext.messages.at(-1)?.timestamp); const stats = getOpenAICodexWebSocketDebugStats(model, { sessionId: "ws-delta-session", @@ -3938,10 +4244,12 @@ describe("openai-codex streaming", () => { let constructorCount = 0; let sendCount = 0; + let prewarmHeaders: WsHeaders | undefined; class ReusableWebSocket extends MockWebSocket { constructor(url: string, options?: { headers?: WsHeaders }) { super(url, options); constructorCount += 1; + prewarmHeaders = options?.headers; this.scheduleOpen(); } @@ -3979,6 +4287,11 @@ describe("openai-codex streaming", () => { sessionId: "ws-reuse-session", providerSessionState, }); + expect(prewarmHeaders?.["session-id"]).toBe("ws-reuse-session"); + expect(prewarmHeaders?.["thread-id"]).toBeDefined(); + expect(prewarmHeaders?.["x-codex-window-id"]).toBeDefined(); + expect(prewarmHeaders?.["x-codex-turn-metadata"]).toBeUndefined(); + expect(prewarmHeaders?.["x-codex-installation-id"]).toBeUndefined(); const firstContext: Context = { systemPrompt: ["You are a helpful assistant."], @@ -4016,6 +4329,26 @@ describe("openai-codex streaming", () => { expect(transportDetails.websocketConnected).toBe(true); expect(transportDetails.prewarmed).toBe(true); expect(transportDetails.canAppend).toBe(true); + resetOpenAICodexHistoryAfterCompaction({ + providerSessionState, + sessionId: "ws-reuse-session", + compaction: { + operationId: "history-rewrite", + trigger: "auto", + reason: "context_limit", + phase: "pre_turn", + strategy: "memento", + }, + }); + expect( + getOpenAICodexTransportDetails(model, { + sessionId: "ws-reuse-session", + providerSessionState, + }), + ).toMatchObject({ + websocketConnected: true, + canAppend: false, + }); }); it("scopes x-codex-turn-state to the current turn on SSE requests", async () => { diff --git a/packages/ai/test/openai-responses-history-payload.test.ts b/packages/ai/test/openai-responses-history-payload.test.ts index afdedd481..9520fabf0 100644 --- a/packages/ai/test/openai-responses-history-payload.test.ts +++ b/packages/ai/test/openai-responses-history-payload.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import { convertCodexResponsesMessages, streamOpenAICodexResponses, @@ -9,6 +9,17 @@ import type { Context, Model, ModelSpec, ProviderSessionState } from "@oh-my-pi/ import { createOpenAIResponsesHistoryPayload, truncateResponseItemId } from "@oh-my-pi/pi-ai/utils"; import { buildModel } from "@oh-my-pi/pi-catalog/build"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import * as piUtils from "@oh-my-pi/pi-utils"; + +const TEST_INSTALLATION_ID = "00000000-0000-4000-8000-000000000001"; + +beforeEach(() => { + vi.spyOn(piUtils, "getInstallId").mockReturnValue(TEST_INSTALLATION_ID); +}); + +afterEach(() => { + vi.restoreAllMocks(); +}); function createAbortedSignal(): AbortSignal { const controller = new AbortController(); diff --git a/packages/ai/test/openai-responses-stream-terminal.test.ts b/packages/ai/test/openai-responses-stream-terminal.test.ts index 63bc5f331..757a7f31f 100644 --- a/packages/ai/test/openai-responses-stream-terminal.test.ts +++ b/packages/ai/test/openai-responses-stream-terminal.test.ts @@ -288,7 +288,13 @@ describe("processResponsesStream: lost output_item.added recovery", () => { { type: "response.output_item.done", output_index: 0, - item: { type: "reasoning", summary: [{ type: "summary_text", text: "first" }] }, + item: { + type: "reasoning", + summary: [ + { type: "summary_text", text: "Plan" }, + { type: "summary_text", text: "Planning details" }, + ], + }, }, { type: "response.output_item.done", @@ -305,7 +311,7 @@ describe("processResponsesStream: lost output_item.added recovery", () => { expect(output.content).toHaveLength(2); const [first, second] = output.content; if (first?.type !== "thinking" || second?.type !== "thinking") throw new Error("expected thinking blocks"); - expect(first.thinking).toBe("first"); + expect(first.thinking).toBe("Plan\n\nPlanning details"); expect(second.thinking).toBe("second"); expect(first.thinkingSignature).toBeDefined(); expect(second.thinkingSignature).toBeDefined(); diff --git a/packages/catalog/src/models.json b/packages/catalog/src/models.json index 706e3b509..735714ef3 100644 --- a/packages/catalog/src/models.json +++ b/packages/catalog/src/models.json @@ -31524,7 +31524,7 @@ }, "openai/gpt-5.6-sol": { "id": "openai/gpt-5.6-sol", - "name": "GPT-5.6 Sol", + "name": "GPT-5.6 Sol (new)", "api": "openai-completions", "provider": "kilo", "baseUrl": "https://api.kilo.ai/api/gateway", @@ -73972,13 +73972,13 @@ "text" ], "cost": { - "input": 0.9199999999999999, - "output": 3, - "cacheRead": 0.18, + "input": 0.84, + "output": 2.64, + "cacheRead": 0.156, "cacheWrite": 0 }, "contextWindow": 1048576, - "maxTokens": 1048576, + "maxTokens": 128000, "thinking": { "mode": "effort", "efforts": [ diff --git a/packages/catalog/src/wire/codex.ts b/packages/catalog/src/wire/codex.ts index c2336b6d2..329ace70f 100644 --- a/packages/catalog/src/wire/codex.ts +++ b/packages/catalog/src/wire/codex.ts @@ -10,6 +10,13 @@ export const OPENAI_HEADERS = { ORIGINATOR: "originator", SESSION_ID: "session_id", CONVERSATION_ID: "conversation_id", + SCOPED_SESSION_ID: "session-id", + THREAD_ID: "thread-id", + INSTALLATION_ID: "x-codex-installation-id", + WINDOW_ID: "x-codex-window-id", + TURN_METADATA: "x-codex-turn-metadata", + PARENT_THREAD_ID: "x-codex-parent-thread-id", + SUBAGENT: "x-openai-subagent", /** Responses Lite transport marker (codex-rs `add_responses_lite_header`); value is always `"true"`. */ RESPONSES_LITE: "x-openai-internal-codex-responses-lite", } as const; diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 3b4274e5f..edac6e685 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -84,6 +84,7 @@ import type { AssistantMessageEvent, AssistantRetryRecovery, AssistantRetryRecoveryKind, + CodexCompactionContext, Context, ImageContent, Message, @@ -117,6 +118,7 @@ import { streamSimple, } from "@oh-my-pi/pi-ai"; import * as AIError from "@oh-my-pi/pi-ai/error"; +import { resetOpenAICodexHistoryAfterCompaction } from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { toolWireSchema } from "@oh-my-pi/pi-ai/utils/schema"; import { GeminiHeaderRunDetector, isGeminiThinkingModel } from "@oh-my-pi/pi-ai/utils/thinking-loop"; import { type RepeatedToolCallDetection, ToolCallLoopGuard } from "@oh-my-pi/pi-ai/utils/tool-call-loop-guard"; @@ -602,6 +604,20 @@ function compactionDeadEndWarning(remedies: string): string { ); } +function createCodexCompactionContext(options: { + trigger: CodexCompactionContext["trigger"]; + reason: CodexCompactionContext["reason"]; + phase: CodexCompactionContext["phase"]; +}): CodexCompactionContext { + return { + operationId: crypto.randomUUID(), + trigger: options.trigger, + reason: options.reason, + phase: options.phase, + strategy: "memento", + }; +} + /** * Per-turn prune cache window. A tool result whose all-message suffix exceeds * this is in the warm, already-sent prompt-cache prefix: re-writing it costs the @@ -2871,6 +2887,12 @@ export class AgentSession { // compaction path so the advisor model's maintenance call also emits spans. const telemetry = resolveTelemetry(agent.telemetry, advisorSessionId); + const codexCompaction = createCodexCompactionContext({ + trigger: "auto", + reason: "context_limit", + phase: "pre_turn", + }); + for (const candidate of candidates) { const apiKey = await this.#modelRegistry.getApiKey(candidate, advisorSessionId); if (!apiKey) continue; @@ -2889,6 +2911,8 @@ export class AgentSession { tools: agent.state.tools, sessionId: advisorSessionId, promptCacheKey: advisorSessionId, + providerSessionState: this.#providerSessionState, + codexCompaction, }, ); break; @@ -9751,6 +9775,7 @@ export class AgentSession { let firstKeptEntryId: string; let tokensBefore: number; let details: unknown; + let codexCompaction: CodexCompactionContext | undefined; // Snapcompact runs locally first. The frame cap is sized from the live // model window via #computeSnapcompactMaxFrames so the post-render context @@ -9830,6 +9855,11 @@ export class AgentSession { details = snapcompactResult.details; preserveData = { ...(compactionPrep.preserveData ?? {}), ...(snapcompactResult.preserveData ?? {}) }; } else { + codexCompaction = createCodexCompactionContext({ + trigger: "manual", + reason: "user_requested", + phase: "standalone_turn", + }); // Generate compaction result. Only convert known abort-shaped // rejections (AbortError raised while the abort signal is set, // or an already-typed sentinel) into `CompactionCancelledError` @@ -9851,6 +9881,7 @@ export class AgentSession { extraContext: compactionPrep.hookContext, remoteInstructions: this.#baseSystemPrompt.join("\n\n"), convertToLlm: messages => this.#convertToLlmForSideRequest(messages), + codexCompaction, }, compactionCandidates, ); @@ -9893,7 +9924,11 @@ export class AgentSession { this.#planReferenceSent = false; this.#resetAllAdvisorRuntimes(); this.#syncTodoPhasesFromBranch(); - this.#closeCodexProviderSessionsForHistoryRewrite(); + if (codexCompaction) { + this.#resetCodexProviderAfterCompaction(codexCompaction); + } else { + this.#closeCodexProviderSessionsForHistoryRewrite(); + } // Get the saved compaction entry for the hook const savedCompactionEntry = newEntries.find(e => e.type === "compaction" && e.summary === summary) as @@ -10259,6 +10294,7 @@ export class AgentSession { await this.#runAutoCompaction("threshold", false, false, false, { autoContinue: false, triggerContextTokens: contextTokens, + phase: "pre_turn", }); } @@ -10333,6 +10369,7 @@ export class AgentSession { suppressContinuation: true, suppressHandoff: true, triggerContextTokens: contextTokens, + phase: "mid_turn", }); if (signal?.aborted) return; @@ -10579,6 +10616,7 @@ export class AgentSession { return await this.#runAutoCompaction("threshold", false, false, allowDefer, { autoContinue, triggerContextTokens: postMaintenanceContextTokens, + phase: "pre_turn", }); } logger.debug("Auto-compaction threshold satisfied but context promotion took over", { @@ -10939,7 +10977,11 @@ export class AgentSession { ): Promise { const compactionEntryBefore = getLatestCompactionEntry(this.sessionManager.getBranch()); await this.#dropPersistedAssistantTurn(assistantMessage); - const result = await this.#runAutoCompaction(reason, true, false, allowDefer, options); + const result = await this.#runAutoCompaction(reason, true, false, allowDefer, { + autoContinue: options.autoContinue, + triggerContextTokens: options.triggerContextTokens, + phase: "mid_turn", + }); const compactionEntryAfter = getLatestCompactionEntry(this.sessionManager.getBranch()); if (result.historyRewritten !== true && compactionEntryAfter === compactionEntryBefore) { this.#restoreFailedAssistantTurn(assistantMessage); @@ -11551,6 +11593,14 @@ export class AgentSession { this.#closeProviderSessionsForModelSwitch(currentModel, currentModel); } + #resetCodexProviderAfterCompaction(compaction: CodexCompactionContext): void { + resetOpenAICodexHistoryAfterCompaction({ + providerSessionState: this.#providerSessionState, + sessionId: this.sessionId, + compaction, + }); + } + #resetCurrentResponsesProviderSession(reason: string): void { const currentModel = this.model; if (currentModel?.api !== "openai-responses" && currentModel?.api !== "openai-codex-responses") { @@ -11982,6 +12032,7 @@ export class AgentSession { tools: this.agent.state.tools, sessionId: this.sessionId, promptCacheKey: this.sessionId, + providerSessionState: this.#providerSessionState, // Route every summarization HTTP request through the // session's side-stream transport so the provider // concurrency cap (e.g. providers.ollama-cloud.maxConcurrency) @@ -12320,6 +12371,7 @@ export class AgentSession { triggerContextTokens?: number; suppressContinuation?: boolean; suppressHandoff?: boolean; + phase?: CodexCompactionContext["phase"]; } = {}, ): Promise { const compactionSettings = this.settings.getGroup("compaction"); @@ -12362,7 +12414,7 @@ export class AgentSession { async signal => { await Promise.resolve(); if (signal.aborted) return; - await this.#runAutoCompaction(reason, willRetry, true); + await this.#runAutoCompaction(reason, willRetry, true, true, { phase: options.phase }); }, { generation }, ); @@ -12509,6 +12561,7 @@ export class AgentSession { let hookCompaction: CompactionResult | undefined; let fromExtension = false; let preserveData: Record | undefined; + let codexCompaction: CodexCompactionContext | undefined; if (this.#extensionRunner?.hasHandlers("session_before_compact")) { const hookResult = (await this.#extensionRunner.emit({ @@ -12647,6 +12700,13 @@ export class AgentSession { const telemetry = resolveTelemetry(this.agent.telemetry, this.sessionId); let compactResult: CompactionResult | undefined; let lastError: unknown; + codexCompaction = createCodexCompactionContext({ + trigger: "auto", + reason: "context_limit", + phase: + options.phase ?? + (reason === "threshold" ? "pre_turn" : reason === "idle" ? "standalone_turn" : "mid_turn"), + }); for (let candidateIndex = 0; candidateIndex < candidates.length; candidateIndex++) { const candidate = candidates[candidateIndex]; @@ -12679,6 +12739,8 @@ export class AgentSession { tools: this.agent.state.tools, sessionId: this.sessionId, promptCacheKey: this.sessionId, + providerSessionState: this.#providerSessionState, + codexCompaction, }, ); break; @@ -12797,7 +12859,11 @@ export class AgentSession { this.#planReferenceSent = false; this.#resetAllAdvisorRuntimes(); this.#syncTodoPhasesFromBranch(); - this.#closeCodexProviderSessionsForHistoryRewrite(); + if (codexCompaction) { + this.#resetCodexProviderAfterCompaction(codexCompaction); + } else { + this.#closeCodexProviderSessionsForHistoryRewrite(); + } // Get the saved compaction entry for the hook const savedCompactionEntry = newEntries.find(e => e.type === "compaction" && e.summary === summary) as diff --git a/packages/coding-agent/test/agent-session-eager-compaction.test.ts b/packages/coding-agent/test/agent-session-eager-compaction.test.ts index a4d0c2db2..cd8528ac9 100644 --- a/packages/coding-agent/test/agent-session-eager-compaction.test.ts +++ b/packages/coding-agent/test/agent-session-eager-compaction.test.ts @@ -2,7 +2,8 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import * as path from "node:path"; import { Agent, type AgentMessage, type AgentTool } from "@oh-my-pi/pi-agent-core"; import * as compactionModule from "@oh-my-pi/pi-agent-core/compaction"; -import type { TextContent } from "@oh-my-pi/pi-ai"; +import type { Model, TextContent } from "@oh-my-pi/pi-ai"; +import * as codexResponses from "@oh-my-pi/pi-ai/providers/openai-codex-responses"; import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; @@ -135,22 +136,23 @@ describe("AgentSession eager prelude re-injection after compaction", () => { async function createHarness( settingsOverride: Record = {}, - opts: { agentId?: string; agentKind?: "main" | "sub" } = {}, + opts: { agentId?: string; agentKind?: "main" | "sub"; model?: Model } = {}, ): Promise { const observedCalls: ObservedPromptCall[] = []; const waiters: Array<{ predicate: (call: ObservedPromptCall) => boolean; resolve: (call: ObservedPromptCall) => void; }> = []; - const bundled = getBundledModel("anthropic", "claude-sonnet-4-5"); - if (!bundled) throw new Error("Expected claude-sonnet-4-5 model to exist"); + const defaultModel = getBundledModel("anthropic", "claude-sonnet-4-5"); + if (!defaultModel) throw new Error("Expected claude-sonnet-4-5 model to exist"); + const selectedModel = opts.model ?? defaultModel; // Pin the window and output reservation: usage figures below trip the // context-full strategy at a 200k/64k threshold; catalog regeneration must // not shift the headroom math. - const model = { ...bundled, contextWindow: 200_000, maxTokens: 64_000 }; + const model = { ...selectedModel, contextWindow: 200_000, maxTokens: 64_000 }; const authStorage = await AuthStorage.create(path.join(tempDir.path(), `testauth-${cleanups.length}.db`)); - authStorage.setRuntimeApiKey("anthropic", "test-key"); + authStorage.setRuntimeApiKey(model.provider, "test-key"); const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir.path(), `models-${cleanups.length}.yml`)); const settings = Settings.isolated({ "compaction.enabled": true, @@ -359,4 +361,25 @@ describe("AgentSession eager prelude re-injection after compaction", () => { expect(continuation.messageTexts.some(text => text.includes("Consider calling"))).toBe(false); expect(continuation.messageTexts.some(text => text.includes("You MUST call"))).toBe(false); }); + + it("resets Codex provider history after successful auto-compaction", async () => { + const model = getBundledModel("openai-codex", "gpt-5.6-terra"); + if (!model) throw new Error("Expected gpt-5.6-terra model to exist"); + const resetSpy = vi.spyOn(codexResponses, "resetOpenAICodexHistoryAfterCompaction"); + const { session, waitForCall } = await createHarness({}, { model }); + stubCompaction(); + + await runToContinuation(session, waitForCall); + + expect(resetSpy).toHaveBeenCalledTimes(1); + const reset = resetSpy.mock.calls[0]?.[0]; + if (!reset) throw new Error("Expected Codex compaction reset"); + expect(reset.providerSessionState).toBe(session.providerSessionState); + expect(reset.sessionId).toBe(session.sessionId); + expect(reset.compaction).toMatchObject({ + trigger: "auto", + reason: "context_limit", + phase: "pre_turn", + }); + }); });