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/work/src/checks.rs

232 lines8,127 bytesCodeBlame
1//! Check runs: the record of the merge queue taking a pull request out,
2//! and of commands written on issues, which g1t ran before a pull
3//! request's checks were the workflows run on it. Those runs are kept so
4//! their history still reads; a sandbox still finishing one reports
5//! through the API with the job's one-time token.
6
7use g1t_contracts::events::ChecksEvent;
8use g1t_contracts::time::rfc3339;
9use g1t_contracts::work::*;
10use g1t_contracts::{FailureCode, Outcome};
11use g1t_kit::now_ms;
12use serde::Deserialize;
13use sha2::{Digest, Sha256};
14use worker::Result;
15use worker::wasm_bindgen::JsValue;
16
17use crate::rows::{PULL_COLUMNS, PullRow};
18use crate::{Work, optional};
19
20const MAX_OUTPUT_CHARS: usize = 16_000;
21const MAX_RESULTS: usize = 20;
22/// How many earlier runs a pull request shows.
23const EARLIER_RUNS: u32 = 10;
24
25#[derive(Deserialize)]
26struct RunRow {
27 id: String,
28 pull_id: String,
29 head_commit: String,
30 status: CheckStatus,
31 /// JSON array of results.
32 results: String,
33 error: Option<String>,
34 token_hash: String,
35 created_at: String,
36 finished_at: Option<String>,
37}
38
39impl From<RunRow> for CheckRun {
40 fn from(row: RunRow) -> Self {
41 CheckRun {
42 id: row.id,
43 head_commit: row.head_commit,
44 status: row.status,
45 results: serde_json::from_str(&row.results).unwrap_or_default(),
46 error: row.error,
47 created_at: row.created_at,
48 finished_at: row.finished_at,
49 }
50 }
51}
52
53pub(crate) fn hash(token: &str) -> String {
54 hex::encode(Sha256::digest(token.as_bytes()))
55}
56
57pub(crate) fn new_token() -> String {
58 let mut bytes = [0u8; 32];
59 getrandom::getrandom(&mut bytes).expect("no source of randomness");
60 hex::encode(bytes)
61}
62
63fn refused<T>(message: &str) -> Outcome<T> {
64 Outcome::fail(FailureCode::Conflict, message)
65}
66
67impl Work {
68 /// The most recent check run of a pull request.
69 pub(crate) async fn latest_checks(&self, pull_id: &str) -> Result<Option<CheckRun>> {
70 Ok(self
71 .db
72 .prepare("SELECT * FROM check_runs WHERE pull_id = ? ORDER BY id DESC LIMIT 1")
73 .bind(&[pull_id.into()])?
74 .first::<RunRow>(None)
75 .await?
76 .map(CheckRun::from))
77 }
78
79 /// The runs before the latest, newest first, without what each command
80 /// printed: enough to see how the checks went over time.
81 pub(crate) async fn earlier_checks(&self, pull_id: &str) -> Result<Vec<CheckRun>> {
82 let rows = self
83 .db
84 .prepare(
85 "SELECT * FROM check_runs WHERE pull_id = ? ORDER BY id DESC LIMIT ? OFFSET 1",
86 )
87 .bind(&[pull_id.into(), EARLIER_RUNS.into()])?
88 .all()
89 .await?
90 .results::<RunRow>()?;
91 Ok(rows
92 .into_iter()
93 .map(|row| {
94 let mut run = CheckRun::from(row);
95 for result in &mut run.results {
96 result.output = String::new();
97 }
98 run
99 })
100 .collect())
101 }
102
103 /// Commands written on issues are no longer run: a pull request's
104 /// checks are the workflows run on it, and the default branch's
105 /// protection says which must pass. Refused, for a runner from before.
106 pub(crate) async fn start_checks(&self, _: StartChecksArgs) -> Result<Outcome<CheckJob>> {
107 Ok(refused(
108 "Checks are the workflows run on a pull request; there are no commands to run.",
109 ))
110 }
111
112 /// Records what a sandbox reports for its run: that it has started, its
113 /// results, that it could not run, or that the run should be forgotten.
114 /// The run's token is the only credential, and a finished run accepts
115 /// nothing more.
116 pub(crate) async fn report_checks(&self, a: ReportChecksArgs) -> Result<Outcome<CheckRun>> {
117 let run = self
118 .db
119 .prepare("SELECT * FROM check_runs WHERE id = ?")
120 .bind(&[a.run_id.as_str().into()])?
121 .first::<RunRow>(None)
122 .await?;
123 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
124 return Ok(Outcome::fail(FailureCode::NotFound, "Check run not found."));
125 };
126 if run.finished_at.is_some() {
127 return Ok(refused("This check run has already finished."));
128 }
129 let latest = "id = ? AND check_run_id = ?";
130 let pull_keys =
131 || -> [JsValue; 2] { [run.pull_id.as_str().into(), run.id.as_str().into()] };
132
133 if a.skip {
134 self.db
135 .batch(vec![
136 self.db
137 .prepare("DELETE FROM check_runs WHERE id = ?")
138 .bind(&[run.id.as_str().into()])?,
139 self.db
140 .prepare(format!(
141 "UPDATE pulls SET check_status = NULL, check_run_id = NULL WHERE {latest}"
142 ))
143 .bind(&pull_keys())?,
144 ])
145 .await?;
146 return Ok(Outcome::Ok(run.into()));
147 }
148
149 let results: Vec<CheckResult> = a
150 .results
151 .into_iter()
152 .take(MAX_RESULTS)
153 .map(|mut result| {
154 // Keep the end of long output: that is where failures are.
155 let length = result.output.chars().count();
156 if length > MAX_OUTPUT_CHARS {
157 result.output = result
158 .output
159 .chars()
160 .skip(length - MAX_OUTPUT_CHARS)
161 .collect();
162 }
163 result
164 })
165 .collect();
166 let status = if a.error.is_some() {
167 CheckStatus::Errored
168 } else if results.is_empty() {
169 CheckStatus::Running
170 } else if results.iter().all(|result| result.passed) {
171 CheckStatus::Passed
172 } else {
173 CheckStatus::Failed
174 };
175 let finished = (status != CheckStatus::Running).then(|| rfc3339(now_ms()));
176 self.db
177 .batch(vec![
178 self.db
179 .prepare(
180 "UPDATE check_runs SET status = ?, results = ?, error = ?, finished_at = ?
181 WHERE id = ?",
182 )
183 .bind(&[
184 status.as_str().into(),
185 serde_json::to_string(&results)?.into(),
186 optional(&a.error),
187 optional(&finished),
188 run.id.as_str().into(),
189 ])?,
190 self.db
191 .prepare(format!("UPDATE pulls SET check_status = ? WHERE {latest}"))
192 .bind(&[
193 status.as_str().into(),
194 run.pull_id.as_str().into(),
195 run.id.as_str().into(),
196 ])?,
197 ])
198 .await?;
199
200 if finished.is_some() {
201 let pull = self
202 .db
203 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE id = ?"))
204 .bind(&[run.pull_id.as_str().into()])?
205 .first::<PullRow>(None)
206 .await?
207 .map(Pull::from);
208 if let Some(pull) = pull {
209 self.publish_as(
210 "checks.completed",
211 &pull.repo_id,
212 None,
213 ChecksEvent {
214 pull_id: pull.id.clone(),
215 repo_id: pull.repo_id.clone(),
216 number: pull.number,
217 status: status.as_str(),
218 commit: run.head_commit.clone(),
219 },
220 )
221 .await?;
222 }
223 }
224 Ok(Outcome::Ok(CheckRun {
225 status,
226 results,
227 error: a.error,
228 finished_at: finished,
229 ..run.into()
230 }))
231 }
232}