1#[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#[derive(Debug, Clone, Copy, Eq, PartialEq)]
40pub enum PackReadTier {
41 Hot,
43 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
74enum 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 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 pub fn list_ids(&self) -> Result<Vec<PackObjectId>> {
188 self.index.ids()
189 }
190
191 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 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 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 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 pub fn representation_hash(&self) -> PackRepresentationHash {
299 PackRepresentationHash::compute(self.data.as_slice())
300 }
301
302 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
960fn 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}