Skip to content
Closed
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
38 changes: 27 additions & 11 deletions plugins/platforms/telegram/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down
108 changes: 108 additions & 0 deletions tests/gateway/test_telegram_network_reconnect.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()."""
Expand Down
Loading