From d3510b309476a87cd0c89fccc4e03dd828ccf767 Mon Sep 17 00:00:00 2001 From: Anthonios Partheniou Date: Thu, 1 Oct 2026 19:16:34 +0000 Subject: [PATCH 1/4] tests: add system tests for resumable uploads --- .../gapic-generator/tests/system/conftest.py | 145 ++++ .../system/test_resumable_upload_progress.py | 256 +++++++ .../system/test_resumable_upload_resume.py | 502 ++++++++++++ .../system/test_resumable_upload_scenarios.py | 723 ++++++++++++++++++ .../system/test_resumable_upload_stall.py | 221 ++++++ 5 files changed, 1847 insertions(+) create mode 100644 packages/gapic-generator/tests/system/test_resumable_upload_progress.py create mode 100644 packages/gapic-generator/tests/system/test_resumable_upload_resume.py create mode 100644 packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py create mode 100644 packages/gapic-generator/tests/system/test_resumable_upload_stall.py diff --git a/packages/gapic-generator/tests/system/conftest.py b/packages/gapic-generator/tests/system/conftest.py index 73169dd8a79f..ed0a98a4cef2 100644 --- a/packages/gapic-generator/tests/system/conftest.py +++ b/packages/gapic-generator/tests/system/conftest.py @@ -37,6 +37,12 @@ from google.showcase import EchoClient from google.showcase import IdentityClient from google.showcase import MessagingClient +try: + from google.showcase import ResumableUploadServiceClient + + HAS_RESUMABLE_UPLOAD_CLIENT = True +except ImportError: + HAS_RESUMABLE_UPLOAD_CLIENT = False if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true": from grpc.experimental import aio @@ -288,6 +294,33 @@ def post_expand_with_metadata(self, request, metadata): return request, metadata +if HAS_RESUMABLE_UPLOAD_CLIENT: + try: + from google.showcase_v1beta1.services.resumable_upload_service.transports import ( + ResumableUploadServiceRestInterceptor, + ) + + class ResumableUploadMetadataClientRestInterceptor( + ResumableUploadServiceRestInterceptor + ): + request_metadata: Sequence[Tuple[str, str]] = [] + response_metadata: Sequence[Tuple[str, str]] = [] + + def pre_upload_media(self, request, metadata): + self.request_metadata = metadata + return request, metadata + + def post_upload_media_with_metadata(self, request, metadata): + self.response_metadata = metadata + return request, metadata + + HAS_RESUMABLE_UPLOAD_INTERCEPTOR = True + except ImportError: + HAS_RESUMABLE_UPLOAD_INTERCEPTOR = False +else: + HAS_RESUMABLE_UPLOAD_INTERCEPTOR = False + + if HAS_ASYNC_REST_ECHO_TRANSPORT: class EchoMetadataClientRestAsyncInterceptor(AsyncEchoRestInterceptor): @@ -516,3 +549,115 @@ def intercepted_echo_rest_async(): ) return EchoAsyncClient(transport=transport), interceptor + + +@pytest.fixture +def intercepted_resumable_upload_rest(use_mtls, use_tls): + if not HAS_RESUMABLE_UPLOAD_CLIENT or not HAS_RESUMABLE_UPLOAD_INTERCEPTOR: + pytest.skip("ResumableUploadServiceClient not available.") + + transport_name = "rest" + transport_cls = ResumableUploadServiceClient.get_transport_class(transport_name) + interceptor = ResumableUploadMetadataClientRestInterceptor() + + url_scheme = "https" if (use_mtls or use_tls) else "http" + transport = transport_cls( + credentials=ga_credentials.AnonymousCredentials(), + host="localhost:7469", + url_scheme=url_scheme, + interceptor=interceptor, + ) + if use_mtls or use_tls: + transport._session.verify = CERT_PATH + transport._session.mount("https://", HostNameIgnoringAdapter()) + if use_mtls: + transport._session.cert = (CERT_PATH, KEY_PATH) + + return ResumableUploadServiceClient(transport=transport), interceptor + + +try: + from google.api_core.resumable_transfer import ( + ResumableUploadConfig, + ResumableUploadSession, + ) +except ImportError: + ResumableUploadConfig = None + ResumableUploadSession = None + + +def make_resumable_upload( + transport, + request_body, + stream, + upload_url, + size=None, + config=None, + **kwargs, +): + content_type = kwargs.pop("content_type", "application/octet-stream") + response_type = kwargs.pop("response_type", None) + retry = kwargs.pop("retry", None) + timeout = kwargs.pop("timeout", None) + + if config is None: + config = ResumableUploadConfig(**kwargs) + elif kwargs: + for k, v in kwargs.items(): + if hasattr(config, k): + setattr(config, k, v) + + session = ResumableUploadSession( + upload_url=upload_url, + config=config, + content_type=content_type, + response_type=response_type, + transport=transport, + ) + return session.upload( + stream=stream, + request_body=request_body, + content_type=content_type, + size=size, + transport=transport, + retry=retry, + timeout=timeout, + ) + + +def resume_resumable_upload( + transport, + upload_url, + stream, + size=None, + config=None, + **kwargs, +): + content_type = kwargs.pop("content_type", None) + response_type = kwargs.pop("response_type", None) + retry = kwargs.pop("retry", None) + timeout = kwargs.pop("timeout", None) + + if config is None: + config = ResumableUploadConfig(**kwargs) + elif kwargs: + for k, v in kwargs.items(): + if hasattr(config, k): + setattr(config, k, v) + + session = ResumableUploadSession( + upload_url=upload_url, + config=config, + transport=transport, + content_type=content_type, + response_type=response_type, + ) + return session.resume( + upload_url=upload_url, + stream=stream, + size=size, + transport=transport, + retry=retry, + timeout=timeout, + ) + diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_progress.py b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py new file mode 100644 index 000000000000..2a22ad51febe --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py @@ -0,0 +1,256 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import io + +from google.api_core.resumable_transfer import ( + DEFAULT_CHUNK_SIZE, + ProgressState, + ResumableUploadConfig, + ResumableUploadSession, + UploadProgress, +) +from google.showcase import UploadMediaResponse + +from conftest import make_resumable_upload + + +class _UnknownSizeUnseekableStream: + """Stream wrapper with unknown total size and non-seekable semantics.""" + + def __init__(self, data: bytes) -> None: + self._buf = io.BytesIO(data) + + def read(self, size: int = -1) -> bytes: + return self._buf.read(size) + + def seekable(self) -> bool: + return False + + +def test_make_resumable_upload_end_to_end(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + stream = io.BytesIO(b"0123456789" * 100) + + # Use make_resumable_upload from start to finish + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "full_e2e_upload.txt"}' + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + chunk_size=256, + ) + assert isinstance(response, bytes) + + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "full_e2e_upload.txt" + assert final_response.size == len(stream.getvalue()) + + +def test_resumable_upload_generator_progress_tracking(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"0123456789" * 100 + stream = io.BytesIO(payload) + + scenario_headers = [("X-Goog-Test-Scenario", "chunk_granularity")] + config = ResumableUploadConfig( + chunk_size=256, + headers=scenario_headers, + ) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_list = [] + # PEP 255 generator progress tracking + for progress in session.iter_upload(stream, request_body='{"name": "generator_upload.txt"}'): + progress_list.append(progress) + assert isinstance(progress, UploadProgress) + assert "sid=" in progress.upload_url + assert progress.chunk_size == 256 + assert progress.total_bytes == len(payload) + + # Verify yielded snapshots + assert len(progress_list) >= 3 + assert progress_list[0].state == ProgressState.STARTED + assert progress_list[-1].state == ProgressState.FINALIZED + assert progress_list[-1].bytes_uploaded == len(payload) + + # Verify response populated on session after generator exhaustion + assert session.response is not None + assert session.response.name == "generator_upload.txt" + assert session.response.size == len(payload) + + +def test_resumable_upload_unseekable_stream_recovery(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "unseekable_stream_upload.txt"}' + + class UnseekableStream(io.BytesIO): + def seekable(self): + return False + + def seek(self, offset, whence=io.SEEK_SET): + raise io.UnsupportedOperation("Stream is not seekable") + + data = b"B" * 1024 + stream = UnseekableStream(data) + + # Injects 503 error on first chunk attempt, which ResumableUploadSession recovers via in-memory buffer + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ("X-Goog-Test-Scenario-Config", '{"error_code":503,"failure_count":1,"after_offset":0}'), + ] + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + size=len(data), + chunk_size=512, + headers=scenario_headers, + ) + assert isinstance(response, bytes) + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "unseekable_stream_upload.txt" + assert final_response.size == len(data) + + +def test_multi_chunk_known_size(intercepted_resumable_upload_rest): + """Golden Path Case 1: Multi-chunk upload with known size (1.5 MB, 512 KiB chunks).""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + total_size = 1_500_000 + chunk_size = 524_288 # 512 KiB + payload = b"M" * total_size + stream = io.BytesIO(payload) + + config = ResumableUploadConfig(chunk_size=chunk_size) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "multi_chunk_known_size.mp4"}', + size=total_size, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.name == "multi_chunk_known_size.mp4" + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.FINALIZED, + ] + assert offsets == [0, 524_288, 1_048_576, 1_500_000] + assert all(p.total_bytes == total_size for p in progress_records) + + +def test_small_upload_default_chunk_size(intercepted_resumable_upload_rest): + """Golden Path Case 2: Default chunk size on small upload (~100 KB).""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + total_size = 100_000 + payload = b"S" * total_size + stream = io.BytesIO(payload) + + session = ResumableUploadSession( + upload_url=initial_url, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "small_default_chunk.bin"}', + size=total_size, + ) + ) + + assert session.chunk_size <= DEFAULT_CHUNK_SIZE + assert isinstance(session.response, UploadMediaResponse) + assert session.response.name == "small_default_chunk.bin" + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.FINALIZED, + ] + assert offsets == [0, 100_000] + assert all(p.total_bytes == total_size for p in progress_records) + + +def test_standalone_finalize_unseekable_stream(intercepted_resumable_upload_rest): + """Golden Path Case 3: Unseekable stream with unknown size and exact chunk multiple.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 # 256 KiB + total_size = 3 * chunk_size # 786_432 bytes + payload = b"U" * total_size + stream = _UnknownSizeUnseekableStream(payload) + + config = ResumableUploadConfig(chunk_size=chunk_size) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "unseekable_exact_multiple.bin"}', + size=None, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.name == "unseekable_exact_multiple.bin" + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.FINALIZED, + ] + assert offsets == [0, 262_144, 524_288, 786_432, 786_432] + assert all(p.total_bytes is None for p in progress_records) diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_resume.py b/packages/gapic-generator/tests/system/test_resumable_upload_resume.py new file mode 100644 index 000000000000..c6f294a7b2f9 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_resume.py @@ -0,0 +1,502 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import io +import pytest + +from google.api_core import exceptions +from google.api_core.resumable_transfer import ( + ProgressState, + ResumableUploadConfig, + ResumableUploadSession, +) +from google.showcase import UploadMediaResponse + +from conftest import resume_resumable_upload + + +class _UnseekableBytesStream: + """File-like stream wrapper that disables seek().""" + + def __init__(self, data: bytes): + self._buf = io.BytesIO(data) + + def read(self, size: int = -1) -> bytes: + return self._buf.read(size) + + def seekable(self) -> bool: + return False + + +def test_resumable_upload_resume_direct(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "resume_direct.txt"}' + data = b"R" * 2048 + stream = io.BytesIO(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + + # Session 1: Initiate and upload first chunk (512 bytes) + config1 = ResumableUploadConfig( + chunk_size=512, + headers=scenario_headers, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session1._initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + saved_url = session1.upload_url + assert saved_url is not None + + # Transmit only the first chunk + session1._transmit_chunk(client.transport._session, stream, len(data)) + assert session1.bytes_uploaded == 512 + assert not session1.finished + + # Session 2: Fresh session simulating resumption across process boundaries + config2 = ResumableUploadConfig( + chunk_size=512, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + # Rewind stream to simulate providing full file stream on resume + stream.seek(0) + final_response = session2.resume( + upload_url=saved_url, + stream=stream, + size=len(data), + transport=client.transport._session, + ) + + assert isinstance(final_response, UploadMediaResponse) + assert final_response.name == "resume_direct.txt" + assert final_response.size == len(data) + assert session2.bytes_uploaded == len(data) + assert session2.finished + + +def test_resumable_upload_iter_resume_generator(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "iter_resume.txt"}' + data = b"I" * 1536 + stream = io.BytesIO(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + + config1 = ResumableUploadConfig( + chunk_size=512, + headers=scenario_headers, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session1._initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + saved_url = session1.upload_url + assert saved_url is not None + + # Transmit first chunk + session1._transmit_chunk(client.transport._session, stream, len(data)) + assert session1.bytes_uploaded == 512 + + # Session 2: Resuming with PEP 255 generator iter_resume + config2 = ResumableUploadConfig( + chunk_size=512, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + stream.seek(0) + progress_snapshots = list( + session2.iter_resume( + upload_url=saved_url, + stream=stream, + size=len(data), + transport=client.transport._session, + ) + ) + + # First event should be OFFSET_RECEIVED recovering to 512 bytes + assert len(progress_snapshots) >= 2 + offset_event = progress_snapshots[0] + assert offset_event.state == ProgressState.OFFSET_RECEIVED + assert offset_event.bytes_uploaded == 512 + + # Final event should be FINALIZED at 1536 bytes + final_event = progress_snapshots[-1] + assert final_event.state == ProgressState.FINALIZED + assert final_event.bytes_uploaded == len(data) + + assert isinstance(session2.response, UploadMediaResponse) + assert session2.response.name == "iter_resume.txt" + assert session2.response.size == len(data) + + +def test_resumable_upload_resume_chunk_size_override( + intercepted_resumable_upload_rest, +): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "chunk_override.txt"}' + data = b"C" * 2048 + stream = io.BytesIO(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + + # Session 1: 256-byte chunks + config1 = ResumableUploadConfig( + chunk_size=256, + headers=scenario_headers, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session1._initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + saved_url = session1.upload_url + + # Transmit 256 bytes + session1._transmit_chunk(client.transport._session, stream, len(data)) + assert session1.bytes_uploaded == 256 + + # Session 2: Resumes overriding chunk_size to 512 (valid multiple of 256) + config2 = ResumableUploadConfig( + chunk_size=512, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + stream.seek(0) + response = session2.resume( + upload_url=saved_url, + stream=stream, + chunk_size=512, + transport=client.transport._session, + ) + + assert session2.chunk_size == 512 + assert isinstance(response, UploadMediaResponse) + assert response.name == "chunk_override.txt" + assert response.size == len(data) + + +def test_resumable_upload_resume_helper_with_raw_bytes( + intercepted_resumable_upload_rest, +): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "helper_bytes.bin"}' + data = b"B" * 1024 + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + + # Session 1: Start and upload first chunk + session1 = ResumableUploadSession( + upload_url=initial_url, + config=ResumableUploadConfig( + chunk_size=256, + headers=scenario_headers, + ), + transport=client.transport._session, + ) + session1._initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + saved_url = session1.upload_url + + stream1 = io.BytesIO(data) + session1._transmit_chunk(client.transport._session, stream1, len(data)) + assert session1.bytes_uploaded == 256 + + # Resume directly using resume_resumable_upload helper with raw bytes + config2 = ResumableUploadConfig( + chunk_size=512, + ) + final_response = resume_resumable_upload( + transport=client.transport._session, + upload_url=saved_url, + stream=data, + config=config2, + response_type=UploadMediaResponse, + ) + + assert isinstance(final_response, UploadMediaResponse) + assert final_response.name == "helper_bytes.bin" + assert final_response.size == len(data) + + +def test_resume_in_progress_upload(intercepted_resumable_upload_rest): + """2.5 Resumption Suite - Case 1: Resume an in-progress multi-chunk upload.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "resume_in_progress.mp4"}' + total_size = 1_500_000 + chunk_size = 524_288 # 512 KiB + data = b"P" * total_size + stream = io.BytesIO(data) + + session1 = ResumableUploadSession( + upload_url=initial_url, + config=ResumableUploadConfig(chunk_size=chunk_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session1._initiate( + transport=client.transport._session, + request_body=request_body, + size=total_size, + ) + saved_url = session1.upload_url + assert saved_url is not None + + # Upload only the first 512 KiB chunk + session1._transmit_chunk(client.transport._session, stream, total_size) + assert session1.bytes_uploaded == 524_288 + assert not session1.finished + + # Resume from a fresh session with a rewound stream + stream.seek(0) + session2 = ResumableUploadSession( + config=ResumableUploadConfig(chunk_size=chunk_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + snapshots = list( + session2.iter_resume( + upload_url=saved_url, + stream=stream, + size=total_size, + chunk_size=chunk_size, + transport=client.transport._session, + ) + ) + + assert [(p.state, p.bytes_uploaded) for p in snapshots] == [ + (ProgressState.OFFSET_RECEIVED, 524_288), + (ProgressState.UPLOADING, 1_048_576), + (ProgressState.FINALIZED, 1_500_000), + ] + assert isinstance(session2.response, UploadMediaResponse) + assert session2.response.name == "resume_in_progress.mp4" + assert session2.response.size == total_size + + +def test_resume_finalized_upload(intercepted_resumable_upload_rest): + """2.5 Resumption Suite - Case 2: Querying/recovering an already-finalized upload.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "already_finalized.txt"}' + total_size = 524_288 + data = b"F" * total_size + + # Complete the upload in Session 1 + session1 = ResumableUploadSession( + upload_url=initial_url, + config=ResumableUploadConfig(chunk_size=total_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + resp1 = session1.upload( + stream=io.BytesIO(data), + request_body=request_body, + size=total_size, + ) + assert isinstance(resp1, UploadMediaResponse) + assert resp1.size == total_size + saved_url = session1.upload_url + assert saved_url is not None + + # Session 2 recovers against the already-finalized session URL + session2 = ResumableUploadSession( + config=ResumableUploadConfig(chunk_size=total_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session2._state._resumable_url = saved_url + session2._state._upload_url = saved_url + session2._needs_recovery = True + progress_queue = [] + snapshots = list( + session2._transmit_all_chunks( + client.transport._session, + io.BytesIO(data), + total_size, + progress_queue=progress_queue, + ) + ) + + assert session2.finished + assert [p.state for p in snapshots] == [ + ProgressState.RECOVERING, + ProgressState.OFFSET_RECEIVED, + ] + assert isinstance(session2.response, UploadMediaResponse) + assert session2.response.name + assert session2.response.size == total_size + + +def test_resume_unseekable_stream_raises(intercepted_resumable_upload_rest): + """2.5 Resumption Suite - Case 5: Resuming with an unseekable stream raises UnseekableStreamError.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "unseekable_resume.bin"}' + total_size = 1_048_576 + chunk_size = 524_288 + data = b"U" * total_size + stream1 = io.BytesIO(data) + + # Session 1 uploads first chunk of 512 KiB + session1 = ResumableUploadSession( + upload_url=initial_url, + config=ResumableUploadConfig(chunk_size=chunk_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session1._initiate( + transport=client.transport._session, + request_body=request_body, + size=total_size, + ) + saved_url = session1.upload_url + session1._transmit_chunk(client.transport._session, stream1, total_size) + assert session1.bytes_uploaded == chunk_size + + # Attempting to resume from offset 524,288 with an unseekable stream must raise UnseekableStreamError + unseekable_stream = _UnseekableBytesStream(data) + session2 = ResumableUploadSession( + config=ResumableUploadConfig(chunk_size=chunk_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + with pytest.raises(exceptions.UnseekableStreamError) as exc_info: + session2.resume( + upload_url=saved_url, + stream=unseekable_stream, + size=total_size, + chunk_size=chunk_size, + transport=client.transport._session, + ) + + assert exc_info.value.upload_url == saved_url + assert exc_info.value.chunk_size == chunk_size + + +def test_golden_user_style_resume_seekable(intercepted_resumable_upload_rest): + """Test Case 3 / 2.5 Case 6: Mid-stream client interruption and resumption with a seekable stream.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "golden_user_resume.mp4"}' + total_size = 1_500_000 + chunk_size = 524_288 # 512 KiB + data = b"G" * total_size + stream = io.BytesIO(data) + + class _ClientPauseError(Exception): + pass + + # Phase 1: Start Session 1 and intentionally abort after the first 512 KiB chunk + session1 = ResumableUploadSession( + upload_url=initial_url, + config=ResumableUploadConfig(chunk_size=chunk_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session1_snapshots = [] + with pytest.raises(_ClientPauseError): + for progress in session1.iter_upload( + stream=stream, + request_body=request_body, + size=total_size, + ): + session1_snapshots.append(progress) + if progress.bytes_uploaded >= chunk_size: + raise _ClientPauseError("Simulated user pause after first chunk") + + # Phase 2: State verification + saved_url = session1.upload_url + saved_chunk_size = session1.chunk_size + assert saved_url is not None + assert saved_chunk_size == chunk_size + assert [(p.state, p.bytes_uploaded) for p in session1_snapshots] == [ + (ProgressState.STARTED, 0), + (ProgressState.UPLOADING, 524_288), + ] + + # Phase 3: Rewind stream to byte 0 and resume in Session 2 + stream.seek(0) + session2 = ResumableUploadSession( + config=ResumableUploadConfig(chunk_size=saved_chunk_size), + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session2_snapshots = list( + session2.iter_resume( + upload_url=saved_url, + stream=stream, + size=total_size, + chunk_size=saved_chunk_size, + transport=client.transport._session, + ) + ) + + assert [(p.state, p.bytes_uploaded) for p in session2_snapshots] == [ + (ProgressState.OFFSET_RECEIVED, 524_288), + (ProgressState.UPLOADING, 1_048_576), + (ProgressState.FINALIZED, 1_500_000), + ] + assert session2.finished + assert isinstance(session2.response, UploadMediaResponse) + assert session2.response.name == "golden_user_resume.mp4" + assert session2.response.size == total_size diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py b/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py new file mode 100644 index 000000000000..6333c0b2d5c7 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py @@ -0,0 +1,723 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import datetime +import io +import json +import time +import uuid +import pytest + +from google.api_core import exceptions as core_exceptions +from google.api_core import retry as retries +from google.api_core.resumable_transfer import ( + ProgressState, + ResumableUploadConfig, + ResumableUploadSession, +) +from google.showcase import UploadMediaResponse + +from conftest import make_resumable_upload + + +FAST_RETRY = retries.Retry(initial=0.05, maximum=0.2, multiplier=1.5, timeout=5.0) +FAST_STREAMING_RETRY = retries.StreamingRetry( + initial=0.05, maximum=0.2, multiplier=1.5, timeout=5.0 +) + + +def test_resumable_upload_scenario_non_fatal_start_error(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "retry_start_upload.txt"}' + stream = io.BytesIO(b"Hello world!") + + # Injects 503 error on start attempt, which gets automatically retried + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps({"client_uuid": str(uuid.uuid4()), "error_code": 503, "failure_count": 1}), + ), + ] + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + chunk_size=256, + headers=scenario_headers, + ) + assert isinstance(response, bytes) + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "retry_start_upload.txt" + assert final_response.size == len(stream.getvalue()) + + +def test_resumable_upload_scenario_fatal_start_error(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "fatal_start_upload.txt"}' + stream = io.BytesIO(b"Hello fatal error!") + + # Injects 403 Forbidden error on start attempt (Category 3 unretriable error) + scenario_headers = [ + ("X-Goog-Test-Scenario", "fatal_error_on_start"), + ("X-Goog-Test-Scenario-Config", '{"error_code":403}'), + ] + + with pytest.raises(core_exceptions.Forbidden): + make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + headers=scenario_headers, + ) + + +@pytest.mark.skip(reason="https://github.com/googleapis/gapic-showcase/issues/1685") +def test_resumable_upload_scenario_missing_status_header_start_retry(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "missing_status_header_upload.txt"}' + stream = io.BytesIO(b"Retrying on missing status header!") + + # Intercept first start response and strip X-Goog-Upload-Status header to test Category 1 retry + original_send = client.transport._session.send + attempt_count = [0] + + def intercepting_send(request, **kwargs): + resp = original_send(request, **kwargs) + if request.headers.get("X-Goog-Upload-Command") == "start": + attempt_count[0] += 1 + if attempt_count[0] == 1: + # Strip X-Goog-Upload-Status on first attempt + resp.headers.pop("X-Goog-Upload-Status", None) + resp.headers.pop("x-goog-upload-status", None) + return resp + + client.transport._session.send = intercepting_send + try: + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + ) + assert isinstance(response, bytes) + assert attempt_count[0] >= 2 # Verified that start was retried upon missing status header + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "missing_status_header_upload.txt" + assert final_response.size == len(stream.getvalue()) + finally: + client.transport._session.send = original_send + + +def test_resumable_upload_scenario_non_fatal_chunk_error(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "recovered_chunk_upload.txt"}' + stream = io.BytesIO(b"A" * 1024) + + # Injects 503 error on first chunk attempt, which ResumableUploadSession recovers via in-memory buffer + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ("X-Goog-Test-Scenario-Config", '{"error_code":503,"failure_count":1,"after_offset":0}'), + ] + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + chunk_size=512, + headers=scenario_headers, + ) + assert isinstance(response, bytes) + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "recovered_chunk_upload.txt" + assert final_response.size == len(stream.getvalue()) + + +def test_resumable_upload_partial_commit_recovery(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "partial_commit_upload.txt"}' + data = b"0123456789" * 50 + stream = io.BytesIO(data) + + # Injects 503 error after server commits only 100 bytes of the chunk + scenario_headers = [ + ("X-Goog-Test-Scenario", "partial_commit_on_chunk_upload"), + ("X-Goog-Test-Scenario-Config", '{"error_code":503,"failure_count":1,"after_offset":0,"partial_bytes":100}'), + ] + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + chunk_size=256, + headers=scenario_headers, + ) + assert isinstance(response, bytes) + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "partial_commit_upload.txt" + assert final_response.size == len(data) + + +def test_resumable_upload_chunk_granularity_alignment(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "granularity_upload.txt"}' + data = b"X" * 1000 + stream = io.BytesIO(data) + + # Server enforces 256 byte chunk granularity + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + chunk_size=300, # Request unaligned chunk size (300) -> state machine aligns up to 512 + headers=scenario_headers, + ) + assert isinstance(response, bytes) + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "granularity_upload.txt" + assert final_response.size == len(data) + + +def test_chunk_granularity_alignment_1mb(intercepted_resumable_upload_rest): + """Chunk Granularity Suite Case 1: 1 MB upload with unaligned chunk_size=300_000.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + total_size = 1_000_000 + payload = b"G" * total_size + stream = io.BytesIO(payload) + + config = ResumableUploadConfig( + chunk_size=300_000, + headers=[("X-Goog-Test-Scenario", "chunk_granularity")], + ) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "granularity_1mb.bin"}', + size=total_size, + timeout=5.0, + ) + ) + + assert session.chunk_size == 300_032 + assert isinstance(session.response, UploadMediaResponse) + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.FINALIZED, + ] + assert offsets == [0, 300_032, 600_064, 900_096, 1_000_000] + + +def test_cat1_error_retried(intercepted_resumable_upload_rest): + """Error Recovery Suite Case 1: 503 transient error on chunk 1 at offset 0.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + total_size = 3 * chunk_size # 786_432 + payload = b"C" * total_size + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ( + "X-Goog-Test-Scenario-Config", + '{"error_code":503,"failure_count":1,"after_offset":0}', + ), + ] + config = ResumableUploadConfig(chunk_size=chunk_size, headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "cat1_503.bin"}', + size=total_size, + retry=FAST_STREAMING_RETRY, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.size == total_size + assert progress_records[-1].state == ProgressState.FINALIZED + assert progress_records[-1].bytes_uploaded == total_size + + +def test_simple_cat2_error_recovery(intercepted_resumable_upload_rest): + """Error Recovery Suite Case 2: Category 2 error recovery at offset 0.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + total_size = 3 * chunk_size # 786_432 + payload = b"D" * total_size + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ( + "X-Goog-Test-Scenario-Config", + '{"error_code":412,"failure_count":1,"after_offset":0}', + ), + ] + config = ResumableUploadConfig(chunk_size=chunk_size, headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "cat2_offset_0.bin"}', + size=total_size, + retry=FAST_STREAMING_RETRY, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.RECOVERING, + ProgressState.OFFSET_RECEIVED, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.FINALIZED, + ] + assert offsets == [0, 0, 0, 262_144, 524_288, 786_432] + + +def test_two_consecutive_cat2_recoveries_on_chunk_2(intercepted_resumable_upload_rest): + """Error Recovery Suite Case 3: Two consecutive Category 2 recoveries on chunk 2.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + total_size = 3 * chunk_size # 786_432 + payload = b"E" * total_size + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ( + "X-Goog-Test-Scenario-Config", + '{"error_code":412,"failure_count":2,"after_offset":262144}', + ), + ] + config = ResumableUploadConfig(chunk_size=chunk_size, headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "cat2_two_on_chunk_2.bin"}', + size=total_size, + retry=FAST_STREAMING_RETRY, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.UPLOADING, + ProgressState.RECOVERING, + ProgressState.OFFSET_RECEIVED, + ProgressState.RECOVERING, + ProgressState.OFFSET_RECEIVED, + ProgressState.UPLOADING, + ProgressState.FINALIZED, + ] + assert offsets == [0, 262_144, 262_144, 262_144, 262_144, 262_144, 524_288, 786_432] + + +def test_cat2_failure_on_finalizing_chunk(intercepted_resumable_upload_rest): + """Error Recovery Suite Case 4: Category 2 failure on the finalizing chunk.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + total_size = 3 * chunk_size - 100 # 786_332 + payload = b"F" * total_size + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ( + "X-Goog-Test-Scenario-Config", + '{"error_code":412,"failure_count":1,"after_offset":524288}', + ), + ] + config = ResumableUploadConfig(chunk_size=chunk_size, headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = list( + session.iter_upload( + stream, + request_body='{"name": "cat2_finalizing_chunk.bin"}', + size=total_size, + retry=FAST_STREAMING_RETRY, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.RECOVERING, + ProgressState.OFFSET_RECEIVED, + ProgressState.FINALIZED, + ] + assert offsets == [0, 262_144, 524_288, 524_288, 524_288, 786_332] + + +def test_no_headers_failure_recovers_until_deadline_exceeded(intercepted_resumable_upload_rest): + """Error Recovery Suite Case 5: Repeated no-header failures until global deadline exceeded.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + payload = b"T" * chunk_size + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ( + "X-Goog-Test-Scenario-Config", + '{"failure_count":0,"action_after_failures":"terminate"}', + ), + ] + deadline = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(seconds=1.0) + config = ResumableUploadConfig( + chunk_size=chunk_size, + deadline=deadline, + headers=scenario_headers, + ) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = [] + with pytest.raises((core_exceptions.DeadlineExceeded, core_exceptions.RetryError)): + for progress in session.iter_upload( + stream, + request_body='{"name": "terminate_until_deadline.bin"}', + size=chunk_size, + retry=FAST_STREAMING_RETRY, + ): + progress_records.append(progress) + + recovering_count = sum(1 for p in progress_records if p.state == ProgressState.RECOVERING) + assert recovering_count >= 1 + + +def test_non_fatal_error_on_start_503(intercepted_resumable_upload_rest): + """Error on Start Suite Case 1: Single 503 on start with fast retry.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"s" * 100 + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps( + {"client_uuid": str(uuid.uuid4()), "error_code": 503, "failure_count": 1} + ), + ), + ] + config = ResumableUploadConfig(headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + start_retry=FAST_RETRY, + ) + + progress_records = list( + session.iter_upload(stream, request_body='{"name": "start_503.bin"}', size=100) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.size == 100 + phases = [p.state for p in progress_records] + assert phases.count(ProgressState.STARTED) == 1 + assert phases[-1] == ProgressState.FINALIZED + + +def test_upload_rejection_400_on_start(intercepted_resumable_upload_rest): + """Test Case 2: Protocol-level 400 Bad Request rejection during session initiation.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"r" * 100 + stream = io.BytesIO(payload) + + # Invalid JSON in X-Goog-Test-Scenario-Config or injected 400 on start + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps( + {"client_uuid": str(uuid.uuid4()), "error_code": 400, "failure_count": 1} + ), + ), + ] + config = ResumableUploadConfig(headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + with pytest.raises(core_exceptions.BadRequest) as exc_info: + session.upload(stream, request_body="{}", size=100) + + assert exc_info.value.code == 400 + + +def test_retry_exhaustion_on_start_times_out(intercepted_resumable_upload_rest): + """Error on Start Suite Case 3: Repeated 503 on start until session deadline expires.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"t" * 100 + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps( + { + "client_uuid": str(uuid.uuid4()), + "error_code": 503, + "failure_count": 10_000, + } + ), + ), + ] + deadline = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(seconds=2.0) + config = ResumableUploadConfig(headers=scenario_headers, deadline=deadline) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + start_retry=retries.Retry(initial=0.2, maximum=0.5, multiplier=1.5, timeout=2.0), + ) + + progress_records = [] + t_start = time.monotonic() + with pytest.raises((core_exceptions.GoogleAPICallError, core_exceptions.RetryError)): + for progress in session.iter_upload( + stream, request_body='{"name": "start_exhaustion.bin"}', size=100 + ): + progress_records.append(progress) + elapsed = time.monotonic() - t_start + + assert elapsed >= 1.5 + phases = [p.state for p in progress_records] + assert ProgressState.UPLOADING not in phases + + +@pytest.mark.parametrize( + "status_code,expected_exc", + [ + (403, core_exceptions.Forbidden), + (404, core_exceptions.NotFound), + ], +) +def test_fatal_error_on_start_raises_immediately( + intercepted_resumable_upload_rest, status_code, expected_exc +): + """Error on Start Suite Case 4: Fatal 403 and 404 on start fail immediately without retry.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"f" * 100 + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "fatal_error_on_start"), + ("X-Goog-Test-Scenario-Config", json.dumps({"error_code": status_code})), + ] + config = ResumableUploadConfig(headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + + progress_records = [] + t_start = time.monotonic() + with pytest.raises(expected_exc) as exc_info: + for progress in session.iter_upload( + stream, request_body='{"name": "fatal_start.bin"}', size=100 + ): + progress_records.append(progress) + elapsed = time.monotonic() - t_start + + assert exc_info.value.code == status_code + assert elapsed < 0.5 + phases = [p.state for p in progress_records] + assert ProgressState.UPLOADING not in phases + + +def test_sequential_runs_session_isolation(intercepted_resumable_upload_rest): + """Error on Start Suite Case 5: Sequential runs with distinct client_uuid values.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"i" * 100 + + for run_idx in range(2): + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps( + { + "client_uuid": str(uuid.uuid4()), + "error_code": 503, + "failure_count": 1, + } + ), + ), + ] + config = ResumableUploadConfig(headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + start_retry=FAST_RETRY, + ) + resp = session.upload( + io.BytesIO(payload), + request_body=f'{{"name": "isolated_run_{run_idx}.bin"}}', + size=100, + ) + assert isinstance(resp, UploadMediaResponse) + assert resp.name == f"isolated_run_{run_idx}.bin" + assert resp.size == 100 + + +def test_non_fatal_error_on_query_recovery(intercepted_resumable_upload_rest): + """Query retry scenario: transient 503 on query during session recovery.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + total_size = 2 * chunk_size # 524_288 + payload = b"Q" * total_size + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_query"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps({"error_code": 503, "failure_count": 1}), + ), + ] + config = ResumableUploadConfig(chunk_size=chunk_size, headers=scenario_headers) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + response_type=UploadMediaResponse, + ) + session._initiate( + transport=client.transport._session, + request_body='{"name": "query_retry.bin"}', + size=total_size, + ) + session._transmit_chunk(client.transport._session, stream, total_size) + assert session.bytes_uploaded == chunk_size + + # Trigger recovery during _transmit_all_chunks; first query returns 503, second succeeds + session._needs_recovery = True + progress_records = list( + session._transmit_all_chunks( + client.transport._session, + stream, + total_size, + retry=FAST_STREAMING_RETRY, + ) + ) + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.name == "query_retry.bin" + assert session.response.size == total_size + phases = [p.state for p in progress_records] + assert ProgressState.RECOVERING in phases + assert ProgressState.OFFSET_RECEIVED in phases + assert phases[-1] == ProgressState.FINALIZED + diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_stall.py b/packages/gapic-generator/tests/system/test_resumable_upload_stall.py new file mode 100644 index 000000000000..ee4e70ac44c5 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_stall.py @@ -0,0 +1,221 @@ +# Copyright 2026 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import datetime +import io +import pytest + +from google.api_core import exceptions +from google.api_core.resumable_transfer import ( + ResumableUploadConfig, + TransferStalledError, +) +from google.showcase import UploadMediaResponse + + +from conftest import make_resumable_upload, resume_resumable_upload + + +def test_resumable_upload_stall_control_success(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "stall_control_success.txt"}' + data = b"S" * 1024 + stream = io.BytesIO(data) + + # Stall control: 100 bytes/s minimum rate, 5s timeout -> fast upload succeeds easily + config = ResumableUploadConfig( + chunk_size=512, + stall_minimum_rate=100.0, + stall_timeout=5.0, + ) + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + config=config, + ) + assert isinstance(response, bytes) + final_response = UploadMediaResponse.from_json(response) + assert final_response.name == "stall_control_success.txt" + assert final_response.size == len(data) + + +def test_resumable_upload_stall_control_triggers_abort(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "stall_abort.txt"}' + data = b"S" * 2048 + stream = io.BytesIO(data) + + # Server injects 600ms delay per chunk; client requires high rate (10000000 bytes/s) with 0.5s stall timeout + scenario_headers = [ + ("X-Goog-Test-Scenario", "happy_path"), + ("X-Goog-Test-Scenario-Config", '{"delay_ms":600}'), + ] + + config = ResumableUploadConfig( + chunk_size=512, + stall_minimum_rate=10_000_000.0, # 10 MB/s minimum + stall_timeout=0.5, # Abort if lagging for > 500ms + headers=scenario_headers, + ) + + with pytest.raises(TransferStalledError) as exc_info: + make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + config=config, + ) + + assert "Upload stalled" in str(exc_info.value) + + +def test_resumable_upload_overall_deadline_exceeded(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "deadline_exceeded.txt"}' + data = b"D" * 2048 + stream = io.BytesIO(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "happy_path"), + ("X-Goog-Test-Scenario-Config", '{"delay_ms":600}'), + ] + + # Deadline is 300ms from now; 600ms server chunk delay forces deadline expiration + deadline = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(milliseconds=300) + config = ResumableUploadConfig( + chunk_size=512, + deadline=deadline, + headers=scenario_headers, + ) + + with pytest.raises(exceptions.DeadlineExceeded) as exc_info: + make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + config=config, + ) + + assert "deadline" in str(exc_info.value).lower() + + +def test_resumable_upload_stall_vs_deadline_conversion(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + + scenario_headers = [ + ("X-Goog-Test-Scenario", "happy_path"), + ("X-Goog-Test-Scenario-Config", '{"delay_ms":600}'), + ] + + # Case A: Stall occurs before deadline -> raises TransferStalledError + future_deadline = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(seconds=10.0) + config_stall = ResumableUploadConfig( + chunk_size=512, + stall_minimum_rate=10_000_000.0, + stall_timeout=0.4, + deadline=future_deadline, + headers=scenario_headers, + ) + with pytest.raises(TransferStalledError) as exc_stall: + make_resumable_upload( + transport=client.transport._session, + request_body='{"name": "stall_before_deadline.txt"}', + stream=io.BytesIO(b"A" * 1024), + upload_url=initial_url, + config=config_stall, + ) + assert "Upload stalled" in str(exc_stall.value) + + # Case B: Server delay causes overall deadline to expire -> raises DeadlineExceeded + near_deadline = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(milliseconds=300) + config_deadline = ResumableUploadConfig( + chunk_size=512, + stall_minimum_rate=10_000_000.0, + stall_timeout=0.4, + deadline=near_deadline, + headers=scenario_headers, + ) + with pytest.raises(exceptions.DeadlineExceeded) as exc_dead: + make_resumable_upload( + transport=client.transport._session, + request_body='{"name": "deadline_before_stall.txt"}', + stream=io.BytesIO(b"B" * 1024), + upload_url=initial_url, + config=config_deadline, + ) + assert "deadline" in str(exc_dead.value).lower() + + +def test_resumable_upload_deadline_timeout_and_resume( + intercepted_resumable_upload_rest, +): + """Timeout on chunk 2 via after_offset + delay_ms, followed by manual resumption via upload_url.""" + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + chunk_size = 262_144 + total_size = 3 * chunk_size + data = b"R" * total_size + + # Inject 600ms delay starting on chunk 2 (offset 262144) so chunk 1 commits before deadline + scenario_headers = [ + ("X-Goog-Test-Scenario", "happy_path"), + ( + "X-Goog-Test-Scenario-Config", + f'{{"delay_ms":600,"after_offset":{chunk_size}}}', + ), + ] + deadline = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta( + milliseconds=350 + ) + config1 = ResumableUploadConfig( + chunk_size=chunk_size, + deadline=deadline, + headers=scenario_headers, + ) + + with pytest.raises(exceptions.DeadlineExceeded) as exc_info: + make_resumable_upload( + transport=client.transport._session, + request_body='{"name": "timeout_then_resume.bin"}', + stream=io.BytesIO(data), + upload_url=initial_url, + config=config1, + response_type=UploadMediaResponse, + ) + + saved_url = getattr(exc_info.value, "upload_url", None) + saved_chunk_size = getattr(exc_info.value, "chunk_size", chunk_size) + assert saved_url is not None + + # Resume with a fresh session and no tight deadline + final_response = resume_resumable_upload( + transport=client.transport._session, + upload_url=saved_url, + stream=io.BytesIO(data), + config=ResumableUploadConfig(chunk_size=saved_chunk_size), + response_type=UploadMediaResponse, + ) + assert isinstance(final_response, UploadMediaResponse) + assert final_response.name == "timeout_then_resume.bin" + assert final_response.size == total_size + From b952b6ca8d59df63d63ed3c1b5a7ab44d2652538 Mon Sep 17 00:00:00 2001 From: Anthonios Partheniou Date: Thu, 1 Oct 2026 19:20:45 +0000 Subject: [PATCH 2/4] add test for bug where request body is not passed --- .../system/test_resumable_upload_progress.py | 22 ++++++++++++++++++- 1 file changed, 21 insertions(+), 1 deletion(-) diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_progress.py b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py index 2a22ad51febe..c22d3cf90225 100644 --- a/packages/gapic-generator/tests/system/test_resumable_upload_progress.py +++ b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py @@ -21,7 +21,7 @@ ResumableUploadSession, UploadProgress, ) -from google.showcase import UploadMediaResponse +from google.showcase import UploadMediaRequest, UploadMediaResponse from conftest import make_resumable_upload @@ -254,3 +254,23 @@ def test_standalone_finalize_unseekable_stream(intercepted_resumable_upload_rest ] assert offsets == [0, 262_144, 524_288, 786_432, 786_432] assert all(p.total_bytes is None for p in progress_records) + + +def test_client_upload_media_passes_request_body(intercepted_resumable_upload_rest): + """ Verify that invoking `client.upload_media(request=UploadMediaRequest(...))` + and calling `session.upload(stream=...)` forwards the serialized request body + on the initial upload session start request so the server echoes back the + expected resource name in `UploadMediaResponse.name`.""" + client, _ = intercepted_resumable_upload_rest + payload = b"0123456789" * 100 + stream = io.BytesIO(payload) + + session = client.upload_media( + request=UploadMediaRequest(name="client_upload_media.txt"), + ) + response = session.upload(stream=stream) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "client_upload_media.txt" + assert response.size == len(payload) + From 6a15b5300fc3b66877306b5c4c224e452524e225 Mon Sep 17 00:00:00 2001 From: Anthonios Partheniou Date: Thu, 1 Oct 2026 21:36:07 +0000 Subject: [PATCH 3/4] add async system tests for resumable uploads --- .../gapic-generator/tests/system/conftest.py | 45 +++++++ .../system/test_resumable_upload_progress.py | 58 +++++++++ .../system/test_resumable_upload_resume.py | 110 +++++++++++++++++- 3 files changed, 212 insertions(+), 1 deletion(-) diff --git a/packages/gapic-generator/tests/system/conftest.py b/packages/gapic-generator/tests/system/conftest.py index ed0a98a4cef2..184caf6c3119 100644 --- a/packages/gapic-generator/tests/system/conftest.py +++ b/packages/gapic-generator/tests/system/conftest.py @@ -44,6 +44,7 @@ except ImportError: HAS_RESUMABLE_UPLOAD_CLIENT = False +HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = False if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true": from grpc.experimental import aio import asyncio @@ -67,6 +68,16 @@ HAS_ASYNC_REST_IDENTITY_TRANSPORT = True except: HAS_ASYNC_REST_IDENTITY_TRANSPORT = False + try: + from google.showcase import ResumableUploadServiceAsyncClient + from google.showcase_v1beta1.services.resumable_upload_service.transports import ( + AsyncResumableUploadServiceRestTransport, + AsyncResumableUploadServiceRestInterceptor, + ) + + HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = True + except: + HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = False _GRPC_VERSION = grpc.__version__ @@ -344,6 +355,19 @@ async def post_expand_with_metadata(self, request, metadata): return request, metadata +if HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT: + + class ResumableUploadMetadataClientRestAsyncInterceptor( + AsyncResumableUploadServiceRestInterceptor + ): + request_metadata: Sequence[Tuple[str, str]] = [] + response_metadata: Sequence[Tuple[str, str]] = [] + + async def pre_upload_media(self, request, metadata): + self.request_metadata = metadata + return request, metadata + + class EchoMetadataClientGrpcInterceptor( grpc.UnaryUnaryClientInterceptor, grpc.UnaryStreamClientInterceptor, @@ -576,6 +600,27 @@ def intercepted_resumable_upload_rest(use_mtls, use_tls): return ResumableUploadServiceClient(transport=transport), interceptor +@pytest.fixture +def intercepted_resumable_upload_rest_async(): + if not HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT: + pytest.skip("Skipping test with async rest.") + + transport_name = "rest_asyncio" + transport_cls = ResumableUploadServiceAsyncClient.get_transport_class( + transport_name + ) + interceptor = ResumableUploadMetadataClientRestAsyncInterceptor() + + transport = transport_cls( + credentials=async_anonymous_credentials(), + host="localhost:7469", + url_scheme="http", + interceptor=interceptor, + ) + + return ResumableUploadServiceAsyncClient(transport=transport), interceptor + + try: from google.api_core.resumable_transfer import ( ResumableUploadConfig, diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_progress.py b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py index c22d3cf90225..0b628c5bf661 100644 --- a/packages/gapic-generator/tests/system/test_resumable_upload_progress.py +++ b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py @@ -13,6 +13,8 @@ # limitations under the License. import io +import os +import pytest from google.api_core.resumable_transfer import ( DEFAULT_CHUNK_SIZE, @@ -274,3 +276,59 @@ def test_client_upload_media_passes_request_body(intercepted_resumable_upload_re assert response.name == "client_upload_media.txt" assert response.size == len(payload) + +if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true": + + @pytest.mark.asyncio + async def test_async_client_upload_media_end_to_end( + intercepted_resumable_upload_rest_async, + ): + """Verify end-to-end async resumable upload via ResumableUploadServiceAsyncClient.""" + client, _ = intercepted_resumable_upload_rest_async + payload = b"0123456789" * 100 + stream = io.BytesIO(payload) + + session = await client.upload_media( + request=UploadMediaRequest(name="async_client_upload_media.txt"), + config=ResumableUploadConfig(chunk_size=256), + ) + response = await session.upload(stream=stream) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "async_client_upload_media.txt" + assert response.size == len(payload) + + @pytest.mark.asyncio + async def test_async_client_upload_media_progress_tracking( + intercepted_resumable_upload_rest_async, + ): + """Verify async progress iteration over ResumableUploadServiceAsyncClient.upload_media.""" + client, _ = intercepted_resumable_upload_rest_async + chunk_size = 262_144 # 256 KiB + total_size = 600_000 + payload = b"A" * total_size + stream = io.BytesIO(payload) + + session = await client.upload_media( + request=UploadMediaRequest(name="async_multi_chunk.bin"), + config=ResumableUploadConfig(chunk_size=chunk_size), + ) + + progress_records = [ + p async for p in session.upload(stream=stream, size=total_size) + ] + + assert isinstance(session.response, UploadMediaResponse) + assert session.response.name == "async_multi_chunk.bin" + assert session.response.size == total_size + + phases = [p.state for p in progress_records] + offsets = [p.bytes_uploaded for p in progress_records] + assert phases == [ + ProgressState.STARTED, + ProgressState.UPLOADING, + ProgressState.UPLOADING, + ProgressState.FINALIZED, + ] + assert offsets == [0, 262_144, 524_288, 600_000] + assert all(p.total_bytes == total_size for p in progress_records) diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_resume.py b/packages/gapic-generator/tests/system/test_resumable_upload_resume.py index c6f294a7b2f9..cee2bb17532d 100644 --- a/packages/gapic-generator/tests/system/test_resumable_upload_resume.py +++ b/packages/gapic-generator/tests/system/test_resumable_upload_resume.py @@ -13,6 +13,7 @@ # limitations under the License. import io +import os import pytest from google.api_core import exceptions @@ -21,7 +22,7 @@ ResumableUploadConfig, ResumableUploadSession, ) -from google.showcase import UploadMediaResponse +from google.showcase import UploadMediaRequest, UploadMediaResponse from conftest import resume_resumable_upload @@ -500,3 +501,110 @@ class _ClientPauseError(Exception): assert isinstance(session2.response, UploadMediaResponse) assert session2.response.name == "golden_user_resume.mp4" assert session2.response.size == total_size + + +if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true": + + @pytest.mark.asyncio + async def test_async_client_upload_media_resume_direct( + intercepted_resumable_upload_rest_async, + ): + """Verify direct async resumption via ResumableUploadServiceAsyncClient.""" + client, _ = intercepted_resumable_upload_rest_async + total_size = 1_500_000 + chunk_size = 524_288 # 512 KiB + data = b"R" * total_size + stream = io.BytesIO(data) + + class _ClientPauseError(Exception): + pass + + session1 = await client.upload_media( + request=UploadMediaRequest(name="async_resume_direct.txt"), + config=ResumableUploadConfig(chunk_size=chunk_size), + ) + with pytest.raises(_ClientPauseError): + async for progress in session1.upload(stream=stream, size=total_size): + if progress.bytes_uploaded >= chunk_size: + raise _ClientPauseError("Simulated user pause after first chunk") + + saved_url = session1.upload_url + saved_chunk_size = session1.chunk_size + assert saved_url is not None + assert saved_chunk_size == chunk_size + + stream.seek(0) + session2 = await client.upload_media( + request=UploadMediaRequest(name="async_resume_direct.txt"), + config=ResumableUploadConfig(chunk_size=saved_chunk_size), + ) + response = await session2.resume( + upload_url=saved_url, + stream=stream, + size=total_size, + chunk_size=saved_chunk_size, + ) + + assert session2.finished + assert isinstance(response, UploadMediaResponse) + assert response.name == "async_resume_direct.txt" + assert response.size == total_size + + @pytest.mark.asyncio + async def test_async_client_upload_media_resume_seekable( + intercepted_resumable_upload_rest_async, + ): + """Verify async progress iteration during resumption via ResumableUploadServiceAsyncClient.""" + client, _ = intercepted_resumable_upload_rest_async + total_size = 1_500_000 + chunk_size = 524_288 # 512 KiB + data = b"G" * total_size + stream = io.BytesIO(data) + + class _ClientPauseError(Exception): + pass + + session1 = await client.upload_media( + request=UploadMediaRequest(name="async_golden_user_resume.mp4"), + config=ResumableUploadConfig(chunk_size=chunk_size), + ) + session1_snapshots = [] + with pytest.raises(_ClientPauseError): + async for progress in session1.upload(stream=stream, size=total_size): + session1_snapshots.append(progress) + if progress.bytes_uploaded >= chunk_size: + raise _ClientPauseError("Simulated user pause after first chunk") + + saved_url = session1.upload_url + saved_chunk_size = session1.chunk_size + assert saved_url is not None + assert saved_chunk_size == chunk_size + assert [(p.state, p.bytes_uploaded) for p in session1_snapshots] == [ + (ProgressState.STARTED, 0), + (ProgressState.UPLOADING, 524_288), + ] + + stream.seek(0) + session2 = await client.upload_media( + request=UploadMediaRequest(name="async_golden_user_resume.mp4"), + config=ResumableUploadConfig(chunk_size=saved_chunk_size), + ) + session2_snapshots = [ + p + async for p in session2.resume( + upload_url=saved_url, + stream=stream, + size=total_size, + chunk_size=saved_chunk_size, + ) + ] + + assert [(p.state, p.bytes_uploaded) for p in session2_snapshots] == [ + (ProgressState.OFFSET_RECEIVED, 524_288), + (ProgressState.UPLOADING, 1_048_576), + (ProgressState.FINALIZED, 1_500_000), + ] + assert session2.finished + assert isinstance(session2.response, UploadMediaResponse) + assert session2.response.name == "async_golden_user_resume.mp4" + assert session2.response.size == total_size From c0f9005020ca52a96dbb32059e98fba15545f100 Mon Sep 17 00:00:00 2001 From: Anthonios Partheniou Date: Sat, 3 Oct 2026 15:23:22 +0000 Subject: [PATCH 4/4] test: add system tests for bug where start_retry is not passed from client.upload_media --- .../system/test_resumable_upload_scenarios.py | 92 ++++++++++++++++++- 1 file changed, 91 insertions(+), 1 deletion(-) diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py b/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py index 6333c0b2d5c7..946ae3fa0c73 100644 --- a/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py +++ b/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py @@ -15,6 +15,7 @@ import datetime import io import json +import os import time import uuid import pytest @@ -26,7 +27,7 @@ ResumableUploadConfig, ResumableUploadSession, ) -from google.showcase import UploadMediaResponse +from google.showcase import UploadMediaRequest, UploadMediaResponse from conftest import make_resumable_upload @@ -721,3 +722,92 @@ def test_non_fatal_error_on_query_recovery(intercepted_resumable_upload_rest): assert ProgressState.OFFSET_RECEIVED in phases assert phases[-1] == ProgressState.FINALIZED + +def test_client_upload_media_passes_start_retry(intercepted_resumable_upload_rest): + """Verify that `retry` passed to `client.upload_media` is forwarded as `start_retry`.""" + client, _ = intercepted_resumable_upload_rest + payload = b"s" * 100 + stream = io.BytesIO(payload) + + # Injects a single 409 Conflict on start, which is not retried by the default + # start retry policy unless the custom `retry` passed to `client.upload_media` + # reaches `ResumableUploadSession(start_retry=...)`. + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps( + {"client_uuid": str(uuid.uuid4()), "error_code": 409, "failure_count": 1} + ), + ), + ] + retried_errors = [] + custom_retry = retries.Retry( + predicate=retries.if_exception_type(core_exceptions.Conflict), + initial=0.05, + maximum=0.2, + multiplier=1.5, + timeout=5.0, + on_error=retried_errors.append, + ) + + session = client.upload_media( + request=UploadMediaRequest(name="client_start_retry.bin"), + config=ResumableUploadConfig(headers=scenario_headers), + retry=custom_retry, + ) + response = session.upload(stream=stream, size=len(payload)) + + assert len(retried_errors) == 1 + assert isinstance(retried_errors[0], core_exceptions.Conflict) + assert isinstance(response, UploadMediaResponse) + assert response.name == "client_start_retry.bin" + assert response.size == len(payload) + + +if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true": + + @pytest.mark.asyncio + async def test_async_client_upload_media_passes_start_retry( + intercepted_resumable_upload_rest_async, + ): + """Verify that `retry` passed to `async_client.upload_media` is forwarded as `start_retry`.""" + client, _ = intercepted_resumable_upload_rest_async + payload = b"s" * 100 + stream = io.BytesIO(payload) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_start"), + ( + "X-Goog-Test-Scenario-Config", + json.dumps( + { + "client_uuid": str(uuid.uuid4()), + "error_code": 409, + "failure_count": 1, + } + ), + ), + ] + retried_errors = [] + custom_retry = retries.AsyncRetry( + predicate=retries.if_exception_type(core_exceptions.Conflict), + initial=0.05, + maximum=0.2, + multiplier=1.5, + timeout=5.0, + on_error=retried_errors.append, + ) + + session = await client.upload_media( + request=UploadMediaRequest(name="async_client_start_retry.bin"), + config=ResumableUploadConfig(headers=scenario_headers), + retry=custom_retry, + ) + response = await session.upload(stream=stream, size=len(payload)) + + assert len(retried_errors) == 1 + assert isinstance(retried_errors[0], core_exceptions.Conflict) + assert isinstance(response, UploadMediaResponse) + assert response.name == "async_client_start_retry.bin" + assert response.size == len(payload)