Skip to content

Commit be137fe

Browse files
committed
add async system tests for resumable uploads
1 parent 3477800 commit be137fe

3 files changed

Lines changed: 212 additions & 1 deletion

File tree

‎packages/gapic-generator/tests/system/conftest.py‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@
4444
except ImportError:
4545
HAS_RESUMABLE_UPLOAD_CLIENT = False
4646

47+
HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = False
4748
if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true":
4849
from grpc.experimental import aio
4950
import asyncio
@@ -67,6 +68,16 @@
6768
HAS_ASYNC_REST_IDENTITY_TRANSPORT = True
6869
except:
6970
HAS_ASYNC_REST_IDENTITY_TRANSPORT = False
71+
try:
72+
from google.showcase import ResumableUploadServiceAsyncClient
73+
from google.showcase_v1beta1.services.resumable_upload_service.transports import (
74+
AsyncResumableUploadServiceRestTransport,
75+
AsyncResumableUploadServiceRestInterceptor,
76+
)
77+
78+
HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = True
79+
except:
80+
HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT = False
7081

7182
_GRPC_VERSION = grpc.__version__
7283

@@ -344,6 +355,19 @@ async def post_expand_with_metadata(self, request, metadata):
344355
return request, metadata
345356

346357

358+
if HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT:
359+
360+
class ResumableUploadMetadataClientRestAsyncInterceptor(
361+
AsyncResumableUploadServiceRestInterceptor
362+
):
363+
request_metadata: Sequence[Tuple[str, str]] = []
364+
response_metadata: Sequence[Tuple[str, str]] = []
365+
366+
async def pre_upload_media(self, request, metadata):
367+
self.request_metadata = metadata
368+
return request, metadata
369+
370+
347371
class EchoMetadataClientGrpcInterceptor(
348372
grpc.UnaryUnaryClientInterceptor,
349373
grpc.UnaryStreamClientInterceptor,
@@ -576,6 +600,27 @@ def intercepted_resumable_upload_rest(use_mtls, use_tls):
576600
return ResumableUploadServiceClient(transport=transport), interceptor
577601

578602

603+
@pytest.fixture
604+
def intercepted_resumable_upload_rest_async():
605+
if not HAS_ASYNC_REST_RESUMABLE_UPLOAD_TRANSPORT:
606+
pytest.skip("Skipping test with async rest.")
607+
608+
transport_name = "rest_asyncio"
609+
transport_cls = ResumableUploadServiceAsyncClient.get_transport_class(
610+
transport_name
611+
)
612+
interceptor = ResumableUploadMetadataClientRestAsyncInterceptor()
613+
614+
transport = transport_cls(
615+
credentials=async_anonymous_credentials(),
616+
host="localhost:7469",
617+
url_scheme="http",
618+
interceptor=interceptor,
619+
)
620+
621+
return ResumableUploadServiceAsyncClient(transport=transport), interceptor
622+
623+
579624
try:
580625
from google.api_core.resumable_transfer import (
581626
ResumableUploadConfig,

‎packages/gapic-generator/tests/system/test_resumable_upload_progress.py‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@
1313
# limitations under the License.
1414

1515
import io
16+
import os
17+
import pytest
1618

1719
from google.api_core.resumable_transfer import (
1820
DEFAULT_CHUNK_SIZE,
@@ -274,3 +276,59 @@ def test_client_upload_media_passes_request_body(intercepted_resumable_upload_re
274276
assert response.name == "client_upload_media.txt"
275277
assert response.size == len(payload)
276278

279+
280+
if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true":
281+
282+
@pytest.mark.asyncio
283+
async def test_async_client_upload_media_end_to_end(
284+
intercepted_resumable_upload_rest_async,
285+
):
286+
"""Verify end-to-end async resumable upload via ResumableUploadServiceAsyncClient."""
287+
client, _ = intercepted_resumable_upload_rest_async
288+
payload = b"0123456789" * 100
289+
stream = io.BytesIO(payload)
290+
291+
session = await client.upload_media(
292+
request=UploadMediaRequest(name="async_client_upload_media.txt"),
293+
config=ResumableUploadConfig(chunk_size=256),
294+
)
295+
response = await session.upload(stream=stream)
296+
297+
assert isinstance(response, UploadMediaResponse)
298+
assert response.name == "async_client_upload_media.txt"
299+
assert response.size == len(payload)
300+
301+
@pytest.mark.asyncio
302+
async def test_async_client_upload_media_progress_tracking(
303+
intercepted_resumable_upload_rest_async,
304+
):
305+
"""Verify async progress iteration over ResumableUploadServiceAsyncClient.upload_media."""
306+
client, _ = intercepted_resumable_upload_rest_async
307+
chunk_size = 262_144 # 256 KiB
308+
total_size = 600_000
309+
payload = b"A" * total_size
310+
stream = io.BytesIO(payload)
311+
312+
session = await client.upload_media(
313+
request=UploadMediaRequest(name="async_multi_chunk.bin"),
314+
config=ResumableUploadConfig(chunk_size=chunk_size),
315+
)
316+
317+
progress_records = [
318+
p async for p in session.upload(stream=stream, size=total_size)
319+
]
320+
321+
assert isinstance(session.response, UploadMediaResponse)
322+
assert session.response.name == "async_multi_chunk.bin"
323+
assert session.response.size == total_size
324+
325+
phases = [p.state for p in progress_records]
326+
offsets = [p.bytes_uploaded for p in progress_records]
327+
assert phases == [
328+
ProgressState.STARTED,
329+
ProgressState.UPLOADING,
330+
ProgressState.UPLOADING,
331+
ProgressState.FINALIZED,
332+
]
333+
assert offsets == [0, 262_144, 524_288, 600_000]
334+
assert all(p.total_bytes == total_size for p in progress_records)

‎packages/gapic-generator/tests/system/test_resumable_upload_resume.py‎

Lines changed: 109 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@
1313
# limitations under the License.
1414

1515
import io
16+
import os
1617
import pytest
1718

1819
from google.api_core import exceptions
@@ -21,7 +22,7 @@
2122
ResumableUploadConfig,
2223
ResumableUploadSession,
2324
)
24-
from google.showcase import UploadMediaResponse
25+
from google.showcase import UploadMediaRequest, UploadMediaResponse
2526

2627
from conftest import resume_resumable_upload
2728

@@ -500,3 +501,110 @@ class _ClientPauseError(Exception):
500501
assert isinstance(session2.response, UploadMediaResponse)
501502
assert session2.response.name == "golden_user_resume.mp4"
502503
assert session2.response.size == total_size
504+
505+
506+
if os.environ.get("GAPIC_PYTHON_ASYNC", "true") == "true":
507+
508+
@pytest.mark.asyncio
509+
async def test_async_client_upload_media_resume_direct(
510+
intercepted_resumable_upload_rest_async,
511+
):
512+
"""Verify direct async resumption via ResumableUploadServiceAsyncClient."""
513+
client, _ = intercepted_resumable_upload_rest_async
514+
total_size = 1_500_000
515+
chunk_size = 524_288 # 512 KiB
516+
data = b"R" * total_size
517+
stream = io.BytesIO(data)
518+
519+
class _ClientPauseError(Exception):
520+
pass
521+
522+
session1 = await client.upload_media(
523+
request=UploadMediaRequest(name="async_resume_direct.txt"),
524+
config=ResumableUploadConfig(chunk_size=chunk_size),
525+
)
526+
with pytest.raises(_ClientPauseError):
527+
async for progress in session1.upload(stream=stream, size=total_size):
528+
if progress.bytes_uploaded >= chunk_size:
529+
raise _ClientPauseError("Simulated user pause after first chunk")
530+
531+
saved_url = session1.upload_url
532+
saved_chunk_size = session1.chunk_size
533+
assert saved_url is not None
534+
assert saved_chunk_size == chunk_size
535+
536+
stream.seek(0)
537+
session2 = await client.upload_media(
538+
request=UploadMediaRequest(name="async_resume_direct.txt"),
539+
config=ResumableUploadConfig(chunk_size=saved_chunk_size),
540+
)
541+
response = await session2.resume(
542+
upload_url=saved_url,
543+
stream=stream,
544+
size=total_size,
545+
chunk_size=saved_chunk_size,
546+
)
547+
548+
assert session2.finished
549+
assert isinstance(response, UploadMediaResponse)
550+
assert response.name == "async_resume_direct.txt"
551+
assert response.size == total_size
552+
553+
@pytest.mark.asyncio
554+
async def test_async_client_upload_media_resume_seekable(
555+
intercepted_resumable_upload_rest_async,
556+
):
557+
"""Verify async progress iteration during resumption via ResumableUploadServiceAsyncClient."""
558+
client, _ = intercepted_resumable_upload_rest_async
559+
total_size = 1_500_000
560+
chunk_size = 524_288 # 512 KiB
561+
data = b"G" * total_size
562+
stream = io.BytesIO(data)
563+
564+
class _ClientPauseError(Exception):
565+
pass
566+
567+
session1 = await client.upload_media(
568+
request=UploadMediaRequest(name="async_golden_user_resume.mp4"),
569+
config=ResumableUploadConfig(chunk_size=chunk_size),
570+
)
571+
session1_snapshots = []
572+
with pytest.raises(_ClientPauseError):
573+
async for progress in session1.upload(stream=stream, size=total_size):
574+
session1_snapshots.append(progress)
575+
if progress.bytes_uploaded >= chunk_size:
576+
raise _ClientPauseError("Simulated user pause after first chunk")
577+
578+
saved_url = session1.upload_url
579+
saved_chunk_size = session1.chunk_size
580+
assert saved_url is not None
581+
assert saved_chunk_size == chunk_size
582+
assert [(p.state, p.bytes_uploaded) for p in session1_snapshots] == [
583+
(ProgressState.STARTED, 0),
584+
(ProgressState.UPLOADING, 524_288),
585+
]
586+
587+
stream.seek(0)
588+
session2 = await client.upload_media(
589+
request=UploadMediaRequest(name="async_golden_user_resume.mp4"),
590+
config=ResumableUploadConfig(chunk_size=saved_chunk_size),
591+
)
592+
session2_snapshots = [
593+
p
594+
async for p in session2.resume(
595+
upload_url=saved_url,
596+
stream=stream,
597+
size=total_size,
598+
chunk_size=saved_chunk_size,
599+
)
600+
]
601+
602+
assert [(p.state, p.bytes_uploaded) for p in session2_snapshots] == [
603+
(ProgressState.OFFSET_RECEIVED, 524_288),
604+
(ProgressState.UPLOADING, 1_048_576),
605+
(ProgressState.FINALIZED, 1_500_000),
606+
]
607+
assert session2.finished
608+
assert isinstance(session2.response, UploadMediaResponse)
609+
assert session2.response.name == "async_golden_user_resume.mp4"
610+
assert session2.response.size == total_size

0 commit comments

Comments
 (0)