mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-24 16:54:43 +00:00
fix(slack): edit status bubbles in place instead of posting new ones
Progress/status callbacks (context-pressure, compression retries,
model fallback) route through _send_or_update_status_coro, which
edits the previous bubble for the same status_key when the adapter
implements send_or_update_status — but only Telegram did. On Slack
every status event posted a fresh thread message, so a compression
retry loop spammed a dozen out-of-order bubbles into the thread
('Context too large 1/3... 2/3... 3/3', fallback switches, etc.).
Implement send_or_update_status on the Slack adapter following the
Telegram pattern (#30045): first call posts and caches the message ts
per (channel, thread, status_key); subsequent calls edit that message
via chat.update. Edit failure drops the cached ts and falls back to a
fresh send. Cache is FIFO-bounded.
This commit is contained in:
parent
93a47dd466
commit
f716b876a8
2 changed files with 195 additions and 0 deletions
|
|
@ -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,
|
||||
|
|
|
|||
144
tests/gateway/test_slack_status_update.py
Normal file
144
tests/gateway/test_slack_status_update.py
Normal file
|
|
@ -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
|
||||
Loading…
Add table
Add a link
Reference in a new issue