Skip to content
Closed
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
63 changes: 51 additions & 12 deletions open-sse/services/combo/validateQuality.ts
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,21 @@ function responsesApiOutputHasContent(output: unknown): boolean {
);
}

function isRecord(value: unknown): value is Record<string, unknown> {
return !!value && typeof value === "object" && !Array.isArray(value);
}

function isStreamingUpstreamError(parsed: unknown, eventType: string): boolean {
if (eventType === "response.failed" || eventType === "error") return true;
if (!isRecord(parsed)) return false;
if (parsed.error != null) return true;

const nestedResponse = isRecord(parsed.response) ? parsed.response : null;
return nestedResponse?.status === "failed" && nestedResponse.error != null;
}

type StreamingPeekOutcome = "content" | "error" | null;

/**
* Validate that a successful (HTTP 200) non-streaming response actually contains
* meaningful content. Returns { valid: true } or { valid: false, reason }.
Expand Down Expand Up @@ -138,10 +153,10 @@ export async function validateResponseQuality(
* flags in the closure. The last (potentially incomplete) line is kept in
* `decodedSoFar` for the next iteration.
*
* Returns true when a content_block_* event is detected — the caller
* should stop peeking and treat the stream as non-empty.
* Returns "content" when a valuable event is detected or "error" when the
* upstream reports a failure before content. Otherwise peeking continues.
*/
function parseAccumulatedSse(): boolean {
function parseAccumulatedSse(): StreamingPeekOutcome {
const lines = decodedSoFar.split(/\r?\n/);
// Retain the potentially-incomplete trailing fragment.
decodedSoFar = lines[lines.length - 1];
Expand Down Expand Up @@ -173,8 +188,12 @@ export async function validateResponseQuality(
(typeof parsed.type === "string" ? parsed.type : null) || pendingEventType || "";
pendingEventType = "";

if (isStreamingUpstreamError(parsed, eventType)) {
return "error";
}

if (isKnownNonClaudeStreamPayload(parsed, eventType)) {
return true;
return "content";
}

switch (eventType) {
Expand All @@ -186,7 +205,7 @@ export async function validateResponseQuality(
case "content_block_stop":
hasContentBlock = true;
// Signal caller to stop buffering immediately.
return true;
return "content";
case "message_stop":
hasLifecycleEnd = true;
break;
Expand All @@ -205,7 +224,7 @@ export async function validateResponseQuality(
break;
}
}
return false;
return null;
}

/**
Expand Down Expand Up @@ -256,7 +275,15 @@ export async function validateResponseQuality(
const tail = decoder.decode(undefined, { stream: false });
if (tail) decodedSoFar += tail;
if (decodedSoFar.trim()) decodedSoFar += "\n\n";
parseAccumulatedSse();
const terminalOutcome = parseAccumulatedSse();

if (terminalOutcome === "error") {
log.warn?.(
"COMBO",
"Streaming response reported an upstream error before content — marking as invalid for combo failover"
);
return { valid: false, reason: "streaming upstream error" };
}

if (hasMessageStart && hasLifecycleEnd && !hasContentBlock) {
// Complete Claude lifecycle with zero content blocks → failover.
Expand Down Expand Up @@ -294,9 +321,20 @@ export async function validateResponseQuality(

// Decode incrementally (stream:true keeps multi-byte char state).
decodedSoFar += decoder.decode(value, { stream: true });
const foundContent = parseAccumulatedSse();
const outcome = parseAccumulatedSse();

if (outcome === "error") {
// Do not await cancellation of a Response.clone() tee branch: the
// promise may remain pending until the client-facing branch drains.
reader.cancel().catch(() => {});
log.warn?.(
"COMBO",
"Streaming response reported an upstream error before content — marking as invalid for combo failover"
);
return { valid: false, reason: "streaming upstream error" };
}

if (foundContent) {
if (outcome === "content") {
anyContentFound = true;
// A content_block_* event was found — stop peeking. Return a
// clonedResponse that replays all buffered bytes (the current chunk
Expand Down Expand Up @@ -379,16 +417,17 @@ export async function validateResponseQuality(
if (errorIsMeaningful) {
const envelopeText = extractEnvelopeErrorText(json);
const errMsg =
rawError && typeof rawError === "object" && typeof (rawError as Record<string, unknown>).message === "string"
rawError &&
typeof rawError === "object" &&
typeof (rawError as Record<string, unknown>).message === "string"
? ((rawError as Record<string, unknown>).message as string)
: envelopeText || JSON.stringify(rawError).substring(0, 200);
return { valid: false, reason: `upstream error in 200 body: ${errMsg}` };
}
{
const envelopeText = extractEnvelopeErrorText(json);
if (envelopeText && EXHAUSTION_MARKER_PATTERN.test(envelopeText)) {
const snippet =
envelopeText.length > 80 ? `${envelopeText.slice(0, 80)}…` : envelopeText;
const snippet = envelopeText.length > 80 ? `${envelopeText.slice(0, 80)}…` : envelopeText;
return { valid: false, reason: `upstream exhaustion marker in 200 body: ${snippet}` };
}
}
Expand Down
118 changes: 118 additions & 0 deletions tests/unit/combo-responses-sse-failure-fallback.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
import test from "node:test";
import assert from "node:assert/strict";

import { handleComboChat, validateResponseQuality } from "../../open-sse/services/combo.ts";

const encoder = new TextEncoder();

function sseResponse(body: string): Response {
return new Response(
new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(encoder.encode(body));
controller.close();
},
}),
{ status: 200, headers: { "Content-Type": "text/event-stream" } }
);
}

function silentLog() {
return { info() {}, warn() {}, error() {}, debug() {} };
}

function failedResponsesSse(): string {
return [
"event: response.failed",
`data: ${JSON.stringify({
type: "response.failed",
response: {
status: "failed",
error: { code: "no_capacity", message: "peak capacity" },
},
})}`,
"",
"",
].join("\n");
}

test("streaming quality rejects a pre-content response.failed event", async () => {
const result = await validateResponseQuality(
sseResponse(failedResponsesSse()),
true,
silentLog()
);

assert.equal(result.valid, false);
assert.equal(result.reason, "streaming upstream error");
});

test("streaming quality rejects a pre-content top-level error envelope", async () => {
const body = [
"event: error",
`data: ${JSON.stringify({
error: { type: "server_error", message: "temporarily unavailable" },
})}`,
"",
"",
].join("\n");

const result = await validateResponseQuality(sseResponse(body), true, silentLog());

assert.equal(result.valid, false);
assert.equal(result.reason, "streaming upstream error");
});

test("combo advances to the next target after a pre-content Responses SSE failure", async () => {
const calls: string[] = [];
const healthy = [
"event: response.output_text.delta",
`data: ${JSON.stringify({ type: "response.output_text.delta", delta: "fallback ok" })}`,
"",
"",
].join("\n");

const result = await handleComboChat({
body: { stream: true, messages: [{ role: "user", content: "hello" }] },
combo: {
name: "responses-sse-failure-fallback",
strategy: "priority",
models: [
{ model: "openai/primary", weight: 0 },
{ model: "openai/secondary", weight: 0 },
],
config: { maxRetries: 0, retryDelayMs: 0 },
},
handleSingleModel: async (_body: unknown, model: string) => {
calls.push(model);
return model.endsWith("/primary") ? sseResponse(failedResponsesSse()) : sseResponse(healthy);
},
isModelAvailable: async () => true,
log: silentLog(),
settings: null,
allCombos: null,
relayOptions: null as never,
});

assert.equal(result.ok, true);
assert.deepEqual(calls, ["openai/primary", "openai/secondary"]);
assert.match(await result.text(), /fallback ok/);
});

test("streaming quality still replays normal Responses lifecycle and content", async () => {
const body = [
"event: response.created",
`data: ${JSON.stringify({ type: "response.created", response: { id: "resp_1" } })}`,
"",
"event: response.output_text.delta",
`data: ${JSON.stringify({ type: "response.output_text.delta", delta: "hello" })}`,
"",
"",
].join("\n");

const result = await validateResponseQuality(sseResponse(body), true, silentLog());

assert.equal(result.valid, true);
assert.ok(result.clonedResponse);
assert.equal(await result.clonedResponse.text(), body);
});
Loading