1pub mod window;
71
72pub mod rewrite;
73use crate::delta;
74use crate::hash::{self, Hash};
75use crate::object::{MkitError, Object};
76use crate::store::{MAX_RAW_OBJECT_SIZE, ObjectStore};
77pub use rewrite::{Rewritten, rewrite_excluding};
78use std::borrow::Cow;
79use std::ops::Range;
80use std::sync::atomic::{AtomicU64, Ordering};
81
82pub const MAGIC: &[u8; 4] = b"MKIT";
84pub const VERSION: u32 = 1;
88pub const VERSION_V2: u32 = 2;
93
94pub const MAX_ENTRIES: u32 = 10_000_000;
96pub const MAX_TOTAL_PAYLOAD: u64 = 4 * 1024 * 1024 * 1024;
98pub const TRAILER_LEN: usize = 32;
100
101pub const HEADER_LEN: usize = 4 + 4 + 4;
103pub const ENTRY_FRAME_LEN: usize = 1 + 4;
105pub const VERSION_OFFSET: usize = 4;
110pub const ENTRY_COUNT_OFFSET: usize = 8;
116
117#[cfg(feature = "pack-zstd")]
122const MIN_COMPRESS_LEN: usize = 64;
123#[cfg(feature = "pack-zstd")]
127const ZSTD_LEVEL: i32 = 3;
128const ZSTD_LEN_PREFIX: usize = 4;
132
133#[derive(Debug, thiserror::Error)]
137pub enum PackError {
138 #[error("packfile is shorter than the {HEADER_LEN}-byte header + {TRAILER_LEN}-byte trailer")]
139 PackfileTooShort,
140 #[error("first 4 bytes are not ASCII \"MKIT\"")]
141 InvalidMagic,
142 #[error("version {0} is not supported (v1 or v2 only)")]
143 UnsupportedVersion(u32),
144 #[error(
145 "entry_type {0:#04x} is not 0x00 (raw), 0x02 (delta), 0x03 (zstd-raw), or 0x04 \
146 (zstd-delta) — or is a v2-only entry type inside a version-1 pack"
147 )]
148 InvalidEntryType(u8),
149 #[error("entry_count {0} exceeds the {MAX_ENTRIES} cap")]
150 TooManyObjects(u32),
151 #[error("sum of payload_len exceeds {MAX_TOTAL_PAYLOAD} bytes")]
152 PackfileTooLarge,
153 #[error("entry payload extends past the trailer offset")]
154 UnexpectedEof,
155 #[error("trailer BLAKE3 mismatch — packfile is corrupt or truncated")]
156 PackfileCorrupted,
157 #[error("delta entry references base hash {0} which is not in this pack or the store")]
158 DeltaBaseMissing(String),
159 #[error("delta entry payload is shorter than the 32-byte base hash prefix")]
160 DeltaEntryTruncated,
161 #[error("delta reconstruction failed: {0}")]
162 DeltaApply(#[from] MkitError),
163 #[error("pack entry is not a canonical storable object: {0}")]
164 InvalidObject(MkitError),
165 #[error("pack entry resolves to pack-only delta object")]
166 NonStorableObject,
167 #[error("pack contains trailing bytes after declared entries")]
168 TrailingData,
169 #[error("store I/O failure: {0}")]
170 Store(#[from] crate::store::StoreError),
171 #[error("zstd entry payload is shorter than its length-prefix header")]
176 ZstdEntryTruncated,
177 #[error(
181 "zstd entry's claimed decompressed size {0} exceeds the {MAX_RAW_OBJECT_SIZE}-byte cap"
182 )]
183 DecompressedSizeOverCap(usize),
184 #[error("zstd entry claims {0} decompressed bytes but produced {1}")]
188 DecompressedSizeMismatch(usize, usize),
189 #[error("zstd decompression failed: {0}")]
191 ZstdDecompress(String),
192 #[error("pack writer is in raw-only mode and does not accept delta entries")]
196 RawOnly,
197}
198
199#[derive(Debug, Clone, Default, PartialEq, Eq)]
202pub struct UnpackReport {
203 pub raw_count: u32,
204 pub delta_count: u32,
205 pub stored: Vec<Hash>,
207}
208
209#[derive(Debug)]
215pub struct PreparedRaw {
216 hash: Hash,
217 bytes: Vec<u8>,
218 frame: Option<Vec<u8>>,
219}
220
221impl PreparedRaw {
222 #[must_use]
227 pub fn hash(&self) -> Hash {
228 self.hash
229 }
230
231 #[must_use]
236 pub fn conservative_len(&self) -> usize {
237 self.bytes.len()
238 }
239}
240
241#[derive(Debug)]
245pub struct PreparedDelta {
246 base: Hash,
247 stream: Vec<u8>,
248 frame: Option<Vec<u8>>,
249}
250
251impl PreparedDelta {
252 #[must_use]
255 pub fn base(&self) -> Hash {
256 self.base
257 }
258
259 #[must_use]
262 pub fn conservative_len(&self) -> usize {
263 hash::HASH_LEN + self.stream.len()
264 }
265}
266
267#[derive(Debug)]
274pub struct PackWriter {
275 buf: Vec<u8>,
282 entry_count: u32,
283 total_payload: u64,
284 has_compressed_entry: bool,
289 raw_only: bool,
294}
295
296impl Default for PackWriter {
297 fn default() -> Self {
298 Self::new()
299 }
300}
301
302impl PackWriter {
303 #[must_use]
305 pub fn new() -> Self {
306 let mut buf = Vec::with_capacity(HEADER_LEN);
307 buf.extend_from_slice(MAGIC);
308 buf.extend_from_slice(&VERSION.to_le_bytes());
309 buf.extend_from_slice(&0u32.to_le_bytes()); Self {
311 buf,
312 entry_count: 0,
313 total_payload: 0,
314 has_compressed_entry: false,
315 raw_only: false,
316 }
317 }
318
319 #[must_use]
328 pub fn new_raw_only() -> Self {
329 let mut w = Self::new();
330 w.raw_only = true;
331 w
332 }
333
334 pub fn push_raw(&mut self, hash_of_bytes: Hash, bytes: &[u8]) -> Result<Hash, PackError> {
350 let frame = if self.raw_only {
351 None
352 } else {
353 maybe_compress(bytes)
354 };
355 self.append_raw_frame(hash_of_bytes, bytes, frame)
356 }
357
358 #[must_use]
373 pub fn prepare_raw(hash_of_bytes: Hash, bytes: Vec<u8>) -> PreparedRaw {
374 let frame = maybe_compress(&bytes);
375 PreparedRaw {
376 hash: hash_of_bytes,
377 bytes,
378 frame,
379 }
380 }
381
382 pub fn push_prepared_raw(&mut self, entry: PreparedRaw) -> Result<Hash, PackError> {
387 let frame = if self.raw_only { None } else { entry.frame };
388 self.append_raw_frame(entry.hash, &entry.bytes, frame)
389 }
390
391 fn append_raw_frame(
397 &mut self,
398 hash_of_bytes: Hash,
399 bytes: &[u8],
400 frame: Option<Vec<u8>>,
401 ) -> Result<Hash, PackError> {
402 if let Some(frame) = frame {
403 let uncompressed_len: u32 = bytes
404 .len()
405 .try_into()
406 .map_err(|_| PackError::PackfileTooLarge)?;
407 let payload_len = ZSTD_LEN_PREFIX + frame.len();
408 self.check_caps_for(payload_len)?;
409 self.total_payload += payload_len as u64;
410 self.append_entry(0x03, &[&uncompressed_len.to_le_bytes(), &frame])?;
411 self.has_compressed_entry = true;
412 } else {
413 self.check_caps_for(bytes.len())?;
414 self.total_payload += bytes.len() as u64;
415 self.append_entry(0x00, &[bytes])?;
416 }
417 self.entry_count += 1;
418 Ok(hash_of_bytes)
419 }
420
421 pub fn push_delta(&mut self, base_hash: &Hash, delta_stream: &[u8]) -> Result<(), PackError> {
433 if self.raw_only {
434 return Err(PackError::RawOnly);
435 }
436 let frame = maybe_compress(delta_stream);
437 self.append_delta_frame(base_hash, delta_stream, frame)
438 }
439
440 #[must_use]
446 pub fn prepare_delta(base_hash: Hash, delta_stream: Vec<u8>) -> PreparedDelta {
447 let frame = maybe_compress(&delta_stream);
448 PreparedDelta {
449 base: base_hash,
450 stream: delta_stream,
451 frame,
452 }
453 }
454
455 pub fn push_prepared_delta(&mut self, entry: PreparedDelta) -> Result<(), PackError> {
459 if self.raw_only {
460 return Err(PackError::RawOnly);
461 }
462 self.append_delta_frame(&entry.base, &entry.stream, entry.frame)
463 }
464
465 fn append_delta_frame(
469 &mut self,
470 base_hash: &Hash,
471 delta_stream: &[u8],
472 frame: Option<Vec<u8>>,
473 ) -> Result<(), PackError> {
474 if let Some(frame) = frame {
475 let uncompressed_len: u32 = delta_stream
476 .len()
477 .try_into()
478 .map_err(|_| PackError::PackfileTooLarge)?;
479 let payload_len = hash::HASH_LEN + ZSTD_LEN_PREFIX + frame.len();
480 self.check_caps_for(payload_len)?;
481 self.total_payload += payload_len as u64;
482 self.append_entry(
483 0x04,
484 &[
485 base_hash.as_slice(),
486 &uncompressed_len.to_le_bytes(),
487 &frame,
488 ],
489 )?;
490 self.has_compressed_entry = true;
491 } else {
492 let payload_len = hash::HASH_LEN + delta_stream.len();
493 self.check_caps_for(payload_len)?;
494 self.total_payload += payload_len as u64;
495 self.append_entry(0x02, &[base_hash.as_slice(), delta_stream])?;
496 }
497 self.entry_count += 1;
498 Ok(())
499 }
500
501 fn append_entry(&mut self, etype: u8, parts: &[&[u8]]) -> Result<(), PackError> {
507 let payload_len: usize = parts.iter().map(|p| p.len()).sum();
508 let plen: u32 = payload_len
509 .try_into()
510 .map_err(|_| PackError::PackfileTooLarge)?;
511 self.buf.push(etype);
512 self.buf.extend_from_slice(&plen.to_le_bytes());
513 for p in parts {
514 self.buf.extend_from_slice(p);
515 }
516 Ok(())
517 }
518
519 fn check_caps_for(&self, add_len: usize) -> Result<(), PackError> {
520 let next_count = u64::from(self.entry_count) + 1;
521 if next_count > u64::from(MAX_ENTRIES) {
522 return Err(PackError::TooManyObjects(MAX_ENTRIES + 1));
523 }
524 let next_total = self.total_payload.saturating_add(add_len as u64);
525 if next_total > MAX_TOTAL_PAYLOAD {
526 return Err(PackError::PackfileTooLarge);
527 }
528 Ok(())
529 }
530
531 #[must_use]
533 pub fn entry_count(&self) -> usize {
534 self.entry_count as usize
535 }
536
537 #[must_use]
548 pub fn total_payload(&self) -> u64 {
549 self.total_payload
550 }
551
552 pub fn finish(self) -> Result<Vec<u8>, PackError> {
560 self.finish_inner(None)
561 }
562
563 #[cfg(test)]
572 pub(crate) fn finish_tracking_bytes_copied(
573 self,
574 bytes_copied: &AtomicU64,
575 ) -> Result<Vec<u8>, PackError> {
576 self.finish_inner(Some(bytes_copied))
577 }
578
579 fn finish_inner(mut self, bytes_copied: Option<&AtomicU64>) -> Result<Vec<u8>, PackError> {
580 if self.entry_count > MAX_ENTRIES {
581 return Err(PackError::TooManyObjects(self.entry_count));
582 }
583 let version = if self.has_compressed_entry {
584 VERSION_V2
585 } else {
586 VERSION
587 };
588 self.buf[VERSION_OFFSET..VERSION_OFFSET + 4].copy_from_slice(&version.to_le_bytes());
589 self.buf[ENTRY_COUNT_OFFSET..ENTRY_COUNT_OFFSET + 4]
590 .copy_from_slice(&self.entry_count.to_le_bytes());
591 let trailer = hash::hash(&self.buf);
592 if let Some(c) = bytes_copied {
593 c.fetch_add(trailer.len() as u64, Ordering::Relaxed);
594 }
595 self.buf.extend_from_slice(&trailer);
596 Ok(self.buf)
597 }
598}
599
600#[must_use]
604pub fn pack_key(pack_bytes: &[u8]) -> Hash {
605 hash::hash(pack_bytes)
606}
607
608#[cfg(feature = "pack-zstd")]
619fn maybe_compress(data: &[u8]) -> Option<Vec<u8>> {
620 maybe_compress_capped(data, MAX_RAW_OBJECT_SIZE)
621}
622
623#[cfg(feature = "pack-zstd")]
624fn maybe_compress_capped(data: &[u8], max_len: usize) -> Option<Vec<u8>> {
625 if data.len() > max_len {
628 return None;
629 }
630 if data.len() < MIN_COMPRESS_LEN {
631 return None;
632 }
633 let compressed = ZSTD_COMPRESSOR
634 .with(|c| c.borrow_mut().compress(data))
635 .ok()?;
636 if ZSTD_LEN_PREFIX + compressed.len() < data.len() {
637 Some(compressed)
638 } else {
639 None
640 }
641}
642
643#[cfg(feature = "pack-zstd")]
644thread_local! {
645 static ZSTD_COMPRESSOR: std::cell::RefCell<zstd::bulk::Compressor<'static>> =
659 std::cell::RefCell::new(
660 zstd::bulk::Compressor::new(ZSTD_LEVEL).expect("ZSTD_LEVEL is a valid zstd level"),
661 );
662}
663
664#[cfg(not(feature = "pack-zstd"))]
665fn maybe_compress(_data: &[u8]) -> Option<Vec<u8>> {
666 None
671}
672
673fn decompress_zstd_entry(payload: &[u8]) -> Result<Vec<u8>, PackError> {
680 let len = zstd_entry_len(payload)?;
681 let mut output = Vec::new();
682 output
683 .try_reserve_exact(len)
684 .map_err(|_| PackError::PackfileTooLarge)?;
685 decompress_zstd_into(payload, &mut output)?;
686 Ok(output)
687}
688
689#[cfg(all(test, feature = "pack-ruzstd"))]
693fn decompress_zstd_entry_with(
694 payload: &[u8],
695 backend: fn(&[u8], usize) -> Result<Vec<u8>, PackError>,
696) -> Result<Vec<u8>, PackError> {
697 let (uncompressed_len, frame) = zstd_claim(payload)?;
698 let decompressed = backend(frame, uncompressed_len)?;
699 if decompressed.len() != uncompressed_len {
700 return Err(PackError::DecompressedSizeMismatch(
701 uncompressed_len,
702 decompressed.len(),
703 ));
704 }
705 Ok(decompressed)
706}
707
708fn zstd_claim(payload: &[u8]) -> Result<(usize, &[u8]), PackError> {
709 if payload.len() < ZSTD_LEN_PREFIX {
710 return Err(PackError::ZstdEntryTruncated);
711 }
712 let uncompressed_len =
713 u32::from_le_bytes(payload[..ZSTD_LEN_PREFIX].try_into().expect("4 bytes")) as usize;
714 if uncompressed_len > MAX_RAW_OBJECT_SIZE {
715 return Err(PackError::DecompressedSizeOverCap(uncompressed_len));
716 }
717 let frame = &payload[ZSTD_LEN_PREFIX..];
718 Ok((uncompressed_len, frame))
719}
720
721fn zstd_entry_len(payload: &[u8]) -> Result<usize, PackError> {
722 zstd_claim(payload).map(|(len, _)| len)
723}
724
725fn decompress_zstd_into(payload: &[u8], output: &mut Vec<u8>) -> Result<(), PackError> {
728 let (len, frame) = zstd_claim(payload)?;
729 output.clear();
730 zstd_decompress_into(frame, len, output)?;
731 if output.len() != len {
732 return Err(PackError::DecompressedSizeMismatch(len, output.len()));
733 }
734 Ok(())
735}
736
737#[cfg(any(feature = "pack-zstd", feature = "pack-ruzstd"))]
742const ZSTD_FRAME_MAGIC: [u8; 4] = [0x28, 0xB5, 0x2F, 0xFD];
743
744#[cfg(any(feature = "pack-zstd", feature = "pack-ruzstd"))]
747fn require_zstd_frame_magic(frame: &[u8]) -> Result<(), PackError> {
748 if frame.starts_with(&ZSTD_FRAME_MAGIC) {
749 Ok(())
750 } else {
751 Err(PackError::ZstdDecompress(
752 "entry payload does not start with a Zstandard frame magic \
753 (skippable and legacy frames are not allowed)"
754 .to_string(),
755 ))
756 }
757}
758
759#[cfg(feature = "pack-zstd")]
765fn zstd_decompress_into(
766 frame: &[u8],
767 capacity: usize,
768 output: &mut Vec<u8>,
769) -> Result<(), PackError> {
770 require_zstd_frame_magic(frame)?;
771 match zstd::zstd_safe::find_frame_compressed_size(frame) {
774 Ok(n) if n == frame.len() => {}
775 Ok(n) => {
776 return Err(PackError::ZstdDecompress(format!(
777 "{} byte(s) after the entry's single zstd frame",
778 frame.len() - n
779 )));
780 }
781 Err(code) => {
782 return Err(PackError::ZstdDecompress(
783 zstd::zstd_safe::get_error_name(code).to_string(),
784 ));
785 }
786 }
787 let mut decoder =
788 zstd::bulk::Decompressor::new().map_err(|e| PackError::ZstdDecompress(e.to_string()))?;
789 decoder
790 .decompress_to_buffer(frame, output)
791 .map_err(|e| PackError::ZstdDecompress(e.to_string()))?;
792 if output.len() > capacity {
793 return Err(PackError::ZstdDecompress(
794 "zstd frame exceeds its claim".to_string(),
795 ));
796 }
797 Ok(())
798}
799
800#[cfg(all(feature = "pack-zstd", test))]
801fn zstd_decompress_capped(frame: &[u8], capacity: usize) -> Result<Vec<u8>, PackError> {
802 let mut output = Vec::new();
803 output
804 .try_reserve_exact(capacity)
805 .map_err(|_| PackError::PackfileTooLarge)?;
806 zstd_decompress_into(frame, capacity, &mut output)?;
807 Ok(output)
808}
809
810#[cfg(all(not(feature = "pack-zstd"), feature = "pack-ruzstd"))]
812fn zstd_decompress_into(
813 frame: &[u8],
814 capacity: usize,
815 output: &mut Vec<u8>,
816) -> Result<(), PackError> {
817 ruzstd_decompress_into(frame, capacity, output)
818}
819
820#[cfg(not(any(feature = "pack-zstd", feature = "pack-ruzstd")))]
821fn zstd_decompress_into(
822 _frame: &[u8],
823 _capacity: usize,
824 _output: &mut Vec<u8>,
825) -> Result<(), PackError> {
826 Err(PackError::ZstdDecompress(
827 "this build was compiled without the `pack-zstd` or `pack-ruzstd` feature".to_string(),
828 ))
829}
830
831#[cfg(feature = "pack-ruzstd")]
835#[cfg_attr(all(feature = "pack-zstd", not(test)), allow(dead_code))]
836const RUZSTD_WINDOW_LIMIT: u64 = 8 << 20;
837
838pub fn peek_delta_header(frame: &[u8]) -> Result<(u32, u32), PackError> {
855 let mut header = [0; delta::HEADER_LEN];
856 read_delta_header_prefix(frame, &mut header)?;
857 if header[0] != delta::STREAM_VERSION {
858 return Err(PackError::DeltaApply(MkitError::UnsupportedObjectVersion));
859 }
860 Ok((
861 u32::from_le_bytes([header[1], header[2], header[3], header[4]]),
862 u32::from_le_bytes([header[5], header[6], header[7], header[8]]),
863 ))
864}
865
866#[cfg(feature = "pack-ruzstd")]
867fn read_delta_header_prefix(
868 frame: &[u8],
869 header: &mut [u8; delta::HEADER_LEN],
870) -> Result<(), PackError> {
871 use ruzstd::decoding::{FrameDecoder, StreamingDecoder};
872 use std::io::Read as _;
873 require_zstd_frame_magic(frame)?;
874 let mut decoder = FrameDecoder::new();
875 decoder.set_max_window_size(RUZSTD_WINDOW_LIMIT);
876 let mut stream = StreamingDecoder::new_with_decoder(frame, decoder)
877 .map_err(|error| PackError::ZstdDecompress(error.to_string()))?;
878 stream
879 .read_exact(header)
880 .map_err(|error| PackError::ZstdDecompress(error.to_string()))
881}
882
883#[cfg(all(feature = "pack-zstd", not(feature = "pack-ruzstd")))]
884fn read_delta_header_prefix(
885 frame: &[u8],
886 header: &mut [u8; delta::HEADER_LEN],
887) -> Result<(), PackError> {
888 use std::io::Read as _;
889 require_zstd_frame_magic(frame)?;
890 let mut stream = zstd::stream::read::Decoder::with_buffer(frame)
891 .map_err(|error| PackError::ZstdDecompress(error.to_string()))?
892 .single_frame();
893 stream
894 .window_log_max(23)
895 .map_err(|error| PackError::ZstdDecompress(error.to_string()))?;
896 stream
897 .read_exact(header)
898 .map_err(|error| PackError::ZstdDecompress(error.to_string()))
899}
900
901#[cfg(not(any(feature = "pack-zstd", feature = "pack-ruzstd")))]
902fn read_delta_header_prefix(
903 frame: &[u8],
904 _header: &mut [u8; delta::HEADER_LEN],
905) -> Result<(), PackError> {
906 zstd_decompress_into(frame, 0, &mut Vec::new())
907}
908
909#[cfg(feature = "pack-ruzstd")]
932#[cfg_attr(not(test), allow(dead_code))]
933pub(crate) fn ruzstd_decompress_capped(
934 frame: &[u8],
935 capacity: usize,
936) -> Result<Vec<u8>, PackError> {
937 let mut output = Vec::new();
938 output
939 .try_reserve_exact(capacity)
940 .map_err(|_| PackError::PackfileTooLarge)?;
941 ruzstd_decompress_into(frame, capacity, &mut output)?;
942 Ok(output)
943}
944
945#[cfg(feature = "pack-ruzstd")]
946#[cfg_attr(all(feature = "pack-zstd", not(test)), allow(dead_code))]
947fn ruzstd_decompress_into(
948 frame: &[u8],
949 capacity: usize,
950 output: &mut Vec<u8>,
951) -> Result<(), PackError> {
952 use ruzstd::decoding::{FrameDecoder, StreamingDecoder};
953 use std::io::Read as _;
954
955 fn fail(msg: impl std::fmt::Display) -> PackError {
956 PackError::ZstdDecompress(msg.to_string())
957 }
958
959 require_zstd_frame_magic(frame)?;
960 let cap = u64::try_from(capacity).unwrap_or(u64::MAX);
961 let mut decoder = FrameDecoder::new();
962 decoder.set_max_window_size(RUZSTD_WINDOW_LIMIT);
963 let mut src = frame;
964 let mut stream = StreamingDecoder::new_with_decoder(&mut src, decoder).map_err(fail)?;
965
966 let descriptor = frame[ZSTD_FRAME_MAGIC.len()];
970 if descriptor & 0x08 != 0 {
971 return Err(fail("zstd frame descriptor has its reserved bit set"));
972 }
973 let declared =
976 (descriptor >> 6 != 0 || descriptor & 0x20 != 0).then(|| stream.decoder.content_size());
977 if let Some(n) = declared
978 && n > cap
979 {
980 return Err(fail(format_args!(
981 "frame content size {n} exceeds the claimed {capacity} bytes"
982 )));
983 }
984
985 let out = output;
986 out.clear();
987 let mut chunk = [0; 8192];
988 loop {
990 let room = capacity.saturating_sub(out.len());
992 let take = room.saturating_add(1).min(chunk.len());
993 let n = stream.read(&mut chunk[..take]).map_err(fail)?;
994 if n == 0 {
995 break;
996 }
997 if n > room {
998 return Err(fail(format_args!(
999 "zstd frame decompresses past the claimed {capacity} bytes"
1000 )));
1001 }
1002 out.extend_from_slice(&chunk[..n]);
1003 }
1004 let decoder = &stream.decoder;
1005 if !decoder.is_finished() {
1006 return Err(fail("zstd frame ended before its last block"));
1007 }
1008 if let Some(n) = declared
1009 && n != out.len() as u64
1010 {
1011 return Err(fail(format_args!(
1012 "frame content size {n} does not match the {} decoded bytes",
1013 out.len()
1014 )));
1015 }
1016 if let Some(stored) = decoder.get_checksum_from_data()
1017 && decoder.get_calculated_checksum() != Some(stored)
1018 {
1019 return Err(fail("zstd frame content checksum mismatch"));
1020 }
1021 drop(stream);
1022 if !src.is_empty() {
1023 return Err(fail(format_args!(
1024 "{} byte(s) after the entry's single zstd frame",
1025 src.len()
1026 )));
1027 }
1028 ruzstd_check_reserved_fields(frame).map_err(fail)?;
1029 Ok(())
1030}
1031
1032#[cfg(feature = "pack-ruzstd")]
1038#[cfg_attr(all(feature = "pack-zstd", not(test)), allow(dead_code))]
1039fn ruzstd_check_reserved_fields(frame: &[u8]) -> Result<(), &'static str> {
1040 const MALFORMED: &str = "malformed zstd frame";
1041 let byte = |i: usize| frame.get(i).copied().ok_or(MALFORMED);
1042 let descriptor = byte(ZSTD_FRAME_MAGIC.len())?;
1043 let single_segment = descriptor & 0x20 != 0;
1044 let fcs_len = match descriptor >> 6 {
1045 0 => usize::from(single_segment),
1046 1 => 2,
1047 2 => 4,
1048 _ => 8,
1049 };
1050 let mut pos = ZSTD_FRAME_MAGIC.len()
1051 + 1
1052 + usize::from(!single_segment)
1053 + [0, 1, 2, 4][usize::from(descriptor & 3)]
1054 + fcs_len;
1055 loop {
1056 let header = u32::from_le_bytes([byte(pos)?, byte(pos + 1)?, byte(pos + 2)?, 0]);
1057 pos += 3;
1058 let size = (header >> 3) as usize;
1059 let body_len = match (header >> 1) & 3 {
1060 1 => 1, _ => size,
1062 };
1063 let body = frame.get(pos..pos + body_len).ok_or(MALFORMED)?;
1064 if (header >> 1) & 3 == 2 && ruzstd_sequence_modes(body)? & 3 != 0 {
1065 return Err("zstd sequences section has its reserved mode bits set");
1066 }
1067 pos += body_len;
1068 if header & 1 == 1 {
1069 return Ok(());
1070 }
1071 }
1072}
1073
1074#[cfg(feature = "pack-ruzstd")]
1077#[cfg_attr(all(feature = "pack-zstd", not(test)), allow(dead_code))]
1078fn ruzstd_sequence_modes(block: &[u8]) -> Result<u8, &'static str> {
1079 const MALFORMED: &str = "malformed zstd block";
1080 let byte = |i: usize| block.get(i).copied().ok_or(MALFORMED).map(u64::from);
1084 let b0 = byte(0)?;
1085 let (header_len, content_len) = match (b0 & 3, (b0 >> 2) & 3) {
1087 (kind @ (0 | 1), format) => {
1089 let (len, regen) = match format {
1090 0 | 2 => (1, b0 >> 3),
1091 1 => (2, (b0 >> 4) | (byte(1)? << 4)),
1092 _ => (3, (b0 >> 4) | (byte(1)? << 4) | (byte(2)? << 12)),
1093 };
1094 (len, if kind == 0 { regen } else { 1 })
1095 }
1096 (_, format) => {
1098 let (len, bits) = match format {
1099 0 | 1 => (3, 10),
1100 2 => (4, 14),
1101 _ => (5, 18),
1102 };
1103 let mut h = 0u64;
1104 for i in (0..len).rev() {
1105 h = (h << 8) | byte(i)?;
1106 }
1107 (len, (h >> (4 + bits)) & ((1 << bits) - 1))
1108 }
1109 };
1110 let content_len = usize::try_from(content_len).map_err(|_| MALFORMED)?;
1112 let seq = header_len + content_len;
1113 let modes_at = match byte(seq)? {
1114 0 => return Ok(0),
1115 n if n < 128 => seq + 1,
1116 255 => seq + 3,
1117 _ => seq + 2,
1118 };
1119 block.get(modes_at).copied().ok_or(MALFORMED)
1120}
1121
1122pub fn delta_base_hashes(pack_bytes: &[u8]) -> Result<Vec<Hash>, PackError> {
1145 if pack_bytes.len() < HEADER_LEN + TRAILER_LEN {
1146 return Err(PackError::PackfileTooShort);
1147 }
1148 if &pack_bytes[..4] != MAGIC.as_slice() {
1149 return Err(PackError::InvalidMagic);
1150 }
1151 let version = u32::from_le_bytes(pack_bytes[4..8].try_into().expect("4 bytes"));
1152 if version != VERSION && version != VERSION_V2 {
1153 return Err(PackError::UnsupportedVersion(version));
1154 }
1155 let count = u32::from_le_bytes(
1156 pack_bytes[ENTRY_COUNT_OFFSET..ENTRY_COUNT_OFFSET + 4]
1157 .try_into()
1158 .expect("4 bytes"),
1159 );
1160 if count > MAX_ENTRIES {
1161 return Err(PackError::TooManyObjects(count));
1162 }
1163 let split = pack_bytes.len() - TRAILER_LEN;
1165
1166 let mut bases = Vec::new();
1167 let mut seen = std::collections::HashSet::new();
1168 let mut pos = HEADER_LEN;
1169 for _ in 0..count {
1170 if ENTRY_FRAME_LEN > split - pos {
1171 return Err(PackError::UnexpectedEof);
1172 }
1173 let etype = pack_bytes[pos];
1174 pos = pos.checked_add(1).ok_or(PackError::UnexpectedEof)?;
1175 let payload_len = u32::from_le_bytes(
1176 pack_bytes[pos..pos.checked_add(4).ok_or(PackError::UnexpectedEof)?]
1177 .try_into()
1178 .expect("4 bytes"),
1179 ) as usize;
1180 pos = pos.checked_add(4).ok_or(PackError::UnexpectedEof)?;
1181 if payload_len > split - pos {
1182 return Err(PackError::UnexpectedEof);
1183 }
1184 if etype == 0x02 || etype == 0x04 {
1189 if payload_len < TRAILER_LEN {
1190 return Err(PackError::DeltaEntryTruncated);
1191 }
1192 let mut base = [0u8; 32];
1193 base.copy_from_slice(
1194 &pack_bytes[pos..pos
1195 .checked_add(TRAILER_LEN)
1196 .ok_or(PackError::UnexpectedEof)?],
1197 );
1198 if seen.insert(base) {
1199 bases.push(base);
1200 }
1201 }
1202 pos = pos
1203 .checked_add(payload_len)
1204 .ok_or(PackError::UnexpectedEof)?;
1205 }
1206 Ok(bases)
1207}
1208
1209#[derive(Debug)]
1213pub struct PackReader;
1214
1215impl PackReader {
1216 pub fn read(pack_bytes: &[u8], store: &ObjectStore) -> Result<UnpackReport, PackError> {
1233 Self::read_with_payload_cap(pack_bytes, store, MAX_TOTAL_PAYLOAD)
1234 }
1235
1236 pub(crate) fn read_with_payload_cap(
1243 pack_bytes: &[u8],
1244 store: &ObjectStore,
1245 payload_cap: u64,
1246 ) -> Result<UnpackReport, PackError> {
1247 Self::read_inner(pack_bytes, store, payload_cap, None)
1248 }
1249
1250 #[cfg(test)]
1258 pub(crate) fn read_tracking_owned_bytes(
1259 pack_bytes: &[u8],
1260 store: &ObjectStore,
1261 owned_bytes: &AtomicU64,
1262 ) -> Result<UnpackReport, PackError> {
1263 Self::read_inner(pack_bytes, store, MAX_TOTAL_PAYLOAD, Some(owned_bytes))
1264 }
1265
1266 fn read_inner(
1267 pack_bytes: &[u8],
1268 store: &ObjectStore,
1269 payload_cap: u64,
1270 owned_bytes: Option<&AtomicU64>,
1271 ) -> Result<UnpackReport, PackError> {
1272 let budget = ResidentBudget::new(resident_bytes_cap(pack_bytes.len()), owned_bytes);
1273 Self::read_with_budget(pack_bytes, store, payload_cap, &budget)
1274 }
1275
1276 fn read_with_budget(
1277 pack_bytes: &[u8],
1278 store: &ObjectStore,
1279 payload_cap: u64,
1280 budget: &ResidentBudget<'_>,
1281 ) -> Result<UnpackReport, PackError> {
1282 let mut parser = PackEntries::new_with_payload_cap(pack_bytes, payload_cap)?;
1283 let mut entries = Vec::new();
1286 entries
1287 .try_reserve_exact(parser.entry_count())
1288 .map_err(|_| PackError::PackfileTooLarge)?;
1289 let mut uses: std::collections::HashMap<Hash, BaseUses> = std::collections::HashMap::new();
1290 for position in 0..parser.entry_count() {
1291 let entry = parser.next_encoded_entry()?;
1292 if let Entry::Delta { base, .. } = entry {
1293 let usage = uses.entry(base).or_default();
1294 usage.remaining += 1;
1295 usage.last_position = position;
1296 }
1297 entries.push(entry);
1298 }
1299 let batch = store.batch();
1300 let raw_frames: Vec<_> = entries
1303 .iter()
1304 .enumerate()
1305 .filter_map(|(position, entry)| match entry {
1306 Entry::Raw(payload) => Some((position, *payload)),
1307 Entry::Delta { .. } => None,
1308 })
1309 .collect();
1310 let raw_results = stage_raw_entries(&batch, &raw_frames, &uses, budget);
1311 finish_pack_read(entries, raw_results, uses, budget, store, batch)
1312 }
1313}
1314
1315pub trait DeltaBaseSource {
1328 const VERIFIED: bool = false;
1339
1340 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError>;
1349
1350 fn base_with_admission(
1359 &mut self,
1360 id: &Hash,
1361 admit: impl FnOnce(usize) -> Result<(), PackError>,
1362 ) -> Result<Option<Vec<u8>>, PackError> {
1363 let Some(bytes) = self.base(id)? else {
1364 return Ok(None);
1365 };
1366 admit(bytes.len())?;
1367 Ok(Some(bytes))
1368 }
1369}
1370
1371#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1375pub struct NoExternalBases;
1376
1377impl DeltaBaseSource for NoExternalBases {
1378 fn base(&mut self, _id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
1379 Ok(None)
1380 }
1381}
1382
1383impl DeltaBaseSource for &ObjectStore {
1384 const VERIFIED: bool = true;
1387
1388 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
1389 if self.contains(id) {
1390 Ok(Some(self.read(id)?))
1391 } else {
1392 Ok(None)
1393 }
1394 }
1395
1396 fn base_with_admission(
1397 &mut self,
1398 id: &Hash,
1399 admit: impl FnOnce(usize) -> Result<(), PackError>,
1400 ) -> Result<Option<Vec<u8>>, PackError> {
1401 if !self.contains(id) {
1402 return Ok(None);
1403 }
1404 self.read_with_allocator(id, |len| {
1405 admit(len)?;
1406 let mut bytes = Vec::new();
1407 bytes
1408 .try_reserve_exact(len)
1409 .map_err(|_| PackError::PackfileTooLarge)?;
1410 Ok(bytes)
1411 })
1412 .map(Some)
1413 }
1414}
1415
1416#[derive(Debug)]
1418#[non_exhaustive]
1419pub struct DecodedEntry<'a> {
1420 pub id: Hash,
1423 pub bytes: &'a [u8],
1426 pub object: Object,
1428 pub from_delta: bool,
1430 pub frame_offset: u64,
1432 pub frame_length: u64,
1434 pub wire_type: u8,
1436 pub delta_base: Option<Hash>,
1438}
1439
1440#[derive(Debug, Default, Clone, PartialEq, Eq)]
1442#[non_exhaustive]
1443pub struct DecodeReport {
1444 pub raw_count: usize,
1446 pub delta_count: usize,
1448 pub ids: Vec<Hash>,
1451}
1452
1453#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1472#[non_exhaustive]
1473pub struct DecodeLimits {
1474 pub max_decoded_bytes: u64,
1482 entry_geometry: Option<(u64, u64)>,
1483}
1484
1485impl DecodeLimits {
1486 pub const DEFAULT_MAX_DECODED_BYTES: u64 = MAX_RAW_OBJECT_SIZE as u64;
1489
1490 #[must_use]
1492 pub const fn with_max_decoded_bytes(mut self, bytes: u64) -> Self {
1493 self.max_decoded_bytes = bytes;
1494 self
1495 }
1496 #[must_use]
1498 pub const fn with_entry_geometry(mut self, frame: u64, delta_stream: u64) -> Self {
1499 self.entry_geometry = Some((frame, delta_stream));
1500 self
1501 }
1502
1503 fn check_frame(&self, kind: u8, payload: &[u8]) -> Result<(), PackError> {
1504 if let Some((frame, stream)) = self.entry_geometry {
1505 let bytes = payload.len() as u64;
1506 if bytes > frame || (kind == 2 && bytes > stream.saturating_add(32)) {
1507 return Err(PackError::PackfileTooLarge);
1508 }
1509 if kind == 4
1510 && zstd_claim(payload.get(32..).ok_or(PackError::DeltaEntryTruncated)?)?.0 as u64
1511 > stream
1512 {
1513 return Err(PackError::PackfileTooLarge);
1514 }
1515 }
1516 Ok(())
1517 }
1518}
1519
1520impl Default for DecodeLimits {
1521 fn default() -> Self {
1522 Self {
1523 max_decoded_bytes: Self::DEFAULT_MAX_DECODED_BYTES,
1524 entry_geometry: None,
1525 }
1526 }
1527}
1528
1529#[derive(Debug)]
1535struct DecodeBudget {
1536 used: u64,
1537 max: u64,
1538}
1539
1540impl DecodeBudget {
1541 fn charge(&mut self, bytes: u64) -> Result<(), PackError> {
1542 self.used = self.used.saturating_add(bytes);
1543 if self.used > self.max {
1544 return Err(PackError::PackfileTooLarge);
1545 }
1546 Ok(())
1547 }
1548}
1549
1550struct ChargedBases<'b, B> {
1555 inner: &'b mut B,
1556 budget: &'b mut DecodeBudget,
1557 charged: &'b mut std::collections::HashMap<Hash, u64>,
1558}
1559
1560impl<B: DeltaBaseSource> DeltaBaseSource for ChargedBases<'_, B> {
1561 const VERIFIED: bool = B::VERIFIED;
1562
1563 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
1564 self.base_with_admission(id, |_| Ok(()))
1565 }
1566
1567 fn base_with_admission(
1568 &mut self,
1569 id: &Hash,
1570 admit: impl FnOnce(usize) -> Result<(), PackError>,
1571 ) -> Result<Option<Vec<u8>>, PackError> {
1572 self.inner.base_with_admission(id, |len| {
1573 let charged_len = u64::try_from(len).map_err(|_| PackError::PackfileTooLarge)?;
1574 self.budget.charge(charged_len)?;
1575 admit(len)?;
1576 let held = self.charged.entry(*id).or_default();
1577 *held = held.saturating_add(charged_len);
1578 Ok(())
1579 })
1580 }
1581}
1582
1583impl<B> ChargedBases<'_, B> {
1584 fn release(&mut self, id: &Hash) {
1586 if let Some(len) = self.charged.remove(id) {
1587 self.budget.used = self.budget.used.saturating_sub(len);
1588 }
1589 }
1590}
1591
1592fn le_u32_at(bytes: &[u8], at: usize) -> Option<u64> {
1594 let field: [u8; 4] = bytes.get(at..at.checked_add(4)?)?.try_into().ok()?;
1595 Some(u64::from(u32::from_le_bytes(field)))
1596}
1597
1598fn charge_compressed_claims(pack: &[u8], budget: &mut DecodeBudget) -> Result<(), PackError> {
1605 let split = pack.len() - TRAILER_LEN;
1606 let mut pos = HEADER_LEN;
1607 while pos < split {
1608 let etype = pack[pos];
1609 let payload_len = le_u32_at(pack, pos + 1)
1610 .and_then(|len| usize::try_from(len).ok())
1611 .ok_or(PackError::UnexpectedEof)?;
1612 let start = pos + ENTRY_FRAME_LEN;
1613 if payload_len > split.saturating_sub(start) {
1614 return Err(PackError::UnexpectedEof);
1615 }
1616 let payload = &pack[start..start + payload_len];
1617 let claim = match etype {
1618 0x03 => le_u32_at(payload, 0),
1619 0x04 => le_u32_at(payload, hash::HASH_LEN),
1620 _ => None,
1621 };
1622 if let Some(claim) = claim {
1623 budget.charge(claim)?;
1624 }
1625 pos = start + payload_len;
1626 }
1627 Ok(())
1628}
1629
1630#[derive(Debug)]
1631struct CursorFrame<'a> {
1632 entry: PackEntry<'a>,
1633 frame_offset: u64,
1634 frame_length: u64,
1635 wire_type: u8,
1636 delta_base: Option<Hash>,
1637}
1638
1639#[derive(Debug)]
1647pub struct PackDecodeCursor<'a> {
1648 entries: Vec<Option<CursorFrame<'a>>>,
1649 next: usize,
1650 budget: DecodeBudget,
1651 charged: std::collections::HashMap<Hash, u64>,
1652 uses: std::collections::HashMap<Hash, usize>,
1653 in_pack: std::collections::HashMap<Hash, Cow<'a, [u8]>>,
1654 report: Option<DecodeReport>,
1655}
1656
1657impl<'a> PackDecodeCursor<'a> {
1658 pub fn new(pack: &'a [u8], limits: DecodeLimits) -> Result<Self, PackError> {
1664 let mut pack_entries = PackEntries::new(pack)?;
1665 let mut budget = DecodeBudget {
1666 used: 0,
1667 max: limits.max_decoded_bytes,
1668 };
1669 charge_compressed_claims(pack, &mut budget)?;
1670
1671 let mut entries = Vec::with_capacity(pack_entries.entry_count());
1672 while let Some(entry) = pack_entries.next() {
1673 let entry = entry?;
1674 let payload = pack_entries
1675 .last_payload_range()
1676 .ok_or(PackError::UnexpectedEof)?;
1677 limits.check_frame(
1678 pack[payload.start - ENTRY_FRAME_LEN],
1679 &pack[payload.clone()],
1680 )?;
1681 let offset = payload
1682 .start
1683 .checked_sub(ENTRY_FRAME_LEN)
1684 .ok_or(PackError::UnexpectedEof)?;
1685 let delta_base = match &entry {
1686 PackEntry::Delta { base, .. } => Some(*base),
1687 PackEntry::Raw { .. } => None,
1688 };
1689 entries.push(Some(CursorFrame {
1690 entry,
1691 frame_offset: offset as u64,
1692 frame_length: (payload.end - offset) as u64,
1693 wire_type: pack[offset],
1694 delta_base,
1695 }));
1696 }
1697
1698 let mut uses = std::collections::HashMap::new();
1701 for frame in entries.iter().flatten() {
1702 if let PackEntry::Delta { base, stream } = &frame.entry {
1703 if let Some(result_len) = le_u32_at(stream, 5) {
1704 budget.charge(result_len)?;
1705 }
1706 let n = uses.entry(*base).or_insert(0usize);
1707 *n = n.saturating_add(1);
1708 }
1709 }
1710 Ok(Self {
1711 report: Some(DecodeReport {
1712 ids: Vec::with_capacity(entries.len()),
1713 ..DecodeReport::default()
1714 }),
1715 entries,
1716 next: 0,
1717 budget,
1718 charged: std::collections::HashMap::new(),
1719 uses,
1720 in_pack: std::collections::HashMap::new(),
1721 })
1722 }
1723
1724 pub fn set_max_decoded_bytes(&mut self, max: u64) -> Result<(), PackError> {
1732 if self.budget.used > max {
1733 return Err(PackError::PackfileTooLarge);
1734 }
1735 self.budget.max = max;
1736 Ok(())
1737 }
1738
1739 #[allow(clippy::too_many_lines)] pub fn resume<B: DeltaBaseSource>(
1747 &mut self,
1748 bases: &mut B,
1749 mut sink: impl FnMut(DecodedEntry<'_>) -> Result<(), PackError>,
1750 ) -> Result<DecodeReport, PackError> {
1751 if self.report.is_none() {
1752 return Err(PackError::PackfileCorrupted);
1753 }
1754 let mut bases = ChargedBases {
1755 inner: bases,
1756 budget: &mut self.budget,
1757 charged: &mut self.charged,
1758 };
1759 while self.next < self.entries.len() {
1760 let frame = self.entries[self.next]
1761 .take()
1762 .ok_or(PackError::PackfileCorrupted)?;
1763 let CursorFrame {
1764 entry,
1765 frame_offset,
1766 frame_length,
1767 wire_type,
1768 delta_base,
1769 } = frame;
1770 match entry {
1771 PackEntry::Raw { bytes } => {
1772 let object = validate_storable_object(&bytes)?;
1773 let id = crate::object::id_from_object(&object, &bytes);
1774 sink(DecodedEntry {
1775 id,
1776 bytes: bytes.as_ref(),
1777 object,
1778 from_delta: false,
1779 frame_offset,
1780 frame_length,
1781 wire_type,
1782 delta_base,
1783 })?;
1784 if self.uses.contains_key(&id) {
1785 self.in_pack.insert(id, bytes);
1786 }
1787 let report = self.report.as_mut().ok_or(PackError::PackfileCorrupted)?;
1788 report.raw_count += 1;
1789 report.ids.push(id);
1790 }
1791 PackEntry::Delta { base, stream } => {
1792 let before_used = bases.budget.used;
1796 let before_charged = bases.charged.get(&base).copied();
1797 let resolved = match resolve_delta_target(
1798 &mut bases,
1799 &mut self.in_pack,
1800 base,
1801 stream.as_ref(),
1802 ) {
1803 Ok(resolved) => resolved,
1804 Err(error @ PackError::DeltaBaseMissing(_)) => {
1805 bases.budget.used = before_used;
1806 match before_charged {
1807 Some(len) => {
1808 bases.charged.insert(base, len);
1809 }
1810 None => {
1811 bases.charged.remove(&base);
1812 }
1813 }
1814 self.entries[self.next] = Some(CursorFrame {
1815 entry: PackEntry::Delta { base, stream },
1816 frame_offset,
1817 frame_length,
1818 wire_type,
1819 delta_base,
1820 });
1821 return Err(error);
1822 }
1823 Err(error) => return Err(error),
1824 };
1825 drop(stream);
1826 if let Some(left) = self.uses.get_mut(&base) {
1827 *left = left.saturating_sub(1);
1828 if *left == 0 {
1829 self.uses.remove(&base);
1830 self.in_pack.remove(&base);
1831 bases.release(&base);
1832 }
1833 }
1834 let object = validate_storable_object(&resolved)?;
1835 let id = crate::object::id_from_object(&object, &resolved);
1836 sink(DecodedEntry {
1837 id,
1838 bytes: &resolved,
1839 object,
1840 from_delta: true,
1841 frame_offset,
1842 frame_length,
1843 wire_type,
1844 delta_base,
1845 })?;
1846 if self.uses.contains_key(&id) {
1847 self.in_pack.insert(id, Cow::Owned(resolved));
1848 }
1849 let report = self.report.as_mut().ok_or(PackError::PackfileCorrupted)?;
1850 report.delta_count += 1;
1851 report.ids.push(id);
1852 }
1853 }
1854 self.next += 1;
1855 }
1856 self.report.take().ok_or(PackError::PackfileCorrupted)
1857 }
1858}
1859
1860pub fn decode_entries_with<B: DeltaBaseSource>(
1896 pack: &[u8],
1897 bases: &mut B,
1898 limits: DecodeLimits,
1899 sink: impl FnMut(DecodedEntry<'_>) -> Result<(), PackError>,
1900) -> Result<DecodeReport, PackError> {
1901 PackDecodeCursor::new(pack, limits)?.resume(bases, sink)
1902}
1903
1904pub fn decode_frame_with<B: DeltaBaseSource>(
1912 frame: &[u8],
1913 version: u32,
1914 bases: &mut B,
1915 limits: DecodeLimits,
1916) -> Result<(Hash, Vec<u8>), PackError> {
1917 if version != VERSION && version != VERSION_V2 {
1918 return Err(PackError::UnsupportedVersion(version));
1919 }
1920 if frame.len() < ENTRY_FRAME_LEN {
1921 return Err(PackError::UnexpectedEof);
1922 }
1923 let payload_len = u32::from_le_bytes(
1924 frame[1..5]
1925 .try_into()
1926 .map_err(|_| PackError::UnexpectedEof)?,
1927 ) as usize;
1928 if Some(frame.len()) != ENTRY_FRAME_LEN.checked_add(payload_len) {
1929 return Err(PackError::UnexpectedEof);
1930 }
1931 let payload = &frame[ENTRY_FRAME_LEN..];
1932 limits.check_frame(frame[0], payload)?;
1933 match frame[0] {
1934 0x00 if payload.len() as u64 > limits.max_decoded_bytes => {
1935 return Err(PackError::PackfileTooLarge);
1936 }
1937 0x03 if zstd_claim(payload)?.0 as u64 > limits.max_decoded_bytes => {
1938 return Err(PackError::PackfileTooLarge);
1939 }
1940 0x04 if payload.len() >= hash::HASH_LEN
1941 && zstd_claim(&payload[hash::HASH_LEN..])?.0 as u64
1942 > limits
1943 .entry_geometry
1944 .map_or(limits.max_decoded_bytes, |(_, stream)| stream) =>
1945 {
1946 return Err(PackError::PackfileTooLarge);
1947 }
1948 _ => {}
1949 }
1950 decode_entry_with(
1951 decode_payload(frame[0], version, &frame[ENTRY_FRAME_LEN..])?,
1952 bases,
1953 limits,
1954 )
1955}
1956
1957pub fn decode_entry_with<B: DeltaBaseSource>(
1965 entry: PackEntry<'_>,
1966 bases: &mut B,
1967 limits: DecodeLimits,
1968) -> Result<(Hash, Vec<u8>), PackError> {
1969 let bytes = match entry {
1970 PackEntry::Raw { bytes } => bytes.into_owned(),
1971 PackEntry::Delta { base, stream } => {
1972 if limits
1973 .entry_geometry
1974 .is_some_and(|(_, cap)| stream.len() as u64 > cap)
1975 || validate_delta_result_size(stream.as_ref())? as u64 > limits.max_decoded_bytes
1976 {
1977 return Err(PackError::PackfileTooLarge);
1978 }
1979 resolve_delta_target(
1980 bases,
1981 &mut std::collections::HashMap::new(),
1982 base,
1983 stream.as_ref(),
1984 )?
1985 }
1986 };
1987 if bytes.len() as u64 > limits.max_decoded_bytes {
1988 return Err(PackError::PackfileTooLarge);
1989 }
1990 let object = validate_storable_object(&bytes)?;
1991 let id = crate::object::id_from_object(&object, &bytes);
1992 Ok((id, bytes))
1993}
1994
1995fn finish_pack_read<'b>(
1997 entries: Vec<Entry<'_>>,
1998 raw_results: RawStageResults<'b>,
1999 mut uses: std::collections::HashMap<Hash, BaseUses>,
2000 budget: &'b ResidentBudget<'_>,
2001 store: &ObjectStore,
2002 batch: crate::batch::WriteBatch<'_>,
2003) -> Result<UnpackReport, PackError> {
2004 let mut bases = store;
2005 let mut raw_results = raw_results.into_iter();
2006 let mut in_pack = std::collections::HashMap::new();
2007 let mut report = UnpackReport::default();
2008 for (position, entry) in entries.into_iter().enumerate() {
2011 match entry {
2012 Entry::Raw(payload) => {
2013 let (stored_hash, retained) = raw_results
2014 .next()
2015 .expect("raw frame result")
2016 .unwrap_or_else(|| {
2017 prepare_and_stage_raw(&batch, position, payload, &uses, budget)
2018 })?;
2019 if has_remaining_uses(&uses, &stored_hash) {
2020 if let Some(bytes) = retained {
2021 in_pack.insert(stored_hash, ResidentBytes::Owned(bytes));
2022 } else if let EncodedPayload::Plain(bytes) = payload {
2023 in_pack.insert(stored_hash, ResidentBytes::Borrowed(bytes));
2024 }
2025 }
2026 report.raw_count += 1;
2027 report.stored.push(stored_hash);
2028 }
2029 Entry::Delta { base, stream } => {
2030 let stream = stream.decode(budget)?;
2031 let stored_hash = stage_delta_target(
2032 &mut bases,
2033 &batch,
2034 &mut in_pack,
2035 &mut uses,
2036 budget,
2037 base,
2038 stream.as_ref(),
2039 )?;
2040 report.delta_count += 1;
2041 report.stored.push(stored_hash);
2042 }
2043 }
2044 }
2045 batch.commit()?;
2046 Ok(report)
2047}
2048
2049fn resident_bytes_cap(pack_len: usize) -> usize {
2055 MAX_RAW_OBJECT_SIZE
2056 .saturating_mul(2)
2057 .max(pack_len.saturating_mul(16))
2058}
2059
2060struct ResidentBudget<'a> {
2064 cap: usize,
2065 used: std::sync::atomic::AtomicUsize,
2066 #[cfg(test)]
2067 peak: std::sync::atomic::AtomicUsize,
2068 owned_bytes: Option<&'a AtomicU64>,
2069}
2070
2071impl<'a> ResidentBudget<'a> {
2072 fn new(cap: usize, owned_bytes: Option<&'a AtomicU64>) -> Self {
2073 Self {
2074 cap,
2075 used: std::sync::atomic::AtomicUsize::new(0),
2076 #[cfg(test)]
2077 peak: std::sync::atomic::AtomicUsize::new(0),
2078 owned_bytes,
2079 }
2080 }
2081
2082 fn charge(&self, len: usize) -> Result<Reservation<'_>, PackError> {
2083 let previous = self
2084 .used
2085 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |used| {
2086 used.checked_add(len).filter(|&total| total <= self.cap)
2087 })
2088 .map_err(|_| PackError::PackfileTooLarge)?;
2089 #[cfg(test)]
2090 self.peak.fetch_max(previous + len, Ordering::Relaxed);
2091 #[cfg(not(test))]
2092 let _ = previous;
2093 Ok(Reservation { budget: self, len })
2094 }
2095
2096 fn allocate(&self, len: usize) -> Result<OwnedBytes<'_>, PackError> {
2097 self.charge(len)?.allocate()
2098 }
2099
2100 fn record_owned(&self, len: usize) {
2101 if let Some(counter) = self.owned_bytes {
2102 counter.fetch_add(len as u64, Ordering::Relaxed);
2103 }
2104 }
2105}
2106
2107struct Reservation<'a> {
2108 budget: &'a ResidentBudget<'a>,
2109 len: usize,
2110}
2111
2112impl<'a> Reservation<'a> {
2113 fn allocate(self) -> Result<OwnedBytes<'a>, PackError> {
2114 let mut bytes = Vec::new();
2115 bytes
2116 .try_reserve_exact(self.len)
2117 .map_err(|_| PackError::PackfileTooLarge)?;
2118 Ok(OwnedBytes {
2119 bytes,
2120 _reservation: self,
2121 })
2122 }
2123}
2124
2125impl Drop for Reservation<'_> {
2126 fn drop(&mut self) {
2127 self.budget.used.fetch_sub(self.len, Ordering::Relaxed);
2128 }
2129}
2130
2131struct OwnedBytes<'a> {
2132 bytes: Vec<u8>,
2133 _reservation: Reservation<'a>,
2134}
2135
2136enum ResidentBytes<'p, 'b> {
2137 Borrowed(&'p [u8]),
2138 Owned(OwnedBytes<'b>),
2139}
2140
2141impl AsRef<[u8]> for ResidentBytes<'_, '_> {
2142 fn as_ref(&self) -> &[u8] {
2143 match self {
2144 Self::Borrowed(bytes) => bytes,
2145 Self::Owned(bytes) => &bytes.bytes,
2146 }
2147 }
2148}
2149
2150#[derive(Clone, Copy)]
2151enum EncodedPayload<'p> {
2152 Plain(&'p [u8]),
2153 Zstd(&'p [u8]),
2154}
2155
2156impl<'p> EncodedPayload<'p> {
2157 fn decode<'b>(
2158 self,
2159 budget: &'b ResidentBudget<'_>,
2160 ) -> Result<ResidentBytes<'p, 'b>, PackError> {
2161 match self {
2162 Self::Plain(bytes) => Ok(ResidentBytes::Borrowed(bytes)),
2163 Self::Zstd(payload) => {
2164 let len = zstd_entry_len(payload)?;
2165 let mut output = budget.allocate(len)?;
2166 decompress_zstd_into(payload, &mut output.bytes)?;
2167 Ok(ResidentBytes::Owned(output))
2168 }
2169 }
2170 }
2171
2172 fn decode_for_staging<'b>(
2174 self,
2175 budget: &'b ResidentBudget<'_>,
2176 ) -> Result<Option<ResidentBytes<'p, 'b>>, PackError> {
2177 match self {
2178 Self::Plain(bytes) => Ok(Some(ResidentBytes::Borrowed(bytes))),
2179 Self::Zstd(payload) => {
2180 let len = zstd_entry_len(payload)?;
2181 let Ok(reservation) = budget.charge(len) else {
2182 return Ok(None);
2183 };
2184 let mut output = reservation.allocate()?;
2185 decompress_zstd_into(payload, &mut output.bytes)?;
2186 Ok(Some(ResidentBytes::Owned(output)))
2187 }
2188 }
2189 }
2190}
2191
2192#[derive(Clone, Copy)]
2193enum Entry<'p> {
2194 Raw(EncodedPayload<'p>),
2195 Delta {
2196 base: Hash,
2197 stream: EncodedPayload<'p>,
2198 },
2199}
2200
2201#[derive(Default)]
2202struct BaseUses {
2203 remaining: usize,
2204 last_position: usize,
2205}
2206
2207fn has_remaining_uses(uses: &std::collections::HashMap<Hash, BaseUses>, hash: &Hash) -> bool {
2208 uses.get(hash).is_some_and(|usage| usage.remaining != 0)
2209}
2210
2211type RawStageResult<'b> = Option<Result<(Hash, Option<OwnedBytes<'b>>), PackError>>;
2213type RawStageResults<'b> = Vec<RawStageResult<'b>>;
2214
2215fn stage_raw_entries<'b>(
2218 batch: &crate::batch::WriteBatch<'_>,
2219 frames: &[(usize, EncodedPayload<'_>)],
2220 uses: &std::collections::HashMap<Hash, BaseUses>,
2221 budget: &'b ResidentBudget<'_>,
2222) -> RawStageResults<'b> {
2223 #[cfg(not(target_arch = "wasm32"))]
2224 {
2225 const ENTRIES_PER_THREAD: usize = 8;
2226 let threads = std::thread::available_parallelism().map_or(1, std::num::NonZeroUsize::get);
2227 if threads > 1 && frames.len() >= ENTRIES_PER_THREAD.saturating_mul(threads) {
2228 return stage_raw_entries_parallel(batch, frames, uses, budget, threads);
2229 }
2230 }
2231 let first_failure = std::sync::atomic::AtomicUsize::new(usize::MAX);
2232 frames
2233 .iter()
2234 .map(|&(position, payload)| {
2235 stage_raw_in_phase_two(batch, position, payload, uses, budget, &first_failure)
2236 })
2237 .collect()
2238}
2239
2240#[cfg(not(target_arch = "wasm32"))]
2241fn stage_raw_entries_parallel<'b>(
2242 batch: &crate::batch::WriteBatch<'_>,
2243 frames: &[(usize, EncodedPayload<'_>)],
2244 uses: &std::collections::HashMap<Hash, BaseUses>,
2245 budget: &'b ResidentBudget<'_>,
2246 threads: usize,
2247) -> RawStageResults<'b> {
2248 let first_failure = &std::sync::atomic::AtomicUsize::new(usize::MAX);
2249 let chunk_size = frames.len().div_ceil(threads).max(1);
2250 let mut out = Vec::with_capacity(frames.len());
2251 std::thread::scope(|scope| {
2252 let handles: Vec<_> = frames
2253 .chunks(chunk_size)
2254 .map(|chunk| {
2255 scope.spawn(move || {
2256 chunk
2257 .iter()
2258 .map(|&(position, payload)| {
2259 stage_raw_in_phase_two(
2260 batch,
2261 position,
2262 payload,
2263 uses,
2264 budget,
2265 first_failure,
2266 )
2267 })
2268 .collect::<Vec<_>>()
2269 })
2270 })
2271 .collect();
2272 for handle in handles {
2273 out.extend(handle.join().expect("pack unpack worker thread panicked"));
2274 }
2275 });
2276 out
2277}
2278
2279fn stage_raw_in_phase_two<'b>(
2280 batch: &crate::batch::WriteBatch<'_>,
2281 position: usize,
2282 payload: EncodedPayload<'_>,
2283 uses: &std::collections::HashMap<Hash, BaseUses>,
2284 budget: &'b ResidentBudget<'_>,
2285 first_failure: &std::sync::atomic::AtomicUsize,
2286) -> RawStageResult<'b> {
2287 if position > first_failure.load(Ordering::Relaxed) {
2288 return None;
2289 }
2290 let result = match payload.decode_for_staging(budget) {
2291 Ok(Some(payload)) => stage_decoded_raw(batch, position, payload, uses, budget),
2292 Ok(None) => return None,
2293 Err(error) => Err(error),
2294 };
2295 if result.is_err() {
2296 first_failure.fetch_min(position, Ordering::Relaxed);
2297 }
2298 Some(result)
2299}
2300
2301fn prepare_and_stage_raw<'b>(
2302 batch: &crate::batch::WriteBatch<'_>,
2303 position: usize,
2304 payload: EncodedPayload<'_>,
2305 uses: &std::collections::HashMap<Hash, BaseUses>,
2306 budget: &'b ResidentBudget<'_>,
2307) -> Result<(Hash, Option<OwnedBytes<'b>>), PackError> {
2308 stage_decoded_raw(batch, position, payload.decode(budget)?, uses, budget)
2309}
2310
2311fn stage_decoded_raw<'b>(
2312 batch: &crate::batch::WriteBatch<'_>,
2313 position: usize,
2314 payload: ResidentBytes<'_, 'b>,
2315 uses: &std::collections::HashMap<Hash, BaseUses>,
2316 budget: &'b ResidentBudget<'_>,
2317) -> Result<(Hash, Option<OwnedBytes<'b>>), PackError> {
2318 let obj = validate_storable_object(payload.as_ref())?;
2319 let stored_hash = crate::object::id_from_object(&obj, payload.as_ref());
2320 batch.write_prehashed(stored_hash, &[payload.as_ref()])?;
2321 let retained = if let ResidentBytes::Owned(bytes) = payload {
2322 budget.record_owned(bytes.bytes.len());
2323 if uses
2324 .get(&stored_hash)
2325 .is_some_and(|usage| usage.last_position > position)
2326 {
2327 Some(bytes)
2328 } else {
2329 None
2330 }
2331 } else {
2332 None
2333 };
2334 Ok((stored_hash, retained))
2335}
2336
2337fn validate_pack_header(pack_bytes: &[u8]) -> Result<(u32, usize, u32), PackError> {
2346 if pack_bytes.len() < HEADER_LEN + TRAILER_LEN {
2348 return Err(PackError::PackfileTooShort);
2349 }
2350 if &pack_bytes[..4] != MAGIC.as_slice() {
2352 return Err(PackError::InvalidMagic);
2353 }
2354 let version = u32::from_le_bytes(pack_bytes[4..8].try_into().expect("4 bytes"));
2357 if version != VERSION && version != VERSION_V2 {
2358 return Err(PackError::UnsupportedVersion(version));
2359 }
2360 let split = pack_bytes.len() - TRAILER_LEN;
2370 let body = &pack_bytes[..split];
2371 let trailer = &pack_bytes[split..];
2372 let computed = hash::hash(body);
2373 if computed.as_slice() != trailer {
2374 return Err(PackError::PackfileCorrupted);
2375 }
2376 let count = u32::from_le_bytes(
2378 pack_bytes[ENTRY_COUNT_OFFSET..ENTRY_COUNT_OFFSET + 4]
2379 .try_into()
2380 .expect("4 bytes"),
2381 );
2382 if count > MAX_ENTRIES {
2383 return Err(PackError::TooManyObjects(count));
2384 }
2385 let body_after_header = body.len() - HEADER_LEN;
2387 if u64::from(count) * ENTRY_FRAME_LEN as u64 > body_after_header as u64 {
2388 return Err(PackError::TooManyObjects(count));
2389 }
2390 Ok((version, split, count))
2391}
2392
2393#[derive(Debug)]
2400pub enum PackEntry<'a> {
2401 Raw { bytes: Cow<'a, [u8]> },
2405 Delta { base: Hash, stream: Cow<'a, [u8]> },
2407}
2408
2409#[derive(Debug)]
2423pub struct PackEntries<'a> {
2424 bytes: &'a [u8],
2425 version: u32,
2426 split: usize,
2427 count: u32,
2428 pos: usize,
2429 yielded: u32,
2430 raw_only: bool,
2431 first_non_raw: Option<u32>,
2432 last_payload_range: Option<Range<usize>>,
2433 done: bool,
2434}
2435
2436impl<'a> PackEntries<'a> {
2437 pub(crate) fn entry_count(&self) -> usize {
2444 self.count as usize
2445 }
2446
2447 pub fn new(bytes: &'a [u8]) -> Result<Self, PackError> {
2455 Self::new_with_payload_cap(bytes, MAX_TOTAL_PAYLOAD)
2456 }
2457
2458 pub(crate) fn new_with_payload_cap(
2459 bytes: &'a [u8],
2460 payload_cap: u64,
2461 ) -> Result<Self, PackError> {
2462 let (version, split, count) = validate_pack_header(bytes)?;
2463 let mut pos = HEADER_LEN;
2464 let mut total_payload: u64 = 0;
2465 let mut raw_only = true;
2466 let mut first_non_raw = None;
2467 for i in 0..count {
2468 if ENTRY_FRAME_LEN > split - pos {
2469 return Err(PackError::UnexpectedEof);
2470 }
2471 let etype = bytes[pos];
2472 pos = pos.checked_add(1).ok_or(PackError::UnexpectedEof)?;
2473 let payload_len = u32::from_le_bytes(
2474 bytes[pos..pos.checked_add(4).ok_or(PackError::UnexpectedEof)?]
2475 .try_into()
2476 .expect("4 bytes"),
2477 ) as usize;
2478 pos = pos.checked_add(4).ok_or(PackError::UnexpectedEof)?;
2479 total_payload = total_payload.saturating_add(payload_len as u64);
2480 if total_payload > payload_cap {
2481 return Err(PackError::PackfileTooLarge);
2482 }
2483 if payload_len > split - pos {
2484 return Err(PackError::UnexpectedEof);
2485 }
2486 match etype {
2487 0x00 => {}
2488 0x02 => {
2489 if payload_len < hash::HASH_LEN {
2490 return Err(PackError::DeltaEntryTruncated);
2491 }
2492 if first_non_raw.is_none() {
2493 first_non_raw = Some(i);
2494 }
2495 raw_only = false;
2496 }
2497 0x03 if version == VERSION_V2 => {
2498 if first_non_raw.is_none() {
2499 first_non_raw = Some(i);
2500 }
2501 raw_only = false;
2502 }
2503 0x04 if version == VERSION_V2 => {
2504 if payload_len < hash::HASH_LEN {
2505 return Err(PackError::DeltaEntryTruncated);
2506 }
2507 if first_non_raw.is_none() {
2508 first_non_raw = Some(i);
2509 }
2510 raw_only = false;
2511 }
2512 0x01 => return Err(PackError::InvalidEntryType(0x01)),
2513 other => return Err(PackError::InvalidEntryType(other)),
2514 }
2515 pos = pos
2516 .checked_add(payload_len)
2517 .ok_or(PackError::UnexpectedEof)?;
2518 }
2519 if pos != split {
2520 return Err(PackError::TrailingData);
2521 }
2522 Ok(Self {
2523 bytes,
2524 version,
2525 split,
2526 count,
2527 pos: HEADER_LEN,
2528 yielded: 0,
2529 raw_only,
2530 first_non_raw,
2531 last_payload_range: None,
2532 done: false,
2533 })
2534 }
2535
2536 #[must_use]
2541 pub fn is_raw_only(&self) -> bool {
2542 self.raw_only
2543 }
2544
2545 #[must_use]
2547 pub fn first_non_raw_index(&self) -> Option<u32> {
2548 self.first_non_raw
2549 }
2550
2551 #[must_use]
2557 pub(crate) fn last_payload_range(&self) -> Option<Range<usize>> {
2558 self.last_payload_range.clone()
2559 }
2560
2561 fn next_encoded_entry(&mut self) -> Result<Entry<'a>, PackError> {
2562 if ENTRY_FRAME_LEN > self.split - self.pos {
2563 return Err(PackError::UnexpectedEof);
2564 }
2565 let etype = self.bytes[self.pos];
2566 self.pos = self.pos.checked_add(1).ok_or(PackError::UnexpectedEof)?;
2567 let payload_len = u32::from_le_bytes(
2568 self.bytes[self.pos..self.pos.checked_add(4).ok_or(PackError::UnexpectedEof)?]
2569 .try_into()
2570 .expect("4 bytes"),
2571 ) as usize;
2572 self.pos = self.pos.checked_add(4).ok_or(PackError::UnexpectedEof)?;
2573 if payload_len > self.split - self.pos {
2574 return Err(PackError::UnexpectedEof);
2575 }
2576 let payload_start = self.pos;
2577 let payload_end = self
2578 .pos
2579 .checked_add(payload_len)
2580 .ok_or(PackError::UnexpectedEof)?;
2581 let payload = &self.bytes[payload_start..payload_end];
2582 self.last_payload_range = Some(payload_start..payload_end);
2583 self.pos = payload_end;
2584 self.yielded += 1;
2585 encoded_payload(etype, self.version, payload)
2586 }
2587
2588 fn next_entry(&mut self) -> Result<PackEntry<'a>, PackError> {
2589 decode_encoded_payload(self.next_encoded_entry()?)
2590 }
2591}
2592
2593fn encoded_payload(etype: u8, version: u32, payload: &[u8]) -> Result<Entry<'_>, PackError> {
2595 match etype {
2596 0x00 => Ok(Entry::Raw(EncodedPayload::Plain(payload))),
2597 0x03 if version == VERSION_V2 => Ok(Entry::Raw(EncodedPayload::Zstd(payload))),
2598 0x02 | 0x04 if etype == 0x02 || version == VERSION_V2 => {
2599 if payload.len() < hash::HASH_LEN {
2600 return Err(PackError::DeltaEntryTruncated);
2601 }
2602 let base = payload[..hash::HASH_LEN].try_into().expect("32 bytes");
2603 let bytes = &payload[hash::HASH_LEN..];
2604 let stream = if etype == 0x02 {
2605 EncodedPayload::Plain(bytes)
2606 } else {
2607 EncodedPayload::Zstd(bytes)
2608 };
2609 Ok(Entry::Delta { base, stream })
2610 }
2611 other => Err(PackError::InvalidEntryType(other)),
2612 }
2613}
2614
2615fn decode_encoded_payload(entry: Entry<'_>) -> Result<PackEntry<'_>, PackError> {
2616 fn decode(payload: EncodedPayload<'_>) -> Result<Cow<'_, [u8]>, PackError> {
2617 match payload {
2618 EncodedPayload::Plain(bytes) => Ok(Cow::Borrowed(bytes)),
2619 EncodedPayload::Zstd(bytes) => Ok(Cow::Owned(decompress_zstd_entry(bytes)?)),
2620 }
2621 }
2622 match entry {
2623 Entry::Raw(bytes) => Ok(PackEntry::Raw {
2624 bytes: decode(bytes)?,
2625 }),
2626 Entry::Delta { base, stream } => Ok(PackEntry::Delta {
2627 base,
2628 stream: decode(stream)?,
2629 }),
2630 }
2631}
2632
2633fn decode_payload(etype: u8, version: u32, payload: &[u8]) -> Result<PackEntry<'_>, PackError> {
2635 decode_encoded_payload(encoded_payload(etype, version, payload)?)
2636}
2637
2638impl<'a> Iterator for PackEntries<'a> {
2639 type Item = Result<PackEntry<'a>, PackError>;
2640
2641 fn next(&mut self) -> Option<Self::Item> {
2642 if self.done || self.yielded >= self.count {
2643 self.done = true;
2644 return None;
2645 }
2646 match self.next_entry() {
2647 Ok(entry) => Some(Ok(entry)),
2648 Err(e) => {
2649 self.done = true;
2650 Some(Err(e))
2651 }
2652 }
2653 }
2654}
2655
2656fn stage_delta_target<'b, B: DeltaBaseSource>(
2660 bases: &mut B,
2661 batch: &crate::batch::WriteBatch<'_>,
2662 in_pack: &mut std::collections::HashMap<Hash, ResidentBytes<'_, 'b>>,
2663 uses: &mut std::collections::HashMap<Hash, BaseUses>,
2664 budget: &'b ResidentBudget<'_>,
2665 base_hash: Hash,
2666 stream: &[u8],
2667) -> Result<Hash, PackError> {
2668 if let std::collections::hash_map::Entry::Vacant(entry) = in_pack.entry(base_hash) {
2669 let mut reservation = None;
2670 let bytes = bases
2671 .base_with_admission(&base_hash, |len| {
2672 reservation = Some(budget.charge(len)?);
2673 Ok(())
2674 })?
2675 .ok_or_else(|| PackError::DeltaBaseMissing(hash::to_hex(&base_hash)))?;
2676 if !external_base_matches::<B>(&bytes, &base_hash)? {
2677 return Err(PackError::DeltaBaseMissing(hash::to_hex(&base_hash)));
2678 }
2679 entry.insert(ResidentBytes::Owned(OwnedBytes {
2680 bytes,
2681 _reservation: reservation.expect("source admitted its buffer"),
2682 }));
2683 }
2684 let result_len = validate_delta_result_size(stream)?;
2685 let mut resolved = budget.allocate(result_len)?;
2686 resolved.bytes =
2687 delta::decode_preallocated(in_pack[&base_hash].as_ref(), stream, resolved.bytes)?;
2688 let obj = validate_storable_object(&resolved.bytes)?;
2689 let stored_hash = crate::object::id_from_object(&obj, &resolved.bytes);
2690 batch.write_prehashed(stored_hash, &[&resolved.bytes])?;
2691 budget.record_owned(resolved.bytes.len());
2692 let usage = uses.get_mut(&base_hash).expect("counted delta base");
2693 usage.remaining -= 1;
2694 if usage.remaining == 0 {
2695 in_pack.remove(&base_hash);
2696 }
2697 if has_remaining_uses(uses, &stored_hash) {
2698 in_pack.insert(stored_hash, ResidentBytes::Owned(resolved));
2699 }
2700 Ok(stored_hash)
2701}
2702
2703fn resolve_delta_target<B: DeltaBaseSource>(
2710 bases: &mut B,
2711 in_pack: &mut std::collections::HashMap<Hash, Cow<'_, [u8]>>,
2712 base_hash: Hash,
2713 stream: &[u8],
2714) -> Result<Vec<u8>, PackError> {
2715 let base_bytes: Cow<'_, [u8]> = if let Some(b) = in_pack.get(&base_hash) {
2729 Cow::Borrowed(b.as_ref())
2730 } else if let Some(bytes) = external_base(bases, &base_hash)? {
2731 in_pack.insert(base_hash, Cow::Owned(bytes.clone()));
2732 Cow::Owned(bytes)
2733 } else {
2734 return Err(PackError::DeltaBaseMissing(hash::to_hex(&base_hash)));
2735 };
2736 validate_delta_result_size(stream)?;
2737 let resolved = delta::decode(base_bytes.as_ref(), stream)?;
2738 Ok(resolved)
2739}
2740
2741fn external_base<B: DeltaBaseSource>(
2752 bases: &mut B,
2753 id: &Hash,
2754) -> Result<Option<Vec<u8>>, PackError> {
2755 let Some(bytes) = bases.base(id)? else {
2756 return Ok(None);
2757 };
2758 Ok(external_base_matches::<B>(&bytes, id)?.then_some(bytes))
2759}
2760
2761fn external_base_matches<B: DeltaBaseSource>(bytes: &[u8], id: &Hash) -> Result<bool, PackError> {
2762 if B::VERIFIED {
2763 validate_storable_object(bytes)?;
2764 return Ok(true);
2765 }
2766 Ok(matches!(validate_storable_object(bytes),
2767 Ok(obj) if crate::object::id_from_object(&obj, bytes) == *id))
2768}
2769
2770fn validate_storable_object(bytes: &[u8]) -> Result<Object, PackError> {
2773 if bytes.len() > MAX_RAW_OBJECT_SIZE {
2774 return Err(PackError::Store(crate::store::StoreError::ObjectTooLarge));
2775 }
2776 match crate::serialize::deserialize(bytes).map_err(PackError::InvalidObject)? {
2777 Object::Delta(_) => Err(PackError::NonStorableObject),
2778 obj @ (Object::Blob(_)
2779 | Object::Tree(_)
2780 | Object::Commit(_)
2781 | Object::Remix(_)
2782 | Object::ChunkedBlob(_)
2783 | Object::Tag(_)) => Ok(obj),
2784 }
2785}
2786
2787fn validate_delta_result_size(stream: &[u8]) -> Result<usize, PackError> {
2788 if stream.len() < delta::HEADER_LEN {
2789 return Err(PackError::DeltaApply(MkitError::UnexpectedEof));
2790 }
2791 let result_len = u32::from_le_bytes(stream[5..9].try_into().expect("4 bytes")) as usize;
2792 if result_len > MAX_RAW_OBJECT_SIZE {
2793 return Err(PackError::Store(crate::store::StoreError::ObjectTooLarge));
2794 }
2795 Ok(result_len)
2796}
2797
2798#[cfg(test)]
2803mod zstd_tests;
2804
2805#[cfg(test)]
2806mod tests {
2807 use super::*;
2808 use tempfile::TempDir;
2809
2810 fn fresh_store() -> (TempDir, ObjectStore) {
2811 let dir = TempDir::new().unwrap();
2812 let store = ObjectStore::init(&crate::layout::RepoLayout::single(dir.path())).unwrap();
2813 (dir, store)
2814 }
2815
2816 fn write_blob_via_serialize(payload: &[u8]) -> Vec<u8> {
2817 let blob = crate::object::Object::Blob(crate::object::Blob {
2821 data: payload.to_vec(),
2822 });
2823 crate::serialize::serialize(&blob).expect("serialize blob")
2824 }
2825
2826 fn finish_pack_body(mut body: Vec<u8>) -> Vec<u8> {
2827 let trailer = hash::hash(&body);
2828 body.extend_from_slice(&trailer);
2829 body
2830 }
2831
2832 fn incompressible_bytes(seed: u64, len: usize) -> Vec<u8> {
2845 let mut buf = vec![0u8; len];
2846 let mut state = seed | 1; for chunk in buf.chunks_mut(8) {
2848 state = state
2849 .wrapping_mul(6_364_136_223_846_793_005)
2850 .wrapping_add(1_442_695_040_888_963_407);
2851 let bytes = state.to_le_bytes();
2852 chunk.copy_from_slice(&bytes[..chunk.len()]);
2853 }
2854 buf
2855 }
2856
2857 #[test]
2858 fn unreferenced_delta_targets_do_not_accumulate() {
2859 let base = write_blob_via_serialize(&vec![b'a'; 64 * 1024]);
2860 let base_hash = hash::hash(&base);
2861 let mut writer = PackWriter::new();
2862 writer.push_raw(base_hash, &base).unwrap();
2863 for i in 0..8 {
2864 let mut content = vec![b'a'; 256 * 1024];
2865 content[0] = b'b' + i;
2866 let target = write_blob_via_serialize(&content);
2867 writer
2868 .push_delta(&base_hash, &delta::encode(&base, &target).unwrap())
2869 .unwrap();
2870 }
2871 let pack = writer.finish().unwrap();
2872 let (_dir, store) = fresh_store();
2873 let budget = ResidentBudget::new(resident_bytes_cap(pack.len()), None);
2874 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget).unwrap();
2875 assert!(budget.peak.load(Ordering::Relaxed) <= 512 * 1024);
2876 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
2877 }
2878
2879 #[test]
2880 fn maximum_wire_lengths_return_framing_errors() {
2881 for len in [u32::MAX, u32::MAX - 4] {
2882 let mut body = Vec::from(MAGIC.as_slice());
2883 body.extend_from_slice(&VERSION.to_le_bytes());
2884 body.extend_from_slice(&1u32.to_le_bytes());
2885 body.push(0x00);
2886 body.extend_from_slice(&len.to_le_bytes());
2887 let pack = finish_pack_body(body);
2888 assert!(matches!(
2889 PackEntries::new(&pack),
2890 Err(PackError::UnexpectedEof)
2891 ));
2892 assert!(matches!(
2893 delta_base_hashes(&pack),
2894 Err(PackError::UnexpectedEof)
2895 ));
2896 let (_dir, store) = fresh_store();
2897 assert!(matches!(
2898 PackReader::read(&pack, &store),
2899 Err(PackError::UnexpectedEof)
2900 ));
2901 let mut parser = PackEntries {
2904 bytes: &pack,
2905 version: VERSION,
2906 split: pack.len() - TRAILER_LEN,
2907 count: 1,
2908 pos: HEADER_LEN,
2909 yielded: 0,
2910 raw_only: true,
2911 first_non_raw: None,
2912 last_payload_range: None,
2913 done: false,
2914 };
2915 assert!(matches!(parser.next(), Some(Err(PackError::UnexpectedEof))));
2916 }
2917 }
2918
2919 #[test]
2920 #[cfg(any(feature = "pack-zstd", feature = "pack-ruzstd"))]
2921 fn junk_zstd_claims_fail_without_zero_filling() {
2922 for etype in [0x03, 0x04] {
2923 for count in [16u32, 64, 256] {
2924 let mut body = Vec::from(MAGIC.as_slice());
2925 body.extend_from_slice(&VERSION_V2.to_le_bytes());
2926 body.extend_from_slice(&count.to_le_bytes());
2927 for _ in 0..count {
2928 body.push(etype);
2929 let base_len = if etype == 0x04 { hash::HASH_LEN } else { 0 };
2930 body.extend_from_slice(&u32::try_from(base_len + 5).unwrap().to_le_bytes());
2931 if etype == 0x04 {
2932 body.extend_from_slice(&[0; hash::HASH_LEN]);
2933 }
2934 body.extend_from_slice(
2935 &u32::try_from(MAX_RAW_OBJECT_SIZE).unwrap().to_le_bytes(),
2936 );
2937 body.push(0xAA);
2938 }
2939 let pack = finish_pack_body(body);
2940 let (_dir, store) = fresh_store();
2941 let started = std::time::Instant::now();
2942 assert!(matches!(
2943 PackReader::read(&pack, &store),
2944 Err(PackError::ZstdDecompress(_))
2945 ));
2946 let elapsed = started.elapsed();
2947 println!("junk zstd type=0x{etype:02x} N={count}: {elapsed:?}");
2948 assert!(
2949 elapsed < std::time::Duration::from_secs(5),
2950 "junk frames must not initialize the claimed buffers: {elapsed:?}"
2951 );
2952 assert!(store.iter_object_hashes().unwrap().is_empty());
2953 }
2954 }
2955 }
2956
2957 #[test]
2958 fn resident_and_decode_limits_are_independent() {
2959 let base = write_blob_via_serialize(b"base payload");
2960 let target = write_blob_via_serialize(b"target payload");
2961 let mut writer = PackWriter::new();
2962 writer.push_raw(hash::hash(&base), &base).unwrap();
2963 let stream = delta::encode(&base, &target).unwrap();
2964 for _ in 0..3 {
2965 writer.push_delta(&hash::hash(&base), &stream).unwrap();
2966 }
2967 let pack = writer.finish().unwrap();
2968 let (_dir, store) = fresh_store();
2969 let resident = ResidentBudget::new(target.len(), None);
2970 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &resident).unwrap();
2971 assert_eq!(resident.used.load(Ordering::Relaxed), 0);
2972 let low = DecodeLimits::default().with_max_decoded_bytes(target.len() as u64);
2975 let mut seen = 0;
2976 assert!(matches!(
2977 decode_entries_with(&pack, &mut NoExternalBases, low, |_| {
2978 seen += 1;
2979 Ok(())
2980 }),
2981 Err(PackError::PackfileTooLarge)
2982 ));
2983 assert_eq!(seen, 0);
2984 let high = DecodeLimits::default().with_max_decoded_bytes((3 * target.len()) as u64);
2985 decode_entries_with(&pack, &mut NoExternalBases, high, |_| Ok(())).unwrap();
2986 let (_dir, store) = fresh_store();
2987 let resident = ResidentBudget::new(target.len() - 1, None);
2988 assert!(matches!(
2989 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &resident),
2990 Err(PackError::PackfileTooLarge)
2991 ));
2992 assert_eq!(resident.peak.load(Ordering::Relaxed), 0);
2993 assert!(store.iter_object_hashes().unwrap().is_empty());
2994 }
2995
2996 #[test]
2997 fn external_base_admission_charges_both_budgets() {
2998 let (_dir, store) = fresh_store();
2999 let bytes = write_blob_via_serialize(b"external base");
3000 let id = store.write(&bytes).unwrap();
3001 for (decoded_cap, resident_cap) in [
3002 (bytes.len() - 1, bytes.len()),
3003 (bytes.len(), bytes.len() - 1),
3004 (bytes.len(), bytes.len()),
3005 ] {
3006 let mut source = &store;
3007 let mut budget = DecodeBudget {
3008 used: 0,
3009 max: decoded_cap as u64,
3010 };
3011 let mut charged = std::collections::HashMap::new();
3012 let mut bases = ChargedBases {
3013 inner: &mut source,
3014 budget: &mut budget,
3015 charged: &mut charged,
3016 };
3017 let resident = ResidentBudget::new(resident_cap, None);
3018 let mut reservation = None;
3019 let result = bases.base_with_admission(&id, |len| {
3020 reservation = Some(resident.charge(len)?);
3021 Ok(())
3022 });
3023 if decoded_cap < bytes.len() || resident_cap < bytes.len() {
3024 assert!(matches!(result, Err(PackError::PackfileTooLarge)));
3025 assert!(reservation.is_none());
3026 assert_eq!(resident.peak.load(Ordering::Relaxed), 0);
3027 } else {
3028 assert_eq!(result.unwrap(), Some(bytes.clone()));
3029 assert_eq!(bases.budget.used, bytes.len() as u64);
3030 assert_eq!(resident.used.load(Ordering::Relaxed), bytes.len());
3031 bases.release(&id);
3032 drop(reservation);
3033 assert_eq!(bases.budget.used, 0);
3034 assert_eq!(resident.used.load(Ordering::Relaxed), 0);
3035 }
3036 }
3037 }
3038
3039 #[test]
3040 fn provided_base_admission_preserves_existing_sources() {
3041 struct Source(Vec<u8>);
3042 impl DeltaBaseSource for Source {
3043 fn base(&mut self, _id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
3044 Ok(Some(self.0.clone()))
3045 }
3046 }
3047 let mut source = Source(vec![1, 2, 3]);
3048 let mut admitted = None;
3049 assert_eq!(
3050 source
3051 .base_with_admission(&[0; 32], |len| {
3052 admitted = Some(len);
3053 Ok(())
3054 })
3055 .unwrap(),
3056 Some(vec![1, 2, 3])
3057 );
3058 assert_eq!(admitted, Some(3));
3059 assert!(matches!(
3060 source.base_with_admission(&[0; 32], |_| Err(PackError::PackfileTooLarge)),
3061 Err(PackError::PackfileTooLarge)
3062 ));
3063 }
3064
3065 #[test]
3066 #[cfg(all(feature = "pack-zstd", not(target_arch = "wasm32")))]
3067 fn phase_two_budget_contention_retries_sequentially() {
3068 let bytes = write_blob_via_serialize(&vec![b'a'; 256 * 1024]);
3069 let mut writer = PackWriter::new();
3070 for _ in 0..16 {
3071 writer.push_raw(hash::hash(&bytes), &bytes).unwrap();
3072 }
3073 let pack = writer.finish().unwrap();
3074 let mut parser = PackEntries::new(&pack).unwrap();
3075 let frames: Vec<_> = (0..parser.entry_count())
3076 .map(|position| match parser.next_encoded_entry().unwrap() {
3077 Entry::Raw(payload @ EncodedPayload::Zstd(_)) => (position, payload),
3078 _ => panic!("expected compressed raw frame"),
3079 })
3080 .collect();
3081 for threads in [1, 2, 4] {
3082 let (_dir, store) = fresh_store();
3083 let batch = store.batch();
3084 let budget = ResidentBudget::new(bytes.len(), None);
3085 let uses = std::collections::HashMap::new();
3086 let in_flight = budget.allocate(bytes.len()).unwrap();
3089 let first_failure = std::sync::atomic::AtomicUsize::new(usize::MAX);
3090 assert!(
3091 stage_raw_in_phase_two(
3092 &batch,
3093 frames[0].0,
3094 frames[0].1,
3095 &uses,
3096 &budget,
3097 &first_failure
3098 )
3099 .is_none()
3100 );
3101 assert_eq!(first_failure.load(Ordering::Relaxed), usize::MAX);
3102 let results = stage_raw_entries_parallel(&batch, &frames, &uses, &budget, threads);
3103 assert_eq!(results.len(), frames.len());
3104 assert!(results.iter().all(Option::is_none));
3105 drop(in_flight);
3106
3107 let entries = frames
3108 .iter()
3109 .map(|(_, payload)| Entry::Raw(*payload))
3110 .collect();
3111 let report = finish_pack_read(entries, results, uses, &budget, &store, batch).unwrap();
3112 assert_eq!(report.raw_count, 16);
3113 assert_eq!(report.delta_count, 0);
3114 assert_eq!(report.stored, vec![hash::hash(&bytes); 16]);
3115 assert_eq!(store.read(&hash::hash(&bytes)).unwrap(), bytes);
3116 assert_eq!(budget.peak.load(Ordering::Relaxed), bytes.len());
3117 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3118 }
3119 }
3120
3121 #[test]
3122 fn phase_two_skips_after_the_first_permanent_failure() {
3123 let (_dir, store) = fresh_store();
3124 let batch = store.batch();
3125 let uses = std::collections::HashMap::new();
3126 let budget = ResidentBudget::new(1024, None);
3127 let first_failure = std::sync::atomic::AtomicUsize::new(usize::MAX);
3128 assert!(matches!(
3129 stage_raw_in_phase_two(
3130 &batch,
3131 7,
3132 EncodedPayload::Plain(b"garbage"),
3133 &uses,
3134 &budget,
3135 &first_failure
3136 ),
3137 Some(Err(PackError::InvalidObject(_)))
3138 ));
3139 let invalid_claim = u32::MAX.to_le_bytes();
3140 assert!(
3141 stage_raw_in_phase_two(
3142 &batch,
3143 19,
3144 EncodedPayload::Zstd(&invalid_claim),
3145 &uses,
3146 &budget,
3147 &first_failure
3148 )
3149 .is_none()
3150 );
3151 assert_eq!(budget.peak.load(Ordering::Relaxed), 0);
3152 let valid = write_blob_via_serialize(b"earlier pack position");
3153 assert!(matches!(
3154 stage_raw_in_phase_two(
3155 &batch,
3156 3,
3157 EncodedPayload::Plain(&valid),
3158 &uses,
3159 &budget,
3160 &first_failure
3161 ),
3162 Some(Ok(_))
3163 ));
3164 assert_eq!(first_failure.load(Ordering::Relaxed), 7);
3165 for _ in 0..2 {
3166 assert!(matches!(
3167 stage_raw_in_phase_two(
3168 &batch,
3169 2,
3170 EncodedPayload::Plain(b"garbage"),
3171 &uses,
3172 &budget,
3173 &first_failure
3174 ),
3175 Some(Err(PackError::InvalidObject(_)))
3176 ));
3177 assert_eq!(first_failure.load(Ordering::Relaxed), 2);
3178 }
3179 }
3180
3181 #[test]
3182 #[cfg(feature = "pack-zstd")]
3183 fn deferred_raw_error_precedes_a_later_staging_failure() {
3184 let mut payload = 7u32.to_le_bytes().to_vec();
3185 payload.extend_from_slice(&zstd::bulk::compress(b"garbage", 3).unwrap());
3186 let over_cap = u32::try_from(MAX_RAW_OBJECT_SIZE + 1)
3187 .unwrap()
3188 .to_le_bytes();
3189 let (_dir, store) = fresh_store();
3190 let batch = store.batch();
3191 let uses = std::collections::HashMap::new();
3192 let budget = ResidentBudget::new(7, None);
3193 let first_failure = std::sync::atomic::AtomicUsize::new(usize::MAX);
3194 let in_flight = budget.allocate(7).unwrap();
3195 let earlier = EncodedPayload::Zstd(&payload);
3196 let later = EncodedPayload::Zstd(&over_cap);
3197 let deferred = stage_raw_in_phase_two(&batch, 0, earlier, &uses, &budget, &first_failure);
3198 assert!(deferred.is_none());
3199 let failed = stage_raw_in_phase_two(&batch, 1, later, &uses, &budget, &first_failure);
3200 assert!(matches!(
3201 failed,
3202 Some(Err(PackError::DecompressedSizeOverCap(_)))
3203 ));
3204 assert_eq!(first_failure.load(Ordering::Relaxed), 1);
3205 drop(in_flight);
3206 assert!(matches!(
3207 finish_pack_read(
3208 vec![Entry::Raw(earlier), Entry::Raw(later)],
3209 vec![deferred, failed],
3210 uses,
3211 &budget,
3212 &store,
3213 batch
3214 ),
3215 Err(PackError::InvalidObject(MkitError::InvalidObjectType(103)))
3216 ));
3217 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3218 assert!(!store.contains(&hash::hash(b"garbage")));
3219 }
3220
3221 #[test]
3222 fn resident_cap_saturates_and_allocation_failure_is_an_error() {
3223 assert_eq!(resident_bytes_cap(0), 2 * MAX_RAW_OBJECT_SIZE);
3224 assert_eq!(
3225 resident_bytes_cap(MAX_RAW_OBJECT_SIZE),
3226 MAX_RAW_OBJECT_SIZE.saturating_mul(16)
3227 );
3228 assert_eq!(resident_bytes_cap(usize::MAX), usize::MAX);
3229 let budget = ResidentBudget::new(usize::MAX, None);
3230 assert!(matches!(
3231 budget.allocate(usize::MAX),
3232 Err(PackError::PackfileTooLarge)
3233 ));
3234 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3235 let _full = budget.charge(usize::MAX).unwrap();
3236 assert!(matches!(budget.charge(1), Err(PackError::PackfileTooLarge)));
3237 }
3238
3239 #[cfg(feature = "pack-zstd")]
3242 fn repeated_blob_delta(base_len: usize, target_len: usize, marker: u8) -> (Vec<u8>, Hash) {
3243 let prologue = crate::serialize::blob_prologue(target_len - 10).unwrap();
3244 let mut stream = vec![delta::STREAM_VERSION];
3245 stream.extend_from_slice(&u32::try_from(base_len).unwrap().to_le_bytes());
3246 stream.extend_from_slice(&u32::try_from(target_len).unwrap().to_le_bytes());
3247 stream.push(10);
3248 stream.extend_from_slice(&prologue);
3249 let mut remaining = target_len - prologue.len() - 1;
3250 let mut hasher = blake3::Hasher::new();
3251 hasher.update(&prologue);
3252 let block = vec![b'a'; usize::from(u16::MAX)];
3253 while remaining != 0 {
3254 let len = remaining.min(block.len());
3255 stream.push(0x80);
3256 stream.extend_from_slice(&10u32.to_le_bytes());
3257 stream.extend_from_slice(&u16::try_from(len).unwrap().to_le_bytes());
3258 hasher.update(&block[..len]);
3259 remaining -= len;
3260 }
3261 stream.extend_from_slice(&[1, marker]);
3262 hasher.update(&[marker]);
3263 (stream, *hasher.finalize().as_bytes())
3264 }
3265
3266 #[test]
3267 #[cfg(feature = "pack-zstd")]
3268 fn compressed_delta_bomb_releases_targets_and_checks_cap_before_allocation() {
3269 const TARGET_LEN: usize = 128 * 1024 * 1024;
3270 let base = write_blob_via_serialize(&vec![b'a'; 64 * 1024]);
3271 let base_hash = hash::hash(&base);
3272 for consume_targets in [false, true] {
3273 let mut writer = PackWriter::new();
3274 writer.push_raw(base_hash, &base).unwrap();
3275 let mut expected = vec![base_hash];
3276 for marker in 0..3 {
3277 let (stream, target_hash) = repeated_blob_delta(base.len(), TARGET_LEN, marker);
3278 writer.push_delta(&base_hash, &stream).unwrap();
3279 expected.push(target_hash);
3280 if consume_targets {
3281 let small = write_blob_via_serialize(&[marker]);
3284 let mut stream = vec![delta::STREAM_VERSION];
3285 stream.extend_from_slice(&u32::try_from(TARGET_LEN).unwrap().to_le_bytes());
3286 stream.extend_from_slice(&u32::try_from(small.len()).unwrap().to_le_bytes());
3287 stream.push(u8::try_from(small.len()).unwrap());
3288 stream.extend_from_slice(&small);
3289 writer.push_delta(&target_hash, &stream).unwrap();
3290 expected.push(hash::hash(&small));
3291 }
3292 }
3293 let pack = writer.finish().unwrap();
3294 assert!(
3295 pack.len() < 2048,
3296 "fixture must remain a small compressed pack"
3297 );
3298 let (_dir, store) = fresh_store();
3299 let budget = ResidentBudget::new(resident_bytes_cap(pack.len()), None);
3300 let report =
3301 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget).unwrap();
3302 assert_eq!(report.stored, expected);
3303 assert_eq!(report.raw_count, 1);
3304 assert_eq!(report.delta_count, if consume_targets { 6 } else { 3 });
3305 assert!(budget.peak.load(Ordering::Relaxed) < TARGET_LEN + 128 * 1024);
3306 assert!(budget.peak.load(Ordering::Relaxed) <= budget.cap);
3307 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3308 for hash in expected {
3309 assert!(store.contains(&hash));
3310 }
3311 let (_dir, store) = fresh_store();
3314 let budget = ResidentBudget::new(TARGET_LEN - 1, None);
3315 assert!(matches!(
3316 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget),
3317 Err(PackError::PackfileTooLarge)
3318 ));
3319 assert!(
3320 budget.peak.load(Ordering::Relaxed) < 128 * 1024,
3321 "target allocation was not admitted"
3322 );
3323 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3324 assert!(!store.contains(&base_hash));
3325 }
3326 }
3327
3328 #[test]
3329 #[cfg(feature = "pack-zstd")]
3330 fn compressed_raw_bomb_is_bounded_by_workers_and_retained_bases() {
3331 const CONTENT_LEN: usize = 4 * 1024 * 1024;
3332 let threads = std::thread::available_parallelism().map_or(1, std::num::NonZeroUsize::get);
3333 let count = 8 * threads;
3334 let mut writer = PackWriter::new();
3335 let mut content = vec![b'a'; CONTENT_LEN];
3336 for i in 0..count {
3337 content[..8].copy_from_slice(&(i as u64).to_le_bytes());
3338 let blob = write_blob_via_serialize(&content);
3339 writer.push_raw(hash::hash(&blob), &blob).unwrap();
3340 }
3341 let pack = writer.finish().unwrap();
3342 assert!(pack.len() < count * 1024);
3343 let (_dir, store) = fresh_store();
3344 let budget = ResidentBudget::new(resident_bytes_cap(pack.len()), None);
3345 let report =
3346 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget).unwrap();
3347 assert_eq!(report.raw_count as usize, count);
3348 assert!(budget.peak.load(Ordering::Relaxed) <= threads * (CONTENT_LEN + 10));
3349 assert!(budget.peak.load(Ordering::Relaxed) <= budget.cap);
3350 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3351 let (_dir, store) = fresh_store();
3352 let budget = ResidentBudget::new(CONTENT_LEN - 1, None);
3353 assert!(matches!(
3354 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget),
3355 Err(PackError::PackfileTooLarge)
3356 ));
3357 assert_eq!(
3358 budget.peak.load(Ordering::Relaxed),
3359 0,
3360 "zstd claim must be rejected before allocation"
3361 );
3362 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3363 }
3364
3365 fn parent_reader(pack: &[u8], store: &ObjectStore) -> Result<UnpackReport, PackError> {
3368 let entries: Vec<_> = PackEntries::new(pack)?.collect::<Result<_, _>>()?;
3369 let batch = store.batch();
3370 let mut in_pack: std::collections::HashMap<Hash, Cow<'_, [u8]>> =
3371 std::collections::HashMap::new();
3372 let mut report = UnpackReport::default();
3373 for entry in entries {
3374 let (bytes, is_delta) = match entry {
3375 PackEntry::Raw { bytes } => (bytes, false),
3376 PackEntry::Delta { base, stream } => {
3377 if let std::collections::hash_map::Entry::Vacant(entry) = in_pack.entry(base) {
3378 if !store.contains(&base) {
3379 return Err(PackError::DeltaBaseMissing(hash::to_hex(&base)));
3380 }
3381 let bytes = store.read(&base)?;
3382 validate_storable_object(&bytes)?;
3383 entry.insert(Cow::Owned(bytes));
3384 }
3385 validate_delta_result_size(&stream)?;
3386 (
3387 Cow::Owned(delta::decode(in_pack[&base].as_ref(), &stream)?),
3388 true,
3389 )
3390 }
3391 };
3392 let obj = validate_storable_object(&bytes)?;
3393 let hash = crate::object::id_from_object(&obj, &bytes);
3394 batch.write_prehashed(hash, &[bytes.as_ref()])?;
3395 in_pack.insert(hash, bytes);
3396 if is_delta {
3397 report.delta_count += 1;
3398 } else {
3399 report.raw_count += 1;
3400 }
3401 report.stored.push(hash);
3402 }
3403 batch.commit()?;
3404 Ok(report)
3405 }
3406
3407 #[test]
3408 fn retention_matches_parent_for_chains_duplicates_and_shared_bases() {
3409 let raw_base = write_blob_via_serialize(&vec![b'a'; 1024]);
3410 let first_target = write_blob_via_serialize(&vec![b'b'; 1024]);
3411 let second_target = write_blob_via_serialize(&vec![b'c'; 1024]);
3412 let shared_target = write_blob_via_serialize(&vec![b'd'; 1024]);
3413 let external = write_blob_via_serialize(b"external base");
3414 let object_hash = |bytes: &[u8]| hash::hash(bytes);
3415 let mut writer = PackWriter::new();
3416 writer.push_raw(object_hash(&raw_base), &raw_base).unwrap();
3417 writer
3418 .push_delta(
3419 &object_hash(&raw_base),
3420 &delta::encode(&raw_base, &first_target).unwrap(),
3421 )
3422 .unwrap();
3423 writer
3424 .push_delta(
3425 &object_hash(&raw_base),
3426 &delta::encode(&raw_base, &second_target).unwrap(),
3427 )
3428 .unwrap();
3429 writer
3430 .push_delta(
3431 &object_hash(&first_target),
3432 &delta::encode(&first_target, &shared_target).unwrap(),
3433 )
3434 .unwrap();
3435 writer.push_raw(object_hash(&raw_base), &raw_base).unwrap();
3437 writer
3439 .push_delta(
3440 &object_hash(&second_target),
3441 &delta::encode(&second_target, &shared_target).unwrap(),
3442 )
3443 .unwrap();
3444 writer
3445 .push_delta(
3446 &object_hash(&shared_target),
3447 &delta::encode(&shared_target, &shared_target).unwrap(),
3448 )
3449 .unwrap();
3450 writer
3451 .push_delta(
3452 &object_hash(&shared_target),
3453 &delta::encode(&shared_target, &first_target).unwrap(),
3454 )
3455 .unwrap();
3456 writer
3457 .push_delta(
3458 &object_hash(&external),
3459 &delta::encode(&external, &raw_base).unwrap(),
3460 )
3461 .unwrap();
3462 writer
3463 .push_delta(
3464 &object_hash(&external),
3465 &delta::encode(&external, &second_target).unwrap(),
3466 )
3467 .unwrap();
3468 let pack = writer.finish().unwrap();
3469 let (_old_dir, old_store) = fresh_store();
3470 let (_new_dir, new_store) = fresh_store();
3471 old_store.write(&external).unwrap();
3472 new_store.write(&external).unwrap();
3473 let old = parent_reader(&pack, &old_store).unwrap();
3474 let budget = ResidentBudget::new(resident_bytes_cap(pack.len()), None);
3475 let new =
3476 PackReader::read_with_budget(&pack, &new_store, MAX_TOTAL_PAYLOAD, &budget).unwrap();
3477 assert_eq!(new, old);
3478 assert_eq!(new.stored.len(), 10);
3479 assert_eq!(new_store.read_call_count(), 1);
3480 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3481 for bytes in [
3482 raw_base,
3483 first_target,
3484 second_target,
3485 shared_target,
3486 external,
3487 ] {
3488 assert_eq!(
3489 new_store.read(&object_hash(&bytes)).unwrap(),
3490 old_store.read(&object_hash(&bytes)).unwrap()
3491 );
3492 }
3493 }
3494
3495 #[test]
3496 fn store_base_is_charged_once_and_released_on_last_use() {
3497 let (_dir, store) = fresh_store();
3498 let base = write_blob_via_serialize(b"external base payload");
3499 let target = write_blob_via_serialize(b"target payload");
3500 let base_hash = store.write(&base).unwrap();
3501 let target_hash = hash::hash(&target);
3502 let stream = delta::encode(&base, &target).unwrap();
3503 let budget = ResidentBudget::new(base.len() + target.len(), None);
3504 let batch = store.batch();
3505 let mut in_pack = std::collections::HashMap::new();
3506 let mut uses = std::collections::HashMap::from([(
3507 base_hash,
3508 BaseUses {
3509 remaining: 2,
3510 last_position: 1,
3511 },
3512 )]);
3513 for remaining in [1, 0] {
3514 assert_eq!(
3515 stage_delta_target(
3516 &mut &store,
3517 &batch,
3518 &mut in_pack,
3519 &mut uses,
3520 &budget,
3521 base_hash,
3522 &stream
3523 )
3524 .unwrap(),
3525 target_hash
3526 );
3527 assert_eq!(uses[&base_hash].remaining, remaining);
3528 assert_eq!(in_pack.contains_key(&base_hash), remaining != 0);
3529 assert!(!in_pack.contains_key(&target_hash));
3530 assert_eq!(
3531 budget.used.load(Ordering::Relaxed),
3532 if remaining == 0 { 0 } else { base.len() }
3533 );
3534 }
3535 assert_eq!(store.read_call_count(), 1);
3536 assert_eq!(budget.peak.load(Ordering::Relaxed), budget.cap);
3537 drop(in_pack);
3540 let budget = ResidentBudget::new(base.len() - 1, None);
3541 let mut in_pack = std::collections::HashMap::new();
3542 let mut uses = std::collections::HashMap::from([(
3543 base_hash,
3544 BaseUses {
3545 remaining: 1,
3546 last_position: 0,
3547 },
3548 )]);
3549 assert!(matches!(
3550 stage_delta_target(
3551 &mut &store,
3552 &batch,
3553 &mut in_pack,
3554 &mut uses,
3555 &budget,
3556 base_hash,
3557 &stream
3558 ),
3559 Err(PackError::PackfileTooLarge)
3560 ));
3561 assert!(in_pack.is_empty());
3562 assert_eq!(budget.peak.load(Ordering::Relaxed), 0);
3563 }
3564
3565 #[test]
3566 #[cfg(feature = "pack-zstd")]
3567 fn retained_compressed_raw_bases_share_the_resident_cap() {
3568 let first = write_blob_via_serialize(&vec![b'a'; 1024 * 1024]);
3569 let second = write_blob_via_serialize(&vec![b'b'; 1024 * 1024]);
3570 let mut writer = PackWriter::new();
3571 for bytes in [&first, &second] {
3572 writer.push_raw(hash::hash(bytes), bytes).unwrap();
3573 }
3574 for bytes in [&first, &second] {
3575 writer
3576 .push_delta(&hash::hash(bytes), &delta::encode(bytes, bytes).unwrap())
3577 .unwrap();
3578 }
3579 let pack = writer.finish().unwrap();
3580 let (_dir, store) = fresh_store();
3581 let budget = ResidentBudget::new(first.len() + second.len() - 1, None);
3582 assert!(matches!(
3583 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget),
3584 Err(PackError::PackfileTooLarge)
3585 ));
3586 assert_eq!(budget.peak.load(Ordering::Relaxed), first.len());
3587 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3588 assert!(!store.contains(&hash::hash(&first)));
3589 assert!(!store.contains(&hash::hash(&second)));
3590 }
3591
3592 #[test]
3593 #[ignore = "decodes more than 200 MiB; run in the serial ignored-lane"]
3594 #[cfg(feature = "pack-zstd")]
3595 fn large_mixed_pack_decodes_under_production_resident_cap() {
3596 const GROUPS: u32 = 13;
3597 const RAW_LEN: usize = 16 * 1024 * 1024;
3598 const COMPRESSED_LEN: usize = 4 * 1024 * 1024;
3599 let started = std::time::Instant::now();
3600 let mut writer = PackWriter::new();
3601 let mut expected = Vec::new();
3602 for group in 0..GROUPS {
3605 for content in [
3606 incompressible_bytes(0xA000_0000 + u64::from(group), RAW_LEN),
3607 vec![u8::try_from(group).unwrap(); COMPRESSED_LEN],
3608 ] {
3609 let base = write_blob_via_serialize(&content);
3610 let base_hash = hash::hash(&base);
3611 writer.push_raw(base_hash, &base).unwrap();
3612 expected.push(base_hash);
3613 let mut target = base.clone();
3614 *target.last_mut().unwrap() ^= 0x80;
3615 let target_hash = hash::hash(&target);
3616 writer
3617 .push_delta(&base_hash, &delta::encode(&base, &target).unwrap())
3618 .unwrap();
3619 expected.push(target_hash);
3620 }
3621 }
3622 let pack = writer.finish().unwrap();
3623 assert!(pack.len() >= 200 * 1024 * 1024);
3626 let mut parser = PackEntries::new(&pack).unwrap();
3627 let (mut raw, mut zstd, mut deltas) = (0, 0, 0);
3628 for _ in 0..parser.entry_count() {
3629 match parser.next_encoded_entry().unwrap() {
3630 Entry::Raw(EncodedPayload::Plain(_)) => raw += 1,
3631 Entry::Raw(EncodedPayload::Zstd(_)) => zstd += 1,
3632 Entry::Delta { .. } => deltas += 1,
3633 }
3634 }
3635 assert_eq!((raw, zstd, deltas), (GROUPS, GROUPS, 2 * GROUPS));
3636 let (_dir, store) = fresh_store();
3637 let budget = ResidentBudget::new(resident_bytes_cap(pack.len()), None);
3638 let report =
3639 PackReader::read_with_budget(&pack, &store, MAX_TOTAL_PAYLOAD, &budget).unwrap();
3640 assert_eq!(report.raw_count, 2 * GROUPS);
3641 assert_eq!(report.delta_count, 2 * GROUPS);
3642 assert_eq!(report.stored, expected);
3643 assert_eq!(budget.used.load(Ordering::Relaxed), 0);
3644 assert!(budget.peak.load(Ordering::Relaxed) <= budget.cap);
3645 for id in report.stored {
3647 store.read(&id).unwrap();
3648 }
3649 eprintln!(
3650 "large mixed pack: {} wire bytes, {} peak owned bytes, {:?}",
3651 pack.len(),
3652 budget.peak.load(Ordering::Relaxed),
3653 started.elapsed()
3654 );
3655 }
3656
3657 #[test]
3658 fn empty_pack_is_44_bytes() {
3659 let pack = PackWriter::new().finish().unwrap();
3660 assert_eq!(pack.len(), HEADER_LEN + TRAILER_LEN);
3661 assert_eq!(&pack[..4], MAGIC);
3662 assert_eq!(u32::from_le_bytes(pack[4..8].try_into().unwrap()), VERSION);
3663 assert_eq!(
3664 u32::from_le_bytes(
3665 pack[ENTRY_COUNT_OFFSET..ENTRY_COUNT_OFFSET + 4]
3666 .try_into()
3667 .unwrap()
3668 ),
3669 0
3670 );
3671
3672 let (_dir, store) = fresh_store();
3673 let report = PackReader::read(&pack, &store).unwrap();
3674 assert_eq!(report.raw_count, 0);
3675 assert_eq!(report.delta_count, 0);
3676 assert!(report.stored.is_empty());
3677 }
3678
3679 #[test]
3680 fn unpack_writes_objects_via_single_batch_flush() {
3681 use crate::batch::testing::{Ev, RecordingSyncer};
3684 use std::sync::Arc;
3685
3686 let mut w = PackWriter::new();
3687 let mut blobs = Vec::new();
3688 for i in 0u32..30 {
3689 let blob = write_blob_via_serialize(format!("pack object {i}").as_bytes());
3690 w.push_raw(hash::hash(&blob), &blob).unwrap();
3691 blobs.push(blob);
3692 }
3693 let pack = w.finish().unwrap();
3694
3695 let (_dir, mut store) = fresh_store();
3696 let rec = Arc::new(RecordingSyncer::default());
3697 store.set_syncer(rec.clone());
3698
3699 let report = PackReader::read(&pack, &store).unwrap();
3700 assert_eq!(report.raw_count, 30);
3701
3702 let fulls = rec
3703 .events()
3704 .iter()
3705 .filter(|e| matches!(e, Ev::Full(_)))
3706 .count();
3707 assert_eq!(
3708 fulls, 2,
3709 "unpack flush cost must be constant, not O(objects)"
3710 );
3711 for blob in &blobs {
3712 assert_eq!(store.read(&hash::hash(blob)).unwrap(), *blob);
3713 }
3714 }
3715
3716 #[test]
3717 fn single_raw_roundtrip() {
3718 let blob = write_blob_via_serialize(b"hello packfile");
3719 let h = hash::hash(&blob);
3720
3721 let mut w = PackWriter::new();
3722 w.push_raw(h, &blob).unwrap();
3723 let pack = w.finish().unwrap();
3724
3725 let (_dir, store) = fresh_store();
3726 let report = PackReader::read(&pack, &store).unwrap();
3727 assert_eq!(report.raw_count, 1);
3728 assert_eq!(report.delta_count, 0);
3729 assert_eq!(report.stored, vec![h]);
3730 assert_eq!(store.read(&h).unwrap(), blob);
3731 }
3732
3733 #[test]
3746 fn prepared_raw_produces_identical_pack_bytes_to_push_raw() {
3747 let blob = write_blob_via_serialize(b"hello prepared packfile");
3748 let h = hash::hash(&blob);
3749
3750 let mut direct = PackWriter::new();
3751 direct.push_raw(h, &blob).unwrap();
3752 let direct_pack = direct.finish().unwrap();
3753
3754 let prepared = PackWriter::prepare_raw(h, blob.clone());
3755 assert_eq!(prepared.hash(), h);
3756 assert_eq!(prepared.conservative_len(), blob.len());
3757 let mut via_prepared = PackWriter::new();
3758 via_prepared.push_prepared_raw(prepared).unwrap();
3759 let prepared_pack = via_prepared.finish().unwrap();
3760
3761 assert_eq!(
3762 direct_pack, prepared_pack,
3763 "push_raw and prepare_raw+push_prepared_raw must produce byte-identical packs"
3764 );
3765 }
3766
3767 #[test]
3768 fn prepared_delta_produces_identical_pack_bytes_to_push_delta() {
3769 let base = write_blob_via_serialize(&incompressible_bytes(0xD00D_0000, 2048));
3770 let target = write_blob_via_serialize(&incompressible_bytes(0xFEED_0000, 2048));
3771 let base_hash = hash::hash(&base);
3772 let stream = delta::encode(&base, &target).unwrap();
3773
3774 let mut direct = PackWriter::new();
3775 direct.push_delta(&base_hash, &stream).unwrap();
3776 let direct_pack = direct.finish().unwrap();
3777
3778 let prepared = PackWriter::prepare_delta(base_hash, stream.clone());
3779 assert_eq!(prepared.base(), base_hash);
3780 let mut via_prepared = PackWriter::new();
3781 via_prepared.push_prepared_delta(prepared).unwrap();
3782 let prepared_pack = via_prepared.finish().unwrap();
3783
3784 assert_eq!(
3785 direct_pack, prepared_pack,
3786 "push_delta and prepare_delta+push_prepared_delta must produce byte-identical packs"
3787 );
3788 }
3789
3790 #[test]
3791 fn prepared_and_direct_entries_interleave_in_push_order() {
3792 let a = write_blob_via_serialize(b"first entry, pushed directly");
3798 let ha = hash::hash(&a);
3799 let b = write_blob_via_serialize(b"second entry, pushed via prepare");
3800 let hb = hash::hash(&b);
3801
3802 let mut w = PackWriter::new();
3803 w.push_raw(ha, &a).unwrap();
3804 let prepared_b = PackWriter::prepare_raw(hb, b.clone());
3805 w.push_prepared_raw(prepared_b).unwrap();
3806 let pack = w.finish().unwrap();
3807
3808 let (_dir, store) = fresh_store();
3809 let report = PackReader::read(&pack, &store).unwrap();
3810 assert_eq!(
3811 report.stored,
3812 vec![ha, hb],
3813 "entries must appear in push order regardless of which path prepared them"
3814 );
3815 assert_eq!(store.read(&ha).unwrap(), a);
3816 assert_eq!(store.read(&hb).unwrap(), b);
3817 }
3818
3819 #[test]
3820 fn total_payload_tracks_wire_sum_for_mixed_raw_and_delta() {
3821 let mut w = PackWriter::new();
3828 assert_eq!(w.total_payload(), 0);
3829
3830 let raw = write_blob_via_serialize(&incompressible_bytes(0xA11C_E000, 2048));
3833 let raw_hash = hash::hash(&raw);
3834 w.push_raw(raw_hash, &raw).unwrap();
3835 assert_eq!(w.total_payload(), raw.len() as u64);
3836
3837 let base = write_blob_via_serialize(&incompressible_bytes(0xB0BA_1000, 2048));
3838 let base_hash = hash::hash(&base);
3839 let target = write_blob_via_serialize(&incompressible_bytes(0xC0FF_EE00, 2048));
3840 let stream = delta::encode(&base, &target).unwrap();
3841 let before_delta = w.total_payload();
3842 w.push_delta(&base_hash, &stream).unwrap();
3843 let delta_wire_len = w.total_payload() - before_delta;
3844
3845 assert!(delta_wire_len <= (hash::HASH_LEN + stream.len()) as u64);
3850 assert_eq!(w.total_payload(), before_delta + delta_wire_len);
3851 assert!(w.total_payload() <= raw.len() as u64 + (hash::HASH_LEN + stream.len()) as u64);
3852 }
3853
3854 #[test]
3855 fn raw_then_delta_resolves_in_pack() {
3856 let mut content_base = vec![0u8; 1024];
3858 for (i, b) in content_base.iter_mut().enumerate() {
3859 *b = u8::try_from(i % 251).expect("modulo < 256");
3860 }
3861 let mut content_target = content_base.clone();
3862 content_target[500] = 0xFF;
3863 content_target[501] = 0xFE;
3864
3865 let base_obj = write_blob_via_serialize(&content_base);
3866 let target_obj = write_blob_via_serialize(&content_target);
3867 let base_hash = hash::hash(&base_obj);
3868 let target_hash = hash::hash(&target_obj);
3869
3870 let stream = delta::encode(&base_obj, &target_obj).unwrap();
3871
3872 let mut w = PackWriter::new();
3873 w.push_raw(base_hash, &base_obj).unwrap();
3874 w.push_delta(&base_hash, &stream).unwrap();
3875 let pack = w.finish().unwrap();
3876
3877 let (_dir, store) = fresh_store();
3878 let report = PackReader::read(&pack, &store).unwrap();
3879 assert_eq!(report.raw_count, 1);
3880 assert_eq!(report.delta_count, 1);
3881 assert_eq!(report.stored, vec![base_hash, target_hash]);
3882 assert_eq!(store.read(&target_hash).unwrap(), target_obj);
3883 }
3884
3885 #[test]
3886 fn delta_before_its_base_in_pack_order_is_rejected() {
3887 let mut content_base = vec![0u8; 1024];
3898 for (i, b) in content_base.iter_mut().enumerate() {
3899 *b = u8::try_from(i % 251).expect("modulo < 256");
3900 }
3901 let mut content_target = content_base.clone();
3902 content_target[500] = 0xFF;
3903 content_target[501] = 0xFE;
3904
3905 let base_obj = write_blob_via_serialize(&content_base);
3906 let target_obj = write_blob_via_serialize(&content_target);
3907 let base_hash = hash::hash(&base_obj);
3908
3909 let stream = delta::encode(&base_obj, &target_obj).unwrap();
3910
3911 let mut w = PackWriter::new();
3912 w.push_delta(&base_hash, &stream).unwrap();
3913 w.push_raw(base_hash, &base_obj).unwrap();
3914 let pack = w.finish().unwrap();
3915
3916 let (_dir, store) = fresh_store();
3917 let err = PackReader::read(&pack, &store).unwrap_err();
3918 assert!(matches!(err, PackError::DeltaBaseMissing(_)), "got {err:?}");
3919 assert!(!store.contains(&base_hash));
3922 }
3923
3924 #[test]
3925 fn delta_before_its_base_is_rejected_under_parallel_raw_fanout() {
3926 let mut content_base = vec![0u8; 256];
3935 for (i, b) in content_base.iter_mut().enumerate() {
3936 *b = u8::try_from(i % 251).expect("modulo < 256");
3937 }
3938 let mut content_target = content_base.clone();
3939 content_target[10] = 0xFF;
3940
3941 let base_obj = write_blob_via_serialize(&content_base);
3942 let target_obj = write_blob_via_serialize(&content_target);
3943 let base_hash = hash::hash(&base_obj);
3944 let stream = delta::encode(&base_obj, &target_obj).unwrap();
3945
3946 let mut w = PackWriter::new();
3947 w.push_delta(&base_hash, &stream).unwrap();
3948 for i in 0..200u32 {
3949 let mut filler = vec![0u8; 64];
3950 for (j, b) in filler.iter_mut().enumerate() {
3951 *b = u8::try_from((i as usize + j) % 251).expect("modulo < 256");
3952 }
3953 let obj = write_blob_via_serialize(&filler);
3954 w.push_raw(hash::hash(&obj), &obj).unwrap();
3955 }
3956 w.push_raw(base_hash, &base_obj).unwrap();
3957 let pack = w.finish().unwrap();
3958
3959 let (_dir, store) = fresh_store();
3960 let err = PackReader::read(&pack, &store).unwrap_err();
3961 assert!(matches!(err, PackError::DeltaBaseMissing(_)), "got {err:?}");
3962 assert!(!store.contains(&base_hash));
3963 }
3964
3965 #[test]
3966 fn earlier_delta_base_missing_wins_over_later_malformed_raw_entry() {
3967 let base_obj = write_blob_via_serialize(&[0u8; 64]);
3978 let target_obj = write_blob_via_serialize(&[1u8; 64]);
3979 let base_hash = hash::hash(&base_obj);
3980 let stream = delta::encode(&base_obj, &target_obj).unwrap();
3981
3982 let mut w = PackWriter::new();
3983 w.push_delta(&base_hash, &stream).unwrap();
3984 for i in 0..200u32 {
3985 if i == 100 {
3986 w.push_raw([0xEE; 32], b"not a valid mkit object").unwrap();
3989 continue;
3990 }
3991 let mut filler = vec![0u8; 64];
3992 for (j, b) in filler.iter_mut().enumerate() {
3993 *b = u8::try_from((i as usize + j) % 251).expect("modulo < 256");
3994 }
3995 let obj = write_blob_via_serialize(&filler);
3996 w.push_raw(hash::hash(&obj), &obj).unwrap();
3997 }
3998 w.push_raw(base_hash, &base_obj).unwrap();
3999 let pack = w.finish().unwrap();
4000
4001 let (_dir, store) = fresh_store();
4002 let err = PackReader::read(&pack, &store).unwrap_err();
4003 assert!(
4004 matches!(err, PackError::DeltaBaseMissing(_)),
4005 "position 0's missing-base error must win over the malformed raw \
4006 entry at a later position, got {err:?}"
4007 );
4008 assert!(!store.contains(&base_hash));
4009 }
4010
4011 #[test]
4012 fn delta_base_hashes_lists_delta_bases_only() {
4013 let base_a = write_blob_via_serialize(b"base alpha content here padding");
4016 let base_b = write_blob_via_serialize(b"base bravo content here padding");
4017 let ha = hash::hash(&base_a);
4018 let hb = hash::hash(&base_b);
4019 let target_a = write_blob_via_serialize(b"base alpha content here PADDED!");
4020 let target_b = write_blob_via_serialize(b"base bravo content here PADDED!");
4021 let stream_a = delta::encode(&base_a, &target_a).unwrap();
4022 let stream_b = delta::encode(&base_b, &target_b).unwrap();
4023
4024 let mut w = PackWriter::new();
4025 w.push_raw(ha, &base_a).unwrap(); w.push_delta(&ha, &stream_a).unwrap();
4027 w.push_delta(&hb, &stream_b).unwrap();
4028 w.push_delta(&ha, &stream_a).unwrap(); let pack = w.finish().unwrap();
4030
4031 let mut bases = delta_base_hashes(&pack).unwrap();
4032 bases.sort_unstable();
4033 let mut expected = vec![ha, hb];
4034 expected.sort_unstable();
4035 assert_eq!(bases, expected);
4036 }
4037
4038 #[test]
4039 fn delta_base_hashes_rejects_bad_magic() {
4040 let mut pack = PackWriter::new().finish().unwrap();
4041 pack[0] = b'X';
4042 assert!(matches!(
4043 delta_base_hashes(&pack),
4044 Err(PackError::InvalidMagic)
4045 ));
4046 }
4047
4048 #[test]
4049 fn decoded_frame_metadata_reconstructs_the_same_object() {
4050 let blob = write_blob_via_serialize(b"frame metadata");
4051 let id = hash::hash(&blob);
4052 let mut writer = PackWriter::new_raw_only();
4053 writer.push_raw(id, &blob).unwrap();
4054 let pack = writer.finish().unwrap();
4055 let mut frames = Vec::new();
4056 decode_entries_with(
4057 &pack,
4058 &mut NoExternalBases,
4059 DecodeLimits::default(),
4060 |entry| {
4061 frames.push((
4062 entry.id,
4063 entry.frame_offset,
4064 entry.frame_length,
4065 entry.wire_type,
4066 entry.delta_base,
4067 ));
4068 Ok(())
4069 },
4070 )
4071 .unwrap();
4072 assert_eq!(frames.len(), 1);
4073 let (found, offset, length, kind, base) = frames[0];
4074 assert_eq!(found, id);
4075 assert_eq!(kind, 0);
4076 assert_eq!(base, None);
4077 let frame =
4078 &pack[usize::try_from(offset).unwrap()..usize::try_from(offset + length).unwrap()];
4079 let (decoded, bytes) = decode_frame_with(
4080 frame,
4081 VERSION,
4082 &mut NoExternalBases,
4083 DecodeLimits::default(),
4084 )
4085 .unwrap();
4086 assert_eq!((decoded, bytes), (id, blob));
4087 }
4088
4089 #[test]
4090 fn rejects_raw_payload_that_is_not_canonical_object_without_store_write() {
4091 let payload = b"not a serialized mkit object".to_vec();
4092 let payload_hash = hash::hash(&payload);
4093 let mut body = Vec::new();
4094 body.extend_from_slice(MAGIC);
4095 body.extend_from_slice(&VERSION.to_le_bytes());
4096 body.extend_from_slice(&1u32.to_le_bytes());
4097 body.push(0x00);
4098 let payload_len = u32::try_from(payload.len()).unwrap();
4099 body.extend_from_slice(&payload_len.to_le_bytes());
4100 body.extend_from_slice(&payload);
4101 let pack = finish_pack_body(body);
4102
4103 let (_dir, store) = fresh_store();
4104 let err = PackReader::read(&pack, &store).unwrap_err();
4105 assert!(matches!(err, PackError::InvalidObject(_)), "got {err:?}");
4106 assert!(!store.contains(&payload_hash));
4107 }
4108
4109 #[test]
4110 fn rejects_raw_delta_object_without_store_write() {
4111 let delta = crate::object::Object::Delta(crate::object::Delta {
4112 base_hash: [0xAB; 32],
4113 result_size: 0,
4114 instructions: Vec::new(),
4115 });
4116 let payload = crate::serialize::serialize(&delta).unwrap();
4117 let payload_hash = hash::hash(&payload);
4118 let mut w = PackWriter::new();
4119 w.push_raw(payload_hash, &payload).unwrap();
4120 let pack = w.finish().unwrap();
4121
4122 let (_dir, store) = fresh_store();
4123 let err = PackReader::read(&pack, &store).unwrap_err();
4124 assert!(matches!(err, PackError::NonStorableObject), "got {err:?}");
4125 assert!(!store.contains(&payload_hash));
4126 }
4127
4128 #[test]
4129 fn rejects_delta_resolving_to_non_object_without_partial_store_write() {
4130 let base_obj = write_blob_via_serialize(b"base bytes");
4131 let base_hash = hash::hash(&base_obj);
4132 let invalid_target = b"not a serialized object".to_vec();
4133 let invalid_hash = hash::hash(&invalid_target);
4134 let stream = delta::encode(&base_obj, &invalid_target).unwrap();
4135
4136 let mut w = PackWriter::new();
4137 w.push_raw(base_hash, &base_obj).unwrap();
4138 w.push_delta(&base_hash, &stream).unwrap();
4139 let pack = w.finish().unwrap();
4140
4141 let (_dir, store) = fresh_store();
4142 let err = PackReader::read(&pack, &store).unwrap_err();
4143 assert!(matches!(err, PackError::InvalidObject(_)), "got {err:?}");
4144 assert!(!store.contains(&base_hash));
4145 assert!(!store.contains(&invalid_hash));
4146 }
4147
4148 #[test]
4149 fn rejects_delta_result_over_object_cap_without_partial_store_write() {
4150 let base_obj = write_blob_via_serialize(b"base bytes");
4151 let base_hash = hash::hash(&base_obj);
4152 let mut stream = Vec::new();
4153 stream.push(delta::STREAM_VERSION);
4154 stream.extend_from_slice(&u32::try_from(base_obj.len()).unwrap().to_le_bytes());
4155 stream.extend_from_slice(
4156 &u32::try_from(MAX_RAW_OBJECT_SIZE + 1)
4157 .unwrap()
4158 .to_le_bytes(),
4159 );
4160
4161 let mut w = PackWriter::new();
4162 w.push_raw(base_hash, &base_obj).unwrap();
4163 w.push_delta(&base_hash, &stream).unwrap();
4164 let pack = w.finish().unwrap();
4165
4166 let (_dir, store) = fresh_store();
4167 let err = PackReader::read(&pack, &store).unwrap_err();
4168 assert!(
4169 matches!(
4170 err,
4171 PackError::Store(crate::store::StoreError::ObjectTooLarge)
4172 ),
4173 "got {err:?}"
4174 );
4175 assert!(!store.contains(&base_hash));
4176 }
4177
4178 #[test]
4179 fn rejects_trailing_bytes_after_declared_entries_without_store_write() {
4180 let blob = write_blob_via_serialize(b"trailing bytes test");
4181 let blob_hash = hash::hash(&blob);
4182 let mut body = Vec::new();
4183 body.extend_from_slice(MAGIC);
4184 body.extend_from_slice(&VERSION.to_le_bytes());
4185 body.extend_from_slice(&1u32.to_le_bytes());
4186 body.push(0x00);
4187 let blob_len = u32::try_from(blob.len()).unwrap();
4188 body.extend_from_slice(&blob_len.to_le_bytes());
4189 body.extend_from_slice(&blob);
4190 body.extend_from_slice(b"junk");
4191 let pack = finish_pack_body(body);
4192
4193 let (_dir, store) = fresh_store();
4194 let err = PackReader::read(&pack, &store).unwrap_err();
4195 assert!(matches!(err, PackError::TrailingData), "got {err:?}");
4196 assert!(!store.contains(&blob_hash));
4197 }
4198
4199 #[test]
4200 fn rejects_invalid_magic() {
4201 let mut pack = PackWriter::new().finish().unwrap();
4204 pack[0] = b'X';
4205 pack[1] = b'X';
4206 pack[2] = b'X';
4207 pack[3] = b'X';
4208 let (_dir, store) = fresh_store();
4209 let err = PackReader::read(&pack, &store).unwrap_err();
4210 assert!(matches!(err, PackError::InvalidMagic));
4211 }
4212
4213 #[test]
4214 fn rejects_unknown_version() {
4215 let mut pack = PackWriter::new().finish().unwrap();
4216 pack[4] = 99;
4218 let (_dir, store) = fresh_store();
4224 let err = PackReader::read(&pack, &store).unwrap_err();
4225 assert!(matches!(err, PackError::UnsupportedVersion(99)));
4226 }
4227
4228 #[test]
4229 fn rejects_truncated_pack() {
4230 let pack = vec![b'M', b'K']; let (_dir, store) = fresh_store();
4232 let err = PackReader::read(&pack, &store).unwrap_err();
4233 assert!(matches!(err, PackError::PackfileTooShort));
4234 }
4235
4236 #[test]
4237 fn rejects_bit_flipped_trailer() {
4238 let blob = write_blob_via_serialize(b"trailer test");
4239 let h = hash::hash(&blob);
4240 let mut w = PackWriter::new();
4241 w.push_raw(h, &blob).unwrap();
4242 let mut pack = w.finish().unwrap();
4243 let last = pack.len() - 1;
4244 pack[last] ^= 0x01; let (_dir, store) = fresh_store();
4246 let err = PackReader::read(&pack, &store).unwrap_err();
4247 assert!(matches!(err, PackError::PackfileCorrupted));
4248 }
4249
4250 #[test]
4251 fn rejects_reserved_entry_type_0x01() {
4252 let mut buf = Vec::new();
4254 buf.extend_from_slice(MAGIC);
4255 buf.extend_from_slice(&VERSION.to_le_bytes());
4256 buf.extend_from_slice(&1u32.to_le_bytes());
4257 buf.push(0x01); buf.extend_from_slice(&0u32.to_le_bytes()); let trailer = hash::hash(&buf);
4260 buf.extend_from_slice(&trailer);
4261
4262 let (_dir, store) = fresh_store();
4263 let err = PackReader::read(&buf, &store).unwrap_err();
4264 assert!(matches!(err, PackError::InvalidEntryType(0x01)));
4265 }
4266
4267 #[test]
4268 fn rejects_unknown_entry_type() {
4269 let mut buf = Vec::new();
4270 buf.extend_from_slice(MAGIC);
4271 buf.extend_from_slice(&VERSION.to_le_bytes());
4272 buf.extend_from_slice(&1u32.to_le_bytes());
4273 buf.push(0x77); buf.extend_from_slice(&0u32.to_le_bytes());
4275 let trailer = hash::hash(&buf);
4276 buf.extend_from_slice(&trailer);
4277
4278 let (_dir, store) = fresh_store();
4279 let err = PackReader::read(&buf, &store).unwrap_err();
4280 assert!(matches!(err, PackError::InvalidEntryType(0x77)));
4281 }
4282
4283 #[test]
4284 fn delta_base_missing_is_loud() {
4285 let mut fake_base = [0u8; 32];
4286 fake_base[0] = 0xAB;
4287 let mut stream = Vec::new();
4289 stream.push(0x01); stream.extend_from_slice(&0u32.to_le_bytes()); stream.extend_from_slice(&0u32.to_le_bytes()); let mut w = PackWriter::new();
4293 w.push_delta(&fake_base, &stream).unwrap();
4294 let pack = w.finish().unwrap();
4295
4296 let (_dir, store) = fresh_store();
4297 let err = PackReader::read(&pack, &store).unwrap_err();
4298 assert!(matches!(err, PackError::DeltaBaseMissing(_)), "got {err:?}");
4299 }
4300
4301 #[test]
4302 fn entry_payload_past_trailer_rejected() {
4303 let mut buf = Vec::new();
4304 buf.extend_from_slice(MAGIC);
4305 buf.extend_from_slice(&VERSION.to_le_bytes());
4306 buf.extend_from_slice(&1u32.to_le_bytes());
4307 buf.push(0x00);
4308 buf.extend_from_slice(&1_000_000u32.to_le_bytes());
4309 let trailer = hash::hash(&buf);
4311 buf.extend_from_slice(&trailer);
4312
4313 let (_dir, store) = fresh_store();
4314 let err = PackReader::read(&buf, &store).unwrap_err();
4315 assert!(matches!(err, PackError::UnexpectedEof));
4316 }
4317
4318 #[test]
4319 fn entry_count_over_cap_rejected() {
4320 let mut buf = Vec::new();
4321 buf.extend_from_slice(MAGIC);
4322 buf.extend_from_slice(&VERSION.to_le_bytes());
4323 buf.extend_from_slice(&u32::MAX.to_le_bytes());
4324 let trailer = hash::hash(&buf);
4329 buf.extend_from_slice(&trailer);
4330
4331 let (_dir, store) = fresh_store();
4332 let err = PackReader::read(&buf, &store).unwrap_err();
4333 assert!(
4336 matches!(err, PackError::TooManyObjects(_)),
4337 "expected TooManyObjects, got {err:?}"
4338 );
4339 }
4340
4341 #[test]
4342 fn payload_sum_over_cap_is_rejected_before_bounds_or_decode() {
4343 let blob_a = write_blob_via_serialize(&incompressible_bytes(0xA5A5, 64));
4353 let blob_b = write_blob_via_serialize(&incompressible_bytes(0xB6B6, 64));
4354 let mut w = PackWriter::new();
4355 w.push_raw(hash::hash(&blob_a), &blob_a).unwrap();
4356 w.push_raw(hash::hash(&blob_b), &blob_b).unwrap();
4357 let pack = w.finish().unwrap();
4358
4359 let (_dir, store) = fresh_store();
4360
4361 let cap = (blob_a.len() as u64) + 10;
4365 let err = PackReader::read_with_payload_cap(&pack, &store, cap).unwrap_err();
4366 assert!(
4367 matches!(err, PackError::PackfileTooLarge),
4368 "expected PackfileTooLarge, got {err:?}"
4369 );
4370
4371 let report = PackReader::read(&pack, &store).unwrap();
4374 assert_eq!(report.raw_count, 2);
4375 }
4376
4377 #[test]
4378 fn pack_key_is_blake3_of_pack_bytes() {
4379 let blob = write_blob_via_serialize(b"key test");
4380 let h = hash::hash(&blob);
4381 let mut w = PackWriter::new();
4382 w.push_raw(h, &blob).unwrap();
4383 let pack = w.finish().unwrap();
4384 assert_eq!(pack_key(&pack), hash::hash(&pack));
4385 }
4386
4387 #[test]
4388 fn unpack_does_not_recopy_raw_payloads_into_a_second_buffer() {
4389 let mut w = PackWriter::new();
4404 for i in 0u32..64 {
4405 let payload = incompressible_bytes(0x1000_0000 + u64::from(i), 16 * 1024);
4406 let blob = write_blob_via_serialize(&payload);
4407 w.push_raw(hash::hash(&blob), &blob).unwrap();
4408 }
4409 let pack = w.finish().unwrap();
4410 assert!(
4411 pack.len() > 512 * 1024,
4412 "sanity: synthetic pack should be substantial, got {}",
4413 pack.len()
4414 );
4415 assert_eq!(
4416 u32::from_le_bytes(pack[VERSION_OFFSET..VERSION_OFFSET + 4].try_into().unwrap()),
4417 VERSION,
4418 "sanity: incompressible filler must stay an uncompressed v1 pack"
4419 );
4420
4421 let (_dir, store) = fresh_store();
4422 let owned_bytes = AtomicU64::new(0);
4423 let report = PackReader::read_tracking_owned_bytes(&pack, &store, &owned_bytes).unwrap();
4424 assert_eq!(report.raw_count, 64);
4425
4426 assert_eq!(
4427 owned_bytes.load(Ordering::Relaxed),
4428 0,
4429 "an all-raw pack must not allocate a second copy of any entry's payload"
4430 );
4431 }
4432
4433 #[test]
4434 fn unpack_owned_bytes_for_deltas_is_exactly_the_delta_targets_not_the_whole_pack() {
4435 let content_base = incompressible_bytes(0x2BAD_2BAD, 4096);
4443 let base_obj = write_blob_via_serialize(&content_base);
4444 let base_hash = hash::hash(&base_obj);
4445
4446 let mut w = PackWriter::new();
4447 w.push_raw(base_hash, &base_obj).unwrap();
4448 let mut expected_owned = 0u64;
4449 for i in 0u32..10 {
4450 let mut target = content_base.clone();
4451 target[i as usize] ^= 0xFF;
4452 let target_obj = write_blob_via_serialize(&target);
4453 let stream = delta::encode(&base_obj, &target_obj).unwrap();
4454 w.push_delta(&base_hash, &stream).unwrap();
4455 expected_owned += target_obj.len() as u64;
4456 }
4457 let pack = w.finish().unwrap();
4458
4459 let (_dir, store) = fresh_store();
4460 let owned_bytes = AtomicU64::new(0);
4461 let report = PackReader::read_tracking_owned_bytes(&pack, &store, &owned_bytes).unwrap();
4462 assert_eq!(report.raw_count, 1);
4463 assert_eq!(report.delta_count, 10);
4464
4465 assert_eq!(
4466 owned_bytes.load(Ordering::Relaxed),
4467 expected_owned,
4468 "owned bytes must equal exactly the sum of delta target sizes — \
4469 no extra copy of the raw base"
4470 );
4471 }
4472
4473 #[test]
4474 fn pack_writer_finish_does_not_recopy_pushed_payloads() {
4475 let mut w = PackWriter::new();
4490 for i in 0u32..64 {
4491 let payload = incompressible_bytes(0x2000_0000 + u64::from(i), 16 * 1024);
4492 let blob = write_blob_via_serialize(&payload);
4493 w.push_raw(hash::hash(&blob), &blob).unwrap();
4494 }
4495 let bytes_copied = AtomicU64::new(0);
4496 let pack = w.finish_tracking_bytes_copied(&bytes_copied).unwrap();
4497 assert!(pack.len() > 512 * 1024);
4498
4499 assert_eq!(
4500 bytes_copied.load(Ordering::Relaxed),
4501 TRAILER_LEN as u64,
4502 "finish() must only append the trailer, not re-copy every pushed entry"
4503 );
4504 }
4505
4506 #[test]
4507 fn delta_resolves_against_pre_existing_store_object() {
4508 let (_dir, store) = fresh_store();
4509 let mut content_base = vec![0u8; 256];
4511 for (i, b) in content_base.iter_mut().enumerate() {
4512 *b = u8::try_from(i % 251).expect("modulo < 256");
4513 }
4514 let base_obj = write_blob_via_serialize(&content_base);
4515 let base_hash = store.write(&base_obj).unwrap();
4516
4517 let mut content_target = content_base.clone();
4519 content_target[100] = 0xAA;
4520 let target_obj = write_blob_via_serialize(&content_target);
4521 let target_hash = hash::hash(&target_obj);
4522 let stream = delta::encode(&base_obj, &target_obj).unwrap();
4523
4524 let mut w = PackWriter::new();
4525 w.push_delta(&base_hash, &stream).unwrap();
4526 let pack = w.finish().unwrap();
4527
4528 let report = PackReader::read(&pack, &store).unwrap();
4529 assert_eq!(report.delta_count, 1);
4530 assert_eq!(report.raw_count, 0);
4531 assert_eq!(store.read(&target_hash).unwrap(), target_obj);
4532 }
4533
4534 #[test]
4535 fn multiple_deltas_against_shared_external_base_read_store_once() {
4536 const N: usize = 5;
4542
4543 let (_dir, store) = fresh_store();
4544
4545 let mut content_base = vec![0u8; 512];
4546 for (i, b) in content_base.iter_mut().enumerate() {
4547 *b = u8::try_from(i % 251).expect("modulo < 256");
4548 }
4549 let base_obj = write_blob_via_serialize(&content_base);
4550 let base_hash = store.write(&base_obj).unwrap();
4551
4552 let mut w = PackWriter::new();
4554 let mut expected_targets = Vec::new();
4555 for i in 0..N {
4556 let mut content_target = content_base.clone();
4557 content_target[100] = u8::try_from(i).unwrap();
4558 let target_obj = write_blob_via_serialize(&content_target);
4559 let target_hash = hash::hash(&target_obj);
4560 let stream = delta::encode(&base_obj, &target_obj).unwrap();
4561 w.push_delta(&base_hash, &stream).unwrap();
4562 expected_targets.push((target_hash, target_obj));
4563 }
4564 let pack = w.finish().unwrap();
4565
4566 let reads_before = store.read_call_count();
4567 let report = PackReader::read(&pack, &store).unwrap();
4568 let reads_after_for_base = store.read_call_count() - reads_before;
4569
4570 assert_eq!(report.delta_count, u32::try_from(N).unwrap());
4571 assert_eq!(
4572 reads_after_for_base, 1,
4573 "base object must be read from the store exactly once for {N} deltas sharing it, got {reads_after_for_base}"
4574 );
4575
4576 for (target_hash, target_obj) in expected_targets {
4580 assert_eq!(store.read(&target_hash).unwrap(), target_obj);
4581 }
4582 }
4583
4584 #[cfg(feature = "pack-zstd")]
4592 fn compressible_bytes(len: usize) -> Vec<u8> {
4593 vec![0x42u8; len]
4594 }
4595
4596 #[test]
4597 #[cfg(feature = "pack-zstd")]
4598 fn compressed_raw_entry_roundtrips() {
4599 let payload = compressible_bytes(4096);
4600 let blob = write_blob_via_serialize(&payload);
4601 let h = hash::hash(&blob);
4602
4603 let mut w = PackWriter::new();
4604 w.push_raw(h, &blob).unwrap();
4605 let pack = w.finish().unwrap();
4606
4607 assert_eq!(
4608 u32::from_le_bytes(pack[VERSION_OFFSET..VERSION_OFFSET + 4].try_into().unwrap()),
4609 VERSION_V2,
4610 "a pack containing a compressed entry must be emitted as version 2"
4611 );
4612 assert_eq!(
4613 pack[HEADER_LEN], 0x03,
4614 "a highly-compressible raw payload must be emitted as 0x03 zstd-raw"
4615 );
4616
4617 let (_dir, store) = fresh_store();
4618 let report = PackReader::read(&pack, &store).unwrap();
4619 assert_eq!(report.raw_count, 1);
4620 assert_eq!(report.delta_count, 0);
4621 assert_eq!(report.stored, vec![h]);
4622 assert_eq!(
4623 store.read(&h).unwrap(),
4624 blob,
4625 "recovered object must be byte-identical to the pre-compression original"
4626 );
4627 }
4628
4629 #[test]
4630 #[cfg(feature = "pack-zstd")]
4631 fn compressed_delta_entry_roundtrips() {
4632 let base_obj =
4637 write_blob_via_serialize(b"delta base filler bytes, not compressible-target-shaped");
4638 let base_hash = hash::hash(&base_obj);
4639 let target_content = compressible_bytes(4096);
4640 let target_obj = write_blob_via_serialize(&target_content);
4641 let target_hash = hash::hash(&target_obj);
4642 let stream = delta::encode(&base_obj, &target_obj).unwrap();
4643 assert!(
4644 stream.len() >= 64,
4645 "sanity: delta stream must clear the writer's compression-candidate floor, got {}",
4646 stream.len()
4647 );
4648
4649 let mut w = PackWriter::new();
4650 w.push_raw(base_hash, &base_obj).unwrap();
4651 w.push_delta(&base_hash, &stream).unwrap();
4652 let pack = w.finish().unwrap();
4653
4654 assert_eq!(
4655 u32::from_le_bytes(pack[VERSION_OFFSET..VERSION_OFFSET + 4].try_into().unwrap()),
4656 VERSION_V2,
4657 "a pack containing a compressed entry must be emitted as version 2"
4658 );
4659 let base_payload_len =
4662 u32::from_le_bytes(pack[HEADER_LEN + 1..HEADER_LEN + 5].try_into().unwrap()) as usize;
4663 let second_entry_type_offset = HEADER_LEN + ENTRY_FRAME_LEN + base_payload_len;
4664 assert_eq!(
4665 pack[second_entry_type_offset], 0x04,
4666 "a highly-compressible delta stream must be emitted as 0x04 zstd-delta"
4667 );
4668
4669 let (_dir, store) = fresh_store();
4670 let report = PackReader::read(&pack, &store).unwrap();
4671 assert_eq!(report.raw_count, 1);
4672 assert_eq!(report.delta_count, 1);
4673 assert_eq!(report.stored, vec![base_hash, target_hash]);
4674 assert_eq!(
4675 store.read(&target_hash).unwrap(),
4676 target_obj,
4677 "recovered delta target must be byte-identical to the pre-compression original"
4678 );
4679 }
4680
4681 #[test]
4682 fn rejects_v2_entry_type_in_v1_pack() {
4683 let mut buf = Vec::new();
4688 buf.extend_from_slice(MAGIC);
4689 buf.extend_from_slice(&VERSION.to_le_bytes()); buf.extend_from_slice(&1u32.to_le_bytes()); buf.push(0x03);
4692 let inner_payload = 0u32.to_le_bytes(); buf.extend_from_slice(&u32::try_from(inner_payload.len()).unwrap().to_le_bytes());
4694 buf.extend_from_slice(&inner_payload);
4695 let pack = finish_pack_body(buf);
4696
4697 let (_dir, store) = fresh_store();
4698 let err = PackReader::read(&pack, &store).unwrap_err();
4699 assert!(
4700 matches!(err, PackError::InvalidEntryType(0x03)),
4701 "got {err:?}"
4702 );
4703 }
4704
4705 #[test]
4706 #[cfg(feature = "pack-zstd")]
4707 fn rejects_decompressed_len_mismatch() {
4708 let payload = compressible_bytes(4096);
4712 let blob = write_blob_via_serialize(&payload);
4713 let h = hash::hash(&blob);
4714 let mut w = PackWriter::new();
4715 w.push_raw(h, &blob).unwrap();
4716 let mut pack = w.finish().unwrap();
4717
4718 assert_eq!(pack[HEADER_LEN], 0x03, "sanity: must be a zstd-raw entry");
4719 let len_prefix_offset = HEADER_LEN + ENTRY_FRAME_LEN;
4720 let claimed_len = u32::from_le_bytes(
4721 pack[len_prefix_offset..len_prefix_offset + 4]
4722 .try_into()
4723 .unwrap(),
4724 );
4725 pack[len_prefix_offset..len_prefix_offset + 4]
4731 .copy_from_slice(&(claimed_len + 1).to_le_bytes());
4732 let split = pack.len() - TRAILER_LEN;
4733 let new_trailer = hash::hash(&pack[..split]);
4734 pack[split..].copy_from_slice(&new_trailer);
4735
4736 let (_dir, store) = fresh_store();
4737 let err = PackReader::read(&pack, &store).unwrap_err();
4738 assert!(
4739 matches!(err, PackError::DecompressedSizeMismatch(_, _)),
4740 "got {err:?}"
4741 );
4742 assert!(!store.contains(&h));
4743 }
4744
4745 #[test]
4746 #[cfg(feature = "pack-zstd")]
4747 fn rejects_decompressed_len_over_object_cap() {
4748 let claimed_len = u32::try_from(MAX_RAW_OBJECT_SIZE + 1).unwrap();
4753 let mut buf = Vec::new();
4754 buf.extend_from_slice(MAGIC);
4755 buf.extend_from_slice(&VERSION_V2.to_le_bytes());
4756 buf.extend_from_slice(&1u32.to_le_bytes());
4757 buf.push(0x03);
4758 let mut inner = Vec::new();
4762 inner.extend_from_slice(&claimed_len.to_le_bytes());
4763 inner.extend_from_slice(&[0u8; 8]);
4764 buf.extend_from_slice(&u32::try_from(inner.len()).unwrap().to_le_bytes());
4765 buf.extend_from_slice(&inner);
4766 let pack = finish_pack_body(buf);
4767
4768 let (_dir, store) = fresh_store();
4769 let err = PackReader::read(&pack, &store).unwrap_err();
4770 assert!(
4771 matches!(err, PackError::DecompressedSizeOverCap(n) if n == claimed_len as usize),
4772 "got {err:?}"
4773 );
4774 }
4775
4776 #[test]
4777 #[cfg(feature = "pack-zstd")]
4778 fn raw_only_writer_emits_v1_raw_for_compressible_payload() {
4779 let payload = compressible_bytes(1024 * 1024);
4780 let blob = write_blob_via_serialize(&payload);
4781 let h = hash::hash(&blob);
4782 let mut w = PackWriter::new_raw_only();
4783 w.push_raw(h, &blob).unwrap();
4784 let pack = w.finish().unwrap();
4785
4786 assert_eq!(
4787 u32::from_le_bytes(pack[VERSION_OFFSET..VERSION_OFFSET + 4].try_into().unwrap()),
4788 VERSION,
4789 "raw-only writer must finish as v1"
4790 );
4791 assert_eq!(pack[HEADER_LEN], 0x00, "every entry must be 0x00");
4792 let entries = PackEntries::new(&pack).unwrap();
4793 assert!(entries.is_raw_only());
4794 assert_eq!(entries.first_non_raw_index(), None);
4795
4796 let (_dir, store) = fresh_store();
4797 let report = PackReader::read(&pack, &store).unwrap();
4798 assert_eq!(report.raw_count, 1);
4799 assert_eq!(store.read(&h).unwrap(), blob);
4800 }
4801
4802 #[test]
4803 fn raw_only_writer_rejects_deltas() {
4804 let mut w = PackWriter::new_raw_only();
4805 let err = w.push_delta(&[0u8; 32], &[0u8; 16]).unwrap_err();
4806 assert!(matches!(err, PackError::RawOnly));
4807 let prepared = PackWriter::prepare_delta([1u8; 32], vec![0u8; 16]);
4808 let err = w.push_prepared_delta(prepared).unwrap_err();
4809 assert!(matches!(err, PackError::RawOnly));
4810 }
4811
4812 #[test]
4813 fn pack_entries_agrees_with_reader_on_empty_and_raw() {
4814 let mut w = PackWriter::new_raw_only();
4815 let blob = write_blob_via_serialize(b"pack-entries");
4816 let h = hash::hash(&blob);
4817 w.push_raw(h, &blob).unwrap();
4818 let pack = w.finish().unwrap();
4819
4820 let entries: Vec<_> = PackEntries::new(&pack)
4821 .unwrap()
4822 .collect::<Result<Vec<_>, _>>()
4823 .unwrap();
4824 assert_eq!(entries.len(), 1);
4825 match &entries[0] {
4826 PackEntry::Raw { bytes } => assert_eq!(bytes.as_ref(), blob.as_slice()),
4827 PackEntry::Delta { .. } => panic!("expected raw"),
4828 }
4829
4830 let (_dir, store) = fresh_store();
4831 PackReader::read(&pack, &store).unwrap();
4832 assert_eq!(store.read(&h).unwrap(), blob);
4833 }
4834
4835 type Variant = (Hash, Vec<u8>, Vec<u8>);
4841 type Seen = (Hash, Vec<u8>, bool);
4843
4844 fn base_and_variants(n: usize) -> (Vec<u8>, Hash, Vec<Variant>) {
4847 let mut content = vec![0u8; 512];
4848 for (i, b) in content.iter_mut().enumerate() {
4849 *b = u8::try_from(i % 251).expect("modulo < 256");
4850 }
4851 let base = write_blob_via_serialize(&content);
4852 let base_hash = hash::hash(&base);
4853 let variants = (0..n)
4854 .map(|i| {
4855 let mut c = content.clone();
4856 c[100] = u8::try_from(i).unwrap() ^ 0xA5;
4857 let target = write_blob_via_serialize(&c);
4858 let stream = delta::encode(&base, &target).unwrap();
4859 (hash::hash(&target), target, stream)
4860 })
4861 .collect();
4862 (base, base_hash, variants)
4863 }
4864
4865 fn decode_collect<B: DeltaBaseSource>(
4867 pack: &[u8],
4868 bases: &mut B,
4869 ) -> Result<(DecodeReport, Vec<Seen>), PackError> {
4870 let mut seen = Vec::new();
4871 let report = decode_entries_with(pack, bases, DecodeLimits::default(), |e| {
4872 assert_eq!(e.id, crate::object::id_from_object(&e.object, e.bytes));
4873 seen.push((e.id, e.bytes.to_vec(), e.from_delta));
4874 Ok(())
4875 })?;
4876 Ok((report, seen))
4877 }
4878
4879 #[test]
4880 fn decode_cursor_resumes_at_two_external_bases_without_replaying_entries() {
4881 #[derive(Default)]
4882 struct Supplied(std::collections::HashMap<Hash, Vec<u8>>);
4883 impl DeltaBaseSource for Supplied {
4884 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
4885 Ok(self.0.get(id).cloned())
4886 }
4887 }
4888
4889 let raw_before = write_blob_via_serialize(b"before first base");
4890 let raw_between = write_blob_via_serialize(b"between bases");
4891 let base_a = write_blob_via_serialize(b"external base a");
4892 let base_b = write_blob_via_serialize(b"external base b");
4893 let target_a = write_blob_via_serialize(b"target a");
4894 let target_b = write_blob_via_serialize(b"target b");
4895 let base_a_id = hash::hash(&base_a);
4896 let second_id = hash::hash(&base_b);
4897 let mut writer = PackWriter::new();
4898 writer
4899 .push_raw(hash::hash(&raw_before), &raw_before)
4900 .unwrap();
4901 writer
4902 .push_delta(&base_a_id, &delta::encode(&base_a, &target_a).unwrap())
4903 .unwrap();
4904 writer
4905 .push_raw(hash::hash(&raw_between), &raw_between)
4906 .unwrap();
4907 writer
4908 .push_delta(&second_id, &delta::encode(&base_b, &target_b).unwrap())
4909 .unwrap();
4910 let pack = writer.finish().unwrap();
4911
4912 let mut cursor = PackDecodeCursor::new(&pack, DecodeLimits::default()).unwrap();
4913 let mut supplied = Supplied::default();
4914 let mut seen = Vec::new();
4915 let first = cursor
4916 .resume(&mut supplied, |entry| {
4917 seen.push((entry.id, entry.bytes.to_vec(), entry.from_delta));
4918 Ok(())
4919 })
4920 .unwrap_err();
4921 assert!(
4922 matches!(first, PackError::DeltaBaseMissing(ref id) if *id == hash::to_hex(&base_a_id))
4923 );
4924 assert_eq!(seen.len(), 1);
4925 assert!(matches!(
4926 cursor.set_max_decoded_bytes(0),
4927 Err(PackError::PackfileTooLarge)
4928 ));
4929 cursor.set_max_decoded_bytes(4096).unwrap();
4930
4931 supplied.0.insert(base_a_id, base_a.clone());
4932 let second = cursor
4933 .resume(&mut supplied, |entry| {
4934 seen.push((entry.id, entry.bytes.to_vec(), entry.from_delta));
4935 Ok(())
4936 })
4937 .unwrap_err();
4938 assert!(
4939 matches!(second, PackError::DeltaBaseMissing(ref id) if *id == hash::to_hex(&second_id))
4940 );
4941 assert_eq!(seen.len(), 3);
4942
4943 supplied.0.insert(second_id, base_b.clone());
4944 let report = cursor
4945 .resume(&mut supplied, |entry| {
4946 seen.push((entry.id, entry.bytes.to_vec(), entry.from_delta));
4947 Ok(())
4948 })
4949 .unwrap();
4950 let (baseline_report, baseline_seen) = decode_collect(&pack, &mut supplied).unwrap();
4951 assert_eq!(report, baseline_report);
4952 assert_eq!(seen, baseline_seen);
4953 assert_eq!(report.ids.len(), 4);
4954 }
4955
4956 fn self_contained_packs() -> Vec<Vec<u8>> {
4959 let mut packs = vec![PackWriter::new().finish().unwrap()];
4960
4961 let mut w = PackWriter::new();
4962 for i in 0..20u32 {
4963 let blob = write_blob_via_serialize(&incompressible_bytes(u64::from(i), 256));
4964 w.push_raw(hash::hash(&blob), &blob).unwrap();
4965 }
4966 packs.push(w.finish().unwrap());
4967
4968 let (base, base_hash, variants) = base_and_variants(3);
4969 let mut w = PackWriter::new();
4970 w.push_raw(base_hash, &base).unwrap();
4971 for (_, _, stream) in &variants {
4972 w.push_delta(&base_hash, stream).unwrap();
4973 }
4974 let (t0_hash, t0, _) = &variants[0];
4976 let mut chained = t0.clone();
4977 let last = chained.len() - 1;
4978 chained[last] ^= 0x01;
4979 w.push_delta(t0_hash, &delta::encode(t0, &chained).unwrap())
4980 .unwrap();
4981 packs.push(w.finish().unwrap());
4982
4983 let tree = crate::serialize::serialize(&Object::Tree(crate::object::Tree {
4984 entries: vec![crate::object::TreeEntry {
4985 name: b"f".to_vec(),
4986 mode: crate::object::EntryMode::Blob,
4987 object_hash: base_hash,
4988 }],
4989 }))
4990 .unwrap();
4991 let mut w = PackWriter::new_raw_only();
4992 w.push_raw(crate::object::object_id_from_bytes(&tree), &tree)
4993 .unwrap();
4994 w.push_raw(base_hash, &base).unwrap();
4995 packs.push(w.finish().unwrap());
4996
4997 #[cfg(feature = "pack-zstd")]
4998 {
4999 let base = write_blob_via_serialize(b"delta base filler, not target-shaped");
5000 let base_hash = hash::hash(&base);
5001 let target = write_blob_via_serialize(&compressible_bytes(4096));
5002 let raw = write_blob_via_serialize(&compressible_bytes(8192));
5003 let mut w = PackWriter::new();
5004 w.push_raw(hash::hash(&raw), &raw).unwrap();
5005 w.push_raw(base_hash, &base).unwrap();
5006 w.push_delta(&base_hash, &delta::encode(&base, &target).unwrap())
5007 .unwrap();
5008 let pack = w.finish().unwrap();
5009 assert_eq!(pack[HEADER_LEN], 0x03, "sanity: zstd-raw entry");
5010 packs.push(pack);
5011 }
5012 packs
5013 }
5014
5015 #[test]
5016 fn decode_with_no_external_bases_matches_reader() {
5017 for pack in self_contained_packs() {
5018 let (_dir, store) = fresh_store();
5019 let unpacked = PackReader::read(&pack, &store).unwrap();
5020 let (report, seen) = decode_collect(&pack, &mut NoExternalBases).unwrap();
5021
5022 assert_eq!(report.ids, unpacked.stored);
5023 assert_eq!(report.raw_count, unpacked.raw_count as usize);
5024 assert_eq!(report.delta_count, unpacked.delta_count as usize);
5025 assert_eq!(
5026 seen.iter().filter(|s| s.2).count(),
5027 unpacked.delta_count as usize
5028 );
5029 for (id, bytes, _) in &seen {
5030 assert_eq!(&store.read(id).unwrap(), bytes, "same bytes under {id:?}");
5031 }
5032 assert_eq!(seen.iter().map(|s| s.0).collect::<Vec<_>>(), report.ids);
5033 }
5034 }
5035
5036 #[test]
5037 fn decode_rejects_exactly_what_reader_rejects() {
5038 let unbounded = DecodeLimits::default().with_max_decoded_bytes(u64::MAX);
5041 let base_obj = write_blob_via_serialize(&[0u8; 64]);
5044 let target_obj = write_blob_via_serialize(&[1u8; 64]);
5045 let base_hash = hash::hash(&base_obj);
5046 let stream = delta::encode(&base_obj, &target_obj).unwrap();
5047 let mut w = PackWriter::new();
5048 w.push_delta(&base_hash, &stream).unwrap();
5049 for i in 0..200u32 {
5050 if i == 100 {
5051 w.push_raw([0xEE; 32], b"not a valid mkit object").unwrap();
5052 continue;
5053 }
5054 let obj = write_blob_via_serialize(&i.to_le_bytes());
5055 w.push_raw(hash::hash(&obj), &obj).unwrap();
5056 }
5057 w.push_raw(base_hash, &base_obj).unwrap();
5058 let delta_first = w.finish().unwrap();
5059
5060 let mut w = PackWriter::new();
5061 w.push_raw([0xEE; 32], b"not a valid mkit object").unwrap();
5062 w.push_delta(&base_hash, &stream).unwrap();
5063 let raw_first = w.finish().unwrap();
5064
5065 let mut w = PackWriter::new();
5066 w.push_raw(base_hash, &base_obj).unwrap();
5067 let oversized = {
5068 let mut s = stream.clone();
5069 s[5..9].copy_from_slice(
5070 &u32::try_from(MAX_RAW_OBJECT_SIZE + 1)
5071 .unwrap()
5072 .to_le_bytes(),
5073 );
5074 s
5075 };
5076 w.push_delta(&base_hash, &oversized).unwrap();
5077 let over_cap = w.finish().unwrap();
5078
5079 let mut flipped = PackWriter::new().finish().unwrap();
5080 flipped[HEADER_LEN] ^= 0x01;
5081
5082 for pack in [delta_first, raw_first, over_cap, flipped] {
5083 let (_dir, store) = fresh_store();
5084 let reader = PackReader::read(&pack, &store).unwrap_err();
5085 let decoder = decode_entries_with(&pack, &mut NoExternalBases, unbounded, |_| Ok(()))
5086 .unwrap_err();
5087 assert_eq!(reader.to_string(), decoder.to_string());
5088 }
5089 }
5090
5091 #[test]
5092 fn external_base_outside_source_is_delta_base_missing() {
5093 struct Membership<'s> {
5098 store: &'s ObjectStore,
5099 members: std::collections::BTreeSet<Hash>,
5100 }
5101 impl DeltaBaseSource for Membership<'_> {
5102 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
5103 if !self.members.contains(id) {
5104 return Ok(None);
5105 }
5106 Ok(Some(self.store.read(id)?))
5107 }
5108 }
5109
5110 let (_other_dir, other_repo) = fresh_store();
5111 let (base, base_hash, variants) = base_and_variants(1);
5112 other_repo.write(&base).unwrap();
5113 let mut w = PackWriter::new();
5114 w.push_delta(&base_hash, &variants[0].2).unwrap();
5115 let pack = w.finish().unwrap();
5116
5117 let (_dir, empty) = fresh_store();
5118 let nowhere = PackReader::read(&pack, &empty).unwrap_err().to_string();
5119 assert_eq!(
5120 nowhere,
5121 PackError::DeltaBaseMissing(hash::to_hex(&base_hash)).to_string()
5122 );
5123
5124 let mut sink_calls = 0usize;
5125 let no_ext =
5126 decode_entries_with(&pack, &mut NoExternalBases, DecodeLimits::default(), |_| {
5127 sink_calls += 1;
5128 Ok(())
5129 })
5130 .unwrap_err();
5131 assert_eq!(no_ext.to_string(), nowhere);
5132
5133 let mut scoped = Membership {
5134 store: &other_repo,
5135 members: std::collections::BTreeSet::new(),
5136 };
5137 let not_member = decode_entries_with(&pack, &mut scoped, DecodeLimits::default(), |_| {
5138 sink_calls += 1;
5139 Ok(())
5140 })
5141 .unwrap_err();
5142 assert_eq!(not_member.to_string(), nowhere);
5143 assert_eq!(sink_calls, 0);
5144
5145 scoped.members.insert(base_hash);
5147 let (report, seen) = decode_collect(&pack, &mut scoped).unwrap();
5148 assert_eq!(report.ids, vec![variants[0].0]);
5149 assert_eq!(seen[0].1, variants[0].1);
5150 assert!(seen[0].2);
5151 }
5152
5153 #[test]
5154 fn untrusted_source_returning_wrong_bytes_is_rejected() {
5155 struct Lying {
5156 answer: Vec<u8>,
5157 }
5158 impl DeltaBaseSource for Lying {
5159 fn base(&mut self, _id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
5160 Ok(Some(self.answer.clone()))
5161 }
5162 }
5163
5164 let (base, base_hash, variants) = base_and_variants(1);
5165 let mut w = PackWriter::new();
5166 w.push_delta(&base_hash, &variants[0].2).unwrap();
5167 let pack = w.finish().unwrap();
5168 let expected = PackError::DeltaBaseMissing(hash::to_hex(&base_hash)).to_string();
5169
5170 let impostor = write_blob_via_serialize(b"a different valid object");
5173 let delta_obj = crate::serialize::serialize(&Object::Delta(crate::object::Delta {
5174 base_hash,
5175 result_size: 0,
5176 instructions: vec![],
5177 }))
5178 .unwrap();
5179 for answer in [impostor, b"garbage".to_vec(), delta_obj] {
5180 let mut emitted = Vec::new();
5181 let err =
5182 decode_entries_with(&pack, &mut Lying { answer }, DecodeLimits::default(), |e| {
5183 emitted.push(e.id);
5184 Ok(())
5185 })
5186 .unwrap_err();
5187 assert_eq!(err.to_string(), expected);
5188 assert!(emitted.is_empty(), "no target may be emitted");
5189 }
5190
5191 let (report, _) = decode_collect(&pack, &mut Lying { answer: base }).unwrap();
5193 assert_eq!(report.ids, vec![variants[0].0]);
5194 }
5195
5196 #[test]
5197 fn store_source_is_verified_once() {
5198 const N: usize = 5;
5201 let (_dir, store) = fresh_store();
5202 let (base, base_hash, variants) = base_and_variants(N);
5203 store.write(&base).unwrap();
5204 let mut w = PackWriter::new();
5205 for (_, _, stream) in &variants {
5206 w.push_delta(&base_hash, stream).unwrap();
5207 }
5208 let pack = w.finish().unwrap();
5209
5210 let reads_before = store.read_call_count();
5211 let mut source = &store;
5212 let (report, seen) = decode_collect(&pack, &mut source).unwrap();
5213 assert_eq!(store.read_call_count() - reads_before, 1);
5214 assert_eq!(report.delta_count, N);
5215 for ((id, bytes, _), (want_id, want, _)) in seen.iter().zip(&variants) {
5216 assert_eq!(id, want_id);
5217 assert_eq!(bytes, want);
5218 }
5219 for (id, _, _) in &variants {
5221 assert!(!store.contains(id));
5222 }
5223 }
5224
5225 #[test]
5226 fn corrupt_store_base_is_a_loud_store_error() {
5227 use std::io::{Seek, Write};
5231 let (_dir, store) = fresh_store();
5232 let (base, base_hash, variants) = base_and_variants(1);
5233 store.write(&base).unwrap();
5234 let mut f = std::fs::OpenOptions::new()
5235 .write(true)
5236 .open(store.path_for(&base_hash))
5237 .unwrap();
5238 f.seek(std::io::SeekFrom::End(-1)).unwrap();
5239 f.write_all(&[base[base.len() - 1] ^ 0xFF]).unwrap();
5240 drop(f);
5241 let mut w = PackWriter::new();
5242 w.push_delta(&base_hash, &variants[0].2).unwrap();
5243 let pack = w.finish().unwrap();
5244
5245 let reader = PackReader::read(&pack, &store).unwrap_err();
5246 let mut source = &store;
5247 let decoder = decode_entries_with(&pack, &mut source, DecodeLimits::default(), |_| Ok(()))
5248 .unwrap_err();
5249 assert!(matches!(reader, PackError::Store(_)), "{reader:?}");
5250 assert_eq!(reader.to_string(), decoder.to_string());
5251 }
5252
5253 fn delta_bomb(n: u32, target_len: usize) -> (Vec<u8>, Vec<Hash>) {
5259 let base = write_blob_via_serialize(&vec![0u8; 65536]);
5260 let base_hash = hash::hash(&base);
5261 let zeros_at = u32::try_from(base.len() - 65536).unwrap();
5262 let mut prologue = write_blob_via_serialize(&[]);
5264 let len_at = prologue.len() - 4;
5265 prologue[len_at..].copy_from_slice(&u32::try_from(target_len).unwrap().to_le_bytes());
5266 let total = prologue.len() + target_len;
5267
5268 let mut w = PackWriter::new();
5269 w.push_raw(base_hash, &base).unwrap();
5270 let mut ids = Vec::new();
5271 for i in 0..n {
5272 let mut s = vec![delta::STREAM_VERSION];
5273 s.extend_from_slice(&u32::try_from(base.len()).unwrap().to_le_bytes());
5274 s.extend_from_slice(&u32::try_from(total).unwrap().to_le_bytes());
5275 s.push(u8::try_from(prologue.len()).unwrap());
5276 s.extend_from_slice(&prologue);
5277 s.push(4);
5279 s.extend_from_slice(&i.to_le_bytes());
5280 let mut left = target_len - 4;
5281 while left > 0 {
5282 let len = left.min(65535);
5283 s.push(delta::OP_COPY);
5284 s.extend_from_slice(&zeros_at.to_le_bytes());
5285 s.extend_from_slice(&u16::try_from(len).unwrap().to_le_bytes());
5286 left -= len;
5287 }
5288 w.push_delta(&base_hash, &s).unwrap();
5289 let mut data = vec![0u8; target_len];
5290 data[..4].copy_from_slice(&i.to_le_bytes());
5291 ids.push(hash::hash(&write_blob_via_serialize(&data)));
5292 }
5293 (w.finish().unwrap(), ids)
5294 }
5295
5296 #[test]
5297 fn delta_bomb_is_rejected_before_any_delta_is_applied() {
5298 let (pack, _) = delta_bomb(16, 128 << 20);
5302 assert!(pack.len() < 512 * 1024, "pack is {} bytes", pack.len());
5303 let mut sink_calls = 0usize;
5304 let err = decode_entries_with(&pack, &mut NoExternalBases, DecodeLimits::default(), |_| {
5305 sink_calls += 1;
5306 Ok(())
5307 })
5308 .unwrap_err();
5309 assert!(matches!(err, PackError::PackfileTooLarge), "{err:?}");
5310 assert_eq!(sink_calls, 0, "rejected before the first entry is judged");
5311 }
5312
5313 #[test]
5314 fn decode_budget_is_caller_set() {
5315 const TARGET: usize = 256 * 1024;
5321 let (pack, ids) = delta_bomb(3, TARGET);
5322 let declared = 3 * (TARGET as u64 + 10);
5323
5324 let fits = DecodeLimits::default().with_max_decoded_bytes(2 * declared);
5325 let (_dir, store) = fresh_store();
5326 let unpacked = PackReader::read(&pack, &store).unwrap();
5327 let mut seen = Vec::new();
5328 let report = decode_entries_with(&pack, &mut NoExternalBases, fits, |e| {
5329 seen.push(e.id);
5330 Ok(())
5331 })
5332 .unwrap();
5333 assert_eq!(report.ids, unpacked.stored);
5334 assert_eq!(seen[1..], ids[..]);
5335
5336 let short = DecodeLimits::default().with_max_decoded_bytes(declared - 1);
5337 let err = decode_entries_with(&pack, &mut NoExternalBases, short, |_| Ok(())).unwrap_err();
5338 assert!(matches!(err, PackError::PackfileTooLarge), "{err:?}");
5339 }
5340
5341 #[test]
5342 fn compressed_claims_are_charged_before_decompression() {
5343 let claimed = u32::try_from(512usize << 20).unwrap();
5347 let mut buf = Vec::new();
5348 buf.extend_from_slice(MAGIC);
5349 buf.extend_from_slice(&VERSION_V2.to_le_bytes());
5350 buf.extend_from_slice(&1u32.to_le_bytes());
5351 buf.push(0x03);
5352 let mut inner = claimed.to_le_bytes().to_vec();
5353 inner.extend_from_slice(&[0u8; 8]);
5354 buf.extend_from_slice(&u32::try_from(inner.len()).unwrap().to_le_bytes());
5355 buf.extend_from_slice(&inner);
5356 let pack = finish_pack_body(buf);
5357
5358 let small = DecodeLimits::default().with_max_decoded_bytes(1 << 20);
5359 let err = decode_entries_with(&pack, &mut NoExternalBases, small, |_| Ok(())).unwrap_err();
5360 assert!(matches!(err, PackError::PackfileTooLarge), "{err:?}");
5361 let err = decode_entries_with(&pack, &mut NoExternalBases, DecodeLimits::default(), |_| {
5363 Ok(())
5364 })
5365 .unwrap_err();
5366 assert!(!matches!(err, PackError::PackfileTooLarge), "{err:?}");
5367 }
5368
5369 #[test]
5370 fn payload_len_u32_max_is_a_clean_error() {
5371 let blob = write_blob_via_serialize(&[1, 2, 3]);
5376 let mut w = PackWriter::new_raw_only();
5377 w.push_raw(hash::hash(&blob), &blob).unwrap();
5378 let mut pack = w.finish().unwrap();
5379 pack[HEADER_LEN + 1..HEADER_LEN + 5].copy_from_slice(&u32::MAX.to_le_bytes());
5380 let split = pack.len() - TRAILER_LEN;
5381 let trailer = hash::hash(&pack[..split]);
5382 pack[split..].copy_from_slice(&trailer);
5383
5384 assert!(matches!(
5385 PackEntries::new(&pack).unwrap_err(),
5386 PackError::UnexpectedEof
5387 ));
5388 assert!(matches!(
5389 delta_base_hashes(&pack).unwrap_err(),
5390 PackError::UnexpectedEof
5391 ));
5392 let err = decode_entries_with(&pack, &mut NoExternalBases, DecodeLimits::default(), |_| {
5393 Ok(())
5394 })
5395 .unwrap_err();
5396 assert!(matches!(err, PackError::UnexpectedEof), "{err:?}");
5397 }
5398
5399 #[test]
5400 fn only_named_bases_stay_resident() {
5401 let (base, base_hash, variants) = base_and_variants(2);
5405 let unrelated = write_blob_via_serialize(b"never a base");
5406 let (t0_hash, t0, _) = &variants[0];
5407 let mut chained = t0.clone();
5408 let last = chained.len() - 1;
5409 chained[last] ^= 0x01;
5410 let mut w = PackWriter::new();
5411 w.push_raw(hash::hash(&unrelated), &unrelated).unwrap();
5412 w.push_raw(base_hash, &base).unwrap();
5413 w.push_delta(&base_hash, &variants[0].2).unwrap();
5414 w.push_delta(t0_hash, &delta::encode(t0, &chained).unwrap())
5415 .unwrap();
5416 let pack = w.finish().unwrap();
5417 let (_dir, store) = fresh_store();
5418 let unpacked = PackReader::read(&pack, &store).unwrap();
5419 let (report, _) = decode_collect(&pack, &mut NoExternalBases).unwrap();
5420 assert_eq!(report.ids, unpacked.stored);
5421 }
5422
5423 struct Members {
5426 objects: std::collections::HashMap<Hash, Vec<u8>>,
5427 fetches: usize,
5428 }
5429
5430 impl DeltaBaseSource for Members {
5431 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
5432 self.fetches += 1;
5433 Ok(self.objects.get(id).cloned())
5434 }
5435 }
5436
5437 fn large_member_bases() -> (Members, Vec<(Hash, Vec<u8>)>) {
5441 let mut objects = std::collections::HashMap::new();
5442 let mut deltas = Vec::new();
5443 for seed in [1u64, 2] {
5444 let base = write_blob_via_serialize(&incompressible_bytes(seed, 600 * 1024));
5445 let id = hash::hash(&base);
5446 let target = write_blob_via_serialize(&base[10..300]);
5447 deltas.push((id, delta::encode(&base, &target).unwrap()));
5448 objects.insert(id, base);
5449 }
5450 (
5451 Members {
5452 objects,
5453 fetches: 0,
5454 },
5455 deltas,
5456 )
5457 }
5458
5459 fn pack_of_deltas(deltas: &[&(Hash, Vec<u8>)]) -> Vec<u8> {
5460 let mut w = PackWriter::new();
5461 for (base, stream) in deltas {
5462 w.push_delta(base, stream).unwrap();
5463 }
5464 w.finish().unwrap()
5465 }
5466
5467 #[test]
5468 fn external_bases_are_charged_against_the_budget() {
5469 let one_mib = DecodeLimits::default().with_max_decoded_bytes(1 << 20);
5470 let (mut members, deltas) = large_member_bases();
5471
5472 let pack = pack_of_deltas(&[&deltas[0], &deltas[1], &deltas[0]]);
5476 let mut sink_calls = 0usize;
5477 let err = decode_entries_with(&pack, &mut members, one_mib, |_| {
5478 sink_calls += 1;
5479 Ok(())
5480 })
5481 .unwrap_err();
5482 assert!(matches!(err, PackError::PackfileTooLarge), "{err:?}");
5483 assert_eq!(sink_calls, 1);
5484
5485 let half_mib = DecodeLimits::default().with_max_decoded_bytes(512 * 1024);
5487 let pack = pack_of_deltas(&[&deltas[0]]);
5488 let mut sink_calls = 0usize;
5489 let err = decode_entries_with(&pack, &mut members, half_mib, |_| {
5490 sink_calls += 1;
5491 Ok(())
5492 })
5493 .unwrap_err();
5494 assert!(matches!(err, PackError::PackfileTooLarge), "{err:?}");
5495 assert_eq!(sink_calls, 0);
5496 }
5497
5498 #[test]
5499 fn external_base_charge_is_released_after_its_last_use() {
5500 let one_mib = DecodeLimits::default().with_max_decoded_bytes(1 << 20);
5504 let (mut members, deltas) = large_member_bases();
5505 let pack = pack_of_deltas(&[&deltas[0], &deltas[0], &deltas[1]]);
5506 let report = decode_entries_with(&pack, &mut members, one_mib, |_| Ok(())).unwrap();
5507 assert_eq!(report.delta_count, 3);
5508 assert_eq!(members.fetches, 2);
5509 }
5510
5511 #[test]
5512 fn sink_error_stops_decode() {
5513 let mut w = PackWriter::new();
5514 let mut ids = Vec::new();
5515 for i in 0..5u32 {
5516 let blob = write_blob_via_serialize(&i.to_le_bytes());
5517 ids.push(w.push_raw(hash::hash(&blob), &blob).unwrap());
5518 }
5519 let pack = w.finish().unwrap();
5520
5521 let mut seen = Vec::new();
5522 let err = decode_entries_with(&pack, &mut NoExternalBases, DecodeLimits::default(), |e| {
5523 seen.push(e.id);
5524 if seen.len() == 2 {
5525 return Err(PackError::TrailingData);
5526 }
5527 Ok(())
5528 })
5529 .unwrap_err();
5530 assert!(matches!(err, PackError::TrailingData), "{err:?}");
5531 assert_eq!(seen, ids[..2]);
5532 }
5533
5534 #[test]
5535 fn pack_entries_is_raw_only_false_for_delta() {
5536 let base = write_blob_via_serialize(b"base-for-delta-scan");
5537 let base_hash = hash::hash(&base);
5538 let target = write_blob_via_serialize(b"target-for-delta-scan!");
5539 let stream = delta::encode(&base, &target).unwrap();
5540 let mut w = PackWriter::new();
5541 w.push_raw(base_hash, &base).unwrap();
5542 w.push_delta(&base_hash, &stream).unwrap();
5543 let pack = w.finish().unwrap();
5544 let entries = PackEntries::new(&pack).unwrap();
5545 assert!(!entries.is_raw_only());
5546 assert_eq!(entries.first_non_raw_index(), Some(1));
5547 }
5548}
5549
5550#[cfg(kani)]
5576mod kani_proofs {
5577 use super::*;
5578
5579 fn toy_hash(data: &[u8]) -> Hash {
5583 let mut h = [0u8; hash::HASH_LEN];
5584 let n = data.len().min(16);
5585 h[..n].copy_from_slice(&data[..n]);
5586 h[16..16 + n].copy_from_slice(&data[data.len() - n..]);
5587 #[allow(clippy::cast_possible_truncation)]
5588 {
5589 h[31] ^= data.len() as u8;
5590 }
5591 h
5592 }
5593
5594 fn stub_zstd(_frame: &[u8], _capacity: usize) -> Result<Vec<u8>, PackError> {
5595 if kani::any() {
5596 return Err(PackError::ZstdDecompress(String::new()));
5597 }
5598 let buf: [u8; 2] = kani::any();
5599 let n: usize = kani::any_where(|&n| n <= 2);
5600 Ok(buf[..n].to_vec())
5601 }
5602
5603 fn any_pack<const BODY: usize>() -> Vec<u8> {
5607 let head: [u8; HEADER_LEN] = kani::any();
5608 let body: [u8; BODY] = kani::any();
5609 let trailer: [u8; TRAILER_LEN] = kani::any();
5610 let mut v = Vec::with_capacity(HEADER_LEN + BODY + TRAILER_LEN);
5611 v.extend_from_slice(&head);
5612 v.extend_from_slice(&body);
5613 v.extend_from_slice(&trailer);
5614 v
5615 }
5616
5617 fn entries_at<const BODY: usize>() {
5618 let bytes = any_pack::<BODY>();
5619 let parsed = PackEntries::new(&bytes);
5620 let Ok(mut entries) = parsed else {
5621 kani::cover!(
5622 matches!(parsed, Err(PackError::PackfileCorrupted)),
5623 "bad_trailer"
5624 );
5625 kani::cover!(matches!(parsed, Err(PackError::TrailingData)), "trailing");
5626 return;
5627 };
5628 let split = bytes.len() - TRAILER_LEN;
5629 assert_eq!(&bytes[..4], MAGIC.as_slice());
5630 assert!(entries.version == VERSION || entries.version == VERSION_V2);
5631 assert_eq!(toy_hash(&bytes[..split]).as_slice(), &bytes[split..]);
5632 let count = entries.count;
5633 let raw_only = entries.is_raw_only();
5634 let v1 = entries.version == VERSION;
5635 let mut ok = 0u32;
5636 let mut failed = false;
5637 while let Some(item) = entries.next() {
5638 let r = entries.last_payload_range().expect("set after an item");
5639 assert!(HEADER_LEN + ENTRY_FRAME_LEN <= r.start && r.end <= split);
5640 match item {
5641 Ok(PackEntry::Raw { bytes: b }) => {
5642 assert!(!v1 || matches!(b, Cow::Borrowed(_)));
5643 ok += 1;
5644 }
5645 Ok(PackEntry::Delta { .. }) => {
5646 assert!(!raw_only);
5647 ok += 1;
5648 }
5649 Err(_) => {
5650 assert!(!v1, "a v1 pack accepted by new() must iterate cleanly");
5651 failed = true;
5652 }
5653 }
5654 }
5655 if !failed {
5656 assert_eq!(ok, count);
5657 assert_eq!(entries.pos, split);
5658 }
5659 kani::cover!(count == 2 && ok == 2, "two_entries");
5660 kani::cover!(!raw_only && ok >= 1, "delta_or_zstd_entry");
5661 }
5662
5663 #[kani::proof]
5681 #[kani::stub(crate::hash::hash, toy_hash)]
5682 #[kani::stub(zstd_decompress_capped, stub_zstd)]
5683 #[kani::unwind(4)]
5684 fn pack_entries_one_frame() {
5685 entries_at::<5>();
5686 }
5687
5688 #[kani::proof]
5691 #[kani::stub(crate::hash::hash, toy_hash)]
5692 #[kani::stub(zstd_decompress_capped, stub_zstd)]
5693 #[kani::unwind(4)]
5694 fn pack_entries_two_entries() {
5695 entries_at::<10>();
5696 }
5697
5698 fn writer_rt<const R: usize, const S: usize>(with_delta: bool) {
5699 let raw: [u8; R] = kani::any();
5700 let base: Hash = kani::any();
5701 let stream: [u8; S] = kani::any();
5702
5703 let mut w = PackWriter::new();
5704 w.push_raw(hash::ZERO, &raw).expect("raw fits caps");
5705 if with_delta {
5706 w.push_delta(&base, &stream).expect("delta fits caps");
5707 }
5708 let pack = w.finish().expect("finish");
5709
5710 let mut it = PackEntries::new(&pack).expect("own pack parses");
5711 assert_eq!(it.version, VERSION);
5712 assert_eq!(it.is_raw_only(), !with_delta);
5713 match it.next() {
5714 Some(Ok(PackEntry::Raw { bytes })) => {
5715 assert_eq!(bytes.as_ref(), raw);
5716 }
5717 _ => panic!("first entry must be the pushed raw payload"),
5718 }
5719 if with_delta {
5720 match it.next() {
5721 Some(Ok(PackEntry::Delta {
5722 base: b,
5723 stream: st,
5724 })) => {
5725 assert_eq!(b, base);
5726 assert_eq!(st.as_ref(), stream);
5727 }
5728 _ => panic!("second entry must be the pushed delta"),
5729 }
5730 assert_eq!(it.first_non_raw_index(), Some(1));
5731 }
5732 assert!(it.next().is_none());
5733 }
5734
5735 #[kani::proof]
5744 #[kani::stub(crate::hash::hash, toy_hash)]
5745 #[kani::unwind(4)]
5748 fn pack_writer_roundtrip_raw() {
5749 writer_rt::<1, 0>(false);
5750 }
5751
5752 #[kani::proof]
5756 #[kani::stub(crate::hash::hash, toy_hash)]
5757 #[kani::unwind(4)]
5759 #[kani::should_panic]
5760 fn pack_canary_mutation_still_parses() {
5761 let mut w = PackWriter::new_raw_only();
5762 w.push_raw(hash::ZERO, b"ab").expect("raw fits caps");
5763 let mut pack = w.finish().expect("finish");
5764 let i: usize = kani::any_where(|&i| i < pack.len());
5765 let flip: u8 = kani::any_where(|&f| f != 0);
5766 pack[i] ^= flip;
5767 assert!(PackEntries::new(&pack).is_ok());
5768 }
5769
5770 fn cursor_bytes<const N: usize>(tags: &[(usize, bool)], depths: [usize; 2]) -> Vec<u8> {
5782 let mut body: [u8; N] = kani::any();
5783 for &(at, present) in tags {
5784 body[at] = u8::from(present);
5785 }
5786 for at in depths {
5787 body[at] = 0;
5788 }
5789 let mut out = Vec::with_capacity(N + hash::HASH_LEN);
5790 out.extend_from_slice(&body);
5791 out.extend_from_slice(&toy_hash(&body));
5792 out
5793 }
5794
5795 fn cursor_anchor() -> Vec<u8> {
5802 cursor_bytes::<125>(
5803 &[
5804 (45, false),
5805 (54, false),
5806 (55, true),
5807 (88, true),
5808 (122, false),
5809 (124, false),
5810 ],
5811 [121, 123],
5812 )
5813 }
5814
5815 #[kani::proof]
5825 #[kani::stub(crate::hash::hash, toy_hash)]
5826 #[kani::unwind(7)]
5829 fn pack_window_cursor_decode() {
5830 let got = window::WindowCursor::from_bytes(&cursor_anchor());
5831 let ok = got.is_ok();
5832 core::mem::forget(got);
5835 kani::cover!(ok, "accepted_anchor_cursor");
5836 }
5837
5838 fn flipped_cursor(mask: u8) -> Vec<u8> {
5845 let mut body = [0u8; 125];
5846 body[0] = 1; body[1..9].copy_from_slice(&100u64.to_le_bytes()); body[9..17].copy_from_slice(&(64u64 << 10).to_le_bytes()); body[17..25].copy_from_slice(&12u64.to_le_bytes()); body[33..37].copy_from_slice(&1u32.to_le_bytes()); body[55] = 1; body[56..88].copy_from_slice(&kani::any::<Hash>());
5853 body[88] = 1; body[89..121].copy_from_slice(&kani::any::<Hash>());
5855 let mut out = body.to_vec();
5856 out.extend_from_slice(&toy_hash(&body));
5857 out[1] ^= mask;
5858 out
5859 }
5860
5861 #[kani::proof]
5869 #[kani::stub(crate::hash::hash, toy_hash)]
5870 #[kani::unwind(7)]
5872 fn pack_window_cursor_flip_rejected() {
5873 let mask: u8 = kani::any();
5874 let got = window::WindowCursor::from_bytes(&flipped_cursor(mask));
5875 let ok = got.is_ok();
5876 core::mem::forget(got);
5877 assert_eq!(ok, mask == 0);
5878 }
5879
5880 #[kani::proof]
5883 #[kani::stub(crate::hash::hash, toy_hash)]
5884 #[kani::unwind(7)]
5886 #[kani::should_panic]
5887 fn pack_window_cursor_canary_flip_accepted() {
5888 let got =
5889 window::WindowCursor::from_bytes(&flipped_cursor(kani::any_where(|&m: &u8| m != 0)));
5890 let ok = got.is_ok();
5891 core::mem::forget(got);
5892 assert!(ok);
5893 }
5894}