Skip to content
61 linesCodeBlameRaw
1/**
2 * An agent's desk (docs/WORKSPACE.md, "The desk"): one Durable Object per
3 * agent, by its id, where everything addressed to it arrives.
4 *
5 * Today everything that arrives is a message to answer. The desk keeps
6 * them in a queue in its storage and works them from an alarm, at most the
7 * agent's capacity at once, so a burst of messages never runs more replies
8 * in parallel than the agent is allowed, and nothing handed over is lost
9 * if the object is moved. While it works, the agent's row says so
10 * (`busy_until`), which is how the Agents page shows it as working without
11 * asking every desk.
12 *
13 * Tasks, steering and triggers arrive here later; replies are the first
14 * kind of work.
15 */
16import { DurableObject } from "cloudflare:workers";
17
18
19import { DEFAULT_CAPACITY, MAX_CAPACITY } from "./definition.ts";
20import { type DeskWork, type ReplyEnv, reply } from "./reply.ts";
21
22/** Messages a desk holds at most; past this the oldest are dropped, as nobody is waiting on them any more. */
23const MAX_QUEUE = 50;
24/** How long a batch of replies says the agent is working, renewed per batch. */
25const BUSY_MS = 3 * 60_000;
26
27export class Desk extends DurableObject<ReplyEnv> {
28 /** Queues a message for the agent and makes sure the desk is working. Returns at once. */
29 async take(delivery: DeskWork): Promise<void> {
30 const queue = (await this.ctx.storage.get<DeskWork[]>("queue")) ?? [];
31 if (queue.some((held) => held.message_id === delivery.message_id)) return;
32 queue.push(delivery);
33 await this.ctx.storage.put("queue", queue.slice(-MAX_QUEUE));
34 if ((await this.ctx.storage.getAlarm()) === null) await this.ctx.storage.setAlarm(Date.now());
35 }
36
37 /** Works the queue until it is empty, a capacity's worth at a time. */
38 async alarm(): Promise<void> {
39 let agent: string | null = null;
40 for (;;) {
41 const queue = (await this.ctx.storage.get<DeskWork[]>("queue")) ?? [];
42 if (!queue.length) break;
43 agent = queue[0].agent_id;
44 const row = await this.env.DB.prepare("SELECT capacity FROM agents WHERE id = ?").bind(agent).first<{ capacity: number }>();
45 const capacity = Math.min(MAX_CAPACITY, Math.max(1, row?.capacity ?? DEFAULT_CAPACITY));
46 const batch = queue.slice(0, capacity);
47 await this.busy(agent, new Date(Date.now() + BUSY_MS).toISOString());
48 // `reply` records every way it ends and never throws; this is a last guard.
49 await Promise.allSettled(batch.map((delivery) => reply(this.env, delivery)));
50 // Taken off only once worked: what arrived meanwhile stays queued.
51 const done = new Set(batch.map((delivery) => delivery.message_id));
52 const left = ((await this.ctx.storage.get<DeskWork[]>("queue")) ?? []).filter((held) => !done.has(held.message_id));
53 await this.ctx.storage.put("queue", left);
54 }
55 if (agent) await this.busy(agent, null);
56 }
57
58 private async busy(agent: string, until: string | null): Promise<void> {
59 await this.env.DB.prepare("UPDATE agents SET busy_until = ? WHERE id = ?").bind(until, agent).run().catch(() => undefined);
60 }
61}