pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/checks.rs

349 lines12,290 bytesCodeBlame
1//! Check runs: an issue's acceptance checks, run against a pull request's
2//! head in a clean sandbox.
3//!
4//! This service keeps the record. The runner service starts the sandbox:
5//! it asks for a job with `start_checks`, and the sandbox reports back
6//! through the API with the job's one-time token. Nothing else can write a
7//! result, including the agent whose work is being checked.
8
9use g1t_contracts::events::ChecksEvent;
10use g1t_contracts::repos::{GetByIdArgs, HeadArgs, Repo, RepoPath};
11use g1t_contracts::time::rfc3339;
12use g1t_contracts::work::*;
13use g1t_contracts::{FailureCode, Outcome, new_id};
14use g1t_kit::now_ms;
15use serde::Deserialize;
16use sha2::{Digest, Sha256};
17use worker::Result;
18use worker::wasm_bindgen::JsValue;
19
20use crate::rows::PullRow;
21use crate::{Work, optional};
22
23const MAX_OUTPUT_CHARS: usize = 16_000;
24const MAX_RESULTS: usize = 20;
25/// How many earlier runs a pull request shows.
26const EARLIER_RUNS: u32 = 10;
27
28#[derive(Deserialize)]
29struct RunRow {
30 id: String,
31 pull_id: String,
32 head_commit: String,
33 status: CheckStatus,
34 /// JSON array of results.
35 results: String,
36 error: Option<String>,
37 token_hash: String,
38 created_at: String,
39 finished_at: Option<String>,
40}
41
42impl From<RunRow> for CheckRun {
43 fn from(row: RunRow) -> Self {
44 CheckRun {
45 id: row.id,
46 head_commit: row.head_commit,
47 status: row.status,
48 results: serde_json::from_str(&row.results).unwrap_or_default(),
49 error: row.error,
50 created_at: row.created_at,
51 finished_at: row.finished_at,
52 }
53 }
54}
55
56pub(crate) fn hash(token: &str) -> String {
57 hex::encode(Sha256::digest(token.as_bytes()))
58}
59
60pub(crate) fn new_token() -> String {
61 let mut bytes = [0u8; 32];
62 getrandom::getrandom(&mut bytes).expect("no source of randomness");
63 hex::encode(bytes)
64}
65
66fn refused<T>(message: &str) -> Outcome<T> {
67 Outcome::fail(FailureCode::Conflict, message)
68}
69
70impl Work {
71 /// The most recent check run of a pull request.
72 pub(crate) async fn latest_checks(&self, pull_id: &str) -> Result<Option<CheckRun>> {
73 Ok(self
74 .db
75 .prepare("SELECT * FROM check_runs WHERE pull_id = ? ORDER BY id DESC LIMIT 1")
76 .bind(&[pull_id.into()])?
77 .first::<RunRow>(None)
78 .await?
79 .map(CheckRun::from))
80 }
81
82 /// The runs before the latest, newest first, without what each command
83 /// printed: enough to see how the checks went over time.
84 pub(crate) async fn earlier_checks(&self, pull_id: &str) -> Result<Vec<CheckRun>> {
85 let rows = self
86 .db
87 .prepare(
88 "SELECT * FROM check_runs WHERE pull_id = ? ORDER BY id DESC LIMIT ? OFFSET 1",
89 )
90 .bind(&[pull_id.into(), EARLIER_RUNS.into()])?
91 .all()
92 .await?
93 .results::<RunRow>()?;
94 Ok(rows
95 .into_iter()
96 .map(|row| {
97 let mut run = CheckRun::from(row);
98 for result in &mut run.results {
99 result.output = String::new();
100 }
101 run
102 })
103 .collect())
104 }
105
106 /// Begins a check run for a pull request that is ready for review, and
107 /// returns what a sandbox needs to carry it out. Any run still in
108 /// progress for the pull request is abandoned.
109 pub(crate) async fn start_checks(&self, a: StartChecksArgs) -> Result<Outcome<CheckJob>> {
110 let pull = self
111 .db
112 .prepare("SELECT * FROM pulls WHERE id = ?")
113 .bind(&[a.pull_id.as_str().into()])?
114 .first::<PullRow>(None)
115 .await?
116 .map(Pull::from);
117 let Some(pull) = pull else {
118 return Ok(Outcome::fail(
119 FailureCode::NotFound,
120 "Pull request not found.",
121 ));
122 };
123 if pull.status != PullStatus::Open {
124 return Ok(refused(
125 "Checks run once a pull request is ready for review.",
126 ));
127 }
128 let issue = match pull.issue {
129 Some(number) => self.issue(&pull.repo_id, number).await?,
130 None => None,
131 };
132 let Some(issue) = issue.filter(|issue| !issue.checks.is_empty()) else {
133 return Ok(refused(
134 "This pull request's issue has no acceptance checks.",
135 ));
136 };
137
138 // The author can read both the repository and the pull request's source.
139 let viewer = self.author_viewer(&pull).await?;
140 let repo: Outcome<Repo> = g1t_kit::call(
141 &self.repos,
142 "get_by_id",
143 &GetByIdArgs {
144 id: pull.repo_id.clone(),
145 viewer,
146 },
147 )
148 .await?;
149 let Outcome::Ok(repo) = crate::retired::unless_archived(repo) else {
150 return Ok(Outcome::fail(
151 FailureCode::NotFound,
152 "Pull request not found.",
153 ));
154 };
155 // Asked of the store, since the recorded head can lag a push.
156 let head: Option<String> = g1t_kit::call(
157 &self.repos,
158 "head",
159 &HeadArgs {
160 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
161 branch: pull
162 .branch
163 .clone()
164 .unwrap_or_else(|| repo.default_branch.clone()),
165 },
166 )
167 .await?;
168 let Some(commit) = head else {
169 return Ok(refused("This pull request has no commits to check."));
170 };
171
172 let now = now_ms();
173 let timestamp = rfc3339(now);
174 let run_id = new_id("chk", now);
175 let token = new_token();
176 self.db
177 .batch(vec![
178 self.db
179 .prepare(
180 "UPDATE check_runs
181 SET status = 'errored', error = 'Replaced by a newer run.', finished_at = ?
182 WHERE pull_id = ? AND finished_at IS NULL",
183 )
184 .bind(&[timestamp.as_str().into(), pull.id.as_str().into()])?,
185 self.db
186 .prepare(
187 "INSERT INTO check_runs (id, pull_id, head_commit, token_hash, created_at)
188 VALUES (?, ?, ?, ?, ?)",
189 )
190 .bind(&[
191 run_id.as_str().into(),
192 pull.id.as_str().into(),
193 commit.as_str().into(),
194 hash(&token).into(),
195 timestamp.as_str().into(),
196 ])?,
197 self.db
198 .prepare(
199 "UPDATE pulls SET check_status = 'queued', check_run_id = ?, head_commit = ?
200 WHERE id = ?",
201 )
202 .bind(&[
203 run_id.as_str().into(),
204 commit.as_str().into(),
205 pull.id.as_str().into(),
206 ])?,
207 ])
208 .await?;
209
210 Ok(Outcome::Ok(CheckJob {
211 run_id,
212 token,
213 commands: issue.checks,
214 source: pull.fork.clone().unwrap_or(RepoPath {
215 namespace: repo.namespace.clone(),
216 name: repo.name.clone(),
217 }),
218 commit,
219 author: pull.author,
220 requested_by: issue.author.username,
221 repo: RepoPath {
222 namespace: repo.namespace,
223 name: repo.name,
224 },
225 number: pull.number,
226 }))
227 }
228
229 /// Records what a sandbox reports for its run: that it has started, its
230 /// results, that it could not run, or that the run should be forgotten.
231 /// The run's token is the only credential, and a finished run accepts
232 /// nothing more.
233 pub(crate) async fn report_checks(&self, a: ReportChecksArgs) -> Result<Outcome<CheckRun>> {
234 let run = self
235 .db
236 .prepare("SELECT * FROM check_runs WHERE id = ?")
237 .bind(&[a.run_id.as_str().into()])?
238 .first::<RunRow>(None)
239 .await?;
240 let Some(run) = run.filter(|run| run.token_hash == hash(&a.token)) else {
241 return Ok(Outcome::fail(FailureCode::NotFound, "Check run not found."));
242 };
243 if run.finished_at.is_some() {
244 return Ok(refused("This check run has already finished."));
245 }
246 let latest = "id = ? AND check_run_id = ?";
247 let pull_keys =
248 || -> [JsValue; 2] { [run.pull_id.as_str().into(), run.id.as_str().into()] };
249
250 if a.skip {
251 self.db
252 .batch(vec![
253 self.db
254 .prepare("DELETE FROM check_runs WHERE id = ?")
255 .bind(&[run.id.as_str().into()])?,
256 self.db
257 .prepare(format!(
258 "UPDATE pulls SET check_status = NULL, check_run_id = NULL WHERE {latest}"
259 ))
260 .bind(&pull_keys())?,
261 ])
262 .await?;
263 return Ok(Outcome::Ok(run.into()));
264 }
265
266 let results: Vec<CheckResult> = a
267 .results
268 .into_iter()
269 .take(MAX_RESULTS)
270 .map(|mut result| {
271 // Keep the end of long output: that is where failures are.
272 let length = result.output.chars().count();
273 if length > MAX_OUTPUT_CHARS {
274 result.output = result
275 .output
276 .chars()
277 .skip(length - MAX_OUTPUT_CHARS)
278 .collect();
279 }
280 result
281 })
282 .collect();
283 let status = if a.error.is_some() {
284 CheckStatus::Errored
285 } else if results.is_empty() {
286 CheckStatus::Running
287 } else if results.iter().all(|result| result.passed) {
288 CheckStatus::Passed
289 } else {
290 CheckStatus::Failed
291 };
292 let finished = (status != CheckStatus::Running).then(|| rfc3339(now_ms()));
293 self.db
294 .batch(vec![
295 self.db
296 .prepare(
297 "UPDATE check_runs SET status = ?, results = ?, error = ?, finished_at = ?
298 WHERE id = ?",
299 )
300 .bind(&[
301 status.as_str().into(),
302 serde_json::to_string(&results)?.into(),
303 optional(&a.error),
304 optional(&finished),
305 run.id.as_str().into(),
306 ])?,
307 self.db
308 .prepare(format!("UPDATE pulls SET check_status = ? WHERE {latest}"))
309 .bind(&[
310 status.as_str().into(),
311 run.pull_id.as_str().into(),
312 run.id.as_str().into(),
313 ])?,
314 ])
315 .await?;
316
317 if finished.is_some() {
318 let pull = self
319 .db
320 .prepare("SELECT * FROM pulls WHERE id = ?")
321 .bind(&[run.pull_id.as_str().into()])?
322 .first::<PullRow>(None)
323 .await?
324 .map(Pull::from);
325 if let Some(pull) = pull {
326 self.publish_as(
327 "checks.completed",
328 &pull.repo_id,
329 None,
330 ChecksEvent {
331 pull_id: pull.id.clone(),
332 repo_id: pull.repo_id.clone(),
333 number: pull.number,
334 status: status.as_str(),
335 commit: run.head_commit.clone(),
336 },
337 )
338 .await?;
339 }
340 }
341 Ok(Outcome::Ok(CheckRun {
342 status,
343 results,
344 error: a.error,
345 finished_at: finished,
346 ..run.into()
347 }))
348 }
349}