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