1#[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#[derive(Debug, Clone, Copy, Eq, PartialEq)]
41pub enum PackReadTier {
42 Hot,
44 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
75enum 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 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 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 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 let _ = self.offsets.set(index);
226 self.offsets
227 .get()
228 .ok_or_else(|| StoreError::InvalidObject("physical pack index unavailable".into()))
229 }
230
231 pub fn list_ids(&self) -> Result<Vec<PackObjectId>> {
233 self.index.ids()
234 }
235
236 #[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 #[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 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 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 pub fn representation_hash(&self) -> PackRepresentationHash {
351 PackRepresentationHash::compute(self.data.as_slice())
352 }
353
354 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
998fn 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}