Conversation
There was a problem hiding this comment.
Code Review
This pull request adds support for resumable upload methods in the GAPIC generator, wrapping responses in ResumableUploadSession or AsyncResumableUploadSession and ensuring these methods route over REST even when gRPC transport is initialized. It updates templates, transports, and unit tests to support this protocol. The reviewer identified a critical issue in grpc_asyncio.py.j2 where AsyncRestTransport is instantiated on every call, which is highly inefficient and causes resource leaks; they suggested caching the transport instance instead.
daniel-sanche
left a comment
There was a problem hiding this comment.
It seems like there could be some runtime issues here. Can you make sure we have full coverage?
b7f0140 to
5282533
Compare
No region tags are edited in this PR.This comment is generated by snippet-bot.
|
ce31648 to
2c60a8c
Compare
a33d1fe to
389dbfb
Compare
1ddea43 to
8e785c8
Compare
|
|
||
| # Save the credentials. | ||
| {% if service.has_resumable_upload_methods %} | ||
| self._credentials: Any = credentials |
There was a problem hiding this comment.
nit: why are we declaring a type hint for a var and setting it to Any? is mypy complaining?
There was a problem hiding this comment.
Yes, mypy complains when self._credentials (typed as Optional[google.auth.credentials.Credentials] ) is passed to
Async...RestTransport in grpc_asyncio.py (which expects Optional[google.auth.aio.credentials.Credentials]), due to the existing async
credential type-hint mismatch tracked in #16268:
I pushed a commit to show the mypy failure in b4b8c20. See build log https://github.com/googleapis/google-cloud-python/actions/runs/36883974364/job/110442747465?pr=18495. The specific failure is
nox > mypy -p google --check-untyped-defs
google/showcase_v1beta1/services/resumable_upload_service/transports/grpc_asyncio.py:376: error: Argument "credentials" to "AsyncResumableUploadServiceRestTransport" has incompatible type "google.auth.credentials.Credentials | None"; expected "google.auth.aio.credentials.Credentials | None" [arg-type]
Found 1 error in 1 file (checked 74 source files)
I found an existing workarounds here
As far as I can tell, this is not a trivial fix. I'll add a comment with a link to the TODO comment to follow up.
Gemini highlighted a clash in GrpcAsyncIOTransport / AsyncRestTransport where grpc_helpers_async.create_channel takes sync credentials while AsyncRestTransport takes async credentials.
| from .grpc import {{ service.name }}GrpcTransport | ||
| {% if service.has_resumable_upload_methods %} | ||
| try: | ||
| from .rest_asyncio import Async{{ service.name }}RestTransport |
There was a problem hiding this comment.
What if it has resumable upload methods but does not have AsyncRestTransport? or have you eliminated this case with your changes to the generator.py?
There was a problem hiding this comment.
I pushed a fix and updated test_get_response_resumable_upload_generates_async_client_and_rest_asyncio to ensure that at code generation time, generator.py will produce rest_asyncio.py (Async{{ service.name }}RestTransport), rest.py, and rest_base.py whenever service.has_resumable_upload_methods is True (even when rest_async_io_enabled is False or transport=grpc is passed).
At runtime, however, importing .rest_asyncio still raises ImportError if the library is installed without the optional [async_rest] (same as existing behavior for async rest SDKs)
See 43d0e42
| self._operations_client: Optional[operations_v1.OperationsClient] = None | ||
| {% endif %} | ||
| {% if service.has_resumable_upload_methods %} | ||
| self._rest_transport: Optional[{{ service.name }}RestTransport] = None |
There was a problem hiding this comment.
it doesn't look clean that we are instantiating a rest transport within grpc code.
In another thread, you've mentioned that there is a strict requirement for grpc to have scotty methods (and I assume we need rest for that).
-
Maybe we should document that requirement.
-
Maybe there's a cleaner way to support this requirement. I haven't look into this in too much depth yet.
There was a problem hiding this comment.
In another thread, you've mentioned that there is a strict requirement for grpc to have scotty methods (and I assume we need rest for that).
This is correct. We have a hard requirement as mentioned in #18495 (comment)
Maybe we should document that requirement.
I added a comment in grpc.py.j2 and grpc_asyncio.py.j2 documenting why REST is needed on gRPC transports in ce47c3c.
Maybe there's a cleaner way to support this requirement. I haven't look into this in too much depth yet.
Why delegation lives in the transport layer rather than the client layer:
{{ service.name }}Transport._prep_wrapped_messages()wrapsself.{{ method }}for every method in service.methods during__init__. Implementing the method property onGrpcTransportandGrpcAsyncIOTransportsatisfies the base Transport interface.- Client and AsyncClient remain transport-agnostic, dispatching all RPCs uniformly via
self._transport._wrapped_methods[self._transport.<method>]and closing resources viaself.transport.close().
| self._stubs['{{ method.transport_safe_name|snake_case }}'] = _ErrorStub() | ||
| else: | ||
| if self._rest_transport is None: | ||
| self._rest_transport = {{ service.name }}RestTransport( |
There was a problem hiding this comment.
what if rest.py file is missing? or can that not be the case?
There was a problem hiding this comment.
Addressed in 43d0e42. rest.py won't be missing. See the test test_get_response_resumable_upload_generates_async_client_and_rest_asyncio
| {% endif %} | ||
| """ | ||
| {% if method.is_resumable_upload %} | ||
| http_options = _Base{{ service.name }}RestTransport._Base{{ method.name }}._get_http_options() |
There was a problem hiding this comment.
Can we use a macro for this base transport for sync and async code?
There was a problem hiding this comment.
Good point. Addressed in 4588ebd. There is no change to the generated output
39c26aa to
f8fa453
Compare
daniel-sanche
left a comment
There was a problem hiding this comment.
I wanted to bring all my remaining comments into one place, although I know some of these are already being addressed
There was a problem hiding this comment.
The temporary change is in a07a705, which will be reverted once google-api-core and google-auth are released
| config=config, | ||
| transport=self._session, | ||
| response_type={{ method.output.ident }}, | ||
| start_retry=retry if isinstance(retry, {{ retry_class }}) else None, |
There was a problem hiding this comment.
Gemini pointed out that this will always be None, because the retry argument expected here is swallowed by _GapicCallable, so it never reaches here. It will always use the default retry
If we pass start_retry=retry in the client's rpc(...) call and accept start_retry on the transport method here, it will be passed through cleanly.
| {# TODO(https://github.com/googleapis/gapic-generator-python/issues/2121): Remove this condition when async rest is GA. #} | ||
| {% if rest_async_io_enabled %} | ||
| if HAS_ASYNC_REST_DEPENDENCIES: # pragma: NO COVER | ||
| _transport_registry["rest_asyncio"] = Async{{ service.name }}RestTransport |
There was a problem hiding this comment.
nit: This will add rest_asyncio to the transport registry, even if it is only used for resumable uploads. Most methods would raise NotImplementedErrors
Can we gate this, so it only shows up for libraries that fully enable rest_asyncio?
| transport._rest_transport = {{ service.name }}RestTransport( | ||
| host=transport._host, | ||
| credentials=transport._credentials, | ||
| client_info=transport._client_info, |
There was a problem hiding this comment.
I think we'd need to pass down client_cert_source_for_mtls for mTLS support
(also for the async side, but it's not currently supported there. Maybe we should expose that limitation with an exception)
Are there other arguments we need to pass through?
| 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 {})} |
There was a problem hiding this comment.
IIUC, should `"Content-Type": "application/json" be part of these headers as well?
|
|
||
|
|
||
| @dataclasses.dataclass | ||
| class _AsyncResumableUploadSessionAdapter: |
| start_retry=retry if isinstance(retry, retries.AsyncRetry) else None, | ||
| **session_kwargs, | ||
| ) | ||
| session.upload = functools.partial(session.upload, request_body=body or "") # type: ignore[method-assign] |
There was a problem hiding this comment.
Note: we can remove this after integrating #18543
| self._stubs['{{ method.transport_safe_name|snake_case }}'] = _ErrorStub() | ||
| elif HAS_ASYNC_REST: | ||
| # TODO(https://github.com/googleapis/google-cloud-python/issues/16268): | ||
| # Resolve the credential type mismatch between {{ service.grpc_asyncio_transport_name }} |
There was a problem hiding this comment.
note: the credential mismatch should go away with #18542
d6b2985 to
7f0f479
Compare
aee4eef to
e20d380
Compare
e20d380 to
a07a705
Compare
Add client, transport, and unit test template support for resumable upload RPCs.
Resumable upload methods accept an optional
ResumableUploadConfigon the client and returnResumableUploadSession/AsyncResumableUploadSessionfromgoogle.api_core.resumable_transfer. Since resumable uploads use HTTP, the gRPC and gRPC-async transports delegate resumable upload calls to the corresponding REST transport.Removed
compliance.protofrom generated showcase goldens. See follow up issue Generated tests for showcase fail withcompliance.proto#16312