chapter-tgz 0.1.0

Specially crafted .tar.gz with embedded chapter boundary information
Documentation
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();

    // This could just be TgzReader::open(Cursor::new(tgz)) but we use a more
    // fancy reader to update the 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 {
            // Option 1: We can directly enqueue chapters into the thread pool.
            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("__")
        {
            // Option 2: We can examine chapter entries to decide whether to
            // process or skip a chapter. Reading the first entry's tar header
            // (path, size, PAX extensions) from a chapter is fast.
            chapters.push(tgz.independent_read_chapter(i)?);
            total_compressed_size += tgz.compressed_size_of_chapter(i);
        } else {
            // Option 3: Also fine to pass i and call independent_read_chapter
            // on the other thread.
        }
    }

    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,
        })
    }
}