#![warn(missing_docs)]
use std::io::{self, Read};
use std::iter::FusedIterator;
use std::ops::Deref;
const DEFAULT_MULTIPLIER: u32 = 0x915f77f5;
const DEFAULT_ADDEND: u32 = 0x34636463;
const MIN_BUFFER_SIZE: usize = 1024 * 1024 * 4;
pub(crate) mod scalar;
#[cfg(all(target_arch = "aarch64", target_feature = "neon"))]
#[path = "neon.rs"]
mod simd;
#[cfg(target_arch = "x86_64")]
#[path = "x86_64.rs"]
mod simd;
#[cfg(not(any(
target_arch = "x86_64",
all(target_arch = "aarch64", target_feature = "neon")
)))]
use scalar as simd;
pub trait Cdc {
fn window_size(&self) -> usize;
fn best_splitpoint(&self, bytes: &[u8]) -> usize;
}
#[non_exhaustive]
#[derive(Copy, Clone, Default, Debug)]
pub struct MinCdc4;
impl MinCdc4 {
#[deprecated = "Unless you have a specific reason to use MinCdc4, prefer MinCdcHash4 instead. MinCdc4 is less robust to certain input patterns, and can easily create skewed chunk sizes. It is not recommended for general use, but is kept for academic purposes."]
pub const fn new() -> Self {
Self
}
}
impl Cdc for MinCdc4 {
#[inline(always)]
fn window_size(&self) -> usize {
4
}
#[inline(always)]
fn best_splitpoint(&self, bytes: &[u8]) -> usize {
if bytes.len() < 4 {
return bytes.len();
}
4 + simd::argmin_u32_overlapping_hashed::<false>(bytes, 1, 0)
}
}
#[derive(Copy, Clone, Debug)]
pub struct MinCdcHash4 {
multiplier: u32,
addend: u32,
}
impl Default for MinCdcHash4 {
fn default() -> Self {
Self::new()
}
}
impl MinCdcHash4 {
pub const fn new() -> Self {
Self::with_params(DEFAULT_MULTIPLIER, DEFAULT_ADDEND)
}
pub const fn with_params(multiplier: u32, addend: u32) -> Self {
assert!(multiplier % 2 == 1, "the MinCDCHash multiplier must be odd");
Self { multiplier, addend }
}
}
impl Cdc for MinCdcHash4 {
#[inline(always)]
fn window_size(&self) -> usize {
4
}
#[inline(always)]
fn best_splitpoint(&self, bytes: &[u8]) -> usize {
if bytes.len() < 4 {
return bytes.len();
}
4 + simd::argmin_u32_overlapping_hashed::<true>(bytes, self.multiplier, self.addend)
}
}
#[derive(Clone)]
pub struct SliceChunker<'a, C> {
min_size: usize,
max_size: usize,
cdc: C,
bytes: &'a [u8],
offset: usize,
}
impl<'a, C> SliceChunker<'a, C> {
pub const fn new(bytes: &'a [u8], min_size: usize, max_size: usize, cdc: C) -> Self {
assert!(min_size <= max_size && max_size > 0);
Self {
min_size,
max_size,
cdc,
bytes,
offset: 0,
}
}
}
impl<'a, C: Cdc> Iterator for SliceChunker<'a, C> {
type Item = Chunk<'a>;
fn next(&mut self) -> Option<Self::Item> {
let bytes_left = self.bytes.len() - self.offset;
if bytes_left == 0 {
return None;
}
if bytes_left <= self.min_size {
let ret = Chunk::new(&self.bytes[self.offset..], self.offset);
self.offset = self.bytes.len();
return Some(ret);
}
let start_search_offset =
self.offset + self.min_size.saturating_sub(self.cdc.window_size());
let stop_search_offset = self.offset + self.max_size;
let search = &self.bytes[start_search_offset..stop_search_offset.min(self.bytes.len())];
let ideal_split = self.cdc.best_splitpoint(search);
let splitpoint = start_search_offset + ideal_split;
let ret = Chunk::new(&self.bytes[self.offset..splitpoint], self.offset);
self.offset = splitpoint;
Some(ret)
}
}
impl<'a, C: Cdc> FusedIterator for SliceChunker<'a, C> {}
#[derive(Clone)]
pub struct ReadChunker<R, C> {
min_size: usize,
max_size: usize,
cdc: C,
reader: R,
buf: Vec<u8>,
buf_offset: usize,
unread_bytes_in_buf: usize,
stream_offset: usize,
}
impl<R, C: Cdc> ReadChunker<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 bytes_needed_for_decision = max_size + 1;
let buf_size = MIN_BUFFER_SIZE + bytes_needed_for_decision + min_size * 4;
Self {
min_size,
max_size,
cdc,
reader,
buf: vec![0; buf_size],
buf_offset: 0,
unread_bytes_in_buf: 0,
stream_offset: 0,
}
}
}
impl<R: Read, C: Cdc> ReadChunker<R, C> {
#[allow(clippy::should_implement_trait)]
pub fn next(&mut self) -> io::Result<Option<Chunk<'_>>> {
if self.stream_offset == usize::MAX {
return Ok(None);
}
let bytes_needed_for_decision = self.max_size + 1;
while self.unread_bytes_in_buf < bytes_needed_for_decision {
if self.buf.len() - self.buf_offset < bytes_needed_for_decision {
self.buf.copy_within(
self.buf_offset..self.buf_offset + self.unread_bytes_in_buf,
0,
);
self.buf_offset = 0;
}
let bytes_read = self
.reader
.read(&mut self.buf[self.buf_offset + self.unread_bytes_in_buf..])?;
if bytes_read == 0 {
break;
}
self.unread_bytes_in_buf += bytes_read;
}
if self.unread_bytes_in_buf <= self.min_size {
let ret = Chunk::new(
&self.buf[self.buf_offset..self.buf_offset + self.unread_bytes_in_buf],
self.stream_offset,
);
self.stream_offset = usize::MAX;
return if ret.bytes.is_empty() {
Ok(None)
} else {
Ok(Some(ret))
};
}
let start_search_offset = self.min_size.saturating_sub(self.cdc.window_size());
let stop_search_offset = self.max_size;
let search = &self.buf[self.buf_offset + start_search_offset
..self.buf_offset + stop_search_offset.min(self.unread_bytes_in_buf)];
let ideal_split = self.cdc.best_splitpoint(search);
let splitpoint = start_search_offset + ideal_split;
let ret = Chunk::new(
&self.buf[self.buf_offset..self.buf_offset + splitpoint],
self.stream_offset,
);
let len = ret.bytes.len();
self.stream_offset += len;
self.buf_offset += len;
self.unread_bytes_in_buf -= len;
Ok(Some(ret))
}
}
#[derive(Copy, Clone, Debug)]
pub struct Chunk<'a> {
bytes: &'a [u8],
offset: usize,
}
impl<'a> Chunk<'a> {
pub const fn new(bytes: &'a [u8], offset: usize) -> Self {
Self { bytes, offset }
}
pub const fn offset(&self) -> usize {
self.offset
}
}
impl<'a> Deref for Chunk<'a> {
type Target = [u8];
fn deref(&self) -> &Self::Target {
self.bytes
}
}
#[cfg(test)]
mod test {
use std::io::Cursor;
use rand::distr::StandardUniform;
use rand::prelude::*;
use crate::{
DEFAULT_ADDEND, DEFAULT_MULTIPLIER, MinCdc4, MinCdcHash4, ReadChunker, SliceChunker,
scalar, simd,
};
#[test]
fn test_argmin_overlapped() {
for size in 0..4096 {
let rng = SmallRng::seed_from_u64(size);
let bytes: Vec<u8> = rng
.sample_iter(StandardUniform)
.take(size as usize)
.collect();
assert_eq!(
simd::argmin_u32_overlapping_hashed::<false>(&bytes, 1, 0),
scalar::argmin_u32_overlapping_hashed::<false>(&bytes, 1, 0)
);
assert_eq!(
simd::argmin_u32_overlapping_hashed::<true>(
&bytes,
DEFAULT_MULTIPLIER,
DEFAULT_ADDEND
),
scalar::argmin_u32_overlapping_hashed::<true>(
&bytes,
DEFAULT_MULTIPLIER,
DEFAULT_ADDEND
)
);
}
}
#[test]
fn test_read_slice_equiv() {
let bounds = [1, 2, 3, 4, 6, 8, 15, 27, 62, 90, 120, 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<u8> = rng
.sample_iter(StandardUniform)
.take(size as usize)
.collect();
let reader = Cursor::new(&bytes);
let mut read_chunker = ReadChunker::new(reader, *min_size, *max_size, MinCdc4);
let slice_chunker = SliceChunker::new(&bytes, *min_size, *max_size, MinCdc4);
for slice_chunk in slice_chunker {
let read_chunk = read_chunker.next().unwrap().unwrap();
assert_eq!(slice_chunk.offset(), read_chunk.offset());
assert_eq!(&slice_chunk[..], &read_chunk[..]);
}
assert!(read_chunker.next().unwrap().is_none());
let reader = Cursor::new(&bytes);
let mut read_chunker =
ReadChunker::new(reader, *min_size, *max_size, MinCdcHash4::new());
let slice_chunker =
SliceChunker::new(&bytes, *min_size, *max_size, MinCdcHash4::new());
for slice_chunk in slice_chunker {
let read_chunk = read_chunker.next().unwrap().unwrap();
assert_eq!(slice_chunk.offset(), read_chunk.offset());
assert_eq!(&slice_chunk[..], &read_chunk[..]);
}
assert!(read_chunker.next().unwrap().is_none());
}
}
}
}
}