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