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