From dc4de143778ab2e156063513fe64d6b41a675860 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Sat, 30 May 2026 22:58:16 -0700 Subject: [PATCH] fix(telegram): retry on httpx pool timeout instead of dropping the send (#35664) When PTB's general httpx pool is exhausted, it converts httpx.PoolTimeout into telegram.error.TimedOut whose message states the request was *not* sent to Telegram. The send retry loop treated all non-connect TimedOut as non-retryable, so a pool timeout raised immediately, skipped all 3 retry attempts, and was returned as retryable=False -- silently dropping the message (agent responses, cron reports, etc.). A pool timeout means the request never left the process, making it the safest case to retry. Add _looks_like_pool_timeout() and treat it like a connect timeout in both the in-loop retry decision and the outer retryable determination, so pool timeouts flow through the existing backoff loop and stay retryable on exhaustion. Reported-by: q3874758 (#35610) --- gateway/platforms/telegram.py | 49 +++++++++++++++-- .../gateway/test_telegram_thread_fallback.py | 54 +++++++++++++++++++ 2 files changed, 99 insertions(+), 4 deletions(-) diff --git a/gateway/platforms/telegram.py b/gateway/platforms/telegram.py index 7b4d00e81..4ab36be4a 100644 --- a/gateway/platforms/telegram.py +++ b/gateway/platforms/telegram.py @@ -840,6 +840,41 @@ class TelegramAdapter(BasePlatformAdapter): stack.append(context) return False + @staticmethod + def _looks_like_pool_timeout(error: Exception) -> bool: + """Return True when a Telegram TimedOut wraps an httpx pool timeout. + + PTB converts ``httpx.PoolTimeout`` into ``telegram.error.TimedOut`` with + a message that explicitly states the request was *not* sent + (``"Pool timeout: All connections in the connection pool are occupied. + Request was *not* sent to Telegram."``). Because the request never left + the process, re-sending is safe and cannot duplicate -- the opposite of + a generic TimedOut, which may have reached Telegram. We match the + wrapped ``httpx.PoolTimeout`` class as well as the message string so the + check survives PTB message-wording changes. + """ + seen: set[int] = set() + stack: list[BaseException] = [error] + while stack: + cur = stack.pop() + ident = id(cur) + if ident in seen: + continue + seen.add(ident) + name = cur.__class__.__name__.lower() + text = str(cur).lower() + if "pooltimeout" in name or "pool timeout" in text or ( + "connection pool" in text and "occupied" in text + ): + return True + cause = getattr(cur, "__cause__", None) + context = getattr(cur, "__context__", None) + if cause is not None: + stack.append(cause) + if context is not None: + stack.append(context) + return False + def _coerce_bool_extra(self, key: str, default: bool = False) -> bool: value = self.config.extra.get(key) if getattr(self.config, "extra", None) else None if value is None: @@ -2001,11 +2036,15 @@ class TelegramAdapter(BasePlatformAdapter): # TimedOut is also a subclass of NetworkError. A # generic timeout may have reached Telegram, so don't # retry; a wrapped ConnectTimeout means no connection - # was established, so retrying is safe. + # was established, so retrying is safe. A pool timeout + # (httpx pool exhausted) is explicitly "not sent to + # Telegram" -- retrying through the loop is safe and + # prevents silent drops when the pool frees up. if ( _TimedOut and isinstance(send_err, _TimedOut) and not self._looks_like_connect_timeout(send_err) + and not self._looks_like_pool_timeout(send_err) ): raise if _send_attempt < 2: @@ -2065,12 +2104,14 @@ class TelegramAdapter(BasePlatformAdapter): return SendResult(success=False, error="message_too_long") # TimedOut usually means the request may have reached Telegram — # mark as non-retryable so _send_with_retry() doesn't re-send. - # Exception: wrapped ConnectTimeout, where no connection was - # established; retrying is safe and prevents silent drops. + # Exceptions: a wrapped ConnectTimeout (no connection established) + # and an httpx pool timeout (request explicitly not sent) -- both + # are safe to re-send and must not be silently dropped. _to = locals().get("_TimedOut") is_timeout = (_to and isinstance(e, _to)) or "timed out" in err_str is_connect_timeout = self._looks_like_connect_timeout(e) - return SendResult(success=False, error=str(e), retryable=(is_connect_timeout or not is_timeout)) + is_pool_timeout = self._looks_like_pool_timeout(e) + return SendResult(success=False, error=str(e), retryable=(is_connect_timeout or is_pool_timeout or not is_timeout)) async def send_or_update_status( self, diff --git a/tests/gateway/test_telegram_thread_fallback.py b/tests/gateway/test_telegram_thread_fallback.py index ddbd8a459..5f56baebc 100644 --- a/tests/gateway/test_telegram_thread_fallback.py +++ b/tests/gateway/test_telegram_thread_fallback.py @@ -1278,6 +1278,60 @@ async def test_send_marks_wrapped_connect_timeout_retryable_after_exhaustion(): assert attempt[0] == 3 +@pytest.mark.asyncio +async def test_send_retries_pool_timeout(): + """Retry TimedOut when it is an httpx pool-timeout (request not sent). + + PTB wraps ``httpx.PoolTimeout`` into ``TimedOut`` with a message that + explicitly states the request was *not* sent to Telegram. Re-sending is + safe and prevents a silent drop when the pool frees up. + """ + adapter = _make_adapter() + + attempt = [0] + + async def mock_send_message(**kwargs): + attempt[0] += 1 + if attempt[0] < 3: + raise FakeTimedOut( + "Pool timeout: All connections in the connection pool are " + "occupied. Request was *not* sent to Telegram. Consider " + "adjusting the connection pool size or the pool timeout." + ) + return SimpleNamespace(message_id=202) + + adapter._bot = SimpleNamespace(send_message=mock_send_message) + + result = await adapter.send(chat_id="123", content="test message") + + assert result.success is True + assert result.message_id == "202" + assert attempt[0] == 3 + + +@pytest.mark.asyncio +async def test_send_marks_pool_timeout_retryable_after_exhaustion(): + """Pool timeout that never clears stays retryable for outer retry handling.""" + adapter = _make_adapter() + + attempt = [0] + + async def mock_send_message(**kwargs): + attempt[0] += 1 + raise FakeTimedOut( + "Pool timeout: All connections in the connection pool are occupied. " + "Request was *not* sent to Telegram." + ) + + adapter._bot = SimpleNamespace(send_message=mock_send_message) + + result = await adapter.send(chat_id="123", content="test message") + + assert result.success is False + assert result.retryable is True + assert attempt[0] == 3 + + @pytest.mark.asyncio async def test_thread_fallback_only_fires_once(): """After clearing thread_id, subsequent chunks should also use None."""