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