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/views.rs

219 lines8,645 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: 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;
GitHub Actions on g1t, part two: running workflows5use g1t_contracts::actions::{
6 Annotation, Job, JobLog, LogChunk, LogsArgs, RunArgs, RunDetail, RunsArgs, SetWorkflowEnabledArgs, StepState, Workflow, WorkflowNote,
7 WorkflowRun, WorkflowsArgs,
8};
9use g1t_contracts::{FailureCode, Outcome};
10use serde::Deserialize;
11use worker::Result;
12
13use crate::plan::{JobRow, RunRow};
14use crate::sync::WorkflowRow;
15use crate::{Actions, check, fail};
16
17const RUNS_SHOWN: u32 = 50;
18
19fn notes(source: &str) -> Vec<WorkflowNote> {
20 workflow::parse(source)
21 .map(|w| {
22 w.notes
23 .into_iter()
24 .map(|note| WorkflowNote {
25 severity: match note.severity {
26 Severity::Info => "info",
27 Severity::Warning => "warning",
28 Severity::Unsupported => "unsupported",
29 }
30 .to_owned(),
31 job: note.job,
32 message: note.message,
33 })
34 .collect()
35 })
36 .unwrap_or_default()
37}
38
39fn job_view(row: JobRow) -> Job {
40 let needs = row.needs();
41 Job {
42 id: row.id,
43 run_id: row.run_id,
44 key: row.key,
45 name: row.name,
46 needs,
47 status: row.status,
48 conclusion: row.conclusion,
49 steps: serde_json::from_str::<Vec<StepState>>(&row.steps).unwrap_or_default(),
50 annotations: serde_json::from_str::<Vec<Annotation>>(&row.annotations).unwrap_or_default(),
51 reason: row.reason,
52 started_at: row.started_at,
53 finished_at: row.finished_at,
Fast pages, required checks on the branch, self-hosted runners, honest incidents54 self_hosted: row.labels.is_some(),
55 runner: row.runner_name,
GitHub Actions on g1t, part two: running workflows56 }
57}
58
59impl Actions {
60 async fn summary(&self, row: &WorkflowRow) -> Result<Workflow> {
61 let last_run = self
62 .db
63 .prepare("SELECT * FROM runs WHERE workflow_id = ? ORDER BY id DESC LIMIT 1")
64 .bind(&[row.id.as_str().into()])?
65 .first::<RunRow>(None)
66 .await?
67 .map(|run| run.summary());
68 let parsed = workflow::parse(&row.source).ok();
69 Ok(Workflow {
70 id: row.id.clone(),
71 path: row.path.clone(),
72 name: row.name.clone(),
73 events: serde_json::from_str(&row.events).unwrap_or_default(),
74 state: row.state.clone(),
75 error: row.error.clone(),
76 notes: notes(&row.source),
77 dispatch: parsed
78 .as_ref()
79 .filter(|_| row.error.is_none())
80 .and_then(|w| w.trigger("workflow_dispatch"))
81 .map(|t| serde_json::Value::Object(t.inputs.clone())),
82 last_run,
83 })
84 }
85
86 pub async fn workflows(&self, a: WorkflowsArgs) -> Result<Outcome<Vec<Workflow>>> {
87 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
88 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
89 };
90 if !self.synced(&repo.id).await?
91 && let Some(ws) = self.workspace_actor(&repo.namespace).await?
92 {
93 self.sync(&repo, &ws).await?;
94 }
95 let rows = self
96 .db
97 .prepare(
98 "SELECT * FROM workflows WHERE repo_id = ?
99 AND (error IS NULL OR error NOT LIKE 'Its file is%' OR id IN (SELECT workflow_id FROM runs WHERE repo_id = ?))
100 ORDER BY name",
101 )
102 .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
103 .all()
104 .await?
105 .results::<WorkflowRow>()?;
106 let mut out = Vec::with_capacity(rows.len());
107 for row in &rows {
108 out.push(self.summary(row).await?);
109 }
110 Ok(Outcome::Ok(out))
111 }
112
113 pub async fn runs(&self, a: RunsArgs) -> Result<Outcome<Vec<WorkflowRun>>> {
114 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
115 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
116 };
117 let mut sql = "SELECT * FROM runs WHERE repo_id = ?".to_owned();
118 let mut binds: Vec<worker::wasm_bindgen::JsValue> = vec![repo.id.as_str().into()];
119 if let Some(workflow) = &a.workflow {
120 sql.push_str(" AND (workflow_id = ? OR path = ? OR path = ?)");
121 binds.push(workflow.as_str().into());
122 binds.push(workflow.as_str().into());
123 binds.push(format!("{}/{workflow}", g1t_actions::workflow::FOLDER).into());
124 }
125 if let Some(branch) = &a.branch {
126 sql.push_str(" AND (git_ref = ? OR json_extract(info, '$.headRef') = ?)");
127 binds.push(format!("refs/heads/{branch}").into());
128 binds.push(branch.as_str().into());
129 }
130 if let Some(event) = &a.event {
131 sql.push_str(" AND event = ?");
132 binds.push(event.as_str().into());
133 }
134 if let Some(pull) = a.pull {
135 sql.push_str(" AND pull = ?");
136 binds.push(pull.into());
137 }
138 if let Some(sha) = &a.sha {
139 sql.push_str(" AND sha = ?");
140 binds.push(sha.as_str().into());
141 }
142 sql.push_str(" ORDER BY id DESC LIMIT ?");
143 binds.push(a.limit.unwrap_or(RUNS_SHOWN).clamp(1, 100).into());
144 let rows = self.db.prepare(sql).bind(&binds)?.all().await?.results::<RunRow>()?;
145 Ok(Outcome::Ok(rows.iter().map(RunRow::summary).collect()))
146 }
147
148 pub async fn run(&self, a: RunArgs) -> Result<Outcome<RunDetail>> {
149 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
150 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
151 }
152 let run = check!(self.run_in(&a.repo, &a.id).await?);
153 let jobs = self.job_rows(&run.id).await?.into_iter().map(job_view).collect();
154 Ok(Outcome::Ok(RunDetail {
155 notes: notes(&run.source),
156 run: run.summary(),
157 jobs,
158 }))
159 }
160
161 pub async fn logs(&self, a: LogsArgs) -> Result<Outcome<JobLog>> {
162 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
163 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
164 }
165 let job = self
166 .db
167 .prepare("SELECT jobs.* FROM jobs JOIN runs ON runs.id = jobs.run_id WHERE jobs.id = ? AND lower(runs.repo) = lower(?)")
168 .bind(&[a.job.as_str().into(), format!("{}/{}", a.repo.namespace, a.repo.name).into()])?
169 .first::<JobRow>(None)
170 .await?;
171 let Some(job) = job else {
172 return Ok(fail(FailureCode::NotFound, "No such job."));
173 };
174 #[derive(Deserialize)]
175 struct Row {
176 seq: u64,
177 step: u32,
178 text: String,
179 }
180 let chunks = self
181 .db
182 .prepare("SELECT seq, step, text FROM logs WHERE job_id = ? AND seq > ? ORDER BY seq LIMIT 500")
183 .bind(&[job.id.as_str().into(), (a.after as f64).into()])?
184 .all()
185 .await?
186 .results::<Row>()?
187 .into_iter()
188 .map(|row| LogChunk { seq: row.seq, step: row.step, text: row.text })
189 .collect();
190 Ok(Outcome::Ok(JobLog {
191 chunks,
192 done: job.status == "completed",
193 }))
194 }
195
196 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 look197 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::ManageSettings).await? {
GitHub Actions on g1t, part two: running workflows198 return Ok(Outcome::Fail(refused));
199 }
GitHub Actions on g1t, part three: .g1t/workflows, the pages, the docs200 let wanted = a.workflow.trim_start_matches(".g1t/workflows/");
GitHub Actions on g1t, part two: running workflows201 let row = self
202 .db
203 .prepare("SELECT * FROM workflows WHERE lower(repo) = lower(?) AND (id = ? OR path = ?)")
204 .bind(&[
205 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
206 a.workflow.as_str().into(),
207 format!("{}/{wanted}", g1t_actions::workflow::FOLDER).into(),
208 ])?
209 .first::<WorkflowRow>(None)
210 .await?;
211 let Some(row) = row else {
212 return Ok(fail(FailureCode::NotFound, "No such workflow."));
213 };
214 let state = if a.enabled { "active" } else { "disabled" };
215 self.db.prepare("UPDATE workflows SET state = ? WHERE id = ?").bind(&[state.into(), row.id.as_str().into()])?.run().await?;
216 let row = WorkflowRow { state: state.to_owned(), ..row };
217 Ok(Outcome::Ok(self.summary(&row).await?))
218 }
219}