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