anvilsign in

collin/anvil

main / crates / anvil-git / src / push.rs
1//! Client side of `git push` (the send-pack protocol), built on gix plumbing.
2//!
3//! anvil is pure gitoxide by design — no `git` binary anywhere in the product
4//! (see CLAUDE.md). gix doesn't implement push yet, so this module speaks
5//! receive-pack's wire format directly: parse the ref advertisement, compute
6//! mirror update commands, send them with a pack from
7//! [`gitserver_core::pack::build_raw_pack`], and read the report-status
8//! reply. The serving counterparts live in `gitserver-core`, which makes the
9//! whole path round-trippable in-process for tests.
10//!
11//! Two transports:
12//! - smart HTTP(S) via reqwest (rustls/ring — see the workspace manifest),
13//! with Basic credentials taken from the URL's userinfo
14//! (`https://x-access-token:<token>@github.com/you/repo.git`);
15//! - a local filesystem path, served by `gitserver-core` in-process.
16
17use std::{
18 collections::BTreeMap,
19 path::{
20 Path,
21 PathBuf,
22 },
23};
24
25use gitserver_core::{
26 pack,
27 pktline,
28 receive_pack,
29};
30
31const ZERO_OID: &str = "0000000000000000000000000000000000000000";
32
33/// What a mirror push did. `up_to_date` means no commands were needed.
34#[derive(Debug, Default)]
35pub struct MirrorOutcome {
36 pub updated: usize,
37 pub deleted: usize,
38 /// Deletions the remote should have seen but doesn't support
39 /// (no `delete-refs` capability).
40 pub skipped_deletes: usize,
41 pub up_to_date: bool,
42}
43
44/// Mirror all local `refs/heads/*` and `refs/tags/*` to `remote`: create or
45/// force-update every local ref, delete remote heads/tags with no local
46/// counterpart. (Unlike `git push --mirror`, other remote namespaces — e.g.
47/// GitHub's read-only `refs/pull/*` — are left alone instead of generating
48/// rejected deletes.)
49pub async fn mirror(repo_path: &Path, remote: &str) -> Result<MirrorOutcome, String> {
50 let remote = Remote::parse(remote);
51
52 let advertisement = remote.advertise(repo_path).await?;
53 let (remote_refs, capabilities) = parse_advertisement(&advertisement)?;
54 let can_delete = capabilities.iter().any(|c| c == "delete-refs");
55 let local_refs = local_refs(repo_path)?;
56
57 // Mirror semantics over heads + tags.
58 let mirrored = |name: &str| name.starts_with("refs/heads/") || name.starts_with("refs/tags/");
59 let mut commands: Vec<(String, String, String)> = Vec::new(); // (old, new, name)
60 let mut skipped_deletes = 0;
61 for (name, new) in &local_refs {
62 match remote_refs.get(name) {
63 Some(old) if old == new => {}
64 Some(old) => commands.push((old.clone(), new.clone(), name.clone())),
65 None => commands.push((ZERO_OID.to_string(), new.clone(), name.clone())),
66 }
67 }
68 for (name, old) in remote_refs.iter().filter(|(name, _)| mirrored(name)) {
69 if !local_refs.contains_key(name) {
70 if can_delete {
71 commands.push((old.clone(), ZERO_OID.to_string(), name.clone()));
72 } else {
73 skipped_deletes += 1;
74 }
75 }
76 }
77 if commands.is_empty() {
78 return Ok(MirrorOutcome {
79 up_to_date: true,
80 skipped_deletes,
81 ..Default::default()
82 });
83 }
84
85 // The pack covers what the remote is missing: wants are the new tips,
86 // haves are whatever advertised tips exist in our object database.
87 let repo = gix::open(repo_path).map_err(|e| format!("open {}: {e}", repo_path.display()))?;
88 let parse_oid =
89 |hex: &str| gix::ObjectId::from_hex(hex.as_bytes()).map_err(|e| format!("bad oid: {e}"));
90 let mut wants = Vec::new();
91 for (_, new, _) in &commands {
92 if new != ZERO_OID {
93 wants.push(parse_oid(new)?);
94 }
95 }
96 let mut haves = Vec::new();
97 for old in remote_refs.values() {
98 let oid = parse_oid(old)?;
99 if repo.try_find_object(oid).ok().flatten().is_some() {
100 haves.push(oid);
101 }
102 }
103
104 // Request body: command pkt-lines (capabilities ride on the first one),
105 // flush, then the pack — omitted when every command is a delete.
106 let mut body = Vec::new();
107 for (i, (old, new, name)) in commands.iter().enumerate() {
108 let line = if i == 0 {
109 format!("{old} {new} {name}\0report-status agent=anvil\n")
110 } else {
111 format!("{old} {new} {name}\n")
112 };
113 body.extend_from_slice(&pktline::encode(line.as_bytes()));
114 }
115 body.extend_from_slice(pktline::flush());
116 if !wants.is_empty() {
117 body.extend_from_slice(
118 &pack::build_raw_pack(repo_path, &wants, &haves).map_err(|e| e.to_string())?,
119 );
120 }
121
122 let response = remote.send(repo_path, body).await?;
123 let (updated, deleted) = commands.iter().fold((0, 0), |(u, d), (_, new, _)| {
124 if new == ZERO_OID {
125 (u, d + 1)
126 } else {
127 (u + 1, d)
128 }
129 });
130 check_report_status(&response)?;
131 Ok(MirrorOutcome {
132 updated,
133 deleted,
134 skipped_deletes,
135 up_to_date: false,
136 })
137}
138
139/// Local refs to mirror: full name → target id (tag refs keep the tag
140/// object's id, matching what git pushes for `refs/tags/*`).
141fn local_refs(repo_path: &Path) -> Result<BTreeMap<String, String>, String> {
142 let repo = gix::open(repo_path).map_err(|e| format!("open {}: {e}", repo_path.display()))?;
143 let platform = repo.references().map_err(|e| e.to_string())?;
144 let mut out = BTreeMap::new();
145 for r in platform.all().map_err(|e| e.to_string())?.flatten() {
146 let name = r.name().as_bstr().to_string();
147 if !(name.starts_with("refs/heads/") || name.starts_with("refs/tags/")) {
148 continue;
149 }
150 if let Some(id) = r.try_id() {
151 out.insert(name, id.to_string());
152 }
153 }
154 Ok(out)
155}
156
157/// Parse a receive-pack ref advertisement (v0, with or without the smart-HTTP
158/// `# service=` preamble) into (ref name → oid hex, server capabilities).
159#[allow(clippy::type_complexity)]
160fn parse_advertisement(bytes: &[u8]) -> Result<(BTreeMap<String, String>, Vec<String>), String> {
161 let mut refs = BTreeMap::new();
162 let mut capabilities = Vec::new();
163 for line in pkt_lines(bytes)? {
164 let line = String::from_utf8_lossy(&line);
165 let line = line.trim_end_matches('\n');
166 if line.starts_with('#') {
167 continue; // smart-HTTP service preamble
168 }
169 // The first ref line carries capabilities after a NUL.
170 let (line, caps) = match line.split_once('\0') {
171 Some((line, caps)) => (line, Some(caps)),
172 None => (line, None),
173 };
174 if let Some(caps) = caps {
175 capabilities.extend(caps.split_whitespace().map(str::to_string));
176 }
177 let Some((oid, name)) = line.split_once(' ') else {
178 continue;
179 };
180 // An empty repo advertises `<zero-oid> capabilities^{}`.
181 if oid == ZERO_OID && name == "capabilities^{}" {
182 continue;
183 }
184 // Peeled-tag annotations are informational.
185 if name.ends_with("^{}") {
186 continue;
187 }
188 refs.insert(name.to_string(), oid.to_string());
189 }
190 Ok((refs, capabilities))
191}
192
193/// The data pkt-lines of a buffer (flush/delim packets skipped).
194fn pkt_lines(mut bytes: &[u8]) -> Result<Vec<Vec<u8>>, String> {
195 let mut out = Vec::new();
196 while bytes.len() >= 4 {
197 let len = usize::from_str_radix(
198 std::str::from_utf8(&bytes[..4]).map_err(|_| "bad pkt length")?,
199 16,
200 )
201 .map_err(|_| "bad pkt length")?;
202 if len < 4 {
203 bytes = &bytes[4..]; // flush (0000) / delim (0001) / response-end
204 continue;
205 }
206 if len > bytes.len() {
207 return Err("truncated pkt-line".into());
208 }
209 out.push(bytes[4..len].to_vec());
210 bytes = &bytes[len..];
211 }
212 Ok(out)
213}
214
215/// Fail on any `unpack`/per-ref error in a report-status reply.
216fn check_report_status(bytes: &[u8]) -> Result<(), String> {
217 let mut errors = Vec::new();
218 for line in pkt_lines(bytes)? {
219 let line = String::from_utf8_lossy(&line);
220 let line = line.trim_end();
221 if let Some(rest) = line.strip_prefix("unpack ") {
222 if rest != "ok" {
223 errors.push(format!("unpack failed: {rest}"));
224 }
225 } else if let Some(rest) = line.strip_prefix("ng ") {
226 errors.push(format!("ref rejected: {rest}"));
227 }
228 }
229 if errors.is_empty() {
230 Ok(())
231 } else {
232 Err(errors.join(" / "))
233 }
234}
235
236/// Where a push goes: a smart-HTTP(S) remote, or a bare repository on the
237/// local filesystem (served by `gitserver-core` in-process).
238enum Remote {
239 Http(String),
240 Local(PathBuf),
241}
242
243impl Remote {
244 fn parse(remote: &str) -> Self {
245 if remote.starts_with("http://") || remote.starts_with("https://") {
246 Remote::Http(remote.trim_end_matches('/').to_string())
247 } else {
248 Remote::Local(PathBuf::from(
249 remote.strip_prefix("file://").unwrap_or(remote),
250 ))
251 }
252 }
253
254 async fn advertise(&self, _local: &Path) -> Result<Vec<u8>, String> {
255 match self {
256 Remote::Local(path) => {
257 receive_pack::advertise_receive_refs(path).map_err(|e| e.to_string())
258 }
259 Remote::Http(url) => {
260 let (client, url, auth) = http_parts(url)?;
261 let mut req = client.get(format!("{url}/info/refs?service=git-receive-pack"));
262 if let Some((user, pass)) = &auth {
263 req = req.basic_auth(user, pass.as_deref());
264 }
265 let resp = req.send().await.map_err(|e| e.to_string())?;
266 if !resp.status().is_success() {
267 return Err(format!("advertisement request: HTTP {}", resp.status()));
268 }
269 Ok(resp.bytes().await.map_err(|e| e.to_string())?.to_vec())
270 }
271 }
272 }
273
274 async fn send(&self, _local: &Path, body: Vec<u8>) -> Result<Vec<u8>, String> {
275 match self {
276 Remote::Local(path) => {
277 let path = path.clone();
278 // receive_pack applies the pack synchronously; keep the
279 // executor unblocked.
280 tokio::task::spawn_blocking(move || {
281 receive_pack::receive_pack(&path, std::io::Cursor::new(body))
282 .map_err(|e| e.to_string())
283 })
284 .await
285 .map_err(|e| e.to_string())?
286 }
287 Remote::Http(url) => {
288 let (client, url, auth) = http_parts(url)?;
289 let mut req = client
290 .post(format!("{url}/git-receive-pack"))
291 .header("Content-Type", "application/x-git-receive-pack-request")
292 .body(body);
293 if let Some((user, pass)) = &auth {
294 req = req.basic_auth(user, pass.as_deref());
295 }
296 let resp = req.send().await.map_err(|e| e.to_string())?;
297 if !resp.status().is_success() {
298 return Err(format!("receive-pack request: HTTP {}", resp.status()));
299 }
300 Ok(resp.bytes().await.map_err(|e| e.to_string())?.to_vec())
301 }
302 }
303 }
304}
305
306/// Basic credentials lifted from a URL's userinfo.
307type UrlAuth = Option<(String, Option<String>)>;
308
309/// Split userinfo credentials out of `url` and build a client. The rustls
310/// *ring* provider is installed process-wide on first use (the workspace
311/// deliberately has no default provider — aws-lc-rs breaks the musl build).
312fn http_parts(url: &str) -> Result<(reqwest::Client, String, UrlAuth), String> {
313 static TLS_PROVIDER: std::sync::Once = std::sync::Once::new();
314 TLS_PROVIDER.call_once(|| {
315 let _ = rustls::crypto::ring::default_provider().install_default();
316 });
317
318 let mut parsed = reqwest::Url::parse(url).map_err(|e| format!("bad mirror url: {e}"))?;
319 let auth = (!parsed.username().is_empty() || parsed.password().is_some()).then(|| {
320 (
321 parsed.username().to_string(),
322 parsed.password().map(str::to_string),
323 )
324 });
325 parsed
326 .set_username("")
327 .and_then(|()| parsed.set_password(None))
328 .map_err(|()| "bad mirror url".to_string())?;
329 let client = reqwest::Client::builder()
330 .user_agent("anvil")
331 .build()
332 .map_err(|e| e.to_string())?;
333 Ok((
334 client,
335 parsed.to_string().trim_end_matches('/').to_string(),
336 auth,
337 ))
338}
339
340#[cfg(test)]
341mod tests {
342 use super::*;
343
344 /// Pure-Rust round trip: build commits with the git CLI as a fixture,
345 /// then mirror them through this client into a bare repo served by
346 /// gitserver-core — no git binary anywhere on the push path.
347 #[tokio::test]
348 async fn mirror_creates_updates_and_deletes() {
349 let tmp = tempfile::tempdir().unwrap();
350 let src = tmp.path().join("src.git");
351 let dst = tmp.path().join("dst.git");
352 let work = tmp.path().join("w");
353 gix::init_bare(&dst).unwrap();
354
355 let git = |args: &[&str], dir: &Path| {
356 let out = std::process::Command::new("git")
357 .args(args)
358 .current_dir(dir)
359 .env("GIT_AUTHOR_NAME", "t")
360 .env("GIT_AUTHOR_EMAIL", "t@example.com")
361 .env("GIT_COMMITTER_NAME", "t")
362 .env("GIT_COMMITTER_EMAIL", "t@example.com")
363 .output()
364 .expect("run git");
365 assert!(out.status.success(), "git {args:?}: {out:?}");
366 };
367 git(&["init", "-q", "--bare", src.to_str().unwrap()], tmp.path());
368 git(
369 &["init", "-q", "-b", "main", work.to_str().unwrap()],
370 tmp.path(),
371 );
372 std::fs::write(work.join("f"), "one").unwrap();
373 git(&["add", "."], &work);
374 git(&["commit", "-qm", "c1"], &work);
375 git(&["tag", "-a", "-m", "annotated", "v1"], &work);
376 git(
377 &[
378 "push",
379 "-q",
380 src.to_str().unwrap(),
381 "main",
382 "main:extra",
383 "v1",
384 ],
385 &work,
386 );
387
388 // First mirror: creates two branches and the annotated tag.
389 let outcome = mirror(&src, dst.to_str().unwrap()).await.unwrap();
390 assert_eq!((outcome.updated, outcome.deleted), (3, 0));
391 let heads = gitserver_core::receive_pack::advertise_receive_refs(&dst).unwrap();
392 let (refs, caps) = parse_advertisement(&heads).unwrap();
393 assert!(caps.iter().any(|c| c == "delete-refs"));
394 assert!(refs.contains_key("refs/heads/main"));
395 assert!(refs.contains_key("refs/heads/extra"));
396 assert!(refs.contains_key("refs/tags/v1"));
397
398 // Mirroring again is a no-op.
399 assert!(
400 mirror(&src, dst.to_str().unwrap())
401 .await
402 .unwrap()
403 .up_to_date
404 );
405
406 // New commit on main + branch deletion both propagate.
407 std::fs::write(work.join("f"), "two").unwrap();
408 git(&["commit", "-qam", "c2"], &work);
409 git(
410 &["push", "-q", src.to_str().unwrap(), "main", ":extra"],
411 &work,
412 );
413 let outcome = mirror(&src, dst.to_str().unwrap()).await.unwrap();
414 assert_eq!((outcome.updated, outcome.deleted), (1, 1));
415 let (refs, _) = parse_advertisement(
416 &gitserver_core::receive_pack::advertise_receive_refs(&dst).unwrap(),
417 )
418 .unwrap();
419 assert!(!refs.contains_key("refs/heads/extra"));
420 assert_eq!(
421 refs["refs/heads/main"],
422 parse_advertisement(
423 &gitserver_core::receive_pack::advertise_receive_refs(&src).unwrap()
424 )
425 .unwrap()
426 .0["refs/heads/main"],
427 "destination main matches source after the second mirror"
428 );
429
430 // The mirrored objects are genuinely usable on the destination side:
431 // read the new file content back out of dst with gix.
432 let dst_repo = gix::open(&dst).unwrap();
433 let head = dst_repo
434 .rev_parse_single("refs/heads/main")
435 .unwrap()
436 .object()
437 .unwrap()
438 .peel_to_commit()
439 .unwrap();
440 let blob = head
441 .tree()
442 .unwrap()
443 .lookup_entry_by_path("f")
444 .unwrap()
445 .unwrap()
446 .object()
447 .unwrap();
448 assert_eq!(blob.data.as_slice(), b"two");
449 }
450
451 #[test]
452 fn advertisement_parsing_handles_preamble_and_empty() {
453 // Smart-HTTP preamble + first-line capabilities.
454 let mut bytes = Vec::new();
455 bytes.extend_from_slice(&pktline::encode_comment("service=git-receive-pack"));
456 bytes.extend_from_slice(pktline::flush());
457 bytes.extend_from_slice(&pktline::encode(
458 b"1111111111111111111111111111111111111111 refs/heads/main\0report-status\n",
459 ));
460 bytes.extend_from_slice(&pktline::encode(
461 b"2222222222222222222222222222222222222222 refs/tags/v1\n",
462 ));
463 bytes.extend_from_slice(pktline::flush());
464 let (refs, caps) = parse_advertisement(&bytes).unwrap();
465 assert!(caps.iter().any(|c| c == "report-status"));
466 assert_eq!(refs.len(), 2);
467 assert_eq!(refs["refs/heads/main"], "1".repeat(40));
468
469 // Empty repository advertisement.
470 let empty =
471 pktline::encode(format!("{ZERO_OID} capabilities^{{}}\0report-status\n").as_bytes());
472 assert!(parse_advertisement(&empty).unwrap().0.is_empty());
473 }
474
475 #[test]
476 fn report_status_failures_surface() {
477 let mut ok = Vec::new();
478 ok.extend_from_slice(&pktline::encode(b"unpack ok\n"));
479 ok.extend_from_slice(&pktline::encode(b"ok refs/heads/main\n"));
480 ok.extend_from_slice(pktline::flush());
481 assert!(check_report_status(&ok).is_ok());
482
483 let mut bad = Vec::new();
484 bad.extend_from_slice(&pktline::encode(b"unpack ok\n"));
485 bad.extend_from_slice(&pktline::encode(b"ng refs/heads/main denied\n"));
486 bad.extend_from_slice(pktline::flush());
487 let err = check_report_status(&bad).unwrap_err();
488 assert!(err.contains("refs/heads/main denied"));
489 }
490}