use std::io::{self, Read};
use crate::mincdc::{Cdc, Chunk, MinCdcHash4, SliceChunker};
use crate::simd;
pub(crate) fn packed_repeats(tail: &[u8], u: usize, max_size: usize) -> usize {
let n = tail.len();
if u == 0 || u >= n {
return 0;
}
let unit = &tail[..u];
let e = if unit[u - 1] == unit[0] && simd::byte_run_len(unit, unit[0]) == u {
u + simd::byte_run_len(&tail[u..], unit[0])
} else {
if 2 * u > n {
return 0;
}
let probe = simd::common_prefix_len(&tail[u..2 * u], unit);
if probe < u {
return 0;
}
2 * u + simd::common_prefix_len(&tail[2 * u..], &tail[u..n - u])
};
let lim = e.min(n - 1);
lim.saturating_sub(max_size) / u
}
#[derive(Debug)]
pub struct Segment<'a> {
kind: SegmentKind<'a>,
}
#[derive(Debug)]
enum SegmentKind<'a> {
Solo(Chunk<'a>),
Caterpillar {
offset: u64,
unit: &'a [u8],
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,
},
}
}
pub fn offset(&self) -> u64 {
match &self.kind {
SegmentKind::Solo(c) => c.offset(),
SegmentKind::Caterpillar { offset, .. } => *offset,
}
}
pub fn chunk_count(&self) -> u64 {
match &self.kind {
SegmentKind::Solo(_) => 1,
SegmentKind::Caterpillar { count, .. } => *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"),
}
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn is_caterpillar(&self) -> bool {
matches!(self.kind, SegmentKind::Caterpillar { .. })
}
pub fn dedup_key(&self) -> &[u8] {
match &self.kind {
SegmentKind::Solo(c) => c,
SegmentKind::Caterpillar { unit, .. } => unit,
}
}
pub fn reconstruct_into(&self, out: &mut Vec<u8>) {
let key = self.dedup_key();
let total = self.len();
let mut written = 0u64;
while written < total {
let take = (key.len() as u64).min(total - written) as usize;
out.extend_from_slice(&key[..take]);
written += take as u64;
}
}
}
pub struct MothChunker<'a, C = MinCdcHash4> {
data: &'a [u8],
inner: SliceChunker<'a, C>,
carry: Option<Chunk<'a>>,
}
impl<'a> MothChunker<'a, MinCdcHash4> {
pub fn new(bytes: &'a [u8], min_size: usize, max_size: usize) -> Self {
Self::with_cdc(bytes, min_size, max_size, MinCdcHash4::new())
}
}
impl<'a, C: Cdc> MothChunker<'a, C> {
pub fn with_cdc(bytes: &'a [u8], min_size: usize, max_size: usize, cdc: C) -> Self {
Self {
data: bytes,
inner: SliceChunker::new(bytes, min_size, max_size, cdc),
carry: None,
}
}
}
impl<'a, C: Cdc> Iterator for MothChunker<'a, C> {
type Item = Segment<'a>;
fn next(&mut self) -> Option<Segment<'a>> {
let first = self.carry.take().or_else(|| self.inner.next())?;
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];
let mut count = 1usize;
let repeats = packed_repeats(&self.data[start..], first_len, self.inner.max_size);
if repeats > 0 {
count += repeats;
self.inner.offset = start + count * first_len;
}
let pending: Option<Chunk<'a>> = loop {
match self.inner.next() {
Some(c) if &*c == unit => count += 1,
other => break other,
}
};
self.carry = pending;
if count >= 2 {
Some(Segment::caterpillar(start as u64, unit, count as u64))
} else {
Some(Segment::solo(first))
}
}
}
pub struct MothReadChunker<R, C = MinCdcHash4> {
min_size: usize,
max_size: usize,
cdc: C,
reader: R,
buf: Vec<u8>,
buf_offset: usize,
unread: usize,
stream_offset: u64,
eof: bool,
done: bool,
pending_len: Option<usize>,
carry_run: Option<(Vec<u8>, u64, u64)>,
emit_unit: Vec<u8>,
}
impl<R> MothReadChunker<R, MinCdcHash4> {
pub fn new(reader: R, min_size: usize, max_size: usize) -> Self {
Self::with_cdc(reader, min_size, max_size, MinCdcHash4::new())
}
pub fn try_new(reader: R, min_size: usize, max_size: usize) -> io::Result<Self> {
Self::try_with_cdc(reader, min_size, max_size, MinCdcHash4::new())
}
}
impl<R, C: Cdc> MothReadChunker<R, C> {
pub fn with_cdc(reader: R, min_size: usize, max_size: usize, cdc: C) -> Self {
Self::try_with_cdc(reader, min_size, max_size, cdc)
.expect("invalid MothReadChunker configuration or buffer allocation failed")
}
pub fn try_with_cdc(reader: R, min_size: usize, max_size: usize, cdc: C) -> io::Result<Self> {
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,
buf_offset: 0,
unread: 0,
stream_offset: 0,
eof: false,
done: false,
pending_len: None,
carry_run: None,
emit_unit: Vec::new(),
})
}
pub fn get_ref(&self) -> &R {
&self.reader
}
pub fn get_mut(&mut self) -> &mut R {
&mut self.reader
}
pub fn into_inner(self) -> R {
self.reader
}
fn chunk_len(&self, p: usize, avail: usize) -> Option<usize> {
crate::mincdc::next_chunk_len(
&self.buf[p..p + avail],
self.min_size,
self.max_size,
self.eof,
&self.cdc,
)
}
fn coalesce_from(&self, base: usize, u: usize) -> (usize, usize, Option<usize>) {
let repeats = packed_repeats(&self.buf[base..base + self.unread], u, self.max_size);
let mut run_len = u + repeats * u;
let mut count = 1 + repeats;
loop {
let cur = base + run_len;
let avail = self.unread - run_len;
match self.chunk_len(cur, avail) {
Some(nl) if nl == u && self.buf[cur..cur + nl] == self.buf[base..base + u] => {
count += 1;
run_len += nl;
},
other => return (run_len, count, other),
}
}
}
}
impl<R: Read, C: Cdc> MothReadChunker<R, C> {
fn ensure(&mut self) -> io::Result<()> {
let need = self.max_size + 1;
while !self.eof && self.unread < need {
if self.buf.len() - self.buf_offset < need {
self.buf
.copy_within(self.buf_offset..self.buf_offset + self.unread, 0);
self.buf_offset = 0;
}
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;
}
self.unread += n;
}
Ok(())
}
#[allow(clippy::should_implement_trait)]
pub fn next(&mut self) -> io::Result<Option<Segment<'_>>> {
if self.done {
return Ok(None);
}
loop {
self.ensure()?;
let base = self.buf_offset;
if self.unread == 0 {
self.done = true;
return Ok(self.carry_run.take().map(|(unit, off, count)| {
self.emit_unit = unit;
if count >= 2 {
Segment::caterpillar(off, &self.emit_unit, count)
} else {
Segment::solo(Chunk::new(&self.emit_unit, off))
}
}));
}
if let Some((unit, run_off, carried)) = self.carry_run.take() {
let u = unit.len();
match self.chunk_len(base, self.unread) {
Some(l) if l == u && self.buf[base..base + l] == unit[..] => {
let (run_len, more, pend) = self.coalesce_from(base, u);
self.buf_offset += run_len;
self.unread -= run_len;
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 {
self.carry_run = Some((unit, run_off, count));
continue;
}
self.pending_len = pend;
self.emit_unit = unit;
return Ok(Some(Segment::caterpillar(run_off, &self.emit_unit, count)));
},
boundary => {
self.pending_len = boundary;
self.emit_unit = unit;
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));
},
}
}
let base_stream = self.stream_offset;
let unit_len = match self.pending_len.take() {
Some(l) => l,
None => self
.chunk_len(base, self.unread)
.expect("buffer was ensured but the first chunk could not be decided"),
};
let (run_len, count, pend) = self.coalesce_from(base, unit_len);
self.buf_offset += run_len;
self.unread -= run_len;
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 {
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(base_stream, &self.buf[base..base + unit_len], count as u64)
} else {
Segment::solo(Chunk::new(&self.buf[base..base + unit_len], base_stream))
};
return Ok(Some(seg));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
const MIN: usize = 2 * 1024;
const MAX: usize = 14 * 1024;
fn xorshift(seed: u64, n: usize) -> Vec<u8> {
let mut s = seed | 1;
(0..n)
.map(|_| {
s ^= s << 13;
s ^= s >> 7;
s ^= s << 17;
(s >> 33) as u8
})
.collect()
}
fn measure(label: &str, data: &[u8], min: usize, max: usize) -> (usize, usize) {
let cdc = MinCdcHash4::new();
let plain = SliceChunker::new(data, min, max, cdc).count();
let mut records = 0usize;
let mut expanded = 0u64;
let mut next_off = 0u64;
let mut rebuilt: Vec<u8> = 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");
records += 1;
expanded += s.chunk_count();
s.reconstruct_into(&mut rebuilt);
next_off += s.len();
}
assert_eq!(rebuilt, data, "{label}: must reconstruct input exactly");
assert_eq!(
expanded, plain as u64,
"{label}: must represent same chunk count"
);
(plain, records)
}
#[test]
fn collapses_low_entropy_and_is_noop_on_random() {
let (plain, records) = measure("zero-fill 1MiB", &vec![0u8; 1024 * 1024], MIN, MAX);
assert!(records * 10 < plain, "zero-fill should collapse");
let data = xorshift(1234, 1024 * 1024);
let (plain, records) = measure("random 1MiB", &data, MIN, MAX);
assert!(
records as f64 >= plain as f64 * 0.98,
"must not coalesce random"
);
}
fn reference_segments(data: &[u8], min: usize, max: usize) -> Vec<(u64, Vec<u8>, u64)> {
let mut out: Vec<(u64, Vec<u8>, 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,
_ => out.push((c.offset(), c.to_vec(), 1)),
}
}
out
}
fn assert_matches_reference(label: &str, data: &[u8], min: usize, max: usize) {
let got: Vec<(u64, Vec<u8>, 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})");
}
#[test]
fn packed_repeats_claims_only_provable_chunks() {
let cdc = MinCdcHash4::new();
for period in [1usize, 2, 3, 5, 8, 13] {
let unit = xorshift(period as u64 + 40, period);
let mut base = Vec::new();
while base.len() < 2000 {
base.extend_from_slice(&unit);
}
base.truncate(2000);
for break_pos in 0..=base.len() {
let mut data = base.clone();
if break_pos < data.len() {
data[break_pos] ^= 0xA5;
}
for (min, max) in [(16usize, 40usize), (16, 16), (20, 24)] {
let u0 = crate::mincdc::next_chunk_len(&data, min, max, true, &cdc)
.expect("non-empty input at eof always chunks");
let claimed = packed_repeats(&data, u0, max);
for k in 1..=claimed {
let tail = &data[k * u0..];
let ctx = format!(
"period={period} break={break_pos} min={min} max={max} k={k}/{claimed} u0={u0}"
);
assert_eq!(
crate::mincdc::next_chunk_len(tail, min, max, true, &cdc),
Some(u0),
"claimed chunk diverges from slow path (eof): {ctx}"
);
assert_eq!(
crate::mincdc::next_chunk_len(tail, min, max, false, &cdc),
Some(u0),
"claimed chunk diverges from slow path (streaming): {ctx}"
);
assert_eq!(
&tail[..u0],
&data[..u0],
"claimed chunk bytes differ from unit: {ctx}"
);
}
}
}
}
}
#[test]
fn packed_fast_path_matches_reference() {
const N: usize = 256 * 1024;
let mut corpora: Vec<(String, Vec<u8>)> = vec![
("zeros".into(), vec![0u8; N]),
("const-0xAB".into(), vec![0xABu8; N]),
("random".into(), xorshift(7, N)),
("tiny-zeros".into(), vec![0u8; 100]),
("small-zeros".into(), vec![0u8; 20_000]),
];
for period in [1usize, 3, 4, 5, 100, 777, 2048, 2500, 5000] {
let unit = xorshift(period as u64, period);
let mut data = Vec::with_capacity(N + period);
while data.len() < N {
data.extend_from_slice(&unit);
}
data.truncate(N);
corpora.push((format!("periodic-{period}"), data));
}
{
let unit = xorshift(99, 777);
let mut data = Vec::new();
while data.len() < N {
data.extend_from_slice(&unit);
}
data.truncate(N);
for pos in [
N - 1,
N - 100,
N - MAX,
N - MAX - 1,
N - MAX + 1,
N / 2,
MAX,
MAX + 1,
MIN,
] {
let mut d = data.clone();
d[pos] ^= 0xFF;
corpora.push((format!("periodic-777-break@{pos}"), d));
}
}
{
let mut d = xorshift(3, N);
for b in &mut d[64 * 1024..192 * 1024] {
*b = 0;
}
corpora.push(("random+zero-hole".into(), d));
}
for (label, data) in &corpora {
for (min, max) in [(MIN, MAX), (2048, 2200), (64, 256), (16, 16), (4, 20)] {
assert_matches_reference(label, data, min, max);
}
}
}
proptest::proptest! {
#![proptest_config(proptest::prelude::ProptestConfig {
cases: 256, ..proptest::prelude::ProptestConfig::default()
})]
#[test]
fn prop_fast_path_matches_reference(
period in 1usize..600,
reps in 2usize..64,
seed in proptest::prelude::any::<u64>(),
min in 4usize..300,
extra in 0usize..300,
prefix in 0usize..500,
break_at in proptest::option::of(0.0f64..1.0),
) {
let unit = xorshift(seed | 1, period);
let mut data = xorshift(seed ^ 0xDEAD, prefix);
for _ in 0..reps {
data.extend_from_slice(&unit);
}
if let Some(f) = break_at {
let pos = ((data.len() as f64 - 1.0) * f) as usize;
data[pos] ^= 0xFF;
}
assert_matches_reference("prop", &data, min, min + extra);
}
}
#[test]
fn streaming_run_crossing_refills_is_one_record() {
use std::io::Cursor;
let n = 24 * 1024 * 1024;
let data = vec![0u8; n];
let mut rc = MothReadChunker::new(Cursor::new(&data), MIN, MAX);
let mut records = 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 as u64);
let plain = SliceChunker::new(&data, MIN, MAX, MinCdcHash4::new()).count();
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)"
);
let mut data = xorshift(11, 64 * 1024);
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, 0u64, 0u64);
let mut rebuilt = Vec::with_capacity(data.len());
while let Some(s) = rc.next().unwrap() {
records += 1;
chunks += s.chunk_count();
covered += s.len();
s.reconstruct_into(&mut rebuilt);
}
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 as u64);
assert!(
records < 120,
"zero run appears to split at refills (got {records} records)"
);
}
struct ReadOnceThenEof {
data: Vec<u8>,
sent: bool,
}
impl std::io::Read for ReadOnceThenEof {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
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 = [
vec![0xA5; 64],
[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() {
let period = xorshift(42, 777);
let mut data = Vec::new();
while data.len() < 4 * 1024 * 1024 {
data.extend_from_slice(&period);
}
let (plain, records) = measure("period777 wide", &data, MIN, MAX);
assert!(
records * 20 < plain,
"tier-1 collapses self-aligned periodic data"
);
}
}