fix(ai): require OpenAI Responses terminal events

This commit is contained in:
Mario Zechner
2026-06-23 16:35:45 +02:00
parent 2285f87964
commit cd95c2749f
12 changed files with 378 additions and 39 deletions
+1
View File
@@ -67,6 +67,7 @@ Migration guide:
### Fixed
- Fixed OpenAI Responses streams to fail when they end before a terminal response event and to treat `response.incomplete` as a length stop ([#5526](https://github.com/earendil-works/pi/pull/5526) by [@dmmulroy](https://github.com/dmmulroy)).
- Fixed Amazon Bedrock endpoint resolution to honor scoped `AWS_PROFILE` values.
- Fixed Cloudflare providers to require account/gateway configuration and route built-in `/compat` requests through provider auth.
- Fixed OpenAI Codex Responses WebSocket sessions to reconnect once when OpenAI's connection limit is reached before output starts ([#5973](https://github.com/earendil-works/pi/issues/5973)).
+39 -29
View File
@@ -294,8 +294,41 @@ export async function processResponsesStream<TApi extends Api>(
): Promise<void> {
let currentItem: ResponseReasoningItem | ResponseOutputMessage | ResponseFunctionToolCall | null = null;
let currentBlock: ThinkingContent | TextContent | (ToolCall & { partialJson: string }) | null = null;
let sawTerminalResponseEvent = false;
const blocks = output.content;
const blockIndex = () => blocks.length - 1;
const finalizeResponse = (
response: Extract<ResponseStreamEvent, { type: "response.completed" | "response.incomplete" }>["response"],
): void => {
sawTerminalResponseEvent = true;
if (response?.id) {
output.responseId = response.id;
}
if (response?.usage) {
const cachedTokens = response.usage.input_tokens_details?.cached_tokens || 0;
output.usage = {
// OpenAI includes cached tokens in input_tokens, so subtract to get non-cached input
input: (response.usage.input_tokens || 0) - cachedTokens,
output: response.usage.output_tokens || 0,
cacheRead: cachedTokens,
cacheWrite: 0,
totalTokens: response.usage.total_tokens || 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
}
calculateCost(model, output.usage);
if (options?.applyServiceTierPricing) {
const serviceTier = options.resolveServiceTier
? options.resolveServiceTier(response?.service_tier, options.serviceTier)
: (response?.service_tier ?? options.serviceTier);
options.applyServiceTierPricing(output.usage, serviceTier);
}
// Map status to stop reason
output.stopReason = mapStopReason(response?.status);
if (output.content.some((b) => b.type === "toolCall") && output.stopReason === "stop") {
output.stopReason = "toolUse";
}
};
for await (const event of openaiStream) {
if (event.type === "response.created") {
@@ -491,38 +524,12 @@ export async function processResponsesStream<TApi extends Api>(
currentBlock = null;
stream.push({ type: "toolcall_end", contentIndex: blockIndex(), toolCall, partial: output });
}
} else if (event.type === "response.completed") {
const response = event.response;
if (response?.id) {
output.responseId = response.id;
}
if (response?.usage) {
const cachedTokens = response.usage.input_tokens_details?.cached_tokens || 0;
output.usage = {
// OpenAI includes cached tokens in input_tokens, so subtract to get non-cached input
input: (response.usage.input_tokens || 0) - cachedTokens,
output: response.usage.output_tokens || 0,
cacheRead: cachedTokens,
cacheWrite: 0,
totalTokens: response.usage.total_tokens || 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
};
}
calculateCost(model, output.usage);
if (options?.applyServiceTierPricing) {
const serviceTier = options.resolveServiceTier
? options.resolveServiceTier(response?.service_tier, options.serviceTier)
: (response?.service_tier ?? options.serviceTier);
options.applyServiceTierPricing(output.usage, serviceTier);
}
// Map status to stop reason
output.stopReason = mapStopReason(response?.status);
if (output.content.some((b) => b.type === "toolCall") && output.stopReason === "stop") {
output.stopReason = "toolUse";
}
} else if (event.type === "response.completed" || event.type === "response.incomplete") {
finalizeResponse(event.response);
} else if (event.type === "error") {
throw new Error(`Error Code ${event.code}: ${event.message}` || "Unknown error");
} else if (event.type === "response.failed") {
sawTerminalResponseEvent = true;
const error = event.response?.error;
const details = event.response?.incomplete_details;
const msg = error
@@ -533,6 +540,9 @@ export async function processResponsesStream<TApi extends Api>(
throw new Error(msg);
}
}
if (!sawTerminalResponseEvent) {
throw new Error("OpenAI Responses stream ended before a terminal response event");
}
}
function mapStopReason(status: OpenAI.Responses.ResponseStatus | undefined): StopReason {
@@ -57,6 +57,11 @@ async function* createFunctionCallEvents(argumentsJson: string): AsyncIterable<R
arguments: argumentsJson,
},
} as ResponseStreamEvent;
yield {
type: "response.completed",
sequence_number: 5,
response: { id: "resp_test", status: "completed" },
} as ResponseStreamEvent;
}
describe("openai responses partialJson cleanup", () => {
@@ -0,0 +1,233 @@
import type { ResponseStreamEvent } from "openai/resources/responses/responses.js";
import { describe, expect, it, vi } from "vitest";
import { stream as streamOpenAIResponses } from "../src/api/openai-responses.ts";
import { processResponsesStream } from "../src/api/openai-responses-shared.ts";
import type { AssistantMessage, AssistantMessageEvent, Context, Model } from "../src/types.ts";
import { AssistantMessageEventStream } from "../src/utils/event-stream.ts";
vi.mock("openai", () => {
async function* createMockResponsesStream(): AsyncIterable<ResponseStreamEvent> {
yield {
type: "response.created",
sequence_number: 0,
response: { id: "resp_wrapper_early_eof" },
} as ResponseStreamEvent;
yield {
type: "response.output_item.added",
sequence_number: 1,
output_index: 0,
item: { type: "reasoning", id: "rs_wrapper_early_eof", summary: [] },
} as ResponseStreamEvent;
yield {
type: "response.reasoning_text.delta",
sequence_number: 2,
output_index: 0,
content_index: 0,
item_id: "rs_wrapper_early_eof",
delta: "partial reasoning before the wrapper stream ends",
} as ResponseStreamEvent;
}
class FakeOpenAI {
responses = {
create: () => {
const responseStream = createMockResponsesStream();
const promise = Promise.resolve(responseStream) as Promise<AsyncIterable<ResponseStreamEvent>> & {
withResponse: () => Promise<{
data: AsyncIterable<ResponseStreamEvent>;
response: { status: number; headers: Headers };
}>;
};
promise.withResponse = async () => ({
data: responseStream,
response: { status: 200, headers: new Headers() },
});
return promise;
},
};
}
return { default: FakeOpenAI };
});
function createModel(): Model<"openai-responses"> {
return {
id: "gpt-5-mini",
name: "GPT-5 Mini",
api: "openai-responses",
provider: "openai",
baseUrl: "https://api.openai.com/v1",
reasoning: true,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 400000,
maxTokens: 128000,
};
}
function createOutput(model: Model<"openai-responses">): AssistantMessage {
return {
role: "assistant",
content: [],
api: model.api,
provider: model.provider,
model: model.id,
usage: {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "stop",
timestamp: Date.now(),
};
}
async function* createEarlyEofEvents(): AsyncIterable<ResponseStreamEvent> {
yield {
type: "response.created",
sequence_number: 0,
response: { id: "resp_early_eof" },
} as ResponseStreamEvent;
yield {
type: "response.output_item.added",
sequence_number: 1,
output_index: 0,
item: { type: "reasoning", id: "rs_early_eof", summary: [] },
} as ResponseStreamEvent;
yield {
type: "response.reasoning_text.delta",
sequence_number: 2,
output_index: 0,
content_index: 0,
item_id: "rs_early_eof",
delta: "partial reasoning before the stream ends",
} as ResponseStreamEvent;
}
async function* createCompletedEvents(): AsyncIterable<ResponseStreamEvent> {
yield {
type: "response.completed",
sequence_number: 0,
response: {
id: "resp_completed",
status: "completed",
usage: {
input_tokens: 20,
output_tokens: 7,
total_tokens: 27,
input_tokens_details: { cached_tokens: 2 },
},
},
} as ResponseStreamEvent;
}
async function* createIncompleteEvents(): AsyncIterable<ResponseStreamEvent> {
yield {
type: "response.incomplete",
sequence_number: 0,
response: {
id: "resp_incomplete",
status: "incomplete",
usage: {
input_tokens: 30,
output_tokens: 12,
total_tokens: 42,
input_tokens_details: { cached_tokens: 5 },
},
},
} as ResponseStreamEvent;
}
async function* createFailedEvents(): AsyncIterable<ResponseStreamEvent> {
yield {
type: "response.failed",
sequence_number: 0,
response: {
id: "resp_failed",
status: "failed",
error: { code: "server_error", message: "boom" },
},
} as ResponseStreamEvent;
}
describe("OpenAI Responses terminal event handling", () => {
it("rejects streams that end before a terminal response event", async () => {
const model = createModel();
const output = createOutput(model);
const stream = new AssistantMessageEventStream();
await expect(processResponsesStream(createEarlyEofEvents(), output, stream, model)).rejects.toThrow(
"OpenAI Responses stream ended before a terminal response event",
);
});
it("emits an error final result when the wrapper stream ends before a terminal response event", async () => {
const model = createModel();
const context: Context = {
systemPrompt: "",
messages: [{ role: "user", content: [{ type: "text", text: "hi" }], timestamp: 0 }],
tools: [],
};
const stream = streamOpenAIResponses(model, context, { apiKey: "test" });
const events: AssistantMessageEvent[] = [];
for await (const event of stream) {
events.push(event);
}
const result = await stream.result();
const lastEvent = events.at(-1);
expect(lastEvent?.type).toBe("error");
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toBe("OpenAI Responses stream ended before a terminal response event");
});
it("finalizes completed terminal events as stop", async () => {
const model = createModel();
const output = createOutput(model);
const stream = new AssistantMessageEventStream();
await processResponsesStream(createCompletedEvents(), output, stream, model);
expect(output.responseId).toBe("resp_completed");
expect(output.stopReason).toBe("stop");
expect(output.usage).toMatchObject({
input: 18,
output: 7,
cacheRead: 2,
cacheWrite: 0,
totalTokens: 27,
});
});
it("finalizes incomplete terminal events as length stops", async () => {
const model = createModel();
const output = createOutput(model);
const stream = new AssistantMessageEventStream();
await processResponsesStream(createIncompleteEvents(), output, stream, model);
expect(output.responseId).toBe("resp_incomplete");
expect(output.stopReason).toBe("length");
expect(output.usage).toMatchObject({
input: 25,
output: 12,
cacheRead: 5,
cacheWrite: 0,
totalTokens: 42,
});
});
it("rejects failed terminal events with the provider error", async () => {
const model = createModel();
const output = createOutput(model);
const stream = new AssistantMessageEventStream();
await expect(processResponsesStream(createFailedEvents(), output, stream, model)).rejects.toThrow(
"server_error: boom",
);
});
});