Routines run when something happens: a pull request ready for review or merged, checks or a deploy failing, an issue opened
- The events service sends agents these events (SUBSCRIBER_AGENTS, queue g1t-events-agents); each matching routine runs once per thing that happened, only in repositories its sponsor can read, at most 20 runs an hour - A routine runs on a schedule, on events, or both, optionally for named repositories - Routines suggested from an agent's responsibilities: reviewing pull requests runs when one is ready for review, chasing flaky checks when checks fail, test plans when an issue is opened
14 files+524−340/14 viewed
| 189 | 189 | "team.deleted", | |
| 190 | 190 | ], | |
| 191 | 191 | ), | |
| 192 | + | // Agents' routines that run on events: a pull request ready for | |
| 193 | + | // review or merged, checks or a deploy failing, an issue opened | |
| 194 | + | // (agents/src/triggers.ts). | |
| 195 | + | ( | |
| 196 | + | "SUBSCRIBER_AGENTS", | |
| 197 | + | &[ | |
| 198 | + | "pull.opened", | |
| 199 | + | "pull.ready", | |
| 200 | + | "pull.merged", | |
| 201 | + | "checks.completed", | |
| 202 | + | "issue.opened", | |
| 203 | + | "deployment.failed", | |
| 204 | + | ], | |
| 205 | + | ), | |
| 192 | 206 | ]; | |
| 193 | 207 | ||
| 194 | 208 | /// Whether `pattern` (see the module's notes) matches the type `kind`. | |
| ⋯ | |||
| 265 | 279 | names.sort_unstable(); | |
| 266 | 280 | names.dedup(); | |
| 267 | 281 | assert_eq!(names.len(), count); | |
| 268 | − | assert_eq!(count, 13); | |
| 282 | + | assert_eq!(count, 14); | |
| 269 | 283 | } | |
| 270 | 284 | ||
| 271 | 285 | #[test] | |
| 432 | 432 | }; | |
| 433 | 433 | ||
| 434 | 434 | /** | |
| 435 | − | * A routine: work an agent does on a schedule, such as Izzy's Monday digest | |
| 435 | + | * Things that happen in the workspace a routine can run on | |
| 436 | + | * (docs/WORKSPACE.md, "Routines"). Each run is one session about the one | |
| 437 | + | * thing that happened, in a repository its sponsor can read. | |
| 438 | + | */ | |
| 439 | + | export const ROUTINE_EVENTS = [ | |
| 440 | + | { key: "pull_ready", label: "A pull request is ready for review", hint: "Opened ready, or moved out of draft." }, | |
| 441 | + | { key: "pull_merged", label: "A pull request is merged", hint: "On any branch it targets." }, | |
| 442 | + | { key: "checks_failed", label: "Checks fail on a pull request", hint: "Its required checks failed or errored." }, | |
| 443 | + | { key: "issue_opened", label: "An issue is opened", hint: "By a person or an agent." }, | |
| 444 | + | { key: "deploy_failed", label: "A deploy fails", hint: "A production or preview deploy." }, | |
| 445 | + | ] as const; | |
| 446 | + | ||
| 447 | + | export type RoutineEvent = (typeof ROUTINE_EVENTS)[number]["key"]; | |
| 448 | + | ||
| 449 | + | /** | |
| 450 | + | * A routine: work an agent does on a schedule or when something happens, such as Izzy's Monday digest | |
| 436 | 451 | * of support themes. Each run is a session posted in the routine's channel, | |
| 437 | 452 | * paid from the agent's budget, and run with the access of the person who | |
| 438 | 453 | * set it up (its sponsor), never more. | |
| ⋯ | |||
| 442 | 457 | agent_id: string; | |
| 443 | 458 | name: string; | |
| 444 | 459 | instructions: string; | |
| 445 | − | schedule: RoutineSchedule; | |
| 460 | + | /** When it runs on a clock; null when it runs only on events. */ | |
| 461 | + | schedule: RoutineSchedule | null; | |
| 462 | + | /** What it runs on; empty when it runs only on its schedule. */ | |
| 463 | + | events: RoutineEvent[]; | |
| 464 | + | /** Which repositories its events come from, by `workspace/name`; empty: every one its sponsor can read. */ | |
| 465 | + | repos: string[]; | |
| 446 | 466 | channel_id: string; | |
| 447 | 467 | channel_name: string | null; | |
| 448 | 468 | sponsor: string; | |
| ⋯ | |||
| 458 | 478 | updated_at: string; | |
| 459 | 479 | }; | |
| 460 | 480 | ||
| 481 | + | /** A routine suggested from an agent's responsibilities, for an owner to add in one step. */ | |
| 482 | + | export type RoutineSuggestion = { responsibility: string; routine: Omit<NewRoutine, "channel_id"> }; | |
| 483 | + | ||
| 461 | 484 | export type NewRoutine = { | |
| 462 | 485 | name: string; | |
| 463 | 486 | instructions: string; | |
| 464 | − | schedule: RoutineSchedule; | |
| 487 | + | /** A schedule, events, or both; at least one. */ | |
| 488 | + | schedule: RoutineSchedule | null; | |
| 489 | + | events?: RoutineEvent[]; | |
| 490 | + | repos?: string[]; | |
| 465 | 491 | /** A channel the agent is in, by id. */ | |
| 466 | 492 | channel_id: string; | |
| 467 | 493 | enabled?: boolean; | |
| ⋯ | |||
| 517 | 543 | waiting_on_you: AgentSession[]; | |
| 518 | 544 | /** Recently finished sessions the viewer can see. */ | |
| 519 | 545 | recent: AgentSession[]; | |
| 520 | − | /** The next routines to run. */ | |
| 546 | + | /** The next routines to run on a schedule. */ | |
| 521 | 547 | upcoming: (AgentRoutine & { agent_handle: string; agent_name: string })[]; | |
| 522 | 548 | spend: AgentSpendBreakdown; | |
| 523 | 549 | can_manage: boolean; | |
| ⋯ | |||
| 594 | 620 | changes: { body?: string; pinned?: boolean }, | |
| 595 | 621 | ): Promise<Result<AgentMemory>>; | |
| 596 | 622 | forget(workspace: string, handle: string, viewer: User, id: string): Promise<Result<null>>; | |
| 597 | − | routines(workspace: string, handle: string, viewer: User): Promise<Result<AgentRoutine[]>>; | |
| 623 | + | routines(workspace: string, handle: string, viewer: User): Promise<Result<{ routines: AgentRoutine[]; suggestions: RoutineSuggestion[] }>>; | |
| 598 | 624 | saveRoutine(workspace: string, handle: string, viewer: User, input: NewRoutine, id?: string | null): Promise<Result<AgentRoutine>>; | |
| 599 | 625 | deleteRoutine(workspace: string, handle: string, viewer: User, id: string): Promise<Result<null>>; | |
| 600 | 626 | /** Runs a routine now, as a session. */ | |
| 118 | 118 | name TEXT NOT NULL, | |
| 119 | 119 | instructions TEXT NOT NULL, | |
| 120 | 120 | -- JSON: every (hour, day, weekday, week), minute, hour, weekday. UTC. | |
| 121 | − | schedule TEXT NOT NULL, | |
| 121 | + | -- Null when it runs only on events. | |
| 122 | + | schedule TEXT, | |
| 123 | + | -- JSON lists: what it runs on (pull_ready, checks_failed...), and which | |
| 124 | + | -- repositories those come from (workspace/name; empty: any its sponsor can read). | |
| 125 | + | events TEXT NOT NULL DEFAULT '[]', | |
| 126 | + | repos TEXT NOT NULL DEFAULT '[]', | |
| 122 | 127 | channel_id TEXT NOT NULL, | |
| 123 | 128 | channel_name TEXT, | |
| 124 | 129 | sponsor TEXT NOT NULL, | |
| ⋯ | |||
| 135 | 140 | ||
| 136 | 141 | CREATE INDEX agent_routines_agent ON agent_routines (agent_id); | |
| 137 | 142 | CREATE INDEX agent_routines_due ON agent_routines (enabled, next_run_at); | |
| 143 | + | CREATE INDEX agent_routines_workspace ON agent_routines (workspace_id, enabled); | |
| 144 | + | ||
| 145 | + | -- What each routine ran on, so one thing that happened runs it once. | |
| 146 | + | CREATE TABLE agent_routine_runs ( | |
| 147 | + | routine_id TEXT NOT NULL, | |
| 148 | + | -- The event's kind and what it was about: pull_ready:rep_1:12. | |
| 149 | + | run_key TEXT NOT NULL, | |
| 150 | + | session_id TEXT, | |
| 151 | + | created_at TEXT NOT NULL, | |
| 152 | + | PRIMARY KEY (routine_id, run_key) | |
| 153 | + | ); | |
| 154 | + | CREATE INDEX agent_routine_runs_recent ON agent_routine_runs (routine_id, created_at); | |
| 138 | 155 | ||
| 139 | 156 | -- The workspace's say over all its agents together. | |
| 140 | 157 | CREATE TABLE agent_policies ( | |
| 1 | + | /** | |
| 2 | + | * Which routine event a system event is, and what makes it the same thing | |
| 3 | + | * twice. Pure, so it is tested on its own (triggers.ts runs routines). | |
| 4 | + | */ | |
| 5 | + | import type { G1tEvent, RoutineEvent } from "@g1t/contracts"; | |
| 6 | + | ||
| 7 | + | export type Happened = { kind: RoutineEvent; repoId: string; key: string; number: number | null; data: Record<string, unknown> }; | |
| 8 | + | ||
| 9 | + | /** The routine event an event is, if any, and what makes it the same thing twice. */ | |
| 10 | + | export function happened(event: Pick<G1tEvent, "type" | "repoId" | "data">): Happened | null { | |
| 11 | + | const data = (event.data ?? {}) as Record<string, unknown>; | |
| 12 | + | const repoId = (typeof data.repoId === "string" ? data.repoId : null) ?? event.repoId; | |
| 13 | + | if (!repoId) return null; | |
| 14 | + | const number = typeof data.number === "number" ? data.number : null; | |
| 15 | + | switch (event.type) { | |
| 16 | + | case "pull.opened": | |
| 17 | + | case "pull.ready": | |
| 18 | + | return number == null ? null : { kind: "pull_ready", repoId, key: `pull_ready:${repoId}:${number}`, number, data }; | |
| 19 | + | case "pull.merged": | |
| 20 | + | return number == null ? null : { kind: "pull_merged", repoId, key: `pull_merged:${repoId}:${number}`, number, data }; | |
| 21 | + | case "checks.completed": { | |
| 22 | + | const status = data.status; | |
| 23 | + | if (number == null || (status !== "failed" && status !== "errored")) return null; | |
| 24 | + | return { kind: "checks_failed", repoId, key: `checks_failed:${repoId}:${number}:${String(data.commit ?? "")}`, number, data }; | |
| 25 | + | } | |
| 26 | + | case "issue.opened": | |
| 27 | + | return number == null ? null : { kind: "issue_opened", repoId, key: `issue_opened:${repoId}:${number}`, number, data }; | |
| 28 | + | case "deployment.failed": | |
| 29 | + | return { kind: "deploy_failed", repoId, key: `deploy_failed:${String(data.deploymentId ?? "")}`, number, data }; | |
| 30 | + | default: | |
| 31 | + | return null; | |
| 32 | + | } | |
| 33 | + | } | |
| 34 | + | ||
| 35 | + | /** Whether a routine follows this repository: all of them, or it is on its list. */ | |
| 36 | + | export function follows(repos: string[], full: string): boolean { | |
| 37 | + | return !repos.length || repos.includes(full.toLowerCase()); | |
| 38 | + | } | |
| 39 | + |
| 10 | 10 | */ | |
| 11 | 11 | import { | |
| 12 | 12 | type AgentDelivery, | |
| 13 | + | type G1tEvent, | |
| 13 | 14 | type NewWorkspaceAgent, | |
| 14 | 15 | type Result, | |
| 15 | 16 | type ServiceBinding, | |
| ⋯ | |||
| 36 | 37 | import { TEMPLATES, TEMPLATE_IDS } from "./templates.ts"; | |
| 37 | 38 | import { readPolicy } from "./policy.ts"; | |
| 38 | 39 | import { runDue } from "./routines.ts"; | |
| 40 | + | import { onEvents } from "./triggers.ts"; | |
| 39 | 41 | import { type SessionEnv, sweep } from "./sessions.ts"; | |
| 40 | 42 | import * as views from "./views.ts"; | |
| 41 | 43 | import { monthKey } from "./budget.ts"; | |
| ⋯ | |||
| 454 | 456 | return opened.finish(await answer(service, match[1], args)); | |
| 455 | 457 | }, | |
| 456 | 458 | ||
| 459 | + | /** Events routines run on, from the events service (SUBSCRIBER_AGENTS). */ | |
| 460 | + | async queue(batch: MessageBatch<unknown>, env: Env): Promise<void> { | |
| 461 | + | await onEvents(env as unknown as SessionEnv, batch.messages.map((message) => message.body as G1tEvent)); | |
| 462 | + | batch.ackAll(); | |
| 463 | + | }, | |
| 464 | + | ||
| 457 | 465 | /** Every few minutes: routines that are due, and session steps a desk lost. */ | |
| 458 | 466 | async scheduled(_controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> { | |
| 459 | 467 | const sessions = env as unknown as SessionEnv; | |
| 9 | 9 | * The schedule is in UTC, every hour, day, weekday or week, at a minute | |
| 10 | 10 | * (and hour, and day). `nextRun` is pure, so it is tested on its own. | |
| 11 | 11 | */ | |
| 12 | − | import { type AgentRoutine, askerAccess, chatClient, identityClient, newId } from "@g1t/contracts"; | |
| 12 | + | import { type AgentRoutine, type RoutineEvent, type RoutineSchedule, askerAccess, chatClient, identityClient, newId } from "@g1t/contracts"; | |
| 13 | 13 | ||
| 14 | 14 | import { type SessionEnv, startSession } from "./sessions.ts"; | |
| 15 | 15 | import { checkSchedule, describeSchedule, nextRun } from "./schedule.ts"; | |
| 16 | + | import { EVENT_KEYS, describeEvents } from "./suggest.ts"; | |
| 16 | 17 | ||
| 17 | 18 | export { checkRoutine, checkSchedule, describeSchedule, nextRun } from "./schedule.ts"; | |
| 18 | 19 | import type { Row } from "./store.ts"; | |
| ⋯ | |||
| 28 | 29 | workspace: string; | |
| 29 | 30 | name: string; | |
| 30 | 31 | instructions: string; | |
| 31 | − | schedule: string; | |
| 32 | + | schedule: string | null; | |
| 33 | + | events: string; | |
| 34 | + | repos: string; | |
| 32 | 35 | channel_id: string; | |
| 33 | 36 | channel_name: string | null; | |
| 34 | 37 | sponsor: string; | |
| ⋯ | |||
| 43 | 46 | updated_at: string; | |
| 44 | 47 | }; | |
| 45 | 48 | ||
| 49 | + | function list<T>(raw: string | null): T[] { | |
| 50 | + | try { | |
| 51 | + | const value = JSON.parse(raw || "[]") as unknown; | |
| 52 | + | return Array.isArray(value) ? (value as T[]) : []; | |
| 53 | + | } catch { | |
| 54 | + | return []; | |
| 55 | + | } | |
| 56 | + | } | |
| 57 | + | ||
| 58 | + | export function scheduleOf(row: Pick<RoutineRow, "schedule">): RoutineSchedule | null { | |
| 59 | + | if (!row.schedule) return null; | |
| 60 | + | try { | |
| 61 | + | const checked = checkSchedule(JSON.parse(row.schedule)); | |
| 62 | + | return checked.ok ? checked.value : null; | |
| 63 | + | } catch { | |
| 64 | + | return null; | |
| 65 | + | } | |
| 66 | + | } | |
| 67 | + | ||
| 68 | + | export function eventsOf(row: Pick<RoutineRow, "events">): RoutineEvent[] { | |
| 69 | + | return list<RoutineEvent>(row.events).filter((e) => EVENT_KEYS.includes(e)); | |
| 70 | + | } | |
| 71 | + | ||
| 46 | 72 | export function toRoutine(row: RoutineRow): AgentRoutine { | |
| 47 | − | const schedule = checkSchedule(JSON.parse(row.schedule || "{}")); | |
| 73 | + | const schedule = scheduleOf(row); | |
| 48 | 74 | return { | |
| 49 | 75 | id: row.id, | |
| 50 | 76 | agent_id: row.agent_id, | |
| 51 | 77 | name: row.name, | |
| 52 | 78 | instructions: row.instructions, | |
| 53 | − | schedule: schedule.ok ? schedule.value : { every: "day", minute: 0, hour: 9, weekday: 1 }, | |
| 79 | + | schedule, | |
| 80 | + | events: eventsOf(row), | |
| 81 | + | repos: list<string>(row.repos), | |
| 54 | 82 | channel_id: row.channel_id, | |
| 55 | 83 | channel_name: row.channel_name, | |
| 56 | 84 | sponsor: row.sponsor, | |
| ⋯ | |||
| 76 | 104 | * channel and that the agent is still in it, then starts the session and | |
| 77 | 105 | * moves its next run on. Returns the session's id, or why it couldn't. | |
| 78 | 106 | */ | |
| 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 }> { | |
| 107 | + | export type RoutineOccasion = { | |
| 108 | + | /** What happened, in a line, and where: "Pull request acme/web#12 is ready for review: Fix CSV export". */ | |
| 109 | + | what: string; | |
| 110 | + | /** Where to read more, relative to the site. */ | |
| 111 | + | href: string | null; | |
| 112 | + | /** The repository it happened in, by id, which the sponsor must be able to read. */ | |
| 113 | + | repo_id: string; | |
| 114 | + | }; | |
| 115 | + | ||
| 116 | + | export async function runRoutine( | |
| 117 | + | env: SessionEnv, | |
| 118 | + | routine: RoutineRow, | |
| 119 | + | agent: Row, | |
| 120 | + | workspace: string, | |
| 121 | + | now = new Date(), | |
| 122 | + | occasion: RoutineOccasion | null = null, | |
| 123 | + | ): Promise<{ ok: true; session: string } | { ok: false; message: string }> { | |
| 80 | 124 | const db = env.DB; | |
| 81 | 125 | const [sponsor] = await identityClient(env.IDENTITY) | |
| 82 | 126 | .usersForAudience([routine.sponsor]) | |
| ⋯ | |||
| 95 | 139 | await pause(db, routine.id, `Paused: @${sponsor.username} is no longer in its channel.`); | |
| 96 | 140 | return { ok: false, message: "Its sponsor is no longer in its channel." }; | |
| 97 | 141 | } | |
| 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. | |
| 142 | + | const schedule = scheduleOf(routine); | |
| 143 | + | // A scheduled run moves the clock on first, so a slow start never runs it twice; an event's run leaves it. | |
| 144 | + | const next = occasion ? routine.next_run_at : schedule ? nextRun(schedule, now).toISOString() : null; | |
| 101 | 145 | await db | |
| 102 | 146 | .prepare("UPDATE agent_routines SET next_run_at = ?, last_run_at = ?, runs = runs + 1, updated_at = ? WHERE id = ?") | |
| 103 | 147 | .bind(next, now.toISOString(), now.toISOString(), routine.id) | |
| 104 | 148 | .run(); | |
| 149 | + | const when = [schedule ? describeSchedule(schedule) : null, describeEvents(eventsOf(routine)) || null].filter(Boolean).join("; "); | |
| 105 | 150 | const session = await startSession(env, { | |
| 106 | 151 | agent, | |
| 107 | 152 | 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}`, | |
| 153 | + | title: occasion ? `${routine.name}: ${occasion.what}`.slice(0, 120) : routine.name, | |
| 154 | + | goal: [ | |
| 155 | + | `This is your routine "${routine.name}" (${when || "run by hand"}), set up by @${sponsor.username}. Post its report for #${routine.channel_name ?? "the channel"}.`, | |
| 156 | + | occasion ? `It runs now because: ${occasion.what}${occasion.href ? ` (${occasion.href})` : ""}. Work on that one thing.` : "", | |
| 157 | + | routine.instructions, | |
| 158 | + | ] | |
| 159 | + | .filter(Boolean) | |
| 160 | + | .join("\n\n"), | |
| 110 | 161 | workspace, | |
| 111 | 162 | channel_id: routine.channel_id, | |
| 112 | 163 | channel_kind: audience.value.kind === "dm" ? "dm" : "channel", | |
| 2 | 2 | * When routines run, and what a routine must have: pure, so it is tested | |
| 3 | 3 | * on its own (routines.ts runs them). | |
| 4 | 4 | */ | |
| 5 | − | import type { NewRoutine, RoutineSchedule } from "@g1t/contracts"; | |
| 5 | + | import type { NewRoutine, RoutineEvent, RoutineSchedule } from "@g1t/contracts"; | |
| 6 | + | ||
| 7 | + | import { EVENT_KEYS } from "./suggest.ts"; | |
| 6 | 8 | ||
| 7 | 9 | const EVERY = ["hour", "day", "weekday", "week"] as const; | |
| 8 | 10 | ||
| ⋯ | |||
| 57 | 59 | } | |
| 58 | 60 | } | |
| 59 | 61 | ||
| 60 | − | /** A routine as given, checked: a name, instructions and a schedule. */ | |
| 61 | − | export function checkRoutine(input: NewRoutine): { ok: true; value: NewRoutine & { schedule: RoutineSchedule } } | { ok: false; message: string } { | |
| 62 | + | export type CheckedRoutine = NewRoutine & { schedule: RoutineSchedule | null; events: RoutineEvent[]; repos: string[] }; | |
| 63 | + | ||
| 64 | + | /** A routine as given, checked: a name, instructions, and a schedule, events, or both. */ | |
| 65 | + | export function checkRoutine(input: NewRoutine): { ok: true; value: CheckedRoutine } | { ok: false; message: string } { | |
| 62 | 66 | const name = typeof input?.name === "string" ? input.name.trim().slice(0, 80) : ""; | |
| 63 | 67 | const instructions = typeof input?.instructions === "string" ? input.instructions.trim().slice(0, 8000) : ""; | |
| 64 | 68 | if (!name) return { ok: false, message: "Give the routine a name." }; | |
| 65 | 69 | if (instructions.length < 10) return { ok: false, message: "Say what the routine does, in a sentence or more." }; | |
| 66 | 70 | if (typeof input.channel_id !== "string" || !input.channel_id) return { ok: false, message: "Choose the channel it posts in." }; | |
| 67 | − | const schedule = checkSchedule(input.schedule); | |
| 68 | − | if (!schedule.ok) return schedule; | |
| 69 | − | return { ok: true, value: { ...input, name, instructions, schedule: schedule.value, enabled: input.enabled !== false } }; | |
| 71 | + | const given = Array.isArray(input.events) ? input.events : []; | |
| 72 | + | if (given.some((e) => !EVENT_KEYS.includes(e))) return { ok: false, message: "That isn't something a routine can run on." }; | |
| 73 | + | const events = [...new Set(given)]; | |
| 74 | + | const repos = [...new Set((Array.isArray(input.repos) ? input.repos : []).filter((r): r is string => typeof r === "string").map((r) => r.trim().toLowerCase()).filter(Boolean))]; | |
| 75 | + | if (repos.some((r) => !/^[a-z0-9._-]+\/[a-z0-9._-]+$/.test(r))) return { ok: false, message: "Name repositories as workspace/name." }; | |
| 76 | + | if (repos.length > 20) return { ok: false, message: "A routine follows at most 20 repositories; leave it empty for all of them." }; | |
| 77 | + | let schedule: RoutineSchedule | null = null; | |
| 78 | + | if (input.schedule) { | |
| 79 | + | const checked = checkSchedule(input.schedule); | |
| 80 | + | if (!checked.ok) return checked; | |
| 81 | + | schedule = checked.value; | |
| 82 | + | } | |
| 83 | + | if (!schedule && !events.length) return { ok: false, message: "A routine runs on a schedule, when something happens, or both." }; | |
| 84 | + | return { ok: true, value: { ...input, name, instructions, schedule, events, repos, enabled: input.enabled !== false } }; | |
| 70 | 85 | } | |
| 71 | − | ||
| 202 | 202 | for (let i = 0; i < 3; i++) assert.equal((await session.run("post_update", { text: `step ${i}` })).outcome, "allowed"); | |
| 203 | 203 | assert.equal((await session.run("post_update", { text: "again" })).outcome, "refused"); | |
| 204 | 204 | }); | |
| 205 | + | ||
| 206 | + | // ── Routines that run when something happens ───────────────────────────── | |
| 207 | + | ||
| 208 | + | import { describeEvents, suggestRoutines } from "./suggest.ts"; | |
| 209 | + | import { checkRoutine as checkEventRoutine } from "./schedule.ts"; | |
| 210 | + | ||
| 211 | + | test("Margo's responsibilities suggest the routines that bind to events", () => { | |
| 212 | + | const margo = ["Reviewing pull requests for risk and test coverage", "Test plans for new features", "Chasing flaky checks"]; | |
| 213 | + | const suggested = suggestRoutines(margo, []); | |
| 214 | + | assert.deepEqual( | |
| 215 | + | suggested.map((s) => [s.routine.name, s.routine.events]), | |
| 216 | + | [ | |
| 217 | + | ["Review pull requests", ["pull_ready"]], | |
| 218 | + | ["Test plans for new work", ["issue_opened"]], | |
| 219 | + | ["Chase failing checks", ["checks_failed"]], | |
| 220 | + | ], | |
| 221 | + | ); | |
| 222 | + | assert.match(suggested[0].routine.instructions, /risk and test coverage/); | |
| 223 | + | // What it already runs on isn't suggested again. | |
| 224 | + | assert.deepEqual( | |
| 225 | + | suggestRoutines(margo, [{ name: "PR reviews", events: ["pull_ready"] }]).map((s) => s.routine.name), | |
| 226 | + | ["Test plans for new work", "Chase failing checks"], | |
| 227 | + | ); | |
| 228 | + | assert.deepEqual(suggestRoutines(["Keep the office plants alive"], []), []); | |
| 229 | + | }); | |
| 230 | + | ||
| 231 | + | test("an event routine needs no schedule, but a routine needs one or the other", () => { | |
| 232 | + | const base = { name: "Review", instructions: "Review pull requests for risk", channel_id: "chn_qa" }; | |
| 233 | + | assert.equal(checkEventRoutine({ ...base, schedule: null, events: ["pull_ready"] }).ok, true); | |
| 234 | + | assert.equal(checkEventRoutine({ ...base, schedule: null, events: [] }).ok, false); | |
| 235 | + | assert.equal(checkEventRoutine({ ...base, schedule: null, events: ["pull_exploded" as never] }).ok, false); | |
| 236 | + | assert.equal(checkEventRoutine({ ...base, schedule: null, events: ["pull_ready"], repos: ["not a repo"] }).ok, false); | |
| 237 | + | const ok = checkEventRoutine({ ...base, schedule: null, events: ["pull_ready", "pull_ready"], repos: ["Acme/Web"] }); | |
| 238 | + | assert.ok(ok.ok && ok.value.events.length === 1 && ok.value.repos[0] === "acme/web"); | |
| 239 | + | assert.equal(describeEvents(["pull_ready", "checks_failed"]), "When a pull request is ready for review or checks fail on a pull request"); | |
| 240 | + | }); |
| 1 | + | /** | |
| 2 | + | * What a routine can run on, in words, and routines that fit an agent's | |
| 3 | + | * responsibilities, for owners to add in one step: Margo's "Reviewing pull | |
| 4 | + | * requests for risk and test coverage" runs when a pull request is ready | |
| 5 | + | * for review. Matched on words, never by a model, so it is instant and the | |
| 6 | + | * same every time. Pure, so it is tested on its own. | |
| 7 | + | */ | |
| 8 | + | import type { NewRoutine, RoutineEvent, RoutineSuggestion } from "@g1t/contracts"; | |
| 9 | + | ||
| 10 | + | /** The events a routine can run on, by key. */ | |
| 11 | + | export const EVENT_KEYS: RoutineEvent[] = ["pull_ready", "pull_merged", "checks_failed", "issue_opened", "deploy_failed"]; | |
| 12 | + | ||
| 13 | + | const WORDS: Record<RoutineEvent, string> = { | |
| 14 | + | pull_ready: "a pull request is ready for review", | |
| 15 | + | pull_merged: "a pull request is merged", | |
| 16 | + | checks_failed: "checks fail on a pull request", | |
| 17 | + | issue_opened: "an issue is opened", | |
| 18 | + | deploy_failed: "a deploy fails", | |
| 19 | + | }; | |
| 20 | + | ||
| 21 | + | /** A routine's events in words: "When a pull request is ready for review or checks fail on a pull request". */ | |
| 22 | + | export function describeEvents(events: RoutineEvent[]): string { | |
| 23 | + | const list = events.map((e) => WORDS[e]).filter(Boolean); | |
| 24 | + | if (!list.length) return ""; | |
| 25 | + | const joined = list.length === 1 ? list[0] : `${list.slice(0, -1).join(", ")} or ${list.at(-1)}`; | |
| 26 | + | return `When ${joined}`; | |
| 27 | + | } | |
| 28 | + | ||
| 29 | + | const trimmed = (duty: string) => duty.trim().replace(/\.$/, ""); | |
| 30 | + | ||
| 31 | + | type Rule = { match: RegExp; routine: (duty: string) => Omit<NewRoutine, "channel_id"> }; | |
| 32 | + | ||
| 33 | + | const RULES: Rule[] = [ | |
| 34 | + | { | |
| 35 | + | match: /\b(review|reviewing|reviews)\b.*\b(pull requests?|prs?|changes?|code)\b|\b(pull requests?|prs?)\b.*\breview/i, | |
| 36 | + | routine: (duty) => ({ | |
| 37 | + | name: "Review pull requests", | |
| 38 | + | instructions: `When a pull request is ready for review, review it (${trimmed(duty)}). Read the change and its checks, then post a short review: what it changes, the risks, which tests cover it and what they miss, and a verdict: looks good, or what to fix first. Link the pull request.`, | |
| 39 | + | schedule: null, | |
| 40 | + | events: ["pull_ready"], | |
| 41 | + | repos: [], | |
| 42 | + | }), | |
| 43 | + | }, | |
| 44 | + | { | |
| 45 | + | match: /\b(flaky|failing|broken)\b.*\b(checks?|tests?|builds?|ci)\b|\b(checks?|ci|builds?)\b.*\b(flaky|fail)/i, | |
| 46 | + | routine: (duty) => ({ | |
| 47 | + | name: "Chase failing checks", | |
| 48 | + | instructions: `When checks fail on a pull request, look into it (${trimmed(duty)}). Read the failure, say whether it looks real or flaky and why, and what to do next. If a flaky check keeps coming back, offer to file an issue.`, | |
| 49 | + | schedule: null, | |
| 50 | + | events: ["checks_failed"], | |
| 51 | + | repos: [], | |
| 52 | + | }), | |
| 53 | + | }, | |
| 54 | + | { | |
| 55 | + | match: /\btest plans?\b|\bnew features?\b.*\btest/i, | |
| 56 | + | routine: (duty) => ({ | |
| 57 | + | name: "Test plans for new work", | |
| 58 | + | instructions: `When an issue is opened for new work, draft a test plan (${trimmed(duty)}): what to test, the edge cases, and what can be automated. Skip bug reports and questions.`, | |
| 59 | + | schedule: null, | |
| 60 | + | events: ["issue_opened"], | |
| 61 | + | repos: [], | |
| 62 | + | }), | |
| 63 | + | }, | |
| 64 | + | { | |
| 65 | + | match: /\b(triage|triaging)\b|\bincoming (issues|bugs|requests)\b/i, | |
| 66 | + | routine: (duty) => ({ | |
| 67 | + | name: "Triage new issues", | |
| 68 | + | instructions: `When an issue is opened, triage it (${trimmed(duty)}): what it is (bug, feature, question), how urgent, which area of the code it touches, and who should own it.`, | |
| 69 | + | schedule: null, | |
| 70 | + | events: ["issue_opened"], | |
| 71 | + | repos: [], | |
| 72 | + | }), | |
| 73 | + | }, | |
| 74 | + | { | |
| 75 | + | match: /\b(deploys?|deployments?|incidents?|on-?call|outages?|rollbacks?)\b/i, | |
| 76 | + | routine: (duty) => ({ | |
| 77 | + | name: "Watch deploys", | |
| 78 | + | instructions: `When a deploy fails, find out why (${trimmed(duty)}): what failed, what changed since the last good deploy, and whether to retry, roll back or fix forward.`, | |
| 79 | + | schedule: null, | |
| 80 | + | events: ["deploy_failed"], | |
| 81 | + | repos: [], | |
| 82 | + | }), | |
| 83 | + | }, | |
| 84 | + | { | |
| 85 | + | match: /\b(release notes|changelog|what shipped|weekly (summary|update|digest)|digests?)\b/i, | |
| 86 | + | routine: (duty) => ({ | |
| 87 | + | name: "Weekly summary", | |
| 88 | + | instructions: `Every Friday (${trimmed(duty)}): what merged and shipped this week, what's still open and what's at risk, in a short post people can skim.`, | |
| 89 | + | schedule: { every: "week", minute: 0, hour: 16, weekday: 5 }, | |
| 90 | + | events: [], | |
| 91 | + | repos: [], | |
| 92 | + | }), | |
| 93 | + | }, | |
| 94 | + | { | |
| 95 | + | match: /\b(support|customers?|tickets?|complaints?|feedback)\b/i, | |
| 96 | + | routine: (duty) => ({ | |
| 97 | + | name: "Support themes", | |
| 98 | + | instructions: `Every weekday morning (${trimmed(duty)}): read yesterday's conversations you can see, group questions and complaints into themes, and post the top ones with how often each came up and whether an issue already covers it.`, | |
| 99 | + | schedule: { every: "weekday", minute: 0, hour: 9, weekday: 1 }, | |
| 100 | + | events: [], | |
| 101 | + | repos: [], | |
| 102 | + | }), | |
| 103 | + | }, | |
| 104 | + | ]; | |
| 105 | + | ||
| 106 | + | /** | |
| 107 | + | * Routines for an agent's responsibilities that it doesn't have yet: one | |
| 108 | + | * per kind, skipping those whose name it already uses or whose events its | |
| 109 | + | * routines already run on. | |
| 110 | + | */ | |
| 111 | + | export function suggestRoutines(responsibilities: string[], existing: { name: string; events: RoutineEvent[] }[]): RoutineSuggestion[] { | |
| 112 | + | const out: RoutineSuggestion[] = []; | |
| 113 | + | const taken = new Set(existing.map((r) => r.name.toLowerCase())); | |
| 114 | + | const covered = new Set(existing.flatMap((r) => r.events)); | |
| 115 | + | for (const duty of responsibilities) { | |
| 116 | + | const rule = RULES.find((r) => r.match.test(duty)); | |
| 117 | + | if (!rule) continue; | |
| 118 | + | const routine = rule.routine(duty); | |
| 119 | + | if (taken.has(routine.name.toLowerCase())) continue; | |
| 120 | + | if (routine.events?.length && routine.events.every((e) => covered.has(e))) continue; | |
| 121 | + | taken.add(routine.name.toLowerCase()); | |
| 122 | + | out.push({ responsibility: duty, routine }); | |
| 123 | + | } | |
| 124 | + | return out; | |
| 125 | + | } |
| 1 | + | import assert from "node:assert/strict"; | |
| 2 | + | import { test } from "node:test"; | |
| 3 | + | ||
| 4 | + | import { follows, happened } from "./happened.ts"; | |
| 5 | + | ||
| 6 | + | const ev = (type: string, data: Record<string, unknown>) => ({ type, repoId: null, data }) as never; | |
| 7 | + | ||
| 8 | + | test("events become the routine events they are, keyed so one thing runs once", () => { | |
| 9 | + | assert.deepEqual(happened(ev("pull.ready", { repoId: "rep_1", number: 12 }))?.key, "pull_ready:rep_1:12"); | |
| 10 | + | assert.equal(happened(ev("pull.opened", { repoId: "rep_1", number: 12 }))?.key, "pull_ready:rep_1:12", "opened and ready are the same thing"); | |
| 11 | + | assert.equal(happened(ev("checks.completed", { repoId: "rep_1", number: 3, status: "passed", commit: "abc" })), null); | |
| 12 | + | assert.equal(happened(ev("checks.completed", { repoId: "rep_1", number: 3, status: "failed", commit: "abc" }))?.key, "checks_failed:rep_1:3:abc"); | |
| 13 | + | assert.notEqual( | |
| 14 | + | happened(ev("checks.completed", { repoId: "rep_1", number: 3, status: "errored", commit: "def" }))?.key, | |
| 15 | + | "checks_failed:rep_1:3:abc", | |
| 16 | + | "a new commit failing is a new thing", | |
| 17 | + | ); | |
| 18 | + | assert.equal(happened(ev("issue.opened", { repoId: "rep_1", number: 9, title: "x" }))?.kind, "issue_opened"); | |
| 19 | + | assert.equal(happened(ev("deployment.failed", { repoId: "rep_1", deploymentId: "dpl_1" }))?.key, "deploy_failed:dpl_1"); | |
| 20 | + | assert.equal(happened(ev("pull.updated", { repoId: "rep_1", number: 1 })), null); | |
| 21 | + | assert.equal(happened(ev("pull.ready", { number: 1 })), null, "no repository, nothing to check access on"); | |
| 22 | + | }); | |
| 23 | + | ||
| 24 | + | test("a routine follows every repository, or only those it names", () => { | |
| 25 | + | assert.equal(follows([], "acme/web"), true); | |
| 26 | + | assert.equal(follows(["acme/web"], "Acme/Web"), true); | |
| 27 | + | assert.equal(follows(["acme/api"], "acme/web"), false); | |
| 28 | + | }); |
| 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 | + | } |
| 22 | 22 | type AgentMemoryScope, | |
| 23 | 23 | type NewRoutine, | |
| 24 | 24 | type Result, | |
| 25 | + | type RoutineSuggestion, | |
| 25 | 26 | type SessionEvent, | |
| 26 | 27 | type SpendSlice, | |
| 27 | 28 | type User, | |
| ⋯ | |||
| 37 | 38 | import { DEFAULT_POLICY, checkPolicy, readPolicy } from "./policy.ts"; | |
| 38 | 39 | import { type RoutineRow, MAX_ROUTINES, checkRoutine, newRoutineId, nextRun, runRoutine, toRoutine } from "./routines.ts"; | |
| 39 | 40 | import { type SessionEnv, type SessionRow, LIVE, approve, sessionRow, steer, stop, toSession } from "./sessions.ts"; | |
| 40 | − | import { type Row, periods, selectAgents, toAgent } from "./store.ts"; | |
| 41 | + | import { type Row, definitionOf, periods, selectAgents, toAgent } from "./store.ts"; | |
| 42 | + | import { suggestRoutines } from "./suggest.ts"; | |
| 41 | 43 | ||
| 42 | 44 | export type ViewContext = { | |
| 43 | 45 | env: SessionEnv; | |
| ⋯ | |||
| 267 | 269 | ||
| 268 | 270 | // ── Routines ────────────────────────────────────────────────────────────── | |
| 269 | 271 | ||
| 270 | − | export async function routines(ctx: ViewContext, handle: string): Promise<Result<AgentRoutine[]>> { | |
| 272 | + | /** An agent's routines, and routines its responsibilities suggest that it doesn't have yet. */ | |
| 273 | + | export async function routines(ctx: ViewContext, handle: string): Promise<Result<{ routines: AgentRoutine[]; suggestions: RoutineSuggestion[] }>> { | |
| 271 | 274 | const agent = await agentByHandle(ctx, handle); | |
| 272 | 275 | if (!agent) return fail("not_found", `There is no agent called @${handle}.`); | |
| 273 | 276 | const rows = await ctx.db.prepare("SELECT * FROM agent_routines WHERE agent_id = ? ORDER BY created_at").bind(agent.id).all<RoutineRow>(); | |
| 274 | − | return ok(rows.results.map(toRoutine)); | |
| 277 | + | const list = rows.results.map(toRoutine); | |
| 278 | + | const duties = definitionOf(agent).responsibilities; | |
| 279 | + | return ok({ routines: list, suggestions: suggestRoutines(duties, list.map((r) => ({ name: r.name, events: r.events }))) }); | |
| 275 | 280 | } | |
| 276 | 281 | ||
| 277 | 282 | /** | |
| ⋯ | |||
| 293 | 298 | const channels = await chatClient(ctx.env.CHAT).sidebar(ctx.slug, ctx.viewer).catch(() => null); | |
| 294 | 299 | const channelName = channels?.ok ? (channelNameIn(channels.value, r.channel_id) ?? null) : null; | |
| 295 | 300 | const now = new Date(); | |
| 296 | − | const next = r.enabled !== false ? nextRun(r.schedule, now).toISOString() : null; | |
| 301 | + | const next = r.enabled !== false && r.schedule ? nextRun(r.schedule, now).toISOString() : null; | |
| 302 | + | const schedule = r.schedule ? JSON.stringify(r.schedule) : null; | |
| 297 | 303 | if (id) { | |
| 298 | 304 | const existing = await ctx.db.prepare("SELECT id FROM agent_routines WHERE id = ? AND agent_id = ?").bind(id, agent.id).first(); | |
| 299 | 305 | if (!existing) return fail("not_found", "There is no such routine."); | |
| 300 | 306 | await ctx.db | |
| 301 | 307 | .prepare( | |
| 302 | − | `UPDATE agent_routines SET name = ?, instructions = ?, schedule = ?, channel_id = ?, channel_name = ?, sponsor = ?, sponsor_username = ?, | |
| 308 | + | `UPDATE agent_routines SET name = ?, instructions = ?, schedule = ?, events = ?, repos = ?, channel_id = ?, channel_name = ?, sponsor = ?, sponsor_username = ?, | |
| 303 | 309 | enabled = ?, paused_note = NULL, next_run_at = ?, workspace = ?, updated_at = ? WHERE id = ?`, | |
| 304 | 310 | ) | |
| 305 | − | .bind(r.name, r.instructions, JSON.stringify(r.schedule), r.channel_id, channelName, ctx.viewer.id, ctx.viewer.username, r.enabled !== false ? 1 : 0, next, ctx.slug, now.toISOString(), id) | |
| 311 | + | .bind(r.name, r.instructions, schedule, JSON.stringify(r.events), JSON.stringify(r.repos), r.channel_id, channelName, ctx.viewer.id, ctx.viewer.username, r.enabled !== false ? 1 : 0, next, ctx.slug, now.toISOString(), id) | |
| 306 | 312 | .run(); | |
| 307 | 313 | } else { | |
| 308 | 314 | const count = await ctx.db.prepare("SELECT COUNT(*) AS n FROM agent_routines WHERE agent_id = ?").bind(agent.id).first<{ n: number }>(); | |
| ⋯ | |||
| 310 | 316 | id = newRoutineId(); | |
| 311 | 317 | await ctx.db | |
| 312 | 318 | .prepare( | |
| 313 | − | `INSERT INTO agent_routines (id, agent_id, workspace_id, workspace, name, instructions, schedule, channel_id, channel_name, sponsor, sponsor_username, enabled, next_run_at, created_at, updated_at) | |
| 314 | − | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, | |
| 319 | + | `INSERT INTO agent_routines (id, agent_id, workspace_id, workspace, name, instructions, schedule, events, repos, channel_id, channel_name, sponsor, sponsor_username, enabled, next_run_at, created_at, updated_at) | |
| 320 | + | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, | |
| 315 | 321 | ) | |
| 316 | − | .bind(id, agent.id, ctx.workspaceId, ctx.slug, r.name, r.instructions, JSON.stringify(r.schedule), r.channel_id, channelName, ctx.viewer.id, ctx.viewer.username, r.enabled !== false ? 1 : 0, next, now.toISOString(), now.toISOString()) | |
| 322 | + | .bind(id, agent.id, ctx.workspaceId, ctx.slug, r.name, r.instructions, schedule, JSON.stringify(r.events), JSON.stringify(r.repos), r.channel_id, channelName, ctx.viewer.id, ctx.viewer.username, r.enabled !== false ? 1 : 0, next, now.toISOString(), now.toISOString()) | |
| 317 | 323 | .run(); | |
| 318 | 324 | } | |
| 319 | 325 | const row = await ctx.db.prepare("SELECT * FROM agent_routines WHERE id = ?").bind(id).first<RoutineRow>(); | |
| 53 | 53 | // Every five minutes: routines that are due, and session steps a desk | |
| 54 | 54 | // lost (src/routines.ts, src/sessions.ts). | |
| 55 | 55 | "triggers": { "crons": ["*/5 * * * *"] }, | |
| 56 | + | // Events routines run on: pull requests ready for review or merged, | |
| 57 | + | // checks and deploys failing, issues opened (src/triggers.ts). Create it | |
| 58 | + | // once: npx wrangler queues create g1t-events-agents | |
| 59 | + | "queues": { | |
| 60 | + | "consumers": [{ "queue": "g1t-events-agents", "max_batch_size": 20, "max_batch_timeout": 2, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }] | |
| 61 | + | }, | |
| 56 | 62 | "vars": { | |
| 57 | 63 | // Who may use g1t's hosted models while billing takes no real money: | |
| 58 | 64 | // the same list as the runner's (services/runner/wrangler.jsonc, and |
| 35 | 35 | { "binding": "SUBSCRIBER_SECURITY", "queue": "g1t-events-security" }, | |
| 36 | 36 | { "binding": "SUBSCRIBER_CONTEXT", "queue": "g1t-events-context" }, | |
| 37 | 37 | { "binding": "SUBSCRIBER_SEARCH", "queue": "g1t-events-search" }, | |
| 38 | − | { "binding": "SUBSCRIBER_PACKAGES", "queue": "g1t-events-packages" } | |
| 38 | + | { "binding": "SUBSCRIBER_PACKAGES", "queue": "g1t-events-packages" }, | |
| 39 | + | // Agents' routines that run on events (services/agents/src/triggers.ts). | |
| 40 | + | { "binding": "SUBSCRIBER_AGENTS", "queue": "g1t-events-agents" } | |
| 39 | 41 | ], | |
| 40 | 42 | "consumers": [{ "queue": "g1t-events", "max_batch_size": 100, "max_batch_timeout": 1, "max_retries": 3, "dead_letter_queue": "g1t-events-dlq" }] | |
| 41 | 43 | }, |