Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions DEV.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand All @@ -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.

Expand Down Expand Up @@ -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.
27 changes: 20 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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)
Expand All @@ -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

Expand All @@ -189,7 +202,7 @@ fn main() -> fbz::Result<()> {
}
```

`decompress` returns a `Vec<u8>`. `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<u8>` 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:

Expand Down Expand Up @@ -243,7 +256,7 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
}
```

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.

Expand Down
50 changes: 37 additions & 13 deletions python/fbz/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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."
Expand Down Expand Up @@ -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:
Expand All @@ -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"
]
2 changes: 1 addition & 1 deletion src/bin/fbz.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?;
Expand Down
10 changes: 6 additions & 4 deletions src/decode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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.
Expand All @@ -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 }
}
}

Expand Down Expand Up @@ -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);
}

Expand Down
2 changes: 0 additions & 2 deletions src/gzip.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1243,8 +1243,6 @@ fn invalid_bit(bit: usize, message: impl Into<String>) -> Error {

#[cfg(test)]
mod tests {
use std::io::Write as _;

use flate2::{
Compression, GzBuilder,
write::{DeflateEncoder, GzEncoder},
Expand Down
Loading