fix(telegram): trim the liveness heartbeat, fix the polling error-callback log lines

- Drop the 900s "inbound liveness" INFO line from the salvage: a periodic
  log heartbeat is a feature with a separate scoping call (#111211 item 3);
  the existing stall watchdog already escalates when getUpdates stops.
- The polling error_callback interpolated the raw exception and left the
  redaction call inside the format string, so its lines read
  "Telegram network _redact_telegram_error_text(error), scheduling
  reconnect: ..." and leaked unredacted text. Redact for real.
- Tests: keep two invariants (recovered wording after network errors;
  clean bootstrap stays "confirmed healthy") and the transport test.
This commit is contained in:
teknium1
2026-09-14 18:30:30 -07:00
committed by Teknium
parent ff307aea5e
commit 1847ad2883
2 changed files with 22 additions and 46 deletions
+6 -19
View File
@@ -356,7 +356,6 @@ _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
_POLLING_LIVENESS_LOG_INTERVAL = 900.0
# Ingress dispatch stall (#102260): the transport probes prove getUpdates round-trips complete, not
# that PTB's dispatcher ever handed the fetched updates to a handler. Two heartbeats (180s) with a
# backlog and no dispatch progress: diagnostic only, never drives recovery (#71240 owns that).
@@ -489,7 +488,6 @@ 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
self._polling_last_liveness_log_monotonic: Optional[float] = None
# Ingress accounting (#102260): received (getUpdates wire) vs dispatched (PTB group-99 catch-all).
self._updates_received_total: int = 0
self._updates_dispatched_total: int = 0
@@ -1627,25 +1625,13 @@ class TelegramAdapter(BasePlatformAdapter):
"""Record successful getUpdates I/O for the current generation only; True when accepted."""
if self._teardown_started or not self._polling_progress_accepting or generation != self._polling_generation:
return False
now = time.monotonic()
first_progress = not self._polling_progress_event.is_set()
if first_progress:
if not self._polling_progress_event.is_set():
# First confirmed round-trip resolves the "health pending" line both reconnect paths end on.
# After network-error WARNINGs the line must read as the matching recovery event (#111211).
state = "recovered" if self._polling_network_error_count else "confirmed healthy"
logger.info("[%s] Telegram polling %s: getUpdates progressing (generation %d)", self.name, state, generation)
self._polling_last_liveness_log_monotonic = now
elif (
self._polling_last_liveness_log_monotonic is None
or now - self._polling_last_liveness_log_monotonic >= _POLLING_LIVENESS_LOG_INTERVAL
):
logger.info(
"[%s] Telegram inbound liveness: getUpdates progressing (generation %d)",
self.name,
generation,
)
self._polling_last_liveness_log_monotonic = now
self._polling_progress_event.set()
self._polling_last_progress_monotonic = now
self._polling_last_progress_monotonic = time.monotonic()
self._polling_network_error_count = 0
if generation == self._polling_conflict_recovery_generation:
self._polling_conflict_recovery_generation = None
@@ -2972,10 +2958,11 @@ class TelegramAdapter(BasePlatformAdapter):
self._disarm_ptb_retry_loop()
self._spawn_polling_recovery(loop, self._handle_polling_conflict(error))
elif self._looks_like_network_error(error):
logger.warning("[%s] Telegram network _redact_telegram_error_text(error), scheduling reconnect: %s", self.name, error)
logger.warning(
"[%s] Telegram network error, scheduling reconnect: %s", self.name, _redact_telegram_error_text(error))
self._spawn_polling_recovery(loop, self._handle_polling_network_error(error))
else:
logger.error("[%s] Telegram polling _redact_telegram_error_text(error): %s", self.name, error, exc_info=True)
logger.error("[%s] Telegram polling error: %s", self.name, _redact_telegram_error_text(error), exc_info=True)
self._polling_error_callback_ref = _polling_error_callback # reused by _handle_polling_conflict
polling_started = await self._start_polling_resilient(
@@ -10,7 +10,6 @@ reconnect is a reliable hung-poll signature.
import asyncio
import logging
from unittest.mock import patch
from gateway.config import Platform # noqa: E402
from plugins.platforms.telegram.adapter import TelegramAdapter # noqa: E402
@@ -29,7 +28,6 @@ def _bare_adapter():
a._polling_conflict_count = 3
a._polling_conflict_recovery_generation = None
a._send_path_degraded = True
a._polling_last_liveness_log_monotonic = None
return a
@@ -43,22 +41,25 @@ class TestPollingHealthConfirmation:
assert "generation 1" in rendered
assert a._polling_progress_event.is_set()
def test_subsequent_progress_before_liveness_interval_is_silent(self, caplog):
"""Progress before the liveness interval must not spam INFO logs."""
def test_subsequent_progress_is_silent(self, caplog):
"""Only the FIRST round-trip of a generation logs — a quiet evening
must not spam one INFO per getUpdates poll."""
a = _bare_adapter()
with patch(
"plugins.platforms.telegram.adapter.time.monotonic",
side_effect=[10.0, 100.0, 909.0],
):
a._record_polling_progress(1) # first — logs
caplog.clear()
with caplog.at_level(logging.INFO, logger="plugins.platforms.telegram.adapter"):
a._record_polling_progress(1) # second — silent
a._record_polling_progress(1) # third — still below 900 seconds
a._record_polling_progress(1) # first — logs
with caplog.at_level(logging.INFO, logger="plugins.platforms.telegram.adapter"):
a._record_polling_progress(1) # second — silent
a._record_polling_progress(1) # third — silent
assert not [rec for rec in caplog.records if "Telegram polling" in rec.getMessage()]
def test_clean_bootstrap_stays_confirmed_healthy(self, caplog):
"""No preceding network errors → no fake "recovered" wording."""
a = _bare_adapter()
a._polling_network_error_count = 0
with caplog.at_level(logging.INFO, logger="plugins.platforms.telegram.adapter"):
a._record_polling_progress(1)
rendered = " | ".join(rec.getMessage() for rec in caplog.records)
assert "Telegram polling" not in rendered
assert "Telegram inbound liveness" not in rendered
assert "confirmed healthy" in rendered
assert "recovered" not in rendered
def test_new_generation_logs_again(self, caplog):
"""A reconnect starts a new generation with a fresh event; its first
@@ -77,18 +78,6 @@ class TestPollingHealthConfirmation:
assert "polling recovered" in rendered
assert "generation 2" in rendered
def test_established_generation_emits_low_frequency_inbound_liveness(self, caplog):
a = _bare_adapter()
with patch("plugins.platforms.telegram.adapter.time.monotonic", return_value=10.0):
a._record_polling_progress(1)
caplog.clear()
with caplog.at_level(logging.INFO, logger="plugins.platforms.telegram.adapter"):
with patch("plugins.platforms.telegram.adapter.time.monotonic", return_value=910.0):
a._record_polling_progress(1)
assert "Telegram inbound liveness: getUpdates progressing" in caplog.text
def test_stale_generation_progress_stays_silent(self, caplog):
"""Progress from an abandoned generation must neither log nor set the
current event (pre-existing guard, pinned here because the log line