diff --git a/plugins/platforms/slack/adapter.py b/plugins/platforms/slack/adapter.py index b2c22bd55951..ef27b760d319 100644 --- a/plugins/platforms/slack/adapter.py +++ b/plugins/platforms/slack/adapter.py @@ -815,6 +815,14 @@ class SlackAdapter(BasePlatformAdapter): # cache bridges lifecycle and message delivery ordering. self._agent_view_contexts: Dict[Tuple[str, str], Dict[str, str]] = {} self._AGENT_VIEW_CONTEXTS_MAX = 5000 + # Status-bubble dedup (issue #30045, extended to Slack): remember the + # message ts of the last status bubble per (channel, thread, status + # key) so repeated progress callbacks (compression retries, fallback + # switches, ...) edit ONE message in place instead of appending a new + # bubble per event — long retry loops used to spam threads with + # dozens of out-of-order status messages. + self._status_message_ids: Dict[Tuple[str, str, str], str] = {} + self._STATUS_MESSAGE_IDS_MAX = 2000 # Cache for _fetch_thread_context results: cache_key → _ThreadContextCache self._thread_context_cache: Dict[str, _ThreadContextCache] = {} self._THREAD_CACHE_TTL = 60.0 @@ -2339,6 +2347,49 @@ class SlackAdapter(BasePlatformAdapter): logger.error("[Slack] Ephemeral send error: %s", e, exc_info=True) return SendResult(success=False, error=str(e)) + async def send_or_update_status( + self, + chat_id: str, + status_key: str, + content: str, + *, + metadata: Optional[Dict[str, Any]] = None, + ) -> SendResult: + """Send a status message, or edit the previous one with the same key. + + Issue #30045 (Telegram) extended to Slack: progress/status callbacks + (context-pressure, compression retries, model fallback, lifecycle) + used to append a fresh bubble on every call, spamming threads during + long retry loops. The first call posts and the message ts is + remembered; subsequent calls with the same (channel, thread, + status_key) edit that message in place via ``chat.update``. If the + edit fails (message deleted, too old, ...) the cached ts is dropped + and a fresh message is sent. + """ + thread_ts = self._resolve_thread_ts(None, metadata) or "" + key = (str(chat_id), str(thread_ts), str(status_key)) + cached_id = self._status_message_ids.get(key) + if cached_id is not None: + result = await self.edit_message( + chat_id, cached_id, content, finalize=False, metadata=metadata, + ) + if result.success: + if result.message_id: + self._status_message_ids[key] = str(result.message_id) + return result + # Edit failed — clear the cached ts and fall through to a fresh send. + self._status_message_ids.pop(key, None) + result = await self.send(chat_id, content, metadata=metadata) + if result.success and result.message_id: + if len(self._status_message_ids) >= self._STATUS_MESSAGE_IDS_MAX: + # Simple FIFO trim: drop the oldest half to bound memory. + for stale in list(self._status_message_ids)[ + : self._STATUS_MESSAGE_IDS_MAX // 2 + ]: + self._status_message_ids.pop(stale, None) + self._status_message_ids[key] = str(result.message_id) + return result + async def edit_message( self, chat_id: str, diff --git a/tests/gateway/test_slack_status_update.py b/tests/gateway/test_slack_status_update.py new file mode 100644 index 000000000000..841d863c51be --- /dev/null +++ b/tests/gateway/test_slack_status_update.py @@ -0,0 +1,144 @@ +"""Tests for SlackAdapter.send_or_update_status (issue #30045, Slack). + +The status-update path must: + 1. Send a fresh message on the first call for a (channel, thread, key). + 2. Edit that same message on subsequent calls with the same key. + 3. Fall back to sending fresh when the cached message edit fails. + 4. Keep distinct keys and distinct threads independent. +""" + +from __future__ import annotations + +import sys +from unittest.mock import AsyncMock, MagicMock + +import pytest + + +def _ensure_slack_mock(): + if "slack_bolt" in sys.modules and hasattr(sys.modules["slack_bolt"], "__file__"): + return + slack_bolt = MagicMock() + slack_bolt.async_app.AsyncApp = MagicMock + slack_bolt.adapter.socket_mode.async_handler.AsyncSocketModeHandler = MagicMock + slack_sdk = MagicMock() + slack_sdk.web.async_client.AsyncWebClient = MagicMock + for name, mod in [ + ("slack_bolt", slack_bolt), + ("slack_bolt.async_app", slack_bolt.async_app), + ("slack_bolt.adapter", slack_bolt.adapter), + ("slack_bolt.adapter.socket_mode", slack_bolt.adapter.socket_mode), + ( + "slack_bolt.adapter.socket_mode.async_handler", + slack_bolt.adapter.socket_mode.async_handler, + ), + ("slack_sdk", slack_sdk), + ("slack_sdk.web", slack_sdk.web), + ("slack_sdk.web.async_client", slack_sdk.web.async_client), + ]: + sys.modules.setdefault(name, mod) + sys.modules.setdefault("aiohttp", MagicMock()) + + +_ensure_slack_mock() + +import plugins.platforms.slack.adapter as _slack_mod # noqa: E402 + +_slack_mod.SLACK_AVAILABLE = True + +from gateway.config import PlatformConfig # noqa: E402 +from plugins.platforms.slack.adapter import SlackAdapter # noqa: E402 + + +@pytest.fixture() +def adapter(): + config = PlatformConfig(enabled=True, token="***") + a = SlackAdapter(config) + a._app = MagicMock() + client = AsyncMock() + client.chat_postMessage = AsyncMock( + side_effect=lambda **kw: {"ok": True, "ts": f"ts_{client.chat_postMessage.call_count}"} + ) + client.chat_update = AsyncMock(return_value={"ok": True}) + a._get_client = MagicMock(return_value=client) + a._bot_user_id = "U_BOT" + a._running = True + a.stop_typing = AsyncMock() + return a + + +METADATA = {"thread_id": "1784585355.415219"} + + +@pytest.mark.asyncio +async def test_first_call_sends_fresh(adapter): + result = await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing 1/3", metadata=METADATA + ) + assert result.success + client = adapter._get_client.return_value + assert client.chat_postMessage.call_count == 1 + assert client.chat_update.call_count == 0 + + +@pytest.mark.asyncio +async def test_second_call_edits_same_message(adapter): + r1 = await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing 1/3", metadata=METADATA + ) + r2 = await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing 2/3", metadata=METADATA + ) + assert r1.success and r2.success + client = adapter._get_client.return_value + assert client.chat_postMessage.call_count == 1 + assert client.chat_update.call_count == 1 + # The edit must target the ts of the first send. + assert client.chat_update.call_args.kwargs["ts"] == r1.message_id + + +@pytest.mark.asyncio +async def test_edit_failure_falls_back_to_fresh_send(adapter): + await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing 1/3", metadata=METADATA + ) + client = adapter._get_client.return_value + client.chat_update = AsyncMock(side_effect=RuntimeError("message_not_found")) + r2 = await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing 2/3", metadata=METADATA + ) + assert r2.success + assert client.chat_postMessage.call_count == 2 + # Cached id was replaced: a third call edits the NEW message. + client.chat_update = AsyncMock(return_value={"ok": True}) + r3 = await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing 3/3", metadata=METADATA + ) + assert r3.success + assert client.chat_update.call_args.kwargs["ts"] == r2.message_id + + +@pytest.mark.asyncio +async def test_distinct_keys_do_not_crosstalk(adapter): + await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing", metadata=METADATA + ) + await adapter.send_or_update_status( + "C_CHAN", "model_fallback", "falling back", metadata=METADATA + ) + client = adapter._get_client.return_value + assert client.chat_postMessage.call_count == 2 + assert client.chat_update.call_count == 0 + + +@pytest.mark.asyncio +async def test_distinct_threads_do_not_crosstalk(adapter): + await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing", metadata={"thread_id": "111.1"} + ) + await adapter.send_or_update_status( + "C_CHAN", "context_pressure", "compressing", metadata={"thread_id": "222.2"} + ) + client = adapter._get_client.return_value + assert client.chat_postMessage.call_count == 2 + assert client.chat_update.call_count == 0