diff --git a/tests/test_activity.py b/tests/test_activity.py index a5682f221..53ffe835e 100644 --- a/tests/test_activity.py +++ b/tests/test_activity.py @@ -752,7 +752,7 @@ async def test_manual_completion(client: Client, env: WorkflowEnvironment): ActivityInput(event_workflow_id=event_workflow_id), id=activity_id, task_queue=task_queue, - start_to_close_timeout=timedelta(seconds=5), + start_to_close_timeout=timedelta(minutes=1), ) async with Worker( @@ -794,7 +794,7 @@ async def test_manual_cancellation(client: Client, env: WorkflowEnvironment): ActivityInput(event_workflow_id=event_workflow_id), id=activity_id, task_queue=task_queue, - start_to_close_timeout=timedelta(seconds=5), + start_to_close_timeout=timedelta(minutes=1), ) async with Worker( @@ -855,7 +855,7 @@ async def test_manual_failure(client: Client, env: WorkflowEnvironment): ActivityInput(event_workflow_id=event_workflow_id), id=activity_id, task_queue=task_queue, - start_to_close_timeout=timedelta(seconds=5), + start_to_close_timeout=timedelta(minutes=1), ) async with Worker( client, @@ -932,7 +932,7 @@ async def test_manual_heartbeat(client: Client, env: WorkflowEnvironment): ), id=activity_id, task_queue=task_queue, - start_to_close_timeout=timedelta(seconds=5), + start_to_close_timeout=timedelta(minutes=1), ) wait_for_activity_start_wf_handle = await client.start_workflow( EventWorkflow.wait, diff --git a/tests/worker/test_workflow.py b/tests/worker/test_workflow.py index d751fdb8b..dcb71a231 100644 --- a/tests/worker/test_workflow.py +++ b/tests/worker/test_workflow.py @@ -947,7 +947,7 @@ async def run(self, params: CancelActivityWorkflowParams) -> None: if params.local: handle = workflow.start_local_activity_method( ActivityWaitCancelNotify.wait_cancel, - schedule_to_close_timeout=timedelta(seconds=5), + schedule_to_close_timeout=timedelta(minutes=1), cancellation_type=workflow.ActivityCancellationType[ params.cancellation_type ], @@ -956,7 +956,7 @@ async def run(self, params: CancelActivityWorkflowParams) -> None: handle = workflow.start_activity_method( ActivityWaitCancelNotify.wait_cancel, schedule_to_close_timeout=timedelta(seconds=5), - heartbeat_timeout=timedelta(seconds=1), + heartbeat_timeout=timedelta(seconds=5), cancellation_type=workflow.ActivityCancellationType[ params.cancellation_type ], @@ -977,14 +977,24 @@ def activity_result(self) -> str: @pytest.mark.parametrize("local", [True, False]) async def test_workflow_cancel_activity(client: Client, local: bool): - # Need short task timeout to timeout LA task and longer assert timeout - # so the task can timeout - task_timeout = timedelta(seconds=1) - assert_timeout = timedelta(seconds=10) + # Core completes the task holding a local activity at 80% of this timeout, and + # the cancel reaches the activity on the second such cycle (~8s), so the assert + # budget needs headroom beyond that on loaded runners + task_timeout = timedelta(seconds=5) + assert_timeout = timedelta(seconds=30) activity_inst = ActivityWaitCancelNotify() + async def wait_cancel_complete() -> None: + await asyncio.wait_for( + activity_inst.wait_cancel_complete.wait(), assert_timeout.total_seconds() + ) + activity_inst.wait_cancel_complete.clear() + async with new_worker( - client, CancelActivityWorkflow, activities=[activity_inst.wait_cancel] + client, + CancelActivityWorkflow, + activities=[activity_inst.wait_cancel], + max_heartbeat_throttle_interval=timedelta(milliseconds=300), ) as worker: # Try cancel - confirm error and activity was sent the cancel handle = await client.start_workflow( @@ -1004,7 +1014,7 @@ async def activity_result() -> str: await assert_eq_eventually( "Error: CancelledError", activity_result, timeout=assert_timeout ) - await activity_inst.wait_cancel_complete.wait() + await wait_cancel_complete() await handle.cancel() # Wait cancel - confirm no error due to graceful cancel handling @@ -1023,7 +1033,7 @@ async def activity_result() -> str: activity_result, timeout=assert_timeout, ) - await activity_inst.wait_cancel_complete.wait() + await wait_cancel_complete() await handle.cancel() # Abandon - confirm error and that activity stays running @@ -1043,7 +1053,7 @@ async def activity_result() -> str: await asyncio.sleep(0.5) assert not activity_inst.wait_cancel_complete.is_set() await handle.cancel() - await activity_inst.wait_cancel_complete.wait() + await wait_cancel_complete() @workflow.defn