Skip to content
682 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 ExtensionInstall,
18 type InstallRequest,
19 type InstallRequests,
20 type NewWorkspaceAgent,
21 type Result,
22 type ServiceBinding,
23 type User,
24 type WorkspaceAgent,
25 UNVERIFIED,
26 askerAccess,
27 awaitsConfirmation,
28 chatClient,
29 cleanRequestNote,
30 extensionById,
31 fail,
32 identityClient,
33 newId,
34 notifyClient,
35 ok,
36 openD1,
37} from "@g1t/contracts";
38
39import { MANAGE_REFUSAL, canManage, canSee } from "./access.ts";
40import { type Definition, applyChanges } from "./definition.ts";
41import { builtinChanges } from "./orchestrator.ts";
42import type { Desk } from "./desk.ts";
43import type { ReplyEnv } from "./reply.ts";
44import { ensureBuiltin } from "./builtin.ts";
45import { type Row, definitionOf, insertAgent, periods, selectAgents, toAgent, updateAgent, versionStatement } from "./store.ts";
46import { TEMPLATES, TEMPLATE_IDS } from "./templates.ts";
47import { readPolicy } from "./policy.ts";
48import { runDue } from "./routines.ts";
49import { onEvents } from "./triggers.ts";
50import { cardAction } from "./cards.ts";
51import { type SessionEnv, sweep } from "./sessions.ts";
52import * as views from "./views.ts";
53import { monthKey } from "./budget.ts";
54import * as extensions from "./extensions.ts";
55import { type Answered, type Person, answerLine, findListing, listRequests, listingPath, openRequest, requestsPath, resolveListing, resolveRequest } from "./installs.ts";
56
57export { Desk } from "./desk.ts";
58
59type Env = ReplyEnv & {
60 IDENTITY: ServiceBinding;
61 /** The audit log. */
62 EVENTS: ServiceBinding;
63 DESKS: DurableObjectNamespace<Desk>;
64};
65
66/** Workspace ids by slug, kept a minute: every call names a workspace by slug. */
67const workspaceIds = new Map<string, { id: string | null; until: number }>();
68
69class Agents {
70 private readonly env: Env;
71 private readonly db: D1Database;
72 private readonly defer: (work: Promise<unknown>) => void;
73
74 // Plain fields, not parameter properties: Node's type stripping does not take those.
75 constructor(env: Env, defer: (work: Promise<unknown>) => void) {
76 this.env = env;
77 this.db = env.DB;
78 this.defer = defer;
79 }
80
81 private async workspaceId(slug: string): Promise<string | null> {
82 const key = slug.toLowerCase();
83 const kept = workspaceIds.get(key);
84 if (kept && kept.until > Date.now()) return kept.id;
85 const workspace = await identityClient(this.env.IDENTITY).getWorkspace(key);
86 if (workspaceIds.size > 5_000) workspaceIds.clear();
87 workspaceIds.set(key, { id: workspace?.id ?? null, until: Date.now() + 60_000 });
88 return workspace?.id ?? null;
89 }
90
91 /** The workspace's id, if the viewer may see its agents. */
92 private async seen(workspace: string, viewer: User | null): Promise<Result<string>> {
93 if (!canSee(viewer, workspace)) return fail("not_found", "There is no such workspace.");
94 const id = await this.workspaceId(workspace);
95 return id ? ok(id) : fail("not_found", "There is no such workspace.");
96 }
97
98 /** Internal, from chat: a person pressed an action on one of agents' cards. */
99 async cardAction(a: AgentCardAction): Promise<Result<CardActionResult>> {
100 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
101 if (!a?.viewer || !canSee(a.viewer, a.workspace ?? "")) return fail("not_found", "No such card.");
102 return cardAction(this.env as unknown as SessionEnv, a);
103 }
104
105 /** What the views need: the workspace, the viewer, and whether they own it. */
106 private async context(workspace: string, viewer: User | null): Promise<Result<views.ViewContext>> {
107 if (viewer && awaitsConfirmation(viewer)) return UNVERIFIED;
108 const seen = await this.seen(workspace, viewer);
109 if (!seen.ok) return seen;
110 return ok({
111 env: this.env as unknown as SessionEnv,
112 db: this.db,
113 slug: workspace.toLowerCase(),
114 workspaceId: seen.value,
115 viewer: viewer!,
116 owner: canManage(viewer, workspace),
117 });
118 }
119
120 /** Runs `view` with the context, or answers why it can't. */
121 async view<T>(a: { workspace: string; viewer: User | null }, view: (ctx: views.ViewContext) => Promise<Result<T>>): Promise<Result<T>> {
122 const ctx = await this.context(a?.workspace ?? "", a?.viewer ?? null);
123 if (!ctx.ok) return ctx;
124 return view(ctx.value);
125 }
126
127 /** The workspace's id, if the viewer may change its agents. */
128 private async managed(workspace: string, viewer: User | null): Promise<Result<string>> {
129 if (awaitsConfirmation(viewer)) return UNVERIFIED;
130 const seen = await this.seen(workspace, viewer);
131 if (!seen.ok) return seen;
132 return canManage(viewer, workspace) ? seen : fail("forbidden", MANAGE_REFUSAL);
133 }
134
135 private async row(workspaceId: string, handle: unknown): Promise<Row | null> {
136 if (typeof handle !== "string") return null;
137 const now = new Date();
138 return this.db
139 .prepare(selectAgents("a.workspace_id = ?3 AND a.handle = ?4 AND a.archived_at IS NULL"))
140 .bind(...periods(now), workspaceId, handle.trim().replace(/^@/, "").toLowerCase())
141 .first<Row>();
142 }
143
144 /** Whether `team` (a slug, or none) is one of the workspace's teams, as the person changing the agent sees them. */
145 private async teamExists(workspace: string, viewer: User, team: string | null): Promise<Result<null>> {
146 if (!team) return ok(null);
147 const found = await identityClient(this.env.IDENTITY)
148 .getTeam(viewer, workspace.toLowerCase(), team)
149 .catch(() => null);
150 return found?.ok ? ok(null) : fail("invalid", `${workspace} has no team called ${team}.`);
151 }
152
153 private async handleTaken(workspaceId: string, handle: string, except: string | null): Promise<boolean> {
154 const found = await this.db
155 .prepare("SELECT id FROM agents WHERE workspace_id = ? AND handle = ? AND archived_at IS NULL")
156 .bind(workspaceId, handle)
157 .first<{ id: string }>();
158 return !!found && found.id !== except;
159 }
160
161 async list(a: { workspace: string; viewer: User | null }): Promise<Result<WorkspaceAgent[]>> {
162 const seen = await this.seen(a.workspace, a.viewer);
163 if (!seen.ok) return seen;
164 const now = new Date();
165 await this.ensureBuiltin(seen.value);
166 const rows = await this.db
167 .prepare(`${selectAgents("a.workspace_id = ?3 AND a.archived_at IS NULL")} ORDER BY a.builtin DESC, a.handle`)
168 .bind(...periods(now), seen.value)
169 .all<Row>();
170 return ok(rows.results.map((row) => toAgent(row, now)));
171 }
172
173 async get(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<WorkspaceAgent>> {
174 const seen = await this.seen(a.workspace, a.viewer);
175 if (!seen.ok) return seen;
176 await this.ensureBuiltin(seen.value);
177 const row = await this.row(seen.value, a.handle);
178 return row ? ok(toAgent(row, new Date())) : fail("not_found", `There is no agent called @${a.handle}.`);
179 }
180
181 /** Internal: agents by id, archived ones too, so old messages still show who wrote them. */
182 async byIds(a: { ids: string[] }): Promise<WorkspaceAgent[]> {
183 const ids = [...new Set((Array.isArray(a.ids) ? a.ids : []).filter((id) => typeof id === "string"))].slice(0, 100);
184 if (!ids.length) return [];
185 const now = new Date();
186 const rows = await this.db
187 .prepare(selectAgents(`a.id IN (${ids.map((_, i) => `?${i + 3}`).join(", ")})`))
188 .bind(...periods(now), ...ids)
189 .all<Row>();
190 return rows.results.map((row) => toAgent(row, now));
191 }
192
193 async create(a: { workspace: string; viewer: User | null; input: NewWorkspaceAgent }): Promise<Result<WorkspaceAgent>> {
194 const managed = await this.managed(a.workspace, a.viewer);
195 if (!managed.ok) return managed;
196 // A new agent starts with the workspace's default monthly budget, unless one was given.
197 const policy = await readPolicy(this.db, managed.value, monthKey(new Date()));
198 const input =
199 policy.default_agent_monthly_micros && a.input?.budget?.monthly_micros === undefined
200 ? { ...a.input, budget: { ...(a.input?.budget ?? {}), monthly_micros: policy.default_agent_monthly_micros } }
201 : a.input;
202 const checked = applyChanges(null, input, TEMPLATE_IDS);
203 if (!checked.ok) return fail("invalid", checked.message);
204 const definition = checked.value;
205 const team = await this.teamExists(a.workspace, a.viewer!, definition.team);
206 if (!team.ok) return team;
207 if (await this.handleTaken(managed.value, definition.handle, null)) {
208 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
209 }
210 const now = new Date().toISOString();
211 const id = newId("agt");
212 const by = a.viewer!.username;
213 try {
214 await this.db.batch([
215 insertAgent(this.db, id, managed.value, definition, { version: 1, created_by: by, created_at: now, updated_at: now }),
216 this.versionStatement(id, 1, definition, by, now),
217 ]);
218 } catch (error) {
219 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
220 throw error;
221 }
222 this.audit(a.viewer!, a.workspace, "create_agent", definition.handle, `Created @${definition.handle} (version 1)`);
223 this.defer(this.hello(a.workspace, managed.value, id, a.viewer!));
224 // Hiring from a catalog role answers every request for that role.
225 if (definition.template) this.defer(this.answerListing(a.workspace, managed.value, `agent:${definition.template}`, a.viewer!));
226 const row = await this.row(managed.value, definition.handle);
227 return ok(toAgent(row!, new Date()));
228 }
229
230 // ── The Marketplace's install requests (./installs.ts) ─────────────────
231
232 /** Requests as the viewer sees them: every one for an owner, their own for anyone else. */
233 async installRequests(a: { workspace: string; viewer: User | null }): Promise<Result<InstallRequests>> {
234 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
235 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
236 if (!seen.ok) return seen;
237 const owner = canManage(a.viewer, a.workspace);
238 const requests = await listRequests(this.db, seen.value, a.viewer!, owner);
239 return ok({ requests, can_resolve: owner });
240 }
241
242 /** A member asks the owners to add something; each owner is notified. */
243 async requestInstall(a: { workspace: string; viewer: User | null; listing: unknown; note?: unknown }): Promise<Result<InstallRequest>> {
244 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
245 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
246 if (!seen.ok) return seen;
247 const viewer = a.viewer!;
248 if ((viewer.kind ?? "user") !== "user") return fail("forbidden", "Only people ask the workspace's owners to add things.");
249 if (canManage(viewer, a.workspace)) return fail("invalid", "You're an owner of this workspace: add it yourself.");
250 const listing = findListing(a.listing);
251 if (!listing) return fail("not_found", "That isn't something a workspace can add yet.");
252 const opened = await openRequest(this.db, seen.value, newId("ins"), listing, viewer, cleanRequestNote(a.note));
253 if (!opened.ok) return opened;
254 const slug = a.workspace.toLowerCase();
255 this.audit(viewer, slug, "request_install", listing.ref, `Asked the owners to add ${listing.name}`, "marketplace");
256 this.defer(this.tellOwners(slug, viewer, opened.value));
257 return opened;
258 }
259
260 /** An owner adds or turns down a request; whoever asked is told. */
261 async resolveInstallRequest(a: { workspace: string; viewer: User | null; id: unknown; status: unknown }): Promise<Result<InstallRequest>> {
262 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
263 if (!managed.ok) {
264 return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners answer requests.") : managed;
265 }
266 if (a.status !== "done" && a.status !== "declined") return fail("invalid", "Answer a request with done or declined.");
267 const answered = await resolveRequest(this.db, managed.value, String(a.id ?? ""), a.status, a.viewer!);
268 if (!answered.ok) return answered;
269 const slug = a.workspace.toLowerCase();
270 const { request } = answered.value;
271 this.audit(a.viewer!, slug, "answer_install_request", request.listing, `${request.status === "done" ? "Added" : "Turned down"} ${request.name} for @${request.requested_by}`, "marketplace");
272 this.defer(this.tellRequester(slug, a.viewer!, answered.value));
273 return ok(request);
274 }
275
276 // ── Extensions installed in the workspace (./extensions.ts) ──────────────
277
278 /** The workspace's installs; any member sees them. */
279 async extensionInstalls(a: { workspace: string; viewer: User | null }): Promise<Result<ExtensionInstall[]>> {
280 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
281 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
282 if (!seen.ok) return seen;
283 return ok(await extensions.listInstalls(this.db, seen.value));
284 }
285
286 /** Owners install a published extension; whoever asked for it is told. */
287 async installExtension(a: { workspace: string; viewer: User | null; extension: unknown }): Promise<Result<ExtensionInstall>> {
288 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
289 if (!managed.ok) return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners install extensions.") : managed;
290 const installed = await extensions.install(this.db, extensionById, managed.value, newId("ins"), String(a.extension ?? ""), a.viewer!.username);
291 if (!installed.ok) return installed;
292 const slug = a.workspace.toLowerCase();
293 this.audit(a.viewer!, slug, "install_extension", installed.value.listing, `Installed ${installed.value.listing} ${installed.value.version}`, "marketplace");
294 this.defer(this.answerListing(slug, managed.value, installed.value.listing, a.viewer!));
295 return installed;
296 }
297
298 /** Owners switch an install on or off. */
299 async setExtensionEnabled(a: { workspace: string; viewer: User | null; listing: unknown; enabled: unknown }): Promise<Result<ExtensionInstall>> {
300 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
301 if (!managed.ok) return managed;
302 const listing = String(a.listing ?? "");
303 if (!extensions.extensionIdOf(listing)) return fail("invalid", "Name an extension as extension:<id>.");
304 const changed = await extensions.setEnabled(this.db, managed.value, listing, a.enabled === true, a.viewer!.username);
305 if (changed.ok) this.audit(a.viewer!, a.workspace, a.enabled === true ? "enable_extension" : "disable_extension", listing, `${a.enabled === true ? "Switched on" : "Switched off"} ${listing}`, "marketplace");
306 return changed;
307 }
308
309 /** Owners cap what an install spends a month. */
310 async setExtensionBudget(a: { workspace: string; viewer: User | null; listing: unknown; monthly_micros: unknown }): Promise<Result<ExtensionInstall>> {
311 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
312 if (!managed.ok) return managed;
313 const listing = String(a.listing ?? "");
314 const micros = a.monthly_micros == null ? null : Number(a.monthly_micros);
315 const changed = await extensions.setBudget(this.db, managed.value, listing, micros);
316 if (changed.ok) this.audit(a.viewer!, a.workspace, "set_extension_budget", listing, `Set ${listing}'s monthly budget`, "marketplace");
317 return changed;
318 }
319
320 /** Owners remove an install. */
321 async uninstallExtension(a: { workspace: string; viewer: User | null; listing: unknown }): Promise<Result<null>> {
322 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
323 if (!managed.ok) return managed;
324 const listing = String(a.listing ?? "");
325 const removed = await extensions.uninstall(this.db, managed.value, listing, a.viewer!.username);
326 if (removed.ok) this.audit(a.viewer!, a.workspace, "uninstall_extension", listing, `Uninstalled ${listing}`, "marketplace");
327 return removed;
328 }
329
330 /** Marks every open request for `listing` added, and tells each person who asked. */
331 private async answerListing(slug: string, workspaceId: string, listing: string, by: User): Promise<void> {
332 try {
333 const answered = await resolveListing(this.db, workspaceId, listing, by as Person);
334 await Promise.all(answered.map((one) => this.tellRequester(slug.toLowerCase(), by, one)));
335 } catch (error) {
336 console.error("agents: requests for a listing were not answered", listing, String(error));
337 }
338 }
339
340 /** Every owner hears of a new request (at most 20 of them). */
341 private async tellOwners(slug: string, asker: User, request: InstallRequest): Promise<void> {
342 if (!this.env.NOTIFY) return;
343 try {
344 const members = await identityClient(this.env.IDENTITY).listMembers(slug, asker);
345 if (!members.ok) throw new Error(members.error.message);
346 const owners = members.value.filter((member) => member.role === "owner").slice(0, 20);
347 const notify = notifyClient(this.env.NOTIFY);
348 await Promise.all(
349 owners.map((owner) =>
350 notify
351 .notify(
352 { username: owner.username },
353 {
354 id: `install-request:${request.id}:${owner.username}`,
355 kind: "approval",
356 workspace: slug,
357 title: `@${asker.username} asks you to add ${request.name}`,
358 body: request.note ?? (request.kind === "agent" ? "An agent from the catalog. Add it, or turn the request down." : request.kind === "extension" ? "An extension. Install it, or turn the request down." : "An integration. Connect it, or turn the request down."),
359 href: requestsPath(slug),
360 actor: { kind: "user", id: asker.id, name: asker.username, avatar: asker.avatar ?? null, avatar_seed: null },
361 created_at: new Date().toISOString(),
362 },
363 )
364 .catch(() => undefined),
365 ),
366 );
367 } catch (error) {
368 console.error("agents: owners were not told of an install request", request.id, String(error));
369 }
370 }
371
372 /** The person who asked hears their request was answered. */
373 private async tellRequester(slug: string, by: User, answered: Answered): Promise<void> {
374 if (!this.env.NOTIFY) return;
375 const { request } = answered;
376 const line = answerLine(request, `@${by.username}`);
377 await notifyClient(this.env.NOTIFY)
378 .notify(
379 { user_id: answered.requested_by_id },
380 {
381 id: `install-answer:${request.id}`,
382 kind: "inbox",
383 workspace: slug,
384 title: line.title,
385 body: line.body,
386 href: request.status === "done" ? listingPath(slug, request.listing) : requestsPath(slug),
387 actor: { kind: "user", id: by.id, name: by.username, avatar: by.avatar ?? null, avatar_seed: null },
388 created_at: new Date().toISOString(),
389 },
390 )
391 .catch((error: unknown) => console.error("agents: a requester was not told", request.id, String(error)));
392 }
393
394 /**
395 * A new agent's first words: the DM with the person who made it is
396 * opened, and the agent says hello there in its own voice, as a reply
397 * billed like any other (a fixed hello when no model can be used).
398 * After the answer; a failure only logs.
399 */
400 private async hello(workspace: string, workspaceId: string, agentId: string, creator: User): Promise<void> {
401 if ((creator.kind ?? "user") !== "user") return;
402 try {
403 const dm = await chatClient(this.env.CHAT).openDm(workspace, creator, [{ kind: "agent", id: agentId }]);
404 if (!dm.ok) throw new Error(dm.error.message);
405 const desk = this.env.DESKS.get(this.env.DESKS.idFromName(agentId));
406 await desk.take({
407 workspace,
408 workspace_id: workspaceId,
409 channel_id: dm.value.id,
410 channel_kind: "dm",
411 channel_name: null,
412 agent_id: agentId,
413 // One hello per agent, however often this runs.
414 message_id: `hello:${agentId}`,
415 thread_root: null,
416 asked_by: creator.id,
417 hops: 0,
418 asker: askerAccess(creator, workspace),
419 hello: true,
420 });
421 } catch (error) {
422 console.error("agents: a new agent's hello was not sent", agentId, String(error));
423 }
424 }
425
426 async update(a: { workspace: string; handle: string; viewer: User | null; changes: Partial<NewWorkspaceAgent> }): Promise<Result<WorkspaceAgent>> {
427 const managed = await this.managed(a.workspace, a.viewer);
428 if (!managed.ok) return managed;
429 const row = await this.row(managed.value, a.handle);
430 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
431 const before = definitionOf(row);
432 // The built-in @g1t keeps who it is and its job; the rest is the workspace's.
433 const allowed = row.builtin ? builtinChanges(before, a.changes) : { ok: true as const, value: a.changes };
434 if (!allowed.ok) return fail("invalid", allowed.message);
435 const checked = applyChanges(before, allowed.value, TEMPLATE_IDS, { builtin: !!row.builtin });
436 if (!checked.ok) return fail("invalid", checked.message);
437 const definition = checked.value;
438 if (JSON.stringify(definition) === JSON.stringify(before)) return ok(toAgent(row, new Date()));
439 if (definition.team !== before.team) {
440 const team = await this.teamExists(a.workspace, a.viewer!, definition.team);
441 if (!team.ok) return team;
442 }
443 if (definition.handle !== before.handle && (await this.handleTaken(managed.value, definition.handle, row.id))) {
444 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
445 }
446 const now = new Date().toISOString();
447 const version = row.version + 1;
448 const by = a.viewer!.username;
449 try {
450 const [updated] = await this.db.batch([
451 // Only from the version read: two owners saving at once never lose
452 // one's change silently; the second is told to look again.
453 updateAgent(this.db, row.id, row.version, definition, { version, updated_at: now }),
454 this.db
455 .prepare(
456 `INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at)
457 SELECT ?1, ?2, ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM agents WHERE id = ?1 AND version = ?2 AND updated_at = ?5)`,
458 )
459 .bind(row.id, version, JSON.stringify(definition), by, now),
460 ]);
461 if (!updated.meta.changes) return fail("conflict", `@${before.handle} was changed meanwhile. Reload it and try again.`);
462 } catch (error) {
463 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
464 throw error;
465 }
466 const renamed = definition.handle !== before.handle ? ` (was @${before.handle})` : "";
467 this.audit(a.viewer!, a.workspace, "update_agent", definition.handle, `Changed @${definition.handle}${renamed} to version ${version}`);
468 const saved = await this.row(managed.value, definition.handle);
469 return ok(toAgent(saved!, new Date()));
470 }
471
472 async archive(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<null>> {
473 const managed = await this.managed(a.workspace, a.viewer);
474 if (!managed.ok) return managed;
475 const row = await this.row(managed.value, a.handle);
476 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
477 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.");
478 const now = new Date().toISOString();
479 await this.db.prepare("UPDATE agents SET archived_at = ?, updated_at = ? WHERE id = ? AND archived_at IS NULL").bind(now, now, row.id).run();
480 this.audit(a.viewer!, a.workspace, "archive_agent", row.handle, `Archived @${row.handle}`);
481 return ok(null);
482 }
483
484 /**
485 * Makes the workspace's built-in @g1t if it does not exist yet: an
486 * ordinary agent row, marked builtin, at version 1. Safe to call on
487 * every request; once seen, an isolate does not write again.
488 */
489 private async ensureBuiltin(workspaceId: string): Promise<void> {
490 await ensureBuiltin(this.db, workspaceId);
491 }
492
493 /** Internal: the workspace's @g1t, made if need be, for the chat service. */
494 async builtin(a: { workspace: string; workspace_id: string }): Promise<Result<WorkspaceAgent>> {
495 if (typeof a?.workspace_id !== "string" || !a.workspace_id) return fail("invalid", "Name the workspace by id.");
496 await this.ensureBuiltin(a.workspace_id);
497 const now = new Date();
498 const row = await this.db
499 .prepare(selectAgents("a.workspace_id = ?3 AND a.builtin = 1"))
500 .bind(...periods(now), a.workspace_id)
501 .first<Row>();
502 return row ? ok(toAgent(row, now)) : fail("not_found", "This workspace's @g1t could not be made.");
503 }
504
505 /**
506 * Hands a message to the agent's desk, which answers it in the
507 * background. Returns as soon as the desk holds it.
508 */
509 async deliver(delivery: AgentDelivery): Promise<Result<null>> {
510 const fields = ["workspace", "workspace_id", "channel_id", "agent_id", "message_id", "asked_by"] as const;
511 if (!delivery || fields.some((field) => typeof delivery[field] !== "string" || !delivery[field])) {
512 return fail("invalid", "A delivery names the workspace, channel, agent, message and who asked.");
513 }
514 if (delivery.channel_kind !== "channel" && delivery.channel_kind !== "dm") return fail("invalid", "channel_kind is channel or dm.");
515 await this.ensureBuiltin(delivery.workspace_id);
516 const desk = this.env.DESKS.get(this.env.DESKS.idFromName(delivery.agent_id));
517 await desk.take({ ...delivery, hops: Math.max(0, Math.floor(Number(delivery.hops) || 0)), thread_root: delivery.thread_root ?? null });
518 return ok(null);
519 }
520
521 private versionStatement(agentId: string, version: number, d: Definition, by: string, at: string): D1PreparedStatement {
522 return versionStatement(this.db, agentId, version, d, by, at);
523 }
524
525 /**
526 * Records a change to an agent in the workspace's audit log, after the
527 * answer. Never fails the change: a log that cannot be written is logged.
528 * The events service's audit contract speaks camelCase (Rust's
529 * `NewAuditEntry`).
530 */
531 private audit(actor: User, workspace: string, action: string, handle: string, message: string, area: "agents" | "marketplace" = "agents"): void {
532 const kind = actor.kind === "workspace" ? "workspace" : actor.kind === "agent" ? "agent" : actor.kind === "system" ? "system" : "person";
533 const entry = {
534 actorKind: kind,
535 actor: actor.username,
536 actorId: actor.id,
537 agent: null,
538 onBehalfOf: null,
539 runId: null,
540 runKind: null,
541 credentialId: null,
542 action,
543 surface: "web",
544 workspace: workspace.toLowerCase(),
545 repo: null,
546 number: null,
547 gitRef: null,
548 path: `${area}/${handle}`,
549 outcome: "allowed",
550 rule: area === "marketplace" && action === "request_install" ? "member" : "owner",
551 result: "ok",
552 message,
553 requestId: `req_${crypto.randomUUID()}`,
554 };
555 this.defer(
556 this.env.EVENTS.fetch("https://service/rpc/audit_record", {
557 method: "POST",
558 headers: { "content-type": "application/json" },
559 body: JSON.stringify({ entries: [entry] }),
560 })
561 .then((response) => {
562 if (!response.ok) throw new Error(`status ${response.status}`);
563 })
564 .catch((error: unknown) => console.error("agents: audit entry not recorded", action, handle, String(error))),
565 );
566 }
567}
568
569/** One RPC method's answer. */
570async function answer(service: Agents, method: string, args: any): Promise<Response> {
571 switch (method) {
572 case "list":
573 return Response.json(await service.list(args));
574 case "get":
575 return Response.json(await service.get(args));
576 case "by_ids":
577 return Response.json(await service.byIds(args));
578 case "create":
579 return Response.json(await service.create(args));
580 case "update":
581 return Response.json(await service.update(args));
582 case "archive":
583 return Response.json(await service.archive(args));
584 case "builtin":
585 return Response.json(await service.builtin(args));
586 case "templates":
587 return Response.json(TEMPLATES);
588 case "deliver":
589 return Response.json(await service.deliver(args));
590 case "overview":
591 return Response.json(await service.view(args, (ctx) => views.overview(ctx)));
592 case "sessions":
593 return Response.json(await service.view(args, (ctx) => views.listSessions(ctx, args)));
594 case "session":
595 return Response.json(await service.view(args, (ctx) => views.sessionDetail(ctx, args.id)));
596 case "stop_session":
597 return Response.json(await service.view(args, (ctx) => views.stopSession(ctx, args.id)));
598 case "approve_session":
599 return Response.json(await service.view(args, (ctx) => views.approveSession(ctx, args.id, args.cap_micros)));
600 case "steer_session":
601 return Response.json(await service.view(args, (ctx) => views.steerSession(ctx, args.id, args.body)));
602 case "memories":
603 return Response.json(await service.view(args, (ctx) => views.memories(ctx, args.handle)));
604 case "remember":
605 return Response.json(await service.view(args, (ctx) => views.remember(ctx, args.handle, args.input)));
606 case "update_memory":
607 return Response.json(await service.view(args, (ctx) => views.updateMemory(ctx, args.handle, args.id, args.changes)));
608 case "forget":
609 return Response.json(await service.view(args, (ctx) => views.forget(ctx, args.handle, args.id)));
610 case "routines":
611 return Response.json(await service.view(args, (ctx) => views.routines(ctx, args.handle)));
612 case "save_routine":
613 return Response.json(await service.view(args, (ctx) => views.saveRoutine(ctx, args.handle, args.input, args.id ?? null)));
614 case "delete_routine":
615 return Response.json(await service.view(args, (ctx) => views.deleteRoutine(ctx, args.handle, args.id)));
616 case "run_routine":
617 return Response.json(await service.view(args, (ctx) => views.runRoutineNow(ctx, args.handle, args.id)));
618 case "spend":
619 return Response.json(await service.view(args, (ctx) => views.spend(ctx, args.handle ?? null, { period: args.period, person: args.person })));
620 case "person_budgets":
621 return Response.json(await service.view(args, (ctx) => views.personBudgetsView(ctx)));
622 case "set_person_budget":
623 return Response.json(await service.view(args, (ctx) => views.setPersonBudget(ctx, args.username, args.monthly_micros)));
624 case "activity":
625 return Response.json(await service.view(args, (ctx) => views.activity(ctx, args.handle)));
626 case "versions":
627 return Response.json(await service.view(args, (ctx) => views.versions(ctx, args.handle)));
628 case "card_action":
629 return Response.json(await service.cardAction(args));
630 case "policy":
631 return Response.json(await service.view(args, (ctx) => views.policy(ctx)));
632 case "install_requests":
633 return Response.json(await service.installRequests(args));
634 case "request_install":
635 return Response.json(await service.requestInstall(args));
636 case "resolve_install_request":
637 return Response.json(await service.resolveInstallRequest(args));
638 case "extension_installs":
639 return Response.json(await service.extensionInstalls(args));
640 case "install_extension":
641 return Response.json(await service.installExtension(args));
642 case "set_extension_enabled":
643 return Response.json(await service.setExtensionEnabled(args));
644 case "set_extension_budget":
645 return Response.json(await service.setExtensionBudget(args));
646 case "uninstall_extension":
647 return Response.json(await service.uninstallExtension(args));
648 case "set_policy":
649 return Response.json(await service.view(args, (ctx) => views.setPolicy(ctx, args.policy)));
650 default:
651 return new Response("Unknown method\n", { status: 404 });
652 }
653}
654
655export default {
656 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
657 const match = new URL(request.url).pathname.match(/^\/rpc\/([a-z_]+)$/);
658 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
659 // A replica near the caller when it asks for one (@g1t/contracts d1.ts).
660 const opened = openD1(env.DB, request);
661 const service = new Agents(Object.create(env, { DB: { value: opened.db } }) as Env, (work) => ctx.waitUntil(work));
662 const args = (await request.json().catch(() => ({}))) as any;
663 return opened.finish(await answer(service, match[1], args));
664 },
665
666 /** Events routines run on, from the events service (SUBSCRIBER_AGENTS). */
667 async queue(batch: MessageBatch<unknown>, env: Env): Promise<void> {
668 await onEvents(env as unknown as SessionEnv, batch.messages.map((message) => message.body as G1tEvent));
669 batch.ackAll();
670 },
671
672 /** Every few minutes: routines that are due, and session steps a desk lost. */
673 async scheduled(_controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
674 const sessions = env as unknown as SessionEnv;
675 ctx.waitUntil(
676 Promise.all([
677 runDue(sessions).catch((error: unknown) => console.error("agents: routines did not run", String(error))),
678 sweep(sessions).catch((error: unknown) => console.error("agents: the session sweep failed", String(error))),
679 ]),
680 );
681 },
682} satisfies ExportedHandler<Env>;