Skip to content
336 linesCodeBlameRaw
1/**
2 * The agents service: a workspace's own agents (docs/WORKSPACE.md,
3 * "Agents"). Their definitions and every version of them, the templates
4 * they start from, each agent's desk, and their replies in chat.
5 *
6 * `@g1t`, the platform's own agent, is not one of these.
7 *
8 * Reached through service bindings: `POST /rpc/<method>`, with snake_case
9 * JSON (@g1t/contracts workspace-agents.ts).
10 */
11import {
12 type AgentDelivery,
13 type NewWorkspaceAgent,
14 type Result,
15 type ServiceBinding,
16 type User,
17 type WorkspaceAgent,
18 UNVERIFIED,
19 awaitsConfirmation,
20 fail,
21 identityClient,
22 newId,
23 ok,
24 openD1,
25} from "@g1t/contracts";
26
27import { MANAGE_REFUSAL, canManage, canSee } from "./access.ts";
28import { type Definition, applyChanges } from "./definition.ts";
29import type { Desk } from "./desk.ts";
30import type { ReplyEnv } from "./reply.ts";
31import { type Row, definitionOf, periods, selectAgents, toAgent } from "./store.ts";
32import { TEMPLATES, TEMPLATE_IDS } from "./templates.ts";
33
34export { Desk } from "./desk.ts";
35
36type Env = ReplyEnv & {
37 IDENTITY: ServiceBinding;
38 /** The audit log. */
39 EVENTS: ServiceBinding;
40 DESKS: DurableObjectNamespace<Desk>;
41};
42
43/** Workspace ids by slug, kept a minute: every call names a workspace by slug. */
44const workspaceIds = new Map<string, { id: string | null; until: number }>();
45
46class Agents {
47 private readonly env: Env;
48 private readonly db: D1Database;
49 private readonly defer: (work: Promise<unknown>) => void;
50
51 // Plain fields, not parameter properties: Node's type stripping does not take those.
52 constructor(env: Env, defer: (work: Promise<unknown>) => void) {
53 this.env = env;
54 this.db = env.DB;
55 this.defer = defer;
56 }
57
58 private async workspaceId(slug: string): Promise<string | null> {
59 const key = slug.toLowerCase();
60 const kept = workspaceIds.get(key);
61 if (kept && kept.until > Date.now()) return kept.id;
62 const workspace = await identityClient(this.env.IDENTITY).getWorkspace(key);
63 if (workspaceIds.size > 5_000) workspaceIds.clear();
64 workspaceIds.set(key, { id: workspace?.id ?? null, until: Date.now() + 60_000 });
65 return workspace?.id ?? null;
66 }
67
68 /** The workspace's id, if the viewer may see its agents. */
69 private async seen(workspace: string, viewer: User | null): Promise<Result<string>> {
70 if (!canSee(viewer, workspace)) return fail("not_found", "There is no such workspace.");
71 const id = await this.workspaceId(workspace);
72 return id ? ok(id) : fail("not_found", "There is no such workspace.");
73 }
74
75 /** The workspace's id, if the viewer may change its agents. */
76 private async managed(workspace: string, viewer: User | null): Promise<Result<string>> {
77 if (awaitsConfirmation(viewer)) return UNVERIFIED;
78 const seen = await this.seen(workspace, viewer);
79 if (!seen.ok) return seen;
80 return canManage(viewer, workspace) ? seen : fail("forbidden", MANAGE_REFUSAL);
81 }
82
83 private async row(workspaceId: string, handle: unknown): Promise<Row | null> {
84 if (typeof handle !== "string") return null;
85 const now = new Date();
86 return this.db
87 .prepare(selectAgents("a.workspace_id = ?3 AND a.handle = ?4 AND a.archived_at IS NULL"))
88 .bind(...periods(now), workspaceId, handle.trim().replace(/^@/, "").toLowerCase())
89 .first<Row>();
90 }
91
92 private async handleTaken(workspaceId: string, handle: string, except: string | null): Promise<boolean> {
93 const found = await this.db
94 .prepare("SELECT id FROM agents WHERE workspace_id = ? AND handle = ? AND archived_at IS NULL")
95 .bind(workspaceId, handle)
96 .first<{ id: string }>();
97 return !!found && found.id !== except;
98 }
99
100 async list(a: { workspace: string; viewer: User | null }): Promise<Result<WorkspaceAgent[]>> {
101 const seen = await this.seen(a.workspace, a.viewer);
102 if (!seen.ok) return seen;
103 const now = new Date();
104 const rows = await this.db
105 .prepare(`${selectAgents("a.workspace_id = ?3 AND a.archived_at IS NULL")} ORDER BY a.handle`)
106 .bind(...periods(now), seen.value)
107 .all<Row>();
108 return ok(rows.results.map((row) => toAgent(row, now)));
109 }
110
111 async get(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<WorkspaceAgent>> {
112 const seen = await this.seen(a.workspace, a.viewer);
113 if (!seen.ok) return seen;
114 const row = await this.row(seen.value, a.handle);
115 return row ? ok(toAgent(row, new Date())) : fail("not_found", `There is no agent called @${a.handle}.`);
116 }
117
118 /** Internal: agents by id, archived ones too, so old messages still show who wrote them. */
119 async byIds(a: { ids: string[] }): Promise<WorkspaceAgent[]> {
120 const ids = [...new Set((Array.isArray(a.ids) ? a.ids : []).filter((id) => typeof id === "string"))].slice(0, 100);
121 if (!ids.length) return [];
122 const now = new Date();
123 const rows = await this.db
124 .prepare(selectAgents(`a.id IN (${ids.map((_, i) => `?${i + 3}`).join(", ")})`))
125 .bind(...periods(now), ...ids)
126 .all<Row>();
127 return rows.results.map((row) => toAgent(row, now));
128 }
129
130 async create(a: { workspace: string; viewer: User | null; input: NewWorkspaceAgent }): Promise<Result<WorkspaceAgent>> {
131 const managed = await this.managed(a.workspace, a.viewer);
132 if (!managed.ok) return managed;
133 const checked = applyChanges(null, a.input, TEMPLATE_IDS);
134 if (!checked.ok) return fail("invalid", checked.message);
135 const definition = checked.value;
136 if (await this.handleTaken(managed.value, definition.handle, null)) {
137 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
138 }
139 const now = new Date().toISOString();
140 const id = newId("agt");
141 const by = a.viewer!.username;
142 try {
143 await this.db.batch([
144 this.db
145 .prepare(
146 `INSERT INTO agents (id, workspace_id, handle, display_name, avatar, role, instructions, personality_preset, personality,
147 routing, budget, autonomy, capacity, template, version, created_by, created_at, updated_at)
148 VALUES (?, ?, ?, ?, NULL, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1, ?, ?, ?)`,
149 )
150 .bind(...this.columns(id, managed.value, definition), by, now, now),
151 this.versionStatement(id, 1, definition, by, now),
152 ]);
153 } catch (error) {
154 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
155 throw error;
156 }
157 this.audit(a.viewer!, a.workspace, "create_agent", definition.handle, `Created @${definition.handle} (version 1)`);
158 const row = await this.row(managed.value, definition.handle);
159 return ok(toAgent(row!, new Date()));
160 }
161
162 async update(a: { workspace: string; handle: string; viewer: User | null; changes: Partial<NewWorkspaceAgent> }): Promise<Result<WorkspaceAgent>> {
163 const managed = await this.managed(a.workspace, a.viewer);
164 if (!managed.ok) return managed;
165 const row = await this.row(managed.value, a.handle);
166 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
167 const before = definitionOf(row);
168 const checked = applyChanges(before, a.changes, TEMPLATE_IDS);
169 if (!checked.ok) return fail("invalid", checked.message);
170 const definition = checked.value;
171 if (JSON.stringify(definition) === JSON.stringify(before)) return ok(toAgent(row, new Date()));
172 if (definition.handle !== before.handle && (await this.handleTaken(managed.value, definition.handle, row.id))) {
173 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
174 }
175 const now = new Date().toISOString();
176 const version = row.version + 1;
177 const by = a.viewer!.username;
178 try {
179 const [updated] = await this.db.batch([
180 // Only from the version read: two owners saving at once never lose
181 // one's change silently; the second is told to look again.
182 this.db
183 .prepare(
184 `UPDATE agents SET handle = ?, display_name = ?, role = ?, instructions = ?, personality_preset = ?, personality = ?,
185 routing = ?, budget = ?, autonomy = ?, capacity = ?, template = ?, version = ?, updated_at = ?
186 WHERE id = ? AND version = ?`,
187 )
188 .bind(...this.columns(null, null, definition), version, now, row.id, row.version),
189 this.db
190 .prepare(
191 `INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at)
192 SELECT ?1, ?2, ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM agents WHERE id = ?1 AND version = ?2 AND updated_at = ?5)`,
193 )
194 .bind(row.id, version, JSON.stringify(definition), by, now),
195 ]);
196 if (!updated.meta.changes) return fail("conflict", `@${before.handle} was changed meanwhile. Reload it and try again.`);
197 } catch (error) {
198 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
199 throw error;
200 }
201 const renamed = definition.handle !== before.handle ? ` (was @${before.handle})` : "";
202 this.audit(a.viewer!, a.workspace, "update_agent", definition.handle, `Changed @${definition.handle}${renamed} to version ${version}`);
203 const saved = await this.row(managed.value, definition.handle);
204 return ok(toAgent(saved!, new Date()));
205 }
206
207 async archive(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<null>> {
208 const managed = await this.managed(a.workspace, a.viewer);
209 if (!managed.ok) return managed;
210 const row = await this.row(managed.value, a.handle);
211 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
212 const now = new Date().toISOString();
213 await this.db.prepare("UPDATE agents SET archived_at = ?, updated_at = ? WHERE id = ? AND archived_at IS NULL").bind(now, now, row.id).run();
214 this.audit(a.viewer!, a.workspace, "archive_agent", row.handle, `Archived @${row.handle}`);
215 return ok(null);
216 }
217
218 /**
219 * Hands a message to the agent's desk, which answers it in the
220 * background. Returns as soon as the desk holds it.
221 */
222 async deliver(delivery: AgentDelivery): Promise<Result<null>> {
223 const fields = ["workspace", "workspace_id", "channel_id", "agent_id", "message_id", "asked_by"] as const;
224 if (!delivery || fields.some((field) => typeof delivery[field] !== "string" || !delivery[field])) {
225 return fail("invalid", "A delivery names the workspace, channel, agent, message and who asked.");
226 }
227 if (delivery.channel_kind !== "channel" && delivery.channel_kind !== "dm") return fail("invalid", "channel_kind is channel or dm.");
228 const desk = this.env.DESKS.get(this.env.DESKS.idFromName(delivery.agent_id));
229 await desk.take({ ...delivery, hops: Math.max(0, Math.floor(Number(delivery.hops) || 0)), thread_root: delivery.thread_root ?? null });
230 return ok(null);
231 }
232
233 /** A definition's columns, in the order the statements above take them. */
234 private columns(id: string | null, workspaceId: string | null, d: Definition): (string | number | null)[] {
235 const head = id && workspaceId ? [id, workspaceId] : [];
236 return [
237 ...head,
238 d.handle,
239 d.display_name,
240 d.role,
241 d.instructions,
242 d.personality_preset,
243 d.personality,
244 JSON.stringify(d.routing),
245 JSON.stringify(d.budget),
246 JSON.stringify(d.autonomy),
247 d.capacity,
248 d.template,
249 ];
250 }
251
252 private versionStatement(agentId: string, version: number, d: Definition, by: string, at: string): D1PreparedStatement {
253 return this.db
254 .prepare("INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at) VALUES (?, ?, ?, ?, ?)")
255 .bind(agentId, version, JSON.stringify(d), by, at);
256 }
257
258 /**
259 * Records a change to an agent in the workspace's audit log, after the
260 * answer. Never fails the change: a log that cannot be written is logged.
261 * The events service's audit contract speaks camelCase (Rust's
262 * `NewAuditEntry`).
263 */
264 private audit(actor: User, workspace: string, action: string, handle: string, message: string): void {
265 const kind = actor.kind === "workspace" ? "workspace" : actor.kind === "agent" ? "agent" : actor.kind === "system" ? "system" : "person";
266 const entry = {
267 actorKind: kind,
268 actor: actor.username,
269 actorId: actor.id,
270 agent: null,
271 onBehalfOf: null,
272 runId: null,
273 runKind: null,
274 credentialId: null,
275 action,
276 surface: "web",
277 workspace: workspace.toLowerCase(),
278 repo: null,
279 number: null,
280 gitRef: null,
281 path: `agents/${handle}`,
282 outcome: "allowed",
283 rule: "owner",
284 result: "ok",
285 message,
286 requestId: `req_${crypto.randomUUID()}`,
287 };
288 this.defer(
289 this.env.EVENTS.fetch("https://service/rpc/audit_record", {
290 method: "POST",
291 headers: { "content-type": "application/json" },
292 body: JSON.stringify({ entries: [entry] }),
293 })
294 .then((response) => {
295 if (!response.ok) throw new Error(`status ${response.status}`);
296 })
297 .catch((error: unknown) => console.error("agents: audit entry not recorded", action, handle, String(error))),
298 );
299 }
300}
301
302/** One RPC method's answer. */
303async function answer(service: Agents, method: string, args: any): Promise<Response> {
304 switch (method) {
305 case "list":
306 return Response.json(await service.list(args));
307 case "get":
308 return Response.json(await service.get(args));
309 case "by_ids":
310 return Response.json(await service.byIds(args));
311 case "create":
312 return Response.json(await service.create(args));
313 case "update":
314 return Response.json(await service.update(args));
315 case "archive":
316 return Response.json(await service.archive(args));
317 case "templates":
318 return Response.json(TEMPLATES);
319 case "deliver":
320 return Response.json(await service.deliver(args));
321 default:
322 return new Response("Unknown method\n", { status: 404 });
323 }
324}
325
326export default {
327 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
328 const match = new URL(request.url).pathname.match(/^\/rpc\/([a-z_]+)$/);
329 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
330 // A replica near the caller when it asks for one (@g1t/contracts d1.ts).
331 const opened = openD1(env.DB, request);
332 const service = new Agents(Object.create(env, { DB: { value: opened.db } }) as Env, (work) => ctx.waitUntil(work));
333 const args = (await request.json().catch(() => ({}))) as any;
334 return opened.finish(await answer(service, match[1], args));
335 },
336} satisfies ExportedHandler<Env>;