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