g1t/services/actions/src/sync.rs
| 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 | |
| 43 | impl Actions { |
| 44 | /// The workflow files of `path` as of `git_ref` (the default branch |
| 45 | /// when absent). |
| 46 | pub async fn read_workflows(&self, path: &RepoPath, actor: &User, git_ref: Option<&str>) -> Result<Read> { |
| 47 | let viewer = Some(actor.clone()); |
| 48 | let tree: Outcome<TreeView> = g1t_kit::call( |
| 49 | &self.repos, |
| 50 | "tree", |
| 51 | &TreeArgs { |
| 52 | path: path.clone(), |
| 53 | viewer: viewer.clone(), |
| 54 | git_ref: git_ref.map(str::to_owned), |
| 55 | tree_path: FOLDER.to_owned(), |
| 56 | }, |
| 57 | ) |
| 58 | .await?; |
| 59 | let (entries, head, resolved) = match tree { |
| 60 | Outcome::Ok(tree) => (tree.entries, tree.head.map(|commit| commit.hash), tree.git_ref), |
| 61 | // No folder: no workflows. |
| 62 | Outcome::Fail(_) => return Ok(Read { files: Vec::new(), head: None }), |
| 63 | }; |
| 64 | let at = head.clone().unwrap_or(resolved); |
| 65 | let mut files = Vec::new(); |
| 66 | for entry in entries |
| 67 | .into_iter() |
| 68 | .filter(|entry| matches!(entry.kind, EntryKind::Blob | EntryKind::Exec)) |
| 69 | .filter(|entry| entry.name.ends_with(".yml") || entry.name.ends_with(".yaml")) |
| 70 | .take(MAX_WORKFLOWS) |
| 71 | { |
| 72 | let file_path = format!("{FOLDER}/{}", entry.name); |
| 73 | let blob: Outcome<BlobView> = g1t_kit::call( |
| 74 | &self.repos, |
| 75 | "blob", |
| 76 | &BlobArgs { |
| 77 | path: path.clone(), |
| 78 | viewer: viewer.clone(), |
| 79 | git_ref: at.clone(), |
| 80 | file_path: file_path.clone(), |
| 81 | }, |
| 82 | ) |
| 83 | .await?; |
| 84 | if let Outcome::Ok(BlobView { text: Some(source), .. }) = blob { |
| 85 | files.push(WorkflowFile { path: file_path, source }); |
| 86 | } |
| 87 | } |
| 88 | Ok(Read { files, head }) |
| 89 | } |
| 90 | |
| 91 | /// Keeps the `workflows` table in step with the default branch. |
| 92 | pub async fn sync(&self, repo: &Repo, actor: &User) -> Result<()> { |
| 93 | let path = RepoPath { |
| 94 | namespace: repo.namespace.clone(), |
| 95 | name: repo.name.clone(), |
| 96 | }; |
| 97 | let read = self.read_workflows(&path, actor, None).await?; |
| 98 | let full_name = format!("{}/{}", repo.namespace, repo.name); |
| 99 | let now = rfc3339(now_ms()); |
| 100 | let mut statements = Vec::new(); |
| 101 | for file in &read.files { |
| 102 | let parsed = workflow::parse(&file.source); |
| 103 | let (name, error, events, crons) = match &parsed { |
| 104 | Ok(parsed) => ( |
| 105 | parsed.display_name(&file.path), |
| 106 | None, |
| 107 | parsed.triggers.iter().map(|t| t.event.clone()).collect::<Vec<_>>(), |
| 108 | parsed.trigger("schedule").map(|t| t.crons.clone()).unwrap_or_default(), |
| 109 | ), |
| 110 | Err(problem) => (file.path.clone(), Some(problem.clone()), Vec::new(), Vec::new()), |
| 111 | }; |
| 112 | statements.push( |
| 113 | self.db |
| 114 | .prepare( |
| 115 | "INSERT INTO workflows (id, repo_id, repo, path, name, source, events, crons, error, updated_at) |
| 116 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) |
| 117 | ON CONFLICT (repo_id, path) DO UPDATE SET |
| 118 | repo = excluded.repo, name = excluded.name, source = excluded.source, |
| 119 | events = excluded.events, crons = excluded.crons, error = excluded.error, |
| 120 | updated_at = excluded.updated_at", |
| 121 | ) |
| 122 | .bind(&[ |
| 123 | new_id("wfl", now_ms()).into(), |
| 124 | repo.id.as_str().into(), |
| 125 | full_name.as_str().into(), |
| 126 | file.path.as_str().into(), |
| 127 | name.into(), |
| 128 | file.source.as_str().into(), |
| 129 | serde_json::to_string(&events)?.into(), |
| 130 | serde_json::to_string(&crons)?.into(), |
| 131 | optional(error.as_deref()), |
| 132 | now.as_str().into(), |
| 133 | ])?, |
| 134 | ); |
| 135 | } |
| 136 | // A workflow whose file is gone keeps its runs, but no longer runs |
| 137 | // on schedule or by hand: it is listed only while it has runs. |
| 138 | let kept: Vec<String> = read.files.iter().map(|file| file.path.clone()).collect(); |
| 139 | let existing = self |
| 140 | .db |
| 141 | .prepare("SELECT * FROM workflows WHERE repo_id = ?") |
| 142 | .bind(&[repo.id.as_str().into()])? |
| 143 | .all() |
| 144 | .await? |
| 145 | .results::<WorkflowRow>()?; |
| 146 | for row in existing.iter().filter(|row| !kept.contains(&row.path)) { |
| 147 | statements.push( |
| 148 | self.db |
| 149 | .prepare("UPDATE workflows SET crons = '[]', error = 'Its file is no longer on the default branch.' WHERE id = ?") |
| 150 | .bind(&[row.id.as_str().into()])?, |
| 151 | ); |
| 152 | } |
| 153 | statements.push( |
| 154 | self.db |
| 155 | .prepare("INSERT OR REPLACE INTO synced (repo_id, at) VALUES (?, ?)") |
| 156 | .bind(&[repo.id.as_str().into(), now.into()])?, |
| 157 | ); |
| 158 | self.db.batch(statements).await?; |
| 159 | Ok(()) |
| 160 | } |
| 161 | |
| 162 | pub async fn synced(&self, repo_id: &str) -> Result<bool> { |
| 163 | Ok(self |
| 164 | .db |
| 165 | .prepare("SELECT COUNT(*) AS n FROM synced WHERE repo_id = ?") |
| 166 | .bind(&[repo_id.into()])? |
| 167 | .first::<Count>(None) |
| 168 | .await? |
| 169 | .is_some_and(|count| count.n > 0)) |
| 170 | } |
| 171 | |
| 172 | /// The row for a workflow file, made if it is new (a file that exists |
| 173 | /// only on a branch still gets its runs counted and listed). |
| 174 | pub async fn workflow_row(&self, repo: &Repo, path: &str, name: &str, source: &str) -> Result<WorkflowRow> { |
| 175 | let now = rfc3339(now_ms()); |
| 176 | self.db |
| 177 | .prepare( |
| 178 | "INSERT INTO workflows (id, repo_id, repo, path, name, source, events, error, updated_at) |
| 179 | VALUES (?, ?, ?, ?, ?, ?, '[]', 'Its file is not on the default branch.', ?) |
| 180 | ON CONFLICT (repo_id, path) DO NOTHING", |
| 181 | ) |
| 182 | .bind(&[ |
| 183 | new_id("wfl", now_ms()).into(), |
| 184 | repo.id.as_str().into(), |
| 185 | format!("{}/{}", repo.namespace, repo.name).into(), |
| 186 | path.into(), |
| 187 | name.into(), |
| 188 | source.into(), |
| 189 | now.into(), |
| 190 | ])? |
| 191 | .run() |
| 192 | .await?; |
| 193 | self.db |
| 194 | .prepare("SELECT * FROM workflows WHERE repo_id = ? AND path = ?") |
| 195 | .bind(&[repo.id.as_str().into(), path.into()])? |
| 196 | .first::<WorkflowRow>(None) |
| 197 | .await? |
| 198 | .ok_or_else(|| worker::Error::RustError("the workflow was not recorded".into())) |
| 199 | } |
| 200 | } |