diff --git a/DEV.md b/DEV.md index 0d8e9df..11ef5c7 100644 --- a/DEV.md +++ b/DEV.md @@ -49,7 +49,7 @@ The bzip2 encoder incrementally forms RLE1 blocks and schedules BWT, MTF/RLE2, a The decoder remains independent of files, Python, and the CLI. Parallel scanning/decoding and indexed seeking are layered over it. Native workers never call Python. Large offsets use explicit 64-bit bit/byte types, and speculative block-marker hits are accepted only when they form an exact stream chain with valid block and combined stream CRCs. -Core decode APIs report completed compressed and decoded byte counts without knowing anything about terminals. The CLI selects bzip2, gzip, LZ4, or ZIP by a recognised extension and falls back to magic for stdin or unknown names. It layers delayed, rate-limited TTY progress rendering over the shared callbacks; redirected stderr and `--quiet` produce no progress output. Decoded files use same-directory temporary files and atomic persistence, then inherit the compressed input's modification time and permissions. `--rm` removes an input only after decode, persistence, and metadata copying all succeed. An `OutputSink` wrapper enforces output-size limits, so each decoder has one code path for files, stdout, validation, listing, and archive extraction. +Core decode APIs report completed compressed and decoded byte counts without knowing anything about terminals. The public `decompress` and writer variants detect bzip2, gzip, or LZ4 magic through one dispatcher; `DecodeOptions::format` provides an explicit override. Bzip2 indexing remains a separate API because its preliminary validation, seeking, and cache controls have different semantics. The CLI selects bzip2, gzip, LZ4, or ZIP by a recognised extension and falls back to magic for stdin or unknown names. It layers delayed, rate-limited TTY progress rendering over the shared callbacks; redirected stderr and `--quiet` produce no progress output. Decoded files use same-directory temporary files and atomic persistence, then inherit the compressed input's modification time and permissions. `--rm` removes an input only after decode, persistence, and metadata copying all succeed. An `OutputSink` wrapper enforces output-size limits, so each decoder has one code path for files, stdout, validation, listing, and archive extraction. The LZ4 decoder parses current frames, concatenation, and skippable frames itself. It validates descriptor bits and XXH32 header, block, and content checksums, and bounds every literal and match before writing. A frame header creates an incremental block cursor rather than a complete layout. Independent blocks are gathered into at most 64-entry batches—only until there is enough work to amortize the pool—and become ordinary `pipeline::Job`s. A parse failure discovered during bounded look-ahead is held until every earlier valid block has decoded and emitted, preserving stream error order. Compressed jobs reserve the frame's declared maximum decoded block size; stored blocks borrow their source bytes and reserve no decoded allocation. Frames containing only stored blocks without block checksums bypass the worker pool because their only remaining work is ordered output and optional content hashing. Retained results remain charged at their allocation capacity rather than logical length, so highly compressible blocks cannot understate memory use. If fewer than two natural blocks fit the speculative budget, decoding proceeds incrementally on the coordinator instead of rejecting the frame. The pool is created lazily and reused across concatenated frames. Linked frames use the same parser, block decoder, output sink, progress, and report path, but decode serially with a rolling 64 KiB history. This is one code path with a scheduling branch, not separate Reader and writer implementations. @@ -65,7 +65,7 @@ ZIP creation follows the same policy over uncompressed sizes. One file of at lea The shared `pipeline.rs` scheduler provides ordered results, byte-budgeted admission, cancellation, and a staged priority queue. Bzip2 uses the rolling candidate path: workers reserve the maximum possible decoded block size, then shrink that reservation to actual retained output until ordered validation consumes or rejects it. Gzip decoding uses the staged path for speculative DEFLATE and priority marker resolution. LZ4 decoding and multi-entry ZIP use ordinary ordered jobs. Compression's long-lived `StreamingOrdered` pool accepts work as bytes arrive, retains output in key order, catches worker panics, and cancels queued work when dropped. Gzip, LZ4, and bzip2 encoders all use it rather than maintaining codec-specific thread/channel machinery. Reservations conservatively cover owned input, temporary working state, and retained output until the coordinator consumes it. -The shared `OutputSink` boundary accepts owned decoder chunks. Direct writer APIs adapt those chunks to `Write::write_all` without a channel or allocation copy. `Reader` and streaming tar extraction instead use the same zero-capacity owned-chunk pipe, so completed bzip2 blocks, resolved parallel-gzip segments, and decoded LZ4 blocks move into the consumer rather than being copied into an intermediate pipe buffer. The pipe adds at most the consumer's current chunk and the producer's next blocked chunk beyond the scheduler budget. Its receiver owns a cancellation flag; bzip2 scanning checks that flag between bounded waves, while LZ4 checks it before parsing each serial block or parallel batch. +The shared `OutputSink` boundary accepts owned decoder chunks. Direct writer APIs adapt those chunks to `Write::write_all` without a channel or allocation copy. Rust `Reader` and streaming tar extraction instead use the same zero-capacity owned-chunk pipe, so completed bzip2 blocks, resolved parallel-gzip segments, and decoded LZ4 blocks move into the consumer rather than being copied into an intermediate pipe buffer. Python's `Reader` is a thin `io.RawIOBase` wrapper over that Rust reader and exposes native `read`/`readinto` operations while releasing the GIL during reads. Its safe buffer-protocol bridge copies at most 1 MiB per `readinto` call rather than retaining an allocation as large as an arbitrary caller-provided buffer. The pipe adds at most the consumer's current chunk and the producer's next blocked chunk beyond the scheduler budget. Its receiver owns a cancellation flag; bzip2 scanning checks that flag between bounded waves, while LZ4 checks it before parsing each serial block or parallel batch. The non-indexing parallel bzip2 path decodes and emits its first CRC-validated block before scanning the remainder of the compressed source. That candidate is then reused by normal ordered assembly rather than decoded twice. Index construction retains the full scan-first path because it produces no output. `Reader` sends decoder success or failure separately from the byte pipe and treats a missing terminal status as an error, so a worker panic or corrupt trailer cannot appear as EOF. Reader errors are sticky across subsequent calls, and drop disconnects the pipe before joining the decoder thread. @@ -379,4 +379,4 @@ The thin PEP 517 backend delegates to Maturin after building and staging the nat 4. For the first crates.io release only, run `cargo publish`, then configure the `ci.yml` trusted publisher for `AnswerDotAI/fbz`; crates.io requires the crate to exist before trusted publishing can be configured. 5. Run `ship-release`. -Fastship pushes the version tag for GitHub Actions, then bumps and pushes `Cargo.toml`. Tagged CI publishes the crate through crates.io trusted publishing as well as building the GitHub release and PyPI packages. +Fastship pushes the version tag for GitHub Actions, then bumps and pushes `Cargo.toml`. Tagged CI publishes the crate through crates.io trusted publishing as well as building the GitHub release and PyPI packages. The AnswerDotAI Homebrew tap discovers the new crate version in its daily bump workflow; green bot-created fbz-only formula PRs publish their bottles automatically through the tap's generated `brew pr-pull` workflow. diff --git a/README.md b/README.md index 72ffcc1..865c967 100644 --- a/README.md +++ b/README.md @@ -120,7 +120,7 @@ fbz --list --json dump.xml.bz2 # emit the complete layout as JSON ## Python -The Python API exposes one-shot compression for all three stream formats. Decompression, validation, scanning, and indexed seeking currently expose the bzip2 backend. +The Python API exposes compression, decompression, validation, and streaming reads for all three stream formats. Bzip2 additionally supports indexed seeking and structural scanning. ### One-shot compression @@ -141,18 +141,31 @@ plain = fbz.decompress(compressed_bytes) fbz.test("dump.xml.bz2") # returns None after successful validation ``` -`decompress` accepts a bytes-like object and returns `bytes`. `test` accepts either compressed bytes or a path and avoids retaining the decoded result. +`decompress` accepts a bytes-like object, detects bzip2, gzip, or LZ4 from its magic, and returns `bytes`. Pass `format="gzip"` (or `"bzip2"`/`"lz4"`) to override detection. `test` accepts either compressed bytes or a path and avoids retaining the decoded result. + +### Streaming reads + +`fbz.open` and `fbz.Reader` start decoding a file immediately and implement Python's raw binary I/O interface. They do not index first or retain the plaintext; wrap them in `io.BufferedReader` when buffering is useful: + +```python +import io, fbz + +with io.BufferedReader(fbz.open("events.json.gz")) as f: + header = f.read(4096) +``` + +They validate the complete stream and report a late checksum failure as `fbz.BadCompressedFile` rather than EOF. Dropping or closing a reader early cancels and joins its workers. Compressed tar inputs yield tar bytes; ZIP remains an archive rather than a byte-stream reader. ### Seekable reads and persistent indexes -`fbz.open` returns a seekable binary `io.RawIOBase`. Opening without an index performs a complete validation pass and builds an in-memory block index; `build_index` can persist that work for later processes: +`fbz.open_indexed` returns a seekable bzip2 `io.RawIOBase`. Opening without an index performs a complete validation pass and builds an in-memory block index; `build_index` can persist that work for later processes: ```python import fbz fbz.build_index("dump.xml.bz2", "dump.xml.bz2.fbz2i") -with fbz.open("dump.xml.bz2", index="dump.xml.bz2.fbz2i") as f: +with fbz.open_indexed("dump.xml.bz2", index="dump.xml.bz2.fbz2i") as f: f.seek(1_000_000_000) chunk = f.read(64 * 1024) print(f.tell(), f.size) @@ -172,7 +185,7 @@ result = scan(bz2.compress(b"hello")) assert result.blocks[0].bit_offset == 32 ``` -Scan results are deliberately untrusted candidates. Use `test`, `decompress`, `build_index`, or `open` when validation is required. +Scan results are deliberately untrusted candidates. Use `test`, `decompress`, `build_index`, or successfully exhaust a streaming reader when validation is required. ## Rust @@ -189,7 +202,7 @@ fn main() -> fbz::Result<()> { } ``` -`decompress` returns a `Vec`. `decode_to_writer` returns a validated `Index` while streaming output, `build_index` validates into a sink, and their `*_with_progress` variants report completed compressed and decoded byte counts. `IndexedReader` implements `Read` and `Seek`; it can build an index itself or load a persisted one with `open_with_index`. +`decompress` returns a `Vec` and detects bzip2, gzip, or LZ4 from magic. Set `DecodeOptions::format` to a `DecodeFormat` variant to override detection. `decompress_to_writer` and its progress variant use the same dispatcher and codec implementations. The bzip2-specific `decode_to_writer` returns a validated `Index`, while `build_index` validates into a sink. `IndexedReader` implements `Read` and `Seek`; it can build an index itself or load a persisted one with `open_with_index`. The unified compression API covers bzip2, gzip, and LZ4 and writes incrementally: @@ -243,7 +256,7 @@ fn main() -> Result<(), Box> { } ``` -Magic takes priority over the filename extension, with the extension used as a fallback for damaged headers. The decoder runs on an owned worker thread and transfers completed decoder allocations through a zero-capacity pipe; it neither materializes the plaintext nor writes an intermediate file. `DecodeOptions` controls decoder threads and speculative memory. Dropping early disconnects the pipe, cancels outstanding work, and joins the worker. +Magic takes priority over the filename extension, with the extension used as a fallback for damaged headers; an explicit `DecodeOptions::format` overrides both. The decoder runs on an owned worker thread and transfers completed decoder allocations through a zero-capacity pipe; it neither materializes the plaintext nor writes an intermediate file. `DecodeOptions` controls format selection, decoder threads, and speculative memory. Dropping early disconnects the pipe, cancels outstanding work, and joins the worker. Checksum errors discovered after output has begun are returned by a later `read()` call. Therefore only successful EOF establishes that the complete stream was valid; dropping early deliberately does not finish validation. Compressed tar inputs yield the decoded tar byte stream rather than extracting it. ZIP is not exposed through `Reader` because an archive has no single decoded byte stream. diff --git a/python/fbz/__init__.py b/python/fbz/__init__.py index d4bbf5c..f1ff468 100644 --- a/python/fbz/__init__.py +++ b/python/fbz/__init__.py @@ -2,7 +2,7 @@ import io, os from pathlib import Path -from ._core import BadBzip2File, _IndexedReader, __version__, _build_index, _compress, _decompress, _scan, _test, bz2_crc32 +from ._core import BadCompressedFile, _IndexedReader, _Reader, __version__, _build_index, _compress, _decompress, _scan, _test_bytes, _test_path, bz2_crc32 DEFAULT_MEMORY_LIMIT = 1024 * 1024 * 1024 DEFAULT_CACHE_LIMIT = 64 * 1024 * 1024 @@ -24,9 +24,9 @@ def scan(data: bytes) -> ScanResult: stream_ends = [EndCandidate(*item) for item in stream_ends] return ScanResult(streams, blocks, stream_ends) -def decompress(data: bytes, *, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT) -> bytes: - "Decompress and fully CRC-validate one or more concatenated bzip2 streams." - return _decompress(data, threads, memory_limit) +def decompress(data: bytes, format=None, *, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT) -> bytes: + "Decompress and validate bzip2, gzip, or LZ4 bytes, selected by magic unless *format* is given." + return _decompress(data, format, threads, memory_limit) def compress(data: bytes, format: str, *, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT, level=None) -> bytes: "Compress *data* as bzip2, gzip, or LZ4." @@ -79,8 +79,33 @@ def close(self): self._reader = None super().close() -def open(source, *, threads=0, index=None, memory_limit=DEFAULT_MEMORY_LIMIT, cache_limit=DEFAULT_CACHE_LIMIT): - "Open a path or bytes object as a seekable bzip2 binary file." +class Reader(io.RawIOBase): + "Streaming binary reader for a bzip2, gzip, or LZ4 file." + + def __init__(self, path, *, format=None, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT): + super().__init__() + self._reader = _Reader.from_path(os.fspath(path), format, threads, memory_limit) + + def readable(self): return True + + def read(self, size=-1): + self._checkClosed() + return self._reader.read(size) + + def readinto(self, buffer): + self._checkClosed() + return self._reader.readinto(buffer) + + def close(self): + self._reader = None + super().close() + +def open(path, *, format=None, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT): + "Open a bzip2, gzip, or LZ4 path as a streaming binary file." + return Reader(path, format=format, threads=threads, memory_limit=memory_limit) + +def open_indexed(source, *, threads=0, index=None, memory_limit=DEFAULT_MEMORY_LIMIT, cache_limit=DEFAULT_CACHE_LIMIT): + "Open a path or bytes object as a seekable indexed bzip2 binary file." return IndexedBzip2File(source, threads=threads, index=index, memory_limit=memory_limit, cache_limit=cache_limit) def build_index(source, path=None, *, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT) -> bytes: @@ -92,13 +117,12 @@ def build_index(source, path=None, *, threads=0, memory_limit=DEFAULT_MEMORY_LIM if path is not None: Path(path).write_bytes(encoded) return encoded -def test(source, *, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT): - "Fully decode and CRC-validate *source*, returning ``None`` on success." - if isinstance(source, (bytes, bytearray, memoryview)): - _IndexedReader.from_bytes(bytes(source), threads, memory_limit, None, DEFAULT_CACHE_LIMIT) - else: _test(os.fspath(source), threads, memory_limit) +def test(source, format=None, *, threads=0, memory_limit=DEFAULT_MEMORY_LIMIT): + "Fully decode and validate a bzip2, gzip, or LZ4 *source*, returning ``None`` on success." + if isinstance(source, (bytes, bytearray, memoryview)): _test_bytes(bytes(source), format, threads, memory_limit) + else: _test_path(os.fspath(source), format, threads, memory_limit) __all__ = [ - "__version__", "BadBzip2File", "BlockCandidate", "DEFAULT_CACHE_LIMIT", "DEFAULT_MEMORY_LIMIT", "EndCandidate", - "IndexedBzip2File", "ScanResult", "StreamHeaderCandidate", "build_index", "bz2_crc32", "compress", "decompress", "open", "scan", "test" + "__version__", "BadCompressedFile", "BlockCandidate", "DEFAULT_CACHE_LIMIT", "DEFAULT_MEMORY_LIMIT", "EndCandidate", + "IndexedBzip2File", "Reader", "ScanResult", "StreamHeaderCandidate", "build_index", "bz2_crc32", "compress", "decompress", "open", "open_indexed", "scan", "test" ] diff --git a/src/bin/fbz.rs b/src/bin/fbz.rs index d9d7cdc..ac968eb 100644 --- a/src/bin/fbz.rs +++ b/src/bin/fbz.rs @@ -129,7 +129,7 @@ fn run(cli: Cli) -> fbz::Result<()> { if cli.compress { return run_compress(&cli); } - let options = DecodeOptions { threads: cli.threads, memory_limit: cli.memory_limit }; + let options = DecodeOptions { threads: cli.threads, memory_limit: cli.memory_limit, ..DecodeOptions::default() }; if cli.test { for input in &cli.inputs { test_input(input, options, cli.max_output, cli.quiet)?; diff --git a/src/decode.rs b/src/decode.rs index 9720427..9652525 100644 --- a/src/decode.rs +++ b/src/decode.rs @@ -5,8 +5,8 @@ use rayon::{ThreadPool, ThreadPoolBuilder}; use crate::format::scan_with_pool; use crate::pipeline::{Job, OrderedResults, PipelineLimits, run_ordered}; use crate::{ - BlockCandidate, BlockIndex, DecodeError, EndCandidate, Error, Index, MAX_DECODED_BLOCK, OutputSink, Result, StreamIndex, WriterSink, combine_stream_crc, - decoder, + BlockCandidate, BlockIndex, DecodeError, DecodeFormat, EndCandidate, Error, Index, MAX_DECODED_BLOCK, OutputSink, Result, StreamIndex, WriterSink, + combine_stream_crc, decoder, }; pub const DEFAULT_MEMORY_LIMIT: usize = 1024 * 1024 * 1024; @@ -19,6 +19,8 @@ pub struct DecodeProgress { #[derive(Clone, Copy, Debug)] pub struct DecodeOptions { + /// Compression format, or automatic magic/filename detection. + pub format: DecodeFormat, /// Zero selects the process's available parallelism. pub threads: usize, /// Maximum decoded bytes reserved for in-flight and completed speculative blocks. @@ -27,7 +29,7 @@ pub struct DecodeOptions { impl Default for DecodeOptions { fn default() -> Self { - Self { threads: 0, memory_limit: DEFAULT_MEMORY_LIMIT } + Self { format: DecodeFormat::Auto, threads: 0, memory_limit: DEFAULT_MEMORY_LIMIT } } } @@ -435,7 +437,7 @@ mod tests { compressed.extend_from_slice(&compress(&plain, Level::BEST)); expected.extend_from_slice(&plain); } - let options = DecodeOptions { threads: 4, memory_limit: MAX_DECODED_BLOCK * 2 }; + let options = DecodeOptions { threads: 4, memory_limit: MAX_DECODED_BLOCK * 2, ..DecodeOptions::default() }; assert_eq!(decompress(&compressed, options).unwrap(), expected); } diff --git a/src/gzip.rs b/src/gzip.rs index 80d7cfc..01e40b4 100644 --- a/src/gzip.rs +++ b/src/gzip.rs @@ -1243,8 +1243,6 @@ fn invalid_bit(bit: usize, message: impl Into) -> Error { #[cfg(test)] mod tests { - use std::io::Write as _; - use flate2::{ Compression, GzBuilder, write::{DeflateEncoder, GzEncoder}, diff --git a/src/lib.rs b/src/lib.rs index f852562..244d4db 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -30,8 +30,9 @@ pub use block::{MAX_DECODED_BLOCK, MAX_ENCODED_BLOCK, decode_block}; pub use bz2_encode::{EncodeReport as Bzip2EncodeReport, Encoder as Bzip2Encoder, compress as compress_bzip2, compress_to_writer as compress_bzip2_to_writer}; pub use crc::{bz2_crc32, combine_stream_crc}; pub use decode::{ - DEFAULT_MEMORY_LIMIT, DecodeOptions, DecodeProgress, build_index, build_index_with_progress, decode_to_writer, decode_to_writer_with_progress, decompress, - decompress_to_sink_with_progress, decompress_to_writer, decompress_to_writer_with_progress, + DEFAULT_MEMORY_LIMIT, DecodeOptions, DecodeProgress, build_index, build_index_with_progress, decode_to_writer, decode_to_writer_with_progress, + decompress as decompress_bzip2, decompress_to_sink_with_progress, decompress_to_writer as decompress_bzip2_to_writer, + decompress_to_writer_with_progress as decompress_bzip2_to_writer_with_progress, }; pub use encode::{EncodeFormat, EncodeOptions, EncodeProgress, EncodeReport, Encoder, compress, compress_to_writer, compress_to_writer_with_progress}; pub use error::{DecodeError, Error, Result}; @@ -41,7 +42,7 @@ pub use indexed::{DEFAULT_CACHE_LIMIT, IndexedReader}; pub use output::{OutputSink, PipeReader, PipeWriter, WriterSink, output_pipe}; pub use reader::Reader; pub use source::Source; -pub use stream::{Format, decode_stream_to_sink_with_progress}; +pub use stream::{DecodeFormat, Format, decode_stream_to_sink_with_progress, decompress, decompress_to_writer, decompress_to_writer_with_progress}; #[cfg(feature = "python")] mod python { @@ -51,19 +52,42 @@ mod python { }; use pyo3::{ + buffer::PyBuffer, create_exception, - exceptions::{PyOSError, PyValueError}, + exceptions::{PyOSError, PyTypeError, PyValueError}, prelude::*, types::PyBytes, }; - create_exception!(fbz, BadBzip2File, PyOSError); + create_exception!(fbz, BadCompressedFile, PyOSError); + + fn decode_format(format: Option<&str>) -> PyResult { + match format { + None | Some("auto") => Ok(crate::DecodeFormat::Auto), + Some("bzip2") => Ok(crate::DecodeFormat::Bzip2), + Some("gzip") => Ok(crate::DecodeFormat::Gzip), + Some("lz4") => Ok(crate::DecodeFormat::Lz4), + Some(_) => Err(PyValueError::new_err("format must be 'bzip2', 'gzip', 'lz4', or None")), + } + } + + fn decode_options(format: Option<&str>, threads: usize, memory_limit: usize) -> PyResult { + Ok(crate::DecodeOptions { format: decode_format(format)?, threads, memory_limit }) + } fn python_error(error: crate::Error) -> PyErr { match error { crate::Error::Io(source) => PyOSError::new_err(source.to_string()), crate::Error::InvalidConfiguration(message) | crate::Error::InvalidIndex(message) => PyValueError::new_err(message), - error => BadBzip2File::new_err(error.to_string()), + error => BadCompressedFile::new_err(error.to_string()), + } + } + + fn python_read_error(error: std::io::Error) -> PyErr { + match error.kind() { + std::io::ErrorKind::InvalidData => BadCompressedFile::new_err(error.to_string()), + std::io::ErrorKind::InvalidInput => PyValueError::new_err(error.to_string()), + _ => PyOSError::new_err(error.to_string()), } } @@ -82,10 +106,11 @@ mod python { crate::bz2_crc32(data) } - #[pyfunction(name = "_decompress", signature = (data, threads=0, memory_limit=crate::DEFAULT_MEMORY_LIMIT))] - fn py_decompress(py: Python<'_>, data: &[u8], threads: usize, memory_limit: usize) -> PyResult> { + #[pyfunction(name = "_decompress", signature = (data, format=None, threads=0, memory_limit=crate::DEFAULT_MEMORY_LIMIT))] + fn py_decompress(py: Python<'_>, data: &[u8], format: Option<&str>, threads: usize, memory_limit: usize) -> PyResult> { + let options = decode_options(format, threads, memory_limit)?; let data = data.to_vec(); - let output = py.detach(move || crate::decompress(&data, crate::DecodeOptions { threads, memory_limit })).map_err(python_error)?; + let output = py.detach(move || crate::decompress(&data, options)).map_err(python_error)?; Ok(PyBytes::new(py, &output).unbind()) } @@ -107,22 +132,88 @@ mod python { let encoded = py .detach(move || { let source = crate::Source::open(path)?; - Ok::<_, crate::Error>(crate::build_index(source.as_slice(), crate::DecodeOptions { threads, memory_limit })?.to_bytes()) + Ok::<_, crate::Error>( + crate::build_index(source.as_slice(), crate::DecodeOptions { threads, memory_limit, ..crate::DecodeOptions::default() })?.to_bytes(), + ) }) .map_err(python_error)?; Ok(PyBytes::new(py, &encoded).unbind()) } - #[pyfunction(name = "_test", signature = (path, threads=0, memory_limit=crate::DEFAULT_MEMORY_LIMIT))] - fn py_test(py: Python<'_>, path: String, threads: usize, memory_limit: usize) -> PyResult<()> { + #[pyfunction(name = "_test_path", signature = (path, format=None, threads=0, memory_limit=crate::DEFAULT_MEMORY_LIMIT))] + fn py_test_path(py: Python<'_>, path: String, format: Option<&str>, threads: usize, memory_limit: usize) -> PyResult<()> { + let options = decode_options(format, threads, memory_limit)?; py.detach(move || { - let source = crate::Source::open(path)?; - crate::build_index(source.as_slice(), crate::DecodeOptions { threads, memory_limit })?; - Ok::<_, crate::Error>(()) + let source = crate::Source::open(&path)?; + let stream_format = options.format.detect_path(std::path::Path::new(&path), source.as_slice())?; + let mut output = crate::WriterSink::new(std::io::sink()); + crate::decode_stream_to_sink_with_progress(stream_format, source.as_slice(), &mut output, options, |_| {}) }) .map_err(python_error) } + #[pyfunction(name = "_test_bytes", signature = (data, format=None, threads=0, memory_limit=crate::DEFAULT_MEMORY_LIMIT))] + fn py_test_bytes(py: Python<'_>, data: &[u8], format: Option<&str>, threads: usize, memory_limit: usize) -> PyResult<()> { + let options = decode_options(format, threads, memory_limit)?; + let data = data.to_vec(); + py.detach(move || crate::decompress_to_writer(&data, &mut std::io::sink(), options)).map_err(python_error) + } + + #[pyclass(name = "_Reader")] + struct PyReader { + inner: Mutex, + } + + const PY_READINTO_CHUNK: usize = 1024 * 1024; + + #[pymethods] + impl PyReader { + #[staticmethod] + #[pyo3(signature = (path, format=None, threads=0, memory_limit=crate::DEFAULT_MEMORY_LIMIT))] + fn from_path(path: String, format: Option<&str>, threads: usize, memory_limit: usize) -> PyResult { + let options = decode_options(format, threads, memory_limit)?; + Ok(Self { inner: Mutex::new(crate::Reader::open(path, options).map_err(python_error)?) }) + } + + fn read(&self, py: Python<'_>, size: i64) -> PyResult> { + let output = py + .detach(|| { + let mut reader = self.inner.lock().map_err(|_| std::io::Error::other("reader lock poisoned"))?; + if size < 0 { + let mut output = Vec::new(); + reader.read_to_end(&mut output)?; + Ok(output) + } else { + let count = usize::try_from(size) + .map_err(|_| std::io::Error::new(std::io::ErrorKind::InvalidInput, "requested read does not fit this platform"))?; + let mut output = vec![0; count]; + let count = reader.read(&mut output)?; + output.truncate(count); + Ok(output) + } + }) + .map_err(python_read_error)?; + Ok(PyBytes::new(py, &output).unbind()) + } + + fn readinto(&self, py: Python<'_>, buffer: PyBuffer) -> PyResult { + if buffer.as_mut_slice(py).is_none() { + return Err(PyTypeError::new_err("readinto() requires a writable, contiguous byte buffer")); + } + let mut output = vec![0; buffer.item_count().min(PY_READINTO_CHUNK)]; + let count = py + .detach(|| { + let mut reader = self.inner.lock().map_err(|_| std::io::Error::other("reader lock poisoned"))?; + reader.read(&mut output) + }) + .map_err(python_read_error)?; + for (destination, byte) in buffer.as_mut_slice(py).unwrap()[..count].iter().zip(&output) { + destination.set(*byte); + } + Ok(count) + } + } + #[pyclass(name = "_IndexedReader")] struct PyIndexedReader { inner: Mutex, @@ -136,7 +227,12 @@ mod python { let inner = py .detach(move || match index_path { Some(index_path) => crate::IndexedReader::open_with_index(path, index_path, cache_limit), - None => crate::IndexedReader::from_source(crate::Source::open(path)?, None, crate::DecodeOptions { threads, memory_limit }, cache_limit), + None => crate::IndexedReader::from_source( + crate::Source::open(path)?, + None, + crate::DecodeOptions { threads, memory_limit, ..crate::DecodeOptions::default() }, + cache_limit, + ), }) .map_err(python_error)?; Ok(Self { inner: Mutex::new(inner) }) @@ -150,9 +246,12 @@ mod python { let inner = py .detach(move || match index { Some(index) => crate::IndexedReader::from_bytes_with_index(data, &index, cache_limit), - None => { - crate::IndexedReader::from_source(crate::Source::from_bytes(data), None, crate::DecodeOptions { threads, memory_limit }, cache_limit) - } + None => crate::IndexedReader::from_source( + crate::Source::from_bytes(data), + None, + crate::DecodeOptions { threads, memory_limit, ..crate::DecodeOptions::default() }, + cache_limit, + ), }) .map_err(python_error)?; Ok(Self { inner: Mutex::new(inner) }) @@ -211,9 +310,11 @@ mod python { m.add_function(wrap_pyfunction!(py_decompress, m)?)?; m.add_function(wrap_pyfunction!(py_compress, m)?)?; m.add_function(wrap_pyfunction!(py_build_index, m)?)?; - m.add_function(wrap_pyfunction!(py_test, m)?)?; + m.add_function(wrap_pyfunction!(py_test_path, m)?)?; + m.add_function(wrap_pyfunction!(py_test_bytes, m)?)?; + m.add_class::()?; m.add_class::()?; - m.add("BadBzip2File", m.py().get_type::())?; + m.add("BadCompressedFile", m.py().get_type::())?; m.add("__version__", env!("CARGO_PKG_VERSION"))?; Ok(()) } diff --git a/src/reader.rs b/src/reader.rs index 32872ec..74862da 100644 --- a/src/reader.rs +++ b/src/reader.rs @@ -52,7 +52,7 @@ impl Reader { let options = options.validate()?; let path = path.as_ref(); let source = Source::open(path)?; - let format = Format::detect(path, source.as_slice())?; + let format = options.format.detect_path(path, source.as_slice())?; if !matches!(format, Format::Bzip2 | Format::Gzip | Format::Lz4) { return Err(Error::UnsupportedFormat("fbz::Reader supports bzip2, gzip, and LZ4 streams".into())); } diff --git a/src/stream.rs b/src/stream.rs index a9f2981..fab4702 100644 --- a/src/stream.rs +++ b/src/stream.rs @@ -1,6 +1,41 @@ -use std::path::Path; +use std::{io::Write, path::Path}; -use crate::{DecodeOptions, DecodeProgress, Error, OutputSink, Result, gzip, lz4}; +use crate::{DecodeOptions, DecodeProgress, Error, OutputSink, Result, WriterSink, gzip, lz4}; + +/// A byte-stream compression format accepted by the unified decoders. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub enum DecodeFormat { + /// Detect the format from magic, falling back to the filename for path inputs. + #[default] + Auto, + Bzip2, + Gzip, + Lz4, +} + +impl DecodeFormat { + fn explicit(self) -> Option { + match self { + Self::Auto => None, + Self::Bzip2 => Some(Format::Bzip2), + Self::Gzip => Some(Format::Gzip), + Self::Lz4 => Some(Format::Lz4), + } + } + + pub(crate) fn detect_data(self, data: &[u8]) -> Result { + self.explicit().or_else(|| Format::from_magic(data)).ok_or_else(|| { + Error::UnsupportedFormat("cannot determine compression format; expected bzip2, gzip, or LZ4 magic, or an explicit DecodeFormat".into()) + }) + } + + pub(crate) fn detect_path(self, path: &Path, data: &[u8]) -> Result { + match self.explicit() { + Some(format) => Ok(format), + None => Format::detect(path, data), + } + } +} /// A compression or archive format recognized by fbz. #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -71,3 +106,22 @@ pub fn decode_stream_to_sink_with_progress( Format::Zip => Err(Error::UnsupportedFormat("ZIP archives do not have one decoded byte stream".into())), } } + +/// Decode bzip2, gzip, or LZ4 bytes selected by magic or `options.format`. +pub fn decompress(data: &[u8], options: DecodeOptions) -> Result> { + let mut output = Vec::new(); + decompress_to_writer(data, &mut output, options)?; + Ok(output) +} + +/// Decode bzip2, gzip, or LZ4 bytes to a streaming output. +pub fn decompress_to_writer(data: &[u8], output: &mut impl Write, options: DecodeOptions) -> Result<()> { + decompress_to_writer_with_progress(data, output, options, |_| {}) +} + +/// Decode a byte stream and report completed compressed and decoded bytes. +pub fn decompress_to_writer_with_progress(data: &[u8], output: &mut impl Write, options: DecodeOptions, progress: impl FnMut(DecodeProgress)) -> Result<()> { + let format = options.format.detect_data(data)?; + let mut output = WriterSink::new(output); + decode_stream_to_sink_with_progress(format, data, &mut output, options, progress) +} diff --git a/tests/encode.rs b/tests/encode.rs index de5b17b..e185b41 100644 --- a/tests/encode.rs +++ b/tests/encode.rs @@ -1,13 +1,9 @@ use std::io::Write; -use fbz::{DecodeOptions, EncodeFormat, EncodeOptions, Encoder, Format, compress, decompress, gzip, lz4}; +use fbz::{DecodeOptions, EncodeFormat, EncodeOptions, Encoder, Format, compress, decompress}; -fn decode(format: EncodeFormat, encoded: &[u8]) -> Vec { - match format { - EncodeFormat::Bzip2 => decompress(encoded, DecodeOptions::default()).unwrap(), - EncodeFormat::Gzip => gzip::decompress(encoded).unwrap(), - EncodeFormat::Lz4 => lz4::decompress(encoded).unwrap(), - } +fn decode(encoded: &[u8]) -> Vec { + decompress(encoded, DecodeOptions::default()).unwrap() } #[test] @@ -16,7 +12,7 @@ fn unified_encoder_covers_every_stream_format() { for (format, magic) in [(EncodeFormat::Bzip2, Format::Bzip2), (EncodeFormat::Gzip, Format::Gzip), (EncodeFormat::Lz4, Format::Lz4)] { let encoded = compress(&input, format, EncodeOptions { threads: 3, memory_limit: 128 * 1024 * 1024, level: None }).unwrap(); assert_eq!(Format::from_magic(&encoded), Some(magic)); - assert_eq!(decode(format, &encoded), input); + assert_eq!(decode(&encoded), input); let mut incremental = Vec::new(); let mut encoder = Encoder::new(&mut incremental, format, EncodeOptions { threads: 2, memory_limit: 128 * 1024 * 1024, level: None }).unwrap(); @@ -26,6 +22,6 @@ fn unified_encoder_covers_every_stream_format() { let (_, report) = encoder.finish().unwrap(); assert_eq!(report.format, format); assert_eq!(report.input_len, input.len() as u64); - assert_eq!(decode(format, &incremental), input); + assert_eq!(decode(&incremental), input); } } diff --git a/tests/gzip_oracle.rs b/tests/gzip_oracle.rs index d286157..63b2a2b 100644 --- a/tests/gzip_oracle.rs +++ b/tests/gzip_oracle.rs @@ -61,7 +61,7 @@ fn large_stored_gzip_matches_oracle() { let plain = random_bytes(20 * 1024 * 1024); let encoded = compress(&plain); assert!(encoded.len() >= 16 * 1024 * 1024); - let options = DecodeOptions { threads: 4, memory_limit: 256 * 1024 * 1024 }; + let options = DecodeOptions { threads: 4, memory_limit: 256 * 1024 * 1024, ..DecodeOptions::default() }; assert_eq!(gzip::decompress_with_options(&encoded, options).unwrap(), oracle_decompress(&encoded)); } @@ -74,7 +74,7 @@ fn parallel_dynamic_gzip_matches_oracle() { encoded.extend(compress(&stored)); plain.extend(stored); assert!(encoded.len() >= 32 * 1024 * 1024); - let options = DecodeOptions { threads: 4, memory_limit: 256 * 1024 * 1024 }; + let options = DecodeOptions { threads: 4, memory_limit: 256 * 1024 * 1024, ..DecodeOptions::default() }; let mut decoded = Vec::new(); let report = gzip::decompress_to_writer_with_options(&encoded, &mut decoded, options).unwrap(); assert!(report.speculative_chunks > 0); diff --git a/tests/lz4_oracle.rs b/tests/lz4_oracle.rs index 133f456..b80e6dc 100644 --- a/tests/lz4_oracle.rs +++ b/tests/lz4_oracle.rs @@ -79,5 +79,5 @@ fn incompressible_multiblock_input_exercises_parallel_scheduler() { fn a_small_speculative_budget_falls_back_to_incremental_serial_blocks() { let data = b"large declared blocks do not require speculative memory ".repeat(120_000); let encoded = encode(&data, BlockSize::Max4MB, BlockMode::Independent, true, true, true); - assert_matches(&data, &encoded, DecodeOptions { threads: 4, memory_limit: MAX_DECODED_BLOCK }); + assert_matches(&data, &encoded, DecodeOptions { threads: 4, memory_limit: MAX_DECODED_BLOCK, ..DecodeOptions::default() }); } diff --git a/tests/reader.rs b/tests/reader.rs index b0d13c4..f798cbe 100644 --- a/tests/reader.rs +++ b/tests/reader.rs @@ -5,7 +5,7 @@ use std::{ }; use crabz2::{Level, compress}; -use fbz::{DecodeOptions, Error, Reader}; +use fbz::{DecodeFormat, DecodeOptions, Error, Reader}; use flate2::{Compression, write::GzEncoder}; use lz4_flex::frame::{BlockMode, BlockSize, FrameEncoder, FrameInfo}; @@ -52,6 +52,20 @@ fn reader_detects_stream_formats_without_prevalidation() { assert_eq!(reader.read(&mut [0; 1]).unwrap_err().kind(), std::io::ErrorKind::InvalidData); } +#[test] +fn explicit_reader_format_overrides_magic_and_extension_detection() { + let directory = tempfile::tempdir().unwrap(); + let path = directory.path().join("misleading.bz2"); + let plain = b"explicit gzip selection".repeat(10_000); + fs::write(&path, gzip(&plain)).unwrap(); + let options = DecodeOptions { format: DecodeFormat::Gzip, threads: 2, ..DecodeOptions::default() }; + assert_eq!(read(&path, options).unwrap(), plain); + + let options = DecodeOptions { format: DecodeFormat::Bzip2, ..DecodeOptions::default() }; + let mut reader = Reader::open(path, options).unwrap(); + assert_eq!(reader.read(&mut [0; 1]).unwrap_err().kind(), std::io::ErrorKind::InvalidData); +} + #[test] fn lz4_reader_emits_valid_blocks_before_a_later_parse_error() { let directory = tempfile::tempdir().unwrap(); diff --git a/tests/test_api.py b/tests/test_api.py index d506449..ca716fe 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -12,6 +12,13 @@ def test_decompress_matches_libbz2_at_every_level(level): plain = patterned(20_000) assert fbz.decompress(bz2.compress(plain, compresslevel=level), threads=2) == plain +def test_decompress_detects_every_stream_format(): + plain = patterned(300_000) + encoded = [bz2.compress(plain), gzip.compress(plain), fbz.compress(plain, "lz4")] + for item in encoded: assert fbz.decompress(item, threads=2) == plain + assert fbz.decompress(encoded[1], "gzip", threads=2) == plain + with pytest.raises(ValueError, match="format"): fbz.decompress(encoded[0], "zip") + def test_compress_stream_formats_and_options(): plain = patterned(300_000) assert bz2.decompress(fbz.compress(plain, "bzip2", threads=2, level=3)) == plain @@ -28,7 +35,7 @@ def test_parallel_multiblock_and_concatenated_are_deterministic(): def test_bad_crc_raises_package_error(): compressed = bytearray(bz2.compress(b"integrity matters")) compressed[-2] ^= 1 - with pytest.raises(fbz.BadBzip2File): fbz.decompress(bytes(compressed)) + with pytest.raises(fbz.BadCompressedFile): fbz.decompress(bytes(compressed)) def test_seekable_file_and_persisted_index(tmp_path): plain = patterned(350_000) @@ -38,7 +45,7 @@ def test_seekable_file_and_persisted_index(tmp_path): source.write_bytes(compressed) encoded = fbz.build_index(source, index_path, threads=2) - with fbz.open(source, index=index_path) as handle: + with fbz.open_indexed(source, index=index_path) as handle: assert handle.size == len(plain) assert handle.seek(99_990) == 99_990 assert handle.read(40) == plain[99_990:100_030] @@ -53,14 +60,46 @@ def test_index_is_bound_to_source(): first = bz2.compress(b"first") second = bz2.compress(b"other") index = fbz.build_index(first) - with pytest.raises(ValueError, match="source identity mismatch"): fbz.open(second, index=index) + with pytest.raises(ValueError, match="source identity mismatch"): fbz.open_indexed(second, index=index) -def test_buffered_reader_compatibility(): +@pytest.mark.parametrize("format", ["bzip2", "gzip", "lz4"]) +def test_buffered_reader_compatibility(tmp_path, format): plain = patterned(180_000) - with io.BufferedReader(fbz.open(bz2.compress(plain, compresslevel=1))) as handle: + source = tmp_path / f"data.{format}" + source.write_bytes(fbz.compress(plain, format, threads=2)) + with io.BufferedReader(fbz.open(source, threads=2)) as handle: assert handle.read(1234) == plain[:1234] - handle.seek(100_000) - assert handle.read() == plain[100_000:] + assert handle.read() == plain[1234:] + assert handle.closed + +def test_stream_reader_readinto_and_explicit_format(tmp_path): + plain = patterned(2_000_000) + source = tmp_path / "misleading.bz2" + source.write_bytes(gzip.compress(plain)) + target = bytearray(len(plain) + 100) + with fbz.Reader(source, format="gzip", threads=2) as handle: + count = handle.readinto(memoryview(target)[17:]) + assert count <= 1024 * 1024 + assert target[17:17 + count] == plain[:count] + assert bytes(target[17:17 + count]) + handle.read() == plain + +def test_stream_reader_reports_late_checksum_failure(tmp_path): + plain = patterned(1_000_000) + encoded = bytearray(gzip.compress(plain)) + encoded[-8] ^= 1 + source = tmp_path / "corrupt.gz" + source.write_bytes(encoded) + with fbz.open(source, threads=2) as handle: + with pytest.raises(fbz.BadCompressedFile): handle.read() + +def test_validates_all_stream_formats_from_paths_and_bytes(tmp_path): + plain = patterned(300_000) + for format in ("bzip2", "gzip", "lz4"): + encoded = fbz.compress(plain, format, threads=2) + source = tmp_path / f"data.{format}" + source.write_bytes(encoded) + assert fbz.test(encoded, threads=2) is None + assert fbz.test(source, threads=2) is None def test_native_decode_releases_gil(): plain = patterned(2_000_000)