g1t/services/work/src/statuses.rs
| 1 | //! Statuses on commits: what workflow runs (and other tools, such as |
| 2 | //! deployments) say about a pull request's head. These are its checks. |
| 3 | //! |
| 4 | //! The default branch's protection names the checks that must pass |
| 5 | //! (`RepoSettings::required_checks`): a required check that failed, is |
| 6 | //! still running or has not reported refuses the merge, for everyone and |
| 7 | //! for the merge queue. Where g1t sees an agent's pull request through, |
| 8 | //! any check that failed sends the agent back to fix it, with what the |
| 9 | //! failing jobs printed; once it is out of revisions, only a required |
| 10 | //! check holds the pull request for a person. |
| 11 | |
| 12 | use g1t_contracts::events::ChecksEvent; |
| 13 | use g1t_contracts::time::rfc3339; |
| 14 | use g1t_contracts::work::{ |
| 15 | CommitStatus, RequiredCheck, RequiredState, SeenCheck, SeenChecksArgs, SetCommitStatusArgs, check_name, |
| 16 | required_checks, |
| 17 | }; |
| 18 | use g1t_contracts::{FailureCode, Outcome}; |
| 19 | use g1t_kit::now_ms; |
| 20 | use serde::Deserialize; |
| 21 | use worker::Result; |
| 22 | |
| 23 | use crate::Work; |
| 24 | |
| 25 | #[derive(Deserialize)] |
| 26 | struct StatusRow { |
| 27 | context: String, |
| 28 | state: String, |
| 29 | description: Option<String>, |
| 30 | target_url: Option<String>, |
| 31 | updated_at: String, |
| 32 | } |
| 33 | |
| 34 | impl From<StatusRow> for CommitStatus { |
| 35 | fn from(row: StatusRow) -> Self { |
| 36 | CommitStatus { |
| 37 | context: row.context, |
| 38 | state: row.state, |
| 39 | description: row.description, |
| 40 | target_url: row.target_url, |
| 41 | updated_at: row.updated_at, |
| 42 | } |
| 43 | } |
| 44 | } |
| 45 | |
| 46 | #[derive(Deserialize)] |
| 47 | struct HeadRow { |
| 48 | id: String, |
| 49 | number: u32, |
| 50 | } |
| 51 | |
| 52 | /// What a commit's checks say: every status still running and every one |
| 53 | /// that failed, by context, and where each required check stands. |
| 54 | #[derive(Clone, Debug, Default, PartialEq)] |
| 55 | pub(crate) struct WorkflowFacts { |
| 56 | pub(crate) pending: Vec<String>, |
| 57 | pub(crate) failed: Vec<String>, |
| 58 | pub(crate) required: Vec<RequiredCheck>, |
| 59 | } |
| 60 | |
| 61 | impl WorkflowFacts { |
| 62 | /// `required` names the checks the default branch's protection requires. |
| 63 | pub(crate) fn of(statuses: &[CommitStatus], required: &[String]) -> WorkflowFacts { |
| 64 | WorkflowFacts { |
| 65 | pending: statuses.iter().filter(|s| s.state == "pending").map(|s| s.context.clone()).collect(), |
| 66 | failed: statuses.iter().filter(|s| s.state == "failure" || s.state == "error").map(|s| s.context.clone()).collect(), |
| 67 | required: required_checks(required, statuses), |
| 68 | } |
| 69 | } |
| 70 | |
| 71 | fn required_in(&self, state: RequiredState) -> Vec<String> { |
| 72 | self.required.iter().filter(|check| check.state == state).map(|check| check.name.clone()).collect() |
| 73 | } |
| 74 | |
| 75 | /// The required checks that failed, by name. |
| 76 | pub(crate) fn required_failed(&self) -> Vec<String> { |
| 77 | self.required_in(RequiredState::Failure) |
| 78 | } |
| 79 | |
| 80 | /// The required checks nothing has reported on the commit yet. |
| 81 | pub(crate) fn expected(&self) -> Vec<String> { |
| 82 | self.required_in(RequiredState::Expected) |
| 83 | } |
| 84 | |
| 85 | /// Why a merge has to wait, if it does: a required check that failed, |
| 86 | /// is still running, or has not reported. Other checks never hold it. |
| 87 | pub(crate) fn refusal(&self) -> Option<String> { |
| 88 | let failed = self.required_failed(); |
| 89 | if !failed.is_empty() { |
| 90 | return Some(format!("The required {} {} failed.", checks_word(&failed), list(&failed))); |
| 91 | } |
| 92 | let running = self.required_in(RequiredState::Pending); |
| 93 | if !running.is_empty() { |
| 94 | let verb = if running.len() == 1 { "is" } else { "are" }; |
| 95 | return Some(format!("The required {} {} {verb} still running.", checks_word(&running), list(&running))); |
| 96 | } |
| 97 | let expected = self.expected(); |
| 98 | if !expected.is_empty() { |
| 99 | let verb = if expected.len() == 1 { "has" } else { "have" }; |
| 100 | return Some(format!( |
| 101 | "The required {} {} {verb} not reported on this commit yet.", |
| 102 | checks_word(&expected), |
| 103 | list(&expected) |
| 104 | )); |
| 105 | } |
| 106 | None |
| 107 | } |
| 108 | } |
| 109 | |
| 110 | fn checks_word(names: &[String]) -> &'static str { |
| 111 | if names.len() == 1 { "check" } else { "checks" } |
| 112 | } |
| 113 | |
| 114 | /// The check names in `(context, last reported)` rows, most recent first: |
| 115 | /// each name once, with the events it was reported for. |
| 116 | pub(crate) fn seen(rows: Vec<(String, String)>) -> Vec<SeenCheck> { |
| 117 | let mut rows = rows; |
| 118 | rows.sort_by(|a, b| b.1.cmp(&a.1)); |
| 119 | let mut out: Vec<SeenCheck> = Vec::new(); |
| 120 | for (context, at) in rows { |
| 121 | let (name, event) = check_name(&context); |
| 122 | match out.iter_mut().find(|seen| seen.name.eq_ignore_ascii_case(name)) { |
| 123 | Some(seen) => { |
| 124 | if let Some(event) = event |
| 125 | && !seen.events.iter().any(|known| known == event) |
| 126 | { |
| 127 | seen.events.push(event.to_owned()); |
| 128 | } |
| 129 | } |
| 130 | None => out.push(SeenCheck { |
| 131 | name: name.to_owned(), |
| 132 | events: event.map(|event| vec![event.to_owned()]).unwrap_or_default(), |
| 133 | last_seen: at, |
| 134 | }), |
| 135 | } |
| 136 | } |
| 137 | out |
| 138 | } |
| 139 | |
| 140 | /// How far back `seen_checks` looks, and the most contexts it reads. |
| 141 | const SEEN_DAYS: u64 = 30; |
| 142 | const SEEN_LIMIT: u32 = 200; |
| 143 | |
| 144 | #[derive(Deserialize)] |
| 145 | struct SeenRow { |
| 146 | context: String, |
| 147 | at: String, |
| 148 | } |
| 149 | |
| 150 | pub(crate) fn list(names: &[String]) -> String { |
| 151 | match names { |
| 152 | [] => String::new(), |
| 153 | [one] => one.clone(), |
| 154 | [rest @ .., last] => format!("{} and {last}", rest.join(", ")), |
| 155 | } |
| 156 | } |
| 157 | |
| 158 | impl Work { |
| 159 | /// Where a commit's checks stand, against the repository's required ones. |
| 160 | pub(crate) async fn facts(&self, repo_id: &str, sha: Option<&str>) -> Result<WorkflowFacts> { |
| 161 | let (statuses, settings) = |
| 162 | futures_util::future::try_join(self.statuses(repo_id, sha), self.settings(repo_id)).await?; |
| 163 | Ok(WorkflowFacts::of(&statuses, &settings.required_checks)) |
| 164 | } |
| 165 | |
| 166 | /// The check names reported on a repository's commits lately, for |
| 167 | /// choosing which to require. |
| 168 | pub(crate) async fn seen_checks(&self, a: SeenChecksArgs) -> Result<Outcome<Vec<SeenCheck>>> { |
| 169 | let repo = match self.repo(&a.repo, &a.viewer).await? { |
| 170 | Outcome::Ok(repo) => repo, |
| 171 | Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)), |
| 172 | }; |
| 173 | let since = rfc3339(now_ms().saturating_sub(SEEN_DAYS * 24 * 60 * 60 * 1000)); |
| 174 | let rows = self |
| 175 | .db |
| 176 | .prepare( |
| 177 | "SELECT context, MAX(updated_at) AS at FROM commit_statuses |
| 178 | WHERE repo_id = ? AND updated_at >= ? GROUP BY context ORDER BY at DESC LIMIT ?", |
| 179 | ) |
| 180 | .bind(&[repo.id.as_str().into(), since.into(), SEEN_LIMIT.into()])? |
| 181 | .all() |
| 182 | .await? |
| 183 | .results::<SeenRow>()?; |
| 184 | Ok(Outcome::Ok(seen(rows.into_iter().map(|row| (row.context, row.at)).collect()))) |
| 185 | } |
| 186 | |
| 187 | pub(crate) async fn statuses(&self, repo_id: &str, sha: Option<&str>) -> Result<Vec<CommitStatus>> { |
| 188 | let Some(sha) = sha else { return Ok(Vec::new()) }; |
| 189 | Ok(self |
| 190 | .db |
| 191 | .prepare("SELECT context, state, description, target_url, updated_at FROM commit_statuses WHERE repo_id = ? AND sha = ? ORDER BY context") |
| 192 | .bind(&[repo_id.into(), sha.into()])? |
| 193 | .all() |
| 194 | .await? |
| 195 | .results::<StatusRow>()? |
| 196 | .into_iter() |
| 197 | .map(CommitStatus::from) |
| 198 | .collect()) |
| 199 | } |
| 200 | |
| 201 | pub(crate) async fn set_commit_status(&self, a: SetCommitStatusArgs) -> Result<Outcome<bool>> { |
| 202 | if !matches!(a.state.as_str(), "pending" | "success" | "failure" | "error") { |
| 203 | return Ok(Outcome::fail(FailureCode::Invalid, "`state` is pending, success, failure or error.")); |
| 204 | } |
| 205 | self.db |
| 206 | .prepare( |
| 207 | "INSERT INTO commit_statuses (repo_id, sha, context, state, description, target_url, updated_at) |
| 208 | VALUES (?, ?, ?, ?, ?, ?, ?) |
| 209 | ON CONFLICT (repo_id, sha, context) DO UPDATE SET |
| 210 | state = excluded.state, description = excluded.description, |
| 211 | target_url = excluded.target_url, updated_at = excluded.updated_at", |
| 212 | ) |
| 213 | .bind(&[ |
| 214 | a.repo_id.as_str().into(), |
| 215 | a.sha.as_str().into(), |
| 216 | a.context.as_str().into(), |
| 217 | a.state.as_str().into(), |
| 218 | a.description.as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into), |
| 219 | a.target_url.as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into), |
| 220 | rfc3339(now_ms()).into(), |
| 221 | ])? |
| 222 | .run() |
| 223 | .await?; |
| 224 | if a.state == "pending" { |
| 225 | return Ok(Outcome::Ok(true)); |
| 226 | } |
| 227 | // Once every workflow on a pull request's head has finished, its |
| 228 | // lifecycle moves on, as it does when its checks finish. |
| 229 | let facts = self.facts(&a.repo_id, Some(&a.sha)).await?; |
| 230 | if !facts.pending.is_empty() { |
| 231 | return Ok(Outcome::Ok(true)); |
| 232 | } |
| 233 | // A merge queue state waiting on its merge_group workflows. |
| 234 | self.merge_group_finished(&a.repo_id, &a.sha, &facts).await?; |
| 235 | let heads = self |
| 236 | .db |
| 237 | .prepare("SELECT id, number FROM pulls WHERE repo_id = ? AND head_commit = ? AND status IN ('draft', 'open')") |
| 238 | .bind(&[a.repo_id.as_str().into(), a.sha.as_str().into()])? |
| 239 | .all() |
| 240 | .await? |
| 241 | .results::<HeadRow>()?; |
| 242 | for head in heads { |
| 243 | // A pull request g1t stopped on picks back up once what stopped |
| 244 | // it passes: its workflows, and its checks if it has any. The |
| 245 | // lifecycle then decides again, within its usual limits. |
| 246 | if facts.failed.is_empty() { |
| 247 | let resumed = self |
| 248 | .db |
| 249 | .prepare( |
| 250 | "UPDATE pulls SET stalled = NULL WHERE id = ? AND managed = 1 AND stalled IS NOT NULL |
| 251 | AND (check_status IS NULL OR check_status = 'passed') RETURNING id AS value", |
| 252 | ) |
| 253 | .bind(&[head.id.as_str().into()])? |
| 254 | .first::<crate::rows::ValueRow>(None) |
| 255 | .await?; |
| 256 | if resumed.is_some() { |
| 257 | self.note( |
| 258 | &a.repo_id, |
| 259 | head.number, |
| 260 | (crate::lifecycle::POLICY_ACTOR_ID, crate::lifecycle::POLICY_ACTOR_NAME), |
| 261 | "picked this back up: its workflows pass now", |
| 262 | ) |
| 263 | .await?; |
| 264 | } |
| 265 | } |
| 266 | self.publish_as( |
| 267 | "checks.completed", |
| 268 | &a.repo_id, |
| 269 | None, |
| 270 | ChecksEvent { |
| 271 | pull_id: head.id, |
| 272 | repo_id: a.repo_id.clone(), |
| 273 | number: head.number, |
| 274 | status: if facts.failed.is_empty() { "passed" } else { "failed" }, |
| 275 | commit: a.sha.clone(), |
| 276 | }, |
| 277 | ) |
| 278 | .await?; |
| 279 | } |
| 280 | Ok(Outcome::Ok(true)) |
| 281 | } |
| 282 | } |
| 283 | |
| 284 | #[cfg(test)] |
| 285 | mod tests { |
| 286 | use super::*; |
| 287 | |
| 288 | fn status(context: &str, state: &str) -> CommitStatus { |
| 289 | CommitStatus { |
| 290 | context: context.into(), |
| 291 | state: state.into(), |
| 292 | description: None, |
| 293 | target_url: None, |
| 294 | updated_at: String::new(), |
| 295 | } |
| 296 | } |
| 297 | |
| 298 | fn names(list: &[&str]) -> Vec<String> { |
| 299 | list.iter().map(|name| (*name).to_owned()).collect() |
| 300 | } |
| 301 | |
| 302 | #[test] |
| 303 | fn only_required_checks_hold_a_merge() { |
| 304 | let statuses = [status("CI / push", "pending"), status("Lint / pull_request", "failure"), status("Docs", "success")]; |
| 305 | // Nothing required: nothing holds it, whatever failed. |
| 306 | let free = WorkflowFacts::of(&statuses, &[]); |
| 307 | assert_eq!(free.pending, ["CI / push"]); |
| 308 | assert_eq!(free.failed, ["Lint / pull_request"]); |
| 309 | assert!(free.refusal().is_none()); |
| 310 | // Failures come before waiting. |
| 311 | let both = WorkflowFacts::of(&statuses, &names(&["CI", "Lint"])); |
| 312 | assert_eq!(both.refusal().unwrap(), "The required check Lint failed."); |
| 313 | let waiting = WorkflowFacts::of(&[status("A / pull_request", "pending"), status("B", "pending")], &names(&["A", "B"])); |
| 314 | assert_eq!(waiting.refusal().unwrap(), "The required checks A and B are still running."); |
| 315 | assert!(WorkflowFacts::of(&[status("Docs", "success")], &names(&["Docs"])).refusal().is_none()); |
| 316 | } |
| 317 | |
| 318 | #[test] |
| 319 | fn a_required_check_nothing_reported_holds_a_merge() { |
| 320 | let facts = WorkflowFacts::of(&[status("CI / pull_request", "success")], &names(&["CI", "Deploy"])); |
| 321 | assert_eq!(facts.expected(), ["Deploy"]); |
| 322 | assert_eq!(facts.refusal().unwrap(), "The required check Deploy has not reported on this commit yet."); |
| 323 | } |
| 324 | |
| 325 | #[test] |
| 326 | fn seen_checks_are_named_once_with_their_events() { |
| 327 | let rows = vec![ |
| 328 | ("CI / push".to_owned(), "2026-10-01T00:00:00Z".to_owned()), |
| 329 | ("CI / pull_request".to_owned(), "2026-10-03T00:00:00Z".to_owned()), |
| 330 | ("g1t / deploy".to_owned(), "2026-10-02T00:00:00Z".to_owned()), |
| 331 | ]; |
| 332 | let seen = seen(rows); |
| 333 | assert_eq!(seen.len(), 2); |
| 334 | assert_eq!(seen[0].name, "CI"); |
| 335 | assert_eq!(seen[0].events, ["pull_request", "push"]); |
| 336 | assert_eq!(seen[0].last_seen, "2026-10-03T00:00:00Z"); |
| 337 | assert_eq!(seen[1].name, "g1t / deploy"); |
| 338 | assert!(seen[1].events.is_empty()); |
| 339 | } |
| 340 | } |