Skip to main content

heddle_pack/store/pack/
pack_reader.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Pack reader for extracting objects from packfiles.
3
4#[cfg(test)]
5use std::sync::atomic::{AtomicUsize, Ordering};
6use std::{
7    collections::{HashMap, HashSet},
8    fs::File,
9    io::Read,
10    path::{Path, PathBuf},
11    sync::OnceLock,
12};
13
14use bytes::Bytes;
15use heddle_format::delta::{DeltaDecoder, MAX_DELTA_OUTPUT_SIZE};
16
17use super::{
18    ObjectType, PackLogicalId, PackObjectId, PackObjectRecord, PackRepresentationHash,
19    append_container_checksum, decode_tagged_entry_header, decompress_pack_payload, has_zstd_magic,
20    pack_container_spec, pack_identity::LogicalIdBuilder, pack_index::PackIndex, varint,
21    verify_supported_container, verify_supported_container_layout, write_container_header,
22};
23use crate::{
24    object::ContentHash,
25    store::{Result, StoreError},
26};
27
28const MAX_PACK_DELTA_OUTPUT_SIZE: usize = MAX_DELTA_OUTPUT_SIZE;
29const MAX_DELTA_CHAIN_DEPTH: usize = 50;
30const MMAP_THRESHOLD_BYTES: u64 = 256 * 1024;
31
32type DecodedCompactObject = (PackObjectId, ObjectType, Vec<u8>);
33type DecodedCompactObjects = Vec<DecodedCompactObject>;
34
35/// Physical read tier for an indexed pack object.
36///
37/// Hot records have one independently addressable record per logical object.
38/// Solid-frame records share one payload across several logical objects and
39/// therefore require a frame decompression on a cache-cold read.
40#[derive(Debug, Clone, Copy, Eq, PartialEq)]
41pub enum PackReadTier {
42    /// Independently addressable, random-access record.
43    Hot,
44    /// Compact payload shared by several tree or state ids.
45    SolidFrame,
46}
47
48fn read_file_bytes_for_pack(path: &Path) -> Result<Bytes> {
49    let file = File::open(path)?;
50    let len = file.metadata()?.len();
51    if len == 0 {
52        return Ok(Bytes::new());
53    }
54    if len >= MMAP_THRESHOLD_BYTES {
55        let mmap = unsafe { memmap2::MmapOptions::new().map(&file)? };
56        if mmap.len() != checked_file_len_to_usize(len)? {
57            return Err(StoreError::InvalidObject(
58                "pack file size changed during memory mapping".to_string(),
59            ));
60        }
61        return Ok(Bytes::from_owner(mmap));
62    }
63    let mut data = Vec::with_capacity(checked_file_len_to_usize(len)?);
64    let mut reader = file;
65    reader.read_to_end(&mut data)?;
66    Ok(Bytes::from(data))
67}
68
69fn checked_file_len_to_usize(len: u64) -> Result<usize> {
70    usize::try_from(len).map_err(|_| {
71        StoreError::InvalidObject(format!("file length {len} exceeds platform limits"))
72    })
73}
74
75/// Pack reader for extracting objects.
76///
77/// `data` is a refcounted [`Bytes`] view of the pack file. For
78/// uncompressed entries we hand back a zero-copy `Bytes::slice` into
79/// this buffer — no per-blob memcpy, no per-blob allocation. Mmap-
80/// backed `Bytes` (via [`Bytes::from_owner`] on the
81/// `memmap2::Mmap`) survives across reads without copying the
82/// whole pack into the heap.
83enum PackData<'a> {
84    Borrowed(&'a [u8]),
85    Owned(Bytes),
86}
87
88impl<'a> PackData<'a> {
89    fn as_slice(&self) -> &[u8] {
90        match self {
91            Self::Borrowed(data) => data,
92            Self::Owned(data) => data,
93        }
94    }
95
96    fn slice(&self, range: std::ops::Range<usize>) -> Bytes {
97        match self {
98            Self::Borrowed(data) => Bytes::copy_from_slice(&data[range]),
99            Self::Owned(data) => data.slice(range),
100        }
101    }
102}
103
104pub struct PackReader<'a> {
105    data: PackData<'a>,
106    index: PackIndex,
107    offsets: OnceLock<super::offset_index::OffsetIndex>,
108    scratch_root: Option<PathBuf>,
109    content_end: usize,
110    #[cfg(test)]
111    compact_frame_reads: AtomicUsize,
112}
113
114#[derive(Debug, Clone)]
115pub struct EncodedPackSubset {
116    pub pack_data: Vec<u8>,
117    pub index_data: Vec<u8>,
118    pub encoded_bytes_copied: u64,
119}
120
121impl PackReader<'static> {
122    /// Open a pack file. mmap-backed when the pack is large enough
123    /// to benefit (the same threshold the loose-blob path uses for
124    /// its own mmap decision); read-into-heap otherwise.
125    pub fn open(pack_path: &Path, index_path: &Path, scratch_root: &Path) -> Result<Self> {
126        Self::open_with_verification(pack_path, index_path, scratch_root, true)
127    }
128
129    pub(super) fn open_lazy(
130        pack_path: &Path,
131        index_path: &Path,
132        scratch_root: &Path,
133    ) -> Result<Self> {
134        Self::open_with_verification(pack_path, index_path, scratch_root, false)
135    }
136
137    fn open_with_verification(
138        pack_path: &Path,
139        index_path: &Path,
140        scratch_root: &Path,
141        verify_checksum: bool,
142    ) -> Result<Self> {
143        let pack_bytes = read_file_bytes_for_pack(pack_path)?;
144        let index_data = read_file_bytes_for_pack(index_path)?;
145        let (_, _, content_end) = if verify_checksum {
146            verify_supported_container(&pack_bytes)?
147        } else {
148            verify_supported_container_layout(&pack_bytes)?
149        };
150        let index = PackIndex::from_owned_bytes(index_data)?;
151        let scratch_root = Some(scratch_root.to_path_buf());
152        let offsets = OnceLock::new();
153        Ok(Self {
154            data: PackData::Owned(pack_bytes),
155            index,
156            offsets,
157            scratch_root,
158            content_end,
159            #[cfg(test)]
160            compact_frame_reads: AtomicUsize::new(0),
161        })
162    }
163
164    pub fn from_bytes(
165        pack_data: impl Into<Bytes>,
166        index_data: impl AsRef<[u8]>,
167        scratch_root: &Path,
168    ) -> Result<Self> {
169        let pack_data = pack_data.into();
170        let (_, _, content_end) = verify_supported_container(&pack_data)?;
171        let index = PackIndex::from_bytes(index_data.as_ref())?;
172        let scratch_root = Some(scratch_root.to_path_buf());
173        let offsets = OnceLock::new();
174        Ok(Self {
175            data: PackData::Owned(pack_data),
176            index,
177            offsets,
178            scratch_root,
179            content_end,
180            #[cfg(test)]
181            compact_frame_reads: AtomicUsize::new(0),
182        })
183    }
184}
185
186impl<'a> PackReader<'a> {
187    /// Buffered stores without a filesystem keep the physical index in memory.
188    /// Source-transfer validation requires a constructor with a scratch root.
189    pub fn from_slice_in_memory(pack_data: &'a [u8], index_data: impl AsRef<[u8]>) -> Result<Self> {
190        let mut reader = Self::from_slice(pack_data, index_data, Path::new(""))?;
191        reader.scratch_root = None;
192        Ok(reader)
193    }
194
195    pub fn from_slice(
196        pack_data: &'a [u8],
197        index_data: impl AsRef<[u8]>,
198        scratch_root: &Path,
199    ) -> Result<Self> {
200        let (_, _, content_end) = verify_supported_container(pack_data)?;
201        let index = PackIndex::from_bytes(index_data.as_ref())?;
202        let scratch_root = Some(scratch_root.to_path_buf());
203        let offsets = OnceLock::new();
204        Ok(Self {
205            data: PackData::Borrowed(pack_data),
206            index,
207            offsets,
208            scratch_root,
209            content_end,
210            #[cfg(test)]
211            compact_frame_reads: AtomicUsize::new(0),
212        })
213    }
214
215    // Metadata point reads need only the mapped identity index. Build the
216    // physical order lazily for all-object visits and shared blob-frame probes.
217    fn offset_index(&self) -> Result<&super::offset_index::OffsetIndex> {
218        if let Some(index) = self.offsets.get() {
219            return Ok(index);
220        }
221        let index =
222            super::offset_index::OffsetIndex::new(&self.index, self.scratch_root.as_deref())?;
223        // A concurrent reader may win initialization. Dropping our unused
224        // candidate closes its mapping and removes its temporary file.
225        let _ = self.offsets.set(index);
226        self.offsets
227            .get()
228            .ok_or_else(|| StoreError::InvalidObject("physical pack index unavailable".into()))
229    }
230
231    /// List all object ids in this pack.
232    pub fn list_ids(&self) -> Result<Vec<PackObjectId>> {
233        self.index.ids()
234    }
235
236    /// Verify the exact selected closure under a decoded byte budget.
237    #[cfg(feature = "source-transfer")]
238    pub fn validate_source_closure(
239        &self,
240        selected: &crate::object::State,
241        max_decoded_bytes: u64,
242    ) -> Result<()> {
243        self.validate_source_closure_with_metadata(selected, &[], None, max_decoded_bytes)
244    }
245    #[cfg(feature = "source-transfer")]
246    pub fn validate_source_closure_with_metadata(
247        &self,
248        selected: &crate::object::State,
249        references: &[crate::object::source_target::capture::ReferenceProof],
250        visibility: Option<&crate::object::thread_replication::CaptureVisibility>,
251        max_decoded_bytes: u64,
252    ) -> Result<()> {
253        self.validate_source_layout(max_decoded_bytes)?;
254        super::source_pack::validate(self, selected, max_decoded_bytes, references, visibility)
255    }
256    /// Preserve verified partial proofs without retaining decoded trees.
257    #[cfg(feature = "source-transfer")]
258    pub fn validate_visible_source_closure(
259        &self,
260        selected: &crate::object::State,
261        max_decoded_bytes: u64,
262    ) -> Result<super::VisibleSourceClosure> {
263        self.validate_source_layout(max_decoded_bytes)?;
264        super::source_pack::validate_disclosure(self, selected, max_decoded_bytes, &[], None, true)
265    }
266    #[cfg(feature = "source-transfer")]
267    pub(super) fn scratch_root(&self) -> Result<&Path> {
268        self.scratch_root.as_deref().ok_or_else(|| {
269            StoreError::InvalidObject("source validation requires a scratch root".into())
270        })
271    }
272    pub fn object_count(&self) -> usize {
273        self.index.len()
274    }
275    // Point reads support generic packs with repeated identities. Whole-pack
276    // visits need one offset per logical ID so a point lookup cannot substitute
277    // another record's payload. Sorted IDs make this a scalar membership check.
278    fn validate_unique_ids(&self) -> Result<()> {
279        let mut previous = None;
280        for entry in self.index.iter() {
281            let id = entry?.id;
282            if previous == Some(id) {
283                return Err(StoreError::InvalidObject(
284                    "duplicate pack object identity".into(),
285                ));
286            }
287            previous = Some(id);
288        }
289        Ok(())
290    }
291    #[cfg(feature = "source-transfer")]
292    fn validate_source_layout(&self, max_decoded_bytes: u64) -> Result<()> {
293        if self.index.is_empty() {
294            return Err(StoreError::InvalidObject("source pack is empty".into()));
295        }
296        self.validate_unique_ids()?;
297        let (_, mut next, end) = verify_supported_container_layout(self.data.as_slice())?;
298        let mut decoded = 0_u64;
299        self.offset_index()?.visit(|offset, _, _| {
300            if checked_index_offset(offset)? != next {
301                return Err(StoreError::InvalidObject(
302                    "source pack has unindexed or overlapping records".into(),
303                ));
304            }
305            let header = decode_tagged_entry_header(self.content_from(next)?)?;
306            decoded = decoded
307                .checked_add(header.uncompressed_size as u64)
308                .ok_or_else(|| StoreError::InvalidObject("source pack size overflow".into()))?;
309            if decoded > max_decoded_bytes {
310                return Err(StoreError::InvalidObject(
311                    "source pack decoded byte budget exceeded".into(),
312                ));
313            }
314            next = next
315                .checked_add(header.header_len)
316                .and_then(|n| n.checked_add(header.compressed_size))
317                .ok_or_else(|| {
318                    StoreError::InvalidObject("source pack record length overflow".into())
319                })?;
320            if next > end {
321                return Err(StoreError::InvalidObject(
322                    "source pack record exceeds container".into(),
323                ));
324            }
325            Ok(())
326        })?;
327        if next != end {
328            return Err(StoreError::InvalidObject(
329                "source pack has unindexed trailing records".into(),
330            ));
331        }
332        Ok(())
333    }
334
335    /// Compute this pack's root-spool-scoped logical identity.
336    ///
337    /// Every logical object is decoded so delta and compact-frame physical
338    /// choices cannot affect the result. The visit also validates exact index
339    /// membership before an identity is returned.
340    pub fn logical_id(&self) -> Result<PackLogicalId> {
341        let mut identity = LogicalIdBuilder::new();
342        self.visit_objects(|id, object_type, data| {
343            identity.push(id, object_type, data);
344            Ok(())
345        })?;
346        Ok(identity.finish())
347    }
348
349    /// Hash the exact finalized pack bytes used by this reader.
350    pub fn representation_hash(&self) -> PackRepresentationHash {
351        PackRepresentationHash::compute(self.data.as_slice())
352    }
353
354    /// List logical ids together with their physical read tier.
355    ///
356    /// A shared compact frame is represented by several index aliases at one
357    /// record offset. Direct records have a unique offset and form the hot,
358    /// random-access tier. A one-object compact frame is intentionally treated
359    /// as hot here: it has no read amplification over a direct record.
360    pub(super) fn indexed_read_tiers(&self) -> Result<Vec<(PackObjectId, PackReadTier)>> {
361        let entries = self.index.entries()?;
362        let mut aliases = HashMap::<u64, usize>::with_capacity(entries.len());
363        for entry in &entries {
364            *aliases.entry(entry.offset).or_default() += 1;
365        }
366        Ok(entries
367            .into_iter()
368            .map(|entry| {
369                let tier = if aliases[&entry.offset] > 1 {
370                    PackReadTier::SolidFrame
371                } else {
372                    PackReadTier::Hot
373                };
374                (entry.id, tier)
375            })
376            .collect())
377    }
378
379    /// Point-membership probe backed by the sorted pack index. This avoids
380    /// enumerating every object merely to locate one hot-path tree or state.
381    pub(super) fn contains_object(&self, id: &PackObjectId) -> Result<bool> {
382        Ok(self.index.find(id)?.is_some())
383    }
384
385    #[cfg(test)]
386    pub(super) fn compact_frame_read_count(&self) -> usize {
387        self.compact_frame_reads.load(Ordering::Relaxed)
388    }
389
390    #[cfg(test)]
391    fn record_compact_frame_read(&self) {
392        self.compact_frame_reads.fetch_add(1, Ordering::Relaxed);
393    }
394
395    pub fn list_hashes(&self) -> Result<Vec<ContentHash>> {
396        Ok(self
397            .list_ids()?
398            .into_iter()
399            .filter_map(|id| match id {
400                PackObjectId::Hash(hash) => Some(hash),
401                PackObjectId::StateId(_) | PackObjectId::AnnotatedTag(_) => None,
402            })
403            .collect())
404    }
405
406    pub fn has_object(&self, id: &PackObjectId) -> Result<bool> {
407        Ok(self.index.find(id)?.is_some())
408    }
409
410    /// Compressed payload bytes used by unique physical records of `obj_type`.
411    ///
412    /// Shared compact frames have one record offset indexed by many logical
413    /// ids, so offsets are deduplicated before bytes are counted.
414    pub fn encoded_payload_bytes(&self, obj_type: ObjectType) -> Result<u64> {
415        let mut bytes = 0u64;
416        self.offset_index()?.visit(|offset, _, _| {
417            let header =
418                decode_tagged_entry_header(self.content_from(checked_index_offset(offset)?)?)?;
419            if header.obj_type == obj_type {
420                bytes = bytes.saturating_add(header.compressed_size as u64);
421            }
422            Ok(())
423        })?;
424        Ok(bytes)
425    }
426
427    /// Visit every indexed physical record once, checking exact frame membership.
428    pub fn visit_objects(
429        &self,
430        mut visitor: impl FnMut(PackObjectId, ObjectType, &[u8]) -> Result<()>,
431    ) -> Result<()> {
432        self.validate_unique_ids()?;
433        self.offset_index()?.visit(|offset, ordinal, aliases| {
434            let offset = checked_index_offset(offset)?;
435            if let Some(objects) = self.read_compact_objects_at(offset)? {
436                if objects.len() != aliases {
437                    return Err(StoreError::InvalidObject(
438                        "compact frame object set differs from its index".into(),
439                    ));
440                }
441                let mut unique = HashSet::with_capacity(objects.len());
442                for (id, _, _) in &objects {
443                    if !unique.insert(*id) || self.index.find(id)? != Some(offset as u64) {
444                        return Err(StoreError::InvalidObject(
445                            "compact frame object set differs from its index".into(),
446                        ));
447                    }
448                }
449                for (id, kind, data) in objects {
450                    visitor(id, kind, &data)?;
451                }
452                return Ok(());
453            }
454            if aliases != 1 {
455                return Err(StoreError::InvalidObject(
456                    "ordinary pack record is indexed by multiple object ids".into(),
457                ));
458            }
459            let id = self.index.entry(ordinal)?.id;
460            let (kind, data) = self
461                .get_object(&id)?
462                .ok_or_else(|| StoreError::InvalidObject("indexed object is missing".into()))?;
463            visitor(id, kind, &data)
464        })
465    }
466
467    /// Copy a validated subset of non-delta encoded entries into a standalone
468    /// hosted transport pack without decoding or recompressing their bodies.
469    ///
470    /// `Ok(None)` is a safe fallback signal: an expected object is absent,
471    /// duplicated, has a different type/size, is delta encoded, or names a
472    /// repository-local entry that may not cross the hosted pack boundary.
473    pub fn copy_hosted_encoded_subset(
474        &self,
475        expected: &[(PackObjectId, ObjectType, u64)],
476    ) -> Result<Option<EncodedPackSubset>> {
477        if expected.is_empty() {
478            return Ok(None);
479        }
480        let mut unique = HashSet::with_capacity(expected.len());
481        if expected.iter().any(|(id, obj_type, _)| {
482            !unique.insert(*id)
483                || matches!(
484                    obj_type,
485                    ObjectType::Delta | ObjectType::StateAttachment | ObjectType::SnapshotCommit
486                )
487        }) {
488            return Ok(None);
489        }
490
491        let mut pack_data = Vec::new();
492        write_container_header(&mut pack_data, pack_container_spec(), expected.len() as u64);
493        let mut index = PackIndex::new();
494        let mut encoded_bytes_copied = 0u64;
495        for (expected_id, expected_type, expected_size) in expected {
496            let Some(offset) = self.index.find(expected_id)? else {
497                return Ok(None);
498            };
499            let offset = checked_index_offset(offset)?;
500            if offset >= self.content_end {
501                return Err(StoreError::InvalidObject(
502                    "Entry offset out of bounds".to_string(),
503                ));
504            }
505            let header = decode_tagged_entry_header(self.content_from(offset)?)?;
506            if self.read_compact_objects_at(offset)?.is_some() {
507                return Ok(None);
508            }
509            let expected_size = usize::try_from(*expected_size).ok();
510            if header.id != *expected_id
511                || header.obj_type != *expected_type
512                || Some(header.uncompressed_size) != expected_size
513                || matches!(
514                    header.obj_type,
515                    ObjectType::Delta | ObjectType::StateAttachment | ObjectType::SnapshotCommit
516                )
517            {
518                return Ok(None);
519            }
520            let encoded_len = header
521                .header_len
522                .checked_add(header.compressed_size)
523                .ok_or_else(|| {
524                    StoreError::InvalidObject("pack entry length overflow".to_string())
525                })?;
526            let encoded_end = offset
527                .checked_add(encoded_len)
528                .ok_or_else(|| StoreError::InvalidObject("pack entry end overflow".to_string()))?;
529            if encoded_end > self.content_end {
530                return Err(StoreError::InvalidObject(
531                    "pack entry extends beyond content boundary".to_string(),
532                ));
533            }
534            let output_offset = u64::try_from(pack_data.len()).map_err(|_| {
535                StoreError::InvalidObject("reused pack offset exceeds u64".to_string())
536            })?;
537            index.add(*expected_id, output_offset);
538            pack_data.extend_from_slice(&self.data.as_slice()[offset..encoded_end]);
539            encoded_bytes_copied = encoded_bytes_copied
540                .checked_add(u64::try_from(encoded_len).map_err(|_| {
541                    StoreError::InvalidObject("encoded pack entry length exceeds u64".to_string())
542                })?)
543                .ok_or_else(|| {
544                    StoreError::InvalidObject("encoded reused byte count overflow".to_string())
545                })?;
546        }
547        index.sort();
548        append_container_checksum(&mut pack_data);
549        Ok(Some(EncodedPackSubset {
550            pack_data,
551            index_data: index.to_bytes(),
552            encoded_bytes_copied,
553        }))
554    }
555
556    /// Get an object from the pack.
557    ///
558    /// Verifies that the tagged id at the indexed offset matches
559    /// `id` before returning. A stale `.idx` file (e.g., overwritten
560    /// in place after a pack rebuild) can otherwise route a request
561    /// for hash `A` to a record physically located at hash `B`'s
562    /// offset — same shape, different content, no error signal.
563    /// This cheap 32-byte id comparison catches that without paying
564    /// a full content-hash recompute on every read; corruption
565    /// strictly *inside* the record body is a separate failure mode
566    /// surfaced via the consumer-side hash verify (see
567    /// `FsStore::loose_blob_path` for the blob equivalent).
568    pub fn get_object(&self, id: &PackObjectId) -> Result<Option<(ObjectType, Vec<u8>)>> {
569        let offset = match self.index.find(id)? {
570            Some(offset) => checked_index_offset(offset)?,
571            None => return Ok(None),
572        };
573
574        let record = self.read_record_at_depth(id, offset, 0)?;
575        Ok(Some((record.obj_type, record.data)))
576    }
577
578    pub fn get_hashed_object(&self, hash: &ContentHash) -> Result<Option<(ObjectType, Vec<u8>)>> {
579        self.get_object(&PackObjectId::Hash(*hash))
580    }
581
582    /// Read an object's logical type from pack headers without reading or
583    /// decoding its payload.
584    ///
585    /// Delta entries inherit the type of their base, so this follows only the
586    /// tagged base ids and headers until it reaches a non-delta entry. No
587    /// compressed or delta payload bytes are decoded. Missing hashes return
588    /// `Ok(None)`.
589    pub fn get_hashed_object_type(&self, hash: &ContentHash) -> Result<Option<ObjectType>> {
590        let id = PackObjectId::Hash(*hash);
591        let Some(offset) = self.index.find(&id)? else {
592            return Ok(None);
593        };
594        self.read_object_type_at_depth(&id, checked_index_offset(offset)?, 0)
595            .map(Some)
596    }
597
598    /// Zero-copy fast path: when the entry is non-delta and stored
599    /// uncompressed, returns `Bytes::slice` into the pack's
600    /// (mmap-backed) buffer — no allocation, no memcpy. Compressed
601    /// or delta entries fall back to `get_object` and wrap the
602    /// resulting `Vec<u8>` in a `Bytes` (one Arc, no body copy).
603    ///
604    /// Use this from the hot read path. The 10 MB benchmark gap
605    /// between the mount and vanilla FS at the 1 MB+ tier is the
606    /// per-blob memcpy this method eliminates.
607    pub fn get_object_bytes(&self, id: &PackObjectId) -> Result<Option<(ObjectType, Bytes)>> {
608        let Some(offset) = self.index.find(id)? else {
609            return Ok(None);
610        };
611        let offset = checked_index_offset(offset)?;
612        if offset >= self.content_end {
613            return Err(StoreError::InvalidObject(
614                "Entry offset out of bounds".to_string(),
615            ));
616        }
617
618        // Verify the tagged id at the indexed offset matches the
619        // requested id — guards against stale-index misrouting (see
620        // `get_object` for the long-form rationale). 32-byte
621        // compare; cheaper than the size+varint decode that follows.
622        let (record_id, id_len) = PackObjectId::decode_tagged(self.content_from(offset)?)?;
623        let header_start = checked_index_add(offset, id_len, "record header start")?;
624        let (encoded_type, uncompressed_size, type_len) =
625            varint::decode_type_and_size(self.content_from(header_start)?).ok_or_else(|| {
626                StoreError::InvalidObject("Truncated type+size varint".to_string())
627            })?;
628        let obj_type = decoded_entry_type(record_id, encoded_type)?;
629        let uncompressed_size = checked_decoded_size("uncompressed_size", uncompressed_size)?;
630        let varint_start = checked_index_add(header_start, type_len, "compressed_size start")?;
631        let (compressed_size, comp_len) = varint::decode_varint(self.content_from(varint_start)?)
632            .ok_or_else(truncated_compressed_size_varint)?;
633        let compressed_size = checked_decoded_size("compressed_size", compressed_size)?;
634
635        // Fast path: non-delta entry stored uncompressed. The most
636        // common shape for snapshot-time packs (the builder skips
637        // the delta search for unrelated blobs).
638        if record_id == *id && obj_type != ObjectType::Delta && compressed_size == uncompressed_size
639        {
640            let data_start = checked_index_add(varint_start, comp_len, "entry data start")?;
641            let data_end = checked_data_end(data_start, compressed_size, self.content_end)?;
642            let data = &self.data.as_slice()[data_start..data_end];
643            if !is_compact_frame(data) {
644                return Ok(Some((obj_type, self.data.slice(data_start..data_end))));
645            }
646        }
647
648        // Slow path: defer to the full record reader (it handles
649        // decompression + delta chains) and Bytes-wrap the Vec.
650        // Bytes::from(Vec) is a single Arc allocation, no body copy.
651        let record = self.read_record_at_depth(id, offset, 0)?;
652        Ok(Some((record.obj_type, Bytes::from(record.data))))
653    }
654
655    pub fn get_hashed_object_bytes(
656        &self,
657        hash: &ContentHash,
658    ) -> Result<Option<(ObjectType, Bytes)>> {
659        self.get_object_bytes(&PackObjectId::Hash(*hash))
660    }
661
662    /// Read just the type+size header for an object without
663    /// decompressing its payload. Returns `Ok(None)` when the object
664    /// isn't in this pack.
665    ///
666    /// For non-delta entries this is one varint decode at the indexed
667    /// offset — much cheaper than `get_object`. Delta entries fall
668    /// back to a full read because their *resolved* size requires
669    /// chasing the base; in practice deltas are rare in the directory
670    /// listing hot path so the fallback is acceptable.
671    pub fn get_hashed_object_size(&self, hash: &ContentHash) -> Result<Option<u64>> {
672        let id = PackObjectId::Hash(*hash);
673        let Some(offset) = self.index.find(&id)? else {
674            return Ok(None);
675        };
676        let offset = checked_index_offset(offset)?;
677        if offset >= self.content_end {
678            return Err(StoreError::InvalidObject(
679                "Entry offset out of bounds".to_string(),
680            ));
681        }
682        let (record_id, id_len) = PackObjectId::decode_tagged(self.content_from(offset)?)?;
683        let header_start = checked_index_add(offset, id_len, "record header start")?;
684        let (obj_type, uncompressed_size, _type_len) = super::varint::decode_type_and_size(
685            self.content_from(header_start)?,
686        )
687        .ok_or_else(|| StoreError::InvalidObject("Truncated type+size varint".to_string()))?;
688        if (obj_type == ObjectType::Blob && self.offset_index()?.aliases(offset as u64)?)
689            || matches!(obj_type, ObjectType::Tree | ObjectType::State)
690        {
691            let Some((_, data)) = self.get_object(&id)? else {
692                return Ok(None);
693            };
694            return Ok(Some(data.len() as u64));
695        }
696        verify_record_id_matches(&id, &record_id)?;
697        if obj_type == ObjectType::Delta {
698            // Delta entries record the *resolved* output size in the
699            // type+size varint already (see `read_record_at_depth`'s
700            // size-mismatch check), so we can still return without
701            // decompressing the payload.
702            return Ok(Some(uncompressed_size));
703        }
704        Ok(Some(uncompressed_size))
705    }
706
707    fn read_object_type_at_depth(
708        &self,
709        requested_id: &PackObjectId,
710        offset: usize,
711        depth: usize,
712    ) -> Result<ObjectType> {
713        if depth > MAX_DELTA_CHAIN_DEPTH {
714            return Err(StoreError::InvalidObject(format!(
715                "Delta chain depth {depth} exceeds max {MAX_DELTA_CHAIN_DEPTH}"
716            )));
717        }
718        if offset >= self.content_end {
719            return Err(StoreError::InvalidObject(
720                "Entry offset out of bounds".to_string(),
721            ));
722        }
723
724        let header = decode_tagged_entry_header(self.content_from(offset)?)?;
725        if header.id != *requested_id {
726            return self
727                .read_record_at_depth(requested_id, offset, depth)
728                .map(|record| record.obj_type);
729        }
730        if header.obj_type != ObjectType::Delta {
731            return Ok(header.obj_type);
732        }
733
734        let base_hash = Self::require_delta_base_hash(header.delta_base)?;
735        let base_id = PackObjectId::Hash(base_hash);
736        let base_offset = self
737            .index
738            .find(&base_id)?
739            .ok_or_else(|| StoreError::NotFound(base_hash.to_string()))?;
740        self.read_object_type_at_depth(&base_id, checked_index_offset(base_offset)?, depth + 1)
741    }
742
743    fn read_record_at_depth(
744        &self,
745        requested_id: &PackObjectId,
746        offset: usize,
747        depth: usize,
748    ) -> Result<PackObjectRecord> {
749        if offset >= self.content_end {
750            return Err(StoreError::InvalidObject(
751                "Entry offset out of bounds".to_string(),
752            ));
753        }
754
755        let (id, id_len) = PackObjectId::decode_tagged(self.content_from(offset)?)?;
756        let header_start = checked_index_add(offset, id_len, "record header start")?;
757
758        let (encoded_type, uncompressed_size, type_len) =
759            varint::decode_type_and_size(self.content_from(header_start)?).ok_or_else(|| {
760                StoreError::InvalidObject("Truncated type+size varint".to_string())
761            })?;
762        let obj_type = decoded_entry_type(id, encoded_type)?;
763        let uncompressed_size = checked_decoded_size("uncompressed_size", uncompressed_size)?;
764
765        let varint_start = checked_index_add(header_start, type_len, "compressed_size start")?;
766        let (compressed_size, comp_len) = varint::decode_varint(self.content_from(varint_start)?)
767            .ok_or_else(truncated_compressed_size_varint)?;
768        let compressed_size = checked_decoded_size("compressed_size", compressed_size)?;
769
770        let mut data_start = checked_index_add(varint_start, comp_len, "entry data start")?;
771
772        // Delta entries carry a tagged base id in pack v2.
773        let base_id = if obj_type == ObjectType::Delta {
774            let (base_id, base_len) = PackObjectId::decode_tagged(self.content_from(data_start)?)?;
775            data_start = checked_index_add(data_start, base_len, "delta data start")?;
776            Some(base_id)
777        } else {
778            None
779        };
780
781        let data_end = checked_data_end(data_start, compressed_size, self.content_end)?;
782
783        let stored_data = &self.data.as_slice()[data_start..data_end];
784
785        // Raw zstd (no wrapper). For non-delta entries, decompress
786        // if sizes differ. For delta entries, the stored data IS the delta
787        // payload (possibly zstd-compressed); check for zstd magic.
788        let decompressed = if obj_type == ObjectType::Delta {
789            if has_zstd_magic(stored_data) {
790                decompress_pack_payload(stored_data, 0)?
791            } else {
792                stored_data.to_vec()
793            }
794        } else if compressed_size != uncompressed_size {
795            decompress_pack_payload(stored_data, uncompressed_size)?
796        } else {
797            stored_data.to_vec()
798        };
799
800        let shared_blob =
801            obj_type == ObjectType::Blob && self.offset_index()?.aliases(offset as u64)?;
802        if obj_type != ObjectType::Delta && (shared_blob || is_compact_frame(&decompressed)) {
803            #[cfg(test)]
804            self.record_compact_frame_read();
805            if let Some(data) =
806                decode_compact_object(requested_id, obj_type, &decompressed, shared_blob)?
807            {
808                return Ok(PackObjectRecord {
809                    id: *requested_id,
810                    obj_type,
811                    data,
812                    delta_base: None,
813                    path_hint: None,
814                });
815            }
816        }
817        verify_record_id_matches(requested_id, &id)?;
818        let (resolved_type, final_data) = if obj_type == ObjectType::Delta {
819            self.read_delta_record(base_id, &decompressed, uncompressed_size, depth)?
820        } else {
821            (obj_type, decompressed)
822        };
823
824        if final_data.len() != uncompressed_size {
825            return Err(StoreError::InvalidObject(format!(
826                "Size mismatch: expected {}, got {}",
827                uncompressed_size,
828                final_data.len()
829            )));
830        }
831
832        Ok(PackObjectRecord {
833            id,
834            obj_type: resolved_type,
835            data: final_data,
836            delta_base: None,
837            path_hint: None,
838        })
839    }
840
841    fn read_compact_objects_at(&self, offset: usize) -> Result<Option<DecodedCompactObjects>> {
842        if offset >= self.content_end {
843            return Err(StoreError::InvalidObject(
844                "Entry offset out of bounds".to_string(),
845            ));
846        }
847        let header = decode_tagged_entry_header(self.content_from(offset)?)?;
848        if !matches!(
849            header.obj_type,
850            ObjectType::Blob | ObjectType::Tree | ObjectType::State
851        ) {
852            return Ok(None);
853        }
854        let shared_blob =
855            header.obj_type == ObjectType::Blob && self.offset_index()?.aliases(offset as u64)?;
856        if header.obj_type == ObjectType::Blob && !shared_blob {
857            return Ok(None);
858        }
859        let data_start = checked_index_add(offset, header.header_len, "entry data start")?;
860        let data_end = checked_data_end(data_start, header.compressed_size, self.content_end)?;
861        let stored = &self.data.as_slice()[data_start..data_end];
862        let data = if header.compressed_size != header.uncompressed_size {
863            decompress_pack_payload(stored, header.uncompressed_size)?
864        } else {
865            stored.to_vec()
866        };
867        if data.len() != header.uncompressed_size {
868            return Err(StoreError::InvalidObject(format!(
869                "Size mismatch: expected {}, got {}",
870                header.uncompressed_size,
871                data.len()
872            )));
873        }
874        #[cfg(test)]
875        if shared_blob || is_compact_frame(&data) {
876            self.record_compact_frame_read();
877        }
878        decode_compact_objects(header.obj_type, &data, shared_blob)
879    }
880
881    fn read_delta_record(
882        &self,
883        base_id: Option<PackObjectId>,
884        delta: &[u8],
885        uncompressed_size: usize,
886        depth: usize,
887    ) -> Result<(ObjectType, Vec<u8>)> {
888        if depth > MAX_DELTA_CHAIN_DEPTH {
889            return Err(StoreError::InvalidObject(format!(
890                "Delta chain depth {} exceeds max {}",
891                depth, MAX_DELTA_CHAIN_DEPTH
892            )));
893        }
894
895        if uncompressed_size > MAX_PACK_DELTA_OUTPUT_SIZE {
896            return Err(StoreError::InvalidObject(format!(
897                "Delta output size {} exceeds max {}",
898                uncompressed_size, MAX_PACK_DELTA_OUTPUT_SIZE
899            )));
900        }
901
902        let base_hash = Self::require_delta_base_hash(base_id)?;
903        let base_offset = self
904            .index
905            .find(&PackObjectId::Hash(base_hash))?
906            .ok_or_else(|| StoreError::NotFound(base_hash.to_string()))?;
907        let base_offset = checked_index_offset(base_offset)?;
908        let base_id = PackObjectId::Hash(base_hash);
909        let base_record = self.read_record_at_depth(&base_id, base_offset, depth + 1)?;
910        let base_type = base_record.obj_type;
911        let base_data = base_record.data;
912
913        let decoded = DeltaDecoder::decode(&base_data, delta, uncompressed_size)
914            .map_err(|error| StoreError::InvalidObject(format!("Delta decode failed: {error}")))?;
915
916        Ok((base_type, decoded))
917    }
918
919    fn require_delta_base_hash(base_id: Option<PackObjectId>) -> Result<ContentHash> {
920        match base_id {
921            Some(PackObjectId::Hash(hash)) => Ok(hash),
922            Some(PackObjectId::StateId(_) | PackObjectId::AnnotatedTag(_)) => Err(
923                StoreError::InvalidObject("pack delta base must be hash-backed content".into()),
924            ),
925            None => Err(StoreError::InvalidObject(
926                "pack object type is Delta but base hash is missing".into(),
927            )),
928        }
929    }
930
931    fn content_from(&self, offset: usize) -> Result<&[u8]> {
932        if offset > self.content_end {
933            return Err(StoreError::InvalidObject(
934                "Entry header out of bounds".to_string(),
935            ));
936        }
937        Ok(&self.data.as_slice()[offset..self.content_end])
938    }
939}
940
941fn checked_index_offset(offset: u64) -> Result<usize> {
942    usize::try_from(offset)
943        .map_err(|_| StoreError::InvalidObject("Entry offset exceeds platform limits".to_string()))
944}
945
946fn checked_decoded_size(field: &str, size: u64) -> Result<usize> {
947    let size = usize::try_from(size).map_err(|_| {
948        StoreError::InvalidObject(format!("Decoded {field} exceeds platform limits"))
949    })?;
950    if field == "uncompressed_size" && size > super::shared::MAX_PACK_OBJECT_OUTPUT_SIZE {
951        return Err(StoreError::InvalidObject(format!(
952            "Pack object output size {size} exceeds max {}",
953            super::shared::MAX_PACK_OBJECT_OUTPUT_SIZE
954        )));
955    }
956    Ok(size)
957}
958
959fn checked_index_add(start: usize, len: usize, field: &str) -> Result<usize> {
960    start.checked_add(len).ok_or_else(|| {
961        StoreError::InvalidObject(format!("{field} offset overflows platform limits"))
962    })
963}
964
965fn checked_data_end(
966    data_start: usize,
967    compressed_size: usize,
968    content_end: usize,
969) -> Result<usize> {
970    let data_end = data_start.checked_add(compressed_size).ok_or_else(|| {
971        StoreError::InvalidObject("Entry data range overflows platform limits".to_string())
972    })?;
973    if data_end > content_end {
974        return Err(StoreError::InvalidObject(
975            "Entry data out of bounds".to_string(),
976        ));
977    }
978    Ok(data_end)
979}
980
981fn truncated_compressed_size_varint() -> StoreError {
982    StoreError::InvalidObject("Truncated compressed_size varint".to_string())
983}
984
985fn decoded_entry_type(id: PackObjectId, encoded: ObjectType) -> Result<ObjectType> {
986    if matches!(id, PackObjectId::AnnotatedTag(_)) {
987        if encoded != ObjectType::Blob {
988            return Err(StoreError::InvalidObject(
989                "annotated-tag pack entry has invalid encoded type".to_string(),
990            ));
991        }
992        Ok(ObjectType::AnnotatedTag)
993    } else {
994        Ok(encoded)
995    }
996}
997
998/// Reject a record whose tagged id at the indexed offset doesn't
999/// match the id the caller asked for. The pack format stores its
1000/// records `[tagged_id, type+size, compressed_size, payload]` so the
1001/// tagged id is the cheapest available authenticator of "we landed
1002/// on the right record"; a stale or hand-edited `.idx` that points
1003/// at the *wrong* record produces a mismatch here and we surface it
1004/// as a real error instead of silently routing the caller to whatever
1005/// bytes happened to be at the bad offset.
1006fn verify_record_id_matches(requested: &PackObjectId, found: &PackObjectId) -> Result<()> {
1007    if requested == found {
1008        return Ok(());
1009    }
1010    Err(StoreError::InvalidObject(format!(
1011        "pack index routed lookup for {requested:?} to record tagged {found:?} \
1012         — index is stale or corrupt; the loose-store path will re-promote on \
1013         the next read"
1014    )))
1015}
1016
1017fn is_compact_frame(data: &[u8]) -> bool {
1018    heddle_object_model::compact::is_blob_frame(data)
1019        || heddle_object_model::compact::is_tree_frame(data)
1020        || heddle_object_model::compact::is_state_frame(data)
1021}
1022
1023fn decode_compact_object(
1024    requested_id: &PackObjectId,
1025    obj_type: ObjectType,
1026    data: &[u8],
1027    require_blob_frame: bool,
1028) -> Result<Option<Vec<u8>>> {
1029    match (obj_type, requested_id) {
1030        (ObjectType::Tree, PackObjectId::Hash(hash))
1031            if heddle_object_model::compact::is_tree_frame(data) =>
1032        {
1033            let tree = heddle_object_model::compact::extract_tree(data, *hash)
1034                .map_err(|error| compact_extract_error(requested_id, error))?;
1035            tree.encode_canonical()
1036                .map(Some)
1037                .map_err(|error| StoreError::InvalidObject(error.to_string()))
1038        }
1039        (ObjectType::State, PackObjectId::StateId(id))
1040            if heddle_object_model::compact::is_state_frame(data) =>
1041        {
1042            let state = heddle_object_model::compact::extract_state(data, *id)
1043                .map_err(|error| compact_extract_error(requested_id, error))?;
1044            state
1045                .encode_current_msgpack()
1046                .map(Some)
1047                .map_err(|error| StoreError::InvalidObject(error.to_string()))
1048        }
1049        _ => {
1050            let Some(objects) = decode_compact_objects(obj_type, data, require_blob_frame)? else {
1051                return Ok(None);
1052            };
1053            objects
1054                .into_iter()
1055                .find_map(|(id, _, bytes)| (id == *requested_id).then_some(bytes))
1056                .map(Some)
1057                .ok_or_else(|| compact_index_miss(requested_id))
1058        }
1059    }
1060}
1061
1062fn compact_extract_error(
1063    id: &PackObjectId,
1064    error: heddle_object_model::compact::CompactError,
1065) -> StoreError {
1066    if matches!(error, heddle_object_model::compact::CompactError::Missing) {
1067        compact_index_miss(id)
1068    } else {
1069        StoreError::InvalidObject(error.to_string())
1070    }
1071}
1072
1073fn decode_compact_objects(
1074    obj_type: ObjectType,
1075    data: &[u8],
1076    require_blob_frame: bool,
1077) -> Result<Option<DecodedCompactObjects>> {
1078    match obj_type {
1079        ObjectType::Blob if require_blob_frame => {
1080            heddle_object_model::compact::decode_blob_frame(data)
1081                .map_err(|error| StoreError::InvalidObject(error.to_string()))?
1082                .into_iter()
1083                .map(|(hash, body)| Ok((PackObjectId::Hash(hash), ObjectType::Blob, body.to_vec())))
1084                .collect::<Result<Vec<_>>>()
1085                .map(Some)
1086        }
1087        ObjectType::Blob => Ok(None),
1088        ObjectType::Tree if heddle_object_model::compact::is_tree_frame(data) => {
1089            heddle_object_model::compact::decode_tree_frame(data)
1090                .map_err(|error| StoreError::InvalidObject(error.to_string()))?
1091                .into_iter()
1092                .map(|tree| {
1093                    let id = PackObjectId::Hash(tree.hash());
1094                    let bytes = tree
1095                        .encode_canonical()
1096                        .map_err(|error| StoreError::InvalidObject(error.to_string()))?;
1097                    Ok((id, ObjectType::Tree, bytes))
1098                })
1099                .collect::<Result<Vec<_>>>()
1100                .map(Some)
1101        }
1102        ObjectType::State if heddle_object_model::compact::is_state_frame(data) => {
1103            heddle_object_model::compact::decode_state_frame(data)
1104                .map_err(|error| StoreError::InvalidObject(error.to_string()))?
1105                .into_iter()
1106                .map(|state| {
1107                    let id = PackObjectId::StateId(state.state_id);
1108                    let bytes = state
1109                        .encode_current_msgpack()
1110                        .map_err(|error| StoreError::InvalidObject(error.to_string()))?;
1111                    Ok((id, ObjectType::State, bytes))
1112                })
1113                .collect::<Result<Vec<_>>>()
1114                .map(Some)
1115        }
1116        _ if is_compact_frame(data) => Err(StoreError::InvalidObject(
1117            "compact frame magic does not match its pack object type".into(),
1118        )),
1119        _ => Ok(None),
1120    }
1121}
1122
1123fn compact_index_miss(id: &PackObjectId) -> StoreError {
1124    StoreError::InvalidObject(format!(
1125        "compact frame does not contain indexed object {id:?}"
1126    ))
1127}
1128
1129#[cfg(test)]
1130mod tests {
1131    use super::{PackObjectId, PackReader, verify_record_id_matches};
1132    use crate::{object::ContentHash, store::StoreError};
1133
1134    #[test]
1135    fn test_require_delta_base_hash_rejects_missing_hash() {
1136        let error =
1137            PackReader::require_delta_base_hash(None).expect_err("missing hash should fail");
1138
1139        assert!(
1140            matches!(error, StoreError::InvalidObject(message) if message == "pack object type is Delta but base hash is missing")
1141        );
1142    }
1143
1144    #[test]
1145    fn verify_record_id_matches_accepts_identical_ids() {
1146        let id = PackObjectId::Hash(ContentHash::from_bytes([7u8; 32]));
1147        verify_record_id_matches(&id, &id).expect("matching ids must verify");
1148    }
1149
1150    #[test]
1151    fn verify_record_id_matches_rejects_mismatched_ids() {
1152        let asked = PackObjectId::Hash(ContentHash::from_bytes([7u8; 32]));
1153        let found = PackObjectId::Hash(ContentHash::from_bytes([8u8; 32]));
1154        let error = verify_record_id_matches(&asked, &found)
1155            .expect_err("mismatched record id must error rather than silently route");
1156        assert!(
1157            matches!(&error, StoreError::InvalidObject(message) if message.contains("stale or corrupt")),
1158            "stale-index mismatch must surface as InvalidObject with the diagnostic phrase, got: {error:?}",
1159        );
1160    }
1161}