Skip to main content

mkit_core/
pack.rs

1//! Packfile writer / reader — conformant to `docs/specs/SPEC-PACKFILE.md`.
2//!
3//! Layout (SPEC-PACKFILE §1, §2, §3, §8):
4//!
5//! ```text
6//! [4B  magic            "MKIT"]                       offset 0
7//! [4B  version u32 LE  == 1 or 2]
8//! [4B  entry_count u32 LE     ]
9//!   for each entry:
10//!     [u8  entry_type]           0x00 raw | 0x02 delta | 0x03 zstd-raw | 0x04 zstd-delta
11//!     [u32 LE payload_len]                            length of payload only
12//!     [payload_len bytes payload]
13//! [32B trailer = BLAKE3 of all preceding bytes]
14//! ```
15//!
16//! Entry types (SPEC-PACKFILE §3):
17//!
18//! * `0x00` raw        — payload is a fully serialised mkit object.
19//! * `0x01`             — RESERVED, MUST be rejected.
20//! * `0x02` delta       — payload is `[32B base_hash][SPEC-DELTA stream]`.
21//! * `0x03` zstd-raw    — v2 only. payload is `[4B uncompressed_len LE][zstd frame]`;
22//!   the frame decompresses to exactly what a `0x00` payload would be.
23//! * `0x04` zstd-delta  — v2 only. payload is `[32B base_hash][4B uncompressed_len LE][zstd frame]`;
24//!   `base_hash` stays uncompressed, the frame decompresses to exactly
25//!   what a `0x02` entry's post-base-hash bytes would be.
26//!
27//! **Version selection is writer policy, not caller policy**
28//! (SPEC-PACKFILE §1): [`PackWriter`] emits `version = 1` when the
29//! finished pack contains no `0x03`/`0x04` entries, and `version = 2`
30//! the moment it contains at least one — even in an otherwise-mixed
31//! pack. `0x03`/`0x04` are illegal inside a `version = 1` pack; a
32//! reader seeing one there rejects with `InvalidEntryType` exactly as
33//! it would for any other unrecognized type (SPEC-PACKFILE §3).
34//!
35//! **Compression is per-entry**, not a whole-pack stream: every
36//! `0x03`/`0x04` entry carries its own independent zstd frame, so
37//! existing framing/caps/trailer semantics (§2, §5, §8) are unchanged
38//! and decompression memory is bounded to one entry at a time.
39//! Decoding a `0x03`/`0x04` entry is bomb-guarded: the claimed
40//! `uncompressed_len` is checked against [`MAX_RAW_OBJECT_SIZE`]
41//! *before* any decompression allocation, decompression itself is
42//! capacity-bounded to that claim, and the actual decompressed length
43//! is re-checked against the claim afterward (`DecompressedSizeMismatch`
44//! / `DecompressedSizeOverCap`). The payload must be exactly one
45//! Zstandard frame. The C decoder (`pack-zstd`) serves reads when it is
46//! compiled in; otherwise the pure-Rust, decode-only `pack-ruzstd`
47//! backend does, under the same checks; with neither, a compressed entry
48//! fails closed.
49//!
50//! Caps (SPEC-PACKFILE §5, unchanged by v2 — measured on the *wire*
51//! size; the decompressed-side cap above is separate and new):
52//!
53//! * `entry_count <= 10_000_000`
54//! * total `payload_len` sum `<= 4 GiB`
55//!
56//! Delta-base ordering rule (SPEC-PACKFILE §4): every delta entry's
57//! (`0x02` or `0x04`) `base_hash` MUST appear earlier in the same pack
58//! as a raw entry, OR already exist in the destination object store.
59//! `0x04`'s `base_hash` is never compressed, so this never requires
60//! decompression to evaluate. The "destination object store" is the
61//! [`DeltaBaseSource`] the decoder is given: the local [`ObjectStore`]
62//! for [`PackReader::read`], or a repository-scoped source for the
63//! store-less [`decode_entries_with`].
64//!
65//! The pack key (SPEC-PACKFILE §7) is `packs/<lower-hex BLAKE3 of entire
66//! pack>`. The trailer is then redundant w.r.t. that key, but it lets a
67//! streaming reader detect bit-rot before the whole pack has been
68//! hashed end-to-end.
69
70pub 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
82/// ASCII magic ("MKIT") at the start of every pack, v1 or v2.
83pub const MAGIC: &[u8; 4] = b"MKIT";
84/// Packfile version emitted when a pack contains no compressed
85/// (`0x03`/`0x04`) entries. Also the minimum version any reader
86/// accepts.
87pub const VERSION: u32 = 1;
88/// Packfile version emitted the moment a pack contains at least one
89/// compressed (`0x03`/`0x04`) entry (SPEC-PACKFILE §1, §9). Readers
90/// accept both `VERSION` and `VERSION_V2`; only `VERSION_V2` packs may
91/// contain `0x03`/`0x04` entries.
92pub const VERSION_V2: u32 = 2;
93
94/// Hard cap on entries (SPEC-PACKFILE §5).
95pub const MAX_ENTRIES: u32 = 10_000_000;
96/// Hard cap on the sum of payload bytes across all entries.
97pub const MAX_TOTAL_PAYLOAD: u64 = 4 * 1024 * 1024 * 1024;
98/// Trailer is a 32-byte raw BLAKE3 digest.
99pub const TRAILER_LEN: usize = 32;
100
101/// Header is `[4B magic][4B version][4B entry_count]`.
102pub const HEADER_LEN: usize = 4 + 4 + 4;
103/// Per-entry framing overhead is `[1B type][4B payload_len]`.
104pub const ENTRY_FRAME_LEN: usize = 1 + 4;
105/// Byte offset of the 4-byte `version` field within the header —
106/// right after the 4-byte magic. `PackWriter::finish` patches this
107/// once the final v1-vs-v2 decision is known (mirrors
108/// `ENTRY_COUNT_OFFSET` below).
109pub const VERSION_OFFSET: usize = 4;
110/// Byte offset of the 4-byte `entry_count` field within the header —
111/// after the 4-byte magic and 4-byte version fields. Found hardcoded
112/// as the literal range `8..12` at four call sites during the
113/// epic-#634 code review; named here instead, consistent with this
114/// file's existing `HEADER_LEN`/`TRAILER_LEN` convention.
115pub const ENTRY_COUNT_OFFSET: usize = 8;
116
117/// Compression candidates shorter than this are never compressed
118/// (SPEC-PACKFILE §3.3) — per-entry zstd framing overhead and CPU
119/// cost isn't worth it for tiny payloads. Only meaningful when the
120/// `pack-zstd` feature is compiled in (see `maybe_compress`).
121#[cfg(feature = "pack-zstd")]
122const MIN_COMPRESS_LEN: usize = 64;
123/// zstd compression level `PackWriter` uses for `0x03`/`0x04` entries.
124/// The library default (`ZSTD_CLEVEL_DEFAULT`); no benchmark evidence
125/// in issue #646 justified deviating from it.
126#[cfg(feature = "pack-zstd")]
127const ZSTD_LEVEL: i32 = 3;
128/// Byte length of a `0x03`/`0x04` entry's `uncompressed_len` length
129/// prefix. Part of the wire format regardless of whether this build
130/// can itself produce/consume `0x03`/`0x04` entries.
131const ZSTD_LEN_PREFIX: usize = 4;
132
133/// Packfile errors. Distinct from [`MkitError`] so callers can match on
134/// pack-specific failures (trailer mismatch, base-missing) without
135/// catching every object decode error.
136#[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    /// `0x03`/`0x04` payload shorter than its `[uncompressed_len]`
172    /// length-prefix header (SPEC-PACKFILE §3.3, §3.4) — distinct from
173    /// `DeltaEntryTruncated`, which covers the 32-byte `base_hash`
174    /// prefix a `0x04` entry has in front of this.
175    #[error("zstd entry payload is shorter than its length-prefix header")]
176    ZstdEntryTruncated,
177    /// Claimed `uncompressed_len` exceeds [`MAX_RAW_OBJECT_SIZE`] —
178    /// rejected before any decompression allocation is attempted
179    /// (SPEC-PACKFILE §3.3 bomb-guarding).
180    #[error(
181        "zstd entry's claimed decompressed size {0} exceeds the {MAX_RAW_OBJECT_SIZE}-byte cap"
182    )]
183    DecompressedSizeOverCap(usize),
184    /// The zstd frame decompressed successfully but produced a
185    /// different byte count than the entry's claimed
186    /// `uncompressed_len` (SPEC-PACKFILE §3.3 bomb-guarding).
187    #[error("zstd entry claims {0} decompressed bytes but produced {1}")]
188    DecompressedSizeMismatch(usize, usize),
189    /// The zstd frame itself is corrupt / not a valid zstd stream.
190    #[error("zstd decompression failed: {0}")]
191    ZstdDecompress(String),
192    /// [`PackWriter::new_raw_only`] refuses delta entries (`push_delta` /
193    /// `push_prepared_delta`). Compression is skipped rather than
194    /// rejected on the raw path.
195    #[error("pack writer is in raw-only mode and does not accept delta entries")]
196    RawOnly,
197}
198
199/// Result of an unpack: which entries were stored, plus a count of
200/// delta resolutions vs raw writes. Useful for transport/CLI summaries.
201#[derive(Debug, Clone, Default, PartialEq, Eq)]
202pub struct UnpackReport {
203    pub raw_count: u32,
204    pub delta_count: u32,
205    /// Hashes inserted into the store this unpack call.
206    pub stored: Vec<Hash>,
207}
208
209/// A raw entry whose pack-compression decision has already been made by
210/// [`PackWriter::prepare_raw`], off any [`PackWriter`] instance — the
211/// CPU-bound step, safe to run in parallel across entries before any of
212/// them touch the writer's sequential state. Feed it to
213/// [`PackWriter::push_prepared_raw`] to append it in order.
214#[derive(Debug)]
215pub struct PreparedRaw {
216    hash: Hash,
217    bytes: Vec<u8>,
218    frame: Option<Vec<u8>>,
219}
220
221impl PreparedRaw {
222    /// The object hash this entry was prepared for — lets a caller
223    /// batching many entries (and a test asserting order-preservation
224    /// across that batching) identify which input produced which
225    /// prepared output without re-deriving it.
226    #[must_use]
227    pub fn hash(&self) -> Hash {
228        self.hash
229    }
230
231    /// Conservative (uncompressed) wire-size bound for this entry —
232    /// the same quantity a caller would use, pre-compression, to decide
233    /// whether pushing it risks exceeding a payload cap (compression
234    /// only ever shrinks the actual wire size, never grows it).
235    #[must_use]
236    pub fn conservative_len(&self) -> usize {
237        self.bytes.len()
238    }
239}
240
241/// The delta-entry counterpart of [`PreparedRaw`], produced by
242/// [`PackWriter::prepare_delta`] and appended via
243/// [`PackWriter::push_prepared_delta`].
244#[derive(Debug)]
245pub struct PreparedDelta {
246    base: Hash,
247    stream: Vec<u8>,
248    frame: Option<Vec<u8>>,
249}
250
251impl PreparedDelta {
252    /// The delta base hash this entry was prepared against — see
253    /// [`PreparedRaw::hash`].
254    #[must_use]
255    pub fn base(&self) -> Hash {
256        self.base
257    }
258
259    /// Conservative (uncompressed) wire-size bound for this entry — see
260    /// [`PreparedRaw::conservative_len`].
261    #[must_use]
262    pub fn conservative_len(&self) -> usize {
263        hash::HASH_LEN + self.stream.len()
264    }
265}
266
267/// Builds a packfile, enforcing entry/payload caps and streaming each
268/// pushed entry's frame directly into the final output buffer as it
269/// arrives. [`Self::finish`] only patches the header's entry count
270/// (unknown up front from a streaming writer) and appends the trailer —
271/// it never re-copies the pushed entries into a second, same-sized
272/// buffer (issue #647).
273#[derive(Debug)]
274pub struct PackWriter {
275    // The final packfile bytes, built incrementally: `new` writes the
276    // header with a zero entry-count placeholder (patched by `finish`
277    // once the final count is known); `push_raw`/`push_delta` append
278    // each entry's `[type][len][payload]` frame directly here. There is
279    // no separate per-entry collection copied a second time at
280    // `finish`.
281    buf: Vec<u8>,
282    entry_count: u32,
283    total_payload: u64,
284    // Set the first time `push_raw`/`push_delta` emits a `0x03`/`0x04`
285    // entry. `finish` reads this to decide the header's `version`
286    // field (SPEC-PACKFILE §1's writer version-selection rule) — v2
287    // the moment ANY entry ended up compressed, v1 otherwise.
288    has_compressed_entry: bool,
289    // When true, `push_raw` never compresses, `push_delta` /
290    // `push_prepared_delta` return [`PackError::RawOnly`], and `finish`
291    // always emits a v1 pack of `0x00` entries. The closure profile
292    // (SPEC-DISCLOSURE) is the consumer: a wasm verifier has no zstd.
293    raw_only: bool,
294}
295
296impl Default for PackWriter {
297    fn default() -> Self {
298        Self::new()
299    }
300}
301
302impl PackWriter {
303    /// Create an empty writer.
304    #[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()); // entry_count placeholder; `finish` patches it in.
310        Self {
311            buf,
312            entry_count: 0,
313            total_payload: 0,
314            has_compressed_entry: false,
315            raw_only: false,
316        }
317    }
318
319    /// A writer that never compresses and never accepts deltas.
320    ///
321    /// Every `push_raw` emits a `0x00` entry even when the `pack-zstd`
322    /// feature is on and the payload is highly compressible.
323    /// `push_delta` / `push_prepared_delta` return [`PackError::RawOnly`].
324    /// [`Self::finish`] always emits a v1 pack. This is the closure
325    /// profile's carrier (SPEC-DISCLOSURE): a wasm verifier is built
326    /// without `pack-zstd` and can only consume raw v1 packs.
327    #[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    /// Append a raw object entry. `bytes` is the fully serialised object
335    /// payload; `hash_of_bytes` is the BLAKE3 of those same bytes —
336    /// callers usually have it on hand from the object store, so we take
337    /// it explicitly to avoid an extra BLAKE3 pass over the same buffer.
338    /// Takes `bytes` by reference (not by value): the streaming writer
339    /// copies it straight into the output buffer as it's pushed, so it
340    /// never needs to own the caller's copy (issue #647). Returns the
341    /// same hash for chaining.
342    ///
343    /// Applies the SPEC-PACKFILE §3.3 compression policy transparently:
344    /// if `bytes` compresses under zstd strictly smaller on the wire
345    /// (and is long enough to bother, see `MIN_COMPRESS_LEN`), this
346    /// emits a `0x03` zstd-raw entry instead of `0x00` raw — callers
347    /// never need to opt in. Either way the returned/stored identity
348    /// (`hash_of_bytes`) is unchanged; only the wire encoding differs.
349    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    /// Pure (no `&self`) compression step for a raw entry, split out of
359    /// [`Self::push_raw`] so a caller building a pack out of many
360    /// independent objects (e.g. a push serialising a whole plan) can
361    /// run the CPU-bound zstd compression for each entry off the
362    /// thread-pool of its choosing — in parallel, since one entry's
363    /// compression never depends on another's — and only then replay
364    /// the writer's sequential, order-preserving append via
365    /// [`Self::push_prepared_raw`].
366    ///
367    /// Takes `bytes` by value (unlike `push_raw`'s `&[u8]`): callers
368    /// with a single object already own the bytes fresh out of the
369    /// object store, and a fan-out across a thread pool needs an owned,
370    /// `'static` value to move into each task anyway, so there is no
371    /// zero-copy path to preserve here the way `push_raw` does.
372    #[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    /// Append a [`PreparedRaw`] produced by [`Self::prepare_raw`].
383    /// Identical wire result and cap-check semantics to `push_raw`
384    /// called on the same bytes — compression already happened, so
385    /// this only replays the cheap bookkeeping + buffer append.
386    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    /// Shared tail of `push_raw`/`push_prepared_raw`: given `bytes` and
392    /// an already-decided `frame` (zstd output, or `None` when
393    /// compression wasn't worth it), do the cap check, bookkeeping, and
394    /// buffer append. The only difference between the two public
395    /// entry points is where `frame` was computed.
396    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    /// Append a delta entry. `base_hash` MUST refer to an earlier raw
422    /// entry in this pack OR an object already in the destination store.
423    /// `delta_stream` MUST be a valid SPEC-DELTA stream — we don't
424    /// re-validate here (the writer is trusted), but the reader will.
425    ///
426    /// Applies the same §3.3 compression policy as [`Self::push_raw`],
427    /// but ONLY to `delta_stream` — `base_hash` is always written
428    /// uncompressed (SPEC-PACKFILE §3.4), so ordering/base-discovery
429    /// logic never needs to decompress anything. Emits `0x04`
430    /// zstd-delta when the stream compresses strictly smaller on the
431    /// wire and is long enough to bother; `0x02` delta otherwise.
432    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    /// Pure (no `&self`) compression step for a delta entry — the
441    /// delta-entry counterpart of [`Self::prepare_raw`]; see its doc
442    /// comment for why this exists and how it's meant to be used (fan
443    /// out `prepare_delta` across a thread pool, then replay results in
444    /// order via [`Self::push_prepared_delta`]).
445    #[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    /// Append a [`PreparedDelta`] produced by [`Self::prepare_delta`].
456    /// Identical wire result and cap-check semantics to `push_delta`
457    /// called on the same base/stream.
458    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    /// Shared tail of `push_delta`/`push_prepared_delta` — see
466    /// [`Self::append_raw_frame`]'s doc comment for the analogous raw-entry
467    /// split.
468    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    /// Append one entry's frame — `[1B type][4B payload_len][payload]`
502    /// — straight onto the output buffer. `parts` is the payload split
503    /// into its logical pieces (a delta entry is `[base_hash][stream]`)
504    /// so no intermediate concatenated buffer is ever built just to
505    /// hand a single contiguous slice to `finish`.
506    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    /// Number of entries pushed so far. Useful for sizing diagnostics.
532    #[must_use]
533    pub fn entry_count(&self) -> usize {
534        self.entry_count as usize
535    }
536
537    /// Sum of wire payload bytes pushed so far — the quantity the
538    /// writer's own internal cap check compares against
539    /// [`MAX_TOTAL_PAYLOAD`]. Measured post-compression (SPEC-PACKFILE
540    /// §5): each `push_raw`/`push_delta` call adds the *wire* payload
541    /// length, not the caller's uncompressed input length. Callers
542    /// deciding whether to seal a pack before pushing another entry can
543    /// use a conservative uncompressed-length estimate against this
544    /// value — the actual wire cost is never more than that estimate,
545    /// since compression is only ever applied when it's strictly
546    /// smaller (see `maybe_compress`).
547    #[must_use]
548    pub fn total_payload(&self) -> u64 {
549        self.total_payload
550    }
551
552    /// Serialise the pack: header + entries + trailer. Entries are
553    /// already in `self.buf` (streamed in by `push_raw`/`push_delta`);
554    /// `finish` patches the header's `version` (SPEC-PACKFILE §1: v2
555    /// iff at least one entry ended up compressed, v1 otherwise) and
556    /// `entry_count`, then appends the trailer,
557    /// `BLAKE3(everything_before_trailer)`. The whole pack's BLAKE3 is
558    /// the on-disk pack key — see [`pack_key`].
559    pub fn finish(self) -> Result<Vec<u8>, PackError> {
560        self.finish_inner(None)
561    }
562
563    /// Test-only variant of [`Self::finish`] that also reports, via
564    /// `bytes_copied`, how many payload bytes it copies WHILE finishing
565    /// (as opposed to while entries were pushed). Proves `finish`
566    /// streams rather than double-buffers (issue #647): the unpatched
567    /// writer re-copied every pushed entry's payload into a fresh
568    /// same-size buffer inside `finish`, so this counter would track
569    /// the whole pack; the streaming writer only ever appends the
570    /// 32-byte trailer here.
571    #[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/// Compute the on-disk pack key: BLAKE3 of the entire packfile bytes
601/// (including the trailer). SPEC-PACKFILE §7. Returns the bare digest;
602/// callers prepend `packs/` and lower-hex-encode for the storage path.
603#[must_use]
604pub fn pack_key(pack_bytes: &[u8]) -> Hash {
605    hash::hash(pack_bytes)
606}
607
608/// SPEC-PACKFILE §3.3 writer compression policy: compress `data` with
609/// zstd and return the frame ONLY if doing so is worth it — `data` is
610/// at least [`MIN_COMPRESS_LEN`] bytes AND the compressed frame plus
611/// its `ZSTD_LEN_PREFIX`-byte length prefix is strictly smaller than
612/// `data` itself. Mirrors `transfer.rs`'s `try_delta` gate's
613/// "strictly smaller or don't bother" posture. Returns `None` (never
614/// an error) on any compression failure or when compression isn't
615/// worth it — compression is a pure wire-size optimization, so a
616/// writer always has a correct fallback (emit the entry uncompressed)
617/// rather than a new failure mode to propagate.
618#[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    // SPEC-PACKFILE §3.3: a reader rejects any `uncompressed_len` over
626    // MAX_RAW_OBJECT_SIZE, so a larger payload is written uncompressed.
627    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    // Per-thread reused `zstd::bulk::Compressor`, keyed off the same
646    // `ZSTD_LEVEL` every call in this build uses. `zstd::bulk::compress`
647    // (the plain free function) allocates and initializes a fresh
648    // `Compressor` — and its underlying `CCtx` — on every single call; on
649    // the push-path's per-entry compression fan-out
650    // (`build_and_upload_packs`'s rayon fan-out over
651    // `prepare_raw`/`prepare_delta`, `pack_build_fanout`) that means one
652    // CCtx alloc/init per object compressed. A one-shot
653    // `Compressor::compress` call carries no state across calls (it is
654    // not a streaming encoder), so reusing the same context across every
655    // object a given thread compresses is behavior-preserving — same
656    // level, same output bytes — and turns that per-object setup cost
657    // into a one-time cost per worker thread.
658    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    // No compression backend compiled in (e.g. mkit-wasm, which opts
667    // out of `pack-zstd` for wasm32-buildability — see mkit-core's
668    // Cargo.toml). Every pack this build writes is a valid v1 pack;
669    // it just never uses the v2-only entry types.
670    None
671}
672
673/// Parse and decompress a `0x03`/`0x04`-style `[uncompressed_len][zstd
674/// frame]` payload, enforcing SPEC-PACKFILE §3.3's bomb-guarding
675/// before any decompression allocation: the claimed length is checked
676/// against [`MAX_RAW_OBJECT_SIZE`] first, decompression is bounded to
677/// that claim, and the actual decompressed length is re-checked
678/// against the claim afterward.
679fn 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/// [`decompress_zstd_entry`] over an explicit backend, so the
690/// differential tests can drive the C and pure-Rust decoders through the
691/// exact same claim / length checks.
692#[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
725/// Reserve and charge before calling this; neither backend grows or zero-fills
726/// the output. Shared by lazy unpack and buffered/window entry decoding.
727fn 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/// RFC 8878 §3.1.1 Zstandard frame magic number, as it appears on the
738/// wire (`0xFD2FB528` little-endian). SPEC-PACKFILE §3.3 allows exactly
739/// one such frame per entry: a skippable frame (`0x184D2A5?`) or a
740/// legacy pre-RFC frame magic is rejected by both backends.
741#[cfg(any(feature = "pack-zstd", feature = "pack-ruzstd"))]
742const ZSTD_FRAME_MAGIC: [u8; 4] = [0x28, 0xB5, 0x2F, 0xFD];
743
744/// Both backends' first check: the payload must open with the one
745/// Zstandard frame magic SPEC-PACKFILE §3.3 permits.
746#[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/// Decompress `frame` — exactly one Zstandard frame, nothing before or
760/// after it — bounding the allocation to `capacity` bytes (already
761/// checked against [`MAX_RAW_OBJECT_SIZE`] by the caller) so a corrupt
762/// or hostile frame can't force an over-large allocation. The C decoder
763/// (`pack-zstd`) is selected whenever it is compiled in.
764#[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    // `ZSTD_decompressDCtx` (behind `bulk::decompress`) would otherwise
772    // decode concatenated frames and skip skippable ones.
773    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/// Without the C library, the pure-Rust decoder serves every read.
811#[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/// Fixed pure-Rust decode window cap: 8 MiB (`windowLog` 23).
832/// This covers the default windows of non-ultra zstd levels. Larger windows
833/// fail closed even if the output claim is larger; native C keeps its policy.
834#[cfg(feature = "pack-ruzstd")]
835#[cfg_attr(all(feature = "pack-zstd", not(test)), allow(dead_code))]
836const RUZSTD_WINDOW_LIMIT: u64 = 8 << 20;
837
838/// Read the base and result lengths from a zstd-compressed SPEC-DELTA v1 header.
839///
840/// `frame` is the zstd frame alone, without the pack frame, base hash or
841/// uncompressed-length prefix. Only the nine-byte delta header is read out;
842/// no allocation is sized from the delta stream's claim or result length.
843/// The pure-Rust backend uses the patched block bound and fixed 8 MiB window.
844/// Its working memory remains within the existing 28 MiB allowance, even when
845/// it retains window history before emitting the prefix. C-only builds use a
846/// streaming decoder with the same window cap. The decoder is dropped here,
847/// before the caller performs its ordinary budgeted decode.
848///
849/// This is a prefix inspection, not full frame or delta validation. Callers
850/// must still decode and verify the complete object after checking metadata.
851///
852/// # Errors
853/// Unsupported compression, malformed/truncated compression or delta header.
854pub 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/// Pure-Rust (`ruzstd`) decode of exactly one Zstandard frame, bounded
910/// to `capacity` output bytes. Compiled whenever `pack-ruzstd` is on,
911/// including alongside `pack-zstd`, so the differential tests can run
912/// both backends over the same inputs.
913///
914/// Matches the C path's accept/reject decisions and error variants:
915/// a declared frame content size must not exceed `capacity` and must
916/// equal the decoded length; output past `capacity` is an error, not a
917/// short read; a content checksum, when present, must match; nothing
918/// may follow the frame. At most `capacity + 1` bytes are ever read out.
919///
920/// Memory: the output is reserved fallibly to the claim without zero-filling.
921/// Unpack charges that reservation before allocation. The pure-Rust decoder's
922/// ring buffer is separate working memory. With the fixed 8 MiB window and
923/// vendored RFC block preflight, working allocations stay below 28 MiB,
924/// including old/new ring allocations during growth and bounded block scratch.
925/// This fixed allowance is outside the owned-payload resident cap, as for the window reader's
926/// carry/output budget. No output allocation grows past the admitted claim.
927///
928/// Allocating adapter for differential tests. Production decode supplies its
929/// already reserved output to `ruzstd_decompress_into` when `pack-zstd` is off,
930/// and to the C `zstd_decompress_into` when it is on.
931#[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    // RFC 8878 §3.1.1.1.1. The header parsed, so the descriptor byte
967    // after the magic exists. ruzstd ignores the reserved bit, which a
968    // decoder must refuse (the C decoder does).
969    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    // A content size is present iff the FCS flag or the single-segment
974    // flag is set.
975    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    // ruzstd's in-memory decoder never yields Interrupted.
989    loop {
990        // Probe one extra byte without allocating past the claim.
991        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/// Re-walk an already-decoded frame's blocks (RFC 8878 §3.1.1.2–3) for
1033/// the reserved fields ruzstd ignores and the C decoder rejects: the
1034/// sequences section's `Symbol_Compression_Modes` reserved bits. Without
1035/// this, a frame the C path refuses would decode here — a fail-open
1036/// divergence between the native and Workers readers.
1037#[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, // RLE: one byte, repeated `size` times
1061            _ => 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/// The `Symbol_Compression_Modes` byte of a compressed block, or `0` when
1075/// the block has no sequences (RFC 8878 §3.1.1.3.1.1, §3.1.1.3.2.1).
1076#[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    // Header fields are assembled in `u64`, never `usize`: the 5-byte
1081    // literals header spans 40 bits, which overflows a 32-bit `usize`
1082    // (wasm32) and silently drops the compressed size's top bits.
1083    let byte = |i: usize| block.get(i).copied().ok_or(MALFORMED).map(u64::from);
1084    let b0 = byte(0)?;
1085    // Literals section: header, then its content.
1086    let (header_len, content_len) = match (b0 & 3, (b0 >> 2) & 3) {
1087        // Raw / RLE literals: 5-, 12- or 20-bit regenerated size.
1088        (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        // Compressed / treeless literals: 10-, 14- or 18-bit sizes.
1097        (_, 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    // At most 20 bits, so it fits any `usize`.
1111    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
1122/// Collect the `base_hash` of every `0x02` delta entry in `pack_bytes`,
1123/// without resolving or storing anything.
1124///
1125/// This lets a caller pre-fetch bases that may live OUTSIDE the pack (e.g.
1126/// objects a legacy per-object remote stored individually) before calling
1127/// [`PackReader::read`], so delta resolution never fails part-way through a
1128/// pack. Raw entries are skipped; duplicates are de-duplicated.
1129///
1130/// Only the header (magic/version) and entry framing are validated — the
1131/// trailer is intentionally NOT verified here, because [`PackReader::read`]
1132/// re-verifies the whole pack (trailer included) before storing anything.
1133///
1134/// # Errors
1135///
1136/// Returns the same framing [`PackError`] variants as [`PackReader::read`]
1137/// for a malformed header or out-of-bounds entry.
1138///
1139/// # Panics
1140///
1141/// The `try_into` calls on fixed 4-byte slices are statically guaranteed by
1142/// the preceding bounds checks; they `expect`-panic only if slice-bounds
1143/// elision is wrong.
1144pub 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    // Entries live between the header and the 32-byte trailer.
1164    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        // `0x04`'s base_hash sits at the same offset (byte 0 of the
1185        // payload) as `0x02`'s and is always uncompressed
1186        // (SPEC-PACKFILE §3.4), so both entry types are scanned the
1187        // same way here — no decompression needed to pre-fetch bases.
1188        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/// Streaming-style packfile reader. Verifies header, trailer, entry
1210/// framing, and the base-before-delta ordering rule. Reconstructs delta
1211/// targets and writes every resolved object to `store`.
1212#[derive(Debug)]
1213pub struct PackReader;
1214
1215impl PackReader {
1216    /// Verify and unpack `pack_bytes` into `store`. Returns counts of
1217    /// raw vs. delta entries plus the list of stored hashes (in pack
1218    /// order, deduped within this call).
1219    ///
1220    /// # Errors
1221    ///
1222    /// Returns the matching [`PackError`] variant on any malformed
1223    /// input or trailer mismatch. The store is not modified if the
1224    /// trailer fails verification.
1225    ///
1226    /// # Panics
1227    ///
1228    /// The internal `try_into` calls on fixed-size byte slices are
1229    /// statically guaranteed to succeed (we slice exactly 4 bytes for
1230    /// every `u32::from_le_bytes`). They `expect`-panic only if the
1231    /// compiler's slice-bounds elision is wrong.
1232    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    /// Same as [`Pack::read`], but with a caller-supplied running-total
1237    /// payload cap instead of the hardcoded [`MAX_TOTAL_PAYLOAD`] (4
1238    /// GiB). `pub(crate)`, not part of the public API — test-only
1239    /// injection point so `PackfileTooLarge` can be exercised without
1240    /// constructing a multi-gigabyte pack. Real callers MUST use
1241    /// [`Pack::read`] instead.
1242    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    /// Test-only variant of [`Self::read`] that also reports, via
1251    /// `owned_bytes`, the total number of payload bytes it allocates
1252    /// freshly (as opposed to borrowing straight from the
1253    /// already-resident `pack_bytes`). Proves the streaming reader
1254    /// (issue #647) never re-copies a raw entry's bytes: only delta
1255    /// targets and successfully staged compressed raw payloads increment
1256    /// this counter. Borrowed raw entries never increment it.
1257    #[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        // Phase 1 retains only borrowed wire frames. Count uses without
1284        // decompressing, and remember the final position naming each base.
1285        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        // Phase 2 keeps the native fan-out. Each worker decompresses its
1301        // current raw frame, stages it, and drops unneeded owned bytes.
1302        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
1315/// Where a delta's *external* base — one that is not an earlier entry
1316/// of the same pack — may come from.
1317///
1318/// SPEC-PACKFILE §3.2's "destination object store" is the implementor:
1319/// the local [`ObjectStore`] on a client, and on a server the pushing
1320/// repository's own membership. A server MUST NOT back this with a
1321/// global content store or another repository's objects: whether a
1322/// delta resolves would then reveal that some other repository holds
1323/// the base (PRD §6.5 repository isolation, no existence oracles).
1324///
1325/// The seam is generic, never `dyn`: [`PackReader::read`] instantiates
1326/// it with `&ObjectStore` so the hot unpack path stays monomorphic.
1327pub trait DeltaBaseSource {
1328    /// Whether [`Self::base`] returns bytes already verified to be the
1329    /// object named by the requested id (the store's `read` verifies).
1330    ///
1331    /// When `false` (the default), the decoder deserializes the bytes
1332    /// and re-derives the id itself; anything that is not a storable
1333    /// canonical object with exactly the requested id is treated as
1334    /// absent ([`PackError::DeltaBaseMissing`]), so an untrusted source
1335    /// can never smuggle a different object in as a base. Set it to
1336    /// `true` only when `base` itself guarantees that identity, as
1337    /// `ObjectStore::read` does; the decoder then skips the re-derive.
1338    const VERIFIED: bool = false;
1339
1340    /// Canonical object bytes for `id`, or `None` when this source may
1341    /// not provide it. "Not permitted" and "does not exist" MUST both
1342    /// be `None`: the decoder reports them identically, by design.
1343    ///
1344    /// # Errors
1345    ///
1346    /// A source failure (I/O, backend error) that is not an answer
1347    /// about `id`. It aborts the decode.
1348    fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError>;
1349
1350    /// Fetch with allocation admission. The default charges after `base`
1351    /// returns; sources that know the length first should override this to
1352    /// call `admit` before allocating, as `&ObjectStore` does. Before
1353    /// returning `Some(bytes)`, every override must successfully invoke
1354    /// `admit` exactly once with `bytes.len()` and propagate its error.
1355    ///
1356    /// # Errors
1357    /// Returns a source failure or the error from `admit`.
1358    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/// A [`DeltaBaseSource`] with no external bases: a pack must be
1372/// self-contained (every delta's base an earlier entry). Used for
1373/// closure-profile packs, pack rewrites and tests.
1374#[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    /// `ObjectStore::read` BLAKE3/BMT-verifies the bytes against `id`
1385    /// and fails with `HashMismatch` otherwise.
1386    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/// One entry handed to [`decode_entries_with`]'s sink, in pack order.
1417#[derive(Debug)]
1418#[non_exhaustive]
1419pub struct DecodedEntry<'a> {
1420    /// The entry's object id, re-derived from `bytes` (BLAKE3, or the
1421    /// BMT root for `Tree`/`ChunkedBlob`).
1422    pub id: Hash,
1423    /// Canonical object bytes: the raw payload (decompressed for
1424    /// `0x03`), or the reconstructed delta target.
1425    pub bytes: &'a [u8],
1426    /// `bytes` decoded, so the consumer need not deserialize again.
1427    pub object: Object,
1428    /// `true` for a `0x02`/`0x04` delta target, `false` for a raw entry.
1429    pub from_delta: bool,
1430    /// Offset of the complete encoded frame in the source pack.
1431    pub frame_offset: u64,
1432    /// Length of the complete encoded frame, including type and length.
1433    pub frame_length: u64,
1434    /// Encoded frame type (`0x00`, `0x02`, `0x03`, or `0x04`).
1435    pub wire_type: u8,
1436    /// Delta base id, if this frame is a delta.
1437    pub delta_base: Option<Hash>,
1438}
1439
1440/// Summary of a successful [`decode_entries_with`] call.
1441#[derive(Debug, Default, Clone, PartialEq, Eq)]
1442#[non_exhaustive]
1443pub struct DecodeReport {
1444    /// Raw (`0x00`/`0x03`) entries decoded.
1445    pub raw_count: usize,
1446    /// Delta (`0x02`/`0x04`) entries reconstructed.
1447    pub delta_count: usize,
1448    /// Every entry's id, in pack order (a repeated entry repeats here,
1449    /// as in [`UnpackReport::stored`]).
1450    pub ids: Vec<Hash>,
1451}
1452
1453/// Resource limits for [`decode_entries_with`].
1454///
1455/// A pack's wire size says little about what decoding it allocates: a
1456/// `0x03`/`0x04` entry may claim up to [`MAX_RAW_OBJECT_SIZE`] bytes from a
1457/// few-byte zstd frame, a delta may declare a result of the same size
1458/// from a short run of copies, and a tiny delta may name a huge external
1459/// base. A decoder that held all of that until the pack ends would let a
1460/// sub-kilobyte pack pin gigabytes. [`decode_entries_with`] therefore
1461/// charges every such allocation against [`Self::max_decoded_bytes`] and
1462/// rejects the pack with [`PackError::PackfileTooLarge`] once the total
1463/// passes it.
1464///
1465/// The budget is **per call**. It bounds one decode, not a process: a
1466/// server running several decodes at once (WP-4.7) must size it from its
1467/// isolate's memory limit divided by its decode concurrency. The default,
1468/// [`Self::DEFAULT_MAX_DECODED_BYTES`] (1 GiB), equals the largest object
1469/// the format allows ([`MAX_RAW_OBJECT_SIZE`]) so that any single valid
1470/// object decodes; it is not sized for a constrained isolate.
1471#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1472#[non_exhaustive]
1473pub struct DecodeLimits {
1474    /// Cap on the bytes one decode may hold beyond the pack itself: every
1475    /// `0x03`/`0x04` entry's claimed `uncompressed_len`, every delta's
1476    /// declared result length, and every external base fetched from the
1477    /// [`DeltaBaseSource`] while it is cached. `0x00` payloads are borrowed
1478    /// from the pack and do not count. An external base is measured once
1479    /// fetched, so the peak can pass the cap by the one base whose fetch
1480    /// trips it.
1481    pub max_decoded_bytes: u64,
1482    entry_geometry: Option<(u64, u64)>,
1483}
1484
1485impl DecodeLimits {
1486    /// Default [`Self::max_decoded_bytes`]: 1 GiB, the maximum object size
1487    /// ([`MAX_RAW_OBJECT_SIZE`]). Servers set their own (see the type docs).
1488    pub const DEFAULT_MAX_DECODED_BYTES: u64 = MAX_RAW_OBJECT_SIZE as u64;
1489
1490    /// These limits with [`Self::max_decoded_bytes`] set to `bytes`.
1491    #[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    /// Bound encoded payloads and delta streams separately from canonical bytes.
1497    #[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/// Running total of a decode's claimed allocations against
1530/// [`DecodeLimits::max_decoded_bytes`].
1531// Distinct from ResidentBudget: compressed and delta claims accumulate over
1532// the entire storeless decode; only external-base charges are credited on last
1533// use, matching DecodeLimits. PackReader uses the peak resident cap instead.
1534#[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
1550/// The decode path's [`DeltaBaseSource`]: forwards to the caller's source
1551/// and charges each base it returns against the decode budget, remembering
1552/// the charge so [`Self::release`] can credit it back once no later delta
1553/// needs that base. [`PackReader::read`] never uses this.
1554struct 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    /// Credit back the charge for external base `id`, if it had one.
1585    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
1592/// Little-endian `u32` at `at` in `bytes`, if all four bytes are there.
1593fn 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
1598/// Charge every `0x03`/`0x04` entry's claimed decompressed size, reading
1599/// only its uncompressed-length prefix. Runs after [`PackEntries::new`]
1600/// has validated the framing (every frame lies inside the body, which is
1601/// why the bounds here can only fall short on a malformed pack) and
1602/// before anything is decompressed. A payload too short to carry its
1603/// prefix is left for the drain to reject.
1604fn 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/// A store-less pack decoder that pauses at an unresolved external base.
1640///
1641/// [`Self::new`] performs the same framing and allocation preflight as
1642/// [`decode_entries_with`]. Each successful entry reaches the sink exactly
1643/// once across calls to [`Self::resume`]; a missing base leaves its frame at
1644/// the cursor so the caller can supply that base and resume. Other errors
1645/// end the decode. The same cumulative budget spans every resume call.
1646#[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    /// Preflight one content-addressed pack without fetching external bases.
1659    ///
1660    /// # Errors
1661    /// Invalid framing, unsupported content, or a decoded-size claim over
1662    /// `limits` is returned before the sink observes any entry.
1663    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        // As in decode_entries_with, all delta result claims precede the
1699        // first object validation; base charges are added only on use.
1700        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    /// Set the remaining decode's allocation cap. A caller retaining
1725    /// external objects beside this cursor can lower the cap as those
1726    /// objects accumulate. Earlier claims and live bases stay charged.
1727    ///
1728    /// # Errors
1729    /// [`PackError::PackfileTooLarge`] if bytes already charged exceed
1730    /// `max`; the previous cap is kept in that case.
1731    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    /// Continue from the first unprocessed frame.
1740    ///
1741    /// # Errors
1742    /// [`PackError::DeltaBaseMissing`] leaves the current frame ready to
1743    /// retry after `bases` gains that id. Any other error is terminal. The
1744    /// sink may have seen earlier entries when either error is returned.
1745    #[allow(clippy::too_many_lines)] // One stateful loop preserves frame order and charges.
1746    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                    // A source may return invalid bytes, which is publicly
1793                    // indistinguishable from absence. Undo that attempted
1794                    // base's charge before a caller retries with good bytes.
1795                    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
1860/// Store-less decode of `pack` over an explicit [`DeltaBaseSource`].
1861///
1862/// Validates the pack exactly as [`PackReader::read`] does — header,
1863/// trailer, caps and framing via [`PackEntries`], every entry drained
1864/// (and `0x03`/`0x04` decompressed) before any entry is judged — then,
1865/// in pack order, validates each raw payload as a storable canonical
1866/// object, resolves each delta against an earlier entry or else
1867/// `bases`, validates the reconstructed target, and hands every entry to
1868/// `sink`, without writing to the store. Given `&ObjectStore` as `bases`
1869/// and non-binding limits, it decodes the same valid objects as
1870/// `PackReader::read`. Its cumulative-budget and decompression preflight
1871/// precede object validation; lazy unpack instead reports errors as their
1872/// pack positions are reached, so malformed packs can differ in error order.
1873///
1874/// Memory is bounded by `limits` (see [`DecodeLimits`]): every
1875/// compressed entry's claimed size is charged before anything is
1876/// decompressed, every delta's declared result length before any delta
1877/// is applied, and every external base as it is fetched. An entry or an
1878/// external base stays resident only until the last delta that names it
1879/// as its base; an external base's charge is then credited back.
1880///
1881/// External bases are fetched only for a delta whose base is not an
1882/// earlier entry, at most once per distinct base while it is needed, and
1883/// never recursively: a base is a canonical object, never a pack-only
1884/// delta.
1885///
1886/// # Errors
1887///
1888/// [`PackError::PackfileTooLarge`] when the decoded size passes
1889/// `limits.max_decoded_bytes`: claims are checked ahead of every
1890/// per-entry error except framing, an external base when the delta that
1891/// names it is reached. After preflight, the first [`PackError`] in pack
1892/// order, or the first error `sink` returns. `sink` may already have seen earlier
1893/// entries when an error is returned; a consumer staging them must
1894/// discard that staging.
1895pub 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
1904/// Decode one complete encoded frame, using the same payload parser and
1905/// repository-supplied base contract as [`decode_entries_with`]. The caller
1906/// authenticates the containing pack and its frame offset separately.
1907///
1908/// # Errors
1909/// Invalid framing, unsupported version, missing or invalid base, decoded
1910/// object, or a result beyond `limits`.
1911pub 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
1957/// Resolve one already-framed entry, such as a [`window::WindowReader`]
1958/// yield, into its canonical bytes and id, with the payload rules of
1959/// [`decode_frame_with`]: a delta's base comes only from `bases`, and the
1960/// result is a storable canonical object bounded by `limits`.
1961///
1962/// # Errors
1963/// A missing or invalid base, an invalid object, or a result beyond `limits`.
1964pub 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
1995/// Replay staging results in pack order; deferred admissions are strict here.
1996fn 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    // Phase 3 exposes bases only at their pack position and releases
2009    // them immediately after their last use, including cached store bases.
2010    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
2049/// Owned payload cap: `max(2 * MAX_RAW_OBJECT_SIZE, 16 * pack_len)`.
2050/// The 2 GiB floor covers a maximum-size base plus target; a compressed delta
2051/// also charges its decoded stream, so a 1 GiB store base plus 1 GiB target
2052/// from a small pack is refused while that stream is resident.
2053/// Saturation avoids arithmetic traps on 32-bit hosts.
2054fn 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
2060/// Peak owned payload bytes, released with each reservation. Unlike
2061/// `DecodeBudget`, this applies to `PackReader::read` and bounds live
2062/// compressed output, delta targets and external bases rather than total work.
2063struct 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    /// Only a failed admission is deferred; allocation and decode errors are permanent.
2173    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
2211// None marks deferred admission or work skipped after an earlier failure.
2212type RawStageResult<'b> = Option<Result<(Hash, Option<OwnedBytes<'b>>), PackError>>;
2213type RawStageResults<'b> = Vec<RawStageResult<'b>>;
2214
2215/// Validate, hash and stage independent raw frames, preserving result order.
2216/// Errors stay in the result queue until phase 3 reaches that pack position.
2217fn 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
2337/// SPEC-PACKFILE §1/§5/§8 steps 1-5: length sanity, magic, version,
2338/// trailer verification (BEFORE anything touches the store), and
2339/// entry-count cap. Returns `(version, split, count)` — `split` is the
2340/// byte offset where the trailer begins (entries live in
2341/// `pack_bytes[HEADER_LEN..split]`). Split out of
2342/// [`PackReader::read_inner`] purely to keep that function's
2343/// entry-parsing loop under clippy's line-count cap; no check moves,
2344/// reorders, or changes behavior.
2345fn validate_pack_header(pack_bytes: &[u8]) -> Result<(u32, usize, u32), PackError> {
2346    // 1. Length sanity: must fit header + trailer at minimum.
2347    if pack_bytes.len() < HEADER_LEN + TRAILER_LEN {
2348        return Err(PackError::PackfileTooShort);
2349    }
2350    // 2. Magic.
2351    if &pack_bytes[..4] != MAGIC.as_slice() {
2352        return Err(PackError::InvalidMagic);
2353    }
2354    // 3. Version. v1 and v2 both decode here; only v2 packs may
2355    // contain `0x03`/`0x04` entries (enforced per-entry by the caller).
2356    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    // 4. Trailer must match BEFORE we touch the store. SPEC-PACKFILE §8.
2361    // Every entry parsed by the caller is staged into its batch as
2362    // it's seen (not buffered up front), but that's only safe to do
2363    // BECAUSE the pack's own integrity is already established here,
2364    // first — a corrupt/truncated pack is rejected before a single
2365    // byte is staged, so the "abort leaves the store untouched"
2366    // guarantee does not depend on holding the whole pack's staged
2367    // output in memory at once (see `WriteBatch`'s module docs: a
2368    // dropped, uncommitted batch unlinks its temp files for free).
2369    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    // 5. Entry count + cap.
2377    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    // Quick lower bound sanity: each entry is at least ENTRY_FRAME_LEN bytes.
2386    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/// One decoded pack entry, with `0x03`/`0x04` already decompressed into
2394/// the matching uncompressed variant when the `pack-zstd` feature is
2395/// compiled in.
2396///
2397/// The closure profile (SPEC-DISCLOSURE) accepts only [`Self::Raw`]
2398/// produced from a `0x00` wire type — see [`PackEntries::is_raw_only`].
2399#[derive(Debug)]
2400pub enum PackEntry<'a> {
2401    /// A fully serialised mkit object (`0x00`, or decompressed `0x03`).
2402    /// Raw (`0x00`) payloads borrow the pack bytes; decompressed `0x03`
2403    /// payloads own a buffer.
2404    Raw { bytes: Cow<'a, [u8]> },
2405    /// A delta (`0x02`, or decompressed `0x04`) against `base`.
2406    Delta { base: Hash, stream: Cow<'a, [u8]> },
2407}
2408
2409/// Store-less iterator over a packfile's entries.
2410///
2411/// [`Self::new`] validates the header, trailer, entry-count cap, and
2412/// the running payload-sum cap *before* yielding anything, and scans
2413/// entry types (without decompressing) so [`Self::is_raw_only`] is
2414/// known up front. Iteration then walks the same frames, decompressing
2415/// `0x03`/`0x04` when `pack-zstd` is compiled in. Canonical-object
2416/// validation of raw payloads is left to the consumer
2417/// ([`PackReader::read`] still rejects non-storable objects before
2418/// they touch the store).
2419///
2420/// [`PackReader::read`] consumes the same private frame parser, deferring
2421/// decompression until validation/staging or delta application.
2422#[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    /// Total entry count declared in the pack header — the header
2438    /// field this iterator validates every position against (`self.pos
2439    /// == self.count`⇒ done), exposed so a caller collecting every
2440    /// entry into a `Vec` up front (as [`PackReader::read_inner`]'s
2441    /// phase 1 does) can size it exactly instead of growing it one
2442    /// `push` at a time.
2443    pub(crate) fn entry_count(&self) -> usize {
2444        self.count as usize
2445    }
2446
2447    /// Validate `bytes` as a packfile and prepare to iterate its entries.
2448    ///
2449    /// # Errors
2450    ///
2451    /// The same framing [`PackError`] variants as [`PackReader::read`]:
2452    /// short input, bad magic/version/trailer, over-cap entry count or
2453    /// payload sum, unknown entry type, trailing data.
2454    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    /// True iff every entry is wire type `0x00` (including the empty
2537    /// pack). Computed during [`Self::new`] by scanning entry types
2538    /// without decompressing — a wasm verifier can reject a non-raw
2539    /// pack before touching zstd.
2540    #[must_use]
2541    pub fn is_raw_only(&self) -> bool {
2542        self.raw_only
2543    }
2544
2545    /// Index of the first non-`0x00` entry, if any.
2546    #[must_use]
2547    pub fn first_non_raw_index(&self) -> Option<u32> {
2548        self.first_non_raw
2549    }
2550
2551    /// Byte range of the payload returned by the most recent successful
2552    /// iteration, relative to the original pack buffer. `None` before the
2553    /// first item. For a raw-only pack, this is the borrowed object slice;
2554    /// compressed entries, which are rejected by the closure profile, still
2555    /// report the encoded payload range rather than the decompressed buffer.
2556    #[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
2593// One type/payload parser shared by buffered iteration, lazy unpack and windows.
2594fn 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
2633// Framing is checked by each caller, payload interpretation is shared.
2634fn 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
2656/// Apply a delta while charging the stream, cached base and declared target
2657/// before allocating. The use count is consumed before retaining the target,
2658/// which also handles identity deltas whose target hash equals their base.
2659fn 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
2703/// Resolve a delta's base (in-pack first, then the external
2704/// [`DeltaBaseSource`]) and decode `stream` against it. Shared by the
2705/// `0x02` and `0x04` branches of [`PackReader::read_inner`] and by
2706/// [`decode_entries_with`] — `0x04` differs only in how `stream` was
2707/// sourced (decompressed vs. borrowed straight from `pack_bytes`), not
2708/// in how base resolution or delta decoding work.
2709fn 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    // Resolve base: in-pack first, then the external source. An
2716    // externally resolved base is cached into `in_pack` under its own
2717    // hash so a later delta entry referencing the same out-of-pack base
2718    // hits the cache-hit branch above instead of paying another full
2719    // read + verify + decode (#643). This is safe because
2720    // `external_base` has already established that the bytes are the
2721    // object `base_hash` names (the store's `read` hash-verifies; an
2722    // unverified source is re-derived there). Cloning once here (vs.
2723    // #643's original Arc::clone) is the cost of composing with #647's
2724    // Cow-based `in_pack`, which trades that one-time clone for
2725    // zero-copy borrows on the far more common raw-entry path — a net
2726    // win, and this clone only happens once per unique out-of-pack
2727    // base, not per delta entry.
2728    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
2741/// Fetch an external delta base from `bases` and establish that it is the
2742/// storable canonical object `id` names, or report it absent.
2743///
2744/// A [`DeltaBaseSource::VERIFIED`] source (the local [`ObjectStore`])
2745/// already hash-verified the bytes, so this runs exactly the pre-seam
2746/// store path: a non-storable object is a loud error. Any other source
2747/// is untrusted for identity: bytes that do not deserialize to a
2748/// storable object whose re-derived id is `id` are treated as absent,
2749/// so the caller reports [`PackError::DeltaBaseMissing`] — the same
2750/// error bytes as a base the source does not have at all.
2751fn 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
2770/// Decode `bytes`, enforce the size and storability invariants, and hand back
2771/// the decoded [`Object`] so callers can address it without decoding twice.
2772fn 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// =========================================================================
2799// Tests
2800// =========================================================================
2801
2802#[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        // Use the serialize/object stack so the bytes are a real mkit
2818        // object — important because `store.write` accepts any bytes
2819        // but unpack-time delta apply produces what serialize would.
2820        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    /// Deterministic, high-entropy filler for tests that specifically
2833    /// exercise UNCOMPRESSED-entry behavior (zero-copy borrows, exact
2834    /// on-wire cap arithmetic) and therefore need payloads the §3.3
2835    /// writer policy will decline to compress. A simple LCG byte
2836    /// stream is enough: zstd's LZ+entropy stages find no exploitable
2837    /// redundancy in it, unlike a repeated-byte or short-period
2838    /// pattern, so `4 + compressed_len < raw_len` never holds and
2839    /// these payloads stay `0x00`/`0x02` exactly as before this
2840    /// change. (Payloads elsewhere in this file that use a short
2841    /// repeating pattern are fine to leave as-is — those tests assert
2842    /// only functional round-trip correctness, not wire-format byte
2843    /// counts, so whether they end up compressed doesn't affect them.)
2844    fn incompressible_bytes(seed: u64, len: usize) -> Vec<u8> {
2845        let mut buf = vec![0u8; len];
2846        let mut state = seed | 1; // odd seed keeps the LCG full-period
2847        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            // Also exercise the iterator's defensive check independently
2902            // of the eager framing validation in its public constructor.
2903            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        // Each target fits simultaneously resident memory, but all three
2973        // declared targets together exceed the cumulative DecodeLimits.
2974        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            // Hold an in-flight worker's entire allowance until every other
3087            // worker has attempted admission, avoiding scheduler-dependent races.
3088            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    /// Build a canonical large Blob using a tiny COPY stream without ever
3240    /// allocating its target in the fixture. The last byte distinguishes IDs.
3241    #[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                    // Consume each large target exactly once into a tiny
3282                    // object, so it is released before the next expansion.
3283                    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            // Inject a lower cap to exercise the same production admission
3312            // check without allocating several GiB in the test suite.
3313            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    /// Copy of the parent reader's ownership/ordering algorithm, deliberately
3366    /// simple and unbounded; used only with small equivalence fixtures.
3367    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        // Raw duplicate after raw_base's final delta use must not re-retain raw_base.
3436        writer.push_raw(object_hash(&raw_base), &raw_base).unwrap();
3437        // Duplicate delta targets and identity targets share raw_base single key.
3438        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        // An insufficient base budget fails inside the store allocator,
3538        // before any buffer exists, and leaves the cache empty.
3539        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        // Model transfer-planner order: each raw base is immediately
3603        // followed by its delta, alternating binary and text-like files.
3604        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        // Require at least 200 MiB on the wire as well as reconstructed
3624        // content; highly compressed entries cannot satisfy this by claim.
3625        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        // Verify every stored object through the real store identity check.
3646        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        // clone/fetch receive N objects per pack; durability must cost
3682        // O(1) full flushes per pack, not O(N).
3683        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    // =================================================================
3734    // `prepare_raw`/`push_prepared_raw` and `prepare_delta`/
3735    // `push_prepared_delta` — the pure-compression / sequential-append
3736    // split added so a caller (mkit-cli's `build_and_upload_packs`) can
3737    // run compression across a thread pool before replaying the
3738    // append in order. Must produce byte-identical output to the
3739    // original `push_raw`/`push_delta` and preserve call order when
3740    // interleaved with them — a future refactor of `append_raw_frame`/
3741    // `append_delta_frame` (the shared tail both paths funnel through)
3742    // must not let the two paths drift apart.
3743    // =================================================================
3744
3745    #[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        // A caller mixing a compression-fan-out batch's prepared
3793        // entries with directly-pushed ones (e.g. across a
3794        // pack-splitting boundary) must see them land in the pack in
3795        // exactly the order they were pushed, whichever path each one
3796        // took.
3797        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        // issue #831: push-side pack-splitting decisions read
3822        // `total_payload()` to decide when to seal a pack, so it must
3823        // track the writer's real running wire-payload sum — checked
3824        // here against both a plain raw entry and a delta entry, and
3825        // bounded by the uncompressed input sizes (compression only
3826        // ever makes the wire payload smaller, never larger).
3827        let mut w = PackWriter::new();
3828        assert_eq!(w.total_payload(), 0);
3829
3830        // Incompressible (random-ish) bytes so `maybe_compress` doesn't
3831        // shrink them — keeps the assertion exact rather than "at most".
3832        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        // The writer never emits more wire bytes than the caller handed
3846        // it (delta payload = base_hash + stream, uncompressed worst
3847        // case), and `total_payload` must equal the sum of what was
3848        // actually appended so far.
3849        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        // Two near-identical blobs. Delta should reconstruct the second.
3857        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        // SPEC-PACKFILE §4: a delta's base MUST appear earlier in the
3888        // pack as a raw entry (or already exist in the destination
3889        // store) — never later. Same fixture as
3890        // `raw_then_delta_resolves_in_pack`, but with the delta and its
3891        // raw base swapped so the base comes *after* the delta that
3892        // references it. The base is genuinely absent from both the
3893        // pack-so-far and the (empty) store at the point the delta is
3894        // read, so this must fail exactly like a base that's missing
3895        // outright — never silently succeed by resolving against a
3896        // same-pack entry the reader hasn't reached yet.
3897        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        // Nothing from this rejected pack should be visible — not even
3920        // the raw base entry that appeared after the bad delta.
3921        assert!(!store.contains(&base_hash));
3922    }
3923
3924    #[test]
3925    fn delta_before_its_base_is_rejected_under_parallel_raw_fanout() {
3926        // Same defect as `delta_before_its_base_in_pack_order_is_rejected`,
3927        // but with enough raw entries ahead of the base (200 — comfortably
3928        // over `stage_raw_entries`'s `ENTRIES_PER_THREAD * threads` on any
3929        // CI runner up to dozens of cores) to force phase 2's parallel
3930        // `std::thread::scope` fan-out rather than its small-pack
3931        // sequential fallback. The base-before-delta rule must hold
3932        // regardless of how phase 2 schedules raw-entry work across
3933        // threads.
3934        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        // Phase 2 validates every raw entry in the pack up front
3968        // (independent per entry, so it can fan out across threads),
3969        // but phase 3 decides *in pack order* which single problem is
3970        // actually reported. A delta at position 0 with a missing base
3971        // must win over a malformed raw entry at a later position —
3972        // the pre-fan-out single-loop reader would have rejected
3973        // position 0 and never even looked at the later one. 200
3974        // well-formed raw fillers keep phase 2 on its parallel branch
3975        // (comfortably over `stage_raw_entries`'s threshold on any CI
3976        // runner), with one malformed entry mixed in among them.
3977        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                // Not a canonical storable object — fails phase 2's
3987                // `validate_storable_object` with `PackError::InvalidObject`.
3988                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        // One raw blob + two deltas against two different bases. The scan
4014        // must return exactly the two (deduped) base hashes, ignoring raw.
4015        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(); // a raw entry — must be ignored
4026        w.push_delta(&ha, &stream_a).unwrap();
4027        w.push_delta(&hb, &stream_b).unwrap();
4028        w.push_delta(&ha, &stream_a).unwrap(); // duplicate base — deduped
4029        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        // Use an arbitrary invalid 4-byte sequence; the rename gate
4202        // forbids spelling out the upstream pre-rename magic literally.
4203        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        // version is u32 LE at offset 4
4217        pack[4] = 99;
4218        // Corrupt trailer so the version check fires first — but
4219        // SPEC-PACKFILE §8 says trailer is checked before entries,
4220        // and we want UnsupportedVersion. Trailer check happens after
4221        // version check in our impl (see read()), so just leave the
4222        // trailer; it will fail UnsupportedVersion on byte 4.
4223        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']; // only 2 bytes
4231        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; // flip one bit
4245        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        // Hand-build a pack with one entry of type 0x01.
4253        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); // RESERVED type
4258        buf.extend_from_slice(&0u32.to_le_bytes()); // payload_len = 0
4259        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); // unknown
4274        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        // Build a minimal SPEC-DELTA stream that targets a nonexistent base.
4288        let mut stream = Vec::new();
4289        stream.push(0x01); // version
4290        stream.extend_from_slice(&0u32.to_le_bytes()); // base_len
4291        stream.extend_from_slice(&0u32.to_le_bytes()); // result_len
4292        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        // No payload bytes follow.
4310        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        // Add a fake trailer so trailer-check passes — wait, it can't
4325        // pass since the body is bogus. Compute it correctly so the
4326        // trailer is the not-the-failure path; then the count cap must
4327        // fire first per read() ordering.
4328        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        // count cap fires after trailer verify in our impl. Either is
4334        // acceptable; assert one of them.
4335        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        // `PackfileTooLarge` on the reader's running-payload-total is
4344        // enforced against MAX_TOTAL_PAYLOAD (4 GiB) in production —
4345        // impractical to trip directly in a unit test without
4346        // allocating gigabytes. `read_with_payload_cap` is the
4347        // test-only injection point: same check, caller-supplied cap.
4348        // Incompressible filler (not `[0xAA; 64]`/`[0xBB; 64]`, which
4349        // the §3.3 writer policy would shrink to a handful of
4350        // compressed bytes): this test's cap arithmetic below assumes
4351        // `blob_a.len()`/`blob_b.len()` ARE the on-wire sizes.
4352        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        // Cap smaller than the combined payload but big enough that
4362        // the first entry alone fits — the SECOND entry's running
4363        // total must trip the cap, not an entry-count or bounds check.
4364        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        // Sanity: the same pack with a generous cap (the real
4372        // MAX_TOTAL_PAYLOAD) unpacks normally.
4373        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        // Issue #647: `PackReader::read` used to copy EVERY raw entry's
4390        // payload into a fresh `Arc<[u8]>` retained in `in_pack` for the
4391        // whole call — redundant given `pack_bytes` (the pack's own
4392        // bytes) is already resident in the caller's memory the whole
4393        // time. A streaming reader only needs a BORROW into
4394        // `pack_bytes` for raw entries; nothing about a raw entry's
4395        // bytes should ever be copied a second time. `owned_bytes`
4396        // tracks the exact production code path that would otherwise
4397        // do that copy (see `read_tracking_owned_bytes`), so this is a
4398        // precise, allocator-free proof rather than a fuzzy proxy.
4399        // Incompressible filler: a repeated-byte 16 KiB payload would
4400        // trip the §3.3 writer policy into emitting `0x03` zstd-raw
4401        // instead of `0x00` raw, which genuinely DOES need an owned
4402        // decompression buffer — that's not what this test is about.
4403        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        // Complements the all-raw test above: the raw base must still
4436        // be a zero-copy borrow, and each delta's "owned" cost must be
4437        // exactly its reconstructed target size — never the base's size
4438        // too, and never the whole pack's. Incompressible base content,
4439        // same rationale as that test: a compressible base would
4440        // legitimately need an owned decompression buffer, which is
4441        // not what this test measures.
4442        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        // Issue #647: `PackWriter::finish()` used to hold every pushed
4476        // entry in a separate `entries` list and then copy ALL of them
4477        // a second time into a freshly `Vec::with_capacity`'d output
4478        // buffer. A streaming writer appends each entry's frame
4479        // directly into the one output buffer as it's pushed, so
4480        // `finish()` itself should only ever append the 32-byte
4481        // trailer — `bytes_copied` tracks exactly that production code
4482        // path (see `finish_tracking_bytes_copied`).
4483        // Incompressible filler (see the comment in
4484        // `unpack_does_not_recopy_raw_payloads_into_a_second_buffer`):
4485        // a repeated-byte payload would shrink dramatically under the
4486        // §3.3 writer policy, invalidating the `pack.len() > 512 KiB`
4487        // sanity check below, which has nothing to do with what this
4488        // test is proving.
4489        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        // Plant the base in the store first.
4510        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        // Pack contains ONLY a delta; the base must be resolved from disk.
4518        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        // Regression for #643: N deltas in one pack all referencing the
4537        // SAME out-of-pack (already-in-store) base object must resolve
4538        // that base with exactly one physical store read, not N — the
4539        // first store-resolved base should be cached into `in_pack` for
4540        // subsequent deltas to hit in memory.
4541        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        // Five distinct deltas against the one shared external base.
4553        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        // Correctness: caching the store-resolved base must not change
4577        // the decoded result for any of the N deltas — every target
4578        // still comes out byte-identical to the uncached decode.
4579        for (target_hash, target_obj) in expected_targets {
4580            assert_eq!(store.read(&target_hash).unwrap(), target_obj);
4581        }
4582    }
4583
4584    // =====================================================================
4585    // SPEC-PACKFILE v2: zstd-compressed entries (issue #646)
4586    // =====================================================================
4587
4588    /// Highly-compressible synthetic payload: `MIN_COMPRESS_LEN` (64) is
4589    /// the writer's floor, so use something well past it — a single
4590    /// repeated byte is the easiest thing for zstd to shrink hard.
4591    #[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        // Base and target share (almost) nothing, so `delta::encode`
4633        // emits a stream dominated by one big INSERT of the target's
4634        // highly-compressible content — long and repetitive enough for
4635        // the §3.3 writer policy to compress it into a 0x04 entry.
4636        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        // Walk past the first (raw base) entry's frame to find the
4660        // second entry's type byte.
4661        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        // Hand-build a version-1-declared pack with one 0x03 entry —
4684        // even though 0x03's byte layout is otherwise well-formed, it
4685        // MUST be rejected because the header says version 1
4686        // (SPEC-PACKFILE §3).
4687        let mut buf = Vec::new();
4688        buf.extend_from_slice(MAGIC);
4689        buf.extend_from_slice(&VERSION.to_le_bytes()); // version = 1
4690        buf.extend_from_slice(&1u32.to_le_bytes()); // entry_count = 1
4691        buf.push(0x03);
4692        let inner_payload = 0u32.to_le_bytes(); // uncompressed_len = 0, no frame bytes
4693        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        // Build a real compressed pack, then tamper the payload so the
4709        // claimed `uncompressed_len` no longer matches what the frame
4710        // actually decompresses to.
4711        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        // Lie about the length: claim one byte more than the frame
4726        // actually decompresses to. The trailer no longer matches the
4727        // tampered body, so recompute it (this test targets the
4728        // length-mismatch check specifically, not trailer verification,
4729        // which is already covered by `rejects_bit_flipped_trailer`).
4730        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        // A `0x03` entry claiming a decompressed size over
4749        // MAX_RAW_OBJECT_SIZE must be rejected before any decompression
4750        // is attempted — hand-build the entry rather than actually
4751        // producing a >1 GiB frame.
4752        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        // Payload: [4B claimed uncompressed_len][tiny bogus "frame"].
4759        // The over-cap check must fire before the (bogus) frame is
4760        // ever touched, so its content doesn't need to be valid zstd.
4761        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    // =====================================================================
4836    // `DeltaBaseSource` seam / `decode_entries_with` (WP-4.2)
4837    // =====================================================================
4838
4839    /// `(target id, target bytes, delta stream against the base)`.
4840    type Variant = (Hash, Vec<u8>, Vec<u8>);
4841    /// `(id, bytes, from_delta)` as a decode sink saw it.
4842    type Seen = (Hash, Vec<u8>, bool);
4843
4844    /// A 512-byte base blob and `n` single-byte variants of it, with
4845    /// their delta streams against the base.
4846    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    /// Decode `pack` through `bases`, collecting what the sink saw.
4866    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    /// Self-contained test packs of every entry shape the writer emits
4957    /// (raw, raw + in-pack delta, and — under `pack-zstd` — `0x03`/`0x04`).
4958    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        // Chain: a delta whose base is itself an earlier delta target.
4975        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        // The reader has no decode budget, so compare against an unbounded
5039        // decoder (the over-cap delta would otherwise trip the budget first).
5040        let unbounded = DecodeLimits::default().with_max_decoded_bytes(u64::MAX);
5041        // Every error-precedence pack pinned above must fail the same way
5042        // through the store-less decoder.
5043        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        // The base exists — in a store the decoder is not given. A
5094        // membership-scoped source that does not list it, and a source
5095        // with no external bases at all, must both fail exactly as a base
5096        // that exists nowhere does.
5097        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        // Once the base IS a member, the same pack decodes.
5146        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        // A different, perfectly valid object; bytes that are not an
5171        // object; and a pack-only delta object: none may stand in.
5172        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        // The genuine bytes from the same untrusted source are accepted.
5192        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        // `multiple_deltas_against_shared_external_base_read_store_once`,
5199        // through the store-less decoder with `&ObjectStore` as source.
5200        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        // Store-less: nothing was written.
5220        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        // The verified store path keeps its pre-seam behavior: a base whose
5228        // on-disk bytes no longer hash to its id is `Store(HashMismatch)`,
5229        // not a silent "missing".
5230        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    /// The WP-4.2 review's memory probe: one 64 KiB raw base, then `n`
5254    /// deltas, each a valid SPEC-DELTA stream that rebuilds a distinct
5255    /// `target_len`-byte blob purely by copying base bytes (a few bytes of
5256    /// wire per 64 KiB of result; `pack-zstd` shrinks it further).
5257    /// Returns the pack and the ids of the delta targets.
5258    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        // Blob prologue + length for a `target_len`-byte blob.
5263        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            // Four distinguishing data bytes, then zeros copied from base.
5278            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        // 16 deltas declaring 128 MiB each (2 GiB in all) from a pack of a
5299        // few hundred KiB at most: over the default 1 GiB budget, so the
5300        // decode fails before it applies (allocates) a single target.
5301        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        // The same construction at a small size decodes to real objects
5316        // under a budget that covers it, and fails under one just short of
5317        // the declared results alone. (Under `pack-zstd` the writer also
5318        // compresses the streams, whose claims are charged too, so the
5319        // exact boundary depends on the build.)
5320        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        // A v2 `0x03` entry claiming 512 MiB behind a bogus frame: a
5344        // decoder that decompressed first would fail on the frame; the
5345        // budget must reject on the claim alone.
5346        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        // Within budget, the frame itself is what fails.
5362        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        // `pos + payload_len` used to overflow a 32-bit `usize` (wasm32)
5372        // here; the checks are now subtraction-based. 64-bit hosts get the
5373        // same clean `UnexpectedEof`; `mkit-core-wasm-check` runs the
5374        // wasm32 lane.
5375        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        // Behavioral pin for the retention rule: an entry no delta names
5402        // is dropped after the sink, yet a later delta against an earlier
5403        // *named* entry (raw or delta target) still resolves.
5404        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    /// A repository-membership source over in-memory objects, counting
5424    /// fetches (the review's `bomb ext` shape: bases live only here).
5425    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    /// Two 600 KiB member blobs and a tiny delta against each: `(source,
5438    /// [(base id, delta stream)])`. Each delta rebuilds a small blob, so
5439    /// its declared result is a few hundred bytes.
5440    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        // A then B then A: A must stay cached while B is fetched, so the
5473        // two 600 KiB bases are held at once and pass the 1 MiB budget at
5474        // B's fetch. Only the first delta reaches the sink.
5475        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        // A single base larger than the budget fails before any sink call.
5486        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        // A, A, then B: A's last use comes before B is fetched, so A is
5501        // dropped and its charge credited back; the peak is one base, and
5502        // the same 1 MiB budget accepts the pack. A is fetched once.
5503        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/// Kani proof harnesses (`cargo kani -p mkit-core --no-default-features
5551/// -Z stubbing --harness pack_`, see `delta.rs`; the reader-side zstd
5552/// call is stubbed either way and the writer never compresses below
5553/// `MIN_COMPRESS_LEN`, so the feature does not change what is checked),
5554/// the model-checked counterparts of the `pack` and `pack_entries` fuzz
5555/// targets.
5556///
5557/// BLAKE3 is stubbed with a cheap deterministic `toy_hash`: the trailer
5558/// bytes stay fully symbolic, so both the §8 trailer-match and -mismatch
5559/// paths are explored, and the writer/reader round-trip still agrees on
5560/// one function. zstd (C FFI, not modelable) is stubbed with a
5561/// nondeterministic "fail, or return any <= 2-byte buffer" — a sound
5562/// over-approximation for panic-freedom of the surrounding framing code.
5563/// `PackReader::read` itself needs an on-disk `ObjectStore`, which Kani
5564/// cannot model. A store-free harness over its per-entry steps (entry
5565/// parsing + the storability gate) ran out of memory even for one 5-byte
5566/// entry frame: CBMC's symbolic execution of the `PackError` →
5567/// `StoreError` → `std::io::Error` drop glue dominates (~13 min), so
5568/// that composition is left to the `pack` fuzz target.
5569///
5570/// The `pack_window_cursor_*` harnesses cover the one new decoder surface
5571/// of the windowed reader (SPEC-PACKFILE §11): the persisted resumable
5572/// cursor (`window::WindowCursor::from_bytes`). Its framing path shares
5573/// `decode_payload` with `PackEntries`; `rewrite::rewrite_excluding`
5574/// parses through `PackEntries`/`decode_entries_with` and adds no decoder.
5575#[cfg(kani)]
5576mod kani_proofs {
5577    use super::*;
5578
5579    /// Loop-free deterministic stand-in for BLAKE3: the first and last
5580    /// 16 bytes of `data` plus its length (loop-free so it does not
5581    /// interact with the global unwind bound).
5582    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    /// Symbolic pack with an entry area of exactly `BODY` bytes (header,
5604    /// body and trailer all symbolic). Harnesses call this once per
5605    /// concrete `BODY` so CBMC can constant-fold the pack length.
5606    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    /// `pack_entries` target, 5-byte entry area (exactly one entry frame
5664    /// — type + length — with an empty payload, or a truncated/oversized
5665    /// one), header, frame and trailer all symbolic (so any magic,
5666    /// version, `entry_count` and trailer): `PackEntries::new` and full
5667    /// iteration never panic/overflow/read OOB. On `Ok`: magic/version
5668    /// valid (§1), trailer equals the hash of the preceding bytes (§8), a
5669    /// v1 pack yields exactly `entry_count` `Ok` items ending at the
5670    /// trailer with no gap (§3, §6), every payload range lies inside the
5671    /// entry area (§2), and `is_raw_only` ⇒ no delta item.
5672    ///
5673    /// Run with `-Z unstable-options --cbmc-args --unwindset memcmp.0:33`
5674    /// (the 32-byte trailer comparison); every other loop is bounded by
5675    /// the global unwind of 4, which keeps CBMC from unrolling the
5676    /// symbolic-`entry_count` loop 33 times. Unwinding assertions stay
5677    /// on, so a too-small bound fails loudly. One entry-area length per
5678    /// harness: 0..=6 (and 0..=12) in one harness did not finish within
5679    /// 15 min, and the empty (0-byte) area alone ran out of memory.
5680    #[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    /// As above for a 10-byte entry area: two empty-payload entries (or
5689    /// one with a 5-byte payload).
5690    #[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    /// Writer → reader round-trip (SPEC-PACKFILE §1–§3): a pack built by
5736    /// `PackWriter` from one raw entry of 1 symbolic byte is accepted by
5737    /// `PackEntries::new` as v1 (no compression below `MIN_COMPRESS_LEN`)
5738    /// and yields exactly the pushed entry. CBMC's symbolic execution of
5739    /// the `PackError` → `StoreError` → `std::io::Error` drop glue costs
5740    /// ~4 min per writer/reader pass here, so the bound is one entry
5741    /// shape (0..=2 bytes in one harness, and a raw + delta pack, ran out
5742    /// of memory; the delta writer path is covered by the unit tests).
5743    #[kani::proof]
5744    #[kani::stub(crate::hash::hash, toy_hash)]
5745    // Run with `-Z unstable-options --cbmc-args --unwindset memcmp.0:33`
5746    // (32-byte trailer / base-hash comparisons); <= 2 entries otherwise.
5747    #[kani::unwind(4)]
5748    fn pack_writer_roundtrip_raw() {
5749        writer_rt::<1, 0>(false);
5750    }
5751
5752    /// Canary: with the trailer check in place, a single flipped body
5753    /// byte must be detectable — the checker has to falsify "every
5754    /// single-byte mutation of a valid pack still parses".
5755    #[kani::proof]
5756    #[kani::stub(crate::hash::hash, toy_hash)]
5757    // As above: `--cbmc-args --unwindset memcmp.0:33`.
5758    #[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    /// A checksum-valid symbolic cursor encoding (canonical v1 field
5771    /// order: see the `window::cursor` module doc) of `N` body bytes, all
5772    /// symbolic except the option tags (0/1) at the concrete `(offset,
5773    /// present)` positions in `tags` and the two tree depth bytes at
5774    /// `depths`, which are 0. A symbolic tag or depth byte (even one
5775    /// restricted to "valid or rejected") makes CBMC explore every
5776    /// layout after it (the reader's `Result` is merged at each return,
5777    /// so later field offsets become symbolic); that did not finish within
5778    /// 10 min. The (toy) checksum is appended
5779    /// so the checksum gate passes and every field parser and the
5780    /// geometry validation run on attacker-chosen values.
5781    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    /// Layout A (125-byte body): first-window boundary bound to the
5796    /// trailer anchor, no pack id requested, no first-non-raw index;
5797    /// current-window prefix commitment present; both trees empty.
5798    /// Offsets: version 0, four u64 + three u32 at 1..45, tags at 45
5799    /// (first-non-raw), 54 (expected), 55 (anchor) + digest, 88 (prefix)
5800    /// + digest, trees at 121/122 and 123/124.
5801    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    /// Resumable-cursor decoder (`window::WindowCursor::from_bytes`, the
5816    /// persisted state of the SPEC-PACKFILE §11 windowed reader) on a
5817    /// checksum-valid layout-A encoding whose every numeric field and
5818    /// digest is symbolic: decoding (field parsing, trailing-byte check,
5819    /// geometry validation) never panics, overflows
5820    /// or reads out of bounds; `cover` shows acceptance is reachable.
5821    /// (Also asserting that an accepted cursor re-encodes byte-exactly
5822    /// via `to_bytes`, and the pack-id-bound layout, each did not finish
5823    /// within 15 min; the round-trip stays with the `window` unit tests.)
5824    #[kani::proof]
5825    #[kani::stub(crate::hash::hash, toy_hash)]
5826    // Run with `-Z unstable-options --cbmc-args --unwindset memcmp.0:33`
5827    // (checksum comparison); the tree loops read no CVs.
5828    #[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        // Skip the harness-side `PackError` drop glue (`Store` →
5833        // `std::io::Error`); nothing about dropping is checked here.
5834        core::mem::forget(got);
5835        kani::cover!(ok, "accepted_anchor_cursor");
5836    }
5837
5838    /// A concrete, valid layout-A cursor (100-byte pack, 64 KiB windows,
5839    /// boundary right after the header, no entries; symbolic anchor and
5840    /// prefix digests) with `mask` XORed into the low byte of `pack_len`
5841    /// (offset 1, inside the checksummed body) after the checksum is
5842    /// computed. Every flipped `pack_len` in 44..=255 is still valid
5843    /// geometry for this cursor, so for those only the checksum rejects.
5844    fn flipped_cursor(mask: u8) -> Vec<u8> {
5845        let mut body = [0u8; 125];
5846        body[0] = 1; // cursor encoding version
5847        body[1..9].copy_from_slice(&100u64.to_le_bytes()); // pack_len
5848        body[9..17].copy_from_slice(&(64u64 << 10).to_le_bytes()); // window
5849        body[17..25].copy_from_slice(&12u64.to_le_bytes()); // pos
5850        body[33..37].copy_from_slice(&1u32.to_le_bytes()); // pack version
5851        body[55] = 1; // anchor present
5852        body[56..88].copy_from_slice(&kani::any::<Hash>());
5853        body[88] = 1; // window prefix present
5854        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    /// Checksum gate: the concrete cursor above decodes iff it is
5862    /// unmodified, i.e. every single-byte change to its `pack_len` field
5863    /// (most of which leave a still-valid cursor) is rejected. One
5864    /// `from_bytes` call per harness: each costs ~7 min of CBMC symbolic
5865    /// execution (~6.8M steps, much of it `PackError` → `StoreError` →
5866    /// `std::io::Error` machinery), so two calls did not finish within
5867    /// 15 min.
5868    #[kani::proof]
5869    #[kani::stub(crate::hash::hash, toy_hash)]
5870    // As above: `--cbmc-args --unwindset memcmp.0:33`.
5871    #[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    /// Canary: the checker must falsify "a flipped `pack_len` byte still
5881    /// decodes" (the negation of the property above).
5882    #[kani::proof]
5883    #[kani::stub(crate::hash::hash, toy_hash)]
5884    // As above: `--cbmc-args --unwindset memcmp.0:33`.
5885    #[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}