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/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/grpc.py index 848bb1096cbe..d8448e6a2a49 100755 --- a/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/rest.py b/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/rest.py index d85aa16473c2..06b9e3b036bc 100755 --- a/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/rest.py +++ b/packages/gapic-generator/tests/integration/goldens/asset/google/cloud/asset_v1/services/asset_service/transports/rest.py @@ -1283,7 +1283,8 @@ def __call__(self, resp = asset_service.AnalyzeIamPolicyResponse() pb_resp = asset_service.AnalyzeIamPolicyResponse.pb(resp) - 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) resp = self._interceptor.post_analyze_iam_policy(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1538,7 +1539,8 @@ def __call__(self, resp = asset_service.AnalyzeMoveResponse() pb_resp = asset_service.AnalyzeMoveResponse.pb(resp) - 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) resp = self._interceptor.post_analyze_move(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1665,7 +1667,8 @@ def __call__(self, resp = asset_service.AnalyzeOrgPoliciesResponse() pb_resp = asset_service.AnalyzeOrgPoliciesResponse.pb(resp) - 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) resp = self._interceptor.post_analyze_org_policies(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1793,7 +1796,8 @@ def __call__(self, resp = asset_service.AnalyzeOrgPolicyGovernedAssetsResponse() pb_resp = asset_service.AnalyzeOrgPolicyGovernedAssetsResponse.pb(resp) - 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) resp = self._interceptor.post_analyze_org_policy_governed_assets(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1921,7 +1925,8 @@ def __call__(self, resp = asset_service.AnalyzeOrgPolicyGovernedContainersResponse() pb_resp = asset_service.AnalyzeOrgPolicyGovernedContainersResponse.pb(resp) - 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) resp = self._interceptor.post_analyze_org_policy_governed_containers(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -2045,7 +2050,8 @@ def __call__(self, resp = asset_service.BatchGetAssetsHistoryResponse() pb_resp = asset_service.BatchGetAssetsHistoryResponse.pb(resp) - 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) resp = self._interceptor.post_batch_get_assets_history(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -2173,7 +2179,8 @@ def __call__(self, resp = asset_service.BatchGetEffectiveIamPoliciesResponse() pb_resp = asset_service.BatchGetEffectiveIamPoliciesResponse.pb(resp) - 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) resp = self._interceptor.post_batch_get_effective_iam_policies(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -2306,7 +2313,8 @@ def __call__(self, resp = asset_service.Feed() pb_resp = asset_service.Feed.pb(resp) - 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) resp = self._interceptor.post_create_feed(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -2433,7 +2441,8 @@ def __call__(self, resp = asset_service.SavedQuery() pb_resp = asset_service.SavedQuery.pb(resp) - 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) resp = self._interceptor.post_create_saved_query(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -2871,7 +2880,8 @@ def __call__(self, resp = asset_service.Feed() pb_resp = asset_service.Feed.pb(resp) - 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) resp = self._interceptor.post_get_feed(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -2997,7 +3007,8 @@ def __call__(self, resp = asset_service.SavedQuery() pb_resp = asset_service.SavedQuery.pb(resp) - 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) resp = self._interceptor.post_get_saved_query(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3121,7 +3132,8 @@ def __call__(self, resp = asset_service.ListAssetsResponse() pb_resp = asset_service.ListAssetsResponse.pb(resp) - 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) resp = self._interceptor.post_list_assets(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3245,7 +3257,8 @@ def __call__(self, resp = asset_service.ListFeedsResponse() pb_resp = asset_service.ListFeedsResponse.pb(resp) - 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) resp = self._interceptor.post_list_feeds(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3369,7 +3382,8 @@ def __call__(self, resp = asset_service.ListSavedQueriesResponse() pb_resp = asset_service.ListSavedQueriesResponse.pb(resp) - 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) resp = self._interceptor.post_list_saved_queries(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3494,7 +3508,8 @@ def __call__(self, resp = asset_service.QueryAssetsResponse() pb_resp = asset_service.QueryAssetsResponse.pb(resp) - 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) resp = self._interceptor.post_query_assets(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3618,7 +3633,8 @@ def __call__(self, resp = asset_service.SearchAllIamPoliciesResponse() pb_resp = asset_service.SearchAllIamPoliciesResponse.pb(resp) - 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) resp = self._interceptor.post_search_all_iam_policies(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3742,7 +3758,8 @@ def __call__(self, resp = asset_service.SearchAllResourcesResponse() pb_resp = asset_service.SearchAllResourcesResponse.pb(resp) - 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) resp = self._interceptor.post_search_all_resources(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -3875,7 +3892,8 @@ def __call__(self, resp = asset_service.Feed() pb_resp = asset_service.Feed.pb(resp) - 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) resp = self._interceptor.post_update_feed(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4002,7 +4020,8 @@ def __call__(self, resp = asset_service.SavedQuery() pb_resp = asset_service.SavedQuery.pb(resp) - 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) resp = self._interceptor.post_update_saved_query(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] diff --git a/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/grpc.py index 18428ad7d6e0..2d8c651a2274 100755 --- a/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import gapic_v1 import google.auth # type: ignore diff --git a/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/rest.py b/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/rest.py index 0cffb09641ed..7b77eb85ece2 100755 --- a/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/rest.py +++ b/packages/gapic-generator/tests/integration/goldens/credentials/google/iam/credentials_v1/services/iam_credentials/transports/rest.py @@ -462,7 +462,8 @@ def __call__(self, resp = common.GenerateAccessTokenResponse() pb_resp = common.GenerateAccessTokenResponse.pb(resp) - 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) resp = self._interceptor.post_generate_access_token(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -587,7 +588,8 @@ def __call__(self, resp = common.GenerateIdTokenResponse() pb_resp = common.GenerateIdTokenResponse.pb(resp) - 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) resp = self._interceptor.post_generate_id_token(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -712,7 +714,8 @@ def __call__(self, resp = common.SignBlobResponse() pb_resp = common.SignBlobResponse.pb(resp) - 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) resp = self._interceptor.post_sign_blob(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -837,7 +840,8 @@ def __call__(self, resp = common.SignJwtResponse() pb_resp = common.SignJwtResponse.pb(resp) - 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) resp = self._interceptor.post_sign_jwt(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] diff --git a/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/grpc.py index ac5d9a0fbe92..955b02345f22 100755 --- a/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/rest.py b/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/rest.py index 1565671cf8d4..a51385d256eb 100755 --- a/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/rest.py +++ b/packages/gapic-generator/tests/integration/goldens/eventarc/google/cloud/eventarc_v1/services/eventarc/transports/rest.py @@ -4029,7 +4029,8 @@ def __call__(self, resp = channel.Channel() pb_resp = channel.Channel.pb(resp) - 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) resp = self._interceptor.post_get_channel(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4161,7 +4162,8 @@ def __call__(self, resp = channel_connection.ChannelConnection() pb_resp = channel_connection.ChannelConnection.pb(resp) - 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) resp = self._interceptor.post_get_channel_connection(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4292,7 +4294,8 @@ def __call__(self, resp = enrollment.Enrollment() pb_resp = enrollment.Enrollment.pb(resp) - 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) resp = self._interceptor.post_get_enrollment(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4420,7 +4423,8 @@ def __call__(self, resp = google_api_source.GoogleApiSource() pb_resp = google_api_source.GoogleApiSource.pb(resp) - 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) resp = self._interceptor.post_get_google_api_source(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4553,7 +4557,8 @@ def __call__(self, resp = google_channel_config.GoogleChannelConfig() pb_resp = google_channel_config.GoogleChannelConfig.pb(resp) - 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) resp = self._interceptor.post_get_google_channel_config(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4686,7 +4691,8 @@ def __call__(self, resp = message_bus.MessageBus() pb_resp = message_bus.MessageBus.pb(resp) - 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) resp = self._interceptor.post_get_message_bus(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4813,7 +4819,8 @@ def __call__(self, resp = pipeline.Pipeline() pb_resp = pipeline.Pipeline.pb(resp) - 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) resp = self._interceptor.post_get_pipeline(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -4940,7 +4947,8 @@ def __call__(self, resp = discovery.Provider() pb_resp = discovery.Provider.pb(resp) - 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) resp = self._interceptor.post_get_provider(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5067,7 +5075,8 @@ def __call__(self, resp = trigger.Trigger() pb_resp = trigger.Trigger.pb(resp) - 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) resp = self._interceptor.post_get_trigger(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5194,7 +5203,8 @@ def __call__(self, resp = eventarc.ListChannelConnectionsResponse() pb_resp = eventarc.ListChannelConnectionsResponse.pb(resp) - 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) resp = self._interceptor.post_list_channel_connections(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5319,7 +5329,8 @@ def __call__(self, resp = eventarc.ListChannelsResponse() pb_resp = eventarc.ListChannelsResponse.pb(resp) - 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) resp = self._interceptor.post_list_channels(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5444,7 +5455,8 @@ def __call__(self, resp = eventarc.ListEnrollmentsResponse() pb_resp = eventarc.ListEnrollmentsResponse.pb(resp) - 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) resp = self._interceptor.post_list_enrollments(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5571,7 +5583,8 @@ def __call__(self, resp = eventarc.ListGoogleApiSourcesResponse() pb_resp = eventarc.ListGoogleApiSourcesResponse.pb(resp) - 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) resp = self._interceptor.post_list_google_api_sources(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5699,7 +5712,8 @@ def __call__(self, resp = eventarc.ListMessageBusEnrollmentsResponse() pb_resp = eventarc.ListMessageBusEnrollmentsResponse.pb(resp) - 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) resp = self._interceptor.post_list_message_bus_enrollments(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5826,7 +5840,8 @@ def __call__(self, resp = eventarc.ListMessageBusesResponse() pb_resp = eventarc.ListMessageBusesResponse.pb(resp) - 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) resp = self._interceptor.post_list_message_buses(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -5953,7 +5968,8 @@ def __call__(self, resp = eventarc.ListPipelinesResponse() pb_resp = eventarc.ListPipelinesResponse.pb(resp) - 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) resp = self._interceptor.post_list_pipelines(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -6078,7 +6094,8 @@ def __call__(self, resp = eventarc.ListProvidersResponse() pb_resp = eventarc.ListProvidersResponse.pb(resp) - 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) resp = self._interceptor.post_list_providers(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -6203,7 +6220,8 @@ def __call__(self, resp = eventarc.ListTriggersResponse() pb_resp = eventarc.ListTriggersResponse.pb(resp) - 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) resp = self._interceptor.post_list_triggers(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -6719,7 +6737,8 @@ def __call__(self, resp = gce_google_channel_config.GoogleChannelConfig() pb_resp = gce_google_channel_config.GoogleChannelConfig.pb(resp) - 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) resp = self._interceptor.post_update_google_channel_config(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] diff --git a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py index d8122989787f..7de54ab790d6 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py index eeb3a8564ee0..3be2e60fc72f 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import gapic_v1 import google.auth # type: ignore diff --git a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py index 2b6003f77476..5baacdcbb32c 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/logging/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import gapic_v1 import google.auth # type: ignore diff --git a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py index d8122989787f..7de54ab790d6 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/config_service_v2/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py index eeb3a8564ee0..3be2e60fc72f 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/logging_service_v2/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import gapic_v1 import google.auth # type: ignore diff --git a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py index 2b6003f77476..5baacdcbb32c 100755 --- a/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/logging_internal/google/cloud/logging_v2/services/metrics_service_v2/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import gapic_v1 import google.auth # type: ignore diff --git a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py index addfbf37e166..b9bd3c9fd101 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest.py b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest.py index ea8778e47a84..dbd9f9965a93 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest.py +++ b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest.py @@ -1495,7 +1495,8 @@ def __call__(self, resp = cloud_redis.Instance() pb_resp = cloud_redis.Instance.pb(resp) - 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) resp = self._interceptor.post_get_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1620,7 +1621,8 @@ def __call__(self, resp = cloud_redis.InstanceAuthString() pb_resp = cloud_redis.InstanceAuthString.pb(resp) - 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) resp = self._interceptor.post_get_instance_auth_string(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1874,7 +1876,8 @@ def __call__(self, resp = cloud_redis.ListInstancesResponse() pb_resp = cloud_redis.ListInstancesResponse.pb(resp) - 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) resp = self._interceptor.post_list_instances(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] diff --git a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py index 4d629a5a8443..5835d5b521f9 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py +++ b/packages/gapic-generator/tests/integration/goldens/redis/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py @@ -1022,7 +1022,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_create_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_create_instance_with_metadata(resp, response_metadata) @@ -1154,7 +1155,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_delete_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_delete_instance_with_metadata(resp, response_metadata) @@ -1287,7 +1289,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_export_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_export_instance_with_metadata(resp, response_metadata) @@ -1420,7 +1423,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_failover_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_failover_instance_with_metadata(resp, response_metadata) @@ -1549,7 +1553,8 @@ async def __call__(self, resp = cloud_redis.Instance() pb_resp = cloud_redis.Instance.pb(resp) 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) resp = await self._interceptor.post_get_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_get_instance_with_metadata(resp, response_metadata) @@ -1678,7 +1683,8 @@ async def __call__(self, resp = cloud_redis.InstanceAuthString() pb_resp = cloud_redis.InstanceAuthString.pb(resp) 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) resp = await self._interceptor.post_get_instance_auth_string(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_get_instance_auth_string_with_metadata(resp, response_metadata) @@ -1811,7 +1817,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_import_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_import_instance_with_metadata(resp, response_metadata) @@ -1942,7 +1949,8 @@ async def __call__(self, resp = cloud_redis.ListInstancesResponse() pb_resp = cloud_redis.ListInstancesResponse.pb(resp) 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) resp = await self._interceptor.post_list_instances(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_list_instances_with_metadata(resp, response_metadata) @@ -2075,7 +2083,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_reschedule_maintenance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_reschedule_maintenance_with_metadata(resp, response_metadata) @@ -2208,7 +2217,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_update_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_update_instance_with_metadata(resp, response_metadata) @@ -2341,7 +2351,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_upgrade_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_upgrade_instance_with_metadata(resp, response_metadata) diff --git a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py index cae682b3d0ae..386b5f1cf163 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest.py b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest.py index 2f972ef00317..1f66b46d2912 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest.py +++ b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest.py @@ -977,7 +977,8 @@ def __call__(self, resp = cloud_redis.Instance() pb_resp = cloud_redis.Instance.pb(resp) - 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) resp = self._interceptor.post_get_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1104,7 +1105,8 @@ def __call__(self, resp = cloud_redis.ListInstancesResponse() pb_resp = cloud_redis.ListInstancesResponse.pb(resp) - 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) resp = self._interceptor.post_list_instances(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] diff --git a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py index 960d9639a214..c3cd5a3dbf23 100755 --- a/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py +++ b/packages/gapic-generator/tests/integration/goldens/redis_selective/google/cloud/redis_v1/services/cloud_redis/transports/rest_asyncio.py @@ -728,7 +728,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_create_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_create_instance_with_metadata(resp, response_metadata) @@ -860,7 +861,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_delete_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_delete_instance_with_metadata(resp, response_metadata) @@ -989,7 +991,8 @@ async def __call__(self, resp = cloud_redis.Instance() pb_resp = cloud_redis.Instance.pb(resp) 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) resp = await self._interceptor.post_get_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_get_instance_with_metadata(resp, response_metadata) @@ -1120,7 +1123,8 @@ async def __call__(self, resp = cloud_redis.ListInstancesResponse() pb_resp = cloud_redis.ListInstancesResponse.pb(resp) 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) resp = await self._interceptor.post_list_instances(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_list_instances_with_metadata(resp, response_metadata) @@ -1253,7 +1257,8 @@ async def __call__(self, resp = operations_pb2.Operation() pb_resp = resp 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) resp = await self._interceptor.post_update_instance(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] resp, _ = await self._interceptor.post_update_instance_with_metadata(resp, response_metadata) diff --git a/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/grpc.py b/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/grpc.py index 1f997d49aabd..fe195106e5f6 100755 --- a/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/grpc.py +++ b/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/grpc.py @@ -19,6 +19,7 @@ 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 from google.api_core import operations_v1 from google.api_core import gapic_v1 diff --git a/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/rest.py b/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/rest.py index 9a4373457926..8750c6f2aff6 100755 --- a/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/rest.py +++ b/packages/gapic-generator/tests/integration/goldens/storagebatchoperations/google/cloud/storagebatchoperations_v1/services/storage_batch_operations/transports/rest.py @@ -739,7 +739,8 @@ def __call__(self, resp = storage_batch_operations.CancelJobResponse() pb_resp = storage_batch_operations.CancelJobResponse.pb(resp) - 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) resp = self._interceptor.post_cancel_job(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1082,7 +1083,8 @@ def __call__(self, resp = storage_batch_operations_types.BucketOperation() pb_resp = storage_batch_operations_types.BucketOperation.pb(resp) - 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) resp = self._interceptor.post_get_bucket_operation(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1208,7 +1210,8 @@ def __call__(self, resp = storage_batch_operations_types.Job() pb_resp = storage_batch_operations_types.Job.pb(resp) - 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) resp = self._interceptor.post_get_job(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1335,7 +1338,8 @@ def __call__(self, resp = storage_batch_operations.ListBucketOperationsResponse() pb_resp = storage_batch_operations.ListBucketOperationsResponse.pb(resp) - 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) resp = self._interceptor.post_list_bucket_operations(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] @@ -1459,7 +1463,8 @@ def __call__(self, resp = storage_batch_operations.ListJobsResponse() pb_resp = storage_batch_operations.ListJobsResponse.pb(resp) - 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) resp = self._interceptor.post_list_jobs(resp) response_metadata = [(k, str(v)) for k, v in response.headers.items()] 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..ba8393852031 --- /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_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 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_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() + 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 +