Automations: rules in .g1t/automations that act when something happens
A repository's automations are YAML files on its default branch. Each names what starts it (events, a cron schedule, or a manual run), optional conditions, and steps: comment, label, unlabel, assign or message the agent, open, close or reopen issues, and notify a chat webhook. The new g1t-automations service reads the files when the default branch moves, runs each matching rule once per event, signs its comments, limits runs per hour, and ignores events its own steps caused so rules cannot loop. Schedules run from a one-minute cron. The API gains list, runs, run and update operations (MCP and REST), and the repository sidebar gains an Automations page with run history and templates.
26 files+2154−50/26 viewed
| 884 | 884 | ] | |
| 885 | 885 | ||
| 886 | 886 | [[package]] | |
| 887 | + | name = "g1t-automations" | |
| 888 | + | version = "0.1.0" | |
| 889 | + | dependencies = [ | |
| 890 | + | "g1t-contracts", | |
| 891 | + | "g1t-kit", | |
| 892 | + | "serde", | |
| 893 | + | "serde_json", | |
| 894 | + | "serde_yaml", | |
| 895 | + | "worker", | |
| 896 | + | ] | |
| 897 | + | ||
| 898 | + | [[package]] | |
| 887 | 899 | name = "g1t-billing" | |
| 888 | 900 | version = "0.1.0" | |
| 889 | 901 | dependencies = [ | |
| 2575 | 2587 | ] | |
| 2576 | 2588 | ||
| 2577 | 2589 | [[package]] | |
| 2590 | + | name = "serde_yaml" | |
| 2591 | + | version = "0.9.34+deprecated" | |
| 2592 | + | source = "registry+https://github.com/rust-lang/crates.io-index" | |
| 2593 | + | checksum = "6a8b1a1a2ebf674015cc02edccce75287f1a0130d394307b36743c2f5d504b47" | |
| 2594 | + | dependencies = [ | |
| 2595 | + | "indexmap", | |
| 2596 | + | "itoa", | |
| 2597 | + | "ryu", | |
| 2598 | + | "serde", | |
| 2599 | + | "unsafe-libyaml", | |
| 2600 | + | ] | |
| 2601 | + | ||
| 2602 | + | [[package]] | |
| 2578 | 2603 | name = "serdect" | |
| 2579 | 2604 | version = "0.4.3" | |
| 2580 | 2605 | source = "registry+https://github.com/rust-lang/crates.io-index" | |
| 3104 | 3129 | ] | |
| 3105 | 3130 | ||
| 3106 | 3131 | [[package]] | |
| 3132 | + | name = "unsafe-libyaml" | |
| 3133 | + | version = "0.2.11" | |
| 3134 | + | source = "registry+https://github.com/rust-lang/crates.io-index" | |
| 3135 | + | checksum = "673aac59facbab8a9007c7f6108d11f63b603f7cabff99fabf650fea5c32b861" | |
| 3136 | + | ||
| 3137 | + | [[package]] | |
| 3107 | 3138 | name = "untrusted" | |
| 3108 | 3139 | version = "0.7.1" | |
| 3109 | 3140 | source = "registry+https://github.com/rust-lang/crates.io-index" |
| 1 | 1 | [workspace] | |
| 2 | 2 | resolver = "3" | |
| 3 | − | members = ["apps/api", "crates/*", "services/billing", "services/events", "services/identity", "services/integrations", "services/webhooks", "services/repos", "services/work"] | |
| 3 | + | members = ["apps/api", "crates/*", "services/automations", "services/billing", "services/events", "services/identity", "services/integrations", "services/webhooks", "services/repos", "services/work"] | |
| 4 | 4 | ||
| 5 | 5 | [workspace.package] | |
| 6 | 6 | edition = "2024" |
| 10 | 10 | let name = op.name(); | |
| 11 | 11 | if name.contains("webhook") { | |
| 12 | 12 | "Webhooks" | |
| 13 | + | } else if name.contains("automation") { | |
| 14 | + | "Automations" | |
| 13 | 15 | } else if name.contains("integration") || name.contains("model_routes") || op == Op::GetContext { | |
| 14 | 16 | "Integrations" | |
| 15 | 17 | } else if op == Op::Whoami || name.contains("workspace") { |
| 25 | 25 | pub billing: Fetcher, | |
| 26 | 26 | pub integrations: Fetcher, | |
| 27 | 27 | pub webhooks: Fetcher, | |
| 28 | + | pub automations: Fetcher, | |
| 28 | 29 | /// Set for a request made with an agent's token: all it may do. | |
| 29 | 30 | pub scope: Option<AgentScope>, | |
| 30 | 31 | } | |
| 40 | 41 | billing: env.service("BILLING")?, | |
| 41 | 42 | integrations: env.service("INTEGRATIONS")?, | |
| 42 | 43 | webhooks: env.service("WEBHOOKS")?, | |
| 44 | + | automations: env.service("AUTOMATIONS")?, | |
| 43 | 45 | scope: None, | |
| 44 | 46 | }) | |
| 45 | 47 | } | |
| 97 | 99 | PingWebhook, | |
| 98 | 100 | ListWebhookDeliveries, | |
| 99 | 101 | RedeliverWebhook, | |
| 102 | + | ListAutomations, | |
| 103 | + | ListAutomationRuns, | |
| 104 | + | RunAutomation, | |
| 105 | + | UpdateAutomation, | |
| 100 | 106 | } | |
| 101 | 107 | ||
| 102 | 108 | fn failed(code: FailureCode, message: &str) -> Result<Outcome<Value>> { | |
| 254 | 260 | } | |
| 255 | 261 | ||
| 256 | 262 | impl Op { | |
| 257 | − | pub const ALL: [Op; 50] = [ | |
| 263 | + | pub const ALL: [Op; 54] = [ | |
| 258 | 264 | Op::Whoami, | |
| 259 | 265 | Op::CreateWorkspace, | |
| 260 | 266 | Op::ListRepos, | |
| 305 | 311 | Op::PingWebhook, | |
| 306 | 312 | Op::ListWebhookDeliveries, | |
| 307 | 313 | Op::RedeliverWebhook, | |
| 314 | + | Op::ListAutomations, | |
| 315 | + | Op::ListAutomationRuns, | |
| 316 | + | Op::RunAutomation, | |
| 317 | + | Op::UpdateAutomation, | |
| 308 | 318 | ]; | |
| 309 | 319 | ||
| 310 | 320 | pub fn by_name(name: &str) -> Option<Op> { | |
| 364 | 374 | Op::PingWebhook => "ping_webhook", | |
| 365 | 375 | Op::ListWebhookDeliveries => "list_webhook_deliveries", | |
| 366 | 376 | Op::RedeliverWebhook => "redeliver_webhook", | |
| 377 | + | Op::ListAutomations => "list_automations", | |
| 378 | + | Op::ListAutomationRuns => "list_automation_runs", | |
| 379 | + | Op::RunAutomation => "run_automation", | |
| 380 | + | Op::UpdateAutomation => "update_automation", | |
| 367 | 381 | } | |
| 368 | 382 | } | |
| 369 | 383 | ||
| 496 | 510 | "A webhook's latest deliveries, newest first: what was sent, how the receiver answered, and when it will be tried again." | |
| 497 | 511 | } | |
| 498 | 512 | Op::RedeliverWebhook => "Send a delivery's payload again, as a new delivery.", | |
| 513 | + | Op::ListAutomations => { | |
| 514 | + | "A repository's automations, read from .g1t/automations/*.yml on its default branch: what starts each, its conditions and steps, whether it is on, any problem with its file, and its last run. To add or change one, commit its file." | |
| 515 | + | } | |
| 516 | + | Op::ListAutomationRuns => { | |
| 517 | + | "A repository's latest automation runs, newest first, of one automation or all: what started each, and how each step went or why it was skipped." | |
| 518 | + | } | |
| 519 | + | Op::RunAutomation => { | |
| 520 | + | "Run an automation now, on an issue or pull request if number is given. Members only." | |
| 521 | + | } | |
| 522 | + | Op::UpdateAutomation => "Turn an automation on or off without changing its file. Members only.", | |
| 499 | 523 | Op::ImportIssue => { | |
| 500 | 524 | "Open an issue from a ticket in Jira or Linear, or from a Sentry issue, by its key or address. The issue is linked to it: agents read the original, and when the work lands the ticket is told. Importing the same ticket again returns the issue already made. With assign, a g1t agent starts on it." | |
| 501 | 525 | } | |
| 838 | 862 | ), | |
| 839 | 863 | Op::GetModelRoutes => object(json!({ "workspace": workspace_schema() }), &["workspace"]), | |
| 840 | 864 | Op::ListWebhooks => object(hook_owner(json!({})), &[]), | |
| 865 | + | Op::ListAutomations => repo_only(), | |
| 866 | + | Op::ListAutomationRuns => object( | |
| 867 | + | json!({ | |
| 868 | + | "repo": repo_schema(), | |
| 869 | + | "automation": { "type": "string", "description": "One automation's id, for its runs only." }, | |
| 870 | + | }), | |
| 871 | + | &["repo"], | |
| 872 | + | ), | |
| 873 | + | Op::RunAutomation => object( | |
| 874 | + | json!({ | |
| 875 | + | "repo": repo_schema(), | |
| 876 | + | "id": { "type": "string", "description": "The automation's id." }, | |
| 877 | + | "number": { "type": "integer", "description": "The issue or pull request to run it on." }, | |
| 878 | + | }), | |
| 879 | + | &["repo", "id"], | |
| 880 | + | ), | |
| 881 | + | Op::UpdateAutomation => object( | |
| 882 | + | json!({ | |
| 883 | + | "repo": repo_schema(), | |
| 884 | + | "id": { "type": "string", "description": "The automation's id." }, | |
| 885 | + | "enabled": { "type": "boolean" }, | |
| 886 | + | }), | |
| 887 | + | &["repo", "id", "enabled"], | |
| 888 | + | ), | |
| 841 | 889 | Op::CreateWebhook => object( | |
| 842 | 890 | hook_owner(json!({ | |
| 843 | 891 | "url": { "type": "string", "description": "An HTTPS address on the public internet." }, | |
| 1035 | 1083 | runner, | |
| 1036 | 1084 | integrations, | |
| 1037 | 1085 | webhooks, | |
| 1086 | + | automations, | |
| 1038 | 1087 | .. | |
| 1039 | 1088 | } = services; | |
| 1040 | 1089 | let workspace = || text(input, "workspace").to_lowercase(); | |
| 1458 | 1507 | ) | |
| 1459 | 1508 | .await | |
| 1460 | 1509 | } | |
| 1510 | + | Op::ListAutomations => pass(automations, "list", &json!({ "repo": repo, "viewer": viewer })).await, | |
| 1511 | + | Op::ListAutomationRuns => { | |
| 1512 | + | pass( | |
| 1513 | + | automations, | |
| 1514 | + | "runs", | |
| 1515 | + | &json!({ "repo": repo, "viewer": viewer, "automation": optional_text(input, "automation") }), | |
| 1516 | + | ) | |
| 1517 | + | .await | |
| 1518 | + | } | |
| 1519 | + | Op::RunAutomation => { | |
| 1520 | + | pass( | |
| 1521 | + | automations, | |
| 1522 | + | "run", | |
| 1523 | + | &json!({ "actor": actor(), "repo": repo, "id": text(input, "id"), "number": integer(input, "number") }), | |
| 1524 | + | ) | |
| 1525 | + | .await | |
| 1526 | + | } | |
| 1527 | + | Op::UpdateAutomation => { | |
| 1528 | + | pass( | |
| 1529 | + | automations, | |
| 1530 | + | "set_enabled", | |
| 1531 | + | &json!({ "actor": actor(), "repo": repo, "id": text(input, "id"), "enabled": input["enabled"].as_bool() == Some(true) }), | |
| 1532 | + | ) | |
| 1533 | + | .await | |
| 1534 | + | } | |
| 1461 | 1535 | Op::ListWebhooks | |
| 1462 | 1536 | | Op::CreateWebhook | |
| 1463 | 1537 | | Op::UpdateWebhook |
| 241 | 241 | Op::TestIntegration, | |
| 242 | 242 | &[], | |
| 243 | 243 | ), | |
| 244 | + | route( | |
| 245 | + | "GET", | |
| 246 | + | "/repos/:owner/:name/automations/runs", | |
| 247 | + | Op::ListAutomationRuns, | |
| 248 | + | &[("automation", "automation")], | |
| 249 | + | ), | |
| 250 | + | route( | |
| 251 | + | "GET", | |
| 252 | + | "/repos/:owner/:name/automations", | |
| 253 | + | Op::ListAutomations, | |
| 254 | + | &[], | |
| 255 | + | ), | |
| 256 | + | route( | |
| 257 | + | "POST", | |
| 258 | + | "/repos/:owner/:name/automations/:id/runs", | |
| 259 | + | Op::RunAutomation, | |
| 260 | + | &[], | |
| 261 | + | ), | |
| 262 | + | route( | |
| 263 | + | "PATCH", | |
| 264 | + | "/repos/:owner/:name/automations/:id", | |
| 265 | + | Op::UpdateAutomation, | |
| 266 | + | &[], | |
| 267 | + | ), | |
| 244 | 268 | route("POST", "/repos/:owner/:name/plans", Op::PlanWork, &[]), | |
| 245 | 269 | route("GET", "/repos/:owner/:name/plans/:plan", Op::GetPlan, &[]), | |
| 246 | 270 | route( |
| 20 | 20 | { "binding": "RUNNER", "service": "g1t-runner" }, | |
| 21 | 21 | { "binding": "BILLING", "service": "g1t-billing" }, | |
| 22 | 22 | { "binding": "INTEGRATIONS", "service": "g1t-integrations" }, | |
| 23 | − | { "binding": "WEBHOOKS", "service": "g1t-webhooks" } | |
| 23 | + | { "binding": "WEBHOOKS", "service": "g1t-webhooks" }, | |
| 24 | + | { "binding": "AUTOMATIONS", "service": "g1t-automations" } | |
| 24 | 25 | ], | |
| 25 | 26 | "observability": { "enabled": true } | |
| 26 | 27 | } |
| 26 | 26 | Settings, | |
| 27 | 27 | Users, | |
| 28 | 28 | Webhook, | |
| 29 | + | Zap, | |
| 29 | 30 | X, | |
| 30 | 31 | } from "lucide-react"; | |
| 31 | 32 | import { type ReactNode, useEffect, useMemo, useRef, useState } from "react"; | |
| 475 | 476 | Plan | |
| 476 | 477 | </SidebarLink> | |
| 477 | 478 | )} | |
| 479 | + | <SidebarLink to={`${repoBase}/automations`} icon={<Zap size={14} />}> | |
| 480 | + | Automations | |
| 481 | + | </SidebarLink> | |
| 478 | 482 | {active.member && ( | |
| 479 | 483 | <SidebarLink to={`${repoBase}/settings`} icon={<Settings size={14} />}> | |
| 480 | 484 | Settings | |
| 510 | 514 | queue: "Merge queue", | |
| 511 | 515 | commits: "Commits", | |
| 512 | 516 | plans: "Plan", | |
| 517 | + | automations: "Automations", | |
| 513 | 518 | settings: "Settings", | |
| 514 | 519 | people: "Members", | |
| 515 | 520 | tokens: "Access tokens", |
| 1 | 1 | import { env } from "cloudflare:workers"; | |
| 2 | 2 | ||
| 3 | 3 | import { | |
| 4 | + | automationsClient, | |
| 4 | 5 | billingClient, | |
| 5 | 6 | eventsClient, | |
| 6 | 7 | identityClient, | |
| 17 | 18 | export const events = eventsClient(env.EVENTS); | |
| 18 | 19 | export const integrations = integrationsClient(env.INTEGRATIONS); | |
| 19 | 20 | export const webhooks = webhooksClient(env.WEBHOOKS); | |
| 21 | + | export const automations = automationsClient(env.AUTOMATIONS); |
| 41 | 41 | route("pulls/new", "routes/repo/pull-new.tsx"), | |
| 42 | 42 | route("pull/:number", "routes/repo/pull.tsx"), | |
| 43 | 43 | route("queue", "routes/repo/queue.tsx"), | |
| 44 | + | route("automations", "routes/repo/automations.tsx"), | |
| 44 | 45 | route("plans", "routes/repo/plans.tsx"), | |
| 45 | 46 | route("plans/:id", "routes/repo/plan.tsx"), | |
| 46 | 47 | route("settings", "routes/repo/settings.tsx"), |
| 1 | + | import { CheckCircle2, ChevronRight, CircleSlash, FileCode2, Play, XCircle, Zap } from "lucide-react"; | |
| 2 | + | import { Form, Link, useNavigation } from "react-router"; | |
| 3 | + | ||
| 4 | + | import type { Automation, AutomationRun } from "@g1t/contracts"; | |
| 5 | + | ||
| 6 | + | import type { Route } from "./+types/automations"; | |
| 7 | + | import { Button, EmptyState, ErrorText, TimeAgo } from "../../components/ui"; | |
| 8 | + | import { automations } from "../../lib/services.server"; | |
| 9 | + | import { assertSameOrigin, getViewer, requireUser, roleIn, unwrap } from "../../lib/session.server"; | |
| 10 | + | ||
| 11 | + | export function meta({ params }: Route.MetaArgs) { | |
| 12 | + | return [{ title: `Automations · ${params.owner}/${params.repo} · g1t` }]; | |
| 13 | + | } | |
| 14 | + | ||
| 15 | + | export async function loader({ params, context }: Route.LoaderArgs) { | |
| 16 | + | const viewer = getViewer(context); | |
| 17 | + | const repo = { namespace: params.owner, name: params.repo }; | |
| 18 | + | const [list, runs] = await Promise.all([automations.list(repo, viewer), automations.runs(repo, viewer)]); | |
| 19 | + | return { | |
| 20 | + | automations: unwrap(list), | |
| 21 | + | runs: runs.ok ? runs.value : [], | |
| 22 | + | member: roleIn(viewer, params.owner) != null, | |
| 23 | + | }; | |
| 24 | + | } | |
| 25 | + | ||
| 26 | + | export async function action({ request, params, context }: Route.ActionArgs) { | |
| 27 | + | assertSameOrigin(request); | |
| 28 | + | const user = requireUser(context, request); | |
| 29 | + | const repo = { namespace: params.owner, name: params.repo }; | |
| 30 | + | const form = await request.formData(); | |
| 31 | + | const id = String(form.get("id") ?? ""); | |
| 32 | + | if (form.get("intent") === "toggle") { | |
| 33 | + | const changed = await automations.setEnabled(user, repo, id, form.get("enabled") === "true"); | |
| 34 | + | return changed.ok ? {} : { error: changed.error.message }; | |
| 35 | + | } | |
| 36 | + | const number = Number(form.get("number")); | |
| 37 | + | const ran = await automations.run(user, repo, id, Number.isInteger(number) && number > 0 ? number : undefined); | |
| 38 | + | return ran.ok ? { ran: ran.value } : { error: ran.error.message }; | |
| 39 | + | } | |
| 40 | + | ||
| 41 | + | /** Ready-made automations to start from. */ | |
| 42 | + | const TEMPLATES: { file: string; title: string; about: string; yaml: string }[] = [ | |
| 43 | + | { | |
| 44 | + | file: "bugs.yml", | |
| 45 | + | title: "Put an agent on every bug", | |
| 46 | + | about: "A new issue labelled bug gets a g1t agent at once, and a note saying so.", | |
| 47 | + | yaml: `name: Put an agent on every bug | |
| 48 | + | on: issue.opened | |
| 49 | + | if: | |
| 50 | + | labels: bug | |
| 51 | + | do: | |
| 52 | + | - comment: "Thanks, {{actor}}. A g1t agent is on it." | |
| 53 | + | - assign_agent`, | |
| 54 | + | }, | |
| 55 | + | { | |
| 56 | + | file: "failed-checks.yml", | |
| 57 | + | title: "Tell the agent when checks fail", | |
| 58 | + | about: "When a pull request's checks fail, the agent working on it hears at once.", | |
| 59 | + | yaml: `name: Tell the agent when checks fail | |
| 60 | + | on: checks.completed | |
| 61 | + | if: | |
| 62 | + | status: failed | |
| 63 | + | do: | |
| 64 | + | - message_agent: "The acceptance checks failed. Read their output on #{{number}} and fix the cause."`, | |
| 65 | + | }, | |
| 66 | + | { | |
| 67 | + | file: "announce-merges.yml", | |
| 68 | + | title: "Announce merges in chat", | |
| 69 | + | about: "Every merge posts to a Slack or Discord channel through its incoming webhook.", | |
| 70 | + | yaml: `name: Announce merges | |
| 71 | + | on: pull.merged | |
| 72 | + | do: | |
| 73 | + | - notify: | |
| 74 | + | url: https://hooks.slack.com/services/... | |
| 75 | + | text: "Merged into {{repo}}: {{title}} {{url}}"`, | |
| 76 | + | }, | |
| 77 | + | { | |
| 78 | + | file: "weekly-tidy.yml", | |
| 79 | + | title: "A weekly tidy-up", | |
| 80 | + | about: "Every Monday morning an issue opens and an agent takes it.", | |
| 81 | + | yaml: `name: Weekly tidy-up | |
| 82 | + | on: | |
| 83 | + | schedule: "0 9 * * mon" | |
| 84 | + | do: | |
| 85 | + | - open_issue: | |
| 86 | + | title: Weekly tidy-up | |
| 87 | + | body: Update dependencies that have new patch releases, and fix any new warnings. | |
| 88 | + | labels: chore | |
| 89 | + | assign_agent: true`, | |
| 90 | + | }, | |
| 91 | + | ]; | |
| 92 | + | ||
| 93 | + | function RunStatus({ run }: { run: AutomationRun }) { | |
| 94 | + | if (run.status === "succeeded") return <CheckCircle2 size={15} className="shrink-0 text-accent" />; | |
| 95 | + | if (run.status === "failed") return <XCircle size={15} className="shrink-0 text-danger" />; | |
| 96 | + | return <CircleSlash size={15} className="shrink-0 text-faint" />; | |
| 97 | + | } | |
| 98 | + | ||
| 99 | + | function RunRow({ run, base }: { run: AutomationRun; base: string }) { | |
| 100 | + | return ( | |
| 101 | + | <details className="group border-t border-line first:border-t-0"> | |
| 102 | + | <summary className="flex cursor-pointer list-none items-center gap-3 px-4 py-2.5 text-sm hover:bg-raised/40"> | |
| 103 | + | <ChevronRight size={14} className="shrink-0 text-faint transition-transform group-open:rotate-90" /> | |
| 104 | + | <RunStatus run={run} /> | |
| 105 | + | <span className="min-w-0 truncate font-medium">{run.name}</span> | |
| 106 | + | <span className="shrink-0 font-mono text-xs text-muted">{run.event}</span> | |
| 107 | + | {run.number != null && ( | |
| 108 | + | <Link to={`${base}/issues/${run.number}`} className="shrink-0 font-mono text-xs text-muted hover:text-fg"> | |
| 109 | + | #{run.number} | |
| 110 | + | </Link> | |
| 111 | + | )} | |
| 112 | + | <span className="ml-auto shrink-0 text-xs text-faint"> | |
| 113 | + | {run.actor && `${run.actor} · `} | |
| 114 | + | <TimeAgo at={run.startedAt} /> | |
| 115 | + | </span> | |
| 116 | + | </summary> | |
| 117 | + | <div className="border-t border-line bg-bg/40 px-4 py-3 text-sm"> | |
| 118 | + | {run.reason && <p className="text-muted">{run.reason}</p>} | |
| 119 | + | {run.steps.length > 0 && ( | |
| 120 | + | <ol className="space-y-1.5"> | |
| 121 | + | {run.steps.map((step, index) => ( | |
| 122 | + | <li key={index} className="flex items-start gap-2"> | |
| 123 | + | {step.ok ? ( | |
| 124 | + | <CheckCircle2 size={14} className="mt-0.5 shrink-0 text-accent" /> | |
| 125 | + | ) : ( | |
| 126 | + | <XCircle size={14} className="mt-0.5 shrink-0 text-danger" /> | |
| 127 | + | )} | |
| 128 | + | <span> | |
| 129 | + | <span className="font-medium">{step.step}</span> | |
| 130 | + | <span className="text-muted"> · {step.detail}</span> | |
| 131 | + | </span> | |
| 132 | + | </li> | |
| 133 | + | ))} | |
| 134 | + | </ol> | |
| 135 | + | )} | |
| 136 | + | </div> | |
| 137 | + | </details> | |
| 138 | + | ); | |
| 139 | + | } | |
| 140 | + | ||
| 141 | + | function AutomationCard({ automation, base, member }: { automation: Automation; base: string; member: boolean }) { | |
| 142 | + | const busy = useNavigation().state === "submitting"; | |
| 143 | + | return ( | |
| 144 | + | <li className="border-t border-line p-4 first:border-t-0"> | |
| 145 | + | <div className="flex flex-wrap items-start gap-3"> | |
| 146 | + | <span | |
| 147 | + | className={`mt-0.5 inline-flex size-8 shrink-0 items-center justify-center rounded-lg ring-1 ${ | |
| 148 | + | automation.error | |
| 149 | + | ? "bg-danger/10 text-danger ring-danger/30" | |
| 150 | + | : automation.enabled | |
| 151 | + | ? "bg-accent/10 text-accent ring-accent/30" | |
| 152 | + | : "bg-surface text-faint ring-line" | |
| 153 | + | }`} | |
| 154 | + | > | |
| 155 | + | <Zap size={15} /> | |
| 156 | + | </span> | |
| 157 | + | <div className="min-w-0 grow"> | |
| 158 | + | <p className="flex flex-wrap items-center gap-2"> | |
| 159 | + | <span className="font-medium">{automation.name}</span> | |
| 160 | + | {!automation.enabled && !automation.error && ( | |
| 161 | + | <span className="rounded-full px-2 py-px text-xs text-muted ring-1 ring-line">Off</span> | |
| 162 | + | )} | |
| 163 | + | </p> | |
| 164 | + | <Link | |
| 165 | + | to={`${base}/blob/HEAD/${automation.path}`} | |
| 166 | + | className="inline-flex items-center gap-1 font-mono text-xs text-faint hover:text-fg" | |
| 167 | + | > | |
| 168 | + | <FileCode2 size={12} /> | |
| 169 | + | {automation.path} | |
| 170 | + | </Link> | |
| 171 | + | {automation.error ? ( | |
| 172 | + | <p className="mt-2 text-sm text-danger">{automation.error}</p> | |
| 173 | + | ) : ( | |
| 174 | + | <p className="mt-2 text-sm text-muted"> | |
| 175 | + | <span className="text-fg">On</span> {automation.trigger} | |
| 176 | + | {automation.conditions.length > 0 && ( | |
| 177 | + | <> | |
| 178 | + | , <span className="text-fg">if</span> {automation.conditions.join(" and ")} | |
| 179 | + | </> | |
| 180 | + | )} | |
| 181 | + | : {automation.steps.join(", then ")}. | |
| 182 | + | </p> | |
| 183 | + | )} | |
| 184 | + | {automation.lastRun && ( | |
| 185 | + | <p className="mt-2 flex items-center gap-1.5 text-xs text-faint"> | |
| 186 | + | <RunStatus run={automation.lastRun} /> | |
| 187 | + | Last run {automation.lastRun.status} <TimeAgo at={automation.lastRun.startedAt} /> | |
| 188 | + | {automation.lastRun.reason && ` · ${automation.lastRun.reason}`} | |
| 189 | + | </p> | |
| 190 | + | )} | |
| 191 | + | </div> | |
| 192 | + | {member && !automation.error && ( | |
| 193 | + | <div className="flex shrink-0 items-center gap-2"> | |
| 194 | + | <Form method="post" className="flex items-center gap-1.5"> | |
| 195 | + | <input type="hidden" name="id" value={automation.id} /> | |
| 196 | + | <input | |
| 197 | + | name="number" | |
| 198 | + | inputMode="numeric" | |
| 199 | + | placeholder="#" | |
| 200 | + | aria-label="Issue or pull request to run it on" | |
| 201 | + | className="w-14 rounded-md border border-line bg-bg px-2 py-1.5 text-center font-mono text-xs outline-none focus:border-accent-dim" | |
| 202 | + | /> | |
| 203 | + | <Button type="submit" variant="quiet" disabled={busy} name="intent" value="run"> | |
| 204 | + | <Play size={13} /> | |
| 205 | + | Run now | |
| 206 | + | </Button> | |
| 207 | + | </Form> | |
| 208 | + | <Form method="post"> | |
| 209 | + | <input type="hidden" name="intent" value="toggle" /> | |
| 210 | + | <input type="hidden" name="id" value={automation.id} /> | |
| 211 | + | <input type="hidden" name="enabled" value={automation.enabled ? "false" : "true"} /> | |
| 212 | + | <Button type="submit" variant="quiet" disabled={busy}> | |
| 213 | + | {automation.enabled ? "Turn off" : "Turn on"} | |
| 214 | + | </Button> | |
| 215 | + | </Form> | |
| 216 | + | </div> | |
| 217 | + | )} | |
| 218 | + | </div> | |
| 219 | + | </li> | |
| 220 | + | ); | |
| 221 | + | } | |
| 222 | + | ||
| 223 | + | function Templates({ repo }: { repo: string }) { | |
| 224 | + | return ( | |
| 225 | + | <section> | |
| 226 | + | <h3 className="text-sm font-medium">Start from one of these</h3> | |
| 227 | + | <p className="mt-1 text-sm text-muted"> | |
| 228 | + | Add the file to <code className="text-fg">.g1t/automations/</code> on {repo}'s default branch. It is read the moment | |
| 229 | + | you push. | |
| 230 | + | </p> | |
| 231 | + | <div className="mt-4 grid gap-3 lg:grid-cols-2"> | |
| 232 | + | {TEMPLATES.map((template) => ( | |
| 233 | + | <div key={template.file} className="rounded-xl border border-line bg-surface p-4"> | |
| 234 | + | <p className="font-medium">{template.title}</p> | |
| 235 | + | <p className="mt-1 text-xs text-muted">{template.about}</p> | |
| 236 | + | <p className="mt-3 font-mono text-xs text-faint">.g1t/automations/{template.file}</p> | |
| 237 | + | <pre className="mt-1.5 overflow-x-auto rounded-lg bg-bg p-3 font-mono text-xs leading-relaxed ring-1 ring-line"> | |
| 238 | + | <code>{template.yaml}</code> | |
| 239 | + | </pre> | |
| 240 | + | </div> | |
| 241 | + | ))} | |
| 242 | + | </div> | |
| 243 | + | </section> | |
| 244 | + | ); | |
| 245 | + | } | |
| 246 | + | ||
| 247 | + | export default function Automations({ loaderData, actionData, params }: Route.ComponentProps) { | |
| 248 | + | const { automations: list, runs, member } = loaderData; | |
| 249 | + | const base = `/${params.owner}/${params.repo}`; | |
| 250 | + | const ran = actionData && "ran" in actionData ? actionData.ran : null; | |
| 251 | + | return ( | |
| 252 | + | <div className="max-w-5xl space-y-10"> | |
| 253 | + | <header className="flex flex-wrap items-start justify-between gap-4"> | |
| 254 | + | <div> | |
| 255 | + | <h2 className="flex items-center gap-2 text-xl font-semibold tracking-tight"> | |
| 256 | + | <Zap size={18} className="text-accent" /> | |
| 257 | + | Automations | |
| 258 | + | </h2> | |
| 259 | + | <p className="mt-1 max-w-2xl text-sm text-muted"> | |
| 260 | + | Rules in <code className="text-fg">.g1t/automations</code> that act when something happens: comment, label, | |
| 261 | + | put an agent on it, tell the agent working on it, open or close issues, post to chat. | |
| 262 | + | </p> | |
| 263 | + | </div> | |
| 264 | + | <a href="https://docs.g1t.sh/guides/automations/" className="text-sm text-muted underline underline-offset-4 hover:text-fg"> | |
| 265 | + | How automations work | |
| 266 | + | </a> | |
| 267 | + | </header> | |
| 268 | + | ||
| 269 | + | {ran && ( | |
| 270 | + | <p className="text-sm text-muted"> | |
| 271 | + | Ran <span className="text-fg">{ran.name}</span>: {ran.status} | |
| 272 | + | {ran.reason ? `. ${ran.reason}` : "."} | |
| 273 | + | </p> | |
| 274 | + | )} | |
| 275 | + | <ErrorText>{actionData && "error" in actionData ? actionData.error : null}</ErrorText> | |
| 276 | + | ||
| 277 | + | {list.length === 0 ? ( | |
| 278 | + | <EmptyState title="No automations yet"> | |
| 279 | + | Commit a file to <code>.g1t/automations/</code> and it shows up here. | |
| 280 | + | </EmptyState> | |
| 281 | + | ) : ( | |
| 282 | + | <ul className="overflow-hidden rounded-xl border border-line bg-surface"> | |
| 283 | + | {list.map((automation) => ( | |
| 284 | + | <AutomationCard key={automation.id} automation={automation} base={base} member={member} /> | |
| 285 | + | ))} | |
| 286 | + | </ul> | |
| 287 | + | )} | |
| 288 | + | ||
| 289 | + | {runs.length > 0 && ( | |
| 290 | + | <section> | |
| 291 | + | <h3 className="mb-3 text-sm font-medium">Recent runs</h3> | |
| 292 | + | <div className="overflow-hidden rounded-xl border border-line bg-surface"> | |
| 293 | + | {runs.map((run) => ( | |
| 294 | + | <RunRow key={run.id} run={run} base={base} /> | |
| 295 | + | ))} | |
| 296 | + | </div> | |
| 297 | + | </section> | |
| 298 | + | )} | |
| 299 | + | ||
| 300 | + | <Templates repo={`${params.owner}/${params.repo}`} /> | |
| 301 | + | </div> | |
| 302 | + | ); | |
| 303 | + | } |
| 12 | 12 | EVENTS: ServiceBinding; | |
| 13 | 13 | INTEGRATIONS: ServiceBinding; | |
| 14 | 14 | WEBHOOKS: ServiceBinding; | |
| 15 | + | AUTOMATIONS: ServiceBinding; | |
| 15 | 16 | } | |
| 16 | 17 | } | |
| 17 | 18 | interface Env extends Cloudflare.Env {} |
| 17 | 17 | { "binding": "BILLING", "service": "g1t-billing" }, | |
| 18 | 18 | { "binding": "EVENTS", "service": "g1t-events" }, | |
| 19 | 19 | { "binding": "INTEGRATIONS", "service": "g1t-integrations" }, | |
| 20 | − | { "binding": "WEBHOOKS", "service": "g1t-webhooks" } | |
| 20 | + | { "binding": "WEBHOOKS", "service": "g1t-webhooks" }, | |
| 21 | + | { "binding": "AUTOMATIONS", "service": "g1t-automations" } | |
| 21 | 22 | ], | |
| 22 | 23 | "observability": { "enabled": true }, | |
| 23 | 24 | "upload_source_maps": true |
| 1 | + | //! The automations service: rules in a repository's `.g1t/automations/` | |
| 2 | + | //! that act when something happens, in the way GitHub Actions' workflows | |
| 3 | + | //! do, but on g1t's events and with g1t's own steps. | |
| 4 | + | //! | |
| 5 | + | //! An automation says *on* which event (or a schedule, or by hand), *if* | |
| 6 | + | //! which conditions hold, *do* which steps: comment, label, put an agent | |
| 7 | + | //! on it, message the agent working on it, open, close or reopen an issue, | |
| 8 | + | //! post to a URL. Files on the default branch are the source of truth; each | |
| 9 | + | //! push reloads them. Every run is kept, step by step. | |
| 10 | + | //! | |
| 11 | + | //! Mirrors `packages/contracts/src/automations.ts`. | |
| 12 | + | ||
| 13 | + | use serde::{Deserialize, Serialize}; | |
| 14 | + | ||
| 15 | + | use crate::repos::RepoPath; | |
| 16 | + | use crate::{User, Viewer}; | |
| 17 | + | ||
| 18 | + | /// One automation, as read from its file. | |
| 19 | + | #[derive(Clone, Debug, Serialize, Deserialize)] | |
| 20 | + | #[serde(rename_all = "camelCase")] | |
| 21 | + | pub struct Automation { | |
| 22 | + | pub id: String, | |
| 23 | + | /// `owner/name`. | |
| 24 | + | pub repo: String, | |
| 25 | + | /// The file it comes from, such as `.g1t/automations/bugs.yml`. | |
| 26 | + | pub path: String, | |
| 27 | + | pub name: String, | |
| 28 | + | /// What starts it, in words: `issue.opened`, `every Monday at 09:00`. | |
| 29 | + | pub trigger: String, | |
| 30 | + | /// Its conditions and steps, in words, one each. | |
| 31 | + | pub conditions: Vec<String>, | |
| 32 | + | pub steps: Vec<String>, | |
| 33 | + | /// Whether it runs. Members can turn one off without changing its file. | |
| 34 | + | pub enabled: bool, | |
| 35 | + | /// Why the file could not be used, if it could not. | |
| 36 | + | pub error: Option<String>, | |
| 37 | + | /// Whether it can be run by hand. | |
| 38 | + | pub manual: bool, | |
| 39 | + | pub last_run: Option<AutomationRun>, | |
| 40 | + | } | |
| 41 | + | ||
| 42 | + | /// One step of a run, and how it went. | |
| 43 | + | #[derive(Clone, Debug, Serialize, Deserialize)] | |
| 44 | + | #[serde(rename_all = "camelCase")] | |
| 45 | + | pub struct StepResult { | |
| 46 | + | pub step: String, | |
| 47 | + | pub ok: bool, | |
| 48 | + | pub detail: String, | |
| 49 | + | } | |
| 50 | + | ||
| 51 | + | #[derive(Clone, Debug, Serialize, Deserialize)] | |
| 52 | + | #[serde(rename_all = "camelCase")] | |
| 53 | + | pub struct AutomationRun { | |
| 54 | + | pub id: String, | |
| 55 | + | pub automation_id: String, | |
| 56 | + | pub name: String, | |
| 57 | + | /// What started it: an event type, `schedule` or `manual`. | |
| 58 | + | pub event: String, | |
| 59 | + | /// The issue or pull request it acted on. | |
| 60 | + | pub number: Option<u32>, | |
| 61 | + | /// `succeeded`, `failed` (a step failed) or `skipped` (and why, in | |
| 62 | + | /// `reason`). | |
| 63 | + | pub status: String, | |
| 64 | + | pub reason: Option<String>, | |
| 65 | + | pub steps: Vec<StepResult>, | |
| 66 | + | /// Who started it by hand, or who caused the event. | |
| 67 | + | pub actor: Option<String>, | |
| 68 | + | /// RFC 3339. | |
| 69 | + | pub started_at: String, | |
| 70 | + | } | |
| 71 | + | ||
| 72 | + | /// `list`: a repository's automations. Returns `Outcome<Vec<Automation>>`. | |
| 73 | + | /// Anyone who can see the repository. | |
| 74 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 75 | + | pub struct ListArgs { | |
| 76 | + | pub repo: RepoPath, | |
| 77 | + | pub viewer: Viewer, | |
| 78 | + | } | |
| 79 | + | ||
| 80 | + | /// `runs`: a repository's latest runs, newest first, of one automation or | |
| 81 | + | /// all. Returns `Outcome<Vec<AutomationRun>>`. | |
| 82 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 83 | + | pub struct RunsArgs { | |
| 84 | + | pub repo: RepoPath, | |
| 85 | + | pub viewer: Viewer, | |
| 86 | + | #[serde(default)] | |
| 87 | + | pub automation: Option<String>, | |
| 88 | + | } | |
| 89 | + | ||
| 90 | + | /// `run`: runs an automation now, on an issue or pull request if given. | |
| 91 | + | /// Returns `Outcome<AutomationRun>`. Members only. | |
| 92 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 93 | + | pub struct RunArgs { | |
| 94 | + | pub actor: User, | |
| 95 | + | pub repo: RepoPath, | |
| 96 | + | pub id: String, | |
| 97 | + | #[serde(default)] | |
| 98 | + | pub number: Option<u32>, | |
| 99 | + | } | |
| 100 | + | ||
| 101 | + | /// `set_enabled`: turns an automation on or off. Returns | |
| 102 | + | /// `Outcome<Automation>`. Members only. | |
| 103 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 104 | + | pub struct SetEnabledArgs { | |
| 105 | + | pub actor: User, | |
| 106 | + | pub repo: RepoPath, | |
| 107 | + | pub id: String, | |
| 108 | + | pub enabled: bool, | |
| 109 | + | } |
| 4 | 4 | //! arguments of each of its methods. Services and their callers depend on | |
| 5 | 5 | //! this crate, never on each other's code. | |
| 6 | 6 | ||
| 7 | + | pub mod automations; | |
| 7 | 8 | pub mod billing; | |
| 8 | 9 | pub mod events; | |
| 9 | 10 | pub mod identity; |
| 131 | 131 | pub viewer: Viewer, | |
| 132 | 132 | } | |
| 133 | 133 | ||
| 134 | + | /// `path_by_id`: where a repository is, whoever may see it. For g1t's own | |
| 135 | + | /// services, which hold a repository's id from an event and act for its | |
| 136 | + | /// workspace; nothing outside reaches it. Returns `Option<RepoPath>`, null | |
| 137 | + | /// for a fork or an unknown id. | |
| 138 | + | #[derive(Debug, Serialize, Deserialize)] | |
| 139 | + | pub struct PathByIdArgs { | |
| 140 | + | pub id: String, | |
| 141 | + | } | |
| 142 | + | ||
| 134 | 143 | /// `get_by_id`. Returns `Outcome<Repo>`. | |
| 135 | 144 | #[derive(Debug, Serialize, Deserialize)] | |
| 136 | 145 | pub struct GetByIdArgs { |
| 1 | + | import type { User, Viewer } from "./identity"; | |
| 2 | + | import type { RepoPath } from "./repos"; | |
| 3 | + | import type { Result } from "./result"; | |
| 4 | + | ||
| 5 | + | /** | |
| 6 | + | * Rules in a repository's `.g1t/automations/` that act when something | |
| 7 | + | * happens. Mirrors `crates/contracts/src/automations.rs`. | |
| 8 | + | */ | |
| 9 | + | ||
| 10 | + | export type StepResult = { step: string; ok: boolean; detail: string }; | |
| 11 | + | ||
| 12 | + | export type AutomationRun = { | |
| 13 | + | id: string; | |
| 14 | + | automationId: string; | |
| 15 | + | name: string; | |
| 16 | + | /** An event type, `schedule` or `manual`. */ | |
| 17 | + | event: string; | |
| 18 | + | number: number | null; | |
| 19 | + | status: "running" | "succeeded" | "failed" | "skipped"; | |
| 20 | + | reason: string | null; | |
| 21 | + | steps: StepResult[]; | |
| 22 | + | actor: string | null; | |
| 23 | + | startedAt: string; | |
| 24 | + | }; | |
| 25 | + | ||
| 26 | + | export type Automation = { | |
| 27 | + | id: string; | |
| 28 | + | repo: string; | |
| 29 | + | path: string; | |
| 30 | + | name: string; | |
| 31 | + | /** What starts it, in words. */ | |
| 32 | + | trigger: string; | |
| 33 | + | conditions: string[]; | |
| 34 | + | steps: string[]; | |
| 35 | + | enabled: boolean; | |
| 36 | + | /** Why its file cannot be used. */ | |
| 37 | + | error: string | null; | |
| 38 | + | manual: boolean; | |
| 39 | + | lastRun: AutomationRun | null; | |
| 40 | + | }; | |
| 41 | + | ||
| 42 | + | export interface AutomationsApi { | |
| 43 | + | list(repo: RepoPath, viewer: Viewer): Promise<Result<Automation[]>>; | |
| 44 | + | runs(repo: RepoPath, viewer: Viewer, automation?: string): Promise<Result<AutomationRun[]>>; | |
| 45 | + | run(actor: User, repo: RepoPath, id: string, number?: number): Promise<Result<AutomationRun>>; | |
| 46 | + | setEnabled(actor: User, repo: RepoPath, id: string, enabled: boolean): Promise<Result<Automation>>; | |
| 47 | + | } |
| 1 | + | import type { AutomationsApi } from "./automations"; | |
| 1 | 2 | import type { BillingApi } from "./billing"; | |
| 2 | 3 | import type { EventsApi } from "./events"; | |
| 3 | 4 | import type { IdentityApi } from "./identity"; | |
| 229 | 230 | redeliver: (actor, owner, deliveryId) => call("redeliver", { actor, ...owner, deliveryId }), | |
| 230 | 231 | }; | |
| 231 | 232 | } | |
| 233 | + | ||
| 234 | + | export function automationsClient(service: ServiceBinding): AutomationsApi { | |
| 235 | + | const call = <T>(method: string, args: object) => rpc<T>(service, method, args); | |
| 236 | + | return { | |
| 237 | + | list: (repo, viewer) => call("list", { repo, viewer }), | |
| 238 | + | runs: (repo, viewer, automation) => call("runs", { repo, viewer, automation }), | |
| 239 | + | run: (actor, repo, id, number) => call("run", { actor, repo, id, number }), | |
| 240 | + | setEnabled: (actor, repo, id, enabled) => call("set_enabled", { actor, repo, id, enabled }), | |
| 241 | + | }; | |
| 242 | + | } |
| 1 | + | export * from "./automations"; | |
| 1 | 2 | export * from "./billing"; | |
| 2 | 3 | export * from "./clients"; | |
| 3 | 4 | export * from "./events"; |
| 1 | + | [package] | |
| 2 | + | name = "g1t-automations" | |
| 3 | + | version = "0.1.0" | |
| 4 | + | edition.workspace = true | |
| 5 | + | license.workspace = true | |
| 6 | + | description = "Rules in a repository's .g1t/automations that act when something happens." | |
| 7 | + | ||
| 8 | + | [lib] | |
| 9 | + | crate-type = ["cdylib"] | |
| 10 | + | ||
| 11 | + | [dependencies] | |
| 12 | + | g1t-contracts.workspace = true | |
| 13 | + | g1t-kit.workspace = true | |
| 14 | + | serde.workspace = true | |
| 15 | + | serde_json.workspace = true | |
| 16 | + | serde_yaml = "0.9" | |
| 17 | + | worker.workspace = true |
| 1 | + | -- Automations, read from repositories' .g1t/automations, and their runs. | |
| 2 | + | -- Every timestamp is RFC 3339 UTC. | |
| 3 | + | ||
| 4 | + | CREATE TABLE automations ( | |
| 5 | + | id TEXT PRIMARY KEY, | |
| 6 | + | repo_id TEXT NOT NULL, | |
| 7 | + | -- owner/name. | |
| 8 | + | repo TEXT NOT NULL, | |
| 9 | + | -- The file it comes from. | |
| 10 | + | path TEXT NOT NULL, | |
| 11 | + | name TEXT NOT NULL, | |
| 12 | + | -- The file as it is on the default branch. | |
| 13 | + | source TEXT NOT NULL, | |
| 14 | + | -- events, schedule, manual, or invalid. | |
| 15 | + | trigger_kind TEXT NOT NULL, | |
| 16 | + | -- For events: the types that start it, as a JSON array. | |
| 17 | + | events TEXT NOT NULL, | |
| 18 | + | -- Why the file cannot be used, if it cannot. | |
| 19 | + | error TEXT, | |
| 20 | + | -- Kept across reloads of the file: a member can turn one off. | |
| 21 | + | enabled INTEGER NOT NULL DEFAULT 1, | |
| 22 | + | updated_at TEXT NOT NULL, | |
| 23 | + | UNIQUE (repo_id, path) | |
| 24 | + | ); | |
| 25 | + | CREATE INDEX automations_by_trigger ON automations (repo_id, enabled, trigger_kind); | |
| 26 | + | CREATE INDEX automations_scheduled ON automations (trigger_kind, enabled); | |
| 27 | + | ||
| 28 | + | CREATE TABLE runs ( | |
| 29 | + | id TEXT PRIMARY KEY, | |
| 30 | + | automation_id TEXT NOT NULL, | |
| 31 | + | repo_id TEXT NOT NULL, | |
| 32 | + | -- What started it, unique per automation: an event's id, the minute of a | |
| 33 | + | -- schedule, or a manual run's own key. An event runs an automation once. | |
| 34 | + | event_key TEXT NOT NULL, | |
| 35 | + | name TEXT NOT NULL, | |
| 36 | + | event TEXT NOT NULL, | |
| 37 | + | number INTEGER, | |
| 38 | + | -- running, succeeded, failed or skipped. | |
| 39 | + | status TEXT NOT NULL, | |
| 40 | + | reason TEXT, | |
| 41 | + | -- Each step's result, as JSON. | |
| 42 | + | steps TEXT NOT NULL, | |
| 43 | + | actor TEXT, | |
| 44 | + | started_at TEXT NOT NULL, | |
| 45 | + | UNIQUE (automation_id, event_key) | |
| 46 | + | ); | |
| 47 | + | CREATE INDEX runs_by_repo ON runs (repo_id, id); | |
| 48 | + | CREATE INDEX runs_by_automation ON runs (automation_id, id); | |
| 49 | + | ||
| 50 | + | -- What an automation just touched, so it does not answer its own doing. | |
| 51 | + | CREATE TABLE effects ( | |
| 52 | + | automation_id TEXT NOT NULL, | |
| 53 | + | repo_id TEXT NOT NULL, | |
| 54 | + | number INTEGER NOT NULL, | |
| 55 | + | at TEXT NOT NULL | |
| 56 | + | ); | |
| 57 | + | CREATE INDEX effects_recent ON effects (automation_id, repo_id, number, at); | |
| 58 | + | ||
| 59 | + | -- Repositories whose files have been read at least once. | |
| 60 | + | CREATE TABLE synced ( | |
| 61 | + | repo_id TEXT PRIMARY KEY, | |
| 62 | + | at TEXT NOT NULL | |
| 63 | + | ); |
| 1 | + | //! Five-field cron schedules, in UTC: minute, hour, day of month, month, | |
| 2 | + | //! day of week. Each field takes `*`, a number, a range `a-b`, a list | |
| 3 | + | //! `a,b`, and a step `*/n` or `a-b/n`; days of the week also take `mon` to | |
| 4 | + | //! `sun`. | |
| 5 | + | ||
| 6 | + | #[derive(Clone, Debug, PartialEq, Eq)] | |
| 7 | + | pub struct Schedule { | |
| 8 | + | minutes: Vec<bool>, | |
| 9 | + | hours: Vec<bool>, | |
| 10 | + | days: Vec<bool>, | |
| 11 | + | months: Vec<bool>, | |
| 12 | + | weekdays: Vec<bool>, | |
| 13 | + | /// Whether day of month and day of week were each restricted: when both | |
| 14 | + | /// are, either matching is enough, as in every cron. | |
| 15 | + | days_restricted: bool, | |
| 16 | + | weekdays_restricted: bool, | |
| 17 | + | } | |
| 18 | + | ||
| 19 | + | const WEEKDAYS: [&str; 7] = ["sun", "mon", "tue", "wed", "thu", "fri", "sat"]; | |
| 20 | + | ||
| 21 | + | fn field(text: &str, low: u32, high: u32, names: &[&str]) -> Result<(Vec<bool>, bool), String> { | |
| 22 | + | let mut set = vec![false; (high + 1) as usize]; | |
| 23 | + | let value = |part: &str| -> Result<u32, String> { | |
| 24 | + | if let Some(at) = names.iter().position(|name| part.eq_ignore_ascii_case(name)) { | |
| 25 | + | return Ok(at as u32); | |
| 26 | + | } | |
| 27 | + | part.parse::<u32>().map_err(|_| format!("`{part}` is not a number")) | |
| 28 | + | }; | |
| 29 | + | for item in text.split(',') { | |
| 30 | + | let (range, step) = match item.split_once('/') { | |
| 31 | + | Some((range, step)) => (range, step.parse::<u32>().map_err(|_| format!("`{step}` is not a step"))?), | |
| 32 | + | None => (item, 1), | |
| 33 | + | }; | |
| 34 | + | if step == 0 { | |
| 35 | + | return Err("a step cannot be 0".to_owned()); | |
| 36 | + | } | |
| 37 | + | let (from, to) = if range == "*" { | |
| 38 | + | (low, high) | |
| 39 | + | } else if let Some((a, b)) = range.split_once('-') { | |
| 40 | + | (value(a)?, value(b)?) | |
| 41 | + | } else { | |
| 42 | + | let at = value(range)?; | |
| 43 | + | (at, if item.contains('/') { high } else { at }) | |
| 44 | + | }; | |
| 45 | + | // Sunday may be written 7. | |
| 46 | + | let (from, to) = if names.len() == 7 && to == 7 { (from.min(6), 6) } else { (from, to) }; | |
| 47 | + | if from < low || to > high || from > to { | |
| 48 | + | return Err(format!("`{item}` is outside {low}-{high}")); | |
| 49 | + | } | |
| 50 | + | let mut at = from; | |
| 51 | + | while at <= to { | |
| 52 | + | set[at as usize] = true; | |
| 53 | + | at += step; | |
| 54 | + | } | |
| 55 | + | if names.len() == 7 && text.split(',').any(|part| part == "7") { | |
| 56 | + | set[0] = true; | |
| 57 | + | } | |
| 58 | + | } | |
| 59 | + | Ok((set, text != "*")) | |
| 60 | + | } | |
| 61 | + | ||
| 62 | + | impl Schedule { | |
| 63 | + | pub fn parse(text: &str) -> Result<Schedule, String> { | |
| 64 | + | let parts: Vec<&str> = text.split_whitespace().collect(); | |
| 65 | + | let [minute, hour, day, month, weekday] = parts[..] else { | |
| 66 | + | return Err("a schedule has five fields: minute hour day month weekday, such as `0 9 * * mon`".to_owned()); | |
| 67 | + | }; | |
| 68 | + | let (minutes, _) = field(minute, 0, 59, &[])?; | |
| 69 | + | let (hours, _) = field(hour, 0, 23, &[])?; | |
| 70 | + | let (days, days_restricted) = field(day, 1, 31, &[])?; | |
| 71 | + | let (months, _) = field(month, 1, 12, &[])?; | |
| 72 | + | let (weekdays, weekdays_restricted) = field(weekday, 0, 6, &WEEKDAYS)?; | |
| 73 | + | Ok(Schedule { | |
| 74 | + | minutes, | |
| 75 | + | hours, | |
| 76 | + | days, | |
| 77 | + | months, | |
| 78 | + | weekdays, | |
| 79 | + | days_restricted, | |
| 80 | + | weekdays_restricted, | |
| 81 | + | }) | |
| 82 | + | } | |
| 83 | + | ||
| 84 | + | /// Whether it fires in the minute starting at `ms` since the epoch, UTC. | |
| 85 | + | pub fn fires_at(&self, ms: u64) -> bool { | |
| 86 | + | let minutes_total = ms / 60_000; | |
| 87 | + | let minute = (minutes_total % 60) as usize; | |
| 88 | + | let hour = (minutes_total / 60 % 24) as usize; | |
| 89 | + | let days_since_epoch = (minutes_total / 60 / 24) as i64; | |
| 90 | + | // 1970-01-01 was a Thursday. | |
| 91 | + | let weekday = ((days_since_epoch + 4) % 7) as usize; | |
| 92 | + | let (_, month, day) = civil_from_days(days_since_epoch); | |
| 93 | + | let day_ok = self.days[day as usize]; | |
| 94 | + | let weekday_ok = self.weekdays[weekday]; | |
| 95 | + | let date_ok = match (self.days_restricted, self.weekdays_restricted) { | |
| 96 | + | (true, true) => day_ok || weekday_ok, | |
| 97 | + | _ => day_ok && weekday_ok, | |
| 98 | + | }; | |
| 99 | + | self.minutes[minute] && self.hours[hour] && self.months[month as usize] && date_ok | |
| 100 | + | } | |
| 101 | + | } | |
| 102 | + | ||
| 103 | + | /// The date of a day counted from 1970-01-01 (Howard Hinnant's algorithm). | |
| 104 | + | fn civil_from_days(days: i64) -> (i64, u32, u32) { | |
| 105 | + | let z = days + 719_468; | |
| 106 | + | let era = z.div_euclid(146_097); | |
| 107 | + | let doe = z.rem_euclid(146_097); | |
| 108 | + | let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; | |
| 109 | + | let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); | |
| 110 | + | let mp = (5 * doy + 2) / 153; | |
| 111 | + | let day = (doy - (153 * mp + 2) / 5 + 1) as u32; | |
| 112 | + | let month = if mp < 10 { mp + 3 } else { mp - 9 } as u32; | |
| 113 | + | (yoe + era * 400 + i64::from(month <= 2), month, day) | |
| 114 | + | } | |
| 115 | + | ||
| 116 | + | #[cfg(test)] | |
| 117 | + | mod tests { | |
| 118 | + | use super::*; | |
| 119 | + | ||
| 120 | + | /// Milliseconds at a UTC date and time. | |
| 121 | + | fn at(days_since_epoch: u64, hour: u64, minute: u64) -> u64 { | |
| 122 | + | ((days_since_epoch * 24 + hour) * 60 + minute) * 60_000 | |
| 123 | + | } | |
| 124 | + | ||
| 125 | + | // 2026-10-05 is a Monday: 20_731 days after 1970-01-01. | |
| 126 | + | const MONDAY: u64 = 20_731; | |
| 127 | + | ||
| 128 | + | #[test] | |
| 129 | + | fn dates_are_worked_out() { | |
| 130 | + | assert_eq!(civil_from_days(0), (1970, 1, 1)); | |
| 131 | + | assert_eq!(civil_from_days(MONDAY as i64), (2026, 10, 5)); | |
| 132 | + | } | |
| 133 | + | ||
| 134 | + | #[test] | |
| 135 | + | fn mondays_at_nine() { | |
| 136 | + | let schedule = Schedule::parse("0 9 * * mon").unwrap(); | |
| 137 | + | assert!(schedule.fires_at(at(MONDAY, 9, 0))); | |
| 138 | + | assert!(!schedule.fires_at(at(MONDAY, 9, 1))); | |
| 139 | + | assert!(!schedule.fires_at(at(MONDAY + 1, 9, 0))); | |
| 140 | + | assert!(Schedule::parse("0 9 * * 1").unwrap().fires_at(at(MONDAY, 9, 0))); | |
| 141 | + | } | |
| 142 | + | ||
| 143 | + | #[test] | |
| 144 | + | fn steps_ranges_and_lists() { | |
| 145 | + | let every_quarter = Schedule::parse("*/15 * * * *").unwrap(); | |
| 146 | + | assert!(every_quarter.fires_at(at(MONDAY, 3, 45))); | |
| 147 | + | assert!(!every_quarter.fires_at(at(MONDAY, 3, 44))); | |
| 148 | + | let weekdays = Schedule::parse("30 8-17/3 * * mon-fri").unwrap(); | |
| 149 | + | assert!(weekdays.fires_at(at(MONDAY, 14, 30))); | |
| 150 | + | assert!(!weekdays.fires_at(at(MONDAY, 15, 30))); | |
| 151 | + | assert!(!weekdays.fires_at(at(MONDAY + 5, 14, 30))); | |
| 152 | + | let sunday = Schedule::parse("0 0 * * 7").unwrap(); | |
| 153 | + | assert!(sunday.fires_at(at(MONDAY + 6, 0, 0))); | |
| 154 | + | } | |
| 155 | + | ||
| 156 | + | #[test] | |
| 157 | + | fn day_of_month_or_week_when_both_are_given() { | |
| 158 | + | // The 1st, or any Monday. | |
| 159 | + | let schedule = Schedule::parse("0 0 1 * mon").unwrap(); | |
| 160 | + | assert!(schedule.fires_at(at(MONDAY, 0, 0))); | |
| 161 | + | assert!(!schedule.fires_at(at(MONDAY + 1, 0, 0))); | |
| 162 | + | } | |
| 163 | + | ||
| 164 | + | #[test] | |
| 165 | + | fn nonsense_is_refused() { | |
| 166 | + | for text in ["", "* * * *", "61 * * * *", "* * * * funday", "*/0 * * * *", "5-1 * * * *"] { | |
| 167 | + | assert!(Schedule::parse(text).is_err(), "{text}"); | |
| 168 | + | } | |
| 169 | + | } | |
| 170 | + | } |
| 1 | + | //! Reading an automation's file, and the rules it gives: what starts it, | |
| 2 | + | //! the conditions, the steps, and the words in `{{ }}` that steps fill in. | |
| 3 | + | ||
| 4 | + | use std::collections::BTreeMap; | |
| 5 | + | ||
| 6 | + | use g1t_contracts::webhooks::EVENT_TYPES; | |
| 7 | + | use serde_yaml::Value; | |
| 8 | + | ||
| 9 | + | use crate::cron::Schedule; | |
| 10 | + | ||
| 11 | + | /// What starts an automation. | |
| 12 | + | #[derive(Clone, Debug)] | |
| 13 | + | pub enum Trigger { | |
| 14 | + | Events(Vec<String>), | |
| 15 | + | Schedule { text: String, schedule: Schedule }, | |
| 16 | + | Manual, | |
| 17 | + | } | |
| 18 | + | ||
| 19 | + | impl Trigger { | |
| 20 | + | pub fn describe(&self) -> String { | |
| 21 | + | match self { | |
| 22 | + | Trigger::Events(events) => events.join(", "), | |
| 23 | + | Trigger::Schedule { text, .. } => format!("on the schedule {text} (UTC)"), | |
| 24 | + | Trigger::Manual => "by hand".to_owned(), | |
| 25 | + | } | |
| 26 | + | } | |
| 27 | + | } | |
| 28 | + | ||
| 29 | + | /// One condition: a field, and the values any of which it may have. | |
| 30 | + | #[derive(Clone, Debug, PartialEq, Eq)] | |
| 31 | + | pub struct Condition { | |
| 32 | + | pub field: String, | |
| 33 | + | pub values: Vec<String>, | |
| 34 | + | } | |
| 35 | + | ||
| 36 | + | impl Condition { | |
| 37 | + | pub fn describe(&self) -> String { | |
| 38 | + | let values = self.values.join(" or "); | |
| 39 | + | match self.field.as_str() { | |
| 40 | + | "labels" => format!("labelled {values}"), | |
| 41 | + | "actor" => format!("caused by {values}"), | |
| 42 | + | "branch" => format!("on {values}"), | |
| 43 | + | field => format!("{field} is {values}"), | |
| 44 | + | } | |
| 45 | + | } | |
| 46 | + | } | |
| 47 | + | ||
| 48 | + | #[derive(Clone, Debug, PartialEq, Eq)] | |
| 49 | + | pub enum Step { | |
| 50 | + | Comment(String), | |
| 51 | + | Label(String), | |
| 52 | + | Unlabel(String), | |
| 53 | + | AssignAgent, | |
| 54 | + | MessageAgent(String), | |
| 55 | + | OpenIssue { title: String, body: String, labels: Vec<String>, assign_agent: bool }, | |
| 56 | + | CloseIssue { not_planned: bool }, | |
| 57 | + | ReopenIssue, | |
| 58 | + | Notify { url: String, text: String }, | |
| 59 | + | } | |
| 60 | + | ||
| 61 | + | impl Step { | |
| 62 | + | pub fn describe(&self) -> String { | |
| 63 | + | match self { | |
| 64 | + | Step::Comment(_) => "comment".to_owned(), | |
| 65 | + | Step::Label(label) => format!("label {label}"), | |
| 66 | + | Step::Unlabel(label) => format!("remove the label {label}"), | |
| 67 | + | Step::AssignAgent => "put a g1t agent on it".to_owned(), | |
| 68 | + | Step::MessageAgent(_) => "message the agent working on it".to_owned(), | |
| 69 | + | Step::OpenIssue { title, assign_agent, .. } => { | |
| 70 | + | format!("open the issue \u{201c}{title}\u{201d}{}", if *assign_agent { " and put an agent on it" } else { "" }) | |
| 71 | + | } | |
| 72 | + | Step::CloseIssue { not_planned } => { | |
| 73 | + | if *not_planned { "close it as not planned".to_owned() } else { "close it".to_owned() } | |
| 74 | + | } | |
| 75 | + | Step::ReopenIssue => "reopen it".to_owned(), | |
| 76 | + | Step::Notify { url, .. } => format!("post to {}", host(url)), | |
| 77 | + | } | |
| 78 | + | } | |
| 79 | + | ||
| 80 | + | /// Whether the step acts on the issue or pull request the event is about. | |
| 81 | + | pub fn needs_target(&self) -> bool { | |
| 82 | + | !matches!(self, Step::OpenIssue { .. } | Step::Notify { .. }) | |
| 83 | + | } | |
| 84 | + | } | |
| 85 | + | ||
| 86 | + | fn host(url: &str) -> &str { | |
| 87 | + | url.trim_start_matches("https://").split('/').next().unwrap_or(url) | |
| 88 | + | } | |
| 89 | + | ||
| 90 | + | #[derive(Clone, Debug)] | |
| 91 | + | pub struct Definition { | |
| 92 | + | pub name: String, | |
| 93 | + | pub trigger: Trigger, | |
| 94 | + | pub conditions: Vec<Condition>, | |
| 95 | + | pub steps: Vec<Step>, | |
| 96 | + | pub per_hour: u32, | |
| 97 | + | } | |
| 98 | + | ||
| 99 | + | /// The default and the most runs an automation makes in an hour. | |
| 100 | + | pub const DEFAULT_PER_HOUR: u32 = 30; | |
| 101 | + | pub const MAX_PER_HOUR: u32 = 200; | |
| 102 | + | ||
| 103 | + | fn text(value: &Value) -> Option<String> { | |
| 104 | + | match value { | |
| 105 | + | Value::String(text) => Some(text.clone()), | |
| 106 | + | Value::Number(number) => Some(number.to_string()), | |
| 107 | + | Value::Bool(flag) => Some(flag.to_string()), | |
| 108 | + | _ => None, | |
| 109 | + | } | |
| 110 | + | } | |
| 111 | + | ||
| 112 | + | fn texts(value: &Value) -> Result<Vec<String>, String> { | |
| 113 | + | match value { | |
| 114 | + | Value::Sequence(items) => items.iter().map(|item| text(item).ok_or_else(|| "a list of words".to_owned())).collect(), | |
| 115 | + | other => text(other).map(|one| vec![one]).ok_or_else(|| "a word or a list of words".to_owned()), | |
| 116 | + | } | |
| 117 | + | } | |
| 118 | + | ||
| 119 | + | fn step(value: &Value) -> Result<Step, String> { | |
| 120 | + | let (name, argument) = match value { | |
| 121 | + | Value::String(name) => (name.as_str(), &Value::Null), | |
| 122 | + | Value::Mapping(map) if map.len() == 1 => { | |
| 123 | + | let (key, argument) = map.iter().next().expect("one entry"); | |
| 124 | + | (key.as_str().ok_or("a step's name is a word")?, argument) | |
| 125 | + | } | |
| 126 | + | _ => return Err("each step is a name, such as `assign_agent`, or a name and what it takes, such as `comment: Thanks!`".to_owned()), | |
| 127 | + | }; | |
| 128 | + | let field = |key: &str| match argument { | |
| 129 | + | Value::Mapping(map) => map.get(key).and_then(text), | |
| 130 | + | _ => None, | |
| 131 | + | }; | |
| 132 | + | let needs_text = |what: &str| text(argument).filter(|t| !t.trim().is_empty()).ok_or(format!("`{name}` needs {what}")); | |
| 133 | + | Ok(match name { | |
| 134 | + | "comment" => Step::Comment(needs_text("the comment's text")?), | |
| 135 | + | "label" => Step::Label(needs_text("a label")?), | |
| 136 | + | "unlabel" => Step::Unlabel(needs_text("a label")?), | |
| 137 | + | "assign_agent" => Step::AssignAgent, | |
| 138 | + | "message_agent" => Step::MessageAgent(needs_text("the message")?), | |
| 139 | + | "close_issue" => Step::CloseIssue { | |
| 140 | + | not_planned: matches!(text(argument).as_deref(), Some("not_planned")), | |
| 141 | + | }, | |
| 142 | + | "reopen_issue" => Step::ReopenIssue, | |
| 143 | + | "open_issue" => Step::OpenIssue { | |
| 144 | + | title: field("title").filter(|t| !t.trim().is_empty()).ok_or("`open_issue` needs a title")?, | |
| 145 | + | body: field("body").unwrap_or_default(), | |
| 146 | + | labels: match argument { | |
| 147 | + | Value::Mapping(map) => map.get("labels").map(texts).transpose()?.unwrap_or_default(), | |
| 148 | + | _ => Vec::new(), | |
| 149 | + | }, | |
| 150 | + | assign_agent: matches!(field("assign_agent").as_deref(), Some("true")), | |
| 151 | + | }, | |
| 152 | + | "notify" => { | |
| 153 | + | let url = field("url").ok_or("`notify` needs a url")?; | |
| 154 | + | if !url.starts_with("https://") { | |
| 155 | + | return Err("`notify` posts only to https:// addresses".to_owned()); | |
| 156 | + | } | |
| 157 | + | Step::Notify { | |
| 158 | + | url, | |
| 159 | + | text: field("text").ok_or("`notify` needs text")?, | |
| 160 | + | } | |
| 161 | + | } | |
| 162 | + | other => { | |
| 163 | + | return Err(format!( | |
| 164 | + | "there is no step called `{other}`; steps are comment, label, unlabel, assign_agent, message_agent, open_issue, close_issue, reopen_issue and notify" | |
| 165 | + | )); | |
| 166 | + | } | |
| 167 | + | }) | |
| 168 | + | } | |
| 169 | + | ||
| 170 | + | /// Reads an automation's file. `Err` says what is wrong with it, for the | |
| 171 | + | /// person who wrote it. | |
| 172 | + | pub fn parse(yaml: &str, file_name: &str) -> Result<Definition, String> { | |
| 173 | + | let root: Value = serde_yaml::from_str(yaml).map_err(|error| format!("It is not valid YAML: {error}"))?; | |
| 174 | + | let Value::Mapping(map) = &root else { | |
| 175 | + | return Err("An automation is a mapping with `on` and `do`.".to_owned()); | |
| 176 | + | }; | |
| 177 | + | let name = map | |
| 178 | + | .get("name") | |
| 179 | + | .and_then(text) | |
| 180 | + | .unwrap_or_else(|| file_name.trim_end_matches(".yml").trim_end_matches(".yaml").replace(['-', '_'], " ")); | |
| 181 | + | let trigger = match map.get("on") { | |
| 182 | + | None => return Err("`on` is missing: say which event starts it, a schedule, or manual.".to_owned()), | |
| 183 | + | Some(Value::String(manual)) if manual == "manual" => Trigger::Manual, | |
| 184 | + | Some(Value::Mapping(on)) if on.contains_key("schedule") => { | |
| 185 | + | let text = on.get("schedule").and_then(text).ok_or("`schedule` takes a cron line, such as \"0 9 * * mon\".")?; | |
| 186 | + | let schedule = Schedule::parse(&text).map_err(|problem| format!("The schedule does not read: {problem}."))?; | |
| 187 | + | Trigger::Schedule { text, schedule } | |
| 188 | + | } | |
| 189 | + | Some(events) => { | |
| 190 | + | let events = texts(events).map_err(|_| "`on` takes an event, a list of events, `manual`, or `schedule:`.".to_owned())?; | |
| 191 | + | if let Some(unknown) = events.iter().find(|event| !EVENT_TYPES.contains(&event.as_str())) { | |
| 192 | + | return Err(format!("There is no event called {unknown}. Events are {}.", EVENT_TYPES.join(", "))); | |
| 193 | + | } | |
| 194 | + | Trigger::Events(events) | |
| 195 | + | } | |
| 196 | + | }; | |
| 197 | + | let mut conditions = Vec::new(); | |
| 198 | + | if let Some(when) = map.get("if") { | |
| 199 | + | let Value::Mapping(when) = when else { | |
| 200 | + | return Err("`if` is a mapping of fields to the values they must have.".to_owned()); | |
| 201 | + | }; | |
| 202 | + | for (field, values) in when { | |
| 203 | + | let field = field.as_str().ok_or("`if` keys are field names")?.to_owned(); | |
| 204 | + | let values = texts(values).map_err(|problem| format!("`if` {field}: give {problem}."))?; | |
| 205 | + | conditions.push(Condition { field, values }); | |
| 206 | + | } | |
| 207 | + | } | |
| 208 | + | let steps = match map.get("do") { | |
| 209 | + | Some(Value::Sequence(steps)) if !steps.is_empty() => steps.iter().map(step).collect::<Result<Vec<_>, _>>()?, | |
| 210 | + | Some(single @ (Value::String(_) | Value::Mapping(_))) => vec![step(single)?], | |
| 211 | + | _ => return Err("`do` is missing: give the steps to take, as a list.".to_owned()), | |
| 212 | + | }; | |
| 213 | + | if matches!(trigger, Trigger::Schedule { .. }) | |
| 214 | + | && let Some(step) = steps.iter().find(|step| step.needs_target()) { | |
| 215 | + | return Err(format!( | |
| 216 | + | "A scheduled automation has no issue or pull request to act on, so it cannot {}. It can open_issue or notify.", | |
| 217 | + | step.describe() | |
| 218 | + | )); | |
| 219 | + | } | |
| 220 | + | let per_hour = match map.get("limits").and_then(|limits| limits.get("per_hour")) { | |
| 221 | + | Some(value) => value | |
| 222 | + | .as_u64() | |
| 223 | + | .filter(|n| (1..=u64::from(MAX_PER_HOUR)).contains(n)) | |
| 224 | + | .ok_or(format!("`limits.per_hour` is a number from 1 to {MAX_PER_HOUR}."))? as u32, | |
| 225 | + | None => DEFAULT_PER_HOUR, | |
| 226 | + | }; | |
| 227 | + | Ok(Definition { | |
| 228 | + | name, | |
| 229 | + | trigger, | |
| 230 | + | conditions, | |
| 231 | + | steps, | |
| 232 | + | per_hour, | |
| 233 | + | }) | |
| 234 | + | } | |
| 235 | + | ||
| 236 | + | /// What a run knows, by name: for conditions and for `{{ }}` in steps. | |
| 237 | + | #[derive(Clone, Debug, Default)] | |
| 238 | + | pub struct Context { | |
| 239 | + | pub vars: BTreeMap<String, String>, | |
| 240 | + | pub labels: Vec<String>, | |
| 241 | + | } | |
| 242 | + | ||
| 243 | + | impl Context { | |
| 244 | + | pub fn set(&mut self, key: &str, value: impl Into<String>) { | |
| 245 | + | self.vars.insert(key.to_owned(), value.into()); | |
| 246 | + | } | |
| 247 | + | ||
| 248 | + | /// An event's data, as `data.<field>` and, where nothing else claims | |
| 249 | + | /// the name, as `<field>`. | |
| 250 | + | pub fn add_data(&mut self, data: &serde_json::Value) { | |
| 251 | + | let Some(fields) = data.as_object() else { return }; | |
| 252 | + | for (key, value) in fields { | |
| 253 | + | let value = match value { | |
| 254 | + | serde_json::Value::String(text) => text.clone(), | |
| 255 | + | serde_json::Value::Number(number) => number.to_string(), | |
| 256 | + | serde_json::Value::Bool(flag) => flag.to_string(), | |
| 257 | + | _ => continue, | |
| 258 | + | }; | |
| 259 | + | self.vars.entry(key.clone()).or_insert_with(|| value.clone()); | |
| 260 | + | self.vars.insert(format!("data.{key}"), value); | |
| 261 | + | } | |
| 262 | + | } | |
| 263 | + | } | |
| 264 | + | ||
| 265 | + | /// Whether every condition holds. `labels` matches when the issue or pull | |
| 266 | + | /// request has any of the labels; other fields compare with the context. | |
| 267 | + | pub fn holds(conditions: &[Condition], context: &Context) -> Result<(), String> { | |
| 268 | + | for condition in conditions { | |
| 269 | + | let ok = if condition.field == "labels" { | |
| 270 | + | condition.values.iter().any(|value| context.labels.iter().any(|label| label.eq_ignore_ascii_case(value))) | |
| 271 | + | } else { | |
| 272 | + | let actual = context.vars.get(&condition.field).or_else(|| context.vars.get(&format!("data.{}", condition.field))); | |
| 273 | + | actual.is_some_and(|actual| condition.values.iter().any(|value| value.eq_ignore_ascii_case(actual))) | |
| 274 | + | }; | |
| 275 | + | if !ok { | |
| 276 | + | return Err(format!("not {}", condition.describe())); | |
| 277 | + | } | |
| 278 | + | } | |
| 279 | + | Ok(()) | |
| 280 | + | } | |
| 281 | + | ||
| 282 | + | /// Fills `{{ name }}` with what the context knows; an unknown name is left | |
| 283 | + | /// empty. | |
| 284 | + | pub fn render(template: &str, context: &Context) -> String { | |
| 285 | + | let mut out = String::with_capacity(template.len()); | |
| 286 | + | let mut rest = template; | |
| 287 | + | while let Some(start) = rest.find("{{") { | |
| 288 | + | out.push_str(&rest[..start]); | |
| 289 | + | let after = &rest[start + 2..]; | |
| 290 | + | let Some(end) = after.find("}}") else { | |
| 291 | + | out.push_str(&rest[start..]); | |
| 292 | + | return out; | |
| 293 | + | }; | |
| 294 | + | let name = after[..end].trim(); | |
| 295 | + | out.push_str(context.vars.get(name).map(String::as_str).unwrap_or_default()); | |
| 296 | + | rest = &after[end + 2..]; | |
| 297 | + | } | |
| 298 | + | out.push_str(rest); | |
| 299 | + | out | |
| 300 | + | } | |
| 301 | + | ||
| 302 | + | #[cfg(test)] | |
| 303 | + | mod tests { | |
| 304 | + | use super::*; | |
| 305 | + | ||
| 306 | + | const BUGS: &str = r#" | |
| 307 | + | name: Put an agent on new bugs | |
| 308 | + | on: issue.opened | |
| 309 | + | if: | |
| 310 | + | labels: [bug, regression] | |
| 311 | + | do: | |
| 312 | + | - comment: "Thanks, {{actor}}. An agent is on it." | |
| 313 | + | - assign_agent | |
| 314 | + | - notify: { url: "https://hooks.slack.com/x", text: "{{repo}}#{{number}}: {{title}}" } | |
| 315 | + | limits: | |
| 316 | + | per_hour: 5 | |
| 317 | + | "#; | |
| 318 | + | ||
| 319 | + | #[test] | |
| 320 | + | fn a_whole_automation_reads() { | |
| 321 | + | let definition = parse(BUGS, "bugs.yml").unwrap(); | |
| 322 | + | assert_eq!(definition.name, "Put an agent on new bugs"); | |
| 323 | + | assert!(matches!(&definition.trigger, Trigger::Events(events) if events == &vec!["issue.opened".to_owned()])); | |
| 324 | + | assert_eq!(definition.conditions, vec![Condition { field: "labels".into(), values: vec!["bug".into(), "regression".into()] }]); | |
| 325 | + | assert_eq!(definition.steps.len(), 3); | |
| 326 | + | assert_eq!(definition.steps[1], Step::AssignAgent); | |
| 327 | + | assert_eq!(definition.per_hour, 5); | |
| 328 | + | assert_eq!(definition.steps[2].describe(), "post to hooks.slack.com"); | |
| 329 | + | } | |
| 330 | + | ||
| 331 | + | #[test] | |
| 332 | + | fn mistakes_are_explained() { | |
| 333 | + | let problem = |yaml: &str| parse(yaml, "x.yml").unwrap_err(); | |
| 334 | + | assert!(problem("on: issue.exploded\ndo: [assign_agent]").contains("no event called issue.exploded")); | |
| 335 | + | assert!(problem("on: issue.opened").contains("`do` is missing")); | |
| 336 | + | assert!(problem("do: [assign_agent]").contains("`on` is missing")); | |
| 337 | + | assert!(problem("on: issue.opened\ndo: [dance]").contains("no step called `dance`")); | |
| 338 | + | assert!(problem("on: issue.opened\ndo: [{notify: {url: 'http://x', text: hi}}]").contains("https://")); | |
| 339 | + | assert!(problem("on: { schedule: '0 9 * * mon' }\ndo: [assign_agent]").contains("cannot put a g1t agent on it")); | |
| 340 | + | assert!(problem("on: { schedule: 'often' }\ndo: [{open_issue: {title: x}}]").contains("schedule does not read")); | |
| 341 | + | assert!(problem("on: issue.opened\ndo: [assign_agent]\nlimits: { per_hour: 0 }").contains("per_hour")); | |
| 342 | + | assert!(problem(": : :").contains("not valid YAML")); | |
| 343 | + | } | |
| 344 | + | ||
| 345 | + | #[test] | |
| 346 | + | fn schedules_and_manual_runs_read() { | |
| 347 | + | let weekly = parse("on: { schedule: '0 9 * * mon' }\ndo:\n - open_issue: { title: Weekly tidy, labels: chore, assign_agent: true }", "weekly.yml").unwrap(); | |
| 348 | + | assert!(matches!(weekly.trigger, Trigger::Schedule { .. })); | |
| 349 | + | assert_eq!(weekly.name, "weekly"); | |
| 350 | + | assert!(matches!(&weekly.steps[0], Step::OpenIssue { labels, assign_agent: true, .. } if labels == &vec!["chore".to_owned()])); | |
| 351 | + | assert!(matches!(parse("on: manual\ndo:\n comment: hi", "m.yml").unwrap().trigger, Trigger::Manual)); | |
| 352 | + | } | |
| 353 | + | ||
| 354 | + | #[test] | |
| 355 | + | fn conditions_check_labels_fields_and_data() { | |
| 356 | + | let mut context = Context::default(); | |
| 357 | + | context.labels = vec!["Bug".into()]; | |
| 358 | + | context.set("actor", "ada"); | |
| 359 | + | context.add_data(&serde_json::json!({ "status": "failed", "number": 7 })); | |
| 360 | + | let condition = |field: &str, values: &[&str]| Condition { field: field.into(), values: values.iter().map(|v| v.to_string()).collect() }; | |
| 361 | + | assert!(holds(&[condition("labels", &["bug"]), condition("status", &["failed", "errored"])], &context).is_ok()); | |
| 362 | + | assert!(holds(&[condition("data.status", &["failed"])], &context).is_ok()); | |
| 363 | + | assert_eq!(holds(&[condition("actor", &["grace"])], &context).unwrap_err(), "not caused by grace"); | |
| 364 | + | assert!(holds(&[condition("verdict", &["approve"])], &context).is_err()); | |
| 365 | + | } | |
| 366 | + | ||
| 367 | + | #[test] | |
| 368 | + | fn templates_fill_in_what_is_known() { | |
| 369 | + | let mut context = Context::default(); | |
| 370 | + | context.set("repo", "acme/web"); | |
| 371 | + | context.set("number", "12"); | |
| 372 | + | assert_eq!(render("{{repo}}#{{ number }} by {{actor}}", &context), "acme/web#12 by "); | |
| 373 | + | assert_eq!(render("no braces", &context), "no braces"); | |
| 374 | + | assert_eq!(render("half {{open", &context), "half {{open"); | |
| 375 | + | } | |
| 376 | + | } |
| 1 | + | //! The automations service: rules in a repository's `.g1t/automations/` | |
| 2 | + | //! that act when something happens. See `g1t_contracts::automations` for | |
| 3 | + | //! the methods, and `definition` for what a file may say. | |
| 4 | + | //! | |
| 5 | + | //! The files on a repository's default branch are read again on every push | |
| 6 | + | //! to it. An event from the bus runs each enabled automation that wants it; | |
| 7 | + | //! the minute's sweep runs scheduled ones; a member can run any by hand. | |
| 8 | + | //! Every run is recorded with each step's result, and every run obeys three | |
| 9 | + | //! rules: an event is handled once per automation, an automation makes at | |
| 10 | + | //! most so many runs an hour, and an automation does not answer what it | |
| 11 | + | //! itself just did. | |
| 12 | + | //! | |
| 13 | + | //! Automations act as their workspace: what they write is the workspace's, | |
| 14 | + | //! and says which automation wrote it. | |
| 15 | + | ||
| 16 | + | mod cron; | |
| 17 | + | mod definition; | |
| 18 | + | ||
| 19 | + | use g1t_contracts::automations::*; | |
| 20 | + | use g1t_contracts::events::Event; | |
| 21 | + | use g1t_contracts::identity::{AGENT_ID, AGENT_NAME, SlugArgs, UsernamesArgs, Workspace}; | |
| 22 | + | use g1t_contracts::repos::{BlobArgs, BlobView, EntryKind, GetArgs, PathByIdArgs, Repo, RepoPath, TreeArgs, TreeView}; | |
| 23 | + | use g1t_contracts::time::rfc3339; | |
| 24 | + | use g1t_contracts::work::{IssueDetail, PullDetail, ViewArgs}; | |
| 25 | + | use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, Role, User, Viewer, new_id}; | |
| 26 | + | use g1t_kit::{args, now_ms, reply, rpc_method}; | |
| 27 | + | use serde::Deserialize; | |
| 28 | + | use serde_json::{Value, json}; | |
| 29 | + | use worker::wasm_bindgen::JsValue; | |
| 30 | + | use worker::{ | |
| 31 | + | Context as WorkerContext, D1Database, Env, Fetch, Fetcher, Headers, MessageBatch, MessageExt, Method, Request, RequestInit, | |
| 32 | + | Response, Result, ScheduleContext, ScheduledEvent, event, | |
| 33 | + | }; | |
| 34 | + | ||
| 35 | + | use definition::{Context, Definition, Step, Trigger}; | |
| 36 | + | ||
| 37 | + | /// Where automations live in a repository. | |
| 38 | + | const FOLDER: &str = ".g1t/automations"; | |
| 39 | + | /// The most automation files read from a repository. | |
| 40 | + | const MAX_FILES: usize = 50; | |
| 41 | + | const RUNS_SHOWN: u32 = 50; | |
| 42 | + | /// How long an automation's own effect is remembered, so it does not answer it. | |
| 43 | + | const OWN_EFFECT_MS: u64 = 10 * 60 * 1000; | |
| 44 | + | const SITE: &str = "https://g1t.sh"; | |
| 45 | + | ||
| 46 | + | #[derive(Deserialize)] | |
| 47 | + | struct AutomationRow { | |
| 48 | + | id: String, | |
| 49 | + | repo_id: String, | |
| 50 | + | repo: String, | |
| 51 | + | path: String, | |
| 52 | + | name: String, | |
| 53 | + | source: String, | |
| 54 | + | error: Option<String>, | |
| 55 | + | enabled: u32, | |
| 56 | + | } | |
| 57 | + | ||
| 58 | + | impl AutomationRow { | |
| 59 | + | fn definition(&self) -> std::result::Result<Definition, String> { | |
| 60 | + | if let Some(error) = &self.error { | |
| 61 | + | return Err(error.clone()); | |
| 62 | + | } | |
| 63 | + | definition::parse(&self.source, self.path.rsplit('/').next().unwrap_or(&self.path)) | |
| 64 | + | } | |
| 65 | + | ||
| 66 | + | fn repo_path(&self) -> RepoPath { | |
| 67 | + | let (namespace, name) = self.repo.split_once('/').unwrap_or((&self.repo, "")); | |
| 68 | + | RepoPath { | |
| 69 | + | namespace: namespace.to_owned(), | |
| 70 | + | name: name.to_owned(), | |
| 71 | + | } | |
| 72 | + | } | |
| 73 | + | } | |
| 74 | + | ||
| 75 | + | #[derive(Deserialize)] | |
| 76 | + | struct RunRow { | |
| 77 | + | id: String, | |
| 78 | + | automation_id: String, | |
| 79 | + | name: String, | |
| 80 | + | event: String, | |
| 81 | + | number: Option<u32>, | |
| 82 | + | status: String, | |
| 83 | + | reason: Option<String>, | |
| 84 | + | steps: String, | |
| 85 | + | actor: Option<String>, | |
| 86 | + | started_at: String, | |
| 87 | + | } | |
| 88 | + | ||
| 89 | + | impl From<RunRow> for AutomationRun { | |
| 90 | + | fn from(row: RunRow) -> Self { | |
| 91 | + | AutomationRun { | |
| 92 | + | id: row.id, | |
| 93 | + | automation_id: row.automation_id, | |
| 94 | + | name: row.name, | |
| 95 | + | event: row.event, | |
| 96 | + | number: row.number, | |
| 97 | + | status: row.status, | |
| 98 | + | reason: row.reason, | |
| 99 | + | steps: serde_json::from_str(&row.steps).unwrap_or_default(), | |
| 100 | + | actor: row.actor, | |
| 101 | + | started_at: row.started_at, | |
| 102 | + | } | |
| 103 | + | } | |
| 104 | + | } | |
| 105 | + | ||
| 106 | + | #[derive(Deserialize)] | |
| 107 | + | struct Count { | |
| 108 | + | n: u32, | |
| 109 | + | } | |
| 110 | + | ||
| 111 | + | /// What started a run. | |
| 112 | + | struct Source { | |
| 113 | + | /// Unique per automation: an event's id, the minute, or a manual run. | |
| 114 | + | key: String, | |
| 115 | + | event: String, | |
| 116 | + | number: Option<u32>, | |
| 117 | + | /// The id of whoever caused it. | |
| 118 | + | actor_id: Option<String>, | |
| 119 | + | /// The event's data, for conditions and `{{data.*}}`. | |
| 120 | + | data: Value, | |
| 121 | + | } | |
| 122 | + | ||
| 123 | + | fn optional(value: Option<&str>) -> JsValue { | |
| 124 | + | value.map_or(JsValue::NULL, JsValue::from) | |
| 125 | + | } | |
| 126 | + | ||
| 127 | + | fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> { | |
| 128 | + | Outcome::fail(code, message) | |
| 129 | + | } | |
| 130 | + | ||
| 131 | + | fn summary(row: &AutomationRow, last_run: Option<AutomationRun>) -> Automation { | |
| 132 | + | let parsed = row.definition(); | |
| 133 | + | Automation { | |
| 134 | + | id: row.id.clone(), | |
| 135 | + | repo: row.repo.clone(), | |
| 136 | + | path: row.path.clone(), | |
| 137 | + | name: row.name.clone(), | |
| 138 | + | trigger: parsed.as_ref().map(|d| d.trigger.describe()).unwrap_or_default(), | |
| 139 | + | conditions: parsed.as_ref().map(|d| d.conditions.iter().map(|c| c.describe()).collect()).unwrap_or_default(), | |
| 140 | + | steps: parsed.as_ref().map(|d| d.steps.iter().map(Step::describe).collect()).unwrap_or_default(), | |
| 141 | + | enabled: row.enabled != 0, | |
| 142 | + | error: parsed.as_ref().err().cloned(), | |
| 143 | + | manual: parsed.is_ok(), | |
| 144 | + | last_run, | |
| 145 | + | } | |
| 146 | + | } | |
| 147 | + | ||
| 148 | + | struct Automations { | |
| 149 | + | db: D1Database, | |
| 150 | + | repos: Fetcher, | |
| 151 | + | work: Fetcher, | |
| 152 | + | identity: Fetcher, | |
| 153 | + | runner: Fetcher, | |
| 154 | + | } | |
| 155 | + | ||
| 156 | + | impl Automations { | |
| 157 | + | fn new(env: &Env) -> Result<Self> { | |
| 158 | + | Ok(Automations { | |
| 159 | + | db: env.d1("DB")?, | |
| 160 | + | repos: env.service("REPOS")?, | |
| 161 | + | work: env.service("WORK")?, | |
| 162 | + | identity: env.service("IDENTITY")?, | |
| 163 | + | runner: env.service("RUNNER")?, | |
| 164 | + | }) | |
| 165 | + | } | |
| 166 | + | ||
| 167 | + | /// The workspace itself, as automations act. | |
| 168 | + | async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> { | |
| 169 | + | let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?; | |
| 170 | + | Ok(workspace.map(|workspace| User { | |
| 171 | + | id: workspace.id, | |
| 172 | + | username: workspace.slug.clone(), | |
| 173 | + | kind: PrincipalKind::Workspace, | |
| 174 | + | verified: true, | |
| 175 | + | workspaces: vec![Membership { | |
| 176 | + | slug: workspace.slug, | |
| 177 | + | role: Role::Member, | |
| 178 | + | }], | |
| 179 | + | })) | |
| 180 | + | } | |
| 181 | + | ||
| 182 | + | async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> { | |
| 183 | + | let found: Outcome<Repo> = g1t_kit::call( | |
| 184 | + | &self.repos, | |
| 185 | + | "get", | |
| 186 | + | &GetArgs { | |
| 187 | + | path: path.clone(), | |
| 188 | + | viewer: viewer.clone(), | |
| 189 | + | }, | |
| 190 | + | ) | |
| 191 | + | .await?; | |
| 192 | + | Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none())) | |
| 193 | + | } | |
| 194 | + | ||
| 195 | + | // --- Reading the files ---------------------------------------------------- | |
| 196 | + | ||
| 197 | + | /// Reads a repository's automation files from its default branch, and | |
| 198 | + | /// keeps the database in step with them. | |
| 199 | + | async fn sync(&self, path: &RepoPath) -> Result<()> { | |
| 200 | + | let Some(actor) = self.workspace_actor(&path.namespace).await? else { | |
| 201 | + | return Ok(()); | |
| 202 | + | }; | |
| 203 | + | let viewer = Some(actor); | |
| 204 | + | let tree: Outcome<TreeView> = g1t_kit::call( | |
| 205 | + | &self.repos, | |
| 206 | + | "tree", | |
| 207 | + | &TreeArgs { | |
| 208 | + | path: path.clone(), | |
| 209 | + | viewer: viewer.clone(), | |
| 210 | + | git_ref: None, | |
| 211 | + | tree_path: FOLDER.to_owned(), | |
| 212 | + | }, | |
| 213 | + | ) | |
| 214 | + | .await?; | |
| 215 | + | let Some(repo) = self.visible_repo(path, &viewer).await? else { | |
| 216 | + | return Ok(()); | |
| 217 | + | }; | |
| 218 | + | let entries = match tree { | |
| 219 | + | Outcome::Ok(tree) => tree.entries, | |
| 220 | + | // No folder: no automations. | |
| 221 | + | Outcome::Fail(_) => Vec::new(), | |
| 222 | + | }; | |
| 223 | + | let files: Vec<String> = entries | |
| 224 | + | .into_iter() | |
| 225 | + | .filter(|entry| matches!(entry.kind, EntryKind::Blob | EntryKind::Exec)) | |
| 226 | + | .map(|entry| entry.name) | |
| 227 | + | .filter(|name| name.ends_with(".yml") || name.ends_with(".yaml")) | |
| 228 | + | .take(MAX_FILES) | |
| 229 | + | .collect(); | |
| 230 | + | let full_name = format!("{}/{}", repo.namespace, repo.name); | |
| 231 | + | let now = rfc3339(now_ms()); | |
| 232 | + | let mut kept = Vec::new(); | |
| 233 | + | for file in &files { | |
| 234 | + | let file_path = format!("{FOLDER}/{file}"); | |
| 235 | + | let blob: Outcome<BlobView> = g1t_kit::call( | |
| 236 | + | &self.repos, | |
| 237 | + | "blob", | |
| 238 | + | &BlobArgs { | |
| 239 | + | path: path.clone(), | |
| 240 | + | viewer: viewer.clone(), | |
| 241 | + | git_ref: repo.default_branch.clone(), | |
| 242 | + | file_path: file_path.clone(), | |
| 243 | + | }, | |
| 244 | + | ) | |
| 245 | + | .await?; | |
| 246 | + | let source = match blob { | |
| 247 | + | Outcome::Ok(BlobView { text: Some(text), .. }) => text, | |
| 248 | + | _ => continue, | |
| 249 | + | }; | |
| 250 | + | let parsed = definition::parse(&source, file); | |
| 251 | + | let (name, error, kind, events) = match &parsed { | |
| 252 | + | Ok(definition) => ( | |
| 253 | + | definition.name.clone(), | |
| 254 | + | None, | |
| 255 | + | match &definition.trigger { | |
| 256 | + | Trigger::Events(_) => "events", | |
| 257 | + | Trigger::Schedule { .. } => "schedule", | |
| 258 | + | Trigger::Manual => "manual", | |
| 259 | + | }, | |
| 260 | + | match &definition.trigger { | |
| 261 | + | Trigger::Events(events) => events.clone(), | |
| 262 | + | _ => Vec::new(), | |
| 263 | + | }, | |
| 264 | + | ), | |
| 265 | + | Err(problem) => (file.clone(), Some(problem.clone()), "invalid", Vec::new()), | |
| 266 | + | }; | |
| 267 | + | self.db | |
| 268 | + | .prepare( | |
| 269 | + | "INSERT INTO automations (id, repo_id, repo, path, name, source, trigger_kind, events, error, updated_at) | |
| 270 | + | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) | |
| 271 | + | ON CONFLICT (repo_id, path) DO UPDATE SET | |
| 272 | + | repo = excluded.repo, name = excluded.name, source = excluded.source, | |
| 273 | + | trigger_kind = excluded.trigger_kind, events = excluded.events, | |
| 274 | + | error = excluded.error, updated_at = excluded.updated_at", | |
| 275 | + | ) | |
| 276 | + | .bind(&[ | |
| 277 | + | new_id("aut", now_ms()).into(), | |
| 278 | + | repo.id.as_str().into(), | |
| 279 | + | full_name.as_str().into(), | |
| 280 | + | file_path.as_str().into(), | |
| 281 | + | name.into(), | |
| 282 | + | source.into(), | |
| 283 | + | kind.into(), | |
| 284 | + | serde_json::to_string(&events)?.into(), | |
| 285 | + | optional(error.as_deref()), | |
| 286 | + | now.as_str().into(), | |
| 287 | + | ])? | |
| 288 | + | .run() | |
| 289 | + | .await?; | |
| 290 | + | kept.push(file_path); | |
| 291 | + | } | |
| 292 | + | // Files that are gone take their automations with them. | |
| 293 | + | let mut statements = Vec::new(); | |
| 294 | + | let existing = self | |
| 295 | + | .db | |
| 296 | + | .prepare("SELECT * FROM automations WHERE repo_id = ?") | |
| 297 | + | .bind(&[repo.id.as_str().into()])? | |
| 298 | + | .all() | |
| 299 | + | .await? | |
| 300 | + | .results::<AutomationRow>()?; | |
| 301 | + | for row in existing.iter().filter(|row| !kept.contains(&row.path)) { | |
| 302 | + | statements.push(self.db.prepare("DELETE FROM automations WHERE id = ?").bind(&[row.id.as_str().into()])?); | |
| 303 | + | } | |
| 304 | + | statements.push( | |
| 305 | + | self.db | |
| 306 | + | .prepare("INSERT OR REPLACE INTO synced (repo_id, at) VALUES (?, ?)") | |
| 307 | + | .bind(&[repo.id.as_str().into(), now.into()])?, | |
| 308 | + | ); | |
| 309 | + | self.db.batch(statements).await?; | |
| 310 | + | Ok(()) | |
| 311 | + | } | |
| 312 | + | ||
| 313 | + | async fn synced(&self, repo_id: &str) -> Result<bool> { | |
| 314 | + | Ok(self | |
| 315 | + | .db | |
| 316 | + | .prepare("SELECT COUNT(*) AS n FROM synced WHERE repo_id = ?") | |
| 317 | + | .bind(&[repo_id.into()])? | |
| 318 | + | .first::<Count>(None) | |
| 319 | + | .await? | |
| 320 | + | .is_some_and(|count| count.n > 0)) | |
| 321 | + | } | |
| 322 | + | ||
| 323 | + | // --- Reading and managing --------------------------------------------------- | |
| 324 | + | ||
| 325 | + | async fn last_run(&self, automation_id: &str) -> Result<Option<AutomationRun>> { | |
| 326 | + | Ok(self | |
| 327 | + | .db | |
| 328 | + | .prepare("SELECT * FROM runs WHERE automation_id = ? ORDER BY id DESC LIMIT 1") | |
| 329 | + | .bind(&[automation_id.into()])? | |
| 330 | + | .first::<RunRow>(None) | |
| 331 | + | .await? | |
| 332 | + | .map(AutomationRun::from)) | |
| 333 | + | } | |
| 334 | + | ||
| 335 | + | async fn list(&self, a: ListArgs) -> Result<Outcome<Vec<Automation>>> { | |
| 336 | + | let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else { | |
| 337 | + | return Ok(fail(FailureCode::NotFound, "There is no such repository.")); | |
| 338 | + | }; | |
| 339 | + | if !self.synced(&repo.id).await? { | |
| 340 | + | self.sync(&RepoPath { | |
| 341 | + | namespace: repo.namespace.clone(), | |
| 342 | + | name: repo.name.clone(), | |
| 343 | + | }) | |
| 344 | + | .await?; | |
| 345 | + | } | |
| 346 | + | let rows = self | |
| 347 | + | .db | |
| 348 | + | .prepare("SELECT * FROM automations WHERE repo_id = ? ORDER BY path") | |
| 349 | + | .bind(&[repo.id.as_str().into()])? | |
| 350 | + | .all() | |
| 351 | + | .await? | |
| 352 | + | .results::<AutomationRow>()?; | |
| 353 | + | let mut out = Vec::with_capacity(rows.len()); | |
| 354 | + | for row in &rows { | |
| 355 | + | out.push(summary(row, self.last_run(&row.id).await?)); | |
| 356 | + | } | |
| 357 | + | Ok(Outcome::Ok(out)) | |
| 358 | + | } | |
| 359 | + | ||
| 360 | + | async fn runs(&self, a: RunsArgs) -> Result<Outcome<Vec<AutomationRun>>> { | |
| 361 | + | let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else { | |
| 362 | + | return Ok(fail(FailureCode::NotFound, "There is no such repository.")); | |
| 363 | + | }; | |
| 364 | + | let rows = match &a.automation { | |
| 365 | + | Some(id) => self | |
| 366 | + | .db | |
| 367 | + | .prepare("SELECT * FROM runs WHERE repo_id = ? AND automation_id = ? ORDER BY id DESC LIMIT ?") | |
| 368 | + | .bind(&[repo.id.as_str().into(), id.as_str().into(), RUNS_SHOWN.into()])?, | |
| 369 | + | None => self | |
| 370 | + | .db | |
| 371 | + | .prepare("SELECT * FROM runs WHERE repo_id = ? ORDER BY id DESC LIMIT ?") | |
| 372 | + | .bind(&[repo.id.as_str().into(), RUNS_SHOWN.into()])?, | |
| 373 | + | } | |
| 374 | + | .all() | |
| 375 | + | .await? | |
| 376 | + | .results::<RunRow>()?; | |
| 377 | + | Ok(Outcome::Ok(rows.into_iter().map(AutomationRun::from).collect())) | |
| 378 | + | } | |
| 379 | + | ||
| 380 | + | /// The automation, if the actor is a member of its repository's workspace. | |
| 381 | + | async fn manageable(&self, actor: &User, repo: &RepoPath, id: &str) -> Result<Outcome<AutomationRow>> { | |
| 382 | + | if actor.kind == PrincipalKind::Agent || !actor.is_member(&repo.namespace.to_lowercase()) { | |
| 383 | + | return Ok(fail(FailureCode::Forbidden, format!("Only members of {} can run or change its automations.", repo.namespace))); | |
| 384 | + | } | |
| 385 | + | let row = self | |
| 386 | + | .db | |
| 387 | + | .prepare("SELECT * FROM automations WHERE id = ? AND lower(repo) = lower(?)") | |
| 388 | + | .bind(&[id.into(), format!("{}/{}", repo.namespace, repo.name).into()])? | |
| 389 | + | .first::<AutomationRow>(None) | |
| 390 | + | .await?; | |
| 391 | + | Ok(row.map_or_else(|| fail(FailureCode::NotFound, "No such automation."), Outcome::Ok)) | |
| 392 | + | } | |
| 393 | + | ||
| 394 | + | async fn set_enabled(&self, a: SetEnabledArgs) -> Result<Outcome<Automation>> { | |
| 395 | + | let row = match self.manageable(&a.actor, &a.repo, &a.id).await? { | |
| 396 | + | Outcome::Ok(row) => row, | |
| 397 | + | Outcome::Fail(refused) => return Ok(Outcome::Fail(refused)), | |
| 398 | + | }; | |
| 399 | + | self.db | |
| 400 | + | .prepare("UPDATE automations SET enabled = ? WHERE id = ?") | |
| 401 | + | .bind(&[(a.enabled as u32).into(), row.id.as_str().into()])? | |
| 402 | + | .run() | |
| 403 | + | .await?; | |
| 404 | + | let row = AutomationRow { | |
| 405 | + | enabled: a.enabled as u32, | |
| 406 | + | ..row | |
| 407 | + | }; | |
| 408 | + | Ok(Outcome::Ok(summary(&row, self.last_run(&row.id).await?))) | |
| 409 | + | } | |
| 410 | + | ||
| 411 | + | async fn run_now(&self, a: RunArgs) -> Result<Outcome<AutomationRun>> { | |
| 412 | + | let row = match self.manageable(&a.actor, &a.repo, &a.id).await? { | |
| 413 | + | Outcome::Ok(row) => row, | |
| 414 | + | Outcome::Fail(refused) => return Ok(Outcome::Fail(refused)), | |
| 415 | + | }; | |
| 416 | + | if let Err(problem) = row.definition() { | |
| 417 | + | return Ok(fail(FailureCode::Invalid, format!("Its file has a problem: {problem}"))); | |
| 418 | + | } | |
| 419 | + | let source = Source { | |
| 420 | + | key: format!("manual:{}", new_id("run", now_ms())), | |
| 421 | + | event: "manual".to_owned(), | |
| 422 | + | number: a.number, | |
| 423 | + | actor_id: Some(a.actor.id.clone()), | |
| 424 | + | data: json!({}), | |
| 425 | + | }; | |
| 426 | + | let id = self.run(&row, source).await?; | |
| 427 | + | Ok(match id { | |
| 428 | + | Some(id) => self | |
| 429 | + | .db | |
| 430 | + | .prepare("SELECT * FROM runs WHERE id = ?") | |
| 431 | + | .bind(&[id.as_str().into()])? | |
| 432 | + | .first::<RunRow>(None) | |
| 433 | + | .await? | |
| 434 | + | .map_or_else(|| fail(FailureCode::NotFound, "The run was not recorded."), |row| Outcome::Ok(row.into())), | |
| 435 | + | None => fail(FailureCode::Conflict, "It did not run."), | |
| 436 | + | }) | |
| 437 | + | } | |
| 438 | + | ||
| 439 | + | // --- Running ------------------------------------------------------------------ | |
| 440 | + | ||
| 441 | + | async fn on_event(&self, event: &Event) -> Result<()> { | |
| 442 | + | let Some(repo_id) = event.repo_id.as_deref() else { | |
| 443 | + | return Ok(()); | |
| 444 | + | }; | |
| 445 | + | // A push to the default branch may have changed the files. | |
| 446 | + | if event.kind == "git.push" && event.data["defaultBranch"].as_bool() == Some(true) { | |
| 447 | + | let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?; | |
| 448 | + | if let Some(path) = path { | |
| 449 | + | self.sync(&path).await?; | |
| 450 | + | } | |
| 451 | + | } | |
| 452 | + | let rows = self | |
| 453 | + | .db | |
| 454 | + | .prepare("SELECT * FROM automations WHERE repo_id = ? AND enabled = 1 AND trigger_kind = 'events'") | |
| 455 | + | .bind(&[repo_id.into()])? | |
| 456 | + | .all() | |
| 457 | + | .await? | |
| 458 | + | .results::<AutomationRow>()?; | |
| 459 | + | for row in rows { | |
| 460 | + | let Ok(definition) = row.definition() else { continue }; | |
| 461 | + | let Trigger::Events(events) = &definition.trigger else { continue }; | |
| 462 | + | if !events.contains(&event.kind) { | |
| 463 | + | continue; | |
| 464 | + | } | |
| 465 | + | let source = Source { | |
| 466 | + | key: event.id.clone(), | |
| 467 | + | event: event.kind.clone(), | |
| 468 | + | number: event.data["number"].as_u64().map(|n| n as u32), | |
| 469 | + | actor_id: event.actor.clone(), | |
| 470 | + | data: event.data.clone(), | |
| 471 | + | }; | |
| 472 | + | self.run(&row, source).await?; | |
| 473 | + | } | |
| 474 | + | Ok(()) | |
| 475 | + | } | |
| 476 | + | ||
| 477 | + | /// Runs scheduled automations whose schedule fires this minute. | |
| 478 | + | async fn on_minute(&self, now: u64) -> Result<()> { | |
| 479 | + | let minute = now / 60_000 * 60_000; | |
| 480 | + | let rows = self | |
| 481 | + | .db | |
| 482 | + | .prepare("SELECT * FROM automations WHERE enabled = 1 AND trigger_kind = 'schedule'") | |
| 483 | + | .all() | |
| 484 | + | .await? | |
| 485 | + | .results::<AutomationRow>()?; | |
| 486 | + | for row in rows { | |
| 487 | + | let Ok(definition) = row.definition() else { continue }; | |
| 488 | + | let Trigger::Schedule { schedule, .. } = &definition.trigger else { continue }; | |
| 489 | + | if schedule.fires_at(minute) { | |
| 490 | + | let source = Source { | |
| 491 | + | key: format!("schedule:{minute}"), | |
| 492 | + | event: "schedule".to_owned(), | |
| 493 | + | number: None, | |
| 494 | + | actor_id: None, | |
| 495 | + | data: json!({}), | |
| 496 | + | }; | |
| 497 | + | self.run(&row, source).await?; | |
| 498 | + | } | |
| 499 | + | } | |
| 500 | + | Ok(()) | |
| 501 | + | } | |
| 502 | + | ||
| 503 | + | /// One run of an automation: recorded once, checked against its rules | |
| 504 | + | /// and conditions, then its steps in order until one fails. | |
| 505 | + | async fn run(&self, row: &AutomationRow, source: Source) -> Result<Option<String>> { | |
| 506 | + | let Ok(definition) = row.definition() else { | |
| 507 | + | return Ok(None); | |
| 508 | + | }; | |
| 509 | + | let now = now_ms(); | |
| 510 | + | let run_id = new_id("arn", now); | |
| 511 | + | let claimed = self | |
| 512 | + | .db | |
| 513 | + | .prepare( | |
| 514 | + | "INSERT OR IGNORE INTO runs (id, automation_id, repo_id, event_key, name, event, number, status, steps, started_at) | |
| 515 | + | VALUES (?, ?, ?, ?, ?, ?, ?, 'running', '[]', ?) RETURNING id", | |
| 516 | + | ) | |
| 517 | + | .bind(&[ | |
| 518 | + | run_id.as_str().into(), | |
| 519 | + | row.id.as_str().into(), | |
| 520 | + | row.repo_id.as_str().into(), | |
| 521 | + | source.key.as_str().into(), | |
| 522 | + | definition.name.as_str().into(), | |
| 523 | + | source.event.as_str().into(), | |
| 524 | + | source.number.map_or(JsValue::NULL, JsValue::from), | |
| 525 | + | rfc3339(now).into(), | |
| 526 | + | ])? | |
| 527 | + | .first::<Value>(None) | |
| 528 | + | .await?; | |
| 529 | + | if claimed.is_none() { | |
| 530 | + | return Ok(None); | |
| 531 | + | } | |
| 532 | + | let repo = row.repo_path(); | |
| 533 | + | let Some(actor) = self.workspace_actor(&repo.namespace).await? else { | |
| 534 | + | self.finish(&run_id, "skipped", Some("The workspace no longer exists."), &[], None).await?; | |
| 535 | + | return Ok(Some(run_id)); | |
| 536 | + | }; | |
| 537 | + | ||
| 538 | + | // Who caused it, by name. | |
| 539 | + | let actor_name = match &source.actor_id { | |
| 540 | + | Some(id) if id == AGENT_ID => Some(AGENT_NAME.to_owned()), | |
| 541 | + | Some(id) => { | |
| 542 | + | let names: std::collections::HashMap<String, String> = | |
| 543 | + | g1t_kit::call(&self.identity, "usernames", &UsernamesArgs { ids: vec![id.clone()] }).await?; | |
| 544 | + | names.get(id).cloned() | |
| 545 | + | } | |
| 546 | + | None => None, | |
| 547 | + | }; | |
| 548 | + | ||
| 549 | + | // Not in answer to its own doing. | |
| 550 | + | if source.actor_id.as_deref() == Some(actor.id.as_str()) | |
| 551 | + | && let Some(number) = source.number | |
| 552 | + | { | |
| 553 | + | let recent = self | |
| 554 | + | .db | |
| 555 | + | .prepare("SELECT COUNT(*) AS n FROM effects WHERE automation_id = ? AND repo_id = ? AND number = ? AND at > ?") | |
| 556 | + | .bind(&[ | |
| 557 | + | row.id.as_str().into(), | |
| 558 | + | row.repo_id.as_str().into(), | |
| 559 | + | number.into(), | |
| 560 | + | rfc3339(now.saturating_sub(OWN_EFFECT_MS)).into(), | |
| 561 | + | ])? | |
| 562 | + | .first::<Count>(None) | |
| 563 | + | .await? | |
| 564 | + | .is_some_and(|count| count.n > 0); | |
| 565 | + | if recent { | |
| 566 | + | self.finish(&run_id, "skipped", Some("This came from its own last run, so it did not answer it."), &[], actor_name.as_deref()) | |
| 567 | + | .await?; | |
| 568 | + | return Ok(Some(run_id)); | |
| 569 | + | } | |
| 570 | + | } | |
| 571 | + | ||
| 572 | + | // At most so many runs an hour. | |
| 573 | + | let this_hour = self | |
| 574 | + | .db | |
| 575 | + | .prepare("SELECT COUNT(*) AS n FROM runs WHERE automation_id = ? AND status IN ('succeeded', 'failed') AND started_at > ?") | |
| 576 | + | .bind(&[row.id.as_str().into(), rfc3339(now.saturating_sub(60 * 60 * 1000)).into()])? | |
| 577 | + | .first::<Count>(None) | |
| 578 | + | .await? | |
| 579 | + | .map_or(0, |count| count.n); | |
| 580 | + | if this_hour >= definition.per_hour { | |
| 581 | + | let reason = format!("It has run {} times in the last hour, its limit.", definition.per_hour); | |
| 582 | + | self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?; | |
| 583 | + | return Ok(Some(run_id)); | |
| 584 | + | } | |
| 585 | + | ||
| 586 | + | // What the run knows. | |
| 587 | + | let mut context = Context::default(); | |
| 588 | + | context.set("repo", &row.repo); | |
| 589 | + | context.set("event", &source.event); | |
| 590 | + | context.set("automation", &definition.name); | |
| 591 | + | if let Some(name) = &actor_name { | |
| 592 | + | context.set("actor", name); | |
| 593 | + | } | |
| 594 | + | if let Some(branch) = source.data["ref"].as_str().and_then(|r| r.strip_prefix("refs/heads/")) { | |
| 595 | + | context.set("branch", branch); | |
| 596 | + | } | |
| 597 | + | context.add_data(&source.data); | |
| 598 | + | // The issue or pull request it is about. | |
| 599 | + | let mut target_issue: Option<u32> = None; | |
| 600 | + | let mut target_pull: Option<u32> = None; | |
| 601 | + | if let Some(number) = source.number { | |
| 602 | + | context.set("number", number.to_string()); | |
| 603 | + | let view = ViewArgs { | |
| 604 | + | repo: repo.clone(), | |
| 605 | + | number, | |
| 606 | + | viewer: Some(actor.clone()), | |
| 607 | + | after_seq: 0, | |
| 608 | + | }; | |
| 609 | + | let issue: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view).await?; | |
| 610 | + | if let Outcome::Ok(detail) = issue { | |
| 611 | + | context.set("title", &detail.issue.title); | |
| 612 | + | context.set("url", format!("{SITE}/{}/issues/{number}", row.repo)); | |
| 613 | + | context.labels = detail.issue.labels.clone(); | |
| 614 | + | target_issue = Some(number); | |
| 615 | + | } else { | |
| 616 | + | let pull: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view).await?; | |
| 617 | + | if let Outcome::Ok(detail) = pull { | |
| 618 | + | context.set("title", &detail.pull.title); | |
| 619 | + | context.set("url", format!("{SITE}/{}/pull/{number}", row.repo)); | |
| 620 | + | if let Some(issue) = &detail.issue { | |
| 621 | + | context.labels = issue.labels.clone(); | |
| 622 | + | } | |
| 623 | + | target_pull = Some(number); | |
| 624 | + | target_issue = detail.pull.issue; | |
| 625 | + | } | |
| 626 | + | } | |
| 627 | + | } | |
| 628 | + | if let Err(reason) = definition::holds(&definition.conditions, &context) { | |
| 629 | + | let reason = format!("Its conditions did not hold: {reason}."); | |
| 630 | + | self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?; | |
| 631 | + | return Ok(Some(run_id)); | |
| 632 | + | } | |
| 633 | + | ||
| 634 | + | let mut results: Vec<StepResult> = Vec::new(); | |
| 635 | + | let mut failed = false; | |
| 636 | + | for step in &definition.steps { | |
| 637 | + | if failed { | |
| 638 | + | results.push(StepResult { | |
| 639 | + | step: step.describe(), | |
| 640 | + | ok: false, | |
| 641 | + | detail: "Not run: an earlier step failed.".to_owned(), | |
| 642 | + | }); | |
| 643 | + | continue; | |
| 644 | + | } | |
| 645 | + | let outcome = self | |
| 646 | + | .step(step, &definition, &repo, &actor, &mut context, target_issue, target_pull) | |
| 647 | + | .await | |
| 648 | + | .unwrap_or_else(|error| Err(format!("g1t could not do it: {error}"))); | |
| 649 | + | failed = outcome.is_err(); | |
| 650 | + | results.push(StepResult { | |
| 651 | + | step: step.describe(), | |
| 652 | + | ok: outcome.is_ok(), | |
| 653 | + | detail: outcome.unwrap_or_else(|problem| problem), | |
| 654 | + | }); | |
| 655 | + | } | |
| 656 | + | // Remember what it touched, so it does not answer itself. | |
| 657 | + | if let Some(number) = target_pull.or(target_issue) { | |
| 658 | + | self.db | |
| 659 | + | .prepare("INSERT INTO effects (automation_id, repo_id, number, at) VALUES (?, ?, ?, ?)") | |
| 660 | + | .bind(&[row.id.as_str().into(), row.repo_id.as_str().into(), number.into(), rfc3339(now_ms()).into()])? | |
| 661 | + | .run() | |
| 662 | + | .await?; | |
| 663 | + | } | |
| 664 | + | self.finish(&run_id, if failed { "failed" } else { "succeeded" }, None, &results, actor_name.as_deref()) | |
| 665 | + | .await?; | |
| 666 | + | Ok(Some(run_id)) | |
| 667 | + | } | |
| 668 | + | ||
| 669 | + | async fn finish(&self, run_id: &str, status: &str, reason: Option<&str>, steps: &[StepResult], actor: Option<&str>) -> Result<()> { | |
| 670 | + | self.db | |
| 671 | + | .prepare("UPDATE runs SET status = ?, reason = ?, steps = ?, actor = ? WHERE id = ?") | |
| 672 | + | .bind(&[status.into(), optional(reason), serde_json::to_string(steps)?.into(), optional(actor), run_id.into()])? | |
| 673 | + | .run() | |
| 674 | + | .await?; | |
| 675 | + | Ok(()) | |
| 676 | + | } | |
| 677 | + | ||
| 678 | + | /// One step. `Ok` with what it did, `Err` with why it could not. | |
| 679 | + | #[allow(clippy::too_many_arguments)] | |
| 680 | + | async fn step( | |
| 681 | + | &self, | |
| 682 | + | step: &Step, | |
| 683 | + | definition: &Definition, | |
| 684 | + | repo: &RepoPath, | |
| 685 | + | actor: &User, | |
| 686 | + | context: &mut Context, | |
| 687 | + | issue: Option<u32>, | |
| 688 | + | pull: Option<u32>, | |
| 689 | + | ) -> Result<std::result::Result<String, String>> { | |
| 690 | + | let render = |text: &str| definition::render(text, context); | |
| 691 | + | let signed = |text: &str| format!("{}\n\n<sub>From the automation **{}**.</sub>", render(text), definition.name); | |
| 692 | + | let target = pull.or(issue); | |
| 693 | + | let outcome = |result: Outcome<Value>, done: String| match result { | |
| 694 | + | Outcome::Ok(_) => Ok(done), | |
| 695 | + | Outcome::Fail(refused) => Err(refused.message), | |
| 696 | + | }; | |
| 697 | + | Ok(match step { | |
| 698 | + | Step::Comment(text) => { | |
| 699 | + | let Some(number) = target else { return Ok(Err("There is no issue or pull request to comment on.".to_owned())) }; | |
| 700 | + | let result = g1t_kit::call( | |
| 701 | + | &self.work, | |
| 702 | + | "add_comment", | |
| 703 | + | &json!({ "actor": actor, "repo": repo, "number": number, "body": signed(text) }), | |
| 704 | + | ) | |
| 705 | + | .await?; | |
| 706 | + | outcome(result, format!("Commented on #{number}.")) | |
| 707 | + | } | |
| 708 | + | Step::Label(label) | Step::Unlabel(label) => { | |
| 709 | + | let Some(number) = issue else { | |
| 710 | + | return Ok(Err("Labels belong to issues, and this is not about one.".to_owned())); | |
| 711 | + | }; | |
| 712 | + | let mut labels = context.labels.clone(); | |
| 713 | + | let adding = matches!(step, Step::Label(_)); | |
| 714 | + | labels.retain(|existing| !existing.eq_ignore_ascii_case(label)); | |
| 715 | + | if adding { | |
| 716 | + | labels.push(label.clone()); | |
| 717 | + | } | |
| 718 | + | let result = g1t_kit::call(&self.work, "update_issue", &json!({ "actor": actor, "repo": repo, "number": number, "labels": labels })) | |
| 719 | + | .await?; | |
| 720 | + | context.labels = labels; | |
| 721 | + | outcome(result, format!("{} {label} on #{number}.", if adding { "Labelled" } else { "Removed the label" })) | |
| 722 | + | } | |
| 723 | + | Step::AssignAgent => { | |
| 724 | + | let Some(number) = issue else { | |
| 725 | + | return Ok(Err("An agent is put on an issue, and this is not about one.".to_owned())); | |
| 726 | + | }; | |
| 727 | + | let result: Outcome<Value> = g1t_kit::call(&self.runner, "run", &json!({ "actor": actor, "repo": repo, "issue": number })).await?; | |
| 728 | + | match result { | |
| 729 | + | Outcome::Ok(pull) => Ok(format!("Put a g1t agent on #{number}: pull request #{}.", pull["number"])), | |
| 730 | + | Outcome::Fail(refused) => Err(refused.message), | |
| 731 | + | } | |
| 732 | + | } | |
| 733 | + | Step::MessageAgent(text) => { | |
| 734 | + | let Some(number) = pull else { | |
| 735 | + | return Ok(Err("Agents are messaged on a pull request, and this is not about one.".to_owned())); | |
| 736 | + | }; | |
| 737 | + | let result = g1t_kit::call( | |
| 738 | + | &self.work, | |
| 739 | + | "message_agent", | |
| 740 | + | &json!({ "actor": actor, "repo": repo, "number": number, "body": render(text) }), | |
| 741 | + | ) | |
| 742 | + | .await?; | |
| 743 | + | outcome(result, format!("Messaged the agent on #{number}.")) | |
| 744 | + | } | |
| 745 | + | Step::OpenIssue { title, body, labels, assign_agent } => { | |
| 746 | + | let opened: Outcome<Value> = g1t_kit::call( | |
| 747 | + | &self.work, | |
| 748 | + | "open_issue", | |
| 749 | + | &json!({ | |
| 750 | + | "actor": actor, "repo": repo, "title": render(title), "body": signed(body), | |
| 751 | + | "labels": labels, "checks": [], | |
| 752 | + | }), | |
| 753 | + | ) | |
| 754 | + | .await?; | |
| 755 | + | let number = match opened { | |
| 756 | + | Outcome::Ok(issue) => issue["number"].as_u64().unwrap_or_default() as u32, | |
| 757 | + | Outcome::Fail(refused) => return Ok(Err(refused.message)), | |
| 758 | + | }; | |
| 759 | + | context.set("opened", number.to_string()); | |
| 760 | + | if *assign_agent { | |
| 761 | + | let started: Outcome<Value> = | |
| 762 | + | g1t_kit::call(&self.runner, "run", &json!({ "actor": actor, "repo": repo, "issue": number })).await?; | |
| 763 | + | if let Outcome::Fail(refused) = started { | |
| 764 | + | return Ok(Err(format!("Opened #{number}, but no agent could start: {}", refused.message))); | |
| 765 | + | } | |
| 766 | + | Ok(format!("Opened #{number} and put a g1t agent on it.")) | |
| 767 | + | } else { | |
| 768 | + | Ok(format!("Opened #{number}.")) | |
| 769 | + | } | |
| 770 | + | } | |
| 771 | + | Step::CloseIssue { not_planned } => { | |
| 772 | + | let Some(number) = issue else { return Ok(Err("There is no issue to close.".to_owned())) }; | |
| 773 | + | let reason = if *not_planned { "not_planned" } else { "completed" }; | |
| 774 | + | let result = g1t_kit::call(&self.work, "close_issue", &json!({ "actor": actor, "repo": repo, "number": number, "reason": reason })) | |
| 775 | + | .await?; | |
| 776 | + | outcome(result, format!("Closed #{number}.")) | |
| 777 | + | } | |
| 778 | + | Step::ReopenIssue => { | |
| 779 | + | let Some(number) = issue else { return Ok(Err("There is no issue to reopen.".to_owned())) }; | |
| 780 | + | let result = g1t_kit::call(&self.work, "reopen_issue", &json!({ "actor": actor, "repo": repo, "number": number })).await?; | |
| 781 | + | outcome(result, format!("Reopened #{number}.")) | |
| 782 | + | } | |
| 783 | + | Step::Notify { url, text } => notify(url, &render(text)).await, | |
| 784 | + | }) | |
| 785 | + | } | |
| 786 | + | } | |
| 787 | + | ||
| 788 | + | /// Posts a message, in the shape Slack's, Discord's and most chat tools' | |
| 789 | + | /// incoming webhooks take. | |
| 790 | + | async fn notify(url: &str, text: &str) -> std::result::Result<String, String> { | |
| 791 | + | let host = url.trim_start_matches("https://").split(['/', ':']).next().unwrap_or_default().to_ascii_lowercase(); | |
| 792 | + | if host == "localhost" || host.ends_with(".local") || host.ends_with(".internal") || host.parse::<std::net::IpAddr>().is_ok() { | |
| 793 | + | return Err("notify posts only to public addresses by name.".to_owned()); | |
| 794 | + | } | |
| 795 | + | let send = async { | |
| 796 | + | let headers = Headers::new(); | |
| 797 | + | headers.set("content-type", "application/json")?; | |
| 798 | + | headers.set("user-agent", "g1t-automations/1")?; | |
| 799 | + | let mut init = RequestInit::new(); | |
| 800 | + | init.with_method(Method::Post) | |
| 801 | + | .with_headers(headers) | |
| 802 | + | .with_body(Some(json!({ "text": text, "content": text }).to_string().into())); | |
| 803 | + | let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?; | |
| 804 | + | Ok::<(u16, String), worker::Error>((response.status_code(), response.text().await.unwrap_or_default())) | |
| 805 | + | }; | |
| 806 | + | match send.await { | |
| 807 | + | Ok((status, _)) if (200..300).contains(&status) => Ok(format!("Posted to {host}.")), | |
| 808 | + | Ok((status, body)) => Err(format!("{host} answered {status}: {}", body.chars().take(200).collect::<String>())), | |
| 809 | + | Err(error) => Err(format!("{host} could not be reached: {error}")), | |
| 810 | + | } | |
| 811 | + | } | |
| 812 | + | ||
| 813 | + | #[event(fetch)] | |
| 814 | + | async fn fetch(mut request: Request, env: Env, _ctx: WorkerContext) -> Result<Response> { | |
| 815 | + | let Some(method) = rpc_method(&request) else { | |
| 816 | + | return Response::error("Not found", 404); | |
| 817 | + | }; | |
| 818 | + | let body: Value = request.json().await?; | |
| 819 | + | let service = Automations::new(&env)?; | |
| 820 | + | match method.as_str() { | |
| 821 | + | "list" => reply(&service.list(args(body)?).await?), | |
| 822 | + | "runs" => reply(&service.runs(args(body)?).await?), | |
| 823 | + | "run" => reply(&service.run_now(args(body)?).await?), | |
| 824 | + | "set_enabled" => reply(&service.set_enabled(args(body)?).await?), | |
| 825 | + | _ => Response::error("Unknown method", 404), | |
| 826 | + | } | |
| 827 | + | } | |
| 828 | + | ||
| 829 | + | /// Events from the bus, on this service's own queue. | |
| 830 | + | #[event(queue)] | |
| 831 | + | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: WorkerContext) -> Result<()> { | |
| 832 | + | let service = Automations::new(&env)?; | |
| 833 | + | for message in batch.messages()? { | |
| 834 | + | service.on_event(message.body()).await?; | |
| 835 | + | message.ack(); | |
| 836 | + | } | |
| 837 | + | Ok(()) | |
| 838 | + | } | |
| 839 | + | ||
| 840 | + | /// Every minute: scheduled automations whose time has come. | |
| 841 | + | #[event(scheduled)] | |
| 842 | + | async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { | |
| 843 | + | match Automations::new(&env) { | |
| 844 | + | Ok(service) => { | |
| 845 | + | if let Err(error) = service.on_minute(now_ms()).await { | |
| 846 | + | worker::console_error!("automations: the minute's sweep failed: {error}"); | |
| 847 | + | } | |
| 848 | + | } | |
| 849 | + | Err(error) => worker::console_error!("automations: could not start: {error}"), | |
| 850 | + | } | |
| 851 | + | } |
| 1 | + | { | |
| 2 | + | "$schema": "../../node_modules/wrangler/config-schema.json", | |
| 3 | + | "name": "g1t-automations", | |
| 4 | + | "account_id": "1e6f2cffa3f445920836e8ebe446bb58", | |
| 5 | + | "compatibility_date": "2026-09-26", | |
| 6 | + | // Runs next to its database: a run makes several queries in turn. | |
| 7 | + | "placement": { "mode": "smart" }, | |
| 8 | + | "main": "build/index.js", | |
| 9 | + | "build": { "command": "cargo install -q worker-build@0.8.7 && worker-build --release" }, | |
| 10 | + | // Reached only through service bindings. | |
| 11 | + | "workers_dev": false, | |
| 12 | + | "d1_databases": [ | |
| 13 | + | { | |
| 14 | + | "binding": "DB", | |
| 15 | + | "database_name": "g1t-automations", | |
| 16 | + | "database_id": "39411186-d3e0-4b8c-8c7f-d6416df64e7a", | |
| 17 | + | "migrations_dir": "migrations" | |
| 18 | + | } | |
| 19 | + | ], | |
| 20 | + | "services": [ | |
| 21 | + | { "binding": "REPOS", "service": "g1t-repos" }, | |
| 22 | + | { "binding": "WORK", "service": "g1t-work" }, | |
| 23 | + | { "binding": "IDENTITY", "service": "g1t-identity" }, | |
| 24 | + | { "binding": "RUNNER", "service": "g1t-runner" } | |
| 25 | + | ], | |
| 26 | + | // Every event on the bus: what starts automations, and pushes that | |
| 27 | + | // change their files. | |
| 28 | + | "queues": { | |
| 29 | + | "consumers": [{ "queue": "g1t-events-automations", "max_batch_size": 20, "max_batch_timeout": 1 }] | |
| 30 | + | }, | |
| 31 | + | // Scheduled automations. | |
| 32 | + | "triggers": { "crons": ["* * * * *"] }, | |
| 33 | + | "observability": { "enabled": true } | |
| 34 | + | } |
| 26 | 26 | { "binding": "SUBSCRIBER_WORK", "queue": "g1t-events-work" }, | |
| 27 | 27 | { "binding": "SUBSCRIBER_RUNNER", "queue": "g1t-events-runner" }, | |
| 28 | 28 | { "binding": "SUBSCRIBER_INTEGRATIONS", "queue": "g1t-events-integrations" }, | |
| 29 | − | { "binding": "SUBSCRIBER_WEBHOOKS", "queue": "g1t-events-webhooks" } | |
| 29 | + | { "binding": "SUBSCRIBER_WEBHOOKS", "queue": "g1t-events-webhooks" }, | |
| 30 | + | { "binding": "SUBSCRIBER_AUTOMATIONS", "queue": "g1t-events-automations" } | |
| 30 | 31 | ], | |
| 31 | 32 | "consumers": [{ "queue": "g1t-events", "max_batch_size": 100, "max_batch_timeout": 1 }] | |
| 32 | 33 | }, |
| 879 | 879 | match method.as_str() { | |
| 880 | 880 | "get" => reply(&repos.get(args(body)?).await?), | |
| 881 | 881 | "get_by_id" => reply(&repos.get_by_id(args(body)?).await?), | |
| 882 | + | "path_by_id" => { | |
| 883 | + | let a: PathByIdArgs = args(body)?; | |
| 884 | + | reply( | |
| 885 | + | &repos | |
| 886 | + | .registry | |
| 887 | + | .by_id(&a.id) | |
| 888 | + | .await? | |
| 889 | + | .filter(|repo| repo.fork_of.is_none()) | |
| 890 | + | .map(|repo| RepoPath { | |
| 891 | + | namespace: repo.namespace, | |
| 892 | + | name: repo.name, | |
| 893 | + | }), | |
| 894 | + | ) | |
| 895 | + | } | |
| 882 | 896 | "list" => { | |
| 883 | 897 | let a: ListArgs = args(body)?; | |
| 884 | 898 | reply( |