Skip to content
62 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
18import type { AgentDelivery } from "@g1t/contracts";
19
20import { DEFAULT_CAPACITY, MAX_CAPACITY } from "./definition.ts";
21import { type ReplyEnv, reply } from "./reply.ts";
22
23/** Messages a desk holds at most; past this the oldest are dropped, as nobody is waiting on them any more. */
24const MAX_QUEUE = 50;
25/** How long a batch of replies says the agent is working, renewed per batch. */
26const BUSY_MS = 3 * 60_000;
27
28export class Desk extends DurableObject<ReplyEnv> {
29 /** Queues a message for the agent and makes sure the desk is working. Returns at once. */
30 async take(delivery: AgentDelivery): Promise<void> {
31 const queue = (await this.ctx.storage.get<AgentDelivery[]>("queue")) ?? [];
32 if (queue.some((held) => held.message_id === delivery.message_id)) return;
33 queue.push(delivery);
34 await this.ctx.storage.put("queue", queue.slice(-MAX_QUEUE));
35 if ((await this.ctx.storage.getAlarm()) === null) await this.ctx.storage.setAlarm(Date.now());
36 }
37
38 /** Works the queue until it is empty, a capacity's worth at a time. */
39 async alarm(): Promise<void> {
40 let agent: string | null = null;
41 for (;;) {
42 const queue = (await this.ctx.storage.get<AgentDelivery[]>("queue")) ?? [];
43 if (!queue.length) break;
44 agent = queue[0].agent_id;
45 const row = await this.env.DB.prepare("SELECT capacity FROM agents WHERE id = ?").bind(agent).first<{ capacity: number }>();
46 const capacity = Math.min(MAX_CAPACITY, Math.max(1, row?.capacity ?? DEFAULT_CAPACITY));
47 const batch = queue.slice(0, capacity);
48 await this.busy(agent, new Date(Date.now() + BUSY_MS).toISOString());
49 // `reply` records every way it ends and never throws; this is a last guard.
50 await Promise.allSettled(batch.map((delivery) => reply(this.env, delivery)));
51 // Taken off only once worked: what arrived meanwhile stays queued.
52 const done = new Set(batch.map((delivery) => delivery.message_id));
53 const left = ((await this.ctx.storage.get<AgentDelivery[]>("queue")) ?? []).filter((held) => !done.has(held.message_id));
54 await this.ctx.storage.put("queue", left);
55 }
56 if (agent) await this.busy(agent, null);
57 }
58
59 private async busy(agent: string, until: string | null): Promise<void> {
60 await this.env.DB.prepare("UPDATE agents SET busy_until = ? WHERE id = ?").bind(until, agent).run().catch(() => undefined);
61 }
62}