diff --git a/docs/design-notes.md b/docs/design-notes.md
index 33ab81cc..eb0f5af0 100644
--- a/docs/design-notes.md
+++ b/docs/design-notes.md
@@ -2767,7 +2767,12 @@ These are the rules both adapters enforce on a model call. The README states the
loop *on purpose*: the AWS SDK already applies its `standard` strategy — also 3 attempts with
exponential backoff — to throttling, 5xx, and node network errors, while failing fast on 4xx.
Verified empirically against a stubbed request handler (3 wire attempts for 503/429/ECONNRESET,
- 1 for a 400). Adding a loop around it would give Bedrock 9 attempts to OpenRouter's 3.
+ 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.
- **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/public/demo.html b/public/demo.html
index 41705180..339d1436 100644
--- a/public/demo.html
+++ b/public/demo.html
@@ -302,6 +302,26 @@
{
+ try {
+ return await this.send(req, system, attempt);
+ } catch (e) {
+ if (!(e instanceof EmptyStreamError)) 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, ` +
+ `a second time.`,
+ );
+ await new Promise((resolve) => setTimeout(resolve, EMPTY_STREAM_RETRY_MS));
+ try {
+ return await this.send(req, system, attempt);
+ } catch (again) {
+ if (!(again instanceof EmptyStreamError)) throw again;
+ // The surviving message says it happened twice. Re-raised rather than rethrown
+ // because the second attempt's own error says "sent once", which would tell an
+ // operator the retry had not been reached.
+ throw new EmptyStreamError({
+ provider: this.name,
+ model: req.model,
+ attempts: 2,
+ detail: EMPTY_STREAM_DETAIL,
+ });
+ }
+ }
+ }
+
// One attempt at one ceiling.
//
// Sent from inside the branch rather than after it, because the SDK's `send` is typed
@@ -957,7 +1063,11 @@ export class BedrockProvider implements ModelProvider {
// arrived yet — see Attempt.spent.
attempt.spent = true;
usage = { ...usage, ...u };
- req.onUsage?.(usage);
+ // `attempt.billed` and not `usage` alone: on a call that is being sent a second time
+ // after an empty stream, the first attempt's prompt was billed and must stay in the
+ // total. Absent on every other call, where `addUsage` returns `usage` untouched.
+ const total = addUsage(attempt.billed, usage);
+ if (total) req.onUsage?.(total);
};
// Which window an event re-arms is decided by whether any text has arrived, not
// by the event's own type. Protocol events (message_start, content_block_start)
@@ -1065,9 +1175,23 @@ export class BedrockProvider implements ModelProvider {
// of a message that had already stopped took nothing from the document.
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
+ // `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.
+ if (!text) {
+ attempt.billed = addUsage(attempt.billed, usage);
+ throw new EmptyStreamError({
+ provider: this.name,
+ model: req.model,
+ attempts: 1,
+ detail: EMPTY_STREAM_DETAIL,
+ });
+ }
throw new Error(
`bedrock: the response stream ended without completing (${text.length} chars received, ` +
- `no message_stop and no stop_reason). Treating a partial document as a whole one would ` +
+ `${EMPTY_STREAM_DETAIL}). Treating a partial document as a whole one would ` +
`deliver content the source never had.`,
);
}
@@ -1105,6 +1229,10 @@ export class BedrockProvider implements ModelProvider {
`whole one.`,
);
}
- return { text, model: req.model, provider: this.name, usage };
+ // `attempt.billed` for the same reason `mergeUsage` reports it: a call that was sent
+ // again after an empty stream paid for both prompts, and the router reads usage off the
+ // result on the surviving path and off the callback on the failing one — the two have to
+ // agree that the abandoned attempt was paid for.
+ return { text, model: req.model, provider: this.name, usage: addUsage(attempt.billed, usage) };
}
}
diff --git a/src/providers/openrouter.ts b/src/providers/openrouter.ts
index 2503bd03..b0a8f9cc 100644
--- a/src/providers/openrouter.ts
+++ b/src/providers/openrouter.ts
@@ -1,5 +1,11 @@
import { DEFAULT_MAX_TOKENS, type Capability, type ProviderBlock } from "../config.ts";
-import { StalledStreamError, TruncatedResponseError, type StallKind } from "./types.ts";
+import {
+ EmptyStreamError,
+ StalledStreamError,
+ TruncatedResponseError,
+ addUsage,
+ type StallKind,
+} from "./types.ts";
import type { CompletionRequest, CompletionResult, ModelProvider, Usage } from "./types.ts";
import {
cacheableSystemPrompt,
@@ -88,26 +94,6 @@ export function normalizeUsage(u?: OpenAIUsage): Usage | undefined {
return Object.keys(usage).length ? usage : undefined;
}
-// Add two usage snapshots. Used across retry attempts, where the counts ADD rather
-// than replace: an attempt that reported tokens and was then abandoned was still
-// billed for them, so reporting only the surviving attempt understates the call — and
-// understates it invisibly, since `tokens.calls_reported` would still say the call was
-// fully accounted for.
-//
-// Absent stays absent when neither side reported: a 0 nobody sent reads as a free
-// half of the call rather than an unreported one.
-function addUsage(a?: Usage, b?: Usage): Usage | undefined {
- if (!a) return b;
- if (!b) return a;
- const sum: Usage = { ...a };
- for (const key of Object.keys(b) as (keyof Usage)[]) {
- const v = b[key];
- if (v == null) continue;
- sum[key] = (sum[key] ?? 0) + v;
- }
- return sum;
-}
-
// OpenRouter adapter. Speaks the OpenAI-compatible chat
// completions API that OpenRouter exposes, including image content parts.
export class OpenRouterProvider implements ModelProvider {
@@ -370,6 +356,18 @@ export class OpenRouterProvider implements ModelProvider {
// TruncatedResponseError exists to prevent, by a different road.
if (expired) throw stalled(expired);
if (!sawDone && !finishReason) {
+ // Nothing arrived at all: a transient failure this loop can answer, rather than a
+ // document cut short, which it cannot. Raised as the shared type so the retry
+ // below recognizes it and so both adapters describe the same event the same way —
+ // see `EmptyStreamError` and providers/bedrock.ts.
+ if (!text) {
+ throw new EmptyStreamError({
+ provider: this.name,
+ model: req.model,
+ attempts: attempt,
+ detail: "no [DONE] and no finish_reason",
+ });
+ }
throw new Error(
`openrouter: the response stream ended without completing (${text.length} chars ` +
`received, no [DONE] and no finish_reason). Treating a partial document as a whole ` +
@@ -410,7 +408,19 @@ export class OpenRouterProvider implements ModelProvider {
// is the case the retry was added for (a proxy resetting a large request
// body, which happens before any output), and it keeps the loop from
// re-billing a long generation that died three quarters of the way through.
- if (attempt < MAX_ATTEMPTS && !text && isTransientNetworkError(e)) {
+ //
+ // `EmptyStreamError` joins the set for issue #480, which was reported on Bedrock:
+ // a stream that opened and closed having sent nothing is the same transient
+ // upstream event as a reset, arriving as a clean 200 instead of a socket error, so
+ // nothing about `isTransientNetworkError` was ever going to recognize it. The `!text`
+ // guard is already exactly its condition — that error is raised only when no
+ // character arrived — and is left in the condition rather than leaned on, because
+ // what makes this retry safe should be readable on the line that decides it.
+ if (
+ attempt < MAX_ATTEMPTS &&
+ !text &&
+ (isTransientNetworkError(e) || e instanceof EmptyStreamError)
+ ) {
await sleep(400 * 2 ** (attempt - 1));
continue;
}
diff --git a/src/providers/types.ts b/src/providers/types.ts
index 3bd22d05..7bb257e8 100644
--- a/src/providers/types.ts
+++ b/src/providers/types.ts
@@ -98,6 +98,33 @@ export interface Usage {
cache_creation_input_tokens?: number;
}
+// Add two usage snapshots. Used across retry attempts, where the counts ADD rather
+// than replace: an attempt that reported tokens and was then abandoned was still
+// billed for them, so reporting only the surviving attempt understates the call — and
+// understates it invisibly, since `tokens.calls_reported` would still say the call was
+// fully accounted for.
+//
+// Absent stays absent when neither side reported: a 0 nobody sent reads as a free
+// half of the call rather than an unreported one.
+//
+// Here rather than in one adapter because both retry now, and the two must not disagree
+// about what a re-sent call cost. Within ONE attempt the counts replace rather than add —
+// the Anthropic stream reports the prompt's half in `message_start` and the output's half
+// at the end, and adding those would double whichever field arrived twice. Across attempts
+// they add. That is the whole distinction, and it is why this is not the only merge in
+// either adapter.
+export function addUsage(a?: Usage, b?: Usage): Usage | undefined {
+ if (!a) return b;
+ if (!b) return a;
+ const sum: Usage = { ...a };
+ for (const key of Object.keys(b) as (keyof Usage)[]) {
+ const v = b[key];
+ if (v == null) continue;
+ sum[key] = (sum[key] ?? 0) + v;
+ }
+ return sum;
+}
+
// A fact an adapter learned while serving one call that only the CALLER can record. Same
// shape of problem as `onUsage`, and unreportable for the same reason a return value cannot
// carry it: it is learned mid-call, it is worth having whether the call then succeeds or
@@ -441,3 +468,52 @@ export class StalledStreamError extends Error {
this.chars = args.chars;
}
}
+
+// A streamed call whose stream opened, delivered NOTHING, and closed without ever saying
+// the message was over. Reported to a user as "Conversion failed" on a whole document,
+// which is what issue #480 was filed about.
+//
+// Its own type, apart from the partial-response failure it used to share a message with,
+// because the two are opposite diagnoses:
+//
+// - A stream that ends after 30,000 characters has a document in hand that is missing
+// its end. Returning it delivers content the source never had, and re-sending it
+// means either discarding what was generated or resuming mid-document. Neither is on
+// offer, so it fails.
+// - A stream that ends after nothing has no document in hand at all. There is nothing to
+// discard, nothing to resume, and nothing that can ship short — so the call can simply
+// be sent again, which is what both adapters now do. It is the upstream ending a 200
+// response before saying anything, and no part of the request is what it objected to.
+//
+// The old message told the second story as the first: an operator reading "treating a
+// partial document as a whole one" about a response of zero characters is being pointed at
+// a truncation that did not happen.
+//
+// `attempts` is what makes the surviving message honest about cost — a re-sent call was
+// billed for its prompt more than once — and is the one line in a run log that says the
+// retry was reached and did not help.
+//
+// Short on purpose. This message is what the demo reads out in a live region, followed by
+// its own "You can try again." (public/demo.html, `failureMessage`), so it says what
+// happened and stops. Advice to send it again would be said twice, and the reasoning above
+// is for whoever reads this file, not for someone whose document just failed.
+export class EmptyStreamError extends Error {
+ readonly provider: string;
+ readonly model: string;
+ readonly attempts: number;
+
+ constructor(args: { provider: string; model: string; attempts: number; detail: string }) {
+ // "ended without completing" is kept from the old message on purpose: it is what anyone
+ // searching run logs for this failure already searches for.
+ super(
+ `${args.provider}: the response stream ended without completing on ${args.model}, ` +
+ `having sent nothing (${args.detail}).` +
+ (args.attempts > 1 ? ` Sent ${args.attempts} times, and each ended the same way.` : ``) +
+ ` Nothing partial was kept.`,
+ );
+ this.name = "EmptyStreamError";
+ this.provider = args.provider;
+ this.model = args.model;
+ this.attempts = args.attempts;
+ }
+}
diff --git a/test/bedrock-converse.test.ts b/test/bedrock-converse.test.ts
index 9a847462..d7a01044 100644
--- a/test/bedrock-converse.test.ts
+++ b/test/bedrock-converse.test.ts
@@ -134,8 +134,19 @@ test("the default deployment still sends the Anthropic body, unchanged", async (
const bedrock = new BedrockProvider({ default_model: MODEL });
const captured = stubConverse(bedrock, script([]));
// Deliberately an empty script: what is asserted is the command, and an empty stream
- // fails the completeness check afterwards, which is the existing path's behaviour.
- await assert.rejects(() => bedrock.complete(req()), /ended without completing/);
+ // fails the completeness check afterwards. Since #480 that means it is sent twice and
+ // then fails, so `captured` holds the second send's command — the same command, built
+ // the same way, which is the whole point of this assertion.
+ //
+ // The warning the retry prints is swallowed here rather than asserted: this test is about
+ // the request body, and empty-stream-retry.test.ts is where that warning is pinned.
+ const warn = console.warn;
+ console.warn = () => {};
+ try {
+ await assert.rejects(() => bedrock.complete(req()), /ended without completing/);
+ } finally {
+ console.warn = warn;
+ }
assert.ok(captured.command instanceof InvokeModelWithResponseStreamCommand);
const body = JSON.parse(String(captured.input.body));
assert.equal(body.anthropic_version, "bedrock-2023-05-31");
diff --git a/test/demo-error-sentence.test.ts b/test/demo-error-sentence.test.ts
new file mode 100644
index 00000000..0aaec785
--- /dev/null
+++ b/test/demo-error-sentence.test.ts
@@ -0,0 +1,120 @@
+// The sentence a user gets when their document failed to convert.
+//
+// It is the only thing a visitor is ever told about a failure, it is put into a live region
+// (`setError`), and it is assembled from two halves that know nothing about each other: an
+// error message written in src/ and a fixed "You can try again." written here. Issue #480
+// quoted the seam — "…content the source never had.. You can try again." — which is what a
+// screen reader reads as a stop, a pause, and a new sentence.
+//
+// The function is lifted out of the inline script rather than copied, the same way
+// test/demo-tally.test.ts lifts `qualityClause`: a copy would keep passing after the page
+// changed, which is the one thing this must not do.
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { readFileSync } from "node:fs";
+import { fileURLToPath } from "node:url";
+import { dirname, join } from "node:path";
+import { EmptyStreamError, StalledStreamError, TruncatedResponseError } from "../src/providers/types.ts";
+
+const repoRoot = join(dirname(fileURLToPath(import.meta.url)), "..");
+const demoHtml = readFileSync(join(repoRoot, "public", "demo.html"), "utf8");
+
+// Take `function failureMessage(...) { ... }` from the page by matching its braces. It
+// touches no DOM and no globals, which is what makes evaluating it in isolation honest
+// rather than a re-implementation.
+function extract(name: string): string {
+ const start = demoHtml.indexOf(`function ${name}(`);
+ assert.notEqual(start, -1, `${name} is no longer in public/demo.html`);
+ let depth = 0;
+ for (let i = demoHtml.indexOf("{", start); i < demoHtml.length; i++) {
+ if (demoHtml[i] === "{") depth++;
+ else if (demoHtml[i] === "}" && --depth === 0) return demoHtml.slice(start, i + 1);
+ }
+ throw new Error(`unbalanced braces reading ${name} from public/demo.html`);
+}
+
+const failureMessage = new Function(`${extract("failureMessage")}; return failureMessage;`)() as (
+ error: unknown,
+) => string;
+
+test("a message that already ends a sentence is not given a second full stop", () => {
+ assert.equal(
+ failureMessage("bedrock: the model refused the request."),
+ "Conversion failed: bedrock: the model refused the request. You can try again.",
+ );
+ assert.ok(!failureMessage("something went wrong.").includes(".."));
+});
+
+test("a message that ends mid-sentence is punctuated, so the next sentence starts cleanly", () => {
+ assert.equal(
+ failureMessage("bedrock: no output arrived within 120s"),
+ "Conversion failed: bedrock: no output arrived within 120s. You can try again.",
+ );
+});
+
+test("an ellipsis, a question and a quoted ending are all already finished", () => {
+ // Deliberately not "is the last character a period": a message can end its sentence in
+ // more than one way, and adding a stop after any of these reads as a typo rather than as
+ // punctuation.
+ for (const why of [
+ "the upstream gave up…",
+ 'the model stopped for "refusal".',
+ "bedrock: decrease input length or `max_tokens` and try again.",
+ "openrouter: is the model name right?",
+ ]) {
+ // The exact string, not "contains no `..`": `…` is one character, so `….` never contains
+ // `..` and that check passed on the very input it names.
+ assert.equal(failureMessage(why), `Conversion failed: ${why} You can try again.`, why);
+ }
+});
+
+test("the #480 failure is read out once, briefly, and says to try again only once", () => {
+ // This is what a screen-reader user hears in a live region. A retry that fails twice used to
+ // announce 97 words ending in two ways of saying "try again".
+ const e = new EmptyStreamError({
+ provider: "bedrock",
+ model: "us.openai.gpt-5.6-luna",
+ attempts: 2,
+ detail: "no message_stop and no stop_reason",
+ });
+ const said = failureMessage(e.message);
+ assert.equal(said.match(/again/gi)?.length, 1, said);
+ const words = said.split(/\s+/).length;
+ assert.ok(words <= 40, `${words} words: ${said}`);
+});
+
+test("no error, an empty one, or a blank one still says something", () => {
+ // A `failed` session with no `error` recorded is the shape this fallback exists for, and
+ // "Conversion failed: . You can try again." would be the alternative.
+ for (const nothing of [undefined, null, "", " "]) {
+ assert.equal(failureMessage(nothing), "Conversion failed: unknown error. You can try again.");
+ }
+});
+
+test("every failure Iris raises for itself lands on the page as one sentence, then another", () => {
+ // The real inputs, from the types that write them, rather than strings invented here: the
+ // seam only stays fixed if the messages the pipeline actually produces are the ones that
+ // pass through it. #480's own message is the first of these.
+ const errors = [
+ new EmptyStreamError({
+ provider: "bedrock",
+ model: "us.anthropic.claude-sonnet-4-6",
+ attempts: 2,
+ detail: "no message_stop and no stop_reason",
+ }),
+ new StalledStreamError({
+ provider: "bedrock",
+ model: "us.anthropic.claude-sonnet-4-6",
+ kind: "first_output",
+ limitMs: 120_000,
+ chars: 0,
+ }),
+ new TruncatedResponseError("openrouter", "m", 32_000, "cut"),
+ ];
+ for (const e of errors) {
+ const said = failureMessage(e.message);
+ assert.ok(!said.includes(".."), said);
+ assert.match(said, /^Conversion failed: /);
+ assert.match(said, /You can try again\.$/);
+ }
+});
diff --git a/test/empty-stream-retry.test.ts b/test/empty-stream-retry.test.ts
new file mode 100644
index 00000000..0c51a943
--- /dev/null
+++ b/test/empty-stream-retry.test.ts
@@ -0,0 +1,343 @@
+// A response stream that opens, sends nothing, and closes (issue #480: a user's document
+// failed with "0 chars received, no message_stop and no stop_reason", and the whole
+// conversion was lost).
+//
+// The distinction every test here turns on is between a stream that ended SHORT and one
+// that ended EMPTY. They used to share one message and one outcome, and they are opposite
+// diagnoses:
+//
+// - Ended short: a document is in hand, missing its end. Sending it again would have to
+// discard what was generated or resume mid-document, so it fails — and must keep
+// failing, which is what the "not retried" tests below pin.
+// - 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.
+import { test } from "node:test";
+import assert from "node:assert/strict";
+import { BedrockProvider } from "../src/providers/bedrock.ts";
+import { OpenRouterProvider } from "../src/providers/openrouter.ts";
+import { EmptyStreamError, StalledStreamError, type Usage } from "../src/providers/types.ts";
+
+const encode = (o: unknown) => new TextEncoder().encode(JSON.stringify(o));
+const messageStart = (usage?: Record) => ({
+ chunk: { bytes: encode({ type: "message_start", message: usage ? { usage } : {} }) },
+});
+const textDelta = (text: string) => ({
+ chunk: { bytes: encode({ type: "content_block_delta", delta: { type: "text_delta", text } }) },
+});
+const messageDelta = (stop_reason: string, usage?: Record) => ({
+ chunk: { bytes: encode({ type: "message_delta", delta: { stop_reason }, ...(usage ? { usage } : {}) }) },
+});
+
+// Replace the adapter's SDK client with one that serves a fresh scripted stream per send
+// and counts the sends. The count IS the assertion in most of these tests: whether the
+// request was sent again is not visible in the result, only in how many times the upstream
+// was asked.
+//
+// `events` is a function of the send number (1-based), because what makes the retry
+// meaningful is that the second attempt can go differently from the first.
+function stubSends(
+ bedrock: BedrockProvider,
+ events: (send: number) => unknown[],
+ key: "body" | "stream" = "body",
+): { count: () => number } {
+ let sends = 0;
+ (bedrock as unknown as { client: unknown }).client = {
+ send: async () => {
+ const script = events(++sends);
+ return {
+ [key]: (async function* () {
+ for (const e of script) yield e;
+ })(),
+ };
+ },
+ };
+ return { count: () => sends };
+}
+
+const bedrockReq = {
+ capability: "vision" as const,
+ model: "us.anthropic.claude-sonnet-4-6",
+ messages: [{ role: "user" as const, content: "fix this document" }],
+};
+
+// Warnings are captured rather than left to print, matching bedrock-output-ceiling.test.ts:
+// the retry says something an operator has to act on, so it is asserted in the first test
+// below and kept out of the suite's output in the rest.
+async function capturingWarnings(body: () => Promise): Promise<[T, string[]]> {
+ const said: string[] = [];
+ const original = console.warn;
+ console.warn = (...args: unknown[]) => said.push(args.join(" "));
+ try {
+ return [await body(), said];
+ } finally {
+ console.warn = original;
+ }
+}
+
+test("a Bedrock stream that closes having sent nothing is sent again, and the retry is delivered", async () => {
+ const bedrock = new BedrockProvider({ default_model: "m" });
+ // The shape of the reported failure: the stream opens and closes with no events at all.
+ const sends = stubSends(bedrock, (send) =>
+ send === 1 ? [] : [messageStart(), textDelta("Whole
"), messageDelta("end_turn")],
+ );
+ const [res, said] = await capturingWarnings(() => bedrock.complete(bedrockReq));
+ assert.equal(res.text, "Whole
");
+ assert.equal(sends.count(), 2);
+ // A call that succeeds on its second attempt is otherwise invisible — the `model_call`
+ // line reports one call with a longer duration — so the warning is the run's only record
+ // that this deployment is meeting the failure at all, and its frequency is the whole
+ // question a reader would have.
+ assert.equal(said.length, 1);
+ assert.match(said[0], /sent nothing at all/);
+ assert.match(said[0], /paid for, a second time/);
+});
+
+test("a Bedrock stream that closes after message_start alone is still empty, and still retried", async () => {
+ // The likelier shape of #480 on this API, and the one that decides whether the retry is
+ // reachable at all in production: `message_start` has arrived, so the prompt has been read
+ // and billed, and only then does the stream close. Nothing was GENERATED, which is the
+ // condition — a guard written on "has this call cost anything" instead would refuse to
+ // retry exactly the case the issue was filed about.
+ const bedrock = new BedrockProvider({ default_model: "m" });
+ const sends = stubSends(bedrock, (send) =>
+ send === 1
+ ? [messageStart({ input_tokens: 900 })]
+ : [messageStart({ input_tokens: 900 }), textDelta("second time
"), messageDelta("end_turn")],
+ );
+ const [res] = await capturingWarnings(() => bedrock.complete(bedrockReq));
+ assert.equal(res.text, "second time
");
+ assert.equal(sends.count(), 2);
+});
+
+test("the abandoned attempt's tokens stay in the call's reported usage", async () => {
+ // What a retry must not do quietly: bill two prompts and report one. `tokens.calls_reported`
+ // would still count this call as fully accounted for, so the undercount would not show up
+ // anywhere as a gap — it would just make the run cheaper than it was.
+ const bedrock = new BedrockProvider({ default_model: "m" });
+ stubSends(bedrock, (send) =>
+ send === 1
+ ? [messageStart({ input_tokens: 900, cache_read_input_tokens: 100 })]
+ : [
+ messageStart({ input_tokens: 900, cache_read_input_tokens: 100 }),
+ textDelta("ok
"),
+ messageDelta("end_turn", { output_tokens: 40 }),
+ ],
+ );
+ const reported: Usage[] = [];
+ const [res] = await capturingWarnings(() =>
+ bedrock.complete({ ...bedrockReq, onUsage: (u) => reported.push(u) }),
+ );
+ // Both prompts, once each, and the output of the attempt that produced some.
+ assert.deepEqual(res.usage, {
+ input_tokens: 1800,
+ cache_read_input_tokens: 200,
+ output_tokens: 40,
+ });
+ // The router reads usage off the callback when a call throws and off the result when it
+ // returns, so the last thing the callback said has to agree with the result. A call whose
+ // retry then truncated would be reported entirely through the callback.
+ assert.deepEqual(reported.at(-1), res.usage);
+});
+
+test("two empty Bedrock streams fail, saying so, and do not describe a document that never arrived", async () => {
+ const bedrock = new BedrockProvider({ default_model: "m" });
+ const sends = stubSends(bedrock, () => []);
+ await capturingWarnings(() =>
+ assert.rejects(
+ () => bedrock.complete(bedrockReq),
+ (e: Error) => {
+ assert.ok(e instanceof EmptyStreamError);
+ assert.equal(e.attempts, 2);
+ // The retry was reached and did not help. Without this the surviving message is the
+ // second attempt's own, which says nothing about the first.
+ assert.match(e.message, /Sent 2 times/);
+ assert.match(e.message, /ended without completing/);
+ assert.match(e.message, /having sent nothing/);
+ // The old message's wording, which was the actual defect in #480: an operator
+ // reading "a partial document" about a response of zero characters is being
+ // pointed at a truncation that did not happen.
+ assert.doesNotMatch(e.message, /partial document/);
+ return true;
+ },
+ ),
+ );
+ assert.equal(sends.count(), 2);
+});
+
+test("a Bedrock stream that ends SHORT is not retried, and still reports what it received", async () => {
+ // The safety pin. Text in hand means a retry would either discard it or deliver the same
+ // passage twice, so this failure stays exactly as it was — one send, and a message naming
+ // the characters that arrived.
+ const bedrock = new BedrockProvider({ default_model: "m" });
+ const sends = stubSends(bedrock, () => [messageStart(), textDelta("| half a document")]);
+ await assert.rejects(
+ () => bedrock.complete(bedrockReq),
+ (e: Error) => {
+ assert.ok(!(e instanceof EmptyStreamError));
+ assert.match(e.message, /30 chars received/);
+ assert.match(e.message, /partial document/);
+ return true;
+ },
+ );
+ 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.
+ 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++;
+ 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.
+ await new Promise((resolve) =>
+ opts.abortSignal.addEventListener("abort", () => resolve(), { once: true }),
+ );
+ })(),
+ };
+ },
+ };
+ await assert.rejects(() => bedrock.complete(bedrockReq), (e: Error) => {
+ assert.ok(e instanceof StalledStreamError);
+ assert.equal(e.kind, "first_output");
+ return true;
+ });
+ assert.equal(sends, 1);
+});
+
+test("the Converse path retries an empty stream too", async () => {
+ // Both APIs go through one `stream`, so this is pinning that the retry sits above the
+ // dialect rather than inside one of them — and the events differ enough between them
+ // (`stream` rather than `body`, `messageStop` rather than `message_stop`) that a retry
+ // wired into the Anthropic path alone would pass every test above and fail here.
+ const bedrock = new BedrockProvider({ default_model: "m", api: "converse" } as never);
+ const sends = stubSends(
+ bedrock,
+ (send) =>
+ send === 1
+ ? []
+ : [
+ { messageStart: { role: "assistant" } },
+ { contentBlockDelta: { delta: { text: " converse " }, contentBlockIndex: 0 } },
+ { messageStop: { stopReason: "end_turn" } },
+ ],
+ "stream",
+ );
+ const [res] = await capturingWarnings(() => bedrock.complete(bedrockReq));
+ assert.equal(res.text, "converse ");
+ assert.equal(sends.count(), 2);
+});
+
+// --- OpenRouter: the same event, arriving as a 200 with an empty body --------
+
+const sseDelta = (content: string) => `data: ${JSON.stringify({ choices: [{ delta: { content } }] })}`;
+const sseFinish = (finish_reason: string) =>
+ `data: ${JSON.stringify({ choices: [{ delta: {}, finish_reason }] })}`;
+const SSE_DONE = "data: [DONE]";
+
+// Swap global fetch for one that serves a fresh canned SSE body per call and counts them.
+async function withFetch(
+ lines: (call: number) => string[],
+ fn: (calls: () => number) => Promise,
+): Promise {
+ const original = globalThis.fetch;
+ let calls = 0;
+ globalThis.fetch = (async () => {
+ const body = lines(++calls).join("\n\n") + "\n\n";
+ return {
+ ok: true,
+ status: 200,
+ text: async () => "",
+ body: (async function* () {
+ yield new TextEncoder().encode(body);
+ })(),
+ };
+ }) as unknown as typeof fetch;
+ try {
+ return await fn(() => calls);
+ } finally {
+ globalThis.fetch = original;
+ }
+}
+
+const openrouter = () =>
+ new OpenRouterProvider({
+ api_key: "test-key",
+ base_url: "http://localhost:1/v1",
+ default_model: "m",
+ });
+
+const openrouterReq = {
+ capability: "text" as const,
+ model: "m",
+ messages: [{ role: "user" as const, content: "hi" }],
+};
+
+test("an OpenRouter stream that ends with no events is retried, not failed", async () => {
+ // #480 was reported on Bedrock, but nothing about it is Bedrock's: a 200 whose body says
+ // nothing is the same transient upstream event as the connection reset this loop already
+ // retried, and `isTransientNetworkError` was never going to recognize it because no socket
+ // error was raised. Fixing one adapter and not the other would leave the two disagreeing
+ // about whether an empty response is fatal.
+ await withFetch(
+ (call) => (call === 1 ? [] : [sseDelta("whole "), sseFinish("stop"), SSE_DONE]),
+ async (calls) => {
+ const res = await openrouter().complete(openrouterReq);
+ assert.equal(res.text, "whole ");
+ assert.equal(calls(), 2);
+ },
+ );
+});
+
+test("an OpenRouter stream that ends SHORT is not retried", async () => {
+ await withFetch(
+ () => [sseDelta("half")],
+ async (calls) => {
+ await assert.rejects(
+ () => openrouter().complete(openrouterReq),
+ (e: Error) => {
+ assert.ok(!(e instanceof EmptyStreamError));
+ assert.match(e.message, /7 chars received/);
+ return true;
+ },
+ );
+ assert.equal(calls(), 1);
+ },
+ );
+});
+
+test("an OpenRouter upstream that answers three times with nothing fails, naming the attempts", async () => {
+ // Three, not two, because this adapter's retry budget is its own (MAX_ATTEMPTS) and the
+ // empty stream joins it rather than bringing a budget of its own.
+ await withFetch(
+ () => [],
+ async (calls) => {
+ await assert.rejects(
+ () => openrouter().complete(openrouterReq),
+ (e: Error) => {
+ assert.ok(e instanceof EmptyStreamError);
+ assert.equal(e.attempts, 3);
+ assert.match(e.message, /Sent 3 times/);
+ assert.doesNotMatch(e.message, /partial document/);
+ return true;
+ },
+ );
+ assert.equal(calls(), 3);
+ },
+ );
+});
|