From 7fd2f4f4e8a3f8b910a89d1700407a6d3e0ca384 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9D=8E=E6=94=BF=E8=BE=BE?= Date: Thu, 30 Jul 2026 11:21:28 +0800 Subject: [PATCH] fix(api): prevent HITL pause from aborting workflow resume (#39485) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: 李政达 <1242427577@qq.com> Co-authored-by: QuantumGhost --- .../app/apps/workflow/app_queue_manager.py | 7 +++++ .../apps/workflow/test_app_queue_manager.py | 27 ++++++++++++++++++- 2 files changed, 33 insertions(+), 1 deletion(-) diff --git a/api/core/app/apps/workflow/app_queue_manager.py b/api/core/app/apps/workflow/app_queue_manager.py index 67df3044fb2..da99d034477 100644 --- a/api/core/app/apps/workflow/app_queue_manager.py +++ b/api/core/app/apps/workflow/app_queue_manager.py @@ -9,6 +9,7 @@ from core.app.entities.queue_entities import ( QueueStopEvent, QueueWorkflowFailedEvent, QueueWorkflowPartialSuccessEvent, + QueueWorkflowPausedEvent, QueueWorkflowSucceededEvent, WorkflowQueueMessage, ) @@ -32,6 +33,11 @@ class WorkflowAppQueueManager(AppQueueManager): self._q.put(message) + # A pause ends only the current listener segment; the workflow stays PAUSED and + # resumes with the same task ID. Without this marker, listen() cleanup calls + # _abort_execution(), whose stop flag and abort command can stop the resumed run. + # This is a compatibility workaround: cancellation policy belongs to the execution + # owner, not the response-stream listener. if isinstance( event, QueueStopEvent @@ -39,6 +45,7 @@ class WorkflowAppQueueManager(AppQueueManager): | QueueMessageEndEvent | QueueWorkflowSucceededEvent | QueueWorkflowFailedEvent + | QueueWorkflowPausedEvent | QueueWorkflowPartialSuccessEvent, ): self.stop_listen(execution_terminal=True) diff --git a/api/tests/unit_tests/core/app/apps/workflow/test_app_queue_manager.py b/api/tests/unit_tests/core/app/apps/workflow/test_app_queue_manager.py index 5e9a69738c4..0867a92f9d4 100644 --- a/api/tests/unit_tests/core/app/apps/workflow/test_app_queue_manager.py +++ b/api/tests/unit_tests/core/app/apps/workflow/test_app_queue_manager.py @@ -5,7 +5,12 @@ from unittest.mock import patch from core.app.apps.base_app_queue_manager import PublishFrom from core.app.apps.workflow.app_queue_manager import WorkflowAppQueueManager from core.app.entities.app_invoke_entities import InvokeFrom -from core.app.entities.queue_entities import QueueMessageEndEvent, QueuePingEvent, QueueStopEvent +from core.app.entities.queue_entities import ( + QueueMessageEndEvent, + QueuePingEvent, + QueueStopEvent, + QueueWorkflowPausedEvent, +) class TestWorkflowAppQueueManager: @@ -99,3 +104,23 @@ class TestWorkflowAppQueueManager: _ = list(manager.listen()) graph_engine_manager.return_value.send_stop_command.assert_not_called() + + def test_workflow_pause_does_not_abort_execution(self): + with ( + patch("core.app.apps.base_app_queue_manager.redis_client") as redis_client, + patch("core.app.apps.base_app_queue_manager.GraphEngineManager") as graph_engine_manager, + ): + redis_client.get.return_value = None + manager = WorkflowAppQueueManager( + task_id="task", + user_id="user", + invoke_from=InvokeFrom.DEBUGGER, + app_mode="workflow", + ) + manager.publish(QueueWorkflowPausedEvent(), PublishFrom.APPLICATION_MANAGER) + listener = manager.listen() + + assert isinstance(next(listener).event, QueueWorkflowPausedEvent) + listener.close() + + graph_engine_manager.return_value.send_stop_command.assert_not_called()