Skip to main content

sley_pack/
parallel_index.rs

1//! Parallel pack discovery, inflate, delta resolution, and index construction.
2//!
3//! Pack entries do not carry their compressed length. The discovery pass scans
4//! possible entry starts in parallel, validates each zlib member without
5//! retaining its output, and then follows the unique entry chain from byte 12.
6//! Selected entries are inflated once more for materialization, also in
7//! parallel. Delta bodies are retained only until their dependency level is
8//! resolved; resolved object bodies are retained only when another entry names
9//! them as a base (or when a caller requests a fully materialized [`PackFile`]).
10
11use super::*;
12use flate2::{Decompress, FlushDecompress};
13
14/// Configuration for the one pack indexing engine.
15///
16/// `threads` is explicit so callers and tests can prove scheduling-independent
17/// output. [`Default`] uses all logical CPUs reported by the operating system.
18#[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}