anvilsign in

collin/anvil

1// This Source Code Form is subject to the terms of the Mozilla Public
2// License, v. 2.0. If a copy of the MPL was not distributed with this
3// file, You can obtain one at https://mozilla.org/MPL/2.0/.
4//
5// Copyright (c) 2026 WJQSERVER
6
7use std::io::{BufRead, BufReader, Read};
8use std::path::Path;
9use std::sync::atomic::AtomicBool;
10
11use gix::objs::bstr::BString;
12use gix::prelude::ObjectIdExt;
13use gix::progress::Discard;
14use gix::refs::Target;
15use gix::refs::transaction::{Change, LogChange, PreviousValue, RefEdit, RefLog};
16
17use crate::error::{Error, Result};
18use crate::pktline;
19
20const ZERO_ID: &str = "0000000000000000000000000000000000000000";
21const CAPABILITIES: &str = concat!(
22 "report-status report-status-v2 side-band-64k quiet ofs-delta object-format=sha1 agent=gitserver/",
23 env!("CARGO_PKG_VERSION")
24);
25
26pub fn advertise_receive_refs(repo_path: &Path) -> Result<Vec<u8>> {
27 let repo = gix::open(repo_path)?;
28 let mut out = Vec::new();
29 let head_name = repo.head_name().ok().flatten();
30 let mut refs: Vec<(String, gix::ObjectId)> = repo
31 .references()
32 .map_err(|e| Error::Protocol(format!("failed to open refs: {e}")))?
33 .all()
34 .map_err(|e| Error::Protocol(format!("failed to iterate refs: {e}")))?
35 .flatten()
36 .filter_map(|mut reference| {
37 reference
38 .peel_to_id()
39 .ok()
40 .map(|id| (reference.name().as_bstr().to_string(), id.detach()))
41 })
42 .collect();
43 refs.sort_by(|a, b| a.0.cmp(&b.0));
44
45 if refs.is_empty() {
46 out.extend_from_slice(&pktline::encode(
47 format!("{ZERO_ID} capabilities^{{}}\0{CAPABILITIES}\n").as_bytes(),
48 ));
49 } else {
50 let (first_name, first_id) = &refs[0];
51 let mut first = format!("{} {}\0{CAPABILITIES}", first_id, first_name);
52 if head_name
53 .as_ref()
54 .is_some_and(|head| head.as_bstr() == first_name.as_str())
55 {
56 first.push_str(&format!(" symref=HEAD:{first_name}"));
57 }
58 first.push('\n');
59 out.extend_from_slice(&pktline::encode(first.as_bytes()));
60
61 for (name, id) in refs.into_iter().skip(1) {
62 out.extend_from_slice(&pktline::encode(format!("{id} {name}\n").as_bytes()));
63 }
64 }
65
66 out.extend_from_slice(pktline::flush());
67 Ok(out)
68}
69
70pub fn receive_pack<R: Read>(repo_path: &Path, request: R) -> Result<Vec<u8>> {
71 let interrupt = AtomicBool::new(false);
72 receive_pack_with_interrupt(repo_path, request, &interrupt)
73}
74
75pub fn receive_pack_with_interrupt<R: Read>(
76 repo_path: &Path,
77 request: R,
78 interrupt: &AtomicBool,
79) -> Result<Vec<u8>> {
80 let repo = gix::open(repo_path)?;
81 let mut parsed = parse_request(request, interrupt)?;
82 let status = apply_commands(&repo, repo_path, &mut parsed, interrupt)?;
83 Ok(encode_report_status(&parsed.capabilities, &status))
84}
85
86#[derive(Default)]
87struct ReceivePackCapabilities {
88 report_status: bool,
89 report_status_v2: bool,
90}
91
92struct ReceivePackRequest<R> {
93 commands: Vec<UpdateCommand>,
94 pack: R,
95 capabilities: ReceivePackCapabilities,
96}
97
98struct UpdateCommand {
99 old_id: String,
100 new_id: String,
101 refname: String,
102}
103
104enum CommandStatus {
105 Ok(String),
106 Ng(String, String),
107}
108
109fn parse_request<R: Read>(
110 request: R,
111 interrupt: &AtomicBool,
112) -> Result<ReceivePackRequest<BufReader<R>>> {
113 let mut request = BufReader::new(request);
114 let mut commands = Vec::new();
115 let mut capabilities = ReceivePackCapabilities::default();
116
117 loop {
118 check_interrupt(interrupt)?;
119 let mut prefix = [0u8; 4];
120 match request.read_exact(&mut prefix) {
121 Ok(()) => {}
122 Err(err) if err.kind() == std::io::ErrorKind::UnexpectedEof => break,
123 Err(err) => return Err(Error::Io(err)),
124 }
125
126 let len_str = std::str::from_utf8(&prefix)
127 .map_err(|_| Error::Protocol("invalid pkt-line length prefix".into()))?;
128
129 if len_str == "0000" {
130 break;
131 }
132
133 let len = usize::from_str_radix(len_str, 16)
134 .map_err(|_| Error::Protocol("invalid pkt-line length".into()))?;
135 if len < 4 {
136 return Err(Error::Protocol("invalid pkt-line frame length".into()));
137 }
138
139 check_interrupt(interrupt)?;
140 let mut payload = vec![0u8; len - 4];
141 request.read_exact(&mut payload)?;
142
143 let (command_bytes, capability_bytes) =
144 if let Some(nul) = payload.iter().position(|b| *b == 0) {
145 (&payload[..nul], Some(&payload[nul + 1..]))
146 } else {
147 (&payload[..], None)
148 };
149
150 if let Some(capability_bytes) = capability_bytes {
151 let capabilities_line = std::str::from_utf8(capability_bytes)
152 .map_err(|_| Error::Protocol("invalid UTF-8 in receive-pack capabilities".into()))?
153 .trim_end_matches('\n');
154 for capability in capabilities_line.split_ascii_whitespace() {
155 match capability {
156 "report-status" => capabilities.report_status = true,
157 "report-status-v2" => {
158 capabilities.report_status = true;
159 capabilities.report_status_v2 = true;
160 }
161 _ => {}
162 }
163 }
164 }
165
166 let line = std::str::from_utf8(command_bytes)
167 .map_err(|_| Error::Protocol("invalid UTF-8 in update command".into()))?
168 .trim_end_matches('\n');
169 let mut parts = line.split_ascii_whitespace();
170 let Some(old_id) = parts.next() else { continue };
171 let Some(new_id) = parts.next() else { continue };
172 let Some(refname) = parts.next() else {
173 continue;
174 };
175
176 commands.push(UpdateCommand {
177 old_id: old_id.to_owned(),
178 new_id: new_id.to_owned(),
179 refname: refname.to_owned(),
180 });
181 }
182
183 Ok(ReceivePackRequest {
184 commands,
185 pack: request,
186 capabilities,
187 })
188}
189
190fn apply_commands<R: BufRead>(
191 repo: &gix::Repository,
192 repo_path: &Path,
193 request: &mut ReceivePackRequest<R>,
194 interrupt: &AtomicBool,
195) -> Result<Vec<CommandStatus>> {
196 check_interrupt(interrupt)?;
197 if request.pack.fill_buf().map(|buf: &[u8]| !buf.is_empty())? {
198 write_pack(repo, repo_path, &mut request.pack, interrupt)?;
199 }
200
201 let mut edits = Vec::with_capacity(request.commands.len());
202 for (index, command) in request.commands.iter().enumerate() {
203 check_interrupt(interrupt)?;
204 match validate_ref_update(repo, command, interrupt) {
205 Ok(edit) => edits.push((command.refname.clone(), edit)),
206 Err(err) => {
207 return Ok(request
208 .commands
209 .iter()
210 .enumerate()
211 .map(|(cmd_index, cmd)| {
212 if cmd_index == index {
213 CommandStatus::Ng(cmd.refname.clone(), err.to_string())
214 } else {
215 CommandStatus::Ng(
216 cmd.refname.clone(),
217 "transaction aborted due to another command failing validation"
218 .into(),
219 )
220 }
221 })
222 .collect());
223 }
224 }
225 }
226
227 check_interrupt(interrupt)?;
228 match repo.edit_references(edits.into_iter().map(|(_, edit)| edit)) {
229 Ok(_) => Ok(request
230 .commands
231 .iter()
232 .map(|cmd| CommandStatus::Ok(cmd.refname.clone()))
233 .collect()),
234 Err(err) => Ok(request
235 .commands
236 .iter()
237 .map(|cmd| CommandStatus::Ng(cmd.refname.clone(), format!("transaction failed: {err}")))
238 .collect()),
239 }
240}
241
242fn write_pack<R: BufRead>(
243 repo: &gix::Repository,
244 repo_path: &Path,
245 pack: &mut R,
246 interrupt: &AtomicBool,
247) -> Result<()> {
248 let mut progress = Discard;
249 // Pushes after the first send a thin pack whose delta bases (ref-deltas) live
250 // only in the existing object database; the lookup lets gix resolve them.
251 let outcome = gix_pack::Bundle::write_to_directory(
252 pack,
253 Some(repo_path.join("objects/pack").as_path()),
254 &mut progress,
255 interrupt,
256 Some(repo),
257 Default::default(),
258 );
259 if interrupt.load(std::sync::atomic::Ordering::Relaxed) {
260 return Err(Error::Io(std::io::Error::new(
261 std::io::ErrorKind::TimedOut,
262 "receive-pack timed out",
263 )));
264 }
265 let outcome =
266 outcome.map_err(|e| Error::Protocol(format!("failed to write incoming pack: {e}")))?;
267
268 if let Some(keep) = outcome.keep_path {
269 let _ = std::fs::remove_file(keep);
270 }
271 Ok(())
272}
273
274fn check_interrupt(interrupt: &AtomicBool) -> Result<()> {
275 if interrupt.load(std::sync::atomic::Ordering::Relaxed) {
276 Err(Error::Io(std::io::Error::new(
277 std::io::ErrorKind::TimedOut,
278 "receive-pack timed out",
279 )))
280 } else {
281 Ok(())
282 }
283}
284
285fn validate_ref_update(
286 repo: &gix::Repository,
287 command: &UpdateCommand,
288 interrupt: &AtomicBool,
289) -> Result<RefEdit> {
290 if command.new_id == ZERO_ID {
291 return Err(Error::Protocol(format!(
292 "deletion prohibited for {}",
293 command.refname
294 )));
295 }
296
297 let is_branch = command.refname.starts_with("refs/heads/");
298 let is_tag = command.refname.starts_with("refs/tags/");
299 let new_id = gix::ObjectId::from_hex(command.new_id.as_bytes())
300 .map_err(|_| Error::Protocol(format!("invalid new object id: {}", command.new_id)))?;
301 let new_header = repo
302 .find_header(new_id)
303 .map_err(|e| Error::Protocol(format!("missing new object {}: {e}", command.new_id)))?;
304 if is_branch && new_header.kind() != gix::objs::Kind::Commit {
305 return Err(Error::Protocol(format!(
306 "updates to {} must point to a commit",
307 command.refname
308 )));
309 }
310
311 let name: gix::refs::FullName = command
312 .refname
313 .as_str()
314 .try_into()
315 .map_err(|e| Error::Protocol(format!("invalid ref name {}: {e}", command.refname)))?;
316
317 let (expected, log_message) = if command.old_id == ZERO_ID {
318 (PreviousValue::MustNotExist, BString::from("push create"))
319 } else {
320 if is_tag {
321 return Err(Error::Protocol(format!(
322 "updating existing tag {} is not allowed",
323 command.refname
324 )));
325 }
326
327 let old_id = gix::ObjectId::from_hex(command.old_id.as_bytes())
328 .map_err(|_| Error::Protocol(format!("invalid old object id: {}", command.old_id)))?;
329 if is_branch {
330 ensure_fast_forward(repo, old_id, new_id, &command.refname, interrupt)?;
331 }
332 (
333 PreviousValue::MustExistAndMatch(Target::Object(old_id)),
334 BString::from("push"),
335 )
336 };
337
338 Ok(RefEdit {
339 change: Change::Update {
340 log: LogChange {
341 mode: RefLog::AndReference,
342 force_create_reflog: false,
343 message: log_message,
344 },
345 expected,
346 new: Target::Object(new_id),
347 },
348 name,
349 deref: false,
350 })
351}
352
353fn ensure_fast_forward(
354 repo: &gix::Repository,
355 old_id: gix::ObjectId,
356 new_id: gix::ObjectId,
357 refname: &str,
358 interrupt: &AtomicBool,
359) -> Result<()> {
360 check_interrupt(interrupt)?;
361 if old_id == new_id {
362 return Ok(());
363 }
364
365 let old_commit_time = repo
366 .find_object(old_id)
367 .map_err(|e| Error::Protocol(format!("failed to inspect current tip for {refname}: {e}")))?
368 .try_into_commit()
369 .map_err(|_| Error::Protocol(format!("current tip of {refname} is not a commit")))?
370 .committer()
371 .map_err(|e| Error::Protocol(format!("failed to read commit metadata for {refname}: {e}")))?
372 .seconds();
373
374 let ancestors = new_id
375 .attach(repo)
376 .ancestors()
377 .sorting(gix::revision::walk::Sorting::ByCommitTimeCutoff {
378 order: Default::default(),
379 seconds: old_commit_time,
380 })
381 .all()
382 .map_err(|e| Error::Protocol(format!("failed to walk commits for {refname}: {e}")))?;
383
384 for id in ancestors {
385 check_interrupt(interrupt)?;
386 if id.is_ok_and(|commit| commit.id == old_id) {
387 return Ok(());
388 }
389 }
390
391 Err(Error::Protocol(format!(
392 "non-fast-forward update to {refname} is not allowed"
393 )))
394}
395
396fn encode_report_status(
397 capabilities: &ReceivePackCapabilities,
398 statuses: &[CommandStatus],
399) -> Vec<u8> {
400 if !capabilities.report_status {
401 return pktline::flush().to_vec();
402 }
403
404 let mut status_lines = Vec::new();
405 status_lines.extend_from_slice(&pktline::encode(b"unpack ok\n"));
406
407 for status in statuses {
408 match status {
409 CommandStatus::Ok(refname) => {
410 status_lines
411 .extend_from_slice(&pktline::encode(format!("ok {refname}\n").as_bytes()));
412 }
413 CommandStatus::Ng(refname, message) => {
414 status_lines.extend_from_slice(&pktline::encode(
415 format!("ng {refname} {message}\n").as_bytes(),
416 ));
417 }
418 }
419 }
420 status_lines.extend_from_slice(pktline::flush());
421
422 if capabilities.report_status_v2 {
423 let mut sideband = Vec::new();
424 const MAX_BAND_PAYLOAD: usize = 65519;
425 for chunk in status_lines.chunks(MAX_BAND_PAYLOAD) {
426 let len = 4 + 1 + chunk.len();
427 sideband.extend_from_slice(format!("{len:04x}").as_bytes());
428 sideband.push(0x01);
429 sideband.extend_from_slice(chunk);
430 }
431 sideband.extend_from_slice(pktline::flush());
432 sideband
433 } else {
434 status_lines
435 }
436}
437
438#[cfg(test)]
439mod tests {
440 use std::process::Command;
441
442 use tempfile::TempDir;
443
444 use super::*;
445
446 fn create_repo_with_commit(root: &std::path::Path) -> std::path::PathBuf {
447 let repo_path = root.join("test.git");
448 let work_dir = root.join("work");
449 std::fs::create_dir(&work_dir).unwrap();
450 Command::new("git")
451 .args(["init", "--bare", repo_path.to_str().unwrap()])
452 .output()
453 .unwrap();
454 Command::new("git")
455 .args(["symbolic-ref", "HEAD", "refs/heads/main"])
456 .current_dir(&repo_path)
457 .output()
458 .unwrap();
459 Command::new("git")
460 .args([
461 "clone",
462 repo_path.to_str().unwrap(),
463 work_dir.to_str().unwrap(),
464 ])
465 .output()
466 .unwrap();
467 Command::new("git")
468 .current_dir(&work_dir)
469 .args(["commit", "--allow-empty", "-m", "init"])
470 .env("GIT_AUTHOR_NAME", "Test")
471 .env("GIT_AUTHOR_EMAIL", "t@t.com")
472 .env("GIT_COMMITTER_NAME", "Test")
473 .env("GIT_COMMITTER_EMAIL", "t@t.com")
474 .output()
475 .unwrap();
476 Command::new("git")
477 .current_dir(&work_dir)
478 .args(["push", "origin", "main"])
479 .output()
480 .unwrap();
481 repo_path
482 }
483
484 #[test]
485 fn advertise_receive_pack_refs() {
486 let root = TempDir::new().unwrap();
487 let repo_path = create_repo_with_commit(root.path());
488 let output = advertise_receive_refs(&repo_path).unwrap();
489 let output_str = String::from_utf8_lossy(&output);
490 assert!(output_str.contains("refs/heads/main"));
491 assert!(output_str.contains("report-status"));
492 }
493
494 #[test]
495 fn parse_receive_pack_request_with_capabilities() {
496 let payload = b"0000000000000000000000000000000000000000 1111111111111111111111111111111111111111 refs/heads/main\0 report-status-v2 side-band-64k\n";
497 let mut body = format!("{:04x}", payload.len() + 4).into_bytes();
498 body.extend_from_slice(payload);
499 body.extend_from_slice(b"0000PACK");
500
501 let interrupt = AtomicBool::new(false);
502 let parsed = parse_request(std::io::Cursor::new(&body), &interrupt).unwrap();
503 assert_eq!(parsed.commands.len(), 1);
504 assert!(parsed.capabilities.report_status);
505 assert!(parsed.capabilities.report_status_v2);
506 let mut pack = String::new();
507 let mut reader = parsed.pack;
508 reader.read_to_string(&mut pack).unwrap();
509 assert_eq!(pack.as_bytes(), b"PACK");
510 }
511
512 #[test]
513 fn branch_updates_require_commit_target() {
514 let root = TempDir::new().unwrap();
515 let repo_path = create_repo_with_commit(root.path());
516 let repo = gix::open(repo_path).unwrap();
517 let tree_id = Command::new("git")
518 .args(["rev-parse", "HEAD^{tree}"])
519 .current_dir(root.path().join("work"))
520 .output()
521 .unwrap();
522 let tree_id = String::from_utf8(tree_id.stdout)
523 .unwrap()
524 .trim()
525 .to_string();
526
527 let err = validate_ref_update(
528 &repo,
529 &UpdateCommand {
530 old_id: ZERO_ID.into(),
531 new_id: tree_id,
532 refname: "refs/heads/feature".into(),
533 },
534 &AtomicBool::new(false),
535 )
536 .unwrap_err();
537
538 assert!(err.to_string().contains("must point to a commit"));
539 }
540
541 #[test]
542 fn receive_thin_pack_with_ref_deltas() {
543 let root = TempDir::new().unwrap();
544 let repo_path = root.path().join("test.git");
545 let work_dir = root.path().join("work");
546 std::fs::create_dir(&work_dir).unwrap();
547 Command::new("git")
548 .args(["init", "--bare", repo_path.to_str().unwrap()])
549 .output()
550 .unwrap();
551 Command::new("git")
552 .args(["symbolic-ref", "HEAD", "refs/heads/main"])
553 .current_dir(&repo_path)
554 .output()
555 .unwrap();
556 Command::new("git")
557 .args([
558 "clone",
559 repo_path.to_str().unwrap(),
560 work_dir.to_str().unwrap(),
561 ])
562 .output()
563 .unwrap();
564
565 let git = |args: &[&str]| {
566 let output = Command::new("git")
567 .current_dir(&work_dir)
568 .args(args)
569 .env("GIT_AUTHOR_NAME", "Test")
570 .env("GIT_AUTHOR_EMAIL", "t@t.com")
571 .env("GIT_COMMITTER_NAME", "Test")
572 .env("GIT_COMMITTER_EMAIL", "t@t.com")
573 .output()
574 .unwrap();
575 assert!(output.status.success(), "git {args:?}: {output:?}");
576 String::from_utf8(output.stdout).unwrap().trim().to_string()
577 };
578
579 // A large repetitive blob so the follow-up commit deltas against it.
580 let base_content = "this line repeats to make the blob delta-friendly\n".repeat(200);
581 std::fs::write(work_dir.join("data.txt"), &base_content).unwrap();
582 git(&["add", "data.txt"]);
583 git(&["commit", "-m", "base"]);
584 git(&["push", "origin", "main"]);
585 let old_id = git(&["rev-parse", "HEAD"]);
586
587 std::fs::write(
588 work_dir.join("data.txt"),
589 format!("{base_content}one more line\n"),
590 )
591 .unwrap();
592 git(&["add", "data.txt"]);
593 git(&["commit", "-m", "append"]);
594 let new_id = git(&["rev-parse", "HEAD"]);
595
596 // Build a thin pack exactly like a push would: bases from old_id stay out.
597 let mut pack_objects = Command::new("git")
598 .current_dir(&work_dir)
599 .args(["pack-objects", "--thin", "--stdout", "--revs", "-q"])
600 .stdin(std::process::Stdio::piped())
601 .stdout(std::process::Stdio::piped())
602 .spawn()
603 .unwrap();
604 use std::io::Write as _;
605 pack_objects
606 .stdin
607 .take()
608 .unwrap()
609 .write_all(format!("{new_id}\n^{old_id}\n").as_bytes())
610 .unwrap();
611 let pack = pack_objects.wait_with_output().unwrap();
612 assert!(pack.status.success());
613 let pack = pack.stdout;
614
615 // The fix only matters if the pack really contains ref deltas.
616 let entries = gix_pack::data::input::BytesToEntriesIter::new_from_header(
617 std::io::Cursor::new(&pack),
618 gix_pack::data::input::Mode::Verify,
619 gix_pack::data::input::EntryDataMode::Ignore,
620 gix::hash::Kind::Sha1,
621 )
622 .unwrap();
623 let has_ref_delta = entries
624 .map(|e| e.unwrap())
625 .any(|entry| matches!(entry.header, gix_pack::data::entry::Header::RefDelta { .. }));
626 assert!(has_ref_delta, "test pack should contain ref delta objects");
627
628 let mut body = Vec::new();
629 body.extend_from_slice(&pktline::encode(
630 format!("{old_id} {new_id} refs/heads/main\0 report-status\n").as_bytes(),
631 ));
632 body.extend_from_slice(pktline::flush());
633 body.extend_from_slice(&pack);
634
635 let response = receive_pack(&repo_path, std::io::Cursor::new(body)).unwrap();
636 let response = String::from_utf8_lossy(&response);
637 assert!(response.contains("unpack ok"), "response: {response}");
638 assert!(
639 response.contains("ok refs/heads/main"),
640 "response: {response}"
641 );
642
643 let repo = gix::open(&repo_path).unwrap();
644 assert_eq!(repo.head_id().unwrap().detach().to_string(), new_id);
645 }
646
647 #[test]
648 fn ensure_fast_forward_respects_interrupt() {
649 let root = TempDir::new().unwrap();
650 let repo_path = create_repo_with_commit(root.path());
651 let repo = gix::open(repo_path).unwrap();
652 let head = repo.head_id().unwrap().detach();
653 let interrupt = AtomicBool::new(true);
654
655 let err =
656 ensure_fast_forward(&repo, head, head, "refs/heads/main", &interrupt).unwrap_err();
657 match err {
658 Error::Io(inner) => assert_eq!(inner.kind(), std::io::ErrorKind::TimedOut),
659 other => panic!("expected timeout io error, got {other}"),
660 }
661 }
662}