Skip to content
38 changes: 35 additions & 3 deletions products/slack_app/backend/slack_thread.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,8 @@
# answer, so the reply is posted plainly instead. The same pair is what the scout delivery in
# signals treats as a block rejection.
_BLOCK_REJECTION_ERROR_CODES = frozenset({"invalid_blocks", "invalid_blocks_format"})
# Slack closed the stream, so every later append and the stop call fail the same way.
_STREAM_ENDED_ERROR_CODE = "message_not_in_streaming_state"


def _split_markdown_text(text: str, limit: int = _MARKDOWN_CHUNK_LIMIT) -> list[str]:
Expand Down Expand Up @@ -224,6 +226,7 @@
self._client: WebClient | None = None
self._bot_user_id: str | None = None
self._fork_flag: bool | None = None
self.stream_ended = False

@classmethod
def for_run(
Expand Down Expand Up @@ -316,6 +319,8 @@
One append per block, and both after the answer's: a request Slack rejects must
cost that control alone, never the reply and never its sibling.
"""
if self.stream_ended:
return
for block, failure in (
(self._fork_menu_actions_block(), "slack_app_fork_menu_append_failed"),
(self._feedback_block(), "slack_app_feedback_buttons_append_failed"),
Expand Down Expand Up @@ -440,12 +445,13 @@
task_updates: list[dict[str, Any]] | None = None,
markdown_text: str | None = None,
plan_title: str | None = None,
) -> None:
"""Append plan-block step transitions and/or markdown_text chunks."""
) -> bool:
"""Append plan-block step transitions and/or markdown_text chunks. Returns whether the stream is still open."""
chunks = _status_chunks(task_updates, markdown_text)
if plan_title:
chunks.insert(0, _plan_update_chunk(plan_title))
self._append_chunks(ts, chunks, "slack_app_status_stream_append_failed")
return not self.stream_ended

def append_status_blocks(self, ts: str, blocks: list[dict[str, Any]]) -> bool:
"""Append Block Kit blocks, such as chart cards, to an open stream. Returns whether Slack took them."""
Expand All @@ -455,7 +461,7 @@
ts, [{"type": "blocks", "blocks": blocks}], "slack_app_status_stream_blocks_append_failed"
)

def stop_status_stream(

Check warning on line 464 in products/slack_app/backend/slack_thread.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`stop_status_stream` has cyclomatic complexity 12 (warn >10)

Check warning on line 464 in products/slack_app/backend/slack_thread.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`stop_status_stream` has cyclomatic complexity 12 (warn >10)
self,
ts: str,
complete_task_id: str | None = None,
Expand All @@ -471,7 +477,10 @@
answer to stream here, the mention closes the message instead, unless ``mention_sent``
says the answer already carried it. ``append_attachments`` runs after the answer, so
chart cards sit under the text that describes them. The provenance footer is a `blocks`
chunk because a `context` block is the only way to get muted text."""
chunk because a `context` block is the only way to get muted text.

When Slack already closed the stream, the answer goes out as a plain thread reply,
and nothing else is sent to the closed stream."""
answer_chunks: list[dict[str, Any]] = []
if plan_title:
answer_chunks.append(_plan_update_chunk(plan_title))
Expand All @@ -481,6 +490,10 @@
for piece in _markdown_text_pieces(self._with_leading_mention(final_markdown)):
answer_chunks.append({"type": "markdown_text", "text": piece})
self._append_chunks(ts, answer_chunks, "slack_app_status_stream_final_append_failed")
if self.stream_ended:
if final_markdown:
self._post_answer_outside_stream(final_markdown)
return
if append_attachments is not None:
try:
append_attachments()
Expand All @@ -498,6 +511,11 @@
self._append_chunks(ts, final_chunks, "slack_app_status_stream_final_append_failed")
if footer:
self._append_trailing_blocks(ts)
self._stop_stream(ts)

def _stop_stream(self, ts: str) -> None:
if self.stream_ended:
return
try:
self._get_client().chat_stopStream(
channel=self.context.channel,
Expand All @@ -506,6 +524,11 @@
except Exception as e:
logger.warning("slack_app_status_stream_stop_failed", error=str(e))

def _post_answer_outside_stream(self, final_markdown: str) -> None:
pieces = _markdown_text_pieces(self._with_leading_mention(final_markdown))
for index, piece in enumerate(pieces):
self.post_thread_message(piece, with_footer=index == len(pieces) - 1, markdown=True)

def _with_leading_mention(self, markdown: str) -> str:
return leading_mention_prefix(markdown, self.actor_slack_user_id) + markdown

Expand All @@ -523,8 +546,17 @@
def _append_chunks(self, ts: str, chunks: list[dict[str, Any]], failure_event: str) -> bool:
if not chunks:
return True
if self.stream_ended:
return False
try:
self._get_client().chat_appendStream(channel=self.context.channel, ts=ts, chunks=chunks)
except SlackApiError as e:
if e.response.get("error") == _STREAM_ENDED_ERROR_CODE:
self.stream_ended = True
logger.info("slack_app_status_stream_ended_by_slack", channel=self.context.channel)
else:
logger.warning(failure_event, error=str(e))
return False
except Exception as e:
logger.warning(failure_event, error=str(e))
return False
Expand Down
29 changes: 29 additions & 0 deletions products/slack_app/backend/tests/test_slack_thread.py
Original file line number Diff line number Diff line change
Expand Up @@ -540,6 +540,35 @@ def test_a_rejected_footer_reposts_the_answer_as_plain_text(self, mock_get_clien
assert not retry.get("blocks")


class TestStreamClosedBySlack(SimpleTestCase):
@patch.object(SlackThreadHandler, "_get_integration")
@patch.object(SlackThreadHandler, "_get_client")
def test_the_answer_is_posted_in_the_thread_and_the_closed_stream_gets_nothing_more(
self, mock_get_client, mock_get_integration
) -> None:
mock_client = MagicMock()
mock_client.chat_appendStream.side_effect = SlackApiError(
"message_not_in_streaming_state", {"error": "message_not_in_streaming_state"}
)
mock_get_client.return_value = mock_client
mock_get_integration.return_value = Integration(id=1, config={}, integration_id="T1")
context = SlackThreadContext(integration_id=1, channel="C001", thread_ts="1234.5678")
handler = SlackThreadHandler(context, RunFooter(model="claude-opus-5"), actor_slack_user_id="U123")

assert (
handler.append_status_chunks(ts="1.0", task_updates=[{"id": "a", "title": "Read", "status": "in_progress"}])
is False
)
handler.stop_status_stream(ts="1.0", final_markdown="Signups grew.")

assert mock_client.chat_appendStream.call_count == 1
mock_client.chat_stopStream.assert_not_called()
posted = mock_client.chat_postMessage.call_args.kwargs
assert posted["thread_ts"] == "1234.5678"
assert "Signups grew." in posted["text"]
assert posted["text"].startswith("<@U123>")


class TestRelayedAnswerFooter(SimpleTestCase):
def _handler(self, footer: RunFooter) -> SlackThreadHandler:
context = SlackThreadContext(integration_id=1, channel="C001", thread_ts="1234.5678")
Expand Down
6 changes: 3 additions & 3 deletions products/tasks/backend/facade/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -3179,7 +3179,7 @@


def update_task_run(
run_id: str | UUID,

Check warning on line 3182 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`update_task_run` has cyclomatic complexity 38 (warn >10)

Check warning on line 3182 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`update_task_run` has cyclomatic complexity 38 (warn >10)
task_id: str | UUID,
team_id: int,
*,
Expand Down Expand Up @@ -4207,7 +4207,7 @@


def finalize_task_run_artifact_uploads(
run_id: str | UUID,

Check warning on line 4210 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`finalize_task_run_artifact_uploads` has cyclomatic complexity 13 (warn >10)

Check warning on line 4210 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`finalize_task_run_artifact_uploads` has cyclomatic complexity 13 (warn >10)
task_id: str | UUID,
team_id: int,
*,
Expand Down Expand Up @@ -4998,7 +4998,7 @@
if run is None:
return None
ensure_subscription_owner(run.state, actor_user_id)
if run.is_terminal or (run.state or {}).get("cancel_requested_at"):

Check warning on line 5001 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`signal_task_run_user_message` has cyclomatic complexity 13 (warn >10)

Check warning on line 5001 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`signal_task_run_user_message` has cyclomatic complexity 13 (warn >10)
if not run.is_terminal:
raise RuntimeError("Task run is still stopping. Try again shortly.")
sandbox_id = (run.state or {}).get("sandbox_id")
Expand Down Expand Up @@ -5406,7 +5406,7 @@
probe = sandbox.execute(_PREVIEW_HEALTH_PROBE, timeout_seconds=_PREVIEW_HEALTH_PROBE_TIMEOUT_SECONDS)
except SandboxNotFoundError:
return _PREVIEW_ENDED
except Exception:

Check warning on line 5409 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`resolve_task_run_preview_redirect` has cyclomatic complexity 12 (warn >10)

Check warning on line 5409 in products/tasks/backend/facade/api.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`resolve_task_run_preview_redirect` has cyclomatic complexity 12 (warn >10)
logger.exception("task_run_preview_health_probe_failed", extra={"run_id": str(run.id)})
return _PREVIEW_UNAVAILABLE
if probe.exit_code != 0:
Expand Down Expand Up @@ -5528,7 +5528,7 @@
trace_id: str | None = None,
) -> tuple[str, str | None]:
"""Queue a Slack relay workflow for a run message, or under the agent-design
flag signal the running task workflow to stream the text inline.
flag give the running task workflow the text as the turn's final answer.

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
Expand All @@ -5548,7 +5548,7 @@
)
from products.tasks.backend.temporal.client import ( # noqa: PLC0415 — keep temporalio off the api import path
execute_posthog_code_agent_relay_workflow,
signal_agent_text_delta,
signal_agent_final_text,
)
from products.tasks.backend.temporal.process_task.activities.feature_flags import ( # noqa: PLC0415 — keep temporal off the api import path
AGENT_DESIGN_STATE_KEY,
Expand All @@ -5567,7 +5567,7 @@

if bool((run.state or {}).get(AGENT_DESIGN_STATE_KEY)):
try:
signal_agent_text_delta(run.workflow_id, trimmed)
signal_agent_final_text(run.workflow_id, trimmed, trace_id)
except Exception:
logger.exception("task_run_relay_text_signal_failed", extra={"run_id": str(run.id)})
return "skipped", None
Expand Down
9 changes: 6 additions & 3 deletions products/tasks/backend/temporal/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -655,11 +655,14 @@ async def signal() -> None:
asyncio.run(signal())


def signal_agent_text_delta(workflow_id: str, text: str) -> None:
"""Push text into the live agent-design plan-block stream for a running task."""
def signal_agent_final_text(workflow_id: str, text: str, trace_id: str | None = None) -> None:
"""Give the live agent-design reply of a running task the turn's whole answer.

The relay replaces its streamed text deltas with this text, so the answer shows one time
even when this signal lands between two deltas."""
client = sync_connect()
handle = client.get_workflow_handle(workflow_id)
asyncio.run(handle.signal("agent_text_delta", text))
asyncio.run(handle.signal("agent_final_text", {"text": text, "trace_id": trace_id}))


def execute_posthog_code_agent_relay_workflow(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,20 +131,21 @@ def start_slack_agent_design_stream(input: StartSlackAgentDesignStreamInput) ->

@activity.defn
@close_db_connections
def append_slack_agent_design_steps(input: AppendSlackAgentDesignStepsInput) -> None:
"""Append plan-block step transitions and a new plan title."""
def append_slack_agent_design_steps(input: AppendSlackAgentDesignStepsInput) -> bool:
"""Append plan-block step transitions and a new plan title. Returns False once Slack has closed the stream."""
from products.slack_app.backend.slack_thread import SlackThreadContext, SlackThreadHandler

try:
context = SlackThreadContext.from_dict(input.slack_thread_context)
handler = SlackThreadHandler(context)
handler.append_status_chunks(
return handler.append_status_chunks(
ts=input.ts,
task_updates=_chunk_dicts(input.task_updates),
plan_title=input.plan_title,
)
except Exception as e:
logger.warning("slack_app_append_agent_design_steps_failed", error=str(e))
return True


@activity.defn
Expand All @@ -155,6 +156,7 @@ def stop_slack_agent_design_stream(input: StopSlackAgentDesignStreamInput) -> No
from products.tasks.backend.logic.services.living_artifacts import (
SlackFileDeliveryResult,
attach_streamed_slack_files,
deliver_pending_slack_file_artifacts,
stream_pending_slack_attachments,
)
from products.tasks.backend.models import TaskRun
Expand Down Expand Up @@ -193,5 +195,8 @@ def _append_attachments() -> None:
attach_streamed_slack_files(
task_run, delivery, attach_files=lambda file_ids: handler.attach_files(input.ts, file_ids)
)
if handler.stream_ended and task_run is not None:
# A closed stream takes no cards, so whatever is still pending posts under the answer as its own message.
deliver_pending_slack_file_artifacts(task_run)
except Exception as e:
logger.warning("slack_app_stop_agent_design_stream_failed", error=str(e))
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,9 @@

The reply is a plan block and the final answer. The plan shows one line per kind of work (see
``products.slack_app.backend.logic.progress_phases``), or the agent's todo list when it keeps one. Each line lists the
descriptions of its calls. The agent's prose does not stream: the last burst before the turn
ends is the answer. While the turn runs, one line always spins and the plan title says what
happens now.
descriptions of its calls. The agent's prose does not stream: the answer is the final text the
agent server reports for the turn, or the last burst before the turn ends when that text does not
arrive in time. While the turn runs, one line always spins and the plan title says what happens now.

The first turn's relay starts before the sandbox exists, so the plan shows the setup steps the
parent sends through ``setup_step``. Slack ends a stream that gets no update for a few minutes,
Expand Down Expand Up @@ -59,6 +59,13 @@
_ACTIVITY_RETRY = RetryPolicy(maximum_attempts=3)
# The parent's progress step that the first setup line shows.
_SANDBOX_SETUP_STEP = "sandbox"
_PATCH_ID_STREAM_ENDED = "tasks-slack-relay-stream-ended"


def _trace_key(trace_id: Optional[str]) -> Optional[str]:
"""The trace id in one form. The relay endpoint sends it as a hyphenated UUID, and the turn-complete
event sends the W3C form, which is the same 32 hex digits without hyphens."""
return trace_id.replace("-", "").lower() if trace_id else None


@frozen
Expand Down Expand Up @@ -97,7 +104,15 @@ def __init__(self) -> None:
# Prose since the last tool call, and the last non-empty burst before it.
self._narrative: str = ""
self._last_burst: str = ""
# The whole answer as the agent server reports it. It can arrive between two text deltas,
# so it is kept apart from the deltas and replaces them at close.
self._final_text: str = ""
self._final_text_trace_id: Optional[str] = None
# Whether this turn's agent sent a tool call or prose yet.
self._turn_has_activity: bool = False
self._stream: Optional[SlackAgentDesignStream] = None
# Slack closed the stream early. Later appends to it can only fail.
self._stream_ended: bool = False
self._last_dispatched_at: float = 0.0
self._last_signal_at: Optional[datetime] = None
self._started_at: datetime = datetime.min
Expand All @@ -111,6 +126,7 @@ def __init__(self) -> None:
async def agent_status_update(self, payload: dict[str, Any] | str) -> None:
"""A tool call or a new agent todo list. Either ends the prose burst before it."""
self._last_signal_at = workflow.now()
self._turn_has_activity = True
narrative, self._narrative = self._narrative, ""
if narrative.strip():
self._last_burst = narrative
Expand Down Expand Up @@ -203,6 +219,19 @@ async def agent_text_delta(self, text: str) -> None:
self._last_signal_at = workflow.now()
if isinstance(text, str) and text:
self._narrative += text
self._turn_has_activity = True

@workflow.signal
async def agent_final_text(self, payload: dict[str, Any]) -> None:
# The agent server sends a turn's final text after the turn ends. A late one can reach the
# relay of the next turn, which has no activity yet, and must not become its answer.
if not self._turn_has_activity:
return
text = payload.get("text")
if isinstance(text, str) and text.strip():
self._final_text = text.strip()
trace_id = payload.get("trace_id")
self._final_text_trace_id = trace_id if isinstance(trace_id, str) else None

@workflow.signal
async def complete_turn(self, trace_id: str | None = None) -> None:
Expand Down Expand Up @@ -289,6 +318,11 @@ def _closing_plan_title(self) -> Optional[str]:
return done_plan_title(workflow.now() - self._started_at)

def _final_answer(self) -> str:
# A trace id that differs from this turn's means the text is a late answer of an earlier turn.
final_trace, turn_trace = _trace_key(self._final_text_trace_id), _trace_key(self._trace_id)
stale = bool(final_trace and turn_trace and final_trace != turn_trace)
if self._final_text and not stale:
return self._final_text
return (self._narrative if self._narrative.strip() else self._last_burst).strip()

def _idle_for(self) -> timedelta:
Expand All @@ -311,7 +345,9 @@ async def _append(
self, input: SlackAgentDesignRelayInput, chunks: list[TaskUpdateChunk], plan_title: Optional[str] = None
) -> None:
assert self._stream is not None
await workflow.execute_activity(
if self._stream_ended:
return
stream_open = await workflow.execute_activity(
append_slack_agent_design_steps,
AppendSlackAgentDesignStepsInput(
slack_thread_context=input.slack_thread_context,
Expand All @@ -322,6 +358,10 @@ async def _append(
start_to_close_timeout=_ACTIVITY_TIMEOUT,
retry_policy=_ACTIVITY_RETRY,
)
# An older activity returns None, which means the stream is still open. The patch keeps
# histories that older workflow code wrote on their recorded command sequence.
if stream_open is False and workflow.patched(_PATCH_ID_STREAM_ENDED):
self._stream_ended = True

@workflow.run
async def run(self, input: SlackAgentDesignRelayInput) -> None:
Expand Down Expand Up @@ -398,6 +438,14 @@ async def _close_stream(self, input: SlackAgentDesignRelayInput) -> None:
self._finish_setup(final_status)
# Lines that never reached Slack because the turn ended inside the debounce window.
pending = self._take_pending_chunks(allow_placeholder=False)
if self._stream_ended:
# The plan is gone with the closed stream, so the answer opens a new message of its own.
self._stream = None
pending = []
self._line_ids = {}
self._agent_plan = []
self._current_key = None
self._placeholder = None
if self._stream is None and (final_answer or pending):
# A turn with no flushed step still streams its answer in a stream of its own.
self._stream = await self._start_stream(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
from django.test import SimpleTestCase, TestCase, override_settings

from parameterized import parameterized
from slack_sdk.errors import SlackApiError

from posthog.models.integration import Integration
from posthog.models.organization import Organization
Expand Down Expand Up @@ -118,6 +119,26 @@ def test_closing_the_stream_hands_the_reply_the_turns_trace_id(self, mock_stop)

assert mock_stop.call_args.args[0].turn_trace_id == trace_id

@patch("products.tasks.backend.logic.services.living_artifacts.deliver_pending_slack_file_artifacts")
@patch.object(SlackThreadHandler, "_get_client")
def test_a_stream_slack_closed_still_delivers_the_turns_attachments(self, mock_get_client, mock_deliver) -> None:
client = mock_get_client.return_value
client.chat_appendStream.side_effect = SlackApiError(
"message_not_in_streaming_state", {"error": "message_not_in_streaming_state"}
)

stop_slack_agent_design_stream(
StopSlackAgentDesignStreamInput(
slack_thread_context={"integration_id": self.integration.id, "channel": "C1", "thread_ts": "1.0"},
ts="2.0",
final_markdown="Signups grew.",
run_id=str(self.task_run.id),
)
)

assert "Signups grew." in client.chat_postMessage.call_args.kwargs["text"]
mock_deliver.assert_called_once_with(self.task_run)

@parameterized.expand(
[
("run_actor", None, "U456"),
Expand Down
Loading
Loading