Stream PDF from archives - #117
Open
lfoppiano wants to merge 5 commits into
Open
Conversation
Add process_archive(), used automatically when --input points to a
.zip/.tar/.tar.gz (.tgz/.tar.bz2/.tbz2) archive. Eligible entries are read out
of the archive in chunks of batch_size, each chunk is extracted to a temporary
directory, sent to GROBID via the existing process_batch (so concurrency,
TEI/JSON/Markdown output, --force/skip and error handling are reused), and the
temporary files are removed before the next chunk is extracted. The archive is
never fully decompressed, so disk usage stays bounded regardless of its size.
- zip via zipfile, tar/tar.gz/tar.bz2 via tarfile ('r:*')
- entries are streamed member-by-member (zipfile.open / tarfile.extractfile)
- archive paths are sanitized to prevent path traversal (zip slip)
- when --output is omitted, results go to a directory named after the archive
- directory-input eligibility check refactored into _is_eligible_input and shared
Documented in the Readme and covered by unit tests (zip, tar.gz, chunking,
cleanup, default output, traversal guard, delegation).
…rsal --input now accepts shell-style glob patterns (with recursive **), e.g. 'paper.zip' (one file), 'paper*.zip' (many), '**/paper*.zip' (subdirectories) or '**/*.pdf'. Each match is dispatched by type: archives are streamed, directories are recursed, eligible files are processed directly, and the results of all matches are aggregated into a single summary. - resolve --input via glob (has_magic + recursive=True, ~ expansion); a plain path is returned unchanged for backward compatibility - directory traversal refactored from os.walk to pathlib.Path.rglob - factor the batching loop, stats summary and archive streaming into reusable helpers (_run_file_batches, _print_processing_summary, _process_archive_core) so directory, loose-file and archive inputs share one code path - loose files matched by a glob are batched together under their common base Documented in the Readme and covered by unit tests (multi-archive glob, recursive **/*.pdf, mixed matches, no-match warning, common-base helper).
lfoppiano
force-pushed
the
feature/stream-archive-input
branch
from
August 11, 2026 21:01
5400cf2 to
87b7e95
Compare
Add s3:// support to --input (and a new --input-list manifest). An s3 zip is range-streamed with smart_open (only the central directory and the requested entries are fetched - the object is never fully downloaded); loose remote PDFs are fetched a batch at a time to a temp dir. Mixed manifests (local + glob + s3) are supported and aggregated into one summary. - refactor process() -> process_paths(list of inputs); process() delegates - _resolve_input_paths handles s3 object / prefix / glob (list_objects_v2 + fnmatch) - _open_archive range-streams s3 zips (smart_open seekable stream); the stream is closed after use; s3 tar is rejected (not range-streamable) - new _process_remote_files streams loose s3 objects in bounded chunks - --input-list reads a file of paths (local/glob/s3, '#' comments) - s3 deps (smart_open[s3], boto3) are an optional 'pip install ...[s3]' extra, lazily imported with a clear install hint if missing - credentials use the standard AWS chain (env / ~/.aws / IAM) Documented in the Readme; covered by moto-backed tests (single object, prefix, glob, zip range-streaming, loose PDFs, mixed manifest, missing-extra error).
process_batch decides a document is already done with os.path.isfile() alone,
so any output truncated by a kill is indistinguishable from a complete one and
is skipped on every subsequent run. The corruption is permanent and silent: it
survives every resume, and `find -name '*.grobid.tei.xml' | wc -l` counts it as
a success. At corpus scale, on jobs that routinely hit a wall clock or a memory
limit, this is not a rare case.
_write_atomic writes to a temp file and renames, so the destination either does
not exist or is the whole document. Applied to all six write sites: TEI, the
error file, and the JSON/Markdown pairs on both the fresh and already-exists
paths.
Three details that are easy to get wrong:
- The temp file goes in the DESTINATION directory, not TMPDIR. os.replace is
only atomic within a filesystem, and on a cluster TMPDIR is usually a
different mount.
- mkstemp hardcodes 0600, where open() respects the umask. Renaming such a
file into place would leave every output private -- unreadable to the group
on shared scratch, which is where the results of a cluster run live. The
mode is taken from a umask read once at import, since reading it means
temporarily setting it and that is not safe from worker threads.
- fd ownership transfers to the file object on a successful fdopen, so the
error path closes the descriptor only when fdopen itself failed. Closing it
unconditionally could close an unrelated descriptor that reused the number,
and process_batch writes from a ThreadPoolExecutor.
The "." prefix and ".tmp" suffix keep an in-flight temp file from matching
*.grobid.tei.xml or *_[0-9]*.txt, so output counting is unaffected mid-write;
there is a test pinning that. The other tests cover the case that matters: a
writer that puts bytes on disk and then dies leaves no destination file and no
temp file, and a failed overwrite preserves the previous content.
Residual risk: a SIGKILL between mkstemp and rename leaks a temp file. That is
visible and harmless, unlike a truncated TEI.
The package ships a py.typed marker since #112, which promises inline types to type checkers of anything importing it; the methods added by this branch had none, so process_paths, process_archive and the S3 helpers came back untyped to every user of the client. The signatures now match the ones on master. _open_archive's handle stays Any on purpose: a ZipFile and a TarFile share no interface here, which is exactly why the function returns a "kind" tag for the callers to dispatch on. mypy --ignore-missing-imports is clean on the module, as it is on master.
lfoppiano
force-pushed
the
feature/stream-archive-input
branch
from
August 11, 2026 21:12
87b7e95 to
ac5b7db
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.