Skip to content

Add streaming support - #28

Open
hoshimura wants to merge 4 commits into
fix_inconsistenciesfrom
add_streaming_support
Open

hoshimura wants to merge 4 commits into
fix_inconsistenciesfrom
add_streaming_support

Conversation

@hoshimura

Copy link
Copy Markdown
Collaborator

Makes large writes and copies bounded-memory instead of buffering whole objects in RAM. Fixes two downstream OOM reports: ShardWriter holding many open("wb") handles, and large local→cloud copy(). This PR lands the foundation plus the first native backend (S3); other backends follow in PR 2.

Changes

Disk-backed write spool — _CloudFile write mode uses a SpooledTemporaryFile (16 MiB resident, then spills to disk) instead of an unbounded BytesIO, so many concurrently-open write handles can't accumulate whole objects in memory. Adds a subclass supplying readable/writable/seekable (Python 3.10's SpooledTemporaryFile lacks them under TextIOWrapper).
_upload_buffer hook — _CloudFile close hands the spooled buffer to path._upload_buffer(buf, size). Default implementation is a single PUT (unchanged for backends without an override).
Native S3 multipart — s3_create_multipart / s3_upload_part / s3_complete_multipart / s3_abort_multipart primitives; S3Path streams the buffer as 8 MiB parts above an 8 MiB threshold (AWS defaults), bounded to one part in memory, abort-on-error. Orchestration stays in Python (per repo convention); Rust provides thin signed requests.
Streamed local→cloud copy — local.py copy hands the source file handle to the destination's _upload_buffer instead of read_bytes() + write_bytes().

Validation

Unit tests for spool spill + multipart orchestration; validated end-to-end against MinIO (single-PUT, multi-part, remainder, 16 MB local→S3 copy — all sha256-verified).

@hoshimura
hoshimura added this pull request to stack #26 October 7, 2026 09:22
@hoshimura
hoshimura force-pushed the add_streaming_support branch from cf70a61 to 64a6b9d Compare October 7, 2026 09:40
- _CloudFile write buffer uses SpooledTemporaryFile (16 MiB resident, then spills to disk) instead of unbounded BytesIO, so many concurrently-open write handles (e.g. sharded writers) can't accumulate whole objects in memory and OOM
- subclass adds readable/writable/seekable (Python 3.10 SpooledTemporaryFile lacks them, needed by TextIOWrapper)
- close the temp file on exit; spill test added
- new s3_create_multipart/s3_upload_part/s3_complete_multipart/s3_abort_multipart primitives (thin signed S3 requests); orchestration stays in Python per repo convention
- _CloudFile close hands the spooled buffer to _upload_buffer; S3 streams it as 8 MiB parts (AWS CLI defaults) above an 8 MiB threshold, bounding memory to one part; abort-on-error
- CloudPathMixin._upload_buffer default = single PUT (other backends unchanged)
- validated end-to-end against MinIO (single-PUT, multi-part, remainder, large text spill); unit tests mock the primitives
- local.py copy/copy-tree now hand the source file handle to the dest's _upload_buffer (S3 streams it as multipart, bounded memory) instead of read_bytes()+write_bytes()
- fixes the OOM where copying a large media file required memory equal to the whole file
- validated against MinIO (16 MB local->s3 multipart copy, sha intact)
- remove duplicate import block in s3.py from fix_inconsistencies merge
- cargo fmt s3.rs multipart helpers
- clear AZURE_STORAGE_CONNECTION_STRING in presign-no-creds test
@hoshimura
hoshimura force-pushed the add_streaming_support branch from 64a6b9d to 65c8edd Compare October 7, 2026 10:04

This branch has not been deployed

No deployments
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