fix(telegram): report a healthy-but-deaf ingress instead of nothing (#102260)

Every Telegram health probe measures the transport. A getUpdates round-trip
that returns 200 proves bytes are moving and nothing else: the stall watchdog
(#92991), the pending-update probe (#42909/#55769), the get_me() heartbeat
(#66377) and the polling-progress instrumentation all stay green while updates
arrive and then die downstream. The adapter then publishes "connected", logs
nothing at all, and is indistinguishable from a bot nobody has messaged.

That is #102260: three weeks of telegram.state "connected" plus "polling
confirmed healthy: getUpdates progressing (generation 1)" with zero inbound
reaching the agent, surviving every restart. Two of the issue's three
hypotheses do not hold on this code — _record_polling_progress fires on every
round-trip (not only at start_polling), and _send_path_degraded is cleared on
the first confirmed round-trip — and the reporter's own observation that fresh
messages are received but not processed places the failure downstream of the
transport, in the one stretch with no instrumentation at all.

Add the missing delivered side of the accounting:

- received: updates Telegram handed the process, read from the getUpdates
  envelope the adapter already parses (an empty result proves the transport,
  not arrival, so only non-empty results count).
- dispatched: updates PTB's dispatcher carried through the whole handler
  chain, stamped in the existing group-99 catch-all before its early returns.
- delivered: inbound events that reached the gateway's message handler,
  stamped in BasePlatformAdapter.handle_message for every platform.

_check_ingress_delivery_gap runs on the existing heartbeat and, when updates
arrived but nothing was delivered for 300s, names the broken hop: received >
dispatched means the dispatcher is not draining, dispatched > delivered means
Hermes is dropping what arrives. Diagnostic only — a received update
legitimately reaches no gateway turn, and reconnecting a healthy transport
cannot repair a dropped update, so this never drives recovery.

Also make the silent discard on the shared funnel speak: handle_message
returned with no log when no message handler was installed, so a mis-wired
adapter discarded 100% of inbound while connected and able to send.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018pT5hFJBRfLj8KMqFhm3qz
This commit is contained in:
joaomarcos
2026-09-03 14:58:43 -03:00
committed by kshitij
parent b05a47b9d2
commit db407dd078
3 changed files with 424 additions and 0 deletions
+35
View File
@@ -1809,6 +1809,11 @@ class BasePlatformAdapter(ABC):
self.config = config
self.platform = platform
self._message_handler: Optional[MessageHandler] = None
# Ingress delivery counter (#102260): a healthy transport can still discard 100% of
# inbound; without a delivered-side counter a deaf gateway looks idle.
self._inbound_delivered_total: int = 0
self._last_inbound_delivered_monotonic: Optional[float] = None
self._no_message_handler_logged: bool = False
self._reaction_handler: Optional[Callable[[Dict[str, Any]], Awaitable[None]]] = None
# Runner-owned boundary for normalized events: auth/profile state never lives in an adapter.
self._platform_event_handler: Optional[Callable[[Dict[str, Any], Any], Awaitable[None]]] = None
@@ -2161,6 +2166,21 @@ class BasePlatformAdapter(ABC):
"""Check if adapter is currently connected."""
return self._running
def note_inbound_delivered(self) -> None:
"""Record that one inbound event reached the gateway message handler.
The delivered side of the ingress accounting introduced for #102260.
Adapters that observe their own wire traffic (Telegram's getUpdates
instrumentation) compare their received counter against this one to
tell "nothing is arriving" apart from "everything arriving is being
dropped downstream" — two failures that look identical from every
transport-level health probe.
"""
self._inbound_delivered_total = (
getattr(self, "_inbound_delivered_total", 0) + 1
)
self._last_inbound_delivered_monotonic = time.monotonic()
def set_message_handler(self, handler: MessageHandler) -> None:
"""Set the incoming-message handler (MessageEvent -> optional response str)."""
self._message_handler = handler
@@ -3617,7 +3637,22 @@ class BasePlatformAdapter(ABC):
task so new messages (and interrupts) can arrive while an agent runs."""
event._gateway_accepted = False
if not self._message_handler:
# A connected adapter with no message handler is silently deaf: it
# polls, publishes "connected", and can still send — while every
# inbound message is discarded here with no log at all. Say so once
# per adapter so the failure is diagnosable (#102260).
if not getattr(self, "_no_message_handler_logged", False):
self._no_message_handler_logged = True
logger.error(
"[%s] Dropping inbound message: no gateway message handler "
"is installed on this adapter. The adapter is connected and "
"can send, but every inbound message is discarded.",
self.name,
)
return
self.note_inbound_delivered()
if event.allow_gateway_control:
coerce_plaintext_gateway_command(event)
expected_session_key = str((event.metadata or {}).get("gateway_session_key") or "").strip()
+101
View File
@@ -344,6 +344,10 @@ _POLLING_PROGRESS_TIMEOUT = 60.0 # generation unhealthy until getUpdates return
# #92991) and no other probe can see it. ~3x the worst-case poll window leaves ample margin against false
# positives while still recovering within a few heartbeat intervals.
_POLLING_STALL_TIMEOUT = 150.0
# Ingress delivery gap (#102260): transport probes prove bytes move, not that any update reached a
# handler. Generous because updates legitimately reach no gateway turn (own messages, reactions,
# unmentioned group chatter); diagnostic only, never drives recovery.
_INGRESS_DELIVERY_GAP_TIMEOUT = 300.0
# sendVideo transcodes before answering, outlasting the 20s read timeout; also how long a user waits
# to hear the attachment failed, so kept modest.
_MEDIA_SEND_READ_TIMEOUT = 60.0
@@ -463,6 +467,12 @@ class TelegramAdapter(BasePlatformAdapter):
# began, and when the last successful getUpdates round-trip completed.
self._polling_generation_started_monotonic: Optional[float] = None
self._polling_last_progress_monotonic: Optional[float] = None
# Ingress accounting (#102260): received (getUpdates wire) vs dispatched (PTB handler chain)
# vs the base adapter's delivered counter.
self._updates_received_total: int = 0
self._updates_dispatched_total: int = 0
self._last_update_received_monotonic: Optional[float] = None
self._ingress_gap_reported_at_received: Optional[int] = None
# Live @username: PTB caches getMe() at initialize() and only rewrites it inside get_me(), so a
# BotFather rename leaves self._bot.username stale; routing reads _current_bot_username().
self._bot_username_observed: Optional[str] = None
@@ -1619,6 +1629,22 @@ class TelegramAdapter(BasePlatformAdapter):
return
if isinstance(envelope, dict) and envelope.get("ok") is True and "result" in envelope:
self._record_polling_progress(generation)
self._record_updates_received(envelope.get("result"))
def _record_updates_received(self, result) -> None:
"""Count updates Telegram handed us on the getUpdates wire (#102260).
``_record_polling_progress`` above proves the transport works; it says
nothing about whether anything arrived. Only a non-empty ``result``
means real updates entered this process, and that is the number the
delivered counter has to be compared against.
"""
if not isinstance(result, list) or not result:
return
self._updates_received_total = (
getattr(self, "_updates_received_total", 0) + len(result)
)
self._last_update_received_monotonic = time.monotonic()
def _instrument_polling_request(self, request):
"""Instrument one dedicated PTB getUpdates request with progress tracking.
@@ -2032,6 +2058,10 @@ class TelegramAdapter(BasePlatformAdapter):
# so a consumer with no successful round-trip past the stall threshold is dead (#92991).
# Pure local-state check — no Bot API call needed.
await self._check_polling_stall()
# Transport health is not delivery health: updates can arrive
# and then die downstream with every probe above still green
# (#102260). Pure local-state check — no Bot API call.
self._check_ingress_delivery_gap()
except asyncio.CancelledError:
return
except (asyncio.TimeoutError, OSError) as probe_err:
@@ -2114,6 +2144,74 @@ class TelegramAdapter(BasePlatformAdapter):
logger.warning(log_message, self.name)
self._polling_error_task = asyncio.get_running_loop().create_task(self._handle_polling_network_error(RuntimeError(reason)))
def _check_ingress_delivery_gap(self) -> None:
"""Report updates that arrived on the wire but reached no gateway turn.
Every other probe in this adapter measures the transport. When the
transport is healthy and updates are still not being acted on, the
gateway publishes ``connected``, logs nothing, and looks exactly like
an idle bot — the failure mode of #102260, which stayed undiagnosed for
three weeks because no counter existed on the delivered side.
Three counters split the failure into its distinguishable causes:
* ``received`` — updates Telegram handed this process (getUpdates).
* ``dispatched`` — updates PTB's dispatcher carried through the whole
handler chain. ``received > dispatched`` means the dispatcher is not
draining; the transport is innocent.
* ``delivered`` — inbound events that reached the gateway's message
handler. ``dispatched > delivered`` means Hermes itself is dropping
what arrives (authorization, mention/topic gating, a missing handler).
Deliberately diagnostic only. A received update legitimately reaches no
gateway turn (the bot's own messages, reactions, unmentioned group
chatter), so this must never drive recovery on its own — reconnecting a
healthy transport would not fix a dropped update anyway. It reports
once per received-count so a persistent gap does not spam the log.
"""
if self._webhook_mode:
return
if getattr(self, "_polling_teardown_started", False):
return
if self.has_fatal_error:
return
received = getattr(self, "_updates_received_total", 0)
if not received:
return
last_received = getattr(self, "_last_update_received_monotonic", None)
if last_received is None:
return
delivered = getattr(self, "_inbound_delivered_total", 0)
last_delivered = getattr(self, "_last_inbound_delivered_monotonic", None)
# Something was delivered after the last update arrived: the chain is
# working end to end. Re-arm so a later gap is reported again.
if last_delivered is not None and last_delivered >= last_received:
self._ingress_gap_reported_at_received = None
return
if time.monotonic() - last_received <= _INGRESS_DELIVERY_GAP_TIMEOUT:
return
if self._ingress_gap_reported_at_received == received:
return
self._ingress_gap_reported_at_received = received
dispatched = getattr(self, "_updates_dispatched_total", 0)
if dispatched < received:
where = (
"PTB's dispatcher is not draining them (received > dispatched)"
)
else:
where = (
"they are being dropped inside Hermes before the gateway "
"handler (dispatched > delivered) — check authorization, "
"mention/topic gating, and the held-inbound queue"
)
logger.error(
"[%s] Telegram ingress is healthy but deaf: %d update(s) received "
"on the getUpdates wire, %d dispatched, %d delivered to the "
"gateway, none delivered in the last %.0fs. Polling is fine — %s.",
self.name, received, dispatched, delivered,
time.monotonic() - last_received, where,
)
async def _check_polling_stall(self) -> None:
"""Watchdog the last successful getUpdates round-trip: a long-poll can wedge without raising
(CLOSE-WAIT after a route flip) while every other probe stays blind; no round-trip for
@@ -2553,6 +2651,9 @@ class TelegramAdapter(BasePlatformAdapter):
async def _on_platform_update(self, update, context) -> None:
"""Catch-all PTB handler (group 99) firing ``gateway_platform_event`` per inbound update with a
stable envelope (no raw SDK objects) and an internal auth source. Never raises into PTB."""
# Last handler group PTB runs: proves the dispatcher carried the update through the chain,
# so stamp before any early return (#102260).
self._updates_dispatched_total = getattr(self, "_updates_dispatched_total", 0) + 1
handler: Optional[Callable[[Dict[str, Any], Any], Awaitable[None]]] = getattr(self, "_platform_event_handler", None)
if handler is None:
return
@@ -0,0 +1,288 @@
"""Telegram ingress delivery accounting (#102260).
Every existing Telegram health probe measures the *transport*: a getUpdates
round-trip that returns 200 proves bytes are moving and nothing else. When
updates arrive and then die downstream — a dispatcher that never drains, a
filter that never matches, an adapter with no gateway handler installed — the
gateway publishes ``connected``, logs nothing at all, and is indistinguishable
from an idle bot. That is #102260: three weeks of a bot reporting
``telegram.state: "connected"`` and ``polling confirmed healthy`` while not one
inbound message reached a handler.
These tests pin the three counters that split the failure into causes:
* ``_updates_received_total`` — updates Telegram handed the process.
* ``_updates_dispatched_total`` — updates PTB carried through the handler chain.
* ``_inbound_delivered_total`` — inbound events that reached the gateway.
"""
import asyncio
import logging
import time as _time
from unittest.mock import AsyncMock, MagicMock
import pytest
from gateway.config import PlatformConfig
from gateway.platforms.base import MessageEvent, MessageType, Platform, SessionSource
from plugins.platforms.telegram.adapter import TelegramAdapter
def _make_adapter() -> TelegramAdapter:
adapter = TelegramAdapter(PlatformConfig(enabled=True, token="***"))
adapter._webhook_mode = False
adapter._app = MagicMock()
adapter._app.updater.running = True
return adapter
def _envelope_request(result):
"""A PTB request double whose parser returns a getUpdates envelope."""
request = MagicMock()
request.parse_json_payload = MagicMock(
return_value={"ok": True, "result": result}
)
return request
# ---------------------------------------------------------------------------
# received counter
# ---------------------------------------------------------------------------
def test_empty_getupdates_result_counts_no_updates():
"""An idle long-poll proves the transport, not that anything arrived."""
adapter = _make_adapter()
adapter._polling_generation = 1
adapter._polling_progress_accepting = True
adapter._polling_progress_event = asyncio.Event()
adapter._observe_polling_request_result(_envelope_request([]), 1, (200, b"{}"))
assert adapter._updates_received_total == 0
assert adapter._last_update_received_monotonic is None
# ...while transport progress is still recorded.
assert adapter._polling_last_progress_monotonic is not None
def test_non_empty_getupdates_result_counts_every_update():
adapter = _make_adapter()
adapter._polling_generation = 1
adapter._polling_progress_accepting = True
adapter._polling_progress_event = asyncio.Event()
adapter._observe_polling_request_result(
_envelope_request([{"update_id": 1}, {"update_id": 2}]), 1, (200, b"{}")
)
assert adapter._updates_received_total == 2
assert adapter._last_update_received_monotonic is not None
# ---------------------------------------------------------------------------
# delivered counter
# ---------------------------------------------------------------------------
def _text_event() -> MessageEvent:
return MessageEvent(
text="hi",
message_type=MessageType.TEXT,
source=SessionSource(
platform=Platform.TELEGRAM,
chat_id="6053106869",
user_id="6053106869",
chat_type="dm",
),
)
@pytest.mark.asyncio
async def test_missing_message_handler_is_logged_once_not_silent(caplog):
"""A connected adapter with no handler discards 100% of inbound.
Before #102260 this returned with no log line at all, so a mis-wired
adapter looked exactly like a quiet chat.
"""
adapter = _make_adapter()
adapter._message_handler = None
with caplog.at_level(logging.ERROR):
await adapter.handle_message(_text_event())
await adapter.handle_message(_text_event())
matches = [r for r in caplog.records if "no gateway message handler" in r.message]
assert len(matches) == 1, "must report the deaf adapter exactly once"
assert adapter._inbound_delivered_total == 0
@pytest.mark.asyncio
async def test_delivery_to_gateway_handler_increments_counter():
adapter = _make_adapter()
adapter.note_inbound_delivered()
assert adapter._inbound_delivered_total == 1
assert adapter._last_inbound_delivered_monotonic is not None
# ---------------------------------------------------------------------------
# gap watchdog
# ---------------------------------------------------------------------------
def _armed_gap_adapter(*, received=3, dispatched=3, delivered=0, age=600.0):
adapter = _make_adapter()
now = _time.monotonic()
adapter._updates_received_total = received
adapter._updates_dispatched_total = dispatched
adapter._inbound_delivered_total = delivered
adapter._last_update_received_monotonic = now - age
adapter._last_inbound_delivered_monotonic = None
return adapter
def test_no_report_while_within_the_gap_window(caplog):
adapter = _armed_gap_adapter(age=10.0)
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
assert not [r for r in caplog.records if "healthy but deaf" in r.message]
def test_no_report_when_nothing_was_ever_received(caplog):
"""An idle bot must never be accused of being deaf."""
adapter = _make_adapter()
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
assert not [r for r in caplog.records if "healthy but deaf" in r.message]
def test_no_report_when_delivery_followed_the_last_update(caplog):
adapter = _armed_gap_adapter()
adapter._inbound_delivered_total = 3
adapter._last_inbound_delivered_monotonic = (
adapter._last_update_received_monotonic + 1
)
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
assert not [r for r in caplog.records if "healthy but deaf" in r.message]
def test_dispatched_shortfall_names_the_dispatcher(caplog):
"""received > dispatched: PTB never carried the updates to the handlers."""
adapter = _armed_gap_adapter(received=5, dispatched=1)
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
(record,) = [r for r in caplog.records if "healthy but deaf" in r.message]
assert "dispatcher is not draining" in record.getMessage()
def test_delivery_shortfall_names_hermes(caplog):
"""dispatched == received but nothing delivered: Hermes drops them."""
adapter = _armed_gap_adapter(received=4, dispatched=4)
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
(record,) = [r for r in caplog.records if "healthy but deaf" in r.message]
message = record.getMessage()
assert "dropped inside Hermes" in message
assert "4 update(s) received" in message
def test_report_is_not_repeated_until_more_updates_arrive(caplog):
adapter = _armed_gap_adapter()
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
adapter._check_ingress_delivery_gap()
assert len([r for r in caplog.records if "healthy but deaf" in r.message]) == 1
# A further update with still no delivery is a new, reportable data point.
adapter._updates_received_total += 1
adapter._updates_dispatched_total += 1
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
assert len([r for r in caplog.records if "healthy but deaf" in r.message]) == 2
def test_gap_check_never_triggers_recovery():
"""Reconnecting a healthy transport cannot fix a dropped update.
The reconnect ladder belongs to the transport probes; this check is
diagnosis only, so it must leave the ladder untouched.
"""
adapter = _armed_gap_adapter()
adapter._check_ingress_delivery_gap()
assert adapter._polling_error_task is None
def test_gap_check_skipped_in_webhook_mode(caplog):
adapter = _armed_gap_adapter()
adapter._webhook_mode = True
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
assert not [r for r in caplog.records if "healthy but deaf" in r.message]
def test_gap_check_skipped_during_teardown(caplog):
adapter = _armed_gap_adapter()
adapter._polling_teardown_started = True
with caplog.at_level(logging.ERROR):
adapter._check_ingress_delivery_gap()
assert not [r for r in caplog.records if "healthy but deaf" in r.message]
# ---------------------------------------------------------------------------
# dispatched counter
# ---------------------------------------------------------------------------
@pytest.mark.asyncio
async def test_catch_all_observer_counts_dispatch_before_early_return():
"""The counter must be stamped even when no plugin hook is installed.
``_on_platform_update`` returns immediately without a platform-event
handler; if the counter sat after that return, every gateway without the
hook would report a permanent dispatcher stall.
"""
adapter = _make_adapter()
adapter._platform_event_handler = None
await adapter._on_platform_update(MagicMock(), MagicMock())
assert adapter._updates_dispatched_total == 1
@pytest.mark.asyncio
async def test_heartbeat_runs_the_gap_check():
"""The check has to be wired into the loop that actually runs."""
adapter = _make_adapter()
bot = MagicMock()
bot.get_me = AsyncMock()
adapter._app.bot = bot
adapter._bot = bot
called = []
adapter._check_ingress_delivery_gap = lambda: called.append(True)
adapter._probe_pending_updates = AsyncMock()
adapter._check_polling_stall = AsyncMock()
real_sleep = asyncio.sleep
async def fast_sleep(_delay, *args, **kwargs):
await real_sleep(0)
task = asyncio.get_running_loop().create_task(adapter._polling_heartbeat_loop())
import plugins.platforms.telegram.adapter as tg_adapter
original = tg_adapter.asyncio.sleep
tg_adapter.asyncio.sleep = fast_sleep
try:
for _ in range(50):
await real_sleep(0)
if called:
break
finally:
tg_adapter.asyncio.sleep = original
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
assert called, "_polling_heartbeat_loop must run the ingress gap check"