From 5495c29cf8f3e5c82aeb2e849a2f0bcff39e5995 Mon Sep 17 00:00:00 2001 From: Alexander Russell Date: Mon, 7 Sep 2026 11:51:52 +0300 Subject: [PATCH] fix(gateway): recognise every definite flood refusal, not only the canonical one The redelivery hook keyed on the canonical flood_control: result, so two real refusals slipped past it and armed no timer, leaving the reply for the next restart. A short wait that outlived the send retries raised instead of failing closed, and an edit refused again after its inline wait returned the platform's raw text. Both now fail closed canonically, the second carrying the new delay rather than the first refusal's. The ledger also accepts a row still carrying the platform's own wording, so a row persisted by an unnormalized path is dated from the delay it states instead of the generic default. Without that a boot sweep claims it at once and spends its one attempt inside the penalty. Matching requires the flood wording as well as a delay, so an unrelated retry suggestion is never read as a flood. Six new assertions fail without this change. 770 passed across the ledger, Telegram, send-retry and queued suites. --- gateway/delivery_ledger.py | 35 ++++++++- plugins/platforms/telegram/adapter.py | 17 +++++ .../test_delivery_ledger_flood_retry.py | 72 +++++++++++++++++++ 3 files changed, 122 insertions(+), 2 deletions(-) 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)