Skip to content

Commit 606a8d3

Browse files
brianstrauch1fanwang
authored andcommitted
Retry cancelled background stream flushes
1 parent 5aa18e4 commit 606a8d3

2 files changed

Lines changed: 8 additions & 5 deletions

File tree

‎temporalio/contrib/workflow_streams/_client.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@
6161
_RETRYABLE_RPC_STATUS_CODES: frozenset[RPCStatusCode] = frozenset(
6262
{
6363
RPCStatusCode.ABORTED,
64+
RPCStatusCode.CANCELLED,
6465
RPCStatusCode.DATA_LOSS,
6566
RPCStatusCode.INTERNAL,
6667
RPCStatusCode.OUT_OF_RANGE,

‎tests/contrib/workflow_streams/test_workflow_streams.py‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1417,7 +1417,10 @@ async def maybe_failing_signal(*args: Any, **kwargs: Any) -> Any:
14171417

14181418

14191419
@pytest.mark.asyncio
1420-
async def test_background_flusher_retries_failed_signal(client: Client) -> None:
1420+
@pytest.mark.parametrize("status", [RPCStatusCode.UNAVAILABLE, RPCStatusCode.CANCELLED])
1421+
async def test_background_flusher_retries_failed_signal(
1422+
client: Client, status: RPCStatusCode
1423+
) -> None:
14211424
async with new_worker(client, BasicWorkflowStreamWorkflow) as worker:
14221425
handle = await client.start_workflow(
14231426
BasicWorkflowStreamWorkflow.run,
@@ -1430,17 +1433,16 @@ async def test_background_flusher_retries_failed_signal(client: Client) -> None:
14301433
first_flush_failed = asyncio.Event()
14311434
retry_succeeded = asyncio.Event()
14321435

1433-
async def fail_first_signal(*args: Any, **kwargs: Any) -> Any:
1436+
async def fail_first_signal(*args: Any, **kwargs: Any) -> None:
14341437
if not first_flush_failed.is_set():
14351438
first_flush_failed.set()
14361439
raise RPCError(
14371440
message="simulated delivery failure",
1438-
status=RPCStatusCode.UNAVAILABLE,
1441+
status=status,
14391442
raw_grpc_status=b"",
14401443
)
1441-
result = await real_signal(*args, **kwargs)
1444+
await real_signal(*args, **kwargs)
14421445
retry_succeeded.set()
1443-
return result
14441446

14451447
with patch.object(handle, "signal", side_effect=fail_first_signal):
14461448
async with stream:

0 commit comments

Comments
 (0)