| 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 | */ |
| 16 | import { DurableObject } from "cloudflare:workers"; |
| 17 | |
| 18 | import type { AgentDelivery } from "@g1t/contracts"; |
| 19 | |
| 20 | import { DEFAULT_CAPACITY, MAX_CAPACITY } from "./definition.ts"; |
| 21 | import { 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. */ |
| 24 | const MAX_QUEUE = 50; |
| 25 | /** How long a batch of replies says the agent is working, renewed per batch. */ |
| 26 | const BUSY_MS = 3 * 60_000; |
| 27 | |
| 28 | export 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 | } |