From 7111db3598bce6acc5028a40d04953a6fcfea9f7 Mon Sep 17 00:00:00 2001 From: Amirsalar Sojoudi Date: Thu, 1 Oct 2026 10:40:33 +0700 Subject: [PATCH] fix(responses): wait for real usage before emitting response.completed --- open-sse/translator/index.js | 5 + .../translator/response/openai-responses.js | 33 +-- open-sse/utils/stream.js | 16 ++ .../openai-responses-completed-usage.test.js | 219 ++++++++++++++++++ .../unit/openai-responses-usage-pivot.test.js | 116 ++++++++++ 5 files changed, 376 insertions(+), 13 deletions(-) create mode 100644 tests/unit/openai-responses-completed-usage.test.js create mode 100644 tests/unit/openai-responses-usage-pivot.test.js diff --git a/open-sse/translator/index.js b/open-sse/translator/index.js index 292ffa36..cfccca84 100644 --- a/open-sse/translator/index.js +++ b/open-sse/translator/index.js @@ -287,6 +287,11 @@ export function initState(sourceFormat) { funcArgsDone: {}, funcItemDone: {}, customToolNames: new Set(), + // Chat Completions usage for response.completed. Not state.usage: other translators in + // the same pipeline overwrite that in their own shapes. + responsesUsage: null, + // finish_reason arrived before usage; response.completed waits for the usage chunk. + completionPending: false, completedSent: false }; } diff --git a/open-sse/translator/response/openai-responses.js b/open-sse/translator/response/openai-responses.js index 6a865d1a..1499a2a5 100644 --- a/open-sse/translator/response/openai-responses.js +++ b/open-sse/translator/response/openai-responses.js @@ -27,17 +27,22 @@ import { ROLE, OPENAI_BLOCK, RESPONSES_ITEM, OPENAI_FINISH, MODEL_FALLBACK } fro function toResponsesUsage(usage) { if (!usage || typeof usage !== "object") return null; - const inputTokens = [usage.input_tokens, usage.prompt_tokens].find(Number.isFinite) ?? 0; - const outputTokens = [usage.output_tokens, usage.completion_tokens].find(Number.isFinite) ?? 0; + const inputTokens = [usage.input_tokens, usage.prompt_tokens].find(Number.isInteger); + const outputTokens = [usage.output_tokens, usage.completion_tokens].find(Number.isInteger); + // Some upstreams attach zeroed placeholders to every chunk. Wait for real counts + // so response.completed cannot freeze the placeholder before the usage trailer. + if (inputTokens === undefined || outputTokens === undefined || inputTokens + outputTokens <= 0) { + return null; + } const responseUsage = { input_tokens: inputTokens, output_tokens: outputTokens, - total_tokens: Number.isFinite(usage.total_tokens) ? usage.total_tokens : inputTokens + outputTokens + total_tokens: inputTokens + outputTokens }; - const cachedTokens = [usage.input_tokens_details?.cached_tokens, usage.prompt_tokens_details?.cached_tokens].find(Number.isFinite); - const reasoningTokens = [usage.output_tokens_details?.reasoning_tokens, usage.completion_tokens_details?.reasoning_tokens].find(Number.isFinite); - if (Number.isFinite(cachedTokens)) responseUsage.input_tokens_details = { cached_tokens: cachedTokens }; - if (Number.isFinite(reasoningTokens)) responseUsage.output_tokens_details = { reasoning_tokens: reasoningTokens }; + const cachedTokens = [usage.input_tokens_details?.cached_tokens, usage.prompt_tokens_details?.cached_tokens].find(Number.isInteger); + const reasoningTokens = [usage.output_tokens_details?.reasoning_tokens, usage.completion_tokens_details?.reasoning_tokens].find(Number.isInteger); + if (Number.isInteger(cachedTokens)) responseUsage.input_tokens_details = { cached_tokens: cachedTokens }; + if (Number.isInteger(reasoningTokens)) responseUsage.output_tokens_details = { reasoning_tokens: reasoningTokens }; return responseUsage; } @@ -47,13 +52,14 @@ export function openaiToOpenAIResponsesResponse(chunk, state) { return flushEvents(state); } - // Capture upstream usage BEFORE the choices guard below: the last OpenAI chunk - // may carry usage together with an empty choices array, and it must not be dropped. - if (chunk.usage) { - state.responsesUsage = toResponsesUsage(chunk.usage); - } + // Capture usage before the choices guard: OpenAI may send it in a trailer + // whose choices array is empty. + const responseUsage = toResponsesUsage(chunk.usage); + if (responseUsage) state.responsesUsage = responseUsage; - if (!chunk.choices?.length) return []; + if (!chunk.choices?.length) { + return state.completionPending && state.responsesUsage ? flushEvents(state) : []; + } const events = []; const nextSeq = () => ++state.seq; @@ -163,6 +169,7 @@ export function openaiToOpenAIResponsesResponse(chunk, state) { // would swallow the terminal event entirely. Keep the old behaviour there. const flushReachesUs = state.targetFormat === FORMATS.OPENAI; if (state.responsesUsage || !flushReachesUs) sendCompleted(state, emit); + else state.completionPending = true; } return events; diff --git a/open-sse/utils/stream.js b/open-sse/utils/stream.js index cc0c6e30..91a041f4 100644 --- a/open-sse/utils/stream.js +++ b/open-sse/utils/stream.js @@ -266,6 +266,22 @@ export function createSSEStream(options = {}) { // For Ollama: done=true is the final chunk with finish_reason/usage, must translate // For other formats: done=true is the [DONE] sentinel, skip if (parsed && parsed.done && targetFormat !== FORMATS.OLLAMA) { + // A direct Chat-to-Responses translation can defer response.completed + // while waiting for a usage trailer. [DONE] ends that opportunity even + // if the upstream keeps the HTTP connection open, so finish now. + if (targetFormat === FORMATS.OPENAI && sourceFormat === FORMATS.OPENAI_RESPONSES && + state.completionPending && !state.completedSent) { + const completed = translateResponse(targetFormat, sourceFormat, null, state); + for (const item of completed || []) { + if (item === null || item === undefined) continue; + const output = formatSSE(item, sourceFormat); + reqLogger?.appendConvertedChunk?.(output); + controller.enqueue(sharedEncoder.encode(output)); + sseEmittedCount++; + } + finalizeStream(); + } + // Synthesize response.failed if the Responses stream never sent a terminal event if (keepsOpenAIResponsesFormat && !openAIResponsesTerminalSeen) { const failedOutput = formatIncompleteOpenAIResponsesStreamFailure(); diff --git a/tests/unit/openai-responses-completed-usage.test.js b/tests/unit/openai-responses-completed-usage.test.js new file mode 100644 index 00000000..f8f8005c --- /dev/null +++ b/tests/unit/openai-responses-completed-usage.test.js @@ -0,0 +1,219 @@ +import { describe, expect, it } from "vitest"; + +import { FORMATS } from "../../open-sse/translator/formats.js"; +import { createSSETransformStreamWithLogger } from "../../open-sse/utils/stream.js"; + +// Codex compacts its history only from the token usage reported on `response.completed` +// (sess.get_total_token_usage). Without it Codex never compacts and eventually sends a +// prompt larger than the model's context window. + +async function runTransform(targetFormat, lines) { + const encoder = new TextEncoder(); + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode(lines.join("\n"))); + controller.close(); + }, + }); + + const output = stream.pipeThrough( + createSSETransformStreamWithLogger(targetFormat, FORMATS.OPENAI_RESPONSES, "test", null, null, "test-model"), + ); + + const reader = output.getReader(); + const decoder = new TextDecoder(); + let text = ""; + while (true) { + const { value, done } = await reader.read(); + if (done) break; + text += decoder.decode(value, { stream: true }); + } + return text + decoder.decode(); +} + +// Parse the client-facing SSE into [{ event, data }], skipping the [DONE] sentinel. +function parseEvents(text) { + return text + .split("\n\n") + .map((block) => { + const event = block.match(/^event: (.+)$/m)?.[1]; + const data = block.match(/^data: (.+)$/m)?.[1]; + if (!event || !data || data === "[DONE]") return null; + return { event, data: JSON.parse(data) }; + }) + .filter(Boolean); +} + +const sse = (data) => [`data: ${JSON.stringify(data)}`, ""]; +const claudeSse = (data) => [`event: ${data.type}`, `data: ${JSON.stringify(data)}`, ""]; + +// Mirrors ResponseCompletedUsage in codex-rs/codex-api/src/sse/responses.rs. The three totals +// are required i64s and the details are optional, but each detail field is a required i64 +// when present. Any deviation fails Codex's parse of the whole event, killing the turn. +function expectCodexUsageShape(usage) { + for (const field of ["input_tokens", "output_tokens", "total_tokens"]) { + expect(Number.isInteger(usage[field]), field).toBe(true); + } + if (usage.input_tokens_details) { + expect(Number.isInteger(usage.input_tokens_details.cached_tokens)).toBe(true); + } + if (usage.output_tokens_details) { + expect(Number.isInteger(usage.output_tokens_details.reasoning_tokens)).toBe(true); + } +} + +describe("Responses response.completed reports token usage", () => { + it("reports Claude upstream usage to a Codex client, counting cached prompt tokens", async () => { + const events = parseEvents(await runTransform(FORMATS.CLAUDE, [ + ...claudeSse({ + type: "message_start", + message: { + id: "msg_1", type: "message", role: "assistant", model: "claude-opus-5", content: [], + usage: { input_tokens: 1000, cache_read_input_tokens: 200, cache_creation_input_tokens: 0, output_tokens: 1 }, + }, + }), + ...claudeSse({ type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }), + ...claudeSse({ type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "Hello" } }), + ...claudeSse({ type: "content_block_stop", index: 0 }), + ...claudeSse({ type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 50 } }), + ...claudeSse({ type: "message_stop" }), + ])); + + const completed = events.filter((e) => e.event === "response.completed"); + expect(completed).toHaveLength(1); + + const usage = completed[0].data.response.usage; + expectCodexUsageShape(usage); + expect(usage).toMatchObject({ + input_tokens: 1200, + output_tokens: 50, + total_tokens: 1250, + input_tokens_details: { cached_tokens: 200 }, + }); + }); + + it("waits for a trailing usage-only chunk instead of completing on finish_reason", async () => { + const events = parseEvents(await runTransform(FORMATS.OPENAI, [ + ...sse({ id: "c1", object: "chat.completion.chunk", choices: [{ index: 0, delta: { role: "assistant", content: "Hi" } }] }), + ...sse({ id: "c1", object: "chat.completion.chunk", choices: [{ index: 0, delta: {}, finish_reason: "stop" }] }), + ...sse({ + id: "c1", object: "chat.completion.chunk", choices: [], + usage: { + prompt_tokens: 300, completion_tokens: 20, total_tokens: 320, + prompt_tokens_details: { cached_tokens: 100 }, + completion_tokens_details: { reasoning_tokens: 5 }, + }, + }), + "data: [DONE]", + "", + ])); + + const completed = events.filter((e) => e.event === "response.completed"); + expect(completed).toHaveLength(1); + // Terminal event last: Codex stops reading at response.completed. + expect(events.at(-1).event).toBe("response.completed"); + + const usage = completed[0].data.response.usage; + expectCodexUsageShape(usage); + expect(usage).toEqual({ + input_tokens: 300, + output_tokens: 20, + total_tokens: 320, + input_tokens_details: { cached_tokens: 100 }, + output_tokens_details: { reasoning_tokens: 5 }, + }); + }); + + it("ignores zeroed placeholder usage on every chunk and reports the real trailing counts", async () => { + const placeholder = { prompt_tokens: 0, completion_tokens: 0, total_tokens: 0 }; + const events = parseEvents(await runTransform(FORMATS.OPENAI, [ + ...sse({ id: "c3", object: "chat.completion.chunk", usage: placeholder, choices: [{ index: 0, delta: { role: "assistant", content: "Hi" } }] }), + ...sse({ id: "c3", object: "chat.completion.chunk", usage: placeholder, choices: [{ index: 0, delta: {}, finish_reason: "stop" }] }), + ...sse({ id: "c3", object: "chat.completion.chunk", choices: [], usage: { prompt_tokens: 300, completion_tokens: 20, total_tokens: 320 } }), + "data: [DONE]", + "", + ])); + + const completed = events.filter((e) => e.event === "response.completed"); + expect(completed).toHaveLength(1); + expect(completed[0].data.response.usage).toEqual({ input_tokens: 300, output_tokens: 20, total_tokens: 320 }); + }); + + it.each([0, 999])("derives the total when upstream reports inconsistent total_tokens=%i", async (totalTokens) => { + const events = parseEvents(await runTransform(FORMATS.OPENAI, [ + ...sse({ id: "c4", object: "chat.completion.chunk", choices: [{ index: 0, delta: { role: "assistant", content: "Hi" } }] }), + ...sse({ + id: "c4", object: "chat.completion.chunk", + choices: [{ index: 0, delta: {}, finish_reason: "stop" }], + usage: { prompt_tokens: 300, completion_tokens: 20, total_tokens: totalTokens }, + }), + "data: [DONE]", + "", + ])); + + const completed = events.filter((e) => e.event === "response.completed"); + expect(completed).toHaveLength(1); + expect(completed[0].data.response.usage).toEqual({ + input_tokens: 300, + output_tokens: 20, + total_tokens: 320, + }); + }); + + it("completes at [DONE] while the upstream connection remains open", async () => { + let upstream; + const input = new ReadableStream({ + start(controller) { + upstream = controller; + }, + }); + const reader = input.pipeThrough( + createSSETransformStreamWithLogger(FORMATS.OPENAI, FORMATS.OPENAI_RESPONSES, "test"), + ).getReader(); + const decoder = new TextDecoder(); + const frames = [ + ...sse({ id: "c5", object: "chat.completion.chunk", choices: [{ index: 0, delta: { role: "assistant", content: "Hi" } }] }), + ...sse({ id: "c5", object: "chat.completion.chunk", choices: [{ index: 0, delta: {}, finish_reason: "stop" }] }), + "data: [DONE]", + "", + ]; + upstream.enqueue(new TextEncoder().encode(frames.join("\n") + "\n")); + + let timer; + try { + let output = ""; + const readUntilCompleted = async () => { + while (!output.includes('"type":"response.completed"')) { + const { value, done } = await reader.read(); + if (done) throw new Error("stream ended before response.completed"); + output += decoder.decode(value, { stream: true }); + } + }; + await Promise.race([ + readUntilCompleted(), + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error("response.completed waited for transport EOF")), 500); + }), + ]); + expect(parseEvents(output).filter((e) => e.event === "response.completed")).toHaveLength(1); + } finally { + clearTimeout(timer); + upstream.close(); + await reader.cancel(); + } + }); + + it("still completes exactly once, without inventing usage, when the upstream reports none", async () => { + const events = parseEvents(await runTransform(FORMATS.OPENAI, [ + ...sse({ id: "c2", object: "chat.completion.chunk", choices: [{ index: 0, delta: { role: "assistant", content: "Hi" } }] }), + ...sse({ id: "c2", object: "chat.completion.chunk", choices: [{ index: 0, delta: {}, finish_reason: "stop" }] }), + "data: [DONE]", + "", + ])); + + const completed = events.filter((e) => e.event === "response.completed"); + expect(completed).toHaveLength(1); + expect(events.at(-1).event).toBe("response.completed"); + expect(completed[0].data.response).not.toHaveProperty("usage"); + }); +}); diff --git a/tests/unit/openai-responses-usage-pivot.test.js b/tests/unit/openai-responses-usage-pivot.test.js new file mode 100644 index 00000000..c9f96af7 --- /dev/null +++ b/tests/unit/openai-responses-usage-pivot.test.js @@ -0,0 +1,116 @@ +import { describe, expect, it } from "vitest"; + +import { FORMATS } from "../../open-sse/translator/formats.js"; +import { createSSETransformStreamWithLogger } from "../../open-sse/utils/stream.js"; + +/** + * Usage must survive the PIVOT, not just the direct openai:openai-responses route. + * + * Codex talks the Responses API, so routing it at a Claude connection runs + * claude -> openai -> openai-responses. The converter that attaches usage to + * response.completed is the second hop, and it only ever sees the intermediate + * OpenAI chunk — so whether Codex learns its context size depends on the first + * hop putting usage on that intermediate chunk. + * + * Signature is (targetFormat, sourceFormat, ...) — targetFormat is what the + * UPSTREAM speaks, sourceFormat is what the CLIENT speaks. + */ +async function runTransform(chunks, targetFormat, provider) { + const encoder = new TextEncoder(); + const input = chunks.map((c) => `data: ${JSON.stringify(c)}\n\n`).join(""); + + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode(input)); + controller.close(); + }, + }); + + const output = stream.pipeThrough( + createSSETransformStreamWithLogger( + targetFormat, + FORMATS.OPENAI_RESPONSES, + provider, + null, + null, + "claude-sonnet-5", + ), + ); + + const reader = output.getReader(); + const decoder = new TextDecoder(); + let text = ""; + + while (true) { + const { value, done } = await reader.read(); + if (done) break; + text += decoder.decode(value, { stream: true }); + } + + text += decoder.decode(); + return text; +} + +function completedResponse(output) { + const lines = output + .split("\n") + .filter((l) => l.startsWith("data: ") && l.includes('"type":"response.completed"')); + expect(lines.length, "expected exactly one response.completed").toBe(1); + return JSON.parse(lines[0].slice(6)).response; +} + +// Anthropic splits the counts across two events: message_start carries the whole +// prompt side (input + both cache buckets), message_delta carries only the output +// side. Neither event alone is the total, which is why the claude converter merges +// them into state before emitting the intermediate chunk. +const CLAUDE_CHUNKS_WITH_USAGE = [ + { + type: "message_start", + message: { + id: "msg_01CfUtmFqMv3Gc5s66ehaTK", + model: "claude-sonnet-5", + usage: { + input_tokens: 1500, + cache_read_input_tokens: 12000, + cache_creation_input_tokens: 300, + output_tokens: 1, + }, + }, + }, + { type: "content_block_start", index: 0, content_block: { type: "text", text: "" } }, + { type: "content_block_delta", index: 0, delta: { type: "text_delta", text: "hi" } }, + { type: "content_block_stop", index: 0 }, + { type: "message_delta", delta: { stop_reason: "end_turn" }, usage: { output_tokens: 42 } }, + { type: "message_stop" }, +]; + +describe("OpenAI Responses usage across the pivot", () => { + // The reported failure: a Codex session on a Claude connection grew unbounded + // (101 -> 503 -> 631 messages) until Anthropic rejected it with + // "prompt is too long: 1676806 tokens > 1000000 maximum", because every + // token_count event Codex recorded had info: null. + it("reports claude usage on response.completed so Codex can auto-compact", async () => { + const output = await runTransform(CLAUDE_CHUNKS_WITH_USAGE, FORMATS.CLAUDE, "claude"); + + // prompt side = input + cache_read + cache_creation = 1500 + 12000 + 300. + expect(completedResponse(output).usage).toEqual({ + input_tokens: 13800, + output_tokens: 42, + total_tokens: 13842, + input_tokens_details: { cached_tokens: 12000 }, + }); + }); + + // Codex deserializes usage into a struct whose three top-level counts are all + // required, so dropping any one of them discards the whole object and leaves the + // context gauge empty — the same end state as reporting nothing. + it("always reports all three top-level counts", async () => { + const output = await runTransform(CLAUDE_CHUNKS_WITH_USAGE, FORMATS.CLAUDE, "claude"); + const usage = completedResponse(output).usage; + + for (const field of ["input_tokens", "output_tokens", "total_tokens"]) { + expect(usage, `missing ${field}`).toHaveProperty(field); + expect(Number.isFinite(usage[field]), `${field} must be a number`).toBe(true); + } + }); +});