Skip to content

Commit b509e2e

Browse files
committed
fix(storage): add retry support for finalize and close in AsyncAppendableObjectWriter
Fixes: b/532527637
1 parent e52b015 commit b509e2e

3 files changed

Lines changed: 273 additions & 42 deletions

File tree

packages/google-cloud-storage/google/cloud/storage/asyncio/async_appendable_object_writer.py

Lines changed: 99 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,6 @@
1919

2020
from google.api_core import exceptions
2121
from google.api_core.retry_async import AsyncRetry
22-
from google.rpc import status_pb2
23-
2422
from google.cloud import _storage_v2
2523
from google.cloud._storage_v2.types import BidiWriteObjectRedirectedError
2624
from google.cloud._storage_v2.types.storage import BidiWriteObjectRequest
@@ -41,6 +39,7 @@
4139
_WriteResumptionStrategy,
4240
_WriteState,
4341
)
42+
from google.rpc import status_pb2
4443

4544
from . import _utils
4645

@@ -299,21 +298,11 @@ def _on_open_error(self, exc):
299298
if redirect_proto.generation:
300299
self.generation = redirect_proto.generation
301300

302-
async def open(
303-
self,
304-
retry_policy: Optional[AsyncRetry] = None,
305-
metadata: Optional[List[Tuple[str, str]]] = None,
306-
) -> None:
307-
"""Opens the underlying bidi-gRPC stream.
308-
309-
:raises ValueError: If the stream is already open.
310-
311-
"""
312-
if self._is_stream_open:
313-
raise ValueError("Underlying bidi-gRPC stream is already open")
314-
301+
def _merge_retry_policy(
302+
self, retry_policy: Optional[AsyncRetry] = None
303+
) -> AsyncRetry:
315304
if retry_policy is None:
316-
retry_policy = AsyncRetry(
305+
return AsyncRetry(
317306
predicate=_is_write_retryable, on_error=self._on_open_error
318307
)
319308
else:
@@ -324,7 +313,7 @@ def combined_on_error(exc):
324313
if original_on_error:
325314
original_on_error(exc)
326315

327-
retry_policy = AsyncRetry(
316+
return AsyncRetry(
328317
predicate=_is_write_retryable,
329318
initial=retry_policy._initial,
330319
maximum=retry_policy._maximum,
@@ -333,6 +322,21 @@ def combined_on_error(exc):
333322
on_error=combined_on_error,
334323
)
335324

325+
async def open(
326+
self,
327+
retry_policy: Optional[AsyncRetry] = None,
328+
metadata: Optional[List[Tuple[str, str]]] = None,
329+
) -> None:
330+
"""Opens the underlying bidi-gRPC stream.
331+
332+
:raises ValueError: If the stream is already open.
333+
334+
"""
335+
if self._is_stream_open:
336+
raise ValueError("Underlying bidi-gRPC stream is already open")
337+
338+
retry_policy = self._merge_retry_policy(retry_policy)
339+
336340
async def _do_open():
337341
current_metadata = list(metadata) if metadata else []
338342

@@ -560,6 +564,7 @@ async def close(
560564
self,
561565
finalize_on_close=False,
562566
full_object_checksum: Optional[int] = None,
567+
retry_policy: Optional[AsyncRetry] = None,
563568
) -> Union[int, _storage_v2.Object]:
564569
"""Closes the underlying bidi-gRPC stream.
565570
@@ -581,6 +586,9 @@ async def close(
581586
crc32c_int = google_crc32c.value(data)
582587
print(crc32c_int)
583588
589+
:type retry_policy: :class:`~google.api_core.retry_async.AsyncRetry`
590+
:param retry_policy: (Optional) The retry policy to use for the operation.
591+
584592
rtype: Union[int, _storage_v2.Object]
585593
returns: Updated `self.persisted_size` by default after closing the
586594
bidi-gRPC stream. However, if `finalize_on_close=True` is passed,
@@ -604,15 +612,47 @@ async def close(
604612
)
605613

606614
if finalize_on_close:
607-
return await self.finalize(full_object_checksum=full_object_checksum)
615+
return await self.finalize(
616+
full_object_checksum=full_object_checksum,
617+
retry_policy=retry_policy,
618+
)
608619

609-
await self.write_obj_stream.close()
620+
retry_policy = self._merge_retry_policy(retry_policy)
610621

611-
self._is_stream_open = False
612-
return self.persisted_size
622+
attempt_count = 0
623+
624+
async def _do_close():
625+
nonlocal attempt_count
626+
attempt_count += 1
627+
628+
if attempt_count > 1:
629+
logger.info(
630+
f"Re-opening the stream for close retry attempt: {attempt_count}"
631+
)
632+
expected_offset = self.offset
633+
self._is_stream_open = False
634+
await self.open()
635+
if (
636+
self.offset is not None
637+
and expected_offset is not None
638+
and self.offset < expected_offset
639+
):
640+
raise exceptions.InternalServerError(
641+
f"Unrecoverable data loss during reconnect. Expected offset {expected_offset}, got {self.offset}"
642+
)
643+
644+
await self.write_obj_stream.close()
645+
return self.persisted_size
646+
647+
try:
648+
return await retry_policy(_do_close)()
649+
finally:
650+
self._is_stream_open = False
613651

614652
async def finalize(
615-
self, full_object_checksum: Optional[int] = None
653+
self,
654+
full_object_checksum: Optional[int] = None,
655+
retry_policy: Optional[AsyncRetry] = None,
616656
) -> _storage_v2.Object:
617657
"""Finalizes the Appendable Object.
618658
@@ -638,6 +678,9 @@ async def finalize(
638678
crc32c_int = google_crc32c.value(data)
639679
print(crc32c_int)
640680
681+
:type retry_policy: :class:`~google.api_core.retry_async.AsyncRetry`
682+
:param retry_policy: (Optional) The retry policy to use for the operation.
683+
641684
rtype: google.cloud.storage_v2.types.Object
642685
returns: The finalized object resource.
643686
@@ -666,14 +709,46 @@ async def finalize(
666709
),
667710
)
668711

669-
try:
712+
retry_policy = self._merge_retry_policy(retry_policy)
713+
714+
attempt_count = 0
715+
716+
async def _do_finalize():
717+
nonlocal attempt_count
718+
attempt_count += 1
719+
720+
if attempt_count > 1:
721+
logger.info(
722+
f"Re-opening the stream for finalize retry attempt: {attempt_count}"
723+
)
724+
expected_offset = self.offset
725+
self._is_stream_open = False
726+
await self.open()
727+
if (
728+
self.offset is not None
729+
and expected_offset is not None
730+
and self.offset < expected_offset
731+
):
732+
raise exceptions.InternalServerError(
733+
f"Unrecoverable data loss during reconnect. Expected offset {expected_offset}, got {self.offset}"
734+
)
735+
670736
await self.write_obj_stream.send(finalize_req)
671737
response = await self.write_obj_stream.recv()
672738
self.object_resource = response.resource
673739
self.persisted_size = self.object_resource.size
674740
return self.object_resource
741+
742+
try:
743+
return await retry_policy(_do_finalize)()
675744
finally:
676-
await self.write_obj_stream.close()
745+
if self.write_obj_stream:
746+
try:
747+
await self.write_obj_stream.close()
748+
except Exception as e:
749+
logger.debug(
750+
f"Stream close during finalize cleanup resulted in: {e}"
751+
)
677752
self._is_stream_open = False
678753
self.offset = None
679754

packages/google-cloud-storage/tests/conformance/test_bidi_writes.py

Lines changed: 38 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -6,10 +6,10 @@
66
import grpc
77
import pytest
88
import requests
9+
910
from google.api_core import client_options, exceptions
1011
from google.api_core.retry_async import AsyncRetry
1112
from google.auth import credentials as auth_credentials
12-
1313
from google.cloud import _storage_v2 as storage_v2
1414
from google.cloud.storage.asyncio.async_appendable_object_writer import (
1515
AsyncAppendableObjectWriter,
@@ -136,31 +136,34 @@ def on_retry_error(exc):
136136
CONTENT, metadata=fault_injection_metadata, retry_policy=policy_to_pass
137137
)
138138
# await writer.finalize()
139-
await writer.close(finalize_on_close=True)
139+
f_o_c = scenario.get("finalize_on_close", True)
140+
await writer.close(finalize_on_close=f_o_c, retry_policy=policy_to_pass)
140141

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

147-
# 4. Verify the object content.
148-
read_request = storage_v2.ReadObjectRequest(
149-
bucket=f"projects/_/buckets/{bucket_name}",
150-
object=object_name,
151-
)
152-
read_stream = await gapic_client.read_object(request=read_request)
153-
data = b""
154-
async for chunk in read_stream:
155-
data += chunk.checksummed_data.content
156-
assert data == CONTENT
148+
# 4. Verify the object content if applicable.
149+
if not scenario.get("skip_verification"):
150+
read_request = storage_v2.ReadObjectRequest(
151+
bucket=f"projects/_/buckets/{bucket_name}",
152+
object=object_name,
153+
)
154+
read_stream = await gapic_client.read_object(request=read_request)
155+
data = b""
156+
async for chunk in read_stream:
157+
data += chunk.checksummed_data.content
158+
assert data == CONTENT
159+
157160
if scenario["expected_error"] is None:
158161
# Scenarios like 503, 500, smarter resumption, and redirects
159162
# SHOULD trigger at least one retry attempt.
160163
if not use_default:
161-
assert retry_count > 0, (
162-
f"Test passed but no retry was actually triggered for {scenario['name']}!"
163-
)
164+
assert (
165+
retry_count > 0
166+
), f"Test passed but no retry was actually triggered for {scenario['name']}!"
164167
else:
165168
print("Successfully recovered using library's default policy.")
166169
print(f"Success: {scenario['name']}")
@@ -235,6 +238,26 @@ async def test_bidi_writes(testbench):
235238
"instruction": "redirect-send-handle-and-token-tokenval",
236239
"expected_error": None,
237240
},
241+
{
242+
"name": "Retry exactly on finalize/close (Redirect Error)",
243+
"method": "storage.objects.insert",
244+
"instruction": "redirect-send-handle-and-token-mytoken-on-finish-write",
245+
"expected_error": None,
246+
},
247+
{
248+
"name": "Retry exactly on finalize/close (503)",
249+
"method": "storage.objects.insert",
250+
"instruction": "return-503-on-finish-write",
251+
"expected_error": None,
252+
},
253+
{
254+
"name": "Retry exactly on close (finalize_on_close=False) (503)",
255+
"method": "storage.objects.insert",
256+
"instruction": "return-503-on-half-close",
257+
"expected_error": None,
258+
"finalize_on_close": False,
259+
"skip_verification": True,
260+
},
238261
]
239262

240263
try:

0 commit comments

Comments
 (0)