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
8 changes: 8 additions & 0 deletions clients/python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -100,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.
Expand Down
2 changes: 2 additions & 0 deletions clients/python/src/objectstore_client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
from objectstore_client.client import (
Client,
GetResponse,
ResumableRetryPolicy,
Session,
Usecase,
)
Expand All @@ -23,6 +24,7 @@
__all__ = [
"Client",
"Usecase",
"ResumableRetryPolicy",
"Session",
"GetResponse",
"RequestError",
Expand Down
115 changes: 115 additions & 0 deletions clients/python/src/objectstore_client/_resumable.py
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -156,6 +160,12 @@ def put(
headers = session._make_headers()
headers["Upload-Offset"] = str(offset)
headers["Content-Length"] = str(length)
# 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
)
with (
storage_span(
"resumable.put",
Expand All @@ -174,6 +184,7 @@ def put(
f"{session._make_url(self.key)}?{query}",
headers=headers,
body=body,
retries=retries,
preload_content=True,
decode_content=True,
)
Expand Down Expand Up @@ -207,3 +218,107 @@ def cancel(self) -> None:
else RequestError
)
raise_for_status(response, error_type=error_type)


def get_size(contents: bytes | IO[bytes]) -> int | None:
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)
Comment thread
lcian marked this conversation as resolved.
return max(0, end - start)


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 is_transient(error.reason)
)
if isinstance(error, RequestError):
return error.status in (408, 429, 502, 503, 504)
return isinstance(
error,
(
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:
policy = session._usecase._resumable_retries
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:
return None

try:
offset = 0
retries = 0
probe = False
while True:
try:
if probe:
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 is_transient(error) or retries >= policy.retries:
raise
time.sleep(policy.delay * 2**retries + random.uniform(0, policy.jitter))
retries += 1
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 probe:
raise ValueError("Upload made no progress")
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
Comment thread
lcian marked this conversation as resolved.
Loading
Loading