| 1 | /** |
| 2 | * Sessions (docs.g1t.sh/guides/agent-sessions/): the work an agent spins off |
| 3 | * from a conversation, a routine's run, or its part in another session. |
| 4 | * |
| 5 | * A conversation with an agent is never one long context. Replies read the |
| 6 | * last few messages, what the agent remembers, and its recent sessions in |
| 7 | * that conversation. Real work happens in a session: |
| 8 | * |
| 9 | * - **Bounded.** A session has a goal, its own working context of turns |
| 10 | * (compacted as it grows: the goal, a summary of earlier steps, and the |
| 11 | * latest turns), a step limit and a spend cap. |
| 12 | * - **Visible.** It posts a live card where it was asked, updates the card |
| 13 | * in place as it works, posts progress in the card's thread, and reports |
| 14 | * back in the conversation when done. Its page has the full transcript. |
| 15 | * - **Steerable.** A reply in the card's thread, or from its page, reaches |
| 16 | * it at its next step, or wakes it again once it is done. |
| 17 | * - **A tree.** It can hand parts to its subagents or bring colleagues in; |
| 18 | * each is a child session whose result comes back to it. Everything in |
| 19 | * the tree is paid by the agent at the root, within the root's cap. |
| 20 | * - **Stoppable.** Anyone who can see it can stop it and everything under it. |
| 21 | * |
| 22 | * Each step runs on the agent's desk (desk.ts) through `metered` (meter.ts), |
| 23 | * so a step is billed, capped and recorded exactly as a reply is. |
| 24 | */ |
| 25 | import { |
| 26 | type AgentRef, |
| 27 | type AgentSession, |
| 28 | type AgentSessionKind, |
| 29 | type AgentSessionStatus, |
| 30 | type AskerAccess, |
| 31 | type MessageCard, |
| 32 | type ModelTier, |
| 33 | type ServiceBinding, |
| 34 | type SessionEvent, |
| 35 | type SessionOutput, |
| 36 | type SubagentDef, |
| 37 | type User, |
| 38 | chatClient, |
| 39 | identityClient, |
| 40 | agentRef, |
| 41 | newId, |
| 42 | workClient, |
| 43 | } from "@g1t/contracts"; |
| 44 | |
| 45 | import { CHAT_MAX_HOPS } from "../../../packages/contracts/src/chat.ts"; |
| 46 | import { Audience } from "./audience.ts"; |
| 47 | import { type MeterEnv, metered } from "./meter.ts"; |
| 48 | import { type RecallPlace, MAX_FACTS, cleanFact, memorySection, recall, scopeFor } from "./memory.ts"; |
| 49 | import { readPolicy } from "./policy.ts"; |
| 50 | import { type PortsEnv, audiencePorts, toolPorts } from "./ports.ts"; |
| 51 | import { systemPrompt } from "./prompt.ts"; |
| 52 | import { conversationFrom } from "./surface.ts"; |
| 53 | import { type Row, definitionOf, periods } from "./store.ts"; |
| 54 | import { type ActionPorts, type ToolCall, ToolBox } from "./tools.ts"; |
| 55 | import { type ModelMessage, SESSION_LIMITS, runTurn } from "./turn.ts"; |
| 56 | import { recallQuery, recallSection } from "./recall.ts"; |
| 57 | import { rosterLines } from "./orchestrator.ts"; |
| 58 | import { dollars } from "./money.ts"; |
| 59 | import { postDraft } from "./cards.ts"; |
| 60 | import { sessionActions } from "./card-views.ts"; |
| 61 | import type { Desk } from "./desk.ts"; |
| 62 | |
| 63 | export type SessionEnv = MeterEnv & |
| 64 | PortsEnv & { |
| 65 | CHAT: ServiceBinding; |
| 66 | IDENTITY: ServiceBinding; |
| 67 | WORK: ServiceBinding; |
| 68 | NOTIFY?: ServiceBinding; |
| 69 | DESKS: DurableObjectNamespace<Desk>; |
| 70 | }; |
| 71 | |
| 72 | /** Steps one session takes at most before it must report. */ |
| 73 | export const MAX_STEPS = 8; |
| 74 | /** Children one session may have running at once. */ |
| 75 | export const MAX_CHILDREN = 4; |
| 76 | /** How deep a tree of sessions may go. */ |
| 77 | export const MAX_DEPTH = 3; |
| 78 | /** Past this many characters of working context, it is compacted. */ |
| 79 | const CONTEXT_LIMIT = 120_000; |
| 80 | /** Turns kept whole when compacting. */ |
| 81 | const KEEP_TURNS = 6; |
| 82 | /** The longest report posted in chat. */ |
| 83 | const MAX_REPORT = 12_000; |
| 84 | |
| 85 | export const LIVE: AgentSessionStatus[] = ["queued", "working", "waiting", "needs_approval"]; |
| 86 | const OVER: AgentSessionStatus[] = ["done", "failed", "stopped"]; |
| 87 | |
| 88 | export type SessionRow = { |
| 89 | id: string; |
| 90 | workspace_id: string; |
| 91 | agent_id: string; |
| 92 | subagent: string | null; |
| 93 | kind: string; |
| 94 | parent_id: string | null; |
| 95 | root_id: string; |
| 96 | payer_agent_id: string; |
| 97 | title: string; |
| 98 | goal: string; |
| 99 | status: string; |
| 100 | status_note: string | null; |
| 101 | summary: string | null; |
| 102 | workspace: string; |
| 103 | channel_id: string; |
| 104 | channel_kind: string; |
| 105 | channel_name: string | null; |
| 106 | thread_root: string | null; |
| 107 | message_id: string | null; |
| 108 | card_message_id: string | null; |
| 109 | asked_by: string | null; |
| 110 | asked_by_username: string | null; |
| 111 | asker: string | null; |
| 112 | routine_id: string | null; |
| 113 | chain: string; |
| 114 | hops: number; |
| 115 | context: string; |
| 116 | inbox: string; |
| 117 | steps: number; |
| 118 | tool_calls: number; |
| 119 | input_tokens: number; |
| 120 | output_tokens: number; |
| 121 | cost_micros: number; |
| 122 | charged_micros: number; |
| 123 | cap_micros: number | null; |
| 124 | model: string | null; |
| 125 | outputs: string; |
| 126 | step_started_at: string | null; |
| 127 | created_at: string; |
| 128 | updated_at: string; |
| 129 | finished_at: string | null; |
| 130 | }; |
| 131 | |
| 132 | /** One turn of a session's working context: plain text, never tool blocks. */ |
| 133 | type Turn = { role: "user" | "assistant"; content: string }; |
| 134 | /** Something that arrived for a session while it worked. */ |
| 135 | type Inbound = { kind: "steer" | "child"; by: string; body: string }; |
| 136 | |
| 137 | function json<T>(raw: string | null | undefined, fallback: T): T { |
| 138 | if (!raw) return fallback; |
| 139 | try { |
| 140 | return JSON.parse(raw) as T; |
| 141 | } catch { |
| 142 | return fallback; |
| 143 | } |
| 144 | } |
| 145 | |
| 146 | const iso = () => new Date().toISOString(); |
| 147 | |
| 148 | /** A session as the contract shows it; `visible` false hides what it was about. */ |
| 149 | export function toSession(row: SessionRow, agent: { handle: string; display_name: string; avatar_seed: string } | null, visible: boolean): AgentSession { |
| 150 | return { |
| 151 | id: row.id, |
| 152 | workspace_id: row.workspace_id, |
| 153 | agent_id: row.agent_id, |
| 154 | agent_handle: agent?.handle ?? "agent", |
| 155 | agent_name: agent?.display_name ?? "An agent", |
| 156 | agent_avatar_seed: agent?.avatar_seed ?? agent?.handle ?? row.agent_id, |
| 157 | subagent: row.subagent, |
| 158 | kind: row.kind as AgentSessionKind, |
| 159 | parent_id: row.parent_id, |
| 160 | root_id: row.root_id, |
| 161 | payer_agent_id: row.payer_agent_id, |
| 162 | title: visible ? row.title : "A private session", |
| 163 | goal: visible ? row.goal : "", |
| 164 | status: row.status as AgentSessionStatus, |
| 165 | status_note: visible ? row.status_note : null, |
| 166 | summary: visible ? row.summary : null, |
| 167 | channel_id: row.channel_id, |
| 168 | channel_kind: row.channel_kind === "dm" ? "dm" : "channel", |
| 169 | channel_name: visible ? row.channel_name : null, |
| 170 | card_message_id: visible ? row.card_message_id : null, |
| 171 | asked_by: row.asked_by, |
| 172 | asked_by_username: visible ? row.asked_by_username : null, |
| 173 | routine_id: row.routine_id, |
| 174 | steps: row.steps, |
| 175 | tool_calls: row.tool_calls, |
| 176 | input_tokens: row.input_tokens, |
| 177 | output_tokens: row.output_tokens, |
| 178 | charged_micros: row.charged_micros, |
| 179 | cost_micros: row.cost_micros, |
| 180 | cap_micros: row.cap_micros, |
| 181 | model: row.model, |
| 182 | outputs: visible ? json<SessionOutput[]>(row.outputs, []) : [], |
| 183 | created_at: row.created_at, |
| 184 | updated_at: row.updated_at, |
| 185 | finished_at: row.finished_at, |
| 186 | visible, |
| 187 | }; |
| 188 | } |
| 189 | |
| 190 | /** A session's card, as its conversation shows it. */ |
| 191 | export function cardFor( |
| 192 | row: Pick<SessionRow, "id" | "title" | "status" | "steps" | "tool_calls" | "charged_micros" | "status_note"> & Partial<Pick<SessionRow, "cap_micros" | "summary" | "goal">>, |
| 193 | slug: string, |
| 194 | handle: string, |
| 195 | children = 0, |
| 196 | ): MessageCard { |
| 197 | const state: Record<string, string> = { |
| 198 | queued: "Queued", |
| 199 | working: "Working", |
| 200 | waiting: children === 1 ? "Waiting on a helper" : "Waiting on helpers", |
| 201 | needs_approval: "Needs approval", |
| 202 | done: "Done", |
| 203 | failed: "Failed", |
| 204 | stopped: "Stopped", |
| 205 | }; |
| 206 | const parts = [ |
| 207 | row.steps ? `Step ${row.steps}` : null, |
| 208 | row.tool_calls ? `${row.tool_calls} tool${row.tool_calls === 1 ? "" : "s"}` : null, |
| 209 | row.charged_micros ? dollars(row.charged_micros) : null, |
| 210 | ].filter(Boolean); |
| 211 | const note = row.status === "needs_approval" || row.status === "failed" || row.status === "stopped" ? row.status_note : null; |
| 212 | const href = `/${slug}/-/agents/${handle}/sessions/${row.id}`; |
| 213 | return { |
| 214 | kind: "session", |
| 215 | title: row.title, |
| 216 | detail: [parts.join(" · ") || "Starting", note].filter(Boolean).join(" — ").slice(0, 480), |
| 217 | state: state[row.status] ?? row.status, |
| 218 | href, |
| 219 | ...(row.status === "done" && row.summary ? { body: row.summary.length > 600 ? `${row.summary.slice(0, 600)}…` : row.summary } : {}), |
| 220 | fields: [ |
| 221 | ...(row.cap_micros ? [{ label: "Spent", value: `${dollars(row.charged_micros)} of ${dollars(row.cap_micros)}` }] : []), |
| 222 | ...(children ? [{ label: "Helpers", value: `${children} working` }] : []), |
| 223 | ], |
| 224 | actions: sessionActions(row.status, row.cap_micros ?? null, row.charged_micros, href), |
| 225 | owner: "agents", |
| 226 | ref: row.id, |
| 227 | }; |
| 228 | } |
| 229 | |
| 230 | /** Appends to a session's transcript. */ |
| 231 | export function eventStatement(db: D1Database, id: string, kind: SessionEvent["kind"], by: string | null, body: string, tool: string | null = null, outcome: string | null = null): D1PreparedStatement { |
| 232 | return db |
| 233 | .prepare( |
| 234 | `INSERT INTO agent_session_events (session_id, seq, kind, by_name, body, tool, outcome, created_at) |
| 235 | SELECT ?1, COALESCE(MAX(seq), 0) + 1, ?2, ?3, ?4, ?5, ?6, ?7 FROM agent_session_events WHERE session_id = ?1`, |
| 236 | ) |
| 237 | .bind(id, kind, by, body.slice(0, 20_000), tool, outcome, iso()); |
| 238 | } |
| 239 | |
| 240 | /** |
| 241 | * The working context, kept bounded: the goal, then a summary of what is |
| 242 | * cut, then the latest turns whole. What is cut stays in the transcript. |
| 243 | */ |
| 244 | export function compact(turns: Turn[], limit = CONTEXT_LIMIT, keep = KEEP_TURNS): Turn[] { |
| 245 | const size = (list: Turn[]) => list.reduce((n, t) => n + t.content.length, 0); |
| 246 | if (size(turns) <= limit || turns.length <= keep + 1) return turns; |
| 247 | const [goal, ...rest] = turns; |
| 248 | let tail = rest.slice(-keep); |
| 249 | // The kept part starts with someone else's turn, as the model needs. |
| 250 | while (tail.length && tail[0].role === "assistant") tail = tail.slice(1); |
| 251 | const cut = rest.slice(0, rest.length - tail.length); |
| 252 | const notes = cut |
| 253 | .filter((t) => t.role === "assistant") |
| 254 | .map((t, i) => `- Step ${i + 1}: ${t.content.replace(/\s+/g, " ").slice(0, 600)}`) |
| 255 | .join("\n"); |
| 256 | const earlier: Turn = { role: "user", content: `${goal.content}\n\n(Earlier in this session, now summarised:\n${notes || "- nothing to note"})` }; |
| 257 | return [earlier, ...tail]; |
| 258 | } |
| 259 | |
| 260 | /** Turns as the Messages API takes them: alternating, someone else's first. */ |
| 261 | function alternate(turns: Turn[]): ModelMessage[] { |
| 262 | const out: Turn[] = []; |
| 263 | for (const turn of turns) { |
| 264 | const last = out[out.length - 1]; |
| 265 | if (last && last.role === turn.role) last.content += `\n\n${turn.content}`; |
| 266 | else out.push({ ...turn }); |
| 267 | } |
| 268 | while (out.length && out[0].role === "assistant") out.shift(); |
| 269 | if (out.length && out[out.length - 1].role === "assistant") out.push({ role: "user", content: "(Go on with the session.)" }); |
| 270 | return out; |
| 271 | } |
| 272 | |
| 273 | async function agentRow(db: D1Database, id: string): Promise<Row | null> { |
| 274 | return db.prepare("SELECT * FROM agents WHERE id = ?").bind(id).first<Row>(); |
| 275 | } |
| 276 | |
| 277 | export async function sessionRow(db: D1Database, id: string): Promise<SessionRow | null> { |
| 278 | return db.prepare("SELECT * FROM agent_sessions WHERE id = ?").bind(id).first<SessionRow>(); |
| 279 | } |
| 280 | |
| 281 | /** Hands a session to its agent's desk to work its next step. */ |
| 282 | export async function wake(env: Pick<SessionEnv, "DESKS">, agentId: string, sessionId: string): Promise<void> { |
| 283 | await env.DESKS.get(env.DESKS.idFromName(agentId)).session(sessionId, agentId); |
| 284 | } |
| 285 | |
| 286 | export type NewSession = { |
| 287 | agent: Row; |
| 288 | kind: AgentSessionKind; |
| 289 | subagent?: SubagentDef | null; |
| 290 | parent?: SessionRow | null; |
| 291 | title: string; |
| 292 | goal: string; |
| 293 | workspace: string; |
| 294 | channel_id: string; |
| 295 | channel_kind: "channel" | "dm"; |
| 296 | channel_name: string | null; |
| 297 | thread_root: string | null; |
| 298 | message_id: string | null; |
| 299 | asked_by: string | null; |
| 300 | asked_by_username: string | null; |
| 301 | asker: AskerAccess | null; |
| 302 | routine_id?: string | null; |
| 303 | chain: string[]; |
| 304 | hops: number; |
| 305 | }; |
| 306 | |
| 307 | /** |
| 308 | * Starts a session: its row, its card where it was asked (a child's card |
| 309 | * goes in its parent's thread), and its first step on the desk. A child's |
| 310 | * cap is what its root has left; a root's is the agent's per-task cap or |
| 311 | * the workspace's default for sessions, whichever is lower. |
| 312 | */ |
| 313 | export async function startSession(env: SessionEnv, input: NewSession): Promise<SessionRow> { |
| 314 | const db = env.DB; |
| 315 | const now = iso(); |
| 316 | const id = newId("asn"); |
| 317 | const parent = input.parent ?? null; |
| 318 | const root = parent ? ((await sessionRow(db, parent.root_id)) ?? parent) : null; |
| 319 | let cap: number | null; |
| 320 | if (root) { |
| 321 | const tree = await db.prepare("SELECT COALESCE(SUM(charged_micros), 0) AS spent FROM agent_sessions WHERE root_id = ?").bind(root.id).first<{ spent: number }>(); |
| 322 | cap = root.cap_micros != null ? Math.max(1, root.cap_micros - (tree?.spent ?? 0)) : null; |
| 323 | } else { |
| 324 | const policy = await readPolicy(db, input.agent.workspace_id, periods(new Date())[0]); |
| 325 | const task = definitionOf(input.agent).budget.task_micros; |
| 326 | cap = Math.min(policy.default_session_micros, task && task > 0 ? task : Number.POSITIVE_INFINITY); |
| 327 | } |
| 328 | const subagent = input.subagent ?? null; |
| 329 | const goal = [ |
| 330 | input.goal, |
| 331 | subagent ? `\n(You are working as ${input.agent.display_name}'s subagent "${subagent.name}": ${subagent.description}\n\n${subagent.instructions})` : "", |
| 332 | ].join(""); |
| 333 | const context: Turn[] = [{ role: "user", content: `Your session: ${input.title}\n\n${goal}` }]; |
| 334 | await db.batch([ |
| 335 | db |
| 336 | .prepare( |
| 337 | `INSERT INTO agent_sessions (id, workspace_id, agent_id, subagent, kind, parent_id, root_id, payer_agent_id, title, goal, status, |
| 338 | workspace, channel_id, channel_kind, channel_name, thread_root, message_id, asked_by, asked_by_username, asker, routine_id, |
| 339 | chain, hops, context, cap_micros, created_at, updated_at) |
| 340 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, |
| 341 | ) |
| 342 | .bind( |
| 343 | id, |
| 344 | input.agent.workspace_id, |
| 345 | input.agent.id, |
| 346 | subagent?.name ?? null, |
| 347 | input.kind, |
| 348 | parent?.id ?? null, |
| 349 | root?.id ?? id, |
| 350 | root?.payer_agent_id ?? input.agent.id, |
| 351 | input.title.slice(0, 120), |
| 352 | input.goal.slice(0, 8000), |
| 353 | input.workspace.toLowerCase(), |
| 354 | input.channel_id, |
| 355 | input.channel_kind, |
| 356 | input.channel_name, |
| 357 | input.thread_root, |
| 358 | input.message_id, |
| 359 | input.asked_by, |
| 360 | input.asked_by_username, |
| 361 | input.asker ? JSON.stringify(input.asker) : null, |
| 362 | input.routine_id ?? null, |
| 363 | JSON.stringify(input.chain), |
| 364 | input.hops, |
| 365 | JSON.stringify(context), |
| 366 | cap === Number.POSITIVE_INFINITY ? null : cap, |
| 367 | now, |
| 368 | now, |
| 369 | ), |
| 370 | eventStatement(db, id, "goal", input.asked_by_username, `${input.title}\n\n${input.goal}`), |
| 371 | ]); |
| 372 | let row = (await sessionRow(db, id))!; |
| 373 | if (parent) { |
| 374 | await addOutput(db, parent.id, { kind: "session", id, agent_handle: subagent ? `${input.agent.handle}/${subagent.name}` : input.agent.handle, title: row.title }); |
| 375 | await db.batch([eventStatement(db, parent.id, "child", input.agent.handle, `${subagent ? `Subagent ${subagent.name}` : `@${input.agent.handle}`} started: ${row.title}`)]); |
| 376 | } else { |
| 377 | // A root session's card, where it was asked. |
| 378 | const posted = await chatClient(env.CHAT) |
| 379 | .postAsAgent(row.workspace, row.channel_id, row.agent_id, { |
| 380 | body: "", |
| 381 | card: cardFor(row, row.workspace, input.agent.handle), |
| 382 | thread_root: row.thread_root, |
| 383 | hops: row.hops, |
| 384 | asked_by: row.asked_by, |
| 385 | asker: input.asker, |
| 386 | chain: input.chain, |
| 387 | }) |
| 388 | .catch(() => null); |
| 389 | if (posted?.ok) { |
| 390 | await db.prepare("UPDATE agent_sessions SET card_message_id = ? WHERE id = ?").bind(posted.value.id, id).run(); |
| 391 | row = { ...row, card_message_id: posted.value.id }; |
| 392 | } |
| 393 | } |
| 394 | await wake(env, row.agent_id, id); |
| 395 | return row; |
| 396 | } |
| 397 | |
| 398 | async function addOutput(db: D1Database, id: string, output: SessionOutput): Promise<void> { |
| 399 | const row = await db.prepare("SELECT outputs FROM agent_sessions WHERE id = ?").bind(id).first<{ outputs: string }>(); |
| 400 | const list = json<SessionOutput[]>(row?.outputs, []); |
| 401 | list.push(output); |
| 402 | await db.prepare("UPDATE agent_sessions SET outputs = ? WHERE id = ?").bind(JSON.stringify(list.slice(-50)), id).run(); |
| 403 | } |
| 404 | |
| 405 | /** The root session of a tree, whose card and agent speak for it in chat. */ |
| 406 | async function speaker(db: D1Database, row: SessionRow): Promise<{ root: SessionRow; agent: Row | null }> { |
| 407 | const root = row.root_id === row.id ? row : ((await sessionRow(db, row.root_id)) ?? row); |
| 408 | return { root, agent: await agentRow(db, root.agent_id) }; |
| 409 | } |
| 410 | |
| 411 | /** Brings the root's card up to date with the tree. Never throws. */ |
| 412 | export async function refreshCard(env: SessionEnv, row: SessionRow): Promise<void> { |
| 413 | try { |
| 414 | const db = env.DB; |
| 415 | const { root, agent } = await speaker(db, row); |
| 416 | if (!root.card_message_id || !agent) return; |
| 417 | const fresh = (await sessionRow(db, root.id)) ?? root; |
| 418 | const tree = await db |
| 419 | .prepare("SELECT COUNT(*) AS n, COALESCE(SUM(charged_micros), 0) AS spent, COALESCE(SUM(tool_calls), 0) AS tools, SUM(CASE WHEN status IN ('queued','working','waiting') AND id <> root_id THEN 1 ELSE 0 END) AS live FROM agent_sessions WHERE root_id = ?") |
| 420 | .bind(root.id) |
| 421 | .first<{ n: number; spent: number; tools: number; live: number }>(); |
| 422 | const shown = { ...fresh, charged_micros: tree?.spent ?? fresh.charged_micros, tool_calls: tree?.tools ?? fresh.tool_calls }; |
| 423 | await chatClient(env.CHAT).updateAsAgent(fresh.workspace, fresh.channel_id, fresh.agent_id, root.card_message_id, { card: cardFor(shown, fresh.workspace, agent.handle, tree?.live ?? 0) }); |
| 424 | } catch (error) { |
| 425 | console.error("agents: a session card was not updated", row.id, String(error)); |
| 426 | } |
| 427 | } |
| 428 | |
| 429 | /** Posts in the root card's thread, as the root's agent; a child's note names who it is from. */ |
| 430 | async function postInThread(env: SessionEnv, row: SessionRow, by: Row, text: string): Promise<boolean> { |
| 431 | const db = env.DB; |
| 432 | const { root } = await speaker(db, row); |
| 433 | if (!root.card_message_id) return false; |
| 434 | const prefix = root.id === row.id ? "" : `**${row.subagent ? `${by.display_name} · ${row.subagent}` : by.display_name}:** `; |
| 435 | const posted = await chatClient(env.CHAT) |
| 436 | .postAsAgent(root.workspace, root.channel_id, root.agent_id, { |
| 437 | body: `${prefix}${text}`.slice(0, 8000), |
| 438 | thread_root: root.card_message_id, |
| 439 | hops: root.hops, |
| 440 | asked_by: root.asked_by, |
| 441 | asker: json<AskerAccess | null>(root.asker, null), |
| 442 | chain: json<string[]>(root.chain, []), |
| 443 | }) |
| 444 | .catch(() => null); |
| 445 | return !!posted?.ok; |
| 446 | } |
| 447 | |
| 448 | /** Posts a card (a draft issue) in the root card's thread, as the root's agent; its id. */ |
| 449 | async function postCardInThread(env: SessionEnv, row: SessionRow, by: Row, card: MessageCard): Promise<string | null> { |
| 450 | const { root } = await speaker(env.DB, row); |
| 451 | const posted = await chatClient(env.CHAT) |
| 452 | .postAsAgent(root.workspace, root.channel_id, root.agent_id, { |
| 453 | body: root.id === row.id ? "" : `**${by.display_name}** drafted this:`, |
| 454 | card, |
| 455 | thread_root: root.card_message_id ?? root.thread_root, |
| 456 | hops: root.hops, |
| 457 | asked_by: root.asked_by, |
| 458 | asker: json<AskerAccess | null>(root.asker, null), |
| 459 | chain: json<string[]>(root.chain, []), |
| 460 | }) |
| 461 | .catch(() => null); |
| 462 | return posted?.ok ? posted.value.id : null; |
| 463 | } |
| 464 | |
| 465 | /** Sets a session's status, records why, and brings its card along. */ |
| 466 | async function setStatus(env: SessionEnv, row: SessionRow, status: AgentSessionStatus, note: string | null, extra: Record<string, string | number | null> = {}): Promise<SessionRow> { |
| 467 | const db = env.DB; |
| 468 | const names = Object.keys(extra); |
| 469 | const finished = OVER.includes(status) ? iso() : null; |
| 470 | await db |
| 471 | .prepare( |
| 472 | `UPDATE agent_sessions SET status = ?, status_note = ?, updated_at = ?, finished_at = COALESCE(?, finished_at)${names.map((n) => `, ${n} = ?`).join("")} WHERE id = ?`, |
| 473 | ) |
| 474 | .bind(status, note, iso(), finished, ...names.map((n) => extra[n]), row.id) |
| 475 | .run(); |
| 476 | const fresh = (await sessionRow(db, row.id))!; |
| 477 | await refreshCard(env, fresh); |
| 478 | return fresh; |
| 479 | } |
| 480 | |
| 481 | /** The person who asked, resolved, for acting on their behalf. */ |
| 482 | async function askerUser(env: SessionEnv, row: SessionRow): Promise<User | null> { |
| 483 | if (!row.asked_by) return null; |
| 484 | const [user] = await identityClient(env.IDENTITY) |
| 485 | .usersForAudience([row.asked_by]) |
| 486 | .catch(() => [] as User[]); |
| 487 | return user ?? null; |
| 488 | } |
| 489 | |
| 490 | /** |
| 491 | * What an agent may do here: remember and forget within where it is, |
| 492 | * file issues as the person who asked, and (in a session) post updates, |
| 493 | * use subagents and bring colleagues in. Shared by replies and sessions. |
| 494 | */ |
| 495 | export function actionPorts( |
| 496 | env: SessionEnv, |
| 497 | input: { |
| 498 | agent: Row; |
| 499 | place: RecallPlace; |
| 500 | source: { kind: "message" | "session"; ref: string; label: string; channel_id: string }; |
| 501 | asker: { id: string | null; username: string | null }; |
| 502 | workspace: string; |
| 503 | session?: SessionRow | null; |
| 504 | /** From a reply: starts a session for the conversation. */ |
| 505 | spinOff?: (title: string, goal: string) => Promise<{ ok: boolean; message: string }>; |
| 506 | /** Posts a card where this work reports (a draft issue); its message id, or null. */ |
| 507 | postCard: (card: MessageCard) => Promise<string | null>; |
| 508 | }, |
| 509 | ): ActionPorts { |
| 510 | const db = env.DB; |
| 511 | const { agent, place } = input; |
| 512 | const session = input.session ?? null; |
| 513 | const ports: ActionPorts = { |
| 514 | async remember(body, wanted, onlyForAsker) { |
| 515 | const fact = cleanFact(body); |
| 516 | if (!fact) return { ok: false, message: "Say what to remember." }; |
| 517 | const count = await db.prepare("SELECT COUNT(*) AS n FROM agent_memories WHERE agent_id = ?").bind(agent.id).first<{ n: number }>(); |
| 518 | if ((count?.n ?? 0) >= MAX_FACTS) return { ok: false, message: "Your memory is full. Forget something out of date first." }; |
| 519 | const privately = !!onlyForAsker && !!input.asker.id; |
| 520 | const { scope, ref } = scopeFor(place, wanted, privately ? input.asker.id : null); |
| 521 | const id = newId("mem"); |
| 522 | const now = iso(); |
| 523 | const label = scope === "person" ? input.asker.username : scope === "channel" ? input.source.label : null; |
| 524 | await db |
| 525 | .prepare( |
| 526 | `INSERT INTO agent_memories (id, agent_id, workspace_id, scope, scope_ref, scope_label, body, source_kind, source_ref, source_label, source_channel_id, created_by, created_by_kind, created_at, updated_at) |
| 527 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'agent', ?, ?)`, |
| 528 | ) |
| 529 | .bind(id, agent.id, agent.workspace_id, scope, ref, label, fact, input.source.kind, input.source.ref, input.source.label, input.source.channel_id, agent.handle, now, now) |
| 530 | .run(); |
| 531 | if (session) await addOutput(db, session.id, { kind: "memory", id, body: fact }); |
| 532 | const where = scope === "workspace" ? "for the whole workspace" : scope === "person" ? "for this person" : "for this conversation"; |
| 533 | const narrowed = privately |
| 534 | ? " (only for them: you read an artifact not everyone in the workspace can)" |
| 535 | : wanted && wanted !== scope |
| 536 | ? ` (${wanted} wasn't allowed from here)` |
| 537 | : ""; |
| 538 | return { ok: true, message: `Remembered ${where}${narrowed}: ${fact}` }; |
| 539 | }, |
| 540 | async forget(id) { |
| 541 | const row = await db.prepare("SELECT scope, scope_ref FROM agent_memories WHERE id = ? AND agent_id = ?").bind(id, agent.id).first<{ scope: string; scope_ref: string }>(); |
| 542 | // Only what could be recalled here can be forgotten from here. |
| 543 | const here = row && (row.scope === "workspace" ? place.kind === "public" : row.scope === "channel" ? row.scope_ref === place.channel_id : place.kind === "dm" && place.people.length === 1 && place.people[0] === row.scope_ref); |
| 544 | if (!row || !here) return { ok: false, message: "There is no such note you can forget here." }; |
| 545 | await db.prepare("DELETE FROM agent_memories WHERE id = ?").bind(id).run(); |
| 546 | return { ok: true, message: "Forgotten." }; |
| 547 | }, |
| 548 | async comment(repo, asker, number, body) { |
| 549 | const made = await workClient(env.WORK).workspaceAgentComment({ namespace: repo.namespace, name: repo.name }, number, refOf(agent), asker, body); |
| 550 | if (!made.ok) return { ok: false, message: `It couldn't be posted: ${made.error.message}` }; |
| 551 | return { ok: true, message: `Commented on ${repo.namespace}/${repo.name}#${number}.` }; |
| 552 | }, |
| 553 | async review(repo, asker, number, verdict, body) { |
| 554 | const made = await workClient(env.WORK).workspaceAgentReview({ namespace: repo.namespace, name: repo.name }, number, refOf(agent), asker, verdict, body); |
| 555 | if (!made.ok) return { ok: false, message: `The review couldn't be posted: ${made.error.message}` }; |
| 556 | const what = verdict === "approve" ? "Approved" : verdict === "request_changes" ? "Requested changes on" : "Reviewed"; |
| 557 | return { ok: true, message: `${what} ${repo.namespace}/${repo.name}#${number} (advisory). Link it in your report: /${repo.namespace}/${repo.name}/pull/${number}` }; |
| 558 | }, |
| 559 | async draftIssue(repo, issue) { |
| 560 | const draft = await postDraft( |
| 561 | env, |
| 562 | { |
| 563 | agent_id: agent.id, |
| 564 | workspace_id: agent.workspace_id, |
| 565 | workspace: input.workspace, |
| 566 | channel_id: input.source.channel_id, |
| 567 | session_id: session?.id ?? null, |
| 568 | repo_id: repo.id, |
| 569 | repo: `${repo.namespace}/${repo.name}`, |
| 570 | title: issue.title, |
| 571 | body: issue.body, |
| 572 | labels: issue.labels, |
| 573 | asked_by: input.asker.id, |
| 574 | }, |
| 575 | input.postCard, |
| 576 | ); |
| 577 | if (!draft) return { ok: false, message: "The draft couldn't be posted; give it in your answer instead." }; |
| 578 | return { ok: true, message: "The draft is in the conversation as a card with File issue and Discard. Tell them in a sentence; don't repeat it." }; |
| 579 | }, |
| 580 | }; |
| 581 | if (input.spinOff) ports.startSession = input.spinOff; |
| 582 | if (session) { |
| 583 | ports.postUpdate = async (text) => { |
| 584 | const posted = await postInThread(env, session, agent, text); |
| 585 | if (posted) await db.batch([eventStatement(db, session.id, "update", agent.handle, text)]); |
| 586 | return posted ? { ok: true, message: "Posted." } : { ok: false, message: "It couldn't be posted; carry on." }; |
| 587 | }; |
| 588 | const child = async (target: Row, subagent: SubagentDef | null, brief: string) => { |
| 589 | const depth = await treeDepth(db, session); |
| 590 | if (depth >= MAX_DEPTH) return { ok: false, message: "This work is already deep enough; do this part yourself." }; |
| 591 | const live = await db |
| 592 | .prepare("SELECT COUNT(*) AS n FROM agent_sessions WHERE parent_id = ? AND status IN ('queued','working','waiting','needs_approval')") |
| 593 | .bind(session.id) |
| 594 | .first<{ n: number }>(); |
| 595 | if ((live?.n ?? 0) >= MAX_CHILDREN) return { ok: false, message: `You already have ${MAX_CHILDREN} helpers working; wait for them.` }; |
| 596 | const title = brief.split("\n")[0].slice(0, 100) || "Helping"; |
| 597 | await startSession(env, { |
| 598 | agent: target, |
| 599 | kind: subagent ? "subagent" : "helper", |
| 600 | subagent, |
| 601 | parent: session, |
| 602 | title, |
| 603 | goal: `${agent.display_name} (@${agent.handle}) asked for your help with part of the session "${session.title}".\n\n${brief}\n\nWhen you're done, answer with your result for ${agent.display_name}: findings, links, and anything left open.`, |
| 604 | workspace: session.workspace, |
| 605 | channel_id: session.channel_id, |
| 606 | channel_kind: session.channel_kind === "dm" ? "dm" : "channel", |
| 607 | channel_name: session.channel_name, |
| 608 | thread_root: session.thread_root, |
| 609 | message_id: session.message_id, |
| 610 | asked_by: session.asked_by, |
| 611 | asked_by_username: session.asked_by_username, |
| 612 | asker: json<AskerAccess | null>(session.asker, null), |
| 613 | chain: [...json<string[]>(session.chain, []), agent.id], |
| 614 | hops: session.hops + 1, |
| 615 | }); |
| 616 | return { ok: true, message: `${subagent ? `Your subagent ${subagent.name}` : `@${target.handle}`} is on it. End this step with what you're waiting for; their result comes back to you before your next step.` }; |
| 617 | }; |
| 618 | ports.useSubagent = async (name, brief) => { |
| 619 | const subagent = definitionOf(agent).subagents.find((s) => s.name === name); |
| 620 | if (!subagent) return { ok: false, message: `You have no subagent called ${name}.` }; |
| 621 | return child(agent, subagent, brief); |
| 622 | }; |
| 623 | ports.bringIn = async (handle, brief) => { |
| 624 | const colleague = await db |
| 625 | .prepare("SELECT * FROM agents WHERE workspace_id = ? AND handle = ? AND archived_at IS NULL") |
| 626 | .bind(agent.workspace_id, handle) |
| 627 | .first<Row>(); |
| 628 | if (!colleague || colleague.id === agent.id) return { ok: false, message: `There is no other agent called @${handle} here.` }; |
| 629 | if (json<string[]>(session.chain, []).includes(colleague.id)) return { ok: false, message: `@${handle} is already part of this work.` }; |
| 630 | if (session.hops + 1 > CHAT_MAX_HOPS) return { ok: false, message: "This work has been passed along too many times; do it yourself." }; |
| 631 | return child(colleague, null, brief); |
| 632 | }; |
| 633 | } |
| 634 | return ports; |
| 635 | } |
| 636 | |
| 637 | /** How work names an agent on issues and pull requests. */ |
| 638 | function refOf(agent: Row): AgentRef { |
| 639 | return agentRef({ id: agent.id, handle: agent.handle, display_name: agent.display_name, avatar_seed: agent.avatar_seed || agent.handle }); |
| 640 | } |
| 641 | |
| 642 | async function treeDepth(db: D1Database, row: SessionRow): Promise<number> { |
| 643 | let depth = 0; |
| 644 | let at: SessionRow | null = row; |
| 645 | while (at?.parent_id && depth < 10) { |
| 646 | depth++; |
| 647 | at = await sessionRow(db, at.parent_id); |
| 648 | } |
| 649 | return depth; |
| 650 | } |
| 651 | |
| 652 | /** The section of the system prompt that says what a session is and how to finish. */ |
| 653 | function sessionSection(row: SessionRow, asker: string): string { |
| 654 | const report = |
| 655 | row.kind === "helper" || row.kind === "subagent" |
| 656 | ? "Your final answer goes back to the agent who asked for your help, not into chat." |
| 657 | : row.kind === "routine" |
| 658 | ? `Your final answer is posted in ${row.channel_kind === "dm" ? "the direct message" : `#${row.channel_name ?? "the channel"}`} as this routine's report.` |
| 659 | : `Your final answer is posted for ${asker} in the conversation where they asked.`; |
| 660 | return [ |
| 661 | "## This session", |
| 662 | "", |
| 663 | `You are working a session: "${row.title}". It is bounded: work through it with your tools, step by step, and finish within ${MAX_STEPS} steps.`, |
| 664 | `- ${report} Make it the report: what you found or did, with links (issues, files, threads), and anything left open.`, |
| 665 | "- Use post_update for real milestones or a question for the people following, not for every step.", |
| 666 | "- When part of the work belongs to a subagent or a colleague, hand it over with use_subagent or bring_in and end your step saying what you're waiting for; their results come back to you.", |
| 667 | "- Never claim to have done or checked something you didn't. If you can't do something from here, say so in the report.", |
| 668 | ].join("\n"); |
| 669 | } |
| 670 | |
| 671 | /** |
| 672 | * Works one step of a session, on its agent's desk. Reads what arrived |
| 673 | * (steering, helpers' results), runs one metered model turn with the |
| 674 | * session's tools, and decides what comes next: done, waiting on helpers, |
| 675 | * another step, or stopped at a limit. Never throws. |
| 676 | */ |
| 677 | export async function advance(env: SessionEnv, id: string): Promise<void> { |
| 678 | const db = env.DB; |
| 679 | let row = await sessionRow(db, id); |
| 680 | if (!row || OVER.includes(row.status as AgentSessionStatus) || row.status === "needs_approval") return; |
| 681 | const agent = await agentRow(db, row.agent_id); |
| 682 | const payer = await agentRow(db, row.payer_agent_id); |
| 683 | if (!agent || !payer || agent.archived_at) { |
| 684 | await setStatus(env, row, "failed", "Its agent was archived."); |
| 685 | return finished(env, row.id); |
| 686 | } |
| 687 | // Waiting on helpers: only once every child is over. |
| 688 | const pending = await db |
| 689 | .prepare("SELECT COUNT(*) AS n FROM agent_sessions WHERE parent_id = ? AND status IN ('queued','working','waiting','needs_approval')") |
| 690 | .bind(row.id) |
| 691 | .first<{ n: number }>(); |
| 692 | if ((pending?.n ?? 0) > 0) { |
| 693 | if (row.status !== "waiting") await setStatus(env, row, "waiting", null); |
| 694 | return; |
| 695 | } |
| 696 | if (row.steps >= MAX_STEPS && json<Inbound[]>(row.inbox, []).length === 0) { |
| 697 | await setStatus(env, row, "done", null, { summary: row.summary ?? "Stopped at the step limit." }); |
| 698 | return finished(env, row.id); |
| 699 | } |
| 700 | // The session's cap, counting its whole tree for a root. |
| 701 | if (row.cap_micros != null && row.charged_micros >= row.cap_micros) { |
| 702 | await setStatus(env, row, "needs_approval", `Reached its cap of ${dollars(row.cap_micros)}.`); |
| 703 | await notifyApproval(env, row, agent); |
| 704 | return; |
| 705 | } |
| 706 | |
| 707 | // What arrived meanwhile becomes the next turn. |
| 708 | const inbox = json<Inbound[]>(row.inbox, []); |
| 709 | let context = json<Turn[]>(row.context, []); |
| 710 | if (inbox.length) { |
| 711 | const lines = inbox.map((item) => (item.kind === "steer" ? `@${item.by} says: ${item.body}` : `Result from ${item.by}:\n${item.body}`)); |
| 712 | context.push({ role: "user", content: lines.join("\n\n") }); |
| 713 | } else if (row.steps > 0) { |
| 714 | context.push({ role: "user", content: "(Go on with the session.)" }); |
| 715 | } |
| 716 | context = compact(context); |
| 717 | const startedAt = iso(); |
| 718 | await db |
| 719 | .prepare("UPDATE agent_sessions SET status = 'working', status_note = NULL, inbox = '[]', context = ?, step_started_at = ?, updated_at = ? WHERE id = ? AND status <> 'stopped'") |
| 720 | .bind(JSON.stringify(context), startedAt, startedAt, row.id) |
| 721 | .run(); |
| 722 | row = (await sessionRow(db, id))!; |
| 723 | if (row.status === "stopped") return; |
| 724 | await refreshCard(env, row); |
| 725 | |
| 726 | const slug = row.workspace; |
| 727 | const definition = definitionOf(agent); |
| 728 | const subagent = row.subagent ? definition.subagents.find((s) => s.name === row!.subagent) ?? null : null; |
| 729 | const events: D1PreparedStatement[] = []; |
| 730 | const calls: ToolCall[] = []; |
| 731 | const startTier: ModelTier = "large"; |
| 732 | const asker = { id: row.asked_by, username: row.asked_by_username }; |
| 733 | const current = row; |
| 734 | |
| 735 | const outcome = await metered( |
| 736 | env, |
| 737 | { |
| 738 | row: agent, |
| 739 | payer, |
| 740 | slug, |
| 741 | task: "session", |
| 742 | start: startTier, |
| 743 | askerName: row.asked_by_username, |
| 744 | person: row.asked_by_username, |
| 745 | leftMicros: row.cap_micros != null ? row.cap_micros - row.charged_micros : null, |
| 746 | limits: subagent ? subagent.routing : null, |
| 747 | }, |
| 748 | async (model) => { |
| 749 | // The audience: who reads what this session posts. Without one it reads nothing but its own context. |
| 750 | let toolbox: ToolBox | null = null; |
| 751 | let place: RecallPlace = { channel_id: current.channel_id, kind: current.channel_kind === "dm" ? "dm" : "private", people: current.asked_by ? [current.asked_by] : [] }; |
| 752 | try { |
| 753 | if (current.asked_by) { |
| 754 | const audience = await Audience.build(slug, current.asked_by, audiencePorts(env, slug, current.channel_id)); |
| 755 | place = { channel_id: current.channel_id, kind: audience.kind, people: audience.shared ? (current.asked_by ? [current.asked_by] : []) : audience.members.map((m) => m.id) }; |
| 756 | const noConsult: ToolPorts["consult"] = async () => ({ ok: false, message: "In a session, bring a colleague in with bring_in instead." }); |
| 757 | const sourceLabel = current.channel_kind === "dm" ? "a direct message" : `#${current.channel_name ?? "a channel"}`; |
| 758 | toolbox = new ToolBox( |
| 759 | audience, |
| 760 | toolPorts(env, slug, agent.workspace_id, current.channel_id, noConsult, agent.id), |
| 761 | { |
| 762 | agentId: agent.id, |
| 763 | notConsult: [agent.handle], |
| 764 | hops: current.hops, |
| 765 | maxHops: CHAT_MAX_HOPS, |
| 766 | session: true, |
| 767 | onCall: (call) => { |
| 768 | calls.push(call); |
| 769 | events.push(eventStatement(db, current.id, "tool", agent.handle, call.args, call.tool, call.outcome)); |
| 770 | }, |
| 771 | }, |
| 772 | [], |
| 773 | actionPorts(env, { |
| 774 | agent, |
| 775 | place, |
| 776 | source: { kind: "session", ref: current.id, label: `the session "${current.title}" in ${sourceLabel}`, channel_id: current.channel_id }, |
| 777 | asker, |
| 778 | workspace: slug, |
| 779 | session: current, |
| 780 | postCard: (card) => postCardInThread(env, current, agent, card), |
| 781 | }), |
| 782 | ); |
| 783 | } |
| 784 | } catch (error) { |
| 785 | console.error("agents: no audience for a session step, so no tools", current.id, String(error)); |
| 786 | } |
| 787 | // What the workspace's artifacts say about the work: its goal, and whatever arrived for this step. |
| 788 | const asked = [current.goal, ...inbox.map((item) => item.body)].reverse(); |
| 789 | const [facts, passages] = await Promise.all([ |
| 790 | recall(db, agent.id, place).catch(() => []), |
| 791 | toolbox ? toolbox.recall(recallQuery(asked, 800), definition.reading ?? []) : Promise.resolve([]), |
| 792 | ]); |
| 793 | const [team, here] = await Promise.all([ |
| 794 | db |
| 795 | .prepare("SELECT handle, display_name, role, title, team, department, responsibilities FROM agents WHERE workspace_id = ? AND archived_at IS NULL AND id <> ? ORDER BY builtin DESC, handle LIMIT 50") |
| 796 | .bind(agent.workspace_id, agent.id) |
| 797 | .all<{ handle: string; display_name: string; role: string; title: string; team: string | null; department: string; responsibilities: string }>(), |
| 798 | // Who reads what this session posts: said every step, as in a reply. A helper may not be a member: then not said. |
| 799 | chatClient(env.CHAT) |
| 800 | .conversationForAgent(slug, current.channel_id, agent.id, current.asked_by) |
| 801 | .then((found) => (found.ok ? conversationFrom(found.value) : null)) |
| 802 | .catch(() => null), |
| 803 | ]); |
| 804 | const roster = rosterLines( |
| 805 | team.results.map((a) => ({ |
| 806 | handle: a.handle, |
| 807 | display_name: a.display_name, |
| 808 | role: a.role, |
| 809 | title: a.title, |
| 810 | team: a.team, |
| 811 | department: a.department, |
| 812 | responsibilities: json<string[]>(a.responsibilities, []), |
| 813 | status: "idle", |
| 814 | spent_month_micros: 0, |
| 815 | monthly_micros: null, |
| 816 | })), |
| 817 | ); |
| 818 | const access = json<AskerAccess | null>(current.asker, null); |
| 819 | const system = [ |
| 820 | systemPrompt({ |
| 821 | agent: { ...definition, id: agent.id }, |
| 822 | workspace: slug, |
| 823 | channel: { kind: current.channel_kind === "dm" ? "dm" : "channel", name: current.channel_name }, |
| 824 | asker: { name: current.asked_by_username ?? "someone", display_name: null, access }, |
| 825 | today: new Date(), |
| 826 | tools: toolbox ? { code: toolbox.definitions().some((tool) => tool.name === "read_file") } : null, |
| 827 | colleagues: roster, |
| 828 | session: true, |
| 829 | conversation: here, |
| 830 | }), |
| 831 | sessionSection(current, current.asked_by_username ? `@${current.asked_by_username}` : "the person who asked"), |
| 832 | memorySection(facts), |
| 833 | recallSection(passages), |
| 834 | ] |
| 835 | .filter(Boolean) |
| 836 | .join("\n\n"); |
| 837 | const result = await runTurn(model.send, { |
| 838 | model: model.model.model, |
| 839 | system, |
| 840 | messages: alternate(context), |
| 841 | tools: toolbox, |
| 842 | price: model.ownModel ? null : model.model.price, |
| 843 | maxRounds: SESSION_LIMITS.rounds, |
| 844 | inputBudget: SESSION_LIMITS.input, |
| 845 | maxOutput: SESSION_LIMITS.output, |
| 846 | onText: (text) => events.push(eventStatement(db, current.id, "text", agent.handle, text)), |
| 847 | stopped: async () => (await db.prepare("SELECT status FROM agent_sessions WHERE id = ?").bind(current.id).first<{ status: string }>())?.status === "stopped", |
| 848 | }); |
| 849 | return { ...result, cost: model.ownModel ? 0 : result.cost }; |
| 850 | }, |
| 851 | ).catch((error: unknown) => ({ ok: false as const, reason: "error", message: error instanceof Error ? error.message : String(error) })); |
| 852 | |
| 853 | if (events.length) await db.batch(events).catch((error: unknown) => console.error("agents: transcript not written", id, String(error))); |
| 854 | row = (await sessionRow(db, id))!; |
| 855 | |
| 856 | if (!outcome.ok) { |
| 857 | if (outcome.reason === "error") { |
| 858 | console.error("agents: a session step failed", id, outcome.message); |
| 859 | await db.batch([eventStatement(db, id, "note", null, `This step failed: ${outcome.message}`)]); |
| 860 | await setStatus(env, row, "failed", "Something went wrong on g1t's side."); |
| 861 | } else { |
| 862 | await db.batch([eventStatement(db, id, "note", null, outcome.message)]); |
| 863 | await setStatus(env, row, "stopped", outcome.message); |
| 864 | } |
| 865 | return finished(env, id); |
| 866 | } |
| 867 | |
| 868 | const answer = outcome.value; |
| 869 | const tokens = outcome.tokens; |
| 870 | context.push({ role: "assistant", content: answer.text || "(no text)" }); |
| 871 | await db |
| 872 | .prepare( |
| 873 | `UPDATE agent_sessions SET steps = steps + 1, tool_calls = tool_calls + ?, input_tokens = input_tokens + ?, output_tokens = output_tokens + ?, |
| 874 | cost_micros = cost_micros + ?, charged_micros = charged_micros + ?, model = ?, context = ?, step_started_at = NULL, updated_at = ? WHERE id = ?`, |
| 875 | ) |
| 876 | .bind( |
| 877 | calls.length, |
| 878 | tokens.input + tokens.cacheRead + tokens.cacheWrite, |
| 879 | tokens.output, |
| 880 | outcome.cost, |
| 881 | outcome.charged, |
| 882 | outcome.model, |
| 883 | JSON.stringify(context), |
| 884 | iso(), |
| 885 | id, |
| 886 | ) |
| 887 | .run(); |
| 888 | // A root's spend counts its tree for the cap: children add theirs to it too. |
| 889 | if (row.root_id !== row.id) { |
| 890 | await db.prepare("UPDATE agent_sessions SET charged_micros = charged_micros + ? WHERE id = ?").bind(outcome.charged, row.root_id).run(); |
| 891 | } |
| 892 | row = (await sessionRow(db, id))!; |
| 893 | if (row.status === "stopped" || answer.stopped) return finished(env, id); |
| 894 | |
| 895 | const children = await db |
| 896 | .prepare("SELECT COUNT(*) AS n FROM agent_sessions WHERE parent_id = ? AND status IN ('queued','working','waiting','needs_approval')") |
| 897 | .bind(id) |
| 898 | .first<{ n: number }>(); |
| 899 | if ((children?.n ?? 0) > 0) { |
| 900 | if (answer.text) await db.batch([eventStatement(db, id, "text", agent.handle, answer.text)]); |
| 901 | await setStatus(env, row, "waiting", null); |
| 902 | return; |
| 903 | } |
| 904 | // Something arrived during the step: another step reads it. |
| 905 | if (json<Inbound[]>(row.inbox, []).length && row.steps < MAX_STEPS + 2) { |
| 906 | if (answer.text) await db.batch([eventStatement(db, id, "text", agent.handle, answer.text)]); |
| 907 | await wake(env, row.agent_id, id); |
| 908 | return; |
| 909 | } |
| 910 | const report = answer.text.trim() || "I finished without anything to report."; |
| 911 | await db.batch([eventStatement(db, id, "result", agent.handle, report)]); |
| 912 | row = await setStatus(env, row, "done", null, { summary: report.slice(0, MAX_REPORT) }); |
| 913 | await finished(env, id); |
| 914 | } |
| 915 | |
| 916 | /** |
| 917 | * After a session is over: a root reports in its conversation; a child |
| 918 | * hands its result to its parent and wakes it once its siblings are done. |
| 919 | */ |
| 920 | async function finished(env: SessionEnv, id: string): Promise<void> { |
| 921 | const db = env.DB; |
| 922 | const row = await sessionRow(db, id); |
| 923 | if (!row) return; |
| 924 | // Everything under a stopped or failed session stops too. |
| 925 | if (row.status === "stopped" || row.status === "failed") await stopChildren(env, row.id, "Its parent session ended."); |
| 926 | const agent = await agentRow(db, row.agent_id); |
| 927 | if (row.parent_id) { |
| 928 | const parent = await sessionRow(db, row.parent_id); |
| 929 | if (!parent || OVER.includes(parent.status as AgentSessionStatus)) return; |
| 930 | const who = row.subagent ? `your subagent ${row.subagent}` : `@${agent?.handle ?? "a colleague"}`; |
| 931 | const body = row.status === "done" ? (row.summary ?? "(no result)") : `They couldn't finish (${row.status}): ${row.status_note ?? "no reason given"}.`; |
| 932 | await pushInbox(db, parent.id, { kind: "child", by: who, body: body.slice(0, 8000) }); |
| 933 | await db.batch([eventStatement(db, parent.id, "child", agent?.handle ?? null, `${who} ${row.status === "done" ? "finished" : row.status}: ${row.title}`)]); |
| 934 | await wake(env, parent.agent_id, parent.id); |
| 935 | return; |
| 936 | } |
| 937 | // A root's report, where it was asked: the request's thread, or the conversation. |
| 938 | if (row.status === "done" && row.summary && agent) { |
| 939 | const mention = row.kind === "chat" && row.asked_by_username ? `@${row.asked_by_username} ` : ""; |
| 940 | await chatClient(env.CHAT) |
| 941 | .postAsAgent(row.workspace, row.channel_id, row.agent_id, { |
| 942 | body: `${mention}${row.summary}`.slice(0, MAX_REPORT), |
| 943 | thread_root: row.thread_root, |
| 944 | hops: row.hops, |
| 945 | asked_by: row.asked_by, |
| 946 | asker: json<AskerAccess | null>(row.asker, null), |
| 947 | chain: json<string[]>(row.chain, []), |
| 948 | }) |
| 949 | .catch((error: unknown) => console.error("agents: a session's report was not posted", row.id, String(error))); |
| 950 | } else if ((row.status === "stopped" || row.status === "failed") && agent && row.status_note) { |
| 951 | await postInThread(env, row, agent, `${row.status === "failed" ? "This session failed" : "This session stopped"}: ${row.status_note}`); |
| 952 | } |
| 953 | await refreshCard(env, row); |
| 954 | } |
| 955 | |
| 956 | async function pushInbox(db: D1Database, id: string, item: Inbound): Promise<void> { |
| 957 | const row = await db.prepare("SELECT inbox FROM agent_sessions WHERE id = ?").bind(id).first<{ inbox: string }>(); |
| 958 | const list = json<Inbound[]>(row?.inbox, []); |
| 959 | list.push(item); |
| 960 | await db.prepare("UPDATE agent_sessions SET inbox = ?, updated_at = ? WHERE id = ?").bind(JSON.stringify(list.slice(-20)), iso(), id).run(); |
| 961 | } |
| 962 | |
| 963 | /** Stops every live session under `id`. */ |
| 964 | async function stopChildren(env: SessionEnv, id: string, note: string): Promise<void> { |
| 965 | const db = env.DB; |
| 966 | const children = await db |
| 967 | .prepare("SELECT * FROM agent_sessions WHERE parent_id = ? AND status IN ('queued','working','waiting','needs_approval')") |
| 968 | .bind(id) |
| 969 | .all<SessionRow>(); |
| 970 | for (const child of children.results) { |
| 971 | await db.prepare("UPDATE agent_sessions SET status = 'stopped', status_note = ?, finished_at = ?, updated_at = ? WHERE id = ?").bind(note, iso(), iso(), child.id).run(); |
| 972 | await stopChildren(env, child.id, note); |
| 973 | } |
| 974 | } |
| 975 | |
| 976 | /** Stops a session and everything under it, by a person. */ |
| 977 | export async function stop(env: SessionEnv, row: SessionRow, by: string): Promise<SessionRow> { |
| 978 | const fresh = await setStatus(env, row, "stopped", `Stopped by @${by}.`); |
| 979 | await env.DB.batch([eventStatement(env.DB, row.id, "note", null, `Stopped by @${by}.`)]); |
| 980 | await stopChildren(env, row.id, `Stopped by @${by}.`); |
| 981 | if (row.parent_id) await finished(env, row.id); |
| 982 | else await refreshCard(env, fresh); |
| 983 | return fresh; |
| 984 | } |
| 985 | |
| 986 | /** A person's message to a session: read at its next step; a finished root goes on again. */ |
| 987 | export async function steer(env: SessionEnv, row: SessionRow, by: string, body: string): Promise<SessionRow> { |
| 988 | const db = env.DB; |
| 989 | await pushInbox(db, row.id, { kind: "steer", by, body: body.slice(0, 4000) }); |
| 990 | await db.batch([eventStatement(db, row.id, "steer", by, body)]); |
| 991 | if (row.status === "working" || row.status === "waiting" || row.status === "queued") { |
| 992 | if (row.status !== "working") await wake(env, row.agent_id, row.id); |
| 993 | return (await sessionRow(db, row.id))!; |
| 994 | } |
| 995 | // Over (or at its cap): it picks up again with its context, a fresh set of steps. |
| 996 | if (row.status === "needs_approval") return (await sessionRow(db, row.id))!; |
| 997 | await db.prepare("UPDATE agent_sessions SET status = 'queued', status_note = NULL, steps = MIN(steps, ?), finished_at = NULL, updated_at = ? WHERE id = ?").bind(Math.max(0, MAX_STEPS - 3), iso(), row.id).run(); |
| 998 | const fresh = (await sessionRow(db, row.id))!; |
| 999 | await refreshCard(env, fresh); |
| 1000 | await wake(env, row.agent_id, row.id); |
| 1001 | return fresh; |
| 1002 | } |
| 1003 | |
| 1004 | /** Raises a session's cap past what it has spent and lets it go on. */ |
| 1005 | export async function approve(env: SessionEnv, row: SessionRow, by: string, capMicros: number): Promise<SessionRow> { |
| 1006 | const db = env.DB; |
| 1007 | await db.prepare("UPDATE agent_sessions SET cap_micros = ?, status = 'queued', status_note = NULL, updated_at = ? WHERE id = ?").bind(Math.floor(capMicros), iso(), row.id).run(); |
| 1008 | await db.batch([eventStatement(db, row.id, "note", null, `@${by} raised its cap to ${dollars(capMicros)}.`)]); |
| 1009 | const fresh = (await sessionRow(db, row.id))!; |
| 1010 | await refreshCard(env, fresh); |
| 1011 | await wake(env, row.agent_id, row.id); |
| 1012 | return fresh; |
| 1013 | } |
| 1014 | |
| 1015 | /** Tells whoever asked, and the agent's maker, that a session waits for more budget. */ |
| 1016 | async function notifyApproval(env: SessionEnv, row: SessionRow, agent: Row): Promise<void> { |
| 1017 | if (!env.NOTIFY) return; |
| 1018 | // The card's own buttons ride along, so it can be approved from the notification. |
| 1019 | const { root } = await speaker(env.DB, row); |
| 1020 | const href = `/${row.workspace}/-/agents/${agent.handle}/sessions/${row.id}`; |
| 1021 | const card = root.card_message_id |
| 1022 | ? { channel_id: root.channel_id, message_id: root.card_message_id, actions: sessionActions("needs_approval", row.cap_micros, row.charged_micros, href) } |
| 1023 | : null; |
| 1024 | const targets = new Set<string>(); |
| 1025 | if (row.asked_by_username) targets.add(row.asked_by_username); |
| 1026 | if (agent.created_by) targets.add(agent.created_by); |
| 1027 | for (const username of targets) { |
| 1028 | await env.NOTIFY.fetch("https://service/rpc/notify", { |
| 1029 | method: "POST", |
| 1030 | headers: { "content-type": "application/json" }, |
| 1031 | body: JSON.stringify({ |
| 1032 | target: { username }, |
| 1033 | notification: { |
| 1034 | id: `approval:${row.id}:${row.cap_micros ?? 0}`, |
| 1035 | kind: "approval", |
| 1036 | workspace: row.workspace, |
| 1037 | title: `${agent.display_name} needs more budget`, |
| 1038 | body: `"${row.title}" reached its cap of ${dollars(row.cap_micros ?? 0)}.`, |
| 1039 | href, |
| 1040 | actor: { kind: "agent", id: agent.id, name: agent.display_name, avatar_seed: agent.avatar_seed || agent.handle }, |
| 1041 | // While that conversation is open the card is there already: no toast. |
| 1042 | channel_id: root.channel_id, |
| 1043 | card, |
| 1044 | created_at: iso(), |
| 1045 | }, |
| 1046 | }), |
| 1047 | }).catch(() => undefined); |
| 1048 | } |
| 1049 | } |
| 1050 | |
| 1051 | /** Sessions stuck mid-step (their desk died): picked up again. */ |
| 1052 | export async function sweep(env: SessionEnv): Promise<number> { |
| 1053 | const before = new Date(Date.now() - 20 * 60_000).toISOString(); |
| 1054 | const stuck = await env.DB.prepare( |
| 1055 | "SELECT id, agent_id FROM agent_sessions WHERE (status = 'working' AND step_started_at < ?) OR (status = 'queued' AND updated_at < ?) LIMIT 50", |
| 1056 | ) |
| 1057 | .bind(before, before) |
| 1058 | .all<{ id: string; agent_id: string }>(); |
| 1059 | for (const s of stuck.results) await wake(env, s.agent_id, s.id).catch(() => undefined); |
| 1060 | return stuck.results.length; |
| 1061 | } |
| 1062 | |
| 1063 | type ToolPorts = import("./tools.ts").ToolPorts; |