fix(tradein/tgbot): ограничение частоты на отправителя в мосте поддержки #2543
2 changed files with 176 additions and 0 deletions
|
|
@ -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)
|
||||||
|
|
|
||||||
|
|
@ -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 ──────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue