g1t/crates/runner/src/queue.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.
| Agents as a team: lifecycle, merge queue, billing and a new shell | 1 | //! 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 | ||
| 26 | use std::path::Path; | |
| 27 | use std::process::Command; | |
| 28 | ||
| 29 | use anyhow::{Context, Result, bail}; | |
| 30 | use serde::Deserialize; | |
| 31 | ||
| 32 | use crate::checks::{redact, run_command}; | |
| 33 | use crate::{WORKDIR, auth_option, env, git}; | |
| 34 | ||
| 35 | #[derive(Deserialize)] | |
| 36 | struct 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. | |
| 45 | enum 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 | ||
| 52 | fn 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 | ||
| 60 | fn 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. | |
| 68 | fn 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 | ||
| 125 | fn 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 | ||
| 133 | pub 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, | |
| Agents and memory, checks and conflicts, profiles, slug renames, custom domains | 166 | "conflicts": files, |
| Agents as a team: lifecycle, merge queue, billing and a new shell | 167 | })); |
| 168 | } | |
| 169 | Err(Stopped::Failed(error)) => { | |
| 170 | return Ok(serde_json::json!({ | |
| 171 | "error": redact(&format!("{error:#}"), &secrets), | |
| 172 | })); | |
| 173 | } | |
| 174 | }; | |
| 175 | let contract: Vec<String> = std::env::var("CONTRACT_CHECKS") | |
| 176 | .ok() | |
| 177 | .and_then(|json| serde_json::from_str(&json).ok()) | |
| 178 | .unwrap_or_default(); | |
| 179 | let workdir = Path::new(WORKDIR); | |
| 180 | let mut results: Vec<_> = commands | |
| 181 | .iter() | |
| 182 | .map(|command| run_command(command, workdir, &secrets)) | |
| 183 | .collect(); | |
| 184 | let contract_results: Vec<_> = contract | |
| 185 | .iter() | |
| 186 | .map(|command| run_command(command, workdir, &secrets)) | |
| 187 | .collect(); | |
| 188 | // A contract check that fails here may have been failing on the base | |
| 189 | // already: try those on the base alone. | |
| 190 | let failing: Vec<usize> = (0..contract_results.len()) | |
| 191 | .filter(|&index| !contract_results[index].passed) | |
| 192 | .collect(); | |
| 193 | let mut contract_results = contract_results; | |
| 194 | if !failing.is_empty() { | |
| 195 | let base = env("BASE_COMMIT")?; | |
| 196 | let combined_head = git(workdir, &["rev-parse", "HEAD"])?.trim().to_owned(); | |
| 197 | git(workdir, &["checkout", "--quiet", "--detach", &base])?; | |
| 198 | for index in failing { | |
| 199 | let on_base = run_command(&contract_results[index].command, workdir, &secrets); | |
| 200 | if !on_base.passed { | |
| 201 | let result = &mut contract_results[index]; | |
| 202 | result.passed = true; | |
| 203 | result.command = format!( | |
| 204 | "{} (already failing on the default branch; not held against this)", | |
| 205 | result.command | |
| 206 | ); | |
| 207 | } | |
| 208 | } | |
| 209 | git(workdir, &["checkout", "--quiet", &combined_head])?; | |
| 210 | } | |
| 211 | results.extend(contract_results); | |
| 212 | // Pushed whether or not it passed, so a failure can be looked at. | |
| 213 | let branch = env("QUEUE_BRANCH")?; | |
| 214 | let base_remote = env("BASE_REMOTE")?; | |
| 215 | git( | |
| 216 | Path::new(WORKDIR), | |
| 217 | &[ | |
| 218 | "-c", | |
| 219 | &auth, | |
| 220 | "push", | |
| 221 | "--quiet", | |
| 222 | "--force", | |
| 223 | &base_remote, | |
| 224 | &format!("HEAD:refs/heads/{branch}"), | |
| 225 | ], | |
| 226 | ) | |
| 227 | .map_err(|error| anyhow::anyhow!("{}", redact(&format!("{error:#}"), &secrets))) | |
| 228 | .context("could not push the tested state")?; | |
| 229 | Ok(serde_json::json!({ "combinedCommit": combined, "results": results })) | |
| 230 | })(); | |
| 231 | let body = result.unwrap_or_else(|error| { | |
| 232 | serde_json::json!({ "error": redact(&format!("{error:#}"), &secrets) }) | |
| 233 | }); | |
| 234 | match report(&api, &entry, &token, body) { | |
| 235 | Ok(()) => 0, | |
| 236 | Err(error) => { | |
| 237 | eprintln!("g1t-runner: {error:#}"); | |
| 238 | 1 | |
| 239 | } | |
| 240 | } | |
| 241 | } |