diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index fd5fc0177eba..0a5ca8de28fc 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -2177,6 +2177,18 @@ async def _polling_heartbeat_loop(self) -> None: HEARTBEAT_INTERVAL = 90 # seconds between probes PROBE_TIMEOUT = 15 # seconds before declaring the path dead + def _trigger_reconnect(err: Exception) -> None: + logger.warning( + "[%s] Polling heartbeat probe failed (%s); triggering reconnect", + self.name, err, + ) + if self._polling_error_task and not self._polling_error_task.done(): + return # reconnect already in progress + loop = asyncio.get_running_loop() + self._polling_error_task = loop.create_task( + self._handle_polling_network_error(err) + ) + while True: try: await asyncio.sleep(HEARTBEAT_INTERVAL) @@ -2204,17 +2216,21 @@ async def _polling_heartbeat_loop(self) -> None: except asyncio.CancelledError: return except (asyncio.TimeoutError, OSError) as probe_err: - logger.warning( - "[%s] Polling heartbeat probe failed (%s); triggering reconnect", - self.name, probe_err, - ) - if self._polling_error_task and not self._polling_error_task.done(): - continue # reconnect already in progress - loop = asyncio.get_running_loop() - self._polling_error_task = loop.create_task( - self._handle_polling_network_error(probe_err) - ) - except Exception: + _trigger_reconnect(probe_err) + except Exception as probe_err: + # python-telegram-bot wraps network failures in TelegramError, + # which is not an OSError. We must explicitly check for PTB + # connectivity errors (NetworkError, TimedOut) and route them + # to the reconnect ladder. + try: + from telegram.error import NetworkError, TimedOut + except ImportError: + # PTB not available or version mismatch; fall through to pass. + pass + else: + if isinstance(probe_err, (NetworkError, TimedOut)): + _trigger_reconnect(probe_err) + continue # Non-connectivity errors (e.g. TelegramError 401) are not # CLOSE-WAIT symptoms — let PTB's own handlers surface them. pass diff --git a/tests/gateway/test_telegram_network_reconnect.py b/tests/gateway/test_telegram_network_reconnect.py index 0016b0e32e75..857b30bb8ca1 100644 --- a/tests/gateway/test_telegram_network_reconnect.py +++ b/tests/gateway/test_telegram_network_reconnect.py @@ -678,6 +678,114 @@ async def telegram_error_wait_for(coro, timeout): adapter._handle_polling_network_error.assert_not_awaited() +@pytest.mark.asyncio +async def test_heartbeat_loop_triggers_reconnect_on_ptb_network_error(): + """python-telegram-bot NetworkError must trigger reconnect, not be swallowed.""" + # Import the real PTB exception classes if available, otherwise mock them. + try: + from telegram.error import NetworkError, TimedOut + except ImportError: + NetworkError = type("NetworkError", (Exception,), {}) + TimedOut = type("TimedOut", (NetworkError,), {}) + + adapter = _make_adapter() + adapter._handle_polling_network_error = AsyncMock() + + mock_app = MagicMock() + adapter._app = mock_app + + sleep_call = 0 + + async def fast_sleep(seconds): + nonlocal sleep_call + sleep_call += 1 + if sleep_call >= 3: + raise asyncio.CancelledError() + + async def network_error_wait_for(coro, timeout): + if asyncio.iscoroutine(coro): + coro.close() + raise NetworkError("Network error from PTB") + + with patch("asyncio.sleep", side_effect=fast_sleep): + with patch("plugins.platforms.telegram.adapter.asyncio.wait_for", side_effect=network_error_wait_for): + await adapter._polling_heartbeat_loop() + + # A reconnect task must have been created for PTB NetworkError. + assert adapter._polling_error_task is not None + + +@pytest.mark.asyncio +async def test_heartbeat_loop_triggers_reconnect_on_ptb_timed_out(): + """python-telegram-bot TimedOut (subclass of NetworkError) must trigger reconnect.""" + try: + from telegram.error import TimedOut + except ImportError: + NetworkError = type("NetworkError", (Exception,), {}) + TimedOut = type("TimedOut", (NetworkError,), {}) + + adapter = _make_adapter() + adapter._handle_polling_network_error = AsyncMock() + + mock_app = MagicMock() + adapter._app = mock_app + + sleep_call = 0 + + async def fast_sleep(seconds): + nonlocal sleep_call + sleep_call += 1 + if sleep_call >= 3: + raise asyncio.CancelledError() + + async def timed_out_wait_for(coro, timeout): + if asyncio.iscoroutine(coro): + coro.close() + raise TimedOut("Timed out from PTB") + + with patch("asyncio.sleep", side_effect=fast_sleep): + with patch("plugins.platforms.telegram.adapter.asyncio.wait_for", side_effect=timed_out_wait_for): + await adapter._polling_heartbeat_loop() + + # A reconnect task must have been created for PTB TimedOut. + assert adapter._polling_error_task is not None + + +@pytest.mark.asyncio +async def test_heartbeat_loop_ignores_ptb_bad_request(): + """PTB BadRequest (non-connectivity) must not trigger reconnect.""" + try: + from telegram.error import BadRequest + except ImportError: + BadRequest = type("BadRequest", (Exception,), {}) + + adapter = _make_adapter() + adapter._handle_polling_network_error = AsyncMock() + + mock_app = MagicMock() + adapter._app = mock_app + + sleep_call = 0 + + async def fast_sleep(seconds): + nonlocal sleep_call + sleep_call += 1 + if sleep_call >= 3: + raise asyncio.CancelledError() + + async def bad_request_wait_for(coro, timeout): + if asyncio.iscoroutine(coro): + coro.close() + raise BadRequest("Bad request from PTB") + + with patch("asyncio.sleep", side_effect=fast_sleep): + with patch("plugins.platforms.telegram.adapter.asyncio.wait_for", side_effect=bad_request_wait_for): + await adapter._polling_heartbeat_loop() + + # No reconnect should have been triggered for BadRequest. + adapter._handle_polling_network_error.assert_not_awaited() + + @pytest.mark.asyncio async def test_heartbeat_loop_exits_on_fatal_error(): """A fatal error short-circuits the loop before probing get_me()."""