| 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 | */ |
| 12 | import { |
| 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 | |
| 47 | import { MANAGE_REFUSAL, NOT_YOURS_REFUSAL, canArchive, canChange, canManage, canSee, canSeeAgent, creatableScope } from "./access.ts"; |
| 48 | import { DRAFT_OUTPUT_TOKENS, draftSystem, jsonIn, proposalFrom, redraftFrom, redraftSystem, startingBudget, trySystem, tryTurns, wordsOf } from "./builder.ts"; |
| 49 | import { callBuilder } from "./builder-call.ts"; |
| 50 | import { type Definition, applyChanges } from "./definition.ts"; |
| 51 | import { listMcpTools } from "./mcp-client.ts"; |
| 52 | import { MAX_MCP_SERVERS, checkMcpName, checkMcpUrl } from "../../../packages/contracts/src/abilities.ts"; |
| 53 | import { builtinChanges } from "./orchestrator.ts"; |
| 54 | import type { Desk } from "./desk.ts"; |
| 55 | import type { ReplyEnv } from "./reply.ts"; |
| 56 | import { ensureBuiltin } from "./builtin.ts"; |
| 57 | import { type Row, definitionOf, insertAgent, isPersonal, periods, selectAgents, toAgent, updateAgent, versionStatement } from "./store.ts"; |
| 58 | import { TEMPLATES, TEMPLATE_IDS } from "./templates.ts"; |
| 59 | import { readPolicy } from "./policy.ts"; |
| 60 | import { runDue } from "./routines.ts"; |
| 61 | import { onEvents } from "./triggers.ts"; |
| 62 | import { cardAction } from "./cards.ts"; |
| 63 | import { type SessionEnv, sweep } from "./sessions.ts"; |
| 64 | import { EFFORT_NAMES, checkDue, effortCostsOf, markResolved, outcomesSince, readRecommendations, recommendationRow, sinceWindow, toRecommendation } from "./recommend.ts"; |
| 65 | import { moveAgentTeams } from "./team-move.ts"; |
| 66 | import { effortOf } from "./routing.ts"; |
| 67 | import * as views from "./views.ts"; |
| 68 | import { monthKey } from "./budget.ts"; |
| 69 | import * as extensions from "./extensions.ts"; |
| 70 | import { skillPushes, skillRpc } from "./skill-rpc.ts"; |
| 71 | import { type Answered, type Person, answerLine, findListing, listRequests, listingPath, openRequest, requestsPath, resolveListing, resolveRequest } from "./installs.ts"; |
| 72 | |
| 73 | export { Desk } from "./desk.ts"; |
| 74 | |
| 75 | type 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. */ |
| 83 | const workspaceIds = new Map<string, { id: string | null; until: number }>(); |
| 84 | |
| 85 | class 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. */ |
| 962 | async 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 | |
| 1072 | export 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>; |