Skip to main content

sley_pack/
stream.rs

1//! Streaming pack reads for index-pack without buffering the whole pack.
2//!
3//! Split out of `lib.rs` in the W21 mechanical refactor: a pure code move
4//! (no function body changed); all items are re-exported from `lib.rs`.
5use super::*;
6
7/// Validate and index a seekable pack stream without writing into a repository.
8///
9/// The reader is consumed from its current position through EOF and left at EOF
10/// on success. Embedders can use the returned bytes as the pack's v2 `.idx`.
11pub fn index_pack_from_reader<R>(
12    reader: &mut R,
13    format: ObjectFormat,
14) -> Result<PackStreamIndexBuild>
15where
16    R: Read + Seek,
17{
18    index_pack_from_reader_with_limits(reader, format, PackReadLimits::default())
19}
20
21/// [`index_pack_from_reader`] with an external ref-delta base resolver.
22pub fn index_pack_from_reader_with_base<R, F>(
23    reader: &mut R,
24    format: ObjectFormat,
25    external_base: F,
26) -> Result<PackStreamIndexBuild>
27where
28    R: Read + Seek,
29    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
30{
31    index_pack_from_reader_with_base_and_limits(
32        reader,
33        format,
34        external_base,
35        PackReadLimits::default(),
36    )
37}
38
39pub fn index_pack_from_reader_with_limits<R>(
40    reader: &mut R,
41    format: ObjectFormat,
42    limits: PackReadLimits,
43) -> Result<PackStreamIndexBuild>
44where
45    R: Read + Seek,
46{
47    let start = reader.stream_position()?;
48    let end = reader.seek(SeekFrom::End(0))?;
49    let pack_len = end
50        .checked_sub(start)
51        .ok_or_else(|| GitError::InvalidFormat("pack stream position overflow".into()))?;
52    reader.seek(SeekFrom::Start(start))?;
53    index_pack_from_reader_with_len_and_limits(reader, format, pack_len, limits)
54}
55
56pub fn index_pack_from_reader_with_base_and_limits<R, F>(
57    reader: &mut R,
58    format: ObjectFormat,
59    external_base: F,
60    limits: PackReadLimits,
61) -> Result<PackStreamIndexBuild>
62where
63    R: Read + Seek,
64    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
65{
66    let start = reader.stream_position()?;
67    let end = reader.seek(SeekFrom::End(0))?;
68    let pack_len = end
69        .checked_sub(start)
70        .ok_or_else(|| GitError::InvalidFormat("pack stream position overflow".into()))?;
71    reader.seek(SeekFrom::Start(start))?;
72    index_pack_from_reader_with_len_base_and_limits(reader, format, pack_len, external_base, limits)
73}
74
75pub(crate) fn index_pack_from_reader_with_len_and_limits<R>(
76    reader: &mut R,
77    format: ObjectFormat,
78    pack_len: u64,
79    limits: PackReadLimits,
80) -> Result<PackStreamIndexBuild>
81where
82    R: Read,
83{
84    index_pack_from_stream_with_base_and_limits(
85        PackReadStream::new(reader, format, Some(pack_len))?,
86        format,
87        |_| Ok(None),
88        limits,
89    )
90}
91
92pub(crate) fn index_pack_from_reader_with_len_base_and_limits<R, F>(
93    reader: &mut R,
94    format: ObjectFormat,
95    pack_len: u64,
96    external_base: F,
97    limits: PackReadLimits,
98) -> Result<PackStreamIndexBuild>
99where
100    R: Read,
101    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
102{
103    index_pack_from_stream_with_base_and_limits(
104        PackReadStream::new(reader, format, Some(pack_len))?,
105        format,
106        external_base,
107        limits,
108    )
109}
110
111/// Validate and index one pack from a stream, stopping after its trailer.
112///
113/// This variant supports non-seekable upload-pack streams whose length is not
114/// known in advance and returns the v2 `.idx` without installing the pack.
115pub fn index_pack_from_reader_to_trailer<R>(
116    reader: &mut R,
117    format: ObjectFormat,
118) -> Result<PackStreamIndexBuild>
119where
120    R: Read,
121{
122    index_pack_from_reader_to_trailer_with_limits(reader, format, PackReadLimits::default())
123}
124
125/// [`index_pack_from_reader_to_trailer`] with an external ref-delta resolver.
126pub fn index_pack_from_reader_to_trailer_with_base<R, F>(
127    reader: &mut R,
128    format: ObjectFormat,
129    external_base: F,
130) -> Result<PackStreamIndexBuild>
131where
132    R: Read,
133    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
134{
135    index_pack_from_reader_to_trailer_with_base_and_limits(
136        reader,
137        format,
138        external_base,
139        PackReadLimits::default(),
140    )
141}
142
143pub fn index_pack_from_reader_to_trailer_with_limits<R>(
144    reader: &mut R,
145    format: ObjectFormat,
146    limits: PackReadLimits,
147) -> Result<PackStreamIndexBuild>
148where
149    R: Read,
150{
151    index_pack_from_stream_with_base_and_limits(
152        PackReadStream::new(reader, format, None)?,
153        format,
154        |_| Ok(None),
155        limits,
156    )
157}
158
159pub fn index_pack_from_reader_to_trailer_with_base_and_limits<R, F>(
160    reader: &mut R,
161    format: ObjectFormat,
162    external_base: F,
163    limits: PackReadLimits,
164) -> Result<PackStreamIndexBuild>
165where
166    R: Read,
167    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
168{
169    index_pack_from_stream_with_base_and_limits(
170        PackReadStream::new(reader, format, None)?,
171        format,
172        external_base,
173        limits,
174    )
175}
176
177/// Index one non-seekable pack through its trailer and report receive progress.
178///
179/// Embedders can use this while streaming upload-pack data to surface byte and
180/// object counters without writing the pack into a repository.
181pub fn index_pack_from_reader_to_trailer_with_progress<R, F>(
182    reader: &mut R,
183    format: ObjectFormat,
184    progress: F,
185) -> Result<PackStreamIndexBuild>
186where
187    R: Read,
188    F: FnMut(PackStreamProgress),
189{
190    index_pack_from_reader_to_trailer_with_progress_and_limits(
191        reader,
192        format,
193        PackReadLimits::default(),
194        progress,
195    )
196}
197
198pub fn index_pack_from_reader_to_trailer_with_progress_and_limits<R, F>(
199    reader: &mut R,
200    format: ObjectFormat,
201    limits: PackReadLimits,
202    progress: F,
203) -> Result<PackStreamIndexBuild>
204where
205    R: Read,
206    F: FnMut(PackStreamProgress),
207{
208    index_pack_from_reader_to_trailer_with_progress_and_cancel_and_limits(
209        reader,
210        format,
211        CancelFlag::never(),
212        limits,
213        progress,
214    )
215}
216
217pub(crate) fn index_pack_from_reader_to_trailer_with_progress_and_cancel_and_limits<R, F>(
218    reader: &mut R,
219    format: ObjectFormat,
220    cancel: CancelFlag<'_>,
221    limits: PackReadLimits,
222    progress: F,
223) -> Result<PackStreamIndexBuild>
224where
225    R: Read,
226    F: FnMut(PackStreamProgress),
227{
228    index_pack_from_stream_with_progress_and_cancel_and_limits(
229        PackReadStream::new(reader, format, None)?,
230        format,
231        cancel,
232        limits,
233        progress,
234    )
235}
236
237/// Index a prepared [`PackReadStream`] without progress callbacks.
238///
239/// This lower-level entry point lets embedders select bounded or trailer-based
240/// streaming before producing the pack's v2 `.idx` bytes.
241pub fn index_pack_from_stream<R>(
242    stream: PackReadStream<'_, R>,
243    format: ObjectFormat,
244) -> Result<PackStreamIndexBuild>
245where
246    R: Read,
247{
248    index_pack_from_stream_with_limits(stream, format, PackReadLimits::default())
249}
250
251/// [`index_pack_from_stream`] with an external ref-delta base resolver.
252pub fn index_pack_from_stream_with_base<R, F>(
253    stream: PackReadStream<'_, R>,
254    format: ObjectFormat,
255    external_base: F,
256) -> Result<PackStreamIndexBuild>
257where
258    R: Read,
259    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
260{
261    index_pack_from_stream_with_base_and_limits(
262        stream,
263        format,
264        external_base,
265        PackReadLimits::default(),
266    )
267}
268
269pub fn index_pack_from_stream_with_limits<R>(
270    stream: PackReadStream<'_, R>,
271    format: ObjectFormat,
272    limits: PackReadLimits,
273) -> Result<PackStreamIndexBuild>
274where
275    R: Read,
276{
277    index_pack_from_stream_with_progress_and_limits(stream, format, limits, |_| {})
278}
279
280pub fn index_pack_from_stream_with_base_and_limits<R, F>(
281    stream: PackReadStream<'_, R>,
282    format: ObjectFormat,
283    external_base: F,
284    limits: PackReadLimits,
285) -> Result<PackStreamIndexBuild>
286where
287    R: Read,
288    F: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
289{
290    index_pack_from_stream_with_base_progress_and_cancel_and_limits(
291        stream,
292        format,
293        external_base,
294        CancelFlag::never(),
295        limits,
296        |_| {},
297    )
298}
299
300/// Approximate cadence for progress emission: report at least every this many
301/// pack bytes, matching how git paces "Receiving objects" (no per-object churn).
302pub(crate) const PROGRESS_BYTE_STEP: u64 = 256 * 1024;
303
304/// Index a prepared [`PackReadStream`] and report receive progress.
305///
306/// The callback receives monotonic byte and object counters while the v2
307/// `.idx` is produced entirely in memory.
308pub fn index_pack_from_stream_with_progress<R, F>(
309    stream: PackReadStream<'_, R>,
310    format: ObjectFormat,
311    progress: F,
312) -> Result<PackStreamIndexBuild>
313where
314    R: Read,
315    F: FnMut(PackStreamProgress),
316{
317    index_pack_from_stream_with_progress_and_limits(
318        stream,
319        format,
320        PackReadLimits::default(),
321        progress,
322    )
323}
324
325pub fn index_pack_from_stream_with_progress_and_limits<R, F>(
326    stream: PackReadStream<'_, R>,
327    format: ObjectFormat,
328    limits: PackReadLimits,
329    progress: F,
330) -> Result<PackStreamIndexBuild>
331where
332    R: Read,
333    F: FnMut(PackStreamProgress),
334{
335    index_pack_from_stream_with_progress_and_cancel_and_limits(
336        stream,
337        format,
338        CancelFlag::never(),
339        limits,
340        progress,
341    )
342}
343
344pub(crate) fn index_pack_from_stream_with_progress_and_cancel_and_limits<R, F>(
345    stream: PackReadStream<'_, R>,
346    format: ObjectFormat,
347    cancel: CancelFlag<'_>,
348    limits: PackReadLimits,
349    progress: F,
350) -> Result<PackStreamIndexBuild>
351where
352    R: Read,
353    F: FnMut(PackStreamProgress),
354{
355    index_pack_from_stream_with_base_progress_and_cancel_and_limits(
356        stream,
357        format,
358        |_| Ok(None),
359        cancel,
360        limits,
361        progress,
362    )
363}
364
365pub(crate) fn index_pack_from_stream_with_base_progress_and_cancel_and_limits<R, B, F>(
366    mut stream: PackReadStream<'_, R>,
367    format: ObjectFormat,
368    mut external_base: B,
369    cancel: CancelFlag<'_>,
370    limits: PackReadLimits,
371    mut progress: F,
372) -> Result<PackStreamIndexBuild>
373where
374    R: Read,
375    B: FnMut(&ObjectId) -> Result<Option<EncodedObject>>,
376    F: FnMut(PackStreamProgress),
377{
378    let mut header = [0u8; 12];
379    stream.read_pack_bytes(&mut header)?;
380    if &header[..4] != b"PACK" {
381        return Err(GitError::InvalidFormat("missing PACK signature".into()));
382    }
383    let version = u32_be(&header[4..8]);
384    if version != 2 && version != 3 {
385        return Err(GitError::Unsupported(format!("pack version {version}")));
386    }
387    let count = u32_be(&header[8..12]) as usize;
388    let total_objects = count as u64;
389    // Emit an initial sample so the consumer learns `total_objects` (the
390    // percentage denominator) as soon as the header is parsed.
391    progress(PackStreamProgress {
392        received_bytes: stream.pack_offset(),
393        received_objects: 0,
394        total_objects,
395    });
396    cancel.check()?;
397    // Throttle per-object emission: every ~1% of objects or `PROGRESS_BYTE_STEP`
398    // bytes, whichever the loop hits first, plus a guaranteed final sample.
399    let object_step = (total_objects / 100).max(1);
400    let mut last_emit_bytes = stream.pack_offset();
401    let mut last_emit_objects = 0u64;
402    // sley#4: the stream has not been read past the header yet, so there is no
403    // length to validate the declared count against — cap the reservation and
404    // let the entry loop's own bounds reject a count the stream cannot honour.
405    let mut parsed_entries = Vec::with_capacity(pack_entry_prealloc(count));
406    let mut raw_entries = Vec::with_capacity(pack_entry_prealloc(count));
407    for index in 0..count {
408        cancel.check()?;
409        let entry_offset = stream.pack_offset();
410        let mut entry_crc = crc32fast::Hasher::new();
411        let header = parse_entry_header_from_stream(&mut stream, &mut entry_crc)?;
412        let base = match header.kind {
413            PackObjectKind::OfsDelta => Some(DeltaBase::Offset(
414                parse_ofs_delta_base_offset_from_stream(&mut stream, &mut entry_crc, entry_offset)?,
415            )),
416            PackObjectKind::RefDelta => {
417                let mut raw = vec![0u8; format.raw_len()];
418                stream.read_entry_bytes(&mut raw, &mut entry_crc)?;
419                Some(DeltaBase::Ref(ObjectId::from_raw(format, &raw)?))
420            }
421            _ => None,
422        };
423        let (body, consumed) = inflate_entry_from_stream(
424            &mut stream,
425            &mut entry_crc,
426            header.size.min(usize::MAX as u64) as usize,
427        )?;
428        if body.len() as u64 != header.size {
429            return Err(GitError::InvalidObject(format!(
430                "pack object declared {} bytes, decoded {}",
431                header.size,
432                body.len()
433            )));
434        }
435        if consumed == 0 {
436            return Err(GitError::InvalidFormat(
437                "empty compressed pack entry".into(),
438            ));
439        }
440        raw_entries.push((entry_offset, entry_crc.finalize()));
441        if let Some(base) = base {
442            parsed_entries.push(ParsedPackEntry::Delta {
443                base,
444                compressed_size: consumed as u64,
445                delta_size: header.size,
446                offset: entry_offset,
447                delta: body,
448            });
449        } else {
450            let object_type = pack_object_kind_to_object_type(header.kind)?;
451            let object = EncodedObject::new(object_type, body);
452            let oid = object.object_id(format)?;
453            parsed_entries.push(ParsedPackEntry::Resolved(PackObject {
454                entry: PackEntry {
455                    oid,
456                    compressed_size: consumed as u64,
457                    uncompressed_size: header.size,
458                    offset: entry_offset,
459                },
460                object,
461            }));
462        }
463        let received_objects = index as u64 + 1;
464        let received_bytes = stream.pack_offset();
465        if received_objects == total_objects
466            || received_objects - last_emit_objects >= object_step
467            || received_bytes - last_emit_bytes >= PROGRESS_BYTE_STEP
468        {
469            last_emit_objects = received_objects;
470            last_emit_bytes = received_bytes;
471            progress(PackStreamProgress {
472                received_bytes,
473                received_objects,
474                total_objects,
475            });
476            cancel.check()?;
477        }
478    }
479    if stream.pack_offset() != stream.trailer_pack_offset() {
480        return Err(GitError::InvalidFormat(format!(
481            "pack has {} trailing bytes before checksum",
482            stream.trailer_pack_offset() - stream.pack_offset()
483        )));
484    }
485    let expected = stream.read_trailer_oid()?;
486    let pack_checksum = stream.finish_digest()?;
487    if pack_checksum != expected {
488        return Err(GitError::InvalidFormat(format!(
489            "pack checksum mismatch: expected {expected}, got {pack_checksum}"
490        )));
491    }
492
493    let resolved = resolve_pack_entries(parsed_entries, format, &mut external_base, limits)?;
494    let entries = resolved
495        .iter()
496        .zip(raw_entries)
497        .map(|(object, (offset, crc32))| PackIndexEntry {
498            oid: object.entry.oid,
499            crc32,
500            offset,
501        })
502        .collect::<Vec<_>>();
503    let objects = resolved
504        .iter()
505        .map(|object| PackIndexedObject {
506            oid: object.entry.oid,
507            object_type: object.object.object_type,
508            size: object.object.body.len() as u64,
509            offset: object.entry.offset,
510        })
511        .collect::<Vec<_>>();
512    let index = PackIndex::write_v2(format, &entries, &pack_checksum)?;
513    Ok(PackStreamIndexBuild {
514        index,
515        pack_checksum,
516        entries,
517        objects,
518    })
519}
520
521pub(crate) fn pack_object_kind_to_object_type(kind: PackObjectKind) -> Result<ObjectType> {
522    match kind {
523        PackObjectKind::Commit => Ok(ObjectType::Commit),
524        PackObjectKind::Tree => Ok(ObjectType::Tree),
525        PackObjectKind::Blob => Ok(ObjectType::Blob),
526        PackObjectKind::Tag => Ok(ObjectType::Tag),
527        PackObjectKind::OfsDelta | PackObjectKind::RefDelta => Err(GitError::InvalidFormat(
528            "delta entry cannot be used as an object type".into(),
529        )),
530    }
531}
532
533/// A pack-byte reader configured for a known length or trailer-delimited input.
534///
535/// Construct this when using [`index_pack_from_stream`] directly; pass a pack
536/// length for bounded seekable input or `None` to stop at the pack trailer.
537pub struct PackReadStream<'a, R> {
538    reader: &'a mut R,
539    position: u64,
540    pack_len: Option<u64>,
541    trailer_position: Option<u64>,
542    digest: StreamingDigest,
543    format: ObjectFormat,
544    pending: VecDeque<u8>,
545}
546
547impl<'a, R> PackReadStream<'a, R>
548where
549    R: Read,
550{
551    /// Wrap `reader` at its current position for streaming pack indexing.
552    pub fn new(reader: &'a mut R, format: ObjectFormat, pack_len: Option<u64>) -> Result<Self> {
553        let trailer_len = format.raw_len() as u64;
554        let trailer_position = pack_len
555            .map(|pack_len| {
556                if pack_len < 12 + trailer_len {
557                    return Err(GitError::InvalidFormat("pack file too short".into()));
558                }
559                Ok(pack_len - trailer_len)
560            })
561            .transpose()?;
562        Ok(Self {
563            reader,
564            position: 0,
565            pack_len,
566            trailer_position,
567            digest: StreamingDigest::new(format),
568            format,
569            pending: VecDeque::new(),
570        })
571    }
572
573    pub(crate) fn pack_offset(&self) -> u64 {
574        self.position
575    }
576
577    pub(crate) fn trailer_pack_offset(&self) -> u64 {
578        self.trailer_position.unwrap_or(self.position)
579    }
580
581    pub(crate) fn read_pack_bytes(&mut self, bytes: &mut [u8]) -> Result<()> {
582        let end = self
583            .position
584            .checked_add(bytes.len() as u64)
585            .ok_or_else(|| GitError::InvalidFormat("pack offset overflow".into()))?;
586        if self
587            .trailer_position
588            .is_some_and(|trailer_position| end > trailer_position)
589        {
590            return Err(GitError::InvalidFormat(
591                "pack entry extends past checksum".into(),
592            ));
593        }
594        self.read_exact_raw(bytes)?;
595        self.position = end;
596        self.digest.update(bytes);
597        Ok(())
598    }
599
600    pub(crate) fn read_exact_raw(&mut self, bytes: &mut [u8]) -> Result<()> {
601        let mut written = 0usize;
602        while written < bytes.len() {
603            if let Some(byte) = self.pending.pop_front() {
604                bytes[written] = byte;
605                written += 1;
606                continue;
607            }
608            self.reader.read_exact(&mut bytes[written..])?;
609            break;
610        }
611        Ok(())
612    }
613
614    pub(crate) fn read_entry_bytes(
615        &mut self,
616        bytes: &mut [u8],
617        crc: &mut crc32fast::Hasher,
618    ) -> Result<()> {
619        self.read_pack_bytes(bytes)?;
620        crc.update(bytes);
621        Ok(())
622    }
623
624    pub(crate) fn read_entry_byte(&mut self, crc: &mut crc32fast::Hasher) -> Result<u8> {
625        let mut byte = [0u8; 1];
626        self.read_entry_bytes(&mut byte, crc)?;
627        Ok(byte[0])
628    }
629
630    pub(crate) fn read_compressed_chunk(&mut self, bytes: &mut [u8]) -> Result<usize> {
631        let len = if let Some(trailer_position) = self.trailer_position {
632            if self.position >= trailer_position {
633                return Ok(0);
634            }
635            let remaining = trailer_position - self.position;
636            if remaining < bytes.len() as u64 {
637                remaining as usize
638            } else {
639                bytes.len()
640            }
641        } else {
642            bytes.len()
643        };
644        let mut read = 0usize;
645        while read < len {
646            let Some(byte) = self.pending.pop_front() else {
647                break;
648            };
649            bytes[read] = byte;
650            read += 1;
651        }
652        if read < len {
653            read += self.reader.read(&mut bytes[read..len])?;
654        }
655        self.position = self
656            .position
657            .checked_add(read as u64)
658            .ok_or_else(|| GitError::InvalidFormat("pack offset overflow".into()))?;
659        Ok(read)
660    }
661
662    pub(crate) fn accept_compressed_bytes(&mut self, bytes: &[u8], crc: &mut crc32fast::Hasher) {
663        self.digest.update(bytes);
664        crc.update(bytes);
665    }
666
667    pub(crate) fn push_back_compressed_bytes(&mut self, bytes: &[u8]) -> Result<()> {
668        if bytes.is_empty() {
669            return Ok(());
670        }
671        self.position = self
672            .position
673            .checked_sub(bytes.len() as u64)
674            .ok_or_else(|| GitError::InvalidFormat("pack offset overflow".into()))?;
675        for byte in bytes.iter().rev() {
676            self.pending.push_front(*byte);
677        }
678        Ok(())
679    }
680
681    pub(crate) fn read_trailer_oid(&mut self) -> Result<ObjectId> {
682        let mut raw = vec![0u8; self.format.raw_len()];
683        self.read_exact_raw(&mut raw)?;
684        self.position = self
685            .position
686            .checked_add(raw.len() as u64)
687            .ok_or_else(|| GitError::InvalidFormat("pack offset overflow".into()))?;
688        if let Some(pack_len) = self.pack_len
689            && self.position != pack_len
690        {
691            return Err(GitError::InvalidFormat(format!(
692                "pack has {} trailing bytes after checksum",
693                pack_len - self.position
694            )));
695        }
696        if self.pack_len.is_none() && !self.pending.is_empty() {
697            return Err(GitError::InvalidFormat(
698                "pack has trailing bytes after checksum".into(),
699            ));
700        }
701        ObjectId::from_raw(self.format, &raw)
702    }
703
704    pub(crate) fn finish_digest(self) -> Result<ObjectId> {
705        self.digest.finalize()
706    }
707}
708
709pub(crate) const STREAM_INFLATE_CHUNK: usize = 32 * 1024;
710
711pub(crate) fn inflate_entry_from_stream<R>(
712    stream: &mut PackReadStream<'_, R>,
713    crc: &mut crc32fast::Hasher,
714    size_hint: usize,
715) -> Result<(Vec<u8>, usize)>
716where
717    R: Read,
718{
719    INFLATE.with(|cell| {
720        let mut decompress = cell.borrow_mut();
721        decompress.reset(true);
722        let mut out = Vec::with_capacity(inflate::bounded_inflate_reserve(
723            size_hint,
724            STREAM_INFLATE_CHUNK,
725        ));
726        let mut compressed_total = 0usize;
727        let mut input = [0u8; STREAM_INFLATE_CHUNK];
728        loop {
729            let read = stream.read_compressed_chunk(&mut input)?;
730            if read == 0 {
731                return Err(GitError::InvalidObject("truncated zlib stream".into()));
732            }
733            let mut cursor = 0usize;
734            while cursor < read {
735                if out.len() == out.capacity() {
736                    out.reserve(out.len().max(64));
737                }
738                let before_in = decompress.total_in();
739                let before_out = decompress.total_out();
740                let status = decompress
741                    .decompress_vec(
742                        &input[cursor..read],
743                        &mut out,
744                        flate2::FlushDecompress::None,
745                    )
746                    .map_err(|err| {
747                        GitError::InvalidObject(format!("zlib inflate failed: {err}"))
748                    })?;
749                let consumed = (decompress.total_in() - before_in) as usize;
750                let produced = decompress.total_out() - before_out;
751                if consumed > 0 {
752                    let consumed_end = cursor + consumed;
753                    stream.accept_compressed_bytes(&input[cursor..consumed_end], crc);
754                    compressed_total = compressed_total
755                        .checked_add(consumed)
756                        .ok_or_else(|| GitError::InvalidFormat("pack offset overflow".into()))?;
757                    cursor = consumed_end;
758                }
759                match status {
760                    flate2::Status::StreamEnd => {
761                        stream.push_back_compressed_bytes(&input[cursor..read])?;
762                        return Ok((out, compressed_total));
763                    }
764                    _ if consumed == 0 && produced == 0 => {
765                        return Err(GitError::InvalidObject("truncated zlib stream".into()));
766                    }
767                    _ => {}
768                }
769            }
770        }
771    })
772}
773
774pub(crate) fn parse_entry_header_from_stream<R>(
775    stream: &mut PackReadStream<'_, R>,
776    crc: &mut crc32fast::Hasher,
777) -> Result<EntryHeader>
778where
779    R: Read,
780{
781    let first = stream.read_entry_byte(crc)?;
782    let mut size = u64::from(first & 0x0f);
783    let kind = match (first >> 4) & 0x07 {
784        1 => PackObjectKind::Commit,
785        2 => PackObjectKind::Tree,
786        3 => PackObjectKind::Blob,
787        4 => PackObjectKind::Tag,
788        6 => PackObjectKind::OfsDelta,
789        7 => PackObjectKind::RefDelta,
790        other => {
791            return Err(GitError::InvalidFormat(format!(
792                "invalid pack object type {other}"
793            )));
794        }
795    };
796    let mut shift = 4;
797    let mut byte = first;
798    while byte & 0x80 != 0 {
799        byte = stream.read_entry_byte(crc)?;
800        let part = u64::from(byte & 0x7f);
801        size = size
802            .checked_add(
803                part.checked_shl(shift)
804                    .ok_or_else(|| GitError::InvalidFormat("pack size overflow".into()))?,
805            )
806            .ok_or_else(|| GitError::InvalidFormat("pack size overflow".into()))?;
807        shift += 7;
808    }
809    Ok(EntryHeader { kind, size })
810}
811
812pub(crate) fn parse_ofs_delta_base_offset_from_stream<R>(
813    stream: &mut PackReadStream<'_, R>,
814    crc: &mut crc32fast::Hasher,
815    entry_offset: u64,
816) -> Result<u64>
817where
818    R: Read,
819{
820    let mut byte = stream.read_entry_byte(crc)?;
821    let mut relative = u64::from(byte & 0x7f);
822    while byte & 0x80 != 0 {
823        byte = stream.read_entry_byte(crc)?;
824        relative = relative
825            .checked_add(1)
826            .and_then(|value| value.checked_shl(7))
827            .and_then(|value| value.checked_add(u64::from(byte & 0x7f)))
828            .ok_or_else(|| GitError::InvalidFormat("ofs-delta offset overflow".into()))?;
829    }
830    entry_offset
831        .checked_sub(relative)
832        .ok_or_else(|| GitError::InvalidFormat("ofs-delta points before pack start".into()))
833}