Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 25 additions & 6 deletions src/server/responses/combo-stream-preflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,24 @@ const TERMINAL_EVENTS = new Set([
"response.incomplete",
]);

const RETRYABLE_ZERO_OUTPUT_INCOMPLETE_REASONS = new Set([
"adapter_eof",
"missing_terminal_event",
"upstream_stall_timeout",
]);

function retryableZeroOutputTerminal(payload: unknown): boolean {
if (!payload || typeof payload !== "object" || Array.isArray(payload)) return false;
const event = payload as {
type?: unknown;
response?: { incomplete_details?: { reason?: unknown } };
};
if (event.type === "response.failed") return true;
if (event.type !== "response.incomplete") return false;
const reason = event.response?.incomplete_details?.reason;
return typeof reason === "string" && RETRYABLE_ZERO_OUTPUT_INCOMPLETE_REASONS.has(reason);
}

/**
* Decide when replaying the request on another combo target would risk duplicating
* client-visible output or a tool-side effect. Unknown event types commit the child
Expand Down Expand Up @@ -128,14 +146,14 @@ export async function preflightComboStreamResponse(
let bufferedBytes = 0;
let outputCommitted = false;
let terminalStatus: ResponsesTerminalStatus | undefined;
let failedPayload: Record<string, unknown> | undefined;
let retryableTerminalPayload: Record<string, unknown> | undefined;
const inspector = createSseInspector({
logCtx,
onParsedPayload: payload => {
if (comboStreamPayloadCommitsOutput(payload)) outputCommitted = true;
if (!payload || typeof payload !== "object" || Array.isArray(payload)) return;
if ((payload as { type?: unknown }).type === "response.failed") {
failedPayload = payload as Record<string, unknown>;
if (retryableZeroOutputTerminal(payload)) {
retryableTerminalPayload = payload as Record<string, unknown>;
}
},
onTerminal: status => { terminalStatus = status; },
Expand All @@ -162,9 +180,10 @@ export async function preflightComboStreamResponse(
inspector.feed(retained);
}

if (terminalStatus === "failed" && !outputCommitted && failedPayload) {
await reader.cancel("retrying zero-output combo stream failure").catch(() => undefined);
return { kind: "failed", response: failedTerminalResponse(response, failedPayload, logCtx) };
if ((terminalStatus === "failed" || terminalStatus === "incomplete")
&& !outputCommitted && retryableTerminalPayload) {
await reader.cancel("retrying zero-output combo stream terminal").catch(() => undefined);
return { kind: "failed", response: failedTerminalResponse(response, retryableTerminalPayload, logCtx) };
}
if (next.done || terminalStatus !== undefined || outputCommitted
|| bufferedBytes >= COMBO_STREAM_PREFLIGHT_MAX_BYTES
Expand Down
67 changes: 67 additions & 0 deletions tests/combo-stream-preflight.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ describe("combo stream preflight", () => {
expect(comboStreamPayloadCommitsOutput({ type: "response.created" })).toBe(false);
expect(comboStreamPayloadCommitsOutput({ type: "response.heartbeat" })).toBe(false);
expect(comboStreamPayloadCommitsOutput({ type: "response.failed" })).toBe(false);
expect(comboStreamPayloadCommitsOutput({ type: "response.incomplete" })).toBe(false);
expect(comboStreamPayloadCommitsOutput({ type: "response.output_text.delta", delta: "x" })).toBe(true);
expect(comboStreamPayloadCommitsOutput({ type: "response.output_item.added", item: { type: "function_call" } })).toBe(true);
expect(comboStreamPayloadCommitsOutput({ type: "provider.future_event" })).toBe(true);
Expand Down Expand Up @@ -49,6 +50,72 @@ describe("combo stream preflight", () => {
expect(JSON.stringify(body)).not.toContain("provider_trace_id");
});

test("converts zero-output transport incompletes into retryable HTTP failures", async () => {
const cases = [
["adapter_eof", "Upstream stream ended unexpectedly without a terminal event"],
["missing_terminal_event", "Upstream incomplete"],
["upstream_stall_timeout", "Upstream stalled"],
] as const;
for (const [reason, message] of cases) {
const result = await preflightComboStreamResponse(sse(
{ type: "response.created", response: { id: "r1", status: "in_progress" } },
{
type: "response.incomplete",
response: {
id: "r1",
status: "incomplete",
incomplete_details: { reason },
usage: { input_tokens: 11, output_tokens: 0, total_tokens: 11 },
},
},
), { model: "m1", provider: "a" });

expect(result.kind).toBe("failed");
expect(result.response.status).toBe(502);
const body = await result.response.json();
expect(body.error).toMatchObject({ type: "upstream_error", code: "upstream_server_error" });
expect(body.error.message).toContain(message);
expect(body.response.usage).toMatchObject({ input_tokens: 11, output_tokens: 0 });
}
});

test("does not replay semantic incompletes that another provider cannot safely repair", async () => {
const source = sse(
{ type: "response.created", response: { id: "r1", status: "in_progress" } },
{
type: "response.incomplete",
response: {
id: "r1",
status: "incomplete",
incomplete_details: { reason: "max_output_tokens" },
},
},
);
const expected = await source.clone().text();
const result = await preflightComboStreamResponse(source, { model: "m1", provider: "a" });
expect(result.kind).toBe("accepted");
expect(await result.response.text()).toBe(expected);
});

test("does not replay transport incompletes after output commits the target", async () => {
const source = sse(
{ type: "response.created", response: { id: "r1", status: "in_progress" } },
{ type: "response.output_text.delta", delta: "visible" },
{
type: "response.incomplete",
response: {
id: "r1",
status: "incomplete",
incomplete_details: { reason: "adapter_eof" },
},
},
);
const expected = await source.clone().text();
const result = await preflightComboStreamResponse(source, { model: "m1", provider: "a" });
expect(result.kind).toBe("accepted");
expect(await result.response.text()).toBe(expected);
});

test("replays buffered bytes unchanged after output commits the target", async () => {
const original = [
{ type: "response.created", response: { id: "r1", status: "in_progress" } },
Expand Down
42 changes: 42 additions & 0 deletions tests/server-combo-failover-e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -204,6 +204,13 @@ function chatStream(text: string): Response {
return new Response(frames, { headers: { "content-type": "text/event-stream" } });
}

function chatTruncatedZeroOutputStream(): Response {
const frames = [
`data: ${JSON.stringify({ choices: [{ index: 0, delta: {}, finish_reason: null }] })}\n\n`,
].join("");
return new Response(frames, { headers: { "content-type": "text/event-stream" } });
}

function chatErrorStream(message: string, prefix?: string): Response {
const frames = [
...(prefix
Expand Down Expand Up @@ -495,6 +502,41 @@ describe("server combo failover 030 activation matrix", () => {
}
});

test("zero-output adapter EOF hops to the next combo target", async () => {
const hits: string[] = [];
const a = serve(() => {
hits.push("a");
return chatTruncatedZeroOutputStream();
});
const b = serve(() => {
hits.push("b");
return chatStream("stream backup after adapter eof");
});
const config = comboConfig({
a: provider("openai-chat", baseUrl(a), "key-a"),
b: provider("openai-chat", baseUrl(b), "key-b"),
});

const response = await postLogged(config, { stream: true });
expect(response.status).toBe(200);
expect(JSON.stringify(await collectSse(response))).toContain("stream backup after adapter eof");
expect(hits).toEqual(["a", "b"]);

const { log, usage } = await latestAttemptReceipts(config);
for (const receipt of [log, usage]) {
expect(receipt).toMatchObject({
provider: "combo",
model: "combo/free",
resolvedModel: "m2",
attempts: [
{ ordinal: 1, provider: "a", model: "m1", status: 502 },
{ ordinal: 2, provider: "b", model: "m2", status: 200 },
],
});
expect(receipt.attempts[0]).not.toHaveProperty("firstOutputMs");
}
});

test("terminal SSE failure after output stays on the first target and never replays", async () => {
const hits: string[] = [];
const a = serve(() => {
Expand Down
Loading