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