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