fix(tradein/tgbot): per-chat_id rate-limit на входящие сообщения бота
All checks were successful
CI Trade-In / changes (pull_request) Successful in 16s
CI / changes (pull_request) Successful in 16s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m11s
CI / frontend-tests (pull_request) Has been skipped

Однопоточный 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 на то же
окно) — молчать нельзя (клиент решит, что доставлено), но повторные
уведомления на каждое превышение сами стали бы источником флуда.
This commit is contained in:
bot-backend 2026-07-27 00:32:45 +03:00
parent a0647a53a9
commit 8e0479c616
2 changed files with 176 additions and 0 deletions

View file

@ -36,6 +36,21 @@
(entrypoint), не здесь. (entrypoint), не здесь.
E) /start клиенту короткое приветствие МЕРЫ, без зеркалирования в топик 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`-протокол маршрутизирующая логика Персистентность вынесена за `BridgeStorage`-протокол маршрутизирующая логика
(`process_update` и приватные `_handle_*`) не завязана на реальную БД, тестируется (`process_update` и приватные `_handle_*`) не завязана на реальную БД, тестируется
@ -58,6 +73,7 @@ from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from app.core.config import settings from app.core.config import settings
from app.core.ratelimit import SlidingWindowLimiter
from app.core.shutdown import shutdown_requested from app.core.shutdown import shutdown_requested
from app.services.tgbot import web_support_storage from app.services.tgbot import web_support_storage
from app.services.tgbot.client import TelegramApiError, TelegramClient 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". # tg_support_messages.kind): "text | photo | document | video | voice | other".
_KNOWN_KINDS = ("text", "photo", "document", "video", "voice") _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: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ) # #tgsupport-web review M2: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ)
# на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор # на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор
# получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил). # получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил).
@ -407,6 +446,30 @@ async def _handle_private_message(
logger.warning("tgbot bridge: приватное сообщение без message_id — игнор") logger.warning("tgbot bridge: приватное сообщение без message_id — игнор")
return 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): 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) header = _format_topic_header(chat_id, username, first_name, last_name)

View file

@ -243,6 +243,25 @@ def _support_chat_settings(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.setattr(bridge.settings, "telegram_support_topic_id", SUPPORT_TOPIC_ID) 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) @pytest.fixture(autouse=True)
def _stop_patches(): def _stop_patches():
"""Останавливает httpx.AsyncClient monkeypatch после каждого теста (unittest.mock.patch.start() """Останавливает 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 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 ────────────────────────────────────────────── # ── B) реплай оператора → user ──────────────────────────────────────────────