g1t/services/automations/src/lib.rs

851 lines35,246 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Automations: rules in .g1t/automations that act when something happens1//! 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
16mod cron;
17mod definition;
18
19use g1t_contracts::automations::*;
20use g1t_contracts::events::Event;
21use g1t_contracts::identity::{AGENT_ID, AGENT_NAME, SlugArgs, UsernamesArgs, Workspace};
22use g1t_contracts::repos::{BlobArgs, BlobView, EntryKind, GetArgs, PathByIdArgs, Repo, RepoPath, TreeArgs, TreeView};
23use g1t_contracts::time::rfc3339;
24use g1t_contracts::work::{IssueDetail, PullDetail, ViewArgs};
25use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, Role, User, Viewer, new_id};
26use g1t_kit::{args, now_ms, reply, rpc_method};
27use serde::Deserialize;
28use serde_json::{Value, json};
29use worker::wasm_bindgen::JsValue;
30use worker::{
31 Context as WorkerContext, D1Database, Env, Fetch, Fetcher, Headers, MessageBatch, MessageExt, Method, Request, RequestInit,
32 Response, Result, ScheduleContext, ScheduledEvent, event,
33};
34
35use definition::{Context, Definition, Step, Trigger};
36
37/// Where automations live in a repository.
38const FOLDER: &str = ".g1t/automations";
39/// The most automation files read from a repository.
40const MAX_FILES: usize = 50;
41const RUNS_SHOWN: u32 = 50;
42/// How long an automation's own effect is remembered, so it does not answer it.
43const OWN_EFFECT_MS: u64 = 10 * 60 * 1000;
44const SITE: &str = "https://g1t.sh";
45
46#[derive(Deserialize)]
47struct 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
58impl 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)]
76struct 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
89impl 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)]
107struct Count {
108 n: u32,
109}
110
111/// What started a run.
112struct 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
123fn optional(value: Option<&str>) -> JsValue {
124 value.map_or(JsValue::NULL, JsValue::from)
125}
126
127fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
128 Outcome::fail(code, message)
129}
130
131fn 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
148struct Automations {
149 db: D1Database,
150 repos: Fetcher,
151 work: Fetcher,
152 identity: Fetcher,
153 runner: Fetcher,
154}
155
156impl 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.
790async 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)]
814async 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)]
831async 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)]
842async 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}