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
8 changes: 8 additions & 0 deletions apps/server/src/db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,14 @@ export class Store {
);
return result.rows.map((row) => row.data as T);
}
async countActiveTasks(owner: string, terminalStatuses: string[]): Promise<number> {
const result = await this.db.query(
"SELECT count(*)::int AS n FROM records WHERE owner=$1 AND kind='tasks' AND NOT (COALESCE(data->>'status','') = ANY($2::text[]))",
[owner, terminalStatuses],
);
const row = result.rows[0] as unknown as { n: number } | undefined;
return row?.n ?? 0;
}
async put<T extends { id: string }>(owner: string, kind: string, value: T): Promise<T> {
await this.db.query(
"INSERT INTO records(owner,kind,id,data) VALUES($1,$2,$3,$4::jsonb) ON CONFLICT(owner,kind,id) DO UPDATE SET data=excluded.data,updated_at=now()",
Expand Down
5 changes: 1 addition & 4 deletions apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -184,10 +184,7 @@ export class AgentService {
const id = idempotencyKey ? hash(`task:${idempotencyKey}`) : randomUUID();
const existing = await this.db.get<AgentTask>(owner, "tasks", id);
if (existing) return existing;
if (
(await this.db.list<AgentTask>(owner, "tasks")).filter((t) => !terminal.has(t.status))
.length >= 100
)
if ((await this.db.countActiveTasks(owner, [...terminal])) >= 100)
throw new AppError("Finish or cancel some tasks before adding more", 409);
const titles =
input.kind === "document"
Expand Down
34 changes: 34 additions & 0 deletions tests/persistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,3 +31,37 @@ test("idle Postgres client errors are logged instead of crashing the process", a
await pool.end();
}
});

test("countActiveTasks matches the createTask cap semantics", async () => {
const store = await createStore();
try {
const terminal = ["succeeded", "failed", "cancelled"];
const statuses = [
"queued",
"running",
"scheduled",
"waiting_input",
"waiting_approval",
"paused",
"succeeded",
"failed",
"cancelled",
];
for (let i = 0; i < statuses.length; i++)
await store.put("owner", "tasks", { id: `t${i}`, status: statuses[i] });
await store.put("owner", "tasks", { id: "t-none" });

const expected = (await store.list<{ id: string; status?: string }>("owner", "tasks")).filter(
(item) => !terminal.includes(item.status ?? ""),
).length;
assert.equal(await store.countActiveTasks("owner", terminal), expected);
assert.equal(expected, 7);

await store.put("other", "tasks", { id: "x", status: "queued" });
await store.put("owner", "goals", { id: "g", status: "queued" });
assert.equal(await store.countActiveTasks("owner", terminal), 7);
assert.equal(await store.countActiveTasks("nobody", terminal), 0);
} finally {
await store.close();
}
});
Loading