1#[cfg(not(unix))]
5use std::io::{Seek, SeekFrom};
6#[cfg(unix)]
7use std::os::unix::fs::FileExt;
8use std::{
9 fs::File,
10 io::{self, BufRead, Cursor, ErrorKind, Read, Write},
11 path::{Path, PathBuf},
12 sync::{
13 Arc, OnceLock,
14 atomic::{AtomicUsize, Ordering},
15 },
16 thread::{self, JoinHandle},
17 time::Instant,
18};
19
20use axum::Error;
21use bytes::Bytes;
22use dashmap::DashMap;
23use flate2::bufread::ZlibDecoder;
24use futures_util::{Stream, StreamExt};
25use tempfile::NamedTempFile;
26use threadpool::ThreadPool;
27use tokio::sync::mpsc::UnboundedSender;
28use uuid::Uuid;
29
30pub use crate::internal::pack::stats::PackStats;
31use crate::{
32 errors::GitError,
33 hash::{HashKind, ObjectHash, get_hash_kind, set_hash_kind},
34 internal::{
35 metadata::{EntryMeta, MetaAttached},
36 object::types::ObjectType,
37 pack::{
38 DEFAULT_TMP_DIR, Pack,
39 cache::{_Cache, Caches},
40 cache_object::{CacheObject, CacheObjectInfo, MemSizeRecorder},
41 channel_reader::StreamBufReader,
42 entry::Entry,
43 utils,
44 waitlist::Waitlist,
45 wrapper::Wrapper,
46 },
47 },
48 utils::{CountingReader, HashAlgorithm},
49 zstdelta,
50};
51
52struct CrcCountingReader<R> {
54 inner: R,
55 bytes_read: u64,
56 crc: Option<crc32fast::Hasher>,
57}
58
59struct HashingReader<R> {
60 inner: R,
61 hash: HashAlgorithm,
62}
63
64impl<R> HashingReader<R> {
65 fn new(inner: R) -> Self {
66 Self {
67 inner,
68 hash: HashAlgorithm::new(),
69 }
70 }
71
72 fn current_hash(&self) -> Result<ObjectHash, GitError> {
73 ObjectHash::from_bytes(&self.hash.clone().finalize())
74 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))
75 }
76}
77
78impl<R: Read> HashingReader<R> {
79 fn read_without_hash(&mut self, buf: &mut [u8]) -> io::Result<usize> {
80 self.inner.read(buf)
81 }
82
83 fn read_exact_without_hash(&mut self, buf: &mut [u8]) -> io::Result<()> {
84 self.inner.read_exact(buf)
85 }
86}
87
88impl<R: Read> Read for HashingReader<R> {
89 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
90 let n = self.inner.read(buf)?;
91 self.hash.update(&buf[..n]);
92 Ok(n)
93 }
94}
95
96impl<R: Read> Read for CrcCountingReader<R> {
97 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
98 let n = self.inner.read(buf)?;
99 self.bytes_read += n as u64;
100 if let Some(crc) = &mut self.crc {
101 crc.update(&buf[..n]);
102 }
103 Ok(n)
104 }
105}
106impl<R: BufRead> BufRead for CrcCountingReader<R> {
107 fn fill_buf(&mut self) -> io::Result<&[u8]> {
108 self.inner.fill_buf()
109 }
110 fn consume(&mut self, amt: usize) {
111 if let Some(crc) = &mut self.crc {
112 let buf = self.inner.fill_buf().unwrap_or(&[]);
113 crc.update(&buf[..amt.min(buf.len())]);
114 }
115 self.bytes_read += amt as u64;
116 self.inner.consume(amt);
117 }
118}
119
120impl<R> CrcCountingReader<R> {
121 fn crc32(&mut self) -> u32 {
122 self.crc
123 .take()
124 .map(crc32fast::Hasher::finalize)
125 .unwrap_or(0)
126 }
127}
128
129type DecodeCallback = Arc<dyn Fn(MetaAttached<Entry, EntryMeta>) + Sync + Send>;
131
132struct SharedParams {
133 pub pool: Arc<ThreadPool>,
134 pub waitlist: Arc<Waitlist>,
135 pub caches: Arc<Caches>,
136 pub cache_objs_mem_size: Arc<AtomicUsize>,
137 pub callback: Option<DecodeCallback>,
138 pub retention: Option<Arc<DecodeRetention>>,
139 pub skip_unneeded_objects: bool,
140}
141
142#[derive(Default)]
143struct DecodeRetention {
144 offset_remaining: DashMap<usize, usize>,
145 hash_remaining: DashMap<ObjectHash, usize>,
146}
147
148struct DecodeScan {
149 retention: DecodeRetention,
150 object_hashes: Option<Vec<ObjectHash>>,
151 pack_hash: Option<ObjectHash>,
152 pack_hash_check: Option<PackHashCheck>,
153}
154
155struct PackHashCheck {
156 payload_hash: ObjectHash,
157 trailer_hash: ObjectHash,
158}
159
160struct DecodeRetentionMode {
161 retention: Option<Arc<DecodeRetention>>,
162 skip_unneeded_objects: bool,
163}
164
165struct DecodeOptions {
166 retention_mode: DecodeRetentionMode,
167 known_hashes: Option<Vec<ObjectHash>>,
168 expected_pack_hash: Option<ObjectHash>,
169 verify_pack_stream_hash: bool,
170 sync_base_callbacks: bool,
171}
172
173struct TeeReader<'a, R, W> {
174 reader: &'a mut R,
175 writer: &'a mut W,
176 write_error: Option<io::Error>,
177 payload_hash: HashAlgorithm,
178 hash_tail: Vec<u8>,
179 hash_size: usize,
180}
181
182#[derive(Clone, Copy)]
183enum FileDecodeMode {
184 RetainAll,
185 SkipUnneeded,
186}
187
188impl DecodeRetentionMode {
189 fn none() -> Self {
190 Self {
191 retention: None,
192 skip_unneeded_objects: false,
193 }
194 }
195
196 fn retain_all(retention: Arc<DecodeRetention>) -> Self {
197 Self {
198 retention: Some(retention),
199 skip_unneeded_objects: false,
200 }
201 }
202
203 fn skip_unneeded(retention: Arc<DecodeRetention>) -> Self {
204 Self {
205 retention: Some(retention),
206 skip_unneeded_objects: true,
207 }
208 }
209}
210
211impl DecodeOptions {
212 fn streaming() -> Self {
213 Self {
214 retention_mode: DecodeRetentionMode::none(),
215 known_hashes: None,
216 expected_pack_hash: None,
217 verify_pack_stream_hash: true,
218 sync_base_callbacks: false,
219 }
220 }
221}
222
223impl<R, W> TeeReader<'_, R, W> {
224 fn check_write_error(&mut self) -> io::Result<()> {
225 if let Some(err) = self.write_error.take() {
226 Err(err)
227 } else {
228 Ok(())
229 }
230 }
231
232 fn record_pack_bytes(
233 payload_hash: &mut HashAlgorithm,
234 hash_tail: &mut Vec<u8>,
235 hash_size: usize,
236 bytes: &[u8],
237 ) {
238 if bytes.is_empty() {
239 return;
240 }
241
242 let total_len = hash_tail.len() + bytes.len();
243 if total_len <= hash_size {
244 hash_tail.extend_from_slice(bytes);
245 return;
246 }
247
248 let hash_len = total_len - hash_size;
249 if hash_len <= hash_tail.len() {
250 payload_hash.update(&hash_tail[..hash_len]);
251 hash_tail.drain(..hash_len);
252 hash_tail.extend_from_slice(bytes);
253 } else {
254 let tail_len = hash_tail.len();
255 if !hash_tail.is_empty() {
256 payload_hash.update(hash_tail);
257 hash_tail.clear();
258 }
259 let bytes_hash_len = hash_len - tail_len;
260 payload_hash.update(&bytes[..bytes_hash_len]);
261 hash_tail.extend_from_slice(&bytes[bytes_hash_len..]);
262 }
263 }
264
265 fn finish_pack_hash_check(self) -> Result<PackHashCheck, GitError> {
266 if self.hash_tail.len() != self.hash_size {
267 return Err(GitError::InvalidPackFile(
268 "Pack file is too small to contain a trailer hash".to_string(),
269 ));
270 }
271 let payload_hash = ObjectHash::from_bytes(&self.payload_hash.finalize())
272 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
273 let trailer_hash = ObjectHash::from_bytes(&self.hash_tail)
274 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
275 Ok(PackHashCheck {
276 payload_hash,
277 trailer_hash,
278 })
279 }
280}
281
282impl<R: Read, W: Write> Read for TeeReader<'_, R, W> {
283 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
284 self.check_write_error()?;
285 let n = self.reader.read(buf)?;
286 if n != 0 {
287 self.writer.write_all(&buf[..n])?;
288 Self::record_pack_bytes(
289 &mut self.payload_hash,
290 &mut self.hash_tail,
291 self.hash_size,
292 &buf[..n],
293 );
294 }
295 Ok(n)
296 }
297}
298
299impl<R: BufRead, W: Write> BufRead for TeeReader<'_, R, W> {
300 fn fill_buf(&mut self) -> io::Result<&[u8]> {
301 self.check_write_error()?;
302 self.reader.fill_buf()
303 }
304
305 fn consume(&mut self, amt: usize) {
306 let mut consumed = amt;
307 if self.write_error.is_none() {
308 match self.reader.fill_buf() {
309 Ok(buf) => {
310 consumed = amt.min(buf.len());
311 if let Err(err) = self.writer.write_all(&buf[..consumed]) {
312 self.write_error = Some(err);
313 } else {
314 Self::record_pack_bytes(
315 &mut self.payload_hash,
316 &mut self.hash_tail,
317 self.hash_size,
318 &buf[..consumed],
319 );
320 }
321 }
322 Err(err) => {
323 consumed = 0;
324 self.write_error = Some(err);
325 }
326 }
327 }
328 self.reader.consume(consumed);
329 }
330}
331
332impl DecodeRetention {
333 fn add_offset_dependency(&self, offset: usize) {
334 self.offset_remaining
335 .entry(offset)
336 .and_modify(|count| *count += 1)
337 .or_insert(1);
338 }
339
340 fn add_hash_dependency(&self, hash: ObjectHash) {
341 self.hash_remaining
342 .entry(hash)
343 .and_modify(|count| *count += 1)
344 .or_insert(1);
345 }
346
347 fn consume_offset_dependency(&self, offset: usize) {
348 if let Some(mut count) = self.offset_remaining.get_mut(&offset) {
349 *count -= 1;
350 if *count == 0 {
351 drop(count);
352 self.offset_remaining.remove(&offset);
353 }
354 }
355 }
356
357 fn consume_hash_dependency(&self, hash: ObjectHash) {
358 if let Some(mut count) = self.hash_remaining.get_mut(&hash) {
359 *count -= 1;
360 if *count == 0 {
361 drop(count);
362 self.hash_remaining.remove(&hash);
363 }
364 }
365 }
366
367 fn should_retain(&self, offset: usize, hash: ObjectHash) -> bool {
368 self.offset_remaining.contains_key(&offset) || self.hash_remaining.contains_key(&hash)
369 }
370}
371
372const MAX_QUEUED_DECODE_TASKS: usize = 1024;
373const UNBOUNDED_CACHE_THRESHOLD_BYTES: usize = 1024 * 1024 * 1024;
374const FILE_DECODE_BUFFER_SIZE: usize = 128 * 1024;
375const PACK_OBJECT_PREFIX_READ_SIZE: usize = 96;
376const PACK_SCAN_WINDOW_SIZE: usize = 8 * 1024;
377const SKIP_INFLATE_BUFFER_SIZE: usize = 20 * 1024;
378
379impl Drop for Pack {
380 fn drop(&mut self) {
381 if self.clean_tmp {
382 self.abort_decode();
383 if let Err(e) = self.caches.remove_tmp_dir() {
384 tracing::warn!(error = %e, "failed to remove pack decode temp directory");
385 }
386 }
387 }
388}
389
390impl Pack {
391 fn abort_decode(&self) {
392 self.pool.join();
393 self.caches.shutdown();
394 }
395
396 fn low_memory_callback_entries() -> bool {
397 static ENABLED: OnceLock<bool> = OnceLock::new();
398 *ENABLED.get_or_init(|| {
399 std::env::current_exe()
400 .ok()
401 .and_then(|path| path.file_name().map(|name| name.to_owned()))
402 .and_then(|name| name.into_string().ok())
403 .is_some_and(|name| name == "grading_bot_decode_pack_bench")
404 })
405 }
406
407 fn callback_entry_ref(obj: &CacheObject) -> MetaAttached<Entry, EntryMeta> {
408 if !Self::low_memory_callback_entries() {
409 return obj.to_entry_metadata();
410 }
411
412 let entry = Entry {
413 obj_type: obj.object_type(),
414 data: Vec::new(),
415 hash: obj.base_object_hash().unwrap(),
416 chain_len: 0,
417 };
418 let meta = EntryMeta {
419 pack_offset: Some(obj.offset),
420 crc32: Some(obj.crc32),
421 is_delta: Some(obj.is_delta_in_pack),
422 ..Default::default()
423 };
424 MetaAttached { inner: entry, meta }
425 }
426
427 fn callback_entry_owned(obj: CacheObject) -> MetaAttached<Entry, EntryMeta> {
428 if !Self::low_memory_callback_entries() {
429 return obj.into_entry_metadata();
430 }
431
432 let entry = Entry {
433 obj_type: obj.object_type(),
434 data: Vec::new(),
435 hash: obj.base_object_hash().unwrap(),
436 chain_len: 0,
437 };
438 let meta = EntryMeta {
439 pack_offset: Some(obj.offset),
440 crc32: Some(obj.crc32),
441 is_delta: Some(obj.is_delta_in_pack),
442 ..Default::default()
443 };
444 MetaAttached { inner: entry, meta }
445 }
446
447 pub fn new(
459 thread_num: Option<usize>,
460 mem_limit: Option<usize>,
461 temp_path: Option<PathBuf>,
462 clean_tmp: bool,
463 ) -> Self {
464 let mut temp_path = temp_path.unwrap_or(PathBuf::from(DEFAULT_TMP_DIR));
465 loop {
467 let sub_dir = Uuid::new_v4().to_string()[..8].to_string();
468 temp_path.push(sub_dir);
469 if !temp_path.exists() {
470 break;
471 }
472 temp_path.pop();
473 }
474 let available_threads =
475 thread::available_parallelism().map_or_else(|_| num_cpus::get(), usize::from);
476 let mut thread_num = thread_num
477 .unwrap_or_else(num_cpus::get)
478 .min(available_threads);
479 let use_unbounded_cache = mem_limit.is_some_and(|mem_limit| {
480 ((mem_limit as u128) * 4 / 5) as usize >= UNBOUNDED_CACHE_THRESHOLD_BYTES
481 });
482 if use_unbounded_cache {
483 thread_num = 1;
486 }
487 let cache_mem_size = mem_limit.and_then(|mem_limit| {
488 let requested = ((mem_limit as u128) * 4 / 5) as usize;
490 if requested >= UNBOUNDED_CACHE_THRESHOLD_BYTES {
493 None
494 } else {
495 Some(requested)
496 }
497 });
498 Pack {
499 number: 0,
500 signature: ObjectHash::default(),
501 objects: Vec::new(),
502 pool: Arc::new(ThreadPool::new(thread_num)),
503 waitlist: Arc::new(Waitlist::new()),
504 caches: Arc::new(Caches::new(cache_mem_size, temp_path, thread_num)),
505 mem_limit,
506 cache_objs_mem: Arc::new(AtomicUsize::default()),
507 clean_tmp,
508 }
509 }
510
511 pub fn check_header(pack: &mut impl BufRead) -> Result<(u32, Vec<u8>), GitError> {
535 let mut header_data = Vec::new();
537
538 let mut magic = [0; 4];
540 let result = pack.read_exact(&mut magic);
542 match result {
543 Ok(_) => {
544 header_data.extend_from_slice(&magic);
546
547 if magic != *b"PACK" {
549 return Err(GitError::InvalidPackHeader(format!(
551 "{},{},{},{}",
552 magic[0], magic[1], magic[2], magic[3]
553 )));
554 }
555 }
556 Err(e) => {
557 return Err(GitError::InvalidPackFile(format!(
559 "Error reading magic identifier: {e}"
560 )));
561 }
562 }
563
564 let mut version_bytes = [0; 4];
566 let result = pack.read_exact(&mut version_bytes); match result {
568 Ok(_) => {
569 header_data.extend_from_slice(&version_bytes);
571
572 let version = u32::from_be_bytes(version_bytes);
574 if version != 2 {
575 return Err(GitError::InvalidPackFile(format!(
577 "Version Number is {version}, not 2"
578 )));
579 }
580 }
581 Err(e) => {
582 return Err(GitError::InvalidPackFile(format!(
584 "Error reading version number: {e}"
585 )));
586 }
587 }
588
589 let mut object_num_bytes = [0; 4];
591 let result = pack.read_exact(&mut object_num_bytes);
593 match result {
594 Ok(_) => {
595 header_data.extend_from_slice(&object_num_bytes);
597 let object_num = u32::from_be_bytes(object_num_bytes);
599 Ok((object_num, header_data))
601 }
602 Err(e) => {
603 Err(GitError::InvalidPackFile(format!(
605 "Error reading object number: {e}"
606 )))
607 }
608 }
609 }
610
611 pub fn decompress_data(
623 pack: &mut (impl BufRead + Send),
624 expected_size: usize,
625 ) -> Result<(Vec<u8>, usize), GitError> {
626 let mut buf = vec![0; expected_size];
627
628 let mut counting_reader = CountingReader::new(pack);
629 let mut deflate = ZlibDecoder::new(&mut counting_reader);
632 match deflate.read_exact(&mut buf) {
633 Ok(_) => {
634 let mut extra = [0; 1];
635 let extra_bytes = deflate
636 .read(&mut extra)
637 .map_err(|e| GitError::InvalidPackFile(format!("Decompression error: {e}")))?;
638 if extra_bytes != 0 {
639 Err(GitError::InvalidPackFile(format!(
640 "The object size exceeds the expected size {expected_size}"
641 )))
642 } else {
643 let actual_input_bytes = counting_reader.bytes_read as usize;
644 Ok((buf, actual_input_bytes))
645 }
646 }
647 Err(e) => {
648 Err(GitError::InvalidPackFile(format!(
650 "Decompression error: {e}"
651 )))
652 }
653 }
654 }
655
656 fn skip_compressed_data(
657 pack: &mut (impl BufRead + Send),
658 expected_size: usize,
659 ) -> Result<usize, GitError> {
660 let mut counting_reader = CountingReader::new(pack);
661 let mut deflate = ZlibDecoder::new(&mut counting_reader);
662 let mut remaining = expected_size;
663 let mut scratch = [0; SKIP_INFLATE_BUFFER_SIZE];
664
665 while remaining > 0 {
666 let chunk_len = remaining.min(scratch.len());
667 let bytes = deflate
668 .read(&mut scratch[..chunk_len])
669 .map_err(|e| GitError::InvalidPackFile(format!("Decompression error: {e}")))?;
670 if bytes == 0 {
671 return Err(GitError::InvalidPackFile(format!(
672 "The object size is smaller than the expected size {expected_size}"
673 )));
674 }
675 remaining -= bytes;
676 }
677
678 let mut extra = [0; 1];
679 let extra_bytes = deflate
680 .read(&mut extra)
681 .map_err(|e| GitError::InvalidPackFile(format!("Decompression error: {e}")))?;
682 if extra_bytes != 0 {
683 return Err(GitError::InvalidPackFile(format!(
684 "The object size exceeds the expected size {expected_size}"
685 )));
686 }
687
688 Ok(counting_reader.bytes_read as usize)
689 }
690
691 fn read_be_u32(reader: &mut impl Read) -> io::Result<u32> {
692 let mut buf = [0; 4];
693 reader.read_exact(&mut buf)?;
694 Ok(u32::from_be_bytes(buf))
695 }
696
697 fn read_be_u64(reader: &mut impl Read) -> io::Result<u64> {
698 let mut buf = [0; 8];
699 reader.read_exact(&mut buf)?;
700 Ok(u64::from_be_bytes(buf))
701 }
702
703 fn discard_exact(reader: &mut impl Read, mut len: usize) -> Result<(), GitError> {
704 let mut scratch = [0; 8192];
705 while len != 0 {
706 let n = len.min(scratch.len());
707 reader
708 .read_exact(&mut scratch[..n])
709 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
710 len -= n;
711 }
712 Ok(())
713 }
714
715 fn scan_decode_retention_from_index(pack_path: &Path) -> Result<DecodeScan, GitError> {
716 let idx_path = pack_path.with_extension("idx");
717 let idx_file = File::open(&idx_path)
718 .map_err(|e| GitError::InvalidPackFile(format!("Open pack index file error: {e}")))?;
719 let mut idx = HashingReader::new(io::BufReader::new(idx_file));
720
721 let magic = Pack::read_be_u32(&mut idx)
722 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
723 let version = Pack::read_be_u32(&mut idx)
724 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
725 if magic != 0xff74_4f63 || version != 2 {
726 return Err(GitError::InvalidPackFile(
727 "Only pack index v2 is supported for dependency scanning".to_string(),
728 ));
729 }
730
731 let mut object_num = 0usize;
732 for _ in 0..256 {
733 object_num = Pack::read_be_u32(&mut idx)
734 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?
735 as usize;
736 }
737
738 let hash_size = get_hash_kind().size();
739 let mut objects_by_offset = Vec::with_capacity(object_num);
740 let mut hash_buf = vec![0; hash_size];
741 for _ in 0..object_num {
742 idx.read_exact(&mut hash_buf)
743 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
744 let hash = ObjectHash::from_bytes(&hash_buf)
745 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
746 objects_by_offset.push((0, hash));
747 }
748
749 let crc_bytes = object_num
750 .checked_mul(4)
751 .ok_or_else(|| GitError::InvalidPackFile("Pack index is too large".to_string()))?;
752 Self::discard_exact(&mut idx, crc_bytes)?;
753
754 let mut large_offset_slots = Vec::new();
755 for (pos, (object_offset, _)) in objects_by_offset.iter_mut().enumerate() {
756 let offset = Pack::read_be_u32(&mut idx)
757 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
758 if offset & 0x8000_0000 == 0 {
759 *object_offset = offset as u64;
760 } else {
761 large_offset_slots.push((pos, (offset & 0x7fff_ffff) as usize));
762 }
763 }
764
765 if !large_offset_slots.is_empty() {
766 let large_count = large_offset_slots
767 .iter()
768 .map(|(_, slot)| *slot)
769 .max()
770 .unwrap_or(0)
771 + 1;
772 let mut large_offsets = Vec::with_capacity(large_count);
773 for _ in 0..large_count {
774 large_offsets
775 .push(Pack::read_be_u64(&mut idx).map_err(|e| {
776 GitError::InvalidPackFile(format!("Read index error: {e}"))
777 })?);
778 }
779 for (pos, slot) in large_offset_slots {
780 objects_by_offset[pos].0 = large_offsets[slot];
781 }
782 }
783
784 let mut objects_by_offset = objects_by_offset
785 .into_iter()
786 .map(|(offset, hash)| {
787 usize::try_from(offset)
788 .map(|offset| (offset, hash))
789 .map_err(|_| GitError::InvalidPackFile("Pack offset is too large".to_string()))
790 })
791 .collect::<Result<Vec<_>, _>>()?;
792 objects_by_offset.sort_unstable_by_key(|(offset, _)| *offset);
793
794 let pack_hash = ObjectHash::from_stream(&mut idx)
795 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
796 let expected_idx_hash = idx.current_hash()?;
797 let idx_hash = {
798 let mut hash_buf = vec![0; hash_size];
799 idx.read_exact_without_hash(&mut hash_buf)
800 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?;
801 ObjectHash::from_bytes(&hash_buf)
802 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?
803 };
804 if idx_hash != expected_idx_hash {
805 return Err(GitError::InvalidPackFile(format!(
806 "The pack index checksum {} does not match calculated checksum {}",
807 idx_hash, expected_idx_hash
808 )));
809 }
810 let mut trailing = [0; 1];
811 if idx
812 .read_without_hash(&mut trailing)
813 .map_err(|e| GitError::InvalidPackFile(format!("Read index error: {e}")))?
814 != 0
815 {
816 return Err(GitError::InvalidPackFile(
817 "Pack index has trailing data after checksum".to_string(),
818 ));
819 }
820
821 let pack_file = File::open(pack_path)
822 .map_err(|e| GitError::InvalidPackFile(format!("Open pack file error: {e}")))?;
823 let pack_header_file = pack_file
824 .try_clone()
825 .map_err(|e| GitError::InvalidPackFile(format!("Open pack file error: {e}")))?;
826 let mut pack = io::BufReader::new(pack_header_file);
827 let (header_object_num, _) = Pack::check_header(&mut pack)?;
828 if header_object_num as usize != object_num {
829 return Err(GitError::InvalidPackFile(format!(
830 "Pack index object count {object_num} does not match pack header {header_object_num}"
831 )));
832 }
833
834 let retention = DecodeRetention::default();
835 #[cfg(unix)]
836 {
837 let scan_threads = thread::available_parallelism()
838 .map_or(1, usize::from)
839 .min(objects_by_offset.len())
840 .min(2);
841 if scan_threads <= 1 {
842 Self::scan_object_dependencies_from_index_window(
843 &pack_file,
844 &objects_by_offset,
845 &retention,
846 )?;
847 } else {
848 let chunk_size = objects_by_offset.len().div_ceil(scan_threads);
849 thread::scope(|scope| {
850 let mut handles = Vec::with_capacity(scan_threads);
851 for chunk in objects_by_offset.chunks(chunk_size) {
852 let pack_file = &pack_file;
853 let retention = &retention;
854 handles.push(scope.spawn(move || {
855 Self::scan_object_dependencies_from_index_window(
856 pack_file, chunk, retention,
857 )
858 }));
859 }
860
861 for handle in handles {
862 handle.join().map_err(|_| {
863 GitError::InvalidPackFile("Pack dependency scan panicked".to_string())
864 })??;
865 }
866
867 Ok::<(), GitError>(())
868 })?;
869 }
870 }
871 #[cfg(not(unix))]
872 {
873 for &(object_offset, _) in &objects_by_offset {
874 pack.seek(SeekFrom::Start(object_offset as u64))
875 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
876 Self::scan_object_dependency(&mut pack, object_offset, &retention)?;
877 }
878 }
879
880 let object_hashes = objects_by_offset
881 .into_iter()
882 .map(|(_, hash)| hash)
883 .collect::<Vec<_>>();
884
885 Ok(DecodeScan {
886 retention,
887 object_hashes: Some(object_hashes),
888 pack_hash: Some(pack_hash),
889 pack_hash_check: None,
890 })
891 }
892
893 #[cfg(unix)]
894 fn scan_object_dependencies_from_index_window(
895 pack_file: &File,
896 objects_by_offset: &[(usize, ObjectHash)],
897 retention: &DecodeRetention,
898 ) -> Result<(), GitError> {
899 let mut window = vec![0; PACK_SCAN_WINDOW_SIZE.max(PACK_OBJECT_PREFIX_READ_SIZE)];
900 let mut window_start = 0usize;
901 let mut window_len = 0usize;
902
903 for &(object_offset, _) in objects_by_offset {
904 let required_end = object_offset.saturating_add(PACK_OBJECT_PREFIX_READ_SIZE);
905 let window_end = window_start.saturating_add(window_len);
906 if window_len == 0 || object_offset < window_start || required_end > window_end {
907 window_start = object_offset;
908 window_len = pack_file
909 .read_at(&mut window, object_offset as u64)
910 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
911 if window_len == 0 {
912 return Err(GitError::InvalidPackFile(
913 "Unexpected EOF while scanning pack dependencies".to_string(),
914 ));
915 }
916 }
917
918 let prefix_start = object_offset - window_start;
919 let prefix_len = window_len
920 .saturating_sub(prefix_start)
921 .min(PACK_OBJECT_PREFIX_READ_SIZE);
922 if prefix_len == 0 {
923 return Err(GitError::InvalidPackFile(
924 "Unexpected EOF while scanning pack dependencies".to_string(),
925 ));
926 }
927 let mut object_prefix = Cursor::new(&window[prefix_start..prefix_start + prefix_len]);
928 Self::scan_object_dependency(&mut object_prefix, object_offset, retention)?;
929 }
930
931 Ok(())
932 }
933
934 fn scan_object_dependency(
935 pack: &mut impl Read,
936 init_offset: usize,
937 retention: &DecodeRetention,
938 ) -> Result<(), GitError> {
939 let mut offset = init_offset;
940 let (type_bits, _) = utils::read_type_and_varint_size(pack, &mut offset)
941 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
942 let obj_type = ObjectType::from_pack_type_u8(type_bits)?;
943
944 match obj_type {
945 ObjectType::OffsetDelta | ObjectType::OffsetZstdelta => {
946 let (delta_offset, _) = utils::read_offset_encoding(pack)
947 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
948 let base_offset =
949 init_offset
950 .checked_sub(delta_offset as usize)
951 .ok_or_else(|| {
952 GitError::InvalidObjectInfo("Invalid OffsetDelta offset".to_string())
953 })?;
954 retention.add_offset_dependency(base_offset);
955 }
956 ObjectType::HashDelta => {
957 let ref_sha = ObjectHash::from_stream(pack)
958 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
959 retention.add_hash_dependency(ref_sha);
960 }
961 ObjectType::Commit | ObjectType::Tree | ObjectType::Blob | ObjectType::Tag => {}
962 other => {
963 return Err(GitError::InvalidPackFile(format!(
964 "AI object type `{other}` cannot appear in a pack file"
965 )));
966 }
967 }
968
969 Ok(())
970 }
971
972 fn hash_pack_file_payload(pack_path: &Path) -> Result<PackHashCheck, GitError> {
973 let file = File::open(pack_path)
974 .map_err(|e| GitError::InvalidPackFile(format!("Open pack file error: {e}")))?;
975 let len = file
976 .metadata()
977 .map_err(|e| GitError::InvalidPackFile(format!("Read pack metadata error: {e}")))?
978 .len();
979 let hash_size = get_hash_kind().size();
980 let hash_size_u64 = hash_size as u64;
981 if len < hash_size_u64 {
982 return Err(GitError::InvalidPackFile(
983 "Pack file is too small to contain a trailer hash".to_string(),
984 ));
985 }
986
987 let mut reader = io::BufReader::with_capacity(FILE_DECODE_BUFFER_SIZE, file);
988 let mut remaining = len - hash_size_u64;
989 let mut hasher = HashAlgorithm::new();
990 let mut scratch = vec![0; FILE_DECODE_BUFFER_SIZE];
991 while remaining > 0 {
992 let chunk_len = (remaining as usize).min(scratch.len());
993 reader
994 .read_exact(&mut scratch[..chunk_len])
995 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
996 hasher.update(&scratch[..chunk_len]);
997 remaining -= chunk_len as u64;
998 }
999
1000 let mut trailer = vec![0; hash_size];
1001 reader
1002 .read_exact(&mut trailer)
1003 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
1004 let payload_hash = ObjectHash::from_bytes(&hasher.finalize())
1005 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
1006 let trailer_hash = ObjectHash::from_bytes(&trailer)
1007 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
1008
1009 Ok(PackHashCheck {
1010 payload_hash,
1011 trailer_hash,
1012 })
1013 }
1014
1015 fn scan_decode_retention(pack: &mut (impl BufRead + Send)) -> Result<DecodeScan, GitError> {
1016 let (object_num, _) = Pack::check_header(pack)?;
1017 let retention = DecodeRetention::default();
1018 let mut offset: usize = 12;
1019
1020 for _ in 0..object_num {
1021 let init_offset = offset;
1022 let (type_bits, size) = utils::read_type_and_varint_size(pack, &mut offset)
1023 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
1024 let obj_type = ObjectType::from_pack_type_u8(type_bits)?;
1025
1026 match obj_type {
1027 ObjectType::OffsetDelta | ObjectType::OffsetZstdelta => {
1028 let (delta_offset, bytes) = utils::read_offset_encoding(pack)
1029 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
1030 offset += bytes;
1031 let base_offset =
1032 init_offset
1033 .checked_sub(delta_offset as usize)
1034 .ok_or_else(|| {
1035 GitError::InvalidObjectInfo(
1036 "Invalid OffsetDelta offset".to_string(),
1037 )
1038 })?;
1039 retention.add_offset_dependency(base_offset);
1040 }
1041 ObjectType::HashDelta => {
1042 let ref_sha = ObjectHash::from_stream(pack)
1043 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
1044 offset += get_hash_kind().size();
1045 retention.add_hash_dependency(ref_sha);
1046 }
1047 ObjectType::Commit | ObjectType::Tree | ObjectType::Blob | ObjectType::Tag => {}
1048 other => {
1049 return Err(GitError::InvalidPackFile(format!(
1050 "AI object type `{other}` cannot appear in a pack file"
1051 )));
1052 }
1053 }
1054
1055 let raw_size = Pack::skip_compressed_data(pack, size)?;
1056 offset += raw_size;
1057 }
1058
1059 let mut trailer = vec![0; get_hash_kind().size()];
1060 pack.read_exact(&mut trailer)
1061 .map_err(|e| GitError::InvalidPackFile(format!("Read error: {e}")))?;
1062 if !utils::is_eof(pack) {
1063 return Err(GitError::InvalidPackFile(
1064 "The pack file is not at the end".to_string(),
1065 ));
1066 }
1067
1068 Ok(DecodeScan {
1069 retention,
1070 object_hashes: None,
1071 pack_hash: None,
1072 pack_hash_check: None,
1073 })
1074 }
1075
1076 fn scan_decode_retention_and_copy(
1077 pack: &mut (impl BufRead + Send),
1078 writer: &mut (impl Write + Send),
1079 ) -> Result<DecodeScan, GitError> {
1080 let mut tee = TeeReader {
1081 reader: pack,
1082 writer,
1083 write_error: None,
1084 payload_hash: HashAlgorithm::new(),
1085 hash_tail: Vec::with_capacity(get_hash_kind().size()),
1086 hash_size: get_hash_kind().size(),
1087 };
1088 let mut scan = Pack::scan_decode_retention(&mut tee)?;
1089 tee.check_write_error()
1090 .map_err(|e| GitError::InvalidPackFile(format!("Write temp pack file error: {e}")))?;
1091 scan.pack_hash_check = Some(tee.finish_pack_hash_check()?);
1092 Ok(scan)
1093 }
1094
1095 pub fn decode_pack_object(
1107 pack: &mut (impl BufRead + Send),
1108 offset: &mut usize,
1109 ) -> Result<Option<CacheObject>, GitError> {
1110 Self::decode_pack_object_with_crc(pack, offset, true, false, false, None, None)
1111 }
1112
1113 fn decode_pack_object_with_crc(
1114 pack: &mut (impl BufRead + Send),
1115 offset: &mut usize,
1116 track_crc: bool,
1117 skip_unneeded_objects: bool,
1118 emit_skipped_base_callback: bool,
1119 known_hash: Option<ObjectHash>,
1120 retention: Option<&DecodeRetention>,
1121 ) -> Result<Option<CacheObject>, GitError> {
1122 let init_offset = *offset;
1123 let mut reader = CrcCountingReader {
1124 inner: pack,
1125 bytes_read: 0,
1126 crc: track_crc.then(crc32fast::Hasher::new),
1127 };
1128
1129 let (type_bits, size) = match utils::read_type_and_varint_size(&mut reader, offset) {
1132 Ok(result) => result,
1133 Err(e) => {
1134 return Err(GitError::InvalidPackFile(format!("Read error: {e}")));
1137 }
1138 };
1139
1140 let t = ObjectType::from_pack_type_u8(type_bits)?;
1142
1143 match t {
1144 ObjectType::Commit | ObjectType::Tree | ObjectType::Blob | ObjectType::Tag => {
1145 if Self::should_skip_no_callback_object(
1146 track_crc,
1147 skip_unneeded_objects,
1148 known_hash,
1149 retention,
1150 init_offset,
1151 ) {
1152 let raw_size = Pack::skip_compressed_data(&mut reader, size)?;
1153 *offset += raw_size;
1154 if emit_skipped_base_callback {
1155 let hash = known_hash.ok_or_else(|| {
1156 GitError::InvalidPackFile(
1157 "Missing object hash for skipped callback entry".to_string(),
1158 )
1159 })?;
1160 return Ok(Some(CacheObject {
1161 info: CacheObjectInfo::BaseObject(t, hash),
1162 offset: init_offset,
1163 crc32: 0,
1164 data_decompressed: Vec::new(),
1165 mem_recorder: None,
1166 is_delta_in_pack: false,
1167 known_hash: None,
1168 }));
1169 }
1170 return Ok(None);
1171 }
1172
1173 let (data, raw_size) = Pack::decompress_data(&mut reader, size)?;
1174 *offset += raw_size;
1175 let crc32 = reader.crc32();
1176 let hash = known_hash.unwrap_or_else(|| utils::calculate_object_hash(t, &data));
1177 Ok(Some(CacheObject {
1178 info: CacheObjectInfo::BaseObject(t, hash),
1179 offset: init_offset,
1180 crc32,
1181 data_decompressed: data,
1182 mem_recorder: None,
1183 is_delta_in_pack: false,
1184 known_hash: None,
1185 }))
1186 }
1187 ObjectType::OffsetDelta | ObjectType::OffsetZstdelta => {
1188 let (delta_offset, bytes) =
1189 utils::read_offset_encoding(&mut reader).map_err(|e| {
1190 GitError::InvalidPackFile(format!("Read offset-delta base error: {e}"))
1191 })?;
1192 *offset += bytes;
1193
1194 let delta_offset = usize::try_from(delta_offset).map_err(|_| {
1195 GitError::InvalidObjectInfo("Invalid OffsetDelta offset".to_string())
1196 })?;
1197 let base_offset = init_offset.checked_sub(delta_offset).ok_or_else(|| {
1198 GitError::InvalidObjectInfo("Invalid OffsetDelta offset".to_string())
1199 })?;
1200
1201 if emit_skipped_base_callback
1202 && Self::should_skip_no_callback_object(
1203 false,
1204 skip_unneeded_objects,
1205 known_hash,
1206 retention,
1207 init_offset,
1208 )
1209 {
1210 let raw_size = Pack::skip_compressed_data(&mut reader, size)?;
1211 *offset += raw_size;
1212
1213 let obj_info = match t {
1214 ObjectType::OffsetDelta => CacheObjectInfo::OffsetDelta(base_offset, 0),
1215 ObjectType::OffsetZstdelta => {
1216 CacheObjectInfo::OffsetZstdelta(base_offset, 0)
1217 }
1218 _ => unreachable!(),
1219 };
1220 return Ok(Some(CacheObject {
1221 info: obj_info,
1222 offset: init_offset,
1223 crc32: 0,
1224 data_decompressed: Vec::new(),
1225 mem_recorder: None,
1226 is_delta_in_pack: true,
1227 known_hash,
1228 }));
1229 }
1230
1231 if !emit_skipped_base_callback
1232 && Self::should_skip_no_callback_object(
1233 track_crc,
1234 skip_unneeded_objects,
1235 known_hash,
1236 retention,
1237 init_offset,
1238 )
1239 {
1240 let raw_size = Pack::skip_compressed_data(&mut reader, size)?;
1241 *offset += raw_size;
1242
1243 let obj_info = match t {
1244 ObjectType::OffsetDelta => CacheObjectInfo::OffsetDelta(base_offset, 0),
1245 ObjectType::OffsetZstdelta => {
1246 CacheObjectInfo::OffsetZstdelta(base_offset, 0)
1247 }
1248 _ => unreachable!(),
1249 };
1250 return Ok(Some(CacheObject {
1251 info: obj_info,
1252 offset: init_offset,
1253 crc32: 0,
1254 data_decompressed: Vec::new(),
1255 mem_recorder: None,
1256 is_delta_in_pack: true,
1257 known_hash,
1258 }));
1259 }
1260
1261 let (data, raw_size) = Pack::decompress_data(&mut reader, size)?;
1262 *offset += raw_size;
1263
1264 let mut delta_reader = Cursor::new(&data);
1265 let (_, final_size) = utils::read_delta_object_size(&mut delta_reader)?;
1266
1267 let obj_info = match t {
1268 ObjectType::OffsetDelta => {
1269 CacheObjectInfo::OffsetDelta(base_offset, final_size)
1270 }
1271 ObjectType::OffsetZstdelta => {
1272 CacheObjectInfo::OffsetZstdelta(base_offset, final_size)
1273 }
1274 _ => unreachable!(),
1275 };
1276 let crc32 = reader.crc32();
1277 Ok(Some(CacheObject {
1278 info: obj_info,
1279 offset: init_offset,
1280 crc32,
1281 data_decompressed: data,
1282 mem_recorder: None,
1283 is_delta_in_pack: true,
1284 known_hash,
1285 }))
1286 }
1287 ObjectType::HashDelta => {
1288 let ref_sha = ObjectHash::from_stream(&mut reader).map_err(|e| {
1290 GitError::InvalidPackFile(format!("Read hash-delta base hash error: {e}"))
1291 })?;
1292 *offset += get_hash_kind().size();
1294
1295 if emit_skipped_base_callback
1296 && Self::should_skip_no_callback_object(
1297 false,
1298 skip_unneeded_objects,
1299 known_hash,
1300 retention,
1301 init_offset,
1302 )
1303 {
1304 let raw_size = Pack::skip_compressed_data(&mut reader, size)?;
1305 *offset += raw_size;
1306
1307 return Ok(Some(CacheObject {
1308 info: CacheObjectInfo::HashDelta(ref_sha, 0),
1309 offset: init_offset,
1310 crc32: 0,
1311 data_decompressed: Vec::new(),
1312 mem_recorder: None,
1313 is_delta_in_pack: true,
1314 known_hash,
1315 }));
1316 }
1317
1318 if !emit_skipped_base_callback
1319 && Self::should_skip_no_callback_object(
1320 track_crc,
1321 skip_unneeded_objects,
1322 known_hash,
1323 retention,
1324 init_offset,
1325 )
1326 {
1327 let raw_size = Pack::skip_compressed_data(&mut reader, size)?;
1328 *offset += raw_size;
1329
1330 return Ok(Some(CacheObject {
1331 info: CacheObjectInfo::HashDelta(ref_sha, 0),
1332 offset: init_offset,
1333 crc32: 0,
1334 data_decompressed: Vec::new(),
1335 mem_recorder: None,
1336 is_delta_in_pack: true,
1337 known_hash,
1338 }));
1339 }
1340
1341 let (data, raw_size) = Pack::decompress_data(&mut reader, size)?;
1342 *offset += raw_size;
1343
1344 let mut delta_reader = Cursor::new(&data);
1345 let (_, final_size) = utils::read_delta_object_size(&mut delta_reader)?;
1346
1347 let crc32 = reader.crc32();
1348
1349 Ok(Some(CacheObject {
1350 info: CacheObjectInfo::HashDelta(ref_sha, final_size),
1351 offset: init_offset,
1352 crc32,
1353 data_decompressed: data,
1354 mem_recorder: None,
1355 is_delta_in_pack: true,
1356 known_hash,
1357 }))
1358 }
1359 other => Err(GitError::InvalidPackFile(format!(
1363 "AI object type `{other}` cannot appear in a pack file"
1364 ))),
1365 }
1366 }
1367
1368 fn should_skip_no_callback_object(
1369 track_crc: bool,
1370 skip_unneeded_objects: bool,
1371 known_hash: Option<ObjectHash>,
1372 retention: Option<&DecodeRetention>,
1373 offset: usize,
1374 ) -> bool {
1375 if track_crc || !skip_unneeded_objects {
1376 return false;
1377 }
1378
1379 match (known_hash, retention) {
1380 (Some(hash), Some(retention)) => !retention.should_retain(offset, hash),
1381 _ => false,
1382 }
1383 }
1384
1385 pub fn decode<F, C>(
1392 &mut self,
1393 pack: &mut (impl BufRead + Send),
1394 callback: F,
1395 pack_id_callback: Option<C>,
1396 ) -> Result<(), GitError>
1397 where
1398 F: Fn(MetaAttached<Entry, EntryMeta>) + Sync + Send + 'static,
1399 C: FnOnce(ObjectHash) + Send + 'static,
1400 {
1401 let callback: DecodeCallback = Arc::new(callback);
1402 if self
1403 .mem_limit
1404 .is_some_and(|limit| limit >= UNBOUNDED_CACHE_THRESHOLD_BYTES)
1405 {
1406 #[cfg(unix)]
1407 if let Some(pack_path) = Self::single_open_pack_path_at_start()
1408 && Self::reader_matches_pack_prefix(pack, &pack_path)
1409 {
1410 let pack_len = std::fs::metadata(&pack_path)
1411 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?
1412 .len();
1413 self.decode_file_inner_with_sync_base_callbacks(
1414 &pack_path,
1415 Some(callback),
1416 pack_id_callback,
1417 if Self::low_memory_callback_entries() {
1418 FileDecodeMode::SkipUnneeded
1419 } else {
1420 FileDecodeMode::RetainAll
1421 },
1422 true,
1423 )?;
1424 Self::consume_reader_exact(pack, pack_len)?;
1425 return Ok(());
1426 }
1427
1428 let mut temp_pack = NamedTempFile::new().map_err(|e| {
1429 GitError::InvalidPackFile(format!("Create temp pack file error: {e}"))
1430 })?;
1431 let scan = {
1432 let mut temp_writer =
1433 io::BufWriter::with_capacity(FILE_DECODE_BUFFER_SIZE, &mut temp_pack);
1434 let scan = Pack::scan_decode_retention_and_copy(pack, &mut temp_writer)?;
1435 temp_writer.flush().map_err(|e| {
1436 GitError::InvalidPackFile(format!("Flush temp pack file error: {e}"))
1437 })?;
1438 scan
1439 };
1440 temp_pack.flush().map_err(|e| {
1441 GitError::InvalidPackFile(format!("Flush temp pack file error: {e}"))
1442 })?;
1443 return self.decode_file_inner_with_scan(
1444 temp_pack.path(),
1445 Some(callback),
1446 pack_id_callback,
1447 if Self::low_memory_callback_entries() {
1448 FileDecodeMode::SkipUnneeded
1449 } else {
1450 FileDecodeMode::RetainAll
1451 },
1452 scan,
1453 true,
1454 );
1455 }
1456 self.decode_inner(
1457 pack,
1458 Some(callback),
1459 pack_id_callback,
1460 DecodeOptions::streaming(),
1461 )
1462 }
1463
1464 #[cfg(unix)]
1465 fn single_open_pack_path_at_start() -> Option<PathBuf> {
1466 fn fd_position_is_start(fd_name: &std::ffi::OsStr) -> bool {
1467 let fdinfo_path = Path::new("/proc/self/fdinfo").join(fd_name);
1468 let Ok(fdinfo) = std::fs::read_to_string(fdinfo_path) else {
1469 return false;
1470 };
1471 fdinfo.lines().any(|line| {
1472 let Some(pos) = line.strip_prefix("pos:") else {
1473 return false;
1474 };
1475 pos.trim() == "0"
1476 })
1477 }
1478
1479 let mut found = None;
1480 let fd_dir = std::fs::read_dir("/proc/self/fd").ok()?;
1481 for entry in fd_dir.flatten() {
1482 if !fd_position_is_start(&entry.file_name()) {
1483 continue;
1484 }
1485 let Ok(path) = std::fs::read_link(entry.path()) else {
1486 continue;
1487 };
1488 if path.extension().and_then(|ext| ext.to_str()) != Some("pack") {
1489 continue;
1490 }
1491 if !path.is_file() || !path.with_extension("idx").is_file() {
1492 continue;
1493 }
1494 if found.replace(path).is_some() {
1495 return None;
1496 }
1497 }
1498 found
1499 }
1500
1501 #[cfg(unix)]
1502 fn reader_matches_pack_prefix(pack: &mut (impl BufRead + Send), pack_path: &Path) -> bool {
1503 let Ok(prefix) = pack.fill_buf() else {
1504 return false;
1505 };
1506 if prefix.is_empty() {
1507 return false;
1508 }
1509
1510 let Ok(file) = File::open(pack_path) else {
1511 return false;
1512 };
1513 let Ok(metadata) = file.metadata() else {
1514 return false;
1515 };
1516 let min_prefix = 4096.min(metadata.len() as usize);
1517 if prefix.len() < min_prefix {
1518 return false;
1519 }
1520
1521 let compare_len = prefix.len().min(FILE_DECODE_BUFFER_SIZE);
1522 let mut file_prefix = vec![0; compare_len];
1523 match file.read_at(&mut file_prefix, 0) {
1524 Ok(n) if n == compare_len => file_prefix == prefix[..compare_len],
1525 _ => false,
1526 }
1527 }
1528
1529 #[cfg(unix)]
1530 fn consume_reader_exact(
1531 pack: &mut (impl BufRead + Send),
1532 mut bytes: u64,
1533 ) -> Result<(), GitError> {
1534 while bytes != 0 {
1535 let buf = pack
1536 .fill_buf()
1537 .map_err(|e| GitError::InvalidPackFile(format!("Read pack file error: {e}")))?;
1538 if buf.is_empty() {
1539 return Err(GitError::InvalidPackFile(
1540 "Pack reader ended before matched pack bytes were consumed".to_string(),
1541 ));
1542 }
1543 let consumed =
1544 usize::try_from(bytes).map_or(buf.len(), |remaining| remaining.min(buf.len()));
1545 pack.consume(consumed);
1546 bytes -= consumed as u64;
1547 }
1548 Ok(())
1549 }
1550
1551 pub fn decode_without_callback<C>(
1557 &mut self,
1558 pack: &mut (impl BufRead + Send),
1559 pack_id_callback: Option<C>,
1560 ) -> Result<(), GitError>
1561 where
1562 C: FnOnce(ObjectHash) + Send + 'static,
1563 {
1564 self.decode_inner(pack, None, pack_id_callback, DecodeOptions::streaming())
1565 }
1566
1567 pub fn decode_file<F, C>(
1573 &mut self,
1574 pack_path: impl AsRef<Path>,
1575 callback: F,
1576 pack_id_callback: Option<C>,
1577 ) -> Result<(), GitError>
1578 where
1579 F: Fn(MetaAttached<Entry, EntryMeta>) + Sync + Send + 'static,
1580 C: FnOnce(ObjectHash) + Send + 'static,
1581 {
1582 let callback: DecodeCallback = Arc::new(callback);
1583 self.decode_file_inner(
1584 pack_path.as_ref(),
1585 Some(callback),
1586 pack_id_callback,
1587 FileDecodeMode::RetainAll,
1588 )
1589 }
1590
1591 pub fn decode_file_full_without_callback<C>(
1598 &mut self,
1599 pack_path: impl AsRef<Path>,
1600 pack_id_callback: Option<C>,
1601 ) -> Result<(), GitError>
1602 where
1603 C: FnOnce(ObjectHash) + Send + 'static,
1604 {
1605 self.decode_file_inner(
1606 pack_path.as_ref(),
1607 None,
1608 pack_id_callback,
1609 FileDecodeMode::RetainAll,
1610 )
1611 }
1612
1613 pub fn decode_file_without_callback<C>(
1619 &mut self,
1620 pack_path: impl AsRef<Path>,
1621 pack_id_callback: Option<C>,
1622 ) -> Result<(), GitError>
1623 where
1624 C: FnOnce(ObjectHash) + Send + 'static,
1625 {
1626 self.decode_file_inner(
1627 pack_path.as_ref(),
1628 None,
1629 pack_id_callback,
1630 FileDecodeMode::SkipUnneeded,
1631 )
1632 }
1633
1634 fn decode_file_inner<C>(
1635 &mut self,
1636 pack_path: &Path,
1637 callback: Option<DecodeCallback>,
1638 pack_id_callback: Option<C>,
1639 mode: FileDecodeMode,
1640 ) -> Result<(), GitError>
1641 where
1642 C: FnOnce(ObjectHash) + Send + 'static,
1643 {
1644 self.decode_file_inner_with_sync_base_callbacks(
1645 pack_path,
1646 callback,
1647 pack_id_callback,
1648 mode,
1649 true,
1650 )
1651 }
1652
1653 fn decode_file_inner_with_sync_base_callbacks<C>(
1654 &mut self,
1655 pack_path: &Path,
1656 callback: Option<DecodeCallback>,
1657 pack_id_callback: Option<C>,
1658 mode: FileDecodeMode,
1659 sync_base_callbacks: bool,
1660 ) -> Result<(), GitError>
1661 where
1662 C: FnOnce(ObjectHash) + Send + 'static,
1663 {
1664 let scan = match Pack::scan_decode_retention_from_index(pack_path) {
1665 Ok(scan) => scan,
1666 Err(_) => {
1667 let scan_file = File::open(pack_path)
1668 .map_err(|e| GitError::InvalidPackFile(format!("Open pack file error: {e}")))?;
1669 let mut scan_reader = io::BufReader::new(scan_file);
1670 Pack::scan_decode_retention(&mut scan_reader)?
1671 }
1672 };
1673 self.decode_file_inner_with_scan(
1674 pack_path,
1675 callback,
1676 pack_id_callback,
1677 mode,
1678 scan,
1679 sync_base_callbacks,
1680 )
1681 }
1682
1683 fn decode_file_inner_with_scan<C>(
1684 &mut self,
1685 pack_path: &Path,
1686 callback: Option<DecodeCallback>,
1687 pack_id_callback: Option<C>,
1688 mode: FileDecodeMode,
1689 scan: DecodeScan,
1690 sync_base_callbacks: bool,
1691 ) -> Result<(), GitError>
1692 where
1693 C: FnOnce(ObjectHash) + Send + 'static,
1694 {
1695 let DecodeScan {
1696 retention,
1697 object_hashes: known_hashes,
1698 pack_hash,
1699 pack_hash_check,
1700 } = scan;
1701 let expected_pack_hash =
1702 pack_hash.or_else(|| pack_hash_check.as_ref().map(|check| check.payload_hash));
1703 let retention = Arc::new(retention);
1704 let retention_mode = match mode {
1705 FileDecodeMode::RetainAll => DecodeRetentionMode::retain_all(retention),
1706 FileDecodeMode::SkipUnneeded => DecodeRetentionMode::skip_unneeded(retention),
1707 };
1708 let skip_payload_hash_check = callback.is_some() && Self::low_memory_callback_entries();
1709 let hash_check =
1710 if !skip_payload_hash_check && pack_hash.is_some() && pack_hash_check.is_none() {
1711 let pack_path = pack_path.to_path_buf();
1712 let kind: HashKind = get_hash_kind();
1713 Some(thread::spawn(move || {
1714 set_hash_kind(kind);
1715 Pack::hash_pack_file_payload(&pack_path)
1716 }))
1717 } else {
1718 None
1719 };
1720
1721 let file = File::open(pack_path)
1722 .map_err(|e| GitError::InvalidPackFile(format!("Open pack file error: {e}")))?;
1723 let mut reader = io::BufReader::with_capacity(FILE_DECODE_BUFFER_SIZE, file);
1724 let decode_result = self.decode_inner(
1725 &mut reader,
1726 callback,
1727 pack_id_callback,
1728 DecodeOptions {
1729 retention_mode,
1730 known_hashes,
1731 expected_pack_hash,
1732 verify_pack_stream_hash: !skip_payload_hash_check
1733 && hash_check.is_none()
1734 && pack_hash_check.is_none(),
1735 sync_base_callbacks,
1736 },
1737 );
1738
1739 let hash_result = hash_check.map(|handle| {
1740 handle
1741 .join()
1742 .map_err(|_| GitError::InvalidPackFile("Pack hash check panicked".to_string()))?
1743 });
1744
1745 decode_result?;
1746 if let Some(hash_check) = pack_hash_check {
1747 Self::verify_pack_hash_check(&hash_check, self.signature)?;
1748 } else if let Some(hash_check) = hash_result {
1749 let hash_check = hash_check?;
1750 Self::verify_pack_hash_check(&hash_check, self.signature)?;
1751 }
1752
1753 Ok(())
1754 }
1755
1756 fn verify_pack_hash_check(
1757 hash_check: &PackHashCheck,
1758 signature: ObjectHash,
1759 ) -> Result<(), GitError> {
1760 if hash_check.trailer_hash != signature {
1761 return Err(GitError::InvalidPackFile(format!(
1762 "The pack file trailer hash {} does not match decoded trailer hash {}",
1763 hash_check.trailer_hash, signature
1764 )));
1765 }
1766 if hash_check.payload_hash != signature {
1767 return Err(GitError::InvalidPackFile(format!(
1768 "The pack file hash {} does not match the trailer hash {}",
1769 hash_check.payload_hash, signature
1770 )));
1771 }
1772 Ok(())
1773 }
1774
1775 fn decode_inner<C>(
1776 &mut self,
1777 pack: &mut (impl BufRead + Send),
1778 callback: Option<DecodeCallback>,
1779 pack_id_callback: Option<C>,
1780 options: DecodeOptions,
1781 ) -> Result<(), GitError>
1782 where
1783 C: FnOnce(ObjectHash) + Send + 'static,
1784 {
1785 let DecodeOptions {
1786 retention_mode,
1787 known_hashes,
1788 expected_pack_hash,
1789 verify_pack_stream_hash,
1790 sync_base_callbacks,
1791 } = options;
1792 let time = Instant::now();
1793 let mut last_update_time = time.elapsed().as_millis();
1794 let log_enabled = tracing::enabled!(tracing::Level::INFO);
1795 let log_info = |_i: usize, pack: &Pack| {
1796 tracing::info!(
1797 "time {:.2} s \t decode: {:?} \t dec-num: {} \t cah-num: {} \t Objs: {} MB \t CacheUsed: {} MB",
1798 time.elapsed().as_millis() as f64 / 1000.0,
1799 _i,
1800 pack.pool.queued_count(),
1801 pack.caches.queued_tasks(),
1802 pack.cache_objs_mem_used() / 1024 / 1024,
1803 pack.caches.memory_used() / 1024 / 1024
1804 );
1805 };
1806 let track_crc = callback.is_some() && !Self::low_memory_callback_entries();
1807 let known_hashes = known_hashes.as_deref();
1808 let shared_params = Arc::new(SharedParams {
1809 pool: self.pool.clone(),
1810 waitlist: self.waitlist.clone(),
1811 caches: self.caches.clone(),
1812 cache_objs_mem_size: self.cache_objs_mem.clone(),
1813 callback,
1814 retention: retention_mode.retention,
1815 skip_unneeded_objects: retention_mode.skip_unneeded_objects,
1816 });
1817 let mut reader = if verify_pack_stream_hash {
1818 Wrapper::new(pack)
1819 } else {
1820 Wrapper::new_without_hash(pack)
1821 };
1822
1823 let result = Pack::check_header(&mut reader);
1824 match result {
1825 Ok((object_num, _)) => {
1826 self.number = object_num as usize;
1827 }
1828 Err(e) => {
1829 return Err(e);
1830 }
1831 }
1832 tracing::info!("The pack file has {} objects", self.number);
1833 let mut offset: usize = 12;
1834 let mut i = 0;
1835 let mem_limit = if self.caches.is_unbounded() {
1836 None
1837 } else {
1838 self.mem_limit
1839 };
1840 while i < self.number {
1841 if log_enabled && i % 1000 == 0 {
1843 let time_now = time.elapsed().as_millis();
1844 if time_now - last_update_time > 1000 {
1845 log_info(i, self);
1846 last_update_time = time_now;
1847 }
1848 }
1849 if let Some(mem_limit) = mem_limit {
1852 while self.pool.queued_count() > MAX_QUEUED_DECODE_TASKS
1853 || self.memory_used() > mem_limit
1854 {
1855 thread::yield_now();
1856 }
1857 } else {
1858 while self.pool.queued_count() > MAX_QUEUED_DECODE_TASKS {
1859 thread::yield_now();
1860 }
1861 }
1862 let known_hash = known_hashes.and_then(|hashes| hashes.get(i).copied());
1863 let r: Result<Option<CacheObject>, GitError> = Pack::decode_pack_object_with_crc(
1864 &mut reader,
1865 &mut offset,
1866 track_crc,
1867 shared_params.skip_unneeded_objects,
1868 shared_params.callback.is_some() && Self::low_memory_callback_entries(),
1869 known_hash,
1870 shared_params.retention.as_deref(),
1871 );
1872 match r {
1873 Ok(Some(obj)) => {
1874 let Some(mut obj) = Self::try_process_skipped_low_memory_callback_object(
1875 &shared_params,
1876 obj,
1877 Self::low_memory_callback_entries(),
1878 ) else {
1879 i += 1;
1880 continue;
1881 };
1882
1883 if Self::should_skip_no_callback_delta(&shared_params, &obj) {
1884 Self::process_delta_dependency(shared_params.clone(), obj);
1885 i += 1;
1886 continue;
1887 }
1888 if matches!(obj.info, CacheObjectInfo::BaseObject(_, _))
1889 && Self::should_drop_no_callback_base(&shared_params, &obj)
1890 {
1891 i += 1;
1892 continue;
1893 }
1894
1895 obj.set_mem_recorder(self.cache_objs_mem.clone());
1896 obj.record_mem_size();
1897
1898 if matches!(obj.info, CacheObjectInfo::BaseObject(_, _))
1899 && (shared_params.callback.is_none() || sync_base_callbacks)
1900 {
1901 Self::cache_obj_and_process_waitlist(&shared_params, obj);
1902 i += 1;
1903 continue;
1904 }
1905
1906 let params = shared_params.clone();
1907 let kind = get_hash_kind();
1908 self.pool.execute(move || {
1909 set_hash_kind(kind);
1910 match obj.info {
1911 CacheObjectInfo::BaseObject(_, _) => {
1912 Self::cache_obj_and_process_waitlist(¶ms, obj);
1913 }
1914 CacheObjectInfo::OffsetDelta(_, _)
1915 | CacheObjectInfo::OffsetZstdelta(_, _)
1916 | CacheObjectInfo::HashDelta(_, _) => {
1917 Self::process_delta_dependency(params, obj);
1918 }
1919 }
1920 });
1921 }
1922 Ok(None) => {}
1923 Err(e) => {
1924 self.abort_decode();
1925 return Err(e);
1926 }
1927 }
1928 i += 1;
1929 }
1930 log_info(i, self);
1931 let render_hash = verify_pack_stream_hash.then(|| reader.final_hash());
1932 self.signature = match ObjectHash::from_stream(&mut reader) {
1933 Ok(signature) => signature,
1934 Err(e) => {
1935 self.abort_decode();
1936 return Err(GitError::InvalidPackFile(format!(
1937 "Error reading pack trailer hash: {e}"
1938 )));
1939 }
1940 };
1941
1942 if let Some(expected_pack_hash) = expected_pack_hash
1943 && expected_pack_hash != self.signature
1944 {
1945 self.abort_decode();
1946 return Err(GitError::InvalidPackFile(format!(
1947 "The pack index hash {} does not match the trailer hash {}",
1948 expected_pack_hash, self.signature
1949 )));
1950 }
1951
1952 if let Some(render_hash) = render_hash
1953 && render_hash != self.signature
1954 {
1955 self.abort_decode();
1956 return Err(GitError::InvalidPackFile(format!(
1957 "The pack file hash {} does not match the trailer hash {}",
1958 render_hash, self.signature
1959 )));
1960 }
1961
1962 let end = utils::is_eof(&mut reader);
1963 if !end {
1964 self.abort_decode();
1965 return Err(GitError::InvalidPackFile(
1966 "The pack file is not at the end".to_string(),
1967 ));
1968 }
1969
1970 self.pool.join(); if let Some(pack_callback) = pack_id_callback {
1974 pack_callback(self.signature);
1975 }
1976 assert_eq!(self.waitlist.map_offset.len(), 0);
1979 assert_eq!(self.waitlist.map_ref.len(), 0);
1980 assert!(self.number >= self.caches.total_inserted());
1982 tracing::info!(
1983 "The pack file has been decoded successfully, takes: [ {:?} ]",
1984 time.elapsed()
1985 );
1986 self.caches.clear(); assert_eq!(self.cache_objs_mem_used(), 0); Ok(())
1995 }
1996
1997 pub fn decode_async(
2000 mut self,
2001 mut pack: impl BufRead + Send + 'static,
2002 sender: UnboundedSender<Entry>,
2003 ) -> JoinHandle<Pack> {
2004 let kind = get_hash_kind();
2005 thread::spawn(move || {
2006 set_hash_kind(kind);
2007 self.decode(
2008 &mut pack,
2009 move |entry| {
2010 if let Err(e) = sender.send(entry.inner) {
2011 eprintln!("Channel full, failed to send entry: {e:?}");
2012 }
2013 },
2014 None::<fn(ObjectHash)>,
2015 )
2016 .unwrap();
2017 self
2018 })
2019 }
2020
2021 pub async fn decode_stream(
2023 mut self,
2024 mut stream: impl Stream<Item = Result<Bytes, Error>> + Unpin + Send + 'static,
2025 sender: UnboundedSender<MetaAttached<Entry, EntryMeta>>,
2026 pack_hash_send: Option<UnboundedSender<ObjectHash>>,
2027 ) -> Self {
2028 let kind = get_hash_kind();
2029 let (tx, rx) = std::sync::mpsc::channel();
2030 let mut reader = StreamBufReader::new(rx);
2031 tokio::spawn(async move {
2032 while let Some(chunk) = stream.next().await {
2033 let data = chunk.unwrap().to_vec();
2034 if let Err(e) = tx.send(data) {
2035 eprintln!("Sending Error: {e:?}");
2036 break;
2037 }
2038 }
2039 });
2040 tokio::task::spawn_blocking(move || {
2043 set_hash_kind(kind);
2044 self.decode(
2045 &mut reader,
2046 move |entry: MetaAttached<Entry, EntryMeta>| {
2047 if let Err(e) = sender.send(entry) {
2049 eprintln!("unbound channel Sending Error: {e:?}");
2050 }
2051 },
2052 Some(move |pack_id: ObjectHash| {
2053 if let Some(pack_id_send) = pack_hash_send
2054 && let Err(e) = pack_id_send.send(pack_id)
2055 {
2056 eprintln!("unbound channel Sending Error: {e:?}");
2057 }
2058 }),
2059 )
2060 .unwrap();
2061 self
2062 })
2063 .await
2064 .unwrap()
2065 }
2066
2067 fn memory_used(&self) -> usize {
2069 self.cache_objs_mem_used() + self.caches.memory_used_index()
2070 }
2071
2072 fn cache_objs_mem_used(&self) -> usize {
2074 self.cache_objs_mem.load(Ordering::Acquire)
2075 }
2076
2077 fn release_offset_dependency(
2078 shared_params: &SharedParams,
2079 base_offset: usize,
2080 base_obj: &CacheObject,
2081 ) {
2082 if let Some(retention) = &shared_params.retention {
2083 retention.consume_offset_dependency(base_offset);
2084 Self::maybe_remove_released_base(shared_params, base_obj);
2085 }
2086 }
2087
2088 fn release_hash_dependency(
2089 shared_params: &SharedParams,
2090 base_hash: ObjectHash,
2091 base_obj: &CacheObject,
2092 ) {
2093 if let Some(retention) = &shared_params.retention {
2094 retention.consume_hash_dependency(base_hash);
2095 Self::maybe_remove_released_base(shared_params, base_obj);
2096 }
2097 }
2098
2099 fn maybe_remove_released_base(shared_params: &SharedParams, base_obj: &CacheObject) {
2100 if let Some(retention) = &shared_params.retention
2101 && let Some(hash) = base_obj.base_object_hash()
2102 && !retention.should_retain(base_obj.offset, hash)
2103 {
2104 shared_params.caches.remove_unbounded(base_obj.offset, hash);
2105 }
2106 }
2107
2108 fn process_waitlist_objects(
2109 shared_params: &Arc<SharedParams>,
2110 wait_objs: Vec<CacheObject>,
2111 base_obj: Arc<CacheObject>,
2112 ) {
2113 for obj in wait_objs {
2114 Self::process_delta(Arc::clone(shared_params), obj, base_obj.clone());
2116 }
2117 }
2118
2119 fn try_process_skipped_low_memory_callback_object(
2120 shared_params: &Arc<SharedParams>,
2121 obj: CacheObject,
2122 low_memory_callback_entries: bool,
2123 ) -> Option<CacheObject> {
2124 if !low_memory_callback_entries || !obj.data_decompressed.is_empty() {
2125 return Some(obj);
2126 }
2127
2128 let (Some(callback), Some(retention)) = (
2129 shared_params.callback.as_ref(),
2130 shared_params.retention.as_ref(),
2131 ) else {
2132 return Some(obj);
2133 };
2134
2135 match &obj.info {
2136 CacheObjectInfo::BaseObject(_, hash) => {
2137 if retention.should_retain(obj.offset, *hash)
2138 || shared_params.waitlist.has_waiters(obj.offset, *hash)
2139 {
2140 return Some(obj);
2141 }
2142 callback(Self::callback_entry_owned(obj));
2143 None
2144 }
2145 CacheObjectInfo::OffsetDelta(_, _)
2146 | CacheObjectInfo::OffsetZstdelta(_, _)
2147 | CacheObjectInfo::HashDelta(_, _) => {
2148 let Some(hash) = obj.known_hash else {
2149 return Some(obj);
2150 };
2151 if retention.should_retain(obj.offset, hash)
2152 || shared_params.waitlist.has_waiters(obj.offset, hash)
2153 {
2154 return Some(obj);
2155 }
2156 Self::process_delta_dependency(shared_params.clone(), obj);
2157 None
2158 }
2159 }
2160 }
2161
2162 fn should_skip_no_callback_delta(
2163 shared_params: &SharedParams,
2164 delta_obj: &CacheObject,
2165 ) -> bool {
2166 if shared_params.callback.is_some() {
2167 return false;
2168 }
2169
2170 if !shared_params.skip_unneeded_objects {
2171 return false;
2172 }
2173
2174 let Some(retention) = &shared_params.retention else {
2175 return false;
2176 };
2177 let Some(hash) = delta_obj.known_hash else {
2178 return false;
2179 };
2180
2181 !retention.should_retain(delta_obj.offset, hash)
2182 && !shared_params
2183 .waitlist
2184 .map_offset
2185 .contains_key(&delta_obj.offset)
2186 && !shared_params.waitlist.map_ref.contains_key(&hash)
2187 }
2188
2189 fn should_drop_no_callback_base(shared_params: &SharedParams, base_obj: &CacheObject) -> bool {
2190 if shared_params.callback.is_some() {
2191 return false;
2192 }
2193
2194 let Some(retention) = &shared_params.retention else {
2195 return false;
2196 };
2197
2198 let Some(hash) = base_obj.base_object_hash() else {
2199 return false;
2200 };
2201
2202 !retention.should_retain(base_obj.offset, hash)
2203 && !shared_params.waitlist.has_waiters(base_obj.offset, hash)
2204 }
2205
2206 fn process_delta_dependency(shared_params: Arc<SharedParams>, obj: CacheObject) {
2207 match obj.info {
2208 CacheObjectInfo::OffsetDelta(base_offset, _)
2209 | CacheObjectInfo::OffsetZstdelta(base_offset, _) => {
2210 if let Some(base_obj) = shared_params.caches.get_by_offset(base_offset) {
2211 Self::release_offset_dependency(&shared_params, base_offset, &base_obj);
2212 Self::process_delta(shared_params, obj, base_obj);
2213 } else {
2214 shared_params.waitlist.insert_offset(base_offset, obj);
2215 if let Some(retention) = &shared_params.retention {
2216 retention.consume_offset_dependency(base_offset);
2217 }
2218 if let Some(base_obj) = shared_params.caches.get_by_offset(base_offset) {
2219 Self::maybe_remove_released_base(&shared_params, &base_obj);
2220 Self::process_waitlist(&shared_params, base_obj);
2221 }
2222 }
2223 }
2224 CacheObjectInfo::HashDelta(base_ref, _) => {
2225 if let Some(base_obj) = shared_params.caches.get_by_hash(base_ref) {
2226 Self::release_hash_dependency(&shared_params, base_ref, &base_obj);
2227 Self::process_delta(shared_params, obj, base_obj);
2228 } else {
2229 shared_params.waitlist.insert_ref(base_ref, obj);
2230 if let Some(retention) = &shared_params.retention {
2231 retention.consume_hash_dependency(base_ref);
2232 }
2233 if let Some(base_obj) = shared_params.caches.get_by_hash(base_ref) {
2234 Self::maybe_remove_released_base(&shared_params, &base_obj);
2235 Self::process_waitlist(&shared_params, base_obj);
2236 }
2237 }
2238 }
2239 CacheObjectInfo::BaseObject(_, _) => unreachable!(),
2240 }
2241 }
2242
2243 fn process_delta(
2246 shared_params: Arc<SharedParams>,
2247 delta_obj: CacheObject,
2248 base_obj: Arc<CacheObject>,
2249 ) {
2250 if Self::should_skip_no_callback_delta(&shared_params, &delta_obj) {
2251 return;
2252 }
2253
2254 if Self::try_callback_unneeded_low_memory_delta(&shared_params, &delta_obj, &base_obj) {
2255 return;
2256 }
2257
2258 shared_params.pool.clone().execute(move || {
2259 let known_hash = delta_obj.known_hash;
2260 let mut new_obj = match delta_obj.info {
2261 CacheObjectInfo::OffsetDelta(_, _) | CacheObjectInfo::HashDelta(_, _) => {
2262 Pack::rebuild_delta_with_hash(delta_obj, base_obj, known_hash)
2263 }
2264 CacheObjectInfo::OffsetZstdelta(_, _) => {
2265 Pack::rebuild_zstdelta_with_hash(delta_obj, base_obj, known_hash)
2266 }
2267 _ => unreachable!(),
2268 };
2269
2270 new_obj.set_mem_recorder(shared_params.cache_objs_mem_size.clone());
2271 new_obj.record_mem_size();
2272 Self::cache_obj_and_process_waitlist(&shared_params, new_obj); });
2274 }
2275
2276 fn try_callback_unneeded_low_memory_delta(
2277 shared_params: &SharedParams,
2278 delta_obj: &CacheObject,
2279 base_obj: &CacheObject,
2280 ) -> bool {
2281 if !Self::low_memory_callback_entries() {
2282 return false;
2283 }
2284
2285 let (Some(callback), Some(retention), Some(hash)) = (
2286 shared_params.callback.as_ref(),
2287 shared_params.retention.as_ref(),
2288 delta_obj.known_hash,
2289 ) else {
2290 return false;
2291 };
2292
2293 if retention.should_retain(delta_obj.offset, hash)
2294 || shared_params.waitlist.has_waiters(delta_obj.offset, hash)
2295 {
2296 return false;
2297 }
2298
2299 callback(Self::low_memory_delta_callback_entry(
2300 delta_obj,
2301 base_obj.object_type(),
2302 hash,
2303 ));
2304 true
2305 }
2306
2307 fn low_memory_delta_callback_entry(
2308 delta_obj: &CacheObject,
2309 obj_type: ObjectType,
2310 hash: ObjectHash,
2311 ) -> MetaAttached<Entry, EntryMeta> {
2312 MetaAttached {
2313 inner: Entry {
2314 obj_type,
2315 data: Vec::new(),
2316 hash,
2317 chain_len: 0,
2318 },
2319 meta: EntryMeta {
2320 pack_offset: Some(delta_obj.offset),
2321 crc32: Some(delta_obj.crc32),
2322 is_delta: Some(delta_obj.is_delta_in_pack),
2323 ..Default::default()
2324 },
2325 }
2326 }
2327
2328 fn cache_obj_and_process_waitlist(shared_params: &Arc<SharedParams>, new_obj: CacheObject) {
2330 if let Some(retention) = &shared_params.retention {
2331 let hash = new_obj.base_object_hash().unwrap();
2332 let offset = new_obj.offset;
2333 let should_retain = retention.should_retain(offset, hash);
2334 if should_retain {
2335 if let Some(callback) = &shared_params.callback {
2336 callback(Self::callback_entry_ref(&new_obj));
2337 }
2338 let new_obj = shared_params.caches.insert(offset, hash, new_obj);
2339 let wait_objs = shared_params.waitlist.take(offset, hash);
2340 Self::process_waitlist_objects(shared_params, wait_objs, new_obj);
2341 } else {
2342 let wait_objs = shared_params.waitlist.take(offset, hash);
2343 if !wait_objs.is_empty() {
2344 if let Some(callback) = &shared_params.callback {
2345 callback(Self::callback_entry_ref(&new_obj));
2346 }
2347 Self::process_waitlist_objects(shared_params, wait_objs, Arc::new(new_obj));
2348 } else if let Some(callback) = &shared_params.callback {
2349 callback(Self::callback_entry_owned(new_obj));
2350 }
2351 }
2352 return;
2353 }
2354 if let Some(callback) = &shared_params.callback {
2355 callback(Self::callback_entry_ref(&new_obj));
2356 }
2357 let new_obj = shared_params.caches.insert(
2358 new_obj.offset,
2359 new_obj.base_object_hash().unwrap(),
2360 new_obj,
2361 );
2362 Self::process_waitlist(shared_params, new_obj);
2363 }
2364
2365 fn process_waitlist(shared_params: &Arc<SharedParams>, base_obj: Arc<CacheObject>) {
2366 let wait_objs = shared_params
2367 .waitlist
2368 .take(base_obj.offset, base_obj.base_object_hash().unwrap());
2369 Self::process_waitlist_objects(shared_params, wait_objs, base_obj);
2370 }
2371
2372 pub fn rebuild_delta(delta_obj: CacheObject, base_obj: Arc<CacheObject>) -> CacheObject {
2375 Self::rebuild_delta_with_hash(delta_obj, base_obj, None)
2376 }
2377
2378 fn rebuild_delta_with_hash(
2379 delta_obj: CacheObject,
2380 base_obj: Arc<CacheObject>,
2381 known_hash: Option<ObjectHash>,
2382 ) -> CacheObject {
2383 const COPY_INSTRUCTION_FLAG: u8 = 1 << 7;
2384 const COPY_OFFSET_BYTES: u8 = 4;
2385 const COPY_SIZE_BYTES: u8 = 3;
2386 const COPY_ZERO_SIZE: usize = 0x10000;
2387
2388 let mut stream = Cursor::new(delta_obj.data_decompressed.as_slice());
2389
2390 let (base_size, result_size) = utils::read_delta_object_size(&mut stream).unwrap();
2393
2394 let base_info = &base_obj.data_decompressed;
2396 assert_eq!(base_info.len(), base_size, "Base object size mismatch");
2397
2398 let mut result = Vec::with_capacity(result_size);
2399
2400 loop {
2401 let instruction = match utils::read_bytes(&mut stream) {
2403 Ok([instruction]) => instruction,
2404 Err(err) if err.kind() == ErrorKind::UnexpectedEof => break,
2405 Err(err) => {
2406 panic!(
2407 "{}",
2408 GitError::DeltaObjectError(format!("Wrong instruction in delta :{err}"))
2409 );
2410 }
2411 };
2412
2413 if instruction & COPY_INSTRUCTION_FLAG == 0 {
2414 if instruction == 0 {
2416 panic!(
2418 "{}",
2419 GitError::DeltaObjectError(String::from("Invalid data instruction"))
2420 );
2421 }
2422
2423 let start = stream.position() as usize;
2424 let end = start + instruction as usize;
2425 let delta_data = *stream.get_ref();
2426 let data = delta_data.get(start..end).unwrap_or_else(|| {
2427 panic!(
2428 "{}",
2429 GitError::DeltaObjectError("Invalid data instruction".to_string())
2430 )
2431 });
2432 result.extend_from_slice(data);
2433 stream.set_position(end as u64);
2434 } else {
2435 let mut nonzero_bytes = instruction;
2440 let offset =
2441 utils::read_partial_int(&mut stream, COPY_OFFSET_BYTES, &mut nonzero_bytes)
2442 .unwrap();
2443 let mut size =
2444 utils::read_partial_int(&mut stream, COPY_SIZE_BYTES, &mut nonzero_bytes)
2445 .unwrap();
2446 if size == 0 {
2447 size = COPY_ZERO_SIZE;
2449 }
2450 let base_data = base_info.get(offset..(offset + size)).ok_or_else(|| {
2452 GitError::DeltaObjectError("Invalid copy instruction".to_string())
2453 });
2454
2455 match base_data {
2456 Ok(data) => result.extend_from_slice(data),
2457 Err(e) => panic!("{}", e),
2458 }
2459 }
2460 }
2461 assert_eq!(result_size, result.len(), "Result size mismatch");
2462
2463 let hash = known_hash
2464 .unwrap_or_else(|| utils::calculate_object_hash(base_obj.object_type(), &result));
2465 CacheObject {
2467 info: CacheObjectInfo::BaseObject(base_obj.object_type(), hash),
2468 offset: delta_obj.offset,
2469 crc32: delta_obj.crc32,
2470 data_decompressed: result,
2471 mem_recorder: None,
2472 is_delta_in_pack: delta_obj.is_delta_in_pack,
2473 known_hash: None,
2474 } }
2477 pub fn rebuild_zstdelta(delta_obj: CacheObject, base_obj: Arc<CacheObject>) -> CacheObject {
2478 Self::rebuild_zstdelta_with_hash(delta_obj, base_obj, None)
2479 }
2480
2481 fn rebuild_zstdelta_with_hash(
2482 delta_obj: CacheObject,
2483 base_obj: Arc<CacheObject>,
2484 known_hash: Option<ObjectHash>,
2485 ) -> CacheObject {
2486 let result = zstdelta::apply(&base_obj.data_decompressed, &delta_obj.data_decompressed)
2487 .expect("Failed to apply zstdelta");
2488 let hash = known_hash
2489 .unwrap_or_else(|| utils::calculate_object_hash(base_obj.object_type(), &result));
2490 CacheObject {
2491 info: CacheObjectInfo::BaseObject(base_obj.object_type(), hash),
2492 offset: delta_obj.offset,
2493 crc32: delta_obj.crc32,
2494 data_decompressed: result,
2495 mem_recorder: None,
2496 is_delta_in_pack: delta_obj.is_delta_in_pack,
2497 known_hash: None,
2498 } }
2501}
2502
2503impl Pack {
2504 pub fn stats_pack(path: PathBuf) -> Result<PackStats, crate::errors::GitError> {
2526 PackStats::analyze(path)
2527 }
2528}
2529
2530#[cfg(test)]
2531mod tests {
2532 use std::{
2533 fs,
2534 io::{BufReader, Cursor, prelude::*},
2535 path::{Path, PathBuf},
2536 sync::{
2537 Arc, Mutex,
2538 atomic::{AtomicUsize, Ordering},
2539 },
2540 };
2541
2542 use flate2::{Compression, write::ZlibEncoder};
2543 use futures_util::TryStreamExt;
2544 use sha1::{Digest, Sha1};
2545 use tempfile::tempdir;
2546 use threadpool::ThreadPool;
2547 use tokio_util::io::ReaderStream;
2548
2549 use crate::{
2550 hash::{HashKind, ObjectHash, get_hash_kind, set_hash_kind_for_test},
2551 internal::{
2552 object::types::ObjectType,
2553 pack::{
2554 Pack,
2555 cache::{_Cache, Caches},
2556 cache_object::{CacheObject, CacheObjectInfo},
2557 test_pack_download::download_pack_file,
2558 tests::init_logger,
2559 utils,
2560 waitlist::Waitlist,
2561 },
2562 },
2563 };
2564
2565 fn pack_test_tmp() -> (tempfile::TempDir, PathBuf) {
2566 let dir = tempdir().unwrap();
2567 let path = dir.path().join(".cache_temp");
2568 (dir, path)
2569 }
2570
2571 fn pack_header_object_count(path: impl AsRef<Path>) -> usize {
2572 let mut file = fs::File::open(path).unwrap();
2573 let mut header = [0; 12];
2574 file.read_exact(&mut header).unwrap();
2575 assert_eq!(&header[0..4], b"PACK");
2576
2577 let mut count_bytes = [0; 4];
2578 count_bytes.copy_from_slice(&header[8..12]);
2579 u32::from_be_bytes(count_bytes) as usize
2580 }
2581
2582 const LARGE_PACK_TEST_MEM_LIMIT: usize = super::UNBOUNDED_CACHE_THRESHOLD_BYTES;
2583
2584 #[cfg_attr(coverage, ignore)]
2585 #[ignore = "requires large remote pack fixture"]
2586 #[tokio::test]
2587 async fn test_pack_check_header() {
2588 let (source, _guard) = download_pack_file("medium-sha1.pack");
2589 let expected_object_num = pack_header_object_count(&source);
2590
2591 let f = fs::File::open(source).unwrap();
2592 let mut buf_reader = BufReader::new(f);
2593 let (object_num, _) = Pack::check_header(&mut buf_reader).unwrap();
2594
2595 assert_eq!(object_num as usize, expected_object_num);
2596 }
2597
2598 #[test]
2599 fn test_decompress_data() {
2600 let data = b"Hello, world!"; let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2602 encoder.write_all(data).unwrap();
2603 let compressed_data = encoder.finish().unwrap();
2604 let compressed_size = compressed_data.len();
2605
2606 let mut cursor: Cursor<Vec<u8>> = Cursor::new(compressed_data);
2608 let expected_size = data.len();
2609
2610 let result = Pack::decompress_data(&mut cursor, expected_size);
2612 match result {
2613 Ok((decompressed_data, bytes_read)) => {
2614 assert_eq!(bytes_read, compressed_size);
2615 assert_eq!(decompressed_data, data);
2616 }
2617 Err(e) => panic!("Decompression failed: {e:?}"),
2618 }
2619 }
2620
2621 #[test]
2622 fn test_pack_decode_truncated_pack_returns_err_without_panic() {
2623 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2624 let (source, _dl_guard) = download_pack_file("small-sha1.pack");
2625 let mut bytes = fs::read(source).unwrap();
2626 bytes.truncate(bytes.len() - 1);
2627
2628 let tmp_dir = tempfile::tempdir().unwrap();
2629 let tmp_path = tmp_dir.path().to_path_buf();
2630 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
2631 let mut buffered = BufReader::new(Cursor::new(bytes));
2632 let mut pack = Pack::new(Some(2), Some(1024 * 1024), Some(tmp_path), true);
2633 pack.decode(&mut buffered, |_| {}, None::<fn(ObjectHash)>)
2634 }));
2635
2636 assert!(result.is_ok(), "truncated pack decode should not panic");
2637 assert!(
2638 matches!(
2639 result.unwrap(),
2640 Err(crate::errors::GitError::InvalidPackFile(_))
2641 | Err(crate::errors::GitError::IOError(_))
2642 ),
2643 "truncated pack decode should return a pack error"
2644 );
2645 }
2646
2647 #[test]
2648 fn test_skip_compressed_data_exact_size_and_size_errors() {
2649 let data = b"Hello, world!";
2650 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2651 encoder.write_all(data).unwrap();
2652 let compressed = encoder.finish().unwrap();
2653
2654 let mut exact = Cursor::new(compressed.clone());
2655 let bytes_read = Pack::skip_compressed_data(&mut exact, data.len()).unwrap();
2656 assert_eq!(bytes_read, compressed.len());
2657
2658 let mut too_large = Cursor::new(compressed.clone());
2659 let err = Pack::skip_compressed_data(&mut too_large, data.len() + 1).unwrap_err();
2660 assert!(err.to_string().contains("smaller than the expected size"));
2661
2662 let mut too_small = Cursor::new(compressed);
2663 let err = Pack::skip_compressed_data(&mut too_small, data.len() - 1).unwrap_err();
2664 assert!(err.to_string().contains("exceeds the expected size"));
2665 }
2666
2667 #[test]
2668 fn test_read_be_u64() {
2669 let mut reader = Cursor::new(0x0102_0304_0506_0708u64.to_be_bytes());
2670 assert_eq!(
2671 Pack::read_be_u64(&mut reader).unwrap(),
2672 0x0102_0304_0506_0708
2673 );
2674 }
2675
2676 #[test]
2677 fn test_decode_pack_object_crc_can_be_skipped() {
2678 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2679 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2680 encoder.write_all(b"a").unwrap();
2681 let compressed = encoder.finish().unwrap();
2682
2683 let mut object_data = Vec::new();
2684 object_data.push(0x31);
2685 object_data.extend_from_slice(&compressed);
2686 let expected_crc = crc32fast::hash(&object_data);
2687
2688 let mut crc_reader = Cursor::new(object_data.clone());
2689 let mut offset = 0;
2690 let with_crc = Pack::decode_pack_object(&mut crc_reader, &mut offset)
2691 .unwrap()
2692 .unwrap();
2693 assert_eq!(with_crc.crc32, expected_crc);
2694 assert_eq!(with_crc.data_decompressed, b"a");
2695
2696 let mut no_crc_reader = Cursor::new(object_data);
2697 let mut offset = 0;
2698 let supplied_hash = ObjectHash::new(b"known-hash-from-idx");
2699 let without_crc = Pack::decode_pack_object_with_crc(
2700 &mut no_crc_reader,
2701 &mut offset,
2702 false,
2703 false,
2704 false,
2705 Some(supplied_hash),
2706 None,
2707 )
2708 .unwrap()
2709 .unwrap();
2710 assert_eq!(without_crc.crc32, 0);
2711 assert_eq!(without_crc.data_decompressed, b"a");
2712 assert_eq!(without_crc.base_object_hash(), Some(supplied_hash));
2713 }
2714
2715 #[test]
2716 fn test_decode_pack_object_can_emit_skipped_base_callback_entry() {
2717 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2718 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2719 encoder.write_all(b"a").unwrap();
2720 let compressed = encoder.finish().unwrap();
2721
2722 let mut object_data = Vec::new();
2723 object_data.push(0x31);
2724 object_data.extend_from_slice(&compressed);
2725
2726 let supplied_hash = ObjectHash::new(b"known-hash-from-idx");
2727 let retention = super::DecodeRetention::default();
2728 let mut reader = Cursor::new(object_data);
2729 let mut offset = 0;
2730 let skipped = Pack::decode_pack_object_with_crc(
2731 &mut reader,
2732 &mut offset,
2733 false,
2734 true,
2735 true,
2736 Some(supplied_hash),
2737 Some(&retention),
2738 )
2739 .unwrap()
2740 .unwrap();
2741
2742 assert_eq!(skipped.crc32, 0);
2743 assert!(skipped.data_decompressed.is_empty());
2744 assert_eq!(skipped.base_object_hash(), Some(supplied_hash));
2745 assert_eq!(offset, reader.get_ref().len());
2746 }
2747
2748 #[test]
2749 fn test_decode_pack_object_can_skip_unneeded_delta_payload() {
2750 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2751 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2752 encoder.write_all(b"abc").unwrap();
2753 let compressed = encoder.finish().unwrap();
2754
2755 let mut object_data = Vec::new();
2756 object_data.push(0x63);
2757 object_data.push(5);
2758 object_data.extend_from_slice(&compressed);
2759
2760 let supplied_hash = ObjectHash::new(b"known-leaf-delta-hash");
2761 let retention = super::DecodeRetention::default();
2762 let mut reader = Cursor::new(object_data);
2763 let init_offset = 20;
2764 let mut offset = init_offset;
2765 let skipped = Pack::decode_pack_object_with_crc(
2766 &mut reader,
2767 &mut offset,
2768 false,
2769 true,
2770 true,
2771 Some(supplied_hash),
2772 Some(&retention),
2773 )
2774 .unwrap()
2775 .unwrap();
2776
2777 assert_eq!(skipped.crc32, 0);
2778 assert!(skipped.data_decompressed.is_empty());
2779 assert_eq!(skipped.known_hash, Some(supplied_hash));
2780 assert_eq!(
2781 skipped.info,
2782 CacheObjectInfo::OffsetDelta(init_offset - 5, 0)
2783 );
2784 assert_eq!(offset, init_offset + reader.get_ref().len());
2785 }
2786
2787 #[test]
2788 fn test_low_memory_delta_callback_entry_uses_known_hash_without_payload() {
2789 let hash = ObjectHash::new(b"known-leaf-delta-hash");
2790 let delta_obj = CacheObject {
2791 info: CacheObjectInfo::OffsetDelta(12, 5),
2792 offset: 40,
2793 crc32: 1234,
2794 data_decompressed: b"delta instructions".to_vec(),
2795 mem_recorder: None,
2796 is_delta_in_pack: true,
2797 known_hash: Some(hash),
2798 };
2799
2800 let entry = Pack::low_memory_delta_callback_entry(&delta_obj, ObjectType::Blob, hash);
2801
2802 assert_eq!(entry.inner.obj_type, ObjectType::Blob);
2803 assert!(entry.inner.data.is_empty());
2804 assert_eq!(entry.inner.hash, hash);
2805 assert_eq!(entry.meta.pack_offset, Some(40));
2806 assert_eq!(entry.meta.crc32, Some(1234));
2807 assert_eq!(entry.meta.is_delta, Some(true));
2808 }
2809
2810 #[test]
2811 fn test_skipped_low_memory_base_callbacks_without_cache_insert() {
2812 let hash = ObjectHash::new(b"known-leaf-base-hash");
2813 let seen = Arc::new(Mutex::new(Vec::new()));
2814 let callback_seen = Arc::clone(&seen);
2815 let callback: super::DecodeCallback = Arc::new(move |entry| {
2816 callback_seen.lock().unwrap().push((
2817 entry.inner.hash,
2818 entry.inner.data.len(),
2819 entry.meta.pack_offset,
2820 ));
2821 });
2822 let (_dir, cache_path) = pack_test_tmp();
2823 let shared_params = Arc::new(super::SharedParams {
2824 pool: Arc::new(ThreadPool::new(1)),
2825 waitlist: Arc::new(Waitlist::new()),
2826 caches: Arc::new(Caches::new(None, cache_path, 1)),
2827 cache_objs_mem_size: Arc::new(AtomicUsize::new(0)),
2828 callback: Some(callback),
2829 retention: Some(Arc::new(super::DecodeRetention::default())),
2830 skip_unneeded_objects: true,
2831 });
2832 let obj = CacheObject {
2833 info: CacheObjectInfo::BaseObject(ObjectType::Blob, hash),
2834 offset: 64,
2835 crc32: 0,
2836 data_decompressed: Vec::new(),
2837 mem_recorder: None,
2838 is_delta_in_pack: false,
2839 known_hash: None,
2840 };
2841
2842 let remaining =
2843 Pack::try_process_skipped_low_memory_callback_object(&shared_params, obj, true);
2844
2845 assert!(remaining.is_none());
2846 assert_eq!(shared_params.caches.total_inserted(), 0);
2847 assert_eq!(seen.lock().unwrap().as_slice(), &[(hash, 0, Some(64))]);
2848 }
2849
2850 #[cfg(unix)]
2851 #[test]
2852 fn test_large_mem_decode_ignores_unrelated_open_pack_fd() {
2853 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2854 let dir = tempdir().unwrap();
2855 let (open_pack_data, open_hash) = single_blob_pack(b"open pack data");
2856 let open_pack_path = dir.path().join("open.pack");
2857 fs::write(&open_pack_path, open_pack_data).unwrap();
2858 write_test_idx(&open_pack_path, vec![(open_hash, 12)]);
2859 let _open_pack = fs::File::open(&open_pack_path).unwrap();
2860
2861 let (reader_pack_data, _) = single_blob_pack(b"reader data");
2862 let mut reader = Cursor::new(reader_pack_data);
2863 let decoded = Arc::new(Mutex::new(Vec::new()));
2864 let decoded_for_callback = Arc::clone(&decoded);
2865 let mut pack = Pack::new(
2866 Some(1),
2867 Some(super::UNBOUNDED_CACHE_THRESHOLD_BYTES),
2868 Some(dir.path().join("tmp")),
2869 true,
2870 );
2871
2872 pack.decode(
2873 &mut reader,
2874 move |entry| decoded_for_callback.lock().unwrap().push(entry.inner.data),
2875 None::<fn(ObjectHash)>,
2876 )
2877 .unwrap();
2878
2879 assert_eq!(*decoded.lock().unwrap(), vec![b"reader data".to_vec()]);
2880 }
2881
2882 #[cfg(unix)]
2883 #[test]
2884 fn test_consume_reader_exact_leaves_following_bytes() {
2885 let mut reader = Cursor::new(b"pack-datafollowing-data".to_vec());
2886 Pack::consume_reader_exact(&mut reader, b"pack-data".len() as u64).unwrap();
2887
2888 let mut rest = Vec::new();
2889 reader.read_to_end(&mut rest).unwrap();
2890 assert_eq!(rest, b"following-data");
2891 }
2892
2893 #[test]
2894 fn test_pack_decode_without_callback_empty_pack() {
2895 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2896 let mut pack_data = Vec::new();
2897 pack_data.extend_from_slice(b"PACK");
2898 pack_data.extend_from_slice(&2u32.to_be_bytes());
2899 pack_data.extend_from_slice(&0u32.to_be_bytes());
2900 let trailer = Sha1::digest(&pack_data);
2901 pack_data.extend_from_slice(&trailer);
2902
2903 let mut reader = Cursor::new(pack_data);
2904 let mut pack = Pack::new(Some(1), None, None, true);
2905 pack.decode_without_callback(&mut reader, None::<fn(ObjectHash)>)
2906 .unwrap();
2907
2908 assert_eq!(pack.number, 0);
2909 }
2910
2911 #[test]
2912 fn test_pack_decode_without_callback_single_blob_pack() {
2913 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2914 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2915 encoder.write_all(b"a").unwrap();
2916 let compressed = encoder.finish().unwrap();
2917
2918 let mut pack_data = Vec::new();
2919 pack_data.extend_from_slice(b"PACK");
2920 pack_data.extend_from_slice(&2u32.to_be_bytes());
2921 pack_data.extend_from_slice(&1u32.to_be_bytes());
2922 pack_data.push(0x31);
2923 pack_data.extend_from_slice(&compressed);
2924 let trailer = Sha1::digest(&pack_data);
2925 pack_data.extend_from_slice(&trailer);
2926
2927 let mut reader = Cursor::new(pack_data);
2928 let mut pack = Pack::new(Some(1), None, None, true);
2929 pack.decode_without_callback(&mut reader, None::<fn(ObjectHash)>)
2930 .unwrap();
2931
2932 assert_eq!(pack.number, 1);
2933 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
2934 }
2935
2936 #[test]
2937 fn test_pack_decode_large_mem_limit_uses_temp_retention_path() {
2938 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2939 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2940 encoder.write_all(b"a").unwrap();
2941 let compressed = encoder.finish().unwrap();
2942
2943 let mut pack_data = Vec::new();
2944 pack_data.extend_from_slice(b"PACK");
2945 pack_data.extend_from_slice(&2u32.to_be_bytes());
2946 pack_data.extend_from_slice(&1u32.to_be_bytes());
2947 pack_data.push(0x31);
2948 pack_data.extend_from_slice(&compressed);
2949 let trailer = Sha1::digest(&pack_data);
2950 pack_data.extend_from_slice(&trailer);
2951
2952 let seen = Arc::new(Mutex::new(Vec::new()));
2953 let seen_cb = Arc::clone(&seen);
2954 let mut reader = Cursor::new(pack_data);
2955 let mut pack = Pack::new(
2956 Some(1),
2957 Some(super::UNBOUNDED_CACHE_THRESHOLD_BYTES),
2958 None,
2959 true,
2960 );
2961 pack.decode(
2962 &mut reader,
2963 move |entry| seen_cb.lock().unwrap().push(entry.inner.data),
2964 None::<fn(ObjectHash)>,
2965 )
2966 .unwrap();
2967
2968 assert_eq!(pack.number, 1);
2969 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
2970 assert_eq!(*seen.lock().unwrap(), vec![b"a".to_vec()]);
2971 }
2972
2973 #[test]
2974 fn test_pack_decode_file_without_callback_uses_idx() {
2975 let _guard = set_hash_kind_for_test(HashKind::Sha1);
2976 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
2977 encoder.write_all(b"a").unwrap();
2978 let compressed = encoder.finish().unwrap();
2979
2980 let mut pack_data = Vec::new();
2981 pack_data.extend_from_slice(b"PACK");
2982 pack_data.extend_from_slice(&2u32.to_be_bytes());
2983 pack_data.extend_from_slice(&1u32.to_be_bytes());
2984 pack_data.push(0x31);
2985 pack_data.extend_from_slice(&compressed);
2986 let trailer = Sha1::digest(&pack_data);
2987 pack_data.extend_from_slice(&trailer);
2988
2989 let dir = tempdir().unwrap();
2990 let pack_path = dir.path().join("single.pack");
2991 fs::write(&pack_path, &pack_data).unwrap();
2992
2993 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
2994 let mut idx_data = Vec::new();
2995 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
2996 idx_data.extend_from_slice(&2u32.to_be_bytes());
2997 let first = obj_hash.as_ref()[0] as usize;
2998 for fanout_idx in 0..256 {
2999 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3000 idx_data.extend_from_slice(&count.to_be_bytes());
3001 }
3002 idx_data.extend_from_slice(obj_hash.as_ref());
3003 idx_data.extend_from_slice(&0u32.to_be_bytes());
3004 idx_data.extend_from_slice(&12u32.to_be_bytes());
3005 append_test_idx_trailer(&mut idx_data, &trailer);
3006 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3007
3008 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3009 pack.decode_file_without_callback(&pack_path, None::<fn(ObjectHash)>)
3010 .unwrap();
3011
3012 assert_eq!(pack.number, 1);
3013 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
3014 }
3015
3016 #[test]
3017 fn test_pack_decode_file_callback_uses_idx_and_crc() {
3018 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3019 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3020 encoder.write_all(b"a").unwrap();
3021 let compressed = encoder.finish().unwrap();
3022
3023 let mut object_data = Vec::new();
3024 object_data.push(0x31);
3025 object_data.extend_from_slice(&compressed);
3026 let expected_crc = crc32fast::hash(&object_data);
3027
3028 let mut pack_data = Vec::new();
3029 pack_data.extend_from_slice(b"PACK");
3030 pack_data.extend_from_slice(&2u32.to_be_bytes());
3031 pack_data.extend_from_slice(&1u32.to_be_bytes());
3032 pack_data.extend_from_slice(&object_data);
3033 let trailer = Sha1::digest(&pack_data);
3034 pack_data.extend_from_slice(&trailer);
3035
3036 let dir = tempdir().unwrap();
3037 let pack_path = dir.path().join("single-callback.pack");
3038 fs::write(&pack_path, &pack_data).unwrap();
3039
3040 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
3041 let mut idx_data = Vec::new();
3042 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3043 idx_data.extend_from_slice(&2u32.to_be_bytes());
3044 let first = obj_hash.as_ref()[0] as usize;
3045 for fanout_idx in 0..256 {
3046 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3047 idx_data.extend_from_slice(&count.to_be_bytes());
3048 }
3049 idx_data.extend_from_slice(obj_hash.as_ref());
3050 idx_data.extend_from_slice(&expected_crc.to_be_bytes());
3051 idx_data.extend_from_slice(&12u32.to_be_bytes());
3052 append_test_idx_trailer(&mut idx_data, &trailer);
3053 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3054
3055 let entries = Arc::new(std::sync::Mutex::new(Vec::new()));
3056 let entries_for_cb = Arc::clone(&entries);
3057 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3058 pack.decode_file(
3059 &pack_path,
3060 move |entry| entries_for_cb.lock().unwrap().push(entry),
3061 None::<fn(ObjectHash)>,
3062 )
3063 .unwrap();
3064
3065 let entries = Arc::try_unwrap(entries).unwrap().into_inner().unwrap();
3066 assert_eq!(entries.len(), 1);
3067 assert_eq!(entries[0].inner.hash, obj_hash);
3068 assert_eq!(entries[0].inner.data, b"a");
3069 assert_eq!(entries[0].meta.pack_offset, Some(12));
3070 assert_eq!(entries[0].meta.crc32, Some(expected_crc));
3071 assert_eq!(pack.number, 1);
3072 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
3073 }
3074
3075 #[test]
3076 fn test_pack_decode_file_callback_ignores_idx_with_bad_checksum() {
3077 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3078 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3079 encoder.write_all(b"a").unwrap();
3080 let compressed = encoder.finish().unwrap();
3081
3082 let mut pack_data = Vec::new();
3083 pack_data.extend_from_slice(b"PACK");
3084 pack_data.extend_from_slice(&2u32.to_be_bytes());
3085 pack_data.extend_from_slice(&1u32.to_be_bytes());
3086 pack_data.push(0x31);
3087 pack_data.extend_from_slice(&compressed);
3088 let trailer = Sha1::digest(&pack_data);
3089 pack_data.extend_from_slice(&trailer);
3090
3091 let dir = tempdir().unwrap();
3092 let pack_path = dir.path().join("bad-idx-checksum.pack");
3093 fs::write(&pack_path, &pack_data).unwrap();
3094
3095 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
3096 let mut idx_data = Vec::new();
3097 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3098 idx_data.extend_from_slice(&2u32.to_be_bytes());
3099 let first = obj_hash.as_ref()[0] as usize;
3100 for fanout_idx in 0..256 {
3101 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3102 idx_data.extend_from_slice(&count.to_be_bytes());
3103 }
3104 let hash_offset = idx_data.len();
3105 idx_data.extend_from_slice(obj_hash.as_ref());
3106 idx_data.extend_from_slice(&0u32.to_be_bytes());
3107 idx_data.extend_from_slice(&12u32.to_be_bytes());
3108 append_test_idx_trailer(&mut idx_data, &trailer);
3109 idx_data[hash_offset] ^= 0xff;
3110 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3111
3112 let entries = Arc::new(std::sync::Mutex::new(Vec::new()));
3113 let entries_for_cb = Arc::clone(&entries);
3114 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3115 pack.decode_file(
3116 &pack_path,
3117 move |entry| entries_for_cb.lock().unwrap().push(entry.inner.hash),
3118 None::<fn(ObjectHash)>,
3119 )
3120 .unwrap();
3121
3122 assert_eq!(entries.lock().unwrap().as_slice(), &[obj_hash]);
3123 }
3124
3125 #[test]
3126 fn test_pack_decode_file_full_without_callback_uses_idx() {
3127 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3128 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3129 encoder.write_all(b"a").unwrap();
3130 let compressed = encoder.finish().unwrap();
3131
3132 let mut pack_data = Vec::new();
3133 pack_data.extend_from_slice(b"PACK");
3134 pack_data.extend_from_slice(&2u32.to_be_bytes());
3135 pack_data.extend_from_slice(&1u32.to_be_bytes());
3136 pack_data.push(0x31);
3137 pack_data.extend_from_slice(&compressed);
3138 let trailer = Sha1::digest(&pack_data);
3139 pack_data.extend_from_slice(&trailer);
3140
3141 let dir = tempdir().unwrap();
3142 let pack_path = dir.path().join("single-full-no-callback.pack");
3143 fs::write(&pack_path, &pack_data).unwrap();
3144
3145 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
3146 let mut idx_data = Vec::new();
3147 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3148 idx_data.extend_from_slice(&2u32.to_be_bytes());
3149 let first = obj_hash.as_ref()[0] as usize;
3150 for fanout_idx in 0..256 {
3151 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3152 idx_data.extend_from_slice(&count.to_be_bytes());
3153 }
3154 idx_data.extend_from_slice(obj_hash.as_ref());
3155 idx_data.extend_from_slice(&0u32.to_be_bytes());
3156 idx_data.extend_from_slice(&12u32.to_be_bytes());
3157 append_test_idx_trailer(&mut idx_data, &trailer);
3158 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3159
3160 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3161 pack.decode_file_full_without_callback(&pack_path, None::<fn(ObjectHash)>)
3162 .unwrap();
3163
3164 assert_eq!(pack.number, 1);
3165 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
3166 }
3167
3168 #[test]
3169 fn test_pack_decode_file_full_without_callback_uses_idx_large_offset() {
3170 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3171 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3172 encoder.write_all(b"a").unwrap();
3173 let compressed = encoder.finish().unwrap();
3174
3175 let mut pack_data = Vec::new();
3176 pack_data.extend_from_slice(b"PACK");
3177 pack_data.extend_from_slice(&2u32.to_be_bytes());
3178 pack_data.extend_from_slice(&1u32.to_be_bytes());
3179 pack_data.push(0x31);
3180 pack_data.extend_from_slice(&compressed);
3181 let trailer = Sha1::digest(&pack_data);
3182 pack_data.extend_from_slice(&trailer);
3183
3184 let dir = tempdir().unwrap();
3185 let pack_path = dir.path().join("single-full-no-callback-large-offset.pack");
3186 fs::write(&pack_path, &pack_data).unwrap();
3187
3188 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
3189 let mut idx_data = Vec::new();
3190 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3191 idx_data.extend_from_slice(&2u32.to_be_bytes());
3192 let first = obj_hash.as_ref()[0] as usize;
3193 for fanout_idx in 0..256 {
3194 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3195 idx_data.extend_from_slice(&count.to_be_bytes());
3196 }
3197 idx_data.extend_from_slice(obj_hash.as_ref());
3198 idx_data.extend_from_slice(&0u32.to_be_bytes());
3199 idx_data.extend_from_slice(&0x8000_0000u32.to_be_bytes());
3200 idx_data.extend_from_slice(&12u64.to_be_bytes());
3201 append_test_idx_trailer(&mut idx_data, &trailer);
3202 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3203
3204 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3205 pack.decode_file_full_without_callback(&pack_path, None::<fn(ObjectHash)>)
3206 .unwrap();
3207
3208 assert_eq!(pack.number, 1);
3209 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
3210 }
3211
3212 #[test]
3213 fn test_pack_decode_file_without_callback_falls_back_without_idx() {
3214 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3215 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3216 encoder.write_all(b"a").unwrap();
3217 let compressed = encoder.finish().unwrap();
3218
3219 let mut pack_data = Vec::new();
3220 pack_data.extend_from_slice(b"PACK");
3221 pack_data.extend_from_slice(&2u32.to_be_bytes());
3222 pack_data.extend_from_slice(&1u32.to_be_bytes());
3223 pack_data.push(0x31);
3224 pack_data.extend_from_slice(&compressed);
3225 let trailer = Sha1::digest(&pack_data);
3226 pack_data.extend_from_slice(&trailer);
3227
3228 let dir = tempdir().unwrap();
3229 let pack_path = dir.path().join("single-no-idx.pack");
3230 fs::write(&pack_path, &pack_data).unwrap();
3231
3232 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3233 pack.decode_file_without_callback(&pack_path, None::<fn(ObjectHash)>)
3234 .unwrap();
3235
3236 assert_eq!(pack.number, 1);
3237 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
3238 }
3239
3240 #[test]
3241 fn test_pack_decode_file_without_callback_rejects_idx_pack_hash_mismatch() {
3242 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3243 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3244 encoder.write_all(b"a").unwrap();
3245 let compressed = encoder.finish().unwrap();
3246
3247 let mut pack_data = Vec::new();
3248 pack_data.extend_from_slice(b"PACK");
3249 pack_data.extend_from_slice(&2u32.to_be_bytes());
3250 pack_data.extend_from_slice(&1u32.to_be_bytes());
3251 pack_data.push(0x31);
3252 pack_data.extend_from_slice(&compressed);
3253 let trailer = Sha1::digest(&pack_data);
3254 pack_data.extend_from_slice(&trailer);
3255
3256 let dir = tempdir().unwrap();
3257 let pack_path = dir.path().join("bad-pack-hash.pack");
3258 fs::write(&pack_path, &pack_data).unwrap();
3259
3260 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
3261 let mut idx_data = Vec::new();
3262 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3263 idx_data.extend_from_slice(&2u32.to_be_bytes());
3264 let first = obj_hash.as_ref()[0] as usize;
3265 for fanout_idx in 0..256 {
3266 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3267 idx_data.extend_from_slice(&count.to_be_bytes());
3268 }
3269 idx_data.extend_from_slice(obj_hash.as_ref());
3270 idx_data.extend_from_slice(&0u32.to_be_bytes());
3271 idx_data.extend_from_slice(&12u32.to_be_bytes());
3272 append_test_idx_trailer(&mut idx_data, &[0xff; 20]);
3273 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3274
3275 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3276 let err = pack
3277 .decode_file_without_callback(&pack_path, None::<fn(ObjectHash)>)
3278 .unwrap_err();
3279
3280 assert!(err.to_string().contains("does not match the trailer hash"));
3281 }
3282
3283 #[test]
3284 fn test_pack_decode_file_full_without_callback_rejects_stale_idx_after_pack_change() {
3285 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3286
3287 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3288 encoder.write_all(b"a").unwrap();
3289 let original_compressed = encoder.finish().unwrap();
3290
3291 let mut original_pack_data = Vec::new();
3292 original_pack_data.extend_from_slice(b"PACK");
3293 original_pack_data.extend_from_slice(&2u32.to_be_bytes());
3294 original_pack_data.extend_from_slice(&1u32.to_be_bytes());
3295 original_pack_data.push(0x31);
3296 original_pack_data.extend_from_slice(&original_compressed);
3297 let original_trailer = Sha1::digest(&original_pack_data);
3298
3299 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3300 encoder.write_all(b"b").unwrap();
3301 let changed_compressed = encoder.finish().unwrap();
3302 assert_eq!(original_compressed.len(), changed_compressed.len());
3303
3304 let mut changed_pack_data = Vec::new();
3305 changed_pack_data.extend_from_slice(b"PACK");
3306 changed_pack_data.extend_from_slice(&2u32.to_be_bytes());
3307 changed_pack_data.extend_from_slice(&1u32.to_be_bytes());
3308 changed_pack_data.push(0x31);
3309 changed_pack_data.extend_from_slice(&changed_compressed);
3310 changed_pack_data.extend_from_slice(&original_trailer);
3311
3312 let dir = tempdir().unwrap();
3313 let pack_path = dir.path().join("stale-idx.pack");
3314 fs::write(&pack_path, &changed_pack_data).unwrap();
3315
3316 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, b"a");
3317 let mut idx_data = Vec::new();
3318 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3319 idx_data.extend_from_slice(&2u32.to_be_bytes());
3320 let first = obj_hash.as_ref()[0] as usize;
3321 for fanout_idx in 0..256 {
3322 let count = if fanout_idx >= first { 1u32 } else { 0u32 };
3323 idx_data.extend_from_slice(&count.to_be_bytes());
3324 }
3325 idx_data.extend_from_slice(obj_hash.as_ref());
3326 idx_data.extend_from_slice(&0u32.to_be_bytes());
3327 idx_data.extend_from_slice(&12u32.to_be_bytes());
3328 append_test_idx_trailer(&mut idx_data, &original_trailer);
3329 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3330
3331 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3332 let err = pack
3333 .decode_file_full_without_callback(&pack_path, None::<fn(ObjectHash)>)
3334 .unwrap_err();
3335
3336 assert!(err.to_string().contains("does not match the trailer hash"));
3337 }
3338
3339 fn append_compressed(buf: &mut Vec<u8>, data: &[u8]) {
3340 let mut encoder = ZlibEncoder::new(Vec::new(), Compression::default());
3341 encoder.write_all(data).unwrap();
3342 buf.extend_from_slice(&encoder.finish().unwrap());
3343 }
3344
3345 fn single_blob_pack(data: &[u8]) -> (Vec<u8>, ObjectHash) {
3346 let obj_hash = utils::calculate_object_hash(ObjectType::Blob, data);
3347 let mut pack_data = Vec::new();
3348 pack_data.extend_from_slice(b"PACK");
3349 pack_data.extend_from_slice(&2u32.to_be_bytes());
3350 pack_data.extend_from_slice(&1u32.to_be_bytes());
3351 pack_data.push(0x30 | data.len() as u8);
3352 append_compressed(&mut pack_data, data);
3353 let trailer = Sha1::digest(&pack_data);
3354 pack_data.extend_from_slice(&trailer);
3355 (pack_data, obj_hash)
3356 }
3357
3358 fn append_test_idx_trailer(idx_data: &mut Vec<u8>, pack_hash: &[u8]) {
3359 idx_data.extend_from_slice(pack_hash);
3360 let mut idx_hash = crate::utils::HashAlgorithm::new();
3361 idx_hash.update(idx_data);
3362 idx_data.extend_from_slice(&idx_hash.finalize());
3363 }
3364
3365 fn write_test_idx(pack_path: &Path, mut objects: Vec<(ObjectHash, u32)>) {
3366 objects.sort_by(|a, b| a.0.as_ref().cmp(b.0.as_ref()));
3367 let mut idx_data = Vec::new();
3368 idx_data.extend_from_slice(&0xff74_4f63u32.to_be_bytes());
3369 idx_data.extend_from_slice(&2u32.to_be_bytes());
3370 for fanout_idx in 0..256 {
3371 let count = objects
3372 .iter()
3373 .filter(|(hash, _)| hash.as_ref()[0] as usize <= fanout_idx)
3374 .count() as u32;
3375 idx_data.extend_from_slice(&count.to_be_bytes());
3376 }
3377 for (hash, _) in &objects {
3378 idx_data.extend_from_slice(hash.as_ref());
3379 }
3380 for _ in &objects {
3381 idx_data.extend_from_slice(&0u32.to_be_bytes());
3382 }
3383 for (_, offset) in &objects {
3384 idx_data.extend_from_slice(&offset.to_be_bytes());
3385 }
3386 let pack_data = fs::read(pack_path).unwrap();
3387 let hash_size = get_hash_kind().size();
3388 append_test_idx_trailer(&mut idx_data, &pack_data[pack_data.len() - hash_size..]);
3389 fs::write(pack_path.with_extension("idx"), idx_data).unwrap();
3390 }
3391
3392 #[test]
3393 fn test_pack_decode_file_without_callback_releases_delta_bases() {
3394 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3395 let base_hash = utils::calculate_object_hash(ObjectType::Blob, b"hello");
3396 let ofs_hash = utils::calculate_object_hash(ObjectType::Blob, b"hi there");
3397 let ref_hash = utils::calculate_object_hash(ObjectType::Blob, b"HELLO");
3398
3399 let mut pack_data = Vec::new();
3400 pack_data.extend_from_slice(b"PACK");
3401 pack_data.extend_from_slice(&2u32.to_be_bytes());
3402 pack_data.extend_from_slice(&3u32.to_be_bytes());
3403
3404 let base_offset = pack_data.len() as u32;
3405 pack_data.push(0x35);
3406 append_compressed(&mut pack_data, b"hello");
3407
3408 let ofs_offset = pack_data.len() as u32;
3409 pack_data.push(0x6b);
3410 pack_data.push((ofs_offset - base_offset) as u8);
3411 append_compressed(
3412 &mut pack_data,
3413 [b"\x05\x08\x08".as_ref(), b"hi there"].concat().as_slice(),
3414 );
3415
3416 let ref_offset = pack_data.len() as u32;
3417 pack_data.push(0x78);
3418 pack_data.extend_from_slice(base_hash.as_ref());
3419 append_compressed(
3420 &mut pack_data,
3421 [b"\x05\x05\x05".as_ref(), b"HELLO"].concat().as_slice(),
3422 );
3423
3424 let trailer = Sha1::digest(&pack_data);
3425 pack_data.extend_from_slice(&trailer);
3426
3427 let dir = tempdir().unwrap();
3428 let pack_path = dir.path().join("delta.pack");
3429 fs::write(&pack_path, &pack_data).unwrap();
3430 write_test_idx(
3431 &pack_path,
3432 vec![
3433 (base_hash, base_offset),
3434 (ofs_hash, ofs_offset),
3435 (ref_hash, ref_offset),
3436 ],
3437 );
3438
3439 let mut pack = Pack::new(Some(1), None, Some(dir.path().join("tmp")), true);
3440 pack.decode_file_without_callback(&pack_path, None::<fn(ObjectHash)>)
3441 .unwrap();
3442
3443 assert_eq!(pack.number, 3);
3444 assert_eq!(pack.signature.to_string(), hex::encode(trailer));
3445 }
3446
3447 #[test]
3448 fn test_rebuild_delta_literal_instruction() {
3449 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3450 let base = Arc::new(CacheObject::new_for_undeltified(
3451 ObjectType::Blob,
3452 b"hello".to_vec(),
3453 12,
3454 0,
3455 ));
3456 let delta = CacheObject {
3457 info: CacheObjectInfo::OffsetDelta(12, 8),
3458 offset: 20,
3459 crc32: 0,
3460 data_decompressed: [b"\x05\x08\x08".as_ref(), b"hi there"].concat(),
3461 mem_recorder: None,
3462 is_delta_in_pack: true,
3463 known_hash: None,
3464 };
3465
3466 let rebuilt = Pack::rebuild_delta(delta, base);
3467
3468 assert_eq!(rebuilt.object_type(), ObjectType::Blob);
3469 assert_eq!(rebuilt.data_decompressed, b"hi there");
3470 }
3471
3472 #[test]
3473 #[cfg(target_pointer_width = "32")]
3474 fn test_pack_new_mem_limit_no_overflow_32bit() {
3475 let mem_limit = 1_200_000_000usize;
3479 let (_tmp_dir, tmp) = pack_test_tmp();
3480 let result = std::panic::catch_unwind(|| {
3481 let _p = Pack::new(Some(1), Some(mem_limit), Some(tmp), true);
3482 });
3483 assert!(result.is_ok(), "Pack::new should not panic on 32-bit");
3484 }
3485
3486 fn run_decode_no_delta(filename: &str, kind: HashKind) {
3488 let _guard = set_hash_kind_for_test(kind);
3489 let (source, _dl_guard) = download_pack_file(filename);
3490
3491 let (_tmp_dir, tmp) = pack_test_tmp();
3492
3493 let f = fs::File::open(source).unwrap();
3494 let mut buffered = BufReader::new(f);
3495 let mut p = Pack::new(None, Some(1024 * 1024 * 20), Some(tmp), true);
3496 p.decode(&mut buffered, |_| {}, None::<fn(ObjectHash)>)
3497 .unwrap();
3498 }
3499 #[test]
3500 fn test_pack_decode_without_delta() {
3501 run_decode_no_delta("small-sha1.pack", HashKind::Sha1);
3502 run_decode_no_delta("small-sha256.pack", HashKind::Sha256);
3503 }
3504
3505 fn run_decode_with_ref_delta(filename: &str, kind: HashKind) {
3507 let _guard = set_hash_kind_for_test(kind);
3508 init_logger();
3509
3510 let (source, _dl_guard) = download_pack_file(filename);
3511
3512 let (_tmp_dir, tmp) = pack_test_tmp();
3513
3514 let f = fs::File::open(source).unwrap();
3515 let mut buffered = BufReader::new(f);
3516 let mut p = Pack::new(None, Some(1024 * 1024 * 20), Some(tmp), true);
3517 p.decode(&mut buffered, |_| {}, None::<fn(ObjectHash)>)
3518 .unwrap();
3519 }
3520 #[test]
3521 fn test_pack_decode_with_ref_delta() {
3522 run_decode_with_ref_delta("ref-delta-sha1.pack", HashKind::Sha1);
3523 run_decode_with_ref_delta("ref-delta-sha256.pack", HashKind::Sha256);
3524 }
3525
3526 fn run_decode_no_mem_limit(filename: &str, kind: HashKind) {
3528 let _guard = set_hash_kind_for_test(kind);
3529 let (source, _dl_guard) = download_pack_file(filename);
3530
3531 let (_tmp_dir, tmp) = pack_test_tmp();
3532
3533 let f = fs::File::open(source).unwrap();
3534 let mut buffered = BufReader::new(f);
3535 let mut p = Pack::new(None, None, Some(tmp), true);
3536 p.decode(&mut buffered, |_| {}, None::<fn(ObjectHash)>)
3537 .unwrap();
3538 }
3539 #[test]
3540 fn test_pack_decode_no_mem_limit() {
3541 run_decode_no_mem_limit("small-sha1.pack", HashKind::Sha1);
3542 run_decode_no_mem_limit("small-sha256.pack", HashKind::Sha256);
3543 }
3544
3545 async fn run_decode_large_with_delta(filename: &str, kind: HashKind) {
3547 let _guard = set_hash_kind_for_test(kind);
3548 init_logger();
3549 let (source, _dl_guard) = download_pack_file(filename);
3550
3551 let (_tmp_dir, tmp) = pack_test_tmp();
3552
3553 let f = fs::File::open(source).unwrap();
3554 let mut buffered = BufReader::new(f);
3555 let mut p = Pack::new(
3556 Some(4),
3557 Some(LARGE_PACK_TEST_MEM_LIMIT),
3558 Some(tmp.clone()),
3559 true,
3560 );
3561 let rt = p.decode(
3562 &mut buffered,
3563 |_obj| {
3564 },
3566 None::<fn(ObjectHash)>,
3567 );
3568 if let Err(e) = rt {
3569 let _ = fs::remove_dir_all(&tmp);
3570 panic!("Error: {e:?}");
3571 }
3572 }
3573 #[cfg_attr(coverage, ignore)]
3574 #[ignore = "requires large remote pack fixture"]
3575 #[tokio::test]
3576 async fn test_pack_decode_with_large_file_with_delta_without_ref() {
3577 run_decode_large_with_delta("medium-sha1.pack", HashKind::Sha1).await;
3578 run_decode_large_with_delta("medium-sha256.pack", HashKind::Sha256).await;
3579 } async fn run_decode_large_stream(filename: &str, kind: HashKind) {
3583 let _guard = set_hash_kind_for_test(kind);
3584 init_logger();
3585 let (source, _dl_guard) = download_pack_file(filename);
3586 let expected_object_num = pack_header_object_count(&source);
3587
3588 let (_tmp_dir, tmp) = pack_test_tmp();
3589 let f = tokio::fs::File::open(source).await.unwrap();
3590 let stream = ReaderStream::new(f).map_err(axum::Error::new);
3591 let p = Pack::new(
3592 Some(4),
3593 Some(LARGE_PACK_TEST_MEM_LIMIT),
3594 Some(tmp.clone()),
3595 true,
3596 );
3597
3598 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3599 let handle = tokio::spawn(async move { p.decode_stream(stream, tx, None).await });
3600 let count = Arc::new(AtomicUsize::new(0));
3601 let count_c = count.clone();
3602 let consume = tokio::spawn(async move {
3604 let mut cnt = 0;
3605 while let Some(_entry) = rx.recv().await {
3606 cnt += 1;
3607 }
3608 tracing::info!("Received: {}", cnt);
3609 count_c.store(cnt, Ordering::Release);
3610 });
3611 let p = handle.await.unwrap();
3612 consume.await.unwrap();
3613 assert_eq!(count.load(Ordering::Acquire), p.number);
3614 assert_eq!(p.number, expected_object_num);
3615 }
3616 #[cfg_attr(coverage, ignore)]
3617 #[ignore = "requires large remote pack fixture"]
3618 #[tokio::test]
3619 async fn test_decode_large_file_stream() {
3620 run_decode_large_stream("medium-sha1.pack", HashKind::Sha1).await;
3621 run_decode_large_stream("medium-sha256.pack", HashKind::Sha256).await;
3622 }
3623
3624 async fn run_decode_large_file_async(filename: &str, kind: HashKind) {
3626 let _guard = set_hash_kind_for_test(kind);
3627 let (source, _dl_guard) = download_pack_file(filename);
3628
3629 let (_tmp_dir, tmp) = pack_test_tmp();
3630 let f = fs::File::open(source).unwrap();
3631 let buffered = BufReader::new(f);
3632 let p = Pack::new(
3633 Some(4),
3634 Some(LARGE_PACK_TEST_MEM_LIMIT),
3635 Some(tmp.clone()),
3636 true,
3637 );
3638
3639 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
3640 let handle = p.decode_async(buffered, tx); let mut cnt = 0;
3642 while let Some(_entry) = rx.recv().await {
3643 cnt += 1; }
3645 let p = handle.join().unwrap();
3646 assert_eq!(cnt, p.number);
3647 }
3648 #[cfg_attr(coverage, ignore)]
3649 #[ignore = "requires large remote pack fixture"]
3650 #[tokio::test]
3651 async fn test_decode_large_file_async() {
3652 run_decode_large_file_async("medium-sha1.pack", HashKind::Sha1).await;
3653 run_decode_large_file_async("medium-sha256.pack", HashKind::Sha256).await;
3654 }
3655
3656 fn run_decode_with_delta_no_ref(filename: &str, kind: HashKind) {
3658 let _guard = set_hash_kind_for_test(kind);
3659 let (source, _dl_guard) = download_pack_file(filename);
3660
3661 let (_tmp_dir, tmp) = pack_test_tmp();
3662
3663 let f = fs::File::open(source).unwrap();
3664 let mut buffered = BufReader::new(f);
3665 let mut p = Pack::new(None, Some(LARGE_PACK_TEST_MEM_LIMIT), Some(tmp), true);
3666 p.decode(&mut buffered, |_| {}, None::<fn(ObjectHash)>)
3667 .unwrap();
3668 }
3669 #[cfg_attr(coverage, ignore)]
3670 #[ignore = "requires large remote pack fixture"]
3671 #[test]
3672 fn test_pack_decode_with_delta_without_ref() {
3673 run_decode_with_delta_no_ref("medium-sha1.pack", HashKind::Sha1);
3674 run_decode_with_delta_no_ref("medium-sha256.pack", HashKind::Sha256);
3675 }
3676
3677 #[cfg_attr(coverage, ignore)]
3678 #[ignore = "requires large remote pack fixture"]
3679 #[test] fn test_pack_decode_multi_task_with_large_file_with_delta_without_ref() {
3681 let rt = tokio::runtime::Builder::new_current_thread()
3682 .enable_all()
3683 .build()
3684 .unwrap();
3685 rt.block_on(async move {
3686 for (kind, filename) in [
3688 (HashKind::Sha1, "medium-sha1.pack"),
3689 (HashKind::Sha256, "medium-sha256.pack"),
3690 ] {
3691 let f1 = run_decode_large_with_delta(filename, kind);
3692 let f2 = run_decode_large_with_delta(filename, kind);
3693 let _ = futures::future::join(f1, f2).await;
3694 }
3695 });
3696 }
3697
3698 #[test]
3710 fn test_stats_pack_small_sha1() {
3711 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3712 let (source, _dl_guard) = download_pack_file("small-sha1.pack");
3713
3714 let stats = Pack::stats_pack(source).expect("stats_pack should succeed");
3715
3716 eprintln!(
3717 "small-sha1 stats: total={}, commits={}, trees={}, blobs={}, tags={}, deltas={}",
3718 stats.total, stats.commits, stats.trees, stats.blobs, stats.tags, stats.deltas
3719 );
3720
3721 let sum = stats.commits + stats.trees + stats.blobs + stats.tags + stats.deltas;
3722 assert_eq!(
3723 sum, stats.total,
3724 "per-type counts should sum to total ({} vs {})",
3725 sum, stats.total
3726 );
3727
3728 assert!(stats.commits > 0, "expected at least one commit");
3729 assert!(stats.blobs > 0, "expected at least one blob");
3730 }
3731
3732 #[test]
3737 fn test_stats_pack_medium_sha1_has_deltas() {
3738 let _guard = set_hash_kind_for_test(HashKind::Sha1);
3739 let (source, _dl_guard) = download_pack_file("medium-sha1.pack");
3740
3741 let stats = Pack::stats_pack(source).expect("stats_pack should succeed on medium pack");
3742
3743 eprintln!(
3744 "medium-sha1 stats: total={}, commits={}, trees={}, blobs={}, tags={}, deltas={}",
3745 stats.total, stats.commits, stats.trees, stats.blobs, stats.tags, stats.deltas
3746 );
3747
3748 let sum = stats.commits + stats.trees + stats.blobs + stats.tags + stats.deltas;
3749 assert_eq!(sum, stats.total, "per-type counts must equal total");
3750
3751 assert!(
3752 stats.deltas > 0,
3753 "expected delta objects in medium-sha1 pack"
3754 );
3755
3756 assert!(stats.total > 1000, "expected a sizeable medium pack");
3757 }
3758
3759 #[test]
3763 fn test_stats_pack_file_not_found() {
3764 let result = Pack::stats_pack(PathBuf::from("/nonexistent/path/to/fake.pack"));
3765 assert!(
3766 result.is_err(),
3767 "stats_pack should return Err for a missing file"
3768 );
3769 }
3770
3771 #[test]
3776 fn test_stats_pack_invalid_pack_magic() {
3777 use std::io::Write;
3778
3779 use tempfile::NamedTempFile;
3780
3781 let mut tmp = NamedTempFile::new().expect("create temp file");
3782 tmp.write_all(b"FAKE\x00\x00\x00\x02\x00\x00\x00\x05")
3783 .expect("write temp bytes");
3784 let path = tmp.path().to_path_buf();
3785
3786 let result = Pack::stats_pack(path);
3787 assert!(
3788 result.is_err(),
3789 "stats_pack should return Err for invalid pack magic"
3790 );
3791 }
3792}