diff --git a/open-sse/services/combo.ts b/open-sse/services/combo.ts index ce1e1dc7e20..12113422ca3 100644 --- a/open-sse/services/combo.ts +++ b/open-sse/services/combo.ts @@ -153,6 +153,7 @@ import { resolveDelayMs, comboModelNotFoundResponse, isStreamReadinessFailureErrorBody, + isStreamEarlyEofErrorBody, isTokenLimitBreachErrorBody, toRecordedTarget, getExhaustedTargetSkipReason, @@ -1511,6 +1512,11 @@ export async function handleComboChat({ const isStreamReadinessFailure = (result.status === 502 || result.status === 504) && isStreamReadinessFailureErrorBody(errorBody); + // An early EOF is an upstream failure, not a readiness probe — the breaker must + // see it even though the transient-retry path below treats both codes alike. + const isStreamEarlyEof = + (result.status === 502 || result.status === 504) && + isStreamEarlyEofErrorBody(errorBody); // FIX 5: a local per-API-key token-limit 429 must not cool shared accounts. const isTokenLimitBreach = @@ -1713,6 +1719,7 @@ export async function handleComboChat({ if ( shouldRecordProviderBreakerFailure({ isStreamReadinessFailure, + isStreamEarlyEof, status: result.status, sameProviderNext, skipProviderBreaker: fallbackResult.skipProviderBreaker, diff --git a/open-sse/services/combo/comboPredicates.ts b/open-sse/services/combo/comboPredicates.ts index dd150e210d2..0b6a939d7bf 100644 --- a/open-sse/services/combo/comboPredicates.ts +++ b/open-sse/services/combo/comboPredicates.ts @@ -133,7 +133,11 @@ const PROVIDER_BREAKER_FAILURE_STATUSES = new Set([408, 500, 502, 503, 504]); * failure (#1731 / #2743 gap-d). This is the consumer side of `skipProviderBreaker`: * * - Stream-readiness failures (pre-flight zombie/ping probes) never count as provider - * failures — they are a connection-readiness signal, not an upstream outage. + * failures — they are a connection-readiness signal, not an upstream outage. EXCEPT a + * STREAM_EARLY_EOF (`isStreamEarlyEof`): there the upstream returned HTTP 200, opened the + * SSE stream and then hung up without a single non-ping event, which is a genuine upstream + * failure. Excluding it made a provider-wide outage invisible to the breaker — see the + * STREAM_EARLY_EOF section of RESILIENCE_GUIDE.md. * - Only whole-provider failure statuses (408/500/502/503/504) count. A plain rate-limit * 429 is deliberately EXCLUDED — it belongs to connection cooldown / model lockout scope * (a genuine quota/token-limit 429 is handled there), NOT the whole-provider breaker. This @@ -163,6 +167,10 @@ const PROVIDER_BREAKER_FAILURE_STATUSES = new Set([408, 500, 502, 503, 504]); */ export function shouldRecordProviderBreakerFailure(args: { isStreamReadinessFailure: boolean; + /** True when the failure is specifically a STREAM_EARLY_EOF (upstream hung up after + * HTTP 200). Overrides the `isStreamReadinessFailure` exemption only; every other + * AND-term below still gates the trip. */ + isStreamEarlyEof?: boolean; status: number; sameProviderNext: boolean; skipProviderBreaker?: boolean; @@ -173,7 +181,7 @@ export function shouldRecordProviderBreakerFailure(args: { isProxyUnreachable?: boolean; }): boolean { return ( - !args.isStreamReadinessFailure && + (!args.isStreamReadinessFailure || args.isStreamEarlyEof === true) && PROVIDER_BREAKER_FAILURE_STATUSES.has(args.status) && (!args.sameProviderNext || args.isProxyUnreachable === true) && !args.skipProviderBreaker && @@ -308,6 +316,28 @@ export function isStreamReadinessFailureErrorBody(errorBody: unknown): boolean { return code === "STREAM_READINESS_TIMEOUT" || code === "STREAM_EARLY_EOF"; } +/** + * A STREAM_EARLY_EOF specifically: the upstream accepted the request (HTTP 200), opened the + * SSE stream, then closed it before emitting a single non-ping event. + * + * This is deliberately NOT the same signal as STREAM_READINESS_TIMEOUT. The readiness probe + * is a pre-flight liveness check on a connection we have not committed to yet, so failing it + * says "this connection looks stale", not "this provider is failing". An early EOF is the + * opposite: the provider took the request and then failed to serve it, which is an upstream + * failure by any reasonable definition. + * + * `isStreamReadinessFailureErrorBody` still covers both codes because the transient-retry and + * semaphore-cooldown paths in combo.ts want identical treatment for both. Only the + * whole-provider circuit breaker needs to tell them apart — see + * `shouldRecordProviderBreakerFailure`. + */ +export function isStreamEarlyEofErrorBody(errorBody: unknown): boolean { + if (!errorBody || typeof errorBody !== "object") return false; + const error = (errorBody as Record).error; + if (!error || typeof error !== "object") return false; + return (error as Record).code === "STREAM_EARLY_EOF"; +} + /** * A local per-API-key token-limit breach surfaces as a 429 tagged with * errorCode "TOKEN_LIMIT_EXCEEDED" (see chatCore.ts Tier 2 early return). This diff --git a/tests/unit/stream-early-eof-breaker.test.ts b/tests/unit/stream-early-eof-breaker.test.ts new file mode 100644 index 00000000000..df6564192f5 --- /dev/null +++ b/tests/unit/stream-early-eof-breaker.test.ts @@ -0,0 +1,176 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + shouldRecordProviderBreakerFailure, + isStreamEarlyEofErrorBody, + isStreamReadinessFailureErrorBody, +} from "../../open-sse/services/combo/comboPredicates.ts"; + +// A STREAM_EARLY_EOF means the upstream returned HTTP 200, opened the SSE stream, then +// closed it without emitting a single non-ping event. It was being classified together +// with STREAM_READINESS_TIMEOUT (a pre-flight liveness probe), and the readiness exemption +// in shouldRecordProviderBreakerFailure meant the whole-provider circuit breaker never saw +// it. During a provider-wide outage that made the breaker blind: every request kept being +// dispatched to the failing provider instead of shedding to the next combo target. +// +// The two codes still share the transient-retry and semaphore-cooldown paths in combo.ts. +// Only the breaker needs to tell them apart. + +const earlyEofBody = { + error: { + message: "Stream ended before producing a non-ping SSE event", + type: "stream_early_eof", + code: "STREAM_EARLY_EOF", + }, +}; + +const readinessBody = { + error: { + message: "Stream readiness timeout", + type: "stream_timeout", + code: "STREAM_READINESS_TIMEOUT", + }, +}; + +test("an early EOF trips the provider breaker", () => { + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: true, + status: 502, + sameProviderNext: false, + skipProviderBreaker: false, + requestScopedFailure: false, + error: "Stream ended before producing a non-ping SSE event", + }), + true + ); +}); + +test("a readiness-probe timeout still does not trip the breaker", () => { + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: false, + status: 502, + sameProviderNext: false, + skipProviderBreaker: false, + requestScopedFailure: false, + error: "Stream readiness timeout", + }), + false + ); +}); + +test("regression: before the fix both codes shared one flag, so the early EOF was exempted", () => { + // isStreamEarlyEof omitted entirely == the old call shape. + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + status: 502, + sameProviderNext: false, + skipProviderBreaker: false, + requestScopedFailure: false, + error: "Stream ended before producing a non-ping SSE event", + }), + false + ); +}); + +// The override is additive: it lifts the readiness exemption and nothing else. Every other +// AND-term in the gate must still be able to veto the trip. + +test("a client abort still does not trip, even on an early EOF", () => { + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: true, + status: 502, + sameProviderNext: false, + skipProviderBreaker: false, + requestScopedFailure: false, + error: "Client disconnected: request_signal_aborted", + }), + false + ); +}); + +test("skipProviderBreaker still wins over an early EOF", () => { + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: true, + status: 502, + sameProviderNext: false, + skipProviderBreaker: true, + requestScopedFailure: false, + error: "Stream ended before producing a non-ping SSE event", + }), + false + ); +}); + +test("a request-scoped failure still does not trip on an early EOF", () => { + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: true, + status: 502, + sameProviderNext: false, + skipProviderBreaker: false, + requestScopedFailure: true, + error: "Stream ended before producing a non-ping SSE event", + }), + false + ); +}); + +test("sameProviderNext still defers the trip on an early EOF", () => { + // Another model on the same provider may still succeed, so the existing policy holds. + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: true, + status: 502, + sameProviderNext: true, + skipProviderBreaker: false, + requestScopedFailure: false, + error: "Stream ended before producing a non-ping SSE event", + }), + false + ); +}); + +test("429 is still excluded from the whole-provider breaker", () => { + assert.equal( + shouldRecordProviderBreakerFailure({ + isStreamReadinessFailure: true, + isStreamEarlyEof: true, + status: 429, + sameProviderNext: false, + skipProviderBreaker: false, + requestScopedFailure: false, + error: "rate limit", + }), + false + ); +}); + +// Body classification: the new predicate must be strictly narrower than the existing one. + +test("isStreamEarlyEofErrorBody matches only the early-EOF code", () => { + assert.equal(isStreamEarlyEofErrorBody(earlyEofBody), true); + assert.equal(isStreamEarlyEofErrorBody(readinessBody), false); +}); + +test("isStreamReadinessFailureErrorBody keeps matching both codes", () => { + // The transient-retry and semaphore paths depend on this staying unchanged. + assert.equal(isStreamReadinessFailureErrorBody(earlyEofBody), true); + assert.equal(isStreamReadinessFailureErrorBody(readinessBody), true); +}); + +test("malformed bodies are not classified as an early EOF", () => { + for (const body of [null, undefined, "STREAM_EARLY_EOF", {}, { error: null }, { error: {} }]) { + assert.equal(isStreamEarlyEofErrorBody(body), false); + } +});