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