Skip to main content

heddle_pack/store/pack/
streaming_builder.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Streaming pack builder for bounded-memory imports.
3//!
4//! `PackBuilder` accumulates every `(id, type, data)` tuple in memory
5//! before producing a pack. That's fine for sync-protocol packets and
6//! small batches, but the import path can produce millions of objects
7//! and would OOM on large repos.
8//!
9//! `StreamingPackBuilder` removes the in-memory buffering by:
10//!
11//! 1. **Streaming pack data to disk** as objects are added. Compression
12//!    runs per-object (the existing zstd path is non-streaming, so the
13//!    one compressed payload is held briefly in a `Vec<u8>` before
14//!    being written), but the writer never holds more than one
15//!    object's worth of data plus its `BufWriter` capacity.
16//!
17//! 2. **External sorting the index** via 512 hash-prefix bucket files
18//!    on disk (256 for `Hash` ids, 256 for `StateId` ids). Each
19//!    `add()` appends one fixed-shape `(id, offset)` record to the
20//!    bucket whose first byte matches the id's first inner byte. At
21//!    finalize, each bucket is externally sorted in 1 MiB runs; the
22//!    concatenation of `Hash` buckets followed by `StateId` buckets
23//!    in byte order produces the exact same global sort `PackBuilder`
24//!    would have via `entries.sort_by_key(|e| e.id)`.
25//!
26//! ## Memory bound
27//!
28//! - Pack data on disk: streamed; only one compressed object held in
29//!   memory at a time.
30//! - Index entries in bucket buffers: at most 32 bucket files are held
31//!   open at once, each behind a default-capacity `BufWriter` (~8 KB),
32//!   so peak buffering is ~256 KB.
33//! - Sort scratch at finalize: 1 MiB of records plus two buffered readers.
34//!   At most 64 merge levels can exist for addressable input files.
35//!
36//! Memory depends on the largest individual object, never the object count.
37//!
38//! ## Trade-offs vs `PackBuilder`
39//!
40//! - **No delta encoding.** Streaming and sliding-window deltas are
41//!   incompatible — delta search needs random access to recently-
42//!   written objects. The import path runs with deltas disabled
43//!   anyway (the cost-benefit is bad on real Heddle history), so this
44//!   is a non-issue for the call site that motivated this builder.
45//! - **No path-grouped reordering.** Entries land in the order added.
46//! - **Output is a pack file at a path** rather than `(Vec<u8>, Vec<u8>)`.
47//!   Callers pair this with `objects::store::ObjectStore::install_pack_from_path`
48//!   which moves/installs the pack without copying it through RAM.
49//! - **Re-reads the pack at finalize** to compute the BLAKE3 trailer
50//!   checksum (the pack format hashes header+body, and the count goes
51//!   in the header — we patch it on finalize via seek-back, then
52//!   re-stream the body to the hasher). 2× sequential disk I/O on the
53//!   pack data is the cost of sticking with the current format. A
54//!   future format change could put the count in the footer to avoid
55//!   the second pass.
56
57use std::{
58    collections::HashSet,
59    fs::{File, OpenOptions},
60    io::{self, BufWriter, Cursor, Read, Seek, SeekFrom, Write},
61    path::PathBuf,
62};
63
64use heddle_format::compression::CompressionConfig;
65
66use super::{ObjectType, PackObjectId, PackStats, pack_container_spec, write_container_header};
67
68/// How many bytes to reserve for the compressed-size varint in the
69/// streaming path. 10 is enough to encode any `u64` (max 9 7-bit
70/// continuation bytes plus 1 terminator). After streaming we patch
71/// the placeholder with a non-canonical varint that pads to exactly
72/// this length. Only the zstd-enabled compress path uses it.
73#[cfg(feature = "zstd")]
74const CSIZE_PLACEHOLDER_LEN: usize = 10;
75use crate::{
76    object::ContentHash,
77    store::{Result, StoreError},
78};
79
80/// Number of buckets per id variant. 256 = one bucket per first byte
81/// of the inner id. We want the bucket boundaries to align with the
82/// `PackObjectId`'s `Ord` derivation (variant tag major, inner bytes
83/// minor) so the concatenated bucket output matches what
84/// `PackIndex::sort()` would have produced.
85const BUCKETS_PER_VARIANT: usize = 256;
86/// 256 each for Hash, StateId, and AnnotatedTag ids.
87const TOTAL_BUCKETS: usize = BUCKETS_PER_VARIANT * 3;
88/// Cap concurrently-open index-bucket files. macOS GUI-launched
89/// processes commonly inherit a 256-fd soft limit; imports also need
90/// room for Git pack/index files, sqlite maps, the output pack, etc.
91const MAX_OPEN_BUCKET_WRITERS: usize = 32;
92
93/// Variant indices into the `bucket_*` arrays. `Hash` ids fill the
94/// lower half (matches the variant order in `PackObjectId` which makes
95/// `Hash(_) < StateId(_)`).
96const HASH_VARIANT: usize = 0;
97const CHANGEID_VARIANT: usize = 1;
98const ANNOTATED_TAG_VARIANT: usize = 2;
99
100/// Fsync staged pack bytes after finalize flush (Wave 5 L7).
101///
102/// Production writers are [`File`]; in-memory [`Cursor`] tests no-op.
103/// Publish still re-fsyncs at `publish_file_durable` install; this closes
104/// the pre-publish window if a caller inspects staged files after finalize.
105pub trait SyncData {
106    fn sync_data_for_durability(&mut self) -> io::Result<()>;
107}
108
109impl SyncData for File {
110    fn sync_data_for_durability(&mut self) -> io::Result<()> {
111        self.sync_all()
112    }
113}
114
115impl SyncData for Cursor<Vec<u8>> {
116    fn sync_data_for_durability(&mut self) -> io::Result<()> {
117        Ok(())
118    }
119}
120
121/// Streaming pack builder. Held generic over the pack writer (`File`
122/// in production, `Cursor<Vec<u8>>` in tests).
123pub struct StreamingPackBuilder<W: Write + Read + Seek> {
124    /// Writer for the pack's `[header][body]` content. The trailer
125    /// checksum is appended to the same writer at `finalize`.
126    /// Wrapped in `Option` so `finalize` can `.take()` it out without
127    /// running afoul of the `Drop` impl's restriction on moving fields.
128    /// `None` after `finalize` succeeds.
129    pack_writer: Option<BufWriter<W>>,
130    /// Position in the pack writer where the header was written, so
131    /// we can seek back at finalize and patch the real `object_count`
132    /// into bytes 8..16.
133    header_offset: u64,
134    /// Logical append position. Avoids flushing the buffered writer before
135    /// every object just to ask the file for its current offset.
136    pack_position: u64,
137    record_count: u64,
138    object_count: u64,
139    declared_object_count: Option<u64>,
140    total_uncompressed: u64,
141    total_compressed: u64,
142    /// Compression knobs. Only consulted when the `zstd` feature is on
143    /// (`enabled` and `min_size` decide whether each entry compresses;
144    /// `level` parameterizes the encoder). Without `zstd` every entry
145    /// takes the raw branch and this field is just along for the ride.
146    #[cfg_attr(not(feature = "zstd"), allow(dead_code))]
147    compression: CompressionConfig,
148    /// Directory holding the 512 bucket files. Owned by the builder
149    /// so we can clean up on `Drop` if `finalize` is never called.
150    bucket_dir: PathBuf,
151    scratch_lease: Option<super::ScratchLease>,
152    /// Buckets `[variant][prefix_byte]` → optional buffered file.
153    /// Lazily opened on first write and capped with LRU eviction so a
154    /// large import cannot exhaust the process fd limit.
155    bucket_writers: Vec<Option<BucketWriter>>,
156    open_bucket_writers: usize,
157    bucket_access_tick: u64,
158    bucket_paths: Vec<PathBuf>,
159    /// File path where the pack index is materialized at `finalize`.
160    /// Bytes are written incrementally as buckets are sorted, so the
161    /// index never sits in memory in its entirety.
162    index_path: PathBuf,
163    /// Set true on `finalize` so `Drop` knows the bucket dir was
164    /// already cleaned and shouldn't be removed again.
165    finalized: bool,
166    durable: bool,
167}
168
169struct BucketWriter {
170    writer: BufWriter<File>,
171    last_used: u64,
172}
173
174#[cfg(feature = "zstd")]
175struct CountingWriter<'a, W: Write> {
176    inner: &'a mut W,
177    written: u64,
178}
179
180#[cfg(feature = "zstd")]
181impl<'a, W: Write> CountingWriter<'a, W> {
182    fn new(inner: &'a mut W) -> Self {
183        Self { inner, written: 0 }
184    }
185}
186
187#[cfg(feature = "zstd")]
188impl<W: Write> Write for CountingWriter<'_, W> {
189    fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
190        let written = self.inner.write(buf)?;
191        self.written = self.written.saturating_add(written as u64);
192        Ok(written)
193    }
194
195    fn flush(&mut self) -> std::io::Result<()> {
196        self.inner.flush()
197    }
198}
199
200impl<W: Write + Read + Seek + SyncData> StreamingPackBuilder<W> {
201    /// Open a streaming builder against `pack_writer`, using
202    /// `bucket_dir` for transient index buckets and writing the
203    /// finalized index to `index_path`. The bucket dir is created if
204    /// it doesn't exist; on a successful `finalize` it's removed
205    /// (along with any bucket files left in it).
206    ///
207    /// `index_path` is *not* created by `new` — opening happens at
208    /// finalize so a misconfigured caller doesn't leave an empty index
209    /// file behind on early failure. It's still recorded here so
210    /// `finalize` can write to a known location and the caller can
211    /// install the file by path.
212    ///
213    /// The `pack_writer` must support `Read` because finalize re-streams
214    /// the body to compute the trailer checksum — see the module-level
215    /// note on the format.
216    #[cfg(feature = "fs")]
217    pub fn new(
218        pack_writer: W,
219        index_path: PathBuf,
220        compression: CompressionConfig,
221        bucket_dir: PathBuf,
222    ) -> Result<Self> {
223        Self::new_inner(pack_writer, index_path, compression, bucket_dir, None, true)
224    }
225
226    /// Open a streaming builder whose final object count is known up front.
227    ///
228    /// This writes the real count into the pack header immediately, so callers
229    /// may safely stream already-flushed pack bytes before `finalize()` appends
230    /// the trailer checksum. `finalize()` still verifies that exactly this many
231    /// objects were added before producing the index.
232    #[cfg(feature = "fs")]
233    pub fn new_with_object_count(
234        pack_writer: W,
235        index_path: PathBuf,
236        compression: CompressionConfig,
237        bucket_dir: PathBuf,
238        object_count: u64,
239    ) -> Result<Self> {
240        Self::new_inner(
241            pack_writer,
242            index_path,
243            compression,
244            bucket_dir,
245            Some(object_count),
246            true,
247        )
248    }
249
250    /// Open a known-count builder for a transient transfer spool.
251    ///
252    /// The source object store remains authoritative, so these files only need
253    /// to be flushed for same-process readers; crash durability would add an
254    /// fsync to the pack, index, and directory for every unary push without
255    /// making the transfer any safer.
256    pub fn new_with_object_count_ephemeral(
257        pack_writer: W,
258        index_path: PathBuf,
259        compression: CompressionConfig,
260        bucket_dir: PathBuf,
261        object_count: u64,
262    ) -> Result<Self> {
263        Self::new_inner(
264            pack_writer,
265            index_path,
266            compression,
267            bucket_dir,
268            Some(object_count),
269            false,
270        )
271    }
272
273    fn new_inner(
274        mut pack_writer: W,
275        index_path: PathBuf,
276        compression: CompressionConfig,
277        bucket_dir: PathBuf,
278        declared_object_count: Option<u64>,
279        durable: bool,
280    ) -> Result<Self> {
281        #[cfg(feature = "fs")]
282        if durable {
283            heddle_fs_prims::fs_atomic::create_dir_all_durable(&bucket_dir)
284                .map_err(StoreError::from)?;
285        } else {
286            std::fs::create_dir_all(&bucket_dir).map_err(StoreError::from)?;
287        }
288        #[cfg(not(feature = "fs"))]
289        {
290            debug_assert!(!durable);
291            std::fs::create_dir_all(&bucket_dir).map_err(StoreError::from)?;
292        }
293        let scratch_lease = super::ScratchLease::acquire(&bucket_dir)?;
294        let header_offset = pack_writer.stream_position().map_err(StoreError::from)?;
295
296        // Write a placeholder header with `count = 0` unless the caller knows
297        // the count up front. The known-count path lets network senders tail
298        // flushed pack bytes before finalize without later mutating bytes they
299        // already sent.
300        let mut header_bytes = Vec::with_capacity(16);
301        write_container_header(
302            &mut header_bytes,
303            pack_container_spec(),
304            declared_object_count.unwrap_or(0),
305        );
306        pack_writer
307            .write_all(&header_bytes)
308            .map_err(StoreError::from)?;
309
310        let bucket_paths: Vec<PathBuf> = (0..TOTAL_BUCKETS)
311            .map(|i| {
312                let variant = match i / BUCKETS_PER_VARIANT {
313                    HASH_VARIANT => 'h',
314                    CHANGEID_VARIANT => 's',
315                    ANNOTATED_TAG_VARIANT => 't',
316                    _ => unreachable!("bucket variant is bounded by TOTAL_BUCKETS"),
317                };
318                let prefix = i % BUCKETS_PER_VARIANT;
319                bucket_dir.join(format!("bucket-{variant}-{prefix:02x}"))
320            })
321            .collect();
322        for path in &bucket_paths {
323            let _ = std::fs::remove_file(path);
324        }
325
326        Ok(Self {
327            pack_writer: Some(BufWriter::new(pack_writer)),
328            header_offset,
329            pack_position: header_offset + header_bytes.len() as u64,
330            record_count: 0,
331            object_count: 0,
332            declared_object_count,
333            total_uncompressed: 0,
334            total_compressed: 0,
335            compression,
336            bucket_dir,
337            scratch_lease: Some(scratch_lease),
338            bucket_writers: (0..TOTAL_BUCKETS).map(|_| None).collect(),
339            open_bucket_writers: 0,
340            bucket_access_tick: 0,
341            bucket_paths,
342            index_path,
343            finalized: false,
344            durable,
345        })
346    }
347
348    /// Flush pack bytes written so far to the underlying writer.
349    ///
350    /// Used by hosted sync's interleaved build/send path: after each complete
351    /// entry is added, the sender flushes and drains full chunks from the file.
352    pub fn flush_pack(&mut self) -> Result<()> {
353        if let Some(writer) = self.pack_writer.as_mut() {
354            writer.flush().map_err(StoreError::from)?;
355        }
356        Ok(())
357    }
358
359    #[cfg(feature = "source-transfer")]
360    pub(super) fn scratch_root(&self) -> &std::path::Path {
361        &self.bucket_dir
362    }
363
364    /// Add an object with a content-hash id.
365    pub fn add(&mut self, hash: ContentHash, obj_type: ObjectType, data: Vec<u8>) -> Result<()> {
366        self.add_id(PackObjectId::Hash(hash), obj_type, data)
367    }
368
369    /// Add an object with an explicit id. Mirrors [`super::PackBuilder::add_id`].
370    ///
371    /// # Memory shape
372    ///
373    /// Per-entry, the only allocations are:
374    ///
375    /// - `data: Vec<u8>` (the input, owned by the caller — comes from
376    ///   gix' `find_object` and isn't ours to stream further).
377    /// - A ~40-byte stack scratch for the entry header.
378    /// - zstd's internal compression context (~128 KB constant).
379    /// - One 50-byte index-bucket entry buffered into the bucket's
380    ///   `BufWriter`.
381    ///
382    /// The compressed payload is **never materialized** as a `Vec<u8>` —
383    /// it streams directly through `zstd::stream::write::Encoder` into
384    /// the pack writer. The pack format requires a `compressed_size`
385    /// varint *before* the compressed bytes, which we don't know yet
386    /// when we write the header; we reserve a 10-byte placeholder and
387    /// seek-back to patch it after the encoder finishes. Heddle's
388    /// varint decoder accepts non-canonical encodings (it walks
389    /// continuation bits without enforcing minimum-byte form), so the
390    /// padded write decodes back to the same value any reader expects.
391    pub fn add_id(
392        &mut self,
393        id: PackObjectId,
394        obj_type: ObjectType,
395        data: impl AsRef<[u8]>,
396    ) -> Result<()> {
397        let data = data.as_ref();
398        // Compute the entry's offset relative to the header from our logical
399        // append cursor. Asking the underlying file for its position would
400        // flush the BufWriter on every object, which defeats the streaming
401        // sender's chunk-sized drain cadence.
402        let pw = self
403            .pack_writer
404            .as_mut()
405            .ok_or_else(|| StoreError::InvalidObject("pack builder is finalized".into()))?;
406        let entry_start = self.pack_position;
407        let offset = entry_start
408            .checked_sub(self.header_offset)
409            .ok_or_else(|| StoreError::InvalidObject("pack position precedes its header".into()))?;
410
411        self.total_uncompressed = self
412            .total_uncompressed
413            .checked_add(data.len() as u64)
414            .ok_or_else(|| StoreError::InvalidObject("pack decoded size overflow".into()))?;
415
416        // Phase 1: write the entry header up to (but not including) the
417        // compressed-size varint. Always small, fits in `entry_header_buf`.
418        let mut header_buf = Vec::with_capacity(40);
419        id.encode_tagged(&mut header_buf);
420        let encoded_type = if obj_type == ObjectType::AnnotatedTag {
421            ObjectType::Blob
422        } else {
423            obj_type
424        };
425        super::varint::encode_type_and_size(encoded_type, data.len() as u64, &mut header_buf);
426        pw.write_all(&header_buf).map_err(StoreError::from)?;
427        self.pack_position = self
428            .pack_position
429            .checked_add(header_buf.len() as u64)
430            .ok_or_else(|| {
431                StoreError::InvalidObject("streaming pack position overflow".to_string())
432            })?;
433        // Only consumed by the zstd-enabled streaming branch below, but
434        // we compute it here while we already have `header_buf`'s length
435        // in scope.
436        #[cfg(feature = "zstd")]
437        let csize_pos = entry_start + header_buf.len() as u64;
438
439        // Phase 2: stream the compressed payload. We branch here on
440        // whether to compress at all — for tiny objects (`< min_size`)
441        // the bulk path traditionally wrote raw bytes to skip zstd
442        // overhead, and the reader's existing `compressed_size ==
443        // uncompressed_size` heuristic in `pack_reader.rs:128` reads
444        // raw entries back unchanged. We preserve that policy.
445        // `want_compress` gates the zstd path. Even with the feature
446        // enabled we fall through to raw for tiny entries (where
447        // zstd's frame overhead dominates) or when the caller
448        // explicitly disabled compression in `CompressionConfig`.
449        // Without the `zstd` Cargo feature, every entry takes the raw
450        // branch — same fallback shape as `compress_pack_payload`.
451        let want_compress: bool;
452        #[cfg(feature = "zstd")]
453        {
454            want_compress = self.compression.enabled && data.len() >= self.compression.min_size;
455        }
456        #[cfg(not(feature = "zstd"))]
457        {
458            want_compress = false;
459        }
460        if !want_compress {
461            // Raw entry: known compressed_size = data.len(). One canonical
462            // varint + the data itself. No seek-back needed.
463            let mut csize_buf = Vec::with_capacity(10);
464            super::varint::encode_varint(data.len() as u64, &mut csize_buf);
465            pw.write_all(&csize_buf).map_err(StoreError::from)?;
466            self.pack_position = self
467                .pack_position
468                .checked_add(csize_buf.len() as u64)
469                .ok_or_else(|| {
470                    StoreError::InvalidObject("streaming pack position overflow".to_string())
471                })?;
472            pw.write_all(data).map_err(StoreError::from)?;
473            self.pack_position = self
474                .pack_position
475                .checked_add(data.len() as u64)
476                .ok_or_else(|| {
477                    StoreError::InvalidObject("streaming pack position overflow".to_string())
478                })?;
479            self.total_compressed += data.len() as u64;
480        } else {
481            #[cfg(feature = "zstd")]
482            {
483                // Streaming entry: reserve 10 bytes for compressed_size,
484                // stream-compress the payload, then seek back to patch.
485                pw.write_all(&[0u8; CSIZE_PLACEHOLDER_LEN])
486                    .map_err(StoreError::from)?;
487                self.pack_position = self
488                    .pack_position
489                    .checked_add(CSIZE_PLACEHOLDER_LEN as u64)
490                    .ok_or_else(|| {
491                        StoreError::InvalidObject("streaming pack position overflow".to_string())
492                    })?;
493                let body_start = self.pack_position;
494                let compressed_size;
495                {
496                    let mut counting = CountingWriter::new(&mut *pw);
497                    let mut enc =
498                        zstd::stream::write::Encoder::new(&mut counting, self.compression.level)
499                            .map_err(StoreError::from)?;
500                    // Pass the source size so the zstd frame's optional
501                    // Frame Content Size field is set — lets decoders
502                    // preallocate output buffers and validates that we
503                    // wrote exactly what we promised at finish().
504                    enc.set_pledged_src_size(Some(data.len() as u64))
505                        .map_err(StoreError::from)?;
506                    enc.write_all(data).map_err(StoreError::from)?;
507                    enc.finish().map_err(StoreError::from)?;
508                    compressed_size = counting.written;
509                }
510                self.pack_position =
511                    self.pack_position
512                        .checked_add(compressed_size)
513                        .ok_or_else(|| {
514                            StoreError::InvalidObject(
515                                "streaming pack position overflow".to_string(),
516                            )
517                        })?;
518                let body_end = body_start.checked_add(compressed_size).ok_or_else(|| {
519                    StoreError::InvalidObject("streaming pack position overflow".to_string())
520                })?;
521                self.total_compressed += compressed_size;
522
523                // Seek back over the placeholder, write a 10-byte
524                // non-canonical varint encoding the actual compressed_size,
525                // then seek forward to where we left off so subsequent
526                // adds append correctly.
527                let mut csize_bytes = [0u8; CSIZE_PLACEHOLDER_LEN];
528                encode_varint_padded_to_10(compressed_size, &mut csize_bytes);
529                pw.flush().map_err(StoreError::from)?;
530                let inner = pw.get_mut();
531                inner
532                    .seek(SeekFrom::Start(csize_pos))
533                    .map_err(StoreError::from)?;
534                inner.write_all(&csize_bytes).map_err(StoreError::from)?;
535                // Pack entries have no explicit compression bit: readers treat
536                // equal stored/logical lengths as raw bytes. A zstd frame can
537                // occasionally be exactly as long as its input, so replace
538                // that ambiguous frame in place with the original payload.
539                if compressed_size == data.len() as u64 {
540                    inner
541                        .seek(SeekFrom::Start(body_start))
542                        .map_err(StoreError::from)?;
543                    inner.write_all(data).map_err(StoreError::from)?;
544                }
545                inner
546                    .seek(SeekFrom::Start(body_end))
547                    .map_err(StoreError::from)?;
548            }
549            #[cfg(not(feature = "zstd"))]
550            {
551                // Unreachable: `want_compress` is forced to `false`
552                // when the `zstd` feature is off.
553                unreachable!("compression branch reached without `zstd` feature");
554            }
555        }
556
557        self.add_index_entry(id, offset)?;
558        self.record_count += 1;
559        self.object_count += 1;
560        Ok(())
561    }
562
563    /// Add one compact frame shared by several logical blob, tree, or state ids.
564    ///
565    /// `stored_data` is either the raw frame or a zstd frame whose decoded
566    /// length is `uncompressed_size`. Every id is indexed at the same physical
567    /// record; [`super::PackReader`] verifies and decodes the complete frame
568    /// before selecting the requested logical object.
569    pub fn add_shared_frame(
570        &mut self,
571        ids: &[PackObjectId],
572        obj_type: ObjectType,
573        uncompressed_size: usize,
574        stored_data: &[u8],
575    ) -> Result<()> {
576        if ids.is_empty() {
577            return Err(StoreError::InvalidObject(
578                "compact frame must contain at least one object".to_string(),
579            ));
580        }
581        if !matches!(
582            obj_type,
583            ObjectType::Blob | ObjectType::Tree | ObjectType::State
584        ) {
585            return Err(StoreError::InvalidObject(
586                "shared compact frames may contain only blobs, trees, or states".to_string(),
587            ));
588        }
589        let unique = ids.iter().copied().collect::<HashSet<_>>();
590        if unique.len() != ids.len() {
591            return Err(StoreError::InvalidObject(
592                "compact frame contains duplicate object ids".to_string(),
593            ));
594        }
595        if ids.iter().any(|id| {
596            !matches!(
597                (obj_type, id),
598                (ObjectType::Blob | ObjectType::Tree, PackObjectId::Hash(_))
599                    | (ObjectType::State, PackObjectId::StateId(_))
600            )
601        }) {
602            return Err(StoreError::InvalidObject(
603                "compact frame id kind does not match its object type".to_string(),
604            ));
605        }
606
607        let entry_start = self.pack_position;
608        let offset = entry_start
609            .checked_sub(self.header_offset)
610            .expect("header offset should precede compact frame");
611        let mut header = Vec::with_capacity(48);
612        ids[0].encode_tagged(&mut header);
613        super::varint::encode_type_and_size(obj_type, uncompressed_size as u64, &mut header);
614        super::varint::encode_varint(stored_data.len() as u64, &mut header);
615        let writer = self
616            .pack_writer
617            .as_mut()
618            .expect("add_shared_frame called after finalize");
619        writer.write_all(&header).map_err(StoreError::from)?;
620        writer.write_all(stored_data).map_err(StoreError::from)?;
621        self.pack_position = self
622            .pack_position
623            .checked_add((header.len() + stored_data.len()) as u64)
624            .ok_or_else(|| {
625                StoreError::InvalidObject("streaming pack position overflow".to_string())
626            })?;
627        self.total_uncompressed = self
628            .total_uncompressed
629            .saturating_add(uncompressed_size as u64);
630        self.total_compressed = self
631            .total_compressed
632            .saturating_add(stored_data.len() as u64);
633        self.record_count = self
634            .record_count
635            .checked_add(1)
636            .ok_or_else(|| StoreError::InvalidObject("pack record count overflow".to_string()))?;
637        for id in ids {
638            self.add_index_entry(*id, offset)?;
639        }
640        self.object_count = self
641            .object_count
642            .checked_add(ids.len() as u64)
643            .ok_or_else(|| StoreError::InvalidObject("pack object count overflow".to_string()))?;
644        Ok(())
645    }
646
647    fn add_index_entry(&mut self, id: PackObjectId, offset: u64) -> Result<()> {
648        let bucket_idx = bucket_index_for(&id);
649        let bucket = self.get_or_open_bucket(bucket_idx)?;
650        let mut entry = Vec::with_capacity(33 + 8);
651        id.encode_tagged(&mut entry);
652        entry.extend_from_slice(&offset.to_be_bytes());
653        bucket.write_all(&entry).map_err(StoreError::from)
654    }
655
656    fn get_or_open_bucket(&mut self, idx: usize) -> Result<&mut BufWriter<File>> {
657        self.bucket_access_tick = self.bucket_access_tick.wrapping_add(1);
658        let last_used = self.bucket_access_tick;
659        if self.bucket_writers[idx].is_none() {
660            if self.open_bucket_writers >= MAX_OPEN_BUCKET_WRITERS {
661                self.evict_lru_bucket()?;
662            }
663            let path = &self.bucket_paths[idx];
664            let f = OpenOptions::new()
665                .create(true)
666                .append(true)
667                .open(path)
668                .map_err(StoreError::from)?;
669            self.bucket_writers[idx] = Some(BucketWriter {
670                writer: BufWriter::new(f),
671                last_used,
672            });
673            self.open_bucket_writers += 1;
674        } else if let Some(bucket) = self.bucket_writers[idx].as_mut() {
675            bucket.last_used = last_used;
676        }
677        Ok(&mut self.bucket_writers[idx]
678            .as_mut()
679            .ok_or_else(|| StoreError::InvalidObject("index bucket writer unavailable".into()))?
680            .writer)
681    }
682
683    fn evict_lru_bucket(&mut self) -> Result<()> {
684        let Some((idx, _)) = self
685            .bucket_writers
686            .iter()
687            .enumerate()
688            .filter_map(|(idx, bucket)| bucket.as_ref().map(|bucket| (idx, bucket.last_used)))
689            .min_by_key(|(_, last_used)| *last_used)
690        else {
691            return Ok(());
692        };
693
694        if let Some(mut bucket) = self.bucket_writers[idx].take() {
695            bucket.writer.flush().map_err(StoreError::from)?;
696            self.open_bucket_writers -= 1;
697        }
698        Ok(())
699    }
700
701    /// Close the pack: patch the header count, append the BLAKE3
702    /// trailer, build the sorted index from bucket files, and clean up
703    /// the bucket directory. Returns `(pack_writer, index_bytes,
704    /// stats)` so the caller can install the pack into its store.
705    ///
706    /// On any failure the bucket dir is left in place; rerunning the
707    /// import will overwrite stale bucket files (they're keyed by
708    /// fixed name, not content) so this isn't a correctness issue —
709    /// just a small amount of disk churn until the next clean
710    /// finalize.
711    pub fn finalize(mut self) -> Result<(W, PackStats)> {
712        // 1. Flush every bucket so reads in the next phase see all
713        //    queued entries. `flatten()` skips the never-opened slots.
714        for bucket in self.bucket_writers.iter_mut().flatten() {
715            bucket.writer.flush().map_err(StoreError::from)?;
716        }
717        // Drop the writers so the OS file handles close before we
718        // re-open the same paths for reading.
719        for slot in self.bucket_writers.iter_mut() {
720            *slot = None;
721        }
722        self.open_bucket_writers = 0;
723
724        // 2. Patch the pack header with the real object count unless the
725        //    caller declared it up front, then re-stream the [header][body]
726        //    bytes to compute the trailer checksum.
727        let bw = self
728            .pack_writer
729            .take()
730            .ok_or_else(|| StoreError::InvalidObject("pack builder is finalized".into()))?;
731        let mut writer = bw
732            .into_inner()
733            .map_err(|e| StoreError::from(std::io::Error::other(e.to_string())))?;
734        if let Some(expected) = self.declared_object_count {
735            if expected != self.record_count {
736                return Err(StoreError::InvalidObject(format!(
737                    "streaming pack declared {expected} record(s) but added {}",
738                    self.record_count
739                )));
740            }
741        } else {
742            writer
743                .seek(SeekFrom::Start(self.header_offset))
744                .map_err(StoreError::from)?;
745            let mut header_bytes = Vec::with_capacity(16);
746            write_container_header(&mut header_bytes, pack_container_spec(), self.record_count);
747            writer.write_all(&header_bytes).map_err(StoreError::from)?;
748        }
749
750        // 3. Hash the on-disk content from header_offset to current
751        //    position (which is just past the body). One sequential
752        //    pass; the BufWriter we drained is gone so this read is
753        //    on the raw writer.
754        writer
755            .seek(SeekFrom::Start(self.header_offset))
756            .map_err(StoreError::from)?;
757        let mut hasher = blake3::Hasher::new();
758        let mut buf = vec![0u8; 64 * 1024];
759        loop {
760            let n = writer.read(&mut buf).map_err(StoreError::from)?;
761            if n == 0 {
762                break;
763            }
764            hasher.update(&buf[..n]);
765        }
766        let checksum = hasher.finalize();
767
768        // 4. Append the trailer checksum.
769        writer.seek(SeekFrom::End(0)).map_err(StoreError::from)?;
770        writer
771            .write_all(checksum.as_bytes())
772            .map_err(StoreError::from)?;
773        writer.flush().map_err(StoreError::from)?;
774        // L7: durable staged pack before return (File fsync; Cursor no-op).
775        if self.durable {
776            writer
777                .sync_data_for_durability()
778                .map_err(StoreError::from)?;
779        }
780
781        // Sort fixed-size runs and merge on disk. Hash-prefix buckets only
782        // partition I/O; even one adversarial bucket uses a fixed working set.
783        let idx_file = File::create(&self.index_path).map_err(StoreError::from)?;
784        let mut idx_writer = BufWriter::new(idx_file);
785        write_index_header(&mut idx_writer, self.object_count)?;
786        let mut entries_written: u64 = 0;
787        for path in self.bucket_paths.iter() {
788            if !path.exists() {
789                continue;
790            }
791            let run = super::disk_sort::sort_records::<41>(File::open(path)?, &self.bucket_dir)?;
792            let mut input = std::io::BufReader::new(run.as_file());
793            while let Some(record) = super::disk_sort::read_record::<41>(&mut input)? {
794                let (id, _) = PackObjectId::decode_tagged(&record)?;
795                let offset = u64::from_be_bytes(record[33..].try_into().map_err(|_| {
796                    StoreError::InvalidObject("index bucket offset truncated".into())
797                })?);
798                write_index_entry(&mut idx_writer, id, offset)?;
799                entries_written += 1;
800            }
801        }
802        idx_writer.flush().map_err(StoreError::from)?;
803        // L7: durable staged index file + parent dirent for rename/read.
804        let _idx_file = idx_writer
805            .into_inner()
806            .map_err(|e| StoreError::from(std::io::Error::other(e.to_string())))?;
807        #[cfg(feature = "fs")]
808        if self.durable {
809            _idx_file.sync_all().map_err(StoreError::from)?;
810            if let Some(parent) = self.index_path.parent() {
811                heddle_fs_prims::fs_atomic::sync_directory(parent).map_err(StoreError::from)?;
812            }
813        }
814        debug_assert_eq!(
815            entries_written, self.object_count,
816            "streaming index entry count drifted from add() count"
817        );
818
819        // 6. Clean up the bucket dir so the heddle store doesn't carry
820        //    transient artifacts. Deletion failures are non-fatal —
821        //    the dir is uniquely named per import so leftovers are at
822        //    worst stale, not corrupting.
823        for path in self.bucket_paths.iter() {
824            let _ = std::fs::remove_file(path);
825        }
826        self.scratch_lease.take();
827        let _ = std::fs::remove_file(self.bucket_dir.join(super::scratch::LEASE_FILE));
828        let _ = std::fs::remove_dir(&self.bucket_dir);
829        self.finalized = true;
830
831        let stats = PackStats {
832            object_count: self.object_count,
833            total_uncompressed: self.total_uncompressed,
834            total_compressed: self.total_compressed,
835            delta_count: 0,
836            compression_ratio: if self.total_uncompressed == 0 {
837                0.0
838            } else {
839                self.total_compressed as f64 / self.total_uncompressed as f64
840            },
841        };
842
843        Ok((writer, stats))
844    }
845}
846
847/// Write the index container header to `out`. Mirrors
848/// [`PackIndex::to_bytes`]'s prefix exactly (4-byte magic, 4-byte
849/// big-endian version, 8-byte big-endian count) so the streaming and
850/// in-memory builders produce the same index container.
851fn write_index_header<W: Write>(out: &mut W, count: u64) -> Result<()> {
852    super::pack_index::index_header().write_to(out, count)
853}
854
855/// Append one `(id, offset)` index entry to `out` using the compact
856/// [`PackIndex::to_bytes`] encoding.
857fn write_index_entry<W: Write>(out: &mut W, id: PackObjectId, offset: u64) -> Result<()> {
858    let buf = super::pack_index::encode_index_entry(id, offset);
859    out.write_all(&buf).map_err(StoreError::from)
860}
861
862/// Encode a `u64` as a non-canonical 10-byte LEB128 varint. The first
863/// 9 bytes always set the continuation bit (`0x80`), the 10th never
864/// does — so the decoder reads exactly 10 bytes regardless of the
865/// value. Used by the streaming path to reserve a fixed-width
866/// placeholder for `compressed_size` before stream-compressing the
867/// payload, then patch the placeholder with the actual size after.
868///
869/// `decode_varint` ignores the canonicalness of the encoding (it
870/// walks continuation bits without checking minimum-byte form), so
871/// the value round-trips exactly. Cost is up to 9 wasted bytes per
872/// entry, ~115 KB on a 13 K-entry import — negligible relative to
873/// the pack body.
874#[cfg(feature = "zstd")]
875fn encode_varint_padded_to_10(value: u64, out: &mut [u8; 10]) {
876    let mut v = value;
877    for slot in out.iter_mut().take(9) {
878        *slot = 0x80 | ((v & 0x7F) as u8);
879        v >>= 7;
880    }
881    out[9] = (v & 0x7F) as u8;
882}
883
884impl<W: Write + Read + Seek> Drop for StreamingPackBuilder<W> {
885    fn drop(&mut self) {
886        if self.finalized {
887            return;
888        }
889        // Best-effort cleanup of bucket dir on abort. Errors here are
890        // suppressed because Drop can't propagate them.
891        for path in self.bucket_paths.iter() {
892            let _ = std::fs::remove_file(path);
893        }
894        self.scratch_lease.take();
895        let _ = std::fs::remove_file(self.bucket_dir.join(super::scratch::LEASE_FILE));
896        let _ = std::fs::remove_dir(&self.bucket_dir);
897    }
898}
899
900/// Map a `PackObjectId` to one of `TOTAL_BUCKETS` buckets. The variant
901/// (Hash vs StateId) picks the upper half; the first byte of the
902/// inner id picks the slot within the half.
903fn bucket_index_for(id: &PackObjectId) -> usize {
904    match id {
905        PackObjectId::Hash(h) => HASH_VARIANT * BUCKETS_PER_VARIANT + h.as_bytes()[0] as usize,
906        PackObjectId::StateId(c) => {
907            CHANGEID_VARIANT * BUCKETS_PER_VARIANT + c.as_bytes()[0] as usize
908        }
909        PackObjectId::AnnotatedTag(hash) => {
910            ANNOTATED_TAG_VARIANT * BUCKETS_PER_VARIANT + hash.as_bytes()[0] as usize
911        }
912    }
913}
914
915// ---------------------- Tests ----------------------
916
917#[cfg(test)]
918mod tests {
919    use std::io::Cursor;
920
921    use super::*;
922    use crate::{
923        object::StateId,
924        store::pack::{PackReader, PackStats},
925    };
926
927    fn deterministic_hash(seed: u8) -> ContentHash {
928        // Spread `seed` across the high byte so different seeds end up
929        // in different hash-prefix buckets. We don't actually want
930        // collisions in the tests that check distribution.
931        let mut bytes = [0u8; 32];
932        bytes[0] = seed;
933        for (i, b) in bytes.iter_mut().enumerate().skip(1) {
934            *b = seed.wrapping_mul(31).wrapping_add(i as u8);
935        }
936        ContentHash::from_bytes(bytes)
937    }
938
939    fn deterministic_state_id(seed: u8) -> StateId {
940        let mut bytes = [0u8; 32];
941        bytes[0] = seed;
942        for (i, b) in bytes.iter_mut().enumerate().skip(1) {
943            *b = seed.wrapping_add(i as u8 * 7);
944        }
945        StateId::from_bytes(bytes)
946    }
947
948    /// Test rig: returns the builder, the bucket dir (for cleanup
949    /// inspection), and the index path the builder will write at
950    /// finalize. The index path lives in the temp dir so it gets
951    /// auto-cleaned with `tmp`.
952    fn fresh_builder(
953        tmp: &tempfile::TempDir,
954    ) -> (StreamingPackBuilder<Cursor<Vec<u8>>>, PathBuf, PathBuf) {
955        let bucket_dir = tmp.path().join("buckets");
956        let index_path = tmp.path().join("test.idx");
957        let cursor = Cursor::new(Vec::<u8>::new());
958        let b = StreamingPackBuilder::new(
959            cursor,
960            index_path.clone(),
961            CompressionConfig::default(),
962            bucket_dir.clone(),
963        )
964        .unwrap();
965        (b, bucket_dir, index_path)
966    }
967
968    /// Finalize the builder and return `(pack_bytes, index_bytes, stats)`.
969    /// The index bytes are read back from the file the builder wrote
970    /// to — verifying that the streaming index path actually produced
971    /// readable bytes.
972    fn finalize_cursor(
973        b: StreamingPackBuilder<Cursor<Vec<u8>>>,
974        index_path: &std::path::Path,
975    ) -> (Vec<u8>, Vec<u8>, PackStats) {
976        let (cursor, stats) = b.finalize().unwrap();
977        let index_bytes = std::fs::read(index_path).unwrap();
978        (cursor.into_inner(), index_bytes, stats)
979    }
980
981    #[test]
982    fn empty_pack_finalizes_to_valid_zero_count_pack() {
983        let tmp = tempfile::TempDir::new().unwrap();
984        let (b, bucket_dir, idx_path) = fresh_builder(&tmp);
985        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
986
987        assert_eq!(stats.object_count, 0);
988        // PackReader can parse the empty pack and reports zero objects.
989        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
990        assert!(reader.list_ids().unwrap().is_empty());
991        // Bucket dir was removed.
992        assert!(
993            !bucket_dir.exists(),
994            "bucket dir should be cleaned on successful finalize"
995        );
996    }
997
998    #[test]
999    fn single_blob_with_hash_id_round_trips() {
1000        let tmp = tempfile::TempDir::new().unwrap();
1001        let (mut b, _, idx_path) = fresh_builder(&tmp);
1002        let hash = deterministic_hash(0x42);
1003        let payload = b"hello, streaming pack".to_vec();
1004        b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1005        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1006
1007        assert_eq!(stats.object_count, 1);
1008        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1009        let id = PackObjectId::Hash(hash);
1010        assert!(reader.has_object(&id).unwrap());
1011        let (got_type, got_data) = reader.get_object(&id).unwrap().unwrap();
1012        assert_eq!(got_type, ObjectType::Blob);
1013        assert_eq!(got_data, payload);
1014    }
1015
1016    #[test]
1017    fn single_state_with_change_id_round_trips() {
1018        let tmp = tempfile::TempDir::new().unwrap();
1019        let (mut b, _, idx_path) = fresh_builder(&tmp);
1020        let cid = deterministic_state_id(0xa5);
1021        let payload = b"serialized-state-bytes".to_vec();
1022        b.add_id(
1023            PackObjectId::StateId(cid),
1024            ObjectType::State,
1025            payload.clone(),
1026        )
1027        .unwrap();
1028        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1029
1030        assert_eq!(stats.object_count, 1);
1031        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1032        let id = PackObjectId::StateId(cid);
1033        let (ty, data) = reader.get_object(&id).unwrap().unwrap();
1034        assert_eq!(ty, ObjectType::State);
1035        assert_eq!(data, payload);
1036    }
1037
1038    #[test]
1039    fn shared_compact_tree_frame_reconstructs_each_indexed_object() {
1040        use crate::object::{Tree, TreeEntry};
1041
1042        let tmp = tempfile::TempDir::new().unwrap();
1043        let (mut builder, _, index_path) = fresh_builder(&tmp);
1044        let blob = deterministic_hash(0x33);
1045        let trees = vec![
1046            Tree::from_entries(vec![TreeEntry::file("a", blob, false).unwrap()]),
1047            Tree::from_entries(vec![TreeEntry::file("b", blob, true).unwrap()]),
1048        ];
1049        let ids = trees
1050            .iter()
1051            .map(|tree| PackObjectId::Hash(tree.hash()))
1052            .collect::<Vec<_>>();
1053        let frame = heddle_object_model::compact::encode_tree_frame(&trees).unwrap();
1054        builder
1055            .add_shared_frame(&ids, ObjectType::Tree, frame.len(), &frame)
1056            .unwrap();
1057        let (pack, index, stats) = finalize_cursor(builder, &index_path);
1058        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1059
1060        assert_eq!(stats.object_count, 2);
1061        assert_eq!(
1062            reader.encoded_payload_bytes(ObjectType::Tree).unwrap(),
1063            frame.len() as u64
1064        );
1065        let hosted = ids
1066            .iter()
1067            .zip(&trees)
1068            .map(|(id, tree)| {
1069                (
1070                    *id,
1071                    ObjectType::Tree,
1072                    tree.encode_canonical().unwrap().len() as u64,
1073                )
1074            })
1075            .collect::<Vec<_>>();
1076        assert!(
1077            reader
1078                .copy_hosted_encoded_subset(&hosted)
1079                .unwrap()
1080                .is_none(),
1081            "repository-local compact frames must use the hosted fallback"
1082        );
1083        for (id, tree) in ids.iter().zip(&trees) {
1084            let (object_type, bytes) = reader.get_object(id).unwrap().unwrap();
1085            assert_eq!(object_type, ObjectType::Tree);
1086            assert_eq!(bytes, tree.encode_canonical().unwrap());
1087        }
1088    }
1089
1090    #[test]
1091    fn compact_tree_extraction_rejects_an_index_alias_with_the_wrong_typed_hash() {
1092        use crate::object::{Tree, TreeEntry};
1093
1094        let tmp = tempfile::TempDir::new().unwrap();
1095        let (mut builder, _, index_path) = fresh_builder(&tmp);
1096        let blob = deterministic_hash(0x34);
1097        let trees = vec![
1098            Tree::from_entries(vec![TreeEntry::file("a", blob, false).unwrap()]),
1099            Tree::from_entries(vec![TreeEntry::file("b", blob, true).unwrap()]),
1100        ];
1101        let wrong_hash = ContentHash::compute_typed("tree", b"not the second tree");
1102        assert_ne!(wrong_hash, trees[1].hash());
1103        let ids = vec![
1104            PackObjectId::Hash(trees[0].hash()),
1105            PackObjectId::Hash(wrong_hash),
1106        ];
1107        let frame = heddle_object_model::compact::encode_tree_frame(&trees).unwrap();
1108        builder
1109            .add_shared_frame(&ids, ObjectType::Tree, frame.len(), &frame)
1110            .unwrap();
1111        let (pack, index, _) = finalize_cursor(builder, &index_path);
1112        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1113
1114        let error = reader.get_object(&ids[1]).unwrap_err();
1115        assert!(
1116            error
1117                .to_string()
1118                .contains("does not contain indexed object"),
1119            "extraction must derive and verify the tree's typed hash: {error}"
1120        );
1121    }
1122
1123    #[test]
1124    fn shared_lineage_blob_frame_reconstructs_each_indexed_object() {
1125        let tmp = tempfile::TempDir::new().unwrap();
1126        let (mut builder, _, index_path) = fresh_builder(&tmp);
1127        let bodies = [b"newest version".as_slice(), b"older version".as_slice()];
1128        let ids = bodies
1129            .iter()
1130            .map(|body| PackObjectId::Hash(ContentHash::compute_typed("blob", body)))
1131            .collect::<Vec<_>>();
1132        let frame = heddle_object_model::compact::encode_blob_frame(&bodies).unwrap();
1133        builder
1134            .add_shared_frame(&ids, ObjectType::Blob, frame.len(), &frame)
1135            .unwrap();
1136        let (pack, index, stats) = finalize_cursor(builder, &index_path);
1137        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1138
1139        assert_eq!(stats.object_count, bodies.len() as u64);
1140        assert_eq!(
1141            reader.encoded_payload_bytes(ObjectType::Blob).unwrap(),
1142            frame.len() as u64
1143        );
1144        for (id, expected) in ids.iter().zip(bodies) {
1145            let (object_type, actual) = reader.get_object(id).unwrap().unwrap();
1146            assert_eq!(object_type, ObjectType::Blob);
1147            assert_eq!(actual, expected);
1148            let PackObjectId::Hash(hash) = id else {
1149                unreachable!("blob ids are hashes")
1150            };
1151            assert_eq!(
1152                reader.get_hashed_object_type(hash).unwrap(),
1153                Some(ObjectType::Blob)
1154            );
1155            assert_eq!(
1156                reader.get_hashed_object_size(hash).unwrap(),
1157                Some(expected.len() as u64)
1158            );
1159        }
1160    }
1161
1162    #[test]
1163    fn ordinary_blob_starting_with_frame_magic_remains_ordinary() {
1164        let tmp = tempfile::TempDir::new().unwrap();
1165        let (mut builder, _, index_path) = fresh_builder(&tmp);
1166        let body = b"HCB2 arbitrary user content".to_vec();
1167        let hash = ContentHash::compute_typed("blob", &body);
1168        builder.add(hash, ObjectType::Blob, body.clone()).unwrap();
1169        let (pack, index, _) = finalize_cursor(builder, &index_path);
1170        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1171
1172        assert_eq!(
1173            reader.get_object(&PackObjectId::Hash(hash)).unwrap(),
1174            Some((ObjectType::Blob, body))
1175        );
1176    }
1177
1178    #[test]
1179    fn corrupt_compact_frame_byte_invalidates_every_contained_object() {
1180        use crate::object::{Tree, TreeEntry};
1181
1182        let tmp = tempfile::TempDir::new().unwrap();
1183        let (mut builder, _, index_path) = fresh_builder(&tmp);
1184        let blob = deterministic_hash(0x44);
1185        let trees = vec![
1186            Tree::from_entries(vec![TreeEntry::file("a", blob, false).unwrap()]),
1187            Tree::from_entries(vec![TreeEntry::file("b", blob, true).unwrap()]),
1188        ];
1189        let ids = trees
1190            .iter()
1191            .map(|tree| PackObjectId::Hash(tree.hash()))
1192            .collect::<Vec<_>>();
1193        let mut frame = heddle_object_model::compact::encode_tree_frame(&trees).unwrap();
1194        let corrupt_at = frame.len() / 2;
1195        frame[corrupt_at] ^= 0x01;
1196        builder
1197            .add_shared_frame(&ids, ObjectType::Tree, frame.len(), &frame)
1198            .unwrap();
1199        let (pack, index, _) = finalize_cursor(builder, &index_path);
1200        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1201
1202        for id in ids {
1203            let error = reader.get_object(&id).unwrap_err();
1204            assert!(
1205                error
1206                    .to_string()
1207                    .contains("compact frame checksum mismatch"),
1208                "unexpected error for {id:?}: {error}"
1209            );
1210        }
1211    }
1212
1213    #[test]
1214    fn corrupt_blob_frame_byte_invalidates_every_contained_object() {
1215        let tmp = tempfile::TempDir::new().unwrap();
1216        let (mut builder, _, index_path) = fresh_builder(&tmp);
1217        let bodies = [b"newest version".as_slice(), b"older version".as_slice()];
1218        let ids = bodies
1219            .iter()
1220            .map(|body| PackObjectId::Hash(ContentHash::compute_typed("blob", body)))
1221            .collect::<Vec<_>>();
1222        let mut frame = heddle_object_model::compact::encode_blob_frame(&bodies).unwrap();
1223        let corrupt_at = frame.len() / 2;
1224        frame[corrupt_at] ^= 0x01;
1225        builder
1226            .add_shared_frame(&ids, ObjectType::Blob, frame.len(), &frame)
1227            .unwrap();
1228        let (pack, index, _) = finalize_cursor(builder, &index_path);
1229        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1230
1231        for id in ids {
1232            let error = reader.get_object(&id).unwrap_err();
1233            assert!(
1234                error
1235                    .to_string()
1236                    .contains("compact frame checksum mismatch"),
1237                "unexpected error for {id:?}: {error}"
1238            );
1239        }
1240    }
1241
1242    #[test]
1243    fn mixed_hash_and_changeid_ids_all_retrievable() {
1244        let tmp = tempfile::TempDir::new().unwrap();
1245        let (mut b, _, idx_path) = fresh_builder(&tmp);
1246        let blob_hash = deterministic_hash(0x10);
1247        let tree_hash = deterministic_hash(0x20);
1248        let state_cid = deterministic_state_id(0x80);
1249
1250        b.add(blob_hash, ObjectType::Blob, b"blob-bytes".to_vec())
1251            .unwrap();
1252        b.add(tree_hash, ObjectType::Tree, b"serialized-tree".to_vec())
1253            .unwrap();
1254        b.add_id(
1255            PackObjectId::StateId(state_cid),
1256            ObjectType::State,
1257            b"serialized-state",
1258        )
1259        .unwrap();
1260
1261        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1262        assert_eq!(stats.object_count, 3);
1263        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1264        assert_eq!(
1265            reader
1266                .get_object(&PackObjectId::Hash(blob_hash))
1267                .unwrap()
1268                .unwrap()
1269                .1,
1270            b"blob-bytes".to_vec()
1271        );
1272        assert_eq!(
1273            reader
1274                .get_object(&PackObjectId::Hash(tree_hash))
1275                .unwrap()
1276                .unwrap()
1277                .1,
1278            b"serialized-tree".to_vec()
1279        );
1280        assert_eq!(
1281            reader
1282                .get_object(&PackObjectId::StateId(state_cid))
1283                .unwrap()
1284                .unwrap()
1285                .1,
1286            b"serialized-state".to_vec()
1287        );
1288    }
1289
1290    #[test]
1291    fn ten_thousand_objects_round_trip_correctly() {
1292        // Stresses the bucket sort: 10K objects spread across
1293        // 256 hash buckets averages 40 entries per bucket — well
1294        // within in-memory sort capacity but covers every bucket.
1295        let tmp = tempfile::TempDir::new().unwrap();
1296        let (mut b, _, idx_path) = fresh_builder(&tmp);
1297        let mut hashes = Vec::with_capacity(10_000);
1298        for i in 0..10_000u32 {
1299            // Use BLAKE3 over the index so first-byte distribution is
1300            // pseudo-uniform across the 256 hash buckets.
1301            let h = blake3::hash(&i.to_le_bytes());
1302            let hash = ContentHash::from_bytes(*h.as_bytes());
1303            hashes.push(hash);
1304            b.add(hash, ObjectType::Blob, format!("payload-{i}").into_bytes())
1305                .unwrap();
1306        }
1307        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1308        assert_eq!(stats.object_count, 10_000);
1309
1310        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1311        assert_eq!(reader.list_ids().unwrap().len(), 10_000);
1312        // Spot-check ten across the range.
1313        for i in [0, 1, 99, 1234, 5_000, 9_999] {
1314            let id = PackObjectId::Hash(hashes[i]);
1315            let (_ty, data) = reader.get_object(&id).unwrap().unwrap();
1316            assert_eq!(data, format!("payload-{i}").into_bytes());
1317        }
1318    }
1319
1320    #[test]
1321    fn bucket_writers_are_lru_capped_below_fd_limit() {
1322        let tmp = tempfile::TempDir::new().unwrap();
1323        let (mut b, _bucket_dir, idx_path) = fresh_builder(&tmp);
1324        let mut ids = Vec::new();
1325
1326        for i in 0..BUCKETS_PER_VARIANT {
1327            let hash = deterministic_hash(i as u8);
1328            ids.push(PackObjectId::Hash(hash));
1329            b.add(hash, ObjectType::Blob, format!("hash-{i}").into_bytes())
1330                .unwrap();
1331            assert!(
1332                b.open_bucket_writers <= MAX_OPEN_BUCKET_WRITERS,
1333                "open bucket writers should stay capped"
1334            );
1335        }
1336
1337        for i in 0..BUCKETS_PER_VARIANT {
1338            let cid = deterministic_state_id(i as u8);
1339            ids.push(PackObjectId::StateId(cid));
1340            b.add_id(
1341                PackObjectId::StateId(cid),
1342                ObjectType::State,
1343                format!("state-{i}").into_bytes(),
1344            )
1345            .unwrap();
1346            assert!(
1347                b.open_bucket_writers <= MAX_OPEN_BUCKET_WRITERS,
1348                "open bucket writers should stay capped"
1349            );
1350        }
1351
1352        for i in 0..BUCKETS_PER_VARIANT {
1353            let hash = deterministic_hash(i as u8);
1354            let id = PackObjectId::AnnotatedTag(hash);
1355            ids.push(id);
1356            b.add_id(
1357                id,
1358                ObjectType::AnnotatedTag,
1359                format!("tag-{i}").into_bytes(),
1360            )
1361            .unwrap();
1362            assert!(
1363                b.open_bucket_writers <= MAX_OPEN_BUCKET_WRITERS,
1364                "open bucket writers should stay capped"
1365            );
1366        }
1367
1368        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1369        assert_eq!(stats.object_count, TOTAL_BUCKETS as u64);
1370        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1371        for id in ids {
1372            assert!(reader.has_object(&id).unwrap(), "missing id {id:?}");
1373        }
1374    }
1375
1376    #[test]
1377    fn index_id_sort_order_matches_packbuilder_output() {
1378        // PackBuilder groups objects by `ObjectType` before encoding,
1379        // which changes the byte offsets relative to a streaming builder
1380        // that writes in added-order. The bytes of the two indices
1381        // therefore can't match exactly. What MUST match is the
1382        // **sort order of ids** — both builders ultimately call
1383        // `PackIndex::sort()` (or the bucket-equivalent), and any
1384        // reader binary-searches against that order.
1385        use crate::store::pack::PackBuilder;
1386        let payloads: Vec<(PackObjectId, ObjectType, Vec<u8>)> = (0..200u32)
1387            .map(|i| {
1388                let h = blake3::hash(&i.to_le_bytes());
1389                (
1390                    PackObjectId::Hash(ContentHash::from_bytes(*h.as_bytes())),
1391                    if i % 3 == 0 {
1392                        ObjectType::Tree
1393                    } else {
1394                        ObjectType::Blob
1395                    },
1396                    format!("body-{i}").into_bytes(),
1397                )
1398            })
1399            .collect();
1400
1401        // Disable delta encoding so the classic builder produces a pack
1402        // shape comparable to the streaming one (which never deltas).
1403        let compression = CompressionConfig {
1404            max_delta_size: 0,
1405            ..CompressionConfig::default()
1406        };
1407        let mut classic = PackBuilder::new(compression);
1408        for (id, ty, data) in payloads.iter() {
1409            classic.add_id(*id, *ty, data.clone());
1410        }
1411        let (classic_pack, classic_index, _) = classic.build().unwrap();
1412        let classic_reader =
1413            PackReader::from_bytes(classic_pack, classic_index, &std::env::temp_dir()).unwrap();
1414
1415        let tmp = tempfile::TempDir::new().unwrap();
1416        let bucket_dir = tmp.path().join("buckets");
1417        let idx_path = tmp.path().join("test.idx");
1418        let cursor = Cursor::new(Vec::<u8>::new());
1419        let mut streaming =
1420            StreamingPackBuilder::new(cursor, idx_path.clone(), compression, bucket_dir).unwrap();
1421        for (id, ty, data) in payloads.iter() {
1422            streaming.add_id(*id, *ty, data.clone()).unwrap();
1423        }
1424        let (streaming_pack, streaming_index, _) = finalize_cursor(streaming, &idx_path);
1425        let streaming_reader =
1426            PackReader::from_bytes(streaming_pack, streaming_index, &std::env::temp_dir()).unwrap();
1427
1428        // Same set of ids in the same sorted order — that's the
1429        // contract for binary search to work.
1430        assert_eq!(
1431            streaming_reader.list_ids().unwrap(),
1432            classic_reader.list_ids().unwrap(),
1433            "streaming and classic indices should report the same id sequence"
1434        );
1435        // Spot-check that each id resolves to a payload that matches
1436        // the classic builder's output (equal bytes after decompression).
1437        for (id, _ty, want) in payloads.iter().take(10).chain(payloads.iter().skip(190)) {
1438            let (_, got) = streaming_reader.get_object(id).unwrap().unwrap();
1439            assert_eq!(&got, want);
1440            let (_, classic_got) = classic_reader.get_object(id).unwrap().unwrap();
1441            assert_eq!(got, classic_got);
1442        }
1443    }
1444
1445    #[test]
1446    fn corrupted_pack_fails_checksum_verification() {
1447        let tmp = tempfile::TempDir::new().unwrap();
1448        let (mut b, _, idx_path) = fresh_builder(&tmp);
1449        b.add(
1450            deterministic_hash(0x01),
1451            ObjectType::Blob,
1452            b"some bytes".to_vec(),
1453        )
1454        .unwrap();
1455        let (mut pack_data, index_data, _) = finalize_cursor(b, &idx_path);
1456        // Flip one byte in the body. The trailer checksum must reject.
1457        let body_byte = 18; // past the 16-byte header
1458        pack_data[body_byte] ^= 0xff;
1459        let result = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir());
1460        assert!(
1461            result.is_err(),
1462            "PackReader should reject pack with mutated body"
1463        );
1464    }
1465
1466    #[test]
1467    fn pack_count_in_header_matches_index_entry_count() {
1468        let tmp = tempfile::TempDir::new().unwrap();
1469        let (mut b, _, idx_path) = fresh_builder(&tmp);
1470        for i in 0..7u8 {
1471            b.add(
1472                deterministic_hash(i),
1473                ObjectType::Blob,
1474                format!("p{i}").into_bytes(),
1475            )
1476            .unwrap();
1477        }
1478        let (pack_data, index_data, _) = finalize_cursor(b, &idx_path);
1479        // Header count is bytes 8..16 (big-endian).
1480        let count = u64::from_be_bytes(pack_data[8..16].try_into().unwrap());
1481        assert_eq!(count, 7);
1482        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1483        assert_eq!(reader.list_ids().unwrap().len(), 7);
1484    }
1485
1486    #[test]
1487    fn declared_pack_count_is_written_before_finalize() {
1488        let tmp = tempfile::TempDir::new().unwrap();
1489        let bucket_dir = tmp.path().join("buckets");
1490        let idx_path = tmp.path().join("test.idx");
1491        let cursor = Cursor::new(Vec::<u8>::new());
1492        let mut b = StreamingPackBuilder::new_with_object_count(
1493            cursor,
1494            idx_path.clone(),
1495            CompressionConfig::default(),
1496            bucket_dir,
1497            2,
1498        )
1499        .unwrap();
1500
1501        b.flush_pack().unwrap();
1502        let initial = b.pack_writer.as_ref().unwrap().get_ref().get_ref().clone();
1503        assert_eq!(u64::from_be_bytes(initial[8..16].try_into().unwrap()), 2);
1504
1505        let hash = deterministic_hash(0x40);
1506        b.add(hash, ObjectType::Blob, b"known-count-entry".to_vec())
1507            .unwrap();
1508        b.flush_pack().unwrap();
1509        let after_add = b.pack_writer.as_ref().unwrap().get_ref().get_ref().clone();
1510        assert_eq!(u64::from_be_bytes(after_add[8..16].try_into().unwrap()), 2);
1511
1512        let second_hash = deterministic_hash(0x41);
1513        b.add(second_hash, ObjectType::Blob, b"second-entry".to_vec())
1514            .unwrap();
1515        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1516
1517        assert_eq!(stats.object_count, 2);
1518        assert_eq!(u64::from_be_bytes(pack_data[8..16].try_into().unwrap()), 2);
1519        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1520        assert!(reader.has_object(&PackObjectId::Hash(hash)).unwrap());
1521        assert!(reader.has_object(&PackObjectId::Hash(second_hash)).unwrap());
1522    }
1523
1524    #[test]
1525    fn declared_pack_count_mismatch_fails_finalize() {
1526        let tmp = tempfile::TempDir::new().unwrap();
1527        let bucket_dir = tmp.path().join("buckets");
1528        let idx_path = tmp.path().join("test.idx");
1529        let cursor = Cursor::new(Vec::<u8>::new());
1530        let mut b = StreamingPackBuilder::new_with_object_count(
1531            cursor,
1532            idx_path,
1533            CompressionConfig::default(),
1534            bucket_dir,
1535            2,
1536        )
1537        .unwrap();
1538
1539        b.add(
1540            deterministic_hash(0x50),
1541            ObjectType::Blob,
1542            b"only-entry".to_vec(),
1543        )
1544        .unwrap();
1545        let error = b.finalize().unwrap_err();
1546
1547        assert!(
1548            error
1549                .to_string()
1550                .contains("streaming pack declared 2 record(s) but added 1")
1551        );
1552    }
1553
1554    #[test]
1555    fn bucket_files_are_cleaned_on_successful_finalize() {
1556        let tmp = tempfile::TempDir::new().unwrap();
1557        let bucket_dir = tmp.path().join("buckets");
1558        let idx_path = tmp.path().join("test.idx");
1559        let cursor = Cursor::new(Vec::<u8>::new());
1560        let mut b = StreamingPackBuilder::new(
1561            cursor,
1562            idx_path.clone(),
1563            CompressionConfig::default(),
1564            bucket_dir.clone(),
1565        )
1566        .unwrap();
1567        for i in 0..50u8 {
1568            b.add(deterministic_hash(i), ObjectType::Blob, vec![i; 32])
1569                .unwrap();
1570        }
1571        // Buckets exist and contain data.
1572        assert!(bucket_dir.exists());
1573        let bucket_count = std::fs::read_dir(&bucket_dir).unwrap().count();
1574        assert!(bucket_count > 0, "bucket dir should hold some files");
1575        let _ = finalize_cursor(b, &idx_path);
1576        assert!(
1577            !bucket_dir.exists(),
1578            "bucket dir should be removed on finalize"
1579        );
1580    }
1581
1582    #[test]
1583    fn bucket_files_are_cleaned_on_drop_without_finalize() {
1584        let tmp = tempfile::TempDir::new().unwrap();
1585        let bucket_dir = tmp.path().join("buckets");
1586        let idx_path = tmp.path().join("test.idx");
1587        {
1588            let cursor = Cursor::new(Vec::<u8>::new());
1589            let mut b = StreamingPackBuilder::new(
1590                cursor,
1591                idx_path.clone(),
1592                CompressionConfig::default(),
1593                bucket_dir.clone(),
1594            )
1595            .unwrap();
1596            for i in 0..10u8 {
1597                b.add(deterministic_hash(i), ObjectType::Blob, vec![0; 32])
1598                    .unwrap();
1599            }
1600            assert!(bucket_dir.exists());
1601            // Drop without finalize — Drop impl should clean up.
1602        }
1603        assert!(
1604            !idx_path.exists(),
1605            "no index file should have been created without finalize"
1606        );
1607        assert!(
1608            !bucket_dir.exists(),
1609            "bucket dir should be removed on Drop when finalize never ran"
1610        );
1611    }
1612
1613    #[test]
1614    fn large_blob_streams_to_disk_without_double_buffering() {
1615        // 4 MiB blob — well under the actual streaming target but big
1616        // enough to confirm we're not buffering the entire pack body in
1617        // RAM. The pack data on disk should be at least 4 MiB; the
1618        // builder's in-memory state is per-object only.
1619        let tmp = tempfile::TempDir::new().unwrap();
1620        let bucket_dir = tmp.path().join("buckets");
1621        let pack_path = tmp.path().join("pack.dat");
1622        let idx_path = tmp.path().join("pack.idx");
1623        let file = std::fs::OpenOptions::new()
1624            .read(true)
1625            .write(true)
1626            .create(true)
1627            .truncate(true)
1628            .open(&pack_path)
1629            .unwrap();
1630        let mut b = StreamingPackBuilder::new(
1631            file,
1632            idx_path.clone(),
1633            CompressionConfig::default(),
1634            bucket_dir,
1635        )
1636        .unwrap();
1637        let payload: Vec<u8> = (0..4 * 1024 * 1024u32).map(|i| (i & 0xff) as u8).collect();
1638        let hash = deterministic_hash(0xff);
1639        b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1640        let (_, stats) = b.finalize().unwrap();
1641        let index_data = std::fs::read(&idx_path).unwrap();
1642        assert_eq!(stats.object_count, 1);
1643        let pack_bytes = std::fs::read(&pack_path).unwrap();
1644        // Pack on disk holds the whole compressed payload + headers
1645        // + trailer. Confirm it round-trips.
1646        let reader = PackReader::from_bytes(pack_bytes, index_data, &std::env::temp_dir()).unwrap();
1647        let (_ty, got) = reader
1648            .get_object(&PackObjectId::Hash(hash))
1649            .unwrap()
1650            .unwrap();
1651        assert_eq!(got, payload);
1652    }
1653
1654    #[test]
1655    fn bucket_distribution_for_random_hashes_is_roughly_uniform() {
1656        // Confirms our sort-time peak memory bound. We accumulate
1657        // 1024 random hashes through a builder and check that no
1658        // single bucket holds more than ~3× the average. (BLAKE3 hash
1659        // first-byte distribution is uniform; this is mostly a
1660        // sanity check that we route to the right bucket and aren't
1661        // accidentally collapsing.)
1662        let tmp = tempfile::TempDir::new().unwrap();
1663        let bucket_dir = tmp.path().join("buckets");
1664        let idx_path = tmp.path().join("test.idx");
1665        let cursor = Cursor::new(Vec::<u8>::new());
1666        let mut b = StreamingPackBuilder::new(
1667            cursor,
1668            idx_path.clone(),
1669            CompressionConfig::default(),
1670            bucket_dir.clone(),
1671        )
1672        .unwrap();
1673        for i in 0..1024u32 {
1674            let h = blake3::hash(&i.to_le_bytes());
1675            let hash = ContentHash::from_bytes(*h.as_bytes());
1676            b.add(hash, ObjectType::Blob, b"x".to_vec()).unwrap();
1677        }
1678        // Inspect bucket file sizes BEFORE finalize (which deletes them).
1679        b.pack_writer.as_mut().unwrap().flush().unwrap();
1680        let mut max_entries = 0usize;
1681        // Bucket files use the temporary tagged ID encoding; only the
1682        // finalized index folds the tag into the offset.
1683        let entry_size = 33 + 8;
1684        for path in b.bucket_paths.iter() {
1685            if path.exists() {
1686                let size = std::fs::metadata(path).unwrap().len() as usize;
1687                let entries = size / entry_size;
1688                if entries > max_entries {
1689                    max_entries = entries;
1690                }
1691            }
1692        }
1693        // Average is 1024 / 256 = 4 entries per bucket. Allow up to 16
1694        // (4× average) — uniformity isn't perfect on small samples.
1695        assert!(
1696            max_entries <= 16,
1697            "max bucket has {max_entries} entries; uniform expected ~4"
1698        );
1699        let _ = finalize_cursor(b, &idx_path);
1700    }
1701
1702    #[test]
1703    fn finalize_returns_correct_stats() {
1704        let tmp = tempfile::TempDir::new().unwrap();
1705        let (mut b, _, idx_path) = fresh_builder(&tmp);
1706        let payload = vec![0xabu8; 1024];
1707        for i in 0..5u8 {
1708            b.add(deterministic_hash(i), ObjectType::Blob, payload.clone())
1709                .unwrap();
1710        }
1711        let (_, _, stats) = finalize_cursor(b, &idx_path);
1712        assert_eq!(stats.object_count, 5);
1713        assert_eq!(stats.total_uncompressed, 5 * 1024);
1714        assert!(stats.total_compressed > 0);
1715        assert!(stats.compression_ratio > 0.0);
1716        assert_eq!(stats.delta_count, 0, "streaming builder never deltas");
1717    }
1718
1719    #[cfg(feature = "zstd")]
1720    #[test]
1721    fn streaming_compression_roundtrips_through_zstd_frame() {
1722        // Force the streaming path with a payload that compresses
1723        // well (long runs of identical bytes). Verifies:
1724        //  1. Streaming output decodes back to the original bytes.
1725        //  2. The compressed body is genuinely smaller than the
1726        //     uncompressed input (proving zstd ran), and
1727        //  3. The non-canonical 10-byte varint patched into the
1728        //     compressed_size slot decodes to the right value.
1729        let tmp = tempfile::TempDir::new().unwrap();
1730        let (mut b, _, idx_path) = fresh_builder(&tmp);
1731        // 64 KiB of zeros — compresses to a tiny zstd frame, well
1732        // above the default `min_size` so we hit the streaming branch.
1733        let payload = vec![0u8; 64 * 1024];
1734        let hash = deterministic_hash(0x77);
1735        b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1736        let (pack_data, index_data, stats) = finalize_cursor(b, &idx_path);
1737        assert!(
1738            stats.total_compressed < stats.total_uncompressed,
1739            "expected compression ratio < 1.0, got {}/{}",
1740            stats.total_compressed,
1741            stats.total_uncompressed
1742        );
1743        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1744        let (_ty, got) = reader
1745            .get_object(&PackObjectId::Hash(hash))
1746            .unwrap()
1747            .unwrap();
1748        assert_eq!(got, payload);
1749    }
1750
1751    #[cfg(feature = "zstd")]
1752    #[test]
1753    fn equal_length_zstd_frame_is_stored_as_raw_payload() {
1754        // The pack format infers compression from stored length != logical
1755        // length. Find a stable boundary payload whose zstd frame is exactly
1756        // as long as its input, then pin the streaming writer's raw fallback.
1757        fn encode_like_streaming_writer(data: &[u8]) -> Vec<u8> {
1758            let mut compressed = Vec::new();
1759            let mut encoder = zstd::stream::write::Encoder::new(&mut compressed, 3).unwrap();
1760            encoder
1761                .set_pledged_src_size(Some(data.len() as u64))
1762                .unwrap();
1763            std::io::Write::write_all(&mut encoder, data).unwrap();
1764            encoder.finish().unwrap();
1765            compressed
1766        }
1767
1768        let mut random = Vec::with_capacity(1024);
1769        for seed in 0u64..32 {
1770            random.extend_from_slice(blake3::hash(&seed.to_le_bytes()).as_bytes());
1771        }
1772        let payload = (256..=random.len())
1773            .find_map(|logical_len| {
1774                (0..=logical_len).find_map(|zero_prefix| {
1775                    let mut candidate = random[..logical_len].to_vec();
1776                    candidate[..zero_prefix].fill(0);
1777                    let compressed = encode_like_streaming_writer(&candidate);
1778                    (compressed.len() == candidate.len()).then_some(candidate)
1779                })
1780            })
1781            .expect("fixture search must find an equal-length zstd frame");
1782        let compressed = encode_like_streaming_writer(&payload);
1783        assert_eq!(compressed.len(), payload.len());
1784        assert_eq!(compressed.first(), Some(&0x28), "fixture must be zstd");
1785
1786        let tmp = tempfile::TempDir::new().unwrap();
1787        let (mut builder, _, index_path) = fresh_builder(&tmp);
1788        let id = PackObjectId::Hash(deterministic_hash(0x78));
1789        builder
1790            .add_id(id, ObjectType::StateAttachment, payload.clone())
1791            .unwrap();
1792        let (pack, index, stats) = finalize_cursor(builder, &index_path);
1793
1794        assert_eq!(stats.total_compressed, stats.total_uncompressed);
1795        let reader = PackReader::from_bytes(pack, index, &std::env::temp_dir()).unwrap();
1796        assert_eq!(
1797            reader.get_object(&id).unwrap(),
1798            Some((ObjectType::StateAttachment, payload))
1799        );
1800    }
1801
1802    #[cfg(feature = "zstd")]
1803    #[test]
1804    fn padded_varint_decodes_to_original_value_for_canonical_decoder() {
1805        // Sanity for the seek-back scheme: for every value we'd want
1806        // to encode (small, mid, large), confirm the existing
1807        // `decode_varint` returns the same `value` from a 10-byte
1808        // padded encoding. If this ever fails the streaming path's
1809        // patched compressed_size would be misread by readers.
1810        let cases: &[u64] = &[0, 1, 127, 128, 4096, 1_000_000, 1_000_000_000_000, u64::MAX];
1811        for &value in cases {
1812            let mut buf = [0u8; 10];
1813            super::encode_varint_padded_to_10(value, &mut buf);
1814            let (decoded, consumed) = super::super::varint::decode_varint(&buf)
1815                .expect("padded varint should always decode");
1816            assert_eq!(decoded, value, "varint roundtrip failed for {value}");
1817            assert_eq!(
1818                consumed, 10,
1819                "padded encoding should consume all 10 bytes for {value}"
1820            );
1821        }
1822    }
1823
1824    #[cfg(feature = "zstd")]
1825    #[test]
1826    fn streaming_path_does_not_buffer_compressed_payload_in_memory() {
1827        // Smoke check: write a single 8 MiB payload, observe the
1828        // pack file size on disk during/after the add. The pack file
1829        // grows incrementally during the streaming compression — if
1830        // we were buffering an intermediate compressed `Vec<u8>` the
1831        // on-disk size would jump by ~8 MiB at finalize, not stay
1832        // bounded as the encoder pumps bytes through.
1833        //
1834        // We can't easily measure peak heap from inside Rust without
1835        // a custom allocator. What we *can* verify is that calling
1836        // `add` returns control with the pack file already at its
1837        // final body size, demonstrating the encoder wrote through
1838        // and didn't accumulate.
1839        let tmp = tempfile::TempDir::new().unwrap();
1840        let bucket_dir = tmp.path().join("buckets");
1841        let pack_path = tmp.path().join("pack.dat");
1842        let idx_path = tmp.path().join("pack.idx");
1843        let file = std::fs::OpenOptions::new()
1844            .read(true)
1845            .write(true)
1846            .create(true)
1847            .truncate(true)
1848            .open(&pack_path)
1849            .unwrap();
1850        let mut b = StreamingPackBuilder::new(
1851            file,
1852            idx_path.clone(),
1853            CompressionConfig::default(),
1854            bucket_dir,
1855        )
1856        .unwrap();
1857        let payload = vec![0xa5u8; 8 * 1024 * 1024];
1858        let hash = deterministic_hash(0x66);
1859        b.add(hash, ObjectType::Blob, payload.clone()).unwrap();
1860        // Pack file already on disk holds at least the entry header +
1861        // compressed payload (excluding the 32-byte trailer the builder
1862        // appends at finalize).
1863        let mid_size = std::fs::metadata(&pack_path).unwrap().len();
1864        assert!(
1865            mid_size > 16 + 40,
1866            "pack file should hold real entry data after add; size={mid_size}"
1867        );
1868        let (_, _) = b.finalize().unwrap();
1869        let pack_bytes = std::fs::read(&pack_path).unwrap();
1870        let index_bytes = std::fs::read(&idx_path).unwrap();
1871        let reader =
1872            PackReader::from_bytes(pack_bytes, index_bytes, &std::env::temp_dir()).unwrap();
1873        let (_ty, got) = reader
1874            .get_object(&PackObjectId::Hash(hash))
1875            .unwrap()
1876            .unwrap();
1877        assert_eq!(got, payload);
1878    }
1879
1880    #[test]
1881    fn list_ids_returns_all_added_ids_sorted() {
1882        let tmp = tempfile::TempDir::new().unwrap();
1883        let (mut b, _, idx_path) = fresh_builder(&tmp);
1884        let mut added: Vec<PackObjectId> = Vec::new();
1885        // Mix of Hash and StateId in a non-sorted order on input.
1886        for seed in [0x05u8, 0xa0, 0x12, 0x9f, 0x33] {
1887            let id = PackObjectId::Hash(deterministic_hash(seed));
1888            b.add_id(id, ObjectType::Blob, vec![seed; 4]).unwrap();
1889            added.push(id);
1890        }
1891        for seed in [0x80u8, 0x10, 0xff] {
1892            let id = PackObjectId::StateId(deterministic_state_id(seed));
1893            b.add_id(id, ObjectType::State, vec![seed; 4]).unwrap();
1894            added.push(id);
1895        }
1896        let (pack_data, index_data, _) = finalize_cursor(b, &idx_path);
1897        let reader = PackReader::from_bytes(pack_data, index_data, &std::env::temp_dir()).unwrap();
1898        let mut got = reader.list_ids().unwrap();
1899        // PackReader's list_ids returns index order — should already be
1900        // sorted because we sort on finalize.
1901        let mut sorted = got.clone();
1902        sorted.sort();
1903        assert_eq!(got, sorted, "list_ids must come back sorted");
1904        // And every added id should appear.
1905        added.sort();
1906        got.sort();
1907        assert_eq!(got, added);
1908    }
1909}