diff --git a/temporalio/worker/_workflow_instance.py b/temporalio/worker/_workflow_instance.py index a8c6d6d47..6a4169134 100644 --- a/temporalio/worker/_workflow_instance.py +++ b/temporalio/worker/_workflow_instance.py @@ -86,6 +86,10 @@ # Set to true to log all cases where we're ignoring things during delete LOG_IGNORE_DURING_DELETE = False +# Core answers a query carrying this id on the query's own task, alone. Held in +# step with LEGACY_QUERY_ID in sdk-core, crates/sdk-core/src/worker/workflow/mod.rs. +_LEGACY_QUERY_ID = "legacy_query" + def _is_workflow_terminal_command( command: temporalio.bridge.proto.workflow_commands.workflow_commands_pb2.WorkflowCommand, @@ -2859,6 +2863,17 @@ def _emit_external_stream_commands(self) -> None: if runtime is None or self._deleting: return + # A legacy query is answered on a task of its own, and Core refuses any + # other command beside that answer. The activation ran no Workflow code, + # so the wait set it would report is the one the retained task already + # holds, and there is no observation delta to commit. + if any( + command.HasField("respond_to_query") + and command.respond_to_query.query_id == _LEGACY_QUERY_ID + for command in self._current_completion.successful.commands + ): + return + # First, because the boundary a Continue-As-New header has to carry is # the one this activation is about to commit, and it only stops moving # here. diff --git a/tests/contrib/external_workflow_streams/test_worker_integration.py b/tests/contrib/external_workflow_streams/test_worker_integration.py index 7c337ffc6..a88c8cc89 100644 --- a/tests/contrib/external_workflow_streams/test_worker_integration.py +++ b/tests/contrib/external_workflow_streams/test_worker_integration.py @@ -359,6 +359,72 @@ async def test_a_blocked_workflow_retains_its_task_rather_than_completing_it( await handle.terminate() +@workflow.defn +class ParkedQueryWorkflow: + """Blocks on an empty stream with a short idle timeout, so the task parks.""" + + def __init__(self) -> None: + self._seen: list[str] = [] + + @workflow.run + async def run(self) -> list[str]: + tokens = external_stream.with_options(idle_timeout=timedelta(seconds=1)).topic( + "tokens", type=str + ) + async for token in tokens.subscribe(): + self._seen.append(token) + break + return self._seen + + @workflow.query + def seen(self) -> list[str]: + return self._seen + + +async def test_a_query_against_a_parked_workflow_is_answered( + client: Client, backend: MemoryStreamBackend +) -> None: + """A parked Run has no open Workflow Task, so its query travels alone. + + The server dispatches it on a task of its own, and Core allows nothing + beside the answer on that task. The wait set is still registered from the + parked task, so reporting it again with the answer made Core refuse the + completion, and the query never returned. + """ + task_queue = f"tq-{uuid.uuid4()}" + async with Worker( + client, + task_queue=task_queue, + workflows=[ParkedQueryWorkflow], + external_stream_backend=backend, + ): + handle = await client.start_workflow( + ParkedQueryWorkflow.run, + id=f"wf-{uuid.uuid4()}", + task_queue=task_queue, + ) + # Past the idle timeout the retained task is reported and the Run + # waits for a wake with no task open. + await asyncio.sleep(3) + + assert await asyncio.wait_for(handle.query(ParkedQueryWorkflow.seen), 10) == [] + + description = await handle.describe() + key = StreamKey( + client.namespace, + handle.id, + description.raw_description.workflow_execution_info.first_run_id, + "tokens", + ) + await publish(backend, key, ["first"]) + assert await asyncio.wait_for(handle.result(), 60) == ["first"] + + events = [e async for e in handle.fetch_history_events()] + assert not any( + e.HasField("workflow_task_failed_event_attributes") for e in events + ), "a Workflow Task failed while the query was outstanding" + + async def test_a_workflow_without_a_configured_backend_says_so( client: Client, ) -> None: