diff --git a/docs/design-notes.md b/docs/design-notes.md index a1a77159..2be0fa3d 100644 --- a/docs/design-notes.md +++ b/docs/design-notes.md @@ -2769,10 +2769,12 @@ These are the rules both adapters enforce on a model call. The README states the Verified empirically against a stubbed request handler (3 wire attempts for 503/429/ECONNRESET, 1 for a 400). A general loop around it would give Bedrock 9 attempts to OpenRouter's 3. - There is one narrow exception (#480): a stream that closes having sent nothing is sent once more. - The SDK cannot retry it, because it arrives as a 200. Both of those sends get the SDK's own 3 wire - attempts, so the worst case is 6 wire attempts instead of 3 (12 across the output-ceiling retry). - That only happens when an empty stream is followed by a throttle or a 5xx. + There is one narrow exception: a call that produced no output is sent once more. That is a stream + that closed empty (#480), or no output within the 120 s first-output window (#484, Bedrock only; + OpenRouter does not retry a stall). The SDK retries neither, because both arrive as a 200. Both + sends get the SDK's own 3 wire attempts, so the worst case is 6 instead of 3 (12 across the + output-ceiling retry), and only when an empty stream or stall is followed by a throttle or a 5xx. + A stall followed by a stall takes 240 s to fail. - **The Bedrock adapter speaks two dialects**, chosen by `providers.bedrock.api`. `invoke` (the default) is `InvokeModelWithResponseStream` carrying an Anthropic-native body, and it is what every published number in this repo was measured through. `converse` is `ConverseStream`, whose request diff --git a/src/providers/bedrock.ts b/src/providers/bedrock.ts index 7a31da5e..77660864 100644 --- a/src/providers/bedrock.ts +++ b/src/providers/bedrock.ts @@ -78,6 +78,11 @@ const EMPTY_STREAM_RETRY_MS = 400; // wording drifted. const EMPTY_STREAM_DETAIL = "no message_stop and no stop_reason"; +// A failed attempt that produced no output, which is the one kind sent again. +function producedNothing(e: unknown): boolean { + return e instanceof EmptyStreamError || (e instanceof StalledStreamError && e.kind === "first_output" && e.chars === 0); +} + // What the upstream actually sends in a usage block, which is a superset of what // `Usage` declares: today `service_tier` and a nested `cache_creation` breakdown ride // along beside the four counts, and a model release can add more without notice. @@ -533,9 +538,10 @@ function outputCeilingRefused(model: string, asked: number, cause: unknown, capp // // ONE case out of that gap is retried here, and it is the one where neither of those // objections applies: a stream that closes having delivered nothing at all -// (`EmptyStreamError`, issue #480 — a user lost a whole document to it). There is no +// (`EmptyStreamError`, issue #480 — a user lost a whole document to it), or one our own +// clock abandoned before any output (a `first_output` stall, issue #484). There is no // streamed output to discard and no document to resume from, so the request is simply -// sent again, once. See `sendRetryingEmptyStream`. +// sent again, once. See `sendRetryingNoOutput`. export class BedrockProvider implements ModelProvider { name = "bedrock"; capabilities: Capability[] = ["text", "vision", "structured_output"]; @@ -676,7 +682,7 @@ export class BedrockProvider implements ModelProvider { spent: false, }; try { - return await this.sendRetryingEmptyStream(req, system, first); + return await this.sendRetryingNoOutput(req, system, first); } catch (e) { // `first.spent` is the guarantee that sending it again costs nothing: a refusal // arrives before generation, so a failure that had already been billed for is not @@ -751,7 +757,7 @@ export class BedrockProvider implements ModelProvider { // deployment's ceiling and below any cap this call carried. const second: Attempt = { maxTokens: stated, ceilingFrom: "model", spent: false }; try { - return await this.sendRetryingEmptyStream(req, system, second); + return await this.sendRetryingNoOutput(req, system, second); } catch (again) { if (!refusedForOutputCeiling(again) || second.spent) throw again; // This page is lost either way — a third attempt is not on offer, since a model that @@ -799,24 +805,22 @@ export class BedrockProvider implements ModelProvider { } } - // One attempt at one ceiling, sent a second time if the stream closed having delivered - // nothing (issue #480: a user's document failed with "0 chars received, no message_stop - // and no stop_reason", and the conversion was lost for a failure that had produced no - // content to protect). + // One attempt at one ceiling, sent a second time if it produced no output at all: the + // stream closed empty (issue #480: "0 chars received, no message_stop and no + // stop_reason"), or no output arrived within the first-output window (issue #484: 10 of + // luna's 308 UIC page calls, each a lost page). // - // Safe in the two ways the note above `BedrockProvider` says a mid-stream retry usually - // is not. Nothing is discarded: `EmptyStreamError` is raised only when not one character - // arrived, so there is no partial document to throw away and no risk of a passage - // shipping twice. And a stalled attempt is never retried: a stall is a `StalledStreamError`, - // checked before the completeness check that raises this, so the attempt this follows is - // one the upstream closed itself. + // Safe in the way the note above `BedrockProvider` says a mid-stream retry usually is + // not: either is retried only when not one character arrived (`producedNothing`), so + // there is no partial document to throw away and no risk of a passage shipping twice. An + // `idle` or `total` stall had output and is not retried. // // "Closed itself" does not mean "closed quickly". On the UIC deployment every empty stream // came from one model, us.openai.gpt-5.6-luna on Converse, after 42, 82 and 83 seconds of // silence (3 of its 308 page calls from 2026-09-01 to 09-24; the same model also hit the // 120 s first-output stall 10 times, and no other model did either). So the retry can add up - // to one more first-output window, 120 s, to a page. That is the price of not losing the - // document. MAX_TOTAL_MS does not cap it: each send arms its own total timer, so the + // to one more first-output window, 120 s, to a page, and a stall followed by a stall + // takes 240 s to fail. That is the price of not losing the page. MAX_TOTAL_MS does not cap it: each send arms its own total timer, so the // retry gets a fresh one. A worst-case call holds its slot for the empty send, then // EMPTY_STREAM_RETRY_MS, then a full MAX_TOTAL_MS. The output-ceiling retry already // works the same way. @@ -828,11 +832,10 @@ export class BedrockProvider implements ModelProvider { // for — and `attempt.billed` keeps the abandoned attempt's counts in the call's reported // usage, so the run log shows what the retry cost rather than hiding it. // - // Once, not until it works. An upstream that answers an identical request with two empty - // streams is not having a blip, and a third attempt would only spend a third prompt to - // say so; the error names the attempt count so a run log can show that this is where it - // ended up. - private async sendRetryingEmptyStream( + // Once, not until it works. An upstream that answers an identical request with nothing + // twice is not having a blip, and a third attempt would only spend a third prompt to + // say so. + private async sendRetryingNoOutput( req: CompletionRequest, system: string, attempt: Attempt, @@ -840,16 +843,16 @@ export class BedrockProvider implements ModelProvider { try { return await this.send(req, system, attempt); } catch (e) { - if (!(e instanceof EmptyStreamError)) throw e; + if (!producedNothing(e)) throw e; // Said on every occurrence rather than once per process, unlike the ceiling warnings // above: those report a standing config error that is the same news however many // pages meet it, while this is a transient upstream event whose FREQUENCY is the // whole question. A deployment seeing it on one page a week and one seeing it on // every page have different problems, and only the count tells them apart. console.warn( - `bedrock: ${req.model} closed a response stream having sent nothing at all, so the ` + - `request is being sent again once after ${EMPTY_STREAM_RETRY_MS}ms. Nothing was ` + - `generated, so nothing is being discarded — but the prompt is read, and paid for, ` + + `bedrock: ${req.model} ${e instanceof EmptyStreamError ? "closed a response stream having sent nothing at all" : "sent no output within the first-output window"}, ` + + `so the request is being sent again once after ${EMPTY_STREAM_RETRY_MS}ms. Nothing was ` + + `generated, so nothing is being discarded — but the prompt may be read, and paid for, ` + `a second time.`, ); await new Promise((resolve) => setTimeout(resolve, EMPTY_STREAM_RETRY_MS)); @@ -1078,8 +1081,10 @@ export class BedrockProvider implements ModelProvider { if (text) arm("idle", this.idleTimeoutMs); else arm("first_output", this.firstOutputTimeoutMs); }; - const stalled = (kind: StallKind): StalledStreamError => - new StalledStreamError({ + const stalled = (kind: StallKind): StalledStreamError => { + // Kept for the retry's report, as an empty stream's is below. + if (kind === "first_output") attempt.billed = addUsage(attempt.billed, usage); + return new StalledStreamError({ provider: this.name, model: req.model, kind, @@ -1091,6 +1096,7 @@ export class BedrockProvider implements ModelProvider { : this.maxTotalMs, chars: text.length, }); + }; // The clock starts before the request: time-to-first-token is exactly as much of // a stall risk as a gap mid-stream, and prompt processing happens in here too. @@ -1176,7 +1182,7 @@ export class BedrockProvider implements ModelProvider { if (expired && !sawStop) throw stalled(expired); if (!sawStop && !stopReason) { // Nothing arrived at all, which is a different failure from a document cut short and - // is the one that can be sent again (see `sendRetryingEmptyStream` and + // is the one that can be sent again (see `sendRetryingNoOutput` and // `EmptyStreamError`). Folded into `attempt.billed` here rather than in the caller, // because this is the only place that knows what this attempt was charged for and the // only failure the caller answers by re-sending. diff --git a/test/empty-stream-retry.test.ts b/test/empty-stream-retry.test.ts index 0c51a943..c715ce5a 100644 --- a/test/empty-stream-retry.test.ts +++ b/test/empty-stream-retry.test.ts @@ -12,13 +12,12 @@ // - Ended empty: nothing is in hand. Nothing to discard, nothing that can ship short, and // nothing about the request the upstream objected to — so it is sent again. // -// Two things are pinned as hard as the retry itself, because both are ways a retry does -// damage rather than good. A stalled call must not become a retried one: it would double -// the time a wedged session takes to fail, and `expired` is checked before the completeness -// check that raises this, which is what makes that true. And the abandoned attempt's token -// counts must survive into the surviving attempt's report: the Anthropic stream reports the -// prompt's counts in `message_start`, so an attempt that got that far and closed was billed, -// and a call that paid for two prompts must not be logged as having paid for one. +// Bedrock also re-sends a call our clock abandoned before any output (a `first_output` +// stall, issue #484), on the same grounds. A stall after output started is not retried. +// +// The abandoned attempt's token counts must survive into the surviving attempt's report: the +// Anthropic stream reports the prompt's counts in `message_start`, so an attempt that got +// that far was billed, and a call that paid for two prompts must not be logged as one. import { test } from "node:test"; import assert from "node:assert/strict"; import { BedrockProvider } from "../src/providers/bedrock.ts"; @@ -190,21 +189,55 @@ test("a Bedrock stream that ends SHORT is not retried, and still reports what it assert.equal(sends.count(), 1); }); -test("a Bedrock call that stalls before any output is a stall, not an empty stream", async () => { - // The other safety pin, and the reason the retry cannot lengthen a wedged session: a call - // abandoned by our own clock has also received 0 characters, so if the completeness check - // were reached first it would look identical to #480 and be sent again — turning a - // 120-second failure into a 240-second one. `expired` is checked first, and this is what - // says so. +// A client whose sends follow `script(send)`: "silent" waits for the abort and ends, "text" +// sends one delta then goes silent, "late" sends one delta just after the abort, "ok" completes. +function stubStalls(bedrock: BedrockProvider, script: (send: number) => "silent" | "text" | "late" | "ok") { + let sends = 0; + (bedrock as unknown as { client: unknown }).client = { + send: async (_cmd: unknown, opts: { abortSignal: AbortSignal }) => { + const kind = script(++sends); + const aborted = () => + new Promise((resolve) => opts.abortSignal.addEventListener("abort", () => resolve(), { once: true })); + return { + body: (async function* () { + if (kind === "ok") { + yield textDelta("

second

"); + yield messageDelta("end_turn"); + return; + } + if (kind === "text") yield textDelta("

partial"); + await aborted(); + if (kind === "late") yield textDelta("

late"); + })(), + }; + }, + }; + return { count: () => sends }; +} + +test("a Bedrock call that stalls before any output is sent again once, and the retry is delivered", async () => { + const bedrock = new BedrockProvider({ default_model: "m" }, { firstOutputTimeoutMs: 50 }); + const sends = stubStalls(bedrock, (n) => (n === 1 ? "silent" : "ok")); + const [res, said] = await capturingWarnings(() => bedrock.complete(bedrockReq)); + assert.equal(res.text, "

second

"); + assert.equal(sends.count(), 2); + assert.match(said.join("\n"), /no output within the first-output window/); +}); + +test("a stalled attempt's prompt tokens stay in the retry's reported usage", async () => { const bedrock = new BedrockProvider({ default_model: "m" }, { firstOutputTimeoutMs: 50 }); let sends = 0; (bedrock as unknown as { client: unknown }).client = { send: async (_cmd: unknown, opts: { abortSignal: AbortSignal }) => { - sends++; + const n = ++sends; return { body: (async function* () { - // Silent until the first-output clock fires, then end without throwing — the - // abort shape that reaches the completeness check rather than the catch. + yield messageStart({ input_tokens: 900 }); + if (n === 2) { + yield textDelta("

ok

"); + yield messageDelta("end_turn", { output_tokens: 40 }); + return; + } await new Promise((resolve) => opts.abortSignal.addEventListener("abort", () => resolve(), { once: true }), ); @@ -212,12 +245,44 @@ test("a Bedrock call that stalls before any output is a stall, not an empty stre }; }, }; + const [res] = await capturingWarnings(() => bedrock.complete(bedrockReq)); + assert.deepEqual(res.usage, { input_tokens: 1800, output_tokens: 40 }); +}); + +test("two first-output stalls fail as a stall, after two sends", async () => { + const bedrock = new BedrockProvider({ default_model: "m" }, { firstOutputTimeoutMs: 50 }); + const sends = stubStalls(bedrock, () => "silent"); + await capturingWarnings(() => + assert.rejects(() => bedrock.complete(bedrockReq), (e: Error) => { + assert.ok(e instanceof StalledStreamError); + assert.equal(e.kind, "first_output"); + return true; + }), + ); + assert.equal(sends.count(), 2); +}); + +test("a first-output stall whose first text lands after the abort is not sent again", async () => { + const bedrock = new BedrockProvider({ default_model: "m" }, { firstOutputTimeoutMs: 50 }); + const sends = stubStalls(bedrock, () => "late"); await assert.rejects(() => bedrock.complete(bedrockReq), (e: Error) => { assert.ok(e instanceof StalledStreamError); assert.equal(e.kind, "first_output"); + assert.ok(e.chars > 0); return true; }); - assert.equal(sends, 1); + assert.equal(sends.count(), 1); +}); + +test("a Bedrock call that stalls after output started is not sent again", async () => { + const bedrock = new BedrockProvider({ default_model: "m" }, { firstOutputTimeoutMs: 50, idleTimeoutMs: 50 }); + const sends = stubStalls(bedrock, () => "text"); + await assert.rejects(() => bedrock.complete(bedrockReq), (e: Error) => { + assert.ok(e instanceof StalledStreamError); + assert.equal(e.kind, "idle"); + return true; + }); + assert.equal(sends.count(), 1); }); test("the Converse path retries an empty stream too", async () => {