diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index 74f78c9725..5b97554620 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -456,6 +456,15 @@ class TelegramAdapter(BasePlatformAdapter): self._telegram_typing_cooldown_until: Dict[str, float] = {} self._telegram_typing_cooldown_seconds: float = self._coerce_float_extra( "typing_cooldown_seconds", 30.0, min_value=1.0, max_value=300.0) + # Post-send typing re-arm: scheduled, deduped and rate-limited per chat. Awaiting a + # sendChatAction round-trip on the send path shares the loop with the getUpdates long-polls, + # and under concurrent streaming it starved them until they rotted into CLOSE-WAIT (#111727). + self._telegram_typing_retrigger_tasks: Dict[str, asyncio.Task] = {} + self._telegram_typing_retrigger_at: Dict[str, float] = {} + # Telegram's bubble lasts ~5s and _keep_typing already refreshes every 2s, so the re-arm only + # has to cover the gap left by a landed message. 0 restores a call per intermediate send. + self._telegram_typing_retrigger_interval: float = self._coerce_float_extra( + "typing_retrigger_min_interval_seconds", 2.0, min_value=0.0, max_value=30.0) # Buffer album/photo bursts into a single MessageEvent instead of self-interrupting turns. self._media_batch_delay_seconds = env_float("HERMES_TELEGRAM_MEDIA_BATCH_DELAY_SECONDS", 0.8) self._pending_photo_batches: Dict[str, MessageEvent] = {} @@ -3383,14 +3392,74 @@ class TelegramAdapter(BasePlatformAdapter): return _flood_cap_result(wait) raise - async def _retrigger_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]]) -> None: - """Re-arm typing after an intermediate send (Telegram clears it when a message lands). Skipped on - the FINAL reply (``metadata["notify"]``): the refresh loop is gone and no API cancels the bubble.""" - if (metadata or {}).get("notify"): - return + def _typing_retrigger_state(self) -> tuple[Dict[str, "asyncio.Task"], Dict[str, float], float]: + """Re-arm bookkeeping, materialised on demand — tests build adapters via ``object.__new__()`` + (no ``__init__``), so these attributes are not guaranteed to exist.""" + tasks = getattr(self, "_telegram_typing_retrigger_tasks", None) + if not isinstance(tasks, dict): + tasks = self._telegram_typing_retrigger_tasks = {} + sent_at = getattr(self, "_telegram_typing_retrigger_at", None) + if not isinstance(sent_at, dict): + sent_at = self._telegram_typing_retrigger_at = {} + try: + interval = float(getattr(self, "_telegram_typing_retrigger_interval", 2.0)) + except (TypeError, ValueError): + interval = 2.0 + return tasks, sent_at, interval + + def _clear_typing_retrigger(self, chat_id: str, finished: "asyncio.Task") -> None: + """Drop the in-flight slot, but only if it is still this task's (a newer one may own it).""" + tasks = getattr(self, "_telegram_typing_retrigger_tasks", None) + if isinstance(tasks, dict) and tasks.get(chat_id) is finished: + tasks.pop(chat_id, None) + + async def _send_typing_quietly(self, chat_id: str, metadata: Optional[Dict[str, Any]]) -> None: + """Best-effort typing send; ``send_typing`` already logs and backs off on its own failures.""" with contextlib.suppress(Exception): await self.send_typing(chat_id, metadata=metadata) + async def _retrigger_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]]) -> None: + """Re-arm typing after an intermediate send (Telegram clears it when a message lands). Skipped on + the FINAL reply (``metadata["notify"]``): the refresh loop is gone and no API cancels the bubble. + + Scheduled rather than awaited. ``sendChatAction`` is a fire-and-forget UI hint whose result nobody + reads, but awaiting it here ran its TLS round-trip on the same event loop as the ``getUpdates`` + long-polls. Streaming re-arms after *every* intermediate send, so with several agents streaming at + once the loop stayed pinned, the polls were never serviced, and they decayed into CLOSE-WAIT while + the adapter still reported ``connected`` (#111727). One in-flight re-arm per chat, at most one per + ``typing_retrigger_min_interval_seconds``.""" + if (metadata or {}).get("notify"): + return + # Only _keep_typing consulted this, so the documented `typing_indicator: false` workaround still + # paid for a sendChatAction on every intermediate send. + if not getattr(getattr(self, "config", None), "typing_indicator", True): + return + chat_key = str(chat_id) + tasks, sent_at, min_interval = self._typing_retrigger_state() + in_flight = tasks.get(chat_key) + if in_flight is not None and not in_flight.done(): + return + try: + loop = asyncio.get_running_loop() + except RuntimeError: + return + now = loop.time() + if min_interval > 0: + previous = sent_at.get(chat_key) + # Stamped at scheduling time, not completion: a burst of chunks must not all pass the + # check while the first round-trip is still open. + if previous is not None and (now - previous) < min_interval: + return + sent_at[chat_key] = now + task = loop.create_task(self._send_typing_quietly(chat_id, metadata)) + tasks[chat_key] = task + task.add_done_callback(lambda finished, key=chat_key: self._clear_typing_retrigger(key, finished)) + # Shutdown cancels _background_tasks, so a detached re-arm cannot outlive the adapter. + tracked = getattr(self, "_background_tasks", None) + if isinstance(tracked, set): + tracked.add(task) + task.add_done_callback(tracked.discard) + async def send( self, chat_id: str, content: str, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a message to a Telegram chat.""" diff --git a/tests/gateway/test_telegram_typing_retrigger.py b/tests/gateway/test_telegram_typing_retrigger.py new file mode 100644 index 0000000000..ef02b83b56 --- /dev/null +++ b/tests/gateway/test_telegram_typing_retrigger.py @@ -0,0 +1,190 @@ +"""Telegram post-send typing re-arm: off the critical path, deduped and rate-limited (#111727). + +Awaiting ``sendChatAction`` after every intermediate send ran its TLS round-trip on the same event +loop as the ``getUpdates`` long-polls; under concurrent streaming the polls were starved until they +rotted into CLOSE-WAIT while the adapter still reported ``connected``. +""" + +import asyncio +import sys +from pathlib import Path +from unittest.mock import AsyncMock + +import pytest + +_repo = str(Path(__file__).resolve().parents[2]) +if _repo not in sys.path: + sys.path.insert(0, _repo) +from gateway.config import PlatformConfig +from plugins.platforms.telegram.adapter import TelegramAdapter + + +def _make_adapter(**config_kwargs): + adapter = TelegramAdapter(PlatformConfig(enabled=True, token="test-token", **config_kwargs)) + adapter._bot = AsyncMock() + adapter._bot.send_chat_action = AsyncMock(return_value=None) + return adapter + + +async def _drain(adapter, chat_id="123"): + """Await the scheduled re-arm so assertions see its effect.""" + task = adapter._telegram_typing_retrigger_tasks.get(str(chat_id)) + if task is not None: + await task + + +@pytest.mark.asyncio +async def test_retrigger_does_not_await_the_round_trip(): + """The send path must return before sendChatAction completes — this is the whole fix.""" + adapter = _make_adapter() + released = asyncio.Event() + started = asyncio.Event() + + async def blocking_action(**kwargs): + started.set() + await released.wait() + + adapter._bot.send_chat_action = AsyncMock(side_effect=blocking_action) + + # Would hang forever if the re-arm were still awaited inline. + await asyncio.wait_for(adapter._retrigger_typing("123", None), timeout=1.0) + + await asyncio.wait_for(started.wait(), timeout=1.0) + assert not adapter._telegram_typing_retrigger_tasks["123"].done() + released.set() + await _drain(adapter) + assert adapter._bot.send_chat_action.await_count == 1 + + +@pytest.mark.asyncio +async def test_streaming_burst_collapses_to_one_chat_action(): + """Every chunk of a streamed reply re-arms typing; within the interval that is one API call.""" + adapter = _make_adapter() + + for _ in range(20): + await adapter._retrigger_typing("123", None) + await _drain(adapter) + + assert adapter._bot.send_chat_action.await_count == 1 + + +@pytest.mark.asyncio +async def test_retrigger_resumes_after_the_interval_elapses(): + adapter = _make_adapter() + adapter._telegram_typing_retrigger_interval = 2.0 + + await adapter._retrigger_typing("123", None) + await _drain(adapter) + + # Simulate the interval having passed rather than sleeping through it. + adapter._telegram_typing_retrigger_at["123"] -= 2.5 + await adapter._retrigger_typing("123", None) + await _drain(adapter) + + assert adapter._bot.send_chat_action.await_count == 2 + + +@pytest.mark.asyncio +async def test_throttle_is_per_chat(): + adapter = _make_adapter() + + await adapter._retrigger_typing("123", None) + await adapter._retrigger_typing("456", None) + await _drain(adapter, "123") + await _drain(adapter, "456") + + assert adapter._bot.send_chat_action.await_count == 2 + + +@pytest.mark.asyncio +async def test_in_flight_rearm_is_not_duplicated(): + """A slow round-trip must not accumulate one task per chunk.""" + adapter = _make_adapter() + adapter._telegram_typing_retrigger_interval = 0.0 # throttle off: the in-flight guard alone + released = asyncio.Event() + + async def blocking_action(**kwargs): + await released.wait() + + adapter._bot.send_chat_action = AsyncMock(side_effect=blocking_action) + + for _ in range(10): + await adapter._retrigger_typing("123", None) + + assert len(adapter._telegram_typing_retrigger_tasks) == 1 + released.set() + await _drain(adapter) + assert adapter._bot.send_chat_action.await_count == 1 + + +@pytest.mark.asyncio +async def test_typing_indicator_disabled_suppresses_rearm(): + """`typing_indicator: false` only gated _keep_typing, so the re-arm still cost a call per send.""" + adapter = _make_adapter(typing_indicator=False) + + await adapter._retrigger_typing("123", None) + await _drain(adapter) + + assert adapter._telegram_typing_retrigger_tasks == {} + adapter._bot.send_chat_action.assert_not_called() + + +@pytest.mark.asyncio +async def test_final_reply_still_does_not_rearm(): + adapter = _make_adapter() + + await adapter._retrigger_typing("123", {"notify": True}) + + assert adapter._telegram_typing_retrigger_tasks == {} + adapter._bot.send_chat_action.assert_not_called() + + +@pytest.mark.asyncio +async def test_rearm_is_tracked_for_shutdown_cancellation(): + """Detached tasks must join _background_tasks or they outlive the adapter.""" + adapter = _make_adapter() + released = asyncio.Event() + + async def blocking_action(**kwargs): + await released.wait() + + adapter._bot.send_chat_action = AsyncMock(side_effect=blocking_action) + + await adapter._retrigger_typing("123", None) + task = adapter._telegram_typing_retrigger_tasks["123"] + assert task in adapter._background_tasks + + await adapter.cancel_background_tasks() + assert task.cancelled() or task.done() + assert adapter._telegram_typing_retrigger_tasks == {} + + +@pytest.mark.asyncio +async def test_rearm_failure_does_not_escape_to_the_send_path(): + adapter = _make_adapter() + adapter._bot.send_chat_action = AsyncMock(side_effect=OSError("telegram network failure")) + + await adapter._retrigger_typing("123", None) + await _drain(adapter) + + assert adapter._telegram_typing_retrigger_tasks == {} + + +@pytest.mark.asyncio +async def test_send_returns_without_waiting_on_typing(): + """End-to-end: an intermediate send completes even while sendChatAction is stalled.""" + adapter = _make_adapter() + adapter._rich_messages_enabled = False + adapter._bot.send_message = AsyncMock(return_value=type("Msg", (), {"message_id": 1})()) + released = asyncio.Event() + + async def blocking_action(**kwargs): + await released.wait() + + adapter._bot.send_chat_action = AsyncMock(side_effect=blocking_action) + + result = await asyncio.wait_for(adapter.send("123", "chunk"), timeout=1.0) + + assert result.success is True + released.set() + await _drain(adapter)