use anyhow::Result;
use chapter_tgz::{Compression, IndependentRead, TgzReader, TgzWriter, tar};
use indicatif::{ProgressBar, ProgressDrawTarget, ProgressStyle};
use rand::RngReader;
use rand::rngs::SmallRng;
use std::io::{self, Cursor, Read, Seek, SeekFrom};
use std::thread;
use std::time::Duration;
fn main() -> Result<()> {
println!("Building tgz...");
let tgz = build_tgz()?;
println!("Extracting in parallel on 8 threads...");
extract_multithreaded(&tgz)?;
Ok(())
}
fn build_tgz() -> Result<Vec<u8>> {
const FILES: u64 = 8;
const SIZE_PER_FILE: u64 = 100_000_000;
let pb = progress_bar();
pb.set_length(FILES * SIZE_PER_FILE);
pb.set_draw_target(ProgressDrawTarget::stderr());
let mut tgz = TgzWriter::new(Vec::new(), Compression::fast());
let mut rng: SmallRng = rand::make_rng();
for i in 0..FILES {
let mut chapter = tgz.create_chapter();
let mut header = tar::Header::new_gnu();
header.set_size(SIZE_PER_FILE);
let path = format!("random/{i}");
let data = RngReader(&mut rng).take(SIZE_PER_FILE);
let progress = pb.wrap_read(data);
chapter.append_data(&mut header, path, progress)?;
}
let compressed = tgz.into_inner()?;
Ok(compressed)
}
fn extract_multithreaded(tgz: &[u8]) -> Result<()> {
let pb = progress_bar();
let tgz = TgzReader::open(ProgressRead {
inner: Cursor::new(tgz),
progress: &pb,
})?;
let n = tgz.chapters();
let mut chapters = Vec::with_capacity(n as usize);
let mut total_compressed_size = 0;
for i in 0..n {
let mut chapter = tgz.independent_read_chapter(i)?;
if i % 2 == 0 {
chapters.push(chapter);
total_compressed_size += tgz.compressed_size_of_chapter(i);
} else if let Some(first_entry) = chapter.entries()?.next()
&& let Some(file_name) = first_entry?.path()?.file_name()
&& let Some(file_name_str) = file_name.to_str()
&& !file_name_str.starts_with("__")
{
chapters.push(tgz.independent_read_chapter(i)?);
total_compressed_size += tgz.compressed_size_of_chapter(i);
} else {
}
}
pb.reset();
pb.set_length(total_compressed_size);
pb.set_draw_target(ProgressDrawTarget::stderr());
thread::scope(|scope| {
for mut chapter in chapters {
scope.spawn(move || {
if let Err(err) = (|| -> Result<()> {
for _entry in chapter.entries()? {}
Ok(())
})() {
eprintln!("Error: {err}");
}
});
}
});
Ok(())
}
fn progress_bar() -> ProgressBar {
let pb = ProgressBar::hidden();
pb.set_style(
ProgressStyle::default_bar()
.template("[{wide_bar:.cyan/blue}] {percent}% ")
.unwrap()
.progress_chars(". "),
);
pb
}
struct ProgressRead<'a, R> {
inner: R,
progress: &'a ProgressBar,
}
impl<R: Read> Read for ProgressRead<'_, R> {
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
thread::sleep(Duration::from_millis(1));
let n = self.inner.read(buf)?;
self.progress.inc(n as u64);
Ok(n)
}
}
impl<R: Seek> Seek for ProgressRead<'_, R> {
fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
self.inner.seek(pos)
}
}
impl<R: IndependentRead> IndependentRead for ProgressRead<'_, R> {
fn independent_clone(&self) -> io::Result<Self> {
Ok(ProgressRead {
inner: self.inner.independent_clone()?,
progress: self.progress,
})
}
}