From 0144767656b9029482b0171ab9076cbfee97b954 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 16:11:54 +0200 Subject: [PATCH 01/14] feat(python): Add automatic resumable uploads --- clients/python/README.md | 8 + .../src/objectstore_client/_resumable.py | 117 +++++++++++++++ .../python/src/objectstore_client/client.py | 142 ++++++++++++++---- .../python/src/objectstore_client/errors.py | 8 +- clients/python/tests/test_e2e.py | 42 ++++++ clients/python/tests/test_upload.py | 137 +++++++++++++++++ 6 files changed, 426 insertions(+), 28 deletions(-) create mode 100644 clients/python/tests/test_upload.py diff --git a/clients/python/README.md b/clients/python/README.md index a37e0c68..ca2cf0e8 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -23,6 +23,14 @@ content = result.payload.read() session.delete(key) ``` +`Session.put()` automatically uses resumable uploads for known remaining source +sizes of at least 32 MiB when the encoded body is seekable; pass +`resumable=False` to opt out. Seekable uncompressed or precompressed streams +upload from their current cursor without staging; streams needing compression +use direct uploads. Any resumable creation failure falls back to direct upload. +Recovery stays within one call, with up to three recovery retries in addition to +the pool's request retries. See `Session.put()` for `RequestError` error behavior. + ## Core Concepts ### Usecases and Scopes diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index 8c467ade..cde856bb 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -1,12 +1,16 @@ from __future__ import annotations +import random +import time from dataclasses import dataclass +from io import SEEK_END from typing import IO, TYPE_CHECKING from urllib.parse import urlencode import urllib3 from objectstore_client.errors import RequestError, raise_for_status +from objectstore_client.metadata import Compression, ExpirationPolicy from objectstore_client.metrics import measure_storage_operation from objectstore_client.tracing import storage_span @@ -207,3 +211,116 @@ def cancel(self) -> None: else RequestError ) raise_for_status(response, error_type=error_type) + + +RESUMABLE_THRESHOLD = 32 * 1024 * 1024 + + +def remaining_size(contents: bytes | IO[bytes]) -> int | None: + """Inspect the remaining size without reading or changing the cursor.""" + if isinstance(contents, bytes): + return len(contents) + try: + if not contents.seekable(): + return None + start = contents.tell() + except (OSError, ValueError): + return None + try: + end = contents.seek(0, SEEK_END) + except (OSError, ValueError): + return None + finally: + contents.seek(start) + return max(0, end - start) + + +def _transient(error: Exception) -> bool: + if isinstance(error, urllib3.exceptions.MaxRetryError): + # Exhausted status retries carry ResponseError rather than the response. + return isinstance(error.reason, urllib3.exceptions.ResponseError) or ( + isinstance(error.reason, Exception) and _transient(error.reason) + ) + if isinstance(error, RequestError): + return error.status in (408, 429, 500, 502, 503, 504) + return isinstance( + error, + ( + urllib3.exceptions.ConnectTimeoutError, + urllib3.exceptions.ReadTimeoutError, + urllib3.exceptions.ProtocolError, + ), + ) + + +def upload( + session: Session, + body: IO[bytes], + encoded_size: int, + key: str | None = None, + compression: Compression | None = None, + content_type: str | None = None, + metadata: dict[str, str] | None = None, + expiration_policy: ExpirationPolicy | None = None, + origin: str | None = None, + filename: str | None = None, +) -> str | None: + """Resume with three recovery retries in addition to the pool's request retries. + + Restore the starting cursor and return None on any creation failure so the + caller can use a direct upload. + Once created, use progress to recover after transient failures without + switching protocols. The recovery budget spans the whole upload, including + failed progress queries, and uses exponential backoff with jitter. + """ + start = body.tell() + try: + handle = session._create_upload( + encoded_size, + key=key, + compression=compression, + content_type=content_type, + metadata=metadata, + expiration_policy=expiration_policy, + origin=origin, + filename=filename, + ) + except Exception: + handle = None + if handle is None: + body.seek(start) + return None + + try: + offset = 0 + retries = 0 + probing = False + while True: + try: + if probing: + result = handle.progress() + else: + body.seek(start + offset) + result = handle.put(offset, (body, encoded_size - offset)) + except UploadOffsetMismatch as error: + result = UploadIncomplete(error.offset) + except Exception as error: + if not _transient(error) or retries == 3: + raise + time.sleep(random.uniform(0, min(0.1 * 2**retries, 2.0))) + retries += 1 + probing = True + continue + + if isinstance(result, UploadComplete): + return handle.key + if not offset <= result.offset <= encoded_size: + raise ValueError("Invalid upload offset") + if result.offset == offset and not probing: + raise ValueError("Upload made no progress") + offset = result.offset + probing = False + except Exception as error: + status = error.status if isinstance(error, RequestError) else None + response = error.response if isinstance(error, RequestError) else "" + raise RequestError("upload failed", status, response) from error diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index 579e9e3b..91cfc7ba 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -37,6 +37,7 @@ from objectstore_client.metrics import ( MetricsBackend, NoOpMetricsBackend, + StorageMetricEmitter, measure_storage_operation, ) from objectstore_client.multipart import MultipartUpload @@ -393,6 +394,7 @@ def put( expiration_policy: ExpirationPolicy | None = None, origin: str | None = None, filename: str | None = None, + resumable: bool = True, ) -> str: """ Uploads the given `contents` to blob storage. @@ -412,6 +414,21 @@ def put( You can use the utility function `objectstore_client.utils.guess_mime_type` to attempt to guess a `content_type` based on magic bytes. + By default, known remaining source sizes of at least 32 MiB use resumable + uploads when the encoded body is seekable: byte payloads are compressed + once in memory, while uncompressed or precompressed streams are uploaded + from their current cursor without staging. Caller-owned streams remain + open. Streams that need on-the-fly compression, smaller or unknown-size + inputs, and calls with ``resumable=False`` use direct uploads. Any session + creation failure falls back to a direct upload of the same bytes. + After creation, recovery stays within this one ``put()`` call and never + switches protocols. Requests honor the pool's retry policy; independently, + transient write or progress-query failures get up to three recovery retries + across the upload, with exponential backoff and jitter. Progress queries + confirm completion or supply the offset to resume from. Execution failures + raise ``RequestError`` with the cause chained; argument, preparation, + and direct-upload errors propagate unchanged. + `compression` is deprecated in favor of `compress`. """ if compression is not None: @@ -431,17 +448,15 @@ def put( if precompressed and precompressed != "zstd": raise ValueError(f"Invalid compression: {precompressed}") - body = BytesIO(contents) if isinstance(contents, bytes) else contents - original_body: IO[bytes] = body - encoding = precompressed or compress or self._usecase._compression compress_with = encoding if precompressed is None else "none" - if compress_with == "zstd": - cctx = zstandard.ZstdCompressor() - body = cctx.stream_reader(original_body) - body = cast(IO[bytes], utils._ZstdCompressionReaderWrapper(body)) - + # On-the-fly compression cannot report its encoded size or seek to an + # encoded offset. Keep those streams on the direct path. + replayable = isinstance(contents, bytes) or compress_with == "none" + source_size = ( + _resumable.remaining_size(contents) if resumable and replayable else None + ) headers = self._metadata_headers( compression=encoding, content_type=content_type, @@ -460,8 +475,89 @@ def put( self._metrics_backend, "put", self._usecase.name ) as metrics, ): - retries = None # by default use the pool's value, set by the Client - if compress_with != "none": + if ( + source_size is not None + and source_size >= _resumable.RESUMABLE_THRESHOLD + ): + if isinstance(contents, bytes): + encoded = ( + zstandard.ZstdCompressor().compress(contents) + if compress_with == "zstd" + else contents + ) + body: IO[bytes] = BytesIO(encoded) + encoded_size = len(encoded) + else: + body = contents + encoded_size = source_size + try: + result_key = _resumable.upload( + self, + body, + encoded_size, + key=key, + compression=encoding, + content_type=content_type, + metadata=metadata, + expiration_policy=expiration_policy, + origin=origin, + filename=filename, + ) + if result_key is None: + headers["Content-Length"] = str(encoded_size) + result_key = self._put_direct( + body, key, headers, compress=False + ) + if precompressed is None: + metrics.record_uncompressed_size(source_size) + if encoding != "none": + metrics.record_compressed_size(encoded_size, encoding) + finally: + if isinstance(contents, bytes): + body.close() + else: + result_key = self._put_direct( + contents, + key, + headers, + compress=compress_with == "zstd", + metrics=metrics, + record_source=precompressed is None, + encoding=encoding, + ) + + # Set after the response, since the key may be server-generated. + span.set_attribute("objectstore.key", result_key) + span.set_attribute("objectstore.compression", encoding) + if metrics.uncompressed_size is not None: + span.set_attribute( + "objectstore.uncompressed_size", metrics.uncompressed_size + ) + if metrics.compressed_size is not None: + span.set_attribute( + "objectstore.compressed_size", metrics.compressed_size + ) + return result_key + + def _put_direct( + self, + contents: bytes | IO[bytes], + key: str | None, + headers: dict[str, str], + compress: bool, + metrics: StorageMetricEmitter | None = None, + record_source: bool = True, + encoding: Compression = "none", + ) -> str: + """Stream a direct upload, optionally compressing on the fly.""" + body = BytesIO(contents) if isinstance(contents, bytes) else contents + original_body: IO[bytes] = body + retries = None # by default use the pool's value, set by the Client + try: + if compress: + cctx = zstandard.ZstdCompressor() + body = cctx.stream_reader(original_body, closefd=False) + body = cast(IO[bytes], utils._ZstdCompressionReaderWrapper(body)) # For on-the-fly compression, don't attempt read retries, # as the stream cannot be rewound after data has been consumed. pool_retries = self._pool.retries @@ -484,23 +580,17 @@ def put( # Must do this after streaming `body` as that's what is responsible # for advancing the seek position in both streams - if precompressed is None: - metrics.record_uncompressed_size(original_body.tell()) - if encoding != "none": - metrics.record_compressed_size(body.tell(), encoding) - - # Set after the response, since the key may be server-generated. - span.set_attribute("objectstore.key", res["key"]) - span.set_attribute("objectstore.compression", encoding) - if metrics.uncompressed_size is not None: - span.set_attribute( - "objectstore.uncompressed_size", metrics.uncompressed_size - ) - if metrics.compressed_size is not None: - span.set_attribute( - "objectstore.compressed_size", metrics.compressed_size - ) + if metrics is not None: + if record_source: + metrics.record_uncompressed_size(original_body.tell()) + if encoding != "none": + metrics.record_compressed_size(body.tell(), encoding) return res["key"] + finally: + if body is not original_body: + body.close() + if isinstance(contents, bytes): + original_body.close() def get( self, diff --git a/clients/python/src/objectstore_client/errors.py b/clients/python/src/objectstore_client/errors.py index 2f94eff6..ecc15245 100644 --- a/clients/python/src/objectstore_client/errors.py +++ b/clients/python/src/objectstore_client/errors.py @@ -4,9 +4,13 @@ class RequestError(Exception): - """Exception raised if an API call to Objectstore fails.""" + """Exception raised if an API call to Objectstore fails. - def __init__(self, message: str, status: int, response: str): + ``status`` is None and ``response`` is empty when no HTTP response is + available. Automatic upload failures chain the underlying cause. + """ + + def __init__(self, message: str, status: int | None, response: str): super().__init__(message) self.status = status self.response = response diff --git a/clients/python/tests/test_e2e.py b/clients/python/tests/test_e2e.py index 8529d57c..fa9126b7 100644 --- a/clients/python/tests/test_e2e.py +++ b/clients/python/tests/test_e2e.py @@ -12,6 +12,7 @@ from datetime import timedelta from io import BytesIO from pathlib import Path +from unittest.mock import Mock import pytest import urllib3 @@ -1153,3 +1154,44 @@ def test_put_stores_under_literal_key(server_url: str) -> None: status, body = _fetch(url) assert status == 200 assert body == payload + + +@pytest.mark.parametrize("precompressed", [False, True]) +def test_compressed_file_upload( + server_url: str, monkeypatch: pytest.MonkeyPatch, precompressed: bool +) -> None: + from objectstore_client import _resumable + + monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD", 1) + session = Client(server_url, token=TestSecretKey.get()).session( + Usecase("test-usecase", expiration_policy=TimeToLive(timedelta(days=1))), org=42 + ) + create = Mock(wraps=session._create_upload) + monkeypatch.setattr(session, "_create_upload", create) + contents = b"file contents\n" * 100 + encoded = ( + zstandard.ZstdCompressor().compress(contents) if precompressed else contents + ) + with tempfile.TemporaryFile() as source: + source.write(b"skip this prefix" + encoded) + source.seek(len(b"skip this prefix")) + key = session.put( + source, + precompressed="zstd" if precompressed else None, + content_type="text/plain", + metadata={"source": "file"}, + origin="203.0.113.42", + filename="example.txt", + ) + assert create.call_count == int(precompressed) + assert not source.closed + stored = session.head(key) + assert stored is not None + assert stored.compression == "zstd" + assert stored.content_type == "text/plain" + assert stored.filename == "example.txt" + assert stored.origin == "203.0.113.42" + assert stored.custom == {"source": "file"} + retrieved = session.get(key) + assert retrieved is not None + assert retrieved.payload.read() == contents diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py new file mode 100644 index 00000000..72ac6ad6 --- /dev/null +++ b/clients/python/tests/test_upload.py @@ -0,0 +1,137 @@ +from io import BytesIO +from typing import Any +from unittest.mock import Mock + +import pytest +import urllib3 +from objectstore_client import Client, RequestError, Session, Usecase, _resumable +from objectstore_client._resumable import UploadComplete + + +@pytest.fixture +def session(monkeypatch: pytest.MonkeyPatch) -> Session: + monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD", 4) + monkeypatch.setattr("objectstore_client._resumable.time.sleep", Mock()) + return Client("http://localhost:8888").session(Usecase("test", compression="none")) + + +@pytest.mark.parametrize( + "size,enabled,declined", + [(3, True, False), (4, True, False), (4, False, False), (4, True, True)], +) +def test_routing( + session: Session, + monkeypatch: pytest.MonkeyPatch, + size: int, + enabled: bool, + declined: bool, +) -> None: + handle = Mock(key="key") + handle.put.return_value = UploadComplete() + create = Mock(return_value=None if declined else handle) + direct = Mock(return_value="key") + monkeypatch.setattr(session, "_create_upload", create) + monkeypatch.setattr(session, "_put_direct", direct) + + assert session.put(b"x" * size, resumable=enabled) == "key" + eligible = enabled and size >= 4 + assert create.call_count == int(eligible) + assert handle.put.call_count == int(eligible and not declined) + assert direct.call_count == int(not eligible or declined) + + +def test_partial_failure_recovery( + session: Session, monkeypatch: pytest.MonkeyPatch +) -> None: + failure = urllib3.exceptions.ReadTimeoutError(session._pool, "/", "lost response") + outcomes = Mock( + side_effect=[ + urllib3.HTTPResponse(status=201, body=b'{"key":"key","session":"token"}'), + failure, + urllib3.HTTPResponse(status=204, headers={"Upload-Offset": "3"}), + urllib3.HTTPResponse(status=201), + ] + ) + sent = [] + + def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: + if body := kwargs.get("body"): + sent.append((kwargs["headers"]["Upload-Offset"], body.read())) + return outcomes() + + monkeypatch.setattr(session._pool, "_make_request", request) + source = BytesIO(b"prefixabcdefgh") + source.seek(len(b"prefix")) + assert session.put(source) == "key" + assert not source.closed + assert sent == [("0", b"abcdefgh"), ("3", b"defgh")] + + +def test_lost_final_response(session: Session, monkeypatch: pytest.MonkeyPatch) -> None: + handle = Mock(key="key") + handle.put.side_effect = urllib3.exceptions.ProtocolError("lost response") + handle.progress.return_value = UploadComplete() + monkeypatch.setattr(session, "_create_upload", Mock(return_value=handle)) + assert session.put(b"complete") == "key" + assert handle.put.call_count == 1 + handle.progress.assert_called_once_with() + + +@pytest.mark.parametrize("pool_retries", [0, 2]) +def test_retry_exhaustion( + session: Session, monkeypatch: pytest.MonkeyPatch, pool_retries: int +) -> None: + policy = urllib3.Retry(total=pool_retries, read=pool_retries) + session._pool.retries = policy + failure = urllib3.exceptions.ReadTimeoutError(session._pool, "/", "lost response") + sleep = Mock() + monkeypatch.setattr("objectstore_client._resumable.time.sleep", sleep) + + def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: + headers = kwargs["headers"] + if "Upload-Length" in headers: + return urllib3.HTTPResponse( + status=201, body=b'{"key":"key","session":"token"}' + ) + if headers["Upload-Offset"] == "*": + return urllib3.HTTPResponse(status=204, headers={"Upload-Offset": "0"}) + raise failure + + make_request = Mock(side_effect=request) + monkeypatch.setattr(session._pool, "_make_request", make_request) + with pytest.raises(RequestError, match="^upload failed$") as raised: + session.put(b"payload") + assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) + assert raised.value.__cause__.reason is failure + # Creation, three progress queries, and four writes with their own pool retries. + assert make_request.call_count == 1 + 3 + 4 * (pool_retries + 1) + assert session._pool.retries is policy + assert sleep.call_count == 3 + for call, limit in zip(sleep.call_args_list, [0.1, 0.2, 0.4], strict=True): + assert 0 <= call.args[0] <= limit + + +@pytest.mark.parametrize( + "error", + [ + RequestError("creation rejected", 403, "forbidden"), + urllib3.exceptions.ProtocolError("connection lost"), + ValueError("invalid creation response"), + ], +) +def test_creation_failure_falls_back( + session: Session, monkeypatch: pytest.MonkeyPatch, error: Exception +) -> None: + monkeypatch.setattr(session, "_create_upload", Mock(side_effect=error)) + sent = [] + + def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: + sent.append(kwargs["body"].read()) + return urllib3.HTTPResponse(status=201, body=b'{"key":"key"}') + + monkeypatch.setattr(session._pool, "request", request) + source = BytesIO(b"prefixpayload") + source.seek(len(b"prefix")) + assert session.put(source) == "key" + assert sent == [b"payload"] + assert not source.closed From 4bf163b9a8af94c8f1cd25d2cb54a81ceaadddd3 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 16:20:38 +0200 Subject: [PATCH 02/14] fix(python): Adjust resumable retry budget and backoff --- clients/python/README.md | 2 +- clients/python/src/objectstore_client/_resumable.py | 8 ++++---- clients/python/src/objectstore_client/client.py | 2 +- clients/python/tests/test_upload.py | 10 +++++----- 4 files changed, 11 insertions(+), 11 deletions(-) diff --git a/clients/python/README.md b/clients/python/README.md index ca2cf0e8..62f3768e 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -28,7 +28,7 @@ sizes of at least 32 MiB when the encoded body is seekable; pass `resumable=False` to opt out. Seekable uncompressed or precompressed streams upload from their current cursor without staging; streams needing compression use direct uploads. Any resumable creation failure falls back to direct upload. -Recovery stays within one call, with up to three recovery retries in addition to +Recovery stays within one call, with up to two recovery retries in addition to the pool's request retries. See `Session.put()` for `RequestError` error behavior. ## Core Concepts diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index cde856bb..0da9889e 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -265,13 +265,13 @@ def upload( origin: str | None = None, filename: str | None = None, ) -> str | None: - """Resume with three recovery retries in addition to the pool's request retries. + """Resume with two recovery retries in addition to the pool's request retries. Restore the starting cursor and return None on any creation failure so the caller can use a direct upload. Once created, use progress to recover after transient failures without switching protocols. The recovery budget spans the whole upload, including - failed progress queries, and uses exponential backoff with jitter. + failed progress queries, and waits 2 then 4 seconds plus up to 1 second of jitter. """ start = body.tell() try: @@ -305,9 +305,9 @@ def upload( except UploadOffsetMismatch as error: result = UploadIncomplete(error.offset) except Exception as error: - if not _transient(error) or retries == 3: + if not _transient(error) or retries == 2: raise - time.sleep(random.uniform(0, min(0.1 * 2**retries, 2.0))) + time.sleep(2 ** (retries + 1) + random.uniform(0, 1)) retries += 1 probing = True continue diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index 91cfc7ba..19b404a0 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -423,7 +423,7 @@ def put( creation failure falls back to a direct upload of the same bytes. After creation, recovery stays within this one ``put()`` call and never switches protocols. Requests honor the pool's retry policy; independently, - transient write or progress-query failures get up to three recovery retries + transient write or progress-query failures get up to two recovery retries across the upload, with exponential backoff and jitter. Progress queries confirm completion or supply the offset to resume from. Execution failures raise ``RequestError`` with the cause chained; argument, preparation, diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index 72ac6ad6..947ef2ee 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -103,12 +103,12 @@ def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: session.put(b"payload") assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) assert raised.value.__cause__.reason is failure - # Creation, three progress queries, and four writes with their own pool retries. - assert make_request.call_count == 1 + 3 + 4 * (pool_retries + 1) + # Creation, two progress queries, and three writes with their own pool retries. + assert make_request.call_count == 1 + 2 + 3 * (pool_retries + 1) assert session._pool.retries is policy - assert sleep.call_count == 3 - for call, limit in zip(sleep.call_args_list, [0.1, 0.2, 0.4], strict=True): - assert 0 <= call.args[0] <= limit + assert sleep.call_count == 2 + for call, delay in zip(sleep.call_args_list, [2, 4], strict=True): + assert delay <= call.args[0] <= delay + 1 @pytest.mark.parametrize( From aa46830fc75956762a146ac6c2b17b113b758a05 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 16:26:29 +0200 Subject: [PATCH 03/14] fix(python): Reserve body retries for resumable recovery --- clients/python/README.md | 6 ++- .../src/objectstore_client/_resumable.py | 12 ++++- .../python/src/objectstore_client/client.py | 8 ++-- clients/python/tests/test_upload.py | 45 ++++++++++++++++--- 4 files changed, 58 insertions(+), 13 deletions(-) diff --git a/clients/python/README.md b/clients/python/README.md index 62f3768e..b3f7b131 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -28,8 +28,10 @@ sizes of at least 32 MiB when the encoded body is seekable; pass `resumable=False` to opt out. Seekable uncompressed or precompressed streams upload from their current cursor without staging; streams needing compression use direct uploads. Any resumable creation failure falls back to direct upload. -Recovery stays within one call, with up to two recovery retries in addition to -the pool's request retries. See `Session.put()` for `RequestError` error behavior. +Recovery stays within one call, with up to two recovery retries. Writes retain +only the pool's connection retries; recovery queries the persisted offset before +resending. Control requests retain the pool's full retry policy. See +`Session.put()` for `RequestError` error behavior. ## Core Concepts diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index 0da9889e..c72cb6c6 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -160,6 +160,11 @@ def put( headers = session._make_headers() headers["Upload-Offset"] = str(offset) headers["Content-Length"] = str(length) + # Only retry connection establishment here. Replaying a body requires + # querying the server's offset first, which the automatic uploader owns. + retries = urllib3.Retry.from_int(session._pool.retries).new( + read=0, status=0, other=0, raise_on_status=False + ) with ( storage_span( "resumable.put", @@ -178,6 +183,8 @@ def put( f"{session._make_url(self.key)}?{query}", headers=headers, body=body, + retries=retries, + redirect=False, preload_content=True, decode_content=True, ) @@ -246,7 +253,6 @@ def _transient(error: Exception) -> bool: return isinstance( error, ( - urllib3.exceptions.ConnectTimeoutError, urllib3.exceptions.ReadTimeoutError, urllib3.exceptions.ProtocolError, ), @@ -265,13 +271,15 @@ def upload( origin: str | None = None, filename: str | None = None, ) -> str | None: - """Resume with two recovery retries in addition to the pool's request retries. + """Resume with two recovery retries; the pool handles connection retries. Restore the starting cursor and return None on any creation failure so the caller can use a direct upload. Once created, use progress to recover after transient failures without switching protocols. The recovery budget spans the whole upload, including failed progress queries, and waits 2 then 4 seconds plus up to 1 second of jitter. + Exhausted connection retries are terminal. Control requests retain the pool's + full retry policy. """ start = body.tell() try: diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index 19b404a0..4d3038ac 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -422,10 +422,12 @@ def put( inputs, and calls with ``resumable=False`` use direct uploads. Any session creation failure falls back to a direct upload of the same bytes. After creation, recovery stays within this one ``put()`` call and never - switches protocols. Requests honor the pool's retry policy; independently, - transient write or progress-query failures get up to two recovery retries + switches protocols. Writes use only the pool's connection retries; control + requests retain its full retry policy. Transient write or progress-query + failures get up to two recovery retries across the upload, with exponential backoff and jitter. Progress queries - confirm completion or supply the offset to resume from. Execution failures + confirm completion or supply the offset to resume from. Exhausted connection + retries are terminal. Execution failures raise ``RequestError`` with the cause chained; argument, preparation, and direct-upload errors propagate unchanged. diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index 947ef2ee..57dd5ef3 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -78,10 +78,19 @@ def test_lost_final_response(session: Session, monkeypatch: pytest.MonkeyPatch) @pytest.mark.parametrize("pool_retries", [0, 2]) +@pytest.mark.parametrize("status_failure", [False, True]) def test_retry_exhaustion( - session: Session, monkeypatch: pytest.MonkeyPatch, pool_retries: int + session: Session, + monkeypatch: pytest.MonkeyPatch, + pool_retries: int, + status_failure: bool, ) -> None: - policy = urllib3.Retry(total=pool_retries, read=pool_retries) + policy = urllib3.Retry( + total=pool_retries, + read=pool_retries, + status=pool_retries, + status_forcelist=[503], + ) session._pool.retries = policy failure = urllib3.exceptions.ReadTimeoutError(session._pool, "/", "lost response") sleep = Mock() @@ -95,22 +104,46 @@ def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: ) if headers["Upload-Offset"] == "*": return urllib3.HTTPResponse(status=204, headers={"Upload-Offset": "0"}) + if status_failure: + return urllib3.HTTPResponse(status=503, body=b"unavailable") raise failure make_request = Mock(side_effect=request) monkeypatch.setattr(session._pool, "_make_request", make_request) with pytest.raises(RequestError, match="^upload failed$") as raised: session.put(b"payload") - assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) - assert raised.value.__cause__.reason is failure - # Creation, two progress queries, and three writes with their own pool retries. - assert make_request.call_count == 1 + 2 + 3 * (pool_retries + 1) + if status_failure: + assert isinstance(raised.value.__cause__, RequestError) + assert raised.value.__cause__.status == 503 + else: + assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) + assert raised.value.__cause__.reason is failure + # Creation, two progress queries, and three writes regardless of pool retries. + assert make_request.call_count == 1 + 2 + 3 assert session._pool.retries is policy assert sleep.call_count == 2 for call, delay in zip(sleep.call_args_list, [2, 4], strict=True): assert delay <= call.args[0] <= delay + 1 +def test_connection_retries_are_not_multiplied( + session: Session, monkeypatch: pytest.MonkeyPatch +) -> None: + policy = urllib3.Retry(total=2, connect=2) + session._pool.retries = policy + handle = session._resume_upload("key", "token") + monkeypatch.setattr(session, "_create_upload", Mock(return_value=handle)) + failure = urllib3.exceptions.ConnectTimeoutError("connection timed out") + request = Mock(side_effect=failure) + monkeypatch.setattr(session._pool, "_make_request", request) + with pytest.raises(RequestError) as raised: + session.put(b"payload") + assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) + assert raised.value.__cause__.reason is failure + assert request.call_count == 3 + assert session._pool.retries is policy + + @pytest.mark.parametrize( "error", [ From f97bab8012ff7b9e71b083b6a8e884d3a7989499 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 16:36:02 +0200 Subject: [PATCH 04/14] fix(python): Preserve resumable upload redirect handling --- clients/python/src/objectstore_client/_resumable.py | 1 - 1 file changed, 1 deletion(-) diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index c72cb6c6..47b23c17 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -184,7 +184,6 @@ def put( headers=headers, body=body, retries=retries, - redirect=False, preload_content=True, decode_content=True, ) From 1b24ba8e8e24f6a18efe6317a0e8ea501d21f94f Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 16:41:01 +0200 Subject: [PATCH 05/14] fix(python): Represent missing error responses as None --- clients/python/src/objectstore_client/_resumable.py | 2 +- clients/python/src/objectstore_client/errors.py | 8 ++------ clients/python/tests/test_upload.py | 2 ++ 3 files changed, 5 insertions(+), 7 deletions(-) diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index 47b23c17..b3408132 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -329,5 +329,5 @@ def upload( probing = False except Exception as error: status = error.status if isinstance(error, RequestError) else None - response = error.response if isinstance(error, RequestError) else "" + response = error.response if isinstance(error, RequestError) else None raise RequestError("upload failed", status, response) from error diff --git a/clients/python/src/objectstore_client/errors.py b/clients/python/src/objectstore_client/errors.py index ecc15245..2f623cbd 100644 --- a/clients/python/src/objectstore_client/errors.py +++ b/clients/python/src/objectstore_client/errors.py @@ -4,13 +4,9 @@ class RequestError(Exception): - """Exception raised if an API call to Objectstore fails. + """Exception raised if an API call to Objectstore fails.""" - ``status`` is None and ``response`` is empty when no HTTP response is - available. Automatic upload failures chain the underlying cause. - """ - - def __init__(self, message: str, status: int | None, response: str): + def __init__(self, message: str, status: int | None, response: str | None): super().__init__(message) self.status = status self.response = response diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index 57dd5ef3..f9dfc02e 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -115,9 +115,11 @@ def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: if status_failure: assert isinstance(raised.value.__cause__, RequestError) assert raised.value.__cause__.status == 503 + assert raised.value.response == "unavailable" else: assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) assert raised.value.__cause__.reason is failure + assert raised.value.response is None # Creation, two progress queries, and three writes regardless of pool retries. assert make_request.call_count == 1 + 2 + 3 assert session._pool.retries is policy From 5dba905d7a66f787cb575c6542c4445ab2f2b02e Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 17:02:01 +0200 Subject: [PATCH 06/14] fix(python): Refine resumable retry handling --- .../src/objectstore_client/_resumable.py | 38 +++++++------------ .../python/src/objectstore_client/client.py | 14 +++---- clients/python/tests/test_e2e.py | 2 +- clients/python/tests/test_upload.py | 2 +- 4 files changed, 23 insertions(+), 33 deletions(-) diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index b3408132..6012a4f3 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -160,8 +160,9 @@ def put( headers = session._make_headers() headers["Upload-Offset"] = str(offset) headers["Content-Length"] = str(length) - # Only retry connection establishment here. Replaying a body requires - # querying the server's offset first, which the automatic uploader owns. + # We don't want the retry policy to retry sending the whole body as that + # defeats the purpose of using resumables, so we disable those retries + # in favor of our retry loop. retries = urllib3.Retry.from_int(session._pool.retries).new( read=0, status=0, other=0, raise_on_status=False ) @@ -219,11 +220,10 @@ def cancel(self) -> None: raise_for_status(response, error_type=error_type) -RESUMABLE_THRESHOLD = 32 * 1024 * 1024 +RESUMABLE_THRESHOLD_BYTES = 32 * 1024 * 1024 -def remaining_size(contents: bytes | IO[bytes]) -> int | None: - """Inspect the remaining size without reading or changing the cursor.""" +def get_size(contents: bytes | IO[bytes]) -> int | None: if isinstance(contents, bytes): return len(contents) try: @@ -241,14 +241,14 @@ def remaining_size(contents: bytes | IO[bytes]) -> int | None: return max(0, end - start) -def _transient(error: Exception) -> bool: +def is_transient(error: Exception) -> bool: if isinstance(error, urllib3.exceptions.MaxRetryError): # Exhausted status retries carry ResponseError rather than the response. return isinstance(error.reason, urllib3.exceptions.ResponseError) or ( - isinstance(error.reason, Exception) and _transient(error.reason) + isinstance(error.reason, Exception) and is_transient(error.reason) ) if isinstance(error, RequestError): - return error.status in (408, 429, 500, 502, 503, 504) + return error.status in (408, 429, 502, 503, 504) return isinstance( error, ( @@ -270,16 +270,6 @@ def upload( origin: str | None = None, filename: str | None = None, ) -> str | None: - """Resume with two recovery retries; the pool handles connection retries. - - Restore the starting cursor and return None on any creation failure so the - caller can use a direct upload. - Once created, use progress to recover after transient failures without - switching protocols. The recovery budget spans the whole upload, including - failed progress queries, and waits 2 then 4 seconds plus up to 1 second of jitter. - Exhausted connection retries are terminal. Control requests retain the pool's - full retry policy. - """ start = body.tell() try: handle = session._create_upload( @@ -301,10 +291,10 @@ def upload( try: offset = 0 retries = 0 - probing = False + probe = False while True: try: - if probing: + if probe: result = handle.progress() else: body.seek(start + offset) @@ -312,21 +302,21 @@ def upload( except UploadOffsetMismatch as error: result = UploadIncomplete(error.offset) except Exception as error: - if not _transient(error) or retries == 2: + if not is_transient(error) or retries == 2: raise time.sleep(2 ** (retries + 1) + random.uniform(0, 1)) retries += 1 - probing = True + probe = True continue if isinstance(result, UploadComplete): return handle.key if not offset <= result.offset <= encoded_size: raise ValueError("Invalid upload offset") - if result.offset == offset and not probing: + if result.offset == offset and not probe: raise ValueError("Upload made no progress") offset = result.offset - probing = False + probe = False except Exception as error: status = error.status if isinstance(error, RequestError) else None response = error.response if isinstance(error, RequestError) else None diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index 4d3038ac..f736c0a7 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -453,12 +453,12 @@ def put( encoding = precompressed or compress or self._usecase._compression compress_with = encoding if precompressed is None else "none" + # On-the-fly compression cannot report its encoded size or seek to an # encoded offset. Keep those streams on the direct path. replayable = isinstance(contents, bytes) or compress_with == "none" - source_size = ( - _resumable.remaining_size(contents) if resumable and replayable else None - ) + body_size = _resumable.get_size(contents) if resumable and replayable else None + headers = self._metadata_headers( compression=encoding, content_type=content_type, @@ -478,8 +478,8 @@ def put( ) as metrics, ): if ( - source_size is not None - and source_size >= _resumable.RESUMABLE_THRESHOLD + body_size is not None + and body_size >= _resumable.RESUMABLE_THRESHOLD_BYTES ): if isinstance(contents, bytes): encoded = ( @@ -491,7 +491,7 @@ def put( encoded_size = len(encoded) else: body = contents - encoded_size = source_size + encoded_size = body_size try: result_key = _resumable.upload( self, @@ -511,7 +511,7 @@ def put( body, key, headers, compress=False ) if precompressed is None: - metrics.record_uncompressed_size(source_size) + metrics.record_uncompressed_size(body_size) if encoding != "none": metrics.record_compressed_size(encoded_size, encoding) finally: diff --git a/clients/python/tests/test_e2e.py b/clients/python/tests/test_e2e.py index fa9126b7..13c296c1 100644 --- a/clients/python/tests/test_e2e.py +++ b/clients/python/tests/test_e2e.py @@ -1162,7 +1162,7 @@ def test_compressed_file_upload( ) -> None: from objectstore_client import _resumable - monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD", 1) + monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD_BYTES", 1) session = Client(server_url, token=TestSecretKey.get()).session( Usecase("test-usecase", expiration_policy=TimeToLive(timedelta(days=1))), org=42 ) diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index f9dfc02e..54922de8 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -10,7 +10,7 @@ @pytest.fixture def session(monkeypatch: pytest.MonkeyPatch) -> Session: - monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD", 4) + monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD_BYTES", 4) monkeypatch.setattr("objectstore_client._resumable.time.sleep", Mock()) return Client("http://localhost:8888").session(Usecase("test", compression="none")) From 03d15f144c2d0ee8212b2dd1cb8b2e76731ee8f4 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 17:50:22 +0200 Subject: [PATCH 07/14] test(python): Consolidate resumable upload coverage --- clients/python/tests/test_e2e.py | 24 +++++++++ clients/python/tests/test_upload.py | 83 +++++------------------------ 2 files changed, 36 insertions(+), 71 deletions(-) diff --git a/clients/python/tests/test_e2e.py b/clients/python/tests/test_e2e.py index 13c296c1..bb87df70 100644 --- a/clients/python/tests/test_e2e.py +++ b/clients/python/tests/test_e2e.py @@ -1195,3 +1195,27 @@ def test_compressed_file_upload( retrieved = session.get(key) assert retrieved is not None assert retrieved.payload.read() == contents + + +@pytest.mark.parametrize( + "error", [None, RequestError("creation rejected", 403, "forbidden")] +) +def test_resumable_creation_fallback( + server_url: str, monkeypatch: pytest.MonkeyPatch, error: RequestError | None +) -> None: + from objectstore_client import _resumable + + monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD_BYTES", 1) + session = Client(server_url, token=TestSecretKey.get()).session( + Usecase("test-usecase") + ) + create = Mock(return_value=None, side_effect=error) + monkeypatch.setattr(session, "_create_upload", create) + source = BytesIO(b"prefixpayload") + source.seek(len(b"prefix")) + key = session.put(source, compress="none") + create.assert_called_once() + assert not source.closed + retrieved = session.get(key) + assert retrieved is not None + assert retrieved.payload.read() == b"payload" diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index 54922de8..b746fd15 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -5,7 +5,6 @@ import pytest import urllib3 from objectstore_client import Client, RequestError, Session, Usecase, _resumable -from objectstore_client._resumable import UploadComplete @pytest.fixture @@ -15,40 +14,18 @@ def session(monkeypatch: pytest.MonkeyPatch) -> Session: return Client("http://localhost:8888").session(Usecase("test", compression="none")) -@pytest.mark.parametrize( - "size,enabled,declined", - [(3, True, False), (4, True, False), (4, False, False), (4, True, True)], -) -def test_routing( - session: Session, - monkeypatch: pytest.MonkeyPatch, - size: int, - enabled: bool, - declined: bool, -) -> None: - handle = Mock(key="key") - handle.put.return_value = UploadComplete() - create = Mock(return_value=None if declined else handle) - direct = Mock(return_value="key") - monkeypatch.setattr(session, "_create_upload", create) - monkeypatch.setattr(session, "_put_direct", direct) - - assert session.put(b"x" * size, resumable=enabled) == "key" - eligible = enabled and size >= 4 - assert create.call_count == int(eligible) - assert handle.put.call_count == int(eligible and not declined) - assert direct.call_count == int(not eligible or declined) - - -def test_partial_failure_recovery( - session: Session, monkeypatch: pytest.MonkeyPatch +@pytest.mark.parametrize("complete", [False, True]) +def test_failure_recovery( + session: Session, monkeypatch: pytest.MonkeyPatch, complete: bool ) -> None: failure = urllib3.exceptions.ReadTimeoutError(session._pool, "/", "lost response") outcomes = Mock( side_effect=[ urllib3.HTTPResponse(status=201, body=b'{"key":"key","session":"token"}'), failure, - urllib3.HTTPResponse(status=204, headers={"Upload-Offset": "3"}), + urllib3.HTTPResponse( + status=201 if complete else 204, headers={"Upload-Offset": "3"} + ), urllib3.HTTPResponse(status=201), ] ) @@ -64,31 +41,21 @@ def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: source.seek(len(b"prefix")) assert session.put(source) == "key" assert not source.closed - assert sent == [("0", b"abcdefgh"), ("3", b"defgh")] - - -def test_lost_final_response(session: Session, monkeypatch: pytest.MonkeyPatch) -> None: - handle = Mock(key="key") - handle.put.side_effect = urllib3.exceptions.ProtocolError("lost response") - handle.progress.return_value = UploadComplete() - monkeypatch.setattr(session, "_create_upload", Mock(return_value=handle)) - assert session.put(b"complete") == "key" - assert handle.put.call_count == 1 - handle.progress.assert_called_once_with() + assert sent == ( + [("0", b"abcdefgh")] if complete else [("0", b"abcdefgh"), ("3", b"defgh")] + ) -@pytest.mark.parametrize("pool_retries", [0, 2]) @pytest.mark.parametrize("status_failure", [False, True]) def test_retry_exhaustion( session: Session, monkeypatch: pytest.MonkeyPatch, - pool_retries: int, status_failure: bool, ) -> None: policy = urllib3.Retry( - total=pool_retries, - read=pool_retries, - status=pool_retries, + total=2, + read=2, + status=2, status_forcelist=[503], ) session._pool.retries = policy @@ -144,29 +111,3 @@ def test_connection_retries_are_not_multiplied( assert raised.value.__cause__.reason is failure assert request.call_count == 3 assert session._pool.retries is policy - - -@pytest.mark.parametrize( - "error", - [ - RequestError("creation rejected", 403, "forbidden"), - urllib3.exceptions.ProtocolError("connection lost"), - ValueError("invalid creation response"), - ], -) -def test_creation_failure_falls_back( - session: Session, monkeypatch: pytest.MonkeyPatch, error: Exception -) -> None: - monkeypatch.setattr(session, "_create_upload", Mock(side_effect=error)) - sent = [] - - def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: - sent.append(kwargs["body"].read()) - return urllib3.HTTPResponse(status=201, body=b'{"key":"key"}') - - monkeypatch.setattr(session._pool, "request", request) - source = BytesIO(b"prefixpayload") - source.seek(len(b"prefix")) - assert session.put(source) == "key" - assert sent == [b"payload"] - assert not source.closed From 46d4611d05087fd711243984c4d346fb1e10fea2 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:18:18 +0200 Subject: [PATCH 08/14] ref(python): Remove redundant upload fallback seek --- clients/python/src/objectstore_client/_resumable.py | 1 - 1 file changed, 1 deletion(-) diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index 6012a4f3..3e67519c 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -285,7 +285,6 @@ def upload( except Exception: handle = None if handle is None: - body.seek(start) return None try: From d971126bafb99150e3ef4b79288f93b9c838972b Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:28:08 +0200 Subject: [PATCH 09/14] feat(py-client): Configure resumable uploads per usecase Refs FS-389 --- clients/python/README.md | 14 ++-- .../python/src/objectstore_client/__init__.py | 2 + .../src/objectstore_client/_resumable.py | 8 +-- .../python/src/objectstore_client/client.py | 69 +++++++++++++++---- clients/python/tests/test_e2e.py | 15 ++-- clients/python/tests/test_upload.py | 7 +- 6 files changed, 78 insertions(+), 37 deletions(-) diff --git a/clients/python/README.md b/clients/python/README.md index b3f7b131..349b004d 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -23,15 +23,11 @@ content = result.payload.read() session.delete(key) ``` -`Session.put()` automatically uses resumable uploads for known remaining source -sizes of at least 32 MiB when the encoded body is seekable; pass -`resumable=False` to opt out. Seekable uncompressed or precompressed streams -upload from their current cursor without staging; streams needing compression -use direct uploads. Any resumable creation failure falls back to direct upload. -Recovery stays within one call, with up to two recovery retries. Writes retain -only the pool's connection retries; recovery queries the persisted offset before -resending. Control requests retain the pool's full retry policy. See -`Session.put()` for `RequestError` error behavior. +`Session.put()` automatically uses resumable uploads for eligible sources of at +least 32 MiB. Set `resumable_threshold_bytes` on `Usecase` or override it on an +individual `put()`; `None` disables resumable uploads. Configure recovery with +`Usecase(resumable_retries=ResumableRetryPolicy(...))`; zero retries also disables +resumable uploads. See `Session.put()` and `ResumableRetryPolicy` for details. ## Core Concepts diff --git a/clients/python/src/objectstore_client/__init__.py b/clients/python/src/objectstore_client/__init__.py index d5c9955f..08fb580c 100644 --- a/clients/python/src/objectstore_client/__init__.py +++ b/clients/python/src/objectstore_client/__init__.py @@ -2,6 +2,7 @@ from objectstore_client.client import ( Client, GetResponse, + ResumableRetryPolicy, Session, Usecase, ) @@ -23,6 +24,7 @@ __all__ = [ "Client", "Usecase", + "ResumableRetryPolicy", "Session", "GetResponse", "RequestError", diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index 3e67519c..b0ac0ee4 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -220,9 +220,6 @@ def cancel(self) -> None: raise_for_status(response, error_type=error_type) -RESUMABLE_THRESHOLD_BYTES = 32 * 1024 * 1024 - - def get_size(contents: bytes | IO[bytes]) -> int | None: if isinstance(contents, bytes): return len(contents) @@ -270,6 +267,7 @@ def upload( origin: str | None = None, filename: str | None = None, ) -> str | None: + policy = session._usecase._resumable_retries start = body.tell() try: handle = session._create_upload( @@ -301,9 +299,9 @@ def upload( except UploadOffsetMismatch as error: result = UploadIncomplete(error.offset) except Exception as error: - if not is_transient(error) or retries == 2: + if not is_transient(error) or retries >= policy.retries: raise - time.sleep(2 ** (retries + 1) + random.uniform(0, 1)) + time.sleep(policy.delay * 2**retries + random.uniform(0, policy.jitter)) retries += 1 probe = True continue diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index f736c0a7..c469cc34 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -56,6 +56,26 @@ class GetResponse(NamedTuple): payload: IO[bytes] +@dataclass(frozen=True) +class ResumableRetryPolicy: + """Recovery limits for one resumable upload, including progress queries. + + ``retries`` counts recovery retries across the upload; zero disables resumable + uploads. + ``delay`` is the initial backoff in seconds, doubling on each retry. ``jitter`` + is the maximum random delay added in seconds. These limits are separate from + the pool's per-request retries; retryable failures are determined internally. + """ + + retries: int = 2 + delay: float = 2.0 + jitter: float = 1.0 + + def __post_init__(self) -> None: + if self.retries < 0 or self.delay < 0 or self.jitter < 0: + raise ValueError("Resumable upload retry settings must be non-negative") + + class Usecase: """ An identifier for a workload in Objectstore, along with defaults to use for all @@ -64,6 +84,10 @@ class Usecase: Usecases need to be statically defined in Objectstore's configuration server-side. Objectstore can make decisions based on the Usecase. For example, choosing the most suitable storage backend. + + ``resumable_threshold_bytes`` defaults to 32 MiB of remaining source bytes, + before any compression. ``None`` disables resumable uploads; zero allows any + eligible size. ``resumable_retries`` configures recovery via `ResumableRetryPolicy`. """ name: str @@ -75,10 +99,17 @@ def __init__( name: str, compression: Compression = "zstd", expiration_policy: ExpirationPolicy | None = None, + *, + resumable_threshold_bytes: int | None = 32 * 1024 * 1024, + resumable_retries: ResumableRetryPolicy = ResumableRetryPolicy(), ): + if resumable_threshold_bytes is not None and resumable_threshold_bytes < 0: + raise ValueError("resumable_threshold_bytes must be non-negative") self.name = name self._compression = compression self._expiration_policy = expiration_policy + self._resumable_threshold_bytes = resumable_threshold_bytes + self._resumable_retries = resumable_retries # Connect timeout used unless overridden in connection parameters. @@ -394,7 +425,7 @@ def put( expiration_policy: ExpirationPolicy | None = None, origin: str | None = None, filename: str | None = None, - resumable: bool = True, + resumable_threshold_bytes: int | None | Literal["unset"] = "unset", ) -> str: """ Uploads the given `contents` to blob storage. @@ -414,18 +445,20 @@ def put( You can use the utility function `objectstore_client.utils.guess_mime_type` to attempt to guess a `content_type` based on magic bytes. - By default, known remaining source sizes of at least 32 MiB use resumable - uploads when the encoded body is seekable: byte payloads are compressed - once in memory, while uncompressed or precompressed streams are uploaded - from their current cursor without staging. Caller-owned streams remain - open. Streams that need on-the-fly compression, smaller or unknown-size - inputs, and calls with ``resumable=False`` use direct uploads. Any session - creation failure falls back to a direct upload of the same bytes. + ``resumable_threshold_bytes`` overrides the Usecase threshold when supplied; + ``None`` disables resumable uploads. A Usecase retry count of zero always + disables them. Eligible source sizes use resumable uploads when the encoded + body is seekable: byte payloads are compressed once in memory, while + uncompressed or precompressed streams upload from their current cursor + without staging. + Caller-owned streams remain open. Direct uploads are used for streams needing + on-the-fly compression, smaller or unknown-size inputs, or when resumable + uploads are disabled. + Any session creation failure falls back to a direct upload of the same bytes. After creation, recovery stays within this one ``put()`` call and never switches protocols. Writes use only the pool's connection retries; control requests retain its full retry policy. Transient write or progress-query - failures get up to two recovery retries - across the upload, with exponential backoff and jitter. Progress queries + failures use the Usecase's ``resumable_retries`` policy. Progress queries confirm completion or supply the offset to resume from. Exhausted connection retries are terminal. Execution failures raise ``RequestError`` with the cause chained; argument, preparation, @@ -450,6 +483,11 @@ def put( if precompressed and precompressed != "zstd": raise ValueError(f"Invalid compression: {precompressed}") + if resumable_threshold_bytes == "unset": + resumable_threshold_bytes = self._usecase._resumable_threshold_bytes + if resumable_threshold_bytes is not None and resumable_threshold_bytes < 0: + raise ValueError("resumable_threshold_bytes must be non-negative") + encoding = precompressed or compress or self._usecase._compression compress_with = encoding if precompressed is None else "none" @@ -457,7 +495,13 @@ def put( # On-the-fly compression cannot report its encoded size or seek to an # encoded offset. Keep those streams on the direct path. replayable = isinstance(contents, bytes) or compress_with == "none" - body_size = _resumable.get_size(contents) if resumable and replayable else None + body_size = ( + _resumable.get_size(contents) + if resumable_threshold_bytes is not None + and self._usecase._resumable_retries.retries > 0 + and replayable + else None + ) headers = self._metadata_headers( compression=encoding, @@ -479,7 +523,8 @@ def put( ): if ( body_size is not None - and body_size >= _resumable.RESUMABLE_THRESHOLD_BYTES + and resumable_threshold_bytes is not None + and body_size >= resumable_threshold_bytes ): if isinstance(contents, bytes): encoded = ( diff --git a/clients/python/tests/test_e2e.py b/clients/python/tests/test_e2e.py index bb87df70..de0f94b3 100644 --- a/clients/python/tests/test_e2e.py +++ b/clients/python/tests/test_e2e.py @@ -1160,11 +1160,13 @@ def test_put_stores_under_literal_key(server_url: str) -> None: def test_compressed_file_upload( server_url: str, monkeypatch: pytest.MonkeyPatch, precompressed: bool ) -> None: - from objectstore_client import _resumable - - monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD_BYTES", 1) session = Client(server_url, token=TestSecretKey.get()).session( - Usecase("test-usecase", expiration_policy=TimeToLive(timedelta(days=1))), org=42 + Usecase( + "test-usecase", + expiration_policy=TimeToLive(timedelta(days=1)), + resumable_threshold_bytes=1, + ), + org=42, ) create = Mock(wraps=session._create_upload) monkeypatch.setattr(session, "_create_upload", create) @@ -1203,11 +1205,8 @@ def test_compressed_file_upload( def test_resumable_creation_fallback( server_url: str, monkeypatch: pytest.MonkeyPatch, error: RequestError | None ) -> None: - from objectstore_client import _resumable - - monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD_BYTES", 1) session = Client(server_url, token=TestSecretKey.get()).session( - Usecase("test-usecase") + Usecase("test-usecase", resumable_threshold_bytes=1) ) create = Mock(return_value=None, side_effect=error) monkeypatch.setattr(session, "_create_upload", create) diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index b746fd15..288873db 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -4,14 +4,15 @@ import pytest import urllib3 -from objectstore_client import Client, RequestError, Session, Usecase, _resumable +from objectstore_client import Client, RequestError, Session, Usecase @pytest.fixture def session(monkeypatch: pytest.MonkeyPatch) -> Session: - monkeypatch.setattr(_resumable, "RESUMABLE_THRESHOLD_BYTES", 4) monkeypatch.setattr("objectstore_client._resumable.time.sleep", Mock()) - return Client("http://localhost:8888").session(Usecase("test", compression="none")) + return Client("http://localhost:8888").session( + Usecase("test", compression="none", resumable_threshold_bytes=4) + ) @pytest.mark.parametrize("complete", [False, True]) From 2efa0f3b4cda754a687ea9c48194c33336a58c62 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:43:44 +0200 Subject: [PATCH 10/14] fix(py-client): Allow positional usecase upload settings Refs FS-389 --- clients/python/src/objectstore_client/client.py | 1 - 1 file changed, 1 deletion(-) diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index c469cc34..9a794846 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -99,7 +99,6 @@ def __init__( name: str, compression: Compression = "zstd", expiration_policy: ExpirationPolicy | None = None, - *, resumable_threshold_bytes: int | None = 32 * 1024 * 1024, resumable_retries: ResumableRetryPolicy = ResumableRetryPolicy(), ): From 333957636e8923fe40f2f52eb52686d2c933e128 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:47:44 +0200 Subject: [PATCH 11/14] docs(py-client): Clarify resumable upload recovery policy --- clients/python/src/objectstore_client/_resumable.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index b0ac0ee4..a470dfb0 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -160,9 +160,9 @@ def put( headers = session._make_headers() headers["Upload-Offset"] = str(offset) headers["Content-Length"] = str(length) - # We don't want the retry policy to retry sending the whole body as that - # defeats the purpose of using resumables, so we disable those retries - # in favor of our retry loop. + # Disable pool retries that replay the body. The resumable upload loop + # uses the Usecase's recovery policy and probes the persisted offset + # before resending data; connection retries retain the pool's policy. retries = urllib3.Retry.from_int(session._pool.retries).new( read=0, status=0, other=0, raise_on_status=False ) From 0af32a3dcdc89de99b9681b2b3360367122f16a0 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 6 Oct 2026 16:57:11 +0200 Subject: [PATCH 12/14] ref(py-client): Simplify resumable upload routing and docs --- clients/python/README.md | 14 ++--- .../python/src/objectstore_client/client.py | 51 +++++++------------ 2 files changed, 26 insertions(+), 39 deletions(-) diff --git a/clients/python/README.md b/clients/python/README.md index 349b004d..d8462fc3 100644 --- a/clients/python/README.md +++ b/clients/python/README.md @@ -23,12 +23,6 @@ content = result.payload.read() session.delete(key) ``` -`Session.put()` automatically uses resumable uploads for eligible sources of at -least 32 MiB. Set `resumable_threshold_bytes` on `Usecase` or override it on an -individual `put()`; `None` disables resumable uploads. Configure recovery with -`Usecase(resumable_retries=ResumableRetryPolicy(...))`; zero retries also disables -resumable uploads. See `Session.put()` and `ResumableRetryPolicy` for details. - ## Core Concepts ### Usecases and Scopes @@ -106,6 +100,14 @@ session.put(video_data, compress="none") session.put(zstd_data, precompressed="zstd") ``` +### Resumable Uploads + +`Session.put()` automatically uses resumable uploads for eligible sources of at +least 32 MiB. Set `resumable_threshold_bytes` on `Usecase` or override it on an +individual `put()`; `None` disables resumable uploads. Configure recovery with +`Usecase(resumable_retries=ResumableRetryPolicy(...))`; zero retries also disables +resumable uploads. See `Session.put()` and `ResumableRetryPolicy` for details. + ### Custom Metadata Arbitrary key-value pairs can be attached to objects and retrieved on download. diff --git a/clients/python/src/objectstore_client/client.py b/clients/python/src/objectstore_client/client.py index 9a794846..810a1dc0 100644 --- a/clients/python/src/objectstore_client/client.py +++ b/clients/python/src/objectstore_client/client.py @@ -62,9 +62,13 @@ class ResumableRetryPolicy: ``retries`` counts recovery retries across the upload; zero disables resumable uploads. - ``delay`` is the initial backoff in seconds, doubling on each retry. ``jitter`` - is the maximum random delay added in seconds. These limits are separate from - the pool's per-request retries; retryable failures are determined internally. + + ``delay`` is the initial backoff in seconds, doubling on each retry. + + ``jitter`` is the maximum random delay added in seconds. + + These limits are separate from the pool's per-request retries; retryable + failures are determined internally. """ retries: int = 2 @@ -86,8 +90,8 @@ class Usecase: suitable storage backend. ``resumable_threshold_bytes`` defaults to 32 MiB of remaining source bytes, - before any compression. ``None`` disables resumable uploads; zero allows any - eligible size. ``resumable_retries`` configures recovery via `ResumableRetryPolicy`. + before any compression. ``None`` disables resumable uploads. + ``resumable_retries`` configures recovery via `ResumableRetryPolicy`. """ name: str @@ -445,23 +449,8 @@ def put( to attempt to guess a `content_type` based on magic bytes. ``resumable_threshold_bytes`` overrides the Usecase threshold when supplied; - ``None`` disables resumable uploads. A Usecase retry count of zero always - disables them. Eligible source sizes use resumable uploads when the encoded - body is seekable: byte payloads are compressed once in memory, while - uncompressed or precompressed streams upload from their current cursor - without staging. - Caller-owned streams remain open. Direct uploads are used for streams needing - on-the-fly compression, smaller or unknown-size inputs, or when resumable - uploads are disabled. - Any session creation failure falls back to a direct upload of the same bytes. - After creation, recovery stays within this one ``put()`` call and never - switches protocols. Writes use only the pool's connection retries; control - requests retain its full retry policy. Transient write or progress-query - failures use the Usecase's ``resumable_retries`` policy. Progress queries - confirm completion or supply the offset to resume from. Exhausted connection - retries are terminal. Execution failures - raise ``RequestError`` with the cause chained; argument, preparation, - and direct-upload errors propagate unchanged. + ``None`` disables resumable uploads. Eligible uploads use the Usecase's + ``resumable_retries`` policy; a retry count of zero disables resumable uploads. `compression` is deprecated in favor of `compress`. """ @@ -491,15 +480,14 @@ def put( compress_with = encoding if precompressed is None else "none" - # On-the-fly compression cannot report its encoded size or seek to an - # encoded offset. Keep those streams on the direct path. replayable = isinstance(contents, bytes) or compress_with == "none" - body_size = ( - _resumable.get_size(contents) - if resumable_threshold_bytes is not None + body_size = _resumable.get_size(contents) + use_resumable = ( + body_size is not None + and resumable_threshold_bytes is not None and self._usecase._resumable_retries.retries > 0 and replayable - else None + and body_size >= resumable_threshold_bytes ) headers = self._metadata_headers( @@ -520,11 +508,8 @@ def put( self._metrics_backend, "put", self._usecase.name ) as metrics, ): - if ( - body_size is not None - and resumable_threshold_bytes is not None - and body_size >= resumable_threshold_bytes - ): + if use_resumable: + assert body_size is not None if isinstance(contents, bytes): encoded = ( zstandard.ZstdCompressor().compress(contents) From 24420d25ab49169a8332ed49d86d6fee405142dd Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 6 Oct 2026 17:25:55 +0200 Subject: [PATCH 13/14] fix(py-client): Cancel resumable uploads after failure --- clients/python/src/objectstore_client/_resumable.py | 4 ++++ clients/python/tests/test_upload.py | 9 ++++++--- 2 files changed, 10 insertions(+), 3 deletions(-) diff --git a/clients/python/src/objectstore_client/_resumable.py b/clients/python/src/objectstore_client/_resumable.py index a470dfb0..6a5934f4 100644 --- a/clients/python/src/objectstore_client/_resumable.py +++ b/clients/python/src/objectstore_client/_resumable.py @@ -315,6 +315,10 @@ def upload( offset = result.offset probe = False except Exception as error: + try: + handle.cancel() + except Exception: + pass status = error.status if isinstance(error, RequestError) else None response = error.response if isinstance(error, RequestError) else None raise RequestError("upload failed", status, response) from error diff --git a/clients/python/tests/test_upload.py b/clients/python/tests/test_upload.py index 288873db..349c2570 100644 --- a/clients/python/tests/test_upload.py +++ b/clients/python/tests/test_upload.py @@ -66,6 +66,8 @@ def test_retry_exhaustion( def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: headers = kwargs["headers"] + if args[1] == "DELETE": + return urllib3.HTTPResponse(status=204) if "Upload-Length" in headers: return urllib3.HTTPResponse( status=201, body=b'{"key":"key","session":"token"}' @@ -88,8 +90,8 @@ def request(*args: Any, **kwargs: Any) -> urllib3.HTTPResponse: assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) assert raised.value.__cause__.reason is failure assert raised.value.response is None - # Creation, two progress queries, and three writes regardless of pool retries. - assert make_request.call_count == 1 + 2 + 3 + # Creation, two progress queries, three writes, and cancellation. + assert make_request.call_count == 1 + 2 + 3 + 1 assert session._pool.retries is policy assert sleep.call_count == 2 for call, delay in zip(sleep.call_args_list, [2, 4], strict=True): @@ -110,5 +112,6 @@ def test_connection_retries_are_not_multiplied( session.put(b"payload") assert isinstance(raised.value.__cause__, urllib3.exceptions.MaxRetryError) assert raised.value.__cause__.reason is failure - assert request.call_count == 3 + # Three connection attempts each for the upload and cancellation. + assert request.call_count == 6 assert session._pool.retries is policy From dd020e029772cc478f48723c50c3546d6d466abd Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 6 Oct 2026 17:53:45 +0200 Subject: [PATCH 14/14] test(py-client): Cover automatic resumable upload sources --- clients/python/tests/test_e2e.py | 26 +++++++++++++++----------- 1 file changed, 15 insertions(+), 11 deletions(-) diff --git a/clients/python/tests/test_e2e.py b/clients/python/tests/test_e2e.py index de0f94b3..51d0ae17 100644 --- a/clients/python/tests/test_e2e.py +++ b/clients/python/tests/test_e2e.py @@ -12,7 +12,7 @@ from datetime import timedelta from io import BytesIO from pathlib import Path -from unittest.mock import Mock +from unittest.mock import Mock, patch import pytest import urllib3 @@ -25,6 +25,7 @@ Usecase, ) from objectstore_client._resumable import ( + ResumableUpload, ResumableUploadUnavailable, UploadComplete, UploadIncomplete, @@ -1156,40 +1157,43 @@ def test_put_stores_under_literal_key(server_url: str) -> None: assert body == payload -@pytest.mark.parametrize("precompressed", [False, True]) -def test_compressed_file_upload( - server_url: str, monkeypatch: pytest.MonkeyPatch, precompressed: bool -) -> None: +@pytest.mark.parametrize("source_kind", ["bytes", "precompressed_bytes", "stream"]) +def test_automatic_resumable_upload(server_url: str, source_kind: str) -> None: session = Client(server_url, token=TestSecretKey.get()).session( Usecase( "test-usecase", + compression="none", expiration_policy=TimeToLive(timedelta(days=1)), resumable_threshold_bytes=1, ), org=42, ) - create = Mock(wraps=session._create_upload) - monkeypatch.setattr(session, "_create_upload", create) contents = b"file contents\n" * 100 + precompressed = source_kind == "precompressed_bytes" encoded = ( zstandard.ZstdCompressor().compress(contents) if precompressed else contents ) - with tempfile.TemporaryFile() as source: + with ( + tempfile.TemporaryFile() as source, + patch.object( + ResumableUpload, "put", autospec=True, side_effect=ResumableUpload.put + ) as put, + ): source.write(b"skip this prefix" + encoded) source.seek(len(b"skip this prefix")) key = session.put( - source, + source if source_kind == "stream" else encoded, precompressed="zstd" if precompressed else None, content_type="text/plain", metadata={"source": "file"}, origin="203.0.113.42", filename="example.txt", ) - assert create.call_count == int(precompressed) + put.assert_called_once() assert not source.closed stored = session.head(key) assert stored is not None - assert stored.compression == "zstd" + assert stored.compression == ("zstd" if precompressed else None) assert stored.content_type == "text/plain" assert stored.filename == "example.txt" assert stored.origin == "203.0.113.42"