flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/actions/src/sync.rs

249 lines9,843 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.

GitHub Actions on g1t, part two: running workflows1//! 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
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 for43/// 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
GitHub Actions on g1t, part two: running workflows56impl 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 for162 .prepare("UPDATE workflows SET crons = '[]', error = ? WHERE id = ?")
163 .bind(&[GONE.into(), row.id.as_str().into()])?,
GitHub Actions on g1t, part two: running workflows164 );
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 for175 /// 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 workflows196 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 for235
236#[cfg(test)]
237mod 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}