Skip to content

g1t/services/actions/src/views.rs

285 lines11,968 bytesCodeBlame
1//! Reading: workflows, runs, a run's jobs, and a job's log.
2
3use g1t_actions::workflow::{self, Severity};
4use g1t_contracts::access::Capability;
5use g1t_contracts::checks::{ActionsChecksArgs, MAX_COMMITS};
6use g1t_contracts::actions::{
7 Annotation, Job, JobLog, LogChunk, LogsArgs, RunArgs, RunDetail, RunsArgs, SetWorkflowEnabledArgs, StepState, Workflow, WorkflowNote,
8 WorkflowRun, WorkflowsArgs,
9};
10use g1t_contracts::{FailureCode, Outcome};
11use serde::Deserialize;
12use worker::Result;
13
14use crate::plan::{JobRow, RunRow};
15use crate::sync::WorkflowRow;
16use crate::{Actions, check, fail};
17
18const RUNS_SHOWN: u32 = 50;
19
20fn notes(source: &str) -> Vec<WorkflowNote> {
21 workflow::parse(source)
22 .map(|w| {
23 w.notes
24 .into_iter()
25 .map(|note| WorkflowNote {
26 severity: match note.severity {
27 Severity::Info => "info",
28 Severity::Warning => "warning",
29 Severity::Unsupported => "unsupported",
30 }
31 .to_owned(),
32 job: note.job,
33 message: note.message,
34 })
35 .collect()
36 })
37 .unwrap_or_default()
38}
39
40fn job_view(row: JobRow) -> Job {
41 let needs = row.needs();
42 Job {
43 id: row.id,
44 run_id: row.run_id,
45 key: row.key,
46 name: row.name,
47 needs,
48 status: row.status,
49 conclusion: row.conclusion,
50 steps: serde_json::from_str::<Vec<StepState>>(&row.steps).unwrap_or_default(),
51 annotations: serde_json::from_str::<Vec<Annotation>>(&row.annotations).unwrap_or_default(),
52 reason: row.reason,
53 started_at: row.started_at,
54 finished_at: row.finished_at,
55 self_hosted: row.labels.is_some(),
56 runner: row.runner_name,
57 }
58}
59
60impl Actions {
61 async fn summary(&self, row: &WorkflowRow) -> Result<Workflow> {
62 let last_run = self
63 .db
64 .prepare("SELECT * FROM runs WHERE workflow_id = ? ORDER BY id DESC LIMIT 1")
65 .bind(&[row.id.as_str().into()])?
66 .first::<RunRow>(None)
67 .await?
68 .map(|run| run.summary());
69 let parsed = workflow::parse(&row.source).ok();
70 Ok(Workflow {
71 id: row.id.clone(),
72 path: row.path.clone(),
73 name: row.name.clone(),
74 events: serde_json::from_str(&row.events).unwrap_or_default(),
75 state: row.state.clone(),
76 error: row.error.clone(),
77 notes: notes(&row.source),
78 dispatch: parsed
79 .as_ref()
80 .filter(|_| row.error.is_none())
81 .and_then(|w| w.trigger("workflow_dispatch"))
82 .map(|t| serde_json::Value::Object(t.inputs.clone())),
83 last_run,
84 })
85 }
86
87 pub async fn workflows(&self, a: WorkflowsArgs) -> Result<Outcome<Vec<Workflow>>> {
88 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
89 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
90 };
91 if !self.synced(&repo.id).await?
92 && let Some(ws) = self.workspace_actor(&repo.namespace).await?
93 {
94 self.sync(&repo, &ws).await?;
95 }
96 // A file that is gone stays listed while it has runs, unless one in
97 // the folder now has its name: then it would read as a second copy.
98 // Files outside the folder are from before g1t stopped reading
99 // `.github`; their runs stay under All workflows.
100 let rows = self
101 .db
102 .prepare(
103 "SELECT * FROM workflows w WHERE w.repo_id = ?1 AND w.path LIKE ?2
104 AND (w.error IS NULL OR w.error NOT LIKE 'Its file is%'
105 OR (w.id IN (SELECT workflow_id FROM runs WHERE repo_id = ?1)
106 AND NOT EXISTS (SELECT 1 FROM workflows o WHERE o.repo_id = ?1 AND o.id <> w.id
107 AND o.name = w.name AND (o.error IS NULL OR o.error NOT LIKE 'Its file is%'))))
108 ORDER BY w.name",
109 )
110 .bind(&[repo.id.as_str().into(), format!("{}/%", g1t_actions::workflow::FOLDER).into()])?
111 .all()
112 .await?
113 .results::<WorkflowRow>()?;
114 let mut out = Vec::with_capacity(rows.len());
115 for row in &rows {
116 out.push(self.summary(row).await?);
117 }
118 Ok(Outcome::Ok(out))
119 }
120
121 pub async fn runs(&self, a: RunsArgs) -> Result<Outcome<Vec<WorkflowRun>>> {
122 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
123 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
124 };
125 let mut sql = "SELECT * FROM runs WHERE repo_id = ?".to_owned();
126 let mut binds: Vec<worker::wasm_bindgen::JsValue> = vec![repo.id.as_str().into()];
127 if let Some(workflow) = &a.workflow {
128 sql.push_str(" AND (workflow_id = ? OR path = ? OR path = ?)");
129 binds.push(workflow.as_str().into());
130 binds.push(workflow.as_str().into());
131 binds.push(format!("{}/{workflow}", g1t_actions::workflow::FOLDER).into());
132 }
133 if let Some(branch) = &a.branch {
134 sql.push_str(" AND (git_ref = ? OR json_extract(info, '$.headRef') = ?)");
135 binds.push(format!("refs/heads/{branch}").into());
136 binds.push(branch.as_str().into());
137 }
138 if let Some(event) = &a.event {
139 sql.push_str(" AND event = ?");
140 binds.push(event.as_str().into());
141 }
142 if let Some(pull) = a.pull {
143 sql.push_str(" AND pull = ?");
144 binds.push(pull.into());
145 }
146 if let Some(sha) = &a.sha {
147 sql.push_str(" AND sha = ?");
148 binds.push(sha.as_str().into());
149 }
150 sql.push_str(" ORDER BY id DESC LIMIT ?");
151 binds.push(a.limit.unwrap_or(RUNS_SHOWN).clamp(1, 100).into());
152 let rows = self.db.prepare(sql).bind(&binds)?.all().await?.results::<RunRow>()?;
153 Ok(Outcome::Ok(rows.iter().map(RunRow::summary).collect()))
154 }
155
156 pub async fn run(&self, a: RunArgs) -> Result<Outcome<RunDetail>> {
157 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
158 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
159 }
160 let run = check!(self.run_in(&a.repo, &a.id).await?);
161 let jobs = self.job_rows(&run.id).await?.into_iter().map(job_view).collect();
162 Ok(Outcome::Ok(RunDetail {
163 notes: notes(&run.source),
164 run: run.summary(),
165 jobs,
166 }))
167 }
168
169 pub async fn logs(&self, a: LogsArgs) -> Result<Outcome<JobLog>> {
170 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
171 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
172 }
173 let job = self
174 .db
175 .prepare("SELECT jobs.* FROM jobs JOIN runs ON runs.id = jobs.run_id WHERE jobs.id = ? AND lower(runs.repo) = lower(?)")
176 .bind(&[a.job.as_str().into(), format!("{}/{}", a.repo.namespace, a.repo.name).into()])?
177 .first::<JobRow>(None)
178 .await?;
179 let Some(job) = job else {
180 return Ok(fail(FailureCode::NotFound, "No such job."));
181 };
182 #[derive(Deserialize)]
183 struct Row {
184 seq: u64,
185 step: u32,
186 text: String,
187 }
188 let chunks = self
189 .db
190 .prepare("SELECT seq, step, text FROM logs WHERE job_id = ? AND seq > ? ORDER BY seq LIMIT 500")
191 .bind(&[job.id.as_str().into(), (a.after as f64).into()])?
192 .all()
193 .await?
194 .results::<Row>()?
195 .into_iter()
196 .map(|row| LogChunk { seq: row.seq, step: row.step, text: row.text })
197 .collect();
198 Ok(Outcome::Ok(JobLog {
199 chunks,
200 done: job.status == "completed",
201 }))
202 }
203
204 pub async fn set_workflow_enabled(&self, a: SetWorkflowEnabledArgs) -> Result<Outcome<Workflow>> {
205 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::ManageSettings).await? {
206 return Ok(Outcome::Fail(refused));
207 }
208 let wanted = a.workflow.trim_start_matches(".g1t/workflows/");
209 let row = self
210 .db
211 .prepare("SELECT * FROM workflows WHERE lower(repo) = lower(?) AND (id = ? OR path = ?)")
212 .bind(&[
213 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
214 a.workflow.as_str().into(),
215 format!("{}/{wanted}", g1t_actions::workflow::FOLDER).into(),
216 ])?
217 .first::<WorkflowRow>(None)
218 .await?;
219 let Some(row) = row else {
220 return Ok(fail(FailureCode::NotFound, "No such workflow."));
221 };
222 let state = if a.enabled { "active" } else { "disabled" };
223 self.db.prepare("UPDATE workflows SET state = ? WHERE id = ?").bind(&[state.into(), row.id.as_str().into()])?.run().await?;
224 let row = WorkflowRow { state: state.to_owned(), ..row };
225 Ok(Outcome::Ok(self.summary(&row).await?))
226 }
227
228 /// `check_runs`: workflow runs with their jobs, by commits, by a run's
229 /// id or by one of its jobs' ids, newest first, for the work service to
230 /// show as check suites and check runs. It decides who may see them.
231 pub async fn check_runs(&self, a: ActionsChecksArgs) -> Result<Vec<RunDetail>> {
232 let runs: Vec<RunRow> = if let Some(job) = &a.job_id {
233 self.db
234 .prepare("SELECT runs.* FROM runs JOIN jobs ON jobs.run_id = runs.id WHERE jobs.id = ? AND runs.repo_id = ?")
235 .bind(&[job.as_str().into(), a.repo_id.as_str().into()])?
236 .all()
237 .await?
238 .results::<RunRow>()?
239 } else if let Some(run) = &a.run_id {
240 self.db
241 .prepare("SELECT * FROM runs WHERE id = ? AND repo_id = ?")
242 .bind(&[run.as_str().into(), a.repo_id.as_str().into()])?
243 .all()
244 .await?
245 .results::<RunRow>()?
246 } else if a.shas.is_empty() {
247 Vec::new()
248 } else {
249 let shas: Vec<&String> = a.shas.iter().take(MAX_COMMITS).collect();
250 let marks = vec!["?"; shas.len()].join(", ");
251 let mut binds: Vec<worker::wasm_bindgen::JsValue> = vec![a.repo_id.as_str().into()];
252 binds.extend(shas.iter().map(|sha| sha.as_str().into()));
253 binds.push(CHECK_RUNS_LIMIT.into());
254 self.db
255 .prepare(format!("SELECT * FROM runs WHERE repo_id = ? AND sha IN ({marks}) ORDER BY id DESC LIMIT ?"))
256 .bind(&binds)?
257 .all()
258 .await?
259 .results::<RunRow>()?
260 };
261 if runs.is_empty() {
262 return Ok(Vec::new());
263 }
264 let marks = vec!["?"; runs.len()].join(", ");
265 let ids: Vec<worker::wasm_bindgen::JsValue> = runs.iter().map(|run| run.id.as_str().into()).collect();
266 let jobs = self
267 .db
268 .prepare(format!("SELECT * FROM jobs WHERE run_id IN ({marks}) ORDER BY rowid"))
269 .bind(&ids)?
270 .all()
271 .await?
272 .results::<JobRow>()?;
273 let mut by_run: std::collections::HashMap<String, Vec<Job>> = std::collections::HashMap::new();
274 for job in jobs {
275 by_run.entry(job.run_id.clone()).or_default().push(job_view(job));
276 }
277 Ok(runs
278 .iter()
279 .map(|run| RunDetail { run: run.summary(), jobs: by_run.remove(&run.id).unwrap_or_default(), notes: Vec::new() })
280 .collect())
281 }
282}
283
284/// The most runs `check_runs` reads for a set of commits.
285const CHECK_RUNS_LIMIT: u32 = 200;