Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
45 commits
Select commit Hold shift + click to select a range
43b0747
feat: add support for resumable uploads
parthea Sep 29, 2026
1977ce0
address review feedback
parthea Sep 29, 2026
db71d9c
update goldens
parthea Sep 29, 2026
a9aefa4
address review feedback
parthea Sep 30, 2026
086f99a
address review feedback
parthea Sep 30, 2026
45f18d0
remove warnings related to rest numberic enums not set
parthea Sep 30, 2026
ee12261
update test
parthea Sep 30, 2026
eac800b
address review feedback
parthea Sep 30, 2026
c0c6465
fix build
parthea Sep 30, 2026
1c0ca3d
fix build
parthea Sep 30, 2026
aa4df26
address review feedback
parthea Sep 30, 2026
8b13205
address review feedback
parthea Sep 30, 2026
e786d2d
update goldens
parthea Sep 30, 2026
fc1aa34
update goldens
parthea Sep 30, 2026
e59e867
address review feedback
parthea Oct 1, 2026
03d618b
address review feedback
parthea Oct 1, 2026
b22f24f
address review feedback
parthea Oct 1, 2026
00c2ce2
address review feedback
parthea Oct 1, 2026
28161a8
address review feedback
parthea Oct 1, 2026
94511d9
address review feedback
parthea Oct 1, 2026
d51b71c
add test for bug where request body is not passed
parthea Oct 1, 2026
485a92b
fix: forward transcoded request body and query params in resumable up…
parthea Oct 1, 2026
bb533ee
address review feedback
parthea Oct 1, 2026
776811e
address review feedback
parthea Oct 1, 2026
b48cfeb
address review feedback
parthea Oct 1, 2026
69cc533
fix build
parthea Oct 1, 2026
4e83566
update unit test
parthea Oct 1, 2026
dafc8de
fix build
parthea Oct 1, 2026
50b0c18
address review feedback
parthea Oct 2, 2026
809a0a7
sync minimum version of google-api-core for resumable uploads
parthea Oct 2, 2026
59fa610
temporary commit to be reverted
parthea Oct 3, 2026
86e88a3
clean up
parthea Oct 3, 2026
c20ac40
test: add unit tests for bug where start_retry is not passed from cli…
parthea Oct 3, 2026
d52a6bf
address review feedback
parthea Oct 3, 2026
9af1477
address review feedback
parthea Oct 3, 2026
723015e
address review feedback
parthea Oct 3, 2026
adab365
fix: set Content-Type on resumable upload start requests
parthea Oct 5, 2026
7e8fb40
temporary change to test deps at head
parthea Oct 5, 2026
aa3041d
Revert "temporary change to test deps at head"
parthea Oct 5, 2026
dd724d6
Revert "temporary commit to be reverted"
parthea Oct 5, 2026
944f07b
address review feedback
parthea Oct 5, 2026
98ed451
address review feedback
parthea Oct 5, 2026
bb8e10c
address review feedback
parthea Oct 5, 2026
87b4c0e
address review feedback
parthea Oct 5, 2026
78ebf32
bump deps
parthea Oct 6, 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
27 changes: 24 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 @@ -358,8 +360,16 @@ 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 (
"grpc" in opts.transport
and service.has_resumable_upload_methods
)
)
or (
"rest_base" in template_name
and "rest" not in opts.transport
and not service.has_resumable_upload_methods
)
or ("rest_base" in template_name and "rest" not in opts.transport)
):
continue

Expand All @@ -386,9 +396,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,10 @@
retry=retry,
timeout=timeout,
metadata=metadata,
{% if method.is_resumable_upload %}
config=config,
start_retry=retry,
{% 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 }}',
Comment thread
daniel-sanche marked this conversation as resolved.
{% 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 {})}
Comment thread
daniel-sanche marked this conversation as resolved.
Comment thread
daniel-sanche marked this conversation as resolved.
headers["Content-Type"] = "application/json"
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 {}
)
# ``start_retry`` is used instead of ``retry`` because ``_GapicCallable``
# consumes the ``retry`` argument before invoking the transport callable
# and only forwards extra keyword arguments such as ``start_retry``.
return {{ client_output_ident }}(
upload_url=upload_url,
config=config,
transport=self._session,
response_type={{ method.output.ident }},
start_retry=start_retry if isinstance(start_retry, {{ retry_class }}) else None,
{% if body_spec %}
request_body=body,
{% endif %}
**session_kwargs,
)
{%- 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
Comment thread
daniel-sanche marked this conversation as resolved.
{% 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,10 @@ class {{ service.async_client_name }}:
retry=retry,
timeout=timeout,
metadata=metadata,
{% if method.is_resumable_upload %}
config=config,
start_retry=retry,
{% endif %}
)
{% if method.lro %}

Expand Down
Original file line number Diff line number Diff line change
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
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,9 @@ from .rest import {{ service.name }}RestInterceptor
ASYNC_REST_CLASSES: Tuple[str, ...]
try:
from .rest_asyncio import Async{{ service.name }}RestTransport
from .rest_asyncio import Async{{ service.name }}RestInterceptor
ASYNC_REST_CLASSES = ('Async{{ service.name }}RestTransport', 'Async{{ service.name }}RestInterceptor')
HAS_REST_ASYNC = True
from .rest_asyncio import Async{{ service.name }}RestInterceptor # pragma: NO COVER
ASYNC_REST_CLASSES = ('Async{{ service.name }}RestTransport', 'Async{{ service.name }}RestInterceptor') # pragma: NO COVER
HAS_REST_ASYNC = True # pragma: NO COVER
except ImportError: # pragma: NO COVER
ASYNC_REST_CLASSES = ()
HAS_REST_ASYNC = False
Expand Down
Loading
Loading