Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
97 changes: 80 additions & 17 deletions plugins/platforms/slack/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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.
Expand All @@ -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(
Expand Down Expand Up @@ -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:
Expand All @@ -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(
Expand Down Expand Up @@ -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,
Expand All @@ -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.
Expand Down
36 changes: 35 additions & 1 deletion tests/gateway/test_clarify_active_session_bypass.py
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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
Loading