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