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
10 changes: 6 additions & 4 deletions docs/design-notes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
62 changes: 34 additions & 28 deletions src/providers/bedrock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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"];
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand All @@ -828,28 +832,27 @@ 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,
): Promise<CompletionResult> {
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));
Expand Down Expand Up @@ -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,
Expand All @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
99 changes: 82 additions & 17 deletions test/empty-stream-retry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -190,34 +189,100 @@ 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<void>((resolve) => opts.abortSignal.addEventListener("abort", () => resolve(), { once: true }));
return {
body: (async function* () {
if (kind === "ok") {
yield textDelta("<p>second</p>");
yield messageDelta("end_turn");
return;
}
if (kind === "text") yield textDelta("<p>partial");
await aborted();
if (kind === "late") yield textDelta("<p>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, "<p>second</p>");
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("<p>ok</p>");
yield messageDelta("end_turn", { output_tokens: 40 });
return;
}
await new Promise<void>((resolve) =>
opts.abortSignal.addEventListener("abort", () => resolve(), { once: true }),
);
})(),
};
},
};
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 () => {
Expand Down
Loading