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