Skip to content

Stream PDF from archives - #117

Open
lfoppiano wants to merge 5 commits into
masterfrom
feature/stream-archive-input
Open

Stream PDF from archives #117
lfoppiano wants to merge 5 commits into
masterfrom
feature/stream-archive-input

Conversation

@lfoppiano

@lfoppiano lfoppiano commented Jul 25, 2026

Copy link
Copy Markdown
Member
  • Add stream of PDF from zip archives
  • Add support for Stream from S3 / HF buckets
  • Add support GLOB on --input parameters

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
lfoppiano force-pushed the feature/stream-archive-input branch from 5400cf2 to 87b7e95 Compare August 11, 2026 21:01
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
lfoppiano force-pushed the feature/stream-archive-input branch from 87b7e95 to ac5b7db Compare August 11, 2026 21:12
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant