Skip to content
272 linesCodeBlameRaw
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
5use g1t_actions::workflow::{self, FOLDER};
6use g1t_contracts::new_id;
7use g1t_contracts::repos::{BlobArgs, BlobView, EntryKind, Repo, RepoPath, TreeArgs, TreeView};
8use g1t_contracts::time::rfc3339;
9use g1t_contracts::{Outcome, User};
10use g1t_kit::now_ms;
11use serde::Deserialize;
12use worker::Result;
13
14use crate::{Actions, Count, MAX_WORKFLOWS, optional};
15
16/// One workflow file as read at a commit.
17pub 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.
23pub struct Read {
24 pub files: Vec<WorkflowFile>,
25 pub head: Option<String>,
26}
27
28#[derive(Deserialize)]
29pub 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/// What a workflow's row says once its file has left the default branch.
44pub 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.
49pub 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
56impl 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 self.read_folders(path, actor, git_ref, &[FOLDER]).await
61 }
62
63 /// The workflow files a repository runs: `.g1t/workflows`, and on a
64 /// mirror in CI failover or taken over, `.github/workflows` too (see
65 /// mirrored.rs), where g1t's own stands in for one of the same name.
66 pub async fn read_for(&self, repo: &Repo, actor: &User, git_ref: Option<&str>, github: bool) -> Result<Read> {
67 let path = RepoPath {
68 namespace: repo.namespace.clone(),
69 name: repo.name.clone(),
70 };
71 if !github {
72 return self.read_workflows(&path, actor, git_ref).await;
73 }
74 let read = self.read_folders(&path, actor, git_ref, &[FOLDER, crate::mirrored::GITHUB_FOLDER]).await?;
75 Ok(Read {
76 files: crate::mirrored::prefer_g1t(read.files),
77 head: read.head,
78 })
79 }
80
81 async fn read_folders(&self, path: &RepoPath, actor: &User, git_ref: Option<&str>, folders: &[&str]) -> Result<Read> {
82 let viewer = Some(actor.clone());
83 let mut files = Vec::new();
84 let mut found_head = None;
85 for folder in folders {
86 let tree: Outcome<TreeView> = g1t_kit::call(
87 &self.repos,
88 "tree",
89 &TreeArgs {
90 path: path.clone(),
91 viewer: viewer.clone(),
92 git_ref: git_ref.map(str::to_owned),
93 tree_path: (*folder).to_owned(),
94 },
95 )
96 .await?;
97 let (entries, head, resolved) = match tree {
98 Outcome::Ok(tree) => (tree.entries, tree.head.map(|commit| commit.hash), tree.git_ref),
99 // No folder: no workflows from it.
100 Outcome::Fail(_) => continue,
101 };
102 let at = head.clone().unwrap_or(resolved);
103 found_head = found_head.or(head);
104 for entry in entries
105 .into_iter()
106 .filter(|entry| matches!(entry.kind, EntryKind::Blob | EntryKind::Exec))
107 .filter(|entry| entry.name.ends_with(".yml") || entry.name.ends_with(".yaml"))
108 .take(MAX_WORKFLOWS.saturating_sub(files.len()))
109 {
110 let file_path = format!("{folder}/{}", entry.name);
111 let blob: Outcome<BlobView> = g1t_kit::call(
112 &self.repos,
113 "blob",
114 &BlobArgs {
115 path: path.clone(),
116 viewer: viewer.clone(),
117 git_ref: at.clone(),
118 file_path: file_path.clone(),
119 },
120 )
121 .await?;
122 if let Outcome::Ok(BlobView { text: Some(source), .. }) = blob {
123 files.push(WorkflowFile { path: file_path, source });
124 }
125 }
126 }
127 Ok(Read { files, head: found_head })
128 }
129
130 /// Keeps the `workflows` table in step with the default branch.
131 pub async fn sync(&self, repo: &Repo, actor: &User) -> Result<()> {
132 let github = crate::mirrored::policy(repo.mirror.as_ref(), false).github;
133 let read = self.read_for(repo, actor, None, github).await?;
134 let full_name = format!("{}/{}", repo.namespace, repo.name);
135 let now = rfc3339(now_ms());
136 let mut statements = Vec::new();
137 for file in &read.files {
138 let parsed = workflow::parse(&file.source);
139 let (name, error, events, crons) = match &parsed {
140 Ok(parsed) => (
141 parsed.display_name(&file.path),
142 None,
143 parsed.triggers.iter().map(|t| t.event.clone()).collect::<Vec<_>>(),
144 parsed.trigger("schedule").map(|t| t.crons.clone()).unwrap_or_default(),
145 ),
146 Err(problem) => (file.path.clone(), Some(problem.clone()), Vec::new(), Vec::new()),
147 };
148 statements.push(
149 self.db
150 .prepare(
151 "INSERT INTO workflows (id, repo_id, repo, path, name, source, events, crons, error, updated_at)
152 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
153 ON CONFLICT (repo_id, path) DO UPDATE SET
154 repo = excluded.repo, name = excluded.name, source = excluded.source,
155 events = excluded.events, crons = excluded.crons, error = excluded.error,
156 updated_at = excluded.updated_at",
157 )
158 .bind(&[
159 new_id("wfl", now_ms()).into(),
160 repo.id.as_str().into(),
161 full_name.as_str().into(),
162 file.path.as_str().into(),
163 name.into(),
164 file.source.as_str().into(),
165 serde_json::to_string(&events)?.into(),
166 serde_json::to_string(&crons)?.into(),
167 optional(error.as_deref()),
168 now.as_str().into(),
169 ])?,
170 );
171 }
172 // A workflow whose file is gone keeps its runs, but no longer runs
173 // on schedule or by hand: it is listed only while it has runs.
174 let kept: Vec<String> = read.files.iter().map(|file| file.path.clone()).collect();
175 let existing = self
176 .db
177 .prepare("SELECT * FROM workflows WHERE repo_id = ?")
178 .bind(&[repo.id.as_str().into()])?
179 .all()
180 .await?
181 .results::<WorkflowRow>()?;
182 for row in existing.iter().filter(|row| !kept.contains(&row.path)) {
183 statements.push(
184 self.db
185 .prepare("UPDATE workflows SET crons = '[]', error = ? WHERE id = ?")
186 .bind(&[GONE.into(), row.id.as_str().into()])?,
187 );
188 }
189 statements.push(
190 self.db
191 .prepare("INSERT OR REPLACE INTO synced (repo_id, at) VALUES (?, ?)")
192 .bind(&[repo.id.as_str().into(), now.into()])?,
193 );
194 self.db.batch(statements).await?;
195 Ok(())
196 }
197
198 /// Whether any of the repository's default-branch workflows could
199 /// start on `event`, from the synced table; None before the first sync.
200 pub async fn listens(&self, repo_id: &str, event: &str) -> Result<Option<bool>> {
201 #[derive(Deserialize)]
202 struct Row {
203 events: String,
204 error: Option<String>,
205 }
206 if !self.synced(repo_id).await? {
207 return Ok(None);
208 }
209 let rows = self
210 .db
211 .prepare("SELECT events, error FROM workflows WHERE repo_id = ? AND state = 'active'")
212 .bind(&[repo_id.into()])?
213 .all()
214 .await?
215 .results::<Row>()?;
216 Ok(Some(rows.iter().any(|row| could_start(&row.events, row.error.as_deref(), event))))
217 }
218
219 pub async fn synced(&self, repo_id: &str) -> Result<bool> {
220 Ok(self
221 .db
222 .prepare("SELECT COUNT(*) AS n FROM synced WHERE repo_id = ?")
223 .bind(&[repo_id.into()])?
224 .first::<Count>(None)
225 .await?
226 .is_some_and(|count| count.n > 0))
227 }
228
229 /// The row for a workflow file, made if it is new (a file that exists
230 /// only on a branch still gets its runs counted and listed).
231 pub async fn workflow_row(&self, repo: &Repo, path: &str, name: &str, source: &str) -> Result<WorkflowRow> {
232 let now = rfc3339(now_ms());
233 self.db
234 .prepare(
235 "INSERT INTO workflows (id, repo_id, repo, path, name, source, events, error, updated_at)
236 VALUES (?, ?, ?, ?, ?, ?, '[]', 'Its file is not on the default branch.', ?)
237 ON CONFLICT (repo_id, path) DO NOTHING",
238 )
239 .bind(&[
240 new_id("wfl", now_ms()).into(),
241 repo.id.as_str().into(),
242 format!("{}/{}", repo.namespace, repo.name).into(),
243 path.into(),
244 name.into(),
245 source.into(),
246 now.into(),
247 ])?
248 .run()
249 .await?;
250 self.db
251 .prepare("SELECT * FROM workflows WHERE repo_id = ? AND path = ?")
252 .bind(&[repo.id.as_str().into(), path.into()])?
253 .first::<WorkflowRow>(None)
254 .await?
255 .ok_or_else(|| worker::Error::RustError("the workflow was not recorded".into()))
256 }
257}
258
259#[cfg(test)]
260mod tests {
261 use super::*;
262
263 #[test]
264 fn a_synced_workflow_starts_on_the_events_it_lists_or_when_broken() {
265 assert!(could_start(r#"["push","issues"]"#, None, "issues"));
266 assert!(!could_start(r#"["push","pull_request"]"#, None, "issue_comment"));
267 // A file that does not parse is read again, so its error shows.
268 assert!(could_start("[]", Some("bad yaml"), "issues"));
269 // A file gone from the branch starts nothing.
270 assert!(!could_start("[]", Some(GONE), "issues"));
271 }
272}