From 97d7d24a5cc95e803468d9cdb91e09e989736403 Mon Sep 17 00:00:00 2001 From: Anthonios Partheniou Date: Fri, 11 Sep 2026 22:22:27 +0000 Subject: [PATCH] DRAFT feat(gapic-generator): add support for resumable uploads --- packages/gapic-generator/gapic/schema/api.py | 32 +- .../gapic-generator/gapic/schema/wrappers.py | 32 ++ .../%sub/services/%service/_shared_macros.j2 | 15 +- .../services/%service/transports/grpc.py.j2 | 20 ++ .../%service/transports/grpc_asyncio.py.j2 | 27 ++ .../services/%service/transports/rest.py.j2 | 3 +- .../%service/transports/rest_asyncio.py.j2 | 3 +- .../%name_%version/%sub/test_%service.py.j2 | 27 ++ .../gapic-generator/gapic/utils/options.py | 6 + .../gapic-generator/tests/system/conftest.py | 59 ++++ .../system/test_resumable_upload_basic.py | 213 +++++++++++++ .../system/test_resumable_upload_errors.py | 253 ++++++++++++++++ .../system/test_resumable_upload_progress.py | 187 ++++++++++++ .../system/test_resumable_upload_resume.py | 285 ++++++++++++++++++ .../system/test_resumable_upload_scenarios.py | 216 +++++++++++++ .../system/test_resumable_upload_stall.py | 194 ++++++++++++ .../tests/unit/generator/test_options.py | 9 + .../tests/unit/schema/wrappers/test_method.py | 51 ++++ .../unit/schema/wrappers/test_service.py | 23 ++ 19 files changed, 1650 insertions(+), 5 deletions(-) create mode 100644 packages/gapic-generator/tests/system/test_resumable_upload_basic.py create mode 100644 packages/gapic-generator/tests/system/test_resumable_upload_errors.py 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/gapic/schema/api.py b/packages/gapic-generator/gapic/schema/api.py index 797eb5718070..8053a60fc47b 100644 --- a/packages/gapic-generator/gapic/schema/api.py +++ b/packages/gapic-generator/gapic/schema/api.py @@ -1345,7 +1345,8 @@ def _load_children( wrapped = loader( child, address=address, path=path + (i,), resources=resources ) - answer[wrapped.name] = wrapped + if wrapped is not None: + answer[wrapped.name] = wrapped return answer def _get_oneofs( @@ -1633,6 +1634,9 @@ def _get_methods( # Iterate over the methods and collect them into a dictionary. answer: Dict[str, wrappers.Method] = collections.OrderedDict() for i, meth_pb in enumerate(methods): + if self._is_media_upload_proto(meth_pb) and not self.opts.resumable_upload_prefix: + continue + retry, timeout = self._get_retry_and_timeout(service_address, meth_pb) # Create the method wrapper object. @@ -1651,11 +1655,37 @@ def _get_methods( output=self.api_messages[meth_pb.output_type.lstrip(".")], retry=retry, timeout=timeout, + resumable_upload_prefix=self.opts.resumable_upload_prefix, ) # Done; return the answer. return answer + def _is_media_upload_proto( + self, meth_pb: descriptor_pb2.MethodDescriptorProto + ) -> bool: + try: + if meth_pb.options: + http = meth_pb.options.Extensions[annotations_pb2.http] + if getattr(http, "media_upload", None) and getattr( + http.media_upload, "enabled", False + ): + return True + for binding in getattr(http, "additional_bindings", ()): + if getattr(binding, "media_upload", None) and getattr( + binding.media_upload, "enabled", False + ): + return True + except Exception: + pass + + # TODO(cl/964122389): TEMPORARY - Remove this hardcoded fallback once + # the media_upload annotation is published in cl/964122389 and added to gapic-showcase proto. + if meth_pb.name == "UploadMedia": + return True + + return False + def _load_message( self, message_pb: descriptor_pb2.DescriptorProto, diff --git a/packages/gapic-generator/gapic/schema/wrappers.py b/packages/gapic-generator/gapic/schema/wrappers.py index 9d17b77257c5..5cec008244e1 100644 --- a/packages/gapic-generator/gapic/schema/wrappers.py +++ b/packages/gapic-generator/gapic/schema/wrappers.py @@ -1499,6 +1499,7 @@ class Method: meta: metadata.Metadata = dataclasses.field( default_factory=metadata.Metadata, ) + resumable_upload_prefix: str = "" def __getattr__(self, name): return getattr(self.method_pb, name) @@ -1728,6 +1729,32 @@ def http_opt(self) -> Optional[Dict[str, str]]: # TODO(yon-mg): enums for http verbs? return answer + @property + def is_resumable_upload(self) -> bool: + """Return True if this method is a resumable upload method.""" + if not self.resumable_upload_prefix: + return False + + try: + if hasattr(self, "options") and self.options: + http = self.options.Extensions[annotations_pb2.http] + if getattr(http, "media_upload", None) and getattr(http.media_upload, "enabled", False): + return True + for binding in getattr(http, "additional_bindings", ()): + if getattr(binding, "media_upload", None) and getattr(binding.media_upload, "enabled", False): + return True + except Exception: + pass + + # TODO(cl/964122389): TEMPORARY - Remove this hardcoded fallback once + # the media_upload annotation is published in cl/964122389 and added to gapic-showcase proto. + pb_name = getattr(self.method_pb, "name", "") + method_name = getattr(self, "name", "") + if pb_name == "UploadMedia" or method_name == "upload_media": + return True + + return False + @property def path_params(self) -> Sequence[str]: """Return the path parameters found in the http annotation path template""" @@ -2208,6 +2235,11 @@ def has_pagers(self) -> bool: """Return whether the service has paged methods.""" return any(m.paged_result_field for m in self.methods.values()) + @property + def has_resumable_upload_methods(self) -> bool: + """Return whether the service has resumable upload methods.""" + return any(m.is_resumable_upload for m in self.methods.values()) + @property def host(self) -> str: """Return the hostname for this service, if specified. diff --git a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/_shared_macros.j2 b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/_shared_macros.j2 index e39425bb8117..a830ae89810a 100644 --- a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/_shared_macros.j2 +++ b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/_shared_macros.j2 @@ -53,6 +53,17 @@ except ImportError: # pragma: NO COVER {% endmacro %} {% macro create_metadata(method) %} + {% if method.is_resumable_upload %} + metadata = () if metadata is None else metadata + resumable_metadata = { + "x-goog-upload-protocol": "resumable", + "x-goog-upload-command": "start", + } + existing_keys = {k.lower() for k, _ in metadata} + metadata = tuple(metadata) + tuple( + (k, v) for k, v in resumable_metadata.items() if k not in existing_keys + ) + {% endif %} {% if method.explicit_routing %} header_params: dict[str, str] = {} {% if not method.client_streaming %} @@ -132,13 +143,13 @@ from google.longrunning import operations_pb2 # type: ignore {% endif %}{# import_ns.has_operations_mixin #} {% endmacro %} -{% macro http_options_method(rules) %} +{% macro http_options_method(rules, is_resumable_upload=False, resumable_upload_prefix="resumable/upload") %} @staticmethod def _get_http_options(): http_options: List[Dict[str, str]] = [ {%- for rule in rules %}{ 'method': '{{ rule.method }}', - 'uri': '{{ rule.uri }}', + 'uri': '{% if is_resumable_upload %}/{{ resumable_upload_prefix }}{% endif %}{{ rule.uri }}', {% if rule.body %} 'body': '{{ rule.body }}', {% endif %}{# rule.body #} diff --git a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc.py.j2 b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc.py.j2 index e906c9d9ea71..0a861a3ed12a 100644 --- a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc.py.j2 +++ b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc.py.j2 @@ -10,6 +10,7 @@ import pickle import warnings from typing import Callable, Dict, Optional, Sequence, Tuple, Union +from google.api_core import exceptions as core_exceptions from google.api_core import grpc_helpers {% if service.has_lro %} from google.api_core import operations_v1 @@ -49,6 +50,9 @@ from google.longrunning import operations_pb2 # type: ignore {% endif %} {% endfilter %} from .base import {{ service.name }}Transport, DEFAULT_CLIENT_INFO +{% if service.has_resumable_upload_methods %} +from .rest import {{ service.name }}RestTransport +{% endif %} try: from google.api_core import client_logging # type: ignore @@ -353,11 +357,27 @@ class {{ service.name }}GrpcTransport({{ service.name }}Transport): # gRPC handles serialization and deserialization, so we just need # to pass in the functions for each. if '{{ method.transport_safe_name|snake_case }}' not in self._stubs: + {% if method.is_resumable_upload %} + if not self._credentials: + def _error_stub(*args, **kwargs): + raise core_exceptions.GoogleAPICallError( + "Resumable upload methods operate over REST and cannot be invoked when the transport is initialized with a pre-constructed gRPC channel. Please supply credentials directly instead of a gRPC channel to use resumable upload functionality." + ) + self._stubs['{{ method.transport_safe_name|snake_case }}'] = _error_stub + else: + rest_transport = {{ service.name }}RestTransport( + host=self._host, + credentials=self._credentials, + client_info=self._client_info, + ) + self._stubs['{{ method.transport_safe_name|snake_case }}'] = rest_transport.{{ method.transport_safe_name|snake_case }} + {% else %} self._stubs['{{ method.transport_safe_name|snake_case }}'] = self._logged_channel.{{ method.grpc_stub_type }}( '/{{ '.'.join(method.meta.address.package) }}.{{ service.name }}/{{ method.name }}', request_serializer={{ method.input.ident }}.{% if method.input.ident.python_import.module.endswith('_pb2') %}SerializeToString{% else %}serialize{% endif %}, response_deserializer={{ method.output.ident }}.{% if method.output.ident.python_import.module.endswith('_pb2') %}FromString{% else %}deserialize{% endif %}, ) + {% endif %} return self._stubs['{{ method.transport_safe_name|snake_case }}'] {% endfor %} diff --git a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc_asyncio.py.j2 b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc_asyncio.py.j2 index 7b8a885d227c..cef9e47eab77 100644 --- a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc_asyncio.py.j2 +++ b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/grpc_asyncio.py.j2 @@ -54,6 +54,13 @@ from google.longrunning import operations_pb2 # type: ignore {% endfilter %} from .base import {{ service.name }}Transport, DEFAULT_CLIENT_INFO from .grpc import {{ service.name }}GrpcTransport +{% if service.has_resumable_upload_methods %} +try: + from .rest_asyncio import Async{{ service.name }}RestTransport + HAS_ASYNC_REST = True +except ImportError: + HAS_ASYNC_REST = False +{% endif %} try: from google.api_core import client_logging # type: ignore @@ -358,11 +365,31 @@ class {{ service.grpc_asyncio_transport_name }}({{ service.name }}Transport): # gRPC handles serialization and deserialization, so we just need # to pass in the functions for each. if '{{ method.transport_safe_name|snake_case }}' not in self._stubs: + {% if method.is_resumable_upload %} + if not self._credentials: + async def _error_stub(*args, **kwargs): + raise core_exceptions.GoogleAPICallError( + "Resumable upload methods operate over REST and cannot be invoked when the transport is initialized with a pre-constructed gRPC channel. Please supply credentials directly instead of a gRPC channel to use resumable upload functionality." + ) + self._stubs['{{ method.transport_safe_name|snake_case }}'] = _error_stub + elif HAS_ASYNC_REST: + rest_transport = Async{{ service.name }}RestTransport( + host=self._host, + credentials=self._credentials, + client_info=self._client_info, + ) + self._stubs['{{ method.transport_safe_name|snake_case }}'] = rest_transport.{{ method.transport_safe_name|snake_case }} + else: + async def _unsupported_stub(*args, **kwargs): + raise NotImplementedError("Async REST transport is required for async resumable upload methods.") + self._stubs['{{ method.transport_safe_name|snake_case }}'] = _unsupported_stub + {% else %} self._stubs['{{ method.transport_safe_name|snake_case }}'] = self._logged_channel.{{ method.grpc_stub_type }}( '/{{ '.'.join(method.meta.address.package) }}.{{ service.name }}/{{ method.name }}', request_serializer={{ method.input.ident }}.{% if method.input.ident.python_import.module.endswith('_pb2') %}SerializeToString{% else %}serialize{% endif %}, response_deserializer={{ method.output.ident }}.{% if method.output.ident.python_import.module.endswith('_pb2') %}FromString{% else %}deserialize{% endif %}, ) + {% endif %} return self._stubs['{{ method.transport_safe_name|snake_case }}'] {% endfor %} diff --git a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest.py.j2 b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest.py.j2 index 1bc499c068ee..1f6e94d63cdb 100644 --- a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest.py.j2 +++ b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest.py.j2 @@ -264,7 +264,8 @@ class {{service.name}}RestTransport(_Base{{ service.name }}RestTransport): pb_resp = resp {% endif %} - json_format.Parse(response.content, pb_resp, ignore_unknown_fields=True) + if response.content and response.content.strip(): + json_format.Parse(response.content, pb_resp, ignore_unknown_fields=True) {% endif %}{# method.lro #} {#- TODO(https://github.com/googleapis/gapic-generator-python/issues/2274): Add debug log before intercepting a request #} resp = self._interceptor.post_{{ method.name|snake_case }}(resp) diff --git a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest_asyncio.py.j2 b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest_asyncio.py.j2 index 0f79d6e1ffef..583880d884f0 100644 --- a/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest_asyncio.py.j2 +++ b/packages/gapic-generator/gapic/templates/%namespace/%name_%version/%sub/services/%service/transports/rest_asyncio.py.j2 @@ -221,7 +221,8 @@ class Async{{service.name}}RestTransport(_Base{{ service.name }}RestTransport): pb_resp = resp {% endif %}{# if method.output.ident.is_proto_plus_type #} content = await response.read() - json_format.Parse(content, pb_resp, ignore_unknown_fields=True) + if content and content.strip(): + json_format.Parse(content, pb_resp, ignore_unknown_fields=True) {% endif %}{# if method.server_streaming #} resp = await self._interceptor.post_{{ method.name|snake_case }}(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] diff --git a/packages/gapic-generator/gapic/templates/tests/unit/gapic/%name_%version/%sub/test_%service.py.j2 b/packages/gapic-generator/gapic/templates/tests/unit/gapic/%name_%version/%sub/test_%service.py.j2 index 68e754caf287..bd573eb1edb1 100644 --- a/packages/gapic-generator/gapic/templates/tests/unit/gapic/%name_%version/%sub/test_%service.py.j2 +++ b/packages/gapic-generator/gapic/templates/tests/unit/gapic/%name_%version/%sub/test_%service.py.j2 @@ -1521,6 +1521,33 @@ def test_{{ service.name|snake_case }}_grpc_asyncio_transport_channel(): assert transport._ssl_channel_credentials == None +{% if service.has_resumable_upload_methods and 'grpc' in opts.transport %} +{% for method in service.methods.values() if method.is_resumable_upload %} +def test_{{ service.name|snake_case }}_{{ method.name|snake_case }}_grpc_channel_without_credentials_error(): + channel = grpc.secure_channel('http://localhost/', grpc.local_channel_credentials()) + transport = transports.{{ service.name }}GrpcTransport( + host="localhost:7469", + channel=channel, + ) + with pytest.raises(core_exceptions.GoogleAPICallError) as exc_info: + transport.{{ method.transport_safe_name|snake_case }}({{ method.input.ident }}()) + assert "operate over REST and cannot be invoked when the transport is initialized with a pre-constructed gRPC channel" in str(exc_info.value) + + +@pytest.mark.asyncio +async def test_{{ service.name|snake_case }}_{{ method.name|snake_case }}_grpc_asyncio_channel_without_credentials_error(): + channel = aio.secure_channel('http://localhost/', grpc.local_channel_credentials()) + transport = transports.{{ service.name }}GrpcAsyncIOTransport( + host="localhost:7469", + channel=channel, + ) + with pytest.raises(core_exceptions.GoogleAPICallError) as exc_info: + await transport.{{ method.transport_safe_name|snake_case }}({{ method.input.ident }}()) + assert "operate over REST and cannot be invoked when the transport is initialized with a pre-constructed gRPC channel" in str(exc_info.value) +{% endfor %} +{% endif %} + + # Remove this test when deprecated arguments (api_mtls_endpoint, client_cert_source) are # removed from grpc/grpc_asyncio transport constructor. @pytest.mark.filterwarnings("ignore::FutureWarning") diff --git a/packages/gapic-generator/gapic/utils/options.py b/packages/gapic-generator/gapic/utils/options.py index 494058d155ab..fbe44b91badd 100644 --- a/packages/gapic-generator/gapic/utils/options.py +++ b/packages/gapic-generator/gapic/utils/options.py @@ -52,6 +52,7 @@ class Options: proto_plus_deps: Tuple[str, ...] = dataclasses.field(default=("",)) gapic_version: str = "0.0.0" resource_name_aliases: Dict[str, str] = dataclasses.field(default_factory=dict) + resumable_upload_prefix: str = "" # Class constants PYTHON_GAPIC_PREFIX: str = "python-gapic-" @@ -78,6 +79,8 @@ class Options: # resource path to a custom TitleCase alias. # Format: resource.path/Name:AliasName "resource-name-alias", + # Prefix for resumable upload requests + "resumable-upload-prefix", ) ) @@ -222,6 +225,8 @@ def tweak_path(p): "Expected format is 'resource.path/Name:AliasName'." ) + resumable_upload_prefix = opts.pop("resumable-upload-prefix", [""])[0] + answer = Options( name=opts.pop("name", [""]).pop(), namespace=tuple(opts.pop("namespace", [])), @@ -245,6 +250,7 @@ def tweak_path(p): proto_plus_deps=proto_plus_deps, gapic_version=opts.pop("gapic-version", ["0.0.0"]).pop(), resource_name_aliases=resource_name_aliases, + resumable_upload_prefix=resumable_upload_prefix, ) # Note: if we ever need to recursively check directories for sample diff --git a/packages/gapic-generator/tests/system/conftest.py b/packages/gapic-generator/tests/system/conftest.py index 73169dd8a79f..0497002271bd 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,29 @@ 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 + diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_basic.py b/packages/gapic-generator/tests/system/test_resumable_upload_basic.py new file mode 100644 index 000000000000..ec4a44f448ad --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_basic.py @@ -0,0 +1,213 @@ +# 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 as core_exceptions +from google.api_core.resumable_transfer import ( + ResumableUploadConfig, + ResumableUploadSession, +) +from google.showcase import ( + ResumableUploadServiceClient, + UploadMediaRequest, + UploadMediaResponse, +) + + +def resume_resumable_upload( + transport, + upload_url, + stream, + size=None, + config=None, + **kwargs, +): + 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( + config=config, + resumable_url=upload_url, + transport=transport, + ) + return session.resume( + upload_url=upload_url, + stream=stream, + size=size, + transport=transport, + ) + + +def test_resumable_upload_start(intercepted_resumable_upload_rest): + client, interceptor = intercepted_resumable_upload_rest + response = client.upload_media(request=UploadMediaRequest(name="test_file.txt")) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "" + assert response.size == 0 + + # Verify that the generated client automatically added the resumable upload protocol headers + req_meta = dict(interceptor.request_metadata) + assert req_meta.get("x-goog-upload-protocol") == "resumable" + assert req_meta.get("x-goog-upload-command") == "start" + + # Verify that the server responded with active status and upload URL + resp_meta = {k.lower(): str(v) for k, v in interceptor.response_metadata} + assert resp_meta.get("x-goog-upload-status") == "active" + assert "x-goog-upload-url" in resp_meta + + +def test_resumable_upload_custom_metadata(intercepted_resumable_upload_rest): + client, interceptor = intercepted_resumable_upload_rest + custom_metadata = [("x-custom-header", "custom-val")] + response = client.upload_media( + request=UploadMediaRequest(name="test_file.txt"), + metadata=custom_metadata, + ) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "" + assert response.size == 0 + + req_meta = dict(interceptor.request_metadata) + assert req_meta.get("x-goog-upload-protocol") == "resumable" + assert req_meta.get("x-goog-upload-command") == "start" + assert req_meta.get("x-custom-header") == "custom-val" + + resp_meta = {k.lower(): str(v) for k, v in interceptor.response_metadata} + assert resp_meta.get("x-goog-upload-status") == "active" + + +def test_resumable_upload_finalize_response(intercepted_resumable_upload_rest): + client, interceptor = intercepted_resumable_upload_rest + + # 1. Start upload session + client.upload_media(request=UploadMediaRequest(name="test_file.txt")) + resp_meta = {k.lower(): str(v) for k, v in interceptor.response_metadata} + upload_url = resp_meta.get("x-goog-upload-url") + + # 2. Upload and finalize using the resumable media helper + stream = io.BytesIO(b"Hello world from resumable upload!") + finalize_response = resume_resumable_upload( + transport=client.transport._session, + upload_url=upload_url, + stream=stream, + ) + assert finalize_response.status_code == 200 + + # 3. Deserialize backend response proto + final_response = UploadMediaResponse.from_json(finalize_response.content) + + # Verify that the backend service returned the resource name and size matching the request + assert final_response.name == "test_file.txt" + assert final_response.size == len(stream.getvalue()) + + +@pytest.mark.parametrize( + "content_type, payload", + [ + ("application/json", b'{"name": "test_file.json"}'), + ("text/plain", b"Hello, this is plain text content!"), + ("application/octet-stream", b"\x00\x01\x02\x03\x04\x05\xff\xfe"), + ("image/png", b"\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00"), + ], +) +def test_resumable_upload_different_content_types( + intercepted_resumable_upload_rest, content_type, payload +): + client, interceptor = intercepted_resumable_upload_rest + + # 1. Start upload session + client.upload_media(request=UploadMediaRequest(name="test_upload")) + resp_meta = {k.lower(): str(v) for k, v in interceptor.response_metadata} + upload_url = resp_meta.get("x-goog-upload-url") + + # 2. Upload and finalize with the specific Content-Type using helper + finalize_response = resume_resumable_upload( + transport=client.transport._session, + upload_url=upload_url, + stream=payload, + content_type=content_type, + ) + assert finalize_response.status_code == 200 + + # 3. Deserialize and verify response name and size + final_response = UploadMediaResponse.from_json(finalize_response.content) + assert final_response.name == "test_upload" + assert final_response.size == len(payload) + + + + +def test_resumable_upload_session_direct_execution(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"Direct execution payload using ResumableUploadSession!" + + config = ResumableUploadConfig( + response_type=UploadMediaResponse, + chunk_size=256, + ) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + ) + + response = session.upload( + stream=io.BytesIO(payload), + request_body='{"name": "direct_session.txt"}', + ) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "direct_session.txt" + assert response.size == len(payload) + + # Verify session properties + assert session.finished is True + assert session.bytes_uploaded == len(payload) + assert session.response is not None + assert session.response.name == "direct_session.txt" + assert session.upload_url is not None + assert "sid=" in session.upload_url + assert session.chunk_size > 0 + + +def test_resumable_upload_session_raw_bytes_payload(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + payload = b"Raw bytes payload directly passed to upload() method" + + config = ResumableUploadConfig(response_type=UploadMediaResponse) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + ) + + response = session.upload( + stream=payload, + request_body='{"name": "raw_bytes_session.txt"}', + ) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "raw_bytes_session.txt" + assert response.size == len(payload) + assert session.response == response diff --git a/packages/gapic-generator/tests/system/test_resumable_upload_errors.py b/packages/gapic-generator/tests/system/test_resumable_upload_errors.py new file mode 100644 index 000000000000..a1b34b376322 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_errors.py @@ -0,0 +1,253 @@ +# 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 ( + ResumableUploadConfig, + ResumableUploadSession, + UnseekableStreamError, + UploadCancelledError, +) +from google.showcase import UploadMediaResponse + + +def make_resumable_upload( + transport, + request_body, + stream, + upload_url, + size=None, + config=None, + **kwargs, +): + 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, + ) + return session.upload( + stream=stream, + request_body=request_body, + size=size, + transport=transport, + ) + + +class StrictlyUnseekableStream(io.RawIOBase): + """Stream that disallows seeking to test unseekable error handling.""" + + def __init__(self, data: bytes): + self._data = data + self._pos = 0 + + def readable(self) -> bool: + return True + + def seekable(self) -> bool: + return False + + def readinto(self, b) -> int: + if self._pos >= len(self._data): + return 0 + n = min(len(b), len(self._data) - self._pos) + b[:n] = self._data[self._pos : self._pos + n] + self._pos += n + return n + + def seek(self, offset, whence=io.SEEK_SET): + raise io.UnsupportedOperation("Stream does not support seeking") + + +def test_resumable_upload_exception_surfaces_attributes(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "error_attributes.txt"}' + data = b"E" * 1024 + stream = io.BytesIO(data) + + # Server scenario terminates chunk upload with non-recoverable error + scenario_headers = [ + ("X-Goog-Test-Scenario", "non_fatal_error_on_chunk_upload"), + ("X-Goog-Test-Scenario-Config", '{"error_code":403,"failure_count":1}'), + ] + + config = ResumableUploadConfig( + chunk_size=512, + headers=scenario_headers, + ) + + with pytest.raises(exceptions.GoogleAPICallError) as exc_info: + make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + config=config, + ) + + # Verify exception attributes + assert hasattr(exc_info.value, "upload_url") + assert exc_info.value.upload_url is not None + assert "/upload?sid=" in exc_info.value.upload_url + assert hasattr(exc_info.value, "chunk_size") + assert exc_info.value.chunk_size == 262144 + + +def test_resumable_upload_crash_recovery_flow(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "crash_recovery.txt"}' + data = b"C" * 1536 + stream = io.BytesIO(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + + # Session 1: Client begins upload + config1 = ResumableUploadConfig( + chunk_size=512, + headers=scenario_headers, + response_type=UploadMediaResponse, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + ) + session1.initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + + # First chunk succeeds + session1._transmit_chunk(client.transport._session, stream, len(data)) + assert session1.bytes_uploaded == 512 + + # Simulate process crash by capturing URL and chunk size + crashed_url = session1.upload_url + crashed_chunk_size = session1.chunk_size + + # Session 2: Fresh process recovers upload from captured URL + config2 = ResumableUploadConfig( + chunk_size=crashed_chunk_size, + response_type=UploadMediaResponse, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + ) + + # Rewind stream to beginning (full file available in new process) + stream.seek(0) + response = session2.resume( + upload_url=crashed_url, + stream=stream, + size=len(data), + transport=client.transport._session, + ) + + assert isinstance(response, UploadMediaResponse) + assert response.name == "crash_recovery.txt" + assert response.size == len(data) + assert session2.bytes_uploaded == len(data) + assert session2.finished + + +def test_resumable_upload_unseekable_stream_beyond_buffer_raises(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "unseekable_error.bin"}' + data = b"U" * 2048 + unseekable = StrictlyUnseekableStream(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + config = ResumableUploadConfig( + chunk_size=512, + headers=scenario_headers, + response_type=UploadMediaResponse, + ) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + ) + session.initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + + # First chunk is transmitted and committed; in-memory buffer is discarded + session._transmit_chunk(client.transport._session, unseekable, len(data)) + assert session.bytes_uploaded == 512 + assert session._buffered_chunk is None + + # Simulating server recovery request to offset 0 (outside discarded buffer) + # on an unseekable stream must raise UnseekableStreamError + with pytest.raises(UnseekableStreamError) as exc_info: + session._reposition_stream_offset(unseekable, 0) + + assert hasattr(exc_info.value, "upload_url") + assert exc_info.value.upload_url == session.upload_url + assert exc_info.value.chunk_size == 512 + + +def test_resumable_upload_session_cancellation(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "cancel_session.txt"}' + data = b"X" * 1024 + stream = io.BytesIO(data) + + scenario_headers = [ + ("X-Goog-Test-Scenario", "chunk_granularity"), + ] + session = ResumableUploadSession( + upload_url=initial_url, + config=ResumableUploadConfig(chunk_size=512, headers=scenario_headers), + transport=client.transport._session, + ) + session.initiate( + transport=client.transport._session, + request_body=request_body, + size=len(data), + ) + assert session.upload_url is not None + + # Transmit first chunk + session._transmit_chunk(client.transport._session, stream, len(data)) + assert session.bytes_uploaded == 512 + + # Cancel session + session.cancel(transport=client.transport._session) + + # Attempting to upload to cancelled session triggers query which discovers cancelled status (410 Gone) or 400 + with pytest.raises(exceptions.GoogleAPICallError) as exc_info: + session._transmit_chunk(client.transport._session, stream, len(data)) + + assert exc_info.value.code in (400, 410) 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..cad2e1f6ad55 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_progress.py @@ -0,0 +1,187 @@ +# 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.resumable_transfer import ( + ProgressState, + ResumableUploadConfig, + ResumableUploadSession, + UploadProgress, +) +from google.showcase import UploadMediaResponse + + +def make_resumable_upload( + transport, + request_body, + stream, + upload_url, + size=None, + config=None, + **kwargs, +): + 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, + ) + return session.upload( + stream=stream, + request_body=request_body, + size=size, + transport=transport, + ) + + +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 response.status_code == 200 + + final_response = UploadMediaResponse.from_json(response.content) + assert final_response.name == "full_e2e_upload.txt" + assert final_response.size == len(stream.getvalue()) + + +def test_resumable_upload_callback_progress_tracking(intercepted_resumable_upload_rest): + client, _ = intercepted_resumable_upload_rest + initial_url = f"{client.transport._host}/resumable/upload/v1beta1/media/upload" + request_body = '{"name": "progress_tracked_upload.txt"}' + payload = b"0123456789" * 100 + stream = io.BytesIO(payload) + + progress_events = [] + + def on_progress(p): + progress_events.append((p.bytes_uploaded, p.state.value, p.upload_url, p.chunk_size)) + + scenario_headers = [("X-Goog-Test-Scenario", "chunk_granularity")] + config = ResumableUploadConfig( + chunk_size=256, + on_progress=on_progress, + headers=scenario_headers, + ) + + response = make_resumable_upload( + transport=client.transport._session, + request_body=request_body, + stream=stream, + upload_url=initial_url, + config=config, + ) + assert response.status_code == 200 + + # Verify progress notifications occurred in order + assert len(progress_events) >= 3 + assert progress_events[0][1] == "started" + assert progress_events[-1][1] == "finalized" + assert progress_events[-1][0] == len(payload) + for bytes_up, state, u, chunk_sz in progress_events: + assert "sid=" in u + assert chunk_sz == 256 + + +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, + response_type=UploadMediaResponse, + headers=scenario_headers, + ) + session = ResumableUploadSession( + upload_url=initial_url, + config=config, + transport=client.transport._session, + ) + + 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, + chunk_size=512, + headers=scenario_headers, + ) + assert response.status_code == 200 + final_response = UploadMediaResponse.from_json(response.content) + assert final_response.name == "unseekable_stream_upload.txt" + assert final_response.size == len(data) 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..d681da68e8e8 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_resume.py @@ -0,0 +1,285 @@ +# 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.resumable_transfer import ( + ProgressState, + ResumableUploadConfig, + ResumableUploadSession, +) +from google.showcase import UploadMediaResponse + + +def resume_resumable_upload( + transport, + upload_url, + stream, + size=None, + config=None, + **kwargs, +): + 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( + config=config, + resumable_url=upload_url, + transport=transport, + ) + return session.resume( + upload_url=upload_url, + stream=stream, + size=size, + transport=transport, + ) + + +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, + response_type=UploadMediaResponse, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + ) + 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, + response_type=UploadMediaResponse, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + ) + + # 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, + response_type=UploadMediaResponse, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + ) + 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, + response_type=UploadMediaResponse, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + ) + + 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, + response_type=UploadMediaResponse, + ) + session1 = ResumableUploadSession( + upload_url=initial_url, + config=config1, + transport=client.transport._session, + ) + 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, + response_type=UploadMediaResponse, + ) + session2 = ResumableUploadSession( + config=config2, + transport=client.transport._session, + ) + + 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, + response_type=UploadMediaResponse, + ) + final_response = resume_resumable_upload( + transport=client.transport._session, + upload_url=saved_url, + stream=data, + config=config2, + ) + + assert isinstance(final_response, UploadMediaResponse) + assert final_response.name == "helper_bytes.bin" + assert final_response.size == len(data) 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..513994eab4f9 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_scenarios.py @@ -0,0 +1,216 @@ +# 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 as core_exceptions +from google.api_core.resumable_transfer import ( + ResumableUploadConfig, + ResumableUploadSession, +) +from google.showcase import UploadMediaResponse + + +def make_resumable_upload( + transport, + request_body, + stream, + upload_url, + size=None, + config=None, + **kwargs, +): + 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, + ) + return session.upload( + stream=stream, + request_body=request_body, + size=size, + transport=transport, + ) + + +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", '{"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 response.status_code == 200 + final_response = UploadMediaResponse.from_json(response.content) + 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, + ) + + +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 response.status_code == 200 + assert attempt_count[0] >= 2 # Verified that start was retried upon missing status header + final_response = UploadMediaResponse.from_json(response.content) + 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 response.status_code == 200 + final_response = UploadMediaResponse.from_json(response.content) + 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 response.status_code == 200 + final_response = UploadMediaResponse.from_json(response.content) + 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 response.status_code == 200 + final_response = UploadMediaResponse.from_json(response.content) + assert final_response.name == "granularity_upload.txt" + assert final_response.size == len(data) 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..b78d7f7226c2 --- /dev/null +++ b/packages/gapic-generator/tests/system/test_resumable_upload_stall.py @@ -0,0 +1,194 @@ +# 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, + ResumableUploadSession, + TransferStalledError, +) +from google.showcase import UploadMediaResponse + + +def make_resumable_upload( + transport, + request_body, + stream, + upload_url, + size=None, + config=None, + **kwargs, +): + 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, + ) + return session.upload( + stream=stream, + request_body=request_body, + size=size, + transport=transport, + ) + + +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_min_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 response.status_code == 200 + final_response = UploadMediaResponse.from_json(response.content) + 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_min_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_min_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_min_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() + diff --git a/packages/gapic-generator/tests/unit/generator/test_options.py b/packages/gapic-generator/tests/unit/generator/test_options.py index df945d68f417..d9090af346d2 100644 --- a/packages/gapic-generator/tests/unit/generator/test_options.py +++ b/packages/gapic-generator/tests/unit/generator/test_options.py @@ -280,3 +280,12 @@ def test_options_resource_name_aliases(): # 3. MissingAlias: # (The empty ' ' string safely 'continues' without warning, as intended) assert warn.call_count == 3 + + +def test_options_resumable_upload_prefix(): + opts_default = Options.build("") + assert opts_default.resumable_upload_prefix == "" + + opts_custom = Options.build("resumable-upload-prefix=custom/upload/prefix") + assert opts_custom.resumable_upload_prefix == "custom/upload/prefix" + diff --git a/packages/gapic-generator/tests/unit/schema/wrappers/test_method.py b/packages/gapic-generator/tests/unit/schema/wrappers/test_method.py index 87eb959cb894..bfc4f85fe27c 100644 --- a/packages/gapic-generator/tests/unit/schema/wrappers/test_method.py +++ b/packages/gapic-generator/tests/unit/schema/wrappers/test_method.py @@ -18,6 +18,7 @@ import pytest from typing import Sequence +from google.api import annotations_pb2 from google.api import field_behavior_pb2 from google.api import http_pb2 from google.api import routing_pb2 @@ -129,6 +130,35 @@ def test_method_client_output_async_empty(): assert method.client_output_async == wrappers.PrimitiveType.build(None) +def test_method_is_resumable_upload_missing_annotation(): + opts = descriptor_pb2.MethodOptions() + http = opts.Extensions[annotations_pb2.http] + http.post = "/v1/upload" + method_pb = descriptor_pb2.MethodDescriptorProto( + name="Upload", + input_type=".foo.bar.v1.Input", + output_type=".foo.bar.v1.Output", + options=opts, + ) + method = wrappers.Method( + method_pb=method_pb, + input=make_message("Input"), + output=make_message("Output"), + resumable_upload_prefix="resumable/upload", + ) + assert method.is_resumable_upload is False + + +def test_method_is_resumable_upload_without_prefix(): + method = make_method("Upload") + assert method.is_resumable_upload is False + + +def test_method_is_resumable_upload_default(): + method_normal = make_method("GetFoo") + assert method_normal.is_resumable_upload is False + + def test_method_paged_result_field_not_first(): paged = make_field(name="foos", message=make_message("Foo"), repeated=True) input_msg = make_message( @@ -1118,3 +1148,24 @@ def test__validate_paged_field_size_type(field_type, pb_type, expected): actual = method._validate_paged_field_size_type(page_field_size=page_size) assert actual == expected + + +def test_method_is_resumable_upload(): + # Without resumable_upload_prefix, method is not resumable upload + method_no_prefix = make_method("UploadMedia") + assert not method_no_prefix.is_resumable_upload + + # With resumable_upload_prefix and UploadMedia method name + method_with_prefix = dataclasses.replace( + make_method("UploadMedia"), + resumable_upload_prefix="resumable/upload", + ) + assert method_with_prefix.is_resumable_upload + + # Non-resumable method with prefix + method_other = dataclasses.replace( + make_method("OtherMethod"), + resumable_upload_prefix="resumable/upload", + ) + assert not method_other.is_resumable_upload + diff --git a/packages/gapic-generator/tests/unit/schema/wrappers/test_service.py b/packages/gapic-generator/tests/unit/schema/wrappers/test_service.py index 2a58ab500c7e..1c7f897511b9 100644 --- a/packages/gapic-generator/tests/unit/schema/wrappers/test_service.py +++ b/packages/gapic-generator/tests/unit/schema/wrappers/test_service.py @@ -13,6 +13,7 @@ # limitations under the License. import collections +import dataclasses import itertools import pytest import typing @@ -749,3 +750,25 @@ def test_resource_messages_raises_on_malformed_typeless_resource(): # 2. Trigger the property and expect it to fail fast with the AIP-123 URL with pytest.raises(ValueError, match="https://google.aip.dev/123"): _ = service.resource_messages + + +def test_service_has_resumable_upload_methods(): + m_upload = dataclasses.replace( + make_method("UploadMedia"), + resumable_upload_prefix="resumable/upload", + ) + m_status = make_method("GetStatus") + + service_with_resumable = make_service( + name="ResumableService", + methods=(m_upload, m_status), + ) + assert service_with_resumable.has_resumable_upload_methods + + m_other = make_method("DoThing") + service_without_resumable = make_service( + name="StandardService", + methods=(m_other,), + ) + assert not service_without_resumable.has_resumable_upload_methods +