Skip to content
1,089 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 AgentProposal,
15 type AgentRedraft,
16 type DraftReply,
17 MAX_PERSONAL_AGENTS,
18 type AgentDelivery,
19 type CardActionResult,
20 type G1tEvent,
21 type ExtensionInstall,
22 type InstallRequest,
23 type InstallRequests,
24 type McpServer,
25 type AgentEffortCosts,
26 type AgentRecommendation,
27 type AgentRecommendations,
28 type NewWorkspaceAgent,
29 type Result,
30 type ServiceBinding,
31 type User,
32 type WorkspaceAgent,
33 UNVERIFIED,
34 askerAccess,
35 awaitsConfirmation,
36 chatClient,
37 cleanRequestNote,
38 extensionById,
39 fail,
40 identityClient,
41 newId,
42 notifyClient,
43 ok,
44 openD1,
45} from "@g1t/contracts";
46
47import { MANAGE_REFUSAL, NOT_YOURS_REFUSAL, canArchive, canChange, canManage, canSee, canSeeAgent, creatableScope } from "./access.ts";
48import { DRAFT_OUTPUT_TOKENS, draftSystem, jsonIn, proposalFrom, redraftFrom, redraftSystem, startingBudget, trySystem, tryTurns, wordsOf } from "./builder.ts";
49import { callBuilder } from "./builder-call.ts";
50import { type Definition, applyChanges } from "./definition.ts";
51import { listMcpTools } from "./mcp-client.ts";
52import { MAX_MCP_SERVERS, checkMcpName, checkMcpUrl } from "../../../packages/contracts/src/abilities.ts";
53import { builtinChanges } from "./orchestrator.ts";
54import type { Desk } from "./desk.ts";
55import type { ReplyEnv } from "./reply.ts";
56import { ensureBuiltin } from "./builtin.ts";
57import { type Row, definitionOf, insertAgent, isPersonal, periods, selectAgents, toAgent, updateAgent, versionStatement } from "./store.ts";
58import { TEMPLATES, TEMPLATE_IDS } from "./templates.ts";
59import { readPolicy } from "./policy.ts";
60import { runDue } from "./routines.ts";
61import { onEvents } from "./triggers.ts";
62import { cardAction } from "./cards.ts";
63import { type SessionEnv, sweep } from "./sessions.ts";
64import { EFFORT_NAMES, checkDue, effortCostsOf, markResolved, outcomesSince, readRecommendations, recommendationRow, sinceWindow, toRecommendation } from "./recommend.ts";
65import { moveAgentTeams } from "./team-move.ts";
66import { effortOf } from "./routing.ts";
67import * as views from "./views.ts";
68import { computerCommands, computerView, forgetComputer, resetComputer, sleepComputer, wakeComputer } from "./computer.ts";
69import { monthKey } from "./budget.ts";
70import * as extensions from "./extensions.ts";
71import { skillPushes, skillRpc } from "./skill-rpc.ts";
72import { type Answered, type Person, answerLine, findListing, listRequests, listingPath, openRequest, requestsPath, resolveListing, resolveRequest } from "./installs.ts";
73
74export { Desk } from "./desk.ts";
75
76type Env = ReplyEnv & {
77 IDENTITY: ServiceBinding;
78 /** The audit log. */
79 EVENTS: ServiceBinding;
80 /** The runner, for agents' own computers (./computer.ts). */
81 RUNNER?: ServiceBinding;
82 DESKS: DurableObjectNamespace<Desk>;
83};
84
85/** Workspace ids by slug, kept a minute: every call names a workspace by slug. */
86const workspaceIds = new Map<string, { id: string | null; until: number }>();
87
88class Agents {
89 private readonly env: Env;
90 private readonly db: D1Database;
91 private readonly defer: (work: Promise<unknown>) => void;
92
93 // Plain fields, not parameter properties: Node's type stripping does not take those.
94 constructor(env: Env, defer: (work: Promise<unknown>) => void) {
95 this.env = env;
96 this.db = env.DB;
97 this.defer = defer;
98 }
99
100 private async workspaceId(slug: string): Promise<string | null> {
101 const key = slug.toLowerCase();
102 const kept = workspaceIds.get(key);
103 if (kept && kept.until > Date.now()) return kept.id;
104 const workspace = await identityClient(this.env.IDENTITY).getWorkspace(key);
105 if (workspaceIds.size > 5_000) workspaceIds.clear();
106 workspaceIds.set(key, { id: workspace?.id ?? null, until: Date.now() + 60_000 });
107 return workspace?.id ?? null;
108 }
109
110 /** The workspace's id, if the viewer may see its agents. */
111 private async seen(workspace: string, viewer: User | null): Promise<Result<string>> {
112 if (!canSee(viewer, workspace)) return fail("not_found", "There is no such workspace.");
113 const id = await this.workspaceId(workspace);
114 return id ? ok(id) : fail("not_found", "There is no such workspace.");
115 }
116
117 /** Internal, from chat: a person pressed an action on one of agents' cards. */
118 async cardAction(a: AgentCardAction): Promise<Result<CardActionResult>> {
119 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
120 if (!a?.viewer || !canSee(a.viewer, a.workspace ?? "")) return fail("not_found", "No such card.");
121 return cardAction(this.env as unknown as SessionEnv, a);
122 }
123
124 /** What the views need: the workspace, the viewer, and whether they own it. */
125 private async context(workspace: string, viewer: User | null): Promise<Result<views.ViewContext>> {
126 if (viewer && awaitsConfirmation(viewer)) return UNVERIFIED;
127 const seen = await this.seen(workspace, viewer);
128 if (!seen.ok) return seen;
129 return ok({
130 env: this.env as unknown as SessionEnv,
131 db: this.db,
132 slug: workspace.toLowerCase(),
133 workspaceId: seen.value,
134 viewer: viewer!,
135 owner: canManage(viewer, workspace),
136 audit: (action, handle, message) => this.audit(viewer!, workspace, action, handle, message),
137 });
138 }
139
140 /** Runs `view` with the context, or answers why it can't. */
141 async view<T>(a: { workspace: string; viewer: User | null }, view: (ctx: views.ViewContext) => Promise<Result<T>>): Promise<Result<T>> {
142 const ctx = await this.context(a?.workspace ?? "", a?.viewer ?? null);
143 if (!ctx.ok) return ctx;
144 return view(ctx.value);
145 }
146
147 /** The workspace's id, if the viewer may change its agents. */
148 private async managed(workspace: string, viewer: User | null): Promise<Result<string>> {
149 if (awaitsConfirmation(viewer)) return UNVERIFIED;
150 const seen = await this.seen(workspace, viewer);
151 if (!seen.ok) return seen;
152 return canManage(viewer, workspace) ? seen : fail("forbidden", MANAGE_REFUSAL);
153 }
154
155 private async row(workspaceId: string, handle: unknown): Promise<Row | null> {
156 if (typeof handle !== "string") return null;
157 const now = new Date();
158 return this.db
159 .prepare(selectAgents("a.workspace_id = ?3 AND a.handle = ?4 AND a.archived_at IS NULL"))
160 .bind(...periods(now), workspaceId, handle.trim().replace(/^@/, "").toLowerCase())
161 .first<Row>();
162 }
163
164 /**
165 * The teams a new agent joins as it is made, by slug: each must be one the
166 * viewer manages (an owner, or the team's maintainer), as adding anyone
167 * to a team asks. Membership lives on the team, in identity.
168 */
169 private async joinable(workspace: string, viewer: User, asked: unknown): Promise<Result<string[]>> {
170 if (asked === undefined || asked === null) return ok([]);
171 if (!Array.isArray(asked) || asked.some((slug) => typeof slug !== "string")) return fail("invalid", "Teams are a list of team slugs.");
172 const slugs = [...new Set((asked as string[]).map((slug) => slug.trim().replace(/^@/, "").toLowerCase()).filter(Boolean))];
173 if (slugs.length > 20) return fail("invalid", "Add a new agent to at most 20 teams.");
174 const identity = identityClient(this.env.IDENTITY);
175 const found = await Promise.all(slugs.map((slug) => identity.getTeam(viewer, workspace.toLowerCase(), slug).catch(() => null)));
176 for (const [i, team] of found.entries()) {
177 if (!team?.ok) return fail("invalid", `${workspace} has no team called ${slugs[i]}.`);
178 if (!team.value.can_manage) return fail("forbidden", `Only owners and ${team.value.name}'s maintainers add agents to it.`);
179 }
180 return ok(slugs);
181 }
182
183 private async handleTaken(workspaceId: string, handle: string, except: string | null): Promise<boolean> {
184 const found = await this.db
185 .prepare("SELECT id FROM agents WHERE workspace_id = ? AND handle = ? AND archived_at IS NULL")
186 .bind(workspaceId, handle)
187 .first<{ id: string }>();
188 return !!found && found.id !== except;
189 }
190
191 /**
192 * The workspace's agents, and personal agents when asked for: `mine`, the
193 * viewer's own; `all`, every member's for an owner.
194 */
195 async list(a: { workspace: string; viewer: User | null; personal?: unknown }): Promise<Result<WorkspaceAgent[]>> {
196 const seen = await this.seen(a.workspace, a.viewer);
197 if (!seen.ok) return seen;
198 const now = new Date();
199 await this.ensureBuiltin(seen.value);
200 const all = a.personal === "all" && canManage(a.viewer, a.workspace);
201 const mine = (a.personal === "mine" || a.personal === "all") && (a.viewer?.kind ?? "user") === "user" ? (a.viewer?.id ?? "") : "";
202 const rows = await this.db
203 .prepare(
204 `${selectAgents(
205 `a.workspace_id = ?3 AND a.archived_at IS NULL AND (a.scope = 'workspace'${all ? " OR a.scope = 'personal'" : mine ? " OR (a.scope = 'personal' AND a.owner_id = ?4)" : ""})`,
206 )} ORDER BY a.builtin DESC, a.handle`,
207 )
208 .bind(...periods(now), seen.value, ...(mine && !all ? [mine] : []))
209 .all<Row>();
210 return ok(rows.results.map((row) => toAgent(row, now)));
211 }
212
213 async get(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<WorkspaceAgent>> {
214 const seen = await this.seen(a.workspace, a.viewer);
215 if (!seen.ok) return seen;
216 await this.ensureBuiltin(seen.value);
217 const row = await this.row(seen.value, a.handle);
218 // Another member's personal agent is theirs: it isn't there for anyone else but owners.
219 return row && canSeeAgent(a.viewer, a.workspace, row) ? ok(toAgent(row, new Date())) : fail("not_found", `There is no agent called @${a.handle}.`);
220 }
221
222 /** The agent by handle, if the viewer may see it. */
223 private async visible(workspace: string, viewer: User | null, handle: unknown): Promise<Result<{ workspaceId: string; row: Row }>> {
224 if (awaitsConfirmation(viewer)) return UNVERIFIED;
225 const seen = await this.seen(workspace, viewer);
226 if (!seen.ok) return seen;
227 const row = await this.row(seen.value, handle);
228 if (!row || !canSeeAgent(viewer, workspace, row)) return fail("not_found", `There is no agent called @${String(handle ?? "")}.`);
229 return ok({ workspaceId: seen.value, row });
230 }
231
232 /** Internal: agents by id, archived ones too, so old messages still show who wrote them. */
233 async byIds(a: { ids: string[] }): Promise<WorkspaceAgent[]> {
234 const ids = [...new Set((Array.isArray(a.ids) ? a.ids : []).filter((id) => typeof id === "string"))].slice(0, 100);
235 if (!ids.length) return [];
236 const now = new Date();
237 const rows = await this.db
238 .prepare(selectAgents(`a.id IN (${ids.map((_, i) => `?${i + 3}`).join(", ")})`))
239 .bind(...periods(now), ...ids)
240 .all<Row>();
241 return rows.results.map((row) => toAgent(row, now));
242 }
243
244 /**
245 * Who may create an agent, and of which scope: the workspace's id and
246 * policy, and the scope it gets (access.ts `creatableScope`).
247 */
248 private async creatable(workspace: string, viewer: User | null, asked: unknown): Promise<Result<{ workspaceId: string; scope: "workspace" | "personal"; policy: Awaited<ReturnType<typeof readPolicy>> }>> {
249 if (awaitsConfirmation(viewer)) return UNVERIFIED;
250 const seen = await this.seen(workspace, viewer);
251 if (!seen.ok) return seen;
252 const policy = await readPolicy(this.db, seen.value, monthKey(new Date()));
253 const scope = creatableScope(viewer, workspace, asked, policy.members_create_agents);
254 if (!scope.ok) return fail("forbidden", scope.message);
255 return ok({ workspaceId: seen.value, scope: scope.scope, policy });
256 }
257
258 async create(a: { workspace: string; viewer: User | null; input: NewWorkspaceAgent; teams?: unknown }): Promise<Result<WorkspaceAgent>> {
259 const allowed = await this.creatable(a.workspace, a.viewer, a.input?.scope);
260 if (!allowed.ok) return allowed;
261 const { workspaceId, scope, policy } = allowed.value;
262 const personal = scope === "personal";
263 const viewer = a.viewer!;
264 if (personal) {
265 const count = await this.db
266 .prepare("SELECT COUNT(*) AS n FROM agents WHERE workspace_id = ? AND owner_id = ? AND scope = 'personal' AND archived_at IS NULL")
267 .bind(workspaceId, viewer.id)
268 .first<{ n: number }>();
269 if ((count?.n ?? 0) >= MAX_PERSONAL_AGENTS) {
270 return fail("invalid", `You have ${MAX_PERSONAL_AGENTS} personal agents here, the most one person keeps. Archive one to make another.`);
271 }
272 }
273 // A new agent starts with a budget, unless one was given: a member's own
274 // $20 a month and $2 a session, or the workspace's default for its agents.
275 const given = a.input?.budget ?? {};
276 const start = startingBudget(scope, policy.default_agent_monthly_micros);
277 const budget = {
278 monthly_micros: given.monthly_micros !== undefined ? given.monthly_micros : start.monthly_micros,
279 daily_micros: given.daily_micros !== undefined ? given.daily_micros : start.daily_micros,
280 task_micros: given.task_micros !== undefined ? given.task_micros : start.task_micros,
281 };
282 const { scope: _scope, ...rest } = a.input ?? ({} as NewWorkspaceAgent);
283 const input = { ...rest, budget };
284 const checked = applyChanges(null, input, TEMPLATE_IDS);
285 if (!checked.ok) return fail("invalid", checked.message);
286 const definition = checked.value;
287 // A personal agent is on no team: teams are shared, and it answers only its member.
288 const teams = personal ? ok([] as string[]) : await this.joinable(a.workspace, viewer, a.teams);
289 if (!teams.ok) return teams;
290 if (await this.handleTaken(workspaceId, definition.handle, null)) {
291 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
292 }
293 const now = new Date().toISOString();
294 const id = newId("agt");
295 const by = viewer.username;
296 const owned: Record<string, string> = personal ? { scope: "personal", owner_id: viewer.id, owner_username: viewer.username.toLowerCase() } : {};
297 try {
298 await this.db.batch([
299 insertAgent(this.db, id, workspaceId, definition, { version: 1, created_by: by, created_at: now, updated_at: now, ...owned }),
300 this.versionStatement(id, 1, definition, by, now),
301 ]);
302 } catch (error) {
303 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
304 throw error;
305 }
306 this.audit(viewer, a.workspace, "create_agent", definition.handle, `Created ${personal ? "the personal agent " : ""}@${definition.handle} (version 1)`, "agents", personal ? "member" : "owner");
307 // On its teams as it is made, so it knows them from its first message.
308 if (teams.value.length) {
309 const identity = identityClient(this.env.IDENTITY);
310 await Promise.all(
311 teams.value.map((slug) =>
312 identity
313 .setTeamAgent(viewer, a.workspace.toLowerCase(), slug, id)
314 .then((added) => (added.ok ? null : console.error("agents: a new agent was not added to its team", id, slug, added.error.message)))
315 .catch((error: unknown) => console.error("agents: a new agent was not added to its team", id, slug, String(error))),
316 ),
317 );
318 }
319 const row = await this.row(workspaceId, definition.handle);
320 return ok(toAgent(row!, new Date()));
321 }
322
323 // ── The builder: describe it, try it, change it in words (./builder.ts) ──
324
325 /** Drafts a whole agent from a description, charged to the viewer. Nothing is saved. */
326 async draft(a: { workspace: string; viewer: User | null; description: unknown; scope?: unknown }): Promise<Result<AgentProposal>> {
327 const allowed = await this.creatable(a?.workspace ?? "", a?.viewer ?? null, a?.scope);
328 if (!allowed.ok) return allowed;
329 const words = wordsOf(a.description, "description");
330 if (!words.ok) return fail("invalid", words.message);
331 const { workspaceId, scope, policy } = allowed.value;
332 const slug = a.workspace.toLowerCase();
333 const [answer, handles] = await Promise.all([
334 callBuilder(this.env, {
335 slug,
336 workspaceId,
337 viewer: a.viewer!,
338 system: draftSystem(slug, scope),
339 messages: [{ role: "user", content: `What the agent should do:\n\n${words.value}` }],
340 maxOutput: DRAFT_OUTPUT_TOKENS,
341 }),
342 this.db.prepare("SELECT handle FROM agents WHERE workspace_id = ? AND archived_at IS NULL").bind(workspaceId).all<{ handle: string }>(),
343 ]);
344 if (!answer.ok) return fail(answer.code, answer.message);
345 const proposal = proposalFrom(jsonIn(answer.text), {
346 scope,
347 taken: new Set(handles.results.map((row) => row.handle)),
348 budget: startingBudget(scope, policy.default_agent_monthly_micros),
349 });
350 if (!proposal.ok) return fail("invalid", proposal.message);
351 return ok({ ...proposal.value, charged_micros: answer.charged });
352 }
353
354 /** Try it: the unsaved definition answers the conversation so far. Charged to the viewer; nothing is saved. */
355 async tryDraft(a: { workspace: string; viewer: User | null; definition: unknown; messages: unknown }): Promise<Result<DraftReply>> {
356 const given = (a?.definition ?? {}) as NewWorkspaceAgent;
357 const allowed = await this.creatable(a?.workspace ?? "", a?.viewer ?? null, given?.scope);
358 if (!allowed.ok) return allowed;
359 const { scope: _scope, ...rest } = given;
360 // A handle someone has already taken is no reason not to try it.
361 const checked = applyChanges(null, { ...rest, handle: rest.handle || "preview" }, TEMPLATE_IDS);
362 if (!checked.ok) return fail("invalid", checked.message);
363 const turns = tryTurns(a.messages);
364 if (!turns.ok) return fail("invalid", turns.message);
365 const viewer = a.viewer!;
366 const answer = await callBuilder(this.env, {
367 slug: a.workspace.toLowerCase(),
368 workspaceId: allowed.value.workspaceId,
369 viewer,
370 system: trySystem({ workspace: a.workspace.toLowerCase(), definition: checked.value, asker: { username: viewer.username, display_name: viewer.display_username ?? null } }),
371 messages: turns.value,
372 routing: { floor: checked.value.routing.floor, ceiling: checked.value.routing.ceiling },
373 });
374 if (!answer.ok) return fail(answer.code, answer.message);
375 return ok({ text: answer.text, charged_micros: answer.charged });
376 }
377
378 /** Drafts changes from a request in words, for whoever may change the agent. Saved only through `update`. */
379 async redraft(a: { workspace: string; handle: string; viewer: User | null; request: unknown }): Promise<Result<AgentRedraft>> {
380 const found = await this.visible(a?.workspace ?? "", a?.viewer ?? null, a?.handle);
381 if (!found.ok) return found;
382 const { workspaceId, row } = found.value;
383 if (!canChange(a.viewer, a.workspace, row)) return fail("forbidden", isPersonal(row) ? NOT_YOURS_REFUSAL : MANAGE_REFUSAL);
384 const words = wordsOf(a.request, "request");
385 if (!words.ok) return fail("invalid", words.message);
386 const before = definitionOf(row);
387 const slug = a.workspace.toLowerCase();
388 const answer = await callBuilder(this.env, {
389 slug,
390 workspaceId,
391 viewer: a.viewer!,
392 system: redraftSystem(slug, before, !!row.builtin),
393 messages: [{ role: "user", content: `The change to make to @${row.handle}:\n\n${words.value}` }],
394 maxOutput: DRAFT_OUTPUT_TOKENS,
395 });
396 if (!answer.ok) return fail(answer.code, answer.message);
397 const drafted = redraftFrom(jsonIn(answer.text), before, !!row.builtin);
398 if (!drafted.ok) return fail("invalid", drafted.message);
399 return ok({ ...drafted.value, from_version: row.version, charged_micros: answer.charged });
400 }
401
402 /**
403 * Owners make a personal agent a workspace agent: a new workspace agent
404 * with the same handle, definition and every version, and one more
405 * version saying so; the personal one is archived in the same step, its
406 * memory and direct messages kept with it (owners can't read a member's
407 * memory to choose from it, so none is carried).
408 */
409 // ── MCP servers (docs.g1t.sh/guides/agent-abilities/, "MCP servers") ─────
410
411 /** The agent, if the viewer may add servers to it: owners only, a workspace agent or a personal one alike. */
412 private async forMcp(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<{ row: Row; workspaceId: string }>> {
413 const found = await this.visible(a?.workspace ?? "", a?.viewer ?? null, a?.handle);
414 if (!found.ok) return found;
415 if (!canManage(a.viewer, a.workspace)) return fail("forbidden", "Only the workspace's owners add MCP servers to an agent.");
416 return ok({ row: found.value.row, workspaceId: found.value.workspaceId });
417 }
418
419 /** Saves `definition` as the agent's next version, from the version read; what the agent is now. */
420 private async saveVersion(row: Row, definition: Definition, by: User, workspace: string, what: string): Promise<Result<WorkspaceAgent>> {
421 const now = new Date().toISOString();
422 const version = row.version + 1;
423 const [updated] = await this.db.batch([
424 updateAgent(this.db, row.id, row.version, definition, { version, updated_at: now }),
425 this.db
426 .prepare(
427 `INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at)
428 SELECT ?1, ?2, ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM agents WHERE id = ?1 AND version = ?2 AND updated_at = ?5)`,
429 )
430 .bind(row.id, version, JSON.stringify(definition), by.username, now),
431 ]);
432 if (!updated.meta.changes) return fail("conflict", `@${row.handle} was changed meanwhile. Reload it and try again.`);
433 this.audit(by, workspace, "update_agent", row.handle, `${what} (version ${version})`, "agents", "owner");
434 const saved = await this.row(row.workspace_id, row.handle);
435 return ok(toAgent(saved!, new Date()));
436 }
437
438 /**
439 * Adds an MCP server to an agent: the name and address are checked, its
440 * tools are listed from it now, and the agent gets a new version with
441 * the server and its tools as abilities (each a write unless the server
442 * says it only reads). Owners only.
443 */
444 async addMcpServer(a: { workspace: string; handle: string; viewer: User | null; name: unknown; url: unknown }): Promise<Result<McpServer>> {
445 const found = await this.forMcp(a);
446 if (!found.ok) return found;
447 const { row } = found.value;
448 const name = checkMcpName(a.name);
449 if (!name.ok) return fail("invalid", name.message);
450 const url = checkMcpUrl(a.url);
451 if (!url.ok) return fail("invalid", url.message);
452 const before = definitionOf(row);
453 const servers = before.abilities.mcp_servers ?? [];
454 if (servers.length >= MAX_MCP_SERVERS) return fail("limit", `An agent has at most ${MAX_MCP_SERVERS} MCP servers.`);
455 if (servers.some((server) => server.name === name.value)) return fail("conflict", `@${row.handle} already has a server called ${name.value}.`);
456 if (servers.some((server) => server.url === url.value)) return fail("conflict", `@${row.handle} already has that server.`);
457 const listed = await listMcpTools(url.value);
458 const now = new Date().toISOString();
459 const server: McpServer = {
460 id: newId("mcp"),
461 name: name.value,
462 url: url.value,
463 tools: listed.ok ? listed.tools : [],
464 added_by: a.viewer!.username,
465 added_at: now,
466 checked_at: now,
467 problem: listed.ok ? null : listed.message,
468 };
469 const checked = applyChanges(before, {}, TEMPLATE_IDS, { builtin: !!row.builtin, mcp_servers: [...servers, server] });
470 if (!checked.ok) return fail("invalid", checked.message);
471 const saved = await this.saveVersion(row, checked.value, a.viewer!, a.workspace, `Added the MCP server ${server.name} (${server.tools.length} tools) to @${row.handle}`);
472 return saved.ok ? ok(server) : saved;
473 }
474
475 /** Takes an MCP server off an agent, with its abilities' settings: a new version. Owners only. */
476 async removeMcpServer(a: { workspace: string; handle: string; viewer: User | null; id: unknown }): Promise<Result<null>> {
477 const found = await this.forMcp(a);
478 if (!found.ok) return found;
479 const { row } = found.value;
480 const before = definitionOf(row);
481 const servers = before.abilities.mcp_servers ?? [];
482 const server = servers.find((s) => s.id === a.id);
483 if (!server) return fail("not_found", "There is no such server on this agent.");
484 const checked = applyChanges(before, {}, TEMPLATE_IDS, { builtin: !!row.builtin, mcp_servers: servers.filter((s) => s.id !== server.id) });
485 if (!checked.ok) return fail("invalid", checked.message);
486 const saved = await this.saveVersion(row, checked.value, a.viewer!, a.workspace, `Removed the MCP server ${server.name} from @${row.handle}`);
487 return saved.ok ? ok(null) : saved;
488 }
489
490 /** Lists a server's tools again; a new version when they changed. Owners only. */
491 async refreshMcpServer(a: { workspace: string; handle: string; viewer: User | null; id: unknown }): Promise<Result<McpServer>> {
492 const found = await this.forMcp(a);
493 if (!found.ok) return found;
494 const { row } = found.value;
495 const before = definitionOf(row);
496 const servers = before.abilities.mcp_servers ?? [];
497 const server = servers.find((s) => s.id === a.id);
498 if (!server) return fail("not_found", "There is no such server on this agent.");
499 const listed = await listMcpTools(server.url);
500 const now = new Date().toISOString();
501 const fresh: McpServer = { ...server, tools: listed.ok ? listed.tools : server.tools, checked_at: now, problem: listed.ok ? null : listed.message };
502 if (JSON.stringify(fresh.tools) === JSON.stringify(server.tools) && fresh.problem === server.problem) return ok(fresh);
503 const checked = applyChanges(before, {}, TEMPLATE_IDS, { builtin: !!row.builtin, mcp_servers: servers.map((s) => (s.id === server.id ? fresh : s)) });
504 if (!checked.ok) return fail("invalid", checked.message);
505 const saved = await this.saveVersion(row, checked.value, a.viewer!, a.workspace, `Listed the MCP server ${server.name}'s tools again for @${row.handle} (${fresh.tools.length} tools)`);
506 return saved.ok ? ok(fresh) : saved;
507 }
508
509 async promote(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<WorkspaceAgent>> {
510 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
511 if (!managed.ok) return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners promote personal agents.") : managed;
512 const row = await this.row(managed.value, a.handle);
513 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
514 if (!isPersonal(row)) return fail("invalid", `@${row.handle} is already a workspace agent.`);
515 const versions = await this.db
516 .prepare("SELECT version, definition, changed_by, created_at FROM agent_versions WHERE agent_id = ? ORDER BY version")
517 .bind(row.id)
518 .all<{ version: number; definition: string; changed_by: string; created_at: string }>();
519 const definition = definitionOf(row);
520 const now = new Date().toISOString();
521 const id = newId("agt");
522 const version = row.version + 1;
523 const by = a.viewer!.username;
524 try {
525 const [archived] = await this.db.batch([
526 // Archived first, from the version read, so the handle is free for the new one in the same step.
527 this.db.prepare("UPDATE agents SET archived_at = ?, updated_at = ? WHERE id = ? AND version = ? AND archived_at IS NULL").bind(now, now, row.id, row.version),
528 insertAgent(this.db, id, managed.value, definition, { version, created_by: row.created_by, created_at: now, updated_at: now }),
529 ...versions.results.map((v) =>
530 this.db.prepare("INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at) VALUES (?, ?, ?, ?, ?)").bind(id, v.version, v.definition, v.changed_by, v.created_at),
531 ),
532 this.versionStatement(id, version, definition, by, now),
533 ]);
534 if (!archived.meta.changes) throw new Error("changed meanwhile");
535 } catch (error) {
536 console.error("agents: a promotion failed", row.id, String(error));
537 return fail("conflict", `@${row.handle} was changed meanwhile. Reload it and try again.`);
538 }
539 this.audit(a.viewer!, a.workspace, "promote_agent", row.handle, `Promoted @${row.owner_username ?? "a member"}'s personal agent @${row.handle} to a workspace agent (version ${version})`);
540 const saved = await this.row(managed.value, row.handle);
541 return ok(toAgent(saved!, new Date()));
542 }
543
544 // ── The Marketplace's install requests (./installs.ts) ─────────────────
545
546 /** Requests as the viewer sees them: every one for an owner, their own for anyone else. */
547 async installRequests(a: { workspace: string; viewer: User | null }): Promise<Result<InstallRequests>> {
548 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
549 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
550 if (!seen.ok) return seen;
551 const owner = canManage(a.viewer, a.workspace);
552 const requests = await listRequests(this.db, seen.value, a.viewer!, owner);
553 return ok({ requests, can_resolve: owner });
554 }
555
556 /** A member asks the owners to add something; each owner is notified. */
557 async requestInstall(a: { workspace: string; viewer: User | null; listing: unknown; note?: unknown }): Promise<Result<InstallRequest>> {
558 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
559 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
560 if (!seen.ok) return seen;
561 const viewer = a.viewer!;
562 if ((viewer.kind ?? "user") !== "user") return fail("forbidden", "Only people ask the workspace's owners to add things.");
563 if (canManage(viewer, a.workspace)) return fail("invalid", "You're an owner of this workspace: add it yourself.");
564 const listing = findListing(a.listing);
565 if (!listing) return fail("not_found", "That isn't something a workspace can add yet.");
566 const opened = await openRequest(this.db, seen.value, newId("ins"), listing, viewer, cleanRequestNote(a.note));
567 if (!opened.ok) return opened;
568 const slug = a.workspace.toLowerCase();
569 this.audit(viewer, slug, "request_install", listing.ref, `Asked the owners to add ${listing.name}`, "marketplace");
570 this.defer(this.tellOwners(slug, viewer, opened.value));
571 return opened;
572 }
573
574 /** An owner adds or turns down a request; whoever asked is told. */
575 async resolveInstallRequest(a: { workspace: string; viewer: User | null; id: unknown; status: unknown }): Promise<Result<InstallRequest>> {
576 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
577 if (!managed.ok) {
578 return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners answer requests.") : managed;
579 }
580 if (a.status !== "done" && a.status !== "declined") return fail("invalid", "Answer a request with done or declined.");
581 const answered = await resolveRequest(this.db, managed.value, String(a.id ?? ""), a.status, a.viewer!);
582 if (!answered.ok) return answered;
583 const slug = a.workspace.toLowerCase();
584 const { request } = answered.value;
585 this.audit(a.viewer!, slug, "answer_install_request", request.listing, `${request.status === "done" ? "Added" : "Turned down"} ${request.name} for @${request.requested_by}`, "marketplace");
586 this.defer(this.tellRequester(slug, a.viewer!, answered.value));
587 return ok(request);
588 }
589
590 // ── Extensions installed in the workspace (./extensions.ts) ──────────────
591
592 /** The workspace's installs; any member sees them. */
593 async extensionInstalls(a: { workspace: string; viewer: User | null }): Promise<Result<ExtensionInstall[]>> {
594 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
595 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
596 if (!seen.ok) return seen;
597 return ok(await extensions.listInstalls(this.db, seen.value));
598 }
599
600 /** Owners install a published extension; whoever asked for it is told. */
601 async installExtension(a: { workspace: string; viewer: User | null; extension: unknown }): Promise<Result<ExtensionInstall>> {
602 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
603 if (!managed.ok) return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners install extensions.") : managed;
604 const installed = await extensions.install(this.db, extensionById, managed.value, newId("ins"), String(a.extension ?? ""), a.viewer!.username);
605 if (!installed.ok) return installed;
606 const slug = a.workspace.toLowerCase();
607 this.audit(a.viewer!, slug, "install_extension", installed.value.listing, `Installed ${installed.value.listing} ${installed.value.version}`, "marketplace");
608 this.defer(this.answerListing(slug, managed.value, installed.value.listing, a.viewer!));
609 return installed;
610 }
611
612 /** Owners switch an install on or off. */
613 async setExtensionEnabled(a: { workspace: string; viewer: User | null; listing: unknown; enabled: unknown }): Promise<Result<ExtensionInstall>> {
614 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
615 if (!managed.ok) return managed;
616 const listing = String(a.listing ?? "");
617 if (!extensions.extensionIdOf(listing)) return fail("invalid", "Name an extension as extension:<id>.");
618 const changed = await extensions.setEnabled(this.db, managed.value, listing, a.enabled === true, a.viewer!.username);
619 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");
620 return changed;
621 }
622
623 /** Owners cap what an install spends a month. */
624 async setExtensionBudget(a: { workspace: string; viewer: User | null; listing: unknown; monthly_micros: unknown }): Promise<Result<ExtensionInstall>> {
625 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
626 if (!managed.ok) return managed;
627 const listing = String(a.listing ?? "");
628 const micros = a.monthly_micros == null ? null : Number(a.monthly_micros);
629 const changed = await extensions.setBudget(this.db, managed.value, listing, micros);
630 if (changed.ok) this.audit(a.viewer!, a.workspace, "set_extension_budget", listing, `Set ${listing}'s monthly budget`, "marketplace");
631 return changed;
632 }
633
634 /** Owners remove an install. */
635 async uninstallExtension(a: { workspace: string; viewer: User | null; listing: unknown }): Promise<Result<null>> {
636 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
637 if (!managed.ok) return managed;
638 const listing = String(a.listing ?? "");
639 const removed = await extensions.uninstall(this.db, managed.value, listing, a.viewer!.username);
640 if (removed.ok) this.audit(a.viewer!, a.workspace, "uninstall_extension", listing, `Uninstalled ${listing}`, "marketplace");
641 return removed;
642 }
643
644 /** Marks every open request for `listing` added, and tells each person who asked. */
645 private async answerListing(slug: string, workspaceId: string, listing: string, by: User): Promise<void> {
646 try {
647 const answered = await resolveListing(this.db, workspaceId, listing, by as Person);
648 await Promise.all(answered.map((one) => this.tellRequester(slug.toLowerCase(), by, one)));
649 } catch (error) {
650 console.error("agents: requests for a listing were not answered", listing, String(error));
651 }
652 }
653
654 /** Every owner hears of a new request (at most 20 of them). */
655 private async tellOwners(slug: string, asker: User, request: InstallRequest): Promise<void> {
656 if (!this.env.NOTIFY) return;
657 try {
658 const members = await identityClient(this.env.IDENTITY).listMembers(slug, asker);
659 if (!members.ok) throw new Error(members.error.message);
660 const owners = members.value.filter((member) => member.role === "owner").slice(0, 20);
661 const notify = notifyClient(this.env.NOTIFY);
662 await Promise.all(
663 owners.map((owner) =>
664 notify
665 .notify(
666 { username: owner.username },
667 {
668 id: `install-request:${request.id}:${owner.username}`,
669 kind: "approval",
670 workspace: slug,
671 title: `@${asker.username} asks you to add ${request.name}`,
672 body: request.note ?? (request.kind === "extension" ? "An extension. Install it, or turn the request down." : "An integration. Connect it, or turn the request down."),
673 href: requestsPath(slug),
674 actor: { kind: "user", id: asker.id, name: asker.username, avatar: asker.avatar ?? null, avatar_seed: null },
675 created_at: new Date().toISOString(),
676 },
677 )
678 .catch(() => undefined),
679 ),
680 );
681 } catch (error) {
682 console.error("agents: owners were not told of an install request", request.id, String(error));
683 }
684 }
685
686 /** The person who asked hears their request was answered. */
687 private async tellRequester(slug: string, by: User, answered: Answered): Promise<void> {
688 if (!this.env.NOTIFY) return;
689 const { request } = answered;
690 const line = answerLine(request, `@${by.username}`);
691 await notifyClient(this.env.NOTIFY)
692 .notify(
693 { user_id: answered.requested_by_id },
694 {
695 id: `install-answer:${request.id}`,
696 kind: "inbox",
697 workspace: slug,
698 title: line.title,
699 body: line.body,
700 href: request.status === "done" ? listingPath(slug, request.listing) : requestsPath(slug),
701 actor: { kind: "user", id: by.id, name: by.username, avatar: by.avatar ?? null, avatar_seed: null },
702 created_at: new Date().toISOString(),
703 },
704 )
705 .catch((error: unknown) => console.error("agents: a requester was not told", request.id, String(error)));
706 }
707
708 async update(a: { workspace: string; handle: string; viewer: User | null; changes: Partial<NewWorkspaceAgent> }): Promise<Result<WorkspaceAgent>> {
709 const found = await this.visible(a.workspace, a.viewer, a.handle);
710 if (!found.ok) return found;
711 const { row } = found.value;
712 // Owners change the workspace's agents; a member changes their own personal one.
713 if (!canChange(a.viewer, a.workspace, row)) return fail("forbidden", isPersonal(row) ? NOT_YOURS_REFUSAL : MANAGE_REFUSAL);
714 const workspaceId = found.value.workspaceId;
715 const before = definitionOf(row);
716 const { scope: _scope, ...asked } = (a.changes ?? {}) as Partial<NewWorkspaceAgent>;
717 // The built-in @g1t keeps who it is and its job; the rest is the workspace's.
718 const allowed = row.builtin ? builtinChanges(before, asked) : { ok: true as const, value: asked };
719 if (!allowed.ok) return fail("invalid", allowed.message);
720 const checked = applyChanges(before, allowed.value, TEMPLATE_IDS, { builtin: !!row.builtin });
721 if (!checked.ok) return fail("invalid", checked.message);
722 const definition = checked.value;
723 if (JSON.stringify(definition) === JSON.stringify(before)) return ok(toAgent(row, new Date()));
724 if (definition.handle !== before.handle && (await this.handleTaken(workspaceId, definition.handle, row.id))) {
725 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
726 }
727 const now = new Date().toISOString();
728 const version = row.version + 1;
729 const by = a.viewer!.username;
730 try {
731 const [updated] = await this.db.batch([
732 // Only from the version read: two owners saving at once never lose
733 // one's change silently; the second is told to look again.
734 updateAgent(this.db, row.id, row.version, definition, { version, updated_at: now }),
735 this.db
736 .prepare(
737 `INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at)
738 SELECT ?1, ?2, ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM agents WHERE id = ?1 AND version = ?2 AND updated_at = ?5)`,
739 )
740 .bind(row.id, version, JSON.stringify(definition), by, now),
741 ]);
742 if (!updated.meta.changes) return fail("conflict", `@${before.handle} was changed meanwhile. Reload it and try again.`);
743 } catch (error) {
744 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
745 throw error;
746 }
747 const renamed = definition.handle !== before.handle ? ` (was @${before.handle})` : "";
748 this.audit(a.viewer!, a.workspace, "update_agent", definition.handle, `Changed @${definition.handle}${renamed} to version ${version}`, "agents", isPersonal(row) ? "member" : "owner");
749 const saved = await this.row(workspaceId, definition.handle);
750 return ok(toAgent(saved!, new Date()));
751 }
752
753 /** What each effort level has cost one agent, from its own finished sessions. */
754 async effortCosts(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<AgentEffortCosts>> {
755 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
756 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
757 if (!seen.ok) return seen;
758 const row = await this.row(seen.value, a.handle);
759 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
760 const outcomes = await outcomesSince(this.db, seen.value, sinceWindow(), row.id);
761 return ok(effortCostsOf(row.handle, effortOf(definitionOf(row).routing), outcomes.get(row.id) ?? []));
762 }
763
764 /** Ways to spend less, checked against past work: every agent's, or one's. */
765 async recommendations(a: { workspace: string; viewer: User | null; handle?: string | null }): Promise<Result<AgentRecommendations>> {
766 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
767 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
768 if (!seen.ok) return seen;
769 let agentId: string | null = null;
770 if (a.handle) {
771 const row = await this.row(seen.value, a.handle);
772 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
773 agentId = row.id;
774 }
775 return ok(await readRecommendations(this.db, seen.value, agentId));
776 }
777
778 /**
779 * Owners apply a suggestion (the agent's effort changes as a new version
780 * of it, and the audit log says so) or dismiss it. Only an open one, and
781 * only while the agent's setting is still the one it was made for.
782 */
783 async resolveRecommendation(a: { workspace: string; viewer: User | null; id: string; action: string }): Promise<Result<AgentRecommendation>> {
784 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
785 if (!managed.ok) return managed;
786 if (a.action !== "apply" && a.action !== "dismiss") return fail("invalid", "Apply or dismiss.");
787 const found = await recommendationRow(this.db, managed.value, String(a.id ?? ""));
788 if (!found) return fail("not_found", "There is no such suggestion.");
789 if (found.status !== "open") return fail("conflict", "This suggestion was already decided, or is no longer current.");
790 const by = a.viewer!.username;
791 if (a.action === "dismiss") {
792 if (!(await markResolved(this.db, found.id, "dismissed", by))) return fail("conflict", "This suggestion was decided meanwhile.");
793 this.audit(a.viewer!, a.workspace, "dismiss_recommendation", found.handle ?? found.agent_id, `Dismissed "${found.title}"`);
794 } else {
795 const agent = await this.row(managed.value, found.handle);
796 if (!agent || agent.id !== found.agent_id) return fail("not_found", "That agent is gone.");
797 if (effortOf(definitionOf(agent).routing) !== found.from_effort) {
798 await this.db.prepare("UPDATE agent_recommendations SET status = 'stale' WHERE id = ? AND status = 'open'").bind(found.id).run();
799 return fail("conflict", `@${agent.handle}'s effort was changed since this was suggested. The next check looks again.`);
800 }
801 // Claimed first, so two owners pressing Apply change the agent once.
802 if (!(await markResolved(this.db, found.id, "applied", by))) return fail("conflict", "This suggestion was decided meanwhile.");
803 const changed = await this.update({ workspace: a.workspace, handle: agent.handle, viewer: a.viewer, changes: { routing: { effort: found.to_effort as never } } });
804 if (!changed.ok) {
805 await this.db.prepare("UPDATE agent_recommendations SET status = 'open', resolved_by = NULL, resolved_at = NULL WHERE id = ?").bind(found.id).run();
806 return changed;
807 }
808 this.audit(
809 a.viewer!,
810 a.workspace,
811 "apply_recommendation",
812 agent.handle,
813 `Applied "${found.title}": @${agent.handle}'s effort from ${EFFORT_NAMES[found.from_effort as keyof typeof EFFORT_NAMES] ?? found.from_effort} to ${EFFORT_NAMES[found.to_effort as keyof typeof EFFORT_NAMES] ?? found.to_effort}`,
814 );
815 }
816 const after = await recommendationRow(this.db, managed.value, found.id);
817 return ok(toRecommendation(after!));
818 }
819
820 async archive(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<null>> {
821 const found = await this.visible(a.workspace, a.viewer, a.handle);
822 if (!found.ok) return found;
823 const { row } = found.value;
824 // Owners archive any agent; a member their own personal one.
825 if (!canArchive(a.viewer, a.workspace, row)) return fail("forbidden", MANAGE_REFUSAL);
826 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.");
827 const now = new Date().toISOString();
828 await this.db.prepare("UPDATE agents SET archived_at = ?, updated_at = ? WHERE id = ? AND archived_at IS NULL").bind(now, now, row.id).run();
829 this.audit(a.viewer!, a.workspace, "archive_agent", row.handle, `Archived @${row.handle}`);
830 // Its computer goes to sleep now; its disk is kept 30 days, then deleted (computer.ts).
831 this.defer(forgetComputer(this.env, row.id));
832 return ok(null);
833 }
834
835 /**
836 * Makes the workspace's built-in @g1t if it does not exist yet: an
837 * ordinary agent row, marked builtin, at version 1. Safe to call on
838 * every request; once seen, an isolate does not write again.
839 */
840 private async ensureBuiltin(workspaceId: string): Promise<void> {
841 await ensureBuiltin(this.db, workspaceId);
842 }
843
844 /** Internal: the workspace's @g1t, made if need be, for the chat service. */
845 async builtin(a: { workspace: string; workspace_id: string }): Promise<Result<WorkspaceAgent>> {
846 if (typeof a?.workspace_id !== "string" || !a.workspace_id) return fail("invalid", "Name the workspace by id.");
847 await this.ensureBuiltin(a.workspace_id);
848 const now = new Date();
849 const row = await this.db
850 .prepare(selectAgents("a.workspace_id = ?3 AND a.builtin = 1"))
851 .bind(...periods(now), a.workspace_id)
852 .first<Row>();
853 return row ? ok(toAgent(row, now)) : fail("not_found", "This workspace's @g1t could not be made.");
854 }
855
856 /**
857 * Hands a message to the agent's desk, which answers it in the
858 * background. Returns as soon as the desk holds it.
859 */
860 async deliver(delivery: AgentDelivery): Promise<Result<null>> {
861 const fields = ["workspace", "workspace_id", "channel_id", "agent_id", "message_id", "asked_by"] as const;
862 if (!delivery || fields.some((field) => typeof delivery[field] !== "string" || !delivery[field])) {
863 return fail("invalid", "A delivery names the workspace, channel, agent, message and who asked.");
864 }
865 if (delivery.channel_kind !== "channel" && delivery.channel_kind !== "dm") return fail("invalid", "channel_kind is channel or dm.");
866 await this.ensureBuiltin(delivery.workspace_id);
867 const desk = this.env.DESKS.get(this.env.DESKS.idFromName(delivery.agent_id));
868 await desk.take({ ...delivery, hops: Math.max(0, Math.floor(Number(delivery.hops) || 0)), thread_root: delivery.thread_root ?? null });
869 return ok(null);
870 }
871
872 private versionStatement(agentId: string, version: number, d: Definition, by: string, at: string): D1PreparedStatement {
873 return versionStatement(this.db, agentId, version, d, by, at);
874 }
875
876 /**
877 * Records a change to an agent in the workspace's audit log, after the
878 * answer. Never fails the change: a log that cannot be written is logged.
879 * The events service's audit contract speaks camelCase (Rust's
880 * `NewAuditEntry`).
881 */
882 /** A change to the skill library, in the audit log as `agents/skills/<name>`. */
883 auditSkill(actor: User, workspace: string, action: string, name: string, message: string): void {
884 this.audit(actor, workspace, action, `skills/${name}`, message);
885 }
886
887 private audit(
888 actor: User,
889 workspace: string,
890 action: string,
891 handle: string,
892 message: string,
893 area: "agents" | "marketplace" = "agents",
894 rule: "owner" | "member" | null = null,
895 ): void {
896 const kind = actor.kind === "workspace" ? "workspace" : actor.kind === "agent" ? "agent" : actor.kind === "system" ? "system" : "person";
897 const entry = {
898 actorKind: kind,
899 actor: actor.username,
900 actorId: actor.id,
901 agent: null,
902 onBehalfOf: null,
903 runId: null,
904 runKind: null,
905 credentialId: null,
906 action,
907 surface: "web",
908 workspace: workspace.toLowerCase(),
909 repo: null,
910 number: null,
911 gitRef: null,
912 path: `${area}/${handle}`,
913 outcome: "allowed",
914 rule: rule ?? (area === "marketplace" && action === "request_install" ? "member" : "owner"),
915 result: "ok",
916 message,
917 requestId: `req_${crypto.randomUUID()}`,
918 };
919 this.defer(
920 this.env.EVENTS.fetch("https://service/rpc/audit_record", {
921 method: "POST",
922 headers: { "content-type": "application/json" },
923 body: JSON.stringify({ entries: [entry] }),
924 })
925 .then((response) => {
926 if (!response.ok) throw new Error(`status ${response.status}`);
927 })
928 .catch((error: unknown) => console.error("agents: audit entry not recorded", action, handle, String(error))),
929 );
930 }
931}
932
933/** One RPC method's answer. */
934async function answer(service: Agents, method: string, args: any): Promise<Response> {
935 // The skill library (./skill-rpc.ts).
936 const skill = await skillRpc(method, args, (a, run) => service.view(a, run), (viewer, workspace, action, name, message) => service.auditSkill(viewer, workspace, action, name, message));
937 if (skill) return Response.json(skill);
938 switch (method) {
939 case "list":
940 return Response.json(await service.list(args));
941 case "get":
942 return Response.json(await service.get(args));
943 case "by_ids":
944 return Response.json(await service.byIds(args));
945 case "create":
946 return Response.json(await service.create(args));
947 case "update":
948 return Response.json(await service.update(args));
949 case "archive":
950 return Response.json(await service.archive(args));
951 case "draft":
952 return Response.json(await service.draft(args));
953 case "try_draft":
954 return Response.json(await service.tryDraft(args));
955 case "redraft":
956 return Response.json(await service.redraft(args));
957 case "add_mcp_server":
958 return Response.json(await service.addMcpServer(args));
959 case "remove_mcp_server":
960 return Response.json(await service.removeMcpServer(args));
961 case "refresh_mcp_server":
962 return Response.json(await service.refreshMcpServer(args));
963 case "promote":
964 return Response.json(await service.promote(args));
965 case "builtin":
966 return Response.json(await service.builtin(args));
967 case "templates":
968 return Response.json(TEMPLATES);
969 case "deliver":
970 return Response.json(await service.deliver(args));
971 case "overview":
972 return Response.json(await service.view(args, (ctx) => views.overview(ctx)));
973 case "sessions":
974 return Response.json(await service.view(args, (ctx) => views.listSessions(ctx, args)));
975 case "session":
976 return Response.json(await service.view(args, (ctx) => views.sessionDetail(ctx, args.id)));
977 case "stop_session":
978 return Response.json(await service.view(args, (ctx) => views.stopSession(ctx, args.id)));
979 case "approve_session":
980 return Response.json(await service.view(args, (ctx) => views.approveSession(ctx, args.id, args.cap_micros)));
981 case "steer_session":
982 return Response.json(await service.view(args, (ctx) => views.steerSession(ctx, args.id, args.body)));
983 // The agent's own computer (./computer.ts).
984 case "computer":
985 return Response.json(await service.view(args, (ctx) => computerView(ctx, args.handle)));
986 case "computer_wake":
987 return Response.json(await service.view(args, (ctx) => wakeComputer(ctx, args.handle)));
988 case "computer_sleep":
989 return Response.json(await service.view(args, (ctx) => sleepComputer(ctx, args.handle)));
990 case "computer_reset":
991 return Response.json(await service.view(args, (ctx) => resetComputer(ctx, args.handle)));
992 case "computer_commands":
993 return Response.json(await service.view(args, (ctx) => computerCommands(ctx, args.handle, args.session_id ?? null)));
994 case "memories":
995 return Response.json(await service.view(args, (ctx) => views.memories(ctx, args.handle)));
996 case "remember":
997 return Response.json(await service.view(args, (ctx) => views.remember(ctx, args.handle, args.input)));
998 case "update_memory":
999 return Response.json(await service.view(args, (ctx) => views.updateMemory(ctx, args.handle, args.id, args.changes)));
1000 case "forget":
1001 return Response.json(await service.view(args, (ctx) => views.forget(ctx, args.handle, args.id)));
1002 case "routines":
1003 return Response.json(await service.view(args, (ctx) => views.routines(ctx, args.handle)));
1004 case "save_routine":
1005 return Response.json(await service.view(args, (ctx) => views.saveRoutine(ctx, args.handle, args.input, args.id ?? null)));
1006 case "delete_routine":
1007 return Response.json(await service.view(args, (ctx) => views.deleteRoutine(ctx, args.handle, args.id)));
1008 case "run_routine":
1009 return Response.json(await service.view(args, (ctx) => views.runRoutineNow(ctx, args.handle, args.id)));
1010 case "spend":
1011 return Response.json(await service.view(args, (ctx) => views.spend(ctx, args.handle ?? null, { period: args.period, person: args.person })));
1012 case "person_budgets":
1013 return Response.json(await service.view(args, (ctx) => views.personBudgetsView(ctx)));
1014 case "set_person_budget":
1015 return Response.json(await service.view(args, (ctx) => views.setPersonBudget(ctx, args.username, args.monthly_micros)));
1016 case "team_context":
1017 return Response.json(await service.view(args, (ctx) => views.teamContext(ctx, args.handle)));
1018 case "activity":
1019 return Response.json(await service.view(args, (ctx) => views.activity(ctx, args.handle)));
1020 case "versions":
1021 return Response.json(await service.view(args, (ctx) => views.versions(ctx, args.handle)));
1022 case "card_action":
1023 return Response.json(await service.cardAction(args));
1024 case "policy":
1025 return Response.json(await service.view(args, (ctx) => views.policy(ctx)));
1026 case "install_requests":
1027 return Response.json(await service.installRequests(args));
1028 case "request_install":
1029 return Response.json(await service.requestInstall(args));
1030 case "resolve_install_request":
1031 return Response.json(await service.resolveInstallRequest(args));
1032 case "extension_installs":
1033 return Response.json(await service.extensionInstalls(args));
1034 case "install_extension":
1035 return Response.json(await service.installExtension(args));
1036 case "set_extension_enabled":
1037 return Response.json(await service.setExtensionEnabled(args));
1038 case "set_extension_budget":
1039 return Response.json(await service.setExtensionBudget(args));
1040 case "uninstall_extension":
1041 return Response.json(await service.uninstallExtension(args));
1042 case "effort_costs":
1043 return Response.json(await service.effortCosts(args));
1044 case "recommendations":
1045 return Response.json(await service.recommendations(args));
1046 case "resolve_recommendation":
1047 return Response.json(await service.resolveRecommendation(args));
1048 case "set_policy":
1049 return Response.json(await service.view(args, (ctx) => views.setPolicy(ctx, args.policy)));
1050 default:
1051 return new Response("Unknown method\n", { status: 404 });
1052 }
1053}
1054
1055export default {
1056 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
1057 const match = new URL(request.url).pathname.match(/^\/rpc\/([a-z_]+)$/);
1058 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
1059 // A replica near the caller when it asks for one (@g1t/contracts d1.ts).
1060 const opened = openD1(env.DB, request);
1061 const service = new Agents(Object.create(env, { DB: { value: opened.db } }) as Env, (work) => ctx.waitUntil(work));
1062 const args = (await request.json().catch(() => ({}))) as any;
1063 return opened.finish(await answer(service, match[1], args));
1064 },
1065
1066 /** Events routines run on, from the events service (SUBSCRIBER_AGENTS). */
1067 async queue(batch: MessageBatch<unknown>, env: Env): Promise<void> {
1068 const events = batch.messages.map((message) => message.body as G1tEvent);
1069 // A push to a default branch: skill libraries that follow that repository read it again (./skill-library.ts).
1070 const pushes = events.flatMap((event) => (event?.type === "git.push" && event.data?.defaultBranch && event.data.repoId ? [{ repoId: event.data.repoId }] : []));
1071 await Promise.all([onEvents(env as unknown as SessionEnv, events.filter((event) => event?.type !== "git.push")), pushes.length ? skillPushes(env as unknown as SessionEnv, pushes) : null]);
1072 batch.ackAll();
1073 },
1074
1075 /** Every few minutes: routines that are due, session steps a desk lost, the spend check, and the one-time team move. */
1076 async scheduled(_controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
1077 const sessions = env as unknown as SessionEnv;
1078 ctx.waitUntil(
1079 Promise.all([
1080 runDue(sessions).catch((error: unknown) => console.error("agents: routines did not run", String(error))),
1081 sweep(sessions).catch((error: unknown) => console.error("agents: the session sweep failed", String(error))),
1082 // Spend less, keep quality: each workspace checked weekly (src/recommend.ts).
1083 checkDue(env.DB).catch((error: unknown) => console.error("agents: the spend check failed", String(error))),
1084 // Once: agents' own old teams become team memberships (src/team-move.ts); after that, one cheap read.
1085 moveAgentTeams(env.DB, (id, agents) => identityClient(env.IDENTITY).adoptAgentTeams(id, agents)).catch((error: unknown) => console.error("agents: moving agents' teams failed", String(error))),
1086 ]),
1087 );
1088 },
1089} satisfies ExportedHandler<Env>;