Skip to main content

mkit_core/
chunker.rs

1//! `FastCDC` content-defined chunker.
2//!
3//! Spec reference: `docs/specs/SPEC-FASTCDC.md`. Frozen v1 parameters:
4//!
5//! * Gear seed: ASCII `"MKITFCDC"` interpreted as a big-endian `u64`
6//!   (`0x4D4B_4954_4643_4443`). Splitmix64-derived 256-entry table.
7//! * `MIN_SIZE = 16 KiB`, `AVG_SIZE = 64 KiB`, `MAX_SIZE = 256 KiB`.
8//! * `MASK_S = 0x0001_FFFF` (strict, used while `i < AVG_SIZE`).
9//! * `MASK   = 0x0000_FFFF` (informative; not used directly — we only
10//!   need the strict and loose masks at runtime).
11//! * `MASK_L = 0x0000_7FFF` (loose, used while `AVG_SIZE <= i < MAX_SIZE`).
12//! * Rolling hash: `h = (h << 1) +% gear[byte]`. Cut when `(h & mask) == 0`,
13//!   else force a cut at `MAX_SIZE`. Never cut before `MIN_SIZE`.
14//!
15//! The spec hard-codes the masks; we pin them as constants here so any
16//! accidental change lights up in the golden tests immediately.
17//!
18//! All chunk boundaries are fully determined by the input bytes plus the
19//! constants above. Any change breaks `chunked_blob` reproducibility
20//! (see `SPEC-FASTCDC.md` §2 determinism contract).
21
22/// Frozen splitmix64 seed for v1 — ASCII `"MKITFCDC"` as big-endian u64.
23pub const SEED: u64 = 0x4D4B_4954_4643_4443;
24
25/// Minimum chunk size (`16 KiB`). `cut` will not return a value below
26/// this for non-final chunks.
27pub const MIN_SIZE: usize = 0x4000;
28/// Average target chunk size (`64 KiB`). The mask transition from strict
29/// to loose happens here.
30pub const AVG_SIZE: usize = 0x10000;
31/// Maximum chunk size (`256 KiB`). `cut` always returns at most this.
32pub const MAX_SIZE: usize = 0x40000;
33
34/// Strict mask (used while `i < AVG_SIZE`). Fewer cuts → bias the chunker
35/// toward `AVG_SIZE`.
36pub const MASK_S: u64 = 0x0001_FFFF;
37/// Loose mask (used while `AVG_SIZE <= i < MAX_SIZE`). More cuts → avoid
38/// running into the `MAX_SIZE` forced boundary.
39pub const MASK_L: u64 = 0x0000_7FFF;
40
41/// 256-entry gear table, derived once at first use from [`SEED`] via
42/// splitmix64. Wrapping arithmetic throughout — see SPEC-FASTCDC §3 for
43/// the exact derivation.
44fn gear_table() -> &'static [u64; 256] {
45    use std::sync::OnceLock;
46    static TABLE: OnceLock<[u64; 256]> = OnceLock::new();
47    TABLE.get_or_init(build_gear_table)
48}
49
50fn build_gear_table() -> [u64; 256] {
51    let mut state: u64 = SEED;
52    let mut table = [0u64; 256];
53    for entry in &mut table {
54        // splitmix64 with the standard gamma. All ops are wrapping.
55        state = state.wrapping_add(0x9e37_79b9_7f4a_7c15);
56        let mut z = state;
57        z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
58        z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
59        z ^= z >> 31;
60        *entry = z;
61    }
62    table
63}
64
65/// One chunk's position in the source buffer.
66#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub struct ChunkBoundary {
68    pub offset: usize,
69    pub length: usize,
70}
71
72/// `FastCDC` chunker parameters. Construct with [`FastCdc::v1`] for the
73/// frozen v1 constants; other constructors exist mainly for tests that
74/// exercise the algorithm at smaller sizes.
75#[derive(Debug, Clone, Copy)]
76pub struct FastCdc {
77    min_size: usize,
78    avg_size: usize,
79    max_size: usize,
80    mask_s: u64,
81    mask_l: u64,
82}
83
84impl FastCdc {
85    /// Construct the v1 chunker (`16 KiB / 64 KiB / 256 KiB`,
86    /// strict/loose `0x0001_FFFF` / `0x0000_7FFF`).
87    #[must_use]
88    pub const fn v1() -> Self {
89        Self {
90            min_size: MIN_SIZE,
91            avg_size: AVG_SIZE,
92            max_size: MAX_SIZE,
93            mask_s: MASK_S,
94            mask_l: MASK_L,
95        }
96    }
97
98    /// Construct a chunker with custom parameters. `avg_size` MUST be a
99    /// power of two; the strict/loose masks are derived as
100    /// `mask = (1 << log2(avg)) - 1; mask_s = mask | (mask << 1); mask_l = mask >> 1`.
101    /// Returns [`crate::object::MkitError::InvalidIdentity`] re-purposed as a generic
102    /// "bad parameter" error if the constraints are violated — but in
103    /// practice this constructor is only used by tests, where the inputs
104    /// are constants, so we panic instead for a clearer failure mode.
105    ///
106    /// # Panics
107    ///
108    /// Panics if `min < avg < max` is not strictly satisfied or `avg` is
109    /// not a power of two.
110    #[must_use]
111    pub fn custom(min: usize, avg: usize, max: usize) -> Self {
112        assert!(min < avg && avg < max, "FastCdc: require min<avg<max");
113        assert!(avg.is_power_of_two(), "FastCdc: avg must be a power of 2");
114        let bits = avg.trailing_zeros();
115        let mask: u64 = (1u64 << bits) - 1;
116        Self {
117            min_size: min,
118            avg_size: avg,
119            max_size: max,
120            mask_s: mask | (mask << 1),
121            mask_l: mask >> 1,
122        }
123    }
124
125    #[must_use]
126    pub const fn min_size(&self) -> usize {
127        self.min_size
128    }
129    #[must_use]
130    pub const fn avg_size(&self) -> usize {
131        self.avg_size
132    }
133    #[must_use]
134    pub const fn max_size(&self) -> usize {
135        self.max_size
136    }
137
138    /// Find the cut point in `data`. Returns the length of the first
139    /// chunk. Semantics:
140    ///
141    /// * `data.len() <= min_size` → returns `data.len()` (no early cut).
142    /// * `min_size < i <= avg_size` → strict mask, fewer boundaries.
143    /// * `avg_size < i <= max_size` → loose mask, more boundaries.
144    /// * Otherwise → forced cut at `min(max_size, data.len())`.
145    #[must_use]
146    pub fn cut(&self, data: &[u8]) -> usize {
147        if data.len() <= self.min_size {
148            return data.len();
149        }
150        let table = gear_table();
151        let mut hash: u64 = 0;
152        let avg_end = self.avg_size.min(data.len());
153        let mut i = self.min_size;
154        while i < avg_end {
155            hash = (hash << 1).wrapping_add(table[data[i] as usize]);
156            if (hash & self.mask_s) == 0 {
157                return i;
158            }
159            i += 1;
160        }
161        let max_end = self.max_size.min(data.len());
162        while i < max_end {
163            hash = (hash << 1).wrapping_add(table[data[i] as usize]);
164            if (hash & self.mask_l) == 0 {
165                return i;
166            }
167            i += 1;
168        }
169        max_end
170    }
171}
172
173/// Iterator over chunk boundaries in a contiguous byte slice.
174#[derive(Debug)]
175pub struct ChunkIterator<'a> {
176    cdc: FastCdc,
177    data: &'a [u8],
178    offset: usize,
179}
180
181impl<'a> ChunkIterator<'a> {
182    #[must_use]
183    pub fn new(cdc: FastCdc, data: &'a [u8]) -> Self {
184        Self {
185            cdc,
186            data,
187            offset: 0,
188        }
189    }
190}
191
192impl Iterator for ChunkIterator<'_> {
193    type Item = ChunkBoundary;
194    fn next(&mut self) -> Option<Self::Item> {
195        if self.offset >= self.data.len() {
196            return None;
197        }
198        let remaining = &self.data[self.offset..];
199        let length = self.cdc.cut(remaining);
200        let boundary = ChunkBoundary {
201            offset: self.offset,
202            length,
203        };
204        self.offset += length;
205        Some(boundary)
206    }
207}
208
209/// Streaming variant of [`ChunkIterator`] that pulls from a [`Read`]
210/// instead of requiring the whole input as an in-memory slice.
211///
212/// [`FastCdc::cut`] only ever looks at most `max_size` bytes ahead of the
213/// current chunk's start (see its doc comment), so a *correct* reader
214/// only strictly needs one `max_size`-sized window resident — but memory
215/// use is bounded, independent of the total input length, either way.
216/// This reader keeps `WINDOW_MULTIPLE * max_size` resident instead of the
217/// bare minimum: buffering further ahead means `fill`'s `read` calls can
218/// pull in more than `max_size` bytes at once when the source allows it
219/// (a `File` or `Cursor` read is not obligated to stop at `max_size`),
220/// amortizing the compaction below across several `max_size` windows
221/// instead of paying it once per chunk. It produces byte-identical chunk
222/// boundaries and content to [`ChunkIterator`] run over the same bytes
223/// regardless of how far ahead it buffers, since [`FastCdc::cut`] itself
224/// clamps its scan to `max_size` no matter how much of the slice handed
225/// to it is valid (pinned by the `streaming_matches_in_memory_*` tests
226/// below).
227///
228/// The window lives in one `WINDOW_MULTIPLE * max_size` buffer allocated
229/// once at construction. `[start, end)` marks the currently valid bytes;
230/// `next_chunk_ref` advances `start` past each cut chunk in place instead
231/// of `split_off`-ing a new `Vec` (an allocation plus a full memcpy of the
232/// remainder) every call, and `fill` reads straight into the buffer's
233/// unused tail instead of through a separate scratch buffer, so a chunk
234/// costs at most one compacting `copy_within` amortized over several
235/// `max_size` windows rather than two copies per chunk.
236///
237/// `WINDOW_MULTIPLE`'s only correctness requirement is `>= 2` (`fill`'s
238/// compaction condition needs at least one full `max_size` of spare tail
239/// room to guarantee progress without compacting on every read); `4` was
240/// chosen so compaction amortizes over 3 windows' worth of consumption
241/// (`(WINDOW_MULTIPLE - 1) * max_size`) — see the `chunker_streaming`
242/// bench for the throughput this buys. For `FastCdc::v1()` this makes the
243/// buffer exactly 1 MiB, the same number as (but not derived from, and
244/// not required to track) `worktree::CHUNK_THRESHOLD` — that constant is
245/// an unrelated policy choice (which files stream through this reader at
246/// all), not a bound this reader depends on.
247const WINDOW_MULTIPLE: usize = 4;
248
249#[derive(Debug)]
250pub struct ChunkReader<R> {
251    cdc: FastCdc,
252    reader: R,
253    buf: Vec<u8>,
254    start: usize,
255    end: usize,
256    eof: bool,
257}
258
259impl<R: std::io::Read> ChunkReader<R> {
260    #[must_use]
261    pub fn new(cdc: FastCdc, reader: R) -> Self {
262        let window = cdc.max_size() * WINDOW_MULTIPLE;
263        Self {
264            cdc,
265            reader,
266            buf: vec![0u8; window],
267            start: 0,
268            end: 0,
269            eof: false,
270        }
271    }
272
273    /// Top up `[start, end)` to at least `max_size` valid bytes (or EOF),
274    /// reading directly into the buffer's tail. Compacts the valid region
275    /// back to offset `0` first whenever the tail no longer has room for a
276    /// full `max_size` read — with a `WINDOW_MULTIPLE * max_size` buffer
277    /// this happens roughly once every `(WINDOW_MULTIPLE - 1) * max_size`
278    /// bytes consumed, not once per chunk. `read` can return short of the
279    /// requested length without signalling EOF, so this loops until either
280    /// the target is met or a `0`-byte read confirms EOF.
281    fn fill(&mut self) -> std::io::Result<()> {
282        debug_assert!(self.start <= self.end && self.end <= self.buf.len());
283        let target = self.cdc.max_size();
284        while !self.eof && self.end - self.start < target {
285            if self.buf.len() - self.end < target {
286                self.buf.copy_within(self.start..self.end, 0);
287                self.end -= self.start;
288                self.start = 0;
289            }
290            let n = self.reader.read(&mut self.buf[self.end..])?;
291            if n == 0 {
292                self.eof = true;
293            } else {
294                self.end += n;
295            }
296        }
297        debug_assert!(self.start <= self.end && self.end <= self.buf.len());
298        Ok(())
299    }
300
301    /// Read the next chunk as a borrow into the internal window, or
302    /// `None` at end of input. The returned slice is only valid until the
303    /// next call — [`next_chunk`](Self::next_chunk) copies it out for
304    /// callers that need an owned `Vec`.
305    ///
306    /// # Errors
307    /// Propagates the underlying reader's I/O errors.
308    pub fn next_chunk_ref(&mut self) -> std::io::Result<Option<&[u8]>> {
309        self.fill()?;
310        if self.start == self.end {
311            return Ok(None);
312        }
313        let length = self.cdc.cut(&self.buf[self.start..self.end]);
314        debug_assert!(length <= self.end - self.start);
315        let chunk_start = self.start;
316        self.start += length;
317        Ok(Some(&self.buf[chunk_start..chunk_start + length]))
318    }
319
320    /// Read and return the next chunk's owned bytes, or `None` at end of
321    /// input. Mirrors `Iterator` exhaustion semantics but is fallible
322    /// because filling the window can fail.
323    ///
324    /// # Errors
325    /// Propagates the underlying reader's I/O errors.
326    pub fn next_chunk(&mut self) -> std::io::Result<Option<Vec<u8>>> {
327        Ok(self.next_chunk_ref()?.map(<[u8]>::to_vec))
328    }
329}
330
331/// Convenience: collect all chunk *end* offsets from `data` using the
332/// frozen v1 chunker. The returned vector starts with the length of the
333/// first chunk and ends with `data.len()`. An empty input yields an
334/// empty vector. Used by goldens to pin boundaries deterministically
335/// without exposing the `ChunkBoundary` struct in serialised form.
336#[must_use]
337pub fn chunk_boundaries(data: &[u8]) -> Vec<usize> {
338    let cdc = FastCdc::v1();
339    let mut out = Vec::new();
340    let mut offset = 0usize;
341    while offset < data.len() {
342        let len = cdc.cut(&data[offset..]);
343        offset += len;
344        out.push(offset);
345    }
346    out
347}
348
349/// Hash the entire 256-entry gear table as little-endian `u64` bytes
350/// (2 048 bytes total) with BLAKE3. Useful as a cheap "did we get the
351/// seed right" check — see SPEC-FASTCDC §8 vector 1.
352#[must_use]
353pub fn gear_table_digest() -> [u8; 32] {
354    let table = gear_table();
355    let mut bytes = [0u8; 256 * 8];
356    for (i, v) in table.iter().enumerate() {
357        bytes[i * 8..i * 8 + 8].copy_from_slice(&v.to_le_bytes());
358    }
359    crate::hash::hash(&bytes)
360}
361
362// =========================================================================
363// Tests
364// =========================================================================
365
366#[cfg(test)]
367mod tests {
368    use super::*;
369
370    /// Deterministic xorshift64* PRNG for test-data generation. We avoid
371    /// `rand` here to keep test inputs explicit and reproducible
372    /// without committing to a specific `rand` version.
373    struct Prng(u64);
374    impl Prng {
375        fn new(seed: u64) -> Self {
376            Self(seed.max(1))
377        }
378        fn next_u64(&mut self) -> u64 {
379            // splitmix64 (same primitive used to derive the gear table —
380            // a different seed gives an unrelated stream).
381            self.0 = self.0.wrapping_add(0x9e37_79b9_7f4a_7c15);
382            let mut z = self.0;
383            z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
384            z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
385            z ^ (z >> 31)
386        }
387        fn fill(&mut self, dst: &mut [u8]) {
388            for chunk in dst.chunks_mut(8) {
389                let bytes = self.next_u64().to_le_bytes();
390                let n = chunk.len();
391                chunk.copy_from_slice(&bytes[..n]);
392            }
393        }
394    }
395
396    #[test]
397    fn gear_table_is_unique_and_nonzero() {
398        let table = gear_table();
399        let mut seen = std::collections::HashSet::new();
400        for &v in table {
401            assert_ne!(v, 0, "gear table entry is zero");
402            assert!(seen.insert(v), "duplicate gear table entry");
403        }
404        assert_eq!(seen.len(), 256);
405    }
406
407    #[test]
408    fn min_max_size_constraints() {
409        let cdc = FastCdc::v1();
410        let mut data = vec![0u8; 512 * 1024];
411        Prng::new(0xCAFE_BABE).fill(&mut data);
412
413        let mut total = 0usize;
414        let mut count = 0usize;
415        for b in ChunkIterator::new(cdc, &data) {
416            assert!(b.length <= MAX_SIZE);
417            // Only the final chunk is allowed to be smaller than min.
418            if b.offset + b.length < data.len() {
419                assert!(b.length >= MIN_SIZE, "chunk under MIN_SIZE: {}", b.length);
420            }
421            total += b.length;
422            count += 1;
423        }
424        assert_eq!(total, data.len());
425        assert!(count > 1);
426    }
427
428    #[test]
429    fn small_input_is_single_chunk() {
430        let cdc = FastCdc::v1();
431        let small = b"hello, this is a tiny file";
432        assert_eq!(cdc.cut(small), small.len());
433        let boundaries: Vec<_> = ChunkIterator::new(cdc, small).collect();
434        assert_eq!(boundaries.len(), 1);
435        assert_eq!(boundaries[0].offset, 0);
436        assert_eq!(boundaries[0].length, small.len());
437    }
438
439    #[test]
440    fn empty_input_iterator_is_empty() {
441        let cdc = FastCdc::v1();
442        assert_eq!(cdc.cut(b""), 0);
443        let none: Vec<_> = ChunkIterator::new(cdc, b"").collect();
444        assert!(none.is_empty());
445    }
446
447    #[test]
448    fn cut_at_exactly_min_size_returns_full() {
449        let cdc = FastCdc::custom(1024, 4096, 16384);
450        let mut data = vec![0u8; 1024];
451        Prng::new(99).fill(&mut data);
452        assert_eq!(cdc.cut(&data), data.len());
453    }
454
455    #[test]
456    fn cut_forces_boundary_at_max_size() {
457        // All-zero buffer, gear hash never trips a natural cut → forced
458        // cut at max_size.
459        let cdc = FastCdc::custom(4, 8, 16);
460        let data = [0u8; 64];
461        let len = cdc.cut(&data);
462        assert!(len <= 16, "cut returned {len} > max=16");
463    }
464
465    #[test]
466    fn boundary_stability_single_byte_insert() {
467        // Insert one byte at offset 32 KiB; expect <=3 chunks differ.
468        let cdc = FastCdc::v1();
469        let mut original = vec![0u8; 64 * 1024];
470        Prng::new(0xBEEF).fill(&mut original);
471        let insert_point = 32 * 1024;
472        let mut modified = Vec::with_capacity(original.len() + 1);
473        modified.extend_from_slice(&original[..insert_point]);
474        modified.push(0xFF);
475        modified.extend_from_slice(&original[insert_point..]);
476
477        let orig_chunks: Vec<_> = ChunkIterator::new(cdc, &original).collect();
478        let mod_chunks: Vec<_> = ChunkIterator::new(cdc, &modified).collect();
479
480        let max_chunks = orig_chunks.len().max(mod_chunks.len());
481        let mut differing = 0usize;
482        for i in 0..max_chunks {
483            match (orig_chunks.get(i), mod_chunks.get(i)) {
484                (Some(o), Some(m)) => {
485                    let os = &original[o.offset..o.offset + o.length];
486                    let ms = &modified[m.offset..m.offset + m.length];
487                    if os != ms {
488                        differing += 1;
489                    }
490                }
491                _ => differing += 1,
492            }
493        }
494        assert!(
495            differing <= 3,
496            "expected <=3 differing chunks, got {differing}"
497        );
498    }
499
500    #[test]
501    fn different_avg_size_yields_different_boundaries() {
502        let mut data = vec![0u8; 100 * 1024];
503        Prng::new(0xABCD).fill(&mut data);
504        let small = FastCdc::custom(8 * 1024, 32 * 1024, 128 * 1024);
505        let large = FastCdc::v1();
506
507        let s: Vec<_> = ChunkIterator::new(small, &data)
508            .map(|b| b.offset + b.length)
509            .collect();
510        let l: Vec<_> = ChunkIterator::new(large, &data)
511            .map(|b| b.offset + b.length)
512            .collect();
513        assert!(
514            s != l,
515            "expected different boundaries for different avg_size"
516        );
517    }
518
519    #[test]
520    fn chunk_boundaries_helper_matches_iterator() {
521        let mut data = vec![0u8; 200 * 1024];
522        Prng::new(0xFEED_FACE).fill(&mut data);
523        let from_helper = chunk_boundaries(&data);
524        let from_iter: Vec<usize> = ChunkIterator::new(FastCdc::v1(), &data)
525            .map(|b| b.offset + b.length)
526            .collect();
527        assert_eq!(from_helper, from_iter);
528    }
529
530    #[test]
531    fn streaming_matches_in_memory_iterator() {
532        let mut data = vec![0u8; 500 * 1024];
533        Prng::new(0x5713_EA31).fill(&mut data);
534
535        let in_memory: Vec<Vec<u8>> = ChunkIterator::new(FastCdc::v1(), &data)
536            .map(|b| data[b.offset..b.offset + b.length].to_vec())
537            .collect();
538
539        let mut reader = ChunkReader::new(FastCdc::v1(), std::io::Cursor::new(&data));
540        let mut streamed = Vec::new();
541        while let Some(chunk) = reader.next_chunk().unwrap() {
542            streamed.push(chunk);
543        }
544
545        assert_eq!(in_memory, streamed);
546    }
547
548    #[test]
549    fn streaming_matches_in_memory_iterator_across_multiple_compactions() {
550        // `streaming_matches_in_memory_iterator` above only pushes 500 KiB
551        // through the reader — under the `WINDOW_MULTIPLE * MAX_SIZE`
552        // (1 MiB) window, so `fill`'s `copy_within` compaction path never
553        // actually runs there. Push several MiB through instead, so `fill`
554        // is forced to compact more than once, and check *content*
555        // (both the owned `next_chunk` and the borrowed `next_chunk_ref`
556        // paths) against `ChunkIterator`, not just that the window stays
557        // bounded (`streaming_reader_buffer_stays_bounded_...` below only
558        // checks buffer size, not chunk correctness). `Cursor` — unlike
559        // `StingyReader` below — happily satisfies a `read` request in one
560        // large memcpy, so this also exercises `fill` buffering far past
561        // `max_size` in a single read before the next `cut`.
562        let mut data = vec![0u8; 6 * 1024 * 1024];
563        Prng::new(0xC0FF_EE01).fill(&mut data);
564
565        let in_memory: Vec<Vec<u8>> = ChunkIterator::new(FastCdc::v1(), &data)
566            .map(|b| data[b.offset..b.offset + b.length].to_vec())
567            .collect();
568        assert!(
569            in_memory.len() > 50,
570            "fixture must be large enough to force several compactions, got {} chunks",
571            in_memory.len()
572        );
573
574        let mut owned_reader = ChunkReader::new(FastCdc::v1(), std::io::Cursor::new(&data));
575        let mut owned = Vec::new();
576        while let Some(chunk) = owned_reader.next_chunk().unwrap() {
577            owned.push(chunk);
578        }
579        assert_eq!(
580            in_memory, owned,
581            "owned next_chunk() diverged from ChunkIterator"
582        );
583
584        let mut ref_reader = ChunkReader::new(FastCdc::v1(), std::io::Cursor::new(&data));
585        let mut borrowed = Vec::new();
586        while let Some(chunk) = ref_reader.next_chunk_ref().unwrap() {
587            borrowed.push(chunk.to_vec());
588        }
589        assert_eq!(
590            in_memory, borrowed,
591            "borrowed next_chunk_ref() diverged from ChunkIterator"
592        );
593    }
594
595    #[test]
596    fn streaming_empty_input_yields_no_chunks() {
597        let mut reader = ChunkReader::new(FastCdc::v1(), std::io::Cursor::new(&[] as &[u8]));
598        assert!(reader.next_chunk().unwrap().is_none());
599    }
600
601    #[test]
602    fn streaming_small_input_is_single_chunk() {
603        let small = b"hello, this is a tiny file";
604        let mut reader = ChunkReader::new(FastCdc::v1(), std::io::Cursor::new(&small[..]));
605        let chunk = reader.next_chunk().unwrap().expect("one chunk");
606        assert_eq!(chunk, small);
607        assert!(reader.next_chunk().unwrap().is_none());
608    }
609
610    /// A reader that only ever yields a handful of bytes per `read` call,
611    /// to exercise `ChunkReader::fill`'s short-read loop (a real `File`
612    /// on a slow filesystem or a network stream behaves this way).
613    struct StingyReader<'a> {
614        data: &'a [u8],
615        pos: usize,
616    }
617    impl std::io::Read for StingyReader<'_> {
618        fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
619            let n = (self.data.len() - self.pos).min(buf.len()).min(7);
620            buf[..n].copy_from_slice(&self.data[self.pos..self.pos + n]);
621            self.pos += n;
622            Ok(n)
623        }
624    }
625
626    #[test]
627    fn streaming_reader_buffer_stays_bounded_regardless_of_input_size() {
628        // The whole point of issue #828: ingest memory must not scale
629        // with file size. Run several MiB through the reader and assert
630        // its internal window is a fixed `WINDOW_MULTIPLE * max_size`
631        // allocation that never grows or reallocates — bounded by the
632        // chunker's own parameters, not by how much input remains. (The
633        // valid region within it can transiently exceed `max_size`: a
634        // single `read` is allowed to fill the whole spare tail in one
635        // syscall, which is the point of buffering ahead — `cut` only
636        // ever consults the first `max_size` bytes of it.)
637        let mut data = vec![0u8; 8 * 1024 * 1024];
638        Prng::new(0x8000_0008).fill(&mut data);
639        let mut reader = ChunkReader::new(FastCdc::v1(), std::io::Cursor::new(&data));
640        let window = WINDOW_MULTIPLE * MAX_SIZE;
641        let mut chunk_count = 0usize;
642        while let Some(_chunk) = reader.next_chunk().unwrap() {
643            assert_eq!(
644                reader.buf.len(),
645                window,
646                "internal window buffer must not grow past its fixed allocation"
647            );
648            assert!(
649                reader.end - reader.start <= window,
650                "valid region grew to {} bytes, expected <= {}",
651                reader.end - reader.start,
652                window
653            );
654            chunk_count += 1;
655        }
656        assert!(chunk_count > 1, "expected multiple chunks from 8 MiB input");
657    }
658
659    #[test]
660    fn streaming_matches_in_memory_iterator_under_short_reads() {
661        let mut data = vec![0u8; 300 * 1024];
662        Prng::new(0x5713_EA32).fill(&mut data);
663
664        let in_memory: Vec<Vec<u8>> = ChunkIterator::new(FastCdc::v1(), &data)
665            .map(|b| data[b.offset..b.offset + b.length].to_vec())
666            .collect();
667
668        let mut reader = ChunkReader::new(
669            FastCdc::v1(),
670            StingyReader {
671                data: &data,
672                pos: 0,
673            },
674        );
675        let mut streamed = Vec::new();
676        while let Some(chunk) = reader.next_chunk().unwrap() {
677            streamed.push(chunk);
678        }
679
680        assert_eq!(in_memory, streamed);
681    }
682
683    /// Pinned v1 gear-table digest, harvested once from the splitmix64
684    /// derivation seeded with "MKITFCDC". Drift = a v2 break.
685    const EXPECTED_GEAR_DIGEST_HEX: &str =
686        "7b238963a8bb10c4dea1bf678aa07d8c3ce94284209c440ca971ff3a97ee5ad4";
687
688    #[test]
689    fn gear_table_digest_is_stable() {
690        // SPEC-FASTCDC §8 vector 1: any change to the seed or splitmix
691        // derivation moves this digest. CI flags drift loud and early.
692        let hex = crate::hash::to_hex(&gear_table_digest());
693        assert_eq!(
694            hex, EXPECTED_GEAR_DIGEST_HEX,
695            "gear table digest changed; refuse to drift silently"
696        );
697    }
698
699    // -- Property tests -------------------------------------------------
700    //
701    // Determinism + cap invariants exercised against arbitrary inputs
702    // via `proptest`. The example tests above cover specific PRNG seeds;
703    // the properties below catch the boundary cases the examples miss
704    // (very-short data, repeating patterns, etc.).
705    proptest::proptest! {
706        /// FastCDC is deterministic: two passes over the same bytes
707        /// produce the same chunk boundaries. This is the core SPEC-
708        /// FASTCDC §2 contract that makes `chunked_blob` content-
709        /// addressable.
710        #[test]
711        fn proptest_determinism(data in proptest::collection::vec(proptest::num::u8::ANY, 0..256 * 1024)) {
712            let cdc = FastCdc::v1();
713            let pass1: Vec<_> = ChunkIterator::new(cdc, &data).collect();
714            let pass2: Vec<_> = ChunkIterator::new(cdc, &data).collect();
715            proptest::prop_assert_eq!(pass1, pass2);
716        }
717
718        /// Boundaries cover the input exactly: sum of lengths equals
719        /// input length; offsets are non-overlapping and monotonic.
720        #[test]
721        fn proptest_boundaries_partition_input(
722            data in proptest::collection::vec(proptest::num::u8::ANY, 0..256 * 1024),
723        ) {
724            let cdc = FastCdc::v1();
725            let boundaries: Vec<_> = ChunkIterator::new(cdc, &data).collect();
726            let mut expected_offset = 0usize;
727            for b in &boundaries {
728                proptest::prop_assert_eq!(b.offset, expected_offset);
729                expected_offset += b.length;
730            }
731            proptest::prop_assert_eq!(expected_offset, data.len());
732        }
733
734        /// `ChunkReader` (streaming) produces byte-identical chunks to
735        /// `ChunkIterator` (in-memory) for arbitrary input — the
736        /// correctness contract issue #828's ingest streaming depends on:
737        /// content addressing must not change based on how a file was
738        /// read.
739        #[test]
740        fn proptest_streaming_matches_in_memory(
741            data in proptest::collection::vec(proptest::num::u8::ANY, 0..256 * 1024),
742        ) {
743            let cdc = FastCdc::v1();
744            let in_memory: Vec<Vec<u8>> = ChunkIterator::new(cdc, &data)
745                .map(|b| data[b.offset..b.offset + b.length].to_vec())
746                .collect();
747
748            let mut reader = ChunkReader::new(cdc, std::io::Cursor::new(&data));
749            let mut streamed = Vec::new();
750            while let Some(chunk) = reader.next_chunk().unwrap() {
751                streamed.push(chunk);
752            }
753            proptest::prop_assert_eq!(in_memory, streamed);
754        }
755    }
756}