Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 22 additions & 6 deletions plugins/platforms/telegram/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -2349,8 +2349,13 @@ async def _verify_polling_after_reconnect(self) -> None:
wedged httpx pool fails this probe; a healthy one returns
well under the timeout.

On any failure, re-enter the reconnect ladder so the existing
MAX_NETWORK_RETRIES path can ultimately escalate to fatal-error.
On connectivity failure, re-enter the reconnect ladder (via
``_schedule_polling_recovery`` so the in-flight guard and task
bookkeeping apply) and let the existing MAX_NETWORK_RETRIES path
ultimately escalate to fatal-error. Auth/validation failures
(``InvalidToken``, ``BadRequest``, ...) are not connectivity symptoms
and must not trigger reconnect churn — same policy as the heartbeat
loop and the pending-update probe (#62098, #63243).
"""
HEARTBEAT_PROBE_DELAY = 60
PROBE_TIMEOUT = 10
Expand All @@ -2364,20 +2369,31 @@ async def _verify_polling_after_reconnect(self) -> None:
"[%s] Updater not running %ds after reconnect — treating as wedged",
self.name, HEARTBEAT_PROBE_DELAY,
)
await self._handle_polling_network_error(
RuntimeError("Updater not running after reconnect heartbeat")
self._schedule_polling_recovery(
RuntimeError("Updater not running after reconnect heartbeat"),
reason="post-reconnect probe: updater not running",
)
return

try:
await asyncio.wait_for(self._app.bot.get_me(), PROBE_TIMEOUT)
self._send_path_degraded = False
except Exception as probe_err:
if not self._looks_like_network_error(probe_err):
logger.warning(
"[%s] Post-reconnect probe hit a non-connectivity error"
" (not retrying): %s",
self.name, _redact_telegram_error_text(probe_err),
)
return
logger.warning(
"[%s] Polling heartbeat probe failed %ds after reconnect: %s",
self.name, HEARTBEAT_PROBE_DELAY, probe_err,
self.name, HEARTBEAT_PROBE_DELAY,
_redact_telegram_error_text(probe_err),
)
self._schedule_polling_recovery(
probe_err, reason="post-reconnect probe failure"
)
await self._handle_polling_network_error(probe_err)

def _disarm_ptb_retry_loop(self) -> None:
"""Synchronously stop PTB's internal polling retry loop.
Expand Down
74 changes: 74 additions & 0 deletions tests/gateway/test_telegram_network_reconnect.py
Original file line number Diff line number Diff line change
Expand Up @@ -385,6 +385,11 @@ async def test_heartbeat_probe_reenters_ladder_when_updater_not_running():
await adapter._verify_polling_after_reconnect()

mock_app.bot.get_me.assert_not_called()
# Recovery is scheduled through _schedule_polling_recovery (#63243), so
# the ladder runs as the tracked _polling_error_task.
task = adapter._polling_error_task
assert task is not None
await task
adapter._handle_polling_network_error.assert_awaited_once()
err = adapter._handle_polling_network_error.await_args.args[0]
assert isinstance(err, RuntimeError)
Expand Down Expand Up @@ -421,6 +426,9 @@ async def fast_wait_for(coro, timeout):
with patch("plugins.platforms.telegram.adapter.asyncio.wait_for", new=fast_wait_for):
await adapter._verify_polling_after_reconnect()

task = adapter._polling_error_task
assert task is not None
await task
adapter._handle_polling_network_error.assert_awaited_once()


Expand All @@ -445,12 +453,78 @@ async def test_heartbeat_probe_reenters_ladder_on_get_me_network_error():
with patch("asyncio.sleep", new_callable=AsyncMock):
await adapter._verify_polling_after_reconnect()

task = adapter._polling_error_task
assert task is not None
# _schedule_polling_recovery must also register the ladder in
# _background_tasks so a failed recovery isn't silently GC'd.
assert task in adapter._background_tasks
await task
adapter._handle_polling_network_error.assert_awaited_once()
assert isinstance(
adapter._handle_polling_network_error.await_args.args[0], ConnectionError
)


@pytest.mark.asyncio
async def test_heartbeat_probe_ignores_auth_errors():
"""
Auth/validation failures from the post-reconnect probe must not enter the
network-reconnect ladder (#63243): a revoked token would otherwise churn
through stop/drain/start_polling cycles that mask the real failure.
"""
adapter = _make_adapter()

mock_updater = MagicMock()
mock_updater.running = True

# Name-shaped like PTB's InvalidToken; _looks_like_network_error excludes
# it by class name, matching real PTB semantics.
invalid_token = type("InvalidToken", (Exception,), {})("token revoked")

mock_app = MagicMock()
mock_app.updater = mock_updater
mock_app.bot.get_me = AsyncMock(side_effect=invalid_token)
adapter._app = mock_app

adapter._handle_polling_network_error = AsyncMock()

with patch("asyncio.sleep", new_callable=AsyncMock):
await adapter._verify_polling_after_reconnect()

assert adapter._polling_error_task is None
adapter._handle_polling_network_error.assert_not_awaited()


@pytest.mark.asyncio
async def test_heartbeat_probe_defers_to_inflight_recovery():
"""
A probe failure while another recovery is mid-flight must not start a
second concurrent stop/drain/start_polling sequence (#63243) — overlapping
recoveries produce dueling getUpdates sessions (self-inflicted 409s).
"""
adapter = _make_adapter()

mock_updater = MagicMock()
mock_updater.running = True

mock_app = MagicMock()
mock_app.updater = mock_updater
mock_app.bot.get_me = AsyncMock(side_effect=ConnectionError("pool wedged"))
adapter._app = mock_app

inflight = MagicMock()
inflight.done.return_value = False
adapter._polling_error_task = inflight

adapter._handle_polling_network_error = AsyncMock()

with patch("asyncio.sleep", new_callable=AsyncMock):
await adapter._verify_polling_after_reconnect()

assert adapter._polling_error_task is inflight
adapter._handle_polling_network_error.assert_not_awaited()


@pytest.mark.asyncio
async def test_heartbeat_probe_skips_when_already_fatal():
"""
Expand Down
Loading