fix(telegram): schedule the post-send typing re-arm instead of awaiting it
`_retrigger_typing` awaited `sendChatAction` inline on the send path, and streaming re-arms after *every* intermediate send. `sendChatAction` is a fire-and-forget UI hint whose result nobody reads, but awaiting it ran its TLS round-trip on the same event loop as the `getUpdates` long-polls. With several agents streaming concurrently the loop stayed pinned, the long-polls were never serviced, and they decayed into CLOSE-WAIT while the adapter still reported `connected` — a gateway that is deaf but healthy, which `Restart=always` cannot recover because the process never exits. py-spy put 20 of 20 MainThread samples in `send_typing` -> `send_chat_action` -> `start_tls`. The handshakes are what cost: with `max_keepalive_connections=4`, a re-arm per chunk churns the pool so most calls pay a fresh TLS handshake on the loop thread. Three changes, all in the re-arm path: - Schedule the re-arm as a tracked task rather than awaiting it, so a round-trip never delays a send or a poll. It joins `_background_tasks`, so shutdown cancels it and it cannot outlive the adapter. - One in-flight re-arm per chat, and at most one per `typing_retrigger_min_interval_seconds` (default 2s, `extra` knob; 0 restores a call per send). 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. - Honour `typing_indicator: false`. Only `_keep_typing` consulted it, so the documented workaround still paid for a `sendChatAction` on every intermediate send. Simulating 200 streamed chunks with a 10ms loop-blocking handshake: 200 `sendChatAction` calls and 2037ms of send-path time before, 1 call and 13ms after. Fixes #111727 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
committed by
Teknium
parent
0b265ef4fc
commit
7c10c249ce
@@ -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."""
|
||||
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user