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.
This commit is contained in:
kshitijk4poor
2026-09-12 23:11:06 +05:30
committed by kshitij
parent 78b98032c5
commit dbd027287e
2 changed files with 48 additions and 39 deletions
+12 -10
View File
@@ -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 "
@@ -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