diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index 3f1d69bc2e..73ef76acf5 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -354,6 +354,7 @@ _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). @@ -486,6 +487,7 @@ 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 @@ -1623,11 +1625,25 @@ 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 - if not self._polling_progress_event.is_set(): + now = time.monotonic() + first_progress = not self._polling_progress_event.is_set() + if first_progress: # First confirmed round-trip resolves the "health pending" line both reconnect paths end on. - logger.info("[%s] Telegram polling confirmed healthy: getUpdates progressing (generation %d)", self.name, generation) + 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 = time.monotonic() + self._polling_last_progress_monotonic = now self._polling_network_error_count = 0 if generation == self._polling_conflict_recovery_generation: self._polling_conflict_recovery_generation = None diff --git a/plugins/platforms/telegram/telegram_network.py b/plugins/platforms/telegram/telegram_network.py index fd784d5a3b..d39bebb3d8 100644 --- a/plugins/platforms/telegram/telegram_network.py +++ b/plugins/platforms/telegram/telegram_network.py @@ -15,6 +15,16 @@ logger = logging.getLogger(__name__) _TELEGRAM_API_HOST = "api.telegram.org" + +def _describe_transport_error(error: Exception) -> str: + """Return a non-empty, secret-safe exception representation for diagnostics.""" + try: + from agent.redact import redact_sensitive_text + + return redact_sensitive_text(repr(error), force=True) + except Exception: + return type(error).__name__ + # TCP keepalive so a half-open/CLOSE-WAIT long-poll errors out instead of blocking getUpdates forever # (Windows leaves SO_KEEPALIVE off). Idle/interval knobs are best-effort per Python/OS combo. # Windows does not enable SO_KEEPALIVE on new sockets by default, so a dead api.telegram.org peer can hang @@ -80,6 +90,7 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport): # ``_UNSET`` / ``None`` / ``str`` = no sticky yet / sticky hostname / sticky IPv4. self._sticky_ip: object = _UNSET self._sticky_lock = asyncio.Lock() + self._last_failure: tuple[str, str] | None = None async def _get_fallback(self, ip: str) -> httpx.AsyncHTTPTransport: async with self._fallback_lock: @@ -136,6 +147,15 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport): transport = self._primary if ip is None else await self._get_fallback(ip) try: response = await transport.handle_async_request(candidate) + if self._last_failure is not None: + failed_path, failure = self._last_failure + self._last_failure = None + logger.info( + "[Telegram] Telegram API transport recovered via %s after %s failed: %s", + ip or _TELEGRAM_API_HOST, + failed_path, + failure, + ) if self._sticky_ip is _UNSET or self._sticky_ip != ip: async with self._sticky_lock: if self._sticky_ip is _UNSET or self._sticky_ip != ip: @@ -148,6 +168,9 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport): last_error = exc if not _is_retryable_connect_error(exc): raise + path = ip or _TELEGRAM_API_HOST + failure = _describe_transport_error(exc) + self._last_failure = (path, failure) if self._sticky_ip is not _UNSET and ip == self._sticky_ip: async with self._sticky_lock: if self._sticky_ip is not _UNSET and self._sticky_ip == ip: @@ -157,9 +180,9 @@ class TelegramFallbackTransport(httpx.AsyncBaseTransport): ip if ip is not None else "api.telegram.org") if ip is None: await self._reset_primary(transport) - logger.warning("[Telegram] Dual-stack api.telegram.org path failed (%s)", exc) + logger.warning("[Telegram] Dual-stack api.telegram.org path failed (%s)", failure) continue - logger.warning("[Telegram] IPv4 Telegram API IP %s failed: %s", ip, exc) + logger.warning("[Telegram] IPv4 Telegram API IP %s failed: %s", ip, failure) await self._reset_fallback(ip) continue if last_error is None: diff --git a/tests/gateway/test_telegram_network.py b/tests/gateway/test_telegram_network.py index 3d624139ed..318b951e7a 100644 --- a/tests/gateway/test_telegram_network.py +++ b/tests/gateway/test_telegram_network.py @@ -177,6 +177,34 @@ class TestFallbackTransport: assert records[0].levelno == logging.WARNING assert "149.154.167.221" in records[0].getMessage() + @pytest.mark.asyncio + async def test_empty_failure_is_diagnostic_and_later_success_logs_recovery( + self, monkeypatch, caplog + ): + import logging + + calls = [] + behavior = { + "149.154.167.220": httpx.ConnectTimeout(""), + "api.telegram.org": httpx.ConnectTimeout(""), + } + monkeypatch.setattr( + tnet.httpx, + "AsyncHTTPTransport", + _fake_transport_factory(calls, behavior), + ) + transport = tnet.TelegramFallbackTransport(["149.154.167.220"]) + + with caplog.at_level(logging.INFO, logger="plugins.platforms.telegram.telegram_network"): + with pytest.raises(httpx.ConnectTimeout): + await transport.handle_async_request(_telegram_request()) + behavior["149.154.167.220"] = "ok" + await transport.handle_async_request(_telegram_request()) + + rendered = " | ".join(record.getMessage() for record in caplog.records) + assert "failed: ConnectTimeout('')" in rendered + assert "transport recovered via 149.154.167.220" in rendered + @pytest.mark.asyncio async def test_sticky_ip_tried_first_but_falls_through_if_stale(self, monkeypatch): diff --git a/tests/gateway/test_telegram_polling_health_confirmation.py b/tests/gateway/test_telegram_polling_health_confirmation.py index 612ad8467d..7f66281c46 100644 --- a/tests/gateway/test_telegram_polling_health_confirmation.py +++ b/tests/gateway/test_telegram_polling_health_confirmation.py @@ -10,6 +10,7 @@ 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 @@ -28,6 +29,7 @@ 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 @@ -37,7 +39,7 @@ class TestPollingHealthConfirmation: 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 "confirmed healthy" in rendered + assert "polling recovered" in rendered assert "generation 1" in rendered assert a._polling_progress_event.is_set() @@ -67,9 +69,21 @@ class TestPollingHealthConfirmation: with caplog.at_level(logging.INFO, logger="plugins.platforms.telegram.adapter"): a._record_polling_progress(2) rendered = " | ".join(rec.getMessage() for rec in caplog.records) - assert "confirmed healthy" in rendered + 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