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/crates/runner/src/queue.rs

240 lines9,801 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.

Agents as a team: lifecycle, merge queue, billing and a new shell1//! Builds one state of a merge queue and checks it: the default branch with
2//! a run of pull requests merged in, in queue order, tested with all of
3//! their acceptance checks. The tested state is pushed to its own branch,
4//! from which g1t lands it if it passed and everything ahead of it has
5//! landed.
6//!
7//! No agent runs here. A merge that does not apply cleanly is reported as a
8//! conflict, naming the pull request ahead whose change it collided with;
9//! resolving it is the job of the pull request's own agent once that one
10//! has landed.
11//!
12//! Configuration comes from the environment:
13//!
14//! - `G1T_API`, `QUEUE_ENTRY`, `QUEUE_TOKEN`: where and how to report.
15//! - `G1T_USER`, `G1T_TOKEN`: a member, to read the changes and push.
16//! - `BASE_REMOTE`, `BASE_COMMIT`: the repository and the commit to build on.
17//! - `QUEUE_BRANCH`: where to push the tested state.
18//! - `STACK`: the pull requests to merge, as JSON `[{number, title, remote,
19//! branch, commit}]`, the entry being tested last.
20//! - `CHECKS`: the stack's own acceptance checks, as a JSON array.
21//! - `CONTRACT_CHECKS`: the checks of issues already completed, which the
22//! default branch has to keep passing. One that fails is run again on the
23//! base alone; if it fails there too, it was broken already and is not
24//! held against the stack.
25
26use std::path::Path;
27use std::process::Command;
28
29use anyhow::{Context, Result, bail};
30use serde::Deserialize;
31
32use crate::checks::{redact, run_command};
33use crate::{WORKDIR, auth_option, env, git};
34
35#[derive(Deserialize)]
36struct Item {
37 number: u32,
38 title: String,
39 remote: String,
40 branch: String,
41 commit: String,
42}
43
44/// What stopped a state from being built.
45enum Stopped {
46 /// Merging this pull request's change did not apply cleanly. The
47 /// pull request ahead whose change it collided with, if one did.
48 Conflict { number: u32, with: Option<u32>, files: Vec<String> },
49 Failed(anyhow::Error),
50}
51
52fn git_ok(dir: &Path, args: &[&str]) -> bool {
53 Command::new("git")
54 .current_dir(dir)
55 .args(args)
56 .output()
57 .is_ok_and(|output| output.status.success())
58}
59
60fn changed_files(dir: &Path, from: &str, to: &str) -> Vec<String> {
61 git(dir, &["diff", "--name-only", from, to])
62 .map(|out| out.lines().map(str::to_owned).collect())
63 .unwrap_or_default()
64}
65
66/// Clones the base, then merges each pull request in order. Returns the
67/// tested state's commit.
68fn build(stack: &[Item], auth: &str) -> std::result::Result<String, Stopped> {
69 let base_remote = env("BASE_REMOTE").map_err(Stopped::Failed)?;
70 let base = env("BASE_COMMIT").map_err(Stopped::Failed)?;
71 std::fs::create_dir_all("/work").map_err(|error| Stopped::Failed(error.into()))?;
72 let workdir = Path::new(WORKDIR);
73 git(
74 Path::new("/work"),
75 &["-c", auth, "clone", "--quiet", &base_remote, WORKDIR],
76 )
77 .and_then(|_| git(workdir, &["checkout", "--quiet", "-B", "g1t-queue", &base]))
78 .and_then(|_| git(workdir, &["config", "user.name", "g1t merge queue"]))
79 .and_then(|_| git(workdir, &["config", "user.email", "queue@g1t.sh"]))
80 .map_err(Stopped::Failed)?;
81
82 for (index, item) in stack.iter().enumerate() {
83 git(
84 workdir,
85 &["-c", auth, "fetch", "--quiet", &item.remote, &item.branch],
86 )
87 .with_context(|| format!("could not fetch #{}", item.number))
88 .map_err(Stopped::Failed)?;
89 // The branch may have moved on since the queue looked; what was
90 // queued is the commit, and it must be there.
91 if !git_ok(workdir, &["cat-file", "-e", &format!("{}^{{commit}}", item.commit)]) {
92 return Err(Stopped::Failed(anyhow::anyhow!(
93 "#{}'s commit {} is no longer on its branch",
94 item.number,
95 &item.commit[..item.commit.len().min(12)]
96 )));
97 }
98 let message = format!("Merge #{}: {}", item.number, item.title);
99 if !git_ok(workdir, &["merge", "--quiet", "--no-edit", "-m", &message, &item.commit]) {
100 let files: Vec<String> = git(workdir, &["diff", "--name-only", "--diff-filter=U"])
101 .map(|out| out.lines().map(str::to_owned).collect())
102 .unwrap_or_default();
103 let _ = git(workdir, &["merge", "--abort"]);
104 // The pull request ahead that changed one of the same files.
105 let with = stack[..index]
106 .iter()
107 .rev()
108 .find(|earlier| {
109 let theirs = changed_files(workdir, &base, &earlier.commit);
110 files.iter().any(|file| theirs.contains(file))
111 })
112 .map(|earlier| earlier.number);
113 return Err(Stopped::Conflict {
114 number: item.number,
115 with,
116 files,
117 });
118 }
119 }
120 git(workdir, &["rev-parse", "HEAD"])
121 .map(|out| out.trim().to_owned())
122 .map_err(Stopped::Failed)
123}
124
125fn report(api: &str, entry: &str, token: &str, mut body: serde_json::Value) -> Result<()> {
126 body["token"] = token.into();
127 ureq::post(&format!("{api}/queue/{entry}"))
128 .send_json(body)
129 .context("could not report the merge queue's result")?;
130 Ok(())
131}
132
133pub fn main() -> i32 {
134 let (api, entry, token) = match (env("G1T_API"), env("QUEUE_ENTRY"), env("QUEUE_TOKEN")) {
135 (Ok(api), Ok(entry), Ok(token)) => (api, entry, token),
136 _ => {
137 eprintln!("g1t-runner: G1T_API, QUEUE_ENTRY and QUEUE_TOKEN must be set");
138 return 2;
139 }
140 };
141 let secrets: Vec<String> = ["G1T_TOKEN", "QUEUE_TOKEN"]
142 .iter()
143 .filter_map(|name| std::env::var(name).ok())
144 .filter(|secret| !secret.is_empty())
145 .collect();
146 let result = (|| -> Result<serde_json::Value> {
147 let auth = auth_option(&env("G1T_USER")?, &env("G1T_TOKEN")?);
148 let stack: Vec<Item> = serde_json::from_str(&env("STACK")?).context("STACK is not valid")?;
149 if stack.is_empty() {
150 bail!("nothing to merge");
151 }
152 let commands: Vec<String> = std::env::var("CHECKS")
153 .ok()
154 .and_then(|json| serde_json::from_str(&json).ok())
155 .unwrap_or_default();
156 let combined = match build(&stack, &auth) {
157 Ok(combined) => combined,
158 Err(Stopped::Conflict { number, with, files }) => {
159 let against = with.map_or("the default branch".to_owned(), |n| format!("#{n}"));
160 return Ok(serde_json::json!({
161 "error": format!(
162 "#{number} does not merge cleanly with {against}: {} conflict.",
163 files.join(", ")
164 ),
165 "conflictWith": with,
166 }));
167 }
168 Err(Stopped::Failed(error)) => {
169 return Ok(serde_json::json!({
170 "error": redact(&format!("{error:#}"), &secrets),
171 }));
172 }
173 };
174 let contract: Vec<String> = std::env::var("CONTRACT_CHECKS")
175 .ok()
176 .and_then(|json| serde_json::from_str(&json).ok())
177 .unwrap_or_default();
178 let workdir = Path::new(WORKDIR);
179 let mut results: Vec<_> = commands
180 .iter()
181 .map(|command| run_command(command, workdir, &secrets))
182 .collect();
183 let contract_results: Vec<_> = contract
184 .iter()
185 .map(|command| run_command(command, workdir, &secrets))
186 .collect();
187 // A contract check that fails here may have been failing on the base
188 // already: try those on the base alone.
189 let failing: Vec<usize> = (0..contract_results.len())
190 .filter(|&index| !contract_results[index].passed)
191 .collect();
192 let mut contract_results = contract_results;
193 if !failing.is_empty() {
194 let base = env("BASE_COMMIT")?;
195 let combined_head = git(workdir, &["rev-parse", "HEAD"])?.trim().to_owned();
196 git(workdir, &["checkout", "--quiet", "--detach", &base])?;
197 for index in failing {
198 let on_base = run_command(&contract_results[index].command, workdir, &secrets);
199 if !on_base.passed {
200 let result = &mut contract_results[index];
201 result.passed = true;
202 result.command = format!(
203 "{} (already failing on the default branch; not held against this)",
204 result.command
205 );
206 }
207 }
208 git(workdir, &["checkout", "--quiet", &combined_head])?;
209 }
210 results.extend(contract_results);
211 // Pushed whether or not it passed, so a failure can be looked at.
212 let branch = env("QUEUE_BRANCH")?;
213 let base_remote = env("BASE_REMOTE")?;
214 git(
215 Path::new(WORKDIR),
216 &[
217 "-c",
218 &auth,
219 "push",
220 "--quiet",
221 "--force",
222 &base_remote,
223 &format!("HEAD:refs/heads/{branch}"),
224 ],
225 )
226 .map_err(|error| anyhow::anyhow!("{}", redact(&format!("{error:#}"), &secrets)))
227 .context("could not push the tested state")?;
228 Ok(serde_json::json!({ "combinedCommit": combined, "results": results }))
229 })();
230 let body = result.unwrap_or_else(|error| {
231 serde_json::json!({ "error": redact(&format!("{error:#}"), &secrets) })
232 });
233 match report(&api, &entry, &token, body) {
234 Ok(()) => 0,
235 Err(error) => {
236 eprintln!("g1t-runner: {error:#}");
237 1
238 }
239 }
240}