| 1 | /** |
| 2 | * What Agents mode reads and changes beyond definitions: sessions, memory, |
| 3 | * routines, spend, activity, versions and the workspace's agent policy. |
| 4 | * Called by the RPC methods in index.ts once they know who is asking and |
| 5 | * that they may see the workspace. |
| 6 | * |
| 7 | * Privacy follows the conversations work came from. A session, a reply or |
| 8 | * a memory from a conversation the viewer is not in shows that it happened |
| 9 | * and what it cost, never what it was about, owners included: owners |
| 10 | * control money and agents, not other people's conversations. |
| 11 | */ |
| 12 | import { |
| 13 | type AgentActivity, |
| 14 | type AgentMemory, |
| 15 | type AgentPolicy, |
| 16 | type AgentRoutine, |
| 17 | type AgentSession, |
| 18 | type AgentSessionDetail, |
| 19 | type AgentSpendBreakdown, |
| 20 | type AgentVersion, |
| 21 | type AgentsOverview, |
| 22 | type AgentMemoryScope, |
| 23 | type NewRoutine, |
| 24 | type PersonBudgets, |
| 25 | type Result, |
| 26 | type RoutineSuggestion, |
| 27 | type SessionEvent, |
| 28 | type SpendSlice, |
| 29 | type User, |
| 30 | type WorkspaceAgent, |
| 31 | chatClient, |
| 32 | fail, |
| 33 | newId, |
| 34 | ok, |
| 35 | } from "@g1t/contracts"; |
| 36 | |
| 37 | import { monthKey, spendSpan } from "./budget.ts"; |
| 38 | import { cleanUsername, personBudgets } from "./person-budget.ts"; |
| 39 | import { type MemoryRow, type MemoryViewer, changeableBy, cleanFact, toMemory, visibleTo } from "./memory.ts"; |
| 40 | import { DEFAULT_POLICY, checkPolicy, readPolicy } from "./policy.ts"; |
| 41 | import { type RoutineRow, MAX_ROUTINES, checkRoutine, newRoutineId, nextRun, runRoutine, toRoutine } from "./routines.ts"; |
| 42 | import { type SessionEnv, type SessionRow, LIVE, approve, sessionRow, steer, stop, toSession } from "./sessions.ts"; |
| 43 | import { type Row, definitionOf, periods, selectAgents, toAgent } from "./store.ts"; |
| 44 | import { suggestRoutines } from "./suggest.ts"; |
| 45 | |
| 46 | export type ViewContext = { |
| 47 | env: SessionEnv; |
| 48 | db: D1Database; |
| 49 | slug: string; |
| 50 | workspaceId: string; |
| 51 | viewer: User; |
| 52 | /** Whether the viewer owns the workspace (or is its token). */ |
| 53 | owner: boolean; |
| 54 | }; |
| 55 | |
| 56 | type AgentFace = { handle: string; display_name: string; avatar_seed: string; team: string | null; department: string }; |
| 57 | |
| 58 | /** The workspace's agents by id, archived ones too, for names on sessions and spend. */ |
| 59 | async function faces(ctx: ViewContext): Promise<Map<string, AgentFace>> { |
| 60 | const rows = await ctx.db |
| 61 | .prepare("SELECT id, handle, display_name, avatar_seed, team, department FROM agents WHERE workspace_id = ?") |
| 62 | .bind(ctx.workspaceId) |
| 63 | .all<{ id: string } & AgentFace>(); |
| 64 | return new Map(rows.results.map((r) => [r.id, { ...r, avatar_seed: r.avatar_seed || r.handle }])); |
| 65 | } |
| 66 | |
| 67 | /** |
| 68 | * Which of these conversations the viewer is in (or can read: a public |
| 69 | * channel). Asked of chat once per conversation, at most 60. |
| 70 | */ |
| 71 | async function readable(ctx: ViewContext, channelIds: string[]): Promise<Set<string>> { |
| 72 | const out = new Set<string>(); |
| 73 | if (ctx.viewer.kind === "workspace") return out; |
| 74 | const chat = chatClient(ctx.env.CHAT); |
| 75 | const ids = [...new Set(channelIds)].slice(0, 60); |
| 76 | await Promise.all( |
| 77 | ids.map(async (id) => { |
| 78 | const audience = await chat.audience(ctx.slug, id).catch(() => null); |
| 79 | if (audience?.ok && (audience.value.kind === "public" || audience.value.member_user_ids.includes(ctx.viewer.id))) out.add(id); |
| 80 | }), |
| 81 | ); |
| 82 | return out; |
| 83 | } |
| 84 | |
| 85 | async function sessionsOut(ctx: ViewContext, rows: SessionRow[], agents?: Map<string, AgentFace>): Promise<AgentSession[]> { |
| 86 | const names = agents ?? (await faces(ctx)); |
| 87 | const can = await readable(ctx, rows.map((r) => r.channel_id)); |
| 88 | return rows.map((row) => toSession(row, names.get(row.agent_id) ?? null, can.has(row.channel_id))); |
| 89 | } |
| 90 | |
| 91 | async function agentByHandle(ctx: ViewContext, handle: string): Promise<Row | null> { |
| 92 | return ctx.db |
| 93 | .prepare("SELECT * FROM agents WHERE workspace_id = ? AND handle = ? AND archived_at IS NULL") |
| 94 | .bind(ctx.workspaceId, String(handle ?? "").trim().replace(/^@/, "").toLowerCase()) |
| 95 | .first<Row>(); |
| 96 | } |
| 97 | |
| 98 | // ── Sessions ────────────────────────────────────────────────────────────── |
| 99 | |
| 100 | export async function listSessions(ctx: ViewContext, filter: { handle?: string | null; status?: "live" | "done" | null; limit?: number | null }): Promise<Result<AgentSession[]>> { |
| 101 | const limit = Math.min(200, Math.max(1, Math.floor(Number(filter.limit) || 50))); |
| 102 | const where = ["workspace_id = ?"]; |
| 103 | const binds: (string | number)[] = [ctx.workspaceId]; |
| 104 | if (filter.handle) { |
| 105 | const agent = await agentByHandle(ctx, filter.handle); |
| 106 | if (!agent) return fail("not_found", `There is no agent called @${filter.handle}.`); |
| 107 | where.push("agent_id = ?"); |
| 108 | binds.push(agent.id); |
| 109 | } |
| 110 | if (filter.status === "live") where.push(`status IN (${LIVE.map(() => "?").join(", ")})`), binds.push(...LIVE); |
| 111 | if (filter.status === "done") where.push(`status NOT IN (${LIVE.map(() => "?").join(", ")})`), binds.push(...LIVE); |
| 112 | const rows = await ctx.db |
| 113 | .prepare(`SELECT * FROM agent_sessions WHERE ${where.join(" AND ")} ORDER BY created_at DESC LIMIT ?`) |
| 114 | .bind(...binds, limit) |
| 115 | .all<SessionRow>(); |
| 116 | return ok(await sessionsOut(ctx, rows.results)); |
| 117 | } |
| 118 | |
| 119 | export async function sessionDetail(ctx: ViewContext, id: string): Promise<Result<AgentSessionDetail>> { |
| 120 | const row = await sessionRow(ctx.db, String(id ?? "")); |
| 121 | if (!row || row.workspace_id !== ctx.workspaceId) return fail("not_found", "There is no such session."); |
| 122 | const agents = await faces(ctx); |
| 123 | const treeRows = await ctx.db.prepare("SELECT * FROM agent_sessions WHERE root_id = ? ORDER BY created_at LIMIT 100").bind(row.root_id).all<SessionRow>(); |
| 124 | const [session] = await sessionsOut(ctx, [row], agents); |
| 125 | const tree = await sessionsOut(ctx, treeRows.results, agents); |
| 126 | let events: SessionEvent[] = []; |
| 127 | if (session.visible) { |
| 128 | const rows = await ctx.db |
| 129 | .prepare("SELECT seq, kind, by_name, body, tool, outcome, created_at FROM agent_session_events WHERE session_id = ? ORDER BY seq LIMIT 1000") |
| 130 | .bind(row.id) |
| 131 | .all<{ seq: number; kind: SessionEvent["kind"]; by_name: string | null; body: string; tool: string | null; outcome: string | null; created_at: string }>(); |
| 132 | events = rows.results.map((e) => ({ seq: e.seq, kind: e.kind, by: e.by_name, body: e.body, tool: e.tool, outcome: e.outcome, created_at: e.created_at })); |
| 133 | } |
| 134 | const live = LIVE.includes(row.status as (typeof LIVE)[number]); |
| 135 | return ok({ |
| 136 | session, |
| 137 | events, |
| 138 | tree, |
| 139 | can_stop: session.visible && live, |
| 140 | can_steer: session.visible, |
| 141 | can_approve: row.status === "needs_approval" && (ctx.owner || (session.visible && row.asked_by === ctx.viewer.id && ctx.owner)), |
| 142 | }); |
| 143 | } |
| 144 | |
| 145 | async function visibleSession(ctx: ViewContext, id: string): Promise<Result<SessionRow>> { |
| 146 | const row = await sessionRow(ctx.db, String(id ?? "")); |
| 147 | if (!row || row.workspace_id !== ctx.workspaceId) return fail("not_found", "There is no such session."); |
| 148 | const can = await readable(ctx, [row.channel_id]); |
| 149 | if (!can.has(row.channel_id)) return fail("not_found", "There is no such session."); |
| 150 | return ok(row); |
| 151 | } |
| 152 | |
| 153 | export async function stopSession(ctx: ViewContext, id: string): Promise<Result<AgentSession>> { |
| 154 | const found = await visibleSession(ctx, id); |
| 155 | if (!found.ok) return found; |
| 156 | if (!LIVE.includes(found.value.status as (typeof LIVE)[number])) return fail("invalid", "That session is already over."); |
| 157 | const row = await stop(ctx.env, found.value, ctx.viewer.username); |
| 158 | return ok((await sessionsOut(ctx, [row]))[0]); |
| 159 | } |
| 160 | |
| 161 | export async function steerSession(ctx: ViewContext, id: string, body: string): Promise<Result<AgentSession>> { |
| 162 | const found = await visibleSession(ctx, id); |
| 163 | if (!found.ok) return found; |
| 164 | const text = typeof body === "string" ? body.trim() : ""; |
| 165 | if (!text) return fail("invalid", "Say something to the session."); |
| 166 | const row = await steer(ctx.env, found.value, ctx.viewer.username, text.slice(0, 4000)); |
| 167 | return ok((await sessionsOut(ctx, [row]))[0]); |
| 168 | } |
| 169 | |
| 170 | /** Owners raise a session's cap; it must end up above what it has spent. */ |
| 171 | export async function approveSession(ctx: ViewContext, id: string, capMicros: number): Promise<Result<AgentSession>> { |
| 172 | if (!ctx.owner) return fail("forbidden", "Only the workspace's owners can approve more spend."); |
| 173 | const row = await sessionRow(ctx.db, String(id ?? "")); |
| 174 | if (!row || row.workspace_id !== ctx.workspaceId) return fail("not_found", "There is no such session."); |
| 175 | if (row.status !== "needs_approval") return fail("invalid", "That session isn't waiting for approval."); |
| 176 | const cap = Math.floor(Number(capMicros)); |
| 177 | if (!Number.isFinite(cap) || cap <= row.charged_micros || cap > 1_000_000_000) return fail("invalid", "The new cap must be above what it has spent."); |
| 178 | const fresh = await approve(ctx.env, row, ctx.viewer.username, cap); |
| 179 | return ok((await sessionsOut(ctx, [fresh]))[0]); |
| 180 | } |
| 181 | |
| 182 | // ── Memory ──────────────────────────────────────────────────────────────── |
| 183 | |
| 184 | async function memoryViewer(ctx: ViewContext, rows: MemoryRow[]): Promise<MemoryViewer> { |
| 185 | const channels = await readable(ctx, rows.filter((r) => r.scope === "channel").map((r) => r.scope_ref)); |
| 186 | return { id: ctx.viewer.id, owner: ctx.owner, inChannel: (id) => channels.has(id) }; |
| 187 | } |
| 188 | |
| 189 | export async function memories(ctx: ViewContext, handle: string): Promise<Result<AgentMemory[]>> { |
| 190 | const agent = await agentByHandle(ctx, handle); |
| 191 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 192 | const rows = await ctx.db.prepare("SELECT * FROM agent_memories WHERE agent_id = ? ORDER BY pinned DESC, updated_at DESC LIMIT 500").bind(agent.id).all<MemoryRow>(); |
| 193 | const viewer = await memoryViewer(ctx, rows.results); |
| 194 | return ok(rows.results.filter((row) => visibleTo(row, viewer)).map(toMemory)); |
| 195 | } |
| 196 | |
| 197 | /** |
| 198 | * A fact a person gives an agent. Workspace facts are the owners'; a |
| 199 | * channel fact needs the person to be in that channel; a person fact is |
| 200 | * always their own. |
| 201 | */ |
| 202 | export async function remember(ctx: ViewContext, handle: string, input: { body: string; scope: AgentMemoryScope; scope_ref?: string | null }): Promise<Result<AgentMemory>> { |
| 203 | const agent = await agentByHandle(ctx, handle); |
| 204 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 205 | const body = cleanFact(input?.body); |
| 206 | if (!body) return fail("invalid", "Say what it should remember."); |
| 207 | let scope: AgentMemoryScope = input?.scope === "workspace" || input?.scope === "channel" ? input.scope : "person"; |
| 208 | let ref = ""; |
| 209 | let label: string | null = null; |
| 210 | if (scope === "workspace" && !ctx.owner) return fail("forbidden", "Only owners give an agent facts for the whole workspace."); |
| 211 | if (scope === "channel") { |
| 212 | ref = String(input.scope_ref ?? ""); |
| 213 | const can = await readable(ctx, [ref]); |
| 214 | if (!can.has(ref)) return fail("forbidden", "You can only give it facts for conversations you're in."); |
| 215 | const audience = await chatClient(ctx.env.CHAT).audience(ctx.slug, ref).catch(() => null); |
| 216 | label = audience?.ok && "name" in audience.value ? ((audience.value as { name?: string | null }).name ?? null) : null; |
| 217 | } |
| 218 | if (scope === "person") { |
| 219 | scope = "person"; |
| 220 | ref = ctx.viewer.id; |
| 221 | label = ctx.viewer.username; |
| 222 | } |
| 223 | const id = newId("mem"); |
| 224 | const now = new Date().toISOString(); |
| 225 | await ctx.db |
| 226 | .prepare( |
| 227 | `INSERT INTO agent_memories (id, agent_id, workspace_id, scope, scope_ref, scope_label, body, source_kind, source_ref, source_label, created_by, created_by_kind, pinned, created_at, updated_at) |
| 228 | VALUES (?, ?, ?, ?, ?, ?, ?, 'person', ?, ?, ?, 'user', 1, ?, ?)`, |
| 229 | ) |
| 230 | .bind(id, agent.id, ctx.workspaceId, scope, ref, label, body, ctx.viewer.username, `@${ctx.viewer.username}`, ctx.viewer.username, now, now) |
| 231 | .run(); |
| 232 | const row = await ctx.db.prepare("SELECT * FROM agent_memories WHERE id = ?").bind(id).first<MemoryRow>(); |
| 233 | return ok(toMemory(row!)); |
| 234 | } |
| 235 | |
| 236 | async function changeable(ctx: ViewContext, handle: string, id: string): Promise<Result<MemoryRow>> { |
| 237 | const agent = await agentByHandle(ctx, handle); |
| 238 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 239 | const row = await ctx.db.prepare("SELECT * FROM agent_memories WHERE id = ? AND agent_id = ?").bind(String(id ?? ""), agent.id).first<MemoryRow>(); |
| 240 | if (!row) return fail("not_found", "There is no such memory."); |
| 241 | const viewer = await memoryViewer(ctx, [row]); |
| 242 | if (!visibleTo(row, viewer)) return fail("not_found", "There is no such memory."); |
| 243 | if (!changeableBy(row, viewer)) return fail("forbidden", "Only owners change what an agent knows for the whole workspace."); |
| 244 | return ok(row); |
| 245 | } |
| 246 | |
| 247 | export async function updateMemory(ctx: ViewContext, handle: string, id: string, changes: { body?: string; pinned?: boolean }): Promise<Result<AgentMemory>> { |
| 248 | const found = await changeable(ctx, handle, id); |
| 249 | if (!found.ok) return found; |
| 250 | const body = changes?.body === undefined ? found.value.body : cleanFact(changes.body); |
| 251 | if (!body) return fail("invalid", "A memory can't be empty; forget it instead."); |
| 252 | const pinned = changes?.pinned === undefined ? found.value.pinned : changes.pinned ? 1 : 0; |
| 253 | const now = new Date().toISOString(); |
| 254 | const edited = body !== found.value.body; |
| 255 | await ctx.db |
| 256 | .prepare( |
| 257 | `UPDATE agent_memories SET body = ?, pinned = ?, updated_at = ?${edited ? ", source_kind = 'person', source_ref = ?, source_label = ?" : ""} WHERE id = ?`, |
| 258 | ) |
| 259 | .bind(...(edited ? [body, pinned, now, ctx.viewer.username, `@${ctx.viewer.username} (corrected)`, found.value.id] : [body, pinned, now, found.value.id])) |
| 260 | .run(); |
| 261 | const row = await ctx.db.prepare("SELECT * FROM agent_memories WHERE id = ?").bind(found.value.id).first<MemoryRow>(); |
| 262 | return ok(toMemory(row!)); |
| 263 | } |
| 264 | |
| 265 | export async function forget(ctx: ViewContext, handle: string, id: string): Promise<Result<null>> { |
| 266 | const found = await changeable(ctx, handle, id); |
| 267 | if (!found.ok) return found; |
| 268 | await ctx.db.prepare("DELETE FROM agent_memories WHERE id = ?").bind(found.value.id).run(); |
| 269 | return ok(null); |
| 270 | } |
| 271 | |
| 272 | // ── Routines ────────────────────────────────────────────────────────────── |
| 273 | |
| 274 | /** An agent's routines, and routines its responsibilities suggest that it doesn't have yet. */ |
| 275 | export async function routines(ctx: ViewContext, handle: string): Promise<Result<{ routines: AgentRoutine[]; suggestions: RoutineSuggestion[] }>> { |
| 276 | const agent = await agentByHandle(ctx, handle); |
| 277 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 278 | const rows = await ctx.db.prepare("SELECT * FROM agent_routines WHERE agent_id = ? ORDER BY created_at").bind(agent.id).all<RoutineRow>(); |
| 279 | const list = rows.results.map(toRoutine); |
| 280 | const duties = definitionOf(agent).responsibilities; |
| 281 | return ok({ routines: list, suggestions: suggestRoutines(duties, list.map((r) => ({ name: r.name, events: r.events }))) }); |
| 282 | } |
| 283 | |
| 284 | /** |
| 285 | * Owners keep an agent's routines. The person who saves one becomes its |
| 286 | * sponsor: it runs with their access, in a channel they and the agent are in. |
| 287 | */ |
| 288 | export async function saveRoutine(ctx: ViewContext, handle: string, input: NewRoutine, id: string | null): Promise<Result<AgentRoutine>> { |
| 289 | if (!ctx.owner) return fail("forbidden", "Only the workspace's owners set up routines."); |
| 290 | if (ctx.viewer.kind === "workspace") return fail("invalid", "A routine runs with a person's access: set it up signed in as yourself."); |
| 291 | const agent = await agentByHandle(ctx, handle); |
| 292 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 293 | const checked = checkRoutine(input); |
| 294 | if (!checked.ok) return fail("invalid", checked.message); |
| 295 | const r = checked.value; |
| 296 | const audience = await chatClient(ctx.env.CHAT).audience(ctx.slug, r.channel_id).catch(() => null); |
| 297 | if (!audience?.ok || (audience.value.kind !== "public" && !audience.value.member_user_ids.includes(ctx.viewer.id))) { |
| 298 | return fail("invalid", "Choose a channel you're in."); |
| 299 | } |
| 300 | const channels = await chatClient(ctx.env.CHAT).sidebar(ctx.slug, ctx.viewer).catch(() => null); |
| 301 | const channelName = channels?.ok ? (channelNameIn(channels.value, r.channel_id) ?? null) : null; |
| 302 | const now = new Date(); |
| 303 | const next = r.enabled !== false && r.schedule ? nextRun(r.schedule, now).toISOString() : null; |
| 304 | const schedule = r.schedule ? JSON.stringify(r.schedule) : null; |
| 305 | if (id) { |
| 306 | const existing = await ctx.db.prepare("SELECT id FROM agent_routines WHERE id = ? AND agent_id = ?").bind(id, agent.id).first(); |
| 307 | if (!existing) return fail("not_found", "There is no such routine."); |
| 308 | await ctx.db |
| 309 | .prepare( |
| 310 | `UPDATE agent_routines SET name = ?, instructions = ?, schedule = ?, events = ?, repos = ?, channel_id = ?, channel_name = ?, sponsor = ?, sponsor_username = ?, |
| 311 | enabled = ?, paused_note = NULL, next_run_at = ?, workspace = ?, updated_at = ? WHERE id = ?`, |
| 312 | ) |
| 313 | .bind(r.name, r.instructions, schedule, JSON.stringify(r.events), JSON.stringify(r.repos), r.channel_id, channelName, ctx.viewer.id, ctx.viewer.username, r.enabled !== false ? 1 : 0, next, ctx.slug, now.toISOString(), id) |
| 314 | .run(); |
| 315 | } else { |
| 316 | const count = await ctx.db.prepare("SELECT COUNT(*) AS n FROM agent_routines WHERE agent_id = ?").bind(agent.id).first<{ n: number }>(); |
| 317 | if ((count?.n ?? 0) >= MAX_ROUTINES) return fail("invalid", `An agent keeps at most ${MAX_ROUTINES} routines.`); |
| 318 | id = newRoutineId(); |
| 319 | await ctx.db |
| 320 | .prepare( |
| 321 | `INSERT INTO agent_routines (id, agent_id, workspace_id, workspace, name, instructions, schedule, events, repos, channel_id, channel_name, sponsor, sponsor_username, enabled, next_run_at, created_at, updated_at) |
| 322 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, |
| 323 | ) |
| 324 | .bind(id, agent.id, ctx.workspaceId, ctx.slug, r.name, r.instructions, schedule, JSON.stringify(r.events), JSON.stringify(r.repos), r.channel_id, channelName, ctx.viewer.id, ctx.viewer.username, r.enabled !== false ? 1 : 0, next, now.toISOString(), now.toISOString()) |
| 325 | .run(); |
| 326 | } |
| 327 | const row = await ctx.db.prepare("SELECT * FROM agent_routines WHERE id = ?").bind(id).first<RoutineRow>(); |
| 328 | return ok(toRoutine(row!)); |
| 329 | } |
| 330 | |
| 331 | function channelNameIn(sidebar: unknown, id: string): string | null { |
| 332 | const seen: unknown[] = [sidebar]; |
| 333 | while (seen.length) { |
| 334 | const value = seen.pop(); |
| 335 | if (Array.isArray(value)) seen.push(...value); |
| 336 | else if (value && typeof value === "object") { |
| 337 | const v = value as Record<string, unknown>; |
| 338 | if (v.id === id && typeof v.name === "string") return v.name; |
| 339 | seen.push(...Object.values(v)); |
| 340 | } |
| 341 | } |
| 342 | return null; |
| 343 | } |
| 344 | |
| 345 | export async function deleteRoutine(ctx: ViewContext, handle: string, id: string): Promise<Result<null>> { |
| 346 | if (!ctx.owner) return fail("forbidden", "Only the workspace's owners change routines."); |
| 347 | const agent = await agentByHandle(ctx, handle); |
| 348 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 349 | await ctx.db.prepare("DELETE FROM agent_routines WHERE id = ? AND agent_id = ?").bind(String(id ?? ""), agent.id).run(); |
| 350 | return ok(null); |
| 351 | } |
| 352 | |
| 353 | export async function runRoutineNow(ctx: ViewContext, handle: string, id: string): Promise<Result<AgentSession>> { |
| 354 | if (!ctx.owner) return fail("forbidden", "Only the workspace's owners run routines by hand."); |
| 355 | const agent = await agentByHandle(ctx, handle); |
| 356 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 357 | const routine = await ctx.db.prepare("SELECT * FROM agent_routines WHERE id = ? AND agent_id = ?").bind(String(id ?? ""), agent.id).first<RoutineRow>(); |
| 358 | if (!routine) return fail("not_found", "There is no such routine."); |
| 359 | const ran = await runRoutine(ctx.env, routine, agent, ctx.slug); |
| 360 | if (!ran.ok) return fail("invalid", ran.message); |
| 361 | const row = await sessionRow(ctx.db, ran.session); |
| 362 | return ok((await sessionsOut(ctx, [row!]))[0]); |
| 363 | } |
| 364 | |
| 365 | // ── Spend ───────────────────────────────────────────────────────────────── |
| 366 | |
| 367 | const KIND_LABELS: Record<string, string> = { reply: "Chat replies", chat: "Sessions", routine: "Routines", helper: "Helping colleagues", subagent: "Subagents" }; |
| 368 | |
| 369 | function slices(rows: { key: string | null; micros: number; n: number }[], label: (key: string) => string): SpendSlice[] { |
| 370 | return rows |
| 371 | .filter((r) => r.micros > 0 || r.n > 0) |
| 372 | .map((r) => ({ key: r.key ?? "", label: label(r.key ?? ""), micros: r.micros, count: r.n })) |
| 373 | .sort((a, b) => b.micros - a.micros); |
| 374 | } |
| 375 | |
| 376 | /** The most channels named in a breakdown; the rest of the channels fold into one slice. */ |
| 377 | const CHANNEL_SLICES = 12; |
| 378 | |
| 379 | /** |
| 380 | * Spend by where it was asked: the costliest channels by name, if the |
| 381 | * viewer can read them; direct messages together; channels they can't |
| 382 | * read, and channels past the first few, folded together without names. |
| 383 | */ |
| 384 | async function channelSlices(ctx: ViewContext, rows: { key: string | null; label: string | null; micros: number; n: number }[]): Promise<SpendSlice[]> { |
| 385 | const ranked = rows.filter((r) => r.micros > 0 || r.n > 0).sort((a, b) => b.micros - a.micros); |
| 386 | const named = ranked.filter((r) => r.key && r.key !== "dm").slice(0, CHANNEL_SLICES); |
| 387 | const can = await readable(ctx, named.map((r) => r.key!)); |
| 388 | const out: SpendSlice[] = []; |
| 389 | const hidden: SpendSlice = { key: "private", label: "Channels you're not in", micros: 0, count: 0 }; |
| 390 | const rest: SpendSlice = { key: "other", label: "Other channels", micros: 0, count: 0 }; |
| 391 | for (const row of ranked) { |
| 392 | if (row.key === "dm") out.push({ key: "dm", label: "Direct messages", micros: row.micros, count: row.n }); |
| 393 | else if (!named.includes(row)) (rest.micros += row.micros), (rest.count += row.n); |
| 394 | else if (can.has(row.key!)) out.push({ key: row.key!, label: row.label ? `#${row.label}` : "A channel", micros: row.micros, count: row.n }); |
| 395 | else (hidden.micros += row.micros), (hidden.count += row.n); |
| 396 | } |
| 397 | for (const folded of [hidden, rest]) if (folded.count > 0) out.push(folded); |
| 398 | return out.sort((a, b) => b.micros - a.micros); |
| 399 | } |
| 400 | |
| 401 | /** |
| 402 | * Budgets per person this month (person-budget.ts): owners see everyone's, |
| 403 | * anyone else only their own. |
| 404 | */ |
| 405 | export async function personBudgetsView(ctx: ViewContext): Promise<Result<PersonBudgets>> { |
| 406 | const now = new Date(); |
| 407 | const policyRow = await readPolicy(ctx.db, ctx.workspaceId, monthKey(now)); |
| 408 | const only = ctx.owner ? null : ctx.viewer.kind === "workspace" ? null : ctx.viewer.username.toLowerCase(); |
| 409 | return ok(await personBudgets(ctx.db, ctx.workspaceId, policyRow.person_monthly_micros, now, only)); |
| 410 | } |
| 411 | |
| 412 | /** Gives a person a budget of their own (0: none at all), or with null puts them back on the default. Owners only. */ |
| 413 | export async function setPersonBudget(ctx: ViewContext, username: unknown, monthly: unknown): Promise<Result<PersonBudgets>> { |
| 414 | if (!ctx.owner) return fail("forbidden", "Only the workspace's owners set budgets for people."); |
| 415 | const name = cleanUsername(username); |
| 416 | if (!name) return fail("invalid", "That isn't a username."); |
| 417 | if (monthly === null || monthly === undefined) { |
| 418 | await ctx.db.prepare("DELETE FROM person_budgets WHERE workspace_id = ? AND username = ?").bind(ctx.workspaceId, name).run(); |
| 419 | } else { |
| 420 | const value = Number(monthly); |
| 421 | if (!Number.isFinite(value) || value < 0 || value > 1_000_000_000_000) return fail("invalid", "A budget is a positive amount, or none."); |
| 422 | await ctx.db |
| 423 | .prepare( |
| 424 | `INSERT INTO person_budgets (workspace_id, username, monthly_micros, updated_by, updated_at) VALUES (?1, ?2, ?3, ?4, ?5) |
| 425 | ON CONFLICT (workspace_id, username) DO UPDATE SET monthly_micros = ?3, updated_by = ?4, updated_at = ?5`, |
| 426 | ) |
| 427 | .bind(ctx.workspaceId, name, Math.floor(value), ctx.viewer.username, new Date().toISOString()) |
| 428 | .run(); |
| 429 | } |
| 430 | return personBudgetsView(ctx); |
| 431 | } |
| 432 | |
| 433 | /** |
| 434 | * Where the spend went over a span of days (this month unless asked), for |
| 435 | * one agent (what it was paid for: its replies and every session it paid |
| 436 | * for, colleagues' help included) or for every agent; for everyone, or only |
| 437 | * the work one person asked for. |
| 438 | */ |
| 439 | export async function spend(ctx: ViewContext, handle: string | null, options: { period?: unknown; person?: unknown } = {}): Promise<Result<AgentSpendBreakdown>> { |
| 440 | const now = new Date(); |
| 441 | const span = spendSpan(options.period, now); |
| 442 | const from = `${span.from}T00:00:00.000Z`; |
| 443 | const end = new Date(`${span.until}T00:00:00.000Z`); |
| 444 | end.setUTCDate(end.getUTCDate() + 1); |
| 445 | const to = end.toISOString(); |
| 446 | const person = options.person == null || options.person === "" ? null : cleanUsername(options.person); |
| 447 | if (options.person != null && options.person !== "" && !person) return fail("invalid", "That isn't a username."); |
| 448 | let agentId: string | null = null; |
| 449 | if (handle) { |
| 450 | const agent = await agentByHandle(ctx, handle); |
| 451 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 452 | agentId = agent.id; |
| 453 | } |
| 454 | const db = ctx.db; |
| 455 | const who = person ? " AND asked_by_username = ?4" : ""; |
| 456 | const rFilter = `${agentId ? "agent_id = ?2" : "workspace_id = ?2"} AND created_at >= ?1 AND created_at < ?3${who}`; |
| 457 | const sFilter = `${agentId ? "payer_agent_id = ?2" : "workspace_id = ?2"} AND created_at >= ?1 AND created_at < ?3${who}`; |
| 458 | const binds: string[] = [from, agentId ?? ctx.workspaceId, to, ...(person ? [person] : [])]; |
| 459 | // Where it was asked: a channel by id; direct messages (a reply's channel has no name) together. |
| 460 | const union = `SELECT 'reply' AS kind, agent_id AS agent, asked_by_username AS person, model, charged_micros AS micros, substr(created_at, 1, 10) AS day, |
| 461 | CASE WHEN channel_name IS NULL THEN 'dm' ELSE channel_id END AS channel, channel_name AS channel_label FROM agent_replies WHERE ${rFilter} |
| 462 | UNION ALL SELECT kind, payer_agent_id AS agent, asked_by_username AS person, model, charged_micros AS micros, substr(created_at, 1, 10) AS day, |
| 463 | CASE WHEN channel_kind = 'dm' THEN 'dm' ELSE channel_id END AS channel, channel_name AS channel_label FROM agent_sessions WHERE ${sFilter} AND parent_id IS NULL |
| 464 | UNION ALL SELECT kind, payer_agent_id AS agent, asked_by_username AS person, model, 0 AS micros, substr(created_at, 1, 10) AS day, |
| 465 | CASE WHEN channel_kind = 'dm' THEN 'dm' ELSE channel_id END AS channel, channel_name AS channel_label FROM agent_sessions WHERE ${sFilter} AND parent_id IS NOT NULL`; |
| 466 | // A child's spend is already counted on its root (sessions.ts), so children add counts, not money. |
| 467 | const group = (column: string) => |
| 468 | db.prepare(`SELECT ${column} AS key, COALESCE(SUM(micros), 0) AS micros, COUNT(*) AS n FROM (${union}) GROUP BY ${column}`).bind(...binds).all<{ key: string | null; micros: number; n: number }>(); |
| 469 | const [byKind, byModel, byPerson, byAgent, byDay, byChannel, top, agents] = await Promise.all([ |
| 470 | group("kind"), |
| 471 | group("model"), |
| 472 | group("person"), |
| 473 | group("agent"), |
| 474 | group("day"), |
| 475 | db |
| 476 | .prepare(`SELECT channel AS key, MAX(channel_label) AS label, COALESCE(SUM(micros), 0) AS micros, COUNT(*) AS n FROM (${union}) GROUP BY channel`) |
| 477 | .bind(...binds) |
| 478 | .all<{ key: string | null; label: string | null; micros: number; n: number }>(), |
| 479 | db.prepare(`SELECT * FROM agent_sessions WHERE ${sFilter} AND parent_id IS NULL ORDER BY charged_micros DESC LIMIT 8`).bind(...binds).all<SessionRow>(), |
| 480 | faces(ctx), |
| 481 | ]); |
| 482 | const channels = await channelSlices(ctx, byChannel.results); |
| 483 | const byTeamMap = new Map<string, SpendSlice>(); |
| 484 | for (const row of byAgent.results) { |
| 485 | const face = agents.get(row.key ?? ""); |
| 486 | const team = face?.team || face?.department || "No team"; |
| 487 | const slice = byTeamMap.get(team) ?? { key: team, label: team, micros: 0, count: 0 }; |
| 488 | slice.micros += row.micros; |
| 489 | slice.count += row.n; |
| 490 | byTeamMap.set(team, slice); |
| 491 | } |
| 492 | const total = byKind.results.reduce((n, r) => n + r.micros, 0); |
| 493 | return ok({ |
| 494 | period: span.period, |
| 495 | span: span.span, |
| 496 | from: span.from, |
| 497 | until: span.until, |
| 498 | person, |
| 499 | total_micros: total, |
| 500 | by_kind: slices(byKind.results, (k) => KIND_LABELS[k] ?? k), |
| 501 | by_model: slices(byModel.results, (k) => k || "No model"), |
| 502 | by_person: slices(byPerson.results, (k) => (k ? `@${k}` : "Routines and agents")), |
| 503 | by_agent: slices(byAgent.results, (k) => { |
| 504 | const face = agents.get(k); |
| 505 | return face ? `${face.display_name} (@${face.handle})` : "An archived agent"; |
| 506 | }), |
| 507 | by_team: [...byTeamMap.values()].sort((a, b) => b.micros - a.micros), |
| 508 | by_channel: channels, |
| 509 | top_sessions: await sessionsOut(ctx, top.results, agents), |
| 510 | days: byDay.results.map((r) => ({ day: r.key ?? "", micros: r.micros })).sort((a, b) => a.day.localeCompare(b.day)), |
| 511 | }); |
| 512 | } |
| 513 | |
| 514 | // ── Activity and versions ──────────────────────────────────────────────── |
| 515 | |
| 516 | export async function activity(ctx: ViewContext, handle: string): Promise<Result<AgentActivity[]>> { |
| 517 | const agent = await agentByHandle(ctx, handle); |
| 518 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 519 | const [replies, sessions] = await Promise.all([ |
| 520 | ctx.db |
| 521 | .prepare( |
| 522 | "SELECT id, status, channel_id, channel_name, asked_by_username, model, tool_count, charged_micros, created_at, reply_id FROM agent_replies WHERE agent_id = ? ORDER BY created_at DESC LIMIT 60", |
| 523 | ) |
| 524 | .bind(agent.id) |
| 525 | .all<{ id: string; status: string; channel_id: string; channel_name: string | null; asked_by_username: string | null; model: string | null; tool_count: number | null; charged_micros: number; created_at: string; reply_id: string | null }>(), |
| 526 | ctx.db.prepare("SELECT * FROM agent_sessions WHERE agent_id = ? ORDER BY created_at DESC LIMIT 40").bind(agent.id).all<SessionRow>(), |
| 527 | ]); |
| 528 | const can = await readable(ctx, [...replies.results.map((r) => r.channel_id), ...sessions.results.map((s) => s.channel_id)]); |
| 529 | const items: AgentActivity[] = [ |
| 530 | ...replies.results.map((r) => { |
| 531 | const visible = can.has(r.channel_id); |
| 532 | return { |
| 533 | id: r.id, |
| 534 | kind: "reply" as const, |
| 535 | status: r.status, |
| 536 | channel_id: r.channel_id, |
| 537 | channel_name: visible ? r.channel_name : null, |
| 538 | title: null, |
| 539 | asked_by_username: visible ? r.asked_by_username : null, |
| 540 | model: r.model, |
| 541 | tools: r.tool_count ?? 0, |
| 542 | charged_micros: r.charged_micros, |
| 543 | created_at: r.created_at, |
| 544 | visible, |
| 545 | ref: visible ? r.reply_id : null, |
| 546 | }; |
| 547 | }), |
| 548 | ...sessions.results.map((s) => { |
| 549 | const visible = can.has(s.channel_id); |
| 550 | return { |
| 551 | id: s.id, |
| 552 | kind: "session" as const, |
| 553 | status: s.status, |
| 554 | channel_id: s.channel_id, |
| 555 | channel_name: visible ? s.channel_name : null, |
| 556 | title: visible ? s.title : null, |
| 557 | asked_by_username: visible ? s.asked_by_username : null, |
| 558 | model: s.model, |
| 559 | tools: s.tool_calls, |
| 560 | charged_micros: s.charged_micros, |
| 561 | created_at: s.created_at, |
| 562 | visible, |
| 563 | ref: s.id, |
| 564 | }; |
| 565 | }), |
| 566 | ]; |
| 567 | return ok(items.sort((a, b) => b.created_at.localeCompare(a.created_at)).slice(0, 80)); |
| 568 | } |
| 569 | |
| 570 | export async function versions(ctx: ViewContext, handle: string): Promise<Result<AgentVersion[]>> { |
| 571 | const agent = await agentByHandle(ctx, handle); |
| 572 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); |
| 573 | const rows = await ctx.db |
| 574 | .prepare("SELECT version, definition, changed_by, created_at FROM agent_versions WHERE agent_id = ? ORDER BY version DESC LIMIT 50") |
| 575 | .bind(agent.id) |
| 576 | .all<{ version: number; definition: string; changed_by: string; created_at: string }>(); |
| 577 | return ok( |
| 578 | rows.results.map((r) => { |
| 579 | let definition = {}; |
| 580 | try { |
| 581 | definition = JSON.parse(r.definition); |
| 582 | } catch { |
| 583 | // An unreadable old version shows as empty. |
| 584 | } |
| 585 | return { version: r.version, changed_by: r.changed_by, created_at: r.created_at, definition }; |
| 586 | }), |
| 587 | ); |
| 588 | } |
| 589 | |
| 590 | // ── Policy and overview ────────────────────────────────────────────────── |
| 591 | |
| 592 | export async function policy(ctx: ViewContext): Promise<Result<AgentPolicy>> { |
| 593 | const row = await readPolicy(ctx.db, ctx.workspaceId, monthKey(new Date())); |
| 594 | return ok(policyOf(row)); |
| 595 | } |
| 596 | |
| 597 | function policyOf(row: AgentPolicy): AgentPolicy { |
| 598 | return { |
| 599 | monthly_micros: row.monthly_micros, |
| 600 | default_agent_monthly_micros: row.default_agent_monthly_micros, |
| 601 | default_session_micros: row.default_session_micros, |
| 602 | person_monthly_micros: row.person_monthly_micros, |
| 603 | }; |
| 604 | } |
| 605 | |
| 606 | export async function setPolicy(ctx: ViewContext, changes: Partial<AgentPolicy>): Promise<Result<AgentPolicy>> { |
| 607 | if (!ctx.owner) return fail("forbidden", "Only the workspace's owners set its agents' budget."); |
| 608 | const current = await policy(ctx); |
| 609 | const checked = checkPolicy(current.ok ? current.value : DEFAULT_POLICY, changes ?? {}); |
| 610 | if (!checked.ok) return fail("invalid", checked.message); |
| 611 | const p = checked.value; |
| 612 | await ctx.db |
| 613 | .prepare( |
| 614 | `INSERT INTO agent_policies (workspace_id, monthly_micros, default_agent_monthly_micros, default_session_micros, person_monthly_micros, updated_by, updated_at) |
| 615 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) |
| 616 | ON CONFLICT (workspace_id) DO UPDATE SET monthly_micros = ?2, default_agent_monthly_micros = ?3, default_session_micros = ?4, person_monthly_micros = ?5, updated_by = ?6, updated_at = ?7`, |
| 617 | ) |
| 618 | .bind(ctx.workspaceId, p.monthly_micros, p.default_agent_monthly_micros, p.default_session_micros, p.person_monthly_micros, ctx.viewer.username, new Date().toISOString()) |
| 619 | .run(); |
| 620 | return ok(p); |
| 621 | } |
| 622 | |
| 623 | export async function overview(ctx: ViewContext): Promise<Result<AgentsOverview>> { |
| 624 | const now = new Date(); |
| 625 | const db = ctx.db; |
| 626 | const [rows, policyRow, live, recent, upcoming, breakdown, agentFaces] = await Promise.all([ |
| 627 | db.prepare(`${selectAgents("a.workspace_id = ?3 AND a.archived_at IS NULL")} ORDER BY a.builtin DESC, a.handle`).bind(...periods(now), ctx.workspaceId).all<Row>(), |
| 628 | readPolicy(db, ctx.workspaceId, monthKey(now)), |
| 629 | db |
| 630 | .prepare(`SELECT * FROM agent_sessions WHERE workspace_id = ? AND status IN (${LIVE.map(() => "?").join(", ")}) ORDER BY created_at DESC LIMIT 60`) |
| 631 | .bind(ctx.workspaceId, ...LIVE) |
| 632 | .all<SessionRow>(), |
| 633 | db |
| 634 | .prepare(`SELECT * FROM agent_sessions WHERE workspace_id = ? AND parent_id IS NULL AND status IN ('done','failed','stopped') ORDER BY finished_at DESC LIMIT 12`) |
| 635 | .bind(ctx.workspaceId) |
| 636 | .all<SessionRow>(), |
| 637 | db.prepare("SELECT * FROM agent_routines WHERE workspace_id = ? AND enabled = 1 AND next_run_at IS NOT NULL ORDER BY next_run_at LIMIT 6").bind(ctx.workspaceId).all<RoutineRow>(), |
| 638 | spend(ctx, null), |
| 639 | faces(ctx), |
| 640 | ]); |
| 641 | const agents: WorkspaceAgent[] = rows.results.map((row) => toAgent(row, now)); |
| 642 | const liveSessions = await sessionsOut(ctx, live.results, agentFaces); |
| 643 | const liveByAgent: Record<string, number> = {}; |
| 644 | for (const s of live.results) liveByAgent[s.agent_id] = (liveByAgent[s.agent_id] ?? 0) + 1; |
| 645 | const level = policyRow.monthly_micros ? [100, 90, 75].find((l) => (policyRow.spent * 100) / policyRow.monthly_micros! >= l) ?? null : null; |
| 646 | return ok({ |
| 647 | policy: policyOf(policyRow), |
| 648 | spent_month_micros: policyRow.spent, |
| 649 | alert: level, |
| 650 | agents, |
| 651 | live_by_agent: liveByAgent, |
| 652 | live: liveSessions.filter((s) => s.visible && !s.parent_id).slice(0, 20), |
| 653 | waiting_on_you: ctx.owner ? liveSessions.filter((s) => s.status === "needs_approval") : liveSessions.filter((s) => s.status === "needs_approval" && s.visible && s.asked_by === ctx.viewer.id), |
| 654 | recent: (await sessionsOut(ctx, recent.results, agentFaces)).filter((s) => s.visible).slice(0, 8), |
| 655 | upcoming: upcoming.results.map((r) => { |
| 656 | const face = agentFaces.get(r.agent_id); |
| 657 | return { ...toRoutine(r), agent_handle: face?.handle ?? "agent", agent_name: face?.display_name ?? "An agent" }; |
| 658 | }), |
| 659 | spend: breakdown.ok ? breakdown.value : emptySpend(now), |
| 660 | can_manage: ctx.owner, |
| 661 | }); |
| 662 | } |
| 663 | |
| 664 | /** A breakdown with nothing in it, for this month. */ |
| 665 | function emptySpend(now: Date): AgentSpendBreakdown { |
| 666 | const span = spendSpan("month", now); |
| 667 | return { ...span, person: null, total_micros: 0, by_kind: [], by_model: [], by_person: [], by_agent: [], by_team: [], by_channel: [], top_sessions: [], days: [] }; |
| 668 | } |