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