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