From dbd027287e1575cbd0acccb8707689f819dfa490 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Sat, 12 Sep 2026 23:11:06 +0530 Subject: [PATCH] refactor(telegram): tighten the dispatch-stall check shape from review - Correct the per-generation reset comment: an in-place updater restart keeps PTB's update_queue, so old-generation dispatches can briefly exceed received; the check already treats that as no backlog. - Make the once-per-stall gate explicit (cap the heartbeat count) instead of relying on `!=`. - Drop the dead getattr in _record_updates_received (only reachable after _record_polling_progress dereferenced the same instance state); keep the fallbacks in the heartbeat check and group-99 handler, which sibling watchdogs share because object.__new__ adapter doubles exist in tests. - Split the single invariant test so a failure names the broken guard: stall-once, re-arm-on-progress, generation-reset; drop the unused _app mock. --- plugins/platforms/telegram/adapter.py | 22 ++++--- .../test_telegram_ingress_delivery_gap.py | 65 ++++++++++--------- 2 files changed, 48 insertions(+), 39 deletions(-) diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index 4b1fda3ca3..a852b7b830 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -1590,7 +1590,8 @@ class TelegramAdapter(BasePlatformAdapter): # See #92991. self._polling_generation_started_monotonic = time.monotonic() self._polling_last_progress_monotonic = None - # A rebuilt consumer drops whatever the old update_queue held; start the backlog from zero. + # Re-base the backlog per generation. On an in-place updater restart PTB keeps the old + # update_queue, so old dispatches can briefly exceed received; the check treats that as no backlog. self._updates_received_total = self._updates_dispatched_total = 0 self._ingress_dispatched_seen = self._ingress_stalled_heartbeats = 0 return self._polling_generation, self._polling_progress_event @@ -1638,7 +1639,7 @@ class TelegramAdapter(BasePlatformAdapter): """Count updates Telegram handed us on the getUpdates wire (#102260). Only reached for the accepted generation, so a late response from a fenced poll cannot inflate the backlog.""" if isinstance(result, list) and result: - self._updates_received_total = getattr(self, "_updates_received_total", 0) + len(result) + self._updates_received_total += len(result) def _instrument_polling_request(self, request): """Instrument one dedicated PTB getUpdates request with progress tracking. @@ -2139,12 +2140,10 @@ class TelegramAdapter(BasePlatformAdapter): def _check_ingress_dispatch_stall(self) -> None: """Report fetched updates PTB's dispatcher is not handing to handlers (#102260). - ``received`` and ``dispatched`` count the same population — every update in a getUpdates - result reaches the group-99 catch-all (no handler raises ApplicationHandlerStop and no error - handler is registered) — so a backlog that does not shrink across - ``_INGRESS_DISPATCH_STALL_HEARTBEATS`` heartbeats is a wedged dispatcher, whatever the traffic - rate. Keyed on dispatcher progress rather than update age: a steady trickle of new updates must - not keep resetting the clock. Reports once per stall and re-arms when dispatch resumes. + ``received`` and ``dispatched`` count the same population (every fetched update reaches the + group-99 catch-all: no handler raises ApplicationHandlerStop, no error handler is registered), + so a backlog with no dispatch progress across ``_INGRESS_DISPATCH_STALL_HEARTBEATS`` heartbeats + is a wedged dispatcher at any traffic rate. Reports once per stall, re-arms on progress. """ if self._webhook_mode or self._teardown_started or self.has_fatal_error: return @@ -2154,8 +2153,11 @@ class TelegramAdapter(BasePlatformAdapter): self._ingress_dispatched_seen = dispatched self._ingress_stalled_heartbeats = 0 return - self._ingress_stalled_heartbeats = getattr(self, "_ingress_stalled_heartbeats", 0) + 1 - if self._ingress_stalled_heartbeats != _INGRESS_DISPATCH_STALL_HEARTBEATS: + stalled = getattr(self, "_ingress_stalled_heartbeats", 0) + if stalled >= _INGRESS_DISPATCH_STALL_HEARTBEATS: + return # already reported this stall + self._ingress_stalled_heartbeats = stalled + 1 + if stalled + 1 < _INGRESS_DISPATCH_STALL_HEARTBEATS: return logger.warning( "[%s] Telegram ingress is healthy but deaf: %d update(s) fetched by getUpdates have not been " diff --git a/tests/gateway/test_telegram_ingress_delivery_gap.py b/tests/gateway/test_telegram_ingress_delivery_gap.py index 7ec4d8b529..da1411601c 100644 --- a/tests/gateway/test_telegram_ingress_delivery_gap.py +++ b/tests/gateway/test_telegram_ingress_delivery_gap.py @@ -19,7 +19,6 @@ _DEAF = "healthy but deaf" def _polling_adapter() -> TelegramAdapter: adapter = TelegramAdapter(PlatformConfig(enabled=True, token="test-token")) adapter._webhook_mode = False - adapter._app = MagicMock() adapter._begin_polling_generation() return adapter @@ -34,56 +33,64 @@ def _receive(adapter: TelegramAdapter, n: int, generation: int | None = None) -> ) +async def _dispatch(adapter: TelegramAdapter, n: int) -> None: + for _ in range(n): + await adapter._on_platform_update(MagicMock(), MagicMock()) + + +def _heartbeats(adapter: TelegramAdapter, n: int) -> None: + for _ in range(n): + adapter._check_ingress_dispatch_stall() + + def _deaf_reports(caplog) -> list[str]: return [r.getMessage() for r in caplog.records if _DEAF in r.message] @pytest.mark.asyncio -async def test_dispatch_stall_is_reported_on_backlog_not_update_age(caplog): - """A wedged dispatcher is reported after two heartbeats even while new updates keep arriving; - dispatch progress re-arms it; a stale generation cannot inflate the backlog.""" +async def test_stall_reported_once_on_backlog_regardless_of_update_age(caplog): + """Healthy dispatch never reports; a wedged dispatcher is reported after two heartbeats even + while new updates keep arriving, and only once per stall.""" adapter = _polling_adapter() caplog.set_level(logging.WARNING) - - # Healthy: every fetched update reaches the group-99 catch-all. _receive(adapter, 2) - for _ in range(2): - await adapter._on_platform_update(MagicMock(), MagicMock()) - for _ in range(3): - adapter._check_ingress_dispatch_stall() + await _dispatch(adapter, 2) + _heartbeats(adapter, 3) assert _deaf_reports(caplog) == [] - # Dispatcher wedged: a fresh update lands before every heartbeat, none dispatched. - _receive(adapter, 1) - adapter._check_ingress_dispatch_stall() - assert _deaf_reports(caplog) == [], "one heartbeat is not a stall" - _receive(adapter, 1) - adapter._check_ingress_dispatch_stall() + for _ in range(3): # a fresh update lands before every heartbeat, none dispatched + _receive(adapter, 1) + adapter._check_ingress_dispatch_stall() (report,) = _deaf_reports(caplog) assert "2 update(s) fetched" in report and "4 received, 2 dispatched" in report - _receive(adapter, 1) - adapter._check_ingress_dispatch_stall() - assert len(_deaf_reports(caplog)) == 1, "a persistent stall is reported once" - # Dispatch resumes but only partially drains the backlog: progress re-arms the check, and the - # next two heartbeats without progress are a new stall. - await adapter._on_platform_update(MagicMock(), MagicMock()) + +@pytest.mark.asyncio +async def test_dispatch_progress_rearms_the_report(caplog): + adapter = _polling_adapter() + caplog.set_level(logging.WARNING) + _receive(adapter, 3) + _heartbeats(adapter, 3) + assert len(_deaf_reports(caplog)) == 1 + + await _dispatch(adapter, 1) # partial drain: progress, backlog remains adapter._check_ingress_dispatch_stall() assert len(_deaf_reports(caplog)) == 1 - for _ in range(2): - adapter._check_ingress_dispatch_stall() + _heartbeats(adapter, 2) assert len(_deaf_reports(caplog)) == 2 - # A rebuilt consumer starts the backlog from zero, and a late response from the fenced - # generation is ignored. + +def test_new_generation_restarts_backlog_and_ignores_fenced_polls(caplog): + adapter = _polling_adapter() + caplog.set_level(logging.WARNING) + _receive(adapter, 3) stale_generation = adapter._polling_generation adapter._begin_polling_generation() assert adapter._record_polling_progress(stale_generation) is False _receive(adapter, 5, generation=stale_generation) assert adapter._updates_received_total == 0 - for _ in range(3): - adapter._check_ingress_dispatch_stall() - assert len(_deaf_reports(caplog)) == 2 + _heartbeats(adapter, 3) + assert _deaf_reports(caplog) == [] @pytest.mark.asyncio