diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7ae2fb6..9d80804 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -22,13 +22,33 @@ jobs: - uses: actions/checkout@v4 - uses: dtolnay/rust-toolchain@stable with: - components: clippy - - name: clippy (lib + tests) - run: cargo clippy --lib --tests -- -D warnings + components: clippy, rustfmt + - name: formatting + run: cargo fmt --all -- --check + - name: clippy (all targets and features) + run: cargo clippy --all-targets --all-features -- -D warnings - name: build (all targets) - run: cargo build --all-targets + run: cargo build --all-targets --all-features - name: test - run: cargo test --release + run: cargo test --release --all-features + + msrv: + name: MSRV (Rust 1.89) + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@1.89 + - run: cargo check --all-targets --all-features + + portable: + name: 32-bit portability + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: dtolnay/rust-toolchain@stable + with: + targets: i686-unknown-linux-gnu + - run: cargo check --target i686-unknown-linux-gnu --lib --all-features # Run the cross-SIMD-width determinism tests under Intel SDE so the AVX-512 # argmin and AVX-512BW packed-scan paths actually execute (and are checked @@ -46,6 +66,6 @@ jobs: # hardware and under SDE 10.8.0 (verified on an Intel Xeon w/ avx512bw). - uses: petarpetrovt/setup-sde@v5.0 - name: Build tests - run: cargo test --release --no-run + run: cargo test --release --all-features --no-run - name: cross-impl determinism under emulated AVX-512 (Sapphire Rapids) - run: ${SDE_PATH}/sde64 -spr -- cargo test --release widths_agree -- --nocapture + run: ${SDE_PATH}/sde64 -spr -- cargo test --release --all-features widths_agree -- --nocapture diff --git a/CHANGELOG.md b/CHANGELOG.md index 3cdfc78..3ddad24 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,8 +1,7 @@ # Changelog -## 0.7.0 — API unification (2026-07-11) +## 0.7.1 — API unification and release hardening (2026-07-12) -Breaking change (0.6.0 has zero dependents, so this ships as a free break). Unifies the public API around the caterpillar layer as the recommended entry point, and separates the inherited MinCDC core from the caterpillar layer at the module level: @@ -20,16 +19,32 @@ layer at the module level: `ReadChunker`) moved into a new `mincdc` submodule and is no longer re-exported at the crate root — reach it as `mothcdc::mincdc::...`. Only `Chunk` stays re-exported at the root (shared by both layers). -- No behavior change to chunk boundary placement, the caterpillar - packed-scanning fast path, or the C API (`mothcdc_next_chunk`'s symbol - and behavior are unchanged). +- `Chunk`/`Segment` offsets, segment lengths, and represented chunk counts are + now `u64`; `Segment` has a private validated representation rather than + publicly constructible invalid variants. +- Reader chunkers retry `Interrupted`, use checked buffer arithmetic, expose + fallible constructors and reader accessors, and report offset/count overflow. +- Streaming caterpillar runs now remain maximal even when a `Read` + implementation returns tiny fragments; a one-byte reader and `Cursor` + produce the same grouping. +- Fixed a panic on x86_64 CPUs without SSE4.1 when the exact scalar argmin was + one of the final three windows. +- The public `Cdc` splitpoint contract is documented and enforced. +- The academic `MinCdc4` unit struct can be constructed directly without using + its deprecated `new()` helper. +- The C API validates size configurations without panicking. A normal Rust + dependency now builds only an `rlib`; the benchmark static library is built + explicitly with `cargo rustc --features capi -- --crate-type staticlib`. +- CI now covers every feature, formatting, the 1.89 MSRV, and a 32-bit target. +- The SIMD prefetch soundness fix was submitted upstream as + [orlp/mincdc#1](https://github.com/orlp/mincdc/pull/1). Also includes the mincatcdc -> mothcdc rename (no functional change). Credits unchanged: MinCDC algorithm (Orson Peters), caterpillar layer inspired by Chonkers (Berger), vector acceleration in the style of VectorCDC (Udayashankar et al.). -## 0.6.0 — optimization campaign (2026-07-11) +## 0.6.0 — optimization campaign and packed scanning (2026-07-10–11) Ten hypothesis-driven loops on a fixed Fly performance-4x machine (AMD EPYC/AVX2), each benched against the same corpus (public Tigris bucket). @@ -65,7 +80,7 @@ AE-Min 55.47%, FastCDC 52.44%, RAM 49.21%) — at 8.3-9.9 GB/s vs VectorCDC-AE-Min's 4.6 GB/s on the same data (only VectorCDC-RAM is faster at 17.2 GB/s, with the worst dedup of the field). -## 0.6.0 — 2026-07-10 +### Packed-scanning caterpillar fast path Packed-scanning caterpillar fast path (VectorCDC-style SIMD). @@ -93,7 +108,7 @@ Packed-scanning caterpillar fast path (VectorCDC-style SIMD). within ~2%. - Synthetic ceiling: zeros 2.4 → 74.3 GiB/s on Intel Xeon/AVX-512BW, 1.8 → 30.2 GiB/s on NEON. -- New `capi` feature: a minimal C API (`mincatcdc_next_chunk`) plus a +- New `capi` feature: a minimal C API (`mothcdc_next_chunk`) plus a dedup-bench fork (github.com/russellromney/dedup-bench, branch `mincatcdc-integration`) that adds `chunking_algo=mincdc`, so mincatcdc is measured by the *same* diff --git a/Cargo.toml b/Cargo.toml index eb22aa9..480862b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "mothcdc" -version = "0.7.0" -authors = ["Orson Peters ", "Russell Romney"] +version = "0.7.1" +authors = ["Russell Romney"] license = "Zlib" repository = "https://github.com/russellromney/mothcdc" readme = "README.md" @@ -14,11 +14,6 @@ edition = "2024" rust-version = "1.89" exclude = ["assets", "tools"] -[lib] -# staticlib is consumed by the dedup-bench integration (feature `capi`); -# rlib is the normal Rust library. -crate-type = ["rlib", "staticlib"] - [features] # Minimal C API for embedding in external harnesses (see src/capi.rs). capi = [] diff --git a/README.md b/README.md index 3e2d65d..aeb3f28 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ deduplicated. To start using `mothcdc` add the following to your `Cargo.toml`: [dependencies] - mothcdc = "0.6" + mothcdc = "0.7" Please refer to [the documentation](https://docs.rs/mothcdc) for more information on usage. @@ -44,6 +44,13 @@ The crate also ships a C API (`--features capi`) and a [dedup-bench fork](https://github.com/russellromney/dedup-bench/tree/MothCDC-integration) so it can be measured by the same harness as every other chunker. +The normal Cargo dependency builds only a Rust `rlib`. To produce the static +library used by C/C++ benchmark harnesses, run: + +```sh +cargo rustc --release --features capi --lib -- --crate-type staticlib +``` + ### Benchmarks All numbers come from [UWASL dedup-bench](https://github.com/UWASL/dedup-bench) @@ -141,6 +148,12 @@ once — at most `max_size` bytes — when it crosses a refill, so even a multi-gigabyte zero region is a single record. Everything else stays zero-copy. +Offsets, segment lengths, and represented chunk counts are `u64`, so streaming +metadata remains correct on 32-bit targets. Individual boundary-search windows +are limited to `mothcdc::mincdc::MAX_CHUNK_SIZE`; the reader chunkers allocate +at least 4 MiB and provide fallible `try_new` constructors for invalid or +unallocatable configurations. + (An experimental second tier — content-defined *period detection* for phase-rotating runs — was evaluated and removed: mincdc self-aligns to most periods so it rarely helped, and it cost 76–99% throughput. See @@ -156,10 +169,6 @@ produced 7,798, with identical deduplicated content. Full method and the other corpora (Linux kernels, containers, SQLite, source trees) are in `examples/REALBENCH_RESULTS.md`. Obviously the benchmark is specific to this case. YMMV. -This fork also fixes a soundness bug in the upstream SIMD prefetch and adds test -coverage (cross-SIMD-width determinism, an invariant/oracle harness). - - ## Algorithm The basic idea of MinCDC is to choose chunk boundaries based on the minimum @@ -240,4 +249,3 @@ smaller. There is still a bias towards smaller chunks as MinCDC breaks ties in the minimum value towards the earlier breakpoint, but this bias is relatively small. - diff --git a/benches/throughput.rs b/benches/throughput.rs index 7380118..f39a227 100644 --- a/benches/throughput.rs +++ b/benches/throughput.rs @@ -166,7 +166,7 @@ fn caterpillar_scalar_reference(data: &[u8]) -> usize { Some((off, len)) if data[off..off + len] == c[..] => {}, _ => { records += 1; - last = Some((c.offset(), c.len())); + last = Some((c.offset() as usize, c.len())); }, } } @@ -184,7 +184,7 @@ fn bench_chunking(c: &mut Criterion) { g.bench_with_input(BenchmarkId::new("plain", &name), &data, |b, d| { b.iter(|| { - let mut acc = 0usize; + let mut acc = 0u64; for c in SliceChunker::new(d, MIN, MAX, MinCdcHash4::new()) { acc ^= c.offset(); } @@ -203,7 +203,7 @@ fn bench_chunking(c: &mut Criterion) { &data, |b, d| { b.iter(|| { - let mut acc = 0usize; + let mut acc = 0u64; for s in MothChunker::new(d, MIN, MAX) { acc ^= s.offset() ^ s.chunk_count(); } diff --git a/examples/catbench.rs b/examples/catbench.rs index b6a6718..1b1b31b 100644 --- a/examples/catbench.rs +++ b/examples/catbench.rs @@ -139,7 +139,7 @@ fn main() { // plain let (secs, _) = best_of(|| { let mut n = 0usize; - let mut s = 0usize; + let mut s = 0u64; for b in &blobs { for c in SliceChunker::new(b, MIN, MAX, cdc) { s ^= c.offset(); @@ -157,7 +157,7 @@ fn main() { { let (secs, _) = best_of(|| { let mut n = 0usize; - let mut s = 0usize; + let mut s = 0u64; for b in &blobs { for seg in MothChunker::with_cdc(b, MIN, MAX, cdc) { s ^= seg.offset(); diff --git a/examples/compare.rs b/examples/compare.rs index 3e9f9ac..1774636 100644 --- a/examples/compare.rs +++ b/examples/compare.rs @@ -40,7 +40,7 @@ struct Stats { /// Accumulates (key, stored_len) records and reports unique/dedup over them. struct Dedup { seen: std::collections::HashMap, - logical: usize, + logical: u64, records: usize, } impl Dedup { @@ -51,7 +51,7 @@ impl Dedup { records: 0, } } - fn add(&mut self, key: u64, stored_len: usize, logical_len: usize) { + fn add(&mut self, key: u64, stored_len: usize, logical_len: u64) { self.seen.entry(key).or_insert(stored_len); self.logical += logical_len; self.records += 1; @@ -124,7 +124,7 @@ fn run_fastcdc(data: &[u8], min: usize, avg: usize, max: usize, d: &mut Dedup) - }); for c in FastCDC::new(data, min, avg, max) { let b = &data[c.offset..c.offset + c.length]; - d.add(fnv1a(b), b.len(), b.len()); + d.add(fnv1a(b), b.len(), b.len() as u64); } Stats { records, @@ -144,7 +144,7 @@ fn run_mincdc(data: &[u8], min: usize, max: usize, d: &mut Dedup) -> Stats { n }); for c in SliceChunker::new(data, min, max, cdc) { - d.add(fnv1a(&c), c.len(), c.len()); + d.add(fnv1a(&c), c.len(), c.len() as u64); } Stats { records, @@ -206,7 +206,7 @@ fn scenario_versioned(min: usize, avg: usize, max: usize) { for data in [&v1[..], &v2[..]] { for c in FastCDC::new(data, min, avg, max) { let b = &data[c.offset..c.offset + c.length]; - d.add(fnv1a(b), b.len(), b.len()); + d.add(fnv1a(b), b.len(), b.len() as u64); } } println!( @@ -222,7 +222,7 @@ fn scenario_versioned(min: usize, avg: usize, max: usize) { let mut d = Dedup::new(); for data in [&v1[..], &v2[..]] { for c in SliceChunker::new(data, min, max, cdc) { - d.add(fnv1a(&c), c.len(), c.len()); + d.add(fnv1a(&c), c.len(), c.len() as u64); } } println!( diff --git a/examples/frontier.rs b/examples/frontier.rs index 179cc46..d3c9821 100644 --- a/examples/frontier.rs +++ b/examples/frontier.rs @@ -47,7 +47,7 @@ fn run(label: &str, blobs: &[Vec], total: u64, min: usize, max: usize, cdc: let mut best = f64::MAX; for _ in 0..3 { let t = Instant::now(); - let mut acc = 0usize; + let mut acc = 0u64; for b in blobs { for s in MothChunker::with_cdc(b, min, max, cdc) { acc ^= s.offset(); @@ -59,7 +59,7 @@ fn run(label: &str, blobs: &[Vec], total: u64, min: usize, max: usize, cdc: // Dedup + metadata (single deterministic pass). let mut store: HashMap = HashMap::new(); - let (mut records, mut chunks) = (0usize, 0usize); + let (mut records, mut chunks) = (0usize, 0u64); for b in blobs { for s in MothChunker::with_cdc(b, min, max, cdc) { let key = s.dedup_key(); diff --git a/examples/realbench.rs b/examples/realbench.rs index 138e714..32c2c8f 100644 --- a/examples/realbench.rs +++ b/examples/realbench.rs @@ -68,11 +68,11 @@ impl Acc { secs: 0.0, } } - fn add(&mut self, key: u64, stored: usize, logical: usize) { + fn add(&mut self, key: u64, stored: usize, logical: u64) { self.seen.entry(key).or_insert(stored); self.records += 1; - self.logical += logical as u64; - self.sizes.push(logical as u32); + self.logical += logical; + self.sizes.push(u32::try_from(logical).unwrap_or(u32::MAX)); } fn report(&self, name: &str, total_bytes: u64) { let gbps = total_bytes as f64 / self.secs.max(1e-9) / 1e9; @@ -147,7 +147,7 @@ fn main() { for b in &blobs { for c in FastCDC::new(b, MIN, AVG, fast_max()) { let s = &b[c.offset..c.offset + c.length]; - a.add(fnv1a(s), s.len(), s.len()); + a.add(fnv1a(s), s.len(), s.len() as u64); } } a.report("fastcdc-v2020", total_bytes); @@ -155,7 +155,7 @@ fn main() { // mincdc-plain let mut a = Acc::new(); let t = Instant::now(); - let mut sink = 0usize; + let mut sink = 0u64; for b in &blobs { for c in SliceChunker::new(b, MIN, MC_MAX, cdc) { sink ^= c.offset(); @@ -165,7 +165,7 @@ fn main() { std::hint::black_box(sink); for b in &blobs { for c in SliceChunker::new(b, MIN, MC_MAX, cdc) { - a.add(fnv1a(&c), c.len(), c.len()); + a.add(fnv1a(&c), c.len(), c.len() as u64); } } a.report("mincdc-plain", total_bytes); @@ -174,7 +174,7 @@ fn main() { { let mut a = Acc::new(); let t = Instant::now(); - let mut sink = 0usize; + let mut sink = 0u64; for b in &blobs { for s in MothChunker::with_cdc(b, MIN, MC_MAX, cdc) { sink ^= s.offset(); diff --git a/rustfmt.toml b/rustfmt.toml index f5c108d..41a35ff 100644 --- a/rustfmt.toml +++ b/rustfmt.toml @@ -1,4 +1,2 @@ -group_imports = "StdExternalCrate" -imports_granularity = "Module" match_block_trailing_comma = true use_field_init_shorthand = true diff --git a/src/capi.rs b/src/capi.rs index 6bd6d44..667a4ba 100644 --- a/src/capi.rs +++ b/src/capi.rs @@ -4,13 +4,13 @@ //! of the stable Rust API. use crate::caterpillar; -use crate::mincdc::{MinCdcHash4, next_chunk_len}; +use crate::mincdc::{MAX_CHUNK_SIZE, MinCdcHash4, next_chunk_len}; /// Length of the next chunk at the front of `data[..len]`, using /// [`MinCdcHash4`] with the default parameters. /// /// `eof != 0` means `data` is all the data there is. With `eof == 0` at -/// least `max_size + 1` bytes must be available, or the boundary is +/// at least `max_size + 1` bytes must be available, or the boundary is /// undecidable and 0 is returned (0 is also returned for `len == 0`). /// /// If `repeats_out` is non-null it receives the number of *additional* @@ -22,8 +22,10 @@ use crate::mincdc::{MinCdcHash4, next_chunk_len}; /// periodic. Pass null to skip the packed scan (plain MinCDC). /// /// # Safety -/// `data` must point to `len` readable bytes. `repeats_out` must be null or -/// point to a writable `usize`. +/// `data` must point to `len` readable bytes in one allocation, and `len` must +/// not exceed `isize::MAX`. `repeats_out` must be null or point to a properly +/// aligned writable `usize`. Invalid size configurations return zero; this +/// function does not panic for valid pointers. #[unsafe(no_mangle)] pub unsafe extern "C" fn mothcdc_next_chunk( data: *const u8, @@ -36,7 +38,13 @@ pub unsafe extern "C" fn mothcdc_next_chunk( if !repeats_out.is_null() { unsafe { *repeats_out = 0 }; } - if data.is_null() || len == 0 || min_size > max_size || max_size == 0 { + if data.is_null() + || len == 0 + || min_size > max_size + || max_size == 0 + || max_size > MAX_CHUNK_SIZE + || (eof == 0 && len <= max_size) + { return 0; } let bytes = unsafe { core::slice::from_raw_parts(data, len) }; @@ -118,4 +126,16 @@ mod tests { "C API stream protocol diverged from SliceChunker" ); } + + #[test] + fn capi_rejects_invalid_sizes_without_panicking() { + let byte = 0u8; + let mut repeats = usize::MAX; + let got = unsafe { mothcdc_next_chunk(&byte, 1, 0, usize::MAX, 0, &mut repeats) }; + assert_eq!(got, 0); + assert_eq!(repeats, 0); + + let got = unsafe { mothcdc_next_chunk(&byte, 1, 0, 8, 0, &mut repeats) }; + assert_eq!(got, 0, "non-EOF input needs max_size + 1 bytes"); + } } diff --git a/src/caterpillar.rs b/src/caterpillar.rs index 946cc79..acbfc63 100644 --- a/src/caterpillar.rs +++ b/src/caterpillar.rs @@ -6,9 +6,9 @@ //! zero-fill or repeated-record region costs metadata out of all proportion to //! its information content. //! -//! [`MothChunker`] wraps a [`SliceChunker`] and run-length-encodes -//! maximal runs of byte-identical adjacent chunks into a single -//! [`Segment::Caterpillar`] record (the unit + a repeat count). It catches +//! [`MothChunker`] wraps a [`SliceChunker`](crate::mincdc::SliceChunker) and run-length-encodes +//! maximal runs of byte-identical adjacent chunks into a single [`Segment`] +//! record (the unit + a repeat count). It catches //! zero-fill, constant bytes, and repeated blocks; it is a no-op (one slice //! compare per chunk) on data with no runs, so it keeps mincdc's speed and //! deduplication everywhere else. The output is lossless. @@ -23,8 +23,9 @@ //! numbers). //! //! [`MothChunker`] works on an in-memory byte slice (it wraps -//! [`SliceChunker`]). For inputs larger than memory, [`MothReadChunker`] -//! does the same coalescing over a streaming [`ReadChunker`] in bounded memory, +//! [`SliceChunker`](crate::mincdc::SliceChunker)). For inputs larger than memory, +//! [`MothReadChunker`] does the same coalescing over a streaming +//! [`ReadChunker`](crate::mincdc::ReadChunker) in bounded memory, //! yielding borrowed [`Segment`]s (valid until the next call, like //! [`ReadChunker::next`](crate::mincdc::ReadChunker)). Runs coalesce across buffer //! refills (a run's unit is copied once — at most `max_size` bytes — when it @@ -40,7 +41,7 @@ //! //! // The whole zero run collapses into one record instead of ~64 chunks. //! assert_eq!(segs.len(), 1); -//! assert!(matches!(segs[0], Segment::Caterpillar { .. })); +//! assert!(segs[0].is_caterpillar()); //! // `dedup_key()` gives the unique bytes to fingerprint/store, regardless of variant. //! assert!(segs[0].dedup_key().iter().all(|&b| b == 0)); //! ``` @@ -128,42 +129,62 @@ pub(crate) fn packed_repeats(tail: &[u8], u: usize, max_size: usize) -> usize { /// **only until the next call**. Process it (or copy [`dedup_key`](Self::dedup_key)) /// before advancing the streaming chunker. #[derive(Debug)] -pub enum Segment<'a> { - /// A single chunk whose neighbor differed — emitted as-is. +pub struct Segment<'a> { + kind: SegmentKind<'a>, +} + +#[derive(Debug)] +enum SegmentKind<'a> { Solo(Chunk<'a>), - /// A run of `count` (>= 2) byte-identical adjacent chunks (tier 1). Caterpillar { - /// Start offset of the run within the input. - offset: usize, - /// The repeated chunk's bytes. + offset: u64, unit: &'a [u8], - /// Number of times `unit` repeats (>= 2). - count: usize, + count: u64, }, } impl<'a> Segment<'a> { + fn solo(chunk: Chunk<'a>) -> Self { + Self { + kind: SegmentKind::Solo(chunk), + } + } + + fn caterpillar(offset: u64, unit: &'a [u8], count: u64) -> Self { + debug_assert!(!unit.is_empty()); + debug_assert!(count >= 2); + Self { + kind: SegmentKind::Caterpillar { + offset, + unit, + count, + }, + } + } + /// Start offset of this segment within the input. - pub fn offset(&self) -> usize { - match self { - Segment::Solo(c) => c.offset(), - Segment::Caterpillar { offset, .. } => *offset, + pub fn offset(&self) -> u64 { + match &self.kind { + SegmentKind::Solo(c) => c.offset(), + SegmentKind::Caterpillar { offset, .. } => *offset, } } - /// Number of underlying chunks represented (1 for [`Segment::Solo`]). - pub fn chunk_count(&self) -> usize { - match self { - Segment::Solo(_) => 1, - Segment::Caterpillar { count, .. } => *count, + /// Number of underlying chunks represented (one for a solo segment). + pub fn chunk_count(&self) -> u64 { + match &self.kind { + SegmentKind::Solo(_) => 1, + SegmentKind::Caterpillar { count, .. } => *count, } } /// Total number of bytes covered by this segment. - pub fn len(&self) -> usize { - match self { - Segment::Solo(c) => c.len(), - Segment::Caterpillar { unit, count, .. } => unit.len() * count, + pub fn len(&self) -> u64 { + match &self.kind { + SegmentKind::Solo(c) => c.len() as u64, + SegmentKind::Caterpillar { unit, count, .. } => (unit.len() as u64) + .checked_mul(*count) + .expect("validated segment length exceeded u64::MAX"), } } @@ -172,9 +193,13 @@ impl<'a> Segment<'a> { self.len() == 0 } + /// Whether this segment represents two or more identical adjacent chunks. + pub fn is_caterpillar(&self) -> bool { + matches!(self.kind, SegmentKind::Caterpillar { .. }) + } + /// The unique content to fingerprint and store for content addressing: the - /// chunk bytes ([`Segment::Solo`]) or the repeated unit - /// ([`Segment::Caterpillar`]). + /// solo chunk bytes or the repeated unit. /// /// This is also exactly the bytes to tile back to [`len`](Self::len) to /// reconstruct the segment, so a store-then-restore round-trip is lossless: @@ -182,24 +207,24 @@ impl<'a> Segment<'a> { /// [`len`](Self::len), and on restore tile the stored bytes to `len`. See /// [`reconstruct_into`](Self::reconstruct_into). pub fn dedup_key(&self) -> &[u8] { - match self { - Segment::Solo(c) => c, - Segment::Caterpillar { unit, .. } => unit, + match &self.kind { + SegmentKind::Solo(c) => c, + SegmentKind::Caterpillar { unit, .. } => unit, } } /// Appends this segment's original bytes to `out` (the inverse of chunking): - /// the chunk for [`Segment::Solo`] or the unit repeated `count` times for - /// [`Segment::Caterpillar`]. Equivalent to tiling [`dedup_key`](Self::dedup_key) + /// the solo chunk or the unit repeated `count` times. Equivalent to tiling + /// [`dedup_key`](Self::dedup_key) /// to [`len`](Self::len). pub fn reconstruct_into(&self, out: &mut Vec) { let key = self.dedup_key(); let total = self.len(); - let mut written = 0; + let mut written = 0u64; while written < total { - let take = key.len().min(total - written); + let take = (key.len() as u64).min(total - written) as usize; out.extend_from_slice(&key[..take]); - written += take; + written += take as u64; } } } @@ -238,7 +263,7 @@ impl<'a, C: Cdc> Iterator for MothChunker<'a, C> { fn next(&mut self) -> Option> { let first = self.carry.take().or_else(|| self.inner.next())?; - let start = first.offset(); + let start = usize::try_from(first.offset()).expect("slice offset does not fit usize"); let first_len = first.len(); let unit: &'a [u8] = &self.data[start..start + first_len]; @@ -264,13 +289,9 @@ impl<'a, C: Cdc> Iterator for MothChunker<'a, C> { }; self.carry = pending; if count >= 2 { - Some(Segment::Caterpillar { - offset: start, - unit, - count, - }) + Some(Segment::caterpillar(start as u64, unit, count as u64)) } else { - Some(Segment::Solo(first)) + Some(Segment::solo(first)) } } } @@ -298,25 +319,32 @@ pub struct MothReadChunker { buf: Vec, buf_offset: usize, unread: usize, - stream_offset: usize, + stream_offset: u64, eof: bool, done: bool, /// Length of the chunk at `buf_offset` if already computed by a previous /// call's lookahead — avoids recomputing the boundary (a SIMD scan) twice. pending_len: Option, /// A run continued across buffer refills: (owned unit bytes, stream offset - /// of the run start, chunks counted so far — always >= 2). - carry_run: Option<(Vec, usize, usize)>, + /// of the run start, chunks counted so far — always >= 1). A singleton is + /// retained only when its neighbor needs another reader refill to decide. + carry_run: Option<(Vec, u64, u64)>, /// Owns the unit of a carried run while its borrowed [`Segment`] is live. emit_unit: Vec, } impl MothReadChunker { /// Creates a zero-copy streaming caterpillar chunker using [`MinCdcHash4`] - /// with the default hash parameters. + /// with the default hash parameters. It allocates at least 4 MiB; use + /// [`MothReadChunker::try_new`] to handle invalid sizes or allocation failure. pub fn new(reader: R, min_size: usize, max_size: usize) -> Self { Self::with_cdc(reader, min_size, max_size, MinCdcHash4::new()) } + + /// Tries to create a streaming chunker without arithmetic or allocation panics. + pub fn try_new(reader: R, min_size: usize, max_size: usize) -> io::Result { + Self::try_with_cdc(reader, min_size, max_size, MinCdcHash4::new()) + } } impl MothReadChunker { @@ -324,14 +352,23 @@ impl MothReadChunker { /// instance (e.g. [`MinCdcHash4::with_params`] or /// [`MinCdc4`](crate::mincdc::MinCdc4)). pub fn with_cdc(reader: R, min_size: usize, max_size: usize, cdc: C) -> Self { - assert!(min_size <= max_size && max_size > 0); - let buf_size = crate::MIN_BUFFER_SIZE + (max_size + 1) + min_size * 4; - Self { + Self::try_with_cdc(reader, min_size, max_size, cdc) + .expect("invalid MothReadChunker configuration or buffer allocation failed") + } + + /// Tries to create a streaming chunker with a custom [`Cdc`] instance. + pub fn try_with_cdc(reader: R, min_size: usize, max_size: usize, cdc: C) -> io::Result { + let buf_size = crate::mincdc::checked_buffer_size(min_size, max_size)?; + let mut buf = Vec::new(); + buf.try_reserve_exact(buf_size) + .map_err(|e| io::Error::new(io::ErrorKind::OutOfMemory, e))?; + buf.resize(buf_size, 0); + Ok(Self { min_size, max_size, cdc, reader, - buf: vec![0; buf_size], + buf, buf_offset: 0, unread: 0, stream_offset: 0, @@ -340,7 +377,24 @@ impl MothReadChunker { pending_len: None, carry_run: None, emit_unit: Vec::new(), - } + }) + } + + /// Returns a shared reference to the underlying reader. + pub fn get_ref(&self) -> &R { + &self.reader + } + + /// Returns a mutable reference to the underlying reader. + /// + /// Reading from it directly can invalidate chunker state. + pub fn get_mut(&mut self) -> &mut R { + &mut self.reader + } + + /// Returns the underlying reader, discarding any bytes already read ahead. + pub fn into_inner(self) -> R { + self.reader } /// Length of the chunk starting at buffer position `p` given `avail` buffered @@ -393,9 +447,15 @@ impl MothReadChunker { .copy_within(self.buf_offset..self.buf_offset + self.unread, 0); self.buf_offset = 0; } - let n = self - .reader - .read(&mut self.buf[self.buf_offset + self.unread..])?; + let n = loop { + match self + .reader + .read(&mut self.buf[self.buf_offset + self.unread..]) + { + Err(e) if e.kind() == io::ErrorKind::Interrupted => continue, + result => break result?, + } + }; if n == 0 { self.eof = true; break; @@ -423,10 +483,10 @@ impl MothReadChunker { // The stream ended exactly at a carried run's boundary. return Ok(self.carry_run.take().map(|(unit, off, count)| { self.emit_unit = unit; - Segment::Caterpillar { - offset: off, - unit: &self.emit_unit, - count, + if count >= 2 { + Segment::caterpillar(off, &self.emit_unit, count) + } else { + Segment::solo(Chunk::new(&self.emit_unit, off)) } })); } @@ -441,8 +501,21 @@ impl MothReadChunker { let (run_len, more, pend) = self.coalesce_from(base, u); self.buf_offset += run_len; self.unread -= run_len; - self.stream_offset += run_len; - let count = carried + more; + self.stream_offset = self + .stream_offset + .checked_add(run_len as u64) + .ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidData, + "stream offset exceeded u64::MAX", + ) + })?; + let count = carried.checked_add(more as u64).ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidData, + "segment chunk count exceeded u64::MAX", + ) + })?; if pend.is_none() && !self.eof { // Ran out of buffer again mid-run: keep carrying. self.carry_run = Some((unit, run_off, count)); @@ -450,22 +523,19 @@ impl MothReadChunker { } self.pending_len = pend; self.emit_unit = unit; - return Ok(Some(Segment::Caterpillar { - offset: run_off, - unit: &self.emit_unit, - count, - })); + return Ok(Some(Segment::caterpillar(run_off, &self.emit_unit, count))); }, // The run ended at the refill boundary: emit it and stash // the just-computed boundary (if any) for the next call. boundary => { self.pending_len = boundary; self.emit_unit = unit; - return Ok(Some(Segment::Caterpillar { - offset: run_off, - unit: &self.emit_unit, - count: carried, - })); + let segment = if carried >= 2 { + Segment::caterpillar(run_off, &self.emit_unit, carried) + } else { + Segment::solo(Chunk::new(&self.emit_unit, run_off)) + }; + return Ok(Some(segment)); }, } } @@ -482,26 +552,33 @@ impl MothReadChunker { let (run_len, count, pend) = self.coalesce_from(base, unit_len); self.buf_offset += run_len; self.unread -= run_len; - self.stream_offset += run_len; - - if pend.is_none() && !self.eof && count >= 2 { - // The run reached the end of the buffered data and may continue - // after a refill: copy the unit (once per crossing run) and - // keep counting instead of splitting the record. - self.carry_run = - Some((self.buf[base..base + unit_len].to_vec(), base_stream, count)); + self.stream_offset = + self.stream_offset + .checked_add(run_len as u64) + .ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidData, + "stream offset exceeded u64::MAX", + ) + })?; + + if pend.is_none() && !self.eof { + // The next boundary needs another refill. Retain even a + // singleton so an arbitrarily fragmented reader cannot split a + // repeated run before its identical neighbor is decidable. + self.carry_run = Some(( + self.buf[base..base + unit_len].to_vec(), + base_stream, + count as u64, + )); continue; } self.pending_len = pend; let seg = if count >= 2 { - Segment::Caterpillar { - offset: base_stream, - unit: &self.buf[base..base + unit_len], - count, - } + Segment::caterpillar(base_stream, &self.buf[base..base + unit_len], count as u64) } else { - Segment::Solo(Chunk::new(&self.buf[base..base + unit_len], base_stream)) + Segment::solo(Chunk::new(&self.buf[base..base + unit_len], base_stream)) }; return Ok(Some(seg)); } @@ -534,8 +611,8 @@ mod tests { let plain = SliceChunker::new(data, min, max, cdc).count(); let mut records = 0usize; - let mut expanded = 0usize; - let mut next_off = 0usize; + let mut expanded = 0u64; + let mut next_off = 0u64; let mut rebuilt: Vec = Vec::with_capacity(data.len()); for s in MothChunker::with_cdc(data, min, max, cdc) { assert_eq!(s.offset(), next_off, "{label}: offset not contiguous"); @@ -545,7 +622,10 @@ mod tests { next_off += s.len(); } assert_eq!(rebuilt, data, "{label}: must reconstruct input exactly"); - assert_eq!(expanded, plain, "{label}: must represent same chunk count"); + assert_eq!( + expanded, plain as u64, + "{label}: must represent same chunk count" + ); (plain, records) } @@ -568,8 +648,8 @@ mod tests { /// chunk from the plain chunker and RLE-coalesce byte-identical neighbors. /// The packed-scanning fast path must produce exactly this segment stream — /// same grouping, same offsets, same unit bytes. - fn reference_segments(data: &[u8], min: usize, max: usize) -> Vec<(usize, Vec, usize)> { - let mut out: Vec<(usize, Vec, usize)> = Vec::new(); + fn reference_segments(data: &[u8], min: usize, max: usize) -> Vec<(u64, Vec, u64)> { + let mut out: Vec<(u64, Vec, u64)> = Vec::new(); for c in SliceChunker::new(data, min, max, MinCdcHash4::new()) { match out.last_mut() { Some((_, unit, count)) if unit[..] == c[..] => *count += 1, @@ -580,15 +660,8 @@ mod tests { } fn assert_matches_reference(label: &str, data: &[u8], min: usize, max: usize) { - let got: Vec<(usize, Vec, usize)> = MothChunker::new(data, min, max) - .map(|s| match s { - Segment::Solo(c) => (c.offset(), c.to_vec(), 1), - Segment::Caterpillar { - offset, - unit, - count, - } => (offset, unit.to_vec(), count), - }) + let got: Vec<(u64, Vec, u64)> = MothChunker::new(data, min, max) + .map(|s| (s.offset(), s.dedup_key().to_vec(), s.chunk_count())) .collect(); let want = reference_segments(data, min, max); assert_eq!(got, want, "{label} (min={min} max={max})"); @@ -756,17 +829,20 @@ mod tests { let data = vec![0u8; n]; let mut rc = MothReadChunker::new(Cursor::new(&data), MIN, MAX); let mut records = 0usize; - let mut chunks = 0usize; - let mut covered = 0usize; + let mut chunks = 0u64; + let mut covered = 0u64; while let Some(s) = rc.next().unwrap() { records += 1; chunks += s.chunk_count(); covered += s.len(); assert!(s.dedup_key().iter().all(|&b| b == 0)); } - assert_eq!(covered, n); + assert_eq!(covered, n as u64); let plain = SliceChunker::new(&data, MIN, MAX, MinCdcHash4::new()).count(); - assert_eq!(chunks, plain, "must represent the same underlying chunks"); + assert_eq!( + chunks, plain as u64, + "must represent the same underlying chunks" + ); assert!( records <= 2, "a single giant run must not split at buffer refills (got {records} records)" @@ -777,7 +853,7 @@ mod tests { data.extend_from_slice(&vec![0u8; 12 * 1024 * 1024]); data.extend_from_slice(&xorshift(12, 64 * 1024)); let mut rc = MothReadChunker::new(Cursor::new(&data), MIN, MAX); - let (mut records, mut chunks, mut covered) = (0usize, 0usize, 0usize); + let (mut records, mut chunks, mut covered) = (0usize, 0u64, 0u64); let mut rebuilt = Vec::with_capacity(data.len()); while let Some(s) = rc.next().unwrap() { records += 1; @@ -785,10 +861,10 @@ mod tests { covered += s.len(); s.reconstruct_into(&mut rebuilt); } - assert_eq!(covered, data.len()); + assert_eq!(covered, data.len() as u64); assert_eq!(rebuilt, data, "must reconstruct exactly"); let plain = SliceChunker::new(&data, MIN, MAX, MinCdcHash4::new()).count(); - assert_eq!(chunks, plain); + assert_eq!(chunks, plain as u64); // ~64 KiB of random on each side is ~40-70 records; the 12 MiB zero // run must contribute ~1, not ~3 (one per 4 MiB refill). assert!( @@ -797,6 +873,54 @@ mod tests { ); } + struct ReadOnceThenEof { + data: Vec, + sent: bool, + } + + impl std::io::Read for ReadOnceThenEof { + fn read(&mut self, buf: &mut [u8]) -> std::io::Result { + if self.sent { + return Ok(0); + } + assert!(buf.len() >= self.data.len()); + buf[..self.data.len()].copy_from_slice(&self.data); + self.sent = true; + Ok(self.data.len()) + } + } + + #[test] + fn streaming_emits_carried_runs_when_the_next_read_is_eof() { + let cases = [ + // Four identical chunks consume the entire first read, so EOF is + // discovered with a repeated run still in carry_run. + vec![0xA5; 64], + // The first full chunk is carried while the short, differing tail + // remains undecidable; EOF then terminates that singleton run. + [vec![0x11; 16], vec![0x22; 8]].concat(), + ]; + + for data in cases { + let want: Vec<_> = MothChunker::new(&data, 16, 16) + .map(|s| (s.offset(), s.len(), s.chunk_count(), s.dedup_key().to_vec())) + .collect(); + let reader = ReadOnceThenEof { data, sent: false }; + let mut chunker = MothReadChunker::new(reader, 16, 16); + let mut got = Vec::new(); + while let Some(segment) = chunker.next().unwrap() { + got.push(( + segment.offset(), + segment.len(), + segment.chunk_count(), + segment.dedup_key().to_vec(), + )); + } + assert_eq!(got, want); + assert!(chunker.next().unwrap().is_none()); + } + } + #[test] fn self_aligns_and_collapses_periodic_data() { // mincdc self-aligns boundaries to periods when a period multiple fits in diff --git a/src/lib.rs b/src/lib.rs index 6fbcd7f..6201b6a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -7,7 +7,13 @@ //! The core chunking algorithm lives in [`mincdc`] (inherited from upstream //! MinCDC); the coalescing layer lives in [`caterpillar`], built on top of it. //! -//! # Examples +//! The default [`MothChunker`] and [`MothReadChunker`] use the robust hashed +//! MinCDC boundary selector and emit validated [`Segment`] values. A segment is +//! either one ordinary chunk or a maximal run of identical adjacent chunks; +//! use [`Segment::dedup_key`] for storage and [`Segment::len`] for its logical +//! byte length. +//! +//! # In-memory example //! //! ```rust //! use mothcdc::MothChunker; @@ -17,6 +23,25 @@ //! //! // The whole zero run collapses into one record instead of ~64 chunks. //! assert_eq!(segs.len(), 1); +//! assert!(segs[0].is_caterpillar()); +//! assert_eq!(segs[0].len(), data.len() as u64); +//! ``` +//! +//! # Streaming example +//! +//! ```rust +//! use std::io::Cursor; +//! use mothcdc::MothReadChunker; +//! +//! let data = vec![0u8; 128 * 4096]; +//! let mut chunker = MothReadChunker::new(Cursor::new(&data), 4096, 12288); +//! let mut restored = Vec::new(); +//! while let Some(segment) = chunker.next()? { +//! // Streaming segments borrow the internal buffer until the next call. +//! segment.reconstruct_into(&mut restored); +//! } +//! assert_eq!(restored, data); +//! # Ok::<(), std::io::Error>(()) //! ``` #![warn(missing_docs)] diff --git a/src/mincdc.rs b/src/mincdc.rs index a426fcc..9f8c7f1 100644 --- a/src/mincdc.rs +++ b/src/mincdc.rs @@ -15,27 +15,30 @@ //! //! This crate provides two implementations of MinCDC, both with a window size //! of 4: -//! - [`MinCdc4`], where the evaluation function is +//! - [`MinCdc4`](crate::mincdc::MinCdc4), where the evaluation function is //! `u32::from_le_bytes(bytes[i - 4..i])`, i.e. a window size of 4 bytes //! interpreting the bytes as a little-endian `u32`, and -//! - [`MinCdcHash4`], where the evaluation function is +//! - [`MinCdcHash4`](crate::mincdc::MinCdcHash4), where the evaluation function is //! `hash(u32::from_le_bytes(bytes[i - 4..i]))`. The hash function used is //! the very simple `hash(x) = x.wrapping_mul(a).wrapping_add(b)`, for //! some constants `a` and `b`. //! -//! **[`MinCdcHash4`] can be slightly (~10%) slower but is far more robust and +//! **[`MinCdcHash4`](crate::mincdc::MinCdcHash4) can be slightly (~10%) slower but is far more robust and //! predictable, it is the recommended default**. //! //! # Usage //! //! This module provides two chunkers: -//! - [`SliceChunker`] for chunking a byte slice, and -//! - [`ReadChunker`] for chunking a reader implementing [`Read`]. +//! - [`SliceChunker`](crate::mincdc::SliceChunker) for chunking a byte slice, and +//! - [`ReadChunker`](crate::mincdc::ReadChunker) for chunking a reader implementing +//! [`Read`](std::io::Read). //! //! Both chunkers take a desired minimum and maximum chunk size as well as a -//! [`Cdc`] instance (either [`MinCdc4`] or -//! [`MinCdcHash4`]). Then by iterating over the chunker -//! (or calling [`next()`](ReadChunker::next) in the case of [`ReadChunker`]) you get chunks of type +//! [`Cdc`](crate::mincdc::Cdc) instance (either +//! [`MinCdc4`](crate::mincdc::MinCdc4) or +//! [`MinCdcHash4`](crate::mincdc::MinCdcHash4)). Then by iterating over the chunker +//! (or calling [`next()`](crate::mincdc::ReadChunker::next) in the case of +//! [`ReadChunker`](crate::mincdc::ReadChunker)) you get chunks of type //! [`Chunk`], which derefs to a byte slice, but also contains the offset of //! that chunk in the input stream. //! @@ -66,16 +69,28 @@ use crate::simd; pub(crate) const DEFAULT_MULTIPLIER: u32 = 0x915f77f5; pub(crate) const DEFAULT_ADDEND: u32 = 0x34636463; +/// Largest supported boundary-search window. +/// +/// SIMD implementations store candidate offsets in `u32` lanes. Keeping this +/// limit consistent across targets also keeps chunk boundaries portable. +pub const MAX_CHUNK_SIZE: usize = u32::MAX as usize - 1; + /// A trait for determining splitpoints in a content-defined way. pub trait Cdc { /// The amount of bytes needed before position `i` to determine if `i` is a /// splitpoint. /// - /// Should return a constant. + /// Must return the same non-zero value for the lifetime of the chunker. fn window_size(&self) -> usize; - /// Returns the best splitpoint `i <= bytes.len()`, indicating bytes is to - /// be split into `bytes[..i]` and `bytes[i..]`. + /// Returns the best splitpoint `i`, indicating bytes is to be split into + /// `bytes[..i]` and `bytes[i..]`. + /// + /// # Contract + /// + /// If `bytes.len() < self.window_size()`, this must return `bytes.len()`. + /// Otherwise it must return a value in + /// `self.window_size()..=bytes.len()`. Chunkers enforce this contract. fn best_splitpoint(&self, bytes: &[u8]) -> usize; } @@ -83,7 +98,6 @@ pub trait Cdc { /// /// This chooses the first splitpoint `i` where /// `u32::from_le_bytes(bytes[i-4..i])` is minimized. -#[non_exhaustive] #[derive(Copy, Clone, Default, Debug)] pub struct MinCdc4; @@ -177,9 +191,10 @@ impl<'a, C> SliceChunker<'a, C> { /// respect the minimum size. /// /// # Panics - /// Panics if `min_size > max_size` or `max_size == 0`. + /// Panics if the sizes are invalid or `max_size` exceeds + /// [`MAX_CHUNK_SIZE`]. pub const fn new(bytes: &'a [u8], min_size: usize, max_size: usize, cdc: C) -> Self { - assert!(min_size <= max_size && max_size > 0); + assert!(min_size <= max_size && max_size > 0 && max_size <= MAX_CHUNK_SIZE); Self { min_size, @@ -209,14 +224,25 @@ pub(crate) fn next_chunk_len( } // Can't reliably place a boundary without the full decision window, unless // this is all the data there is. - if !eof && n < max_size + 1 { + if !eof && n <= max_size { return None; } if n <= min_size { return Some(n); // final short chunk (only reachable at eof) } - let start = min_size.saturating_sub(cdc.window_size()); - Some(start + cdc.best_splitpoint(&avail[start..max_size.min(n)])) + let window = cdc.window_size(); + assert!(window > 0, "Cdc::window_size() must be non-zero"); + let start = min_size.saturating_sub(window); + let search = &avail[start..max_size.min(n)]; + let split = cdc.best_splitpoint(search); + let minimum = window.min(search.len()); + assert!( + split >= minimum && split <= search.len(), + "Cdc::best_splitpoint() returned {split} for {} bytes with window size {window}; expected {minimum}..={}", + search.len(), + search.len() + ); + Some(start + split) } impl<'a, C: Cdc> Iterator for SliceChunker<'a, C> { @@ -235,7 +261,10 @@ impl<'a, C: Cdc> Iterator for SliceChunker<'a, C> { &self.cdc, ) .expect("eof=true always yields a chunk for non-empty input"); - let ret = Chunk::new(&self.bytes[self.offset..self.offset + len], self.offset); + let ret = Chunk::new( + &self.bytes[self.offset..self.offset + len], + self.offset as u64, + ); self.offset += len; Some(ret) } @@ -246,7 +275,8 @@ impl<'a, C: Cdc> FusedIterator for SliceChunker<'a, C> {} /// A chunker for a reader implementing [`Read`]. /// /// Note that unlike [`SliceChunker`] this stores bytes in an internal buffer -/// which is re-used and thus it can not implement [`Iterator`]. +/// which is re-used and thus it can not implement [`Iterator`]. It allocates at +/// least 4 MiB so that ordinary reads amortize buffer movement. #[derive(Clone)] pub struct ReadChunker { min_size: usize, @@ -256,7 +286,8 @@ pub struct ReadChunker { buf: Vec, buf_offset: usize, unread_bytes_in_buf: usize, - stream_offset: usize, + stream_offset: u64, + done: bool, } impl ReadChunker { @@ -267,30 +298,83 @@ impl ReadChunker { /// respect the minimum size. /// /// # Panics - /// Panics if `min_size > max_size` or `max_size == 0`. + /// Panics if the sizes are invalid, arithmetic overflows, or the internal + /// buffer cannot be allocated. Use [`ReadChunker::try_new`] to handle those + /// cases as errors. pub fn new(reader: R, min_size: usize, max_size: usize, cdc: C) -> Self { - assert!(min_size <= max_size && max_size > 0); + Self::try_new(reader, min_size, max_size, cdc) + .expect("invalid ReadChunker configuration or buffer allocation failed") + } - let bytes_needed_for_decision = max_size + 1; - let buf_size = MIN_BUFFER_SIZE + bytes_needed_for_decision + min_size * 4; - Self { + /// Tries to create a reader chunker without arithmetic or allocation panics. + /// Invalid sizes produce [`io::ErrorKind::InvalidInput`]; allocation failure + /// produces [`io::ErrorKind::OutOfMemory`]. + pub fn try_new(reader: R, min_size: usize, max_size: usize, cdc: C) -> io::Result { + let buf_size = checked_buffer_size(min_size, max_size)?; + let mut buf = Vec::new(); + buf.try_reserve_exact(buf_size) + .map_err(|e| io::Error::new(io::ErrorKind::OutOfMemory, e))?; + buf.resize(buf_size, 0); + Ok(Self { min_size, max_size, cdc, reader, - buf: vec![0; buf_size], + buf, buf_offset: 0, unread_bytes_in_buf: 0, stream_offset: 0, - } + done: false, + }) + } + + /// Returns a shared reference to the underlying reader. + pub fn get_ref(&self) -> &R { + &self.reader } + + /// Returns a mutable reference to the underlying reader. + /// + /// Reading from it directly can invalidate chunker state. + pub fn get_mut(&mut self) -> &mut R { + &mut self.reader + } + + /// Returns the underlying reader, discarding any bytes already read ahead. + pub fn into_inner(self) -> R { + self.reader + } +} + +pub(crate) fn checked_buffer_size(min_size: usize, max_size: usize) -> io::Result { + if min_size > max_size || max_size == 0 || max_size > MAX_CHUNK_SIZE { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("expected 0 <= min_size <= max_size <= {MAX_CHUNK_SIZE}, with max_size > 0"), + )); + } + let decision = max_size + .checked_add(1) + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "max_size + 1 overflowed"))?; + let slack = min_size + .checked_mul(4) + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "min_size * 4 overflowed"))?; + MIN_BUFFER_SIZE + .checked_add(decision) + .and_then(|n| n.checked_add(slack)) + .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "buffer size overflowed")) } impl ReadChunker { /// Gets the next [`Chunk`] from the reader, or [`None`] if it is exhausted. + /// + /// A `read` returning `Ok(0)` is treated as end of input. + /// [`io::ErrorKind::Interrupted`] is retried internally; any other error is + /// returned, and progress already made is preserved so the next call can + /// resume. #[allow(clippy::should_implement_trait)] pub fn next(&mut self) -> io::Result>> { - if self.stream_offset == usize::MAX { + if self.done { return Ok(None); } @@ -305,9 +389,15 @@ impl ReadChunker { self.buf_offset = 0; } - let bytes_read = self - .reader - .read(&mut self.buf[self.buf_offset + self.unread_bytes_in_buf..])?; + let bytes_read = loop { + match self + .reader + .read(&mut self.buf[self.buf_offset + self.unread_bytes_in_buf..]) + { + Err(e) if e.kind() == io::ErrorKind::Interrupted => continue, + result => break result?, + } + }; if bytes_read == 0 { break; } @@ -321,7 +411,7 @@ impl ReadChunker { match next_chunk_len(avail, self.min_size, self.max_size, eof, &self.cdc) { None => { // Reader is exhausted. - self.stream_offset = usize::MAX; + self.done = true; Ok(None) }, Some(len) => { @@ -329,7 +419,13 @@ impl ReadChunker { &self.buf[self.buf_offset..self.buf_offset + len], self.stream_offset, ); - self.stream_offset += len; + self.stream_offset = + self.stream_offset.checked_add(len as u64).ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidData, + "stream offset exceeded u64::MAX", + ) + })?; self.buf_offset += len; self.unread_bytes_in_buf -= len; Ok(Some(ret)) @@ -344,17 +440,18 @@ impl ReadChunker { #[derive(Copy, Clone, Debug)] pub struct Chunk<'a> { bytes: &'a [u8], - offset: usize, + offset: u64, } impl<'a> Chunk<'a> { /// Creates a new [`Chunk`] with the given `bytes` and `offset`. - pub const fn new(bytes: &'a [u8], offset: usize) -> Self { + pub const fn new(bytes: &'a [u8], offset: u64) -> Self { Self { bytes, offset } } - /// The start offset of this chunk within the full data. - pub const fn offset(&self) -> usize { + /// The start offset of this chunk within the full data. Stream offsets use + /// `u64` even on 32-bit targets. + pub const fn offset(&self) -> u64 { self.offset } } @@ -369,7 +466,7 @@ impl<'a> Deref for Chunk<'a> { #[cfg(test)] mod test { - use std::io::Cursor; + use std::io::{self, Cursor, Read}; use rand::distr::StandardUniform; use rand::prelude::*; @@ -439,18 +536,18 @@ mod test { #[test] fn test_read_slice_equiv() { - let bounds = [1, 2, 3, 4, 6, 8, 15, 27, 62, 90, 120, 200]; + // The integration invariant suite exercises a much larger matrix + // against an independent oracle. Keep this unit smoke test small so a + // normal debug `cargo test` remains fast. + let bounds = [1, 4, 27, 200]; for min_size in &bounds { for max_size in &bounds { if min_size > max_size { continue; } - for size in 0..4096 { - let rng = SmallRng::seed_from_u64(size); - let bytes: Vec = rng - .sample_iter(StandardUniform) - .take(size as usize) - .collect(); + for size in [0usize, 1, 3, 4, 17, 200, 511, 4096] { + let rng = SmallRng::seed_from_u64(size as u64); + let bytes: Vec = rng.sample_iter(StandardUniform).take(size).collect(); let reader = Cursor::new(&bytes); let mut read_chunker = ReadChunker::new(reader, *min_size, *max_size, MinCdc4); @@ -477,4 +574,217 @@ mod test { } } } + + #[derive(Clone, Copy)] + struct InvalidCdc { + window: usize, + split: usize, + } + + impl super::Cdc for InvalidCdc { + fn window_size(&self) -> usize { + self.window + } + + fn best_splitpoint(&self, _bytes: &[u8]) -> usize { + self.split + } + } + + #[test] + #[should_panic(expected = "best_splitpoint")] + fn invalid_cdc_cannot_yield_empty_chunks() { + let mut chunks = SliceChunker::new( + b"enough input to require a split", + 1, + 8, + InvalidCdc { + window: 1, + split: 0, + }, + ); + let _ = chunks.next(); + } + + #[test] + #[should_panic(expected = "window_size")] + fn invalid_cdc_window_is_rejected() { + let mut chunks = SliceChunker::new( + b"enough input", + 1, + 8, + InvalidCdc { + window: 0, + split: 1, + }, + ); + let _ = chunks.next(); + } + + struct InterruptOnce { + inner: R, + interrupted: bool, + } + + impl Read for InterruptOnce { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if !self.interrupted { + self.interrupted = true; + return Err(io::ErrorKind::Interrupted.into()); + } + self.inner.read(buf) + } + } + + #[test] + fn read_chunker_retries_interrupted() { + let data = vec![7u8; 1024]; + let reader = InterruptOnce { + inner: Cursor::new(&data), + interrupted: false, + }; + let mut chunks = ReadChunker::new(reader, 16, 64, MinCdcHash4::new()); + let mut rebuilt = Vec::new(); + while let Some(chunk) = chunks.next().unwrap() { + rebuilt.extend_from_slice(&chunk); + } + assert_eq!(rebuilt, data); + } + + struct OneByteReader(R); + + impl Read for OneByteReader { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + let n = buf.len().min(1); + self.0.read(&mut buf[..n]) + } + } + + #[test] + fn read_chunker_handles_one_byte_reads() { + let data: Vec = (0..8192).map(|i| (i * 31) as u8).collect(); + let want: Vec<_> = SliceChunker::new(&data, 64, 256, MinCdcHash4::new()) + .map(|c| (c.offset(), c.to_vec())) + .collect(); + let mut chunker = ReadChunker::new( + OneByteReader(Cursor::new(&data)), + 64, + 256, + MinCdcHash4::new(), + ); + let mut got = Vec::new(); + while let Some(chunk) = chunker.next().unwrap() { + got.push((chunk.offset(), chunk.to_vec())); + } + assert_eq!(got, want); + } + + #[test] + fn read_chunker_compacts_its_buffer_without_changing_boundaries() { + let n = crate::MIN_BUFFER_SIZE + 16 * 1024; + let data: Vec = (0..n) + .map(|i| { + let x = (i as u64).wrapping_mul(0x9E37_79B9_7F4A_7C15); + (x ^ (x >> 29)) as u8 + }) + .collect(); + let want: Vec<_> = SliceChunker::new(&data, 64, 256, MinCdcHash4::new()) + .map(|c| (c.offset(), c.to_vec())) + .collect(); + + let mut chunker = ReadChunker::new(Cursor::new(&data), 64, 256, MinCdcHash4::new()); + let mut got = Vec::new(); + while let Some(chunk) = chunker.next().unwrap() { + got.push((chunk.offset(), chunk.to_vec())); + } + assert_eq!(got, want); + assert!(chunker.next().unwrap().is_none()); + } + + struct PartialThenError { + data: Vec, + offset: usize, + state: u8, + } + + impl Read for PartialThenError { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if self.state == 1 { + self.state = 2; + return Err(io::Error::other("transient reader failure")); + } + if self.offset == self.data.len() { + return Ok(0); + } + let limit = if self.state == 0 { 17 } else { buf.len() }; + let n = limit.min(buf.len()).min(self.data.len() - self.offset); + buf[..n].copy_from_slice(&self.data[self.offset..self.offset + n]); + self.offset += n; + self.state = 1.max(self.state); + Ok(n) + } + } + + #[test] + fn read_chunker_resumes_after_error_with_internal_progress() { + let data: Vec = (0..8192).map(|i| (i * 31) as u8).collect(); + let want: Vec<_> = SliceChunker::new(&data, 64, 256, MinCdcHash4::new()) + .map(|c| (c.offset(), c.to_vec())) + .collect(); + let reader = PartialThenError { + data, + offset: 0, + state: 0, + }; + let mut chunker = ReadChunker::new(reader, 64, 256, MinCdcHash4::new()); + + assert_eq!(chunker.next().unwrap_err().kind(), io::ErrorKind::Other); + let mut got = Vec::new(); + while let Some(chunk) = chunker.next().unwrap() { + got.push((chunk.offset(), chunk.to_vec())); + } + assert_eq!(got, want); + } + + struct ErrorAfterFirstRead { + data: Vec, + first: bool, + } + + impl Read for ErrorAfterFirstRead { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if self.first { + return Err(io::Error::other("reader failed")); + } + self.first = true; + let n = self.data.len().min(buf.len()); + buf[..n].copy_from_slice(&self.data[..n]); + Ok(n) + } + } + + #[test] + fn read_chunker_preserves_progress_before_error() { + let reader = ErrorAfterFirstRead { + data: vec![0; 17], + first: false, + }; + let mut chunks = ReadChunker::new(reader, 8, 16, MinCdcHash4::new()); + assert!(chunks.next().unwrap().is_some()); + assert_eq!(chunks.next().unwrap_err().kind(), io::ErrorKind::Other); + } + + #[test] + fn fallible_constructor_rejects_overflowing_configuration() { + let result = ReadChunker::try_new( + Cursor::new(Vec::::new()), + 0, + usize::MAX, + MinCdcHash4::new(), + ); + match result { + Err(e) => assert_eq!(e.kind(), io::ErrorKind::InvalidInput), + Ok(_) => panic!("overflowing configuration was accepted"), + } + } } diff --git a/src/x86_64.rs b/src/x86_64.rs index e407da1..6613328 100644 --- a/src/x86_64.rs +++ b/src/x86_64.rs @@ -287,16 +287,38 @@ pub fn argmin_u32_overlapping_hashed( multiplier: u32, addend: u32, ) -> usize { - const MIN_SIMD_LEN: usize = 16 + 3; - assert!(bytes.len() <= u32::MAX as usize); - - let within_four_offset = if SHOULD_HASH { - unsafe { ARGMIN_HASH_IMPL(bytes, multiplier, addend) } + let implementation = if SHOULD_HASH { + *ARGMIN_HASH_IMPL } else { - unsafe { ARGMIN_IMPL(bytes, multiplier, addend) } + *ARGMIN_IMPL }; + argmin_with_impl::(bytes, multiplier, addend, implementation) +} - if bytes.len() >= MIN_SIMD_LEN { +#[inline(always)] +fn argmin_with_impl( + bytes: &[u8], + multiplier: u32, + addend: u32, + implementation: ArgMinFn, +) -> usize { + assert!(bytes.len() <= u32::MAX as usize); + let within_four_offset = unsafe { implementation(bytes, multiplier, addend) }; + refine_argmin::(bytes, within_four_offset, multiplier, addend) +} + +#[inline(always)] +fn refine_argmin( + bytes: &[u8], + within_four_offset: usize, + multiplier: u32, + addend: u32, +) -> usize { + // SIMD implementations return the start of a four-position block and + // therefore always leave the seven bytes needed by the exact refinement. + // The scalar runtime-dispatch fallback already returns the exact argmin, + // which may be one of the final three windows. Do not refine that result. + if bytes.len() >= 16 + 3 && bytes.len() - within_four_offset >= 7 { let final_bump = scalar::argmin_u32_overlapping_hashed_four::( &bytes[within_four_offset..], multiplier, @@ -565,6 +587,27 @@ mod tests { use rand::distr::StandardUniform; use rand::prelude::*; + #[test] + fn scalar_dispatch_result_at_tail_needs_no_refinement() { + let mut bytes = vec![0xFF; 19]; + bytes[15..].fill(0); + + let got = + argmin_with_impl::(&bytes, 1, 0, scalar::argmin_u32_overlapping_hashed::); + assert_eq!(got, 15); + + let multiplier = 0x9E37_79B1; + let addend = 0x85EB_CA77; + let exact = scalar::argmin_u32_overlapping_hashed::(&bytes, multiplier, addend); + let got = argmin_with_impl::( + &bytes, + multiplier, + addend, + scalar::argmin_u32_overlapping_hashed::, + ); + assert_eq!(got, exact); + } + // Reconstruct the full argmin from a "within-four" SIMD result, mirroring // the dispatch logic in argmin_u32_overlapping_hashed. The *_impl functions // only locate the block containing the global minimizer within the next four @@ -572,19 +615,9 @@ mod tests { // bump, so the test must do the same before comparing to the scalar oracle. macro_rules! full_from_impl { ($within_four:expr, $hash:literal, $mul:expr, $add:expr, $bytes:expr) => {{ - const MIN_SIMD_LEN: usize = 16 + 3; let bytes: &[u8] = $bytes; let within_four = $within_four; - if bytes.len() >= MIN_SIMD_LEN { - within_four - + scalar::argmin_u32_overlapping_hashed_four::<$hash>( - &bytes[within_four..], - $mul, - $add, - ) - } else { - within_four - } + refine_argmin::<$hash>(bytes, within_four, $mul, $add) }}; } diff --git a/tests/invariants.rs b/tests/invariants.rs index 3c0fdcf..2e2c64c 100644 --- a/tests/invariants.rs +++ b/tests/invariants.rs @@ -14,8 +14,6 @@ //! 5. Reader/Slice agreement -- ReadChunker == SliceChunker under a hostile //! (1-byte) reader that stresses buffer refills. -#![allow(deprecated)] // MinCdc4 is deprecated but we still want to test it. - use std::io::{self, Read}; use mothcdc::mincdc::{Cdc, MinCdc4, MinCdcHash4, ReadChunker, SliceChunker}; @@ -112,11 +110,11 @@ fn oracle_chunks( fn slice_chunks(data: &[u8], min: usize, max: usize, mode: Mode) -> Vec<(usize, usize)> { fn collect(data: &[u8], min: usize, max: usize, cdc: C) -> Vec<(usize, usize)> { SliceChunker::new(data, min, max, cdc) - .map(|c| (c.offset(), c.len())) + .map(|c| (c.offset() as usize, c.len())) .collect() } match mode { - Mode::Plain => collect(data, min, max, MinCdc4::new()), + Mode::Plain => collect(data, min, max, MinCdc4), Mode::Hashed => collect(data, min, max, MinCdcHash4::with_params(HASH_M, HASH_A)), } } @@ -156,12 +154,12 @@ fn read_chunks( let mut chunker = ReadChunker::new(reader, min, max, cdc); let mut out = Vec::new(); while let Some(chunk) = chunker.next().unwrap() { - out.push((chunk.offset(), chunk.len())); + out.push((chunk.offset() as usize, chunk.len())); } out } match mode { - Mode::Plain => collect(data, min, max, MinCdc4::new(), step), + Mode::Plain => collect(data, min, max, MinCdc4, step), Mode::Hashed => collect( data, min, @@ -220,7 +218,10 @@ fn check_all(data: &[u8], min: usize, max: usize, mode: Mode) { data.len() ); - // Tier 5: ReadChunker must match SliceChunker, even with a choked reader. + // Tier 5: ReadChunker must match SliceChunker, even with a choked reader + // that fragments the internal buffer refills. This runs on every corpus and + // proptest case so the streaming buffer state machine is fuzzed against the + // in-memory chunker, not just checked on a handful of fixed inputs. for &step in &[1usize, 3, 7, 64, 4096] { let read = read_chunks(data, min, max, mode, step); assert_eq!( @@ -283,6 +284,19 @@ fn corpus_all_invariants() { } } +#[test] +fn hostile_reader_matches_slice_chunker() { + let mut data = (0..8192).map(|i| (i * 31) as u8).collect::>(); + data[2048..6144].fill(0); + for (min, max) in [(1usize, 4usize), (64, 256)] { + for mode in [Mode::Plain, Mode::Hashed] { + let slice = slice_chunks(&data, min, max, mode); + let read = read_chunks(&data, min, max, mode, 3); + assert_eq!(read, slice, "mode={mode:?}, min={min}, max={max}"); + } + } +} + // --------------------------------------------------------------------------- // Tier 3b: large random buffers to exercise the SIMD main loops + tail paths. // --------------------------------------------------------------------------- @@ -324,12 +338,12 @@ fn large_random_differential() { // --------------------------------------------------------------------------- proptest! { - #![proptest_config(ProptestConfig { cases: 512, ..ProptestConfig::default() })] + #![proptest_config(ProptestConfig { cases: 64, ..ProptestConfig::default() })] #[test] fn prop_invariants_and_oracle( data in proptest::collection::vec(any::(), 0..8192), - min in 4usize..400, + min in 1usize..400, extra in 0usize..400, ) { let max = min + extra; diff --git a/tests/roundtrip.rs b/tests/roundtrip.rs index 966ba62..59ef887 100644 --- a/tests/roundtrip.rs +++ b/tests/roundtrip.rs @@ -10,10 +10,10 @@ //! not catch a `dedup_key()` that isn't reconstruction-safe. use std::collections::HashMap; -use std::io::Cursor; +use std::io::{self, Cursor, Read}; use mothcdc::mincdc::{MinCdcHash4, ReadChunker, SliceChunker}; -use mothcdc::{MothChunker, MothReadChunker, Segment}; +use mothcdc::{MothChunker, MothReadChunker}; use proptest::prelude::*; fn fnv1a(b: &[u8]) -> u64 { @@ -49,8 +49,8 @@ fn assert_roundtrip(label: &str, data: &[u8], min: usize, max: usize) { // --- ingest --- let mut store: HashMap> = HashMap::new(); - let mut manifest: Vec<(u64, usize)> = Vec::new(); // (key_hash, logical_len) - let mut next_off = 0usize; + let mut manifest: Vec<(u64, u64)> = Vec::new(); // (key_hash, logical_len) + let mut next_off = 0u64; for seg in build(data, min, max) { assert_eq!(seg.offset(), next_off, "{tag}: non-contiguous offset"); assert!(!seg.is_empty(), "{tag}: empty segment"); @@ -62,14 +62,14 @@ fn assert_roundtrip(label: &str, data: &[u8], min: usize, max: usize) { manifest.push((h, seg.len())); next_off += seg.len(); } - assert_eq!(next_off, data.len(), "{tag}: coverage gap"); + assert_eq!(next_off, data.len() as u64, "{tag}: coverage gap"); // The caterpillar must represent exactly the same underlying chunks as plain // mincdc — it only groups them, never adds or drops any. let plain = SliceChunker::new(data, min, max, MinCdcHash4::new()).count(); - let expanded: usize = build(data, min, max).map(|s| s.chunk_count()).sum(); + let expanded: u64 = build(data, min, max).map(|s| s.chunk_count()).sum(); assert_eq!( - expanded, plain, + expanded, plain as u64, "{tag}: chunk_count {expanded} != plain {plain}" ); @@ -79,11 +79,11 @@ fn assert_roundtrip(label: &str, data: &[u8], min: usize, max: usize) { let bytes = store .get(h) .expect("manifest references a chunk not in the store"); - let mut w = 0; + let mut w = 0u64; while w < *len { - let take = bytes.len().min(*len - w); + let take = (bytes.len() as u64).min(*len - w) as usize; restored.extend_from_slice(&bytes[..take]); - w += take; + w += take as u64; } } assert_eq!(restored, data, "{tag}: store round-trip corrupted the data"); @@ -134,7 +134,7 @@ fn roundtrip_wide_and_narrow() { } /// A reader that returns at most `step` bytes per call, to exercise the streaming -/// chunker's buffer refill/shift logic (and its reader-dependent run splitting). +/// chunker's buffer refill/shift logic. struct ChokedReader<'a> { data: &'a [u8], pos: usize, @@ -149,11 +149,198 @@ impl std::io::Read for ChokedReader<'_> { } } +fn segment_layout(reader: R, min: usize, max: usize) -> Vec<(u64, u64, u64, Vec)> { + let mut chunker = MothReadChunker::new(reader, min, max); + let mut out = Vec::new(); + while let Some(segment) = chunker.next().unwrap() { + out.push(( + segment.offset(), + segment.len(), + segment.chunk_count(), + segment.dedup_key().to_vec(), + )); + } + out +} + +#[test] +fn streaming_grouping_is_independent_of_read_fragmentation() { + let data = vec![0u8; 1024 * 1024]; + let cursor = segment_layout(Cursor::new(&data), 2048, 14336); + let one_byte = segment_layout( + ChokedReader { + data: &data, + pos: 0, + step: 1, + }, + 2048, + 14336, + ); + assert_eq!(one_byte, cursor); + assert_eq!(cursor.len(), 1, "a zero run should be one record"); + assert!(cursor[0].2 >= 2); +} + +/// The in-memory slice caterpillar's grouping: (offset, len, count, dedup_key). +fn slice_grouping(data: &[u8], min: usize, max: usize) -> Vec<(u64, u64, u64, Vec)> { + MothChunker::new(data, min, max) + .map(|s| (s.offset(), s.len(), s.chunk_count(), s.dedup_key().to_vec())) + .collect() +} + +/// The streaming caterpillar must produce *exactly* the slice caterpillar's +/// segment grouping — same offsets, lengths, counts, and unit bytes — for every +/// corpus and every reader fragmentation. This is the strong form of the +/// "grouping is independent of `Read` fragmentation" claim: not merely that two +/// readers agree with each other, but that the stream matches the in-memory +/// reference, so a coalesced run is never split or merged differently by a +/// hostile reader. Guards the singleton-carry logic in `MothReadChunker::next`. +#[test] +fn streaming_grouping_matches_slice_across_fragmentations() { + let configs = [(64usize, 256usize), (2048, 14336), (2048, 2200), (16, 16)]; + let steps = [1usize, 2, 3, 7, 13, 64, 4096, 100_000]; + for (name, data) in corpora() { + for &(min, max) in &configs { + let want = slice_grouping(&data, min, max); + for &step in &steps { + let got = segment_layout( + ChokedReader { + data: &data, + pos: 0, + step, + }, + min, + max, + ); + assert_eq!( + got, + want, + "{name} (min={min} max={max} step={step}): streaming grouping \ + diverges from slice caterpillar ({} vs {} records)", + got.len(), + want.len() + ); + } + } + } +} + +/// A single run longer than the internal buffer must coalesce to the same +/// grouping as the slice caterpillar regardless of how the reader fragments the +/// refills that the run crosses. +#[test] +fn streaming_run_crossing_buffer_matches_slice() { + let mut data = xorshift(1, 100 * 1024); + data.extend_from_slice(&vec![0u8; 10 * 1024 * 1024]); + data.extend_from_slice(&xorshift(2, 100 * 1024)); + let (min, max) = (2048usize, 14336usize); + let want = slice_grouping(&data, min, max); + for &step in &[1usize, 4096, 1_000_000, 5_000_000] { + let got = segment_layout( + ChokedReader { + data: &data, + pos: 0, + step, + }, + min, + max, + ); + assert_eq!(got, want, "large-run grouping diverges at step={step}"); + } +} + +struct InterruptOnce { + inner: R, + interrupted: bool, +} + +impl Read for InterruptOnce { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if !self.interrupted { + self.interrupted = true; + return Err(io::ErrorKind::Interrupted.into()); + } + self.inner.read(buf) + } +} + +#[test] +fn moth_reader_retries_interrupted() { + let data = xorshift(88, 8192); + let reader = InterruptOnce { + inner: Cursor::new(&data), + interrupted: false, + }; + assert_stream_roundtrip("interrupted", reader, &data, 64, 256); +} + +struct ErrorOnceAfterData<'a> { + data: &'a [u8], + pos: usize, + errored: bool, +} + +impl Read for ErrorOnceAfterData<'_> { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + if self.pos >= 257 && !self.errored { + self.errored = true; + return Err(io::Error::other("transient failure")); + } + let remaining_before_error = 257usize.saturating_sub(self.pos); + let n = buf + .len() + .min(self.data.len() - self.pos) + .min(remaining_before_error.max(1)); + buf[..n].copy_from_slice(&self.data[self.pos..self.pos + n]); + self.pos += n; + Ok(n) + } +} + +#[test] +fn moth_reader_can_resume_after_error_with_internal_progress() { + let data = vec![0u8; 4096]; + let reader = ErrorOnceAfterData { + data: &data, + pos: 0, + errored: false, + }; + let mut chunker = MothReadChunker::new(reader, 64, 256); + assert_eq!(chunker.next().unwrap_err().kind(), io::ErrorKind::Other); + + let mut rebuilt = Vec::new(); + while let Some(segment) = chunker.next().unwrap() { + segment.reconstruct_into(&mut rebuilt); + } + assert_eq!(rebuilt, data); +} + +#[test] +fn moth_reader_reports_permanent_errors() { + let mut chunker = MothReadChunker::new(FailingReader, 64, 256); + assert_eq!(chunker.next().unwrap_err().kind(), io::ErrorKind::Other); +} + +#[test] +fn moth_fallible_constructor_rejects_overflowing_configuration() { + let result = MothReadChunker::try_new(Cursor::new(Vec::::new()), 0, usize::MAX); + match result { + Err(e) => assert_eq!(e.kind(), io::ErrorKind::InvalidInput), + Ok(_) => panic!("overflowing configuration was accepted"), + } +} + +struct FailingReader; + +impl Read for FailingReader { + fn read(&mut self, _buf: &mut [u8]) -> io::Result { + Err(io::Error::other("permanent failure")) + } +} + /// Drives the streaming caterpillar over `reader`, asserting contiguity, lossless /// reconstruction, exact coverage, and that it represents the same underlying -/// chunk count as plain mincdc. (It is NOT asserted to match the slice -/// caterpillar's segment grouping — the streaming version splits long runs at -/// buffer boundaries, which depends on the reader.) +/// chunk count as plain mincdc. fn assert_stream_roundtrip( tag: &str, reader: R, @@ -164,7 +351,7 @@ fn assert_stream_roundtrip( let cdc = MinCdcHash4::new(); let plain = SliceChunker::new(data, min, max, cdc).count(); let mut rc = MothReadChunker::with_cdc(reader, min, max, cdc); - let (mut next_off, mut chunks) = (0usize, 0usize); + let (mut next_off, mut chunks) = (0u64, 0u64); let mut rebuilt = Vec::with_capacity(data.len()); while let Some(s) = rc.next().unwrap() { assert_eq!(s.offset(), next_off, "{tag}: non-contiguous"); @@ -174,9 +361,9 @@ fn assert_stream_roundtrip( s.reconstruct_into(&mut rebuilt); } assert_eq!(rebuilt, data, "{tag}: stream round-trip corrupted data"); - assert_eq!(next_off, data.len(), "{tag}: coverage gap"); + assert_eq!(next_off, data.len() as u64, "{tag}: coverage gap"); assert_eq!( - chunks, plain, + chunks, plain as u64, "{tag}: chunk_count {chunks} != plain {plain}" ); } @@ -202,17 +389,9 @@ fn streaming_chunk_boundaries_match_plain_mincdc() { let mut expanded = Vec::new(); let mut sc = MothReadChunker::with_cdc(Cursor::new(&data), min, max, cdc); while let Some(s) = sc.next().unwrap() { - match s { - Segment::Solo(c) => expanded.push((c.offset(), c.len())), - Segment::Caterpillar { - offset, - unit, - count, - } => { - for i in 0..count { - expanded.push((offset + i * unit.len(), unit.len())); - } - }, + let unit_len = s.dedup_key().len(); + for i in 0..s.chunk_count() { + expanded.push((s.offset() + i * unit_len as u64, unit_len)); } } @@ -248,7 +427,7 @@ fn streaming_caterpillar_roundtrips() { } proptest! { - #![proptest_config(ProptestConfig { cases: 400, ..ProptestConfig::default() })] + #![proptest_config(ProptestConfig { cases: 64, ..ProptestConfig::default() })] /// Fuzz the full round-trip over random data and sizes. Shrinks any /// state-machine or reconstruction bug to a minimal case.