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
20 changes: 20 additions & 0 deletions src/sse/handlers/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
81 changes: 81 additions & 0 deletions tests/unit/stream-early-eof-affinity-8928.test.ts
Original file line number Diff line number Diff line change
@@ -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"
);
});