Skip to content
119 linesCodeBlameRaw
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 */
18import { type G1tEvent, type User, identityClient, reposClient, workClient } from "@g1t/contracts";
19
20import { type RoutineOccasion, type RoutineRow, eventsOf, runRoutine } from "./routines.ts";
21import type { SessionEnv } from "./sessions.ts";
22import type { Row } from "./store.ts";
23import { type Happened, follows, happened } from "./happened.ts";
24
25export { follows, happened } from "./happened.ts";
26
27export 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 */
33async 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. */
61export 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
83async 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}