From 80c2eb090ab65f72905fea74591d65573336fa99 Mon Sep 17 00:00:00 2001 From: Jade Guo Date: Wed, 15 Jul 2026 01:21:01 +0800 Subject: [PATCH 1/3] fix(combo): fall back on Responses SSE failures --- open-sse/services/combo/validateQuality.ts | 62 +++++++-- ...mbo-responses-sse-failure-fallback.test.ts | 118 ++++++++++++++++++ 2 files changed, 168 insertions(+), 12 deletions(-) create mode 100644 tests/unit/combo-responses-sse-failure-fallback.test.ts diff --git a/open-sse/services/combo/validateQuality.ts b/open-sse/services/combo/validateQuality.ts index 2b561464795..a1cfeee0d9e 100644 --- a/open-sse/services/combo/validateQuality.ts +++ b/open-sse/services/combo/validateQuality.ts @@ -76,6 +76,20 @@ function responsesApiOutputHasContent(output: unknown): boolean { ); } +function isRecord(value: unknown): value is Record { + return !!value && typeof value === "object" && !Array.isArray(value); +} + +function isStreamingUpstreamError(parsed: Record, eventType: string): boolean { + if (eventType === "response.failed" || eventType === "error") return true; + 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 }. @@ -138,10 +152,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]; @@ -173,8 +187,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) { @@ -186,7 +204,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; @@ -205,7 +223,7 @@ export async function validateResponseQuality( break; } } - return false; + return null; } /** @@ -256,7 +274,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. @@ -294,9 +320,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 @@ -379,7 +416,9 @@ export async function validateResponseQuality( if (errorIsMeaningful) { const envelopeText = extractEnvelopeErrorText(json); const errMsg = - rawError && typeof rawError === "object" && typeof (rawError as Record).message === "string" + rawError && + typeof rawError === "object" && + typeof (rawError as Record).message === "string" ? ((rawError as Record).message as string) : envelopeText || JSON.stringify(rawError).substring(0, 200); return { valid: false, reason: `upstream error in 200 body: ${errMsg}` }; @@ -387,8 +426,7 @@ export async function validateResponseQuality( { 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}` }; } } diff --git a/tests/unit/combo-responses-sse-failure-fallback.test.ts b/tests/unit/combo-responses-sse-failure-fallback.test.ts new file mode 100644 index 00000000000..b723adae5db --- /dev/null +++ b/tests/unit/combo-responses-sse-failure-fallback.test.ts @@ -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({ + 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); +}); From dadab7ddd3b53a8a50dbf939895f857fa1f44e05 Mon Sep 17 00:00:00 2001 From: Jade Guo Date: Wed, 15 Jul 2026 01:27:36 +0800 Subject: [PATCH 2/3] Update open-sse/services/combo/validateQuality.ts Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> --- open-sse/services/combo/validateQuality.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/open-sse/services/combo/validateQuality.ts b/open-sse/services/combo/validateQuality.ts index a1cfeee0d9e..ac365acc3aa 100644 --- a/open-sse/services/combo/validateQuality.ts +++ b/open-sse/services/combo/validateQuality.ts @@ -80,8 +80,9 @@ function isRecord(value: unknown): value is Record { return !!value && typeof value === "object" && !Array.isArray(value); } -function isStreamingUpstreamError(parsed: Record, eventType: string): boolean { +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; From 6edbf3191f8d013b557de10215ab469af3fad3f6 Mon Sep 17 00:00:00 2001 From: Jade Guo Date: Wed, 15 Jul 2026 01:38:58 +0800 Subject: [PATCH 3/3] chore: rerun CI with PR evidence