diff --git a/.changeset/isolate-provider-rate-limit-pauses.md b/.changeset/isolate-provider-rate-limit-pauses.md new file mode 100644 index 0000000000..d550127142 --- /dev/null +++ b/.changeset/isolate-provider-rate-limit-pauses.md @@ -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. diff --git a/packages/dashboard/src/__tests__/provider-health-monitor.test.ts b/packages/dashboard/src/__tests__/provider-health-monitor.test.ts new file mode 100644 index 0000000000..67ce02c428 --- /dev/null +++ b/packages/dashboard/src/__tests__/provider-health-monitor.test.ts @@ -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 { + return { + name: "Claude", + icon: "test", + status: "ok", + windows: [{ + label: "Weekly", + percentUsed: 40, + percentLeft: 60, + resetText: "resets in 4d", + }], + ...overrides, + }; +} + +function createStore(tasks: Array>) { + 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(); + }); +}); diff --git a/packages/dashboard/src/provider-health-monitor.ts b/packages/dashboard/src/provider-health-monitor.ts new file mode 100644 index 0000000000..5172b88698 --- /dev/null +++ b/packages/dashboard/src/provider-health-monitor.ts @@ -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; + +export interface ProviderHealthMonitorOptions { + getStores: () => Iterable; + 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 { + 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(); + private timer: ReturnType | 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 { + const stores = Array.from(new Set(this.options.getStores())); + const providers = new Set(); + + 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 { + 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?.(); + } + } + } +} diff --git a/packages/dashboard/src/server.ts b/packages/dashboard/src/server.ts index 7834d7ee3c..564dbf0dc8 100644 --- a/packages/dashboard/src/server.ts +++ b/packages/dashboard/src/server.ts @@ -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) => { const normalizedArgs = normalizeListenArgsForTests(args) as Parameters; @@ -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([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", { diff --git a/packages/dashboard/src/usage.ts b/packages/dashboard/src/usage.ts index d9c41b2e47..3cf1253277 100644 --- a/packages/dashboard/src/usage.ts +++ b/packages/dashboard/src/usage.ts @@ -898,7 +898,7 @@ async function fetchClaudeUsageViaCli(): Promise { * 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 { +export async function fetchClaudeUsage(authStorage?: AuthStorageLike): Promise { const usage: ProviderUsage = { name: "Claude", icon: "🟠", @@ -1251,7 +1251,7 @@ async function loadCodexCredential(): Promise { return null; } -async function fetchCodexUsage(): Promise { +export async function fetchCodexUsage(): Promise { const usage: ProviderUsage = { name: "Codex", icon: "🟢", diff --git a/packages/engine/src/__tests__/in-process-runtime.pg.test.ts b/packages/engine/src/__tests__/in-process-runtime.pg.test.ts index 1418b8e894..3d1a96e110 100644 --- a/packages/engine/src/__tests__/in-process-runtime.pg.test.ts +++ b/packages/engine/src/__tests__/in-process-runtime.pg.test.ts @@ -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" }); diff --git a/packages/engine/src/__tests__/merger-merge-details.test.ts b/packages/engine/src/__tests__/merger-merge-details.test.ts index 48f8c6703f..af9c960b41 100644 --- a/packages/engine/src/__tests__/merger-merge-details.test.ts +++ b/packages/engine/src/__tests__/merger-merge-details.test.ts @@ -185,6 +185,7 @@ function createMockStore(taskOverrides: Partial = {}, 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(); }); }); - diff --git a/packages/engine/src/__tests__/reviewer.test.ts b/packages/engine/src/__tests__/reviewer.test.ts index 93dbc01eed..14566f4172 100644 --- a/packages/engine/src/__tests__/reviewer.test.ts +++ b/packages/engine/src/__tests__/reviewer.test.ts @@ -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); diff --git a/packages/engine/src/__tests__/usage-limit-detector.test.ts b/packages/engine/src/__tests__/usage-limit-detector.test.ts index ddf90403a3..0ee283ebf7 100644 --- a/packages/engine/src/__tests__/usage-limit-detector.test.ts +++ b/packages/engine/src/__tests__/usage-limit-detector.test.ts @@ -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(); + }); }); diff --git a/packages/engine/src/engine-errors.ts b/packages/engine/src/engine-errors.ts index 8f242d0ed8..f07ce04451 100644 --- a/packages/engine/src/engine-errors.ts +++ b/packages/engine/src/engine-errors.ts @@ -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); } diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index dc405eb96b..be9409c886 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -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; diff --git a/packages/engine/src/merger.ts b/packages/engine/src/merger.ts index 25e9f677a6..272d8f10b0 100644 --- a/packages/engine/src/merger.ts +++ b/packages/engine/src/merger.ts @@ -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; diff --git a/packages/engine/src/rate-limit-retry.ts b/packages/engine/src/rate-limit-retry.ts index 4501b02962..f5e487fad0 100644 --- a/packages/engine/src/rate-limit-retry.ts +++ b/packages/engine/src/rate-limit-retry.ts @@ -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( 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; } diff --git a/packages/engine/src/recovery-policy.ts b/packages/engine/src/recovery-policy.ts index 7e4b8791c0..b396fecb68 100644 --- a/packages/engine/src/recovery-policy.ts +++ b/packages/engine/src/recovery-policy.ts @@ -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) diff --git a/packages/engine/src/retry-with-backoff.ts b/packages/engine/src/retry-with-backoff.ts index 6a313c23b5..faa322236f 100644 --- a/packages/engine/src/retry-with-backoff.ts +++ b/packages/engine/src/retry-with-backoff.ts @@ -227,7 +227,7 @@ function withTimeout( * `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. * diff --git a/packages/engine/src/reviewer.ts b/packages/engine/src/reviewer.ts index 57302a6dbb..5eac7fad84 100644 --- a/packages/engine/src/reviewer.ts +++ b/packages/engine/src/reviewer.ts @@ -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) { diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index 1eb7a54a9a..869f8b3648 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -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) => { diff --git a/packages/engine/src/transient-error-detector.ts b/packages/engine/src/transient-error-detector.ts index accb537935..4d77397c25 100644 --- a/packages/engine/src/transient-error-detector.ts +++ b/packages/engine/src/transient-error-detector.ts @@ -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"; } diff --git a/packages/engine/src/triage.ts b/packages/engine/src/triage.ts index 43c4c95d92..ea90032c41 100644 --- a/packages/engine/src/triage.ts +++ b/packages/engine/src/triage.ts @@ -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) { /* diff --git a/packages/engine/src/usage-limit-detector.ts b/packages/engine/src/usage-limit-detector.ts index d0b87af4a4..f81193e99a 100644 --- a/packages/engine/src/usage-limit-detector.ts +++ b/packages/engine/src/usage-limit-detector.ts @@ -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 { + 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>, + 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 { - // 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 { + 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`); } }