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