| 1 | /** |
| 2 | * Routines (docs/WORKSPACE.md, "Routines"): work an agent does on a |
| 3 | * schedule, such as Izzy's Monday digest of support themes or Bruno's |
| 4 | * morning look at failed deploys. Each run is a session in the routine's |
| 5 | * channel, paid from the agent's budget, with the access of the person who |
| 6 | * set it up (its sponsor), never more: if the sponsor leaves the workspace |
| 7 | * or can no longer read the channel, the routine pauses and says why. |
| 8 | * |
| 9 | * The schedule is in UTC, every hour, day, weekday or week, at a minute |
| 10 | * (and hour, and day). `nextRun` is pure, so it is tested on its own. |
| 11 | */ |
| 12 | import { type AgentRoutine, askerAccess, chatClient, identityClient, newId } from "@g1t/contracts"; |
| 13 | |
| 14 | import { type SessionEnv, startSession } from "./sessions.ts"; |
| 15 | import { checkSchedule, describeSchedule, nextRun } from "./schedule.ts"; |
| 16 | |
| 17 | export { checkRoutine, checkSchedule, describeSchedule, nextRun } from "./schedule.ts"; |
| 18 | import type { Row } from "./store.ts"; |
| 19 | |
| 20 | /** Routines one agent may keep. */ |
| 21 | export const MAX_ROUTINES = 25; |
| 22 | |
| 23 | export type RoutineRow = { |
| 24 | id: string; |
| 25 | agent_id: string; |
| 26 | workspace_id: string; |
| 27 | /** The workspace's slug, for posting and billing. */ |
| 28 | workspace: string; |
| 29 | name: string; |
| 30 | instructions: string; |
| 31 | schedule: string; |
| 32 | channel_id: string; |
| 33 | channel_name: string | null; |
| 34 | sponsor: string; |
| 35 | sponsor_username: string | null; |
| 36 | enabled: number; |
| 37 | paused_note: string | null; |
| 38 | next_run_at: string | null; |
| 39 | last_run_at: string | null; |
| 40 | last_session_id: string | null; |
| 41 | runs: number; |
| 42 | created_at: string; |
| 43 | updated_at: string; |
| 44 | }; |
| 45 | |
| 46 | export function toRoutine(row: RoutineRow): AgentRoutine { |
| 47 | const schedule = checkSchedule(JSON.parse(row.schedule || "{}")); |
| 48 | return { |
| 49 | id: row.id, |
| 50 | agent_id: row.agent_id, |
| 51 | name: row.name, |
| 52 | instructions: row.instructions, |
| 53 | schedule: schedule.ok ? schedule.value : { every: "day", minute: 0, hour: 9, weekday: 1 }, |
| 54 | channel_id: row.channel_id, |
| 55 | channel_name: row.channel_name, |
| 56 | sponsor: row.sponsor, |
| 57 | sponsor_username: row.sponsor_username, |
| 58 | enabled: !!row.enabled, |
| 59 | paused_note: row.paused_note, |
| 60 | next_run_at: row.enabled ? row.next_run_at : null, |
| 61 | last_run_at: row.last_run_at, |
| 62 | last_session_id: row.last_session_id, |
| 63 | runs: row.runs, |
| 64 | created_at: row.created_at, |
| 65 | updated_at: row.updated_at, |
| 66 | }; |
| 67 | } |
| 68 | |
| 69 | /** Pauses a routine and says why, on its page. */ |
| 70 | async function pause(db: D1Database, id: string, note: string): Promise<void> { |
| 71 | await db.prepare("UPDATE agent_routines SET enabled = 0, paused_note = ?, updated_at = ? WHERE id = ?").bind(note, new Date().toISOString(), id).run(); |
| 72 | } |
| 73 | |
| 74 | /** |
| 75 | * Runs one routine now, as a session: checks its sponsor may still read its |
| 76 | * channel and that the agent is still in it, then starts the session and |
| 77 | * moves its next run on. Returns the session's id, or why it couldn't. |
| 78 | */ |
| 79 | export async function runRoutine(env: SessionEnv, routine: RoutineRow, agent: Row, workspace: string, now = new Date()): Promise<{ ok: true; session: string } | { ok: false; message: string }> { |
| 80 | const db = env.DB; |
| 81 | const [sponsor] = await identityClient(env.IDENTITY) |
| 82 | .usersForAudience([routine.sponsor]) |
| 83 | .catch(() => []); |
| 84 | const membership = sponsor?.workspaces?.find((m) => m.slug.toLowerCase() === workspace.toLowerCase()); |
| 85 | if (!sponsor || !membership) { |
| 86 | await pause(db, routine.id, "Paused: the person who set it up is no longer in the workspace. Anyone who can manage the agent can take it over by saving it."); |
| 87 | return { ok: false, message: "Its sponsor is no longer in the workspace." }; |
| 88 | } |
| 89 | const audience = await chatClient(env.CHAT).audience(workspace, routine.channel_id); |
| 90 | if (!audience.ok) { |
| 91 | await pause(db, routine.id, "Paused: its channel is gone, or the agent is no longer in it."); |
| 92 | return { ok: false, message: "Its channel can't be read." }; |
| 93 | } |
| 94 | if (audience.value.kind !== "public" && !audience.value.member_user_ids.includes(sponsor.id)) { |
| 95 | await pause(db, routine.id, `Paused: @${sponsor.username} is no longer in its channel.`); |
| 96 | return { ok: false, message: "Its sponsor is no longer in its channel." }; |
| 97 | } |
| 98 | const schedule = checkSchedule(JSON.parse(routine.schedule || "{}")); |
| 99 | const next = schedule.ok ? nextRun(schedule.value, now).toISOString() : null; |
| 100 | // Moved on first, so a slow start never runs it twice. |
| 101 | await db |
| 102 | .prepare("UPDATE agent_routines SET next_run_at = ?, last_run_at = ?, runs = runs + 1, updated_at = ? WHERE id = ?") |
| 103 | .bind(next, now.toISOString(), now.toISOString(), routine.id) |
| 104 | .run(); |
| 105 | const session = await startSession(env, { |
| 106 | agent, |
| 107 | kind: "routine", |
| 108 | title: routine.name, |
| 109 | goal: `This is your routine "${routine.name}" (${schedule.ok ? describeSchedule(schedule.value) : "scheduled"}), set up by @${sponsor.username}. Post its report for #${routine.channel_name ?? "the channel"}.\n\n${routine.instructions}`, |
| 110 | workspace, |
| 111 | channel_id: routine.channel_id, |
| 112 | channel_kind: audience.value.kind === "dm" ? "dm" : "channel", |
| 113 | channel_name: routine.channel_name, |
| 114 | thread_root: null, |
| 115 | message_id: null, |
| 116 | asked_by: sponsor.id, |
| 117 | asked_by_username: sponsor.username, |
| 118 | asker: askerAccess(sponsor, workspace), |
| 119 | routine_id: routine.id, |
| 120 | chain: [], |
| 121 | hops: 0, |
| 122 | }); |
| 123 | await db.prepare("UPDATE agent_routines SET last_session_id = ? WHERE id = ?").bind(session.id, routine.id).run(); |
| 124 | return { ok: true, session: session.id }; |
| 125 | } |
| 126 | |
| 127 | /** Every routine due now, run. For the cron trigger, every few minutes. */ |
| 128 | export async function runDue(env: SessionEnv, now = new Date()): Promise<number> { |
| 129 | const db = env.DB; |
| 130 | const due = await db |
| 131 | .prepare("SELECT * FROM agent_routines WHERE enabled = 1 AND next_run_at IS NOT NULL AND next_run_at <= ? ORDER BY next_run_at LIMIT 25") |
| 132 | .bind(now.toISOString()) |
| 133 | .all<RoutineRow>(); |
| 134 | let ran = 0; |
| 135 | for (const routine of due.results) { |
| 136 | const agent = await db.prepare("SELECT * FROM agents WHERE id = ?").bind(routine.agent_id).first<Row>(); |
| 137 | if (!agent || agent.archived_at) { |
| 138 | await pause(db, routine.id, "Paused: its agent was archived."); |
| 139 | continue; |
| 140 | } |
| 141 | const slug = routine.workspace; |
| 142 | if (!slug) continue; |
| 143 | const result = await runRoutine(env, routine, agent, slug, now).catch((error: unknown) => ({ ok: false as const, message: String(error) })); |
| 144 | if (result.ok) ran++; |
| 145 | else console.error("agents: a routine did not run", routine.id, result.message); |
| 146 | } |
| 147 | return ran; |
| 148 | } |
| 149 | |
| 150 | export function newRoutineId(): string { |
| 151 | return newId("rtn"); |
| 152 | } |