mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-09-21 14:22:14 +03:00
fix(sse): enforce strict OpenAI Responses API streaming schema compliance (#13956)
* fix(sse): include output in response.in_progress event * fix(sse): include output in translator response.in_progress event * fix(sse): include status in all responses output items * fix(sse): include sequence_number in Responses API error frames * docs(changelog): expand responses streaming schema compliance notes * fix(sse): always include token details in Responses API usage --------- Co-authored-by: TheDemonTuan <nguyenviettuanbp@gmail.com>
This commit is contained in:
5
changelog.d/fixes/responses-in-progress-output.md
Normal file
5
changelog.d/fixes/responses-in-progress-output.md
Normal file
@@ -0,0 +1,5 @@
|
||||
- **fix(responses):** ensure full compliance with the OpenAI Responses API streaming schema for strict deserializers (e.g. OpenAI Responses SDK, Grok CLI / pager):
|
||||
- Include `output: []`, `background: false`, and `error: null` in the `response.in_progress` lifecycle event across both the Responses transformer and response translator.
|
||||
- Include `status` (`in_progress` or `completed`) on all emitted output items (`message`, `reasoning`, `function_call`, `custom_tool_call`) in `response.output_item.added`, `response.output_item.done`, and `response.output[]`.
|
||||
- Include `sequence_number: 0` in in-band Responses stream error frames (`OPENAI_RESPONSES_ERROR_FRAME` and `buildResponsesErrorDataLine`) emitted after early keepalive streams commit.
|
||||
- Ensure `input_tokens_details` (with `cached_tokens: 0`) and `output_tokens_details` (with `reasoning_tokens: 0`) are always populated in `response.usage` even when upstreams (e.g. Gemini) omit reasoning or caching tokens.
|
||||
@@ -297,6 +297,7 @@ export function createResponsesApiTransformStream(
|
||||
id: state.reasoningId,
|
||||
type: "reasoning",
|
||||
summary: [],
|
||||
status: "in_progress",
|
||||
},
|
||||
});
|
||||
|
||||
@@ -347,6 +348,7 @@ export function createResponsesApiTransformStream(
|
||||
id: state.reasoningId,
|
||||
type: "reasoning",
|
||||
summary: [{ type: "summary_text", text: state.reasoningBuf }],
|
||||
status: "completed",
|
||||
};
|
||||
|
||||
emit(controller, "response.output_item.done", {
|
||||
@@ -388,6 +390,7 @@ export function createResponsesApiTransformStream(
|
||||
type: "message",
|
||||
content: [{ type: "output_text", annotations: [], logprobs: [], text: fullText }],
|
||||
role: "assistant",
|
||||
status: "completed",
|
||||
};
|
||||
|
||||
emit(controller, "response.output_item.done", {
|
||||
@@ -438,7 +441,7 @@ export function createResponsesApiTransformStream(
|
||||
...(customTool ? { input: "" } : { arguments: "" }),
|
||||
call_id: state.funcCallIds[idx],
|
||||
name: state.funcNames[idx] || "",
|
||||
...(customTool ? { status: "in_progress" } : {}),
|
||||
status: "in_progress",
|
||||
},
|
||||
});
|
||||
return true;
|
||||
@@ -519,6 +522,7 @@ export function createResponsesApiTransformStream(
|
||||
arguments: args,
|
||||
call_id: callId,
|
||||
name: toolName,
|
||||
status: "completed",
|
||||
};
|
||||
}
|
||||
|
||||
@@ -678,6 +682,9 @@ export function createResponsesApiTransformStream(
|
||||
object: "response",
|
||||
created_at: state.created,
|
||||
status: "in_progress",
|
||||
background: false,
|
||||
error: null,
|
||||
output: [],
|
||||
},
|
||||
});
|
||||
}
|
||||
@@ -763,7 +770,13 @@ export function createResponsesApiTransformStream(
|
||||
emit(controller, "response.output_item.added", {
|
||||
type: "response.output_item.added",
|
||||
output_index: msgIdx,
|
||||
item: { id: msgId, type: "message", content: [], role: "assistant" },
|
||||
item: {
|
||||
id: msgId,
|
||||
type: "message",
|
||||
content: [],
|
||||
role: "assistant",
|
||||
status: "in_progress",
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -171,20 +171,24 @@ export function openaiToOpenAIResponsesResponse(chunk, state) {
|
||||
const u = chunk.usage;
|
||||
const input_tokens = u.input_tokens ?? u.prompt_tokens ?? 0;
|
||||
const output_tokens = u.output_tokens ?? u.completion_tokens ?? 0;
|
||||
const cacheDetails = resolveResponsesCacheUsageDetails(u);
|
||||
const rawReasoning =
|
||||
u.output_tokens_details?.reasoning_tokens ?? u.completion_tokens_details?.reasoning_tokens;
|
||||
const reasoningTokens =
|
||||
typeof rawReasoning === "number" && Number.isFinite(rawReasoning) ? rawReasoning : 0;
|
||||
|
||||
state.usage = {
|
||||
input_tokens,
|
||||
input_tokens_details: {
|
||||
cached_tokens: 0,
|
||||
...(cacheDetails || {}),
|
||||
},
|
||||
output_tokens,
|
||||
output_tokens_details: {
|
||||
reasoning_tokens: reasoningTokens,
|
||||
},
|
||||
total_tokens: u.total_tokens ?? input_tokens + output_tokens,
|
||||
};
|
||||
const cacheDetails = resolveResponsesCacheUsageDetails(u);
|
||||
if (cacheDetails) {
|
||||
state.usage.input_tokens_details = cacheDetails;
|
||||
}
|
||||
const reasoningTokens =
|
||||
u.output_tokens_details?.reasoning_tokens ?? u.completion_tokens_details?.reasoning_tokens;
|
||||
if (reasoningTokens) {
|
||||
state.usage.output_tokens_details = { reasoning_tokens: reasoningTokens };
|
||||
}
|
||||
}
|
||||
|
||||
if (!chunk.choices?.length) {
|
||||
@@ -268,6 +272,9 @@ export function openaiToOpenAIResponsesResponse(chunk, state) {
|
||||
object: "response",
|
||||
created_at: state.created,
|
||||
status: "in_progress",
|
||||
background: false,
|
||||
error: null,
|
||||
output: [],
|
||||
};
|
||||
if (state.model) inProgressResponse.model = state.model;
|
||||
emit("response.in_progress", {
|
||||
@@ -395,7 +402,7 @@ function startReasoning(state, emit, idx) {
|
||||
emit("response.output_item.added", {
|
||||
type: "response.output_item.added",
|
||||
output_index: idx,
|
||||
item: { id: state.reasoningId, type: "reasoning", summary: [] },
|
||||
item: { id: state.reasoningId, type: "reasoning", summary: [], status: "in_progress" },
|
||||
});
|
||||
|
||||
emit("response.reasoning_summary_part.added", {
|
||||
@@ -445,6 +452,7 @@ function closeReasoning(state, emit) {
|
||||
id: state.reasoningId,
|
||||
type: "reasoning",
|
||||
summary: [{ type: "summary_text", text: state.reasoningBuf }],
|
||||
status: "completed",
|
||||
};
|
||||
|
||||
emit("response.output_item.done", {
|
||||
@@ -465,7 +473,7 @@ function emitTextContent(state, emit, idx, content) {
|
||||
emit("response.output_item.added", {
|
||||
type: "response.output_item.added",
|
||||
output_index: idx,
|
||||
item: { id: msgId, type: "message", content: [], role: "assistant" },
|
||||
item: { id: msgId, type: "message", content: [], role: "assistant", status: "in_progress" },
|
||||
});
|
||||
}
|
||||
|
||||
@@ -523,6 +531,7 @@ function closeMessage(state, emit, idx) {
|
||||
type: "message",
|
||||
content: [{ type: "output_text", annotations: [], logprobs: [], text: fullText }],
|
||||
role: "assistant",
|
||||
status: "completed",
|
||||
};
|
||||
|
||||
emit("response.output_item.done", {
|
||||
|
||||
@@ -86,6 +86,7 @@ export const OPENAI_RESPONSES_ERROR_FRAME = ENCODER.encode(
|
||||
code: null,
|
||||
message: "Upstream stream failed before completion.",
|
||||
param: null,
|
||||
sequence_number: 0,
|
||||
})}\n\n`
|
||||
);
|
||||
|
||||
@@ -124,7 +125,7 @@ function buildResponsesErrorDataLine(text: string): string {
|
||||
parsed && typeof parsed.diagnostics === "object" && parsed.diagnostics !== null
|
||||
? { diagnostics: parsed.diagnostics }
|
||||
: {};
|
||||
return JSON.stringify({ type: "error", code, message, param, ...extras });
|
||||
return JSON.stringify({ type: "error", code, message, param, sequence_number: 0, ...extras });
|
||||
}
|
||||
|
||||
export type EarlyStreamKeepaliveOptions = {
|
||||
|
||||
@@ -7,9 +7,8 @@ import {
|
||||
echoModelInSseLine,
|
||||
} from "../../open-sse/services/responseModelEcho.ts";
|
||||
|
||||
const { openaiToOpenAIResponsesResponse } = await import(
|
||||
"../../open-sse/translator/response/openai-responses.ts"
|
||||
);
|
||||
const { openaiToOpenAIResponsesResponse } =
|
||||
await import("../../open-sse/translator/response/openai-responses.ts");
|
||||
const { initState } = await import("../../open-sse/translator/index.ts");
|
||||
const { FORMATS } = await import("../../open-sse/translator/formats.ts");
|
||||
|
||||
@@ -82,6 +81,68 @@ test("OpenAI -> Responses translator omits model when the upstream never sent on
|
||||
assert.equal("model" in (completed!.data.response as Record<string, unknown>), false);
|
||||
});
|
||||
|
||||
test("OpenAI -> Responses translator emits response.in_progress with output: [], background: false, error: null", () => {
|
||||
const events = collectResponsesEvents([
|
||||
{
|
||||
id: "chatcmpl-1",
|
||||
model: "gpt-5.5",
|
||||
choices: [{ index: 0, delta: { content: "hi" }, finish_reason: null }],
|
||||
},
|
||||
null,
|
||||
]);
|
||||
|
||||
const inProgress = events.find((e) => e.event === "response.in_progress");
|
||||
assert.ok(inProgress, "response.in_progress must be emitted");
|
||||
const resp = inProgress!.data.response as Record<string, unknown>;
|
||||
assert.ok(Array.isArray(resp.output), "output must be an array");
|
||||
assert.deepEqual(resp.output, []);
|
||||
assert.equal(resp.background, false);
|
||||
assert.equal(resp.error, null);
|
||||
|
||||
const addedItem = events.find((e) => e.event === "response.output_item.added")?.data
|
||||
.item as Record<string, unknown>;
|
||||
assert.ok(addedItem, "output_item.added must exist");
|
||||
assert.equal(addedItem.status, "in_progress");
|
||||
|
||||
const completed = events.find((e) => e.event === "response.completed")?.data.response as Record<
|
||||
string,
|
||||
unknown
|
||||
>;
|
||||
assert.ok(completed, "response.completed must exist");
|
||||
const completedOutput = completed.output as Array<Record<string, unknown>>;
|
||||
assert.equal(completedOutput[0].status, "completed");
|
||||
});
|
||||
|
||||
test("OpenAI -> Responses translator always populates input_tokens_details and output_tokens_details", () => {
|
||||
const events = collectResponsesEvents([
|
||||
{
|
||||
id: "chatcmpl-1",
|
||||
model: "gemini-3.8-flash",
|
||||
choices: [{ index: 0, delta: { content: "hi" }, finish_reason: null }],
|
||||
usage: { prompt_tokens: 10, completion_tokens: 5, total_tokens: 15 },
|
||||
},
|
||||
{
|
||||
id: "chatcmpl-1",
|
||||
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
|
||||
},
|
||||
null,
|
||||
]);
|
||||
|
||||
const completed = events.find((e) => e.event === "response.completed")?.data?.response as Record<
|
||||
string,
|
||||
unknown
|
||||
>;
|
||||
assert.ok(completed.usage, "usage must be present");
|
||||
const usage = completed.usage as Record<string, unknown>;
|
||||
assert.equal(usage.input_tokens, 10);
|
||||
assert.equal(usage.output_tokens, 5);
|
||||
assert.equal(usage.total_tokens, 15);
|
||||
assert.ok(usage.input_tokens_details, "input_tokens_details must be present");
|
||||
assert.ok(usage.output_tokens_details, "output_tokens_details must be present");
|
||||
assert.deepEqual(usage.input_tokens_details, { cached_tokens: 0 });
|
||||
assert.deepEqual(usage.output_tokens_details, { reasoning_tokens: 0 });
|
||||
});
|
||||
|
||||
test("full shim pipeline: bare upstream model in Responses payloads gets rewritten to the requested effort-suffixed id", () => {
|
||||
const events = collectResponsesEvents([
|
||||
{
|
||||
|
||||
@@ -82,6 +82,7 @@ test("Responses route: post-keepalive JSON error body must carry a `type` field
|
||||
`instead of surfacing the real upstream error.`
|
||||
);
|
||||
assert.equal(lastPayload.type, "error");
|
||||
assert.equal(typeof lastPayload.sequence_number, "number");
|
||||
assert.equal(lastPayload.message, 'Unknown name "encrypted" ... Cannot find field.');
|
||||
assert.equal(lastPayload.code, "bad_request");
|
||||
});
|
||||
@@ -109,6 +110,7 @@ test("Responses route: non-JSON/empty post-keepalive error body falls back to a
|
||||
const lastPayload = lastDataPayload(await readAll(result));
|
||||
|
||||
assert.equal(lastPayload.type, "error");
|
||||
assert.equal(typeof lastPayload.sequence_number, "number");
|
||||
assert.ok(
|
||||
typeof lastPayload.message === "string" && lastPayload.message.length > 0,
|
||||
`fallback frame must never be opaque/empty; got ${JSON.stringify(lastPayload)}`
|
||||
|
||||
@@ -72,9 +72,22 @@ test("createResponsesApiTransformStream converts plain chat deltas into Response
|
||||
);
|
||||
assert.ok(types.includes("response.created"));
|
||||
assert.ok(types.includes("response.in_progress"));
|
||||
|
||||
const inProgress = JSON.parse(
|
||||
events.find((event) => event.event === "response.in_progress").data
|
||||
).response;
|
||||
assert.ok(Array.isArray(inProgress.output), "response.in_progress must include an output array");
|
||||
assert.deepEqual(inProgress.output, []);
|
||||
|
||||
assert.ok(types.includes("response.output_item.added"));
|
||||
const addedItem = JSON.parse(
|
||||
events.find((event) => event.event === "response.output_item.added").data
|
||||
).item;
|
||||
assert.equal(addedItem.status, "in_progress");
|
||||
|
||||
assert.ok(types.includes("response.output_text.done"));
|
||||
assert.equal(completed.output[0].content[0].text, "Hello");
|
||||
assert.equal(completed.output[0].status, "completed");
|
||||
assert.deepEqual(completed.usage, {
|
||||
input_tokens: 1,
|
||||
input_tokens_details: { cached_tokens: 0 },
|
||||
|
||||
@@ -46,7 +46,13 @@ test("BUG #6906: live translator — response.completed carries usage when the u
|
||||
assert.ok(completedEvent, "response.completed event should be emitted");
|
||||
assert.deepEqual(
|
||||
completedEvent.data.response.usage,
|
||||
{ input_tokens: 2249, output_tokens: 123, total_tokens: 2372 },
|
||||
{
|
||||
input_tokens: 2249,
|
||||
input_tokens_details: { cached_tokens: 0 },
|
||||
output_tokens: 123,
|
||||
output_tokens_details: { reasoning_tokens: 0 },
|
||||
total_tokens: 2372,
|
||||
},
|
||||
"response.completed must carry usage even when the usage-only chunk trails finish_reason"
|
||||
);
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user