#![cfg(all(feature = "lzma2", feature = "parallel"))]
use std::io::Cursor;
use zesven::read::Archive;
use zesven::write::{EntryMeta, WriteOptions, Writer};
use zesven::{ArchivePath, MemoryLimit, Threads};
const STREAMING_THRESHOLD: usize = 64 * 1024 * 1024;
const SIGNATURE_HEADER_SIZE: usize = 32;
fn payload(seed: u64, len: usize) -> Vec<u8> {
let mut data = Vec::with_capacity(len);
let mut state = seed | 1;
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 options(threads: usize) -> WriteOptions {
WriteOptions::new()
.level(1)
.expect("level")
.memory_limit(MemoryLimit::bytes_or_auto(1 << 30))
.threads(Threads::count_or_single(threads))
}
fn write_entries(entries: &[(&str, &[u8])], threads: usize) -> Vec<u8> {
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options(threads));
for (name, data) in entries {
writer
.add_bytes(ArchivePath::new(name).expect("path"), data)
.expect("adds");
}
let (_result, cursor) = writer.finish_into_inner().expect("finishes");
cursor.into_inner()
}
fn data_area(archive: &[u8]) -> usize {
let offset = u64::from_le_bytes(archive[12..20].try_into().expect("eight bytes"));
SIGNATURE_HEADER_SIZE + offset as usize
}
#[test]
fn test_an_entry_is_written_the_same_whatever_follows_it() {
let first = payload(1, STREAMING_THRESHOLD + 3 * 1024 * 1024);
let second = payload(2, STREAMING_THRESHOLD + 5 * 1024 * 1024);
let alone = write_entries(&[("a.bin", &first)], 8);
let followed = write_entries(&[("a.bin", &first), ("b.bin", &second)], 8);
let span = data_area(&alone);
assert_eq!(
&followed[SIGNATURE_HEADER_SIZE..span],
&alone[SIGNATURE_HEADER_SIZE..span],
"the first entry's stream changed when a second entry was written after it",
);
assert!(
followed.len() > alone.len(),
"the second entry produced no data at all",
);
}
#[test]
fn test_overlapping_entries_do_not_depend_on_the_worker_count() {
let first = payload(3, STREAMING_THRESHOLD + 1024 * 1024);
let second = payload(4, STREAMING_THRESHOLD + 2 * 1024 * 1024);
let third = payload(5, STREAMING_THRESHOLD + 3 * 1024 * 1024);
let entries: &[(&str, &[u8])] = &[("a.bin", &first), ("b.bin", &second), ("c.bin", &third)];
let reference = write_entries(entries, 2);
for threads in [4, 8, 16] {
assert_eq!(
write_entries(entries, threads),
reference,
"{threads} workers changed the archive",
);
}
let mut archive = Archive::open(Cursor::new(reference)).expect("opens");
assert_eq!(archive.extract_to_vec("a.bin").expect("extracts"), first);
assert_eq!(archive.extract_to_vec("b.bin").expect("extracts"), second);
assert_eq!(archive.extract_to_vec("c.bin").expect("extracts"), third);
}
#[test]
fn test_order_survives_an_entry_still_being_finished() {
let first = payload(6, STREAMING_THRESHOLD + 2 * 1024 * 1024);
let last = payload(7, STREAMING_THRESHOLD + 1024 * 1024);
let small = payload(8, 64 * 1024);
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options(8));
writer
.add_bytes(ArchivePath::new("1_big.bin").expect("path"), &first)
.expect("adds");
writer
.add_directory(
ArchivePath::new("2_dir").expect("path"),
EntryMeta::directory(),
)
.expect("adds");
writer
.add_bytes(ArchivePath::new("3_empty.bin").expect("path"), b"")
.expect("adds");
writer
.add_bytes(ArchivePath::new("4_small.bin").expect("path"), &small)
.expect("adds");
writer
.add_bytes(ArchivePath::new("5_big.bin").expect("path"), &last)
.expect("adds");
let (_result, cursor) = writer.finish_into_inner().expect("finishes");
let mut archive = Archive::open(Cursor::new(cursor.into_inner())).expect("opens");
let names: Vec<String> = archive
.entries()
.iter()
.map(|entry| entry.path.as_str().to_string())
.collect();
assert_eq!(
names,
vec![
"1_big.bin",
"2_dir",
"3_empty.bin",
"4_small.bin",
"5_big.bin",
],
);
assert_eq!(
archive.extract_to_vec("1_big.bin").expect("extracts"),
first
);
assert_eq!(
archive.extract_to_vec("3_empty.bin").expect("extracts"),
b""
);
assert_eq!(
archive.extract_to_vec("4_small.bin").expect("extracts"),
small
);
assert_eq!(archive.extract_to_vec("5_big.bin").expect("extracts"), last);
}
struct RefusesAfter {
seen: std::sync::atomic::AtomicUsize,
limit: usize,
finished_badly: std::sync::Arc<std::sync::Mutex<Vec<String>>>,
}
impl zesven::progress::ProgressReporter for RefusesAfter {
fn on_progress(&mut self, _processed: u64, _total: u64) -> bool {
self.seen.fetch_add(1, std::sync::atomic::Ordering::Relaxed) < self.limit
}
fn on_entry_complete(&mut self, entry_name: &str, success: bool) {
if !success {
if let Ok(mut held) = self.finished_badly.lock() {
held.push(entry_name.to_string());
}
}
}
}
#[test]
fn test_calling_off_a_write_ends_as_a_cancellation() {
let first = payload(9, STREAMING_THRESHOLD + 2 * 1024 * 1024);
let second = payload(10, STREAMING_THRESHOLD + 2 * 1024 * 1024);
let badly = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options(8))
.progress(RefusesAfter {
seen: std::sync::atomic::AtomicUsize::new(0),
limit: 24,
finished_badly: std::sync::Arc::clone(&badly),
});
writer
.add_bytes(ArchivePath::new("a.bin").expect("path"), &first)
.expect("adds");
let stopped = writer.add_bytes(ArchivePath::new("b.bin").expect("path"), &second);
let error = stopped.expect_err("the reporter called the write off");
assert!(
matches!(error, zesven::Error::Cancelled),
"a refusal reached the caller as {error} rather than as a cancellation",
);
writer
.finish()
.expect_err("an archive stopped partway cannot be finished");
let badly = badly.lock().expect("lock");
assert!(
!badly.is_empty(),
"nothing was reported as having failed, so a caller watching would \
still be showing entries as running",
);
}
#[test]
fn test_dropping_a_writer_does_not_wait_for_work_nobody_will_read() {
let data = payload(11, STREAMING_THRESHOLD + 4 * 1024 * 1024);
let written = std::time::Instant::now();
{
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options(8));
writer
.add_bytes(ArchivePath::new("a.bin").expect("path"), &data)
.expect("adds");
let _result = writer.finish().expect("finishes");
}
let writing = written.elapsed();
let (done, finished) = std::sync::mpsc::channel();
let dropped = data.clone();
std::thread::spawn(move || {
let started = std::time::Instant::now();
let mut writer = Writer::create(Cursor::new(Vec::new()))
.expect("writer")
.options(options(8));
writer
.add_bytes(ArchivePath::new("a.bin").expect("path"), &dropped)
.expect("adds");
drop(writer);
let _ = done.send(started.elapsed());
});
let abandoning = finished
.recv_timeout(writing * 4 + std::time::Duration::from_secs(60))
.expect(
"dropping a writer with an entry still being compressed did not return: \
the thread finishing it was left waiting on an area nobody drains",
);
let allowed = writing / 2 + std::time::Duration::from_secs(2);
assert!(
abandoning <= allowed,
"writing the entry took {writing:?} and abandoning it took {abandoning:?}, \
over an allowance of {allowed:?}: dropping the writer is seeing the \
entry through rather than stopping it",
);
}