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
7 changes: 7 additions & 0 deletions open-sse/services/combo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,7 @@ import {
resolveDelayMs,
comboModelNotFoundResponse,
isStreamReadinessFailureErrorBody,
isStreamEarlyEofErrorBody,
isTokenLimitBreachErrorBody,
toRecordedTarget,
getExhaustedTargetSkipReason,
Expand Down Expand Up @@ -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);
Comment on lines +1517 to +1519

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Count early EOFs in round-robin combos

When the configured strategy is round-robin, handleComboChat returns through handleRoundRobinCombo at lines 676–688 before reaching this new classification. The round-robin failure path at lines 2763–2898 still groups STREAM_EARLY_EOF with readiness failures and never calls shouldRecordProviderBreakerFailure or recordProviderFailure, so sustained early EOFs on round-robin combos continue leaving the whole-provider breaker at zero. Apply the same early-EOF breaker decision in that handler as well.

AGENTS.md reference: open-sse/services/AGENTS.md:L9-L13

Useful? React with 👍 / 👎.


// FIX 5: a local per-API-key token-limit 429 must not cool shared accounts.
const isTokenLimitBreach =
Expand Down Expand Up @@ -1713,6 +1719,7 @@ export async function handleComboChat({
if (
shouldRecordProviderBreakerFailure({
isStreamReadinessFailure,
isStreamEarlyEof,
status: result.status,
sameProviderNext,
skipProviderBreaker: fallbackResult.skipProviderBreaker,
Expand Down
34 changes: 32 additions & 2 deletions open-sse/services/combo/comboPredicates.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand All @@ -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 &&
Expand Down Expand Up @@ -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<string, unknown>).error;
if (!error || typeof error !== "object") return false;
return (error as Record<string, unknown>).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
Expand Down
176 changes: 176 additions & 0 deletions tests/unit/stream-early-eof-breaker.test.ts
Original file line number Diff line number Diff line change
@@ -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);
}
});
Loading