diff --git a/open-sse/translator/response/openai-responses.js b/open-sse/translator/response/openai-responses.js index bd435f9c..c3d16448 100644 --- a/open-sse/translator/response/openai-responses.js +++ b/open-sse/translator/response/openai-responses.js @@ -14,13 +14,47 @@ import { ROLE, OPENAI_BLOCK, RESPONSES_ITEM, OPENAI_FINISH, MODEL_FALLBACK } fro * Translate OpenAI chunk to Responses API events * @returns {Array} Array of events with { event, data } structure */ +// Upstream Chat Completions usage -> Responses API usage shape. +// Without this, /v1/responses never reports usage: Responses clients (Codex CLI) +// keep their "context used" gauge pinned at 0 and never auto-compact, so a long +// session grows until the upstream context limit rejects it (9router issue #3432). +// +// Note this is stored under state.responsesUsage, NOT state.usage: state.usage is +// owned by the stream layer, which fills it with normalizeUsage()-shaped counts +// (prompt_tokens/prompt_tokens_details) and hands it to finalizeStream() for +// logging and cost accounting. Overwriting it with this shape silently drops +// cached/reasoning tokens from those stats. +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 responseUsage = { + input_tokens: inputTokens, + output_tokens: outputTokens, + total_tokens: Number.isFinite(usage.total_tokens) ? usage.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 }; + + return responseUsage; +} + export function openaiToOpenAIResponsesResponse(chunk, state) { if (!chunk) { 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); + } + if (!chunk.choices?.length) return []; - + const events = []; const nextSeq = () => ++state.seq; @@ -112,7 +146,19 @@ export function openaiToOpenAIResponsesResponse(chunk, state) { for (const i in state.msgItemAdded) closeMessage(state, emit, i); closeReasoning(state, emit); for (const i in state.funcCallIds) closeToolCall(state, emit, i); - sendCompleted(state, emit); + // Upstreams report usage either on the finish chunk itself or on a trailing chunk + // whose `choices` array is empty (OpenAI does the latter). Emitting + // response.completed here would freeze the payload before that trailing chunk is + // parsed, so when usage is not known yet we leave completion to flushEvents(), + // which runs once the upstream stream ends and by then has seen every chunk. + // + // That only holds on the direct openai:openai-responses route. When this converter + // runs as the second hop of a pivot (Claude/Gemini/Kiro upstream), translateResponse() + // drops the terminal null chunk before reaching us — the first hop returns null for + // it, leaving nothing to iterate — so flushEvents() is never called and deferring + // would swallow the terminal event entirely. Keep the old behaviour there. + const flushReachesUs = state.targetFormat === FORMATS.OPENAI; + if (state.responsesUsage || !flushReachesUs) sendCompleted(state, emit); } return events; @@ -376,7 +422,8 @@ function sendCompleted(state, emit) { created_at: state.created, status: "completed", background: false, - error: null + error: null, + ...(state.responsesUsage ? { usage: state.responsesUsage } : {}) } }); } diff --git a/open-sse/utils/stream.js b/open-sse/utils/stream.js index 15ee37d8..cc0c6e30 100644 --- a/open-sse/utils/stream.js +++ b/open-sse/utils/stream.js @@ -60,7 +60,13 @@ export function createSSEStream(options = {}) { const decoder = new TextDecoder("utf-8", { fatal: false }); const state = mode === STREAM_MODE.TRANSLATE - ? { ...initState(sourceFormat), provider, toolNameMap, customToolNames: new Set(customToolNames || []), model, sessionId: credentials?._clientSessionId || null } + ? { ...initState(sourceFormat), provider, toolNameMap, customToolNames: new Set(customToolNames || []), model, sessionId: credentials?._clientSessionId || null, + // Which upstream format this stream came from. A response translator can be + // reached either directly (target === its registered source) or as the second + // hop of a pivot, and on the terminal null chunk the pivot drops it — so a + // translator that defers closing events until flush needs to know which case + // it is in. Absent/undefined means "unknown", i.e. do not defer. + targetFormat } : null; let totalContentLength = 0; diff --git a/tests/unit/openai-responses-usage-completed.test.js b/tests/unit/openai-responses-usage-completed.test.js new file mode 100644 index 00000000..fe750d02 --- /dev/null +++ b/tests/unit/openai-responses-usage-completed.test.js @@ -0,0 +1,162 @@ +import { describe, expect, it } from "vitest"; + +import { FORMATS } from "../../open-sse/translator/formats.js"; +import { createSSETransformStreamWithLogger } from "../../open-sse/utils/stream.js"; + +/** + * Upstream chunks -> client Responses API events. + * + * The converter under test is openaiToOpenAIResponsesResponse(), reached through + * the registered OPENAI:OPENAI_RESPONSES pair. Without it, /v1/responses never + * reports usage and Responses clients (Codex CLI) keep their context gauge at 0, + * so they never auto-compact and eventually hit the upstream context limit. + * + * Signature is (targetFormat, sourceFormat, ...) — targetFormat is what the + * UPSTREAM speaks, sourceFormat is what the CLIENT speaks. + */ +async function runTransform(chunks, targetFormat = FORMATS.OPENAI) { + 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, + "deepseek", + null, + null, + "deepseek-flash", + ), + ); + + 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 completedEvents(output) { + return output + .split("\n") + .filter((l) => l.startsWith("data: ") && l.includes('"type":"response.completed"')); +} + +function completedResponse(output) { + const lines = completedEvents(output); + expect(lines.length, "expected exactly one response.completed").toBe(1); + return JSON.parse(lines[0].slice(6)).response; +} + +const TEXT_CHUNK = { + id: "chatcmpl-test", + object: "chat.completion.chunk", + created: 1700000000, + model: "deepseek-flash", + choices: [{ index: 0, delta: { role: "assistant", content: "好" } }], +}; + +const FINISH_CHUNK = { + id: "chatcmpl-test", + object: "chat.completion.chunk", + created: 1700000000, + model: "deepseek-flash", + choices: [{ index: 0, delta: {}, finish_reason: "stop" }], +}; + +// Usage-only trailer: `choices` is empty, exactly as OpenAI emits it when +// stream_options.include_usage is set. +const USAGE_ONLY_CHUNK = { + id: "chatcmpl-test", + object: "chat.completion.chunk", + created: 1700000000, + model: "deepseek-flash", + choices: [], + usage: { + prompt_tokens: 884, + completion_tokens: 37, + total_tokens: 921, + prompt_tokens_details: { cached_tokens: 256 }, + }, +}; + +const EXPECTED_USAGE = { + input_tokens: 884, + output_tokens: 37, + total_tokens: 921, + input_tokens_details: { cached_tokens: 256 }, +}; + +// Claude-shaped stream with NO usage anywhere: the only way the client gets a +// terminal event is the finish_reason branch, because the pivot never reaches +// flushEvents() with the terminal null chunk. +const CLAUDE_CHUNKS = [ + { type: "message_start", message: { id: "msg_1", model: "claude-x" } }, + { 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" } }, + { type: "message_stop" }, +]; + +describe("OpenAI Responses usage on response.completed", () => { + it("maps usage reported on the finish chunk", async () => { + const output = await runTransform([ + TEXT_CHUNK, + { + ...FINISH_CHUNK, + usage: { + prompt_tokens: 884, + completion_tokens: 37, + total_tokens: 921, + prompt_tokens_details: { cached_tokens: 256 }, + completion_tokens_details: { reasoning_tokens: 12 }, + }, + }, + ]); + + expect(completedResponse(output).usage).toEqual({ + ...EXPECTED_USAGE, + output_tokens_details: { reasoning_tokens: 12 }, + }); + }); + + it("maps usage reported on a trailing usage-only chunk with empty choices", async () => { + const output = await runTransform([TEXT_CHUNK, FINISH_CHUNK, USAGE_ONLY_CHUNK]); + + expect(completedResponse(output).usage).toEqual(EXPECTED_USAGE); + }); + + it("still completes when the upstream reports no usage at all", async () => { + const output = await runTransform([TEXT_CHUNK, FINISH_CHUNK]); + + const response = completedResponse(output); + expect(response.status).toBe("completed"); + expect(response).not.toHaveProperty("usage"); + }); + + // Regression guard for the pivot: with a Claude upstream the converter runs as + // the second hop, translateResponse() drops the terminal null chunk before it + // reaches this converter, so flushEvents() never runs. Deferring completion + // there would leave the client without any terminal event. + it("completes on a pivoted stream whose upstream never reports usage", async () => { + const output = await runTransform(CLAUDE_CHUNKS, FORMATS.CLAUDE); + + const response = completedResponse(output); + expect(response.status).toBe("completed"); + }); +});