Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
d9b6707
feat: add support for resumable uploads
parthea Sep 29, 2026
c448c74
address review feedback
parthea Sep 29, 2026
2d9db92
update goldens
parthea Sep 29, 2026
f3a255b
address review feedback
parthea Sep 30, 2026
16fbe33
address review feedback
parthea Sep 30, 2026
b50d5b3
remove warnings related to rest numberic enums not set
parthea Sep 30, 2026
4405eec
update test
parthea Sep 30, 2026
52ef3ac
address review feedback
parthea Sep 30, 2026
54fd19c
fix build
parthea Sep 30, 2026
04b91b6
fix build
parthea Sep 30, 2026
f712705
address review feedback
parthea Sep 30, 2026
7ffd24b
address review feedback
parthea Sep 30, 2026
8059ad0
update goldens
parthea Sep 30, 2026
0d9212e
update goldens
parthea Sep 30, 2026
f952d40
address review feedback
parthea Oct 1, 2026
6424adc
address review feedback
parthea Oct 1, 2026
08b304b
address review feedback
parthea Oct 1, 2026
8684167
address review feedback
parthea Oct 1, 2026
96f2f6b
address review feedback
parthea Oct 1, 2026
76103bf
address review feedback
parthea Oct 1, 2026
1d6a0ff
add test for bug where request body is not passed
parthea Oct 1, 2026
cc6b703
fix: forward transcoded request body and query params in resumable up…
parthea Oct 1, 2026
5352808
address review feedback
parthea Oct 1, 2026
2303a13
address review feedback
parthea Oct 1, 2026
8110282
address review feedback
parthea Oct 1, 2026
9e41135
fix build
parthea Oct 1, 2026
4c9067a
update unit test
parthea Oct 1, 2026
7f0f479
fix build
parthea Oct 1, 2026
98b80fd
address review feedback
parthea Oct 2, 2026
ed5d9c9
sync minimum version of google-api-core for resumable uploads
parthea Oct 2, 2026
8bd4fe8
updated noxfile
daniel-sanche Oct 2, 2026
5ea59b4
Revert "sync minimum version of google-api-core for resumable uploads"
daniel-sanche Oct 3, 2026
96549f5
allow later versions
daniel-sanche Oct 3, 2026
825a38b
Revert "address review feedback"
daniel-sanche Oct 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 22 additions & 3 deletions packages/gapic-generator/gapic/generator/generator.py
Original file line number Diff line number Diff line change
Expand Up @@ -342,7 +342,9 @@ def _render_template(
)
or (
"transport" in template_name
and not self._is_desired_transport(template_name, opts)
and not self._is_desired_transport(
template_name, opts, service=service
)
)
or
# TODO(https://github.com/googleapis/gapic-generator-python/issues/2121): Remove this condition when async rest is GA.
Expand All @@ -352,14 +354,20 @@ def _render_template(
and not api_schema.all_library_settings[
api_schema.naming.proto_package
].python_settings.experimental_features.rest_async_io_enabled
and not service.has_resumable_upload_methods
)
or (
"rest_asyncio" in template_name
and not api_schema.all_library_settings[
api_schema.naming.proto_package
].python_settings.experimental_features.rest_async_io_enabled
and not service.has_resumable_upload_methods
)
or (
"rest_base" in template_name
and "rest" not in opts.transport
and not service.has_resumable_upload_methods
)
Comment on lines 354 to 370

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The service parameter is optional and can be None when rendering non-service templates. Accessing service.has_resumable_upload_methods directly without a null check will raise an AttributeError. Please use defensive programming and guard the access with a check for service.

Suggested change
and not api_schema.all_library_settings[
api_schema.naming.proto_package
].python_settings.experimental_features.rest_async_io_enabled
and not service.has_resumable_upload_methods
)
or (
"rest_asyncio" in template_name
and not api_schema.all_library_settings[
api_schema.naming.proto_package
].python_settings.experimental_features.rest_async_io_enabled
and not service.has_resumable_upload_methods
)
or (
"rest_base" in template_name
and "rest" not in opts.transport
and not service.has_resumable_upload_methods
)
and not api_schema.all_library_settings[
api_schema.naming.proto_package
].python_settings.experimental_features.rest_async_io_enabled
and not (service and service.has_resumable_upload_methods)
)
or (
"rest_asyncio" in template_name
and not api_schema.all_library_settings[
api_schema.naming.proto_package
].python_settings.experimental_features.rest_async_io_enabled
and not (service and service.has_resumable_upload_methods)
)
or (
"rest_base" in template_name
and "rest" not in opts.transport
and not (service and service.has_resumable_upload_methods)
)

or ("rest_base" in template_name and "rest" not in opts.transport)
):
continue

Expand All @@ -386,9 +394,20 @@ def _render_template(
)
return answer

def _is_desired_transport(self, template_name: str, opts: Options) -> bool:
def _is_desired_transport(
self,
template_name: str,
opts: Options,
service: Optional[Any] = None,
) -> bool:
"""Returns true if template name contains a desired transport"""
desired_transports = ["__init__", "base", "README"] + opts.transport
if (
service is not None
and service.has_resumable_upload_methods
and "rest" not in desired_transports
):
desired_transports.append("rest")
return any(transport in template_name for transport in desired_transports)

def _get_file(
Expand Down
27 changes: 27 additions & 0 deletions packages/gapic-generator/gapic/schema/wrappers.py
Original file line number Diff line number Diff line change
Expand Up @@ -1640,6 +1640,28 @@ def _client_output(self, enable_asyncio: bool):
)
)

# If this method is a resumable upload, return a PythonType instance
# representing the resumable upload session (while self.output remains
# the underlying protobuf response message for final deserialization).
if self.is_resumable_upload:
return PythonType(
meta=metadata.Metadata(
address=metadata.Address(
name=(
"AsyncResumableUploadSession"
if enable_asyncio
else "ResumableUploadSession"
),
module="resumable_transfer",
package=("google", "api_core"),
collisions=self.input.ident.collisions,
),
documentation=utils.doc(
"An object representing a resumable upload session."
),
),
)

# Return the usual output.
return self.output

Expand Down Expand Up @@ -1943,6 +1965,11 @@ def _ref_types(self, recursive: bool) -> Sequence[Union[MessageType, EnumType]]:
if self.paged_result_field and self.paged_result_field.message:
answer.append(self.paged_result_field.message)

# If this method is a resumable upload, client_output is ResumableUploadSession,
# so explicitly include self.output to ensure the underlying response message is imported.
if self.is_resumable_upload:
answer.append(self.output)

# Done; return the answer.
return tuple(answer)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,9 @@
requests: Optional[Iterator[{{ method.input.ident }}]] = None,
*,
{% endif %}
{% if method.is_resumable_upload %}
config: Optional[ResumableUploadConfig] = None,
{% endif %}
retry: OptionalRetry = gapic_v1.method.DEFAULT,
timeout: Union[float, object] = gapic_v1.method.DEFAULT,
{{ shared_macros.client_method_metadata_argument()|indent(8) }} = {{ shared_macros.client_method_metadata_default_value() }},
Expand Down Expand Up @@ -65,6 +68,10 @@
The request object iterator.{{ " " }}
{{- method.input.meta.doc|rst(width=72, indent=16, nl=False) }}
{% endif %}
{% if method.is_resumable_upload %}
config (Optional[google.api_core.resumable_transfer.ResumableUploadConfig]):
Optional configuration for the resumable upload session.
{% endif %}
retry (google.api_core.retry.Retry): Designation of what errors, if any,
should be retried.
timeout (float): The timeout for this request.
Expand Down Expand Up @@ -161,6 +168,9 @@
retry=retry,
timeout=timeout,
metadata=metadata,
{% if method.is_resumable_upload %}
config=config,
{% endif %}
)
{% if method.lro %}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -132,13 +132,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 #}
Expand Down Expand Up @@ -263,6 +263,68 @@ def _get_http_options():
{% endmacro %}


{# rest_resumable_upload_call_method_common includes the common code for a rest
resumable upload __call__ method to be re-used for sync and async REST
__call__ implementation.

Args:
method: The method.
service: The service.
is_async (bool): Used to determine the code path i.e. whether for sync or async call.
rest_numeric_enums (bool): Used to determine whether to encode enums as numbers. #}
{% macro rest_resumable_upload_call_method_common(method, service, is_async=False, rest_numeric_enums=False) %}
{% set service_name = service.name %}
{% set method_name = method.name %}
{% set body_spec = method.http_options[0].body %}
{% set await_prefix = "await " if is_async else "" %}
{% set client_output_ident = method.client_output_async.ident if is_async else method.client_output.ident %}
{% set retry_class = "retries.AsyncRetry" if is_async else "retries.Retry" %}
http_options = _Base{{ service_name }}RestTransport._Base{{ method_name }}._get_http_options()
request, metadata = {{ await_prefix }}self._interceptor.pre_{{ method_name|snake_case }}(request, metadata)
transcoded_request, body, query_params = transcode_request(
http_options,
request,
required_fields_default_values=getattr(
_Base{{ service_name }}RestTransport._Base{{ method_name }},
"_Base{{ method_name }}__REQUIRED_FIELDS_DEFAULT_VALUES",
None,
),
rest_numeric_enums={{ rest_numeric_enums }},
)

uri = transcoded_request["uri"]
params = rest_helpers.flatten_query_params(query_params, strict=True)
query_string = f"?{urllib.parse.urlencode(params)}" if params else ""
upload_url = f"{self._host}{uri}{query_string}"
headers: Dict[str, Any] = {**dict(metadata), **dict((config.headers or {}) if config else {})}
if config is None:
config = resumable_transfer.ResumableUploadConfig(headers=headers)
else:
config = dataclasses.replace(config, headers=headers)

session_kwargs: Dict[str, Any] = (
{"start_timeout": timeout}
if isinstance(timeout, (int, float))
else {}
)
session = {{ client_output_ident }}(
upload_url=upload_url,
config=config,
transport=self._session,
response_type={{ method.output.ident }},
start_retry=retry if isinstance(retry, {{ retry_class }}) else None,
**session_kwargs,
)
{% if body_spec %}
session.upload = functools.partial(session.upload, request_body=body or "") # type: ignore[method-assign]
{% if not is_async %}
session.iter_upload = functools.partial(session.iter_upload, request_body=body or "") # type: ignore[method-assign]
{% endif %}
{% endif %}
return session
{%- endmacro %}


{% macro unary_request_interceptor_common(service) %}
logging_enabled = CLIENT_LOGGING_SUPPORTED and _LOGGER.isEnabledFor(std_logging.DEBUG)
if logging_enabled: # pragma: NO COVER
Expand Down Expand Up @@ -299,11 +361,12 @@ def _get_http_options():
{%- endmacro %}


{% macro prep_wrapped_messages_async_method(api, service) %}
{% macro prep_wrapped_messages_async_method(api, service, is_rest_asyncio=False) %}
{% set rest_async_io_enabled = api.all_library_settings[api.naming.proto_package].python_settings.experimental_features.rest_async_io_enabled %}
def _prep_wrapped_messages(self, client_info):
""" Precompute the wrapped methods, overriding the base class method to use async wrappers."""
self._wrapped_methods = {
{% for method in service.methods.values() %}
{% for method in service.methods.values() if not is_rest_asyncio or rest_async_io_enabled or method.is_resumable_upload %}
self.{{ method.transport_safe_name|snake_case }}: self._wrap_method(
self.{{ method.transport_safe_name|snake_case }},
{% if method.retry %}
Expand All @@ -329,6 +392,7 @@ def _prep_wrapped_messages(self, client_info):
client_info=client_info,
),
{% endfor %}{# service.methods.values() #}
{% if not is_rest_asyncio or rest_async_io_enabled %}
{% for method_name in api.mixin_api_methods.keys() %}
{# TODO(https://github.com/googleapis/gapic-generator-python/issues/2197): Use `transport_safe_name` similar
# to what we do for non-mixin methods above.
Expand All @@ -339,6 +403,7 @@ def _prep_wrapped_messages(self, client_info):
client_info=client_info,
),
{% endfor %}{# method_name in api.mixin_api_methods.keys() #}
{% endif %}
}
{% endmacro %}

Expand All @@ -362,6 +427,7 @@ def _wrap_method(self, func, *args, **kwargs):
# synchronous and asynchronous rest transports
#}
{% macro create_interceptor_class(api, service, method, is_async=False) %}
{% set rest_async_io_enabled = api.all_library_settings[api.naming.proto_package].python_settings.experimental_features.rest_async_io_enabled %}
{% set async_prefix = "async " if is_async else "" %}
{% set async_method_name_prefix = "Async" if is_async else "" %}
{% set async_docstring = "Asynchronous " if is_async else "" %}
Expand All @@ -382,12 +448,12 @@ class {{ async_method_name_prefix }}{{ service.name }}RestInterceptor:

.. code-block:: python
class MyCustom{{ service.name }}Interceptor({{ service.name }}RestInterceptor):
{% for _, method in service.methods|dictsort if not method.client_streaming %}
{% for _, method in service.methods|dictsort if not method.client_streaming and (not is_async or rest_async_io_enabled or method.is_resumable_upload) %}
{{ async_prefix }}def pre_{{ method.name|snake_case }}(self, request, metadata):
logging.log(f"Received request: {request}")
return request, metadata

{% if not method.void %}
{% if not method.void and not method.is_resumable_upload %}
{{ async_prefix }}def post_{{ method.name|snake_case }}(self, response):
logging.log(f"Received response: {response}")
return response
Expand All @@ -400,7 +466,7 @@ class {{ async_method_name_prefix }}{{ service.name }}RestInterceptor:


"""
{% for method in service.methods.values()|sort(attribute="name") if not method.client_streaming and method.http_options %}
{% for method in service.methods.values()|sort(attribute="name") if not method.client_streaming and method.http_options and (not is_async or rest_async_io_enabled or method.is_resumable_upload) %}
{# TODO(https://github.com/googleapis/gapic-generator-python/issues/2147): Remove the condition below once async rest transport supports the guarded methods. #}
{{ async_prefix }}def pre_{{ method.name|snake_case }}(self, request: {{method.input.ident}}, {{ client_method_metadata_argument() }}) -> Tuple[{{method.input.ident}}, {{ client_method_metadata_type() }}]:
"""Pre-rpc interceptor for {{ method.name|snake_case }}
Expand All @@ -410,7 +476,7 @@ class {{ async_method_name_prefix }}{{ service.name }}RestInterceptor:
"""
return request, metadata

{% if not method.void %}
{% if not method.void and not method.is_resumable_upload %}
{% if not method.server_streaming %}
{{ async_prefix }}def post_{{ method.name|snake_case }}(self, response: {{method.output.ident}}) -> {{method.output.ident}}:
{% else %}
Expand Down Expand Up @@ -450,6 +516,7 @@ class {{ async_method_name_prefix }}{{ service.name }}RestInterceptor:
{% endif %}{# not method.void #}
{% endfor %}

{% if not is_async or rest_async_io_enabled %}
{% for name, signature in api.mixin_api_signatures.items() %}
{{ async_prefix }}def pre_{{ name|snake_case }}(
self, request: {{signature.request_type}}, {{ client_method_metadata_argument() }}
Expand All @@ -473,6 +540,7 @@ class {{ async_method_name_prefix }}{{ service.name }}RestInterceptor:
return response

{% endfor %}
{% endif %}
{% endmacro %}

{% macro generate_mixin_call_method(service, api, name, sig, is_async) %}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ from {{package_path}} import gapic_version as package_version
from google.api_core.client_options import ClientOptions
from google.api_core import exceptions as core_exceptions
from google.api_core import gapic_v1
{% if service.has_resumable_upload_methods %}
from google.api_core.resumable_transfer import ResumableUploadConfig
{% endif %}
{% if has_auto_populated_fields %}
from {{package_path}}._compat import setup_request_id
{% endif %}
Expand Down Expand Up @@ -289,6 +292,9 @@ class {{ service.async_client_name }}:
requests: Optional[AsyncIterator[{{ method.input.ident }}]] = None,
*,
{% endif %}
{% if method.is_resumable_upload %}
config: Optional[ResumableUploadConfig] = None,
{% endif %}
retry: OptionalRetry = gapic_v1.method.DEFAULT,
timeout: Union[float, object] = gapic_v1.method.DEFAULT,
{{ shared_macros.client_method_metadata_argument()|indent(8) }} = {{ shared_macros.client_method_metadata_default_value() }},
Expand Down Expand Up @@ -324,6 +330,10 @@ class {{ service.async_client_name }}:
The request object AsyncIterator.{{ " " }}
{{- method.input.meta.doc|rst(width=72, indent=16, nl=False) }}
{% endif %}
{% if method.is_resumable_upload %}
config (Optional[google.api_core.resumable_transfer.ResumableUploadConfig]):
Optional configuration for the resumable upload session.
{% endif %}
retry (google.api_core.retry_async.AsyncRetry): Designation of what errors, if any,
should be retried.
timeout (float): The timeout for this request.
Expand Down Expand Up @@ -414,6 +424,9 @@ class {{ service.async_client_name }}:
retry=retry,
timeout=timeout,
metadata=metadata,
{% if method.is_resumable_upload %}
config=config,
{% endif %}
)
{% if method.lro %}

Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
{% extends '_base.py.j2' %}
{% block content %}
{# TODO(https://github.com/googleapis/gapic-generator-python/issues/2121): Remove the following variable (and the condition later in this file) for async rest transport once support for it is GA. #}
{% set rest_async_io_enabled = api.all_library_settings[api.naming.proto_package].python_settings.experimental_features.rest_async_io_enabled %}
{% set rest_async_io_enabled = api.all_library_settings[api.naming.proto_package].python_settings.experimental_features.rest_async_io_enabled or service.has_resumable_upload_methods %}
{% set has_auto_populated_fields = api.all_method_settings.values()|map(attribute="auto_populated_fields", default=[])|select|list %}
{% import "%namespace/%name_%version/%sub/services/%service/_client_macros.j2" as macros %}
{% import "%namespace/%name_%version/%sub/services/%service/_shared_macros.j2" as shared_macros %}
Expand Down Expand Up @@ -30,6 +30,9 @@ from google.api_core import exceptions as core_exceptions
from google.api_core import extended_operation
{% endif %}
from google.api_core import gapic_v1
{% if service.has_resumable_upload_methods %}
from google.api_core.resumable_transfer import ResumableUploadConfig
{% endif %}
from {{package_path}}._compat import get_universe_domain, get_api_endpoint, get_default_mtls_endpoint, should_use_client_cert, read_environment_variables
{% if has_auto_populated_fields %}
from {{package_path}}._compat import setup_request_id
Expand Down Expand Up @@ -82,14 +85,12 @@ from .transports.grpc_asyncio import {{ service.grpc_asyncio_transport_name }}
from .transports.rest import {{ service.name }}RestTransport
{# TODO(https://github.com/googleapis/gapic-generator-python/issues/2121): Remove this condition when async rest is GA. #}
{% if rest_async_io_enabled %}
ASYNC_REST_EXCEPTION = None
try:
from .transports.rest_asyncio import Async{{ service.name }}RestTransport
HAS_ASYNC_REST_DEPENDENCIES = True
HAS_ASYNC_REST_DEPENDENCIES = True # pragma: NO COVER
{# NOTE: `pragma: NO COVER` is needed since the coverage for presubmits isn't combined. #}
except ImportError as e: # pragma: NO COVER
except ImportError: # pragma: NO COVER
HAS_ASYNC_REST_DEPENDENCIES = False
ASYNC_REST_EXCEPTION = e

{% endif %}{# if rest_async_io_enabled #}
{% endif %}
Expand Down Expand Up @@ -133,7 +134,9 @@ class {{ service.client_name }}Meta(type):
{% if rest_async_io_enabled %}
{# NOTE: `pragma: NO COVER` is needed since the coverage for presubmits isn't combined. #}
if label == "rest_asyncio" and not HAS_ASYNC_REST_DEPENDENCIES: # pragma: NO COVER
raise ASYNC_REST_EXCEPTION
raise ImportError(
"`rest_asyncio` transport requires the library to be installed with the `async_rest` extra. Install the library with the `async_rest` extra using `pip install {{ api.naming.warehouse_package_name }}[async_rest]`"
)
{% endif %}
if label:
return cls._transport_registry[label]
Expand Down Expand Up @@ -494,7 +497,7 @@ class {{ service.client_name }}(metaclass={{ service.client_name }}Meta):
else cast(Callable[..., {{ service.name }}Transport], transport)
)

if "rest_asyncio" in str(transport_init):
if "rest_asyncio" in str(transport_init): # pragma: NO COVER
{# TODO(https://github.com/googleapis/gapic-generator-python/issues/2136): Support the following parameters in async rest: #}
unsupported_params = {
"google.api_core.client_options.ClientOptions.credentials_file": self._client_options.credentials_file,
Expand Down
Loading
Loading