fix(session): handled subscription-cap retry exhaustion
- Classified subscription and plan rate caps as credential-rotatable usage limits while preserving transient per-minute throttles. - Applied reason-specific backoff to transient rate limits and collapsed exhausted retry attempts behind one budget-labeled terminal error. - Covered classification, delay selection, persisted transcript aggregation, and retry event propagation. Fixes #7767
This commit is contained in:
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Classified subscription/plan-cap 429 responses as credential-rotatable usage limits instead of transient per-minute throttles ([#7767](https://github.com/can1357/oh-my-pi/issues/7767)).
|
||||
|
||||
## [17.2.9] - 2026-08-05
|
||||
|
||||
### Fixed
|
||||
|
||||
@@ -22,6 +22,8 @@ const ACCOUNT_RATE_LIMIT_PATTERN =
|
||||
/\baccount(?:'s)?\b[^\n]{0,80}\brate.?limit\b|\brate.?limit\b[^\n]{0,80}\baccount\b/i;
|
||||
const INSUFFICIENT_BALANCE_PATTERN = /insufficient.?balance/i;
|
||||
const SPEND_LIMIT_PATTERN = /spend.?limit/i;
|
||||
const SUBSCRIPTION_CAP_PATTERN =
|
||||
/\b(?:subscription|plan|membership)\b[^\n]{0,80}\b(?:rate.?limits?|quota|cap)\b|\b(?:rate.?limits?|quota|cap)\b[^\n]{0,80}\b(?:subscription|plan|membership)\b/i;
|
||||
const OPENROUTER_DAILY_FREE_LIMIT_PATTERN = /\bfree[-_ ]models[-_ ]per[-_ ]day\b/i;
|
||||
// gRPC/Connect end-streams carry the status as its name (`resource_exhausted`),
|
||||
// while HTTP bodies use the phrase ("resource exhausted"). Strip either form
|
||||
@@ -80,6 +82,10 @@ export function parseRateLimitReason(errorMessage: string): RateLimitReason {
|
||||
return "QUOTA_EXHAUSTED";
|
||||
}
|
||||
|
||||
if (SUBSCRIPTION_CAP_PATTERN.test(errorMessage)) {
|
||||
return "QUOTA_EXHAUSTED";
|
||||
}
|
||||
|
||||
if (OPENROUTER_DAILY_FREE_LIMIT_PATTERN.test(errorMessage)) {
|
||||
return "QUOTA_EXHAUSTED";
|
||||
}
|
||||
@@ -226,6 +232,7 @@ export function matchesUsageLimitText(errorMessage: string): boolean {
|
||||
USAGE_LIMIT_PATTERN.test(errorMessage) ||
|
||||
SPEND_LIMIT_PATTERN.test(errorMessage) ||
|
||||
ACCOUNT_RATE_LIMIT_PATTERN.test(errorMessage) ||
|
||||
SUBSCRIPTION_CAP_PATTERN.test(errorMessage) ||
|
||||
OPENROUTER_DAILY_FREE_LIMIT_PATTERN.test(errorMessage)
|
||||
);
|
||||
}
|
||||
|
||||
+24
-14
@@ -831,22 +831,32 @@ export interface DeveloperMessage {
|
||||
timestamp: number; // Unix timestamp in milliseconds
|
||||
}
|
||||
|
||||
/** How an automatic retry recovered or ultimately settled a failed attempt. */
|
||||
export type AssistantRetryRecoveryKind = "credential" | "model" | "wait" | "plain";
|
||||
|
||||
export interface AssistantRetryRecovery {
|
||||
kind: "auto-retry";
|
||||
status: "recovered";
|
||||
attempt: number;
|
||||
recoveredAt: string;
|
||||
recovery: AssistantRetryRecoveryKind;
|
||||
note: string;
|
||||
supersededBy?: {
|
||||
timestamp: number;
|
||||
responseId?: string;
|
||||
provider: string;
|
||||
model: string;
|
||||
};
|
||||
}
|
||||
/** Persisted presentation state for an assistant error superseded by an automatic retry saga. */
|
||||
export type AssistantRetryRecovery =
|
||||
| {
|
||||
kind: "auto-retry";
|
||||
status: "recovered";
|
||||
attempt: number;
|
||||
recoveredAt: string;
|
||||
recovery: AssistantRetryRecoveryKind;
|
||||
note: string;
|
||||
supersededBy?: {
|
||||
timestamp: number;
|
||||
responseId?: string;
|
||||
provider: string;
|
||||
model: string;
|
||||
};
|
||||
}
|
||||
| {
|
||||
kind: "auto-retry";
|
||||
status: "superseded";
|
||||
attempt: number;
|
||||
recovery: AssistantRetryRecoveryKind;
|
||||
note: string;
|
||||
};
|
||||
|
||||
export interface ContextSnapshot {
|
||||
promptTokens: number; // authoritative provider prompt/input tokens
|
||||
|
||||
@@ -257,6 +257,19 @@ describe("isUsageLimitOutcome", () => {
|
||||
expect(isUsageLimitOutcome(429, "Please retry in 5s")).toBe(false);
|
||||
});
|
||||
|
||||
it("rotates on subscription caps without treating generic rate limits as usage exhaustion", () => {
|
||||
const subscriptionCap =
|
||||
"429 You've exceeded your subscription rate limits. Upgrade, or try again later. You can view your usage at https://api.synthetic.new/usage";
|
||||
expect(parseRateLimitReason(subscriptionCap)).toBe("QUOTA_EXHAUSTED");
|
||||
expect(isUsageLimitOutcome(429, subscriptionCap)).toBe(true);
|
||||
expect(isUsageLimit(Object.assign(new Error(subscriptionCap), { status: 429 }))).toBe(true);
|
||||
|
||||
const transient = "429 Rate limit exceeded, too many requests";
|
||||
expect(parseRateLimitReason(transient)).toBe("RATE_LIMIT_EXCEEDED");
|
||||
expect(isUsageLimitOutcome(429, transient)).toBe(false);
|
||||
expect(isUsageLimit(Object.assign(new Error(transient), { status: 429 }))).toBe(false);
|
||||
});
|
||||
|
||||
it("still rotates on 429 with explicit account rate-limit framing", () => {
|
||||
expect(
|
||||
isUsageLimitOutcome(
|
||||
|
||||
@@ -2,6 +2,10 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Applied reason-specific backoff to transient rate-limit retries and collapsed exhausted retry sagas into one terminal error naming the spent budget ([#7767](https://github.com/can1357/oh-my-pi/issues/7767)).
|
||||
|
||||
## [17.2.9] - 2026-08-05
|
||||
|
||||
### Breaking Changes
|
||||
|
||||
@@ -30,7 +30,7 @@ import type { LocalProtocolOptions } from "../../internal-urls/local-protocol";
|
||||
import type { Theme } from "../../modes/theme/theme";
|
||||
import type { ReadonlySessionManager } from "../../session/session-manager";
|
||||
import type { TodoItem } from "../../tools/todo";
|
||||
import type { RecoveredRetryError } from "../shared-events";
|
||||
import type { RetryErrorUpdate } from "../shared-events";
|
||||
|
||||
/** Alias for clarity */
|
||||
export type CustomToolUIContext = HookUIContext;
|
||||
@@ -139,7 +139,7 @@ export type CustomToolSessionEvent =
|
||||
success: boolean;
|
||||
attempt: number;
|
||||
finalError?: string;
|
||||
recoveredErrors?: RecoveredRetryError[];
|
||||
retryErrors?: RetryErrorUpdate[];
|
||||
}
|
||||
| {
|
||||
reason: "ttsr_triggered";
|
||||
|
||||
@@ -249,7 +249,8 @@ export interface AutoRetryStartEvent {
|
||||
errorId?: number;
|
||||
}
|
||||
|
||||
export interface RecoveredRetryError {
|
||||
/** Persisted retry error whose transcript presentation changed when the retry saga settled. */
|
||||
export interface RetryErrorUpdate {
|
||||
entryId: string;
|
||||
persistenceKey?: string;
|
||||
note: string;
|
||||
@@ -262,7 +263,7 @@ export interface AutoRetryEndEvent {
|
||||
success: boolean;
|
||||
attempt: number;
|
||||
finalError?: string;
|
||||
recoveredErrors?: RecoveredRetryError[];
|
||||
retryErrors?: RetryErrorUpdate[];
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
|
||||
@@ -1896,21 +1896,19 @@ export class EventController {
|
||||
this.ctx.retryLoader = undefined;
|
||||
this.ctx.statusContainer.disposeChildren();
|
||||
}
|
||||
if (event.success) {
|
||||
let appliedRecovered = false;
|
||||
for (const recovered of event.recoveredErrors ?? []) {
|
||||
const component = this.#takeRetrySupersededAssistantComponent(recovered.persistenceKey);
|
||||
if (!component) continue;
|
||||
component.applyRetryRecovery(recovered.retryRecovery);
|
||||
if (this.#pinnedErrorComponent === component) this.#pinnedErrorComponent = undefined;
|
||||
appliedRecovered = true;
|
||||
}
|
||||
if (appliedRecovered || (event.recoveredErrors?.length ?? 0) > 0) {
|
||||
this.ctx.clearPinnedError();
|
||||
}
|
||||
this.#clearRetrySupersededAssistantComponents();
|
||||
} else {
|
||||
this.#clearRetrySupersededAssistantComponents();
|
||||
let appliedRetryUpdate = false;
|
||||
for (const retryError of event.retryErrors ?? []) {
|
||||
const component = this.#takeRetrySupersededAssistantComponent(retryError.persistenceKey);
|
||||
if (!component) continue;
|
||||
component.applyRetryRecovery(retryError.retryRecovery);
|
||||
if (this.#pinnedErrorComponent === component) this.#pinnedErrorComponent = undefined;
|
||||
appliedRetryUpdate = true;
|
||||
}
|
||||
if (appliedRetryUpdate || (event.retryErrors?.length ?? 0) > 0) {
|
||||
this.ctx.clearPinnedError();
|
||||
}
|
||||
this.#clearRetrySupersededAssistantComponents();
|
||||
if (!event.success) {
|
||||
this.ctx.showError(`Retry failed after ${event.attempt} attempts: ${event.finalError || "Unknown error"}`);
|
||||
}
|
||||
this.#ensureWorkingLoaderWhileStreaming();
|
||||
|
||||
@@ -211,13 +211,15 @@ function sanitizeRecoveredRetryNote(note: string): string {
|
||||
|
||||
/**
|
||||
* Resolve the turn-ending assistant error presentation, if any.
|
||||
* Silent and user-interrupt aborts yield no label. Recovered auto-retry errors
|
||||
* collapse to a single non-error note; terminal errors keep the full red presentation.
|
||||
* Silent and user-interrupt aborts yield no label. Recovered retry attempts
|
||||
* render a compact note; attempts superseded by an exhausted budget are hidden
|
||||
* while the final terminal error keeps its full presentation.
|
||||
*/
|
||||
export function resolveAssistantErrorPresentation(
|
||||
message: AssistantAgentMessage,
|
||||
retryAttempt = 0,
|
||||
): AssistantErrorPresentation {
|
||||
if (message.retryRecovery?.status === "superseded") return { kind: "none" };
|
||||
if (message.retryRecovery?.status === "recovered") {
|
||||
return {
|
||||
kind: "compact-recovered",
|
||||
|
||||
@@ -1039,7 +1039,7 @@ function createCustomToolsExtension(tools: CustomTool[]): ExtensionFactory {
|
||||
success: event.success,
|
||||
attempt: event.attempt,
|
||||
finalError: event.finalError,
|
||||
recoveredErrors: event.recoveredErrors,
|
||||
retryErrors: event.retryErrors,
|
||||
},
|
||||
ctx,
|
||||
),
|
||||
|
||||
@@ -2,7 +2,7 @@ import type { AgentEvent, ThinkingLevel } from "@oh-my-pi/pi-agent-core";
|
||||
import type { CompactionResult } from "@oh-my-pi/pi-agent-core/compaction";
|
||||
import type { Effort } from "@oh-my-pi/pi-ai";
|
||||
import type { Rule } from "../capability/rule";
|
||||
import type { RecoveredRetryError } from "../extensibility/shared-events";
|
||||
import type { RetryErrorUpdate } from "../extensibility/shared-events";
|
||||
import type { Goal, GoalModeState } from "../goals/state";
|
||||
import type { ConfiguredThinkingLevel } from "../thinking";
|
||||
import type { TodoItem } from "../tools/todo";
|
||||
@@ -43,7 +43,7 @@ export type AgentSessionEvent =
|
||||
success: boolean;
|
||||
attempt: number;
|
||||
finalError?: string;
|
||||
recoveredErrors?: RecoveredRetryError[];
|
||||
retryErrors?: RetryErrorUpdate[];
|
||||
}
|
||||
| { type: "retry_fallback_applied"; from: string; to: string; role: string }
|
||||
| { type: "retry_fallback_succeeded"; model: string; role: string }
|
||||
|
||||
@@ -3439,7 +3439,7 @@ export class AgentSession {
|
||||
success: event.success,
|
||||
attempt: event.attempt,
|
||||
finalError: event.finalError,
|
||||
recoveredErrors: event.recoveredErrors,
|
||||
retryErrors: event.retryErrors,
|
||||
});
|
||||
} else if (event.type === "ttsr_triggered") {
|
||||
await this.#extensionRunner.emit({ type: "ttsr_triggered", rules: event.rules });
|
||||
|
||||
@@ -332,11 +332,7 @@ export function buildSessionContext(
|
||||
const appendMessage = (entry: SessionEntry) => {
|
||||
handleEntryResetTracking(entry);
|
||||
if (entry.type === "message") {
|
||||
if (
|
||||
!options?.transcript &&
|
||||
entry.message.role === "assistant" &&
|
||||
entry.message.retryRecovery?.status === "recovered"
|
||||
) {
|
||||
if (!options?.transcript && entry.message.role === "assistant" && entry.message.retryRecovery) {
|
||||
return;
|
||||
}
|
||||
pushMessage(entry.message);
|
||||
|
||||
@@ -28,7 +28,7 @@ import type { ModelRegistry } from "../config/model-registry";
|
||||
import { formatModelStringWithRouting, resolveModelOverride } from "../config/model-resolver";
|
||||
|
||||
import type { Settings } from "../config/settings";
|
||||
import type { RecoveredRetryError } from "../extensibility/shared-events";
|
||||
import type { RetryErrorUpdate } from "../extensibility/shared-events";
|
||||
import emptyStopRetryTemplate from "../prompts/system/empty-stop-retry.md" with { type: "text" };
|
||||
import thinkingLoopRedirectTemplate from "../prompts/system/thinking-loop-redirect.md" with { type: "text" };
|
||||
import unexpectedStopRetryTemplate from "../prompts/system/unexpected-stop-retry.md" with { type: "text" };
|
||||
@@ -161,7 +161,7 @@ export interface TurnRecoveryOptions {
|
||||
initialRetryFallback?: InitialRetryFallbackState;
|
||||
}
|
||||
|
||||
type PendingRecoveredRetryError = {
|
||||
type PendingRetryError = {
|
||||
entryId: string;
|
||||
persistenceKey: string;
|
||||
recovery: AssistantRetryRecoveryKind;
|
||||
@@ -184,7 +184,7 @@ export class TurnRecovery {
|
||||
#retryResolve: (() => void) | undefined;
|
||||
#activeRetryFallback: ActiveRetryFallbackState | undefined;
|
||||
#usageReserveApprovedSelector: string | undefined;
|
||||
#pendingRecoveredRetryErrors: PendingRecoveredRetryError[] = [];
|
||||
#pendingRetryErrors: PendingRetryError[] = [];
|
||||
#usageLimitOutcomes = new WeakMap<AssistantMessage, Promise<UsageLimitOutcome>>();
|
||||
#emptyStopRetryCount = 0;
|
||||
#unexpectedStopRetryCount = 0;
|
||||
@@ -250,14 +250,17 @@ export class TurnRecovery {
|
||||
role: this.#activeRetryFallback.role,
|
||||
});
|
||||
}
|
||||
const recoveredErrors = await this.#markPendingRecoveredRetryErrors(message);
|
||||
const retryErrors = await this.#markPendingRetryErrors({
|
||||
status: "recovered",
|
||||
supersedingMessage: message,
|
||||
});
|
||||
await this.#host.emitSessionEvent({
|
||||
type: "auto_retry_end",
|
||||
success: true,
|
||||
attempt: this.#retryAttempt,
|
||||
recoveredErrors,
|
||||
retryErrors,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
this.#retryAttempt = 0;
|
||||
}
|
||||
|
||||
@@ -272,7 +275,7 @@ export class TurnRecovery {
|
||||
attempt,
|
||||
finalError: message.errorMessage,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
}
|
||||
|
||||
/** Persists an otherwise skipped terminal empty error turn. */
|
||||
@@ -381,8 +384,8 @@ export class TurnRecovery {
|
||||
}
|
||||
}
|
||||
|
||||
#clearPendingRecoveredRetryErrors(): void {
|
||||
this.#pendingRecoveredRetryErrors = [];
|
||||
#clearPendingRetryErrors(): void {
|
||||
this.#pendingRetryErrors = [];
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -432,7 +435,7 @@ export class TurnRecovery {
|
||||
return parts.join("; ");
|
||||
}
|
||||
|
||||
async #recordPendingRecoveredRetryError(
|
||||
async #recordPendingRetryError(
|
||||
message: AssistantMessage,
|
||||
id: number,
|
||||
options: { switchedCredential: boolean; switchedModel: boolean; delayMs: number },
|
||||
@@ -451,11 +454,11 @@ export class TurnRecovery {
|
||||
break;
|
||||
}
|
||||
if (!branchEntry) return;
|
||||
if (this.#pendingRecoveredRetryErrors.some(error => error.entryId === branchEntry.id)) return;
|
||||
if (this.#pendingRetryErrors.some(error => error.entryId === branchEntry.id)) return;
|
||||
const rateLimited = AIError.is(id, AIError.Flag.UsageLimit);
|
||||
const recovery = this.#retryRecoveryKind(id, options.switchedCredential, options.switchedModel, options.delayMs);
|
||||
const note = this.#retryRecoveryNote(recovery, rateLimited);
|
||||
this.#pendingRecoveredRetryErrors.push({
|
||||
this.#pendingRetryErrors.push({
|
||||
entryId: branchEntry.id,
|
||||
persistenceKey,
|
||||
recovery,
|
||||
@@ -464,24 +467,17 @@ export class TurnRecovery {
|
||||
});
|
||||
}
|
||||
|
||||
async #markPendingRecoveredRetryErrors(supersedingMessage: AssistantMessage): Promise<RecoveredRetryError[]> {
|
||||
if (this.#pendingRecoveredRetryErrors.length === 0) return [];
|
||||
async #markPendingRetryErrors(
|
||||
completion: { status: "recovered"; supersedingMessage: AssistantMessage } | { status: "superseded" },
|
||||
): Promise<RetryErrorUpdate[]> {
|
||||
if (this.#pendingRetryErrors.length === 0) return [];
|
||||
const branch = this.#host.sessionManager.getBranch();
|
||||
const branchById = new Map<string, SessionEntry>();
|
||||
for (const entry of branch) {
|
||||
branchById.set(entry.id, entry);
|
||||
}
|
||||
const recoveredAt = new Date().toISOString();
|
||||
const supersededBy: AssistantRetryRecovery["supersededBy"] = {
|
||||
timestamp: supersedingMessage.timestamp,
|
||||
provider: supersedingMessage.provider,
|
||||
model: supersedingMessage.model,
|
||||
};
|
||||
if (supersedingMessage.responseId) {
|
||||
supersededBy.responseId = supersedingMessage.responseId;
|
||||
}
|
||||
const recoveredErrors: RecoveredRetryError[] = [];
|
||||
for (const pending of this.#pendingRecoveredRetryErrors) {
|
||||
const retryErrors: RetryErrorUpdate[] = [];
|
||||
for (const pending of this.#pendingRetryErrors) {
|
||||
let entry = branchById.get(pending.entryId);
|
||||
if (entry?.type !== "message" || entry.message.role !== "assistant") {
|
||||
entry = branch
|
||||
@@ -495,27 +491,45 @@ export class TurnRecovery {
|
||||
);
|
||||
}
|
||||
if (entry?.type !== "message" || entry.message.role !== "assistant") continue;
|
||||
const retryRecovery: AssistantRetryRecovery = {
|
||||
kind: "auto-retry",
|
||||
status: "recovered",
|
||||
attempt: pending.attempt,
|
||||
recoveredAt,
|
||||
recovery: pending.recovery,
|
||||
note: pending.note,
|
||||
supersededBy,
|
||||
};
|
||||
let retryRecovery: AssistantRetryRecovery;
|
||||
if (completion.status === "recovered") {
|
||||
retryRecovery = {
|
||||
kind: "auto-retry",
|
||||
status: "recovered",
|
||||
attempt: pending.attempt,
|
||||
recoveredAt: new Date().toISOString(),
|
||||
recovery: pending.recovery,
|
||||
note: pending.note,
|
||||
supersededBy: {
|
||||
timestamp: completion.supersedingMessage.timestamp,
|
||||
...(completion.supersedingMessage.responseId === undefined
|
||||
? {}
|
||||
: { responseId: completion.supersedingMessage.responseId }),
|
||||
provider: completion.supersedingMessage.provider,
|
||||
model: completion.supersedingMessage.model,
|
||||
},
|
||||
};
|
||||
} else {
|
||||
retryRecovery = {
|
||||
kind: "auto-retry",
|
||||
status: "superseded",
|
||||
attempt: pending.attempt,
|
||||
recovery: pending.recovery,
|
||||
note: pending.note,
|
||||
};
|
||||
}
|
||||
entry.message.retryRecovery = retryRecovery;
|
||||
recoveredErrors.push({
|
||||
retryErrors.push({
|
||||
entryId: entry.id,
|
||||
persistenceKey: pending.persistenceKey,
|
||||
note: retryRecovery.note,
|
||||
retryRecovery,
|
||||
});
|
||||
}
|
||||
if (recoveredErrors.length > 0) {
|
||||
if (retryErrors.length > 0) {
|
||||
await this.#host.sessionManager.rewriteEntries();
|
||||
}
|
||||
return recoveredErrors;
|
||||
return retryErrors;
|
||||
}
|
||||
|
||||
async #handleEmptyAssistantStop(assistantMessage: AssistantMessage): Promise<boolean> {
|
||||
@@ -547,7 +561,7 @@ export class TurnRecovery {
|
||||
attempt: this.#retryAttempt > 0 ? this.#retryAttempt : attempts,
|
||||
finalError,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
this.#retryAttempt = 0;
|
||||
this.resolveRetry();
|
||||
// A zero-content turn carries no transcript value, while its provider usage
|
||||
@@ -1581,27 +1595,24 @@ export class TurnRecovery {
|
||||
|
||||
const errorMessage = message.errorMessage || "Unknown error";
|
||||
const id = this.#classifyRetryMessage(message);
|
||||
const rateLimitReason = parseRateLimitReason(errorMessage);
|
||||
const staleOpenAIResponsesReplayError = AIError.is(id, AIError.Flag.StaleResponsesItem);
|
||||
const recordedUsageLimitOutcome = await this.#usageLimitOutcomes.get(message);
|
||||
const parsedRetryAfterMs = this.#parseRetryAfterMsFromError(errorMessage);
|
||||
let delayMs = staleOpenAIResponsesReplayError
|
||||
? 0
|
||||
: calculateRetryBackoffDelayMs(retrySettings.baseDelayMs, this.#retryAttempt);
|
||||
// Concurrency caps shed-and-backoff (5s) rather than burning a sibling
|
||||
// credential, so the usage-limit rotation branch below is deliberately
|
||||
// skipped for them. Apply the reason-based backoff to the transient
|
||||
// same-model retry path too — otherwise the default exponential base
|
||||
// (≈500ms) re-hits the cap immediately and burns the retry budget while
|
||||
// the concurrency slot stays occupied. A categorical 402 billing cap whose
|
||||
// body merely mentions concurrency is still a usage limit (handled below),
|
||||
// so gate on the flag matching the rotation decision.
|
||||
// Transient rate/concurrency caps stay on the same credential, but must
|
||||
// honor their reason-specific windows. The default exponential base
|
||||
// (≈500ms, capped at 8s) otherwise re-hits the cap and burns the retry
|
||||
// budget before either window can clear.
|
||||
if (
|
||||
!staleOpenAIResponsesReplayError &&
|
||||
!AIError.is(id, AIError.Flag.UsageLimit) &&
|
||||
parseRateLimitReason(errorMessage) === "CONCURRENT_LIMIT"
|
||||
(rateLimitReason === "CONCURRENT_LIMIT" || rateLimitReason === "RATE_LIMIT_EXCEEDED")
|
||||
) {
|
||||
const concurrentBackoffMs = calculateRateLimitBackoffMs("CONCURRENT_LIMIT");
|
||||
if (concurrentBackoffMs > delayMs) delayMs = concurrentBackoffMs;
|
||||
const reasonBackoffMs = calculateRateLimitBackoffMs(rateLimitReason);
|
||||
if (reasonBackoffMs > delayMs) delayMs = reasonBackoffMs;
|
||||
}
|
||||
let switchedCredential = false;
|
||||
let switchedModel = false;
|
||||
@@ -1677,16 +1688,18 @@ export class TurnRecovery {
|
||||
|
||||
if (retryBudgetExhausted) {
|
||||
if (!switchedModel) {
|
||||
const attempt = this.#retryAttempt - 1;
|
||||
message.errorMessage = `Retry budget exhausted after ${attempt} ${attempt === 1 ? "retry" : "retries"}: ${errorMessage}`;
|
||||
await this.persistTerminalEmptyErrorTurn(message);
|
||||
// Max retries exceeded and no fallback model to switch to: emit
|
||||
// final failure and reset.
|
||||
const retryErrors = await this.#markPendingRetryErrors({ status: "superseded" });
|
||||
await this.#host.emitSessionEvent({
|
||||
type: "auto_retry_end",
|
||||
success: false,
|
||||
attempt: this.#retryAttempt - 1,
|
||||
finalError: message.errorMessage,
|
||||
attempt,
|
||||
finalError: errorMessage,
|
||||
retryErrors,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
this.#retryAttempt = 0;
|
||||
this.resolveRetry(); // Resolve so waitForRetry() completes
|
||||
return false;
|
||||
@@ -1711,7 +1724,7 @@ export class TurnRecovery {
|
||||
attempt: this.#retryAttempt - 1,
|
||||
finalError: errorMessage,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
}
|
||||
this.#retryAttempt = 0;
|
||||
this.resolveRetry();
|
||||
@@ -1736,7 +1749,7 @@ export class TurnRecovery {
|
||||
attempt: this.#retryAttempt - 1,
|
||||
finalError: errorMessage,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
}
|
||||
this.#retryAttempt = 0;
|
||||
this.resolveRetry();
|
||||
@@ -1761,12 +1774,12 @@ export class TurnRecovery {
|
||||
attempt,
|
||||
finalError: `Provider requested ${delayMs}ms wait, exceeds retry.maxDelayMs (${maxDelayMs}ms). Original error: ${errorMessage}`,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
this.resolveRetry();
|
||||
return false;
|
||||
}
|
||||
|
||||
await this.#recordPendingRecoveredRetryError(message, id, { switchedCredential, switchedModel, delayMs });
|
||||
await this.#recordPendingRetryError(message, id, { switchedCredential, switchedModel, delayMs });
|
||||
|
||||
await this.#host.emitSessionEvent({
|
||||
type: "auto_retry_start",
|
||||
@@ -1808,7 +1821,7 @@ export class TurnRecovery {
|
||||
attempt,
|
||||
finalError: "Retry cancelled",
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
this.resolveRetry();
|
||||
return false;
|
||||
}
|
||||
@@ -1881,7 +1894,7 @@ export class TurnRecovery {
|
||||
attempt,
|
||||
finalError: `Retry continuation failed locally: ${localError}. Original error: ${message.errorMessage ?? "Unknown error"}`,
|
||||
});
|
||||
this.#clearPendingRecoveredRetryErrors();
|
||||
this.#clearPendingRetryErrors();
|
||||
this.resolveRetry();
|
||||
}
|
||||
|
||||
|
||||
@@ -151,6 +151,61 @@ describe("AgentSession retry delay cap", () => {
|
||||
expect(session.isRetrying).toBe(false);
|
||||
});
|
||||
|
||||
it("honors the reason backoff for a transient rate-limit 429 without a provider hint", async () => {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!model) {
|
||||
throw new Error("Expected bundled Anthropic test model to exist");
|
||||
}
|
||||
|
||||
const mock = createMockModel({
|
||||
responses: [
|
||||
{ throw: "429 Rate limit exceeded, too many requests" },
|
||||
{ content: ["recovered after rate-limit window"], stopReason: "stop" },
|
||||
],
|
||||
});
|
||||
const agent = new Agent({
|
||||
getApiKey: requestedModel => `${requestedModel.provider}-test-key`,
|
||||
initialState: {
|
||||
model,
|
||||
systemPrompt: ["Test"],
|
||||
tools: [],
|
||||
messages: [],
|
||||
},
|
||||
streamFn: (requestedModel, context, options) => mock.stream(requestedModel, context, options),
|
||||
});
|
||||
const settings = Settings.isolated({
|
||||
"compaction.enabled": false,
|
||||
"retry.baseDelayMs": 5,
|
||||
"retry.maxDelayMs": 60_000,
|
||||
"retry.maxRetries": 1,
|
||||
"retry.modelFallback": false,
|
||||
});
|
||||
settings.setModelRole("default", `${model.provider}/${model.id}`);
|
||||
|
||||
session = new AgentSession({
|
||||
agent,
|
||||
sessionManager: SessionManager.inMemory(),
|
||||
settings,
|
||||
modelRegistry,
|
||||
});
|
||||
const waitSpy = vi.spyOn(scheduler, "wait").mockResolvedValue(undefined);
|
||||
const retryStartEvents: AutoRetryStartEvent[] = [];
|
||||
session.subscribe(event => {
|
||||
if (event.type === "auto_retry_start") retryStartEvents.push(event);
|
||||
});
|
||||
|
||||
await session.prompt("Trigger transient rate limit without retry-after");
|
||||
await session.waitForIdle();
|
||||
|
||||
expect(retryStartEvents).toHaveLength(1);
|
||||
expect(retryStartEvents[0].delayMs).toBe(30_000);
|
||||
expect(waitSpy.mock.calls.some(call => call[0] === 30_000)).toBe(true);
|
||||
expect(lastAssistant(session).content).toContainEqual({
|
||||
type: "text",
|
||||
text: "recovered after rate-limit window",
|
||||
});
|
||||
});
|
||||
|
||||
it("auto-retries OpenAI Responses stream_read_error instead of stopping the conversation", async () => {
|
||||
const model = getBundledModel("openai", "gpt-5");
|
||||
if (!model) {
|
||||
@@ -1646,7 +1701,9 @@ describe("AgentSession retry delay cap", () => {
|
||||
expect(retryStartEvents[0]).toMatchObject({ attempt: 1, maxAttempts: 1 });
|
||||
expect(retryEndEvents).toHaveLength(1);
|
||||
expect(retryEndEvents[0]).toMatchObject({ success: false, attempt: 1 });
|
||||
expect(lastAssistant(session).errorMessage).toBe("server_error: stream closed with reason: error");
|
||||
expect(lastAssistant(session).errorMessage).toBe(
|
||||
"Retry budget exhausted after 1 retry: server_error: stream closed with reason: error",
|
||||
);
|
||||
expect(session.isRetrying).toBe(false);
|
||||
});
|
||||
|
||||
|
||||
@@ -216,16 +216,16 @@ describe("AgentSession retry recovery", () => {
|
||||
expect(new Set(requestedKeys)).toEqual(new Set(["anthropic-key-1", "anthropic-key-2"]));
|
||||
expect(retryEndEvents).toHaveLength(1);
|
||||
expect(retryEndEvents[0]).toMatchObject({ success: true, attempt: 1 });
|
||||
expect(retryEndEvents[0].recoveredErrors).toHaveLength(1);
|
||||
expect(retryEndEvents[0].retryErrors).toHaveLength(1);
|
||||
|
||||
const recoveredEntry = recoveredAssistantEntry(sessionManager);
|
||||
const successfulEntry = successfulAssistantEntry(sessionManager, "recovered after credential switch");
|
||||
const recoveredEvent = retryEndEvents[0].recoveredErrors?.[0];
|
||||
const recoveredEvent = retryEndEvents[0].retryErrors?.[0];
|
||||
if (!recoveredEvent) {
|
||||
throw new Error("Expected a recovered error payload on auto_retry_end");
|
||||
}
|
||||
const recoveredMarker = recoveredEntry.message.retryRecovery;
|
||||
if (!recoveredMarker) {
|
||||
if (recoveredMarker?.status !== "recovered") {
|
||||
throw new Error("Expected recovered marker on superseded assistant message");
|
||||
}
|
||||
expect(recoveredEvent.entryId).toBe(recoveredEntry.entry.id);
|
||||
@@ -245,7 +245,7 @@ describe("AgentSession retry recovery", () => {
|
||||
timestamp: successfulEntry.message.timestamp,
|
||||
},
|
||||
});
|
||||
expect(Date.parse(recoveredEntry.message.retryRecovery?.recoveredAt ?? "")).not.toBeNaN();
|
||||
expect(Date.parse(recoveredMarker.recoveredAt)).not.toBeNaN();
|
||||
|
||||
const modelContext = sessionManager.buildSessionContext();
|
||||
expect(modelContext.messages.map(message => message.role)).toEqual(["user", "assistant"]);
|
||||
@@ -266,7 +266,7 @@ describe("AgentSession retry recovery", () => {
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("leaves exhausted retries as terminal errors without recovery presentation", async () => {
|
||||
it("collapses exhausted retries into one terminal error naming the spent budget", async () => {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
|
||||
if (!model) {
|
||||
throw new Error("Expected bundled Anthropic test model to exist");
|
||||
@@ -311,27 +311,40 @@ describe("AgentSession retry recovery", () => {
|
||||
|
||||
await session.prompt("Exhaust retry attempts");
|
||||
await session.waitForIdle();
|
||||
await sessionManager.flush();
|
||||
|
||||
expect(mock.calls).toHaveLength(2);
|
||||
expect(retryEndEvents).toHaveLength(1);
|
||||
expect(retryEndEvents[0]).toMatchObject({ success: false, attempt: 1 });
|
||||
expect(retryEndEvents[0].recoveredErrors).toBeUndefined();
|
||||
expect(retryEndEvents[0].retryErrors).toHaveLength(1);
|
||||
expect(retryEndEvents[0].retryErrors?.[0].retryRecovery).toMatchObject({
|
||||
status: "superseded",
|
||||
attempt: 1,
|
||||
});
|
||||
|
||||
const terminalError = assistantEntries(sessionManager).at(-1)?.message;
|
||||
if (!terminalError) {
|
||||
throw new Error("Expected a terminal assistant error entry");
|
||||
}
|
||||
const errors = assistantEntries(sessionManager).filter(candidate => candidate.message.stopReason === "error");
|
||||
expect(errors).toHaveLength(2);
|
||||
expect(errors[0].message.retryRecovery).toMatchObject({ status: "superseded", attempt: 1 });
|
||||
expect(resolveAssistantErrorPresentation(errors[0].message)).toEqual({ kind: "none" });
|
||||
|
||||
const terminalError = errors[1].message;
|
||||
const terminalErrorText = terminalError.errorMessage;
|
||||
if (!terminalErrorText) {
|
||||
throw new Error("Expected a terminal assistant errorMessage");
|
||||
throw new Error("Expected an aggregated terminal error message");
|
||||
}
|
||||
expect(terminalError).toMatchObject({ role: "assistant", stopReason: "error" });
|
||||
expect(terminalError.retryRecovery).toBeUndefined();
|
||||
expect(terminalErrorText).toBe(`Retry budget exhausted after 1 retry: ${RETRIABLE_SERVER_ERROR}`);
|
||||
expect(resolveAssistantErrorPresentation(terminalError)).toEqual({
|
||||
kind: "full",
|
||||
text: terminalErrorText,
|
||||
isError: true,
|
||||
});
|
||||
|
||||
const visibleErrors = errors
|
||||
.map(candidate => resolveAssistantErrorPresentation(candidate.message))
|
||||
.filter(presentation => presentation.kind !== "none");
|
||||
expect(visibleErrors).toHaveLength(1);
|
||||
expect(sessionManager.buildSessionContext().messages.map(message => message.role)).toEqual(["user"]);
|
||||
});
|
||||
|
||||
it("maps assistant error presentation for recovered, unrecovered, and silent abort turns", () => {
|
||||
|
||||
@@ -183,7 +183,7 @@ describe("AgentSession thinking-loop retry", () => {
|
||||
expect(AIError.is(retryStartEvents[0].errorId, AIError.Flag.ThinkingLoop)).toBe(true);
|
||||
expect(retryEndEvents).toHaveLength(1);
|
||||
expect(retryEndEvents[0]).toMatchObject({ type: "auto_retry_end", success: true, attempt: 1 });
|
||||
expect(retryEndEvents[0].recoveredErrors).toBeArray();
|
||||
expect(retryEndEvents[0].retryErrors).toBeArray();
|
||||
const assistants = session.agent.state.messages.filter(
|
||||
(message): message is AssistantMessage => message.role === "assistant",
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user