g1t/crates/runner/src/backup.rs

445 lines18,437 bytesCodeBlame
1//! Cuts one repository's nightly backup: a `git bundle` of everything it
2//! has, or of what is new since the last one, sent to g1t in parts
3//! (`g1t_contracts::backups` has the flow; services/repos/src/backups.rs
4//! keeps the chain).
5//!
6//! Nothing here is an agent and nothing is pushed. The sandbox holds the
7//! job's token and nothing else: it asks for the job, which comes with a
8//! read-only credential for the repository in the git store that lasts
9//! minutes, clones every ref (`--mirror`, never shallow), and bundles
10//! `--all` but the commits the last bundle ended at, and everything they
11//! reach. Those are written as refs of their own first, so a repository
12//! with thousands of refs never makes a command line too long.
13//!
14//! Configuration comes from the environment:
15//!
16//! - `G1T_API`: where to report.
17//! - `BACKUP_JOB`, `BACKUP_TOKEN`: the job, and the token that does it.
18
19use std::collections::BTreeMap;
20use std::fs::File;
21use std::io::{Read, Write};
22use std::path::{Path, PathBuf};
23use std::process::{Command, Stdio};
24use std::time::Duration;
25
26use anyhow::{Context, Result, bail};
27use serde::Deserialize;
28use serde_json::{Value, json};
29use sha2::{Digest, Sha256};
30
31use crate::checks::redact;
32use crate::env;
33
34/// The header the job's token goes in (`g1t_contracts::backups::TOKEN_HEADER`).
35const TOKEN_HEADER: &str = "x-g1t-backup-token";
36/// Where the prerequisites are written as refs while the bundle is cut.
37const PREREQ_REFS: &str = "refs/g1t-backup-prerequisites";
38/// A transfer slower than this many bytes a second for `LOW_SPEED_SECONDS`
39/// is given up, so a stalled clone does not hold the sandbox for hours.
40const LOW_SPEED_BYTES: &str = "1000";
41const LOW_SPEED_SECONDS: &str = "120";
42const PART_TRIES: u32 = 3;
43
44/// The job, as the API gives it (`BackupSpec`).
45#[derive(Debug, Deserialize)]
46struct Spec {
47 kind: String,
48 remote: String,
49 git_token: String,
50 #[serde(default)]
51 prerequisites: Vec<String>,
52 #[serde(default)]
53 previous_refs: BTreeMap<String, String>,
54 part_bytes: u64,
55}
56
57/// What was cut.
58#[derive(Debug, PartialEq, Eq)]
59pub(crate) struct Cut {
60 /// Every ref the clone has, and `HEAD`.
61 pub refs: BTreeMap<String, String>,
62 /// The bundle, or None when there was nothing new to put in one.
63 pub bundle: Option<PathBuf>,
64}
65
66/// Runs git in `dir`, with `stdin` given to it, and returns its trimmed
67/// output, failing on a non-zero exit with what it said.
68fn git_in(dir: &Path, args: &[&str], stdin: Option<&str>) -> Result<String> {
69 let mut child = Command::new("git")
70 .current_dir(dir)
71 .args(args)
72 .stdin(if stdin.is_some() { Stdio::piped() } else { Stdio::null() })
73 .stdout(Stdio::piped())
74 .stderr(Stdio::piped())
75 .spawn()
76 .context("could not run git")?;
77 if let (Some(input), Some(mut pipe)) = (stdin, child.stdin.take()) {
78 pipe.write_all(input.as_bytes())?;
79 }
80 let output = child.wait_with_output()?;
81 if !output.status.success() {
82 bail!(
83 "git {} failed: {}",
84 args.iter().find(|arg| !arg.starts_with('-') && !arg.contains('=')).unwrap_or(&""),
85 String::from_utf8_lossy(&output.stderr).trim()
86 );
87 }
88 Ok(String::from_utf8_lossy(&output.stdout).trim().to_owned())
89}
90
91/// `git for-each-ref` output, `<hash> <name>` a line, as a map.
92pub(crate) fn parse_refs(listing: &str) -> BTreeMap<String, String> {
93 listing
94 .lines()
95 .filter_map(|line| line.trim().split_once(' '))
96 .filter(|(_, name)| !name.starts_with(PREREQ_REFS))
97 .map(|(hash, name)| (name.trim().to_owned(), hash.trim().to_owned()))
98 .collect()
99}
100
101/// Of `git cat-file --batch-check` output, the objects it found.
102pub(crate) fn present(check: &str) -> Vec<String> {
103 check
104 .lines()
105 .filter(|line| !line.ends_with(" missing"))
106 .filter_map(|line| line.split(' ').next())
107 .filter(|hash| !hash.is_empty())
108 .map(str::to_owned)
109 .collect()
110}
111
112/// Whether git refused because the bundle would hold no objects: every
113/// ref still points where the last bundle left it, or at what it reaches.
114pub(crate) fn is_empty_bundle(error: &str) -> bool {
115 error.contains("empty bundle")
116}
117
118/// Every ref of the clone in `dir`, and `HEAD` when it points somewhere.
119pub(crate) fn refs_of(dir: &Path) -> Result<BTreeMap<String, String>> {
120 let mut refs = parse_refs(&git_in(dir, &["for-each-ref", "--format=%(objectname) %(refname)"], None)?);
121 if let Ok(head) = git_in(dir, &["rev-parse", "--verify", "--quiet", "HEAD"], None)
122 && !head.is_empty()
123 {
124 refs.insert("HEAD".to_owned(), head);
125 }
126 Ok(refs)
127}
128
129/// Cuts the bundle of the clone in `dir` into `out`: every ref, leaving
130/// out `prerequisites` (those the clone still has) and all they reach.
131/// Nothing is cut when the refs are `previous` exactly, or there are none,
132/// or nothing new is there.
133pub(crate) fn cut(dir: &Path, prerequisites: &[String], previous: &BTreeMap<String, String>, out: &Path) -> Result<Cut> {
134 let refs = refs_of(dir)?;
135 if refs.is_empty() || (&refs == previous && !previous.is_empty() && !prerequisites.is_empty()) {
136 return Ok(Cut { refs, bundle: None });
137 }
138 // Only those the clone has: a commit force-pushed away is no longer
139 // there to leave out, and the bundle then carries a little more.
140 let kept = if prerequisites.is_empty() {
141 Vec::new()
142 } else {
143 present(&git_in(dir, &["cat-file", "--batch-check=%(objectname) %(objecttype)"], Some(&format!("{}\n", prerequisites.join("\n"))))?)
144 };
145 if !kept.is_empty() {
146 let updates: String = kept.iter().map(|hash| format!("create {PREREQ_REFS}/{hash} {hash}\n")).collect();
147 git_in(dir, &["update-ref", "--stdin"], Some(&updates))?;
148 }
149 let out_text = out.to_str().context("the bundle's path is not text")?;
150 let glob = format!("--glob={PREREQ_REFS}/*");
151 let exclude = format!("--exclude={PREREQ_REFS}/*");
152 let mut args = vec!["bundle", "create", "--quiet", out_text, exclude.as_str(), "--all"];
153 if !kept.is_empty() {
154 args.extend(["--not", glob.as_str()]);
155 }
156 let made = git_in(dir, &args, None);
157 if !kept.is_empty() {
158 let deletes: String = kept.iter().map(|hash| format!("delete {PREREQ_REFS}/{hash}\n")).collect();
159 git_in(dir, &["update-ref", "--stdin"], Some(&deletes))?;
160 }
161 match made {
162 Ok(_) => {
163 git_in(dir, &["bundle", "verify", "--quiet", out_text], None).context("the bundle does not verify")?;
164 Ok(Cut { refs, bundle: Some(out.to_owned()) })
165 }
166 Err(error) if is_empty_bundle(&error.to_string()) => Ok(Cut { refs, bundle: None }),
167 Err(error) => Err(error),
168 }
169}
170
171/// The bytes of the clone's packs, which is what it read from the store.
172fn pack_bytes(dir: &Path) -> u64 {
173 std::fs::read_dir(dir.join("objects/pack"))
174 .map(|entries| entries.filter_map(|entry| entry.ok()?.metadata().ok()).map(|meta| meta.len()).sum())
175 .unwrap_or(0)
176}
177
178/// Talks to the API about one job.
179struct Job {
180 api: String,
181 id: String,
182 token: String,
183 agent: ureq::Agent,
184}
185
186impl Job {
187 fn url(&self, action: &str) -> String {
188 format!("{}/backups/{}/{action}", self.api, self.id)
189 }
190
191 fn answer(result: std::result::Result<ureq::Response, ureq::Error>) -> Result<Value> {
192 match result {
193 Ok(response) => Ok(response.into_json()?),
194 Err(ureq::Error::Status(status, response)) => {
195 let body: Value = response.into_json().unwrap_or(Value::Null);
196 let said = body["error"]["message"].as_str().unwrap_or("no reason given").to_owned();
197 bail!("g1t answered {status}: {said}")
198 }
199 Err(error) => Err(error.into()),
200 }
201 }
202
203 fn post(&self, action: &str, body: Value) -> Result<Value> {
204 Job::answer(self.agent.post(&self.url(action)).set(TOKEN_HEADER, &self.token).send_json(body))
205 }
206
207 fn put_part(&self, number: u16, bytes: &[u8]) -> Result<Value> {
208 let mut last = None;
209 for _ in 0..PART_TRIES {
210 let sent = self
211 .agent
212 .put(&self.url(&format!("parts/{number}")))
213 .set(TOKEN_HEADER, &self.token)
214 .set("content-type", "application/octet-stream")
215 .send_bytes(bytes);
216 match Job::answer(sent) {
217 Ok(part) => return Ok(part),
218 Err(error) => last = Some(error),
219 }
220 }
221 Err(last.unwrap_or_else(|| anyhow::anyhow!("the part was not sent")))
222 }
223}
224
225/// Sends the bundle in parts of `part_bytes`, hashing it on the way, and
226/// returns its size, SHA-256 and the parts as g1t kept them.
227fn send(job: &Job, bundle: &Path, part_bytes: u64) -> Result<(u64, String, Vec<Value>)> {
228 let mut file = File::open(bundle)?;
229 let mut hasher = Sha256::new();
230 let mut parts = Vec::new();
231 let mut size = 0u64;
232 let mut buffer = vec![0u8; part_bytes as usize];
233 loop {
234 let mut filled = 0;
235 while filled < buffer.len() {
236 let read = file.read(&mut buffer[filled..])?;
237 if read == 0 {
238 break;
239 }
240 filled += read;
241 }
242 if filled == 0 {
243 break;
244 }
245 hasher.update(&buffer[..filled]);
246 size += filled as u64;
247 let number = u16::try_from(parts.len() + 1).context("the bundle has too many parts")?;
248 parts.push(job.put_part(number, &buffer[..filled]).with_context(|| format!("could not send part {number}"))?);
249 crate::abuse::touch();
250 if filled < buffer.len() {
251 break;
252 }
253 }
254 Ok((size, hex::encode(hasher.finalize()), parts))
255}
256
257fn back_up(job: &Job, fetched: &mut u64) -> Result<String> {
258 let spec: Spec = serde_json::from_value(job.post("spec", json!({}))?).context("the job's spec could not be read")?;
259 let work = Path::new("/work");
260 std::fs::create_dir_all(work)?;
261 let mirror = work.join("backup.git");
262 let auth = format!("http.extraHeader=Authorization: Bearer {}", spec.git_token);
263 let mirror_text = mirror.to_str().context("the clone's path is not text")?;
264 git_in(
265 work,
266 &[
267 "-c",
268 &auth,
269 "-c",
270 &format!("http.lowSpeedLimit={LOW_SPEED_BYTES}"),
271 "-c",
272 &format!("http.lowSpeedTime={LOW_SPEED_SECONDS}"),
273 "clone",
274 "--mirror",
275 "--quiet",
276 &spec.remote,
277 mirror_text,
278 ],
279 None,
280 )
281 .context("could not clone the repository")?;
282 *fetched = pack_bytes(&mirror).max(1);
283 let cut = cut(&mirror, &spec.prerequisites, &spec.previous_refs, &work.join("backup.bundle"))?;
284 let (size, sha256, parts) = match &cut.bundle {
285 Some(bundle) => {
286 let (size, sha256, parts) = send(job, bundle, spec.part_bytes)?;
287 (size, Some(sha256), parts)
288 }
289 None => (0, None, Vec::new()),
290 };
291 job.post(
292 "complete",
293 json!({ "refs": cut.refs, "size": size, "sha256": sha256, "parts": parts, "fetched_bytes": *fetched }),
294 )
295 .context("could not report the backup")?;
296 Ok(match cut.bundle {
297 Some(_) => format!("{} bundle of {} refs, {size} bytes", spec.kind, cut.refs.len()),
298 None => format!("nothing new to bundle in {} refs", cut.refs.len()),
299 })
300}
301
302pub fn main() -> i32 {
303 let (api, id, token) = match (env("G1T_API"), env("BACKUP_JOB"), env("BACKUP_TOKEN")) {
304 (Ok(api), Ok(id), Ok(token)) => (api, id, token),
305 _ => {
306 eprintln!("g1t-runner: G1T_API, BACKUP_JOB and BACKUP_TOKEN must be set");
307 return 2;
308 }
309 };
310 let agent = ureq::AgentBuilder::new().timeout_connect(Duration::from_secs(30)).timeout(Duration::from_secs(600)).build();
311 let job = Job { api, id, token: token.clone(), agent };
312 let mut fetched = 0;
313 match back_up(&job, &mut fetched) {
314 Ok(said) => {
315 println!("g1t-runner: backed up: {said}");
316 0
317 }
318 Err(error) => {
319 let said = redact(&format!("{error:#}"), &[token]);
320 eprintln!("g1t-runner: the backup failed: {said}");
321 if let Err(error) = job.post("fail", json!({ "error": said, "fetched_bytes": fetched })) {
322 eprintln!("g1t-runner: could not report the failure: {error:#}");
323 }
324 1
325 }
326 }
327}
328
329#[cfg(test)]
330mod tests {
331 use super::*;
332
333 #[test]
334 fn refs_are_read_and_the_prerequisites_own_left_out() {
335 let listing = "c71546fcd893ef8b0f57388b65e620d759705dda refs/heads/main\n\
336 4807077b296e6edbf410d55e72749d3e1170c291 refs/pull/pr_1/head\n\
337 4807077b296e6edbf410d55e72749d3e1170c291 refs/g1t-backup-prerequisites/4807077b296e6edbf410d55e72749d3e1170c291\n";
338 let refs = parse_refs(listing);
339 assert_eq!(refs.len(), 2);
340 assert_eq!(refs["refs/heads/main"], "c71546fcd893ef8b0f57388b65e620d759705dda");
341 assert!(parse_refs("").is_empty());
342 }
343
344 #[test]
345 fn missing_prerequisites_are_dropped() {
346 let check = "c71546fcd893ef8b0f57388b65e620d759705dda commit\n4807077b296e6edbf410d55e72749d3e1170c291 missing\n";
347 assert_eq!(present(check), ["c71546fcd893ef8b0f57388b65e620d759705dda"]);
348 }
349
350 #[test]
351 fn an_empty_bundle_is_told_apart_from_a_failure() {
352 assert!(is_empty_bundle("git bundle failed: fatal: Refusing to create empty bundle."));
353 assert!(!is_empty_bundle("git bundle failed: fatal: bad revision"));
354 }
355
356 // With git itself: a repository backed up full, then incrementally,
357 // then restored from the chain the way the restore drill does it.
358
359 fn scratch(name: &str) -> PathBuf {
360 let dir = std::env::temp_dir().join(format!("g1t-backup-test-{name}-{}", std::process::id()));
361 let _ = std::fs::remove_dir_all(&dir);
362 std::fs::create_dir_all(&dir).unwrap();
363 dir
364 }
365
366 fn commit(dir: &Path, file: &str, text: &str) {
367 std::fs::write(dir.join(file), text).unwrap();
368 git_in(dir, &["add", "--all"], None).unwrap();
369 git_in(dir, &["-c", "user.name=t", "-c", "user.email=t@example.com", "commit", "--quiet", "-m", text], None).unwrap();
370 }
371
372 fn mirror_of(origin: &Path, into: &Path) {
373 let _ = std::fs::remove_dir_all(into);
374 let parent = into.parent().unwrap();
375 git_in(parent, &["clone", "--mirror", "--quiet", origin.to_str().unwrap(), into.to_str().unwrap()], None).unwrap();
376 }
377
378 fn prerequisites(refs: &BTreeMap<String, String>) -> Vec<String> {
379 let unique: std::collections::BTreeSet<&String> = refs.values().collect();
380 unique.into_iter().cloned().collect()
381 }
382
383 #[test]
384 fn a_chain_of_bundles_restores_every_ref() {
385 let root = scratch("chain");
386 let origin = root.join("origin");
387 std::fs::create_dir_all(&origin).unwrap();
388 git_in(&origin, &["init", "--quiet", "--initial-branch=main"], None).unwrap();
389 commit(&origin, "a.txt", "one");
390 git_in(&origin, &["tag", "-a", "v1", "-m", "v1"], None).unwrap();
391 let mirror = root.join("mirror.git");
392
393 // Full.
394 mirror_of(&origin, &mirror);
395 let full = cut(&mirror, &[], &BTreeMap::new(), &root.join("0-full.bundle")).unwrap();
396 assert!(full.bundle.is_some());
397 assert!(full.refs.contains_key("refs/tags/v1") && full.refs.contains_key("HEAD"));
398
399 // Nothing changed: nothing is cut.
400 mirror_of(&origin, &mirror);
401 let same = cut(&mirror, &prerequisites(&full.refs), &full.refs, &root.join("x.bundle")).unwrap();
402 assert_eq!(same.bundle, None);
403
404 // New work, a new branch and a deleted tag: incremental.
405 commit(&origin, "b.txt", "two");
406 git_in(&origin, &["branch", "feature"], None).unwrap();
407 git_in(&origin, &["tag", "-d", "v1"], None).unwrap();
408 mirror_of(&origin, &mirror);
409 let incr = cut(&mirror, &prerequisites(&full.refs), &full.refs, &root.join("1-incr.bundle")).unwrap();
410 let incremental = incr.bundle.clone().expect("new commits make a bundle");
411 assert!(!incr.refs.contains_key("refs/tags/v1"));
412 // It needs the full one: alone, it does not verify.
413 let lone = root.join("lone");
414 git_in(&root, &["init", "--quiet", "--bare", lone.to_str().unwrap()], None).unwrap();
415 assert!(git_in(&lone, &["bundle", "verify", incremental.to_str().unwrap()], None).is_err());
416
417 // Only a ref moved to a commit already kept: no objects, no bundle.
418 git_in(&origin, &["branch", "-f", "feature", "HEAD~1"], None).unwrap();
419 mirror_of(&origin, &mirror);
420 let moved = cut(&mirror, &prerequisites(&incr.refs), &incr.refs, &root.join("2-incr.bundle")).unwrap();
421 assert_eq!(moved.bundle, None);
422 assert_ne!(moved.refs, incr.refs);
423
424 // Restored: each bundle in order, then the refs the last entry says.
425 let restored = root.join("restored.git");
426 git_in(&root, &["init", "--quiet", "--bare", restored.to_str().unwrap()], None).unwrap();
427 for bundle in [full.bundle.unwrap(), incremental] {
428 git_in(&restored, &["bundle", "verify", "--quiet", bundle.to_str().unwrap()], None).unwrap();
429 git_in(&restored, &["fetch", "--quiet", "--no-tags", bundle.to_str().unwrap(), "+refs/*:refs/backup-staging/*"], None).unwrap();
430 }
431 let updates: String = moved
432 .refs
433 .iter()
434 .filter(|(name, _)| name.as_str() != "HEAD")
435 .map(|(name, hash)| format!("update {name} {hash}\n"))
436 .collect();
437 git_in(&restored, &["update-ref", "--stdin"], Some(&updates)).unwrap();
438 let staging: String = git_in(&restored, &["for-each-ref", "--format=delete %(refname)", "refs/backup-staging/"], None).unwrap();
439 git_in(&restored, &["update-ref", "--stdin"], Some(&format!("{staging}\n"))).unwrap();
440 git_in(&restored, &["symbolic-ref", "HEAD", "refs/heads/main"], None).unwrap();
441 assert_eq!(refs_of(&restored).unwrap(), moved.refs);
442 git_in(&restored, &["fsck", "--no-progress", "--connectivity-only"], None).unwrap();
443 let _ = std::fs::remove_dir_all(&root);
444 }
445}