| 1 | /** |
| 2 | * Routines that run when something happens |
| 3 | * (docs.g1t.sh/guides/agent-routines/, "When something happens"): the |
| 4 | * events service sends this service the events routines can run on |
| 5 | * (`SUBSCRIBER_AGENTS`, crates/contracts subscribers.rs), and each matching |
| 6 | * routine runs once for the one thing that happened. |
| 7 | * |
| 8 | * - **Which routines:** enabled ones in the event's workspace that run on |
| 9 | * its kind, and follow its repository (or every repository). |
| 10 | * - **Whose access:** the routine's sponsor must be able to read the |
| 11 | * repository; otherwise the routine never hears of it, and nothing says |
| 12 | * it happened. |
| 13 | * - **Once:** each routine runs once per pull request ready for review or |
| 14 | * merged, issue opened, deploy failed, and per commit whose checks fail. |
| 15 | * - **Not too often:** at most `MAX_RUNS_PER_HOUR` runs per routine, so a |
| 16 | * burst of events can't run up its agent's budget; the rest are skipped. |
| 17 | */ |
| 18 | import { type G1tEvent, type User, identityClient, reposClient, workClient } from "@g1t/contracts"; |
| 19 | |
| 20 | import { type RoutineOccasion, type RoutineRow, eventsOf, runRoutine } from "./routines.ts"; |
| 21 | import type { SessionEnv } from "./sessions.ts"; |
| 22 | import type { Row } from "./store.ts"; |
| 23 | import { type Happened, follows, happened } from "./happened.ts"; |
| 24 | |
| 25 | export { follows, happened } from "./happened.ts"; |
| 26 | |
| 27 | export const MAX_RUNS_PER_HOUR = 20; |
| 28 | |
| 29 | /** |
| 30 | * What happened, as the routine's session is told it, read as its sponsor: |
| 31 | * a draft pull request is not ready for review, so it is nothing yet. |
| 32 | */ |
| 33 | async function occasionFor(env: SessionEnv, what: Happened, repo: { namespace: string; name: string }, sponsor: User): Promise<RoutineOccasion | null> { |
| 34 | const full = `${repo.namespace}/${repo.name}`; |
| 35 | const work = workClient(env.WORK); |
| 36 | if (what.kind === "pull_ready" || what.kind === "pull_merged" || what.kind === "checks_failed") { |
| 37 | const pull = await work.getPull(repo, what.number!, sponsor).catch(() => null); |
| 38 | if (!pull?.ok) return null; |
| 39 | if (what.kind === "pull_ready" && pull.value.pull.status !== "open") return null; |
| 40 | const title = pull.value.pull.title; |
| 41 | const href = `/${full}/pull/${what.number}`; |
| 42 | const line = |
| 43 | what.kind === "pull_ready" |
| 44 | ? `pull request ${full}#${what.number} is ready for review: ${title}` |
| 45 | : what.kind === "pull_merged" |
| 46 | ? `pull request ${full}#${what.number} was merged: ${title}` |
| 47 | : `checks failed on pull request ${full}#${what.number} (${title}) at ${String(what.data.commit ?? "").slice(0, 8)}`; |
| 48 | return { what: line, href, repo_id: what.repoId }; |
| 49 | } |
| 50 | if (what.kind === "issue_opened") { |
| 51 | const title = typeof what.data.title === "string" ? what.data.title : `#${what.number}`; |
| 52 | return { what: `issue ${full}#${what.number} was opened: ${title}`, href: `/${full}/issues/${what.number}`, repo_id: what.repoId }; |
| 53 | } |
| 54 | const project = typeof what.data.project === "string" ? what.data.project : full; |
| 55 | const kind = what.data.kind === "preview" ? "preview" : "production"; |
| 56 | const branch = typeof what.data.branch === "string" ? ` of ${what.data.branch}` : ""; |
| 57 | return { what: `a ${kind} deploy${branch} of ${project} failed`, href: null, repo_id: what.repoId }; |
| 58 | } |
| 59 | |
| 60 | /** Runs every routine an event is for. Never throws for one routine's sake. */ |
| 61 | export async function onEvents(env: SessionEnv, events: G1tEvent[], now = new Date()): Promise<number> { |
| 62 | const db = env.DB; |
| 63 | let ran = 0; |
| 64 | for (const event of events) { |
| 65 | const what = happened(event); |
| 66 | if (!what) continue; |
| 67 | const routines = await db |
| 68 | .prepare("SELECT * FROM agent_routines WHERE enabled = 1 AND events LIKE ? LIMIT 200") |
| 69 | .bind(`%"${what.kind}"%`) |
| 70 | .all<RoutineRow>(); |
| 71 | for (const routine of routines.results) { |
| 72 | if (!eventsOf(routine).includes(what.kind)) continue; |
| 73 | try { |
| 74 | if (await runFor(env, routine, what, now)) ran++; |
| 75 | } catch (error) { |
| 76 | console.error("agents: a routine didn't run on an event", routine.id, event.type, String(error)); |
| 77 | } |
| 78 | } |
| 79 | } |
| 80 | return ran; |
| 81 | } |
| 82 | |
| 83 | async function runFor(env: SessionEnv, routine: RoutineRow, what: Happened, now: Date): Promise<boolean> { |
| 84 | const db = env.DB; |
| 85 | const [sponsor] = await identityClient(env.IDENTITY) |
| 86 | .usersForAudience([routine.sponsor]) |
| 87 | .catch(() => [] as User[]); |
| 88 | if (!sponsor) return false; |
| 89 | // The repository, as the sponsor can read it, in the routine's own workspace. |
| 90 | const [repo] = await reposClient(env.REPOS) |
| 91 | .readable([what.repoId], sponsor) |
| 92 | .catch(() => []); |
| 93 | if (!repo || repo.namespace.toLowerCase() !== routine.workspace.toLowerCase()) return false; |
| 94 | let repos: string[] = []; |
| 95 | try { |
| 96 | repos = JSON.parse(routine.repos || "[]") as string[]; |
| 97 | } catch { |
| 98 | repos = []; |
| 99 | } |
| 100 | if (!follows(repos, `${repo.namespace}/${repo.name}`)) return false; |
| 101 | const hourAgo = new Date(now.getTime() - 3_600_000).toISOString(); |
| 102 | const recent = await db.prepare("SELECT COUNT(*) AS n FROM agent_routine_runs WHERE routine_id = ? AND created_at > ?").bind(routine.id, hourAgo).first<{ n: number }>(); |
| 103 | if ((recent?.n ?? 0) >= MAX_RUNS_PER_HOUR) return false; |
| 104 | const occasion = await occasionFor(env, what, repo, sponsor); |
| 105 | if (!occasion) return false; |
| 106 | // Once per thing that happened, however often it is told. |
| 107 | const claimed = await db |
| 108 | .prepare("INSERT INTO agent_routine_runs (routine_id, run_key, created_at) VALUES (?, ?, ?) ON CONFLICT DO NOTHING RETURNING routine_id") |
| 109 | .bind(routine.id, what.key, now.toISOString()) |
| 110 | .first(); |
| 111 | if (!claimed) return false; |
| 112 | const agent = await db.prepare("SELECT * FROM agents WHERE id = ? AND archived_at IS NULL").bind(routine.agent_id).first<Row>(); |
| 113 | if (!agent) return false; |
| 114 | const result = await runRoutine(env, routine, agent, routine.workspace, now, occasion); |
| 115 | if (result.ok) { |
| 116 | await db.prepare("UPDATE agent_routine_runs SET session_id = ? WHERE routine_id = ? AND run_key = ?").bind(result.session, routine.id, what.key).run(); |
| 117 | } |
| 118 | return result.ok; |
| 119 | } |