Repository navigation
feat(storage): stream resumable file uploads #59
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
75b6a91
bd90802
5ef28e9
4f8f7d2
2eb8e5f
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -5,11 +5,13 @@ | |
| import base64 | ||
| import binascii | ||
| import json | ||
| from collections.abc import Mapping, Sequence | ||
| from contextlib import suppress | ||
| from collections.abc import Generator, Mapping, Sequence | ||
| from contextlib import contextmanager, suppress | ||
| from dataclasses import dataclass | ||
| from datetime import datetime | ||
| from typing import Any, Protocol, cast | ||
| from io import SEEK_END, BytesIO | ||
| from tempfile import TemporaryFile | ||
| from typing import Any, BinaryIO, Protocol, cast | ||
| from urllib.parse import quote | ||
|
|
||
| from ._transport import ( | ||
|
|
@@ -21,7 +23,6 @@ | |
| invoke, | ||
| response_payload, | ||
| ) | ||
| from .errors import VolcanoError | ||
| from .models import ( | ||
| JSONValue, | ||
| StorageObject, | ||
|
|
@@ -40,6 +41,8 @@ | |
| _INVALID_PUBLIC_URL_PATH = "Public URL paths cannot contain dot segments" | ||
| _JWT_PART_COUNT = 3 | ||
| _HTTP_PARTIAL_CONTENT = 206 | ||
| _UPLOAD_SPOOL_READ_SIZE = 1_048_576 | ||
| _UPLOAD_SOURCE_UNAVAILABLE = "Upload source is temporarily unavailable" | ||
|
|
||
|
|
||
| def _optional_datetime(value: object) -> datetime | None: | ||
|
|
@@ -198,6 +201,65 @@ def _encoded_storage_path(path: str) -> str: | |
| return "/".join(quote(segment, safe="") for segment in segments) | ||
|
|
||
|
|
||
| def _remaining_upload_bytes(source: BinaryIO) -> int | None: | ||
| try: | ||
| if not source.seekable(): | ||
| return None | ||
| position = source.tell() | ||
| except (AttributeError, OSError, ValueError): | ||
| return None | ||
| try: | ||
| try: | ||
| source.seek(0, SEEK_END) | ||
| remaining = max(0, source.tell() - position) | ||
| except (OSError, ValueError): | ||
| remaining = None | ||
| finally: | ||
| source.seek(position) | ||
| return remaining | ||
|
|
||
|
|
||
| def _spool_upload_source(source: BinaryIO, target: BinaryIO) -> None: | ||
| while True: | ||
| chunk = cast("bytes | None", source.read(_UPLOAD_SPOOL_READ_SIZE)) | ||
| if chunk is None: | ||
| raise BlockingIOError(_UPLOAD_SOURCE_UNAVAILABLE) | ||
| if chunk == b"": | ||
| return | ||
| target.write(chunk) | ||
|
|
||
|
|
||
| def _read_upload_part(source: BinaryIO, part_size: int) -> bytes: | ||
| part = bytearray() | ||
| while len(part) < part_size: | ||
| chunk = cast("bytes | None", source.read(part_size - len(part))) | ||
| if chunk is None: | ||
| raise BlockingIOError(_UPLOAD_SOURCE_UNAVAILABLE) | ||
| if chunk == b"": | ||
| break | ||
| part.extend(chunk) | ||
| return bytes(part) | ||
|
|
||
|
|
||
| @contextmanager | ||
| def _resumable_upload_source( | ||
| data: bytes | BinaryIO, | ||
| ) -> Generator[tuple[BinaryIO, int], None, None]: | ||
| if isinstance(data, bytes): | ||
| with BytesIO(data) as source: | ||
| yield source, len(data) | ||
| return | ||
| remaining = _remaining_upload_bytes(data) | ||
| if remaining is not None: | ||
| yield data, remaining | ||
| return | ||
| with TemporaryFile(mode="w+b") as source: | ||
| _spool_upload_source(data, source) | ||
| total_size = source.tell() | ||
| source.seek(0) | ||
| yield source, total_size | ||
|
|
||
|
|
||
| class StorageContext(Protocol): | ||
| """Client capabilities required by object storage.""" | ||
|
|
||
|
|
@@ -498,42 +560,46 @@ def abort_upload_session( | |
| def upload_resumable( | ||
| self, | ||
| path: str, | ||
| data: bytes, | ||
| data: bytes | BinaryIO, | ||
| *, | ||
| content_type: str = "application/octet-stream", | ||
| part_size: int | None = None, | ||
| ) -> StorageObject: | ||
| """Upload bytes through a server-managed resumable session.""" | ||
| session = self.create_upload_session( | ||
| path, | ||
| total_size=len(data), | ||
| content_type=content_type, | ||
| part_size=part_size, | ||
| ) | ||
| try: | ||
| self._upload_session_parts(path, data, session) | ||
| except VolcanoError: | ||
| self._abort_failed_upload(path, session.session_id) | ||
| raise | ||
| return self.complete_upload_session(path, session_id=session.session_id) | ||
| """Upload bytes or a binary stream through a resumable session.""" | ||
| path = _storage_path(path) | ||
| self._client._session_token() | ||
| with _resumable_upload_source(data) as (source, total_size): | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
For a non-seekable input, entering this context drains the entire source into a temporary file before Useful? React with 👍 / 👎.
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Addressed the valid local portion in bd90802. Path and active-session authentication are validated before a non-seekable source is spooled. The server-side active-upload check remains in createUploadSession because total_size is required by that wire operation and is unknowable until spooling completes. |
||
| session = self.create_upload_session( | ||
| path, | ||
| total_size=total_size, | ||
| content_type=content_type, | ||
| part_size=part_size, | ||
| ) | ||
| upload_succeeded = False | ||
| try: | ||
| self._upload_session_parts(path, source, session) | ||
| upload_succeeded = True | ||
| finally: | ||
| if not upload_succeeded: | ||
| self._abort_failed_upload(path, session.session_id) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
When a stream read fails and the best-effort abort raises something other than Useful? React with 👍 / 👎.
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in 11bf76a. Best-effort abort now suppresses ordinary cleanup exceptions, and regression coverage proves the original reader failure remains observable when abort itself raises.
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Correction: the fix is in 4f8f7d2 (the prior reply contained a mistyped short SHA). |
||
| return self.complete_upload_session(path, session_id=session.session_id) | ||
|
|
||
| def _upload_session_parts( | ||
| self, | ||
| path: str, | ||
| data: bytes, | ||
| source: BinaryIO, | ||
| session: UploadSession, | ||
| ) -> None: | ||
| for part_index in range(session.total_parts): | ||
| offset = part_index * session.part_size | ||
| self.upload_part( | ||
| path, | ||
| session_id=session.session_id, | ||
| part_number=part_index + 1, | ||
| data=data[offset : offset + session.part_size], | ||
| data=_read_upload_part(source, session.part_size), | ||
| ) | ||
|
|
||
| def _abort_failed_upload(self, path: str, session_id: str) -> None: | ||
| with suppress(VolcanoError): | ||
| with suppress(Exception): | ||
| self.abort_upload_session(path, session_id=session_id) | ||
|
|
||
| def list( | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
seekable()as non-seekableFor a non-seekable binary file-like object that exposes bounded
read()but does not inheritIOBaseand therefore has noseekable()method, this unconditional probe raisesAttributeErrorinstead of using the advertised spooling path. This affects common read-only stream wrappers such as HTTP response bodies; treat an absent seekability probe as non-seekable so these sources can be spooled.Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Fixed in 6d2f4b8. Missing seekability probes are treated as non-seekable and use the bounded spooling path, with regression coverage for a duck-typed read-only stream.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Correction: the fix is in 5ef28e9 (the prior reply contained a mistyped short SHA).