| 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 | let reason = format!( |
| 567 | "An automation made this change, and this one changed #{number} in the last 10 minutes, so it did not answer: that could loop." |
| 568 | ); |
| 569 | self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?; |
| 570 | return Ok(Some(run_id)); |
| 571 | } |
| 572 | } |
| 573 | |
| 574 | // At most so many runs an hour. |
| 575 | let this_hour = self |
| 576 | .db |
| 577 | .prepare("SELECT COUNT(*) AS n FROM runs WHERE automation_id = ? AND status IN ('succeeded', 'failed') AND started_at > ?") |
| 578 | .bind(&[row.id.as_str().into(), rfc3339(now.saturating_sub(60 * 60 * 1000)).into()])? |
| 579 | .first::<Count>(None) |
| 580 | .await? |
| 581 | .map_or(0, |count| count.n); |
| 582 | if this_hour >= definition.per_hour { |
| 583 | let reason = format!("It has run {} times in the last hour, its limit.", definition.per_hour); |
| 584 | self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?; |
| 585 | return Ok(Some(run_id)); |
| 586 | } |
| 587 | |
| 588 | // What the run knows. |
| 589 | let mut context = Context::default(); |
| 590 | context.set("repo", &row.repo); |
| 591 | context.set("event", &source.event); |
| 592 | context.set("automation", &definition.name); |
| 593 | if let Some(name) = &actor_name { |
| 594 | context.set("actor", name); |
| 595 | } |
| 596 | if let Some(branch) = source.data["ref"].as_str().and_then(|r| r.strip_prefix("refs/heads/")) { |
| 597 | context.set("branch", branch); |
| 598 | } |
| 599 | context.add_data(&source.data); |
| 600 | // The issue or pull request it is about. |
| 601 | let mut target_issue: Option<u32> = None; |
| 602 | let mut target_pull: Option<u32> = None; |
| 603 | if let Some(number) = source.number { |
| 604 | context.set("number", number.to_string()); |
| 605 | let view = ViewArgs { |
| 606 | repo: repo.clone(), |
| 607 | number, |
| 608 | viewer: Some(actor.clone()), |
| 609 | after_seq: 0, |
| 610 | }; |
| 611 | let issue: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view).await?; |
| 612 | if let Outcome::Ok(detail) = issue { |
| 613 | context.set("title", &detail.issue.title); |
| 614 | context.set("url", format!("{SITE}/{}/issues/{number}", row.repo)); |
| 615 | context.labels = detail.issue.labels.clone(); |
| 616 | target_issue = Some(number); |
| 617 | } else { |
| 618 | let pull: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view).await?; |
| 619 | if let Outcome::Ok(detail) = pull { |
| 620 | context.set("title", &detail.pull.title); |
| 621 | context.set("url", format!("{SITE}/{}/pull/{number}", row.repo)); |
| 622 | if let Some(issue) = &detail.issue { |
| 623 | context.labels = issue.labels.clone(); |
| 624 | } |
| 625 | target_pull = Some(number); |
| 626 | target_issue = detail.pull.issue; |
| 627 | } |
| 628 | } |
| 629 | } |
| 630 | if let Err(reason) = definition::holds(&definition.conditions, &context) { |
| 631 | let reason = format!("Its conditions did not hold: {reason}."); |
| 632 | self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?; |
| 633 | return Ok(Some(run_id)); |
| 634 | } |
| 635 | |
| 636 | let mut results: Vec<StepResult> = Vec::new(); |
| 637 | let mut failed = false; |
| 638 | for step in &definition.steps { |
| 639 | if failed { |
| 640 | results.push(StepResult { |
| 641 | step: step.describe(), |
| 642 | ok: false, |
| 643 | detail: "Not run: an earlier step failed.".to_owned(), |
| 644 | }); |
| 645 | continue; |
| 646 | } |
| 647 | let outcome = self |
| 648 | .step(step, &definition, &repo, &actor, &mut context, target_issue, target_pull) |
| 649 | .await |
| 650 | .unwrap_or_else(|error| Err(format!("g1t could not do it: {error}"))); |
| 651 | failed = outcome.is_err(); |
| 652 | results.push(StepResult { |
| 653 | step: step.describe(), |
| 654 | ok: outcome.is_ok(), |
| 655 | detail: outcome.unwrap_or_else(|problem| problem), |
| 656 | }); |
| 657 | } |
| 658 | // Remember what it touched, so it does not answer itself. |
| 659 | if let Some(number) = target_pull.or(target_issue) { |
| 660 | self.db |
| 661 | .prepare("INSERT INTO effects (automation_id, repo_id, number, at) VALUES (?, ?, ?, ?)") |
| 662 | .bind(&[row.id.as_str().into(), row.repo_id.as_str().into(), number.into(), rfc3339(now_ms()).into()])? |
| 663 | .run() |
| 664 | .await?; |
| 665 | } |
| 666 | self.finish(&run_id, if failed { "failed" } else { "succeeded" }, None, &results, actor_name.as_deref()) |
| 667 | .await?; |
| 668 | Ok(Some(run_id)) |
| 669 | } |
| 670 | |
| 671 | async fn finish(&self, run_id: &str, status: &str, reason: Option<&str>, steps: &[StepResult], actor: Option<&str>) -> Result<()> { |
| 672 | self.db |
| 673 | .prepare("UPDATE runs SET status = ?, reason = ?, steps = ?, actor = ? WHERE id = ?") |
| 674 | .bind(&[status.into(), optional(reason), serde_json::to_string(steps)?.into(), optional(actor), run_id.into()])? |
| 675 | .run() |
| 676 | .await?; |
| 677 | Ok(()) |
| 678 | } |
| 679 | |
| 680 | /// One step. `Ok` with what it did, `Err` with why it could not. |
| 681 | #[allow(clippy::too_many_arguments)] |
| 682 | async fn step( |
| 683 | &self, |
| 684 | step: &Step, |
| 685 | definition: &Definition, |
| 686 | repo: &RepoPath, |
| 687 | actor: &User, |
| 688 | context: &mut Context, |
| 689 | issue: Option<u32>, |
| 690 | pull: Option<u32>, |
| 691 | ) -> Result<std::result::Result<String, String>> { |
| 692 | let render = |text: &str| definition::render(text, context); |
| 693 | let signed = |text: &str| format!("{}\n\n<sub>From the automation **{}**.</sub>", render(text), definition.name); |
| 694 | let target = pull.or(issue); |
| 695 | let outcome = |result: Outcome<Value>, done: String| match result { |
| 696 | Outcome::Ok(_) => Ok(done), |
| 697 | Outcome::Fail(refused) => Err(refused.message), |
| 698 | }; |
| 699 | Ok(match step { |
| 700 | Step::Comment(text) => { |
| 701 | let Some(number) = target else { return Ok(Err("There is no issue or pull request to comment on.".to_owned())) }; |
| 702 | let result = g1t_kit::call( |
| 703 | &self.work, |
| 704 | "add_comment", |
| 705 | &json!({ "actor": actor, "repo": repo, "number": number, "body": signed(text) }), |
| 706 | ) |
| 707 | .await?; |
| 708 | outcome(result, format!("Commented on #{number}.")) |
| 709 | } |
| 710 | Step::Label(label) | Step::Unlabel(label) => { |
| 711 | let Some(number) = issue else { |
| 712 | return Ok(Err("Labels belong to issues, and this is not about one.".to_owned())); |
| 713 | }; |
| 714 | let mut labels = context.labels.clone(); |
| 715 | let adding = matches!(step, Step::Label(_)); |
| 716 | labels.retain(|existing| !existing.eq_ignore_ascii_case(label)); |
| 717 | if adding { |
| 718 | labels.push(label.clone()); |
| 719 | } |
| 720 | let result = g1t_kit::call(&self.work, "update_issue", &json!({ "actor": actor, "repo": repo, "number": number, "labels": labels })) |
| 721 | .await?; |
| 722 | context.labels = labels; |
| 723 | outcome(result, format!("{} {label} on #{number}.", if adding { "Labelled" } else { "Removed the label" })) |
| 724 | } |
| 725 | Step::AssignAgent => { |
| 726 | let Some(number) = issue else { |
| 727 | return Ok(Err("An agent is put on an issue, and this is not about one.".to_owned())); |
| 728 | }; |
| 729 | let result: Outcome<Value> = g1t_kit::call(&self.runner, "run", &json!({ "actor": actor, "repo": repo, "issue": number })).await?; |
| 730 | match result { |
| 731 | Outcome::Ok(pull) => Ok(format!("Put a g1t agent on #{number}: pull request #{}.", pull["number"])), |
| 732 | Outcome::Fail(refused) => Err(refused.message), |
| 733 | } |
| 734 | } |
| 735 | Step::MessageAgent(text) => { |
| 736 | let Some(number) = pull else { |
| 737 | return Ok(Err("Agents are messaged on a pull request, and this is not about one.".to_owned())); |
| 738 | }; |
| 739 | let result = g1t_kit::call( |
| 740 | &self.work, |
| 741 | "message_agent", |
| 742 | &json!({ "actor": actor, "repo": repo, "number": number, "body": render(text) }), |
| 743 | ) |
| 744 | .await?; |
| 745 | outcome(result, format!("Messaged the agent on #{number}.")) |
| 746 | } |
| 747 | Step::OpenIssue { title, body, labels, assign_agent } => { |
| 748 | let opened: Outcome<Value> = g1t_kit::call( |
| 749 | &self.work, |
| 750 | "open_issue", |
| 751 | &json!({ |
| 752 | "actor": actor, "repo": repo, "title": render(title), "body": signed(body), |
| 753 | "labels": labels, "checks": [], |
| 754 | }), |
| 755 | ) |
| 756 | .await?; |
| 757 | let number = match opened { |
| 758 | Outcome::Ok(issue) => issue["number"].as_u64().unwrap_or_default() as u32, |
| 759 | Outcome::Fail(refused) => return Ok(Err(refused.message)), |
| 760 | }; |
| 761 | context.set("opened", number.to_string()); |
| 762 | if *assign_agent { |
| 763 | let started: Outcome<Value> = |
| 764 | g1t_kit::call(&self.runner, "run", &json!({ "actor": actor, "repo": repo, "issue": number })).await?; |
| 765 | if let Outcome::Fail(refused) = started { |
| 766 | return Ok(Err(format!("Opened #{number}, but no agent could start: {}", refused.message))); |
| 767 | } |
| 768 | Ok(format!("Opened #{number} and put a g1t agent on it.")) |
| 769 | } else { |
| 770 | Ok(format!("Opened #{number}.")) |
| 771 | } |
| 772 | } |
| 773 | Step::CloseIssue { not_planned } => { |
| 774 | let Some(number) = issue else { return Ok(Err("There is no issue to close.".to_owned())) }; |
| 775 | let reason = if *not_planned { "not_planned" } else { "completed" }; |
| 776 | let result = g1t_kit::call(&self.work, "close_issue", &json!({ "actor": actor, "repo": repo, "number": number, "reason": reason })) |
| 777 | .await?; |
| 778 | outcome(result, format!("Closed #{number}.")) |
| 779 | } |
| 780 | Step::ReopenIssue => { |
| 781 | let Some(number) = issue else { return Ok(Err("There is no issue to reopen.".to_owned())) }; |
| 782 | let result = g1t_kit::call(&self.work, "reopen_issue", &json!({ "actor": actor, "repo": repo, "number": number })).await?; |
| 783 | outcome(result, format!("Reopened #{number}.")) |
| 784 | } |
| 785 | Step::Notify { url, text } => notify(url, &render(text)).await, |
| 786 | }) |
| 787 | } |
| 788 | } |
| 789 | |
| 790 | /// Posts a message, in the shape Slack's, Discord's and most chat tools' |
| 791 | /// incoming webhooks take. |
| 792 | async fn notify(url: &str, text: &str) -> std::result::Result<String, String> { |
| 793 | let host = url.trim_start_matches("https://").split(['/', ':']).next().unwrap_or_default().to_ascii_lowercase(); |
| 794 | if host == "localhost" || host.ends_with(".local") || host.ends_with(".internal") || host.parse::<std::net::IpAddr>().is_ok() { |
| 795 | return Err("notify posts only to public addresses by name.".to_owned()); |
| 796 | } |
| 797 | let send = async { |
| 798 | let headers = Headers::new(); |
| 799 | headers.set("content-type", "application/json")?; |
| 800 | headers.set("user-agent", "g1t-automations/1")?; |
| 801 | let mut init = RequestInit::new(); |
| 802 | init.with_method(Method::Post) |
| 803 | .with_headers(headers) |
| 804 | .with_body(Some(json!({ "text": text, "content": text }).to_string().into())); |
| 805 | let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?; |
| 806 | Ok::<(u16, String), worker::Error>((response.status_code(), response.text().await.unwrap_or_default())) |
| 807 | }; |
| 808 | match send.await { |
| 809 | Ok((status, _)) if (200..300).contains(&status) => Ok(format!("Posted to {host}.")), |
| 810 | Ok((status, body)) => Err(format!("{host} answered {status}: {}", body.chars().take(200).collect::<String>())), |
| 811 | Err(error) => Err(format!("{host} could not be reached: {error}")), |
| 812 | } |
| 813 | } |
| 814 | |
| 815 | #[event(fetch)] |
| 816 | async fn fetch(mut request: Request, env: Env, _ctx: WorkerContext) -> Result<Response> { |
| 817 | let Some(method) = rpc_method(&request) else { |
| 818 | return Response::error("Not found", 404); |
| 819 | }; |
| 820 | let body: Value = request.json().await?; |
| 821 | let service = Automations::new(&env)?; |
| 822 | match method.as_str() { |
| 823 | "list" => reply(&service.list(args(body)?).await?), |
| 824 | "runs" => reply(&service.runs(args(body)?).await?), |
| 825 | "run" => reply(&service.run_now(args(body)?).await?), |
| 826 | "set_enabled" => reply(&service.set_enabled(args(body)?).await?), |
| 827 | _ => Response::error("Unknown method", 404), |
| 828 | } |
| 829 | } |
| 830 | |
| 831 | /// Events from the bus, on this service's own queue. |
| 832 | #[event(queue)] |
| 833 | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: WorkerContext) -> Result<()> { |
| 834 | let service = Automations::new(&env)?; |
| 835 | for message in batch.messages()? { |
| 836 | service.on_event(message.body()).await?; |
| 837 | message.ack(); |
| 838 | } |
| 839 | Ok(()) |
| 840 | } |
| 841 | |
| 842 | /// Every minute: scheduled automations whose time has come. |
| 843 | #[event(scheduled)] |
| 844 | async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) { |
| 845 | match Automations::new(&env) { |
| 846 | Ok(service) => { |
| 847 | if let Err(error) = service.on_minute(now_ms()).await { |
| 848 | worker::console_error!("automations: the minute's sweep failed: {error}"); |
| 849 | } |
| 850 | } |
| 851 | Err(error) => worker::console_error!("automations: could not start: {error}"), |
| 852 | } |
| 853 | } |