Skip to content
Open
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
47 changes: 45 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,49 @@ for child in p.parent.iterdir(): # yields synchronously
Local paths (`file://` or bare `/path`) bypass the proxy entirely and use
`pathlib.Path` directly — zero threading overhead.

### Large uploads stream automatically

Writing or copying a large object never buffers the whole thing in memory — it
streams in bounded-size parts using each backend's native mechanism (S3
multipart, Azure block blobs, GCS resumable, SSH chunked writes, FTP data
connection, Artifactory streamed `PUT`). This is the **default** behaviour — no
flag, no API change:

```python
# open("wb") spills to disk past 16 MiB and uploads in parts on close
async with AsAnyPath("s3://bucket/huge.bin").open("wb") as f:
f.write(data)

# local -> cloud copy streams straight from the source file (async and sync)
await AsAnyPath("/data/movie.mp4").copy("s3://bucket/movie.mp4")
```

Peak memory per upload is one part — **8 MiB by default**. Objects smaller than
one part upload in a single request. (`write_bytes(data)` / `write_text(...)`
take the whole payload in memory by definition — use `open("wb")` or `copy()`
for large data.)

**Tuning the memory footprint.** The part size *is* the max RAM held per
in-flight chunk (parts upload sequentially, so there is no concurrency
multiplier). Set it globally via env, or per call:

```bash
export ASANYPATH_UPLOAD_CHUNK_SIZE=33554432 # 32 MiB parts (bytes)
export ASANYPATH_SPOOL_MAX_SIZE=67108864 # 64 MiB before open("wb") spills to disk
```

```python
# per call: copy(chunk_size=...) or, for writes, open(..., buffering=...)
await src.copy("s3://bucket/huge.bin", chunk_size=32 * 1024 * 1024)
async with dst.open("wb", buffering=32 * 1024 * 1024) as f:
f.write(data)
```

Bigger parts mean fewer requests and a larger single-object ceiling (S3 allows
10,000 parts, so 8 MiB parts cap an object at ~80 GB). Backend minimums are
applied automatically: **S3** clamps to ≥ 5 MiB parts, **GCS** rounds down to a
256 KiB multiple.

Unknown URI schemes now construct as `UnsupportedProtocolPath` placeholders.
Pure path operations (e.g. `.name`, `.parent`, joins) still work, while backend
operations (e.g. `.exists()`, `.read_bytes()`, `.open()`) raise
Expand Down Expand Up @@ -383,8 +426,8 @@ Factory that returns the appropriate path instance based on the URI scheme.

All path implementations support:

- `open(mode, buffering)` — file-like streaming with lazy range reads
- `copy(dst, recursive)` — copy file or tree (cross-backend)
- `open(mode, buffering)` — file-like streaming with lazy range reads; for writes, `buffering` sets the upload part size
- `copy(dst, recursive, chunk_size)` — copy file or tree (cross-backend); `chunk_size` tunes the upload part size
- `read_bytes()` / `write_bytes(data)`
- `read_text(encoding)` / `write_text(data, encoding)`
- `exists()` / `is_file()` / `is_dir()`
Expand Down
20 changes: 16 additions & 4 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@ disable_error_code = ["attr-defined", "arg-type", "assignment"]
module = "asanypath.local"
disallow_untyped_defs = false
warn_return_any = false
disable_error_code = ["type-var", "misc", "override", "return-value"]
disable_error_code = ["type-var", "misc", "override", "return-value", "assignment", "no-redef"]

[[tool.mypy.overrides]]
module = "asanypath.common"
Expand All @@ -120,19 +120,31 @@ disable_error_code = ["override", "type-arg", "return-value", "assignment", "att
module = ["asanypath.s3", "asanypath.gcs", "asanypath.azure", "asanypath.artifactory"]
disallow_untyped_defs = false
warn_return_any = false
disable_error_code = ["no-any-return", "attr-defined", "assignment", "override"]
disable_error_code = ["no-any-return", "attr-defined", "assignment", "override", "arg-type", "misc"]

[[tool.mypy.overrides]]
module = ["asanypath.ssh", "asanypath.ftp"]
disallow_untyped_defs = false
warn_return_any = false
disable_error_code = ["attr-defined", "override", "arg-type", "misc", "no-any-return", "assignment"]

[[tool.mypy.overrides]]
module = "asanypath.sync"
disallow_untyped_defs = false
warn_return_any = false
disable_error_code = ["no-any-return", "misc", "attr-defined", "arg-type"]

[[tool.mypy.overrides]]
module = "asanypath.cloud"
disallow_untyped_defs = false
warn_return_any = false
disable_error_code = ["no-any-return", "arg-type"]
disable_error_code = ["no-any-return", "arg-type", "attr-defined", "misc", "assignment"]

[[tool.mypy.overrides]]
module = "asanypath.asanypath"
disallow_untyped_defs = false
warn_return_any = false
disable_error_code = ["no-any-return", "misc", "union-attr", "type-var"]
disable_error_code = ["no-any-return", "misc", "union-attr", "type-var", "attr-defined"]

[[tool.mypy.overrides]]
module = "asanypath.batcher"
Expand Down
20 changes: 14 additions & 6 deletions src/asanypath/artifactory.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,13 @@
)

from asanypath.cloud import CloudPathMixin
from asanypath.options import AccessGrant, AccessPolicy, AccessPolicyPatch, BackendOptions
from asanypath.options import (
UPLOAD_CHUNK_SIZE,
AccessGrant,
AccessPolicy,
AccessPolicyPatch,
BackendOptions,
)


def _jfrog_conf_token(base_url: str) -> str | None:
Expand Down Expand Up @@ -99,9 +105,9 @@ async def main():

protocol: str = "art"
_supports_range_read: bool = True
# Artifactory has no multipart API: large objects stream a single PUT body
# from disk (bounded memory) instead of buffering the whole object.
_STREAM_THRESHOLD = 8 * 1024 * 1024
# Artifactory has no multipart API: objects at/above this size stream a single
# PUT body from disk (bounded memory) instead of buffering the whole object.
_STREAM_THRESHOLD = UPLOAD_CHUNK_SIZE
_copy_batch_fn = staticmethod(art_copy_batch)

@staticmethod
Expand Down Expand Up @@ -201,9 +207,11 @@ def _native_kwargs(self) -> dict:
# Backend-specific operations
# ------------------------------------------------------------------

async def _upload_buffer(self, fileobj, size: int, *, backend_options=None) -> None:
async def _upload_buffer(
self, fileobj, size: int, *, backend_options=None, chunk_size: int | None = None
) -> None:
"""Upload a spooled write buffer; large objects stream a single PUT from disk."""
if size < self._STREAM_THRESHOLD:
if size < (chunk_size or self._STREAM_THRESHOLD):
data = fileobj.read()
if backend_options is None:
await self.write_bytes(data)
Expand Down
26 changes: 17 additions & 9 deletions src/asanypath/azure.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,13 @@
)

from asanypath.cloud import CloudPathMixin
from asanypath.options import AccessGrant, AccessPolicy, AccessPolicyPatch, BackendOptions
from asanypath.options import (
UPLOAD_CHUNK_SIZE,
AccessGrant,
AccessPolicy,
AccessPolicyPatch,
BackendOptions,
)

AZURE_API_VERSION = "2023-11-03"

Expand Down Expand Up @@ -85,9 +91,8 @@ class AzurePath(CloudPathMixin):

protocol: str = "az"
_supports_range_read: bool = True
# Single PUT below the threshold, else staged block-blob parts.
_MULTIPART_THRESHOLD = 8 * 1024 * 1024
_MULTIPART_PART_SIZE = 8 * 1024 * 1024
# Bytes per staged block; objects below one block upload in a single PUT.
_MULTIPART_PART_SIZE = UPLOAD_CHUNK_SIZE
_is_dir_batch_fn = staticmethod(az_is_dir_batch)
_copy_batch_fn = staticmethod(az_copy_batch)

Expand Down Expand Up @@ -186,26 +191,29 @@ def _native_kwargs(self) -> dict:
# Backend-specific operations
# ------------------------------------------------------------------

async def _upload_buffer(self, fileobj, size: int, *, backend_options=None) -> None:
async def _upload_buffer(
self, fileobj, size: int, *, backend_options=None, chunk_size: int | None = None
) -> None:
"""Upload a spooled write buffer, using block blobs for large objects."""
if size < self._MULTIPART_THRESHOLD:
part = chunk_size or self._MULTIPART_PART_SIZE
if size < part:
data = fileobj.read()
if backend_options is None:
await self.write_bytes(data)
else:
await self.write_bytes(data, backend_options=backend_options)
return
await self._block_upload(fileobj, backend_options=backend_options)
await self._block_upload(fileobj, part, backend_options=backend_options)

async def _block_upload(self, fileobj, *, backend_options=None) -> None:
async def _block_upload(self, fileobj, part_size: int, *, backend_options=None) -> None:
"""Stage fixed-size blocks then commit the list (memory bounded to one block)."""
options_json = (
msgspec.json.encode(backend_options).decode() if backend_options is not None else "{}"
)
block_ids: list[str] = []
index = 0
while True:
chunk = fileobj.read(self._MULTIPART_PART_SIZE)
chunk = fileobj.read(part_size)
if not chunk:
break
# Block ids must be equal-length base64 strings.
Expand Down
31 changes: 26 additions & 5 deletions src/asanypath/cloud.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@
)
from asanypath.batcher import H2_BATCH_THRESHOLD, MicroBatcher, make_batcher
from asanypath.common import SCHEME_SEP, CommonPurePathMixin
from asanypath.options import AccessPolicy, AccessPolicyPatch, BackendOptions
from asanypath.options import SPOOL_MAX_SIZE, AccessPolicy, AccessPolicyPatch, BackendOptions


class CloudPathMixin(CommonPurePathMixin):
Expand Down Expand Up @@ -354,12 +354,18 @@ async def write_text(
return await self.write_bytes(encoded, backend_options=backend_options)

async def _upload_buffer(
self, fileobj: IO[bytes], size: int, *, backend_options: BackendOptions | None = None
self,
fileobj: IO[bytes],
size: int,
*,
backend_options: BackendOptions | None = None,
chunk_size: int | None = None,
) -> None:
"""Upload a spooled write buffer (default: a single PUT of the whole buffer).

Backends with native multipart override this to bound memory for large
objects; ``size`` is the byte length and ``fileobj`` is positioned at 0.
``chunk_size`` tunes the part size for those backends (ignored here).
"""
data = fileobj.read()
if backend_options is None:
Expand Down Expand Up @@ -989,7 +995,7 @@ class _CloudFile:
_DEFAULT_CHUNK = 8 * 1024 * 1024 # 8 MiB — good default for cloud latency
# Write buffer spills to disk past this size so many concurrently-open write
# handles (e.g. sharded writers) can't accumulate whole objects in memory.
_SPOOL_MAX_SIZE = 16 * 1024 * 1024 # 16 MiB resident per open write handle
_SPOOL_MAX_SIZE = SPOOL_MAX_SIZE # resident bytes per open write handle before spilling

def __init__(self, path, mode, buffering, encoding, errors, newline, backend_options=None):
self._path = path
Expand All @@ -1012,6 +1018,11 @@ def _chunk_size(self) -> int:
return self._DEFAULT_CHUNK
return self._buffering

@property
def _write_chunk_size(self) -> int | None:
"""Write-mode: a positive ``buffering`` sets the upload part size."""
return self._buffering if self._buffering and self._buffering > 0 else None

def _make_buffer(self, data: bytes | None = None) -> io.IOBase:
"""Create the backing buffer: in-memory for reads, disk-spilling for writes."""
if "r" in self._mode:
Expand Down Expand Up @@ -1100,7 +1111,12 @@ def __exit__(self, exc_type, exc_val, exc_tb):
if ("w" in self._mode or "a" in self._mode) and exc_type is None:
buf, size = self._spool_for_upload()
_SyncRunner.get().run(
self._path._upload_buffer(buf, size, backend_options=self._backend_options)
self._path._upload_buffer(
buf,
size,
backend_options=self._backend_options,
chunk_size=self._write_chunk_size,
)
)
finally:
if self._reader is not None:
Expand Down Expand Up @@ -1137,7 +1153,12 @@ async def __aexit__(self, exc_type, exc_val, exc_tb):
try:
if ("w" in self._mode or "a" in self._mode) and exc_type is None:
buf, size = self._spool_for_upload()
await self._path._upload_buffer(buf, size, backend_options=self._backend_options)
await self._path._upload_buffer(
buf,
size,
backend_options=self._backend_options,
chunk_size=self._write_chunk_size,
)
finally:
if self._reader is not None:
self._reader.close()
Expand Down
40 changes: 26 additions & 14 deletions src/asanypath/ftp.py
Original file line number Diff line number Diff line change
Expand Up @@ -59,12 +59,19 @@
from datetime import datetime, timezone
from os import getenv
from types import SimpleNamespace
from typing import IO, TYPE_CHECKING, TypeVar
from typing import IO, TYPE_CHECKING, TypeVar, cast

import aioftp
from yarl import URL

from asanypath.cloud import CloudPathMixin
from asanypath.options import AccessGrant, AccessPolicy, AccessPolicyPatch, BackendOptions
from asanypath.options import (
UPLOAD_CHUNK_SIZE,
AccessGrant,
AccessPolicy,
AccessPolicyPatch,
BackendOptions,
)

if TYPE_CHECKING:
from typing import Self
Expand Down Expand Up @@ -155,10 +162,8 @@ class FTPPath(CloudPathMixin):
_tls: bool = False
_default_port: int = 21
_supports_range_read: bool = False
# Large objects stream to the data connection in chunks instead of buffering
# the whole payload in memory.
_STREAM_THRESHOLD = 8 * 1024 * 1024
_STREAM_CHUNK = 1024 * 1024
# Bytes per streamed read; objects below this write in a single call.
_STREAM_CHUNK = UPLOAD_CHUNK_SIZE

def __init__(
self,
Expand All @@ -171,10 +176,11 @@ def __init__(
) -> None:
super().__init__(*parts)
cfg = self.env_config
self._host = host or self._path.host or cfg.host
self._port = port or self._path.port or cfg.port
url_pw = self._path.password
url_user = self._path.user
url = cast(URL, self._path)
self._host = host or url.host or cfg.host
self._port = port or url.port or cfg.port
url_pw = url.password
url_user = url.user
# Username precedence: kwarg > URL > netrc > env default.
# Password precedence: kwarg > URL > netrc > anonymous default.
netrc_user, netrc_pw, netrc_account = _netrc_lookup(self._host)
Expand Down Expand Up @@ -218,7 +224,7 @@ def _create_env_config(cls) -> SimpleNamespace:

@property
def _item_path(self) -> str:
return self._path.path or "/"
return cast(URL, self._path).path or "/"

@property
def _native_kwargs(self) -> dict:
Expand Down Expand Up @@ -404,18 +410,24 @@ async def _write(client: aioftp.Client) -> int:
return await self._op(_write)

async def _upload_buffer(
self, fileobj: IO[bytes], size: int, *, backend_options: BackendOptions | None = None
self,
fileobj: IO[bytes],
size: int,
*,
backend_options: BackendOptions | None = None,
chunk_size: int | None = None,
) -> None:
"""Stream large objects to the FTP data connection (bounded memory)."""
if size < self._STREAM_THRESHOLD:
chunk_bytes = chunk_size or self._STREAM_CHUNK
if size < chunk_bytes:
await self.write_bytes(fileobj.read())
return
path = self._item_path

async def _stream(client: aioftp.Client) -> int:
async with client.upload_stream(path) as stream:
while True:
chunk = fileobj.read(self._STREAM_CHUNK)
chunk = fileobj.read(chunk_bytes)
if not chunk:
break
await stream.write(chunk)
Expand Down
Loading
Loading