anvilsign in

collin/anvil

1use std::{
2 collections::HashSet,
3 path::Path,
4};
5
6use bytes::Bytes;
7use sha1::{
8 Digest,
9 Sha1,
10};
11use tokio::io::AsyncRead;
12use tokio_util::io::StreamReader;
13
14use crate::{
15 error::{
16 Error,
17 Result,
18 },
19 pktline,
20};
21
22#[derive(Clone, Debug, Default)]
23pub struct UploadPackCapabilities {
24 pub ofs_delta: bool,
25 pub multi_ack: bool,
26 pub multi_ack_detailed: bool,
27}
28
29#[derive(Clone, Debug, Default)]
30pub struct ShallowRequest {
31 pub depth: Option<usize>,
32 pub client_shallows: Vec<gix::ObjectId>,
33 pub deepen_relative: bool,
34}
35
36/// A parsed upload-pack request from a Git client.
37pub struct UploadPackRequest {
38 pub wants: Vec<gix::ObjectId>,
39 pub haves: Vec<gix::ObjectId>,
40 pub done: bool,
41 pub capabilities: UploadPackCapabilities,
42 pub shallow: ShallowRequest,
43 pub object_ids: Option<Vec<gix::ObjectId>>,
44}
45
46impl UploadPackRequest {
47 /// Parse a pkt-line encoded upload-pack request body.
48 ///
49 /// The body contains:
50 /// - "want <oid> [capabilities]\n" lines
51 /// - flush packet "0000"
52 /// - "have <oid>\n" lines (optional)
53 /// - "done\n"
54 pub fn parse(body: &[u8]) -> Result<Self> {
55 let mut wants = Vec::new();
56 let mut haves = Vec::new();
57 let mut done = false;
58 let mut capabilities = UploadPackCapabilities::default();
59 let mut shallow = ShallowRequest::default();
60 let mut pos = 0;
61
62 while pos < body.len() {
63 // Check for flush packet
64 if body[pos..].starts_with(b"0000") {
65 pos += 4;
66 continue;
67 }
68
69 // Read 4-byte hex length prefix
70 if pos + 4 > body.len() {
71 break;
72 }
73 let len_str = std::str::from_utf8(&body[pos..pos + 4])
74 .map_err(|_| Error::Protocol("invalid pkt-line length prefix".into()))?;
75 let len = usize::from_str_radix(len_str, 16)
76 .map_err(|_| Error::Protocol("invalid pkt-line length".into()))?;
77
78 if len == 0 {
79 // flush packet already handled above, but just in case
80 pos += 4;
81 continue;
82 }
83
84 if len < 4 || pos + len > body.len() {
85 break;
86 }
87
88 let payload = &body[pos + 4..pos + len];
89 let line = std::str::from_utf8(payload)
90 .map_err(|_| Error::Protocol("invalid UTF-8 in pkt-line".into()))?;
91 let line = line.trim_end_matches('\n');
92
93 if line == "done" {
94 done = true;
95 } else if let Some(rest) = line.strip_prefix("deepen ") {
96 let depth = rest
97 .parse::<usize>()
98 .map_err(|_| Error::Protocol(format!("invalid deepen value: {rest}")))?;
99 shallow.depth = Some(depth);
100 } else if line == "deepen-relative" {
101 shallow.deepen_relative = true;
102 } else if let Some(rest) = line.strip_prefix("shallow ") {
103 let oid = gix::ObjectId::from_hex(rest.as_bytes())
104 .map_err(|_| Error::Protocol(format!("invalid OID in shallow: {rest}")))?;
105 shallow.client_shallows.push(oid);
106 } else if let Some(rest) = line.strip_prefix("want ") {
107 let mut parts = rest.split_ascii_whitespace();
108 let oid_hex = parts
109 .next()
110 .ok_or_else(|| Error::Protocol("missing OID in want".into()))?;
111 let oid = gix::ObjectId::from_hex(oid_hex.as_bytes())
112 .map_err(|_| Error::Protocol(format!("invalid OID in want: {oid_hex}")))?;
113 if wants.is_empty() {
114 for capability in parts {
115 if capability == "ofs-delta" {
116 capabilities.ofs_delta = true;
117 } else if capability == "multi_ack" {
118 capabilities.multi_ack = true;
119 } else if capability == "multi_ack_detailed" {
120 capabilities.multi_ack = true;
121 capabilities.multi_ack_detailed = true;
122 }
123 }
124 }
125 wants.push(oid);
126 } else if let Some(rest) = line.strip_prefix("have ") {
127 let oid_hex = rest
128 .split_ascii_whitespace()
129 .next()
130 .ok_or_else(|| Error::Protocol("missing OID in have".into()))?;
131 let oid = gix::ObjectId::from_hex(oid_hex.as_bytes())
132 .map_err(|_| Error::Protocol(format!("invalid OID in have: {oid_hex}")))?;
133 haves.push(oid);
134 }
135
136 pos += len;
137 }
138
139 Ok(Self {
140 wants,
141 haves,
142 done,
143 capabilities,
144 shallow,
145 object_ids: None,
146 })
147 }
148}
149
150/// Encode the variable-length pack object header.
151///
152/// Format: first byte = MSB continuation + 3-bit type + 4-bit size
153/// Subsequent bytes: 7-bit size chunks with MSB continuation
154fn encode_pack_object_header(obj_type: u8, size: usize) -> Vec<u8> {
155 let mut header = Vec::new();
156 let mut byte = (obj_type << 4) | (size as u8 & 0x0f);
157 let mut remaining = size >> 4;
158
159 if remaining > 0 {
160 byte |= 0x80; // set continuation bit
161 header.push(byte);
162 while remaining > 0 {
163 byte = remaining as u8 & 0x7f;
164 remaining >>= 7;
165 if remaining > 0 {
166 byte |= 0x80;
167 }
168 header.push(byte);
169 }
170 } else {
171 header.push(byte);
172 }
173
174 header
175}
176
177fn encode_ofs_delta_base_distance(mut distance: u64) -> Vec<u8> {
178 debug_assert!(distance > 0, "offset deltas must point backwards");
179
180 let mut buf = [0u8; 10];
181 let mut bytes_written = 1;
182 buf[buf.len() - 1] = distance as u8 & 0x7f;
183
184 for out in buf.iter_mut().rev().skip(1) {
185 distance >>= 7;
186 if distance == 0 {
187 break;
188 }
189 distance -= 1;
190 *out = 0x80 | (distance as u8 & 0x7f);
191 bytes_written += 1;
192 }
193
194 buf[buf.len() - bytes_written..].to_vec()
195}
196
197fn encode_delta_size(mut size: usize, out: &mut Vec<u8>) {
198 loop {
199 let mut byte = (size & 0x7f) as u8;
200 size >>= 7;
201 if size > 0 {
202 byte |= 0x80;
203 }
204 out.push(byte);
205 if size == 0 {
206 break;
207 }
208 }
209}
210
211fn encode_delta_copy_instruction(out: &mut Vec<u8>, offset: usize, size: usize) {
212 debug_assert!(size > 0 && size <= 0x10000);
213
214 let command_pos = out.len();
215 out.push(0x80);
216 let mut command = 0x80;
217
218 if offset & 0xff != 0 {
219 command |= 0x01;
220 out.push(offset as u8);
221 }
222 if (offset >> 8) & 0xff != 0 {
223 command |= 0x02;
224 out.push((offset >> 8) as u8);
225 }
226 if (offset >> 16) & 0xff != 0 {
227 command |= 0x04;
228 out.push((offset >> 16) as u8);
229 }
230 if (offset >> 24) & 0xff != 0 {
231 command |= 0x08;
232 out.push((offset >> 24) as u8);
233 }
234
235 if size != 0x10000 {
236 if size & 0xff != 0 {
237 command |= 0x10;
238 out.push(size as u8);
239 }
240 if (size >> 8) & 0xff != 0 {
241 command |= 0x20;
242 out.push((size >> 8) as u8);
243 }
244 if (size >> 16) & 0xff != 0 {
245 command |= 0x40;
246 out.push((size >> 16) as u8);
247 }
248 }
249
250 out[command_pos] = command;
251}
252
253fn encode_delta_copy(out: &mut Vec<u8>, mut offset: usize, mut size: usize) {
254 while size > 0 {
255 let chunk = size.min(0x10000);
256 encode_delta_copy_instruction(out, offset, chunk);
257 offset += chunk;
258 size -= chunk;
259 }
260}
261
262fn encode_delta_insert(out: &mut Vec<u8>, data: &[u8]) {
263 for chunk in data.chunks(0x7f) {
264 out.push(chunk.len() as u8);
265 out.extend_from_slice(chunk);
266 }
267}
268
269fn encode_blob_delta(base: &[u8], target: &[u8]) -> Option<Vec<u8>> {
270 let mut prefix = 0;
271 let max_prefix = base.len().min(target.len());
272 while prefix < max_prefix && base[prefix] == target[prefix] {
273 prefix += 1;
274 }
275
276 let max_suffix = base
277 .len()
278 .saturating_sub(prefix)
279 .min(target.len().saturating_sub(prefix));
280 let mut suffix = 0;
281 while suffix < max_suffix && base[base.len() - 1 - suffix] == target[target.len() - 1 - suffix]
282 {
283 suffix += 1;
284 }
285
286 if prefix == 0 && suffix == 0 {
287 return None;
288 }
289
290 let mut delta = Vec::new();
291 encode_delta_size(base.len(), &mut delta);
292 encode_delta_size(target.len(), &mut delta);
293
294 if prefix > 0 {
295 encode_delta_copy(&mut delta, 0, prefix);
296 }
297
298 let insert_start = prefix;
299 let insert_end = target.len() - suffix;
300 encode_delta_insert(&mut delta, &target[insert_start..insert_end]);
301
302 if suffix > 0 {
303 encode_delta_copy(&mut delta, base.len() - suffix, suffix);
304 }
305
306 Some(delta)
307}
308
309fn build_base_entry(kind: gix::object::Kind, data: &[u8]) -> Vec<u8> {
310 let type_num = object_type_number(kind);
311 let obj_header = encode_pack_object_header(type_num, data.len());
312 let compressed = miniz_oxide::deflate::compress_to_vec_zlib(data, 6);
313
314 let mut entry = Vec::with_capacity(obj_header.len() + compressed.len());
315 entry.extend_from_slice(&obj_header);
316 entry.extend_from_slice(&compressed);
317 entry
318}
319
320fn build_ofs_delta_entry(
321 pack_offset: u64,
322 base_pack_offset: u64,
323 base_data: &[u8],
324 target_data: &[u8],
325) -> Option<Vec<u8>> {
326 let delta = encode_blob_delta(base_data, target_data)?;
327 let obj_header = encode_pack_object_header(6, delta.len());
328 let base_distance = encode_ofs_delta_base_distance(pack_offset - base_pack_offset);
329 let compressed = miniz_oxide::deflate::compress_to_vec_zlib(&delta, 6);
330
331 let mut entry = Vec::with_capacity(obj_header.len() + base_distance.len() + compressed.len());
332 entry.extend_from_slice(&obj_header);
333 entry.extend_from_slice(&base_distance);
334 entry.extend_from_slice(&compressed);
335 Some(entry)
336}
337
338struct BlobDeltaBase {
339 pack_offset: u64,
340 data: Vec<u8>,
341}
342
343/// Map gix object kind to pack type number.
344fn object_type_number(kind: gix::object::Kind) -> u8 {
345 match kind {
346 gix::object::Kind::Commit => 1,
347 gix::object::Kind::Tree => 2,
348 gix::object::Kind::Blob => 3,
349 gix::object::Kind::Tag => 4,
350 }
351}
352
353/// Send raw bytes through the channel.
354fn send(
355 tx: &tokio::sync::mpsc::Sender<std::result::Result<Bytes, std::io::Error>>,
356 data: &[u8],
357) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> {
358 tx.blocking_send(Ok(Bytes::copy_from_slice(data)))
359 .map_err(|_| "receiver dropped".into())
360}
361
362/// Send pack data through the channel wrapped in side-band-64k framing
363/// (band 1 = pack data).
364///
365/// Respects LARGE_PACKET_MAX: each pkt-line frame carries at most
366/// 65520 - 4 (prefix) - 1 (band byte) = 65515 bytes of payload.
367fn send_sideband(
368 tx: &tokio::sync::mpsc::Sender<std::result::Result<Bytes, std::io::Error>>,
369 data: &[u8],
370) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> {
371 const MAX_DATA_PER_FRAME: usize = 65515;
372
373 for chunk in data.chunks(MAX_DATA_PER_FRAME) {
374 let pkt_len = 4 + 1 + chunk.len();
375 let mut frame = Vec::with_capacity(pkt_len);
376 frame.extend_from_slice(format!("{pkt_len:04x}").as_bytes());
377 frame.push(0x01); // band 1 = pack data
378 frame.extend_from_slice(chunk);
379 send(tx, &frame)?;
380 }
381
382 Ok(())
383}
384
385fn encode_ack_line(oid: gix::ObjectId, suffix: Option<&str>) -> Vec<u8> {
386 let mut line = format!("ACK {oid}");
387 if let Some(suffix) = suffix {
388 line.push(' ');
389 line.push_str(suffix);
390 }
391 line.push('\n');
392 pktline::encode(line.as_bytes())
393}
394
395/// Recursively collect tree and blob OIDs reachable from `tree_oid`.
396///
397/// Uses a single `find_object` call per object and parses raw tree
398/// bytes via `TreeRefIter` to avoid a second ODB lookup.
399fn collect_tree_oids(
400 repo: &gix::Repository,
401 tree_oid: gix::ObjectId,
402 seen: &mut HashSet<gix::ObjectId>,
403 oids: &mut Vec<gix::ObjectId>,
404) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> {
405 if !seen.insert(tree_oid) {
406 return Ok(());
407 }
408
409 let tree_obj = repo.find_object(tree_oid)?;
410 let tree_data = tree_obj.data.to_vec();
411 oids.push(tree_oid);
412
413 for entry_result in gix::objs::TreeRefIter::from_bytes(&tree_data, gix::hash::Kind::Sha1) {
414 let entry = entry_result?;
415 let entry_oid = entry.oid.to_owned();
416 let entry_mode = entry.mode;
417
418 if entry_mode.is_tree() {
419 collect_tree_oids(repo, entry_oid, seen, oids)?;
420 } else if seen.insert(entry_oid) && !entry_mode.is_commit() {
421 oids.push(entry_oid);
422 }
423 }
424
425 Ok(())
426}
427
428/// Peel `wants` for a pack walk: tag wants contribute the tag object itself
429/// (plus the chain down to a non-tag), and the peeled target joins the walk —
430/// commits as walk tips, trees/blobs (rare tag targets) collected directly.
431/// `git fetch`/`push` send tag *object* ids whenever annotated tags are
432/// involved, so every pack path needs this.
433pub(crate) fn peel_wants(
434 repo: &gix::Repository,
435 wants: &[gix::ObjectId],
436) -> std::result::Result<
437 (Vec<gix::ObjectId>, Vec<gix::ObjectId>),
438 Box<dyn std::error::Error + Send + Sync>,
439> {
440 let mut commit_wants = Vec::new();
441 let mut extras = Vec::new();
442 let mut seen = HashSet::new();
443 for want in wants {
444 let mut oid = *want;
445 loop {
446 let obj = repo.find_object(oid)?;
447 match obj.kind {
448 gix::object::Kind::Tag => {
449 if seen.insert(oid) {
450 extras.push(oid);
451 }
452 oid = gix::objs::TagRefIter::from_bytes(&obj.data, gix::hash::Kind::Sha1)
453 .target_id()?;
454 }
455 gix::object::Kind::Commit => {
456 commit_wants.push(oid);
457 break;
458 }
459 gix::object::Kind::Tree => {
460 collect_tree_oids(repo, oid, &mut seen, &mut extras)?;
461 break;
462 }
463 gix::object::Kind::Blob => {
464 if seen.insert(oid) {
465 extras.push(oid);
466 }
467 break;
468 }
469 }
470 }
471 }
472 Ok((commit_wants, extras))
473}
474
475/// Walk commits from `wants` (excluding `haves`) and collect all
476/// reachable ObjectIds (commits, trees, blobs). Tag wants are peeled, with
477/// the tag objects themselves included in the result.
478///
479/// Pass 1 of the two-pass streaming approach: only OIDs are stored,
480/// not object data.
481fn collect_all_oids(
482 repo: &gix::Repository,
483 wants: &[gix::ObjectId],
484 haves: &[gix::ObjectId],
485) -> std::result::Result<Vec<gix::ObjectId>, Box<dyn std::error::Error + Send + Sync>> {
486 let (commit_wants, extras) = peel_wants(repo, wants)?;
487 let have_set: HashSet<gix::ObjectId> = haves.iter().copied().collect();
488 let mut seen = HashSet::new();
489 let mut oids = Vec::new();
490
491 // Mark have objects as already seen so we skip them
492 for have in haves {
493 seen.insert(*have);
494 }
495
496 let walk = repo
497 .rev_walk(commit_wants.iter().copied())
498 .with_hidden(haves.iter().copied())
499 .all()?;
500
501 for info_result in walk {
502 let info = info_result?;
503 let commit_oid = info.id;
504
505 if have_set.contains(&commit_oid) || !seen.insert(commit_oid) {
506 continue;
507 }
508
509 // Extract tree OID from raw commit bytes (single ODB read)
510 let commit_obj = repo.find_object(commit_oid)?;
511 let tree_oid =
512 gix::objs::CommitRefIter::from_bytes(&commit_obj.data, gix::hash::Kind::Sha1)
513 .tree_id()?;
514
515 oids.push(commit_oid);
516
517 collect_tree_oids(repo, tree_oid, &mut seen, &mut oids)?;
518 }
519
520 // Tag objects (and direct tree/blob tag targets) ride along after the
521 // walk; anything already collected or known to the remote is skipped.
522 oids.extend(
523 extras
524 .into_iter()
525 .filter(|oid| !seen.contains(oid) && !have_set.contains(oid)),
526 );
527
528 Ok(oids)
529}
530
531fn common_haves(
532 repo: &gix::Repository,
533 wants: &[gix::ObjectId],
534 haves: &[gix::ObjectId],
535) -> std::result::Result<Vec<gix::ObjectId>, Box<dyn std::error::Error + Send + Sync>> {
536 let want_set: HashSet<gix::ObjectId> =
537 collect_all_oids(repo, wants, &[])?.into_iter().collect();
538
539 Ok(haves
540 .iter()
541 .copied()
542 .filter(|oid| want_set.contains(oid))
543 .collect())
544}
545
546/// Build a complete, self-contained pack (header + entries + SHA-1 trailer)
547/// of everything reachable from `wants` but not from `haves`, returned as one
548/// buffer. Entries are undeltified base objects — used by the *push*
549/// (send-pack) client, where mirror pushes are infrequent enough that pack
550/// size doesn't justify delta search.
551///
552/// `haves` must exist in the local object database (callers filter the remote
553/// advertisement accordingly).
554pub fn build_raw_pack(
555 repo_path: &Path,
556 wants: &[gix::ObjectId],
557 haves: &[gix::ObjectId],
558) -> Result<Vec<u8>> {
559 let repo = gix::open(repo_path)?;
560 // Tag wants (annotated tags) are peeled inside `collect_all_oids`, with
561 // the tag objects themselves included in the result.
562 let oids = collect_all_oids(&repo, wants, haves)
563 .map_err(|e| Error::Protocol(format!("collecting pack objects: {e}")))?;
564
565 let mut out = Vec::new();
566 out.extend_from_slice(b"PACK");
567 out.extend_from_slice(&2u32.to_be_bytes());
568 out.extend_from_slice(&(oids.len() as u32).to_be_bytes());
569 for oid in &oids {
570 let obj = repo
571 .find_object(*oid)
572 .map_err(|e| Error::Protocol(format!("reading {oid}: {e}")))?;
573 out.extend_from_slice(&build_base_entry(obj.kind, &obj.data));
574 }
575 let checksum = Sha1::digest(&out);
576 out.extend_from_slice(&checksum);
577 Ok(out)
578}
579
580/// Generate the complete pack response for a Git upload-pack request.
581///
582/// Returns an `AsyncRead` producing the side-band-64k framed response that
583/// can be streamed as the HTTP response body.
584pub fn generate_pack(
585 repo_path: &Path,
586 request: &UploadPackRequest,
587) -> Result<impl AsyncRead + Send + Unpin + use<>> {
588 let repo_path = repo_path.to_path_buf();
589 let wants: Vec<gix::ObjectId> = request.wants.clone();
590 let haves: Vec<gix::ObjectId> = request.haves.clone();
591 let object_ids = request.object_ids.clone();
592 let done = request.done;
593 let ofs_delta = request.capabilities.ofs_delta;
594 let multi_ack = request.capabilities.multi_ack;
595 let multi_ack_detailed = request.capabilities.multi_ack_detailed;
596
597 let (tx, rx) = tokio::sync::mpsc::channel::<std::result::Result<Bytes, std::io::Error>>(64);
598
599 let handle = tokio::task::spawn_blocking(move || {
600 if let Err(e) = generate_pack_sync(
601 &repo_path,
602 &wants,
603 &haves,
604 object_ids,
605 GeneratePackOptions {
606 done,
607 ofs_delta,
608 multi_ack,
609 multi_ack_detailed,
610 },
611 &tx,
612 ) {
613 let _ = tx.blocking_send(Err(std::io::Error::other(e.to_string())));
614 }
615 });
616
617 // Log panics from the blocking task without blocking the stream
618 tokio::spawn(async move {
619 if let Err(e) = handle.await {
620 tracing::error!("pack generation task panicked: {e}");
621 }
622 });
623
624 let stream = tokio_stream::wrappers::ReceiverStream::new(rx);
625 Ok(StreamReader::new(stream))
626}
627
628/// Synchronous two-pass streaming pack generator.
629///
630/// Pass 1: collect OIDs only (lightweight -- no object data retained).
631/// Pass 2: re-read each object, compress, and stream it through `tx`.
632struct GeneratePackOptions {
633 done: bool,
634 ofs_delta: bool,
635 multi_ack: bool,
636 multi_ack_detailed: bool,
637}
638
639fn generate_pack_sync(
640 repo_path: &Path,
641 wants: &[gix::ObjectId],
642 haves: &[gix::ObjectId],
643 object_ids: Option<Vec<gix::ObjectId>>,
644 options: GeneratePackOptions,
645 tx: &tokio::sync::mpsc::Sender<std::result::Result<Bytes, std::io::Error>>,
646) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync>> {
647 const MAX_DELTA_BASES: usize = 8;
648 const MIN_DELTA_BLOB_SIZE: usize = 1024;
649
650 let repo = gix::open(repo_path)?;
651
652 let common = if !haves.is_empty() {
653 common_haves(&repo, wants, haves)?
654 } else {
655 Vec::new()
656 };
657
658 if options.multi_ack && !haves.is_empty() && !options.done {
659 for oid in &common {
660 let suffix = if options.multi_ack_detailed {
661 "common"
662 } else {
663 "continue"
664 };
665 send(tx, &encode_ack_line(*oid, Some(suffix)))?;
666 }
667 send(tx, &pktline::encode(b"NAK\n"))?;
668 return Ok(());
669 }
670
671 if options.multi_ack && !common.is_empty() {
672 send(tx, &encode_ack_line(*common.last().unwrap(), None))?;
673 } else {
674 // NAK line
675 send(tx, &pktline::encode(b"NAK\n"))?;
676 }
677
678 // Pass 1: collect OIDs only
679 let oids = match object_ids {
680 Some(oids) => oids,
681 None => collect_all_oids(&repo, wants, haves)?,
682 };
683
684 // Pass 2: stream each object
685 let mut hasher = Sha1::new();
686
687 // Pack header
688 let mut header = Vec::with_capacity(12);
689 header.extend_from_slice(b"PACK");
690 header.extend_from_slice(&2u32.to_be_bytes());
691 header.extend_from_slice(&(oids.len() as u32).to_be_bytes());
692 hasher.update(&header);
693 send_sideband(tx, &header)?;
694
695 let mut pack_offset = header.len() as u64;
696 let mut recent_blob_bases = Vec::<BlobDeltaBase>::new();
697
698 // Each object: read, compress, frame, send
699 for oid in &oids {
700 let obj = repo.find_object(*oid)?;
701 let full_entry = build_base_entry(obj.kind, &obj.data);
702 let mut used_delta = false;
703 let entry = if options.ofs_delta
704 && obj.kind == gix::object::Kind::Blob
705 && obj.data.len() >= MIN_DELTA_BLOB_SIZE
706 {
707 recent_blob_bases
708 .iter()
709 .filter(|base| base.data.len() >= MIN_DELTA_BLOB_SIZE)
710 .filter_map(|base| {
711 build_ofs_delta_entry(pack_offset, base.pack_offset, &base.data, &obj.data)
712 })
713 .min_by_key(Vec::len)
714 .filter(|delta_entry| delta_entry.len() < full_entry.len())
715 .inspect(|_| {
716 used_delta = true;
717 })
718 .unwrap_or(full_entry)
719 } else {
720 full_entry
721 };
722
723 hasher.update(&entry);
724 send_sideband(tx, &entry)?;
725
726 if obj.kind == gix::object::Kind::Blob
727 && !used_delta
728 && obj.data.len() >= MIN_DELTA_BLOB_SIZE
729 {
730 recent_blob_bases.push(BlobDeltaBase {
731 pack_offset,
732 data: obj.data.to_vec(),
733 });
734 if recent_blob_bases.len() > MAX_DELTA_BASES {
735 recent_blob_bases.remove(0);
736 }
737 }
738
739 pack_offset += entry.len() as u64;
740 }
741
742 // SHA-1 checksum over raw pack bytes
743 let checksum = hasher.finalize();
744 send_sideband(tx, &checksum)?;
745
746 // Flush
747 send(tx, b"0000")?;
748
749 Ok(())
750}
751
752#[cfg(test)]
753mod tests {
754 use std::{
755 path::{
756 Path,
757 PathBuf,
758 },
759 process::Command,
760 };
761
762 use tempfile::TempDir;
763 use tokio::io::AsyncReadExt;
764
765 use super::*;
766
767 fn make_pktline(data: &str) -> Vec<u8> {
768 let len = data.len() + 4;
769 format!("{len:04x}{data}").into_bytes()
770 }
771
772 /// Create a bare repo with a single commit on the `main` branch.
773 fn create_repo_with_commit(root: &Path) -> PathBuf {
774 let bare_path = root.join("test.git");
775 let clone_path = root.join("workdir");
776
777 let out = Command::new("git")
778 .args(["init", "--bare", bare_path.to_str().unwrap()])
779 .output()
780 .expect("git init --bare failed");
781 assert!(out.status.success(), "git init --bare failed: {:?}", out);
782
783 let out = Command::new("git")
784 .args(["symbolic-ref", "HEAD", "refs/heads/main"])
785 .current_dir(&bare_path)
786 .output()
787 .expect("git symbolic-ref failed");
788 assert!(out.status.success());
789
790 let out = Command::new("git")
791 .args([
792 "clone",
793 bare_path.to_str().unwrap(),
794 clone_path.to_str().unwrap(),
795 ])
796 .output()
797 .expect("git clone failed");
798 assert!(out.status.success(), "git clone failed: {:?}", out);
799
800 for (key, val) in [("user.name", "Test User"), ("user.email", "test@test.com")] {
801 Command::new("git")
802 .args(["config", key, val])
803 .current_dir(&clone_path)
804 .output()
805 .expect("git config failed");
806 }
807
808 // Create a file and commit
809 std::fs::write(clone_path.join("README.md"), "# Test\n").unwrap();
810
811 Command::new("git")
812 .args(["add", "README.md"])
813 .current_dir(&clone_path)
814 .output()
815 .expect("git add failed");
816
817 let out = Command::new("git")
818 .args(["commit", "-m", "initial commit"])
819 .current_dir(&clone_path)
820 .env("GIT_AUTHOR_NAME", "Test User")
821 .env("GIT_AUTHOR_EMAIL", "test@test.com")
822 .env("GIT_COMMITTER_NAME", "Test User")
823 .env("GIT_COMMITTER_EMAIL", "test@test.com")
824 .output()
825 .expect("git commit failed");
826 assert!(out.status.success(), "git commit failed: {:?}", out);
827
828 let out = Command::new("git")
829 .args(["push", "origin", "main"])
830 .current_dir(&clone_path)
831 .output()
832 .expect("git push failed");
833 assert!(out.status.success(), "git push failed: {:?}", out);
834
835 bare_path
836 }
837
838 #[test]
839 fn parse_simple_want() {
840 let hash = "0000000000000000000000000000000000000001";
841 let mut body = make_pktline(&format!("want {hash}\n"));
842 body.extend_from_slice(b"00000009done\n");
843 let req = UploadPackRequest::parse(&body).unwrap();
844 assert_eq!(req.wants.len(), 1);
845 assert!(req.haves.is_empty());
846 assert!(req.done);
847 assert!(!req.capabilities.ofs_delta);
848 assert_eq!(req.shallow.depth, None);
849 }
850
851 #[test]
852 fn parse_wants_and_haves() {
853 let want = "0000000000000000000000000000000000000001";
854 let have = "0000000000000000000000000000000000000002";
855 let mut body = make_pktline(&format!("want {want}\n"));
856 body.extend_from_slice(b"0000");
857 body.extend_from_slice(&make_pktline(&format!("have {have}\n")));
858 body.extend_from_slice(b"0009done\n");
859 let req = UploadPackRequest::parse(&body).unwrap();
860 assert_eq!(req.wants.len(), 1);
861 assert_eq!(req.haves.len(), 1);
862 assert!(req.done);
863 assert!(!req.capabilities.ofs_delta);
864 assert!(req.shallow.client_shallows.is_empty());
865 }
866
867 #[test]
868 fn parse_ofs_delta_capability() {
869 let hash = "0000000000000000000000000000000000000001";
870 let mut body = make_pktline(&format!("want {hash} side-band-64k ofs-delta\n"));
871 body.extend_from_slice(b"0009done\n");
872 let req = UploadPackRequest::parse(&body).unwrap();
873 assert!(req.capabilities.ofs_delta);
874 }
875
876 #[test]
877 fn parse_multi_ack_capability() {
878 let hash = "0000000000000000000000000000000000000001";
879 let mut body = make_pktline(&format!("want {hash} multi_ack side-band-64k\n"));
880 body.extend_from_slice(b"0009done\n");
881 let req = UploadPackRequest::parse(&body).unwrap();
882 assert!(req.capabilities.multi_ack);
883 }
884
885 #[test]
886 fn parse_multi_ack_detailed_capability() {
887 let hash = "0000000000000000000000000000000000000001";
888 let mut body = make_pktline(&format!("want {hash} multi_ack_detailed side-band-64k\n"));
889 body.extend_from_slice(b"0009done\n");
890 let req = UploadPackRequest::parse(&body).unwrap();
891 assert!(req.capabilities.multi_ack);
892 assert!(req.capabilities.multi_ack_detailed);
893 }
894
895 #[test]
896 fn parse_shallow_request() {
897 let hash = "0000000000000000000000000000000000000001";
898 let mut body = make_pktline(&format!("want {hash}\n"));
899 body.extend_from_slice(&make_pktline("deepen 2\n"));
900 body.extend_from_slice(&make_pktline(&format!("shallow {hash}\n")));
901 body.extend_from_slice(&make_pktline("deepen-relative\n"));
902 body.extend_from_slice(b"0009done\n");
903 let req = UploadPackRequest::parse(&body).unwrap();
904 assert_eq!(req.shallow.depth, Some(2));
905 assert_eq!(
906 req.shallow.client_shallows,
907 vec![gix::ObjectId::from_hex(hash.as_bytes()).unwrap()]
908 );
909 assert!(req.shallow.deepen_relative);
910 }
911
912 #[tokio::test]
913 async fn generate_pack_for_clone() {
914 let dir = TempDir::new().unwrap();
915 let repo_path = create_repo_with_commit(dir.path());
916
917 // Get HEAD OID
918 let repo = gix::open(&repo_path).unwrap();
919 let head_oid = repo.head_id().unwrap().detach();
920 drop(repo);
921
922 let request = UploadPackRequest {
923 wants: vec![head_oid],
924 haves: vec![],
925 done: true,
926 capabilities: UploadPackCapabilities::default(),
927 shallow: ShallowRequest::default(),
928 object_ids: None,
929 };
930
931 let mut reader = generate_pack(&repo_path, &request).unwrap();
932 let mut buf = Vec::new();
933 reader.read_to_end(&mut buf).await.unwrap();
934
935 let response = String::from_utf8_lossy(&buf);
936 assert!(
937 response.contains("NAK"),
938 "response should contain NAK: {response:?}"
939 );
940
941 // Find PACK signature in the binary response
942 let pack_found = buf.windows(4).any(|window| window == b"PACK");
943 assert!(pack_found, "response should contain PACK signature");
944 }
945
946 /// Regression: a want may be an annotated tag *object* (protocol-v2
947 /// `ls-refs` advertises unpeeled ids). The walk used to choke on it
948 /// mid-stream, killing `git clone` of any repo with an annotated tag;
949 /// now the tag is peeled and the tag object itself joins the pack.
950 #[tokio::test]
951 async fn tag_object_wants_are_peeled_and_packed() {
952 let dir = TempDir::new().unwrap();
953 let repo_path = create_repo_with_commit(dir.path());
954
955 let clone_path = dir.path().join("workdir");
956 let out = Command::new("git")
957 .args(["tag", "-a", "-m", "annotated", "v1"])
958 .current_dir(&clone_path)
959 .env("GIT_AUTHOR_NAME", "t")
960 .env("GIT_AUTHOR_EMAIL", "t@t")
961 .env("GIT_COMMITTER_NAME", "t")
962 .env("GIT_COMMITTER_EMAIL", "t@t")
963 .output()
964 .unwrap();
965 assert!(out.status.success());
966 let out = Command::new("git")
967 .args(["push", "-q", "origin", "v1"])
968 .current_dir(&clone_path)
969 .output()
970 .unwrap();
971 assert!(out.status.success());
972
973 let repo = gix::open(&repo_path).unwrap();
974 let tag_oid = repo
975 .find_reference("refs/tags/v1")
976 .unwrap()
977 .try_id()
978 .unwrap()
979 .detach();
980 assert_eq!(
981 repo.find_object(tag_oid).unwrap().kind,
982 gix::object::Kind::Tag,
983 "fixture must produce a real tag object"
984 );
985 drop(repo);
986
987 let pack = build_raw_pack(&repo_path, &[tag_oid], &[]).unwrap();
988 assert!(pack.starts_with(b"PACK"));
989 let count = u32::from_be_bytes(pack[8..12].try_into().unwrap());
990 // tag + commit + tree + blob
991 assert_eq!(count, 4, "tag object and full closure packed");
992
993 // The streaming (fetch-serving) path accepts the same want.
994 let request = UploadPackRequest {
995 wants: vec![tag_oid],
996 haves: vec![],
997 done: true,
998 capabilities: UploadPackCapabilities::default(),
999 shallow: ShallowRequest::default(),
1000 object_ids: None,
1001 };
1002 let mut reader = generate_pack(&repo_path, &request).unwrap();
1003 let mut buf = Vec::new();
1004 reader.read_to_end(&mut buf).await.unwrap();
1005 assert!(
1006 buf.windows(4).any(|w| w == b"PACK"),
1007 "fetch response should contain a pack, got: {:?}",
1008 String::from_utf8_lossy(&buf[..buf.len().min(200)])
1009 );
1010 }
1011}