pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/checks.rs
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.
| Acceptance checks in sandboxes, line comments and review verdicts | 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 | ||
| 9 | use g1t_contracts::events::ChecksEvent; | |
| 10 | use g1t_contracts::repos::{GetByIdArgs, HeadArgs, Repo, RepoPath}; | |
| 11 | use g1t_contracts::time::rfc3339; | |
| 12 | use g1t_contracts::work::*; | |
| 13 | use g1t_contracts::{FailureCode, Outcome, new_id}; | |
| 14 | use g1t_kit::now_ms; | |
| 15 | use serde::Deserialize; | |
| 16 | use sha2::{Digest, Sha256}; | |
| 17 | use worker::Result; | |
| 18 | use worker::wasm_bindgen::JsValue; | |
| 19 | ||
| 20 | use crate::rows::PullRow; | |
| 21 | use crate::{Work, optional}; | |
| 22 | ||
| 23 | const MAX_OUTPUT_CHARS: usize = 16_000; | |
| 24 | const MAX_RESULTS: usize = 20; | |
| 25 | ||
| 26 | #[derive(Deserialize)] | |
| 27 | struct 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 | ||
| 40 | impl 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 | ||
| Agents as a team: lifecycle, merge queue, billing and a new shell | 54 | pub(crate) fn hash(token: &str) -> String { |
| Acceptance checks in sandboxes, line comments and review verdicts | 55 | hex::encode(Sha256::digest(token.as_bytes())) |
| 56 | } | |
| 57 | ||
| Agents as a team: lifecycle, merge queue, billing and a new shell | 58 | pub(crate) fn new_token() -> String { |
| Acceptance checks in sandboxes, line comments and review verdicts | 59 | let mut bytes = [0u8; 32]; |
| 60 | getrandom::getrandom(&mut bytes).expect("no source of randomness"); | |
| 61 | hex::encode(bytes) | |
| 62 | } | |
| 63 | ||
| 64 | fn refused<T>(message: &str) -> Outcome<T> { | |
| 65 | Outcome::fail(FailureCode::Conflict, message) | |
| 66 | } | |
| 67 | ||
| 68 | impl 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 = Some(pull.author.clone()); | |
| 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 | } |