Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,6 @@

from google.api_core import exceptions
from google.api_core.retry_async import AsyncRetry
from google.rpc import status_pb2

from google.cloud import _storage_v2
from google.cloud._storage_v2.types import BidiWriteObjectRedirectedError
from google.cloud._storage_v2.types.storage import BidiWriteObjectRequest
Expand All @@ -41,6 +39,7 @@
_WriteResumptionStrategy,
_WriteState,
)
from google.rpc import status_pb2

from . import _utils

Expand Down Expand Up @@ -299,21 +298,11 @@ def _on_open_error(self, exc):
if redirect_proto.generation:
self.generation = redirect_proto.generation

async def open(
self,
retry_policy: Optional[AsyncRetry] = None,
metadata: Optional[List[Tuple[str, str]]] = None,
) -> None:
"""Opens the underlying bidi-gRPC stream.

:raises ValueError: If the stream is already open.

"""
if self._is_stream_open:
raise ValueError("Underlying bidi-gRPC stream is already open")

def _merge_retry_policy(
self, retry_policy: Optional[AsyncRetry] = None
) -> AsyncRetry:
if retry_policy is None:
retry_policy = AsyncRetry(
return AsyncRetry(
predicate=_is_write_retryable, on_error=self._on_open_error
)
else:
Expand All @@ -324,7 +313,7 @@ def combined_on_error(exc):
if original_on_error:
original_on_error(exc)

retry_policy = AsyncRetry(
return AsyncRetry(
predicate=_is_write_retryable,
initial=retry_policy._initial,
maximum=retry_policy._maximum,
Expand All @@ -333,6 +322,21 @@ def combined_on_error(exc):
on_error=combined_on_error,
)

async def open(
self,
retry_policy: Optional[AsyncRetry] = None,
metadata: Optional[List[Tuple[str, str]]] = None,
) -> None:
"""Opens the underlying bidi-gRPC stream.

:raises ValueError: If the stream is already open.

"""
if self._is_stream_open:
raise ValueError("Underlying bidi-gRPC stream is already open")

retry_policy = self._merge_retry_policy(retry_policy)

async def _do_open():
current_metadata = list(metadata) if metadata else []

Expand Down Expand Up @@ -560,6 +564,7 @@ async def close(
self,
finalize_on_close=False,
full_object_checksum: Optional[int] = None,
retry_policy: Optional[AsyncRetry] = None,
) -> Union[int, _storage_v2.Object]:
"""Closes the underlying bidi-gRPC stream.

Expand All @@ -581,6 +586,9 @@ async def close(
crc32c_int = google_crc32c.value(data)
print(crc32c_int)

:type retry_policy: :class:`~google.api_core.retry_async.AsyncRetry`
:param retry_policy: (Optional) The retry policy to use for the operation.

rtype: Union[int, _storage_v2.Object]
returns: Updated `self.persisted_size` by default after closing the
bidi-gRPC stream. However, if `finalize_on_close=True` is passed,
Expand All @@ -604,15 +612,47 @@ async def close(
)

if finalize_on_close:
return await self.finalize(full_object_checksum=full_object_checksum)
return await self.finalize(
full_object_checksum=full_object_checksum,
retry_policy=retry_policy,
)

await self.write_obj_stream.close()
retry_policy = self._merge_retry_policy(retry_policy)

self._is_stream_open = False
return self.persisted_size
attempt_count = 0

async def _do_close():
nonlocal attempt_count
attempt_count += 1

if attempt_count > 1:
logger.info(
f"Re-opening the stream for close retry attempt: {attempt_count}"
)
expected_offset = self.offset
self._is_stream_open = False
await self.open()
if (
self.offset is not None
and expected_offset is not None
and self.offset < expected_offset
):
raise exceptions.InternalServerError(
f"Unrecoverable data loss during reconnect. Expected offset {expected_offset}, got {self.offset}"
)

await self.write_obj_stream.close()
return self.persisted_size

try:
return await retry_policy(_do_close)()
finally:
self._is_stream_open = False

async def finalize(
self, full_object_checksum: Optional[int] = None
self,
full_object_checksum: Optional[int] = None,
retry_policy: Optional[AsyncRetry] = None,
) -> _storage_v2.Object:
"""Finalizes the Appendable Object.

Expand All @@ -638,6 +678,9 @@ async def finalize(
crc32c_int = google_crc32c.value(data)
print(crc32c_int)

:type retry_policy: :class:`~google.api_core.retry_async.AsyncRetry`
:param retry_policy: (Optional) The retry policy to use for the operation.

rtype: google.cloud.storage_v2.types.Object
returns: The finalized object resource.

Expand Down Expand Up @@ -666,14 +709,46 @@ async def finalize(
),
)

try:
retry_policy = self._merge_retry_policy(retry_policy)

attempt_count = 0

async def _do_finalize():
nonlocal attempt_count
attempt_count += 1

if attempt_count > 1:
logger.info(
f"Re-opening the stream for finalize retry attempt: {attempt_count}"
)
expected_offset = self.offset
self._is_stream_open = False
await self.open()
if (
self.offset is not None
and expected_offset is not None
and self.offset < expected_offset
):
raise exceptions.InternalServerError(
f"Unrecoverable data loss during reconnect. Expected offset {expected_offset}, got {self.offset}"
)

await self.write_obj_stream.send(finalize_req)
response = await self.write_obj_stream.recv()
self.object_resource = response.resource
self.persisted_size = self.object_resource.size
return self.object_resource

try:
return await retry_policy(_do_finalize)()
finally:
await self.write_obj_stream.close()
if self.write_obj_stream:
try:
await self.write_obj_stream.close()
except Exception as e:
logger.debug(
f"Stream close during finalize cleanup resulted in: {e}"
)
self._is_stream_open = False
self.offset = None

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,10 @@
import grpc
import pytest
import requests

from google.api_core import client_options, exceptions
from google.api_core.retry_async import AsyncRetry
from google.auth import credentials as auth_credentials

from google.cloud import _storage_v2 as storage_v2
from google.cloud.storage.asyncio.async_appendable_object_writer import (
AsyncAppendableObjectWriter,
Expand Down Expand Up @@ -136,24 +136,27 @@ def on_retry_error(exc):
CONTENT, metadata=fault_injection_metadata, retry_policy=policy_to_pass
)
# await writer.finalize()
await writer.close(finalize_on_close=True)
f_o_c = scenario.get("finalize_on_close", True)
await writer.close(finalize_on_close=f_o_c, retry_policy=policy_to_pass)

# If an exception was expected, this line should not be reached.
if scenario["expected_error"] is not None:
raise AssertionError(
f"Expected exception {scenario['expected_error']} was not raised."
)

# 4. Verify the object content.
read_request = storage_v2.ReadObjectRequest(
bucket=f"projects/_/buckets/{bucket_name}",
object=object_name,
)
read_stream = await gapic_client.read_object(request=read_request)
data = b""
async for chunk in read_stream:
data += chunk.checksummed_data.content
assert data == CONTENT
# 4. Verify the object content if applicable.
if not scenario.get("skip_verification"):
read_request = storage_v2.ReadObjectRequest(
bucket=f"projects/_/buckets/{bucket_name}",
object=object_name,
)
read_stream = await gapic_client.read_object(request=read_request)
data = b""
async for chunk in read_stream:
data += chunk.checksummed_data.content
assert data == CONTENT

if scenario["expected_error"] is None:
# Scenarios like 503, 500, smarter resumption, and redirects
# SHOULD trigger at least one retry attempt.
Expand Down Expand Up @@ -235,6 +238,26 @@ async def test_bidi_writes(testbench):
"instruction": "redirect-send-handle-and-token-tokenval",
"expected_error": None,
},
{
"name": "Retry exactly on finalize/close (Redirect Error)",
"method": "storage.objects.insert",
"instruction": "redirect-send-handle-and-token-mytoken-on-finish-write",
"expected_error": None,
},
{
"name": "Retry exactly on finalize/close (503)",
"method": "storage.objects.insert",
"instruction": "return-503-on-finish-write",
"expected_error": None,
},
{
"name": "Retry exactly on close (finalize_on_close=False) (503)",
"method": "storage.objects.insert",
"instruction": "return-503-on-half-close",
"expected_error": None,
"finalize_on_close": False,
"skip_verification": True,
},
]

try:
Expand Down
Loading
Loading