#![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,
);
}