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/statuses.rs

199 lines7,548 bytesCodeBlame

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.

GitHub Actions on g1t, part two: running workflows1//! Statuses on commits: what workflow runs say about a pull request's head.
2//! A pending status holds the pull request, a failed one sends its agent
3//! back (or, for anyone else's, refuses the merge), as acceptance checks do.
4
5use g1t_contracts::events::ChecksEvent;
6use g1t_contracts::time::rfc3339;
7use g1t_contracts::work::{CommitStatus, SetCommitStatusArgs};
8use g1t_contracts::{FailureCode, Outcome};
9use g1t_kit::now_ms;
10use serde::Deserialize;
11use worker::Result;
12
13use crate::Work;
14
15#[derive(Deserialize)]
16struct StatusRow {
17 context: String,
18 state: String,
19 description: Option<String>,
20 target_url: Option<String>,
21 updated_at: String,
22}
23
24impl From<StatusRow> for CommitStatus {
25 fn from(row: StatusRow) -> Self {
26 CommitStatus {
27 context: row.context,
28 state: row.state,
29 description: row.description,
30 target_url: row.target_url,
31 updated_at: row.updated_at,
32 }
33 }
34}
35
36#[derive(Deserialize)]
37struct HeadRow {
38 id: String,
39 number: u32,
40}
41
42/// The workflows still running and the ones that failed, by name.
43#[derive(Clone, Debug, Default, PartialEq, Eq)]
44pub(crate) struct WorkflowFacts {
45 pub(crate) pending: Vec<String>,
46 pub(crate) failed: Vec<String>,
47}
48
49impl WorkflowFacts {
50 pub(crate) fn of(statuses: &[CommitStatus]) -> WorkflowFacts {
51 WorkflowFacts {
52 pending: statuses.iter().filter(|s| s.state == "pending").map(|s| s.context.clone()).collect(),
53 failed: statuses.iter().filter(|s| s.state == "failure" || s.state == "error").map(|s| s.context.clone()).collect(),
54 }
55 }
56
57 /// Why a merge has to wait, if it does.
58 pub(crate) fn refusal(&self) -> Option<String> {
59 if !self.failed.is_empty() {
60 return Some(format!("{} failed.", list(&self.failed)));
61 }
62 if !self.pending.is_empty() {
63 return Some(format!("{} {} still running.", list(&self.pending), if self.pending.len() == 1 { "is" } else { "are" }));
64 }
65 None
66 }
67}
68
69pub(crate) fn list(names: &[String]) -> String {
70 match names {
71 [] => String::new(),
72 [one] => one.clone(),
73 [rest @ .., last] => format!("{} and {last}", rest.join(", ")),
74 }
75}
76
77impl Work {
78 pub(crate) async fn statuses(&self, repo_id: &str, sha: Option<&str>) -> Result<Vec<CommitStatus>> {
79 let Some(sha) = sha else { return Ok(Vec::new()) };
80 Ok(self
81 .db
82 .prepare("SELECT context, state, description, target_url, updated_at FROM commit_statuses WHERE repo_id = ? AND sha = ? ORDER BY context")
83 .bind(&[repo_id.into(), sha.into()])?
84 .all()
85 .await?
86 .results::<StatusRow>()?
87 .into_iter()
88 .map(CommitStatus::from)
89 .collect())
90 }
91
92 pub(crate) async fn set_commit_status(&self, a: SetCommitStatusArgs) -> Result<Outcome<bool>> {
93 if !matches!(a.state.as_str(), "pending" | "success" | "failure" | "error") {
94 return Ok(Outcome::fail(FailureCode::Invalid, "`state` is pending, success, failure or error."));
95 }
96 self.db
97 .prepare(
98 "INSERT INTO commit_statuses (repo_id, sha, context, state, description, target_url, updated_at)
99 VALUES (?, ?, ?, ?, ?, ?, ?)
100 ON CONFLICT (repo_id, sha, context) DO UPDATE SET
101 state = excluded.state, description = excluded.description,
102 target_url = excluded.target_url, updated_at = excluded.updated_at",
103 )
104 .bind(&[
105 a.repo_id.as_str().into(),
106 a.sha.as_str().into(),
107 a.context.as_str().into(),
108 a.state.as_str().into(),
109 a.description.as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
110 a.target_url.as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
111 rfc3339(now_ms()).into(),
112 ])?
113 .run()
114 .await?;
115 if a.state == "pending" {
116 return Ok(Outcome::Ok(true));
117 }
118 // Once every workflow on a pull request's head has finished, its
119 // lifecycle moves on, as it does when its checks finish.
120 let facts = WorkflowFacts::of(&self.statuses(&a.repo_id, Some(&a.sha)).await?);
121 if !facts.pending.is_empty() {
122 return Ok(Outcome::Ok(true));
123 }
Sidebar: the panels really slide124 // A merge queue state waiting on its merge_group workflows.
125 self.merge_group_finished(&a.repo_id, &a.sha, &facts.failed).await?;
GitHub Actions on g1t, part two: running workflows126 let heads = self
127 .db
128 .prepare("SELECT id, number FROM pulls WHERE repo_id = ? AND head_commit = ? AND status IN ('draft', 'open')")
129 .bind(&[a.repo_id.as_str().into(), a.sha.as_str().into()])?
130 .all()
131 .await?
132 .results::<HeadRow>()?;
133 for head in heads {
A stalled pull request picks back up when its workflows pass134 // A pull request g1t stopped on picks back up once what stopped
135 // it passes: its workflows, and its checks if it has any. The
136 // lifecycle then decides again, within its usual limits.
137 if facts.failed.is_empty() {
138 let resumed = self
139 .db
140 .prepare(
141 "UPDATE pulls SET stalled = NULL WHERE id = ? AND managed = 1 AND stalled IS NOT NULL
142 AND (check_status IS NULL OR check_status = 'passed') RETURNING id AS value",
143 )
144 .bind(&[head.id.as_str().into()])?
145 .first::<crate::rows::ValueRow>(None)
146 .await?;
147 if resumed.is_some() {
148 self.note(
149 &a.repo_id,
150 head.number,
151 (crate::lifecycle::POLICY_ACTOR_ID, crate::lifecycle::POLICY_ACTOR_NAME),
152 "picked this back up: its workflows pass now",
153 )
154 .await?;
155 }
156 }
GitHub Actions on g1t, part two: running workflows157 self.publish_as(
158 "checks.completed",
159 &a.repo_id,
160 None,
161 ChecksEvent {
162 pull_id: head.id,
163 repo_id: a.repo_id.clone(),
164 number: head.number,
165 status: if facts.failed.is_empty() { "passed" } else { "failed" },
166 commit: a.sha.clone(),
167 },
168 )
169 .await?;
170 }
171 Ok(Outcome::Ok(true))
172 }
173}
174
175#[cfg(test)]
176mod tests {
177 use super::*;
178
179 fn status(context: &str, state: &str) -> CommitStatus {
180 CommitStatus {
181 context: context.into(),
182 state: state.into(),
183 description: None,
184 target_url: None,
185 updated_at: String::new(),
186 }
187 }
188
189 #[test]
190 fn failures_come_before_waiting() {
191 let facts = WorkflowFacts::of(&[status("CI / push", "pending"), status("Lint / pull_request", "failure"), status("Docs", "success")]);
192 assert_eq!(facts.pending, ["CI / push"]);
193 assert_eq!(facts.failed, ["Lint / pull_request"]);
194 assert_eq!(facts.refusal().unwrap(), "Lint / pull_request failed.");
195 let waiting = WorkflowFacts::of(&[status("A", "pending"), status("B", "pending")]);
196 assert_eq!(waiting.refusal().unwrap(), "A and B are still running.");
197 assert!(WorkflowFacts::of(&[status("A", "success")]).refusal().is_none());
198 }
199}