Skip to content

Commit 91d5f1b

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

2 files changed

Lines changed: 97 additions & 4 deletions

File tree

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

Lines changed: 21 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -560,6 +560,7 @@ async def close(
560560
self,
561561
finalize_on_close=False,
562562
full_object_checksum: Optional[int] = None,
563+
retry_policy: Optional[AsyncRetry] = None,
563564
) -> Union[int, _storage_v2.Object]:
564565
"""Closes the underlying bidi-gRPC stream.
565566
@@ -581,6 +582,9 @@ async def close(
581582
crc32c_int = google_crc32c.value(data)
582583
print(crc32c_int)
583584
585+
:type retry_policy: :class:`~google.api_core.retry_async.AsyncRetry`
586+
:param retry_policy: (Optional) The retry policy to use for the operation.
587+
584588
rtype: Union[int, _storage_v2.Object]
585589
returns: Updated `self.persisted_size` by default after closing the
586590
bidi-gRPC stream. However, if `finalize_on_close=True` is passed,
@@ -604,15 +608,20 @@ async def close(
604608
)
605609

606610
if finalize_on_close:
607-
return await self.finalize(full_object_checksum=full_object_checksum)
611+
return await self.finalize(
612+
full_object_checksum=full_object_checksum,
613+
retry_policy=retry_policy,
614+
)
608615

609616
await self.write_obj_stream.close()
610617

611618
self._is_stream_open = False
612619
return self.persisted_size
613620

614621
async def finalize(
615-
self, full_object_checksum: Optional[int] = None
622+
self,
623+
full_object_checksum: Optional[int] = None,
624+
retry_policy: Optional[AsyncRetry] = None,
616625
) -> _storage_v2.Object:
617626
"""Finalizes the Appendable Object.
618627
@@ -638,6 +647,9 @@ async def finalize(
638647
crc32c_int = google_crc32c.value(data)
639648
print(crc32c_int)
640649
650+
:type retry_policy: :class:`~google.api_core.retry_async.AsyncRetry`
651+
:param retry_policy: (Optional) The retry policy to use for the operation.
652+
641653
rtype: google.cloud.storage_v2.types.Object
642654
returns: The finalized object resource.
643655
@@ -666,12 +678,18 @@ async def finalize(
666678
),
667679
)
668680

669-
try:
681+
if retry_policy is None:
682+
retry_policy = AsyncRetry(predicate=_is_write_retryable)
683+
684+
async def _do_finalize():
670685
await self.write_obj_stream.send(finalize_req)
671686
response = await self.write_obj_stream.recv()
672687
self.object_resource = response.resource
673688
self.persisted_size = self.object_resource.size
674689
return self.object_resource
690+
691+
try:
692+
return await retry_policy(_do_finalize)()
675693
finally:
676694
await self.write_obj_stream.close()
677695
self._is_stream_open = False

packages/google-cloud-storage/tests/unit/asyncio/test_async_appendable_object_writer.py

Lines changed: 76 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -560,7 +560,9 @@ async def test_close_with_checksum_and_finalize(self, mock_appendable_writer):
560560

561561
checksum = 12345678
562562
await writer.close(finalize_on_close=True, full_object_checksum=checksum)
563-
writer.finalize.assert_awaited_once_with(full_object_checksum=checksum)
563+
writer.finalize.assert_awaited_once_with(
564+
full_object_checksum=checksum, retry_policy=None
565+
)
564566

565567
@pytest.mark.asyncio
566568
async def test_close_with_checksum_without_finalize_raises(
@@ -625,3 +627,76 @@ async def test_finalize_mismatch_closes_stream(self, mock_appendable_writer):
625627
# Assert stream was closed and local state reset despite exception
626628
mock_appendable_writer["mock_stream"].close.assert_awaited()
627629
assert not writer._is_stream_open
630+
631+
@pytest.mark.asyncio
632+
async def test_finalize_retry_on_transient_error(self, mock_appendable_writer):
633+
writer = self._make_one(mock_appendable_writer["mock_client"])
634+
writer._is_stream_open = True
635+
writer.write_obj_stream = mock_appendable_writer["mock_stream"]
636+
637+
resource = storage_type.Object(size=999)
638+
mock_appendable_writer["mock_stream"].recv.side_effect = [
639+
exceptions.InternalServerError("500 Transient Error"),
640+
storage_type.BidiWriteObjectResponse(resource=resource),
641+
]
642+
643+
res = await writer.finalize()
644+
645+
assert res == resource
646+
assert writer.persisted_size == 999
647+
assert mock_appendable_writer["mock_stream"].send.await_count == 2
648+
assert not writer._is_stream_open
649+
650+
@pytest.mark.asyncio
651+
async def test_finalize_custom_retry_policy(self, mock_appendable_writer):
652+
from google.api_core.retry_async import AsyncRetry
653+
654+
writer = self._make_one(mock_appendable_writer["mock_client"])
655+
writer._is_stream_open = True
656+
writer.write_obj_stream = mock_appendable_writer["mock_stream"]
657+
658+
custom_policy = AsyncRetry(predicate=lambda exc: isinstance(exc, exceptions.InternalServerError))
659+
resource = storage_type.Object(size=999)
660+
mock_appendable_writer["mock_stream"].recv.return_value = (
661+
storage_type.BidiWriteObjectResponse(resource=resource)
662+
)
663+
664+
res = await writer.finalize(retry_policy=custom_policy)
665+
assert res == resource
666+
667+
@pytest.mark.asyncio
668+
async def test_close_with_finalize_and_custom_retry_policy(self, mock_appendable_writer):
669+
from google.api_core.retry_async import AsyncRetry
670+
671+
writer = self._make_one(mock_appendable_writer["mock_client"])
672+
writer._is_stream_open = True
673+
writer.finalize = AsyncMock()
674+
675+
custom_policy = AsyncRetry(predicate=lambda exc: False)
676+
await writer.close(finalize_on_close=True, retry_policy=custom_policy)
677+
writer.finalize.assert_awaited_once_with(
678+
full_object_checksum=None,
679+
retry_policy=custom_policy,
680+
)
681+
682+
@pytest.mark.asyncio
683+
async def test_close_retry_on_transient_error(self, mock_appendable_writer):
684+
writer = self._make_one(mock_appendable_writer["mock_client"])
685+
writer._is_stream_open = True
686+
writer.write_obj_stream = mock_appendable_writer["mock_stream"]
687+
688+
resource = storage_type.Object(size=999)
689+
mock_appendable_writer["mock_stream"].recv.side_effect = [
690+
exceptions.InternalServerError("500 Transient Error"),
691+
storage_type.BidiWriteObjectResponse(resource=resource),
692+
]
693+
694+
res = await writer.close(finalize_on_close=True)
695+
696+
assert res == resource
697+
assert writer.persisted_size == 999
698+
assert mock_appendable_writer["mock_stream"].send.await_count == 2
699+
assert not writer._is_stream_open
700+
701+
702+

0 commit comments

Comments
 (0)