use std::io::{Read, Seek, Write};
use crate::codec::Encoder;
use crate::{ArchivePath, Error, Result};
use super::options::{EntryMeta, WriteOptions};
use super::{PendingEntry, Writer};
pub(crate) const STREAMING_THRESHOLD: u64 = 64 * 1024 * 1024;
const READ_CHUNK: usize = 256 * 1024;
#[cfg(feature = "parallel")]
pub(crate) struct InFlightBatch {
handle: std::thread::JoinHandle<Result<Vec<super::entry_compression::BatchOutcome>>>,
options: WriteOptions,
footprint: u64,
}
#[cfg(feature = "parallel")]
fn batch_footprint(batch: &[super::BufferedEntry], options: &WriteOptions, reserved: u64) -> u64 {
let entries: u64 = batch.iter().map(|entry| entry.data.len() as u64).sum();
let largest = batch
.iter()
.map(|entry| entry.data.len())
.max()
.unwrap_or(0);
let workers = super::entry_compression::workers_within_budget(options, largest, reserved)
.min(batch.len())
.max(1) as u64;
entries.saturating_mul(4).saturating_add(
workers.saturating_mul(super::codecs::encoder_memory_usage(options, largest)),
)
}
#[derive(Default)]
struct Held {
bytes: Vec<u8>,
#[cfg(all(feature = "parallel", feature = "lzma2"))]
cap: Option<usize>,
#[cfg(all(feature = "parallel", feature = "lzma2"))]
done: bool,
#[cfg(all(feature = "parallel", feature = "lzma2"))]
abandoned: bool,
}
#[derive(Clone)]
struct HoldingArea {
shared: std::sync::Arc<(std::sync::Mutex<Held>, std::sync::Condvar)>,
}
fn lock_failed() -> Error {
Error::Io(std::io::Error::other("compression thread failed"))
}
impl HoldingArea {
fn new() -> Self {
Self {
shared: std::sync::Arc::new((
std::sync::Mutex::new(Held::default()),
std::sync::Condvar::new(),
)),
}
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn bound_to(&self, cap: usize) {
let (lock, _) = &*self.shared;
if let Ok(mut held) = lock.lock() {
held.cap = Some(cap.max(1));
}
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn no_more(&self) {
let (lock, ready) = &*self.shared;
if let Ok(mut held) = lock.lock() {
held.done = true;
}
ready.notify_all();
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn abandon(&self) {
let (lock, ready) = &*self.shared;
if let Ok(mut held) = lock.lock() {
held.abandoned = true;
held.bytes = Vec::new();
}
ready.notify_all();
}
fn waiting(&self) -> usize {
let (lock, _) = &*self.shared;
lock.lock().map_or(0, |held| held.bytes.len())
}
fn swap_into(&self, taker: &mut Vec<u8>) -> Result<()> {
let (lock, ready) = &*self.shared;
let mut held = lock.lock().map_err(|_| lock_failed())?;
std::mem::swap(&mut held.bytes, taker);
drop(held);
ready.notify_all();
Ok(())
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn take_when_ready(&self, taker: &mut Vec<u8>) -> Result<bool> {
let (lock, ready) = &*self.shared;
let mut held = lock.lock().map_err(|_| lock_failed())?;
while held.bytes.is_empty() && !held.done {
held = ready.wait(held).map_err(|_| lock_failed())?;
}
taker.clear();
std::mem::swap(&mut held.bytes, taker);
let complete = held.done;
drop(held);
ready.notify_all();
Ok(complete)
}
}
impl Write for HoldingArea {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let (lock, ready) = &*self.shared;
#[cfg_attr(not(all(feature = "parallel", feature = "lzma2")), allow(unused_mut))]
let mut held = lock
.lock()
.map_err(|_| std::io::Error::other("compression thread failed"))?;
#[cfg(all(feature = "parallel", feature = "lzma2"))]
{
while !held.abandoned
&& held
.cap
.is_some_and(|cap| held.bytes.len() >= cap && !held.bytes.is_empty())
{
held = ready
.wait(held)
.map_err(|_| std::io::Error::other("compression thread failed"))?;
}
if held.abandoned {
drop(held);
ready.notify_all();
return Err(std::io::Error::other("the archive was abandoned"));
}
}
held.bytes.extend_from_slice(buf);
drop(held);
ready.notify_all();
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
fn time_to_collect(batch_finished: bool, waiting: usize, cap: usize) -> bool {
batch_finished || waiting >= cap
}
pub(crate) struct StreamedFolder {
path: ArchivePath,
meta: EntryMeta,
uncompressed_size: u64,
crc: u32,
packed_size: u64,
properties: Vec<u8>,
method: crate::codec::CodecMethod,
}
struct FolderSoFar {
path: ArchivePath,
meta: EntryMeta,
uncompressed_size: u64,
crc: u32,
properties: Vec<u8>,
method: crate::codec::CodecMethod,
}
impl FolderSoFar {
fn packed(self, packed_size: u64) -> StreamedFolder {
StreamedFolder {
path: self.path,
meta: self.meta,
uncompressed_size: self.uncompressed_size,
crc: self.crc,
packed_size,
properties: self.properties,
method: self.method,
}
}
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
pub(crate) struct StreamedTail {
worker: TailWorker,
produced: std::sync::Arc<std::sync::atomic::AtomicU64>,
folder: FolderSoFar,
called_off: std::sync::Arc<std::sync::atomic::AtomicBool>,
held: std::sync::Arc<std::sync::atomic::AtomicU64>,
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
struct TailWorker {
handle: Option<std::thread::JoinHandle<std::io::Result<()>>>,
holding: HoldingArea,
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
impl TailWorker {
fn join(&mut self) -> std::io::Result<()> {
match self.handle.take() {
None => Ok(()),
Some(handle) => handle
.join()
.unwrap_or_else(|_| Err(std::io::Error::other("finishing an entry panicked"))),
}
}
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
impl Drop for TailWorker {
fn drop(&mut self) {
self.holding.abandon();
drop(self.join());
}
}
#[cfg(feature = "parallel")]
pub(crate) enum Ahead {
Batch(InFlightBatch),
#[cfg(feature = "lzma2")]
Tail(StreamedTail),
}
struct Pumped {
encoder: Box<dyn Encoder>,
produced: std::sync::Arc<std::sync::atomic::AtomicU64>,
properties: Vec<u8>,
held: Option<std::sync::Arc<std::sync::atomic::AtomicU64>>,
}
struct StreamedEntry {
crc: crc32fast::Hasher,
uncompressed_size: u64,
}
struct CountingWriter<W> {
inner: W,
written: std::sync::Arc<std::sync::atomic::AtomicU64>,
watcher: Option<Watcher>,
}
struct Reservation {
fixed: u64,
live: Option<std::sync::Arc<std::sync::atomic::AtomicU64>>,
}
impl Reservation {
fn none() -> Self {
Self {
fixed: 0,
live: None,
}
}
}
struct PumpSetup<'a> {
options: &'a WriteOptions,
reserved: Reservation,
watcher: Option<Watcher>,
}
struct Watcher {
reporter: std::sync::Arc<std::sync::Mutex<Box<dyn crate::progress::ProgressReporter>>>,
declared: u64,
called_off: std::sync::Arc<std::sync::atomic::AtomicBool>,
}
impl<W: Write> CountingWriter<W> {
fn new(inner: W) -> Self {
Self {
inner,
written: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
watcher: None,
}
}
fn watched(mut self, watcher: Option<Watcher>) -> Self {
self.watcher = watcher;
self
}
}
impl<W: Write> Write for CountingWriter<W> {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let n = self.inner.write(buf)?;
let written = self
.written
.fetch_add(n as u64, std::sync::atomic::Ordering::Relaxed)
+ n as u64;
let mut called_off = false;
if let Some(watcher) = self.watcher.as_ref() {
if let Ok(mut held) = watcher.reporter.lock() {
called_off = !held.on_progress(written, watcher.declared);
}
}
if called_off {
if let Some(watcher) = self.watcher.as_ref() {
watcher
.called_off
.store(true, std::sync::atomic::Ordering::Relaxed);
}
return Err(std::io::Error::other(
"the progress reporter called the write off",
));
}
Ok(n)
}
fn flush(&mut self) -> std::io::Result<()> {
self.inner.flush()
}
}
pub(super) fn read_some(source: &mut dyn Read, buffer: &mut [u8]) -> Result<usize> {
let mut filled = 0;
while filled < buffer.len() {
let read = source.read(&mut buffer[filled..]).map_err(Error::Io)?;
if read == 0 {
break;
}
filled += read;
}
Ok(filled)
}
pub(crate) fn can_stream(options: &WriteOptions) -> bool {
if options.solid.is_solid() || options.filter.is_active() {
return false;
}
#[cfg(feature = "aes")]
if options.is_data_encrypted() {
return false;
}
encoder_is_available(options)
}
fn encoder_is_available(options: &WriteOptions) -> bool {
use crate::codec::CodecMethod;
match options.method {
CodecMethod::Copy => true,
#[cfg(feature = "lzma2")]
CodecMethod::Lzma2 => true,
#[cfg(feature = "lzma")]
CodecMethod::Lzma => true,
_ => false,
}
}
struct StoreEncoder<W> {
inner: W,
}
impl<W: Write + Send> Write for StoreEncoder<W> {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.inner.write(buf)
}
fn flush(&mut self) -> std::io::Result<()> {
self.inner.flush()
}
}
impl<W: Write + Send> Encoder for StoreEncoder<W> {
fn method_id(&self) -> &'static [u8] {
crate::codec::method::COPY
}
fn finish(mut self: Box<Self>) -> std::io::Result<()> {
self.inner.flush()
}
}
fn overlap_share(options: &WriteOptions) -> usize {
usize::try_from(options.memory_limit.bytes() / 8).unwrap_or(usize::MAX)
}
fn own_buffering(options: &WriteOptions) -> u64 {
overlap_share(options) as u64
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn tail_share(options: &WriteOptions) -> u64 {
options.memory_limit.bytes() / 2
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn can_afford_a_tail(options: &WriteOptions) -> bool {
let encoder = crate::codec::lzma::encoder_memory_usage(
options.level,
super::codecs::stream_dictionary_size(options),
);
tail_share(options) >= encoder
}
fn holding_cap(options: &WriteOptions) -> usize {
overlap_share(options) / 2
}
struct BuiltEncoder {
encoder: Box<dyn Encoder>,
properties: Vec<u8>,
held: Option<std::sync::Arc<std::sync::atomic::AtomicU64>>,
}
fn encoder_for<W: Write + Send + 'static>(
options: &WriteOptions,
output: W,
reserved: Reservation,
) -> Result<BuiltEncoder> {
use crate::codec::CodecMethod;
let Reservation { fixed, live } = reserved;
#[cfg(not(all(feature = "parallel", feature = "lzma2")))]
let _ = (fixed, live);
let plain = |encoder: Box<dyn Encoder>, properties: Vec<u8>| BuiltEncoder {
encoder,
properties,
held: None,
};
match options.method {
CodecMethod::Copy => Ok(plain(Box::new(StoreEncoder { inner: output }), Vec::new())),
#[cfg(feature = "lzma2")]
CodecMethod::Lzma2 => {
use crate::codec::lzma::{Lzma2Encoder, Lzma2EncoderOptions};
let opts = Lzma2EncoderOptions {
preset: options.level,
dict_size: Some(super::codecs::stream_dictionary_size(options)),
};
let properties = opts.properties();
#[cfg(feature = "parallel")]
if super::codecs::lzma2_is_chunked(options, &opts) {
use crate::codec::lzma2_chunked::ChunkedLzma2Encoder;
let encoder = ChunkedLzma2Encoder::new(
output,
&opts,
options.threads.count(),
options.memory_limit.bytes().saturating_sub(fixed),
live,
)?;
let held = encoder.held();
return Ok(BuiltEncoder {
encoder: Box::new(encoder),
properties,
held: Some(held),
});
}
Ok(plain(
Box::new(Lzma2Encoder::new(output, &opts)),
properties,
))
}
#[cfg(feature = "lzma")]
CodecMethod::Lzma => {
use crate::codec::lzma::{LzmaEncoder, LzmaEncoderOptions};
let opts = LzmaEncoderOptions {
preset: options.level,
dict_size: Some(super::codecs::stream_dictionary_size(options)),
};
let properties = opts.properties();
Ok(plain(
Box::new(LzmaEncoder::new(output, &opts)?),
properties,
))
}
method => Err(Error::UnsupportedMethod {
method_id: method.method_id(),
}),
}
}
impl<W: Write + Seek> Writer<W> {
pub(crate) fn record_streamed_folder(&mut self, folder: StreamedFolder) {
let StreamedFolder {
path,
meta,
uncompressed_size,
crc,
packed_size,
properties,
method,
} = folder;
let pending = PendingEntry {
path,
meta,
uncompressed_size,
};
self.compressed_bytes += packed_size;
self.stream_info.pack_sizes.push(packed_size);
self.stream_info.unpack_sizes.push(uncompressed_size);
self.stream_info.coder_methods.push(method);
self.stream_info.coder_properties.push(properties);
self.stream_info.crcs.push(None);
self.stream_info.substream_sizes.push(uncompressed_size);
self.stream_info.substream_crcs.push(crc);
#[cfg(feature = "aes")]
self.stream_info.encryption_info.push(None);
self.stream_info.filter_info.push(None);
self.stream_info.bcj2_folder_info.push(None);
self.stream_info.num_unpack_streams_per_folder.push(1);
self.record_entry(pending);
}
fn pour(&mut self, holding: &HoldingArea, scratch: &mut Vec<u8>) -> Result<()> {
scratch.clear();
holding.swap_into(scratch)?;
if scratch.is_empty() {
return Ok(());
}
self.sink.write_all(scratch).map_err(Error::Io)
}
#[cfg(feature = "parallel")]
fn ahead_reservation(&self) -> Reservation {
match self.ahead.as_ref() {
None => Reservation::none(),
Some(Ahead::Batch(batch)) => Reservation {
fixed: batch.footprint,
live: None,
},
#[cfg(feature = "lzma2")]
Some(Ahead::Tail(tail)) => Reservation {
fixed: own_buffering(&self.options),
live: Some(std::sync::Arc::clone(&tail.held)),
},
}
}
#[cfg(feature = "parallel")]
pub(crate) fn ahead_reservation_now(&self) -> u64 {
let reservation = self.ahead_reservation();
reservation.fixed.saturating_add(
reservation
.live
.map_or(0, |held| held.load(std::sync::atomic::Ordering::Relaxed)),
)
}
#[cfg(feature = "parallel")]
fn ahead_is_finished(&self) -> bool {
match self.ahead.as_ref() {
None => false,
Some(Ahead::Batch(batch)) => batch.handle.is_finished(),
#[cfg(feature = "lzma2")]
Some(Ahead::Tail(tail)) => tail
.worker
.handle
.as_ref()
.is_some_and(|handle| handle.is_finished()),
}
}
#[cfg(feature = "parallel")]
pub(crate) fn settle_ahead(&mut self) -> Result<()> {
match self.ahead.take() {
None => Ok(()),
Some(Ahead::Batch(batch)) => self.settle_batch(batch),
#[cfg(feature = "lzma2")]
Some(Ahead::Tail(tail)) => self.settle_tail(tail),
}
}
#[cfg(feature = "parallel")]
fn settle_batch(&mut self, batch: InFlightBatch) -> Result<()> {
let InFlightBatch {
handle, options, ..
} = batch;
let outcomes = handle
.join()
.map_err(|_| Error::Io(std::io::Error::other("compressing a batch panicked")))??;
self.write_batch_outcomes(outcomes, &options)
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
fn settle_tail(&mut self, tail: StreamedTail) -> Result<()> {
let StreamedTail {
mut worker,
produced,
folder,
called_off,
..
} = tail;
let mut scratch = Vec::new();
loop {
let complete = match worker.holding.take_when_ready(&mut scratch) {
Ok(complete) => complete,
Err(e) => return self.fail(e),
};
if !scratch.is_empty() {
if let Err(e) = self.sink.write_all(&scratch) {
return self.fail(Error::Io(e));
}
}
if complete {
break;
}
}
if let Err(e) = worker.join() {
let e = if called_off.load(std::sync::atomic::Ordering::Relaxed) {
Error::Cancelled
} else {
Error::Io(e)
};
return self.fail(e);
}
let packed_size = produced.load(std::sync::atomic::Ordering::Relaxed);
self.record_streamed_folder(folder.packed(packed_size));
Ok(())
}
#[cfg(feature = "parallel")]
pub(crate) fn abandon_ahead(&mut self) {
match self.ahead.take() {
None => {}
Some(Ahead::Batch(batch)) => drop(batch.handle.join()),
#[cfg(feature = "lzma2")]
Some(Ahead::Tail(tail)) => drop(tail),
}
}
#[cfg(feature = "parallel")]
fn something_ahead(&self) -> bool {
self.ahead.is_some()
}
}
impl<W: Write + Seek + Send> Writer<W> {
fn send_batch_ahead(&mut self) -> Result<()> {
#[cfg(not(feature = "parallel"))]
return Ok(());
#[cfg(feature = "parallel")]
{
if self.pending_batch.is_empty() {
return Ok(());
}
if self.ahead.is_some() {
return Ok(());
}
if !self.solid_buffer.is_empty() {
return Ok(());
}
let options = self
.pending_batch
.first()
.map(|entry| (*entry.options).clone())
.unwrap_or_else(|| (*self.active_options).clone());
if options.threads.count() <= 1 {
return Ok(());
}
let batch = std::mem::take(&mut self.pending_batch);
self.pending_batch_size = 0;
let footprint = batch_footprint(&batch, &options, 0);
if footprint > overlap_share(&options) as u64 {
self.pending_batch = batch;
self.pending_batch_size =
self.pending_batch.iter().map(|e| e.data.len() as u64).sum();
return Ok(());
}
self.announce_entries(
batch
.iter()
.map(|entry| (entry.path.as_str().to_string(), entry.data.len() as u64))
.collect(),
);
let for_thread = options.clone();
let handle = match std::thread::Builder::new()
.name("zesven-batch".into())
.spawn(move || {
super::entry_compression::compress_batch_owned(batch, &for_thread, 0)
}) {
Ok(handle) => handle,
Err(e) => return self.fail(Error::Io(e)),
};
self.ahead = Some(Ahead::Batch(InFlightBatch {
handle,
options,
footprint,
}));
Ok(())
}
}
fn pump(
source: &mut dyn Read,
prefix: Vec<u8>,
holding: HoldingArea,
setup: PumpSetup<'_>,
state: &mut StreamedEntry,
mut between: impl FnMut(u64) -> Result<()>,
) -> Result<Pumped> {
let PumpSetup {
options,
reserved,
watcher,
} = setup;
let mut buffer = vec![0u8; READ_CHUNK];
let counting = CountingWriter::new(holding).watched(watcher);
let produced = std::sync::Arc::clone(&counting.written);
let BuiltEncoder {
mut encoder,
properties,
held,
} = encoder_for(options, counting, reserved)?;
state.crc.update(&prefix);
state.uncompressed_size += prefix.len() as u64;
let mut result = Ok(());
for piece in prefix.chunks(READ_CHUNK) {
if let Err(e) = encoder.write_all(piece) {
result = Err(Error::Io(e));
break;
}
if let Err(e) = between(produced.load(std::sync::atomic::Ordering::Relaxed)) {
result = Err(e);
break;
}
}
drop(prefix);
while result.is_ok() {
let read = match read_some(source, &mut buffer) {
Ok(n) => n,
Err(e) => {
result = Err(e);
break;
}
};
if read == 0 {
break;
}
state.crc.update(&buffer[..read]);
state.uncompressed_size += read as u64;
if let Err(e) = encoder.write_all(&buffer[..read]) {
result = Err(Error::Io(e));
break;
}
if let Err(e) = between(produced.load(std::sync::atomic::Ordering::Relaxed)) {
result = Err(e);
break;
}
}
if let Err(e) = result {
drop(encoder);
return Err(e);
}
Ok(Pumped {
encoder,
produced,
properties,
held,
})
}
pub(crate) fn compress_entry_streaming(
&mut self,
archive_path: ArchivePath,
prefix: Vec<u8>,
source: &mut dyn Read,
meta: EntryMeta,
) -> Result<()> {
self.send_batch_ahead()?;
if let Err(e) = self.flush_buffered_entries() {
self.abandon_ahead();
return self.fail(e);
}
let mut state = StreamedEntry {
crc: crc32fast::Hasher::new(),
uncompressed_size: 0,
};
let holding = HoldingArea::new();
let mut scratch = Vec::new();
let held_limit = holding_cap(&self.options);
let options = self.options.clone();
let mut reserved = self.ahead_reservation();
reserved.fixed = reserved.fixed.saturating_add(own_buffering(&options));
let declared = meta.size;
self.announce_entries(vec![(archive_path.as_str().to_string(), declared)]);
let shared = self.progress.clone();
let called_off = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let watcher = shared.as_ref().map(|reporter| Watcher {
reporter: std::sync::Arc::clone(reporter),
declared,
called_off: std::sync::Arc::clone(&called_off),
});
let outcome = Self::pump(
source,
prefix,
holding.clone(),
PumpSetup {
options: &options,
reserved,
watcher,
},
&mut state,
|_| {
if self.something_ahead()
&& time_to_collect(self.ahead_is_finished(), holding.waiting(), held_limit)
{
self.settle_ahead()?;
}
if !self.something_ahead() {
self.pour(&holding, &mut scratch)?;
}
Ok(())
},
);
drop(shared);
let pumped = match outcome {
Ok(pumped) => pumped,
Err(e) => {
self.abandon_ahead();
let e = if called_off.load(std::sync::atomic::Ordering::Relaxed) {
Error::Cancelled
} else {
e
};
return self.fail(e);
}
};
let settled = self
.settle_ahead()
.and_then(|()| self.pour(&holding, &mut scratch));
if let Err(e) = settled {
self.abandon_ahead();
return self.fail(e);
}
let StreamedEntry {
crc,
uncompressed_size,
} = state;
let folder = FolderSoFar {
path: archive_path,
meta,
uncompressed_size,
crc: crc.finalize(),
properties: pumped.properties,
method: options.method,
};
#[cfg_attr(not(all(feature = "parallel", feature = "lzma2")), allow(unused_mut))]
let Pumped {
mut encoder,
produced,
held,
..
} = pumped;
let Some(held) = held else {
return self.finish_entry_here(encoder, &produced, folder, &holding, &mut scratch);
};
#[cfg(not(all(feature = "parallel", feature = "lzma2")))]
{
let _ = (held, called_off, &options, held_limit);
return self.finish_entry_here(encoder, &produced, folder, &holding, &mut scratch);
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
{
if !can_afford_a_tail(&options) {
return self.finish_entry_here(encoder, &produced, folder, &holding, &mut scratch);
}
let target = tail_share(&options);
while held.load(std::sync::atomic::Ordering::Relaxed) > target {
match encoder.drain_one_block() {
Ok(true) => {}
Ok(false) => break,
Err(e) => return self.fail(Error::Io(e)),
}
if let Err(e) = self.pour(&holding, &mut scratch) {
return self.fail(e);
}
}
holding.bound_to(held_limit);
let for_thread = holding.clone();
let handle = std::thread::Builder::new()
.name("zesven-tail".into())
.spawn(move || {
let _closing = ClosingArea(for_thread);
encoder.finish()
});
match handle {
Ok(handle) => {
self.ahead = Some(Ahead::Tail(StreamedTail {
worker: TailWorker {
handle: Some(handle),
holding,
},
produced,
folder,
called_off,
held,
}));
Ok(())
}
Err(e) => self.fail(Error::Io(e)),
}
}
}
fn finish_entry_here(
&mut self,
encoder: Box<dyn Encoder>,
produced: &std::sync::atomic::AtomicU64,
folder: FolderSoFar,
holding: &HoldingArea,
scratch: &mut Vec<u8>,
) -> Result<()> {
let finished = encoder.finish().map_err(Error::Io);
let settled = finished.and_then(|()| self.pour(holding, scratch));
if let Err(e) = settled {
return self.fail(e);
}
let packed_size = produced.load(std::sync::atomic::Ordering::Relaxed);
self.record_streamed_folder(folder.packed(packed_size));
Ok(())
}
}
#[cfg(all(feature = "parallel", feature = "lzma2"))]
struct ClosingArea(HoldingArea);
#[cfg(all(feature = "parallel", feature = "lzma2"))]
impl Drop for ClosingArea {
fn drop(&mut self) {
self.0.no_more();
}
}
#[cfg(test)]
mod tests {
use super::time_to_collect;
#[test]
fn test_an_entry_is_charged_for_both_buffers_it_pours_through() {
use super::{overlap_share, own_buffering};
use crate::write::options::WriteOptions;
let options = WriteOptions::new();
let share = overlap_share(&options) as u64;
assert_eq!(super::holding_cap(&options) as u64 * 2, share);
assert_eq!(own_buffering(&options), share);
}
#[cfg(feature = "parallel")]
#[test]
fn test_a_batch_is_charged_for_its_output_as_well_as_its_input() {
use super::super::BufferedEntry;
use super::batch_footprint;
use crate::ArchivePath;
use crate::write::options::{EntryMeta, WriteOptions};
let options = std::sync::Arc::new(WriteOptions::new());
let entry = |name: &str, len: usize| BufferedEntry {
path: ArchivePath::new(name).expect("path"),
data: vec![0u8; len],
meta: EntryMeta::file(len as u64),
crc: 0,
options: options.clone(),
};
let lean = vec![entry("a.bin", 16 << 20), entry("b.bin", 1)];
let full = vec![entry("a.bin", 16 << 20), entry("b.bin", 8 << 20)];
let more = ((8 << 20) - 1) as u64;
let grown = batch_footprint(&full, &options, 0) - batch_footprint(&lean, &options, 0);
assert_eq!(
grown,
more * 4,
"{more} more bytes in a batch grew what it is charged by {grown}: \
its output, or the room the vector holding that output reserves \
past it, is not being charged beside its input",
);
}
#[test]
fn test_a_batch_is_collected_when_it_ends_or_when_there_is_no_room() {
assert!(!time_to_collect(false, 0, 64));
assert!(!time_to_collect(false, 63, 64));
assert!(time_to_collect(true, 0, 64));
assert!(time_to_collect(false, 64, 64));
assert!(time_to_collect(false, 65, 64));
}
}
#[cfg(not(feature = "parallel"))]
impl<W: Write + Seek> Writer<W> {
pub(crate) fn ahead_reservation_now(&self) -> u64 {
0
}
fn ahead_reservation(&self) -> Reservation {
Reservation::none()
}
fn something_ahead(&self) -> bool {
false
}
fn ahead_is_finished(&self) -> bool {
false
}
pub(crate) fn settle_ahead(&mut self) -> Result<()> {
Ok(())
}
pub(crate) fn abandon_ahead(&mut self) {}
}