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