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