From 8e0479c616740f3a2e25a3fe57c9479e87c7da7e Mon Sep 17 00:00:00 2001 From: bot-backend Date: Mon, 27 Jul 2026 00:32:45 +0300 Subject: [PATCH] =?UTF-8?q?fix(tradein/tgbot):=20per-chat=5Fid=20rate-limi?= =?UTF-8?q?t=20=D0=BD=D0=B0=20=D0=B2=D1=85=D0=BE=D0=B4=D1=8F=D1=89=D0=B8?= =?UTF-8?q?=D0=B5=20=D1=81=D0=BE=D0=BE=D0=B1=D1=89=D0=B5=D0=BD=D0=B8=D1=8F?= =?UTF-8?q?=20=D0=B1=D0=BE=D1=82=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Однопоточный long-polling воркер обрабатывал апдейты строго последовательно, и ничто не мешало одному флудящему клиенту слать поток сообщений: каждое зеркалировалось (copyMessage) в support-топик, а групповой Telegram-лимит ~20 msg/min общий на ВСЕХ клиентов сразу — превышение даёт 429 с ожиданием 30-60с, за которое воркер не может обработать ни одного апдейта от кого-либо ещё. Переиспользован app.core.ratelimit.SlidingWindowLimiter (тот же примитив, что уже применён для веб-чата поддержки, app/api/v1/support.py) — ключ здесь TELEGRAM chat_id отправителя, лимит заметно ниже группового Telegram-порога (5 msg/60s). Сообщения сверх бюджета не зеркалируются (и не пишутся в tg_support_messages — маршрутизировать ответ всё равно нечего без topic_message_id), клиент получает явное уведомление о недоставке РОВНО один раз за окно (второй лимитер с limit=1 на то же окно) — молчать нельзя (клиент решит, что доставлено), но повторные уведомления на каждое превышение сами стали бы источником флуда. --- .../backend/app/services/tgbot/bridge.py | 63 ++++++++++ .../tests/services/tgbot/test_bridge.py | 113 ++++++++++++++++++ 2 files changed, 176 insertions(+) diff --git a/tradein-mvp/backend/app/services/tgbot/bridge.py b/tradein-mvp/backend/app/services/tgbot/bridge.py index 6cbc8f3a..fc49aa68 100644 --- a/tradein-mvp/backend/app/services/tgbot/bridge.py +++ b/tradein-mvp/backend/app/services/tgbot/bridge.py @@ -36,6 +36,21 @@ (entrypoint), не здесь. E) /start клиенту → короткое приветствие МЕРЫ, без зеркалирования в топик (команда — не содержательное обращение, не должна засорять топик). + F) Флуд-лимит на отправителя (низкий приоритет, per-chat_id): воркер + long-polling однопоточный и обрабатывает апдейты СТРОГО последовательно, а + Telegram ограничивает саму support-группу ~20 сообщениями/минуту — ОДНИМ + бюджетом на ВСЕХ клиентов разом (зеркала + шапки + ответы оператора). + Превышение — 429 с ожиданием 30-60с, на которые воркер не может обработать + НИ ОДНОГО следующего апдейта — один флудящий клиент подвешивает доставку + всем остальным. `_flood_limiter` (тот же `SlidingWindowLimiter`, что и + веб-чат поддержки, ключ — TELEGRAM chat_id) режет per-sender поток заметно + ниже группового лимита; сообщения сверх бюджета НЕ зеркалируются (иначе + сам факт мирроринга уже съедает групповой бюджет, который мы и защищаем) и + НЕ пишутся в tg_support_messages (нечего маршрутизировать без + topic_message_id). Клиент получает уведомление, что сообщение НЕ + доставлено (молчать нельзя — иначе клиент решит, что оператор его получил), + но не чаще ОДНОГО РАЗА за то же окно (`_flood_notify_limiter`, limit=1) — + иначе само уведомление стало бы вторым источником флуда. Персистентность вынесена за `BridgeStorage`-протокол — маршрутизирующая логика (`process_update` и приватные `_handle_*`) не завязана на реальную БД, тестируется @@ -58,6 +73,7 @@ from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.orm import Session from app.core.config import settings +from app.core.ratelimit import SlidingWindowLimiter from app.core.shutdown import shutdown_requested from app.services.tgbot import web_support_storage from app.services.tgbot.client import TelegramApiError, TelegramClient @@ -87,6 +103,29 @@ SERVICE_UNAVAILABLE_TEXT = ( # tg_support_messages.kind): "text | photo | document | video | voice | other". _KNOWN_KINDS = ("text", "photo", "document", "video", "voice") +# (низкий приоритет, флуд-защита) — см. пункт F) в докстринге модуля. Порог +# НАМЕРЕННО заметно ниже группового лимита Telegram (~20 msg/min): бюджет +# делится с шапками-идентификациями и ответами оператора, и с другими +# одновременными клиентами — щедрый лимит одного отправителя всё равно упёрся +# бы в общий групповой 429. Тот же примитив, что и веб-чат поддержки +# (app/api/v1/support.py `_send_limiter`), ключ здесь — TELEGRAM chat_id +# отправителя (не username — у Telegram-клиента username может отсутствовать). +_FLOOD_LIMIT = 5 +_FLOOD_WINDOW_S = 60.0 +_flood_limiter = SlidingWindowLimiter(limit=_FLOOD_LIMIT, window_s=_FLOOD_WINDOW_S) + +# Уведомление о флуде — не чаще ОДНОГО раза за то же окно, иначе само +# уведомление стало бы вторым источником флуда. Отдельный лимитер с limit=1 на +# то же окно: `check()` возвращает None (и фиксирует попытку) ровно один раз за +# окно, дальше молчит до его истечения — без отдельной структуры "когда в +# последний раз уведомляли". +_flood_notify_limiter = SlidingWindowLimiter(limit=1, window_s=_FLOOD_WINDOW_S) + +FLOOD_LIMITED_TEXT = ( + "Сообщение не доставлено — вы отправляете сообщения слишком часто. " + "Пожалуйста, подождите немного и напишите ещё раз." +) + # #tgsupport-web review M2: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ) # на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор # получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил). @@ -407,6 +446,30 @@ async def _handle_private_message( logger.warning("tgbot bridge: приватное сообщение без message_id — игнор") return + # F) Флуд-лимит на отправителя — peek БЕЗ расхода бюджета (тот же паттерн, + # что `_send_limiter` в app/api/v1/support.py: под лимитом ниже сразу + # `.record()`-им попытку). Над лимитом — НЕ зеркалируем (иначе сам мирроринг + # уже съедает групповой Telegram-бюджет, который лимит и защищает) и НЕ + # пишем в tg_support_messages (без topic_message_id маршрутизировать ответ + # всё равно нечего). + flood_key = str(chat_id) # SlidingWindowLimiter — ключ str (см. app/core/ratelimit.py) + if _flood_limiter.retry_after(flood_key) is not None: + logger.warning( + "tgbot bridge: chat_id=%d превысил флуд-лимит (%d msg/%.0fs) — " + "сообщение НЕ зеркалируется в топик (защита группового Telegram-лимита)", + chat_id, + _FLOOD_LIMIT, + _FLOOD_WINDOW_S, + ) + # Уведомляем клиента, что сообщение НЕ доставлено (молчать нельзя — + # иначе клиент решит, что оператор его получил), но не чаще одного раза + # за окно — `_flood_notify_limiter.check()` возвращает None (и сам + # фиксирует попытку) ровно один раз за окно. + if _flood_notify_limiter.check(flood_key) is None: + await client.send_message(chat_id=chat_id, text=FLOOD_LIMITED_TEXT) + return + _flood_limiter.record(flood_key) + # Шапка — только на первое сообщение клиента за окно, иначе топик засоряется. if not storage.had_recent_inbound(chat_id, window_seconds=_HEADER_THROTTLE_WINDOW_S): header = _format_topic_header(chat_id, username, first_name, last_name) diff --git a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py index fdc82334..601f4f32 100644 --- a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py +++ b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py @@ -243,6 +243,25 @@ def _support_chat_settings(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setattr(bridge.settings, "telegram_support_topic_id", SUPPORT_TOPIC_ID) +@pytest.fixture(autouse=True) +def _reset_flood_limiters(monkeypatch: pytest.MonkeyPatch) -> None: + """F) `_flood_limiter`/`_flood_notify_limiter` — module-level singletons (тот же + паттерн, что `_send_limiter` в app/api/v1/support.py); большинство тестов в + этом файле шлют сообщения от одного и того же chat_id=555, поэтому без сброса + накопленные хиты одного теста бы протекали в следующий и ломали его + предположения (тест флуда должен видеть ЧИСТЫЙ бюджет).""" + monkeypatch.setattr( + bridge, + "_flood_limiter", + bridge.SlidingWindowLimiter(limit=bridge._FLOOD_LIMIT, window_s=bridge._FLOOD_WINDOW_S), + ) + monkeypatch.setattr( + bridge, + "_flood_notify_limiter", + bridge.SlidingWindowLimiter(limit=1, window_s=bridge._FLOOD_WINDOW_S), + ) + + @pytest.fixture(autouse=True) def _stop_patches(): """Останавливает httpx.AsyncClient monkeypatch после каждого теста (unittest.mock.patch.start() @@ -427,6 +446,100 @@ async def test_private_message_notifies_client_when_support_chat_unset( assert storage.get_offset() == 14 +# ── F) флуд-лимит на отправителя (низкий приоритет) ───────────────────────── + + +async def test_private_message_flood_limit_blocks_excess_and_notifies_once() -> None: + """Больше `_FLOOD_LIMIT` сообщений от ОДНОГО chat_id за окно — зеркалирование + сверх лимита отключается (никакого copyMessage, никакой записи в + tg_support_messages — маршрутизировать ответ всё равно нечего без + topic_message_id). Клиент получает уведомление о недоставке РОВНО один раз + за окно, а не на каждое следующее превышение — иначе само уведомление стало + бы вторым источником флуда.""" + calls: list[tuple[str, dict[str, Any]]] = [] + client = _make_client({"copyMessage": {"message_id": 900}}, calls) + storage = FakeBridgeStorage() + + update_id = 100 + for i in range(bridge._FLOOD_LIMIT): + update = {"update_id": update_id, "message": _private_message(message_id=i + 1)} + await bridge.process_update(update, client, storage) + update_id += 1 + + # Ровно _FLOOD_LIMIT сообщений прошли мирроринг: первое — шапка + зеркало, + # остальные — только зеркало. + mirrored_calls = [m for m, _ in calls if m == "copyMessage"] + assert len(mirrored_calls) == bridge._FLOOD_LIMIT + assert len(storage.messages) == bridge._FLOOD_LIMIT + + calls.clear() + over_limit_update = { + "update_id": update_id, + "message": _private_message(message_id=bridge._FLOOD_LIMIT + 1), + } + await bridge.process_update(over_limit_update, client, storage) + update_id += 1 + + # Сверх лимита — НЕ зеркалируется, НЕ пишется в лог переписки, клиент + # получает уведомление о недоставке (не тихий игнор — клиент не должен + # решить, что оператор получил сообщение). + assert len(calls) == 1 + method, payload = calls[0] + assert method == "sendMessage" + assert payload["chat_id"] == 555 + assert payload["text"] == bridge.FLOOD_LIMITED_TEXT + assert len(storage.messages) == bridge._FLOOD_LIMIT + + calls.clear() + second_over_limit_update = { + "update_id": update_id, + "message": _private_message(message_id=bridge._FLOOD_LIMIT + 2), + } + await bridge.process_update(second_over_limit_update, client, storage) + + # Повторное превышение в ТОМ ЖЕ окне — уведомление подавлено (не второй + # источник флуда), никаких Telegram-вызовов вообще. + assert calls == [] + assert len(storage.messages) == bridge._FLOOD_LIMIT + + +async def test_private_message_flood_limit_does_not_block_other_client() -> None: + """Флуд-лимит — per-chat_id: клиент А исчерпал свой бюджет, но клиент Б + (другой chat_id) продолжает получать зеркалирование как обычно — один + флудящий клиент не блокирует доставку сообщений остальным (сама суть + задачи — воркер однопоточный, но лимит не даёт флудеру монополизировать + его через Telegram 429).""" + calls: list[tuple[str, dict[str, Any]]] = [] + client = _make_client({"copyMessage": {"message_id": 901}}, calls) + storage = FakeBridgeStorage() + + flooding_chat_id = 555 + update_id = 300 + for i in range(bridge._FLOOD_LIMIT + 2): + update = { + "update_id": update_id, + "message": _private_message(chat_id=flooding_chat_id, message_id=i + 1), + } + await bridge.process_update(update, client, storage) + update_id += 1 + + calls.clear() + + other_chat_id = 777001 + other_update = { + "update_id": update_id, + "message": _private_message(chat_id=other_chat_id, message_id=1, username="another_client"), + } + await bridge.process_update(other_update, client, storage) + + methods = [m for m, _ in calls] + # Другой клиент получает шапку (первое обращение) + зеркало как обычно — + # флуд первого клиента на него не влияет. + assert methods == ["sendMessage", "copyMessage"] + mirror_call = calls[1][1] + assert mirror_call["from_chat_id"] == other_chat_id + + # ── B) реплай оператора → user ──────────────────────────────────────────────