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