@@ -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,134 @@ 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 (
659+ predicate = lambda exc : isinstance (exc , exceptions .InternalServerError )
660+ )
661+ resource = storage_type .Object (size = 999 )
662+ mock_appendable_writer [
663+ "mock_stream"
664+ ].recv .return_value = storage_type .BidiWriteObjectResponse (resource = resource )
665+
666+ res = await writer .finalize (retry_policy = custom_policy )
667+ assert res == resource
668+
669+ @pytest .mark .asyncio
670+ async def test_close_with_finalize_and_custom_retry_policy (
671+ self , mock_appendable_writer
672+ ):
673+ from google .api_core .retry_async import AsyncRetry
674+
675+ writer = self ._make_one (mock_appendable_writer ["mock_client" ])
676+ writer ._is_stream_open = True
677+ writer .finalize = AsyncMock ()
678+
679+ custom_policy = AsyncRetry (predicate = lambda exc : False )
680+ await writer .close (finalize_on_close = True , retry_policy = custom_policy )
681+ writer .finalize .assert_awaited_once_with (
682+ full_object_checksum = None ,
683+ retry_policy = custom_policy ,
684+ )
685+
686+ @pytest .mark .asyncio
687+ async def test_close_retry_on_transient_error (self , mock_appendable_writer ):
688+ writer = self ._make_one (mock_appendable_writer ["mock_client" ])
689+ writer ._is_stream_open = True
690+ writer .write_obj_stream = mock_appendable_writer ["mock_stream" ]
691+
692+ resource = storage_type .Object (size = 999 )
693+ mock_appendable_writer ["mock_stream" ].recv .side_effect = [
694+ exceptions .InternalServerError ("500 Transient Error" ),
695+ storage_type .BidiWriteObjectResponse (resource = resource ),
696+ ]
697+
698+ res = await writer .close (finalize_on_close = True )
699+
700+ assert res == resource
701+ assert writer .persisted_size == 999
702+ assert mock_appendable_writer ["mock_stream" ].send .await_count == 2
703+ assert not writer ._is_stream_open
704+
705+ @pytest .mark .asyncio
706+ async def test_finalize_retry_on_redirect_error (self , mock_appendable_writer ):
707+ writer = self ._make_one (mock_appendable_writer ["mock_client" ])
708+ writer ._is_stream_open = True
709+ writer .write_obj_stream = mock_appendable_writer ["mock_stream" ]
710+
711+ redirect = BidiWriteObjectRedirectedError (
712+ routing_token = "rt1" ,
713+ write_handle = storage_type .BidiWriteHandle (handle = b"h1" ),
714+ )
715+ exc = exceptions .Aborted ("aborted" , errors = [redirect ])
716+
717+ resource = storage_type .Object (size = 999 )
718+ mock_appendable_writer ["mock_stream" ].recv .side_effect = [
719+ exc ,
720+ storage_type .BidiWriteObjectResponse (resource = resource ),
721+ ]
722+
723+ writer .open = mock .AsyncMock ()
724+
725+ res = await writer .finalize ()
726+
727+ assert res == resource
728+ assert writer .persisted_size == 999
729+ assert mock_appendable_writer ["mock_stream" ].send .await_count == 2
730+ assert writer ._routing_token == "rt1"
731+ assert writer .write_handle .handle == b"h1"
732+ writer .open .assert_awaited_once ()
733+
734+ @pytest .mark .asyncio
735+ async def test_close_retry_on_redirect_error (self , mock_appendable_writer ):
736+ writer = self ._make_one (mock_appendable_writer ["mock_client" ])
737+ writer ._is_stream_open = True
738+ writer .write_obj_stream = mock_appendable_writer ["mock_stream" ]
739+
740+ redirect = BidiWriteObjectRedirectedError (
741+ routing_token = "rt2" ,
742+ write_handle = storage_type .BidiWriteHandle (handle = b"h2" ),
743+ )
744+ exc = exceptions .Aborted ("aborted" , errors = [redirect ])
745+
746+ mock_appendable_writer ["mock_stream" ].close .side_effect = [
747+ exc ,
748+ None ,
749+ ]
750+
751+ writer .open = mock .AsyncMock ()
752+ writer .persisted_size = 999
753+
754+ res = await writer .close ()
755+
756+ assert res == 999
757+ assert mock_appendable_writer ["mock_stream" ].close .await_count == 2
758+ assert writer ._routing_token == "rt2"
759+ assert writer .write_handle .handle == b"h2"
760+ writer .open .assert_awaited_once ()
0 commit comments