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