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