diff --git a/tests/test_worker_inbound_transport.py b/tests/test_worker_inbound_transport.py index 1ee4241b..d1681d35 100644 --- a/tests/test_worker_inbound_transport.py +++ b/tests/test_worker_inbound_transport.py @@ -338,6 +338,13 @@ async def test_reconnect_redelivers_an_inbound_result_not_yet_acknowledged( async def test_lost_result_ack_refetch_keeps_the_committed_inbound_artifact( inbound, monkeypatch ): + from worker.inbound.artifacts import KeyedArtifactTransport + cleanup_done = asyncio.Event() + real_cleanup = KeyedArtifactTransport.result_acked_async + async def observe_cleanup(transport, artifacts): + await real_cleanup(transport, artifacts) + cleanup_done.set() + monkeypatch.setattr(KeyedArtifactTransport, "result_acked_async", observe_cleanup) issued = inbound.keys.issue("Test panel") payload = b"rendered audio" * 24_000 @@ -395,9 +402,10 @@ async def test_lost_result_ack_refetch_keeps_the_committed_inbound_artifact( assert open(committed, "rb").read() == payload # The ACK removes the pending item before awaiting filesystem cleanup. # Observe completion of that asynchronous operation, not its queue marker. - await _until(lambda: inbound.artifacts.open_result( + await asyncio.wait_for(cleanup_done.wait(), timeout=5) + assert inbound.artifacts.open_result( artifact_id, key_id=inbound.panel_key_id - ) is None) + ) is None assert inbound.executed == [task.task_id]