diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index c4c77cd4f..d17d0da63 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -5,6 +5,7 @@ ### Fixed - Fixed `openai-codex-responses` fresh plan execution requests that contained only system/developer guidance by mirroring the final instruction as user input so Codex accepts the first turn. ([#4714](https://github.com/can1357/oh-my-pi/issues/4714)) +- Fixed Codex WebSocket compact/resume delta diagnostics to record request shape and raw-vs-displayed usage buckets, so persistent server-reported uncached suffixes without `orchestration_*` fields are visible in debug stats. ([#4707](https://github.com/can1357/oh-my-pi/issues/4707)) ## [16.3.10] - 2026-07-06 diff --git a/packages/ai/src/providers/openai-codex-responses.ts b/packages/ai/src/providers/openai-codex-responses.ts index 08ca0a434..f1f71167a 100644 --- a/packages/ai/src/providers/openai-codex-responses.ts +++ b/packages/ai/src/providers/openai-codex-responses.ts @@ -37,6 +37,7 @@ import type { Tool, ToolCall, ToolChoice, + Usage, } from "../types"; import { createOpenAIResponsesHistoryPayload, @@ -246,10 +247,66 @@ type CodexOutputBlock = | TextContent | (ToolCall & { [kStreamingPartialJson]: string; [kStreamingLastParseLen]?: number }); +interface CodexResponseUsage { + input_tokens?: number; + output_tokens?: number; + total_tokens?: number; + prompt_cache_hit_tokens?: number; + input_tokens_details?: { + cached_tokens?: number; + cache_write_tokens?: number; + orchestration_input_tokens?: number; + orchestration_input_cached_tokens?: number; + }; + output_tokens_details?: { + reasoning_tokens?: number; + orchestration_output_tokens?: number; + }; +} + +/** Shape of the Codex request sent on the latest provider turn. */ +export interface OpenAICodexTurnRequestDiagnostics { + transport: "sse" | "websocket"; + previousResponseIdPresent: boolean; + inputItemCount: number; + inputItemTypes: string[]; + firstInputItemType?: string; + inputJsonBytes: number; + promptCacheKey?: string; + toolsHash?: string; + optionsHash: string; + canAppendBeforeRequest: boolean; +} + +/** Raw provider usage plus the normalized buckets OMP displays for the latest Codex turn. */ +export interface OpenAICodexTurnUsageDiagnostics { + rawInputTokens: number; + rawCachedTokens: number; + rawUncachedTokens: number; + rawOutputTokens: number; + rawTotalTokens?: number; + rawOrchestrationInputTokens?: number; + rawOrchestrationCachedTokens?: number; + rawOrchestrationOutputTokens?: number; + displayedInputTokens: number; + displayedOutputTokens: number; + displayedCacheReadTokens: number; + displayedCacheWriteTokens: number; + displayedTotalTokens: number; + displayedOrchestrationInputTokens: number; + displayedOrchestrationCacheReadTokens: number; + displayedOrchestrationOutputTokens: number; +} + +/** Latest Codex turn request/usage diagnostics exposed to debug UIs and tests. */ +export interface OpenAICodexTurnDiagnostics { + request: OpenAICodexTurnRequestDiagnostics; + usage?: OpenAICodexTurnUsageDiagnostics; +} + /** - * Per-session request-shape counters. Despite the name, these cover both - * transports: once stateful SSE chaining is enabled, SSE requests are counted - * too (the shared chained-request builder records every request it shapes). + * Per-session request-shape counters and latest turn diagnostics. Despite the + * name, these cover both transports. */ export interface OpenAICodexWebSocketDebugStats { fullContextRequests: number; @@ -257,6 +314,7 @@ export interface OpenAICodexWebSocketDebugStats { lastInputItems: number; lastDeltaInputItems?: number; lastPreviousResponseId?: string; + lastTurn?: OpenAICodexTurnDiagnostics; } /** @@ -1002,6 +1060,7 @@ async function openCodexWebSocketTransport( requestBodyForState: RequestBody; transport: CodexTransport; }> { + 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`. @@ -1021,6 +1080,7 @@ async function openCodexWebSocketTransport( if (replacementWebsocketRequest !== undefined) { websocketRequest = replacementWebsocketRequest as typeof websocketRequest; } + recordCodexTurnRequestDiagnostics(websocketState, websocketRequest, "websocket", canAppendBeforeRequest); const websocketHeaders = createCodexHeaders( requestContext.requestHeaders, requestContext.accountId, @@ -1113,12 +1173,13 @@ async function openCodexSseTransport( ), ); }; + const canAppendBeforeRequest = state?.canAppend === true; let wireBody = body; const replacementWireBody = await options?.onPayload?.(wireBody, model); if (replacementWireBody !== undefined) { wireBody = replacementWireBody as RequestBody; } - recordCodexWebSocketRequestStats(state, wireBody); + recordCodexTurnRequestDiagnostics(state, wireBody, "sse", canAppendBeforeRequest); return { eventStream: await open(wireBody), requestBodyForState: structuredCloneJSON(wireBody), transport: "sse" }; } @@ -1525,34 +1586,19 @@ class CodexStreamProcessor { #handleResponseCompleted(rawEvent: Record): void { const { runtime, model, output } = this; runtime.sawTerminalEvent = true; - const response = ( - rawEvent as { - response?: { - id?: string; - usage?: { - input_tokens?: number; - output_tokens?: number; - total_tokens?: number; - input_tokens_details?: { - cached_tokens?: number; - orchestration_input_tokens?: number; - orchestration_input_cached_tokens?: number; - }; - output_tokens_details?: { - reasoning_tokens?: number; - orchestration_output_tokens?: number; - }; - }; - status?: string; - service_tier?: ServiceTier | "default"; - end_turn?: boolean; - }; - } - ).response; + const rawResponse = rawEvent.response; + const response = rawResponse && typeof rawResponse === "object" ? rawResponse : undefined; + const responseId = response && "id" in response && typeof response.id === "string" ? response.id : undefined; + const usage = response && "usage" in response ? parseCodexResponseUsage(response.usage) : undefined; + const serviceTier = + response && "service_tier" in response ? parseCodexServiceTier(response.service_tier) : undefined; + const status = response && "status" in response ? parseCodexResponseStatus(response.status) : undefined; + const endTurn = response && "end_turn" in response ? response.end_turn : undefined; - populateResponsesUsageFromResponse(output, response?.usage); - if (typeof response?.id === "string" && response.id.length > 0) { - output.responseId = response.id; + populateResponsesUsageFromResponse(output, usage); + recordCodexTurnUsageDiagnostics(runtime.websocketState, usage, output.usage); + if (responseId) { + output.responseId = responseId; } const state = runtime.websocketState; @@ -1564,8 +1610,8 @@ class CodexStreamProcessor { resetCodexWebSocketAppendState(state); } else { state.lastRequest = structuredCloneJSON(runtime.requestBodyForState); - if (typeof response?.id === "string" && response.id.length > 0) { - state.lastResponseId = response.id; + if (responseId) { + state.lastResponseId = responseId; state.lastResponseItems = stripInputItemIds(structuredCloneJSON(runtime.nativeOutputItems)); state.canAppend = rawEvent.type === "response.done" || rawEvent.type === "response.completed"; } else { @@ -1578,14 +1624,9 @@ class CodexStreamProcessor { finalizePendingResponsesToolCalls(output); calculateCost(model, output.usage); - applyCodexServiceTierPricing( - model, - output.usage, - response?.service_tier, - runtime.requestBodyForState.service_tier, - ); - output.stopReason = mapOpenAIResponsesStopReason(response?.status as ResponseStatus | undefined); - promoteResponsesToolUseStopReason(output, response?.end_turn); + applyCodexServiceTierPricing(model, output.usage, serviceTier, runtime.requestBodyForState.service_tier); + output.stopReason = mapOpenAIResponsesStopReason(status); + promoteResponsesToolUseStopReason(output, endTurn === true ? true : endTurn === false ? false : undefined); } async #recoverStreamError(error: unknown): Promise { @@ -2270,22 +2311,230 @@ function stripInputItemIds(items: Array>): InputItem[] { }); } -function recordCodexWebSocketRequestStats( +const codexDiagnosticsTextEncoder = new TextEncoder(); + +function jsonByteLength(value: unknown): number { + const json = JSON.stringify(value); + return codexDiagnosticsTextEncoder.encode(json === undefined ? "undefined" : json).byteLength; +} + +function hashJson(value: unknown): string { + const json = JSON.stringify(value); + return String(Bun.hash(json === undefined ? "undefined" : json)); +} + +function parseCodexServiceTier(value: unknown): ServiceTier | undefined { + switch (value) { + case "auto": + case "default": + case "flex": + case "scale": + case "priority": + return value; + default: + return undefined; + } +} + +function parseCodexResponseStatus(value: unknown): ResponseStatus | undefined { + switch (value) { + case "completed": + case "failed": + case "in_progress": + case "cancelled": + case "queued": + case "incomplete": + return value; + default: + return undefined; + } +} + +function parseCodexResponseUsage(value: unknown): CodexResponseUsage | undefined { + if (!value || typeof value !== "object") return undefined; + const usage: CodexResponseUsage = {}; + let hasUsage = false; + if ("input_tokens" in value && typeof value.input_tokens === "number") { + usage.input_tokens = value.input_tokens; + hasUsage = true; + } + if ("output_tokens" in value && typeof value.output_tokens === "number") { + usage.output_tokens = value.output_tokens; + hasUsage = true; + } + if ("total_tokens" in value && typeof value.total_tokens === "number") { + usage.total_tokens = value.total_tokens; + hasUsage = true; + } + if ("prompt_cache_hit_tokens" in value && typeof value.prompt_cache_hit_tokens === "number") { + usage.prompt_cache_hit_tokens = value.prompt_cache_hit_tokens; + hasUsage = true; + } + if ( + "input_tokens_details" in value && + value.input_tokens_details && + typeof value.input_tokens_details === "object" + ) { + const details = value.input_tokens_details; + const parsedDetails: NonNullable = {}; + let hasDetails = false; + if ("cached_tokens" in details && typeof details.cached_tokens === "number") { + parsedDetails.cached_tokens = details.cached_tokens; + hasDetails = true; + } + if ("cache_write_tokens" in details && typeof details.cache_write_tokens === "number") { + parsedDetails.cache_write_tokens = details.cache_write_tokens; + hasDetails = true; + } + if ("orchestration_input_tokens" in details && typeof details.orchestration_input_tokens === "number") { + parsedDetails.orchestration_input_tokens = details.orchestration_input_tokens; + hasDetails = true; + } + if ( + "orchestration_input_cached_tokens" in details && + typeof details.orchestration_input_cached_tokens === "number" + ) { + parsedDetails.orchestration_input_cached_tokens = details.orchestration_input_cached_tokens; + hasDetails = true; + } + if (hasDetails) { + usage.input_tokens_details = parsedDetails; + hasUsage = true; + } + } + if ( + "output_tokens_details" in value && + value.output_tokens_details && + typeof value.output_tokens_details === "object" + ) { + const details = value.output_tokens_details; + const parsedDetails: NonNullable = {}; + let hasDetails = false; + if ("reasoning_tokens" in details && typeof details.reasoning_tokens === "number") { + parsedDetails.reasoning_tokens = details.reasoning_tokens; + hasDetails = true; + } + if ("orchestration_output_tokens" in details && typeof details.orchestration_output_tokens === "number") { + parsedDetails.orchestration_output_tokens = details.orchestration_output_tokens; + hasDetails = true; + } + if (hasDetails) { + usage.output_tokens_details = parsedDetails; + hasUsage = true; + } + } + return hasUsage ? usage : undefined; +} + +function describeCodexInputItemType(item: unknown): string { + if (item && typeof item === "object") { + if ("type" in item && typeof item.type === "string") return item.type; + if ("role" in item && typeof item.role === "string") return item.role; + } + return typeof item; +} + +function createCodexOptionsHash(request: Record): string { + const options: Record = {}; + for (const key in request) { + if (key === "input" || key === "previous_response_id" || key === "type" || key === "client_metadata") { + continue; + } + options[key] = request[key]; + } + return hashJson(options); +} + +function buildCodexTurnRequestDiagnostics( + request: Record, + transport: CodexTransport, + canAppendBeforeRequest: boolean, +): OpenAICodexTurnRequestDiagnostics { + const input = request.input; + const inputItems = Array.isArray(input) ? input : []; + const inputItemTypes = inputItems.map(describeCodexInputItemType); + const promptCacheKey = typeof request.prompt_cache_key === "string" ? request.prompt_cache_key : undefined; + const toolsHash = request.tools === undefined ? undefined : hashJson(request.tools); + return { + transport, + previousResponseIdPresent: + typeof request.previous_response_id === "string" && request.previous_response_id.length > 0, + inputItemCount: inputItems.length, + inputItemTypes, + ...(inputItemTypes[0] ? { firstInputItemType: inputItemTypes[0] } : {}), + inputJsonBytes: jsonByteLength(inputItems), + ...(promptCacheKey !== undefined ? { promptCacheKey } : {}), + ...(toolsHash !== undefined ? { toolsHash } : {}), + optionsHash: createCodexOptionsHash(request), + canAppendBeforeRequest, + }; +} + +function recordCodexTurnRequestDiagnostics( state: CodexWebSocketSessionState | undefined, request: Record, + transport: CodexTransport, + canAppendBeforeRequest: boolean, ): void { if (!state) return; const input = request.input; state.stats.lastInputItems = Array.isArray(input) ? input.length : 0; - if (typeof request.previous_response_id === "string" && request.previous_response_id.length > 0) { + const previousResponseId = + typeof request.previous_response_id === "string" ? request.previous_response_id : undefined; + if (previousResponseId && previousResponseId.length > 0) { state.stats.deltaRequests += 1; state.stats.lastDeltaInputItems = state.stats.lastInputItems; - state.stats.lastPreviousResponseId = request.previous_response_id; - return; + state.stats.lastPreviousResponseId = previousResponseId; + } else { + state.stats.fullContextRequests += 1; + state.stats.lastDeltaInputItems = undefined; + state.stats.lastPreviousResponseId = undefined; } - state.stats.fullContextRequests += 1; - state.stats.lastDeltaInputItems = undefined; - state.stats.lastPreviousResponseId = undefined; + state.stats.lastTurn = { + request: buildCodexTurnRequestDiagnostics(request, transport, canAppendBeforeRequest), + }; + CODEX_DEBUG && logger.debug("[codex] codex turn request diagnostics", { diagnostics: state.stats.lastTurn.request }); +} + +function recordCodexTurnUsageDiagnostics( + state: CodexWebSocketSessionState | undefined, + rawUsage: CodexResponseUsage | undefined, + displayedUsage: Usage, +): void { + if (!state?.stats.lastTurn || !rawUsage) return; + const details = rawUsage.input_tokens_details; + const outputDetails = rawUsage.output_tokens_details; + const rawInputTokens = rawUsage.input_tokens ?? 0; + const rawCachedTokens = details?.cached_tokens ?? rawUsage.prompt_cache_hit_tokens ?? 0; + const usageDiagnostics: OpenAICodexTurnUsageDiagnostics = { + rawInputTokens, + rawCachedTokens, + rawUncachedTokens: Math.max(0, rawInputTokens - rawCachedTokens), + rawOutputTokens: rawUsage.output_tokens ?? 0, + ...(typeof rawUsage.total_tokens === "number" ? { rawTotalTokens: rawUsage.total_tokens } : {}), + ...(typeof details?.orchestration_input_tokens === "number" + ? { rawOrchestrationInputTokens: details.orchestration_input_tokens } + : {}), + ...(typeof details?.orchestration_input_cached_tokens === "number" + ? { rawOrchestrationCachedTokens: details.orchestration_input_cached_tokens } + : {}), + ...(typeof outputDetails?.orchestration_output_tokens === "number" + ? { rawOrchestrationOutputTokens: outputDetails.orchestration_output_tokens } + : {}), + displayedInputTokens: displayedUsage.input, + displayedOutputTokens: displayedUsage.output, + displayedCacheReadTokens: displayedUsage.cacheRead, + displayedCacheWriteTokens: displayedUsage.cacheWrite, + displayedTotalTokens: displayedUsage.totalTokens, + displayedOrchestrationInputTokens: displayedUsage.orchestration?.input ?? 0, + displayedOrchestrationCacheReadTokens: displayedUsage.orchestration?.cacheRead ?? 0, + displayedOrchestrationOutputTokens: displayedUsage.orchestration?.output ?? 0, + }; + state.stats.lastTurn = { + ...state.stats.lastTurn, + usage: usageDiagnostics, + }; + CODEX_DEBUG && logger.debug("[codex] codex turn diagnostics", { diagnostics: state.stats.lastTurn }); } /** @@ -2306,9 +2555,7 @@ function buildCodexChainedRequestBody( ? buildResponsesDeltaInput(state.lastRequest, state.lastResponseItems, requestBody) : null; if (appendInput && appendInput.length > 0 && state?.lastResponseId) { - const body: RequestBody = { ...requestBody, previous_response_id: state.lastResponseId, input: appendInput }; - recordCodexWebSocketRequestStats(state, body); - return body; + return { ...requestBody, previous_response_id: state.lastResponseId, input: appendInput }; } if (chainable && state) { // Chaining was eligible but the prefix/options check failed: history @@ -2322,7 +2569,6 @@ function buildCodexChainedRequestBody( state.turnState = undefined; state.modelsEtag = undefined; } - recordCodexWebSocketRequestStats(state, requestBody); return requestBody; } diff --git a/packages/ai/test/openai-codex-stream.test.ts b/packages/ai/test/openai-codex-stream.test.ts index b7979ab40..13bd10e96 100644 --- a/packages/ai/test/openai-codex-stream.test.ts +++ b/packages/ai/test/openai-codex-stream.test.ts @@ -140,7 +140,14 @@ function encodeWebSocketMessage(value: Record): Uint8Array { type WsHeaders = Record; type WsEventType = "open" | "message" | "error" | "close"; -const DEFAULT_USAGE = { +type CodexTestUsage = { + input_tokens: number; + output_tokens: number; + total_tokens: number; + input_tokens_details: { cached_tokens: number }; +}; + +const DEFAULT_USAGE: CodexTestUsage = { input_tokens: 5, output_tokens: 3, total_tokens: 8, @@ -210,8 +217,16 @@ class MockWebSocket { text: string; terminalType?: "response.done" | "response.completed"; includeCreated?: boolean; + usage?: CodexTestUsage; }): void { - const { messageId, responseId, text, terminalType = "response.done", includeCreated = false } = opts; + const { + messageId, + responseId, + text, + terminalType = "response.done", + includeCreated = false, + usage = DEFAULT_USAGE, + } = opts; if (includeCreated) { this.sendJson({ type: "response.created", response: { id: responseId } }); } @@ -236,7 +251,7 @@ class MockWebSocket { response: { id: responseId, status: "completed", - usage: DEFAULT_USAGE, + usage, }, }); } @@ -2230,7 +2245,7 @@ describe("openai-codex streaming", () => { expect(result.usage.premiumRequests).toBeUndefined(); }); - it("sends websocket continuation deltas after prior assistant response items and records stats", async () => { + it("records websocket delta request and usage diagnostics", async () => { const tempDir = TempDir.createSync("@pi-codex-stream-"); setAgentDir(tempDir.path()); const payload = Buffer.from( @@ -2242,6 +2257,12 @@ describe("openai-codex streaming", () => { const fetchMock = vi.fn(async () => { throw new Error("SSE fallback should not be called"); }); + const secondTurnUsage: CodexTestUsage = { + input_tokens: 132278, + input_tokens_details: { cached_tokens: 124416 }, + output_tokens: 29, + total_tokens: 132307, + }; class DeltaWebSocket extends MockWebSocket { constructor(url: string, options?: { headers?: WsHeaders }) { @@ -2258,6 +2279,7 @@ describe("openai-codex streaming", () => { text: responseIndex === 1 ? "First answer" : "Second answer", terminalType: "response.completed", includeCreated: true, + usage: responseIndex === 2 ? secondTurnUsage : DEFAULT_USAGE, }); } } @@ -2287,6 +2309,7 @@ describe("openai-codex streaming", () => { sessionId: "ws-delta-session", providerSessionState, }).result(); + expect(firstResponse.stopReason).toBe("stop"); const secondContext: Context = { systemPrompt: ["You are a helpful assistant.", "Use concise answers."], messages: [ @@ -2295,12 +2318,13 @@ describe("openai-codex streaming", () => { { role: "user", content: "Second question", timestamp: Date.now() }, ], }; - await streamOpenAICodexResponses(model, secondContext, { + const secondResponse = await streamOpenAICodexResponses(model, secondContext, { fetch: fetchMock as FetchImpl, apiKey: token, sessionId: "ws-delta-session", providerSessionState, }).result(); + expect(secondResponse.stopReason).toBe("stop"); expect(fetchMock).not.toHaveBeenCalled(); expect(sentRequests).toHaveLength(2); @@ -2330,12 +2354,36 @@ describe("openai-codex streaming", () => { sessionId: "ws-delta-session", providerSessionState, }); - expect(stats).toEqual({ - fullContextRequests: 1, - deltaRequests: 1, - lastInputItems: 1, - lastDeltaInputItems: 1, - lastPreviousResponseId: "resp_1", + expect(stats?.fullContextRequests).toBe(1); + expect(stats?.deltaRequests).toBe(1); + expect(stats?.lastInputItems).toBe(1); + expect(stats?.lastDeltaInputItems).toBe(1); + expect(stats?.lastPreviousResponseId).toBe("resp_1"); + expect(stats?.lastTurn?.request).toMatchObject({ + transport: "websocket", + previousResponseIdPresent: true, + inputItemCount: 1, + inputItemTypes: ["user"], + firstInputItemType: "user", + canAppendBeforeRequest: true, + promptCacheKey: "ws-delta-session", + }); + expect(stats?.lastTurn?.request.inputJsonBytes).toBeGreaterThan(0); + expect(stats?.lastTurn?.request.inputJsonBytes).toBeLessThan(1000); + expect(stats?.lastTurn?.usage).toEqual({ + rawInputTokens: 132278, + rawCachedTokens: 124416, + rawUncachedTokens: 7862, + rawOutputTokens: 29, + rawTotalTokens: 132307, + displayedInputTokens: 7862, + displayedOutputTokens: 29, + displayedCacheReadTokens: 124416, + displayedCacheWriteTokens: 0, + displayedTotalTokens: 132307, + displayedOrchestrationInputTokens: 0, + displayedOrchestrationCacheReadTokens: 0, + displayedOrchestrationOutputTokens: 0, }); }); @@ -2652,7 +2700,7 @@ describe("openai-codex streaming", () => { sessionId: "ws-expired-previous-response-session", providerSessionState, }); - expect(stats).toEqual({ + expect(stats).toMatchObject({ fullContextRequests: 2, deltaRequests: 1, lastInputItems: (retryInput as unknown[]).length,