use std::io::{self, Read};
use crate::{Cdc, Chunk, SliceChunker, 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 enum Segment<'a> {
Solo(Chunk<'a>),
Caterpillar {
offset: usize,
unit: &'a [u8],
count: usize,
},
}
impl<'a> Segment<'a> {
pub fn offset(&self) -> usize {
match self {
Segment::Solo(c) => c.offset(),
Segment::Caterpillar { offset, .. } => *offset,
}
}
pub fn chunk_count(&self) -> usize {
match self {
Segment::Solo(_) => 1,
Segment::Caterpillar { count, .. } => *count,
}
}
pub fn len(&self) -> usize {
match self {
Segment::Solo(c) => c.len(),
Segment::Caterpillar { unit, count, .. } => unit.len() * count,
}
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn dedup_key(&self) -> &[u8] {
match self {
Segment::Solo(c) => c,
Segment::Caterpillar { unit, .. } => unit,
}
}
pub fn reconstruct_into(&self, out: &mut Vec<u8>) {
let key = self.dedup_key();
let total = self.len();
let mut written = 0;
while written < total {
let take = key.len().min(total - written);
out.extend_from_slice(&key[..take]);
written += take;
}
}
}
pub struct CaterpillarChunker<'a, C> {
data: &'a [u8],
inner: SliceChunker<'a, C>,
carry: Option<Chunk<'a>>,
}
impl<'a, C: Cdc> CaterpillarChunker<'a, C> {
pub fn new(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 CaterpillarChunker<'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 = first.offset();
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 {
offset: start,
unit,
count,
})
} else {
Some(Segment::Solo(first))
}
}
}
pub struct CaterpillarReadChunker<R, C> {
min_size: usize,
max_size: usize,
cdc: C,
reader: R,
buf: Vec<u8>,
buf_offset: usize,
unread: usize,
stream_offset: usize,
eof: bool,
done: bool,
pending_len: Option<usize>,
carry_run: Option<(Vec<u8>, usize, usize)>,
emit_unit: Vec<u8>,
}
impl<R, C: Cdc> CaterpillarReadChunker<R, C> {
pub fn new(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 {
min_size,
max_size,
cdc,
reader,
buf: vec![0; buf_size],
buf_offset: 0,
unread: 0,
stream_offset: 0,
eof: false,
done: false,
pending_len: None,
carry_run: None,
emit_unit: Vec::new(),
}
}
fn chunk_len(&self, p: usize, avail: usize) -> Option<usize> {
crate::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> CaterpillarReadChunker<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 = self
.reader
.read(&mut self.buf[self.buf_offset + self.unread..])?;
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;
Segment::Caterpillar {
offset: off,
unit: &self.emit_unit,
count,
}
}));
}
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 += run_len;
let count = carried + more;
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 {
offset: run_off,
unit: &self.emit_unit,
count,
}));
},
boundary => {
self.pending_len = boundary;
self.emit_unit = unit;
return Ok(Some(Segment::Caterpillar {
offset: run_off,
unit: &self.emit_unit,
count: carried,
}));
},
}
}
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 += run_len;
if pend.is_none() && !self.eof && count >= 2 {
self.carry_run =
Some((self.buf[base..base + unit_len].to_vec(), base_stream, count));
continue;
}
self.pending_len = pend;
let seg = if count >= 2 {
Segment::Caterpillar {
offset: base_stream,
unit: &self.buf[base..base + unit_len],
count,
}
} else {
Segment::Solo(Chunk::new(&self.buf[base..base + unit_len], base_stream))
};
return Ok(Some(seg));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::MinCdcHash4;
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 = 0usize;
let mut next_off = 0usize;
let mut rebuilt: Vec<u8> = Vec::with_capacity(data.len());
for s in CaterpillarChunker::new(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, "{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<(usize, Vec<u8>, usize)> {
let mut out: Vec<(usize, Vec<u8>, usize)> = 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<(usize, Vec<u8>, usize)> =
CaterpillarChunker::new(data, min, max, MinCdcHash4::new())
.map(|s| match s {
Segment::Solo(c) => (c.offset(), c.to_vec(), 1),
Segment::Caterpillar {
offset,
unit,
count,
} => (offset, unit.to_vec(), 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::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::next_chunk_len(tail, min, max, true, &cdc),
Some(u0),
"claimed chunk diverges from slow path (eof): {ctx}"
);
assert_eq!(
crate::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 = CaterpillarReadChunker::new(Cursor::new(&data), MIN, MAX, MinCdcHash4::new());
let mut records = 0usize;
let mut chunks = 0usize;
let mut covered = 0usize;
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);
let plain = SliceChunker::new(&data, MIN, MAX, MinCdcHash4::new()).count();
assert_eq!(chunks, plain, "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 = CaterpillarReadChunker::new(Cursor::new(&data), MIN, MAX, MinCdcHash4::new());
let (mut records, mut chunks, mut covered) = (0usize, 0usize, 0usize);
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());
assert_eq!(rebuilt, data, "must reconstruct exactly");
let plain = SliceChunker::new(&data, MIN, MAX, MinCdcHash4::new()).count();
assert_eq!(chunks, plain);
assert!(
records < 120,
"zero run appears to split at refills (got {records} records)"
);
}
#[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"
);
}
}