Skip to content

g1t/services/actions/src/views.rs

300 lines12,611 bytesCodeBlameRaw

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: workflows, runs, a run's jobs, and a job's log.
2
3use g1t_actions::workflow::{self, Severity};
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look4use g1t_contracts::access::Capability;
Merge checks: statuses and check runs on every commit5use g1t_contracts::checks::{ActionsChecksArgs, MAX_COMMITS};
GitHub Actions on g1t, part two: running workflows6use 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,
Merge branch 'worktree-agent-a3abfcce648e87dca'55 environment: row.environment,
Fast pages, required checks on the branch, self-hosted runners, honest incidents56 self_hosted: row.labels.is_some(),
57 runner: row.runner_name,
GitHub Actions on g1t, part two: running workflows58 }
59}
60
61impl Actions {
62 async fn summary(&self, row: &WorkflowRow) -> Result<Workflow> {
63 let last_run = self
64 .db
65 .prepare("SELECT * FROM runs WHERE workflow_id = ? ORDER BY id DESC LIMIT 1")
66 .bind(&[row.id.as_str().into()])?
67 .first::<RunRow>(None)
68 .await?
69 .map(|run| run.summary());
70 let parsed = workflow::parse(&row.source).ok();
71 Ok(Workflow {
72 id: row.id.clone(),
73 path: row.path.clone(),
74 name: row.name.clone(),
75 events: serde_json::from_str(&row.events).unwrap_or_default(),
76 state: row.state.clone(),
77 error: row.error.clone(),
78 notes: notes(&row.source),
79 dispatch: parsed
80 .as_ref()
81 .filter(|_| row.error.is_none())
82 .and_then(|w| w.trigger("workflow_dispatch"))
83 .map(|t| serde_json::Value::Object(t.inputs.clone())),
84 last_run,
85 })
86 }
87
88 pub async fn workflows(&self, a: WorkflowsArgs) -> Result<Outcome<Vec<Workflow>>> {
89 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
90 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
91 };
92 if !self.synced(&repo.id).await?
93 && let Some(ws) = self.workspace_actor(&repo.namespace).await?
94 {
95 self.sync(&repo, &ws).await?;
96 }
One CI in the workflow list: hide files from before .github stopped being read, and gone files a live one replaced97 // A file that is gone stays listed while it has runs, unless one in
98 // the folder now has its name: then it would read as a second copy.
99 // Files outside the folder are from before g1t stopped reading
100 // `.github`; their runs stay under All workflows.
GitHub Actions on g1t, part two: running workflows101 let rows = self
102 .db
103 .prepare(
One CI in the workflow list: hide files from before .github stopped being read, and gone files a live one replaced104 "SELECT * FROM workflows w WHERE w.repo_id = ?1 AND w.path LIKE ?2
105 AND (w.error IS NULL OR w.error NOT LIKE 'Its file is%'
106 OR (w.id IN (SELECT workflow_id FROM runs WHERE repo_id = ?1)
107 AND NOT EXISTS (SELECT 1 FROM workflows o WHERE o.repo_id = ?1 AND o.id <> w.id
108 AND o.name = w.name AND (o.error IS NULL OR o.error NOT LIKE 'Its file is%'))))
109 ORDER BY w.name",
GitHub Actions on g1t, part two: running workflows110 )
One CI in the workflow list: hide files from before .github stopped being read, and gone files a live one replaced111 .bind(&[repo.id.as_str().into(), format!("{}/%", g1t_actions::workflow::FOLDER).into()])?
GitHub Actions on g1t, part two: running workflows112 .all()
113 .await?
114 .results::<WorkflowRow>()?;
115 let mut out = Vec::with_capacity(rows.len());
116 for row in &rows {
117 out.push(self.summary(row).await?);
118 }
119 Ok(Outcome::Ok(out))
120 }
121
122 pub async fn runs(&self, a: RunsArgs) -> Result<Outcome<Vec<WorkflowRun>>> {
123 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
124 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
125 };
126 let mut sql = "SELECT * FROM runs WHERE repo_id = ?".to_owned();
127 let mut binds: Vec<worker::wasm_bindgen::JsValue> = vec![repo.id.as_str().into()];
128 if let Some(workflow) = &a.workflow {
129 sql.push_str(" AND (workflow_id = ? OR path = ? OR path = ?)");
130 binds.push(workflow.as_str().into());
131 binds.push(workflow.as_str().into());
132 binds.push(format!("{}/{workflow}", g1t_actions::workflow::FOLDER).into());
133 }
134 if let Some(branch) = &a.branch {
135 sql.push_str(" AND (git_ref = ? OR json_extract(info, '$.headRef') = ?)");
136 binds.push(format!("refs/heads/{branch}").into());
137 binds.push(branch.as_str().into());
138 }
139 if let Some(event) = &a.event {
140 sql.push_str(" AND event = ?");
141 binds.push(event.as_str().into());
142 }
143 if let Some(pull) = a.pull {
144 sql.push_str(" AND pull = ?");
145 binds.push(pull.into());
146 }
147 if let Some(sha) = &a.sha {
148 sql.push_str(" AND sha = ?");
149 binds.push(sha.as_str().into());
150 }
151 sql.push_str(" ORDER BY id DESC LIMIT ?");
152 binds.push(a.limit.unwrap_or(RUNS_SHOWN).clamp(1, 100).into());
153 let rows = self.db.prepare(sql).bind(&binds)?.all().await?.results::<RunRow>()?;
154 Ok(Outcome::Ok(rows.iter().map(RunRow::summary).collect()))
155 }
156
157 pub async fn run(&self, a: RunArgs) -> Result<Outcome<RunDetail>> {
158 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
159 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
160 }
161 let run = check!(self.run_in(&a.repo, &a.id).await?);
Merge branch 'worktree-agent-a3abfcce648e87dca'162 let jobs: Vec<Job> = self.job_rows(&run.id).await?.into_iter().map(job_view).collect();
163 let pending_deployments = self.pending_for(&run, &a.viewer).await?;
164 let mut summary = run.summary();
165 // Jobs held at an environment's rules: the run waits, as on GitHub.
166 if matches!(summary.status.as_str(), "queued" | "in_progress") && jobs.iter().any(|job| job.status == "pending") {
167 summary.status = "waiting".to_owned();
168 }
GitHub Actions on g1t, part two: running workflows169 Ok(Outcome::Ok(RunDetail {
170 notes: notes(&run.source),
Merge branch 'worktree-agent-a3abfcce648e87dca'171 run: summary,
GitHub Actions on g1t, part two: running workflows172 jobs,
Merge branch 'worktree-agent-a3abfcce648e87dca'173 approval: run.approval(),
174 pending_deployments,
GitHub Actions on g1t, part two: running workflows175 }))
176 }
177
178 pub async fn logs(&self, a: LogsArgs) -> Result<Outcome<JobLog>> {
179 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
180 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
181 }
182 let job = self
183 .db
184 .prepare("SELECT jobs.* FROM jobs JOIN runs ON runs.id = jobs.run_id WHERE jobs.id = ? AND lower(runs.repo) = lower(?)")
185 .bind(&[a.job.as_str().into(), format!("{}/{}", a.repo.namespace, a.repo.name).into()])?
186 .first::<JobRow>(None)
187 .await?;
188 let Some(job) = job else {
189 return Ok(fail(FailureCode::NotFound, "No such job."));
190 };
191 #[derive(Deserialize)]
192 struct Row {
193 seq: u64,
194 step: u32,
195 text: String,
196 }
197 let chunks = self
198 .db
199 .prepare("SELECT seq, step, text FROM logs WHERE job_id = ? AND seq > ? ORDER BY seq LIMIT 500")
200 .bind(&[job.id.as_str().into(), (a.after as f64).into()])?
201 .all()
202 .await?
203 .results::<Row>()?
204 .into_iter()
205 .map(|row| LogChunk { seq: row.seq, step: row.step, text: row.text })
206 .collect();
207 Ok(Outcome::Ok(JobLog {
208 chunks,
209 done: job.status == "completed",
210 }))
211 }
212
213 pub async fn set_workflow_enabled(&self, a: SetWorkflowEnabledArgs) -> Result<Outcome<Workflow>> {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look214 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::ManageSettings).await? {
GitHub Actions on g1t, part two: running workflows215 return Ok(Outcome::Fail(refused));
216 }
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs217 let wanted = a.workflow.trim_start_matches(".g1t/workflows/");
GitHub Actions on g1t, part two: running workflows218 let row = self
219 .db
220 .prepare("SELECT * FROM workflows WHERE lower(repo) = lower(?) AND (id = ? OR path = ?)")
221 .bind(&[
222 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
223 a.workflow.as_str().into(),
224 format!("{}/{wanted}", g1t_actions::workflow::FOLDER).into(),
225 ])?
226 .first::<WorkflowRow>(None)
227 .await?;
228 let Some(row) = row else {
229 return Ok(fail(FailureCode::NotFound, "No such workflow."));
230 };
231 let state = if a.enabled { "active" } else { "disabled" };
232 self.db.prepare("UPDATE workflows SET state = ? WHERE id = ?").bind(&[state.into(), row.id.as_str().into()])?.run().await?;
233 let row = WorkflowRow { state: state.to_owned(), ..row };
234 Ok(Outcome::Ok(self.summary(&row).await?))
235 }
Merge checks: statuses and check runs on every commit236
237 /// `check_runs`: workflow runs with their jobs, by commits, by a run's
238 /// id or by one of its jobs' ids, newest first, for the work service to
239 /// show as check suites and check runs. It decides who may see them.
240 pub async fn check_runs(&self, a: ActionsChecksArgs) -> Result<Vec<RunDetail>> {
241 let runs: Vec<RunRow> = if let Some(job) = &a.job_id {
242 self.db
243 .prepare("SELECT runs.* FROM runs JOIN jobs ON jobs.run_id = runs.id WHERE jobs.id = ? AND runs.repo_id = ?")
244 .bind(&[job.as_str().into(), a.repo_id.as_str().into()])?
245 .all()
246 .await?
247 .results::<RunRow>()?
248 } else if let Some(run) = &a.run_id {
249 self.db
250 .prepare("SELECT * FROM runs WHERE id = ? AND repo_id = ?")
251 .bind(&[run.as_str().into(), a.repo_id.as_str().into()])?
252 .all()
253 .await?
254 .results::<RunRow>()?
255 } else if a.shas.is_empty() {
256 Vec::new()
257 } else {
258 let shas: Vec<&String> = a.shas.iter().take(MAX_COMMITS).collect();
259 let marks = vec!["?"; shas.len()].join(", ");
260 let mut binds: Vec<worker::wasm_bindgen::JsValue> = vec![a.repo_id.as_str().into()];
261 binds.extend(shas.iter().map(|sha| sha.as_str().into()));
262 binds.push(CHECK_RUNS_LIMIT.into());
263 self.db
264 .prepare(format!("SELECT * FROM runs WHERE repo_id = ? AND sha IN ({marks}) ORDER BY id DESC LIMIT ?"))
265 .bind(&binds)?
266 .all()
267 .await?
268 .results::<RunRow>()?
269 };
270 if runs.is_empty() {
271 return Ok(Vec::new());
272 }
273 let marks = vec!["?"; runs.len()].join(", ");
274 let ids: Vec<worker::wasm_bindgen::JsValue> = runs.iter().map(|run| run.id.as_str().into()).collect();
275 let jobs = self
276 .db
277 .prepare(format!("SELECT * FROM jobs WHERE run_id IN ({marks}) ORDER BY rowid"))
278 .bind(&ids)?
279 .all()
280 .await?
281 .results::<JobRow>()?;
282 let mut by_run: std::collections::HashMap<String, Vec<Job>> = std::collections::HashMap::new();
283 for job in jobs {
284 by_run.entry(job.run_id.clone()).or_default().push(job_view(job));
285 }
286 Ok(runs
287 .iter()
Merge branch 'worktree-agent-a3abfcce648e87dca'288 .map(|run| RunDetail {
289 run: run.summary(),
290 jobs: by_run.remove(&run.id).unwrap_or_default(),
291 notes: Vec::new(),
292 approval: run.approval(),
293 pending_deployments: Vec::new(),
294 })
Merge checks: statuses and check runs on every commit295 .collect())
296 }
GitHub Actions on g1t, part two: running workflows297}
Merge checks: statuses and check runs on every commit298
299/// The most runs `check_runs` reads for a set of commits.
300const CHECK_RUNS_LIMIT: u32 = 200;

This file's history is long; its oldest lines are credited to the oldest commit read.