fix(telegram): expose polling transport recovery

This commit is contained in:
fangliquan
2026-09-15 05:05:41 +08:00
committed by Teknium
parent 2298dc8122
commit e1e943a9df
4 changed files with 88 additions and 7 deletions
+19 -3
View File
@@ -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
+25 -2
View File
@@ -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:
+28
View File
@@ -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):
@@ -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