g1t/services/actions/src/sync.rs
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.
| GitHub Actions on g1t, part two: running workflows | 1 | //! Reading workflow files: at any commit for a run, and from the default |
| 2 | //! branch into the `workflows` table, which lists them, holds their | |
| 3 | //! schedules and remembers which are turned off. | |
| 4 | ||
| 5 | use g1t_actions::workflow::{self, FOLDER}; | |
| 6 | use g1t_contracts::new_id; | |
| 7 | use g1t_contracts::repos::{BlobArgs, BlobView, EntryKind, Repo, RepoPath, TreeArgs, TreeView}; | |
| 8 | use g1t_contracts::time::rfc3339; | |
| 9 | use g1t_contracts::{Outcome, User}; | |
| 10 | use g1t_kit::now_ms; | |
| 11 | use serde::Deserialize; | |
| 12 | use worker::Result; | |
| 13 | ||
| 14 | use crate::{Actions, Count, MAX_WORKFLOWS, optional}; | |
| 15 | ||
| 16 | /// One workflow file as read at a commit. | |
| 17 | pub struct WorkflowFile { | |
| 18 | pub path: String, | |
| 19 | pub source: String, | |
| 20 | } | |
| 21 | ||
| 22 | /// The files at a commit, and the commit the ref resolved to. | |
| 23 | pub struct Read { | |
| 24 | pub files: Vec<WorkflowFile>, | |
| 25 | pub head: Option<String>, | |
| 26 | } | |
| 27 | ||
| 28 | #[derive(Deserialize)] | |
| 29 | pub struct WorkflowRow { | |
| 30 | pub id: String, | |
| 31 | pub repo_id: String, | |
| 32 | pub repo: String, | |
| 33 | pub path: String, | |
| 34 | pub name: String, | |
| 35 | pub source: String, | |
| 36 | pub events: String, | |
| 37 | pub crons: String, | |
| 38 | pub error: Option<String>, | |
| 39 | pub state: String, | |
| 40 | pub updated_at: String, | |
| 41 | } | |
| 42 | ||
| Fewer Artifacts reads: the store is asked for a handle only when needed, objects are kept in the isolate, and issue events read no workflows nobody listens for | 43 | /// What a workflow's row says once its file has left the default branch. |
| 44 | pub const GONE: &str = "Its file is no longer on the default branch."; | |
| 45 | ||
| 46 | /// Whether a synced workflow could start on `event`: it lists the event, | |
| 47 | /// or it did not parse (and is still on the branch), so reading it again | |
| 48 | /// reports why. | |
| 49 | pub fn could_start(events: &str, error: Option<&str>, event: &str) -> bool { | |
| 50 | match error { | |
| 51 | Some(error) => error != GONE, | |
| 52 | None => serde_json::from_str::<Vec<String>>(events).is_ok_and(|events| events.iter().any(|e| e == event)), | |
| 53 | } | |
| 54 | } | |
| 55 | ||
| GitHub Actions on g1t, part two: running workflows | 56 | impl Actions { |
| 57 | /// The workflow files of `path` as of `git_ref` (the default branch | |
| 58 | /// when absent). | |
| 59 | pub async fn read_workflows(&self, path: &RepoPath, actor: &User, git_ref: Option<&str>) -> Result<Read> { | |
| 60 | let viewer = Some(actor.clone()); | |
| 61 | let tree: Outcome<TreeView> = g1t_kit::call( | |
| 62 | &self.repos, | |
| 63 | "tree", | |
| 64 | &TreeArgs { | |
| 65 | path: path.clone(), | |
| 66 | viewer: viewer.clone(), | |
| 67 | git_ref: git_ref.map(str::to_owned), | |
| 68 | tree_path: FOLDER.to_owned(), | |
| 69 | }, | |
| 70 | ) | |
| 71 | .await?; | |
| 72 | let (entries, head, resolved) = match tree { | |
| 73 | Outcome::Ok(tree) => (tree.entries, tree.head.map(|commit| commit.hash), tree.git_ref), | |
| 74 | // No folder: no workflows. | |
| 75 | Outcome::Fail(_) => return Ok(Read { files: Vec::new(), head: None }), | |
| 76 | }; | |
| 77 | let at = head.clone().unwrap_or(resolved); | |
| 78 | let mut files = Vec::new(); | |
| 79 | for entry in entries | |
| 80 | .into_iter() | |
| 81 | .filter(|entry| matches!(entry.kind, EntryKind::Blob | EntryKind::Exec)) | |
| 82 | .filter(|entry| entry.name.ends_with(".yml") || entry.name.ends_with(".yaml")) | |
| 83 | .take(MAX_WORKFLOWS) | |
| 84 | { | |
| 85 | let file_path = format!("{FOLDER}/{}", entry.name); | |
| 86 | let blob: Outcome<BlobView> = g1t_kit::call( | |
| 87 | &self.repos, | |
| 88 | "blob", | |
| 89 | &BlobArgs { | |
| 90 | path: path.clone(), | |
| 91 | viewer: viewer.clone(), | |
| 92 | git_ref: at.clone(), | |
| 93 | file_path: file_path.clone(), | |
| 94 | }, | |
| 95 | ) | |
| 96 | .await?; | |
| 97 | if let Outcome::Ok(BlobView { text: Some(source), .. }) = blob { | |
| 98 | files.push(WorkflowFile { path: file_path, source }); | |
| 99 | } | |
| 100 | } | |
| 101 | Ok(Read { files, head }) | |
| 102 | } | |
| 103 | ||
| 104 | /// Keeps the `workflows` table in step with the default branch. | |
| 105 | pub async fn sync(&self, repo: &Repo, actor: &User) -> Result<()> { | |
| 106 | let path = RepoPath { | |
| 107 | namespace: repo.namespace.clone(), | |
| 108 | name: repo.name.clone(), | |
| 109 | }; | |
| 110 | let read = self.read_workflows(&path, actor, None).await?; | |
| 111 | let full_name = format!("{}/{}", repo.namespace, repo.name); | |
| 112 | let now = rfc3339(now_ms()); | |
| 113 | let mut statements = Vec::new(); | |
| 114 | for file in &read.files { | |
| 115 | let parsed = workflow::parse(&file.source); | |
| 116 | let (name, error, events, crons) = match &parsed { | |
| 117 | Ok(parsed) => ( | |
| 118 | parsed.display_name(&file.path), | |
| 119 | None, | |
| 120 | parsed.triggers.iter().map(|t| t.event.clone()).collect::<Vec<_>>(), | |
| 121 | parsed.trigger("schedule").map(|t| t.crons.clone()).unwrap_or_default(), | |
| 122 | ), | |
| 123 | Err(problem) => (file.path.clone(), Some(problem.clone()), Vec::new(), Vec::new()), | |
| 124 | }; | |
| 125 | statements.push( | |
| 126 | self.db | |
| 127 | .prepare( | |
| 128 | "INSERT INTO workflows (id, repo_id, repo, path, name, source, events, crons, error, updated_at) | |
| 129 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) | |
| 130 | ON CONFLICT (repo_id, path) DO UPDATE SET | |
| 131 | repo = excluded.repo, name = excluded.name, source = excluded.source, | |
| 132 | events = excluded.events, crons = excluded.crons, error = excluded.error, | |
| 133 | updated_at = excluded.updated_at", | |
| 134 | ) | |
| 135 | .bind(&[ | |
| 136 | new_id("wfl", now_ms()).into(), | |
| 137 | repo.id.as_str().into(), | |
| 138 | full_name.as_str().into(), | |
| 139 | file.path.as_str().into(), | |
| 140 | name.into(), | |
| 141 | file.source.as_str().into(), | |
| 142 | serde_json::to_string(&events)?.into(), | |
| 143 | serde_json::to_string(&crons)?.into(), | |
| 144 | optional(error.as_deref()), | |
| 145 | now.as_str().into(), | |
| 146 | ])?, | |
| 147 | ); | |
| 148 | } | |
| 149 | // A workflow whose file is gone keeps its runs, but no longer runs | |
| 150 | // on schedule or by hand: it is listed only while it has runs. | |
| 151 | let kept: Vec<String> = read.files.iter().map(|file| file.path.clone()).collect(); | |
| 152 | let existing = self | |
| 153 | .db | |
| 154 | .prepare("SELECT * FROM workflows WHERE repo_id = ?") | |
| 155 | .bind(&[repo.id.as_str().into()])? | |
| 156 | .all() | |
| 157 | .await? | |
| 158 | .results::<WorkflowRow>()?; | |
| 159 | for row in existing.iter().filter(|row| !kept.contains(&row.path)) { | |
| 160 | statements.push( | |
| 161 | self.db | |
| Fewer Artifacts reads: the store is asked for a handle only when needed, objects are kept in the isolate, and issue events read no workflows nobody listens for | 162 | .prepare("UPDATE workflows SET crons = '[]', error = ? WHERE id = ?") |
| 163 | .bind(&[GONE.into(), row.id.as_str().into()])?, | |
| GitHub Actions on g1t, part two: running workflows | 164 | ); |
| 165 | } | |
| 166 | statements.push( | |
| 167 | self.db | |
| 168 | .prepare("INSERT OR REPLACE INTO synced (repo_id, at) VALUES (?, ?)") | |
| 169 | .bind(&[repo.id.as_str().into(), now.into()])?, | |
| 170 | ); | |
| 171 | self.db.batch(statements).await?; | |
| 172 | Ok(()) | |
| 173 | } | |
| 174 | ||
| Fewer Artifacts reads: the store is asked for a handle only when needed, objects are kept in the isolate, and issue events read no workflows nobody listens for | 175 | /// Whether any of the repository's default-branch workflows could |
| 176 | /// start on `event`, from the synced table; None before the first sync. | |
| 177 | pub async fn listens(&self, repo_id: &str, event: &str) -> Result<Option<bool>> { | |
| 178 | #[derive(Deserialize)] | |
| 179 | struct Row { | |
| 180 | events: String, | |
| 181 | error: Option<String>, | |
| 182 | } | |
| 183 | if !self.synced(repo_id).await? { | |
| 184 | return Ok(None); | |
| 185 | } | |
| 186 | let rows = self | |
| 187 | .db | |
| 188 | .prepare("SELECT events, error FROM workflows WHERE repo_id = ? AND state = 'active'") | |
| 189 | .bind(&[repo_id.into()])? | |
| 190 | .all() | |
| 191 | .await? | |
| 192 | .results::<Row>()?; | |
| 193 | Ok(Some(rows.iter().any(|row| could_start(&row.events, row.error.as_deref(), event)))) | |
| 194 | } | |
| 195 | ||
| GitHub Actions on g1t, part two: running workflows | 196 | pub async fn synced(&self, repo_id: &str) -> Result<bool> { |
| 197 | Ok(self | |
| 198 | .db | |
| 199 | .prepare("SELECT COUNT(*) AS n FROM synced WHERE repo_id = ?") | |
| 200 | .bind(&[repo_id.into()])? | |
| 201 | .first::<Count>(None) | |
| 202 | .await? | |
| 203 | .is_some_and(|count| count.n > 0)) | |
| 204 | } | |
| 205 | ||
| 206 | /// The row for a workflow file, made if it is new (a file that exists | |
| 207 | /// only on a branch still gets its runs counted and listed). | |
| 208 | pub async fn workflow_row(&self, repo: &Repo, path: &str, name: &str, source: &str) -> Result<WorkflowRow> { | |
| 209 | let now = rfc3339(now_ms()); | |
| 210 | self.db | |
| 211 | .prepare( | |
| 212 | "INSERT INTO workflows (id, repo_id, repo, path, name, source, events, error, updated_at) | |
| 213 | VALUES (?, ?, ?, ?, ?, ?, '[]', 'Its file is not on the default branch.', ?) | |
| 214 | ON CONFLICT (repo_id, path) DO NOTHING", | |
| 215 | ) | |
| 216 | .bind(&[ | |
| 217 | new_id("wfl", now_ms()).into(), | |
| 218 | repo.id.as_str().into(), | |
| 219 | format!("{}/{}", repo.namespace, repo.name).into(), | |
| 220 | path.into(), | |
| 221 | name.into(), | |
| 222 | source.into(), | |
| 223 | now.into(), | |
| 224 | ])? | |
| 225 | .run() | |
| 226 | .await?; | |
| 227 | self.db | |
| 228 | .prepare("SELECT * FROM workflows WHERE repo_id = ? AND path = ?") | |
| 229 | .bind(&[repo.id.as_str().into(), path.into()])? | |
| 230 | .first::<WorkflowRow>(None) | |
| 231 | .await? | |
| 232 | .ok_or_else(|| worker::Error::RustError("the workflow was not recorded".into())) | |
| 233 | } | |
| 234 | } | |
| Fewer Artifacts reads: the store is asked for a handle only when needed, objects are kept in the isolate, and issue events read no workflows nobody listens for | 235 | |
| 236 | #[cfg(test)] | |
| 237 | mod tests { | |
| 238 | use super::*; | |
| 239 | ||
| 240 | #[test] | |
| 241 | fn a_synced_workflow_starts_on_the_events_it_lists_or_when_broken() { | |
| 242 | assert!(could_start(r#"["push","issues"]"#, None, "issues")); | |
| 243 | assert!(!could_start(r#"["push","pull_request"]"#, None, "issue_comment")); | |
| 244 | // A file that does not parse is read again, so its error shows. | |
| 245 | assert!(could_start("[]", Some("bad yaml"), "issues")); | |
| 246 | // A file gone from the branch starts nothing. | |
| 247 | assert!(!could_start("[]", Some(GONE), "issues")); | |
| 248 | } | |
| 249 | } |