| 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 | * Two kinds of work arrive: messages to answer (replies), and sessions to |
| 6 | * advance by a step. They are worked in two lanes from one alarm, each at |
| 7 | * most the agent's capacity at once, so a long session step never keeps |
| 8 | * the agent from answering a quick question, and a burst of messages never |
| 9 | * runs more replies in parallel than the agent is allowed. Both queues live |
| 10 | * in the desk's storage, so nothing handed over is lost if the object |
| 11 | * moves; a session step that never finished is picked up again by the |
| 12 | * sweep (sessions.ts). While it works, the agent's row says so |
| 13 | * (`busy_until`), which is how the Agents page shows it as working without |
| 14 | * asking every desk. |
| 15 | */ |
| 16 | import { DurableObject } from "cloudflare:workers"; |
| 17 | |
| 18 | import { DEFAULT_CAPACITY, MAX_CAPACITY } from "./definition.ts"; |
| 19 | import { type DeskWork, type ReplyEnv, reply } from "./reply.ts"; |
| 20 | import { type SessionEnv, advance } from "./sessions.ts"; |
| 21 | |
| 22 | /** Messages a desk holds at most; past this the oldest are dropped, as nobody is waiting on them any more. */ |
| 23 | const MAX_QUEUE = 50; |
| 24 | /** Sessions waiting for a step a desk holds at most. */ |
| 25 | const MAX_SESSIONS = 200; |
| 26 | /** How long a batch says the agent is working, renewed per batch. */ |
| 27 | const BUSY_MS = 3 * 60_000; |
| 28 | /** How long a lane with nothing to do waits for the other before it looks again. */ |
| 29 | const IDLE_MS = 400; |
| 30 | |
| 31 | export class Desk extends DurableObject<ReplyEnv> { |
| 32 | /** Queues a message for the agent and makes sure the desk is working. Returns at once. */ |
| 33 | async take(delivery: DeskWork): Promise<void> { |
| 34 | const queue = (await this.ctx.storage.get<DeskWork[]>("queue")) ?? []; |
| 35 | if (queue.some((held) => held.message_id === delivery.message_id)) return; |
| 36 | queue.push(delivery); |
| 37 | await this.ctx.storage.put("queue", queue.slice(-MAX_QUEUE)); |
| 38 | await this.ctx.storage.put("agent", delivery.agent_id); |
| 39 | await this.start(); |
| 40 | } |
| 41 | |
| 42 | /** Queues a session for its next step. Returns at once. */ |
| 43 | async session(id: string, agentId: string): Promise<void> { |
| 44 | const sessions = (await this.ctx.storage.get<string[]>("sessions")) ?? []; |
| 45 | if (!sessions.includes(id)) sessions.push(id); |
| 46 | await this.ctx.storage.put("sessions", sessions.slice(-MAX_SESSIONS)); |
| 47 | await this.ctx.storage.put("agent", agentId); |
| 48 | await this.start(); |
| 49 | } |
| 50 | |
| 51 | private running = false; |
| 52 | |
| 53 | private async start(): Promise<void> { |
| 54 | if (this.running) return; |
| 55 | if ((await this.ctx.storage.getAlarm()) === null) await this.ctx.storage.setAlarm(Date.now()); |
| 56 | } |
| 57 | |
| 58 | /** Works both lanes until both are empty. */ |
| 59 | async alarm(): Promise<void> { |
| 60 | this.running = true; |
| 61 | const agent = (await this.ctx.storage.get<string>("agent")) ?? null; |
| 62 | // A lane with nothing to do waits while the other works, since work can |
| 63 | // arrive for it meanwhile; once both are idle, the alarm ends. |
| 64 | const idle = [false, false]; |
| 65 | const lane = async (index: number, work: () => Promise<boolean>) => { |
| 66 | for (;;) { |
| 67 | if (await work()) { |
| 68 | idle[index] = false; |
| 69 | continue; |
| 70 | } |
| 71 | idle[index] = true; |
| 72 | if (idle.every(Boolean)) break; |
| 73 | await new Promise((resolve) => setTimeout(resolve, IDLE_MS)); |
| 74 | } |
| 75 | }; |
| 76 | try { |
| 77 | await Promise.all([lane(0, () => this.replies()), lane(1, () => this.steps())]); |
| 78 | } finally { |
| 79 | this.running = false; |
| 80 | if (agent) await this.busy(agent, null); |
| 81 | // Anything that arrived as the lanes closed is worked by a fresh alarm. |
| 82 | if (await this.pending()) await this.ctx.storage.setAlarm(Date.now() + 100); |
| 83 | } |
| 84 | } |
| 85 | |
| 86 | private async pending(): Promise<boolean> { |
| 87 | const [queue, sessions] = await Promise.all([this.ctx.storage.get<DeskWork[]>("queue"), this.ctx.storage.get<string[]>("sessions")]); |
| 88 | return (queue?.length ?? 0) > 0 || (sessions?.length ?? 0) > 0; |
| 89 | } |
| 90 | |
| 91 | private async capacity(agent: string): Promise<number> { |
| 92 | const row = await this.env.DB.prepare("SELECT capacity FROM agents WHERE id = ?").bind(agent).first<{ capacity: number }>(); |
| 93 | return Math.min(MAX_CAPACITY, Math.max(1, row?.capacity ?? DEFAULT_CAPACITY)); |
| 94 | } |
| 95 | |
| 96 | /** One batch of replies; false when there were none. */ |
| 97 | private async replies(): Promise<boolean> { |
| 98 | const queue = (await this.ctx.storage.get<DeskWork[]>("queue")) ?? []; |
| 99 | if (!queue.length) return false; |
| 100 | const agent = queue[0].agent_id; |
| 101 | const batch = queue.slice(0, await this.capacity(agent)); |
| 102 | await this.busy(agent, new Date(Date.now() + BUSY_MS).toISOString()); |
| 103 | // `reply` records every way it ends and never throws; this is a last guard. |
| 104 | await Promise.allSettled(batch.map((delivery) => reply(this.env, delivery))); |
| 105 | // Taken off only once worked: what arrived meanwhile stays queued. |
| 106 | const done = new Set(batch.map((delivery) => delivery.message_id)); |
| 107 | const left = ((await this.ctx.storage.get<DeskWork[]>("queue")) ?? []).filter((held) => !done.has(held.message_id)); |
| 108 | await this.ctx.storage.put("queue", left); |
| 109 | return true; |
| 110 | } |
| 111 | |
| 112 | /** One batch of session steps; false when there were none. */ |
| 113 | private async steps(): Promise<boolean> { |
| 114 | const sessions = (await this.ctx.storage.get<string[]>("sessions")) ?? []; |
| 115 | if (!sessions.length) return false; |
| 116 | const agent = (await this.ctx.storage.get<string>("agent")) ?? null; |
| 117 | const batch = sessions.slice(0, agent ? await this.capacity(agent) : DEFAULT_CAPACITY); |
| 118 | // Taken off before the step: a session woken again during it (steering, a helper's result) is queued anew. |
| 119 | const taken = new Set(batch); |
| 120 | await this.ctx.storage.put("sessions", ((await this.ctx.storage.get<string[]>("sessions")) ?? []).filter((id) => !taken.has(id))); |
| 121 | if (agent) await this.busy(agent, new Date(Date.now() + BUSY_MS).toISOString()); |
| 122 | await Promise.allSettled( |
| 123 | batch.map((id) => advance(this.env as unknown as SessionEnv, id).catch((error: unknown) => console.error("agents: a session step threw", id, String(error)))), |
| 124 | ); |
| 125 | return true; |
| 126 | } |
| 127 | |
| 128 | private async busy(agent: string, until: string | null): Promise<void> { |
| 129 | await this.env.DB.prepare("UPDATE agents SET busy_until = ? WHERE id = ?").bind(until, agent).run().catch(() => undefined); |
| 130 | } |
| 131 | } |