FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

feat(gapic-generator): generate code samples for resumable upload RPCs by parthea · Pull Request #18501 · googleapis/google-cloud-python · GitHub

Repository navigation

Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension .j2  (2) .json  (1) .py  (6) All 3 file types selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
17 changes: 12 additions & 5 deletions packages/gapic-generator/gapic/samplegen/samplegen.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -1183,14 +1183,20 @@ def _get_sample_imports(sample: Dict, rpc: wrappers.Method) -> List[str]:
module_name = sample["module_name"]
module_import = f"from {module_namespace} import {module_name}"

imports = [module_import]
address = rpc.input.meta.address
# This checks if the request message is part of the service proto package.
# If not, we should try to include a separate import statement.
if address.proto_package.startswith(address.api_naming.proto_package):
return [module_import]
else:
request_import = str(address.python_import)
return sorted([module_import, request_import])
if not address.proto_package.startswith(address.api_naming.proto_package):
imports.append(str(address.python_import))

if rpc.is_resumable_upload:
imports.append(
"from google.api_core.resumable_transfer import ResumableUploadConfig"
)
imports.append("import io")

return sorted(imports)


def generate_sample(
Expand Down Expand Up @@ -1222,6 +1228,7 @@ def generate_sample(

calling_form = types.CallingForm.method_default(rpc)
sample["is_internal"] = rpc.is_internal
sample["is_resumable_upload"] = rpc.is_resumable_upload

v = Validator(rpc, api_schema)
# Tweak some small aspects of the sample to set defaults for optional
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,11 @@
# It may require modifications to work in your environment.

# To install the latest published package dependency, execute the following:
{% if sample.get("is_resumable_upload") and sample.transport in ("grpc-async", "rest-async") %}
# python3 -m pip install {{ sample.package_name }}[async_rest]
{% else %}
# python3 -m pip install {{ sample.package_name }}
{% endif %}
{% endmacro %}

{% macro print_string_formatting(string_list) %}
Expand Down Expand Up @@ -226,7 +230,7 @@ request=request
{% endmacro %}


{% macro render_method_call(sample, calling_form, calling_form_enum, transport) %}
{% macro render_method_call(sample, calling_form, calling_form_enum, transport, include_config=False) %}
{# Note: this doesn't deal with enums or unions #}
{# LROs return operation objects and paged requests return pager objects #}
{% if transport == "grpc-async" and calling_form not in
Expand All @@ -238,14 +242,71 @@ await{{ " "}}
client.{{ render_method_name(sample)|trim }}(requests=request_generator())
{% else %}{# TODO: deal with flattening #}
{# TODO: set up client streaming once some questions are answered #}
client.{{ render_method_name(sample)|trim }}({{ render_request_params_unary(sample.request)|trim }})
client.{{ render_method_name(sample)|trim }}({{ render_request_params_unary(sample.request)|trim }}{% if include_config %}, config=config{% endif %})
{% endif %}
{% endmacro %}


{% macro render_resumable_upload_samples(sample, calling_form, calling_form_enum) %}
{% set is_async = sample.transport in ("grpc-async", "rest-async") %}
{% set async_prefix = "async " if is_async else "" %}
{% set await_prefix = "await " if is_async else "" %}
{% set fn_name = sample.rpc|snake_case|trim %}
{% set fn_params = print_input_params(sample.request)|trim %}
{% set method_call_with_config = render_method_call(sample, calling_form, calling_form_enum, sample.transport, include_config=True)|trim %}
{{ async_prefix }}def sample_{{ fn_name }}({{ fn_params }}):
{{ render_client_setup(sample.module_name, sample.client_name)|indent }}
{{ render_request_setup(sample.request, sample.request_module_name, sample.request_type, calling_form, calling_form_enum)|indent }}
# Configure optional transfer settings such as chunk size and stall detection
config = ResumableUploadConfig(
chunk_size=8 * 1024 * 1024, # 8 MB
stall_minimum_rate=64 * 1024,
stall_timeout=120,
)

# Make the request
upload_session = {{ method_call_with_config }}

# Option 1: Upload the entire stream directly and {% if is_async %}await{% else %}return{% endif %} the final response
stream = io.BytesIO(b"Example upload data")
{{ "response = " if sample.response else "" }}{{ await_prefix }}upload_session.upload(stream)

# Option 2: Alternatively, iterate over the upload to receive progress updates per chunk
{% if is_async %}
# async for progress in upload_session.upload(stream):
{% else %}
# for progress in upload_session.iter_upload(stream):
{% endif %}
Comment thread
parthea marked this conversation as resolved.
# print(f"Uploaded {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# print(f"Session URL: {progress.upload_url}")
{% if sample.response %}
# response = upload_session.response
{% endif %}

# Option 3: Alternatively, resume an interrupted upload from a saved session URL
# {{ "response = " if sample.response else "" }}{{ await_prefix }}upload_session.resume(upload_url, stream, chunk_size=config.chunk_size)

# Option 4: Alternatively, resume an interrupted upload while receiving progress updates
{% if is_async %}
# async for progress in upload_session.resume(upload_url, stream, chunk_size=config.chunk_size):
{% else %}
# for progress in upload_session.iter_resume(upload_url, stream, chunk_size=config.chunk_size):
{% endif %}
Comment thread
parthea marked this conversation as resolved.
# print(f"Resumed {progress.bytes_uploaded} bytes | State: {progress.state.name}")
{% if sample.response %}
# response = upload_session.response

# Handle the response
{% for statement in sample.response %}
{{ dispatch_statement(statement)|trim }}
{% endfor %}
{% endif %}
{% endmacro %}


{# Setting up the method invocation is the responsibility of the caller: #}
{# it's just easier to set up client side streaming and other things from outside this macro. #}
{% macro render_calling_form(method_invocation_text, calling_form, calling_form_enum, transport, response_statements ) %}
{% macro render_calling_form(method_invocation_text, calling_form, calling_form_enum, transport, response_statements, sample=None ) %}
# Make the request
{% if calling_form in [calling_form_enum.Request, calling_form_enum.RequestStreamingClient] %}
{% if response_statements %}response = {% endif %}{{ method_invocation_text|trim }}
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -32,12 +32,16 @@


{# also need calling form #}
{% if sample.get("is_resumable_upload") %}
{{ frags.render_resumable_upload_samples(sample, calling_form, calling_form_enum) }}
Comment thread
parthea marked this conversation as resolved.
{% else %}
{% if sample.transport == "grpc-async" %}async {% endif %}def sample_{{ sample.rpc|snake_case|trim }}({{ frags.print_input_params(sample.request)|trim }}):
{{ frags.render_client_setup(sample.module_name, sample.client_name)|indent }}
{{ frags.render_request_setup(sample.request, sample.request_module_name, sample.request_type, calling_form, calling_form_enum)|indent }}
{% with method_call = frags.render_method_call(sample, calling_form, calling_form_enum, sample.transport) %}
{{ frags.render_calling_form(method_call, calling_form, calling_form_enum, sample.transport, sample.response)|indent -}}
{{ frags.render_calling_form(method_call, calling_form, calling_form_enum, sample.transport, sample.response, sample)|indent -}}
{% endwith %}
{% endif %}

# [END {{ sample.id }}]
{# TODO: Enable main block (or decide to remove main block from python sample) #}
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,8 @@ async def upload_media(self,
# client as shown in:
# https://googleapis.dev/python/google-api-core/latest/client_options.html
from google import showcase_v1beta1
from google.api_core.resumable_transfer import ResumableUploadConfig
import io

async def sample_upload_media():
# Create a client
Expand All @@ -289,8 +291,33 @@ async def sample_upload_media():
request = showcase_v1beta1.UploadMediaRequest(
)

# Configure optional transfer settings such as chunk size and stall detection
config = ResumableUploadConfig(
chunk_size=8 * 1024 * 1024, # 8 MB
stall_minimum_rate=64 * 1024,
stall_timeout=120,
)

# Make the request
response = await client.upload_media(request=request)
upload_session = await client.upload_media(request=request, config=config)

# Option 1: Upload the entire stream directly and await the final response
stream = io.BytesIO(b"Example upload data")
response = await upload_session.upload(stream)

# Option 2: Alternatively, iterate over the upload to receive progress updates per chunk
# async for progress in upload_session.upload(stream):
# print(f"Uploaded {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# print(f"Session URL: {progress.upload_url}")
# response = upload_session.response

# Option 3: Alternatively, resume an interrupted upload from a saved session URL
# response = await upload_session.resume(upload_url, stream, chunk_size=config.chunk_size)

# Option 4: Alternatively, resume an interrupted upload while receiving progress updates
# async for progress in upload_session.resume(upload_url, stream, chunk_size=config.chunk_size):
# print(f"Resumed {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# response = upload_session.response

Comment thread
daniel-sanche marked this conversation as resolved.
# Handle the response
print(response)
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -514,6 +514,8 @@ def upload_media(self,
# client as shown in:
# https://googleapis.dev/python/google-api-core/latest/client_options.html
from google import showcase_v1beta1
from google.api_core.resumable_transfer import ResumableUploadConfig
import io

def sample_upload_media():
# Create a client
Expand All @@ -523,8 +525,33 @@ def sample_upload_media():
request = showcase_v1beta1.UploadMediaRequest(
)

# Configure optional transfer settings such as chunk size and stall detection
config = ResumableUploadConfig(
chunk_size=8 * 1024 * 1024, # 8 MB
stall_minimum_rate=64 * 1024,
stall_timeout=120,
)

# Make the request
response = client.upload_media(request=request)
upload_session = client.upload_media(request=request, config=config)

# Option 1: Upload the entire stream directly and return the final response
stream = io.BytesIO(b"Example upload data")
response = upload_session.upload(stream)

# Option 2: Alternatively, iterate over the upload to receive progress updates per chunk
# for progress in upload_session.iter_upload(stream):
# print(f"Uploaded {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# print(f"Session URL: {progress.upload_url}")
# response = upload_session.response

# Option 3: Alternatively, resume an interrupted upload from a saved session URL
# response = upload_session.resume(upload_url, stream, chunk_size=config.chunk_size)

# Option 4: Alternatively, resume an interrupted upload while receiving progress updates
# for progress in upload_session.iter_resume(upload_url, stream, chunk_size=config.chunk_size):
# print(f"Resumed {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# response = upload_session.response

# Handle the response
print(response)
Expand Down
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
# It may require modifications to work in your environment.

# To install the latest published package dependency, execute the following:
# python3 -m pip install google-showcase
# python3 -m pip install google-showcase[async_rest]


# [START localhost_v1beta1_generated_ResumableUploadService_UploadMedia_async]
Comment thread
daniel-sanche marked this conversation as resolved.
Expand All @@ -32,6 +32,8 @@
# client as shown in:
# https://googleapis.dev/python/google-api-core/latest/client_options.html
from google import showcase_v1beta1
from google.api_core.resumable_transfer import ResumableUploadConfig
import io


async def sample_upload_media():
Expand All @@ -42,10 +44,36 @@ async def sample_upload_media():
request = showcase_v1beta1.UploadMediaRequest(
)

# Configure optional transfer settings such as chunk size and stall detection
config = ResumableUploadConfig(
chunk_size=8 * 1024 * 1024, # 8 MB
stall_minimum_rate=64 * 1024,
stall_timeout=120,
)

# Make the request
response = await client.upload_media(request=request)
upload_session = await client.upload_media(request=request, config=config)

# Option 1: Upload the entire stream directly and await the final response
stream = io.BytesIO(b"Example upload data")
response = await upload_session.upload(stream)

# Option 2: Alternatively, iterate over the upload to receive progress updates per chunk
# async for progress in upload_session.upload(stream):
# print(f"Uploaded {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# print(f"Session URL: {progress.upload_url}")
# response = upload_session.response

# Option 3: Alternatively, resume an interrupted upload from a saved session URL
# response = await upload_session.resume(upload_url, stream, chunk_size=config.chunk_size)

# Option 4: Alternatively, resume an interrupted upload while receiving progress updates
# async for progress in upload_session.resume(upload_url, stream, chunk_size=config.chunk_size):
# print(f"Resumed {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# response = upload_session.response

# Handle the response
print(response)


# [END localhost_v1beta1_generated_ResumableUploadService_UploadMedia_async]
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@
# client as shown in:
# https://googleapis.dev/python/google-api-core/latest/client_options.html
from google import showcase_v1beta1
from google.api_core.resumable_transfer import ResumableUploadConfig
import io


def sample_upload_media():
Expand All @@ -42,10 +44,36 @@ def sample_upload_media():
request = showcase_v1beta1.UploadMediaRequest(
)

# Configure optional transfer settings such as chunk size and stall detection
config = ResumableUploadConfig(
chunk_size=8 * 1024 * 1024, # 8 MB
stall_minimum_rate=64 * 1024,
stall_timeout=120,
)

# Make the request
response = client.upload_media(request=request)
upload_session = client.upload_media(request=request, config=config)

# Option 1: Upload the entire stream directly and return the final response
stream = io.BytesIO(b"Example upload data")
response = upload_session.upload(stream)

# Option 2: Alternatively, iterate over the upload to receive progress updates per chunk
# for progress in upload_session.iter_upload(stream):
# print(f"Uploaded {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# print(f"Session URL: {progress.upload_url}")
# response = upload_session.response

# Option 3: Alternatively, resume an interrupted upload from a saved session URL
# response = upload_session.resume(upload_url, stream, chunk_size=config.chunk_size)

# Option 4: Alternatively, resume an interrupted upload while receiving progress updates
# for progress in upload_session.iter_resume(upload_url, stream, chunk_size=config.chunk_size):
# print(f"Resumed {progress.bytes_uploaded} bytes | State: {progress.state.name}")
# response = upload_session.response

# Handle the response
print(response)


# [END localhost_v1beta1_generated_ResumableUploadService_UploadMedia_sync]
Loading
Loading

Back | FazBrowse Home | New Git URL