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