diff --git a/gateway/run.py b/gateway/run.py index ccfa8e92c143e..764c23116d377 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -9101,6 +9101,21 @@ async def _handle_message(self, event: MessageEvent) -> Optional[str]: "Gateway intercepted clarify text response (session=%s, id=%s)", _quick_key, _pending_clarify.clarify_id, ) + # The clarify callback pauses the platform typing/status + # indicator while waiting so Slack users can type their + # answer. The active agent resumes as soon as this reply + # resolves the wait, so re-enable its indicator here too. + # Without this, Slack stays silent until the independent + # long-running heartbeat fires (three minutes by default). + _clarify_adapter = self._adapter_for_source(source) + if _clarify_adapter: + try: + _clarify_adapter.resume_typing_for_chat(source.chat_id) + except Exception: + logger.debug( + "Failed to resume typing after clarify response", + exc_info=True, + ) # Acknowledge with empty string so adapters that emit # the agent's response don't double-post. The agent # itself will produce the next user-facing message. diff --git a/plugins/platforms/slack/adapter.py b/plugins/platforms/slack/adapter.py index 05aeaf56f46aa..58870bed3d8dd 100644 --- a/plugins/platforms/slack/adapter.py +++ b/plugins/platforms/slack/adapter.py @@ -472,8 +472,14 @@ def __init__(self, config: PlatformConfig): # Track message IDs that should get reaction lifecycle (DMs / @mentions). self._reacting_message_ids: set = set() # Track active assistant thread status indicators so stop_typing can - # clear them (chat_id → thread_ts). - self._active_status_threads: Dict[str, str] = {} + # clear them. A single Slack channel can have several concurrent + # thread-bound Hermes turns; keying this by chat_id alone lets a newer + # thread overwrite an older thread, leaving the older assistant status + # stuck as "is thinking..." after final delivery. + self._active_status_threads: Dict[str, set[str]] = {} + # Defensive bound for permanent Slack API failures or lost lifecycle + # metadata. Normal concurrency is far below this threshold. + self._ACTIVE_STATUS_THREADS_PER_CHAT_MAX = 128 # Slash-command contexts: stash response_url + user_id so send() # can route the first reply ephemerally. Keyed by # (channel_id, user_id) to avoid cross-user collisions. @@ -1360,7 +1366,11 @@ async def send( if not self._app: return SendResult(success=False, error="Not connected") + thread_ts = None try: + # Resolve once so success and exception cleanup always target the + # same exact thread, including callers that provide reply_to only. + thread_ts = self._resolve_thread_ts(reply_to, metadata) # Check for a pending slash-command context. When the user ran a # native slash command (e.g. /q, /stop, /model), the initial ack # already showed an ephemeral "Running /cmd…" message. If we have @@ -1379,7 +1389,6 @@ async def send( # Split long messages, preserving code block boundaries chunks = self.truncate_message(formatted, self.MAX_MESSAGE_LENGTH) - thread_ts = self._resolve_thread_ts(reply_to, metadata) last_result = None # reply_broadcast: also post thread replies to the main channel. @@ -1409,9 +1418,11 @@ async def send( last_result = await self._get_client(chat_id).chat_postMessage(**kwargs) - # Clear Slack Assistant status as soon as the final message is posted. + # Clear Slack Assistant status for this thread as soon as the final + # message is posted. Preserve any other concurrently-active thread + # statuses in the same channel. if thread_ts: - await self.stop_typing(chat_id) + await self.stop_typing(chat_id, metadata={"thread_id": thread_ts}) # Track the sent message ts so we can auto-respond to thread # replies without requiring @mention. @@ -1434,6 +1445,20 @@ async def send( except Exception as e: # pragma: no cover - defensive logging logger.error("[Slack] Send error: %s", e, exc_info=True) + # A failed final post must not leave Slack's persistent Assistant + # status behind. Clear only this turn's thread; other concurrent + # threads in the same DM/channel must keep their own status. + if thread_ts: + try: + await self.stop_typing( + chat_id, + metadata={"thread_id": thread_ts}, + ) + except Exception: + logger.debug( + "[Slack] Failed to clear Assistant status after send error", + exc_info=True, + ) return SendResult(success=False, error=str(e)) async def send_private_notice( @@ -1479,6 +1504,7 @@ async def edit_message( content: str, *, finalize: bool = False, + metadata: Optional[Dict[str, Any]] = None, ) -> SendResult: """Edit a previously sent Slack message.""" if not self._app: @@ -1500,7 +1526,7 @@ async def edit_message( update_kwargs["blocks"] = blocks await self._get_client(chat_id).chat_update(**update_kwargs) if finalize: - await self.stop_typing(chat_id) + await self.stop_typing(chat_id, metadata=metadata) return SendResult(success=True, message_id=message_id) except Exception as e: # pragma: no cover - defensive logging logger.error( @@ -1529,7 +1555,17 @@ async def send_typing(self, chat_id: str, metadata=None) -> None: if not thread_ts: return # Can only set status in a thread context - self._active_status_threads[chat_id] = thread_ts + tracked = self._active_status_threads.setdefault(chat_id, set()) + thread_ts = str(thread_ts) + tracked.add(thread_ts) + while len(tracked) > self._ACTIVE_STATUS_THREADS_PER_CHAT_MAX: + # Keep the status being refreshed. Slack statuses also expire + # server-side, so dropping an older orphan is safer than allowing + # an unbounded local leak after permanent API/scope failures. + stale_thread_ts = next((ts for ts in tracked if ts != thread_ts), None) + if stale_thread_ts is None: + break + tracked.discard(stale_thread_ts) try: await self._get_client(chat_id).assistant_threads_setStatus( channel_id=chat_id, @@ -1545,17 +1581,44 @@ async def stop_typing(self, chat_id: str, metadata=None) -> None: """Clear the assistant thread status indicator.""" if not self._app: return - thread_ts = self._active_status_threads.pop(chat_id, None) - if not thread_ts: + + requested_thread_ts = None + if metadata: + requested_thread_ts = metadata.get("thread_id") or metadata.get("thread_ts") + + active = self._active_status_threads.get(chat_id) + if not active: return - try: - await self._get_client(chat_id).assistant_threads_setStatus( - channel_id=chat_id, - thread_ts=thread_ts, - status="", - ) - except Exception as e: - logger.debug("[Slack] assistant.threads.setStatus clear failed: %s", e) + + if isinstance(active, str): + # Backward compatibility for tests or long-lived adapters created + # before the per-thread tracking map shape existed. + legacy_thread_ts = str(active) + active = {legacy_thread_ts} + self._active_status_threads[chat_id] = active + + if requested_thread_ts: + thread_ts_values = [str(requested_thread_ts)] + else: + thread_ts_values = list(active) + + for thread_ts in thread_ts_values: + try: + await self._get_client(chat_id).assistant_threads_setStatus( + channel_id=chat_id, + thread_ts=thread_ts, + status="", + ) + except Exception as e: + # Slack still owns a visible status when this call fails. Keep + # the thread tracked so later cleanup attempts can retry rather + # than forgetting a server-side status that remains visible. + logger.debug("[Slack] assistant.threads.setStatus clear failed: %s", e) + continue + active.discard(thread_ts) + + if not active: + self._active_status_threads.pop(chat_id, None) def _dm_top_level_threads_as_sessions(self) -> bool: """Whether top-level Slack DMs get per-message session threads. diff --git a/tests/gateway/test_clarify_active_session_bypass.py b/tests/gateway/test_clarify_active_session_bypass.py index bfa5ffff5617d..0f5e6614251f3 100644 --- a/tests/gateway/test_clarify_active_session_bypass.py +++ b/tests/gateway/test_clarify_active_session_bypass.py @@ -1,7 +1,7 @@ """Regression tests for clarify replies while a gateway session is busy.""" import asyncio -from unittest.mock import AsyncMock +from unittest.mock import AsyncMock, patch import pytest @@ -85,3 +85,37 @@ async def test_active_session_routes_typed_choice_clarify_reply_to_runner_not_bu adapter._message_handler.assert_awaited_once_with(event) adapter._busy_session_handler.assert_not_awaited() assert adapter._pending_messages == {} + + +@pytest.mark.asyncio +async def test_gateway_clarify_reply_resumes_typing_before_returning_empty_ack(): + """A clarify answer must re-enable the active run's typing indicator. + + Clarify pauses typing while waiting so Slack's Assistant API does not + disable the compose box. The typed answer is intercepted by the gateway + and returns an empty acknowledgment instead of starting a second run; that + interception path must therefore resume the original run's indicator. + """ + _clear_clarify_state() + from gateway.run import GatewayRunner + from tools import clarify_gateway as cm + + adapter = _ClarifyBypassAdapter() + adapter.pause_typing_for_chat("12345") + event = _event("the missing details") + + runner = GatewayRunner.__new__(GatewayRunner) + runner._startup_restore_in_progress = False + runner._scale_to_zero_note_real_inbound = lambda: None + runner._is_user_authorized = lambda source: True + runner._session_key_for_source = lambda source: "clarify-session" + runner._adapter_for_source = lambda source: adapter + runner._update_prompt_pending = {} + + cm.register("clarify-2", "clarify-session", "What is missing?", None) + + with patch("hermes_cli.plugins.invoke_hook", return_value=[]): + result = await runner._handle_message(event) + + assert result == "" + assert "12345" not in adapter._typing_paused diff --git a/tests/gateway/test_slack.py b/tests/gateway/test_slack.py index de3edf2db272e..a26eb20c21d97 100644 --- a/tests/gateway/test_slack.py +++ b/tests/gateway/test_slack.py @@ -2028,12 +2028,62 @@ class TestSendTyping: async def test_sets_status_in_thread(self, adapter): adapter._app.client.assistant_threads_setStatus = AsyncMock() await adapter.send_typing("C123", metadata={"thread_id": "parent_ts"}) + assert adapter._active_status_threads == {"C123": {"parent_ts"}} adapter._app.client.assistant_threads_setStatus.assert_called_once_with( channel_id="C123", thread_ts="parent_ts", status="is thinking...", ) + @pytest.mark.asyncio + async def test_status_tracking_is_per_thread_not_per_channel(self, adapter): + adapter._app.client.assistant_threads_setStatus = AsyncMock() + + await adapter.send_typing("C123", metadata={"thread_id": "111.000"}) + await adapter.send_typing("C123", metadata={"thread_id": "222.000"}) + + assert adapter._active_status_threads == {"C123": {"111.000", "222.000"}} + + await adapter.stop_typing("C123", metadata={"thread_id": "111.000"}) + + assert adapter._active_status_threads == {"C123": {"222.000"}} + adapter._app.client.assistant_threads_setStatus.assert_any_await( + channel_id="C123", + thread_ts="111.000", + status="", + ) + + @pytest.mark.asyncio + async def test_active_status_threads_bounded_per_chat(self, adapter): + adapter._app.client.assistant_threads_setStatus = AsyncMock( + side_effect=Exception("missing_scope") + ) + adapter._ACTIVE_STATUS_THREADS_PER_CHAT_MAX = 4 + + for i in range(20): + await adapter.send_typing("C123", metadata={"thread_id": f"ts_{i}"}) + + tracked = adapter._active_status_threads["C123"] + assert len(tracked) <= 4 + assert "ts_19" in tracked + + @pytest.mark.asyncio + async def test_stop_typing_without_metadata_clears_all_channel_threads(self, adapter): + adapter._app.client.assistant_threads_setStatus = AsyncMock() + await adapter.send_typing("C123", metadata={"thread_id": "111.000"}) + await adapter.send_typing("C123", metadata={"thread_id": "222.000"}) + adapter._app.client.assistant_threads_setStatus.reset_mock() + + await adapter.stop_typing("C123") + + assert "C123" not in adapter._active_status_threads + cleared = { + call.kwargs["thread_ts"] + for call in adapter._app.client.assistant_threads_setStatus.await_args_list + if call.kwargs.get("status") == "" + } + assert cleared == {"111.000", "222.000"} + @pytest.mark.asyncio async def test_noop_without_thread(self, adapter): adapter._app.client.assistant_threads_setStatus = AsyncMock() @@ -2096,7 +2146,72 @@ async def test_stop_typing_handles_api_error_gracefully(self, adapter): thread_ts="parent_ts", status="", ) - assert "C123" not in adapter._active_status_threads + assert adapter._active_status_threads["C123"] == {"parent_ts"} + + @pytest.mark.asyncio + async def test_stop_typing_partial_failure_keeps_only_failed_thread_tracked(self, adapter): + adapter._active_status_threads["C123"] = {"thread_ok", "thread_fail"} + + async def set_status(*, channel_id, thread_ts, status): + if thread_ts == "thread_fail": + raise Exception("temporary Slack failure") + + adapter._app.client.assistant_threads_setStatus = AsyncMock( + side_effect=set_status + ) + + await adapter.stop_typing("C123") + + assert adapter._active_status_threads["C123"] == {"thread_fail"} + cleared = { + call.kwargs["thread_ts"] + for call in adapter._app.client.assistant_threads_setStatus.await_args_list + } + assert cleared == {"thread_ok", "thread_fail"} + + @pytest.mark.asyncio + async def test_send_failure_attempts_exact_thread_status_clear(self, adapter): + adapter._app.client.chat_postMessage = AsyncMock( + side_effect=Exception("Slack post failed") + ) + adapter._app.client.assistant_threads_setStatus = AsyncMock() + adapter._active_status_threads["C123"] = {"111.000", "222.000"} + + result = await adapter.send( + "C123", + "done", + metadata={"thread_id": "111.000"}, + ) + + assert not result.success + adapter._app.client.assistant_threads_setStatus.assert_called_once_with( + channel_id="C123", + thread_ts="111.000", + status="", + ) + assert adapter._active_status_threads == {"C123": {"222.000"}} + + @pytest.mark.asyncio + async def test_send_failure_with_reply_to_clears_only_that_thread(self, adapter): + adapter._app.client.chat_postMessage = AsyncMock( + side_effect=Exception("Slack post failed") + ) + adapter._app.client.assistant_threads_setStatus = AsyncMock() + adapter._active_status_threads["C123"] = {"111.000", "222.000"} + + result = await adapter.send( + "C123", + "done", + reply_to="111.000", + ) + + assert not result.success + adapter._app.client.assistant_threads_setStatus.assert_called_once_with( + channel_id="C123", + thread_ts="111.000", + status="", + ) + assert adapter._active_status_threads == {"C123": {"222.000"}} @pytest.mark.asyncio async def test_send_clears_status_after_final_post(self, adapter): @@ -2104,7 +2219,7 @@ async def test_send_clears_status_after_final_post(self, adapter): return_value={"ts": "reply_ts"} ) adapter._app.client.assistant_threads_setStatus = AsyncMock() - adapter._active_status_threads["C123"] = "parent_ts" + adapter._active_status_threads["C123"] = {"parent_ts"} result = await adapter.send("C123", "done", metadata={"thread_id": "parent_ts"}) @@ -2121,7 +2236,7 @@ async def test_send_clears_status_after_final_post(self, adapter): async def test_streaming_final_edit_clears_status(self, adapter): adapter._app.client.chat_update = AsyncMock() adapter._app.client.assistant_threads_setStatus = AsyncMock() - adapter._active_status_threads["C123"] = "parent_ts" + adapter._active_status_threads["C123"] = {"parent_ts"} result = await adapter.edit_message( "C123", @@ -2147,7 +2262,7 @@ async def test_streaming_final_edit_clears_status(self, adapter): async def test_streaming_intermediate_edit_keeps_status(self, adapter): adapter._app.client.chat_update = AsyncMock() adapter._app.client.assistant_threads_setStatus = AsyncMock() - adapter._active_status_threads["C123"] = "parent_ts" + adapter._active_status_threads["C123"] = {"parent_ts"} result = await adapter.edit_message( "C123", @@ -2158,7 +2273,29 @@ async def test_streaming_intermediate_edit_keeps_status(self, adapter): assert result.success adapter._app.client.assistant_threads_setStatus.assert_not_called() - assert adapter._active_status_threads["C123"] == "parent_ts" + assert adapter._active_status_threads["C123"] == {"parent_ts"} + + @pytest.mark.asyncio + async def test_streaming_final_edit_with_metadata_preserves_other_threads(self, adapter): + adapter._app.client.chat_update = AsyncMock() + adapter._app.client.assistant_threads_setStatus = AsyncMock() + adapter._active_status_threads["C123"] = {"111.000", "222.000"} + + result = await adapter.edit_message( + "C123", + "reply_ts", + "done", + finalize=True, + metadata={"thread_id": "111.000"}, + ) + + assert result.success + adapter._app.client.assistant_threads_setStatus.assert_called_once_with( + channel_id="C123", + thread_ts="111.000", + status="", + ) + assert adapter._active_status_threads == {"C123": {"222.000"}} # ---------------------------------------------------------------------------