diff --git a/apps/server/src/engine/worker.ts b/apps/server/src/engine/worker.ts index eeadea40..b69fe634 100644 --- a/apps/server/src/engine/worker.ts +++ b/apps/server/src/engine/worker.ts @@ -25,6 +25,8 @@ export class TaskWorker { private ticking = false; private stopping = false; private active = new Map(); + // Tasks this worker has started and not yet settled, including ones still being claimed. + private inFlight = new Set(); lastTickAt?: string; constructor( private readonly db: Store, @@ -56,7 +58,7 @@ export class TaskWorker { if (this.timer) clearInterval(this.timer); this.timer = undefined; for (const controller of this.active.values()) controller.abort(); - while (this.active.size || this.ticking) await new Promise((r) => setTimeout(r, 10)); + while (this.inFlight.size || this.ticking) await new Promise((r) => setTimeout(r, 10)); } abort(taskId: string) { this.active.get(taskId)?.abort(); @@ -71,11 +73,12 @@ export class TaskWorker { if (this.ticking) return; this.ticking = true; this.lastTickAt = new Date(this.now()).toISOString(); + const started: Promise[] = []; try { const records = await this.db.scan("tasks"); const due = records.filter( ({ value: t }) => - !this.active.has(t.id) && + !this.inFlight.has(t.id) && (t.status === "queued" || (t.status === "scheduled" && Date.parse(t.nextRunAt ?? "") <= this.now()) || (t.status === "running" && Date.parse(t.leaseUntil ?? "") <= this.now()) || @@ -83,6 +86,7 @@ export class TaskWorker { ); const eligible = []; for (const record of due) { + if (this.inFlight.size + eligible.length >= 3) break; if (record.value.status === "waiting_approval") { const action = record.value.actionId ? await this.db.get<{ status: string; expiresAt?: string }>( @@ -106,12 +110,16 @@ export class TaskWorker { else if (action && ["awaiting_review", "executing"].includes(action.status)) continue; } eligible.push(record); - if (eligible.length === 3) break; } - await Promise.all(eligible.map(({ owner, value }) => this.run(owner, value))); + // Runs outlive the tick, so a long task does not stop later ticks from starting queued work. + for (const { owner, value } of eligible) { + this.inFlight.add(value.id); + started.push(this.run(owner, value).finally(() => this.inFlight.delete(value.id))); + } } finally { this.ticking = false; } + await Promise.all(started); } private async run(owner: string, previous: AgentTask) { if (this.stopping) return; diff --git a/tests/engine.test.ts b/tests/engine.test.ts index 947b1007..7b23ab06 100644 --- a/tests/engine.test.ts +++ b/tests/engine.test.ts @@ -168,6 +168,37 @@ test("pending reviews do not starve queued work", async () => { await db.close(); } }); +test("a long-running task does not starve queued work", async () => { + const db = await createStore(); + let finish: () => void = () => {}; + const finished = new Promise((resolve) => { + finish = resolve; + }); + let started: () => void = () => {}; + const running = new Promise((resolve) => { + started = resolve; + }); + const worker = new TaskWorker(db, async (_owner, current) => { + if (current.id === "long") { + started(); + await finished; + } + return { status: "succeeded" }; + }); + await db.put("owner", "tasks", task("long")); + const first = worker.tick(); + try { + await running; + await db.put("owner", "tasks", task("ready")); + await worker.tick(); + assert.equal((await db.get("owner", "tasks", "ready"))?.status, "succeeded"); + assert.equal((await db.get("owner", "tasks", "long"))?.status, "running"); + } finally { + finish(); + await first; + await db.close(); + } +}); test("run history keeps the time the run started", async () => { const db = await createStore(); try {