#[cfg(feature = "async")]
pub mod processor {
use crate::{CompressionMode, CompressionStats, DictionarySize, Result};
use futures::stream::{self, StreamExt, TryStreamExt};
use std::path::{Path, PathBuf};
use tokio::fs::File;
use tokio::io::{AsyncReadExt, BufReader};
#[derive(Debug, Clone)]
pub struct AsyncBatchProcessor {
concurrency_limit: usize,
chunk_size: usize,
memory_limit: usize,
}
impl AsyncBatchProcessor {
pub fn new() -> Self {
Self {
concurrency_limit: num_cpus::get(),
chunk_size: 64 * 1024, memory_limit: 256 * 1024 * 1024, }
}
pub fn with_concurrency(mut self, limit: usize) -> Self {
self.concurrency_limit = limit;
self
}
pub fn with_chunk_size(mut self, size: usize) -> Self {
self.chunk_size = size;
self
}
pub fn with_memory_limit(mut self, limit: usize) -> Self {
self.memory_limit = limit;
self
}
pub async fn compress_files<P: AsRef<Path> + Send + Sync>(
&self,
files: Vec<P>,
mode: CompressionMode,
dict_size: DictionarySize,
) -> Result<Vec<(PathBuf, Vec<u8>)>> {
let results = stream::iter(files.into_iter().map(|path| {
let processor = self.clone();
async move { processor.compress_single_file(path, mode, dict_size).await }
}))
.buffer_unordered(self.concurrency_limit)
.try_collect()
.await?;
Ok(results)
}
pub fn compress_files_streaming<P: AsRef<Path> + Send + Sync + 'static>(
&self,
files: Vec<P>,
) -> impl futures::Stream<Item = Result<(PathBuf, CompressionStats)>> + '_ {
let mode = CompressionMode::Binary;
let dict_size = DictionarySize::Size4K;
stream::iter(files.into_iter().map(move |path| {
let processor = self.clone();
async move {
let (path_buf, _data) = processor
.compress_single_file(path, mode, dict_size)
.await?;
let stats = CompressionStats {
literal_count: 0,
match_count: 0,
bytes_processed: 0,
longest_match: 0,
input_bytes: 0, output_bytes: 0, compression_ratio: 0.0,
};
Ok((path_buf, stats))
}
}))
.buffer_unordered(self.concurrency_limit)
}
async fn compress_single_file<P: AsRef<Path>>(
&self,
path: P,
mode: CompressionMode,
dict_size: DictionarySize,
) -> Result<(PathBuf, Vec<u8>)> {
let path = path.as_ref();
let file = File::open(path).await?;
let reader = BufReader::new(file);
let compressed = self.compress_reader(reader, mode, dict_size).await?;
Ok((path.to_path_buf(), compressed))
}
async fn compress_reader<R: tokio::io::AsyncRead + Unpin>(
&self,
mut reader: R,
mode: CompressionMode,
dict_size: DictionarySize,
) -> Result<Vec<u8>> {
use crate::async_implode::AsyncImplodeWriter;
let mut output = Vec::new();
let mut writer = AsyncImplodeWriter::new(&mut output, mode, dict_size)?;
let mut buffer = vec![0u8; self.chunk_size];
loop {
let bytes_read = reader.read(&mut buffer).await?;
if bytes_read == 0 {
break;
}
writer.write_chunk(&buffer[..bytes_read]).await?;
if bytes_read == self.chunk_size {
tokio::task::yield_now().await;
}
}
writer.finish().await?;
Ok(output)
}
pub async fn compress_files_with_memory_limit<P: AsRef<Path> + Send + Sync>(
&self,
files: Vec<P>,
mode: CompressionMode,
dict_size: DictionarySize,
) -> Result<Vec<(PathBuf, Vec<u8>)>> {
let mut results = Vec::new();
let mut current_memory = 0;
for chunk in files.chunks(self.concurrency_limit) {
let chunk_results = stream::iter(chunk.iter().map(|path| {
let processor = self.clone();
async move { processor.compress_single_file(path, mode, dict_size).await }
}))
.buffer_unordered(self.concurrency_limit)
.try_collect::<Vec<_>>()
.await?;
let chunk_memory: usize = chunk_results.iter().map(|(_, data)| data.len()).sum();
current_memory += chunk_memory;
results.extend(chunk_results);
if current_memory > self.memory_limit * 3 / 4 {
tokio::task::yield_now().await;
current_memory = 0; }
}
Ok(results)
}
}
impl Default for AsyncBatchProcessor {
fn default() -> Self {
Self::new()
}
}
}
#[cfg(feature = "async")]
pub use processor::AsyncBatchProcessor;