fix(responses): wait for real usage before emitting response.completed
This commit is contained in:
1 parent
5e9bd464f2
commit
7111db3598
5 files changed
+376
-13
No files matched your search
@@ -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
|
||||
};
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in new issue
Block a user