fix(engine): isolate provider rate-limit pauses (#2339)

## What changed

- Construct one `UsageLimitPauser` per project runtime and wire it into
both executor and triage.
- Replace the project-wide emergency stop for 429/quota failures with
provider-scoped task parking.
- Resolve execution, planning, validator, and merger providers for
active tasks; park only tasks routed through the unavailable provider.
- Preserve the actual reviewer provider on `ReviewerProviderError`, so a
Claude Plan Review 429 does not stop Codex work.
- Record `provider-rate-limit:<provider>` pause provenance without
storing provider response bodies in pause metadata.
- Run one daemon-owned provider-health monitor that probes only
providers with persisted rate-limit parks.
- Resume exact matching provider parks across every project only after
the existing authenticated usage probe succeeds and all reported
capacity windows are usable.
- Probe at five-minute intervals for the first five checks, then back
off independently per provider to 10/20/40/60 minutes with a one-hour
cap.

## Root cause and impact

The runtime refactor left `usageLimitPauser` undefined for
`TriageProcessor`. In the observed FN-922 incident, Claude Plan Review
returned four explicit 429 responses; Fusion backed off for roughly
60/120/240 seconds and then failed the task, but never invoked its pause
coordinator. The older coordinator also used `globalPause`, which would
terminate healthy sessions on every other provider.

After this change, active tasks using the unavailable provider are
parked while work routed exclusively through healthy providers
continues. Recovery is a provider-health state transition: the daemon
checks Claude/Codex authentication and metered capacity independently of
task execution, including after restart, and clears only exact
`provider-rate-limit:<provider>` parks. Logged-out, errored, exhausted,
manually paused, user-paused, and other-provider tasks remain parked.
Explicit global/engine pause controls remain unchanged.

## Surface enumeration

- executor usage-limit catches
- triage planner and Plan Review catches
- reviewer provider-error propagation
- merger usage-limit catches
- per-project runtime construction and wiring
- task model overrides plus project/global execution, planning,
validator, and merger resolution
- daemon startup/listen and shutdown lifecycle
- multi-project provider-probe deduplication
- Claude and Codex authenticated usage/capacity probes
- done/archived/already-paused task exclusions
- manual, user, generic, and other-provider pause provenance

## Symptom verification

**Original symptom:** Anthropic/Claude 429s retried and failed FN-922
without pausing Claude-routed work; a functioning global pauser would
also have stopped Codex, and provider parks had no positive-health
recovery path.

**Exact reproduction:** Raise `ReviewerProviderError("429
overloaded_error", "usage-limit", { provider: "anthropic" })` during
Plan Review with Anthropic and Codex tasks present, then return
logged-out/error/exhausted and finally healthy Claude usage responses
from the daemon probe.

**Assertion it is gone:** Anthropic-routed active tasks receive
`provider-rate-limit:anthropic`; Codex-only tasks are not paused and
`globalPause` is never changed. Unhealthy probes leave the Anthropic
tasks parked; a positive authenticated response with remaining capacity
resumes only exact Anthropic provider parks without executing a model
call as a probe.

## Validation

- `packages/engine/src/__tests__/usage-limit-detector.test.ts`: 49
passed
- `packages/dashboard/src/__tests__/provider-health-monitor.test.ts`: 8
passed
- Engine TypeScript check passed
- Dashboard server and app TypeScript checks passed
- Scoped ESLint passed
- Changeset strict format check passed
- Reapply script passed `bash -n`, two consecutive fixture applications,
and `node --check`


<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->

## Summary by CodeRabbit

- **New Features**
- Tasks paused due to a provider’s rate limits can now automatically
resume when capacity returns.
- Provider health is monitored in the background, including retry
backoff for unavailable providers.

- **Bug Fixes**
- Rate-limit issues now pause only affected provider-routed tasks
instead of stopping unrelated work.
- Provider failures are handled separately from invalid review results,
improving recovery behavior.
- Healthy providers remain available while another provider is
rate-limited.

<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Co-authored-by: v <v@v.speedport.ip>
This commit is contained in:
flexi767
2026-07-22 02:08:44 +02:00
committed by GitHub
parent 8f7f52784d
commit c71a9545b0
20 changed files with 823 additions and 107 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Keep healthy AI providers running and resume provider-paused tasks when capacity returns.
category: fix
dev: Provider-scoped parks recover from daemon-side authenticated usage and capacity health transitions without task-call probes.

View File

@@ -0,0 +1,224 @@
import { describe, expect, it, vi } from "vitest";
import type { ProviderUsage } from "../usage.js";
import {
hasUsableProviderCapacity,
ProviderHealthMonitor,
hasIndependentProviderHealthProbe,
providerHealthProbeDelayMs,
providerIdFromRateLimitReason,
} from "../provider-health-monitor.js";
function providerUsage(overrides: Partial<ProviderUsage> = {}): ProviderUsage {
return {
name: "Claude",
icon: "test",
status: "ok",
windows: [{
label: "Weekly",
percentUsed: 40,
percentLeft: 60,
resetText: "resets in 4d",
}],
...overrides,
};
}
function createStore(tasks: Array<Record<string, unknown>>) {
return {
listTasks: vi.fn().mockResolvedValue(tasks),
logEntry: vi.fn().mockResolvedValue(undefined),
pauseTask: vi.fn().mockImplementation(async (id: string, paused: boolean) => {
const task = tasks.find((candidate) => candidate.id === id);
if (task) {
task.paused = paused || undefined;
if (!paused) task.pausedReason = undefined;
}
}),
} as any;
}
function createLogger() {
return {
scope: "test",
info: vi.fn(),
warn: vi.fn(),
error: vi.fn(),
child: vi.fn(),
} as any;
}
describe("ProviderHealthMonitor", () => {
it("derives provider-qualified and legacy unqualified rate-limit parks", () => {
expect(providerIdFromRateLimitReason("provider-rate-limit:Anthropic")).toBe("anthropic");
expect(providerIdFromRateLimitReason("provider-rate-limit")).toBe("unknown");
expect(providerIdFromRateLimitReason("manual")).toBeNull();
});
it("identifies providers with independent capacity meters", () => {
expect(hasIndependentProviderHealthProbe("Anthropic")).toBe(true);
expect(hasIndependentProviderHealthProbe("openai-codex")).toBe(true);
expect(hasIndependentProviderHealthProbe("openrouter")).toBe(false);
});
it("requires positive auth and non-exhausted metered capacity", () => {
expect(hasUsableProviderCapacity(providerUsage())).toBe(true);
expect(hasUsableProviderCapacity(providerUsage({ status: "no-auth" }))).toBe(false);
expect(hasUsableProviderCapacity(providerUsage({ status: "error", error: "HTTP 429" }))).toBe(false);
expect(hasUsableProviderCapacity(providerUsage({ windows: [] }))).toBe(false);
expect(hasUsableProviderCapacity(providerUsage({
windows: [{ label: "Weekly", percentUsed: 100, percentLeft: 0, resetText: "resets in 1h" }],
}))).toBe(false);
});
it("checks at five-minute cadence five times, then backs off to a one-hour cap", () => {
expect(providerHealthProbeDelayMs(1)).toBe(300_000);
expect(providerHealthProbeDelayMs(4)).toBe(300_000);
expect(providerHealthProbeDelayMs(5)).toBe(600_000);
expect(providerHealthProbeDelayMs(6)).toBe(1_200_000);
expect(providerHealthProbeDelayMs(7)).toBe(2_400_000);
expect(providerHealthProbeDelayMs(8)).toBe(3_600_000);
expect(providerHealthProbeDelayMs(20)).toBe(3_600_000);
});
it("does not re-probe before a provider's growing backoff expires", async () => {
let now = 0;
const store = createStore([
{ id: "FN-9", paused: true, pausedReason: "provider-rate-limit:anthropic" },
]);
const probe = vi.fn().mockResolvedValue(providerUsage({ status: "error", error: "HTTP 429" }));
const monitor = new ProviderHealthMonitor({
getStores: () => [store],
logger: createLogger(),
probe,
now: () => now,
});
await monitor.checkNow();
now = 299_999;
await monitor.checkNow();
expect(probe).toHaveBeenCalledTimes(1);
now = 300_000;
await monitor.checkNow();
expect(probe).toHaveBeenCalledTimes(2);
});
it("requeues unsupported and legacy-unqualified provider parks after a bounded cooldown", async () => {
let now = 0;
const store = createStore([
{ id: "FN-7", paused: true, pausedReason: "provider-rate-limit:openrouter" },
{ id: "FN-8", paused: true, pausedReason: "provider-rate-limit:unknown" },
{ id: "FN-9", paused: true, pausedReason: "provider-rate-limit" },
]);
const failedStore = createStore([
{ id: "FN-FAILED", paused: true, pausedReason: "provider-rate-limit:openrouter" },
]);
failedStore.listTasks
.mockResolvedValueOnce([{ id: "FN-FAILED", paused: true, pausedReason: "provider-rate-limit:openrouter" }])
.mockRejectedValue(new Error("database unavailable"));
const probe = vi.fn();
const logger = createLogger();
const monitor = new ProviderHealthMonitor({
getStores: () => [failedStore, store],
logger,
probe,
supportsProbe: () => false,
now: () => now,
pollIntervalMs: 1_000,
});
await monitor.checkNow();
expect(store.pauseTask).not.toHaveBeenCalled();
now = 1_000;
await monitor.checkNow();
expect(probe).not.toHaveBeenCalled();
expect(store.pauseTask).toHaveBeenCalledWith("FN-7", false);
expect(store.pauseTask).toHaveBeenCalledWith("FN-8", false);
expect(store.pauseTask).toHaveBeenCalledWith("FN-9", false);
expect(logger.warn).toHaveBeenCalledWith(
"Failed to recover tasks after provider cooldown",
{ providerId: "openrouter", error: "database unavailable" },
);
});
it("probes a persisted unavailable provider once and resumes exact matching parks across projects", async () => {
const claudeStoreA = createStore([
{ id: "FN-1", paused: true, pausedReason: "provider-rate-limit:anthropic" },
{ id: "FN-2", paused: true, pausedReason: "manual" },
]);
const claudeStoreB = createStore([
{ id: "FN-3", paused: true, pausedReason: "provider-rate-limit:anthropic" },
{ id: "FN-4", paused: true, userPaused: true, pausedReason: "provider-rate-limit:anthropic" },
{ id: "FN-5", paused: true, pausedReason: "provider-rate-limit:openai-codex" },
]);
const probe = vi.fn().mockImplementation(async (providerId: string) =>
providerId === "anthropic"
? providerUsage()
: providerUsage({ name: "Codex", status: "error", error: "HTTP 429" }));
const logger = createLogger();
const monitor = new ProviderHealthMonitor({
getStores: () => [claudeStoreA, claudeStoreB],
logger,
probe,
});
await monitor.checkNow();
expect(probe).toHaveBeenCalledTimes(2);
expect(probe).toHaveBeenCalledWith("anthropic", undefined);
expect(probe).toHaveBeenCalledWith("openai-codex", undefined);
expect(claudeStoreA.pauseTask).toHaveBeenCalledWith("FN-1", false);
expect(claudeStoreA.pauseTask).not.toHaveBeenCalledWith("FN-2", false);
expect(claudeStoreB.pauseTask).toHaveBeenCalledWith("FN-3", false);
expect(claudeStoreB.pauseTask).not.toHaveBeenCalledWith("FN-4", false);
expect(claudeStoreB.pauseTask).not.toHaveBeenCalledWith("FN-5", false);
expect(logger.info).toHaveBeenCalledWith(
"Provider available again; resumed provider-paused tasks",
{ providerId: "anthropic", recoveredTasks: 2 },
);
});
it("continues scanning healthy project stores when one store cannot list tasks", async () => {
const failedStore = createStore([]);
failedStore.listTasks.mockRejectedValue(new Error("database unavailable"));
const healthyStore = createStore([
{ id: "FN-6", paused: true, pausedReason: "provider-rate-limit:anthropic" },
]);
const logger = createLogger();
const probe = vi.fn().mockResolvedValue(providerUsage());
const monitor = new ProviderHealthMonitor({
getStores: () => [failedStore, healthyStore],
logger,
probe,
});
await monitor.checkNow();
expect(logger.warn).toHaveBeenCalledWith(
"Failed to list tasks for provider health scan",
{ error: "database unavailable" },
);
expect(probe).toHaveBeenCalledWith("anthropic", undefined);
expect(healthyStore.pauseTask).toHaveBeenCalledWith("FN-6", false);
});
it.each([
providerUsage({ status: "no-auth", error: "login required" }),
providerUsage({ status: "error", error: "HTTP 429" }),
providerUsage({ windows: [{ label: "Weekly", percentUsed: 100, percentLeft: 0, resetText: null }] }),
])("keeps tasks parked without using task execution as a recovery probe", async (usage) => {
const store = createStore([
{ id: "FN-10", paused: true, pausedReason: "provider-rate-limit:anthropic" },
]);
const monitor = new ProviderHealthMonitor({
getStores: () => [store],
logger: createLogger(),
probe: vi.fn().mockResolvedValue(usage),
});
await monitor.checkNow();
expect(store.pauseTask).not.toHaveBeenCalled();
});
});

View File

@@ -0,0 +1,280 @@
import type { TaskStore } from "@fusion/core";
import { UsageLimitPauser } from "@fusion/engine";
import type { RuntimeLogger } from "./runtime-logger.js";
import {
fetchClaudeUsage,
fetchCodexUsage,
type AuthStorageLike,
type ProviderUsage,
} from "./usage.js";
const PROVIDER_RATE_LIMIT_PREFIX = "provider-rate-limit:";
const DEFAULT_POLL_INTERVAL_MS = 300_000;
const DEFAULT_MAX_POLL_INTERVAL_MS = 3_600_000;
const CHECKS_BEFORE_BACKOFF = 5;
const INDEPENDENTLY_METERED_PROVIDERS = new Set([
"anthropic",
"anthropic-subscription",
"openai-codex",
]);
export type ProviderHealthProbe = (
providerId: string,
authStorage?: AuthStorageLike,
) => Promise<ProviderUsage | null>;
export interface ProviderHealthMonitorOptions {
getStores: () => Iterable<TaskStore>;
authStorage?: AuthStorageLike;
logger: RuntimeLogger;
pollIntervalMs?: number;
maxPollIntervalMs?: number;
probe?: ProviderHealthProbe;
supportsProbe?: (providerId: string) => boolean;
now?: () => number;
}
interface ProviderMonitorState {
status: "unavailable" | "available";
failedChecks: number;
nextCheckAt: number;
}
function normalizeProviderId(providerId: string): string {
return providerId.trim().toLowerCase();
}
export function providerIdFromRateLimitReason(reason: string | undefined): string | null {
if (reason === "provider-rate-limit") return "unknown";
if (!reason?.startsWith(PROVIDER_RATE_LIMIT_PREFIX)) return null;
const providerId = normalizeProviderId(reason.slice(PROVIDER_RATE_LIMIT_PREFIX.length));
return providerId || null;
}
export function hasIndependentProviderHealthProbe(providerId: string): boolean {
return INDEPENDENTLY_METERED_PROVIDERS.has(normalizeProviderId(providerId));
}
/**
* A successful auth check is not enough: a provider can serve its usage API
* while an enforced session, weekly, or model-specific window is exhausted.
*/
export function hasUsableProviderCapacity(usage: ProviderUsage | null): boolean {
return usage?.status === "ok"
&& usage.windows.length > 0
&& usage.windows.every((window) => window.percentUsed < 100 && window.percentLeft > 0);
}
export function providerHealthProbeDelayMs(
failedChecks: number,
baseIntervalMs = DEFAULT_POLL_INTERVAL_MS,
maxIntervalMs = DEFAULT_MAX_POLL_INTERVAL_MS,
): number {
if (failedChecks < CHECKS_BEFORE_BACKOFF) return baseIntervalMs;
const exponent = failedChecks - CHECKS_BEFORE_BACKOFF + 1;
return Math.min(baseIntervalMs * (2 ** exponent), maxIntervalMs);
}
export async function probeMeteredProviderHealth(
providerId: string,
authStorage?: AuthStorageLike,
): Promise<ProviderUsage | null> {
switch (normalizeProviderId(providerId)) {
case "anthropic":
case "anthropic-subscription":
return fetchClaudeUsage(authStorage);
case "openai-codex":
return fetchCodexUsage();
default:
return null;
}
}
/**
* Daemon-owned provider recovery monitor.
*
* It probes only providers that currently own persisted rate-limit parks. A
* single provider check is shared across every project store, and task model
* execution is never used as the health probe.
*/
export class ProviderHealthMonitor {
private readonly pollIntervalMs: number;
private readonly maxPollIntervalMs: number;
private readonly probe: ProviderHealthProbe;
private readonly supportsProbe: (providerId: string) => boolean;
private readonly now: () => number;
private readonly states = new Map<string, ProviderMonitorState>();
private timer: ReturnType<typeof setTimeout> | null = null;
private running = false;
constructor(private readonly options: ProviderHealthMonitorOptions) {
this.pollIntervalMs = options.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS;
this.maxPollIntervalMs = options.maxPollIntervalMs ?? DEFAULT_MAX_POLL_INTERVAL_MS;
this.probe = options.probe ?? probeMeteredProviderHealth;
this.supportsProbe = options.supportsProbe
?? (options.probe ? () => true : hasIndependentProviderHealthProbe);
this.now = options.now ?? Date.now;
}
start(): void {
if (this.running) return;
this.running = true;
void this.runAndSchedule();
}
stop(): void {
this.running = false;
if (this.timer) clearTimeout(this.timer);
this.timer = null;
}
async checkNow(): Promise<void> {
const stores = Array.from(new Set(this.options.getStores()));
const providers = new Set<string>();
await Promise.all(stores.map(async (store) => {
try {
const tasks = await store.listTasks();
for (const task of tasks) {
if (task.paused !== true || task.userPaused === true) continue;
const providerId = providerIdFromRateLimitReason(task.pausedReason);
if (providerId) providers.add(providerId);
}
} catch (error: unknown) {
this.options.logger.warn("Failed to list tasks for provider health scan", {
error: error instanceof Error ? error.message : String(error),
});
}
}));
for (const knownProvider of Array.from(this.states.keys())) {
if (!providers.has(knownProvider)) this.states.delete(knownProvider);
}
await Promise.all(Array.from(providers, async (providerId) => {
const state = this.states.get(providerId);
const checkedAt = this.now();
if (state && state.nextCheckAt > checkedAt) return;
/*
FNXC:ProviderRateLimitRecovery 2026-07-21-21:30:
A provider without an independent quota meter must never be parked forever. Give it one normal poll interval as a cooldown, then requeue its exact provider-qualified parks so ordinary execution can confirm that capacity returned. Legacy unqualified parks use the same bounded fallback under the synthetic "unknown" provider id.
FNXC:ProviderRateLimitRecovery 2026-07-21-21:50:
Cooldown recovery is isolated per project store so one unavailable database cannot prevent healthy projects or other providers from resuming.
*/
if (!this.supportsProbe(providerId)) {
if (!state || state.status !== "unavailable") {
this.states.set(providerId, {
status: "unavailable",
failedChecks: 1,
nextCheckAt: checkedAt + this.pollIntervalMs,
});
this.options.logger.warn("Provider has no independent health probe; applying bounded cooldown", {
providerId,
});
return;
}
const recoveredCounts = await Promise.all(stores.map(async (store) => {
try {
return await new UsageLimitPauser(store).onProviderAvailable(providerId);
} catch (error: unknown) {
this.options.logger.warn("Failed to recover tasks after provider cooldown", {
providerId,
error: error instanceof Error ? error.message : String(error),
});
return 0;
}
}));
const recoveredTasks = recoveredCounts.reduce((total, count) => total + count, 0);
this.states.set(providerId, {
status: "available",
failedChecks: 0,
nextCheckAt: checkedAt + this.pollIntervalMs,
});
if (recoveredTasks > 0) {
this.options.logger.info("Provider cooldown elapsed; resumed provider-paused tasks", {
providerId,
recoveredTasks,
});
}
return;
}
const usage = await this.probe(providerId, this.options.authStorage).catch((error: unknown) => {
this.options.logger.warn("Provider health probe failed", {
providerId,
error: error instanceof Error ? error.message : String(error),
});
return null;
});
if (!hasUsableProviderCapacity(usage)) {
const failedChecks = (state?.failedChecks ?? 0) + 1;
const retryDelayMs = providerHealthProbeDelayMs(
failedChecks,
this.pollIntervalMs,
this.maxPollIntervalMs,
);
if (state?.status !== "unavailable") {
this.options.logger.warn("Provider unavailable; rate-limited tasks remain paused", {
providerId,
status: usage?.status ?? "unsupported",
error: usage?.error,
});
}
this.states.set(providerId, {
status: "unavailable",
failedChecks,
nextCheckAt: checkedAt + retryDelayMs,
});
return;
}
/*
FNXC:ProviderRateLimitRecovery 2026-07-19-20:15:
Recovery is driven by an authenticated daemon-side usage/capacity transition, not by waking a parked task and spending a model call as a probe. The persisted pause reason seeds this monitor after restart; one provider probe fans out only to exact matching parks across project engines.
Provider probes start at a five-minute cadence; after five failed checks each provider backs off independently to 10/20/40/60 minutes. This bounds subscription API traffic without permanently stranding parks during a long outage.
*/
const recoveredCounts = await Promise.all(stores.map(async (store) => {
try {
return await new UsageLimitPauser(store).onProviderAvailable(providerId);
} catch (error: unknown) {
this.options.logger.warn("Failed to recover tasks for available provider", {
providerId,
error: error instanceof Error ? error.message : String(error),
});
return 0;
}
}));
const recoveredTasks = recoveredCounts.reduce((total, count) => total + count, 0);
this.states.set(providerId, {
status: "available",
failedChecks: 0,
nextCheckAt: checkedAt + this.pollIntervalMs,
});
if (recoveredTasks > 0) {
this.options.logger.info("Provider available again; resumed provider-paused tasks", {
providerId,
recoveredTasks,
});
}
}));
}
private async runAndSchedule(): Promise<void> {
try {
await this.checkNow();
} catch (error: unknown) {
this.options.logger.warn("Provider health reconciliation failed", {
error: error instanceof Error ? error.message : String(error),
});
} finally {
if (this.running) {
this.timer = setTimeout(() => void this.runAndSchedule(), this.pollIntervalMs);
this.timer.unref?.();
}
}
}
}

View File

@@ -89,6 +89,7 @@ import {
resolveDashboardPostgresLayer,
type DashboardTaskIdIntegrityHealth,
} from "./dashboard-postgres-health.js";
import { ProviderHealthMonitor } from "./provider-health-monitor.js";
const __dirname = dirname(fileURLToPath(import.meta.url));
@@ -2206,6 +2207,7 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
// FUSION_OTEL_METRICS_ENDPOINT is explicitly configured. Held here so the
// server "close" handler can stop its timer.
let otelExporter: OtelExporterHandle | null = null;
let providerHealthMonitor: ProviderHealthMonitor | null = null;
dashboardApp.listen = ((...args: Parameters<typeof dashboardApp.listen>) => {
const normalizedArgs = normalizeListenArgsForTests(args) as Parameters<typeof originalListen>;
@@ -2241,11 +2243,31 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
});
}
if (!providerHealthMonitor && (options?.engineManager || options?.engine)) {
const providerHealthLogger = runtimeLogger.child("provider-health");
providerHealthMonitor = new ProviderHealthMonitor({
authStorage: options?.authStorage,
logger: providerHealthLogger,
getStores: () => {
const stores = new Set<TaskStore>([store]);
const singleEngineStore = options?.engine?.getTaskStore?.();
if (singleEngineStore) stores.add(singleEngineStore);
for (const projectEngine of options?.engineManager?.getAllEngines?.().values() ?? []) {
stores.add(projectEngine.getTaskStore());
}
return stores;
},
});
providerHealthMonitor.start();
}
server.once("close", () => {
clearAiSessionCleanupInterval();
aiSessionStore?.stopScheduledCleanup();
otelExporter?.stop();
otelExporter = null;
providerHealthMonitor?.stop();
providerHealthMonitor = null;
(apiRouter as Router & { dispose?: () => void }).dispose?.();
void stopAllDevServers().catch((error) => {
runtimeLogger.warn("Failed to shutdown dev-server managers", {

View File

@@ -898,7 +898,7 @@ async function fetchClaudeUsageViaCli(): Promise<ProviderUsage> {
* Includes retry logic with exponential backoff for transient 429 responses.
* Falls back to parsing `claude /usage` CLI output when rate limited.
*/
async function fetchClaudeUsage(authStorage?: AuthStorageLike): Promise<ProviderUsage> {
export async function fetchClaudeUsage(authStorage?: AuthStorageLike): Promise<ProviderUsage> {
const usage: ProviderUsage = {
name: "Claude",
icon: "🟠",
@@ -1251,7 +1251,7 @@ async function loadCodexCredential(): Promise<CodexCredential | null> {
return null;
}
async function fetchCodexUsage(): Promise<ProviderUsage> {
export async function fetchCodexUsage(): Promise<ProviderUsage> {
const usage: ProviderUsage = {
name: "Codex",
icon: "🟢",

View File

@@ -75,6 +75,13 @@ pgDescribe("InProcessRuntime PostgreSQL composition", () => {
expect(taskStore.isBackendMode()).toBe(true);
expect(layer?.projectId).toBe("runtime-composition");
expect(runtime.getMissionExecutionLoop()).toBeDefined();
const runtimeInternals = runtime as unknown as {
usageLimitPauser?: unknown;
triageProcessor?: { options?: { usageLimitPauser?: unknown } };
};
expect(runtimeInternals.usageLimitPauser).toBeDefined();
expect(runtimeInternals.triageProcessor?.options?.usageLimitPauser)
.toBe(runtimeInternals.usageLimitPauser);
const missionStore = taskStore.getMissionStore();
const mission = await missionStore.createMission({ title: "Runtime composition" });

View File

@@ -185,6 +185,7 @@ function createMockStore(taskOverrides: Partial<Task> = {}, allTasks: Task[] = [
getTask: vi.fn().mockResolvedValue({ ...baseTask, prompt: "# test" }),
listTasks: vi.fn().mockResolvedValue(allTasks),
updateTask: vi.fn().mockResolvedValue(baseTask),
pauseTask: vi.fn().mockResolvedValue({ ...baseTask, paused: true }),
moveTask: vi.fn().mockResolvedValue({ ...baseTask, column: "done" }),
logEntry: vi.fn().mockResolvedValue(undefined),
appendAgentLog: vi.fn().mockResolvedValue(undefined),
@@ -572,7 +573,7 @@ describe("aiMergeTask — usage limit detection", () => {
setupFailingTheirsStrategy();
});
it("triggers global pause when merger catches a usage-limit error", async () => {
it("parks only the merge task when merger catches a usage-limit error", async () => {
const store = createMockStore(
{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050" },
[{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
@@ -595,14 +596,15 @@ describe("aiMergeTask — usage limit detection", () => {
"merger",
"FN-050",
"rate_limit_error: Rate limit exceeded",
undefined,
);
expect(store.updateSettings).toHaveBeenCalledWith({
globalPause: true,
globalPauseReason: "rate-limit",
expect(store.pauseTask).toHaveBeenCalledWith("FN-050", true, undefined, {
pausedReason: "provider-rate-limit",
});
expect(store.updateSettings).not.toHaveBeenCalled();
});
it("triggers global pause when session.prompt() resolves with exhausted-retry error on state.error", async () => {
it("parks the merge task when session.prompt() resolves with exhausted-retry error on state.error", async () => {
const store = createMockStore(
{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050" },
[{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
@@ -627,6 +629,7 @@ describe("aiMergeTask — usage limit detection", () => {
"merger",
"FN-050",
"429 Too Many Requests",
undefined,
);
// git reset --merge should be called to abort the merge
const resetCalls = mockedExecSync.mock.calls.filter(
@@ -635,7 +638,7 @@ describe("aiMergeTask — usage limit detection", () => {
expect(resetCalls.length).toBeGreaterThan(0);
});
it("does NOT trigger global pause for non-usage-limit errors", async () => {
it("does NOT park the task for non-usage-limit errors", async () => {
const store = createMockStore(
{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050" },
[{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
@@ -676,7 +679,7 @@ describe("aiMergeTask — usage limit detection", () => {
).rejects.toThrow("AI merge failed");
});
it("triggers global pause for overloaded error", async () => {
it("parks the task for an overloaded error", async () => {
const store = createMockStore(
{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050" },
[{ id: "FN-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
@@ -699,6 +702,7 @@ describe("aiMergeTask — usage limit detection", () => {
"merger",
"FN-050",
"overloaded_error: Overloaded",
undefined,
);
});
});
@@ -1170,4 +1174,3 @@ describe("aiMergeTask — merge details collection", () => {
expect(mergeDetails.deletions).toBeUndefined();
});
});

View File

@@ -1725,11 +1725,15 @@ describe("reviewStep — provider errors are not review verdicts", () => {
mockedPromptWithFallback.mockRejectedValue(new Error(RATE_LIMIT_ERROR));
const error = await reviewStep(
"/tmp/worktree", "FN-RL", 2, "Rate limited", "code", "# prompt", "abc123", {},
"/tmp/worktree", "FN-RL", 2, "Rate limited", "code", "# prompt", "abc123", {
defaultProvider: "anthropic",
defaultModelId: "claude-sonnet",
},
).then(() => null, (err: unknown) => err);
expect(error).toBeInstanceOf(ReviewerProviderError);
expect((error as ReviewerProviderError).classification).toBe("usage-limit");
expect((error as ReviewerProviderError).provider).toBe("anthropic");
// The reported bug: with no configured fallback the ladder re-ran the SAME model
// immediately, so a 429 spawned a second session (and a second identical marker).
expect(mockedCreateFnAgent).toHaveBeenCalledTimes(1);

View File

@@ -128,11 +128,23 @@ describe("checkSessionError", () => {
// ── UsageLimitPauser tests ───────────────────────────────────────────
function createMockStore(globalPause = false) {
function createMockStore(tasks: any[] = []) {
return {
getSettings: vi.fn().mockResolvedValue({ globalPause }),
updateSettings: vi.fn().mockResolvedValue({ globalPause: true }),
logEntry: vi.fn().mockResolvedValue(undefined),
pauseTask: vi.fn().mockResolvedValue(undefined),
listTasks: vi.fn().mockResolvedValue(tasks),
getTask: vi.fn().mockImplementation(async (id: string) => tasks.find((task) => task.id === id) ?? {
id,
column: "todo",
dependencies: [],
steps: [],
currentStep: 0,
log: [],
}),
getSettings: vi.fn().mockResolvedValue({
defaultProvider: "openai-codex",
defaultModelId: "gpt-5",
}),
} as any;
}
@@ -141,16 +153,20 @@ describe("UsageLimitPauser", () => {
vi.clearAllMocks();
});
it("calls store.updateSettings({ globalPause: true, globalPauseReason: \"rate-limit\" }) on usage limit hit", async () => {
const store = createMockStore();
it("pauses only the affected task instead of activating global pause", async () => {
const store = createMockStore([
{ id: "FN-001", column: "todo", modelProvider: "anthropic", modelId: "claude-sonnet" },
{ id: "FN-002", column: "todo", modelProvider: "openai-codex", modelId: "gpt-5" },
]);
const pauser = new UsageLimitPauser(store);
await pauser.onUsageLimitHit("executor", "FN-001", "rate_limit_error: Rate limit exceeded");
await pauser.onUsageLimitHit("executor", "FN-001", "rate_limit_error: Rate limit exceeded", "anthropic");
expect(store.updateSettings).toHaveBeenCalledWith({
globalPause: true,
globalPauseReason: "rate-limit",
expect(store.pauseTask).toHaveBeenCalledWith("FN-001", true, undefined, {
pausedReason: "provider-rate-limit:anthropic",
});
expect(store.pauseTask).not.toHaveBeenCalledWith("FN-002", expect.anything(), expect.anything(), expect.anything());
expect(store.updateSettings).toBeUndefined();
});
it("logs the triggering error on the task via store.logEntry", async () => {
@@ -161,51 +177,91 @@ describe("UsageLimitPauser", () => {
expect(store.logEntry).toHaveBeenCalledWith(
"FN-002",
"Usage limit detected (triage): overloaded_error",
"Usage limit detected (triage/unknown): overloaded_error",
);
});
it("is idempotent — calling multiple times only triggers one pause", async () => {
const store = createMockStore();
// After first call, globalPause will be true
store.getSettings.mockResolvedValue({ globalPause: true });
it("parks active executor tasks on the unavailable provider while other lanes and providers continue", async () => {
const store = createMockStore([
{ id: "FN-001", column: "in-progress", modelProvider: "anthropic", modelId: "claude-sonnet" },
{ id: "FN-002", column: "in-progress", modelProvider: "anthropic", modelId: "claude-sonnet" },
{ id: "FN-003", column: "in-progress", modelProvider: "openai-codex", modelId: "gpt-5" },
{ id: "FN-004", column: "triage", planningModelProvider: "anthropic", planningModelId: "claude-sonnet" },
]);
const pauser = new UsageLimitPauser(store);
await pauser.onUsageLimitHit("executor", "FN-001", "rate limit");
await pauser.onUsageLimitHit("triage", "FN-002", "rate limit");
await pauser.onUsageLimitHit("merger", "FN-003", "rate limit");
await pauser.onUsageLimitHit("executor", "FN-001", "rate limit", "anthropic");
// updateSettings should only be called once
expect(store.updateSettings).toHaveBeenCalledTimes(1);
expect(store.pauseTask).toHaveBeenCalledTimes(2);
expect(store.pauseTask).toHaveBeenCalledWith("FN-001", true, undefined, { pausedReason: "provider-rate-limit:anthropic" });
expect(store.pauseTask).toHaveBeenCalledWith("FN-002", true, undefined, { pausedReason: "provider-rate-limit:anthropic" });
expect(store.pauseTask).not.toHaveBeenCalledWith("FN-003", expect.anything(), expect.anything(), expect.anything());
expect(store.pauseTask).not.toHaveBeenCalledWith("FN-004", expect.anything(), expect.anything(), expect.anything());
});
it("re-triggers pause if globalPause was externally reset to false", async () => {
it("parks only active triage tasks whose planning or validator lane uses the unavailable provider", async () => {
const store = createMockStore([
{ id: "FN-010", column: "triage", validatorModelProvider: "anthropic", validatorModelId: "claude-sonnet" },
{ id: "FN-011", column: "triage", planningModelProvider: "openai-codex", planningModelId: "gpt-5", validatorModelProvider: "openai-codex", validatorModelId: "gpt-5" },
{ id: "FN-012", column: "in-progress", validatorModelProvider: "anthropic", validatorModelId: "claude-sonnet" },
]);
const pauser = new UsageLimitPauser(store);
await pauser.onUsageLimitHit("triage", "FN-010", "429", "anthropic");
expect(store.pauseTask).toHaveBeenCalledTimes(1);
expect(store.pauseTask).toHaveBeenCalledWith("FN-010", true, undefined, { pausedReason: "provider-rate-limit:anthropic" });
});
it("uses a recoverable qualified reason when the caller cannot identify the provider", async () => {
const store = createMockStore();
const pauser = new UsageLimitPauser(store);
// First hit — triggers pause
store.getSettings.mockResolvedValue({ globalPause: true });
await pauser.onUsageLimitHit("executor", "FN-001", "rate limit");
expect(store.updateSettings).toHaveBeenCalledTimes(1);
// External reset: globalPause set to false
store.getSettings.mockResolvedValue({ globalPause: false });
// Second hit — should trigger again since it was reset
await pauser.onUsageLimitHit("executor", "FN-004", "rate limit again");
expect(store.updateSettings).toHaveBeenCalledTimes(2);
expect(store.pauseTask).toHaveBeenCalledWith("FN-001", true, undefined, {
pausedReason: "provider-rate-limit:unknown",
});
});
it("includes agent type in the log entry", async () => {
const store = createMockStore();
const pauser = new UsageLimitPauser(store);
await pauser.onUsageLimitHit("merger", "FN-005", "quota exceeded");
await pauser.onUsageLimitHit("merger", "FN-005", "quota exceeded", "Anthropic API");
expect(store.logEntry).toHaveBeenCalledWith(
"FN-005",
expect.stringContaining("merger"),
expect.stringContaining("merger/anthropic-api"),
);
});
it("resumes only exact provider-rate-limit parks after positive provider health", async () => {
const store = createMockStore([
{ id: "FN-101", paused: true, pausedReason: "provider-rate-limit:anthropic" },
{ id: "FN-102", paused: true, pausedReason: "provider-rate-limit:openai-codex" },
{ id: "FN-103", paused: true, pausedReason: "manual" },
{ id: "FN-104", paused: true, userPaused: true, pausedReason: "provider-rate-limit:anthropic" },
{ id: "FN-105", paused: false, pausedReason: "provider-rate-limit:anthropic" },
]);
const pauser = new UsageLimitPauser(store);
await expect(pauser.onProviderAvailable("Anthropic")).resolves.toBe(1);
expect(store.pauseTask).toHaveBeenCalledTimes(1);
expect(store.pauseTask).toHaveBeenCalledWith("FN-101", false);
expect(store.logEntry).toHaveBeenCalledWith(
"FN-101",
"Provider anthropic is available again; resuming task",
);
});
it("does nothing when provider health has no matching persisted parks", async () => {
const store = createMockStore([
{ id: "FN-201", paused: true, pausedReason: "provider-rate-limit:openai-codex" },
]);
const pauser = new UsageLimitPauser(store);
await expect(pauser.onProviderAvailable("anthropic")).resolves.toBe(0);
expect(store.pauseTask).not.toHaveBeenCalled();
});
});

View File

@@ -181,12 +181,12 @@ export class ValidationError extends PermanentError {
}
}
// ── Rate Limit Error (special — global pause, not local retry) ──────────
// ── Rate Limit Error (special — provider-scoped park, not local retry) ──
/**
* Rate-limit / usage-limit error. Unlike transient errors, these should
* NOT be retried locally — instead they trigger a global pause via
* UsageLimitPauser so all agents back off simultaneously.
* NOT be retried indefinitely — after bounded backoff they park only the
* affected provider-routed task via UsageLimitPauser.
*/
export class RateLimitError extends EngineError {
/** Suggested retry-after in milliseconds (from Retry-After header or heuristic). */
@@ -224,7 +224,7 @@ export function classifyThrownError(err: unknown): EngineError {
const message = err instanceof Error ? err.message : String(err ?? "");
// Rate-limit (triggers global pause)
// Rate-limit (triggers a provider-scoped task park)
if (isUsageLimitError(message)) {
return new RateLimitError(message, undefined, undefined, err instanceof Error ? err : undefined);
}

View File

@@ -1541,7 +1541,10 @@ export interface TaskExecutorOptions {
semaphore?: AgentSemaphore;
/** Worktree pool for recycling idle worktrees across tasks. */
pool?: WorktreePool;
/** Usage limit pauser — triggers global pause when API limits are detected. */
/**
* FNXC:ProviderRateLimitIsolation 2026-07-21-18:00:
* Parks only tasks routed through the provider whose API limit was detected.
*/
usageLimitPauser?: UsageLimitPauser;
/** Stuck task detector — monitors agent sessions for stagnation and triggers recovery. */
stuckTaskDetector?: StuckTaskDetector;

View File

@@ -4928,7 +4928,7 @@ export interface MergerOptions {
/** Worktree pool — when provided and `recycleWorktrees` is enabled,
* worktrees are released to the pool instead of being removed. */
pool?: WorktreePool;
/** Usage limit pauser — triggers global pause when API limits are detected. */
/** Usage limit pauser — parks only the affected provider-routed task. */
usageLimitPauser?: UsageLimitPauser;
/** Called with the agent session immediately after creation. Enables the
* caller (e.g. dashboard.ts) to track and externally dispose the session
@@ -11023,7 +11023,12 @@ async function runAiAgentForCommit(params: AiAgentParams): Promise<{ success: bo
mergerLog.error(`Agent failed: ${err.message}`);
if (options.usageLimitPauser && isUsageLimitError(err.message)) {
await options.usageLimitPauser.onUsageLimitHit("merger", taskId, err.message);
await options.usageLimitPauser.onUsageLimitHit(
"merger",
taskId,
err.message,
session.state?.model?.provider ?? mergerSessionModel.provider,
);
}
throw err;

View File

@@ -2,10 +2,11 @@
* Rate Limit Retry — wraps async agent work with exponential backoff
* specifically for rate-limit / usage-limit errors.
*
* FNXC:ProviderRateLimitIsolation 2026-07-21-18:00:
* When an AI model returns a rate limit error (429, overloaded, quota, etc.),
* this utility retries the operation with exponential backoff before letting
* the error propagate to the caller's catch block, which triggers a global
* pause via `UsageLimitPauser`.
* provider-scoped task park via `UsageLimitPauser`.
*
* **Backoff strategy:** `delay = min(baseDelayMs × 2^attempt, maxDelayMs)` with
* ±10 % jitter to avoid thundering-herd effects across concurrent agents.
@@ -75,8 +76,9 @@ export interface RateLimitRetryOptions {
* a short flat delay — this budget is separate and does not consume rate-limit
* attempts. All other errors are re-thrown immediately.
*
* FNXC:ProviderRateLimitIsolation 2026-07-21-18:00:
* After all retries are exhausted, the **original** error is thrown so the
* caller's existing catch block can trigger the global pause via
* caller's existing catch block can park the affected provider-routed task via
* `UsageLimitPauser`.
*
* @example
@@ -138,7 +140,8 @@ export async function withRateLimitRetry<T>(
continue;
}
// All retries exhausted — throw so caller can trigger global pause
// FNXC:ProviderRateLimitIsolation 2026-07-21-18:00: exhaustion parks only
// the affected provider-routed task instead of stopping the project.
if (attempt >= maxRetries) {
throw lastError;
}

View File

@@ -21,7 +21,8 @@
* - Exhausted retry budgets escalate to a real failure (task marked failed or error set).
*
* **Not retried via this policy:**
* - Usage-limit errors (handled by `UsageLimitPauser` with global pause)
* - FNXC:ProviderRateLimitIsolation 2026-07-21-18:00: usage-limit errors
* (handled by `UsageLimitPauser` with a provider-scoped task park)
* - User pauses (handled by pause flow)
* - Stuck-task-detector kills (handled by stuck flow)
* - Dependency-abort cleanups (handled by dep-abort flow)

View File

@@ -227,7 +227,7 @@ function withTimeout<T>(
* `maxRetries` times. Non-retryable errors are re-thrown immediately.
*
* Rate-limit / usage-limit errors are NEVER retried by this function — they
* should be handled by `withRateLimitRetry` or trigger a global pause.
* should be handled by `withRateLimitRetry` or a provider-scoped task park.
*
* After all retries are exhausted, the original error is thrown.
*

View File

@@ -48,7 +48,7 @@ A reviewer provider failure (rate limit / flaky network) is NOT a review verdict
Root cause this type exists to fix: the reviewer was the only AI lane that never classified provider errors. A 429 became `UNAVAILABLE`, which drove the fallback ladder to re-hit the SAME rate-limited model instantly (when no validator fallback is configured the "fallback" is a same-model strict-prompt rerun), and `fn_review_step` then told the model "code review remains blocking; retry once", so the executor's agent re-called the tool indefinitely. Observed symptom: 14 identical "Reviewer using model: umans/umans-kimi-k2.7" markers with no review text, one per spawned session, hammering an already-limited provider.
`UNAVAILABLE` is reserved for its real meaning: the reviewer RAN and could not produce a parseable verdict. Provider failures throw this instead so they reach the machinery that already exists to handle them — `withRateLimitRetry` backoff, `UsageLimitPauser` global pause, and the executor's bounded transient recovery. See `getDeferredReviewerFatal` in executor.ts for why the escape needs a deferred re-raise.
`UNAVAILABLE` is reserved for its real meaning: the reviewer RAN and could not produce a parseable verdict. Provider failures throw this instead so they reach the machinery that already exists to handle them — `withRateLimitRetry` backoff, the provider-scoped `UsageLimitPauser` task park, and the executor's bounded transient recovery. See `getDeferredReviewerFatal` in executor.ts for why the escape needs a deferred re-raise.
*/
/** Bounded local retry budget for transient network blips inside one review attempt. */
const REVIEWER_TRANSIENT_MAX_RETRIES = 3;
@@ -56,13 +56,16 @@ const REVIEWER_TRANSIENT_MAX_RETRIES = 3;
export class ReviewerProviderError extends Error {
constructor(
message: string,
/** `usage-limit` → global pause; `transient` → bounded recovery retry. Never `permanent`. */
/** `usage-limit` → affected-task provider pause; `transient` → bounded recovery retry. Never `permanent`. */
public readonly classification: "usage-limit" | "transient",
options?: { cause?: unknown },
options?: { cause?: unknown; provider?: string },
) {
super(message, options);
this.name = "ReviewerProviderError";
this.provider = options?.provider;
}
readonly provider?: string;
}
export interface ReviewResult {
@@ -688,7 +691,7 @@ export async function reviewStep(
/*
FNXC:ReviewerProviderErrors 2026-07-15-11:20:
Classify BEFORE the fallback ladder. The ladder's premise is "this model produced a bad review, try another prompt/model" — a premise that is false for provider failures and actively harmful for them:
- usage-limit: the ladder's same-model strict-prompt rerun (taken whenever no validator fallback is configured) re-hits the exact model that just rate-limited us, with no delay. Escalate instead so `withRateLimitRetry` backs off and `UsageLimitPauser` pauses every lane.
- usage-limit: the ladder's same-model strict-prompt rerun (taken whenever no validator fallback is configured) re-hits the exact model that just rate-limited us, with no delay. Escalate instead so `withRateLimitRetry` backs off and `UsageLimitPauser` parks only this provider-routed task.
- transient: `runAttempt` already spent its bounded backoff budget above, so the network is genuinely down. Escalate to the executor's bounded recovery (requeue with delay) rather than burning the reviewer fallback budget on a dead link.
Neither burns `reviewerFallbackRetryCount` — that budget exists to bound BAD REVIEWS, and spending it on an outage would fail tasks that have nothing wrong with them.
*/
@@ -700,7 +703,10 @@ export async function reviewStep(
if (options.store && options.taskId) {
await options.store.logEntry(options.taskId, escalationMessage).catch(() => undefined);
}
throw new ReviewerProviderError(providerErrorMessage, classification, { cause: err });
throw new ReviewerProviderError(providerErrorMessage, classification, {
cause: err,
provider: validatorProvider,
});
}
if (hasConfiguredFallback) {

View File

@@ -48,7 +48,7 @@ import type {
import { runtimeLog } from "../logger.js";
import { getActiveNotificationService } from "../notifier.js";
import { StuckTaskDetector } from "../stuck-task-detector.js";
import type { UsageLimitPauser } from "../usage-limit-detector.js";
import { UsageLimitPauser } from "../usage-limit-detector.js";
import { SelfHealingManager, VALIDATOR_RUN_STALE_MAX_AGE_MS } from "../self-healing.js";
import { RestartRecoveryCoordinator } from "../restart-recovery-coordinator.js";
import { MeshLeaseManager } from "../mesh-lease-manager.js";
@@ -356,6 +356,12 @@ export class InProcessRuntime
);
}
/*
FNXC:ProviderRateLimitIsolation 2026-07-19-19:10:
Every project runtime owns one usage-limit coordinator and shares it across executor, triage, reviewer, and merger surfaces. Runtime isolation replaced the old dashboard-level construction site; constructing it here prevents a silently undefined pauser while keeping a provider outage local to the affected project/task.
*/
this.usageLimitPauser ??= new UsageLimitPauser(this.taskStore);
// Initialize MessageStore early so TaskExecutor receives send_message capability.
// FNXC:RuntimeSatelliteAsync 2026-06-24-12:45:
// In backend mode, pass the AsyncDataLayer so MessageStore delegates to the
@@ -1002,6 +1008,7 @@ export class InProcessRuntime
{
semaphore: this.projectSemaphore,
stuckTaskDetector: this.stuckTaskDetector,
usageLimitPauser: this.usageLimitPauser,
agentStore: this.agentStore,
pluginRunner: this.pluginRunner,
onSpecifyStart: (t) => {

View File

@@ -11,7 +11,8 @@
* being incorrectly marked as failed due to temporary infrastructure issues.
*
* Contrast with:
* - Usage limit errors: Systemic conditions (rate limits, quota) → trigger global pause
* - FNXC:ProviderRateLimitIsolation 2026-07-21-18:00: usage limit errors are
* provider-local conditions (rate limits, quota) → park the affected task
* - Permanent errors: Code issues, test failures, logic errors → mark task as failed
*/
@@ -103,7 +104,8 @@ export function isSilentTransientError(errorMessage: string): boolean {
/**
* Comprehensive error classification that distinguishes between:
* - 'usage-limit': Rate limits, quota exceeded, billing issues → triggers global pause
* - FNXC:ProviderRateLimitIsolation 2026-07-21-18:00: 'usage-limit' means rate
* limits, quota exceeded, or billing issues → provider-scoped task park
* - 'transient': Network blips, connection errors → move task to "todo" for retry
* - 'permanent': Code errors, test failures, logic errors → mark task as failed
*
@@ -119,7 +121,8 @@ export function classifyError(errorMessage: string): "transient" | "usage-limit"
return "permanent";
}
// Check usage limits first (highest priority - triggers global pause)
// FNXC:ProviderRateLimitIsolation 2026-07-21-18:00: usage limits have highest
// priority and park only the affected provider-routed task.
if (isUsageLimitError(errorMessage)) {
return "usage-limit";
}

View File

@@ -191,7 +191,10 @@ import { buildAgentGatedActionSummary } from "./permanent-agent-gating.js";
export interface TriageProcessorOptions {
pollIntervalMs?: number;
semaphore?: AgentSemaphore;
/** Usage limit pauser — triggers global pause when API limits are detected. */
/**
* FNXC:ProviderRateLimitIsolation 2026-07-21-18:00:
* Parks only tasks routed through the provider whose API limit was detected.
*/
usageLimitPauser?: UsageLimitPauser;
/** Stuck task detector — monitors triage sessions for stagnation and triggers recovery. */
stuckTaskDetector?: StuckTaskDetector;
@@ -1325,6 +1328,7 @@ export class TriageProcessor {
);
this.options.onSpecifyStart?.(task);
let activePlanningProvider: string | undefined;
try {
const detail = await this.store.getTask(task.id);
const currentTask = detail ?? task;
@@ -1578,6 +1582,7 @@ export class TriageProcessor {
settings,
assignedAgent?.runtimeConfig,
);
activePlanningProvider = planningModel.provider;
const planningSessionModelOptions = {
defaultProvider: planningModel.provider,
@@ -1981,12 +1986,14 @@ export class TriageProcessor {
this.stuckAborted.delete(task.id);
await this.handleStuckAbortRequeue(task, "catch");
} else {
// Check if the error is a usage-limit error and trigger global pause
// FNXC:ProviderRateLimitIsolation 2026-07-21-18:00: preserve the resolved
// planning provider so health recovery can resume only its parked lane.
if (this.options.usageLimitPauser && isUsageLimitError(errorMessage)) {
await this.options.usageLimitPauser.onUsageLimitHit(
"triage",
task.id,
errorMessage,
activePlanningProvider,
);
} else if (err instanceof ModelFallbackExhaustedError) {
/*

View File

@@ -1,15 +1,22 @@
/**
* Usage Limit Detector — classifies API errors as usage-limit-related
* and triggers the global pause mechanism when detected.
* Usage Limit Detector — classifies API errors as usage-limit-related and
* parks only the task routed through the unavailable provider.
*
* Usage-limit errors indicate systemic conditions (rate limits, quota exceeded,
* billing issues, overloaded APIs) where continued retrying across multiple
* agents is wasteful. Transient server errors (500, timeout, connection refused)
* are NOT classified as usage-limit errors — they are temporary and may resolve
* on their own via per-session retry.
* Usage-limit errors indicate provider-local conditions (rate limits, quota
* exceeded, billing issues, overloaded APIs). Continued retrying the affected
* task is wasteful, but unrelated providers must remain available. Transient
* server errors (500, timeout, connection refused) are NOT classified as usage-
* limit errors — they are temporary and may resolve on their own via per-session
* retry.
*/
import type { TaskStore } from "@fusion/core";
import type { Task, TaskStore } from "@fusion/core";
import {
resolveExecutorSessionModel,
resolveMergerSessionModel,
resolvePlanningSessionModel,
resolveValidatorSessionModel,
} from "./agent-session-helpers.js";
import { createLogger } from "./logger.js";
const log = createLogger("usage-limit");
@@ -33,10 +40,9 @@ const USAGE_LIMIT_PATTERNS: RegExp[] = [
/**
* Classify whether an error message indicates a usage-limit condition.
*
* Returns `true` for rate limits, overloaded errors, quota/billing issues —
* conditions where all agents should stop. Returns `false` for transient
* server errors (500/502/503/504, timeout, connection refused) that may
* resolve on their own.
* Returns `true` for rate limits, overloaded errors, and quota/billing issues.
* Returns `false` for transient server errors (500/502/503/504, timeout,
* connection refused) that may resolve on their own.
*/
export function isUsageLimitError(errorMessage: string): boolean {
return USAGE_LIMIT_PATTERNS.some((pattern) => pattern.test(errorMessage));
@@ -44,12 +50,9 @@ export function isUsageLimitError(errorMessage: string): boolean {
/**
* Lightweight coordinator that agents call when they detect usage-limit errors.
* Triggers the global pause mechanism by calling
* `store.updateSettings({ globalPause: true, globalPauseReason: "rate-limit" })`.
*
* **Idempotency:** Tracks an internal `paused` flag so that multiple concurrent
* agents hitting limits only trigger one pause. The flag resets when `globalPause`
* is externally set back to `false` (detected by reading settings before pausing).
* It parks only the task that reached the unavailable provider. A provider-local
* outage must never activate the project-wide emergency stop because doing so
* also kills healthy Codex/Claude/Grok work routed through other providers.
*/
/**
* Check if an agent session resolved with an error after exhausting retries.
@@ -74,45 +77,120 @@ export function checkSessionError(session: { state: { errorMessage?: string; err
}
export class UsageLimitPauser {
private paused = false;
constructor(private store: TaskStore) {}
private normalizeProviderId(provider: string): string {
return provider.trim().toLowerCase().replace(/[^a-z0-9._-]+/g, "-").replace(/^-+|-+$/g, "");
}
/**
* Clear only parks created for a provider whose independent health probe has
* transitioned back to usable. Manual/user pauses and every other provider
* reason remain untouched.
*/
async onProviderAvailable(provider: string): Promise<number> {
const providerId = this.normalizeProviderId(provider);
if (!providerId) return 0;
const pausedReason = `provider-rate-limit:${providerId}`;
const tasks = await this.store.listTasks();
const recoverableTasks = tasks.filter((task) =>
task.paused === true
&& task.userPaused !== true
&& (task.pausedReason === pausedReason
|| (providerId === "unknown" && task.pausedReason === "provider-rate-limit")));
/*
FNXC:ProviderRateLimitRecovery 2026-07-19-20:15:
Provider recovery is a health-state transition, never a task call used as a probe. The daemon's independent authenticated usage/capacity monitor invokes this seam only after positive health, and this exact-reason filter ensures recovery cannot clear manual parks, unrelated failure reasons, or another provider's outage.
*/
await Promise.all(recoverableTasks.map(async (task) => {
await this.store.logEntry(task.id, `Provider ${providerId} is available again; resuming task`);
await this.store.pauseTask(task.id, false);
}));
if (recoverableTasks.length > 0) {
log.log(`Provider ${providerId} recovered; resumed ${recoverableTasks.length} task(s)`);
}
return recoverableTasks.length;
}
private taskUsesProvider(
task: Task,
provider: string,
settings: Awaited<ReturnType<TaskStore["getSettings"]>>,
agentType: string,
): boolean {
const providersByActiveLane = agentType === "triage"
? (task.column === "triage" ? [
resolvePlanningSessionModel(task.planningModelProvider, task.planningModelId, settings).provider,
resolveValidatorSessionModel(task.validatorModelProvider, task.validatorModelId, settings).provider,
] : [])
: agentType === "executor"
? (task.column === "in-progress" ? [
resolveExecutorSessionModel(task.modelProvider, task.modelId, settings).provider,
resolveValidatorSessionModel(task.validatorModelProvider, task.validatorModelId, settings).provider,
] : [])
: agentType === "merger"
? (task.column === "in-review" ? [resolveMergerSessionModel(settings, undefined, task).provider] : [])
: [];
const resolvedProviders = providersByActiveLane;
return resolvedProviders.some((candidate) => candidate?.trim().toLowerCase() === provider);
}
/**
* Called by agents when a usage-limit error is detected after retries are exhausted.
* Triggers global pause if not already paused.
* Parks the affected task while leaving every other provider lane running.
*
* @param agentType - The type of agent that hit the limit (e.g., "executor", "triage", "merger")
* @param taskId - The task that was being processed when the limit was hit
* @param errorMessage - The error message from the API
* @param provider - Best-effort provider identifier used in the pause reason
*/
async onUsageLimitHit(agentType: string, taskId: string, errorMessage: string): Promise<void> {
// If we already triggered a pause, check if it was externally reset
if (this.paused) {
const settings = await this.store.getSettings();
if (settings.globalPause) {
// Still paused — no need to trigger again
log.log(`Global pause already active — ignoring duplicate from ${agentType}/${taskId}`);
return;
}
// External reset detected — allow re-triggering
this.paused = false;
}
async onUsageLimitHit(agentType: string, taskId: string, errorMessage: string, provider?: string): Promise<void> {
const providerId = this.normalizeProviderId(provider ?? "unknown") || "unknown";
const pausedReason = `provider-rate-limit:${providerId}`;
this.paused = true;
log.warn(`${agentType} hit usage limit on ${taskId}: ${errorMessage}`);
/*
FNXC:ProviderRateLimitIsolation 2026-07-19-19:10:
A 429 is provider-local, not a project emergency. Park only the task that exhausted retries on that provider so healthy provider lanes continue executing. Keep the provider id in structured pause provenance when the caller can identify it; never persist the full provider response as pause metadata.
*/
log.warn(`${agentType} hit usage limit${providerId ? ` for ${providerId}` : ""} on ${taskId}: ${errorMessage}`);
log.warn(`Matched pattern in error: "${errorMessage.slice(0, 200)}"`);
// Log the triggering error on the task
await this.store.logEntry(
taskId,
`Usage limit detected (${agentType}): ${errorMessage}`,
`Usage limit detected (${agentType}${providerId ? `/${providerId}` : ""}): ${errorMessage}`,
);
// Activate global pause
await this.store.updateSettings({ globalPause: true, globalPauseReason: "rate-limit" });
const [settings, tasks] = await Promise.all([
this.store.getSettings(),
this.store.listTasks(),
]);
const affectedTasks = tasks.filter((task) =>
task.column !== "done"
&& task.column !== "archived"
&& task.paused !== true
&& providerId !== "unknown"
&& this.taskUsesProvider(task, providerId, settings, agentType));
log.warn("⚠ Global pause activated — all automated activity will halt");
// Always include the task that produced the 429 even if its actual provider
// came from a runtime fallback not represented in persisted task settings.
if (!affectedTasks.some((task) => task.id === taskId)) {
const triggeringTask = await this.store.getTask(taskId).catch(() => null);
if (triggeringTask && triggeringTask.paused !== true) affectedTasks.push(triggeringTask);
}
await Promise.all(affectedTasks.map(async (task) => {
if (task.id !== taskId) {
await this.store.logEntry(
task.id,
`Paused because provider ${providerId} reached a usage limit on ${taskId}`,
);
}
await this.store.pauseTask(task.id, true, undefined, { pausedReason });
}));
log.warn(`Paused ${affectedTasks.length} task(s) routed through ${providerId}; other provider lanes remain active`);
}
}