From f157113639cb8e9e3d4ab4bfd42197ba6b8704cd Mon Sep 17 00:00:00 2001 From: Vojta Bartos Date: Wed, 22 Jul 2026 15:21:47 +0200 Subject: [PATCH] fix(tasks): relay late Slack answers from completed runs A background (task-notification) turn can finish while its run is being torn down, or just after the inactivity timeout marked the run completed. relay_task_run_message skipped every terminal run, so that answer was silently dropped even though the Slack thread mapping still exists. Accept completed runs (failed and cancelled stay silent) and route them to the relay workflow directly, since a closed workflow can no longer receive the agent-design inline-stream signal. --- products/tasks/backend/facade/api.py | 17 ++++-- products/tasks/backend/tests/test_api.py | 78 +++++++++++++++++++++++- 2 files changed, 89 insertions(+), 6 deletions(-) diff --git a/products/tasks/backend/facade/api.py b/products/tasks/backend/facade/api.py index 836482e79989..04605696ef1c 100644 --- a/products/tasks/backend/facade/api.py +++ b/products/tasks/backend/facade/api.py @@ -2852,8 +2852,13 @@ def relay_task_run_message( flag signal the running task workflow to stream the text inline. Returns ``(status, relay_id)`` where status is ``"accepted"`` (relay_id set), ``"skipped"`` - (run not found / terminal / no Slack mapping / empty text / streamed inline under the - agent-design flag), or ``"failed"``. + (run not found / failed or cancelled / no Slack mapping / empty text / streamed inline + under the agent-design flag), or ``"failed"``. + + A completed run still relays: a background (task-notification) turn can finish while + the run is being torn down, or just after the inactivity timeout completed it, and its + answer should reach the thread rather than be dropped. Failed and cancelled runs stay + silent — their terminal Slack card is the last word. When ``text_parts`` is provided the last non-empty entry is used — it's the post-last-tool-use answer, and posting only that keeps the interim narration @@ -2872,7 +2877,9 @@ def relay_task_run_message( ) run = _get_visible_run(run_id, task_id, team_id) - if run is None or run.is_terminal: + if run is None: + return "skipped", None + if run.is_terminal and run.status != TaskRun.Status.COMPLETED: return "skipped", None if not SlackThreadTaskMapping.objects.filter(task_run=run).exists(): return "skipped", None @@ -2882,7 +2889,9 @@ def relay_task_run_message( if not trimmed: return "skipped", None - if bool((run.state or {}).get(AGENT_DESIGN_STATE_KEY)): + # A terminal run's workflow is closed, so the inline-stream signal can't be + # delivered — fall through to the relay workflow even under the flag. + if not run.is_terminal and bool((run.state or {}).get(AGENT_DESIGN_STATE_KEY)): try: signal_agent_text_delta(run.workflow_id, trimmed) except Exception: diff --git a/products/tasks/backend/tests/test_api.py b/products/tasks/backend/tests/test_api.py index 26f05f5b4888..c3bdbd8695c8 100644 --- a/products/tasks/backend/tests/test_api.py +++ b/products/tasks/backend/tests/test_api.py @@ -4901,6 +4901,23 @@ def test_partial_update_state_remove_keys_is_atomic(self): self.assertNotIn("pending_user_message", run.state) self.assertNotIn("pending_user_artifact_ids", run.state) + def _create_slack_mapping(self, task, run): + from posthog.models.integration import Integration + + from products.slack_app.backend.models import SlackThreadTaskMapping + + integration = Integration.objects.create(team=self.team, kind="slack", integration_id="T_SLACK", config={}) + return SlackThreadTaskMapping.objects.create( + team=self.team, + integration=integration, + slack_workspace_id="T_SLACK", + channel="C123", + thread_ts="1234.5678", + task=task, + task_run=run, + mentioning_slack_user_id="U123", + ) + @patch("products.tasks.backend.temporal.client.execute_posthog_code_agent_relay_workflow") def test_relay_message_enqueues_slack_relay_workflow(self, mock_execute_relay): from posthog.models.integration import Integration @@ -5007,10 +5024,17 @@ def test_relay_message_skips_when_no_slack_mapping(self, mock_execute_relay): self.assertEqual(response.json(), {"status": "skipped"}) mock_execute_relay.assert_not_called() + @parameterized.expand( + [ + ("failed", TaskRun.Status.FAILED), + ("cancelled", TaskRun.Status.CANCELLED), + ] + ) @patch("products.tasks.backend.temporal.client.execute_posthog_code_agent_relay_workflow") - def test_relay_message_skips_for_terminal_run(self, mock_execute_relay): + def test_relay_message_skips_for_failed_or_cancelled_run(self, _name, run_status, mock_execute_relay): task = self.create_task() - run = TaskRun.objects.create(task=task, team=self.team, status=TaskRun.Status.COMPLETED) + run = TaskRun.objects.create(task=task, team=self.team, status=run_status) + self._create_slack_mapping(task, run) response = self.client.post( f"/api/projects/@current/tasks/{task.id}/runs/{run.id}/relay_message/", @@ -5022,6 +5046,56 @@ def test_relay_message_skips_for_terminal_run(self, mock_execute_relay): self.assertEqual(response.json(), {"status": "skipped"}) mock_execute_relay.assert_not_called() + @patch("products.tasks.backend.temporal.client.execute_posthog_code_agent_relay_workflow") + def test_relay_message_accepts_completed_run(self, mock_execute_relay): + # A background (task-notification) turn can finish just after the + # inactivity timeout completed the run; its answer must still post. + task = self.create_task() + run = TaskRun.objects.create(task=task, team=self.team, status=TaskRun.Status.COMPLETED) + self._create_slack_mapping(task, run) + mock_execute_relay.return_value = "relay-late" + + response = self.client.post( + f"/api/projects/@current/tasks/{task.id}/runs/{run.id}/relay_message/", + {"text": "Late background answer"}, + format="json", + ) + + self.assertEqual(response.status_code, status.HTTP_200_OK) + self.assertEqual(response.json(), {"status": "accepted", "relay_id": "relay-late"}) + mock_execute_relay.assert_called_once_with( + run_id=str(run.id), + text="Late background answer", + delete_progress=True, + message_id=None, + ) + + @patch("products.tasks.backend.temporal.client.signal_agent_text_delta") + @patch("products.tasks.backend.temporal.client.execute_posthog_code_agent_relay_workflow") + def test_relay_message_completed_run_bypasses_agent_design_signal(self, mock_execute_relay, mock_signal): + # The workflow behind a terminal run is closed, so the agent-design + # inline-stream signal cannot be delivered; the relay workflow posts instead. + task = self.create_task() + run = TaskRun.objects.create( + task=task, + team=self.team, + status=TaskRun.Status.COMPLETED, + state={"slack_app_agent_design_enabled": True}, + ) + self._create_slack_mapping(task, run) + mock_execute_relay.return_value = "relay-late" + + response = self.client.post( + f"/api/projects/@current/tasks/{task.id}/runs/{run.id}/relay_message/", + {"text": "Late background answer"}, + format="json", + ) + + self.assertEqual(response.status_code, status.HTTP_200_OK) + self.assertEqual(response.json(), {"status": "accepted", "relay_id": "relay-late"}) + mock_signal.assert_not_called() + mock_execute_relay.assert_called_once() + @patch("products.tasks.backend.temporal.client.execute_posthog_code_agent_relay_workflow") def test_relay_message_rejects_blank_text(self, mock_execute_relay): task = self.create_task()