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

200 lines7,973 bytesCodeBlame
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
43impl 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}