diff --git a/gateway/delivery_ledger.py b/gateway/delivery_ledger.py index 6e22c7a035..af3352754a 100644 --- a/gateway/delivery_ledger.py +++ b/gateway/delivery_ledger.py @@ -15,6 +15,7 @@ from __future__ import annotations import hashlib import logging import os +import re import sqlite3 import threading import time @@ -62,10 +63,34 @@ FLOOD_RETRY_DEFAULT_SECONDS = 60.0 FLOOD_RETRY_CAP_SECONDS = 15 * 60.0 FLOOD_RETRY_SLACK_SECONDS = 2.0 +# The canonical prefix above is what the adapters produce for a flood they decide not to sleep. It is +# not the only shape that reaches this ledger. PTB raises ``RetryAfter``, whose own text reads +# "Flood control exceeded. Retry in 185 seconds", and that text is what lands in ``last_error`` +# whenever a send fails on a path that has not been normalized, or was written by an older build +# before it was. Such a row has to be recognised too: unrecognised, it is treated as an ordinary +# failure, so no redelivery timer is armed for it and a boot sweep claims it immediately instead of +# adopting it until its deadline, spending the one attempt inside the penalty that caused it. +# Matching requires the "flood control" wording as well as the delay, so an unrelated error that +# merely suggests retrying is never mistaken for a flood. +_RAW_FLOOD_RE = re.compile(r"flood control exceeded.*?retry in\s+(\d+(?:\.\d+)?)", re.IGNORECASE) + + +def _raw_flood_wait(text: str) -> Optional[float]: + """Seconds asked for by a flood error still carrying the platform's own wording, else ``None``.""" + match = _RAW_FLOOD_RE.search(text or "") + if not match: + return None + try: + return float(match.group(1)) + except (TypeError, ValueError): + return None + def is_flood_error(error: Any) -> bool: - """True for the adapters' fail-closed flood result (``flood_control:``).""" - return str(error or "").strip().lower().startswith(FLOOD_ERROR_PREFIX) + """True for a flood refusal: the adapters' fail-closed ``flood_control:`` result, or a + row still carrying the platform's own flood wording (see ``_RAW_FLOOD_RE``).""" + text = str(error or "").strip().lower() + return text.startswith(FLOOD_ERROR_PREFIX) or _raw_flood_wait(text) is not None def flood_wait_seconds(error: Any, default: float = FLOOD_RETRY_DEFAULT_SECONDS) -> float: @@ -77,6 +102,12 @@ def flood_wait_seconds(error: Any, default: float = FLOOD_RETRY_DEFAULT_SECONDS) wait = float(text[len(FLOOD_ERROR_PREFIX):].strip()) except ValueError: wait = default + else: + # A row still carrying the platform's own flood wording states its delay just as precisely, + # and the deadline must use it rather than the generic default. + raw = _raw_flood_wait(text) + if raw is not None: + wait = raw return wait if wait > 0 else default diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index bf29b640cc..5149facc72 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -3302,6 +3302,15 @@ class TelegramAdapter(BasePlatformAdapter): _send_attempt + 1, wait, safe_send_error) await asyncio.sleep(wait) continue + # Retries exhausted and still flooded. Fail closed the same way a long penalty + # does: raising here handed the caller the platform's own wording instead of the + # canonical result, so the delivery ledger did not recognise the row as a flood + # refusal, armed no redelivery timer, and the reply waited for the next restart. + logger.warning( + "[%s] Telegram flood control on send persisted across %d attempts; failing " + "closed so the delivery ledger owns the wait: %s", + self.name, _send_attempt + 1, safe_send_error) + return _flood_cap_result(wait) raise async def _retrigger_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]]) -> None: @@ -3514,6 +3523,14 @@ class TelegramAdapter(BasePlatformAdapter): except Exception as retry_err: safe_retry_error = _redact_telegram_error_text(retry_err) logger.error("[%s] Edit retry failed after flood wait: %s", self.name, safe_retry_error) + retry_wait = getattr(retry_err, "retry_after", None) + if retry_wait is not None or "retry after" in str(retry_err).lower(): + # Still flooded after the inline wait, and typically for much longer than the + # first refusal asked for. Fail closed canonically so the ledger arms its + # timer on this delay rather than storing the platform's raw wording, which + # it would read as an ordinary failure and never redeliver. + return _flood_cap_result( + float(retry_wait) if retry_wait is not None else wait) return SendResult(success=False, error=safe_retry_error) safe_error = _redact_telegram_error_text(e) # Transient network errors must not permanently disable progress-message editing. diff --git a/tests/gateway/test_delivery_ledger_flood_retry.py b/tests/gateway/test_delivery_ledger_flood_retry.py index 579bc754e1..9832374341 100644 --- a/tests/gateway/test_delivery_ledger_flood_retry.py +++ b/tests/gateway/test_delivery_ledger_flood_retry.py @@ -142,6 +142,13 @@ async def _drain_timers(runner): ("Forbidden: bot was blocked by the user", False), ("", False), (None, False), + # PTB's own RetryAfter wording, which is what a row carries when it was written by a send path + # that had not been normalized, or by an older build. Unrecognised, such a row gets no timer. + ("Flood control exceeded. Retry in 185 seconds", True), + ("flood control exceeded. retry in 30 seconds", True), + # The flood wording is required, not merely a delay: an unrelated error that suggests retrying + # must not be read as a flood refusal. + ("Bad Gateway; retry in 5 seconds", False), ]) def test_is_flood_error(error, expected): assert dl.is_flood_error(error) is expected @@ -153,6 +160,10 @@ def test_is_flood_error(error, expected): ("flood_control:abc", 60.0), # unreadable -> default ("flood_control:0", 60.0), ("send_path_degraded", 60.0), + # The platform's own wording states the delay just as precisely as the canonical prefix, so the + # deadline must come from it rather than the generic default. + ("Flood control exceeded. Retry in 185 seconds", 185.0), + ("Bad Gateway; retry in 5 seconds", 60.0), ]) def test_flood_wait_seconds(error, expected): assert dl.flood_wait_seconds(error) == pytest.approx(expected) @@ -639,3 +650,64 @@ async def test_other_outcomes_do_not_arm_the_flood_timer(result): await adapter._finalize_delivery_obligation("ob-y", result, _event(), adapter) runner._schedule_flood_redelivery.assert_not_called() + + +# --------------------------------------------------------------------------- +# The adapter normalizes every definite flood refusal, so the ledger sees one shape. +# --------------------------------------------------------------------------- + +def _real_telegram_adapter(): + """A real Telegram adapter with a stub bot, for driving the send and edit flood paths.""" + from plugins.platforms.telegram.adapter import TelegramAdapter + + adapter = TelegramAdapter(PlatformConfig(enabled=True, token="test-token", extra={})) + adapter._bot = MagicMock() + return adapter + + +def test_a_persisted_row_with_the_platforms_own_wording_is_dated_from_it(): + """The boot sweep decides adopt-or-claim from the row's deadline. Read as an ordinary failure, a + raw flood row is claimed at once and spends its one attempt inside the penalty that caused it.""" + raw = "Flood control exceeded. Retry in 185 seconds" + + assert dl.is_flood_error(raw) is True + assert dl.flood_wait_seconds(raw) == pytest.approx(185.0) + assert dl.flood_not_before(T0, raw) == pytest.approx(T0 + 185.0) + + +@pytest.mark.asyncio +async def test_a_short_flood_that_outlives_the_retries_fails_closed_canonically(): + """A wait under the inline cap is slept and retried. When the flood is still there after the last + attempt the send used to raise, handing the caller the platform's wording instead of the + canonical result, so no redelivery timer was armed and the reply waited for the next restart.""" + from telegram.error import RetryAfter + + adapter = _real_telegram_adapter() + adapter._send_chunk_markdown_or_plain = AsyncMock(side_effect=RetryAfter(1)) + + result = await adapter._send_chunk_with_retries( + "5230977008", "the final answer", 0, None, None, None, False, + adapter._telegram_error_types()) + + assert isinstance(result, SendResult) + assert result.success is False + assert result.error == "flood_control:1.0" + assert dl.is_flood_error(result.error) is True + assert result.retry_after == pytest.approx(1.0) + assert adapter._send_chunk_markdown_or_plain.await_count == 3 + + +@pytest.mark.asyncio +async def test_an_edit_still_flooded_after_the_inline_wait_fails_closed_canonically(): + """The edit path sleeps a short wait once and retries. A retry that is refused again, typically + for far longer, must report the new delay canonically rather than as raw text.""" + from telegram.error import RetryAfter + + adapter = _real_telegram_adapter() + adapter._edit_text = AsyncMock(side_effect=[RetryAfter(1), RetryAfter(185)]) + + result = await adapter.edit_message("5230977008", "900", "the final answer") + + assert result.success is False + assert result.error == "flood_control:185.0" + assert dl.flood_wait_seconds(result.error) == pytest.approx(185.0)