g1t/services/automations/src/lib.rs

852 lines35,368 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 definition;
17
18use g1t_contracts::automations::*;
19use g1t_contracts::events::Event;
20use g1t_contracts::identity::{AGENT_ID, AGENT_NAME, SlugArgs, UsernamesArgs, Workspace};
21use g1t_contracts::repos::{BlobArgs, BlobView, EntryKind, GetArgs, PathByIdArgs, Repo, RepoPath, TreeArgs, TreeView};
22use g1t_contracts::time::rfc3339;
23use g1t_contracts::work::{IssueDetail, PullDetail, ViewArgs};
24use g1t_contracts::{FailureCode, Membership, Outcome, PrincipalKind, Role, User, Viewer, new_id};
25use g1t_kit::{args, now_ms, reply, rpc_method};
26use serde::Deserialize;
27use serde_json::{Value, json};
28use worker::wasm_bindgen::JsValue;
29use worker::{
30 Context as WorkerContext, D1Database, Env, Fetch, Fetcher, Headers, MessageBatch, MessageExt, Method, Request, RequestInit,
31 Response, Result, ScheduleContext, ScheduledEvent, event,
32};
33
34use definition::{Context, Definition, Step, Trigger};
35
36/// Where automations live in a repository.
37const FOLDER: &str = ".g1t/automations";
38/// The most automation files read from a repository.
39const MAX_FILES: usize = 50;
40const RUNS_SHOWN: u32 = 50;
41/// How long an automation's own effect is remembered, so it does not answer it.
42const OWN_EFFECT_MS: u64 = 10 * 60 * 1000;
43const SITE: &str = "https://g1t.sh";
44
45#[derive(Deserialize)]
46struct AutomationRow {
47 id: String,
48 repo_id: String,
49 repo: String,
50 path: String,
51 name: String,
52 source: String,
53 error: Option<String>,
54 enabled: u32,
55}
56
57impl AutomationRow {
58 fn definition(&self) -> std::result::Result<Definition, String> {
59 if let Some(error) = &self.error {
60 return Err(error.clone());
61 }
62 definition::parse(&self.source, self.path.rsplit('/').next().unwrap_or(&self.path))
63 }
64
65 fn repo_path(&self) -> RepoPath {
66 let (namespace, name) = self.repo.split_once('/').unwrap_or((&self.repo, ""));
67 RepoPath {
68 namespace: namespace.to_owned(),
69 name: name.to_owned(),
70 }
71 }
72}
73
74#[derive(Deserialize)]
75struct RunRow {
76 id: String,
77 automation_id: String,
78 name: String,
79 event: String,
80 number: Option<u32>,
81 status: String,
82 reason: Option<String>,
83 steps: String,
84 actor: Option<String>,
85 started_at: String,
86}
87
88impl From<RunRow> for AutomationRun {
89 fn from(row: RunRow) -> Self {
90 AutomationRun {
91 id: row.id,
92 automation_id: row.automation_id,
93 name: row.name,
94 event: row.event,
95 number: row.number,
96 status: row.status,
97 reason: row.reason,
98 steps: serde_json::from_str(&row.steps).unwrap_or_default(),
99 actor: row.actor,
100 started_at: row.started_at,
101 }
102 }
103}
104
105#[derive(Deserialize)]
106struct Count {
107 n: u32,
108}
109
110/// What started a run.
111struct Source {
112 /// Unique per automation: an event's id, the minute, or a manual run.
113 key: String,
114 event: String,
115 number: Option<u32>,
116 /// The id of whoever caused it.
117 actor_id: Option<String>,
118 /// The event's data, for conditions and `{{data.*}}`.
119 data: Value,
120}
121
122fn optional(value: Option<&str>) -> JsValue {
123 value.map_or(JsValue::NULL, JsValue::from)
124}
125
126fn fail<T>(code: FailureCode, message: impl Into<String>) -> Outcome<T> {
127 Outcome::fail(code, message)
128}
129
130fn summary(row: &AutomationRow, last_run: Option<AutomationRun>) -> Automation {
131 let parsed = row.definition();
132 Automation {
133 id: row.id.clone(),
134 repo: row.repo.clone(),
135 path: row.path.clone(),
136 name: row.name.clone(),
137 trigger: parsed.as_ref().map(|d| d.trigger.describe()).unwrap_or_default(),
138 conditions: parsed.as_ref().map(|d| d.conditions.iter().map(|c| c.describe()).collect()).unwrap_or_default(),
139 steps: parsed.as_ref().map(|d| d.steps.iter().map(Step::describe).collect()).unwrap_or_default(),
140 enabled: row.enabled != 0,
141 error: parsed.as_ref().err().cloned(),
142 manual: parsed.is_ok(),
143 last_run,
144 }
145}
146
147struct Automations {
148 db: D1Database,
149 repos: Fetcher,
150 work: Fetcher,
151 identity: Fetcher,
152 runner: Fetcher,
153}
154
155impl Automations {
156 fn new(env: &Env) -> Result<Self> {
157 Ok(Automations {
158 db: env.d1("DB")?,
159 repos: env.service("REPOS")?,
160 work: env.service("WORK")?,
161 identity: env.service("IDENTITY")?,
162 runner: env.service("RUNNER")?,
163 })
164 }
165
166 /// The workspace itself, as automations act.
167 async fn workspace_actor(&self, slug: &str) -> Result<Option<User>> {
168 let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_owned() }).await?;
169 Ok(workspace.map(|workspace| User {
170 id: workspace.id,
171 username: workspace.slug.clone(),
172 kind: PrincipalKind::Workspace,
173 verified: true,
174 workspaces: vec![Membership {
175 slug: workspace.slug,
176 role: Role::Member,
177 }],
178 }))
179 }
180
181 async fn visible_repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
182 let found: Outcome<Repo> = g1t_kit::call(
183 &self.repos,
184 "get",
185 &GetArgs {
186 path: path.clone(),
187 viewer: viewer.clone(),
188 },
189 )
190 .await?;
191 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
192 }
193
194 // --- Reading the files ----------------------------------------------------
195
196 /// Reads a repository's automation files from its default branch, and
197 /// keeps the database in step with them.
198 async fn sync(&self, path: &RepoPath) -> Result<()> {
199 let Some(actor) = self.workspace_actor(&path.namespace).await? else {
200 return Ok(());
201 };
202 let viewer = Some(actor);
203 let tree: Outcome<TreeView> = g1t_kit::call(
204 &self.repos,
205 "tree",
206 &TreeArgs {
207 path: path.clone(),
208 viewer: viewer.clone(),
209 git_ref: None,
210 tree_path: FOLDER.to_owned(),
211 },
212 )
213 .await?;
214 let Some(repo) = self.visible_repo(path, &viewer).await? else {
215 return Ok(());
216 };
217 let entries = match tree {
218 Outcome::Ok(tree) => tree.entries,
219 // No folder: no automations.
220 Outcome::Fail(_) => Vec::new(),
221 };
222 let files: Vec<String> = entries
223 .into_iter()
224 .filter(|entry| matches!(entry.kind, EntryKind::Blob | EntryKind::Exec))
225 .map(|entry| entry.name)
226 .filter(|name| name.ends_with(".yml") || name.ends_with(".yaml"))
227 .take(MAX_FILES)
228 .collect();
229 let full_name = format!("{}/{}", repo.namespace, repo.name);
230 let now = rfc3339(now_ms());
231 let mut kept = Vec::new();
232 for file in &files {
233 let file_path = format!("{FOLDER}/{file}");
234 let blob: Outcome<BlobView> = g1t_kit::call(
235 &self.repos,
236 "blob",
237 &BlobArgs {
238 path: path.clone(),
239 viewer: viewer.clone(),
240 git_ref: repo.default_branch.clone(),
241 file_path: file_path.clone(),
242 },
243 )
244 .await?;
245 let source = match blob {
246 Outcome::Ok(BlobView { text: Some(text), .. }) => text,
247 _ => continue,
248 };
249 let parsed = definition::parse(&source, file);
250 let (name, error, kind, events) = match &parsed {
251 Ok(definition) => (
252 definition.name.clone(),
253 None,
254 match &definition.trigger {
255 Trigger::Events(_) => "events",
256 Trigger::Schedule { .. } => "schedule",
257 Trigger::Manual => "manual",
258 },
259 match &definition.trigger {
260 Trigger::Events(events) => events.clone(),
261 _ => Vec::new(),
262 },
263 ),
264 Err(problem) => (file.clone(), Some(problem.clone()), "invalid", Vec::new()),
265 };
266 self.db
267 .prepare(
268 "INSERT INTO automations (id, repo_id, repo, path, name, source, trigger_kind, events, error, updated_at)
269 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
270 ON CONFLICT (repo_id, path) DO UPDATE SET
271 repo = excluded.repo, name = excluded.name, source = excluded.source,
272 trigger_kind = excluded.trigger_kind, events = excluded.events,
273 error = excluded.error, updated_at = excluded.updated_at",
274 )
275 .bind(&[
276 new_id("aut", now_ms()).into(),
277 repo.id.as_str().into(),
278 full_name.as_str().into(),
279 file_path.as_str().into(),
280 name.into(),
281 source.into(),
282 kind.into(),
283 serde_json::to_string(&events)?.into(),
284 optional(error.as_deref()),
285 now.as_str().into(),
286 ])?
287 .run()
288 .await?;
289 kept.push(file_path);
290 }
291 // Files that are gone take their automations with them.
292 let mut statements = Vec::new();
293 let existing = self
294 .db
295 .prepare("SELECT * FROM automations WHERE repo_id = ?")
296 .bind(&[repo.id.as_str().into()])?
297 .all()
298 .await?
299 .results::<AutomationRow>()?;
300 for row in existing.iter().filter(|row| !kept.contains(&row.path)) {
301 statements.push(self.db.prepare("DELETE FROM automations WHERE id = ?").bind(&[row.id.as_str().into()])?);
302 }
303 statements.push(
304 self.db
305 .prepare("INSERT OR REPLACE INTO synced (repo_id, at) VALUES (?, ?)")
306 .bind(&[repo.id.as_str().into(), now.into()])?,
307 );
308 self.db.batch(statements).await?;
309 Ok(())
310 }
311
312 async fn synced(&self, repo_id: &str) -> Result<bool> {
313 Ok(self
314 .db
315 .prepare("SELECT COUNT(*) AS n FROM synced WHERE repo_id = ?")
316 .bind(&[repo_id.into()])?
317 .first::<Count>(None)
318 .await?
319 .is_some_and(|count| count.n > 0))
320 }
321
322 // --- Reading and managing ---------------------------------------------------
323
324 async fn last_run(&self, automation_id: &str) -> Result<Option<AutomationRun>> {
325 Ok(self
326 .db
327 .prepare("SELECT * FROM runs WHERE automation_id = ? ORDER BY id DESC LIMIT 1")
328 .bind(&[automation_id.into()])?
329 .first::<RunRow>(None)
330 .await?
331 .map(AutomationRun::from))
332 }
333
334 async fn list(&self, a: ListArgs) -> Result<Outcome<Vec<Automation>>> {
335 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
336 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
337 };
338 if !self.synced(&repo.id).await? {
339 self.sync(&RepoPath {
340 namespace: repo.namespace.clone(),
341 name: repo.name.clone(),
342 })
343 .await?;
344 }
345 let rows = self
346 .db
347 .prepare("SELECT * FROM automations WHERE repo_id = ? ORDER BY path")
348 .bind(&[repo.id.as_str().into()])?
349 .all()
350 .await?
351 .results::<AutomationRow>()?;
352 let mut out = Vec::with_capacity(rows.len());
353 for row in &rows {
354 out.push(summary(row, self.last_run(&row.id).await?));
355 }
356 Ok(Outcome::Ok(out))
357 }
358
359 async fn runs(&self, a: RunsArgs) -> Result<Outcome<Vec<AutomationRun>>> {
360 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
361 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
362 };
363 let rows = match &a.automation {
364 Some(id) => self
365 .db
366 .prepare("SELECT * FROM runs WHERE repo_id = ? AND automation_id = ? ORDER BY id DESC LIMIT ?")
367 .bind(&[repo.id.as_str().into(), id.as_str().into(), RUNS_SHOWN.into()])?,
368 None => self
369 .db
370 .prepare("SELECT * FROM runs WHERE repo_id = ? ORDER BY id DESC LIMIT ?")
371 .bind(&[repo.id.as_str().into(), RUNS_SHOWN.into()])?,
372 }
373 .all()
374 .await?
375 .results::<RunRow>()?;
376 Ok(Outcome::Ok(rows.into_iter().map(AutomationRun::from).collect()))
377 }
378
379 /// The automation, if the actor is a member of its repository's workspace.
380 async fn manageable(&self, actor: &User, repo: &RepoPath, id: &str) -> Result<Outcome<AutomationRow>> {
381 if actor.kind == PrincipalKind::Agent || !actor.is_member(&repo.namespace.to_lowercase()) {
382 return Ok(fail(FailureCode::Forbidden, format!("Only members of {} can run or change its automations.", repo.namespace)));
383 }
384 let row = self
385 .db
386 .prepare("SELECT * FROM automations WHERE id = ? AND lower(repo) = lower(?)")
387 .bind(&[id.into(), format!("{}/{}", repo.namespace, repo.name).into()])?
388 .first::<AutomationRow>(None)
389 .await?;
390 Ok(row.map_or_else(|| fail(FailureCode::NotFound, "No such automation."), Outcome::Ok))
391 }
392
393 async fn set_enabled(&self, a: SetEnabledArgs) -> Result<Outcome<Automation>> {
394 let row = match self.manageable(&a.actor, &a.repo, &a.id).await? {
395 Outcome::Ok(row) => row,
396 Outcome::Fail(refused) => return Ok(Outcome::Fail(refused)),
397 };
398 self.db
399 .prepare("UPDATE automations SET enabled = ? WHERE id = ?")
400 .bind(&[(a.enabled as u32).into(), row.id.as_str().into()])?
401 .run()
402 .await?;
403 let row = AutomationRow {
404 enabled: a.enabled as u32,
405 ..row
406 };
407 Ok(Outcome::Ok(summary(&row, self.last_run(&row.id).await?)))
408 }
409
410 async fn run_now(&self, a: RunArgs) -> Result<Outcome<AutomationRun>> {
411 let row = match self.manageable(&a.actor, &a.repo, &a.id).await? {
412 Outcome::Ok(row) => row,
413 Outcome::Fail(refused) => return Ok(Outcome::Fail(refused)),
414 };
415 if let Err(problem) = row.definition() {
416 return Ok(fail(FailureCode::Invalid, format!("Its file has a problem: {problem}")));
417 }
418 let source = Source {
419 key: format!("manual:{}", new_id("run", now_ms())),
420 event: "manual".to_owned(),
421 number: a.number,
422 actor_id: Some(a.actor.id.clone()),
423 data: json!({}),
424 };
425 let id = self.run(&row, source).await?;
426 Ok(match id {
427 Some(id) => self
428 .db
429 .prepare("SELECT * FROM runs WHERE id = ?")
430 .bind(&[id.as_str().into()])?
431 .first::<RunRow>(None)
432 .await?
433 .map_or_else(|| fail(FailureCode::NotFound, "The run was not recorded."), |row| Outcome::Ok(row.into())),
434 None => fail(FailureCode::Conflict, "It did not run."),
435 })
436 }
437
438 // --- Running ------------------------------------------------------------------
439
440 async fn on_event(&self, event: &Event) -> Result<()> {
441 let Some(repo_id) = event.repo_id.as_deref() else {
442 return Ok(());
443 };
444 // A push to the default branch may have changed the files.
445 if event.kind == "git.push" && event.data["defaultBranch"].as_bool() == Some(true) {
446 let path: Option<RepoPath> = g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await?;
447 if let Some(path) = path {
448 self.sync(&path).await?;
449 }
450 }
451 let rows = self
452 .db
453 .prepare("SELECT * FROM automations WHERE repo_id = ? AND enabled = 1 AND trigger_kind = 'events'")
454 .bind(&[repo_id.into()])?
455 .all()
456 .await?
457 .results::<AutomationRow>()?;
458 for row in rows {
459 let Ok(definition) = row.definition() else { continue };
460 let Trigger::Events(events) = &definition.trigger else { continue };
461 if !events.contains(&event.kind) {
462 continue;
463 }
464 let source = Source {
465 key: event.id.clone(),
466 event: event.kind.clone(),
467 number: event.data["number"].as_u64().map(|n| n as u32),
468 actor_id: event.actor.clone(),
469 data: event.data.clone(),
470 };
471 self.run(&row, source).await?;
472 }
473 Ok(())
474 }
475
476 /// Runs scheduled automations whose schedule fires this minute.
477 async fn on_minute(&self, now: u64) -> Result<()> {
478 let minute = now / 60_000 * 60_000;
479 let rows = self
480 .db
481 .prepare("SELECT * FROM automations WHERE enabled = 1 AND trigger_kind = 'schedule'")
482 .all()
483 .await?
484 .results::<AutomationRow>()?;
485 for row in rows {
486 let Ok(definition) = row.definition() else { continue };
487 let Trigger::Schedule { schedule, .. } = &definition.trigger else { continue };
488 if schedule.fires_at(minute) {
489 let source = Source {
490 key: format!("schedule:{minute}"),
491 event: "schedule".to_owned(),
492 number: None,
493 actor_id: None,
494 data: json!({}),
495 };
496 self.run(&row, source).await?;
497 }
498 }
499 Ok(())
500 }
501
502 /// One run of an automation: recorded once, checked against its rules
503 /// and conditions, then its steps in order until one fails.
504 async fn run(&self, row: &AutomationRow, source: Source) -> Result<Option<String>> {
505 let Ok(definition) = row.definition() else {
506 return Ok(None);
507 };
508 let now = now_ms();
509 let run_id = new_id("arn", now);
510 let claimed = self
511 .db
512 .prepare(
513 "INSERT OR IGNORE INTO runs (id, automation_id, repo_id, event_key, name, event, number, status, steps, started_at)
514 VALUES (?, ?, ?, ?, ?, ?, ?, 'running', '[]', ?) RETURNING id",
515 )
516 .bind(&[
517 run_id.as_str().into(),
518 row.id.as_str().into(),
519 row.repo_id.as_str().into(),
520 source.key.as_str().into(),
521 definition.name.as_str().into(),
522 source.event.as_str().into(),
523 source.number.map_or(JsValue::NULL, JsValue::from),
524 rfc3339(now).into(),
525 ])?
526 .first::<Value>(None)
527 .await?;
528 if claimed.is_none() {
529 return Ok(None);
530 }
531 let repo = row.repo_path();
532 let Some(actor) = self.workspace_actor(&repo.namespace).await? else {
533 self.finish(&run_id, "skipped", Some("The workspace no longer exists."), &[], None).await?;
534 return Ok(Some(run_id));
535 };
536
537 // Who caused it, by name.
538 let actor_name = match &source.actor_id {
539 Some(id) if id == AGENT_ID => Some(AGENT_NAME.to_owned()),
540 Some(id) => {
541 let names: std::collections::HashMap<String, String> =
542 g1t_kit::call(&self.identity, "usernames", &UsernamesArgs { ids: vec![id.clone()] }).await?;
543 names.get(id).cloned()
544 }
545 None => None,
546 };
547
548 // Not in answer to its own doing.
549 if source.actor_id.as_deref() == Some(actor.id.as_str())
550 && let Some(number) = source.number
551 {
552 let recent = self
553 .db
554 .prepare("SELECT COUNT(*) AS n FROM effects WHERE automation_id = ? AND repo_id = ? AND number = ? AND at > ?")
555 .bind(&[
556 row.id.as_str().into(),
557 row.repo_id.as_str().into(),
558 number.into(),
559 rfc3339(now.saturating_sub(OWN_EFFECT_MS)).into(),
560 ])?
561 .first::<Count>(None)
562 .await?
563 .is_some_and(|count| count.n > 0);
564 if recent {
Automations: say plainly why a run was skipped to avoid a loop565 let reason = format!(
566 "An automation made this change, and this one changed #{number} in the last 10 minutes, so it did not answer: that could loop."
567 );
568 self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?;
Automations: rules in .g1t/automations that act when something happens569 return Ok(Some(run_id));
570 }
571 }
572
573 // At most so many runs an hour.
574 let this_hour = self
575 .db
576 .prepare("SELECT COUNT(*) AS n FROM runs WHERE automation_id = ? AND status IN ('succeeded', 'failed') AND started_at > ?")
577 .bind(&[row.id.as_str().into(), rfc3339(now.saturating_sub(60 * 60 * 1000)).into()])?
578 .first::<Count>(None)
579 .await?
580 .map_or(0, |count| count.n);
581 if this_hour >= definition.per_hour {
582 let reason = format!("It has run {} times in the last hour, its limit.", definition.per_hour);
583 self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?;
584 return Ok(Some(run_id));
585 }
586
587 // What the run knows.
588 let mut context = Context::default();
589 context.set("repo", &row.repo);
590 context.set("event", &source.event);
591 context.set("automation", &definition.name);
592 if let Some(name) = &actor_name {
593 context.set("actor", name);
594 }
595 if let Some(branch) = source.data["ref"].as_str().and_then(|r| r.strip_prefix("refs/heads/")) {
596 context.set("branch", branch);
597 }
598 context.add_data(&source.data);
599 // The issue or pull request it is about.
600 let mut target_issue: Option<u32> = None;
601 let mut target_pull: Option<u32> = None;
602 if let Some(number) = source.number {
603 context.set("number", number.to_string());
604 let view = ViewArgs {
605 repo: repo.clone(),
606 number,
607 viewer: Some(actor.clone()),
608 after_seq: 0,
609 };
610 let issue: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view).await?;
611 if let Outcome::Ok(detail) = issue {
612 context.set("title", &detail.issue.title);
613 context.set("url", format!("{SITE}/{}/issues/{number}", row.repo));
614 context.labels = detail.issue.labels.clone();
615 target_issue = Some(number);
616 } else {
617 let pull: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view).await?;
618 if let Outcome::Ok(detail) = pull {
619 context.set("title", &detail.pull.title);
620 context.set("url", format!("{SITE}/{}/pull/{number}", row.repo));
621 if let Some(issue) = &detail.issue {
622 context.labels = issue.labels.clone();
623 }
624 target_pull = Some(number);
625 target_issue = detail.pull.issue;
626 }
627 }
628 }
629 if let Err(reason) = definition::holds(&definition.conditions, &context) {
630 let reason = format!("Its conditions did not hold: {reason}.");
631 self.finish(&run_id, "skipped", Some(&reason), &[], actor_name.as_deref()).await?;
632 return Ok(Some(run_id));
633 }
634
635 let mut results: Vec<StepResult> = Vec::new();
636 let mut failed = false;
637 for step in &definition.steps {
638 if failed {
639 results.push(StepResult {
640 step: step.describe(),
641 ok: false,
642 detail: "Not run: an earlier step failed.".to_owned(),
643 });
644 continue;
645 }
646 let outcome = self
647 .step(step, &definition, &repo, &actor, &mut context, target_issue, target_pull)
648 .await
649 .unwrap_or_else(|error| Err(format!("g1t could not do it: {error}")));
650 failed = outcome.is_err();
651 results.push(StepResult {
652 step: step.describe(),
653 ok: outcome.is_ok(),
654 detail: outcome.unwrap_or_else(|problem| problem),
655 });
656 }
657 // Remember what it touched, so it does not answer itself.
658 if let Some(number) = target_pull.or(target_issue) {
659 self.db
660 .prepare("INSERT INTO effects (automation_id, repo_id, number, at) VALUES (?, ?, ?, ?)")
661 .bind(&[row.id.as_str().into(), row.repo_id.as_str().into(), number.into(), rfc3339(now_ms()).into()])?
662 .run()
663 .await?;
664 }
665 self.finish(&run_id, if failed { "failed" } else { "succeeded" }, None, &results, actor_name.as_deref())
666 .await?;
667 Ok(Some(run_id))
668 }
669
670 async fn finish(&self, run_id: &str, status: &str, reason: Option<&str>, steps: &[StepResult], actor: Option<&str>) -> Result<()> {
671 self.db
672 .prepare("UPDATE runs SET status = ?, reason = ?, steps = ?, actor = ? WHERE id = ?")
673 .bind(&[status.into(), optional(reason), serde_json::to_string(steps)?.into(), optional(actor), run_id.into()])?
674 .run()
675 .await?;
676 Ok(())
677 }
678
679 /// One step. `Ok` with what it did, `Err` with why it could not.
680 #[allow(clippy::too_many_arguments)]
681 async fn step(
682 &self,
683 step: &Step,
684 definition: &Definition,
685 repo: &RepoPath,
686 actor: &User,
687 context: &mut Context,
688 issue: Option<u32>,
689 pull: Option<u32>,
690 ) -> Result<std::result::Result<String, String>> {
691 let render = |text: &str| definition::render(text, context);
692 let signed = |text: &str| format!("{}\n\n<sub>From the automation **{}**.</sub>", render(text), definition.name);
693 let target = pull.or(issue);
694 let outcome = |result: Outcome<Value>, done: String| match result {
695 Outcome::Ok(_) => Ok(done),
696 Outcome::Fail(refused) => Err(refused.message),
697 };
698 Ok(match step {
699 Step::Comment(text) => {
700 let Some(number) = target else { return Ok(Err("There is no issue or pull request to comment on.".to_owned())) };
701 let result = g1t_kit::call(
702 &self.work,
703 "add_comment",
704 &json!({ "actor": actor, "repo": repo, "number": number, "body": signed(text) }),
705 )
706 .await?;
707 outcome(result, format!("Commented on #{number}."))
708 }
709 Step::Label(label) | Step::Unlabel(label) => {
710 let Some(number) = issue else {
711 return Ok(Err("Labels belong to issues, and this is not about one.".to_owned()));
712 };
713 let mut labels = context.labels.clone();
714 let adding = matches!(step, Step::Label(_));
715 labels.retain(|existing| !existing.eq_ignore_ascii_case(label));
716 if adding {
717 labels.push(label.clone());
718 }
719 let result = g1t_kit::call(&self.work, "update_issue", &json!({ "actor": actor, "repo": repo, "number": number, "labels": labels }))
720 .await?;
721 context.labels = labels;
722 outcome(result, format!("{} {label} on #{number}.", if adding { "Labelled" } else { "Removed the label" }))
723 }
724 Step::AssignAgent => {
725 let Some(number) = issue else {
726 return Ok(Err("An agent is put on an issue, and this is not about one.".to_owned()));
727 };
728 let result: Outcome<Value> = g1t_kit::call(&self.runner, "run", &json!({ "actor": actor, "repo": repo, "issue": number })).await?;
729 match result {
730 Outcome::Ok(pull) => Ok(format!("Put a g1t agent on #{number}: pull request #{}.", pull["number"])),
731 Outcome::Fail(refused) => Err(refused.message),
732 }
733 }
734 Step::MessageAgent(text) => {
735 let Some(number) = pull else {
736 return Ok(Err("Agents are messaged on a pull request, and this is not about one.".to_owned()));
737 };
738 let result = g1t_kit::call(
739 &self.work,
740 "message_agent",
741 &json!({ "actor": actor, "repo": repo, "number": number, "body": render(text) }),
742 )
743 .await?;
744 outcome(result, format!("Messaged the agent on #{number}."))
745 }
746 Step::OpenIssue { title, body, labels, assign_agent } => {
747 let opened: Outcome<Value> = g1t_kit::call(
748 &self.work,
749 "open_issue",
750 &json!({
751 "actor": actor, "repo": repo, "title": render(title), "body": signed(body),
752 "labels": labels, "checks": [],
753 }),
754 )
755 .await?;
756 let number = match opened {
757 Outcome::Ok(issue) => issue["number"].as_u64().unwrap_or_default() as u32,
758 Outcome::Fail(refused) => return Ok(Err(refused.message)),
759 };
760 context.set("opened", number.to_string());
761 if *assign_agent {
762 let started: Outcome<Value> =
763 g1t_kit::call(&self.runner, "run", &json!({ "actor": actor, "repo": repo, "issue": number })).await?;
764 if let Outcome::Fail(refused) = started {
765 return Ok(Err(format!("Opened #{number}, but no agent could start: {}", refused.message)));
766 }
767 Ok(format!("Opened #{number} and put a g1t agent on it."))
768 } else {
769 Ok(format!("Opened #{number}."))
770 }
771 }
772 Step::CloseIssue { not_planned } => {
773 let Some(number) = issue else { return Ok(Err("There is no issue to close.".to_owned())) };
774 let reason = if *not_planned { "not_planned" } else { "completed" };
775 let result = g1t_kit::call(&self.work, "close_issue", &json!({ "actor": actor, "repo": repo, "number": number, "reason": reason }))
776 .await?;
777 outcome(result, format!("Closed #{number}."))
778 }
779 Step::ReopenIssue => {
780 let Some(number) = issue else { return Ok(Err("There is no issue to reopen.".to_owned())) };
781 let result = g1t_kit::call(&self.work, "reopen_issue", &json!({ "actor": actor, "repo": repo, "number": number })).await?;
782 outcome(result, format!("Reopened #{number}."))
783 }
784 Step::Notify { url, text } => notify(url, &render(text)).await,
785 })
786 }
787}
788
789/// Posts a message, in the shape Slack's, Discord's and most chat tools'
790/// incoming webhooks take.
791async fn notify(url: &str, text: &str) -> std::result::Result<String, String> {
792 let host = url.trim_start_matches("https://").split(['/', ':']).next().unwrap_or_default().to_ascii_lowercase();
793 if host == "localhost" || host.ends_with(".local") || host.ends_with(".internal") || host.parse::<std::net::IpAddr>().is_ok() {
794 return Err("notify posts only to public addresses by name.".to_owned());
795 }
796 let send = async {
797 let headers = Headers::new();
798 headers.set("content-type", "application/json")?;
799 headers.set("user-agent", "g1t-automations/1")?;
800 let mut init = RequestInit::new();
801 init.with_method(Method::Post)
802 .with_headers(headers)
803 .with_body(Some(json!({ "text": text, "content": text }).to_string().into()));
804 let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?;
805 Ok::<(u16, String), worker::Error>((response.status_code(), response.text().await.unwrap_or_default()))
806 };
807 match send.await {
808 Ok((status, _)) if (200..300).contains(&status) => Ok(format!("Posted to {host}.")),
809 Ok((status, body)) => Err(format!("{host} answered {status}: {}", body.chars().take(200).collect::<String>())),
810 Err(error) => Err(format!("{host} could not be reached: {error}")),
811 }
812}
813
814#[event(fetch)]
815async fn fetch(mut request: Request, env: Env, _ctx: WorkerContext) -> Result<Response> {
816 let Some(method) = rpc_method(&request) else {
817 return Response::error("Not found", 404);
818 };
819 let body: Value = request.json().await?;
820 let service = Automations::new(&env)?;
821 match method.as_str() {
822 "list" => reply(&service.list(args(body)?).await?),
823 "runs" => reply(&service.runs(args(body)?).await?),
824 "run" => reply(&service.run_now(args(body)?).await?),
825 "set_enabled" => reply(&service.set_enabled(args(body)?).await?),
826 _ => Response::error("Unknown method", 404),
827 }
828}
829
830/// Events from the bus, on this service's own queue.
831#[event(queue)]
832async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: WorkerContext) -> Result<()> {
833 let service = Automations::new(&env)?;
834 for message in batch.messages()? {
835 service.on_event(message.body()).await?;
836 message.ack();
837 }
838 Ok(())
839}
840
841/// Every minute: scheduled automations whose time has come.
842#[event(scheduled)]
843async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
844 match Automations::new(&env) {
845 Ok(service) => {
846 if let Err(error) = service.on_minute(now_ms()).await {
847 worker::console_error!("automations: the minute's sweep failed: {error}");
848 }
849 }
850 Err(error) => worker::console_error!("automations: could not start: {error}"),
851 }
852}