fix(responses): report usage on response.completed so clients can auto-compact
Map upstream Chat Completions usage to the Responses API shape and attach it to response.completed. Capture chunk.usage before the empty-choices guard so the usage-only trailer chunk survives, and defer completion to flushEvents() when usage is not yet known — only on the direct openai:openai-responses route, since a pivoted stream never reaches flushEvents. Fixes #3432.
This commit is contained in:
@@ -14,11 +14,45 @@ import { ROLE, OPENAI_BLOCK, RESPONSES_ITEM, OPENAI_FINISH, MODEL_FALLBACK } fro
|
|||||||
* Translate OpenAI chunk to Responses API events
|
* Translate OpenAI chunk to Responses API events
|
||||||
* @returns {Array} Array of events with { event, data } structure
|
* @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) {
|
export function openaiToOpenAIResponsesResponse(chunk, state) {
|
||||||
if (!chunk) {
|
if (!chunk) {
|
||||||
return flushEvents(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);
|
||||||
|
}
|
||||||
|
|
||||||
if (!chunk.choices?.length) return [];
|
if (!chunk.choices?.length) return [];
|
||||||
|
|
||||||
const events = [];
|
const events = [];
|
||||||
@@ -112,7 +146,19 @@ export function openaiToOpenAIResponsesResponse(chunk, state) {
|
|||||||
for (const i in state.msgItemAdded) closeMessage(state, emit, i);
|
for (const i in state.msgItemAdded) closeMessage(state, emit, i);
|
||||||
closeReasoning(state, emit);
|
closeReasoning(state, emit);
|
||||||
for (const i in state.funcCallIds) closeToolCall(state, emit, i);
|
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;
|
return events;
|
||||||
@@ -376,7 +422,8 @@ function sendCompleted(state, emit) {
|
|||||||
created_at: state.created,
|
created_at: state.created,
|
||||||
status: "completed",
|
status: "completed",
|
||||||
background: false,
|
background: false,
|
||||||
error: null
|
error: null,
|
||||||
|
...(state.responsesUsage ? { usage: state.responsesUsage } : {})
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -60,7 +60,13 @@ export function createSSEStream(options = {}) {
|
|||||||
const decoder = new TextDecoder("utf-8", { fatal: false });
|
const decoder = new TextDecoder("utf-8", { fatal: false });
|
||||||
|
|
||||||
const state = mode === STREAM_MODE.TRANSLATE
|
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;
|
: null;
|
||||||
|
|
||||||
let totalContentLength = 0;
|
let totalContentLength = 0;
|
||||||
|
|||||||
162
tests/unit/openai-responses-usage-completed.test.js
Normal file
162
tests/unit/openai-responses-usage-completed.test.js
Normal file
@@ -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");
|
||||||
|
});
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user