#![cfg(all(feature = "lzma2", feature = "parallel"))]
use std::io::{Cursor, Read, Write};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use zesven::read::Archive;
use zesven::write::{EntryMeta, WriteOptions, Writer};
use zesven::{ArchivePath, Threads};
const STREAMING_THRESHOLD: usize = 64 * 1024 * 1024;
fn payload(len: usize) -> Vec<u8> {
let mut data = Vec::with_capacity(len);
let mut state = 0x243F_6A88_85A3_08D3u64;
while data.len() < len {
state ^= state << 13;
state ^= state >> 7;
state ^= state << 17;
data.extend_from_slice(&state.to_le_bytes());
if state % 8 == 0 {
data.extend_from_slice(b"........................");
}
}
data.truncate(len);
data
}
fn write_archive(data: &[u8], threads: Threads) -> Vec<u8> {
let options = WriteOptions::new()
.level(1)
.expect("level")
.threads(threads);
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options);
writer
.add_bytes(ArchivePath::new("big.bin").expect("path"), data)
.expect("adds");
let (_result, cursor) = writer.finish_into_inner().expect("finishes");
cursor.into_inner()
}
fn read_back(archive: &[u8]) -> Vec<u8> {
let mut archive = Archive::open(Cursor::new(archive.to_vec())).expect("opens");
archive.extract_to_vec("big.bin").expect("extracts")
}
#[test]
fn test_large_entries_round_trip_around_block_edges() {
let mib = 1024 * 1024;
for len in [
STREAMING_THRESHOLD - 1,
STREAMING_THRESHOLD,
STREAMING_THRESHOLD + 1,
STREAMING_THRESHOLD + 4 * mib,
STREAMING_THRESHOLD + 5 * mib - 1,
STREAMING_THRESHOLD + 5 * mib,
STREAMING_THRESHOLD + 5 * mib + 1,
] {
let data = payload(len);
let archive = write_archive(&data, Threads::count_or_single(4));
assert_eq!(read_back(&archive), data, "len={len}");
}
}
#[test]
fn test_the_declared_size_does_not_reach_the_archive() {
let data = payload(STREAMING_THRESHOLD + 6 * 1024 * 1024);
let write = |declared: u64| {
let options = WriteOptions::new()
.level(1)
.expect("level")
.threads(Threads::count_or_single(4));
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options);
writer
.add_stream(
ArchivePath::new("big.bin").expect("path"),
&mut Cursor::new(data.clone()),
EntryMeta::file(declared),
)
.expect("adds");
let (_result, cursor) = writer.finish_into_inner().expect("finishes");
cursor.into_inner()
};
let truthful = write(data.len() as u64);
for declared in [0, 1, (STREAMING_THRESHOLD - 1) as u64, 4 << 30] {
assert_eq!(
write(declared),
truthful,
"a declared size of {declared} changed the archive",
);
}
assert_eq!(read_back(&truthful), data);
}
#[test]
fn test_output_does_not_depend_on_worker_count() {
let data = payload(STREAMING_THRESHOLD + 3 * 1024 * 1024);
let reference = write_archive(&data, Threads::count_or_single(2));
for threads in [4, 8, 16] {
assert_eq!(
write_archive(&data, Threads::count_or_single(threads)),
reference,
"{threads} workers changed the archive",
);
}
assert_eq!(read_back(&reference), data);
}
#[test]
fn test_single_thread_round_trips() {
let data = payload(STREAMING_THRESHOLD + 1024 * 1024);
let archive = write_archive(&data, Threads::Single);
assert_eq!(read_back(&archive), data);
}
struct CountingSink {
inner: Cursor<Vec<u8>>,
written: Arc<AtomicU64>,
}
impl Write for CountingSink {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let n = self.inner.write(buf)?;
self.written.fetch_add(n as u64, Ordering::Release);
Ok(n)
}
fn flush(&mut self) -> std::io::Result<()> {
self.inner.flush()
}
}
impl std::io::Seek for CountingSink {
fn seek(&mut self, pos: std::io::SeekFrom) -> std::io::Result<u64> {
self.inner.seek(pos)
}
}
struct GapReader<'a> {
data: &'a [u8],
position: usize,
written: Arc<AtomicU64>,
largest_gap: u64,
}
impl Read for GapReader<'_> {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
let n = buf.len().min(self.data.len() - self.position);
buf[..n].copy_from_slice(&self.data[self.position..self.position + n]);
self.position += n;
let gap = (self.position as u64).saturating_sub(self.written.load(Ordering::Acquire));
self.largest_gap = self.largest_gap.max(gap);
Ok(n)
}
}
fn incompressible(len: usize) -> Vec<u8> {
let mut data = Vec::with_capacity(len);
let mut state = 0x2545_F491_4F6C_DD1Du64;
while data.len() < len {
state ^= state << 13;
state ^= state >> 7;
state ^= state << 17;
data.extend_from_slice(&state.to_le_bytes());
}
data.truncate(len);
data
}
fn largest_backlog(data: &[u8]) -> u64 {
let written = Arc::new(AtomicU64::new(0));
let sink = CountingSink {
inner: Cursor::new(Vec::new()),
written: Arc::clone(&written),
};
let options = WriteOptions::new()
.level(1)
.expect("level")
.threads(Threads::count_or_single(4));
let mut writer = Writer::create(sink).expect("writer").options(options);
let mut source = GapReader {
data,
position: 0,
written: Arc::clone(&written),
largest_gap: 0,
};
writer
.add_stream(
ArchivePath::new("big.bin").expect("path"),
&mut source,
EntryMeta::file(data.len() as u64),
)
.expect("adds");
let gap = source.largest_gap;
let _ = writer.finish_into_inner().expect("finishes");
gap
}
#[test]
fn test_memory_follows_the_window_not_the_entry() {
let small = incompressible(STREAMING_THRESHOLD + 8 * 1024 * 1024);
let large = incompressible(STREAMING_THRESHOLD + 136 * 1024 * 1024);
let small_backlog = largest_backlog(&small);
let large_backlog = largest_backlog(&large);
let slack = 16 * 1024 * 1024;
assert!(
large_backlog <= small_backlog + slack,
"backlog went from {small_backlog} bytes on a {} byte entry to \
{large_backlog} on a {} byte one: it is following the entry rather \
than the window",
small.len(),
large.len(),
);
let bound = (STREAMING_THRESHOLD + 64 * 1024 * 1024) as u64;
for (backlog, data) in [(small_backlog, &small), (large_backlog, &large)] {
assert!(
backlog < bound,
"backlog {backlog} on a {} byte entry, over a bound of {bound}",
data.len(),
);
}
}
#[test]
fn test_a_size_declared_too_small_still_bounds_memory() {
let data = incompressible(STREAMING_THRESHOLD + 136 * 1024 * 1024);
let written = Arc::new(AtomicU64::new(0));
let sink = CountingSink {
inner: Cursor::new(Vec::new()),
written: Arc::clone(&written),
};
let options = WriteOptions::new()
.level(1)
.expect("level")
.threads(Threads::count_or_single(4));
let mut writer = Writer::create(sink).expect("writer").options(options);
let mut source = GapReader {
data: &data,
position: 0,
written: Arc::clone(&written),
largest_gap: 0,
};
writer
.add_stream(
ArchivePath::new("big.bin").expect("path"),
&mut source,
EntryMeta::file(data.len() as u64 / 10),
)
.expect("adds");
let backlog = source.largest_gap;
let _ = writer.finish_into_inner().expect("finishes");
assert!(
backlog < (STREAMING_THRESHOLD + 64 * 1024 * 1024) as u64,
"backlog {backlog} on a {} byte entry declared as {}: the declared size \
decided how the entry was held",
data.len(),
data.len() / 10,
);
}
#[test]
fn test_a_batch_compressed_alongside_does_not_reach_the_archive() {
use zesven::MemoryLimit;
let small: Vec<Vec<u8>> = (0..4).map(|i| payload(1 << 20 << (i % 3))).collect();
let large = payload(STREAMING_THRESHOLD + (4 << 20));
let build = |limit: MemoryLimit, threads: Threads| {
let options = WriteOptions::new()
.level(1)
.expect("level")
.memory_limit(limit)
.threads(threads);
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options);
for round in 0..2 {
for (i, data) in small.iter().enumerate() {
writer
.add_bytes(
ArchivePath::new(&format!("{round}-{i}.bin")).expect("path"),
data,
)
.expect("adds");
}
writer
.add_bytes(
ArchivePath::new(&format!("{round}-big.bin")).expect("path"),
&large,
)
.expect("adds");
}
writer.finish_into_inner().expect("finishes").1.into_inner()
};
let cramped = build(
MemoryLimit::bytes_or_auto(16 << 20),
Threads::count_or_single(2),
);
let roomy = build(
MemoryLimit::bytes_or_auto(4 << 30),
Threads::count_or_single(8),
);
assert_eq!(
cramped, roomy,
"when the batch ahead was collected reached the bytes",
);
let mut archive = Archive::open(Cursor::new(cramped)).expect("opens");
let names: Vec<String> = archive
.entries()
.iter()
.map(|entry| entry.path.to_string())
.collect();
assert_eq!(
names,
vec![
"0-0.bin",
"0-1.bin",
"0-2.bin",
"0-3.bin",
"0-big.bin",
"1-0.bin",
"1-1.bin",
"1-2.bin",
"1-3.bin",
"1-big.bin",
],
);
assert_eq!(
archive.extract_to_vec("0-2.bin").expect("extracts"),
small[2]
);
assert_eq!(
archive.extract_to_vec("1-big.bin").expect("extracts"),
large
);
let single = build(MemoryLimit::bytes_or_auto(4 << 30), Threads::Single);
let mut archive = Archive::open(Cursor::new(single)).expect("opens");
assert_eq!(
archive.extract_to_vec("0-2.bin").expect("extracts"),
small[2]
);
assert_eq!(
archive.extract_to_vec("1-big.bin").expect("extracts"),
large
);
}
#[test]
fn test_a_batch_ahead_does_not_hold_the_entry() {
use zesven::MemoryLimit;
let budget = MemoryLimit::bytes_or_auto(1 << 30);
let batch: Vec<Vec<u8>> = (0..3).map(|_| incompressible(2 << 20)).collect();
let backlog_of = |entry: &[u8]| -> u64 {
let written = Arc::new(AtomicU64::new(0));
let sink = CountingSink {
inner: Cursor::new(Vec::new()),
written: Arc::clone(&written),
};
let options = WriteOptions::new()
.level(1)
.expect("level")
.memory_limit(budget)
.threads(Threads::count_or_single(4));
let mut writer = Writer::create(sink).expect("writer").options(options);
for (i, data) in batch.iter().enumerate() {
writer
.add_bytes(ArchivePath::new(&format!("s{i}.bin")).expect("path"), data)
.expect("adds");
}
let mut source = GapReader {
data: entry,
position: 0,
written: Arc::clone(&written),
largest_gap: 0,
};
writer
.add_stream(
ArchivePath::new("big.bin").expect("path"),
&mut source,
EntryMeta::file(entry.len() as u64),
)
.expect("adds");
let gap = source.largest_gap;
let _ = writer.finish_into_inner().expect("finishes");
gap
};
let small = incompressible(STREAMING_THRESHOLD + 8 * 1024 * 1024);
let large = incompressible(STREAMING_THRESHOLD + 136 * 1024 * 1024);
let small_backlog = backlog_of(&small);
let large_backlog = backlog_of(&large);
let slack = 16 * 1024 * 1024;
assert!(
large_backlog <= small_backlog + slack,
"backlog went from {small_backlog} on a {} byte entry to \
{large_backlog} on a {} byte one, behind the same batch: the batch \
ahead turned the entry into a buffer",
small.len(),
large.len(),
);
}
struct FirstWriteWatcher {
inner: Cursor<Vec<u8>>,
read_so_far: Arc<AtomicU64>,
read_at_first_write: Arc<AtomicU64>,
seen_a_write: bool,
}
impl Write for FirstWriteWatcher {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let read = self.read_so_far.load(Ordering::Acquire);
if !self.seen_a_write && !buf.is_empty() && read > 0 {
self.seen_a_write = true;
self.read_at_first_write.store(read, Ordering::Release);
}
self.inner.write(buf)
}
fn flush(&mut self) -> std::io::Result<()> {
self.inner.flush()
}
}
impl std::io::Seek for FirstWriteWatcher {
fn seek(&mut self, pos: std::io::SeekFrom) -> std::io::Result<u64> {
self.inner.seek(pos)
}
}
struct WatchedReader<'a> {
data: &'a [u8],
position: usize,
read_so_far: Arc<AtomicU64>,
}
impl Read for WatchedReader<'_> {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
let taken = (self.data.len() - self.position).min(buf.len());
buf[..taken].copy_from_slice(&self.data[self.position..self.position + taken]);
self.position += taken;
self.read_so_far
.store(self.position as u64, Ordering::Release);
Ok(taken)
}
}
#[test]
fn test_the_batch_is_compressed_while_the_entry_is_read() {
use zesven::MemoryLimit;
let read_so_far = Arc::new(AtomicU64::new(0));
let read_at_first_write = Arc::new(AtomicU64::new(0));
let options = WriteOptions::new()
.level(1)
.expect("level")
.memory_limit(MemoryLimit::bytes_or_auto(4 << 30))
.threads(Threads::count_or_single(4));
let mut writer = Writer::create(FirstWriteWatcher {
inner: Cursor::new(Vec::new()),
read_so_far: Arc::clone(&read_so_far),
read_at_first_write: Arc::clone(&read_at_first_write),
seen_a_write: false,
})
.expect("writer")
.options(options);
let waiting = incompressible(STREAMING_THRESHOLD - (4 << 20));
for name in ["batch-0.bin", "batch-1.bin"] {
writer
.add_bytes(ArchivePath::new(name).expect("path"), &waiting)
.expect("adds");
}
let large = incompressible(110 << 20);
writer
.add_stream(
ArchivePath::new("big.bin").expect("path"),
&mut WatchedReader {
data: &large,
position: 0,
read_so_far: Arc::clone(&read_so_far),
},
EntryMeta::file(large.len() as u64),
)
.expect("adds");
let archive = writer
.finish_into_inner()
.expect("finishes")
.1
.inner
.into_inner();
let at_first_write = read_at_first_write.load(Ordering::Acquire);
let prefix = STREAMING_THRESHOLD as u64;
assert!(
at_first_write > prefix,
"the first bytes reached the sink after {at_first_write} bytes had been \
read, against a {prefix} byte prefix that is read before the entry is \
even recognised: the batch was compressed before the entry rather than \
alongside it",
);
let mut archive = Archive::open(Cursor::new(archive)).expect("opens");
assert_eq!(archive.extract_to_vec("big.bin").expect("extracts"), large);
}