feat(KB-147): auto-pause engine on API usage limit errors
- Add UsageLimitPauser class and isUsageLimitError detector for rate limits, overloaded, and quota errors - Integrate usage limit detection into executor, triage, and merger error handlers - Wire shared UsageLimitPauser instance in dashboard startup across all agents - Export UsageLimitPauser and isUsageLimitError from @kb/engine public API - Add comprehensive tests for detector patterns and agent integration
This commit is contained in:
5
.changeset/usage-limit-auto-pause.md
Normal file
5
.changeset/usage-limit-auto-pause.md
Normal file
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@dustinbyrne/kb": patch
|
||||
---
|
||||
|
||||
Auto-pause engine when API usage limits are detected (rate limits, overloaded, quota exceeded). Prevents wasteful retries across concurrent agents.
|
||||
@@ -2,7 +2,7 @@ import { exec } from "node:child_process";
|
||||
import type { AddressInfo } from "node:net";
|
||||
import { TaskStore } from "@kb/core";
|
||||
import { createServer } from "@kb/dashboard";
|
||||
import { TriageProcessor, TaskExecutor, Scheduler, AgentSemaphore, WorktreePool, aiMergeTask, PRIORITY_MERGE } from "@kb/engine";
|
||||
import { TriageProcessor, TaskExecutor, Scheduler, AgentSemaphore, WorktreePool, aiMergeTask, UsageLimitPauser, PRIORITY_MERGE } from "@kb/engine";
|
||||
import { AuthStorage, ModelRegistry } from "@mariozechner/pi-coding-agent";
|
||||
|
||||
function openBrowser(url: string): void {
|
||||
@@ -47,12 +47,22 @@ export async function runDashboard(port: number, opts: { open?: boolean } = {})
|
||||
//
|
||||
const pool = new WorktreePool();
|
||||
|
||||
// ── Usage limit pauser ──────────────────────────────────────────────
|
||||
//
|
||||
// Shared pauser that triggers globalPause when any agent hits an API
|
||||
// usage limit (rate limits, overloaded, quota exceeded). A single
|
||||
// instance is shared across triage, executor, and merger so that the
|
||||
// pause is deduplicated across concurrent agents.
|
||||
//
|
||||
const usageLimitPauser = new UsageLimitPauser(store);
|
||||
|
||||
// AI-powered merge handler (used by the web UI for manual merges).
|
||||
// Wrapped with the shared semaphore so merges count toward the global
|
||||
// concurrency limit alongside triage and execution agents.
|
||||
const rawMerge = (taskId: string) =>
|
||||
aiMergeTask(store, cwd, taskId, {
|
||||
pool,
|
||||
usageLimitPauser,
|
||||
onAgentText: (delta) => process.stdout.write(delta),
|
||||
onAgentTool: (name) => console.log(`[merger] tool: ${name}`),
|
||||
});
|
||||
@@ -153,6 +163,7 @@ export async function runDashboard(port: number, opts: { open?: boolean } = {})
|
||||
{
|
||||
const triage = new TriageProcessor(store, cwd, {
|
||||
semaphore,
|
||||
usageLimitPauser,
|
||||
onSpecifyStart: (t) => console.log(`[engine] Specifying ${t.id}...`),
|
||||
onSpecifyComplete: (t) => console.log(`[engine] ✓ ${t.id} → todo`),
|
||||
onSpecifyError: (t, e) => console.log(`[engine] ✗ ${t.id}: ${e.message}`),
|
||||
@@ -161,6 +172,7 @@ export async function runDashboard(port: number, opts: { open?: boolean } = {})
|
||||
const executor = new TaskExecutor(store, cwd, {
|
||||
semaphore,
|
||||
pool,
|
||||
usageLimitPauser,
|
||||
onStart: (t, p) => console.log(`[engine] Executing ${t.id} in ${p}`),
|
||||
onComplete: (t) => console.log(`[engine] ✓ ${t.id} → in-review`),
|
||||
onError: (t, e) => console.log(`[engine] ✗ ${t.id}: ${e.message}`),
|
||||
|
||||
@@ -69,6 +69,7 @@ function createMockStore() {
|
||||
moveTask: vi.fn().mockResolvedValue({}),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
parseStepsFromPrompt: vi.fn().mockResolvedValue([]),
|
||||
updateSettings: vi.fn().mockResolvedValue({}),
|
||||
getSettings: vi.fn().mockResolvedValue({
|
||||
maxConcurrent: 2,
|
||||
maxWorktrees: 4,
|
||||
@@ -2331,3 +2332,139 @@ describe("task_add_dep tool", () => {
|
||||
expect(store.updateTask).not.toHaveBeenCalledWith("KB-DEP", { status: "failed" });
|
||||
});
|
||||
});
|
||||
|
||||
// ── Usage limit detection in executor ────────────────────────────────
|
||||
|
||||
import { UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
describe("TaskExecutor usage limit detection", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
mockedExistsSync.mockReturnValue(true);
|
||||
});
|
||||
|
||||
it("triggers global pause when executor catches a usage-limit error", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("rate_limit_error: Rate limit exceeded"));
|
||||
|
||||
const onError = vi.fn();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {
|
||||
onError,
|
||||
usageLimitPauser: pauser,
|
||||
});
|
||||
|
||||
await executor.execute({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "in-progress",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
expect(onUsageLimitHitSpy).toHaveBeenCalledWith(
|
||||
"executor",
|
||||
"KB-001",
|
||||
"rate_limit_error: Rate limit exceeded",
|
||||
);
|
||||
expect(store.updateSettings).toHaveBeenCalledWith({ globalPause: true });
|
||||
// Task should still be marked as failed
|
||||
expect(store.updateTask).toHaveBeenCalledWith("KB-001", { status: "failed" });
|
||||
expect(onError).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("does NOT trigger global pause for non-usage-limit errors", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("connection refused"));
|
||||
|
||||
const onError = vi.fn();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {
|
||||
onError,
|
||||
usageLimitPauser: pauser,
|
||||
});
|
||||
|
||||
await executor.execute({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "in-progress",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
expect(onUsageLimitHitSpy).not.toHaveBeenCalled();
|
||||
// Task should still be marked as failed
|
||||
expect(store.updateTask).toHaveBeenCalledWith("KB-001", { status: "failed" });
|
||||
});
|
||||
|
||||
it("works without usageLimitPauser (backward compatible)", async () => {
|
||||
const store = createMockStore();
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("rate_limit_error: Rate limit exceeded"));
|
||||
|
||||
const onError = vi.fn();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { onError });
|
||||
|
||||
await executor.execute({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "in-progress",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
// Should not crash — just mark as failed
|
||||
expect(store.updateTask).toHaveBeenCalledWith("KB-001", { status: "failed" });
|
||||
expect(onError).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("triggers global pause for overloaded error", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("overloaded_error: Overloaded"));
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", {
|
||||
usageLimitPauser: pauser,
|
||||
});
|
||||
|
||||
await executor.execute({
|
||||
id: "KB-002",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "in-progress",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
expect(onUsageLimitHitSpy).toHaveBeenCalledWith(
|
||||
"executor",
|
||||
"KB-002",
|
||||
"overloaded_error: Overloaded",
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -12,6 +12,7 @@ import { PRIORITY_EXECUTE, type AgentSemaphore } from "./concurrency.js";
|
||||
import type { WorktreePool } from "./worktree-pool.js";
|
||||
import { AgentLogger } from "./agent-logger.js";
|
||||
import { executorLog, reviewerLog } from "./logger.js";
|
||||
import { isUsageLimitError, type UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
// Re-export for backward compatibility (tests import from executor.ts)
|
||||
export { summarizeToolArgs } from "./agent-logger.js";
|
||||
@@ -144,6 +145,8 @@ 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. */
|
||||
usageLimitPauser?: UsageLimitPauser;
|
||||
onStart?: (task: Task, worktreePath: string) => void;
|
||||
onComplete?: (task: Task) => void;
|
||||
onError?: (task: Task, error: Error) => void;
|
||||
@@ -452,6 +455,10 @@ export class TaskExecutor {
|
||||
await this.store.logEntry(task.id, "Execution paused — agent terminated, moved to todo");
|
||||
await this.store.moveTask(task.id, "todo");
|
||||
} else {
|
||||
// Check if the error is a usage-limit error and trigger global pause
|
||||
if (this.options.usageLimitPauser && isUsageLimitError(err.message)) {
|
||||
await this.options.usageLimitPauser.onUsageLimitHit("executor", task.id, err.message);
|
||||
}
|
||||
executorLog.error(`✗ ${task.id} execution failed:`, err.message);
|
||||
await this.store.logEntry(task.id, `Execution failed: ${err.message}`);
|
||||
await this.store.updateTask(task.id, { status: "failed" });
|
||||
|
||||
@@ -8,3 +8,4 @@ export { reviewStep, type ReviewType, type ReviewVerdict, type ReviewResult, typ
|
||||
export { createKbAgent, type AgentOptions, type AgentResult } from "./pi.js";
|
||||
export { WorktreePool } from "./worktree-pool.js";
|
||||
export { createLogger, type Logger } from "./logger.js";
|
||||
export { isUsageLimitError, UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
@@ -46,6 +46,7 @@ function createMockStore(taskOverrides: Partial<Task> = {}, allTasks: Task[] = [
|
||||
moveTask: vi.fn().mockResolvedValue(baseTask),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
appendAgentLog: vi.fn().mockResolvedValue(undefined),
|
||||
updateSettings: vi.fn().mockResolvedValue({}),
|
||||
getSettings: vi.fn().mockResolvedValue({ ...DEFAULT_SETTINGS }),
|
||||
emit: vi.fn(),
|
||||
on: vi.fn(),
|
||||
@@ -426,3 +427,109 @@ describe("aiMergeTask — agent log persistence", () => {
|
||||
expect(store.appendAgentLog).toHaveBeenCalledWith("KB-050", "hi", "text", undefined, "merger");
|
||||
});
|
||||
});
|
||||
|
||||
// ── Usage limit detection in merger ──────────────────────────────────
|
||||
|
||||
import { UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
describe("aiMergeTask — usage limit detection", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
mockedExistsSync.mockReturnValue(true);
|
||||
setupHappyPathExecSync();
|
||||
});
|
||||
|
||||
it("triggers global pause when merger catches a usage-limit error", async () => {
|
||||
const store = createMockStore(
|
||||
{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050" },
|
||||
[{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
|
||||
);
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockResolvedValue({
|
||||
session: {
|
||||
prompt: vi.fn().mockRejectedValue(new Error("rate_limit_error: Rate limit exceeded")),
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
} as any);
|
||||
|
||||
await expect(
|
||||
aiMergeTask(store, "/tmp/root", "KB-050", { usageLimitPauser: pauser }),
|
||||
).rejects.toThrow("AI merge failed");
|
||||
|
||||
expect(onUsageLimitHitSpy).toHaveBeenCalledWith(
|
||||
"merger",
|
||||
"KB-050",
|
||||
"rate_limit_error: Rate limit exceeded",
|
||||
);
|
||||
expect(store.updateSettings).toHaveBeenCalledWith({ globalPause: true });
|
||||
});
|
||||
|
||||
it("does NOT trigger global pause for non-usage-limit errors", async () => {
|
||||
const store = createMockStore(
|
||||
{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050" },
|
||||
[{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
|
||||
);
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockResolvedValue({
|
||||
session: {
|
||||
prompt: vi.fn().mockRejectedValue(new Error("connection refused")),
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
} as any);
|
||||
|
||||
await expect(
|
||||
aiMergeTask(store, "/tmp/root", "KB-050", { usageLimitPauser: pauser }),
|
||||
).rejects.toThrow("AI merge failed");
|
||||
|
||||
expect(onUsageLimitHitSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("works without usageLimitPauser (backward compatible)", async () => {
|
||||
const store = createMockStore(
|
||||
{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050" },
|
||||
[{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
|
||||
);
|
||||
|
||||
mockedCreateHaiAgent.mockResolvedValue({
|
||||
session: {
|
||||
prompt: vi.fn().mockRejectedValue(new Error("rate_limit_error: Rate limit exceeded")),
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
} as any);
|
||||
|
||||
// Should not crash — just re-throw
|
||||
await expect(
|
||||
aiMergeTask(store, "/tmp/root", "KB-050"),
|
||||
).rejects.toThrow("AI merge failed");
|
||||
});
|
||||
|
||||
it("triggers global pause for overloaded error", async () => {
|
||||
const store = createMockStore(
|
||||
{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050" },
|
||||
[{ id: "KB-050", worktree: "/tmp/root/.worktrees/KB-050", column: "in-review" } as Task],
|
||||
);
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockResolvedValue({
|
||||
session: {
|
||||
prompt: vi.fn().mockRejectedValue(new Error("overloaded_error: Overloaded")),
|
||||
dispose: vi.fn(),
|
||||
},
|
||||
} as any);
|
||||
|
||||
await expect(
|
||||
aiMergeTask(store, "/tmp/root", "KB-050", { usageLimitPauser: pauser }),
|
||||
).rejects.toThrow("AI merge failed");
|
||||
|
||||
expect(onUsageLimitHitSpy).toHaveBeenCalledWith(
|
||||
"merger",
|
||||
"KB-050",
|
||||
"overloaded_error: Overloaded",
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -5,6 +5,7 @@ import { createKbAgent } from "./pi.js";
|
||||
import type { WorktreePool } from "./worktree-pool.js";
|
||||
import { AgentLogger } from "./agent-logger.js";
|
||||
import { mergerLog } from "./logger.js";
|
||||
import { isUsageLimitError, type UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
/**
|
||||
* Build the merge system prompt. When `includeTaskId` is true (default),
|
||||
@@ -104,6 +105,8 @@ 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. */
|
||||
usageLimitPauser?: UsageLimitPauser;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -269,6 +272,10 @@ export async function aiMergeTask(
|
||||
} catch (err: any) {
|
||||
// Agent failed — try to abort the merge
|
||||
mergerLog.error(`Agent failed: ${err.message}`);
|
||||
// Check if the error is a usage-limit error and trigger global pause
|
||||
if (options.usageLimitPauser && isUsageLimitError(err.message)) {
|
||||
await options.usageLimitPauser.onUsageLimitHit("merger", taskId, err.message);
|
||||
}
|
||||
try {
|
||||
execSync("git reset --merge", { cwd: rootDir, stdio: "pipe" });
|
||||
} catch { /* */ }
|
||||
|
||||
@@ -38,6 +38,7 @@ function createMockStore(tasks: any[] = []) {
|
||||
parseDependenciesFromPrompt: vi.fn().mockResolvedValue([]),
|
||||
logEntry: vi.fn().mockResolvedValue({}),
|
||||
deleteTask: vi.fn().mockResolvedValue({}),
|
||||
updateSettings: vi.fn().mockResolvedValue({}),
|
||||
getSettings: vi.fn().mockResolvedValue({
|
||||
maxConcurrent: 2,
|
||||
maxWorktrees: 4,
|
||||
@@ -952,3 +953,141 @@ describe("TriageProcessor dependency parsing", () => {
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
// ── Usage limit detection in triage ──────────────────────────────────
|
||||
|
||||
import { UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
describe("TriageProcessor usage limit detection", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
});
|
||||
|
||||
it("triggers global pause when triage catches a usage-limit error", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("rate_limit_error: Rate limit exceeded"));
|
||||
|
||||
const onError = vi.fn();
|
||||
const triage = new TriageProcessor(store, "/tmp/test", {
|
||||
onSpecifyError: onError,
|
||||
usageLimitPauser: pauser,
|
||||
});
|
||||
|
||||
await triage.specifyTask({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "triage",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
expect(onUsageLimitHitSpy).toHaveBeenCalledWith(
|
||||
"triage",
|
||||
"KB-001",
|
||||
"rate_limit_error: Rate limit exceeded",
|
||||
);
|
||||
expect(store.updateSettings).toHaveBeenCalledWith({ globalPause: true });
|
||||
// Error callback should still fire
|
||||
expect(onError).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("does NOT trigger global pause for non-usage-limit errors", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("connection refused"));
|
||||
|
||||
const onError = vi.fn();
|
||||
const triage = new TriageProcessor(store, "/tmp/test", {
|
||||
onSpecifyError: onError,
|
||||
usageLimitPauser: pauser,
|
||||
});
|
||||
|
||||
await triage.specifyTask({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "triage",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
expect(onUsageLimitHitSpy).not.toHaveBeenCalled();
|
||||
expect(onError).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("does NOT trigger global pause for ENOENT errors (deleted tasks)", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
const onUsageLimitHitSpy = vi.spyOn(pauser, "onUsageLimitHit");
|
||||
|
||||
const enoentError = Object.assign(
|
||||
new Error("ENOENT: no such file or directory"),
|
||||
{ code: "ENOENT" },
|
||||
);
|
||||
store.updateTask.mockRejectedValue(enoentError);
|
||||
|
||||
const onError = vi.fn();
|
||||
const triage = new TriageProcessor(store, "/tmp/test", {
|
||||
onSpecifyError: onError,
|
||||
usageLimitPauser: pauser,
|
||||
});
|
||||
|
||||
await triage.specifyTask({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "triage",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
expect(onUsageLimitHitSpy).not.toHaveBeenCalled();
|
||||
// ENOENT errors don't call onSpecifyError
|
||||
expect(onError).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("works without usageLimitPauser (backward compatible)", async () => {
|
||||
const store = createMockStore();
|
||||
|
||||
mockedCreateHaiAgent.mockRejectedValue(new Error("rate_limit_error: Rate limit exceeded"));
|
||||
|
||||
const onError = vi.fn();
|
||||
const triage = new TriageProcessor(store, "/tmp/test", {
|
||||
onSpecifyError: onError,
|
||||
});
|
||||
|
||||
await triage.specifyTask({
|
||||
id: "KB-001",
|
||||
title: "Test",
|
||||
description: "Test",
|
||||
column: "triage",
|
||||
dependencies: [],
|
||||
steps: [],
|
||||
currentStep: 0,
|
||||
log: [],
|
||||
createdAt: new Date().toISOString(),
|
||||
updatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
// Should not crash — just call onError
|
||||
expect(onError).toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -6,6 +6,7 @@ import { createKbAgent } from "./pi.js";
|
||||
import { PRIORITY_SPECIFY, type AgentSemaphore } from "./concurrency.js";
|
||||
import { AgentLogger } from "./agent-logger.js";
|
||||
import { triageLog } from "./logger.js";
|
||||
import { isUsageLimitError, type UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
const TRIAGE_SYSTEM_PROMPT = `You are a task specification agent for "kb", an AI-orchestrated task board.
|
||||
|
||||
@@ -157,6 +158,8 @@ Write the PROMPT.md directly using the write tool. Nothing else.`;
|
||||
export interface TriageProcessorOptions {
|
||||
pollIntervalMs?: number;
|
||||
semaphore?: AgentSemaphore;
|
||||
/** Usage limit pauser — triggers global pause when API limits are detected. */
|
||||
usageLimitPauser?: UsageLimitPauser;
|
||||
onSpecifyStart?: (task: Task) => void;
|
||||
onSpecifyComplete?: (task: Task) => void;
|
||||
onSpecifyError?: (task: Task, error: Error) => void;
|
||||
@@ -361,6 +364,10 @@ export class TriageProcessor {
|
||||
if (err.code === "ENOENT") {
|
||||
triageLog.log(`${task.id} no longer exists — skipping`);
|
||||
} else {
|
||||
// Check if the error is a usage-limit error and trigger global pause
|
||||
if (this.options.usageLimitPauser && isUsageLimitError(err.message)) {
|
||||
await this.options.usageLimitPauser.onUsageLimitHit("triage", task.id, err.message);
|
||||
}
|
||||
await this.store.updateTask(task.id, { status: null }).catch(() => {});
|
||||
triageLog.error(`✗ ${task.id} specification failed:`, err.message);
|
||||
this.options.onSpecifyError?.(task, err);
|
||||
|
||||
155
packages/engine/src/usage-limit-detector.test.ts
Normal file
155
packages/engine/src/usage-limit-detector.test.ts
Normal file
@@ -0,0 +1,155 @@
|
||||
import { describe, it, expect, vi, beforeEach } from "vitest";
|
||||
import { isUsageLimitError, UsageLimitPauser } from "./usage-limit-detector.js";
|
||||
|
||||
// ── isUsageLimitError classification tests ───────────────────────────
|
||||
|
||||
describe("isUsageLimitError", () => {
|
||||
describe("should match usage-limit errors", () => {
|
||||
const usageLimitMessages = [
|
||||
// Anthropic overloaded
|
||||
"overloaded_error: Overloaded",
|
||||
"API is overloaded",
|
||||
// Rate limiting
|
||||
"rate_limit_error: Rate limit exceeded",
|
||||
"rate limit exceeded",
|
||||
"Rate Limit Reached",
|
||||
"Too many requests",
|
||||
"too many requests, please retry after 60s",
|
||||
// HTTP status codes
|
||||
"Request failed with status 429",
|
||||
"HTTP 429: Too Many Requests",
|
||||
"529 overloaded",
|
||||
"Status 529",
|
||||
// Quota / billing
|
||||
"quota exceeded for this billing period",
|
||||
"Quota limit reached",
|
||||
"billing account is inactive",
|
||||
"Billing issue detected",
|
||||
"insufficient credit balance",
|
||||
"Insufficient credits",
|
||||
"credit balance too low",
|
||||
];
|
||||
|
||||
for (const msg of usageLimitMessages) {
|
||||
it(`matches: "${msg}"`, () => {
|
||||
expect(isUsageLimitError(msg)).toBe(true);
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
describe("should NOT match transient server errors", () => {
|
||||
const transientMessages = [
|
||||
"Internal Server Error",
|
||||
"Request failed with status 500",
|
||||
"HTTP 502: Bad Gateway",
|
||||
"503 Service Unavailable",
|
||||
"504 Gateway Timeout",
|
||||
"connection refused",
|
||||
"Connection reset by peer",
|
||||
"ECONNREFUSED",
|
||||
"timeout exceeded",
|
||||
"request timed out",
|
||||
"socket hang up",
|
||||
"network error",
|
||||
"ETIMEDOUT",
|
||||
"DNS lookup failed",
|
||||
"getaddrinfo ENOTFOUND",
|
||||
];
|
||||
|
||||
for (const msg of transientMessages) {
|
||||
it(`does not match: "${msg}"`, () => {
|
||||
expect(isUsageLimitError(msg)).toBe(false);
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
it("returns false for empty string", () => {
|
||||
expect(isUsageLimitError("")).toBe(false);
|
||||
});
|
||||
|
||||
it("returns false for generic error messages", () => {
|
||||
expect(isUsageLimitError("Something went wrong")).toBe(false);
|
||||
expect(isUsageLimitError("Unexpected token in JSON")).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
// ── UsageLimitPauser tests ───────────────────────────────────────────
|
||||
|
||||
function createMockStore(globalPause = false) {
|
||||
return {
|
||||
getSettings: vi.fn().mockResolvedValue({ globalPause }),
|
||||
updateSettings: vi.fn().mockResolvedValue({ globalPause: true }),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
} as any;
|
||||
}
|
||||
|
||||
describe("UsageLimitPauser", () => {
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
});
|
||||
|
||||
it("calls store.updateSettings({ globalPause: true }) on usage limit hit", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
|
||||
await pauser.onUsageLimitHit("executor", "KB-001", "rate_limit_error: Rate limit exceeded");
|
||||
|
||||
expect(store.updateSettings).toHaveBeenCalledWith({ globalPause: true });
|
||||
});
|
||||
|
||||
it("logs the triggering error on the task via store.logEntry", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
|
||||
await pauser.onUsageLimitHit("triage", "KB-002", "overloaded_error");
|
||||
|
||||
expect(store.logEntry).toHaveBeenCalledWith(
|
||||
"KB-002",
|
||||
"Usage limit detected (triage): 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 });
|
||||
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
|
||||
await pauser.onUsageLimitHit("executor", "KB-001", "rate limit");
|
||||
await pauser.onUsageLimitHit("triage", "KB-002", "rate limit");
|
||||
await pauser.onUsageLimitHit("merger", "KB-003", "rate limit");
|
||||
|
||||
// updateSettings should only be called once
|
||||
expect(store.updateSettings).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("re-triggers pause if globalPause was externally reset to false", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
|
||||
// First hit — triggers pause
|
||||
store.getSettings.mockResolvedValue({ globalPause: true });
|
||||
await pauser.onUsageLimitHit("executor", "KB-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", "KB-004", "rate limit again");
|
||||
expect(store.updateSettings).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it("includes agent type in the log entry", async () => {
|
||||
const store = createMockStore();
|
||||
const pauser = new UsageLimitPauser(store);
|
||||
|
||||
await pauser.onUsageLimitHit("merger", "KB-005", "quota exceeded");
|
||||
|
||||
expect(store.logEntry).toHaveBeenCalledWith(
|
||||
"KB-005",
|
||||
expect.stringContaining("merger"),
|
||||
);
|
||||
});
|
||||
});
|
||||
93
packages/engine/src/usage-limit-detector.ts
Normal file
93
packages/engine/src/usage-limit-detector.ts
Normal file
@@ -0,0 +1,93 @@
|
||||
/**
|
||||
* Usage Limit Detector — classifies API errors as usage-limit-related
|
||||
* and triggers the global pause mechanism when detected.
|
||||
*
|
||||
* 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.
|
||||
*/
|
||||
|
||||
import type { TaskStore } from "@kb/core";
|
||||
import { createLogger } from "./logger.js";
|
||||
|
||||
const log = createLogger("usage-limit");
|
||||
|
||||
/**
|
||||
* Patterns that indicate API usage/capacity/billing limits.
|
||||
* These are checked case-insensitively against error messages.
|
||||
*/
|
||||
const USAGE_LIMIT_PATTERNS: RegExp[] = [
|
||||
/overloaded/i,
|
||||
/rate[_\s]?limit/i,
|
||||
/too many requests/i,
|
||||
/\b429\b/,
|
||||
/\b529\b/,
|
||||
/quota/i,
|
||||
/billing/i,
|
||||
/\bcredit/i,
|
||||
/insufficient/i,
|
||||
];
|
||||
|
||||
/**
|
||||
* 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.
|
||||
*/
|
||||
export function isUsageLimitError(errorMessage: string): boolean {
|
||||
return USAGE_LIMIT_PATTERNS.some((pattern) => pattern.test(errorMessage));
|
||||
}
|
||||
|
||||
/**
|
||||
* Lightweight coordinator that agents call when they detect usage-limit errors.
|
||||
* Triggers the global pause mechanism by calling `store.updateSettings({ globalPause: true })`.
|
||||
*
|
||||
* **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).
|
||||
*/
|
||||
export class UsageLimitPauser {
|
||||
private paused = false;
|
||||
|
||||
constructor(private store: TaskStore) {}
|
||||
|
||||
/**
|
||||
* Called by agents when a usage-limit error is detected after retries are exhausted.
|
||||
* Triggers global pause if not already paused.
|
||||
*
|
||||
* @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
|
||||
*/
|
||||
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
|
||||
return;
|
||||
}
|
||||
// External reset detected — allow re-triggering
|
||||
this.paused = false;
|
||||
}
|
||||
|
||||
this.paused = true;
|
||||
|
||||
log.warn(`${agentType} hit usage limit on ${taskId}: ${errorMessage}`);
|
||||
|
||||
// Log the triggering error on the task
|
||||
await this.store.logEntry(
|
||||
taskId,
|
||||
`Usage limit detected (${agentType}): ${errorMessage}`,
|
||||
);
|
||||
|
||||
// Activate global pause
|
||||
await this.store.updateSettings({ globalPause: true });
|
||||
|
||||
log.warn("⚠ Global pause activated — all automated activity will halt");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user