diff --git a/src/sse/handlers/chat.ts b/src/sse/handlers/chat.ts index da10e9a5506..631a63c1132 100644 --- a/src/sse/handlers/chat.ts +++ b/src/sse/handlers/chat.ts @@ -1530,6 +1530,26 @@ async function handleSingleModelChat( continue; } + // #8928: once the bounded same-connection retry is unavailable or exhausted, + // remove only the affinity pin that still points at this failed connection. This lets + // the next client retry select another eligible account without deleting a + // pin that may already have moved to a healthy connection. + const isTerminalStreamEarlyEof = + result.errorCode === "STREAM_EARLY_EOF" || + result.errorType === "stream_early_eof"; + + if (isTerminalStreamEarlyEof && runtimeOptions.sessionAffinityKey) { + try { + evictSessionAccountAffinityForConnection( + runtimeOptions.sessionAffinityKey, + provider, + credentials.connectionId + ); + } catch { + // Best-effort: the current response still surfaces the original 502. + } + } + // Stream readiness timeout is an upstream stall after an HTTP response was received, // not an account/quota failure. Do NOT mark the account unavailable here. return withSelectedConnectionHeader(result.response, credentials?.connectionId); diff --git a/tests/unit/stream-early-eof-affinity-8928.test.ts b/tests/unit/stream-early-eof-affinity-8928.test.ts new file mode 100644 index 00000000000..2af93bd1523 --- /dev/null +++ b/tests/unit/stream-early-eof-affinity-8928.test.ts @@ -0,0 +1,81 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; + +const chatSource = fs.readFileSync( + new URL("../../src/sse/handlers/chat.ts", import.meta.url), + "utf8" +); + +function getNonAntigravityStreamFailureBranch(): string { + const startMarker = [ + " if (", + ' (result.errorType === "stream_timeout" || result.errorType === "stream_early_eof") &&', + " !isAntigravityStreamReadinessFailure", + " ) {", + ].join("\n"); + + const start = chatSource.indexOf(startMarker); + assert.notEqual(start, -1, "non-Antigravity stream-failure branch must exist"); + + const end = chatSource.indexOf( + "\n if (isAntigravityStreamReadinessFailure)", + start + ); + assert.notEqual(end, -1, "stream-failure branch end marker must exist"); + + return chatSource.slice(start, end); +} + +test("terminal STREAM_EARLY_EOF evicts affinity after the bounded retry (#8928)", () => { + const branch = getNonAntigravityStreamFailureBranch(); + + const retryContinue = branch.indexOf("continue;"); + const eviction = branch.indexOf( + "evictSessionAccountAffinityForConnection(" + ); + const terminalReturn = branch.indexOf( + "return withSelectedConnectionHeader(" + ); + + assert.ok(retryContinue >= 0, "the existing bounded retry must remain"); + assert.ok(eviction > retryContinue, "eviction must happen only after retry is exhausted"); + assert.ok(terminalReturn > eviction, "eviction must happen before the terminal 502 is returned"); +}); + +test("early-EOF eviction recognizes both typed error signals (#8928)", () => { + const branch = getNonAntigravityStreamFailureBranch(); + + assert.match( + branch, + /result\.errorCode === "STREAM_EARLY_EOF"/, + "the canonical STREAM_EARLY_EOF code must trigger eviction" + ); + assert.match( + branch, + /result\.errorType === "stream_early_eof"/, + "the typed stream_early_eof fallback must trigger eviction" + ); +}); + +test("early-EOF eviction is session-key and connection guarded (#8928)", () => { + const branch = getNonAntigravityStreamFailureBranch(); + + assert.match( + branch, + /isTerminalStreamEarlyEof && runtimeOptions\.sessionAffinityKey/, + "eviction must run only for terminal early EOF with an affinity key" + ); + + assert.match( + branch, + /evictSessionAccountAffinityForConnection\(\s*runtimeOptions\.sessionAffinityKey,\s*provider,\s*credentials\.connectionId\s*\)/s, + "the connection-matched helper must protect a pin that moved elsewhere" + ); + + assert.doesNotMatch( + branch, + /deleteSessionAccountAffinity\(/, + "the terminal branch must not perform an unguarded affinity deletion" + ); +});