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