1use super::*;
12use flate2::{Decompress, FlushDecompress};
13
14#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub struct PackIndexOptions {
20 pub limits: PackReadLimits,
21 threads: usize,
22}
23
24impl PackIndexOptions {
25 pub fn new(limits: PackReadLimits) -> Self {
26 Self {
27 limits,
28 threads: std::thread::available_parallelism()
29 .map(|count| count.get())
30 .unwrap_or(1),
31 }
32 }
33
34 pub fn with_threads(mut self, threads: usize) -> Self {
35 self.threads = threads.max(1);
36 self
37 }
38
39 pub const fn threads(self) -> usize {
40 self.threads
41 }
42}
43
44impl Default for PackIndexOptions {
45 fn default() -> Self {
46 Self::new(PackReadLimits::default())
47 }
48}
49
50#[derive(Debug, Clone)]
51struct EntryDescriptor {
52 offset: usize,
53 data_offset: usize,
54 end_offset: usize,
55 header: EntryHeader,
56 base: Option<DeltaBase>,
57}
58
59#[derive(Debug)]
60struct ResolvedEntry {
61 oid: ObjectId,
62 object_type: ObjectType,
63 size: u64,
64 crc32: u32,
65 depth: usize,
66 body: Option<Vec<u8>>,
67}
68
69#[derive(Debug, Clone, Copy)]
70enum ReadyBase {
71 Internal(usize),
72 External(ObjectId),
73}
74
75#[derive(Debug, Clone, Copy)]
76struct ReadyDelta {
77 index: usize,
78 base: ReadyBase,
79 depth: usize,
80}
81
82#[derive(Debug, Clone, Copy)]
83struct ResolutionSettings {
84 format: ObjectFormat,
85 options: PackIndexOptions,
86 retain_all: bool,
87}
88
89pub(crate) fn build_parallel_index<F, P>(
90 pack: &[u8],
91 format: ObjectFormat,
92 external_base: &mut F,
93 options: PackIndexOptions,
94 cancel: CancelFlag<'_>,
95 progress: &mut P,
96) -> Result<PackIndexBuild>
97where
98 F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
99 P: FnMut(PackIndexProgress),
100{
101 let (pack_checksum, descriptors) = discover_pack_entries(pack, format, options, cancel)?;
102 let resolved = resolve_entries_parallel(
103 pack,
104 &descriptors,
105 external_base,
106 ResolutionSettings {
107 format,
108 options,
109 retain_all: false,
110 },
111 cancel,
112 progress,
113 )?;
114 finish_index(pack_checksum, &descriptors, resolved, format)
115}
116
117pub(crate) fn parse_parallel_pack<F>(
118 pack: &[u8],
119 format: ObjectFormat,
120 external_base: &mut F,
121 options: PackIndexOptions,
122 cancel: CancelFlag<'_>,
123) -> Result<PackFile>
124where
125 F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
126{
127 let (checksum, descriptors) = discover_pack_entries(pack, format, options, cancel)?;
128 let resolved = resolve_entries_parallel(
129 pack,
130 &descriptors,
131 external_base,
132 ResolutionSettings {
133 format,
134 options,
135 retain_all: true,
136 },
137 cancel,
138 &mut |_| {},
139 )?;
140 let mut entries = Vec::with_capacity(resolved.len());
141 for (descriptor, resolved) in descriptors.iter().zip(resolved) {
142 let body = resolved.body.ok_or_else(|| {
143 GitError::InvalidObject("parallel pack parse discarded an object body".into())
144 })?;
145 entries.push(PackObject {
146 entry: PackEntry {
147 oid: resolved.oid,
148 compressed_size: compressed_size(descriptor)?,
149 uncompressed_size: resolved.size,
150 offset: descriptor.offset as u64,
151 },
152 object: EncodedObject::new(resolved.object_type, body),
153 });
154 }
155 Ok(PackFile {
156 version: u32_be(&pack[4..8]),
157 entries,
158 checksum,
159 })
160}
161
162fn finish_index(
163 pack_checksum: ObjectId,
164 descriptors: &[EntryDescriptor],
165 resolved: Vec<ResolvedEntry>,
166 format: ObjectFormat,
167) -> Result<PackIndexBuild> {
168 let entries = descriptors
169 .iter()
170 .zip(&resolved)
171 .map(|(descriptor, resolved)| PackIndexEntry {
172 oid: resolved.oid,
173 crc32: resolved.crc32,
174 offset: descriptor.offset as u64,
175 })
176 .collect::<Vec<_>>();
177 let objects = descriptors
178 .iter()
179 .zip(&resolved)
180 .map(|(descriptor, resolved)| PackIndexedObject {
181 oid: resolved.oid,
182 object_type: resolved.object_type,
183 size: resolved.size,
184 offset: descriptor.offset as u64,
185 })
186 .collect::<Vec<_>>();
187 let index = PackIndex::write_v2(format, &entries, &pack_checksum)?;
188 Ok(PackIndexBuild {
189 index,
190 pack_checksum,
191 entries,
192 objects,
193 })
194}
195
196fn discover_pack_entries(
197 pack: &[u8],
198 format: ObjectFormat,
199 options: PackIndexOptions,
200 cancel: CancelFlag<'_>,
201) -> Result<(ObjectId, Vec<EntryDescriptor>)> {
202 cancel.check()?;
203 let trailer_len = format.raw_len();
204 if pack.len() < 12 + trailer_len {
205 return Err(GitError::InvalidFormat("pack file too short".into()));
206 }
207 if &pack[..4] != b"PACK" {
208 return Err(GitError::InvalidFormat("missing PACK signature".into()));
209 }
210 let version = u32_be(&pack[4..8]);
211 if version != 2 && version != 3 {
212 return Err(GitError::Unsupported(format!("pack version {version}")));
213 }
214 let trailer_offset = pack.len() - trailer_len;
215 let count = checked_pack_object_count(
216 u32_be(&pack[8..12]),
217 trailer_offset.saturating_sub(12) as u64,
218 )?;
219 let pack_checksum = sley_core::digest_bytes(format, &pack[..trailer_offset])?;
220 let expected = ObjectId::from_raw(format, &pack[trailer_offset..])?;
221 if pack_checksum != expected {
222 return Err(GitError::InvalidFormat(format!(
223 "pack checksum mismatch: expected {expected}, got {pack_checksum}"
224 )));
225 }
226 if count == 0 {
227 if trailer_offset != 12 {
228 return Err(GitError::InvalidFormat(format!(
229 "empty pack has {} trailing bytes before checksum",
230 trailer_offset - 12
231 )));
232 }
233 return Ok((pack_checksum, Vec::new()));
234 }
235
236 #[cfg(feature = "fetch-profile")]
237 let _inflate_span =
238 sley_core::fetch_profile::Span::enter(sley_core::fetch_profile::Stage::Inflate);
239 let candidates = scan_candidates_parallel(
240 pack,
241 format,
242 trailer_offset,
243 options.threads.min(count.max(1)),
244 cancel,
245 )?;
246 #[cfg(feature = "fetch-profile")]
247 drop(_inflate_span);
248
249 let mut descriptors = Vec::with_capacity(pack_entry_prealloc(count));
250 let mut candidate_index = 0usize;
251 let mut offset = 12usize;
252 for _ in 0..count {
253 while candidate_index < candidates.len() && candidates[candidate_index].offset < offset {
254 candidate_index += 1;
255 }
256 let descriptor = candidates
257 .get(candidate_index)
258 .filter(|candidate| candidate.offset == offset)
259 .ok_or_else(|| {
260 GitError::InvalidObject(format!(
261 "pack entry at offset {offset} is not a valid zlib member"
262 ))
263 })?
264 .clone();
265 if descriptor.end_offset <= descriptor.offset {
266 return Err(GitError::InvalidFormat(
267 "empty compressed pack entry".into(),
268 ));
269 }
270 offset = descriptor.end_offset;
271 descriptors.push(descriptor);
272 candidate_index += 1;
273 }
274 if offset != trailer_offset {
275 let detail = if offset < trailer_offset {
276 format!("{} trailing bytes before checksum", trailer_offset - offset)
277 } else {
278 "entry extends past checksum".into()
279 };
280 return Err(GitError::InvalidFormat(format!("pack has {detail}")));
281 }
282 Ok((pack_checksum, descriptors))
283}
284
285fn scan_candidates_parallel(
286 pack: &[u8],
287 format: ObjectFormat,
288 trailer_offset: usize,
289 requested_threads: usize,
290 cancel: CancelFlag<'_>,
291) -> Result<Vec<EntryDescriptor>> {
292 let possible_starts = trailer_offset.saturating_sub(12);
293 let worker_count = requested_threads.max(1).min(possible_starts.max(1));
294 let chunk_len = possible_starts.div_ceil(worker_count);
295 std::thread::scope(|scope| {
296 let mut handles = Vec::with_capacity(worker_count);
297 for worker in 0..worker_count {
298 let start = 12 + worker * chunk_len;
299 let end = (start + chunk_len).min(trailer_offset);
300 if start >= end {
301 continue;
302 }
303 handles.push(scope.spawn(sley_core::diagnostics::inherit(move || {
304 scan_candidate_range(pack, format, trailer_offset, start, end, cancel)
305 })));
306 }
307 let mut candidates = Vec::new();
308 for handle in handles {
309 match handle.join() {
310 Ok(Ok(mut worker_candidates)) => candidates.append(&mut worker_candidates),
311 Ok(Err(err)) => return Err(err),
312 Err(_) => {
313 return Err(GitError::InvalidObject(
314 "parallel pack discovery worker panicked".into(),
315 ));
316 }
317 }
318 }
319 Ok(candidates)
320 })
321}
322
323fn scan_candidate_range(
324 pack: &[u8],
325 format: ObjectFormat,
326 trailer_offset: usize,
327 start: usize,
328 end: usize,
329 cancel: CancelFlag<'_>,
330) -> Result<Vec<EntryDescriptor>> {
331 let mut candidates = Vec::new();
332 let mut decompress = Decompress::new(true);
333 let mut output = vec![0u8; 64 * 1024];
334 for offset in start..end {
335 if offset & 0xfff == 0 {
336 cancel.check()?;
337 }
338 let Some(mut descriptor) = candidate_header(pack, format, trailer_offset, offset) else {
339 continue;
340 };
341 let expected = match usize::try_from(descriptor.header.size) {
342 Ok(expected) => expected,
343 Err(_) => continue,
344 };
345 let Some(consumed) = measure_zlib_member(
346 &mut decompress,
347 &pack[descriptor.data_offset..trailer_offset],
348 expected,
349 &mut output,
350 cancel,
351 )?
352 else {
353 continue;
354 };
355 let Some(end_offset) = descriptor.data_offset.checked_add(consumed) else {
356 continue;
357 };
358 if consumed == 0 || end_offset > trailer_offset {
359 continue;
360 }
361 descriptor.end_offset = end_offset;
362 candidates.push(descriptor);
363 }
364 cancel.check()?;
365 Ok(candidates)
366}
367
368fn candidate_header(
369 pack: &[u8],
370 format: ObjectFormat,
371 trailer_offset: usize,
372 offset: usize,
373) -> Option<EntryDescriptor> {
374 let first = *pack.get(offset)?;
375 let kind = match (first >> 4) & 0x07 {
376 1 => PackObjectKind::Commit,
377 2 => PackObjectKind::Tree,
378 3 => PackObjectKind::Blob,
379 4 => PackObjectKind::Tag,
380 6 => PackObjectKind::OfsDelta,
381 7 => PackObjectKind::RefDelta,
382 _ => return None,
383 };
384 let mut cursor = offset + 1;
385 let mut byte = first;
386 let mut size = u64::from(first & 0x0f);
387 let mut shift = 4u32;
388 while byte & 0x80 != 0 {
389 byte = *pack.get(cursor)?;
390 cursor = cursor.checked_add(1)?;
391 let part = u64::from(byte & 0x7f).checked_shl(shift)?;
392 size = size.checked_add(part)?;
393 shift = shift.checked_add(7)?;
394 if shift > 67 {
395 return None;
396 }
397 }
398 let base = match kind {
399 PackObjectKind::OfsDelta => {
400 let mut base_cursor = cursor;
401 let base = parse_ofs_delta_base_offset(pack, &mut base_cursor, offset as u64).ok()?;
402 cursor = base_cursor;
403 Some(DeltaBase::Offset(base))
404 }
405 PackObjectKind::RefDelta => {
406 let end = cursor.checked_add(format.raw_len())?;
407 if end > trailer_offset {
408 return None;
409 }
410 let oid = ObjectId::from_raw(format, pack.get(cursor..end)?).ok()?;
411 cursor = end;
412 Some(DeltaBase::Ref(oid))
413 }
414 _ => None,
415 };
416 if cursor >= trailer_offset || !is_zlib_header(pack.get(cursor..cursor.checked_add(2)?)?) {
417 return None;
418 }
419 Some(EntryDescriptor {
420 offset,
421 data_offset: cursor,
422 end_offset: 0,
423 header: EntryHeader { kind, size },
424 base,
425 })
426}
427
428fn is_zlib_header(bytes: &[u8]) -> bool {
429 let cmf = bytes[0];
430 let flg = bytes[1];
431 cmf & 0x0f == 8
432 && cmf >> 4 <= 7
433 && flg & 0x20 == 0
434 && u16::from_be_bytes([cmf, flg]).is_multiple_of(31)
435}
436
437fn measure_zlib_member(
438 decompress: &mut Decompress,
439 compressed: &[u8],
440 expected: usize,
441 output: &mut [u8],
442 cancel: CancelFlag<'_>,
443) -> Result<Option<usize>> {
444 decompress.reset(true);
445 let mut input = compressed;
446 loop {
447 cancel.check()?;
448 let before_in = decompress.total_in();
449 let before_out = decompress.total_out();
450 let status = match decompress.decompress(input, output, FlushDecompress::None) {
451 Ok(status) => status,
452 Err(_) => return Ok(None),
453 };
454 let consumed = (decompress.total_in() - before_in) as usize;
455 let produced = (decompress.total_out() - before_out) as usize;
456 if decompress.total_out() > expected as u64 {
457 return Ok(None);
458 }
459 input = match input.get(consumed..) {
460 Some(remaining) => remaining,
461 None => return Ok(None),
462 };
463 match status {
464 flate2::Status::StreamEnd if decompress.total_out() == expected as u64 => {
465 return Ok(Some(decompress.total_in() as usize));
466 }
467 flate2::Status::StreamEnd => return Ok(None),
468 _ if consumed == 0 && produced == 0 => return Ok(None),
469 _ => {}
470 }
471 }
472}
473
474fn resolve_entries_parallel<F, P>(
475 pack: &[u8],
476 descriptors: &[EntryDescriptor],
477 external_base: &mut F,
478 settings: ResolutionSettings,
479 cancel: CancelFlag<'_>,
480 progress: &mut P,
481) -> Result<Vec<ResolvedEntry>>
482where
483 F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
484 P: FnMut(PackIndexProgress),
485{
486 let mut offset_to_index = HashMap::with_capacity(descriptors.len());
487 for (index, descriptor) in descriptors.iter().enumerate() {
488 offset_to_index.insert(descriptor.offset as u64, index);
489 }
490 let mut ofs_bases = HashSet::new();
491 let mut ref_bases = HashSet::new();
492 for descriptor in descriptors {
493 match descriptor.base {
494 Some(DeltaBase::Offset(offset)) => {
495 let index = offset_to_index.get(&offset).copied().ok_or_else(|| {
496 GitError::InvalidFormat(format!("ofs-delta base offset {offset} not found"))
497 })?;
498 ofs_bases.insert(index);
499 }
500 Some(DeltaBase::Ref(oid)) => {
501 ref_bases.insert(oid);
502 }
503 None => {}
504 }
505 }
506
507 let base_indices = descriptors
508 .iter()
509 .enumerate()
510 .filter_map(|(index, descriptor)| descriptor.base.is_none().then_some(index))
511 .collect::<Vec<_>>();
512 let mut resolved: Vec<Option<ResolvedEntry>> = std::iter::repeat_with(|| None)
513 .take(descriptors.len())
514 .collect();
515
516 #[cfg(feature = "fetch-profile")]
517 let _inflate_span =
518 sley_core::fetch_profile::Span::enter(sley_core::fetch_profile::Stage::Inflate);
519 let mut completed = 0u64;
520 let base_results = parallel_chunks(
521 &base_indices,
522 settings.options.threads,
523 |indices| {
524 let mut output = Vec::with_capacity(indices.len());
525 for &index in indices {
526 let descriptor = &descriptors[index];
527 let body = inflate_descriptor(pack, descriptor, cancel)?;
528 let object_type = object_type_for_kind(descriptor.header.kind)?;
529 let oid = cancellable_object_id_bytes(object_type, &body, settings.format, cancel)?;
530 let keep_body =
531 settings.retain_all || ofs_bases.contains(&index) || ref_bases.contains(&oid);
532 output.push((
533 index,
534 ResolvedEntry {
535 oid,
536 object_type,
537 size: body.len() as u64,
538 crc32: crc32fast::hash(&pack[descriptor.offset..descriptor.end_offset]),
539 depth: 0,
540 body: keep_body.then_some(body),
541 },
542 ));
543 }
544 Ok(output)
545 },
546 |batch_len| {
547 completed = completed.saturating_add(batch_len as u64);
548 progress(PackIndexProgress {
549 completed_objects: completed,
550 total_objects: descriptors.len() as u64,
551 });
552 cancel.check()
553 },
554 )?;
555 #[cfg(feature = "fetch-profile")]
556 drop(_inflate_span);
557
558 let mut oid_to_index = HashMap::with_capacity(descriptors.len());
559 for (index, entry) in base_results {
560 oid_to_index.entry(entry.oid).or_insert(index);
561 resolved[index] = Some(entry);
562 }
563 if descriptors.is_empty() {
564 progress(PackIndexProgress::default());
565 cancel.check()?;
566 }
567
568 let mut unresolved = descriptors.len().saturating_sub(base_indices.len());
569 let mut external = HashMap::<ObjectId, EncodedObject>::new();
570 let mut external_missing = HashSet::<ObjectId>::new();
571 while unresolved != 0 {
572 cancel.check()?;
573 let mut ready = ready_internal_deltas(
574 descriptors,
575 &resolved,
576 &offset_to_index,
577 &oid_to_index,
578 settings.options.limits,
579 )?;
580 if ready.is_empty() {
581 ready = ready_external_deltas(
582 descriptors,
583 &resolved,
584 &mut external,
585 &mut external_missing,
586 external_base,
587 settings.format,
588 settings.options.limits,
589 )?;
590 }
591 if ready.is_empty() {
592 return Err(GitError::Unsupported(
593 "unresolved, cyclic, or mis-ordered delta base".into(),
594 ));
595 }
596
597 #[cfg(feature = "fetch-profile")]
598 let _delta_span =
599 sley_core::fetch_profile::Span::enter(sley_core::fetch_profile::Stage::DeltaResolution);
600 let batch_results = parallel_chunks(
601 &ready,
602 settings.options.threads,
603 |tasks| {
604 let mut output = Vec::with_capacity(tasks.len());
605 for task in tasks {
606 output.push((
607 task.index,
608 resolve_delta_entry(
609 pack,
610 settings.format,
611 &descriptors[task.index],
612 *task,
613 &resolved,
614 &external,
615 settings.retain_all,
616 &ofs_bases,
617 &ref_bases,
618 cancel,
619 )?,
620 ));
621 }
622 Ok(output)
623 },
624 |batch_len| {
625 completed = completed.saturating_add(batch_len as u64);
626 progress(PackIndexProgress {
627 completed_objects: completed,
628 total_objects: descriptors.len() as u64,
629 });
630 cancel.check()
631 },
632 )?;
633 #[cfg(feature = "fetch-profile")]
634 {
635 sley_core::fetch_profile::add_count(
636 sley_core::fetch_profile::Stage::DeltaResolution,
637 batch_results.len() as u64,
638 );
639 drop(_delta_span);
640 }
641 for (index, entry) in batch_results {
642 oid_to_index.entry(entry.oid).or_insert(index);
643 resolved[index] = Some(entry);
644 unresolved -= 1;
645 }
646 }
647
648 resolved
649 .into_iter()
650 .map(|entry| entry.ok_or_else(|| GitError::InvalidFormat("unresolved pack entry".into())))
651 .collect()
652}
653
654fn ready_internal_deltas(
655 descriptors: &[EntryDescriptor],
656 resolved: &[Option<ResolvedEntry>],
657 offset_to_index: &HashMap<u64, usize>,
658 oid_to_index: &HashMap<ObjectId, usize>,
659 limits: PackReadLimits,
660) -> Result<Vec<ReadyDelta>> {
661 let mut ready = Vec::new();
662 for (index, descriptor) in descriptors.iter().enumerate() {
663 if resolved[index].is_some() {
664 continue;
665 }
666 let base_index = match descriptor.base {
667 Some(DeltaBase::Offset(offset)) => offset_to_index.get(&offset).copied(),
668 Some(DeltaBase::Ref(oid)) => oid_to_index.get(&oid).copied(),
669 None => None,
670 };
671 let Some(base_index) = base_index else {
672 continue;
673 };
674 let Some(base) = resolved[base_index].as_ref() else {
675 continue;
676 };
677 let depth = base.depth + 1;
678 check_delta_depth(descriptor.offset, depth, limits)?;
679 ready.push(ReadyDelta {
680 index,
681 base: ReadyBase::Internal(base_index),
682 depth,
683 });
684 }
685 Ok(ready)
686}
687
688fn ready_external_deltas<F>(
689 descriptors: &[EntryDescriptor],
690 resolved: &[Option<ResolvedEntry>],
691 external: &mut HashMap<ObjectId, EncodedObject>,
692 external_missing: &mut HashSet<ObjectId>,
693 external_base: &mut F,
694 format: ObjectFormat,
695 limits: PackReadLimits,
696) -> Result<Vec<ReadyDelta>>
697where
698 F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
699{
700 let mut ready = Vec::new();
701 for (index, descriptor) in descriptors.iter().enumerate() {
702 if resolved[index].is_some() {
703 continue;
704 }
705 let Some(DeltaBase::Ref(oid)) = descriptor.base else {
706 continue;
707 };
708 if !external.contains_key(&oid) && !external_missing.contains(&oid) {
709 match external_base(&oid)? {
710 Some(object) => {
711 let actual = object.object_id(format)?;
712 if actual != oid {
713 return Err(GitError::InvalidObject(format!(
714 "external delta base {oid} hashes to {actual}"
715 )));
716 }
717 external.insert(oid, object);
718 }
719 None => {
720 external_missing.insert(oid);
721 }
722 }
723 }
724 if external.contains_key(&oid) {
725 check_delta_depth(descriptor.offset, 1, limits)?;
726 ready.push(ReadyDelta {
727 index,
728 base: ReadyBase::External(oid),
729 depth: 1,
730 });
731 }
732 }
733 Ok(ready)
734}
735
736#[allow(clippy::too_many_arguments)]
737fn resolve_delta_entry(
738 pack: &[u8],
739 format: ObjectFormat,
740 descriptor: &EntryDescriptor,
741 task: ReadyDelta,
742 resolved: &[Option<ResolvedEntry>],
743 external: &HashMap<ObjectId, EncodedObject>,
744 retain_all: bool,
745 ofs_bases: &HashSet<usize>,
746 ref_bases: &HashSet<ObjectId>,
747 cancel: CancelFlag<'_>,
748) -> Result<ResolvedEntry> {
749 let (base_type, base_body) = match task.base {
750 ReadyBase::Internal(index) => {
751 let base = resolved[index]
752 .as_ref()
753 .ok_or_else(|| GitError::InvalidFormat("delta base is not resolved".into()))?;
754 let body = base.body.as_deref().ok_or_else(|| {
755 GitError::InvalidFormat("delta base body was released before use".into())
756 })?;
757 (base.object_type, body)
758 }
759 ReadyBase::External(oid) => {
760 let base = external.get(&oid).ok_or_else(|| {
761 GitError::InvalidFormat("external delta base is not available".into())
762 })?;
763 (base.object_type, base.body.as_slice())
764 }
765 };
766 let delta = inflate_descriptor(pack, descriptor, cancel)?;
767 if delta.len() as u64 != descriptor.header.size {
768 return Err(GitError::InvalidObject(format!(
769 "pack delta declared {} bytes, decoded {}",
770 descriptor.header.size,
771 delta.len()
772 )));
773 }
774 let plan = plan_pack_delta(base_body, &delta)?;
775 let result_size = usize::try_from(plan.result_size)
776 .map_err(|_| GitError::InvalidObject("delta result size overflows usize".into()))?;
777 let mut body = Vec::new();
778 body.try_reserve_exact(result_size)
779 .map_err(|_| GitError::InvalidObject("could not allocate delta result".into()))?;
780 apply_pack_delta_exact(base_body, &delta, plan, &mut body, cancel)?;
781 let oid = cancellable_object_id_bytes(base_type, &body, format, cancel)?;
782 let keep_body = retain_all || ofs_bases.contains(&task.index) || ref_bases.contains(&oid);
783 Ok(ResolvedEntry {
784 oid,
785 object_type: base_type,
786 size: body.len() as u64,
787 crc32: crc32fast::hash(&pack[descriptor.offset..descriptor.end_offset]),
788 depth: task.depth,
789 body: keep_body.then_some(body),
790 })
791}
792
793fn check_delta_depth(offset: usize, depth: usize, limits: PackReadLimits) -> Result<()> {
794 if depth <= limits.max_delta_depth {
795 return Ok(());
796 }
797 Err(GitError::InvalidFormat(format!(
798 "pack delta chain at offset {offset} has observed depth {depth}, which exceeds maximum \
799 depth (configured limit {}); raise PackReadLimits::max_delta_depth or run `git repack \
800 --depth={}`",
801 limits.max_delta_depth, limits.max_delta_depth
802 )))
803}
804
805fn inflate_descriptor(
806 pack: &[u8],
807 descriptor: &EntryDescriptor,
808 cancel: CancelFlag<'_>,
809) -> Result<Vec<u8>> {
810 let expected = usize::try_from(descriptor.header.size)
811 .map_err(|_| GitError::InvalidObject("pack object size overflows usize".into()))?;
812 let (body, consumed) = inflate_exact(
813 &pack[descriptor.data_offset..descriptor.end_offset],
814 expected,
815 cancel,
816 )?;
817 let expected_consumed = descriptor.end_offset - descriptor.data_offset;
818 if consumed != expected_consumed {
819 return Err(GitError::InvalidObject(format!(
820 "pack entry compressed span changed during parallel inflate: expected \
821 {expected_consumed}, consumed {consumed}"
822 )));
823 }
824 Ok(body)
825}
826
827fn inflate_exact(
828 compressed: &[u8],
829 expected: usize,
830 cancel: CancelFlag<'_>,
831) -> Result<(Vec<u8>, usize)> {
832 let mut output = Vec::new();
833 output
834 .try_reserve_exact(expected)
835 .map_err(|_| GitError::InvalidObject("could not allocate pack object".into()))?;
836 let mut decompress = Decompress::new(true);
837 let mut input = compressed;
838 let mut overflow = [0u8; 1];
839 loop {
840 cancel.check()?;
841 let before_in = decompress.total_in();
842 let before_out = decompress.total_out();
843 let checking_overflow = output.len() == expected;
844 let status = if checking_overflow {
845 decompress.decompress(input, &mut overflow, FlushDecompress::None)
846 } else {
847 decompress.decompress_vec(input, &mut output, FlushDecompress::None)
848 }
849 .map_err(|error| GitError::InvalidObject(format!("zlib inflate failed: {error}")))?;
850 let consumed = (decompress.total_in() - before_in) as usize;
851 let produced = (decompress.total_out() - before_out) as usize;
852 if output.len() > expected || (checking_overflow && produced != 0) {
853 return Err(GitError::InvalidObject(format!(
854 "pack object declared {expected} bytes, decoded more than {expected}"
855 )));
856 }
857 input = input.get(consumed..).ok_or_else(|| {
858 GitError::InvalidObject("zlib consumed beyond pack entry input".into())
859 })?;
860 match status {
861 flate2::Status::StreamEnd if output.len() == expected => {
862 return Ok((output, decompress.total_in() as usize));
863 }
864 flate2::Status::StreamEnd => {
865 return Err(GitError::InvalidObject(format!(
866 "pack object declared {expected} bytes, decoded {}",
867 output.len()
868 )));
869 }
870 _ if consumed == 0 && produced == 0 => {
871 return Err(GitError::InvalidObject("truncated zlib stream".into()));
872 }
873 _ => {}
874 }
875 }
876}
877
878fn object_type_for_kind(kind: PackObjectKind) -> Result<ObjectType> {
879 match kind {
880 PackObjectKind::Commit => Ok(ObjectType::Commit),
881 PackObjectKind::Tree => Ok(ObjectType::Tree),
882 PackObjectKind::Blob => Ok(ObjectType::Blob),
883 PackObjectKind::Tag => Ok(ObjectType::Tag),
884 PackObjectKind::OfsDelta | PackObjectKind::RefDelta => Err(GitError::InvalidFormat(
885 "delta entry cannot be used as an object type".into(),
886 )),
887 }
888}
889
890fn cancellable_object_id_bytes(
891 object_type: ObjectType,
892 body: &[u8],
893 format: ObjectFormat,
894 cancel: CancelFlag<'_>,
895) -> Result<ObjectId> {
896 cancel.check()?;
897 #[cfg(feature = "fetch-profile")]
898 let _oid_span = sley_core::fetch_profile::Span::enter(sley_core::fetch_profile::Stage::OidHash);
899 let mut digest = StreamingDigest::new(format);
900 digest.update(object_type.as_str().as_bytes());
901 digest.update(b" ");
902 digest.update(body.len().to_string().as_bytes());
903 digest.update(b"\0");
904 for chunk in body.chunks(256 * 1024) {
905 cancel.check()?;
906 digest.update(chunk);
907 }
908 let oid = digest.finalize()?;
909 #[cfg(feature = "fetch-profile")]
910 {
911 sley_core::fetch_profile::add_count(sley_core::fetch_profile::Stage::OidHash, 1);
912 sley_core::fetch_profile::add_bytes(
913 sley_core::fetch_profile::Stage::OidHash,
914 body.len() as u64,
915 );
916 drop(_oid_span);
917 }
918 Ok(oid)
919}
920
921fn compressed_size(descriptor: &EntryDescriptor) -> Result<u64> {
922 u64::try_from(descriptor.end_offset - descriptor.data_offset)
923 .map_err(|_| GitError::InvalidFormat("compressed pack entry size overflows u64".into()))
924}
925
926fn parallel_chunks<T, U, F, P>(
927 items: &[T],
928 threads: usize,
929 work: F,
930 mut batch_complete: P,
931) -> Result<Vec<U>>
932where
933 T: Sync,
934 U: Send,
935 F: Fn(&[T]) -> Result<Vec<U>> + Sync,
936 P: FnMut(usize) -> Result<()>,
937{
938 if items.is_empty() {
939 return Ok(Vec::new());
940 }
941 let worker_count = threads.max(1).min(items.len());
942 let chunk_len = items.len().div_ceil(worker_count);
943 std::thread::scope(|scope| {
944 let mut handles = Vec::with_capacity(worker_count);
945 for chunk in items.chunks(chunk_len) {
946 let work = &work;
947 handles.push(scope.spawn(sley_core::diagnostics::inherit(move || work(chunk))));
948 }
949 let mut output = Vec::with_capacity(items.len());
950 for handle in handles {
951 match handle.join() {
952 Ok(Ok(mut values)) => {
953 batch_complete(values.len())?;
954 output.append(&mut values);
955 }
956 Ok(Err(err)) => return Err(err),
957 Err(_) => {
958 return Err(GitError::InvalidObject(
959 "parallel pack worker panicked".into(),
960 ));
961 }
962 }
963 }
964 Ok(output)
965 })
966}