feat(ai): add retryTransientCompletion for oneshot LLM calls

completeSimple reports a transient provider failure (Anthropic
overloaded_error / rate_limit_error, HTTP 429/500/502/503/529) by
RESOLVING with stopReason "error" rather than throwing, so callers that
wrapped it in a try/catch-based retry never actually retried.

The helper handles both shapes, classifies with the existing AIError
predicates so the retryable set stays defined in one place, honors
retry-after / x-ratelimit-reset* from response headers or error text, and
returns/rethrows the failure unchanged once attempts are exhausted so
existing caller fallbacks keep working.

Aborts stay aborts: a cancel observed at an attempt boundary surfaces
that attempt's own result, while a cancel during backoff rejects with the
abort reason instead of being relabelled as the provider failure.
This commit is contained in:
wonjun3991
2026-08-13 10:51:58 +09:00
parent 06aecdd51f
commit b55dcd01ff
4 changed files with 494 additions and 0 deletions
+3
View File
@@ -74,6 +74,9 @@
- Fixed OpenAI Codex usage telemetry blocking explicitly allowed ChatGPT Team credentials when a weekly `used_percent` rounded to 100, which could route multi-account sessions to an actually exhausted sibling instead ([#7617](https://github.com/can1357/oh-my-pi/issues/7617)).
- Fixed OpenAI Codex GPT-5.x requests sending optional `reasoning.summary`, `reasoning.context`, and `text.verbosity` controls by default, reducing Codex `server_error` disconnects from unsupported request shapes. ([#4949](https://github.com/can1357/oh-my-pi/issues/4949))
- Classified concurrent-request caps separately from quota exhaustion so they use a short retry backoff without burning a credential, and rotate credentials for account-scoped 403 caps such as Devin's overall message limit.
### Added
- Added `retryTransientCompletion` for oneshot (non-agent-loop) LLM calls. `completeSimple` reports a transient provider failure — Anthropic `overloaded_error` / `rate_limit_error`, HTTP 429/500/502/503/529 — by **resolving** with `stopReason: "error"` rather than throwing, so callers that wrapped it in a try/catch-based retry never actually retried. The helper handles both shapes (resolved error-stop and thrown error), classifies with the existing `AIError` predicates so the retryable set stays defined in one place, honors `retry-after` / `x-ratelimit-reset*` from response headers (supplied via `getResponseHeaders`, since an `AssistantMessage` carries none) or from the error text, and returns/rethrows the failure unchanged once attempts are exhausted so existing caller fallbacks keep working. Aborts surface immediately.
## [17.2.7] - 2026-08-03
+1
View File
@@ -6,6 +6,7 @@ export * from "./auth-gateway/types";
export * from "./auth-retry";
export * from "./auth-storage";
export * from "./error/rate-limit";
export * from "./oneshot-retry";
export * from "./provider-details";
export * from "./providers/anthropic";
export * from "./providers/anthropic-client";
+213
View File
@@ -0,0 +1,213 @@
import { extractRetryHint } from "@oh-my-pi/pi-utils";
import * as AIError from "./error";
import type { AssistantMessage } from "./types";
import { getHeadersFromError, getRetryAfterMsFromHeaders, type HeadersLike } from "./utils/retry-after";
/**
* Transient-failure retry for **oneshot** (non-agent-loop) completions.
*
* Why this exists: `streamSimple`/`completeSimple` retry *auth* failures
* (credential rotation) but deliberately surface *transient* provider failures
* — Anthropic `overloaded_error`, `rate_limit_error`, HTTP 429/500/502/503/529
* — as a **resolved** `AssistantMessage` with `stopReason: "error"`. For the
* main agent turn that is correct: `TurnRecovery` owns recovery there, and it
* must refuse to replay once tool calls or visible text have streamed.
*
* Oneshots have no such hazard. A summary, title, handoff, or image
* description produces no side effects, so re-issuing the whole request is
* safe and is almost always what the caller wants. Before this helper every
* oneshot call site had to re-implement that decision, and most did not —
* failing on the first blip, or swallowing it into `null` so a transient
* overload was indistinguishable from a legitimate empty result.
*
* Classification reuses the existing provider predicates (`AIError`), so the
* set of retryable Anthropic failures stays defined in exactly one place.
* Usage limits are included: unlike the provider loop — which excludes them so
* credential rotation can own them — a oneshot has no rotation layer above it,
* and the retry hint the provider supplies (`retry-after`, "try again in ~5m")
* is honored, so waiting is the correct response.
*/
export interface OneshotRetryOptions {
/** Total attempts, including the first. Default 3. Values < 1 are treated as 1. */
maxAttempts?: number;
/** First backoff step in ms; doubles per attempt. Default 500. */
baseDelayMs?: number;
/**
* Upper bound for a single wait. Default 30_000. A provider retry hint
* longer than this aborts the retry instead of parking the caller — the
* error surfaces so higher-level recovery (or the user) can decide.
*/
maxDelayMs?: number;
/**
* Stops further attempts. Two distinct paths, both preserving the caller's
* intent: an abort already visible when an attempt settles surfaces that
* attempt's own result (`completeSimple` reports `stopReason: "aborted"`),
* while an abort that lands during the backoff wait rejects with the abort
* reason — a user cancel stays a cancel and is never relabelled as the
* provider failure we happened to be waiting on.
*
* This helper does NOT pass the signal into `run` — cancelling the in-flight
* request is the closure's job, because a per-attempt deadline must be
* rebuilt on every attempt. Construct it inside `run`
* (`signal: AbortSignal.timeout(MS)`, or `AbortSignal.any([outer, perAttempt])`);
* a deadline captured outside would fire once and then abort every retry,
* silently turning this helper into a single attempt.
*/
signal?: AbortSignal;
/**
* Headers of the attempt that just failed, used to honor `retry-after`.
*
* Load-bearing: a transient Anthropic failure arrives as a **resolved**
* `AssistantMessage`, and `AssistantMessage` carries no headers — so without
* this the real `retry-after` / `x-ratelimit-reset` values on a 429/529 are
* invisible and only the (usually hint-free) error text is available.
* Callers that already capture headers via `SimpleStreamOptions.onResponse`
* should return the latest capture here; it is read once per failed attempt.
* Thrown errors need no wiring — headers are recovered from the error itself.
*/
getResponseHeaders?: () => HeadersLike;
/** Observability hook. Fires immediately before sleeping. */
onRetry?: (info: OneshotRetryInfo) => void;
}
export interface OneshotRetryInfo {
/** 1-based index of the attempt that just failed. */
attempt: number;
maxAttempts: number;
delayMs: number;
/** True when `delayMs` came from a provider retry hint rather than backoff. */
fromRetryHint: boolean;
errorMessage: string;
/** `AIError` classification bits of the failure. */
errorId: number;
}
const DEFAULT_MAX_ATTEMPTS = 3;
const DEFAULT_BASE_DELAY_MS = 500;
const DEFAULT_MAX_DELAY_MS = 30_000;
/** Cap on pure backoff growth. A provider hint may still exceed this, up to `maxDelayMs`. */
const BACKOFF_CEILING_MS = 8_000;
function backoffDelayMs(attempt: number, baseDelayMs: number): number {
const growth = Math.min(baseDelayMs * 2 ** (attempt - 1), BACKOFF_CEILING_MS);
// 75-100% jitter, matching the provider loop and TurnRecovery, so a fleet of
// concurrent oneshots does not re-converge on the same instant.
return Math.round(growth * (0.75 + Math.random() * 0.25));
}
/** Retryable when the provider says transient, or when it says "wait, then retry". */
function isRetryableOneshotFailure(errorId: number): boolean {
return (
AIError.is(errorId, AIError.Flag.Transient) ||
AIError.is(errorId, AIError.Flag.UsageLimit) ||
AIError.retriable(errorId)
);
}
function sleep(delayMs: number, signal?: AbortSignal): Promise<void> {
if (delayMs <= 0) return Promise.resolve();
const { promise, resolve, reject } = Promise.withResolvers<void>();
const timer = setTimeout(() => {
signal?.removeEventListener("abort", onAbort);
resolve();
}, delayMs);
const onAbort = () => {
clearTimeout(timer);
reject(signal?.reason ?? new AIError.AbortError("oneshot retry aborted"));
};
if (signal) {
if (signal.aborted) {
clearTimeout(timer);
return Promise.reject(signal.reason ?? new AIError.AbortError("oneshot retry aborted"));
}
signal.addEventListener("abort", onAbort, { once: true });
}
return promise;
}
/**
* Run a oneshot completion, retrying transient provider failures.
*
* Handles both failure shapes: a resolved `AssistantMessage` carrying
* `stopReason: "error"` (what `completeSimple` produces) and a thrown error
* (what the raw HTTP helpers produce). A non-retryable failure is returned or
* rethrown unchanged, so existing caller error handling keeps working — this
* only removes the *first-blip* failure mode.
*/
export async function retryTransientCompletion(
run: (attempt: number) => Promise<AssistantMessage>,
options?: OneshotRetryOptions,
): Promise<AssistantMessage> {
const maxAttempts = Math.max(1, options?.maxAttempts ?? DEFAULT_MAX_ATTEMPTS);
const baseDelayMs = options?.baseDelayMs ?? DEFAULT_BASE_DELAY_MS;
const maxDelayMs = options?.maxDelayMs ?? DEFAULT_MAX_DELAY_MS;
const signal = options?.signal;
for (let attempt = 1; ; attempt++) {
let message: AssistantMessage | undefined;
let thrown: unknown;
try {
message = await run(attempt);
if (message.stopReason !== "error") return message;
} catch (error) {
thrown = error;
}
// A caller abort is never a transient failure — surface it immediately so
// cancellation stays responsive.
if (signal?.aborted) {
if (thrown !== undefined) throw thrown;
return message as AssistantMessage;
}
const errorId =
thrown !== undefined ? AIError.classify(thrown) : AIError.classifyMessage(message as AssistantMessage);
if (AIError.is(errorId, AIError.Flag.Abort) || AIError.is(errorId, AIError.Flag.UserInterrupt)) {
if (thrown !== undefined) throw thrown;
return message as AssistantMessage;
}
const errorMessage =
thrown !== undefined
? thrown instanceof Error
? thrown.message
: String(thrown)
: ((message as AssistantMessage).errorMessage ?? "unknown error");
const lastAttempt = attempt >= maxAttempts;
if (lastAttempt || !isRetryableOneshotFailure(errorId)) {
if (thrown !== undefined) throw thrown;
return message as AssistantMessage;
}
// Headers first: a real Anthropic 429/529 carries `retry-after` /
// `x-ratelimit-reset*` in the response, and the resolved AssistantMessage
// has none — the caller supplies them via `getResponseHeaders`. Thrown
// errors (e.g. AnthropicApiError) carry their own headers.
const headers: HeadersLike = thrown !== undefined ? getHeadersFromError(thrown) : options?.getResponseHeaders?.();
const headerHintMs = getRetryAfterMsFromHeaders(headers);
const textHintMs = extractRetryHint(undefined, errorMessage);
const hintMs =
headerHintMs === undefined && textHintMs === undefined
? undefined
: Math.max(headerHintMs ?? 0, textHintMs ?? 0);
// An over-cap hint means "come back much later"; parking a oneshot that
// long is worse than surfacing the failure to the caller.
if (hintMs !== undefined && hintMs > maxDelayMs) {
if (thrown !== undefined) throw thrown;
return message as AssistantMessage;
}
const backoff = backoffDelayMs(attempt, baseDelayMs);
const delayMs = Math.min(Math.max(hintMs ?? 0, backoff), maxDelayMs);
options?.onRetry?.({
attempt,
maxAttempts,
delayMs,
fromRetryHint: hintMs !== undefined && hintMs >= backoff,
errorMessage,
errorId,
});
// Aborting mid-backoff rejects with the caller's abort reason: a user
// cancel must stay a cancel, not get relabelled as the provider failure we
// happened to be waiting on.
await sleep(delayMs, signal);
}
}
+277
View File
@@ -0,0 +1,277 @@
import { describe, expect, it } from "bun:test";
import { retryTransientCompletion } from "@oh-my-pi/pi-ai/oneshot-retry";
import type { AssistantMessage, Usage } from "@oh-my-pi/pi-ai/types";
/**
* Defends the contract every oneshot LLM call site now depends on:
* `completeSimple` reports a transient provider failure by RESOLVING with
* `stopReason: "error"` rather than throwing, so a retry layer that only
* catches exceptions silently never fires. Before this helper, an Anthropic
* `overloaded_error` / 429 / 529 on a summary, title, handoff or image
* description failed on the first blip — or was swallowed into `null`, making a
* transient overload indistinguishable from a legitimate empty result.
*
* These tests pin the four properties the call sites rely on: transient
* error-stops are re-issued, non-transient ones are not, the final failure is
* handed back unchanged (so existing `null`/throw fallbacks still work), and a
* caller abort wins immediately.
*/
const emptyUsage = (): Usage =>
({
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
}) as unknown as Usage;
function message(overrides: Partial<AssistantMessage> = {}): AssistantMessage {
return {
role: "assistant",
content: [],
api: "anthropic-messages",
provider: "anthropic",
model: "claude-sonnet-4-6",
usage: emptyUsage(),
stopReason: "stop",
timestamp: 0,
...overrides,
} as AssistantMessage;
}
const overloaded = (): AssistantMessage =>
message({
stopReason: "error",
errorStatus: 529,
errorMessage: "Anthropic stream error (overloaded_error): Overloaded",
});
const rateLimited = (): AssistantMessage =>
message({ stopReason: "error", errorStatus: 429, errorMessage: "rate_limit_error: too many requests" });
// Keep the suite fast: the helper's real backoff floor is 500ms.
const fast = { baseDelayMs: 1, maxAttempts: 3 } as const;
describe("retryTransientCompletion", () => {
it("re-issues an Anthropic 529 error-stop and returns the eventual success", async () => {
const results = [overloaded(), overloaded(), message({ stopReason: "stop" })];
let calls = 0;
const final = await retryTransientCompletion(() => {
calls += 1;
return Promise.resolve(results.shift()!);
}, fast);
expect(calls).toBe(3);
expect(final.stopReason).toBe("stop");
});
it("re-issues a 429 rate-limit error-stop", async () => {
let calls = 0;
const final = await retryTransientCompletion(() => {
calls += 1;
return Promise.resolve(calls === 1 ? rateLimited() : message());
}, fast);
expect(calls).toBe(2);
expect(final.stopReason).toBe("stop");
});
it("returns the failing message unchanged once attempts are exhausted, so caller fallbacks still apply", async () => {
let calls = 0;
const final = await retryTransientCompletion(() => {
calls += 1;
return Promise.resolve(overloaded());
}, fast);
expect(calls).toBe(3);
expect(final.stopReason).toBe("error");
expect(final.errorMessage).toContain("overloaded_error");
});
it("does not retry a non-transient provider error", async () => {
let calls = 0;
const final = await retryTransientCompletion(
() => {
calls += 1;
return Promise.resolve(
message({
stopReason: "error",
errorStatus: 400,
errorMessage: "invalid_request_error: messages: at least one message is required",
}),
);
},
{ ...fast, maxAttempts: 5 },
);
expect(calls).toBe(1);
expect(final.stopReason).toBe("error");
});
it("retries a thrown transient error and rethrows the last one when exhausted", async () => {
let calls = 0;
const attempt = retryTransientCompletion(() => {
calls += 1;
const error = new Error("529 overloaded_error: Overloaded") as Error & { status?: number };
error.status = 529;
throw error;
}, fast);
await expect(attempt).rejects.toThrow(/overloaded_error/);
expect(calls).toBe(3);
});
it("does not retry a thrown non-transient error", async () => {
let calls = 0;
const attempt = retryTransientCompletion(() => {
calls += 1;
throw new Error("invalid_request_error: bad tool schema");
}, fast);
await expect(attempt).rejects.toThrow(/invalid_request_error/);
expect(calls).toBe(1);
});
it("stops immediately when the caller aborts", async () => {
const controller = new AbortController();
let calls = 0;
const final = await retryTransientCompletion(
() => {
calls += 1;
controller.abort();
return Promise.resolve(overloaded());
},
{ ...fast, signal: controller.signal },
);
expect(calls).toBe(1);
expect(final.stopReason).toBe("error");
});
it("rejects with the abort reason when the caller cancels during backoff", async () => {
// The cancel lands while we are waiting, not while an attempt is in flight:
// it must stay a cancellation rather than being reported as the provider
// failure we happened to be sleeping on.
const controller = new AbortController();
const reason = new Error("user pressed escape");
let calls = 0;
const attempt = retryTransientCompletion(
() => {
calls += 1;
setTimeout(() => controller.abort(reason), 5);
return Promise.resolve(overloaded());
},
{ maxAttempts: 3, baseDelayMs: 200, signal: controller.signal },
);
await expect(attempt).rejects.toThrow("user pressed escape");
expect(calls).toBe(1);
});
it("reports each retry through onRetry so callers can log the wait", async () => {
const seen: number[] = [];
let calls = 0;
await retryTransientCompletion(
() => {
calls += 1;
return Promise.resolve(calls === 1 ? overloaded() : message());
},
{ ...fast, onRetry: info => seen.push(info.attempt) },
);
expect(seen).toEqual([1]);
});
it("surfaces the failure instead of parking when the provider asks for longer than maxDelayMs", async () => {
let calls = 0;
const final = await retryTransientCompletion(
() => {
calls += 1;
return Promise.resolve(
message({
stopReason: "error",
errorStatus: 429,
errorMessage: "rate_limit_error: please retry in 600s",
}),
);
},
{ ...fast, maxDelayMs: 1_000 },
);
expect(calls).toBe(1);
expect(final.stopReason).toBe("error");
});
it("honors a retry-after-ms response header over the backoff floor", async () => {
// The header is the only place a real Anthropic 429 carries its wait: the
// resolved AssistantMessage has no headers, so a helper that reads only the
// error text would silently fall back to plain backoff.
let calls = 0;
let observedDelay = -1;
await retryTransientCompletion(
() => {
calls += 1;
return Promise.resolve(calls === 1 ? rateLimited() : message());
},
{
maxAttempts: 2,
baseDelayMs: 1,
getResponseHeaders: () => ({ "retry-after-ms": "120" }),
onRetry: info => {
observedDelay = info.delayMs;
},
},
);
expect(calls).toBe(2);
expect(observedDelay).toBe(120);
});
it("surfaces the failure when a retry-after header exceeds maxDelayMs", async () => {
let calls = 0;
const final = await retryTransientCompletion(
() => {
calls += 1;
return Promise.resolve(rateLimited());
},
{
maxAttempts: 3,
baseDelayMs: 1,
maxDelayMs: 1_000,
getResponseHeaders: () => ({ "retry-after": "300" }),
},
);
expect(calls).toBe(1);
expect(final.stopReason).toBe("error");
});
it("recovers retry-after from a thrown provider error's own headers", async () => {
let calls = 0;
let observedDelay = -1;
const attempt = retryTransientCompletion(
() => {
calls += 1;
const error = new Error("529 overloaded_error: Overloaded") as Error & {
status?: number;
headers?: Record<string, string>;
};
error.status = 529;
error.headers = { "retry-after-ms": "90" };
throw error;
},
{
maxAttempts: 2,
baseDelayMs: 1,
onRetry: info => {
observedDelay = info.delayMs;
},
},
);
await expect(attempt).rejects.toThrow(/overloaded_error/);
expect(calls).toBe(2);
expect(observedDelay).toBe(90);
});
});