Skip to content
1,122 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 before it says hello, so it knows them.
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 this.defer(this.hello(a.workspace, workspaceId, id, viewer));
320 const row = await this.row(workspaceId, definition.handle);
321 return ok(toAgent(row!, new Date()));
322 }
323
324 // ── The builder: describe it, try it, change it in words (./builder.ts) ──
325
326 /** Drafts a whole agent from a description, charged to the viewer. Nothing is saved. */
327 async draft(a: { workspace: string; viewer: User | null; description: unknown; scope?: unknown }): Promise<Result<AgentProposal>> {
328 const allowed = await this.creatable(a?.workspace ?? "", a?.viewer ?? null, a?.scope);
329 if (!allowed.ok) return allowed;
330 const words = wordsOf(a.description, "description");
331 if (!words.ok) return fail("invalid", words.message);
332 const { workspaceId, scope, policy } = allowed.value;
333 const slug = a.workspace.toLowerCase();
334 const [answer, handles] = await Promise.all([
335 callBuilder(this.env, {
336 slug,
337 workspaceId,
338 viewer: a.viewer!,
339 system: draftSystem(slug, scope),
340 messages: [{ role: "user", content: `What the agent should do:\n\n${words.value}` }],
341 maxOutput: DRAFT_OUTPUT_TOKENS,
342 }),
343 this.db.prepare("SELECT handle FROM agents WHERE workspace_id = ? AND archived_at IS NULL").bind(workspaceId).all<{ handle: string }>(),
344 ]);
345 if (!answer.ok) return fail(answer.code, answer.message);
346 const proposal = proposalFrom(jsonIn(answer.text), {
347 scope,
348 taken: new Set(handles.results.map((row) => row.handle)),
349 budget: startingBudget(scope, policy.default_agent_monthly_micros),
350 });
351 if (!proposal.ok) return fail("invalid", proposal.message);
352 return ok({ ...proposal.value, charged_micros: answer.charged });
353 }
354
355 /** Try it: the unsaved definition answers the conversation so far. Charged to the viewer; nothing is saved. */
356 async tryDraft(a: { workspace: string; viewer: User | null; definition: unknown; messages: unknown }): Promise<Result<DraftReply>> {
357 const given = (a?.definition ?? {}) as NewWorkspaceAgent;
358 const allowed = await this.creatable(a?.workspace ?? "", a?.viewer ?? null, given?.scope);
359 if (!allowed.ok) return allowed;
360 const { scope: _scope, ...rest } = given;
361 // A handle someone has already taken is no reason not to try it.
362 const checked = applyChanges(null, { ...rest, handle: rest.handle || "preview" }, TEMPLATE_IDS);
363 if (!checked.ok) return fail("invalid", checked.message);
364 const turns = tryTurns(a.messages);
365 if (!turns.ok) return fail("invalid", turns.message);
366 const viewer = a.viewer!;
367 const answer = await callBuilder(this.env, {
368 slug: a.workspace.toLowerCase(),
369 workspaceId: allowed.value.workspaceId,
370 viewer,
371 system: trySystem({ workspace: a.workspace.toLowerCase(), definition: checked.value, asker: { username: viewer.username, display_name: viewer.display_username ?? null } }),
372 messages: turns.value,
373 routing: { floor: checked.value.routing.floor, ceiling: checked.value.routing.ceiling },
374 });
375 if (!answer.ok) return fail(answer.code, answer.message);
376 return ok({ text: answer.text, charged_micros: answer.charged });
377 }
378
379 /** Drafts changes from a request in words, for whoever may change the agent. Saved only through `update`. */
380 async redraft(a: { workspace: string; handle: string; viewer: User | null; request: unknown }): Promise<Result<AgentRedraft>> {
381 const found = await this.visible(a?.workspace ?? "", a?.viewer ?? null, a?.handle);
382 if (!found.ok) return found;
383 const { workspaceId, row } = found.value;
384 if (!canChange(a.viewer, a.workspace, row)) return fail("forbidden", isPersonal(row) ? NOT_YOURS_REFUSAL : MANAGE_REFUSAL);
385 const words = wordsOf(a.request, "request");
386 if (!words.ok) return fail("invalid", words.message);
387 const before = definitionOf(row);
388 const slug = a.workspace.toLowerCase();
389 const answer = await callBuilder(this.env, {
390 slug,
391 workspaceId,
392 viewer: a.viewer!,
393 system: redraftSystem(slug, before, !!row.builtin),
394 messages: [{ role: "user", content: `The change to make to @${row.handle}:\n\n${words.value}` }],
395 maxOutput: DRAFT_OUTPUT_TOKENS,
396 });
397 if (!answer.ok) return fail(answer.code, answer.message);
398 const drafted = redraftFrom(jsonIn(answer.text), before, !!row.builtin);
399 if (!drafted.ok) return fail("invalid", drafted.message);
400 return ok({ ...drafted.value, from_version: row.version, charged_micros: answer.charged });
401 }
402
403 /**
404 * Owners make a personal agent a workspace agent: a new workspace agent
405 * with the same handle, definition and every version, and one more
406 * version saying so; the personal one is archived in the same step, its
407 * memory and direct messages kept with it (owners can't read a member's
408 * memory to choose from it, so none is carried).
409 */
410 // ── MCP servers (docs.g1t.sh/guides/agent-abilities/, "MCP servers") ─────
411
412 /** The agent, if the viewer may add servers to it: owners only, a workspace agent or a personal one alike. */
413 private async forMcp(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<{ row: Row; workspaceId: string }>> {
414 const found = await this.visible(a?.workspace ?? "", a?.viewer ?? null, a?.handle);
415 if (!found.ok) return found;
416 if (!canManage(a.viewer, a.workspace)) return fail("forbidden", "Only the workspace's owners add MCP servers to an agent.");
417 return ok({ row: found.value.row, workspaceId: found.value.workspaceId });
418 }
419
420 /** Saves `definition` as the agent's next version, from the version read; what the agent is now. */
421 private async saveVersion(row: Row, definition: Definition, by: User, workspace: string, what: string): Promise<Result<WorkspaceAgent>> {
422 const now = new Date().toISOString();
423 const version = row.version + 1;
424 const [updated] = await this.db.batch([
425 updateAgent(this.db, row.id, row.version, definition, { version, updated_at: now }),
426 this.db
427 .prepare(
428 `INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at)
429 SELECT ?1, ?2, ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM agents WHERE id = ?1 AND version = ?2 AND updated_at = ?5)`,
430 )
431 .bind(row.id, version, JSON.stringify(definition), by.username, now),
432 ]);
433 if (!updated.meta.changes) return fail("conflict", `@${row.handle} was changed meanwhile. Reload it and try again.`);
434 this.audit(by, workspace, "update_agent", row.handle, `${what} (version ${version})`, "agents", "owner");
435 const saved = await this.row(row.workspace_id, row.handle);
436 return ok(toAgent(saved!, new Date()));
437 }
438
439 /**
440 * Adds an MCP server to an agent: the name and address are checked, its
441 * tools are listed from it now, and the agent gets a new version with
442 * the server and its tools as abilities (each a write unless the server
443 * says it only reads). Owners only.
444 */
445 async addMcpServer(a: { workspace: string; handle: string; viewer: User | null; name: unknown; url: unknown }): Promise<Result<McpServer>> {
446 const found = await this.forMcp(a);
447 if (!found.ok) return found;
448 const { row } = found.value;
449 const name = checkMcpName(a.name);
450 if (!name.ok) return fail("invalid", name.message);
451 const url = checkMcpUrl(a.url);
452 if (!url.ok) return fail("invalid", url.message);
453 const before = definitionOf(row);
454 const servers = before.abilities.mcp_servers ?? [];
455 if (servers.length >= MAX_MCP_SERVERS) return fail("limit", `An agent has at most ${MAX_MCP_SERVERS} MCP servers.`);
456 if (servers.some((server) => server.name === name.value)) return fail("conflict", `@${row.handle} already has a server called ${name.value}.`);
457 if (servers.some((server) => server.url === url.value)) return fail("conflict", `@${row.handle} already has that server.`);
458 const listed = await listMcpTools(url.value);
459 const now = new Date().toISOString();
460 const server: McpServer = {
461 id: newId("mcp"),
462 name: name.value,
463 url: url.value,
464 tools: listed.ok ? listed.tools : [],
465 added_by: a.viewer!.username,
466 added_at: now,
467 checked_at: now,
468 problem: listed.ok ? null : listed.message,
469 };
470 const checked = applyChanges(before, {}, TEMPLATE_IDS, { builtin: !!row.builtin, mcp_servers: [...servers, server] });
471 if (!checked.ok) return fail("invalid", checked.message);
472 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}`);
473 return saved.ok ? ok(server) : saved;
474 }
475
476 /** Takes an MCP server off an agent, with its abilities' settings: a new version. Owners only. */
477 async removeMcpServer(a: { workspace: string; handle: string; viewer: User | null; id: unknown }): Promise<Result<null>> {
478 const found = await this.forMcp(a);
479 if (!found.ok) return found;
480 const { row } = found.value;
481 const before = definitionOf(row);
482 const servers = before.abilities.mcp_servers ?? [];
483 const server = servers.find((s) => s.id === a.id);
484 if (!server) return fail("not_found", "There is no such server on this agent.");
485 const checked = applyChanges(before, {}, TEMPLATE_IDS, { builtin: !!row.builtin, mcp_servers: servers.filter((s) => s.id !== server.id) });
486 if (!checked.ok) return fail("invalid", checked.message);
487 const saved = await this.saveVersion(row, checked.value, a.viewer!, a.workspace, `Removed the MCP server ${server.name} from @${row.handle}`);
488 return saved.ok ? ok(null) : saved;
489 }
490
491 /** Lists a server's tools again; a new version when they changed. Owners only. */
492 async refreshMcpServer(a: { workspace: string; handle: string; viewer: User | null; id: unknown }): Promise<Result<McpServer>> {
493 const found = await this.forMcp(a);
494 if (!found.ok) return found;
495 const { row } = found.value;
496 const before = definitionOf(row);
497 const servers = before.abilities.mcp_servers ?? [];
498 const server = servers.find((s) => s.id === a.id);
499 if (!server) return fail("not_found", "There is no such server on this agent.");
500 const listed = await listMcpTools(server.url);
501 const now = new Date().toISOString();
502 const fresh: McpServer = { ...server, tools: listed.ok ? listed.tools : server.tools, checked_at: now, problem: listed.ok ? null : listed.message };
503 if (JSON.stringify(fresh.tools) === JSON.stringify(server.tools) && fresh.problem === server.problem) return ok(fresh);
504 const checked = applyChanges(before, {}, TEMPLATE_IDS, { builtin: !!row.builtin, mcp_servers: servers.map((s) => (s.id === server.id ? fresh : s)) });
505 if (!checked.ok) return fail("invalid", checked.message);
506 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)`);
507 return saved.ok ? ok(fresh) : saved;
508 }
509
510 async promote(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<WorkspaceAgent>> {
511 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
512 if (!managed.ok) return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners promote personal agents.") : managed;
513 const row = await this.row(managed.value, a.handle);
514 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
515 if (!isPersonal(row)) return fail("invalid", `@${row.handle} is already a workspace agent.`);
516 const versions = await this.db
517 .prepare("SELECT version, definition, changed_by, created_at FROM agent_versions WHERE agent_id = ? ORDER BY version")
518 .bind(row.id)
519 .all<{ version: number; definition: string; changed_by: string; created_at: string }>();
520 const definition = definitionOf(row);
521 const now = new Date().toISOString();
522 const id = newId("agt");
523 const version = row.version + 1;
524 const by = a.viewer!.username;
525 try {
526 const [archived] = await this.db.batch([
527 // Archived first, from the version read, so the handle is free for the new one in the same step.
528 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),
529 insertAgent(this.db, id, managed.value, definition, { version, created_by: row.created_by, created_at: now, updated_at: now }),
530 ...versions.results.map((v) =>
531 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),
532 ),
533 this.versionStatement(id, version, definition, by, now),
534 ]);
535 if (!archived.meta.changes) throw new Error("changed meanwhile");
536 } catch (error) {
537 console.error("agents: a promotion failed", row.id, String(error));
538 return fail("conflict", `@${row.handle} was changed meanwhile. Reload it and try again.`);
539 }
540 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})`);
541 const saved = await this.row(managed.value, row.handle);
542 return ok(toAgent(saved!, new Date()));
543 }
544
545 // ── The Marketplace's install requests (./installs.ts) ─────────────────
546
547 /** Requests as the viewer sees them: every one for an owner, their own for anyone else. */
548 async installRequests(a: { workspace: string; viewer: User | null }): Promise<Result<InstallRequests>> {
549 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
550 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
551 if (!seen.ok) return seen;
552 const owner = canManage(a.viewer, a.workspace);
553 const requests = await listRequests(this.db, seen.value, a.viewer!, owner);
554 return ok({ requests, can_resolve: owner });
555 }
556
557 /** A member asks the owners to add something; each owner is notified. */
558 async requestInstall(a: { workspace: string; viewer: User | null; listing: unknown; note?: unknown }): Promise<Result<InstallRequest>> {
559 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
560 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
561 if (!seen.ok) return seen;
562 const viewer = a.viewer!;
563 if ((viewer.kind ?? "user") !== "user") return fail("forbidden", "Only people ask the workspace's owners to add things.");
564 if (canManage(viewer, a.workspace)) return fail("invalid", "You're an owner of this workspace: add it yourself.");
565 const listing = findListing(a.listing);
566 if (!listing) return fail("not_found", "That isn't something a workspace can add yet.");
567 const opened = await openRequest(this.db, seen.value, newId("ins"), listing, viewer, cleanRequestNote(a.note));
568 if (!opened.ok) return opened;
569 const slug = a.workspace.toLowerCase();
570 this.audit(viewer, slug, "request_install", listing.ref, `Asked the owners to add ${listing.name}`, "marketplace");
571 this.defer(this.tellOwners(slug, viewer, opened.value));
572 return opened;
573 }
574
575 /** An owner adds or turns down a request; whoever asked is told. */
576 async resolveInstallRequest(a: { workspace: string; viewer: User | null; id: unknown; status: unknown }): Promise<Result<InstallRequest>> {
577 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
578 if (!managed.ok) {
579 return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners answer requests.") : managed;
580 }
581 if (a.status !== "done" && a.status !== "declined") return fail("invalid", "Answer a request with done or declined.");
582 const answered = await resolveRequest(this.db, managed.value, String(a.id ?? ""), a.status, a.viewer!);
583 if (!answered.ok) return answered;
584 const slug = a.workspace.toLowerCase();
585 const { request } = answered.value;
586 this.audit(a.viewer!, slug, "answer_install_request", request.listing, `${request.status === "done" ? "Added" : "Turned down"} ${request.name} for @${request.requested_by}`, "marketplace");
587 this.defer(this.tellRequester(slug, a.viewer!, answered.value));
588 return ok(request);
589 }
590
591 // ── Extensions installed in the workspace (./extensions.ts) ──────────────
592
593 /** The workspace's installs; any member sees them. */
594 async extensionInstalls(a: { workspace: string; viewer: User | null }): Promise<Result<ExtensionInstall[]>> {
595 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
596 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
597 if (!seen.ok) return seen;
598 return ok(await extensions.listInstalls(this.db, seen.value));
599 }
600
601 /** Owners install a published extension; whoever asked for it is told. */
602 async installExtension(a: { workspace: string; viewer: User | null; extension: unknown }): Promise<Result<ExtensionInstall>> {
603 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
604 if (!managed.ok) return managed.error.code === "forbidden" ? fail("forbidden", "Only the workspace's owners install extensions.") : managed;
605 const installed = await extensions.install(this.db, extensionById, managed.value, newId("ins"), String(a.extension ?? ""), a.viewer!.username);
606 if (!installed.ok) return installed;
607 const slug = a.workspace.toLowerCase();
608 this.audit(a.viewer!, slug, "install_extension", installed.value.listing, `Installed ${installed.value.listing} ${installed.value.version}`, "marketplace");
609 this.defer(this.answerListing(slug, managed.value, installed.value.listing, a.viewer!));
610 return installed;
611 }
612
613 /** Owners switch an install on or off. */
614 async setExtensionEnabled(a: { workspace: string; viewer: User | null; listing: unknown; enabled: unknown }): Promise<Result<ExtensionInstall>> {
615 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
616 if (!managed.ok) return managed;
617 const listing = String(a.listing ?? "");
618 if (!extensions.extensionIdOf(listing)) return fail("invalid", "Name an extension as extension:<id>.");
619 const changed = await extensions.setEnabled(this.db, managed.value, listing, a.enabled === true, a.viewer!.username);
620 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");
621 return changed;
622 }
623
624 /** Owners cap what an install spends a month. */
625 async setExtensionBudget(a: { workspace: string; viewer: User | null; listing: unknown; monthly_micros: unknown }): Promise<Result<ExtensionInstall>> {
626 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
627 if (!managed.ok) return managed;
628 const listing = String(a.listing ?? "");
629 const micros = a.monthly_micros == null ? null : Number(a.monthly_micros);
630 const changed = await extensions.setBudget(this.db, managed.value, listing, micros);
631 if (changed.ok) this.audit(a.viewer!, a.workspace, "set_extension_budget", listing, `Set ${listing}'s monthly budget`, "marketplace");
632 return changed;
633 }
634
635 /** Owners remove an install. */
636 async uninstallExtension(a: { workspace: string; viewer: User | null; listing: unknown }): Promise<Result<null>> {
637 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
638 if (!managed.ok) return managed;
639 const listing = String(a.listing ?? "");
640 const removed = await extensions.uninstall(this.db, managed.value, listing, a.viewer!.username);
641 if (removed.ok) this.audit(a.viewer!, a.workspace, "uninstall_extension", listing, `Uninstalled ${listing}`, "marketplace");
642 return removed;
643 }
644
645 /** Marks every open request for `listing` added, and tells each person who asked. */
646 private async answerListing(slug: string, workspaceId: string, listing: string, by: User): Promise<void> {
647 try {
648 const answered = await resolveListing(this.db, workspaceId, listing, by as Person);
649 await Promise.all(answered.map((one) => this.tellRequester(slug.toLowerCase(), by, one)));
650 } catch (error) {
651 console.error("agents: requests for a listing were not answered", listing, String(error));
652 }
653 }
654
655 /** Every owner hears of a new request (at most 20 of them). */
656 private async tellOwners(slug: string, asker: User, request: InstallRequest): Promise<void> {
657 if (!this.env.NOTIFY) return;
658 try {
659 const members = await identityClient(this.env.IDENTITY).listMembers(slug, asker);
660 if (!members.ok) throw new Error(members.error.message);
661 const owners = members.value.filter((member) => member.role === "owner").slice(0, 20);
662 const notify = notifyClient(this.env.NOTIFY);
663 await Promise.all(
664 owners.map((owner) =>
665 notify
666 .notify(
667 { username: owner.username },
668 {
669 id: `install-request:${request.id}:${owner.username}`,
670 kind: "approval",
671 workspace: slug,
672 title: `@${asker.username} asks you to add ${request.name}`,
673 body: request.note ?? (request.kind === "extension" ? "An extension. Install it, or turn the request down." : "An integration. Connect it, or turn the request down."),
674 href: requestsPath(slug),
675 actor: { kind: "user", id: asker.id, name: asker.username, avatar: asker.avatar ?? null, avatar_seed: null },
676 created_at: new Date().toISOString(),
677 },
678 )
679 .catch(() => undefined),
680 ),
681 );
682 } catch (error) {
683 console.error("agents: owners were not told of an install request", request.id, String(error));
684 }
685 }
686
687 /** The person who asked hears their request was answered. */
688 private async tellRequester(slug: string, by: User, answered: Answered): Promise<void> {
689 if (!this.env.NOTIFY) return;
690 const { request } = answered;
691 const line = answerLine(request, `@${by.username}`);
692 await notifyClient(this.env.NOTIFY)
693 .notify(
694 { user_id: answered.requested_by_id },
695 {
696 id: `install-answer:${request.id}`,
697 kind: "inbox",
698 workspace: slug,
699 title: line.title,
700 body: line.body,
701 href: request.status === "done" ? listingPath(slug, request.listing) : requestsPath(slug),
702 actor: { kind: "user", id: by.id, name: by.username, avatar: by.avatar ?? null, avatar_seed: null },
703 created_at: new Date().toISOString(),
704 },
705 )
706 .catch((error: unknown) => console.error("agents: a requester was not told", request.id, String(error)));
707 }
708
709 /**
710 * A new agent's first words: the DM with the person who made it is
711 * opened, and the agent says hello there in its own voice, as a reply
712 * billed like any other (a fixed hello when no model can be used).
713 * After the answer; a failure only logs.
714 */
715 private async hello(workspace: string, workspaceId: string, agentId: string, creator: User): Promise<void> {
716 if ((creator.kind ?? "user") !== "user") return;
717 try {
718 const dm = await chatClient(this.env.CHAT).openDm(workspace, creator, [{ kind: "agent", id: agentId }]);
719 if (!dm.ok) throw new Error(dm.error.message);
720 const desk = this.env.DESKS.get(this.env.DESKS.idFromName(agentId));
721 await desk.take({
722 workspace,
723 workspace_id: workspaceId,
724 channel_id: dm.value.id,
725 channel_kind: "dm",
726 channel_name: null,
727 agent_id: agentId,
728 // One hello per agent, however often this runs.
729 message_id: `hello:${agentId}`,
730 thread_root: null,
731 asked_by: creator.id,
732 hops: 0,
733 asker: askerAccess(creator, workspace),
734 hello: true,
735 });
736 } catch (error) {
737 console.error("agents: a new agent's hello was not sent", agentId, String(error));
738 }
739 }
740
741 async update(a: { workspace: string; handle: string; viewer: User | null; changes: Partial<NewWorkspaceAgent> }): Promise<Result<WorkspaceAgent>> {
742 const found = await this.visible(a.workspace, a.viewer, a.handle);
743 if (!found.ok) return found;
744 const { row } = found.value;
745 // Owners change the workspace's agents; a member changes their own personal one.
746 if (!canChange(a.viewer, a.workspace, row)) return fail("forbidden", isPersonal(row) ? NOT_YOURS_REFUSAL : MANAGE_REFUSAL);
747 const workspaceId = found.value.workspaceId;
748 const before = definitionOf(row);
749 const { scope: _scope, ...asked } = (a.changes ?? {}) as Partial<NewWorkspaceAgent>;
750 // The built-in @g1t keeps who it is and its job; the rest is the workspace's.
751 const allowed = row.builtin ? builtinChanges(before, asked) : { ok: true as const, value: asked };
752 if (!allowed.ok) return fail("invalid", allowed.message);
753 const checked = applyChanges(before, allowed.value, TEMPLATE_IDS, { builtin: !!row.builtin });
754 if (!checked.ok) return fail("invalid", checked.message);
755 const definition = checked.value;
756 if (JSON.stringify(definition) === JSON.stringify(before)) return ok(toAgent(row, new Date()));
757 if (definition.handle !== before.handle && (await this.handleTaken(workspaceId, definition.handle, row.id))) {
758 return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
759 }
760 const now = new Date().toISOString();
761 const version = row.version + 1;
762 const by = a.viewer!.username;
763 try {
764 const [updated] = await this.db.batch([
765 // Only from the version read: two owners saving at once never lose
766 // one's change silently; the second is told to look again.
767 updateAgent(this.db, row.id, row.version, definition, { version, updated_at: now }),
768 this.db
769 .prepare(
770 `INSERT INTO agent_versions (agent_id, version, definition, changed_by, created_at)
771 SELECT ?1, ?2, ?3, ?4, ?5 WHERE EXISTS (SELECT 1 FROM agents WHERE id = ?1 AND version = ?2 AND updated_at = ?5)`,
772 )
773 .bind(row.id, version, JSON.stringify(definition), by, now),
774 ]);
775 if (!updated.meta.changes) return fail("conflict", `@${before.handle} was changed meanwhile. Reload it and try again.`);
776 } catch (error) {
777 if (String(error).includes("UNIQUE")) return fail("conflict", `${a.workspace} already has an agent called @${definition.handle}.`);
778 throw error;
779 }
780 const renamed = definition.handle !== before.handle ? ` (was @${before.handle})` : "";
781 this.audit(a.viewer!, a.workspace, "update_agent", definition.handle, `Changed @${definition.handle}${renamed} to version ${version}`, "agents", isPersonal(row) ? "member" : "owner");
782 const saved = await this.row(workspaceId, definition.handle);
783 return ok(toAgent(saved!, new Date()));
784 }
785
786 /** What each effort level has cost one agent, from its own finished sessions. */
787 async effortCosts(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<AgentEffortCosts>> {
788 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
789 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
790 if (!seen.ok) return seen;
791 const row = await this.row(seen.value, a.handle);
792 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
793 const outcomes = await outcomesSince(this.db, seen.value, sinceWindow(), row.id);
794 return ok(effortCostsOf(row.handle, effortOf(definitionOf(row).routing), outcomes.get(row.id) ?? []));
795 }
796
797 /** Ways to spend less, checked against past work: every agent's, or one's. */
798 async recommendations(a: { workspace: string; viewer: User | null; handle?: string | null }): Promise<Result<AgentRecommendations>> {
799 if (a?.viewer && awaitsConfirmation(a.viewer)) return UNVERIFIED;
800 const seen = await this.seen(a?.workspace ?? "", a?.viewer ?? null);
801 if (!seen.ok) return seen;
802 let agentId: string | null = null;
803 if (a.handle) {
804 const row = await this.row(seen.value, a.handle);
805 if (!row) return fail("not_found", `There is no agent called @${a.handle}.`);
806 agentId = row.id;
807 }
808 return ok(await readRecommendations(this.db, seen.value, agentId));
809 }
810
811 /**
812 * Owners apply a suggestion (the agent's effort changes as a new version
813 * of it, and the audit log says so) or dismiss it. Only an open one, and
814 * only while the agent's setting is still the one it was made for.
815 */
816 async resolveRecommendation(a: { workspace: string; viewer: User | null; id: string; action: string }): Promise<Result<AgentRecommendation>> {
817 const managed = await this.managed(a?.workspace ?? "", a?.viewer ?? null);
818 if (!managed.ok) return managed;
819 if (a.action !== "apply" && a.action !== "dismiss") return fail("invalid", "Apply or dismiss.");
820 const found = await recommendationRow(this.db, managed.value, String(a.id ?? ""));
821 if (!found) return fail("not_found", "There is no such suggestion.");
822 if (found.status !== "open") return fail("conflict", "This suggestion was already decided, or is no longer current.");
823 const by = a.viewer!.username;
824 if (a.action === "dismiss") {
825 if (!(await markResolved(this.db, found.id, "dismissed", by))) return fail("conflict", "This suggestion was decided meanwhile.");
826 this.audit(a.viewer!, a.workspace, "dismiss_recommendation", found.handle ?? found.agent_id, `Dismissed "${found.title}"`);
827 } else {
828 const agent = await this.row(managed.value, found.handle);
829 if (!agent || agent.id !== found.agent_id) return fail("not_found", "That agent is gone.");
830 if (effortOf(definitionOf(agent).routing) !== found.from_effort) {
831 await this.db.prepare("UPDATE agent_recommendations SET status = 'stale' WHERE id = ? AND status = 'open'").bind(found.id).run();
832 return fail("conflict", `@${agent.handle}'s effort was changed since this was suggested. The next check looks again.`);
833 }
834 // Claimed first, so two owners pressing Apply change the agent once.
835 if (!(await markResolved(this.db, found.id, "applied", by))) return fail("conflict", "This suggestion was decided meanwhile.");
836 const changed = await this.update({ workspace: a.workspace, handle: agent.handle, viewer: a.viewer, changes: { routing: { effort: found.to_effort as never } } });
837 if (!changed.ok) {
838 await this.db.prepare("UPDATE agent_recommendations SET status = 'open', resolved_by = NULL, resolved_at = NULL WHERE id = ?").bind(found.id).run();
839 return changed;
840 }
841 this.audit(
842 a.viewer!,
843 a.workspace,
844 "apply_recommendation",
845 agent.handle,
846 `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}`,
847 );
848 }
849 const after = await recommendationRow(this.db, managed.value, found.id);
850 return ok(toRecommendation(after!));
851 }
852
853 async archive(a: { workspace: string; handle: string; viewer: User | null }): Promise<Result<null>> {
854 const found = await this.visible(a.workspace, a.viewer, a.handle);
855 if (!found.ok) return found;
856 const { row } = found.value;
857 // Owners archive any agent; a member their own personal one.
858 if (!canArchive(a.viewer, a.workspace, row)) return fail("forbidden", MANAGE_REFUSAL);
859 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.");
860 const now = new Date().toISOString();
861 await this.db.prepare("UPDATE agents SET archived_at = ?, updated_at = ? WHERE id = ? AND archived_at IS NULL").bind(now, now, row.id).run();
862 this.audit(a.viewer!, a.workspace, "archive_agent", row.handle, `Archived @${row.handle}`);
863 // Its computer goes to sleep now; its disk is kept 30 days, then deleted (computer.ts).
864 this.defer(forgetComputer(this.env, row.id));
865 return ok(null);
866 }
867
868 /**
869 * Makes the workspace's built-in @g1t if it does not exist yet: an
870 * ordinary agent row, marked builtin, at version 1. Safe to call on
871 * every request; once seen, an isolate does not write again.
872 */
873 private async ensureBuiltin(workspaceId: string): Promise<void> {
874 await ensureBuiltin(this.db, workspaceId);
875 }
876
877 /** Internal: the workspace's @g1t, made if need be, for the chat service. */
878 async builtin(a: { workspace: string; workspace_id: string }): Promise<Result<WorkspaceAgent>> {
879 if (typeof a?.workspace_id !== "string" || !a.workspace_id) return fail("invalid", "Name the workspace by id.");
880 await this.ensureBuiltin(a.workspace_id);
881 const now = new Date();
882 const row = await this.db
883 .prepare(selectAgents("a.workspace_id = ?3 AND a.builtin = 1"))
884 .bind(...periods(now), a.workspace_id)
885 .first<Row>();
886 return row ? ok(toAgent(row, now)) : fail("not_found", "This workspace's @g1t could not be made.");
887 }
888
889 /**
890 * Hands a message to the agent's desk, which answers it in the
891 * background. Returns as soon as the desk holds it.
892 */
893 async deliver(delivery: AgentDelivery): Promise<Result<null>> {
894 const fields = ["workspace", "workspace_id", "channel_id", "agent_id", "message_id", "asked_by"] as const;
895 if (!delivery || fields.some((field) => typeof delivery[field] !== "string" || !delivery[field])) {
896 return fail("invalid", "A delivery names the workspace, channel, agent, message and who asked.");
897 }
898 if (delivery.channel_kind !== "channel" && delivery.channel_kind !== "dm") return fail("invalid", "channel_kind is channel or dm.");
899 await this.ensureBuiltin(delivery.workspace_id);
900 const desk = this.env.DESKS.get(this.env.DESKS.idFromName(delivery.agent_id));
901 await desk.take({ ...delivery, hops: Math.max(0, Math.floor(Number(delivery.hops) || 0)), thread_root: delivery.thread_root ?? null });
902 return ok(null);
903 }
904
905 private versionStatement(agentId: string, version: number, d: Definition, by: string, at: string): D1PreparedStatement {
906 return versionStatement(this.db, agentId, version, d, by, at);
907 }
908
909 /**
910 * Records a change to an agent in the workspace's audit log, after the
911 * answer. Never fails the change: a log that cannot be written is logged.
912 * The events service's audit contract speaks camelCase (Rust's
913 * `NewAuditEntry`).
914 */
915 /** A change to the skill library, in the audit log as `agents/skills/<name>`. */
916 auditSkill(actor: User, workspace: string, action: string, name: string, message: string): void {
917 this.audit(actor, workspace, action, `skills/${name}`, message);
918 }
919
920 private audit(
921 actor: User,
922 workspace: string,
923 action: string,
924 handle: string,
925 message: string,
926 area: "agents" | "marketplace" = "agents",
927 rule: "owner" | "member" | null = null,
928 ): void {
929 const kind = actor.kind === "workspace" ? "workspace" : actor.kind === "agent" ? "agent" : actor.kind === "system" ? "system" : "person";
930 const entry = {
931 actorKind: kind,
932 actor: actor.username,
933 actorId: actor.id,
934 agent: null,
935 onBehalfOf: null,
936 runId: null,
937 runKind: null,
938 credentialId: null,
939 action,
940 surface: "web",
941 workspace: workspace.toLowerCase(),
942 repo: null,
943 number: null,
944 gitRef: null,
945 path: `${area}/${handle}`,
946 outcome: "allowed",
947 rule: rule ?? (area === "marketplace" && action === "request_install" ? "member" : "owner"),
948 result: "ok",
949 message,
950 requestId: `req_${crypto.randomUUID()}`,
951 };
952 this.defer(
953 this.env.EVENTS.fetch("https://service/rpc/audit_record", {
954 method: "POST",
955 headers: { "content-type": "application/json" },
956 body: JSON.stringify({ entries: [entry] }),
957 })
958 .then((response) => {
959 if (!response.ok) throw new Error(`status ${response.status}`);
960 })
961 .catch((error: unknown) => console.error("agents: audit entry not recorded", action, handle, String(error))),
962 );
963 }
964}
965
966/** One RPC method's answer. */
967async function answer(service: Agents, method: string, args: any): Promise<Response> {
968 // The skill library (./skill-rpc.ts).
969 const skill = await skillRpc(method, args, (a, run) => service.view(a, run), (viewer, workspace, action, name, message) => service.auditSkill(viewer, workspace, action, name, message));
970 if (skill) return Response.json(skill);
971 switch (method) {
972 case "list":
973 return Response.json(await service.list(args));
974 case "get":
975 return Response.json(await service.get(args));
976 case "by_ids":
977 return Response.json(await service.byIds(args));
978 case "create":
979 return Response.json(await service.create(args));
980 case "update":
981 return Response.json(await service.update(args));
982 case "archive":
983 return Response.json(await service.archive(args));
984 case "draft":
985 return Response.json(await service.draft(args));
986 case "try_draft":
987 return Response.json(await service.tryDraft(args));
988 case "redraft":
989 return Response.json(await service.redraft(args));
990 case "add_mcp_server":
991 return Response.json(await service.addMcpServer(args));
992 case "remove_mcp_server":
993 return Response.json(await service.removeMcpServer(args));
994 case "refresh_mcp_server":
995 return Response.json(await service.refreshMcpServer(args));
996 case "promote":
997 return Response.json(await service.promote(args));
998 case "builtin":
999 return Response.json(await service.builtin(args));
1000 case "templates":
1001 return Response.json(TEMPLATES);
1002 case "deliver":
1003 return Response.json(await service.deliver(args));
1004 case "overview":
1005 return Response.json(await service.view(args, (ctx) => views.overview(ctx)));
1006 case "sessions":
1007 return Response.json(await service.view(args, (ctx) => views.listSessions(ctx, args)));
1008 case "session":
1009 return Response.json(await service.view(args, (ctx) => views.sessionDetail(ctx, args.id)));
1010 case "stop_session":
1011 return Response.json(await service.view(args, (ctx) => views.stopSession(ctx, args.id)));
1012 case "approve_session":
1013 return Response.json(await service.view(args, (ctx) => views.approveSession(ctx, args.id, args.cap_micros)));
1014 case "steer_session":
1015 return Response.json(await service.view(args, (ctx) => views.steerSession(ctx, args.id, args.body)));
1016 // The agent's own computer (./computer.ts).
1017 case "computer":
1018 return Response.json(await service.view(args, (ctx) => computerView(ctx, args.handle)));
1019 case "computer_wake":
1020 return Response.json(await service.view(args, (ctx) => wakeComputer(ctx, args.handle)));
1021 case "computer_sleep":
1022 return Response.json(await service.view(args, (ctx) => sleepComputer(ctx, args.handle)));
1023 case "computer_reset":
1024 return Response.json(await service.view(args, (ctx) => resetComputer(ctx, args.handle)));
1025 case "computer_commands":
1026 return Response.json(await service.view(args, (ctx) => computerCommands(ctx, args.handle, args.session_id ?? null)));
1027 case "memories":
1028 return Response.json(await service.view(args, (ctx) => views.memories(ctx, args.handle)));
1029 case "remember":
1030 return Response.json(await service.view(args, (ctx) => views.remember(ctx, args.handle, args.input)));
1031 case "update_memory":
1032 return Response.json(await service.view(args, (ctx) => views.updateMemory(ctx, args.handle, args.id, args.changes)));
1033 case "forget":
1034 return Response.json(await service.view(args, (ctx) => views.forget(ctx, args.handle, args.id)));
1035 case "routines":
1036 return Response.json(await service.view(args, (ctx) => views.routines(ctx, args.handle)));
1037 case "save_routine":
1038 return Response.json(await service.view(args, (ctx) => views.saveRoutine(ctx, args.handle, args.input, args.id ?? null)));
1039 case "delete_routine":
1040 return Response.json(await service.view(args, (ctx) => views.deleteRoutine(ctx, args.handle, args.id)));
1041 case "run_routine":
1042 return Response.json(await service.view(args, (ctx) => views.runRoutineNow(ctx, args.handle, args.id)));
1043 case "spend":
1044 return Response.json(await service.view(args, (ctx) => views.spend(ctx, args.handle ?? null, { period: args.period, person: args.person })));
1045 case "person_budgets":
1046 return Response.json(await service.view(args, (ctx) => views.personBudgetsView(ctx)));
1047 case "set_person_budget":
1048 return Response.json(await service.view(args, (ctx) => views.setPersonBudget(ctx, args.username, args.monthly_micros)));
1049 case "team_context":
1050 return Response.json(await service.view(args, (ctx) => views.teamContext(ctx, args.handle)));
1051 case "activity":
1052 return Response.json(await service.view(args, (ctx) => views.activity(ctx, args.handle)));
1053 case "versions":
1054 return Response.json(await service.view(args, (ctx) => views.versions(ctx, args.handle)));
1055 case "card_action":
1056 return Response.json(await service.cardAction(args));
1057 case "policy":
1058 return Response.json(await service.view(args, (ctx) => views.policy(ctx)));
1059 case "install_requests":
1060 return Response.json(await service.installRequests(args));
1061 case "request_install":
1062 return Response.json(await service.requestInstall(args));
1063 case "resolve_install_request":
1064 return Response.json(await service.resolveInstallRequest(args));
1065 case "extension_installs":
1066 return Response.json(await service.extensionInstalls(args));
1067 case "install_extension":
1068 return Response.json(await service.installExtension(args));
1069 case "set_extension_enabled":
1070 return Response.json(await service.setExtensionEnabled(args));
1071 case "set_extension_budget":
1072 return Response.json(await service.setExtensionBudget(args));
1073 case "uninstall_extension":
1074 return Response.json(await service.uninstallExtension(args));
1075 case "effort_costs":
1076 return Response.json(await service.effortCosts(args));
1077 case "recommendations":
1078 return Response.json(await service.recommendations(args));
1079 case "resolve_recommendation":
1080 return Response.json(await service.resolveRecommendation(args));
1081 case "set_policy":
1082 return Response.json(await service.view(args, (ctx) => views.setPolicy(ctx, args.policy)));
1083 default:
1084 return new Response("Unknown method\n", { status: 404 });
1085 }
1086}
1087
1088export default {
1089 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
1090 const match = new URL(request.url).pathname.match(/^\/rpc\/([a-z_]+)$/);
1091 if (request.method !== "POST" || !match) return new Response("Not found\n", { status: 404 });
1092 // A replica near the caller when it asks for one (@g1t/contracts d1.ts).
1093 const opened = openD1(env.DB, request);
1094 const service = new Agents(Object.create(env, { DB: { value: opened.db } }) as Env, (work) => ctx.waitUntil(work));
1095 const args = (await request.json().catch(() => ({}))) as any;
1096 return opened.finish(await answer(service, match[1], args));
1097 },
1098
1099 /** Events routines run on, from the events service (SUBSCRIBER_AGENTS). */
1100 async queue(batch: MessageBatch<unknown>, env: Env): Promise<void> {
1101 const events = batch.messages.map((message) => message.body as G1tEvent);
1102 // A push to a default branch: skill libraries that follow that repository read it again (./skill-library.ts).
1103 const pushes = events.flatMap((event) => (event?.type === "git.push" && event.data?.defaultBranch && event.data.repoId ? [{ repoId: event.data.repoId }] : []));
1104 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]);
1105 batch.ackAll();
1106 },
1107
1108 /** Every few minutes: routines that are due, session steps a desk lost, the spend check, and the one-time team move. */
1109 async scheduled(_controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
1110 const sessions = env as unknown as SessionEnv;
1111 ctx.waitUntil(
1112 Promise.all([
1113 runDue(sessions).catch((error: unknown) => console.error("agents: routines did not run", String(error))),
1114 sweep(sessions).catch((error: unknown) => console.error("agents: the session sweep failed", String(error))),
1115 // Spend less, keep quality: each workspace checked weekly (src/recommend.ts).
1116 checkDue(env.DB).catch((error: unknown) => console.error("agents: the spend check failed", String(error))),
1117 // Once: agents' own old teams become team memberships (src/team-move.ts); after that, one cheap read.
1118 moveAgentTeams(env.DB, (id, agents) => identityClient(env.IDENTITY).adoptAgentTeams(id, agents)).catch((error: unknown) => console.error("agents: moving agents' teams failed", String(error))),
1119 ]),
1120 );
1121 },
1122} satisfies ExportedHandler<Env>;