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