use super::filter_policy::{
compression_info, encode_member_with_filter_policy_candidates_and_progress,
should_store_compressed_payload,
};
use super::FilterPolicy;
use crate::codec::rar50::{encode_lz_streaming_block, EncodeOptions};
use crate::streaming::Spool;
use crate::{EntrySource, Error, Result, WriterResources};
use std::io::{Read, Write};
pub(super) struct CompressedMember {
pub(super) input_size: u64,
pub(super) crc32: u32,
pub(super) hash: [u8; 32],
pub(super) packed: Spool,
pub(super) store: bool,
pub(super) solid_continuation: bool,
}
#[derive(Debug, Clone)]
pub(super) struct CompressPlan {
pub(super) algorithm_version: u8,
pub(super) encode_options: EncodeOptions,
pub(super) dictionary_size: u64,
pub(super) block_size: usize,
pub(super) solid: bool,
pub(super) method: u8,
pub(super) filter_policy: FilterPolicy,
pub(super) candidates: Vec<EncodeOptions>,
}
struct BlockJob {
member: usize,
data: Vec<u8>,
history: Vec<u8>,
is_last: bool,
}
struct MemberStream {
member: usize,
reader: Box<dyn crate::EntryReader>,
remaining: u64,
packed: Spool,
}
pub(super) fn compress_members_reporting(
sources: &[EntrySource],
plan: CompressPlan,
resources: &WriterResources,
advance: &mut dyn FnMut(u64) -> bool,
) -> Result<Vec<CompressedMember>> {
let mut integrity = Vec::with_capacity(sources.len());
for source in sources {
let input_size = source.len()?;
let (crc32, hash) = super::source_integrity(source, input_size, plan.block_size)?;
integrity.push((input_size, crc32, hash));
}
let optimal_parse = plan.encode_options.optimal_parse
|| plan.candidates.iter().any(|options| options.optimal_parse);
let required =
super::streaming_lz_workspace(plan.dictionary_size, plan.block_size, optimal_parse);
let max_jobs_by_memory = resources.memory_limit() / required;
if max_jobs_by_memory == 0 {
resources.acquire(required, plan.dictionary_size)?;
unreachable!("oversized workspace acquisition must fail");
}
let batch_capacity = usize::try_from(max_jobs_by_memory)
.unwrap_or(usize::MAX)
.min(crate::parallel::threads())
.max(1);
let wants_whole_member =
plan.method != 0 && (plan.filter_policy != FilterPolicy::None || plan.candidates.len() > 1);
if wants_whole_member && !plan.solid {
return compress_members_whole(sources, &integrity, &plan, resources, advance);
}
let packed = if plan.method == 0 {
for (input_size, _, _) in &integrity {
if !advance(*input_size) {
return Err(Error::Cancelled);
}
}
integrity
.iter()
.map(|_| Spool::create(resources))
.collect::<Result<Vec<_>>>()?
} else if plan.solid {
compress_solid_chain(
sources,
&integrity,
&plan,
batch_capacity,
required,
resources,
advance,
)?
} else {
compress_independent_members(
sources,
&integrity,
&plan,
batch_capacity,
required,
resources,
advance,
)?
};
Ok(packed
.into_iter()
.zip(&integrity)
.enumerate()
.map(
|(member, (packed, &(input_size, crc32, hash)))| CompressedMember {
input_size,
crc32,
hash,
store: plan.method == 0
|| input_size == 0
|| should_store_compressed_payload(
input_size,
packed.len(),
plan.solid,
&plan.filter_policy,
),
packed,
solid_continuation: plan.solid && member > 0,
},
)
.collect())
}
fn whole_member_workspace(input_size: u64) -> u64 {
input_size.saturating_mul(4).saturating_add(2 * 1024 * 1024)
}
fn compress_members_whole(
sources: &[EntrySource],
integrity: &[(u64, u32, [u8; 32])],
plan: &CompressPlan,
resources: &WriterResources,
advance: &mut dyn FnMut(u64) -> bool,
) -> Result<Vec<CompressedMember>> {
let mut members = Vec::with_capacity(sources.len());
for (index, source) in sources.iter().enumerate() {
let (input_size, crc32, hash) = integrity[index];
let required = whole_member_workspace(input_size);
let mut packed_spool = Spool::create(resources)?;
let mut stored = input_size == 0;
if !stored {
match resources.acquire(required, plan.dictionary_size) {
Ok(_permit) => {
let mut data = Vec::with_capacity(input_size as usize);
source.open()?.read_to_end(&mut data)?;
if data.len() as u64 != input_size {
return Err(Error::InvalidHeader(
"entry source size changed while compressing",
));
}
let walk = super::filter_policy_walk_bytes(
&data,
&plan.filter_policy,
plan.algorithm_version,
plan.candidates.len(),
)
.max(input_size)
.max(1);
let share = |bytes: u64| {
(u128::from(bytes) * u128::from(input_size) / u128::from(walk)) as u64
};
let mut reported = 0u64;
let mut charged = 0u64;
let mut report = |position: usize| {
let position = position as u64;
if position < reported {
reported = 0;
}
let delta = position - reported;
reported = position;
let target = (charged + delta).min(walk);
let scaled = share(target) - share(charged);
charged = target;
advance(scaled)
};
let packed = encode_member_with_filter_policy_candidates_and_progress(
&data,
plan.algorithm_version,
&plan.filter_policy,
&plan.candidates,
Some(&mut report),
)?;
stored = should_store_compressed_payload(
data.len() as u64,
packed.len() as u64,
plan.solid,
&plan.filter_policy,
);
if !stored {
packed_spool.write_all(&packed)?;
}
}
Err(error) => {
if plan.filter_policy != FilterPolicy::Auto {
return Err(error);
}
let streamed = compress_members_reporting(
std::slice::from_ref(source),
CompressPlan {
filter_policy: FilterPolicy::None,
candidates: vec![plan.encode_options],
..plan.clone()
},
resources,
advance,
)?;
members.extend(streamed);
continue;
}
}
}
members.push(CompressedMember {
input_size,
crc32,
hash,
store: stored,
packed: packed_spool,
solid_continuation: false,
});
}
Ok(members)
}
#[allow(clippy::too_many_arguments)]
fn compress_independent_members(
sources: &[EntrySource],
integrity: &[(u64, u32, [u8; 32])],
plan: &CompressPlan,
batch_capacity: usize,
required: u64,
resources: &WriterResources,
advance: &mut dyn FnMut(u64) -> bool,
) -> Result<Vec<Spool>> {
let mut packed = Vec::with_capacity(sources.len());
for (group_index, group) in sources.chunks(batch_capacity).enumerate() {
let group_start = group_index * batch_capacity;
let mut streams = group
.iter()
.enumerate()
.map(|(offset, source)| {
Ok(MemberStream {
member: offset,
reader: source.open()?,
remaining: integrity[group_start + offset].0,
packed: Spool::create(resources)?,
})
})
.collect::<Result<Vec<_>>>()?;
let mut histories = vec![Vec::new(); streams.len()];
let mut cursor = 0usize;
while streams.iter().any(|stream| stream.remaining != 0) {
let reserved = required.saturating_mul(batch_capacity as u64);
let _permit = resources.acquire(reserved, plan.dictionary_size)?;
let mut jobs = Vec::with_capacity(batch_capacity);
let mut misses = 0usize;
while jobs.len() < batch_capacity && misses < streams.len() {
let stream_count = streams.len();
let stream = &mut streams[cursor];
cursor = (cursor + 1) % stream_count;
if stream.remaining == 0 {
misses += 1;
continue;
}
misses = 0;
let member = stream.member;
let data = read_block(stream, plan.block_size)?;
let is_last = stream.remaining == 0;
jobs.push(BlockJob {
member,
history: histories[member].clone(),
is_last,
data,
});
advance_history(
&mut histories[member],
&jobs.last().expect("just pushed").data,
plan.encode_options.max_match_distance,
);
}
compress_wave(jobs, plan, &mut streams, advance)?;
}
packed.extend(streams.into_iter().map(|stream| stream.packed));
}
Ok(packed)
}
#[allow(clippy::too_many_arguments)]
fn compress_solid_chain(
sources: &[EntrySource],
integrity: &[(u64, u32, [u8; 32])],
plan: &CompressPlan,
batch_capacity: usize,
required: u64,
resources: &WriterResources,
advance: &mut dyn FnMut(u64) -> bool,
) -> Result<Vec<Spool>> {
let mut streams = sources
.iter()
.enumerate()
.map(|(member, source)| {
Ok(MemberStream {
member,
reader: source.open()?,
remaining: integrity[member].0,
packed: Spool::create(resources)?,
})
})
.collect::<Result<Vec<_>>>()?;
let mut history: Vec<u8> = Vec::new();
let mut next = 0usize;
loop {
let reserved = required.saturating_mul(batch_capacity as u64);
let _permit = resources.acquire(reserved, plan.dictionary_size)?;
let mut jobs = Vec::with_capacity(batch_capacity);
while jobs.len() < batch_capacity {
while next < streams.len() && streams[next].remaining == 0 {
next += 1;
}
let Some(stream) = streams.get_mut(next) else {
break;
};
let member = stream.member;
let data = read_block(stream, plan.block_size)?;
let is_last = stream.remaining == 0;
jobs.push(BlockJob {
member,
history: history.clone(),
is_last,
data,
});
advance_history(
&mut history,
&jobs.last().expect("just pushed").data,
plan.encode_options.max_match_distance,
);
}
if jobs.is_empty() {
break;
}
compress_wave(jobs, plan, &mut streams, advance)?;
}
Ok(streams.into_iter().map(|stream| stream.packed).collect())
}
fn read_block(stream: &mut MemberStream, block_size: usize) -> Result<Vec<u8>> {
let wanted = usize::try_from(stream.remaining.min(block_size as u64))
.map_err(|_| Error::InvalidHeader("RAR 5 block size overflows usize"))?;
let mut data = vec![0u8; wanted];
stream.reader.read_exact(&mut data)?;
stream.remaining -= wanted as u64;
if stream.remaining == 0 {
let mut trailing = [0u8; 1];
if stream.reader.read(&mut trailing)? != 0 {
return Err(Error::InvalidHeader(
"entry source size changed while compressing",
));
}
}
Ok(data)
}
fn advance_history(history: &mut Vec<u8>, data: &[u8], max_match_distance: usize) {
history.extend_from_slice(data);
let keep_from = history.len().saturating_sub(max_match_distance);
if keep_from != 0 {
history.drain(..keep_from);
}
}
fn compress_wave(
jobs: Vec<BlockJob>,
plan: &CompressPlan,
streams: &mut [MemberStream],
advance: &mut dyn FnMut(u64) -> bool,
) -> Result<()> {
let wave_bytes: u64 = jobs.iter().map(|job| job.data.len() as u64).sum();
let packed_blocks = crate::parallel::map_collect(jobs, |job| {
let packed = encode_lz_streaming_block(
&job.data,
&job.history,
plan.algorithm_version,
plan.encode_options,
job.is_last,
)?;
Ok::<_, crate::codec::Error>((job.member, packed))
})?;
for (member, packed) in packed_blocks {
streams[member].packed.write_all(&packed)?;
}
if !advance(wave_bytes) {
return Err(Error::Cancelled);
}
Ok(())
}
pub(super) fn member_compression_info(
plan: &CompressPlan,
member: &CompressedMember,
method: u8,
) -> Result<u64> {
compression_info(
plan.algorithm_version,
if member.store { 0 } else { method },
plan.dictionary_size,
member.solid_continuation,
)
}