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