Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 12 additions & 4 deletions apps/server/src/engine/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ export class TaskWorker {
private ticking = false;
private stopping = false;
private active = new Map<string, AbortController>();
// Tasks this worker has started and not yet settled, including ones still being claimed.
private inFlight = new Set<string>();
lastTickAt?: string;
constructor(
private readonly db: Store,
Expand Down Expand Up @@ -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();
Expand All @@ -71,18 +73,20 @@ export class TaskWorker {
if (this.ticking) return;
this.ticking = true;
this.lastTickAt = new Date(this.now()).toISOString();
const started: Promise<void>[] = [];
try {
const records = await this.db.scan<AgentTask>("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()) ||
t.status === "waiting_approval"),
);
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 }>(
Expand All @@ -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;
Expand Down
31 changes: 31 additions & 0 deletions tests/engine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((resolve) => {
finish = resolve;
});
let started: () => void = () => {};
const running = new Promise<void>((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<AgentTask>("owner", "tasks", "ready"))?.status, "succeeded");
assert.equal((await db.get<AgentTask>("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 {
Expand Down
Loading