pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/checks.rs

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