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

217 lines8,569 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::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,
54 }
55}
56
57impl Actions {
58 async fn summary(&self, row: &WorkflowRow) -> Result<Workflow> {
59 let last_run = self
60 .db
61 .prepare("SELECT * FROM runs WHERE workflow_id = ? ORDER BY id DESC LIMIT 1")
62 .bind(&[row.id.as_str().into()])?
63 .first::<RunRow>(None)
64 .await?
65 .map(|run| run.summary());
66 let parsed = workflow::parse(&row.source).ok();
67 Ok(Workflow {
68 id: row.id.clone(),
69 path: row.path.clone(),
70 name: row.name.clone(),
71 events: serde_json::from_str(&row.events).unwrap_or_default(),
72 state: row.state.clone(),
73 error: row.error.clone(),
74 notes: notes(&row.source),
75 dispatch: parsed
76 .as_ref()
77 .filter(|_| row.error.is_none())
78 .and_then(|w| w.trigger("workflow_dispatch"))
79 .map(|t| serde_json::Value::Object(t.inputs.clone())),
80 last_run,
81 })
82 }
83
84 pub async fn workflows(&self, a: WorkflowsArgs) -> Result<Outcome<Vec<Workflow>>> {
85 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
86 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
87 };
88 if !self.synced(&repo.id).await?
89 && let Some(ws) = self.workspace_actor(&repo.namespace).await?
90 {
91 self.sync(&repo, &ws).await?;
92 }
93 let rows = self
94 .db
95 .prepare(
96 "SELECT * FROM workflows WHERE repo_id = ?
97 AND (error IS NULL OR error NOT LIKE 'Its file is%' OR id IN (SELECT workflow_id FROM runs WHERE repo_id = ?))
98 ORDER BY name",
99 )
100 .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
101 .all()
102 .await?
103 .results::<WorkflowRow>()?;
104 let mut out = Vec::with_capacity(rows.len());
105 for row in &rows {
106 out.push(self.summary(row).await?);
107 }
108 Ok(Outcome::Ok(out))
109 }
110
111 pub async fn runs(&self, a: RunsArgs) -> Result<Outcome<Vec<WorkflowRun>>> {
112 let Some(repo) = self.visible_repo(&a.repo, &a.viewer).await? else {
113 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
114 };
115 let mut sql = "SELECT * FROM runs WHERE repo_id = ?".to_owned();
116 let mut binds: Vec<worker::wasm_bindgen::JsValue> = vec![repo.id.as_str().into()];
117 if let Some(workflow) = &a.workflow {
118 sql.push_str(" AND (workflow_id = ? OR path = ? OR path = ?)");
119 binds.push(workflow.as_str().into());
120 binds.push(workflow.as_str().into());
121 binds.push(format!("{}/{workflow}", g1t_actions::workflow::FOLDER).into());
122 }
123 if let Some(branch) = &a.branch {
124 sql.push_str(" AND (git_ref = ? OR json_extract(info, '$.headRef') = ?)");
125 binds.push(format!("refs/heads/{branch}").into());
126 binds.push(branch.as_str().into());
127 }
128 if let Some(event) = &a.event {
129 sql.push_str(" AND event = ?");
130 binds.push(event.as_str().into());
131 }
132 if let Some(pull) = a.pull {
133 sql.push_str(" AND pull = ?");
134 binds.push(pull.into());
135 }
136 if let Some(sha) = &a.sha {
137 sql.push_str(" AND sha = ?");
138 binds.push(sha.as_str().into());
139 }
140 sql.push_str(" ORDER BY id DESC LIMIT ?");
141 binds.push(a.limit.unwrap_or(RUNS_SHOWN).clamp(1, 100).into());
142 let rows = self.db.prepare(sql).bind(&binds)?.all().await?.results::<RunRow>()?;
143 Ok(Outcome::Ok(rows.iter().map(RunRow::summary).collect()))
144 }
145
146 pub async fn run(&self, a: RunArgs) -> Result<Outcome<RunDetail>> {
147 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
148 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
149 }
150 let run = check!(self.run_in(&a.repo, &a.id).await?);
151 let jobs = self.job_rows(&run.id).await?.into_iter().map(job_view).collect();
152 Ok(Outcome::Ok(RunDetail {
153 notes: notes(&run.source),
154 run: run.summary(),
155 jobs,
156 }))
157 }
158
159 pub async fn logs(&self, a: LogsArgs) -> Result<Outcome<JobLog>> {
160 if self.visible_repo(&a.repo, &a.viewer).await?.is_none() {
161 return Ok(fail(FailureCode::NotFound, "There is no such repository."));
162 }
163 let job = self
164 .db
165 .prepare("SELECT jobs.* FROM jobs JOIN runs ON runs.id = jobs.run_id WHERE jobs.id = ? AND lower(runs.repo) = lower(?)")
166 .bind(&[a.job.as_str().into(), format!("{}/{}", a.repo.namespace, a.repo.name).into()])?
167 .first::<JobRow>(None)
168 .await?;
169 let Some(job) = job else {
170 return Ok(fail(FailureCode::NotFound, "No such job."));
171 };
172 #[derive(Deserialize)]
173 struct Row {
174 seq: u64,
175 step: u32,
176 text: String,
177 }
178 let chunks = self
179 .db
180 .prepare("SELECT seq, step, text FROM logs WHERE job_id = ? AND seq > ? ORDER BY seq LIMIT 500")
181 .bind(&[job.id.as_str().into(), (a.after as f64).into()])?
182 .all()
183 .await?
184 .results::<Row>()?
185 .into_iter()
186 .map(|row| LogChunk { seq: row.seq, step: row.step, text: row.text })
187 .collect();
188 Ok(Outcome::Ok(JobLog {
189 chunks,
190 done: job.status == "completed",
191 }))
192 }
193
194 pub async fn set_workflow_enabled(&self, a: SetWorkflowEnabledArgs) -> Result<Outcome<Workflow>> {
195 if let Outcome::Fail(refused) = self.may(&a.actor, &a.repo, Capability::ManageSettings).await? {
196 return Ok(Outcome::Fail(refused));
197 }
198 let wanted = a.workflow.trim_start_matches(".g1t/workflows/");
199 let row = self
200 .db
201 .prepare("SELECT * FROM workflows WHERE lower(repo) = lower(?) AND (id = ? OR path = ?)")
202 .bind(&[
203 format!("{}/{}", a.repo.namespace, a.repo.name).into(),
204 a.workflow.as_str().into(),
205 format!("{}/{wanted}", g1t_actions::workflow::FOLDER).into(),
206 ])?
207 .first::<WorkflowRow>(None)
208 .await?;
209 let Some(row) = row else {
210 return Ok(fail(FailureCode::NotFound, "No such workflow."));
211 };
212 let state = if a.enabled { "active" } else { "disabled" };
213 self.db.prepare("UPDATE workflows SET state = ? WHERE id = ?").bind(&[state.into(), row.id.as_str().into()])?.run().await?;
214 let row = WorkflowRow { state: state.to_owned(), ..row };
215 Ok(Outcome::Ok(self.summary(&row).await?))
216 }
217}