From 6433477f7ce4a3a8fbfaa0aef59fb14582b54f67 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 14:59:54 +0300 Subject: [PATCH 1/3] =?UTF-8?q?fix(tgbot):=20=D0=BF=D1=80=D0=BE=D0=B2?= =?UTF-8?q?=D0=B5=D1=80=D0=BA=D0=B0=20=D1=82=D0=B5=D0=BC=D1=8B=20=D0=BF?= =?UTF-8?q?=D1=80=D0=B8=20=D1=81=D1=82=D0=B0=D1=80=D1=82=D0=B5=20+=20?= =?UTF-8?q?=D0=BE=D0=B1=D1=89=D0=B8=D0=B9=20rate=20limit=20=D0=BD=D0=B0=20?= =?UTF-8?q?=D0=B3=D1=80=D1=83=D0=BF=D0=BF=D1=83=20(#3471)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Бот падал молча на каждом сообщении, если тему форума удалили/переименовали: узнавали об этом только по отсутствию сообщений у людей. tgbot_main теперь один раз на старте проверяет getChat + typing-индикатор с message_thread_id (единственный способ Bot API провалидировать message_thread_id без создания видимого сообщения) и громко пишет error при отказе, не роняя процесс. Второе: лимит Telegram (~20 msg/min) общий на всю группу, все темы делят бюджет — всплеск GlitchTip-алертов вместе с потоком поддержки в ту же группу уже давал 429 и терял сообщения. TelegramGroupRateLimiter — скользящее окно per-chat_id (НЕ per-теме) с asyncio.Lock на чат, встроен прямо в TelegramClient._request перед _post, поэтому считает все отправки независимо от relay/прямого пути и без изменений в support.py/glitchtip.py (они уже идут через общий клиент). Порог настраивается через TELEGRAM_GROUP_RATE_LIMIT_PER_MINUTE (дефолт 18, чуть ниже потолка площадки). Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG --- tradein-mvp/backend/app/core/config.py | 9 + .../backend/app/services/tgbot/client.py | 201 +++++++++++++ .../backend/app/services/tgbot/shared.py | 1 + tradein-mvp/backend/app/tgbot_main.py | 18 +- .../test_topic_check_and_group_rate_limit.py | 263 ++++++++++++++++++ 5 files changed, 491 insertions(+), 1 deletion(-) create mode 100644 tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py diff --git a/tradein-mvp/backend/app/core/config.py b/tradein-mvp/backend/app/core/config.py index 46a15b08..4a930171 100644 --- a/tradein-mvp/backend/app/core/config.py +++ b/tradein-mvp/backend/app/core/config.py @@ -1287,6 +1287,15 @@ class Settings(BaseSettings): telegram_alerts_chat_id: int = Field(default=0, validation_alias="TELEGRAM_ALERTS_CHAT_ID") telegram_alerts_topic_id: int = Field(default=0, validation_alias="TELEGRAM_ALERTS_TOPIC_ID") + # Общий (не per-тему) лимит частоты отправки в ОДНУ группу — Telegram + # считает ~20 сообщений/минуту на группу суммарно по всем её темам (#3471: + # всплеск GlitchTip-алертов + поток поддержки в ту же группу давали 429 и + # потерю сообщений). Дефолт — чуть ниже площадочного потолка. 0/отрицательное + # значение выключает лимитер (см. `TelegramGroupRateLimiter.acquire`). + telegram_group_rate_limit_per_minute: int = Field( + default=18, validation_alias="TELEGRAM_GROUP_RATE_LIMIT_PER_MINUTE" + ) + # ── Ретранслятор Bot API через Beget (#3471) ───────────────────────────── # Замер 12.09.2026, оба хоста в одни и те же минуты: `getMe` с Selectel — 9 # успешных из 12 (три ConnectTimeout), TCP-443 до адреса Selectel — 5/6, TCP-443 diff --git a/tradein-mvp/backend/app/services/tgbot/client.py b/tradein-mvp/backend/app/services/tgbot/client.py index eb8c2777..5f0c93ab 100644 --- a/tradein-mvp/backend/app/services/tgbot/client.py +++ b/tradein-mvp/backend/app/services/tgbot/client.py @@ -43,6 +43,7 @@ from __future__ import annotations import asyncio import logging +import time from typing import Any import httpx @@ -54,6 +55,94 @@ _DEFAULT_RETRY_AFTER_S = 5.0 _MAX_BACKOFF_S = 30.0 _DEFAULT_MAX_RETRIES = 5 +# Telegram документирует ~20 сообщений/минуту на ОДНУ группу (общий лимит на +# все темы супергруппы разом, не на тему по отдельности — превышение даёт 429 +# на ЛЮБОЕ следующее сообщение в группу, кто бы его ни отправлял). До #3471 +# лимита на нашей стороне не было вовсе: всплеск GlitchTip-алертов + поток +# сообщений поддержки в ту же группу (разные темы, общий чат) укладывались в +# 429 и теряли сообщения — retry в `_request` уважает `retry_after`, но не +# предотвращает сам всплеск. `_DEFAULT_GROUP_RATE_LIMIT_PER_MINUTE` — чуть +# ниже площадочного потолка, с запасом на неточность скользящего окна и на то, +# что сама площадка не обязана быть педантичной ровно к 20-й отправке. +_DEFAULT_GROUP_RATE_LIMIT_PER_MINUTE = 18 +_RATE_LIMIT_WINDOW_S = 60.0 + + +class TelegramGroupRateLimiter: + """Общий (per-`chat_id`, НЕ per-теме) ограничитель частоты отправки в группу. + + Зачем ключ — `chat_id`, а не `(chat_id, message_thread_id)`: лимит Telegram + считается на группу целиком, все темы супергруппы делят один бюджет. + Ограничитель с ключом по теме позволил бы двум темам суммарно превысить + лимит группы и всё равно поймать 429 — ровно баг, который здесь чинится. + + Реализация — скользящее окно (список меток времени последних отправок за + `_RATE_LIMIT_WINDOW_S`), а не токен-бакет с фиксированным пополнением: + окно точнее соответствует тому, как Telegram считает лимит («N сообщений + за последние 60 секунд», а не «N сообщений в календарную минуту»). + + Конкурентность (asyncio, один процесс, несколько отправителей): на каждый + `chat_id` — свой `asyncio.Lock`. Лок держится ВКЛЮЧАЯ время ожидания + (`asyncio.sleep`), а не только на чтение/запись счётчика — это осознанно: + цель не просто «не гонять счётчик без гонки», а ФАКТИЧЕСКИ сериализовать + отправителей в этот чат, чтобы они не просыпались все разом по истечении + окна и не били по лимиту повторно. Другие `chat_id` (например, тема + поддержки и тема алертов лежат в РАЗНЫХ группах) не блокируют друг друга — + у каждого свой лок и свой список меток. + + Регистр локов защищён отдельным `asyncio.Lock` только на момент + создания записи — сам подсчёт/сон идёт уже под персональным локом чата. + """ + + def __init__( + self, + max_per_window: int = _DEFAULT_GROUP_RATE_LIMIT_PER_MINUTE, + window_s: float = _RATE_LIMIT_WINDOW_S, + ) -> None: + self._max_per_window = max_per_window + self._window_s = window_s + self._registry_lock = asyncio.Lock() + self._locks: dict[int, asyncio.Lock] = {} + self._sent_at: dict[int, list[float]] = {} + + async def _lock_for(self, chat_id: int) -> asyncio.Lock: + async with self._registry_lock: + lock = self._locks.get(chat_id) + if lock is None: + lock = asyncio.Lock() + self._locks[chat_id] = lock + return lock + + async def acquire(self, chat_id: int) -> None: + """Блокируется, пока в окне `_window_s` для `chat_id` есть свободный слот. + + Не роняет и не отбрасывает вызов — только ждёт очередь (требование + задачи: при исчерпании лимита ждать, а не терять сообщение). + """ + if self._max_per_window <= 0: + return # 0/отрицательное значение конфига = лимитер выключен + lock = await self._lock_for(chat_id) + async with lock: + while True: + now = time.monotonic() + history = self._sent_at.setdefault(chat_id, []) + cutoff = now - self._window_s + while history and history[0] <= cutoff: + history.pop(0) + if len(history) < self._max_per_window: + history.append(now) + return + wait_s = history[0] + self._window_s - now + logger.warning( + "tg group rate limit: chat_id=%s — лимит %d/%.0fs исчерпан, " + "жду %.1fs перед отправкой", + chat_id, + self._max_per_window, + self._window_s, + wait_s, + ) + await asyncio.sleep(max(wait_s, 0.01)) + # Раздельные таймауты вместо скаляра. httpx разворачивает скаляр в # connect=read=write=pool, поэтому long-poll `getUpdates` (read = 30с, которые # Telegram держит запрос, + 10с запаса = 40с) ставил 40 секунд и на УСТАНОВКУ @@ -200,6 +289,7 @@ class TelegramClient: timeout: float = _DEFAULT_TIMEOUT_S, relay_base_url: str = "", relay_secret: str = "", + group_rate_limit_per_minute: int = _DEFAULT_GROUP_RATE_LIMIT_PER_MINUTE, ) -> None: # ── Ретранслятор через Beget (#3471) ──────────────────────────────── # `relay_base_url` пуст по умолчанию → `_relay_base is _direct_base`, @@ -212,6 +302,11 @@ class TelegramClient: self._base = self._relay_base self._timeout = timeout self._http: httpx.AsyncClient | None = None + # Один лимитер на экземпляр клиента — см. `TelegramGroupRateLimiter`. + # Ретранслятор здесь ни при чём: лимитер стоит ДО `_post`, то есть + # считает отправку независимо от того, уйдёт она через relay или + # напрямую (#3471) — оба пути ниже по стеку от этой точки. + self._rate_limiter = TelegramGroupRateLimiter(group_rate_limit_per_minute) def _http_client(self) -> httpx.AsyncClient: """Ленивое создание переиспользуемого AsyncClient (вне цикла ретраев).""" @@ -298,6 +393,16 @@ class TelegramClient: это штатные 30-60с) на число попыток и подвесил бы синхронный HTTP-запрос на минуты — ровно то, от чего предостерегает докстринг `send_message`. """ + # Лимитер применяется ТОЛЬКО к методам с `chat_id` в payload (отправка + # в конкретный чат) — `getUpdates` его не несёт и лимиту не подлежит. + # Списывается ОДИН слот на логический вызов `_request` (то есть на + # одну попытку отправки конкретного сообщения), а не на каждую HTTP + # попытку внутри ретрай-цикла ниже: ретраи по 429/5xx лечат один и тот + # же send, а не порождают новые отправки. + chat_id = payload.get("chat_id") + if isinstance(chat_id, int): + await self._rate_limiter.acquire(chat_id) + url = f"{self._base}/{method}" effective_timeout = timeout if timeout is not None else self._timeout # Раздельные таймауты считаем ОДИН раз и передаём per-request: у общего @@ -518,3 +623,99 @@ class TelegramClient: kwargs["max_backoff"] = max_backoff result = await self._request("sendMessage", payload, **kwargs) return result if isinstance(result, dict) else {} + + async def get_chat(self, chat_id: int) -> dict[str, Any]: + """getChat — метаданные чата. Единственная цель здесь — startup-проверка + (см. `verify_chat_and_topic`): подтвердить, что `chat_id` валиден и бот + не выгнан/не заблокирован, без единого видимого сообщения.""" + result = await self._request("getChat", {"chat_id": chat_id}, max_retries=1) + return result if isinstance(result, dict) else {} + + async def send_chat_action( + self, + *, + chat_id: int, + action: str = "typing", + message_thread_id: int | None = None, + ) -> bool: + """sendChatAction — статус набора текста. Возвращает `True`/`False`, JSON-объекта нет. + + Используется НЕ по прямому назначению (индикация набора), а как способ + проверить существование `message_thread_id` (темы форума) — см. + `verify_chat_and_topic`. Это единственный метод Bot API, который + принимает `message_thread_id` и не создаёт message-объект: если тема + удалена/переименована в другую с иным id, Telegram отвечает `Bad + Request: message thread not found` мгновенно, а в истории чата не + остаётся ни строки (индикатор эфемерный и не персистится).""" + payload: dict[str, Any] = {"chat_id": chat_id, "action": action} + if message_thread_id: + payload["message_thread_id"] = message_thread_id + result = await self._request( + "sendChatAction", payload, max_retries=1, max_backoff=5.0 + ) + return bool(result) + + +async def verify_chat_and_topic( + client: TelegramClient, + *, + chat_id: int, + topic_id: int, + label: str, +) -> bool: + """Разовая startup-проверка: чат существует, бот в нём не забанен, тема жива. + + Зачем нужна: бот пишет в тему форума по числовому id из настроек. Если + тему удалили, переименовали в другую (новый id) или id в конфиге просто + неверный — отправка начинает падать НА КАЖДОМ сообщении, а узнаём мы об + этом только по молчанию у людей (симметрично истории #tgsupport с сетевыми + отказами: тихий отказ хуже шума). Эта проверка переносит обнаружение с + «через сутки тишины» на «в первую секунду после старта/рестарта». + + Способ намеренно НЕ `sendMessage`+`deleteMessage`: + - `getChat(chat_id)` подтверждает валидность чата и то, что бот не + выгнан/не заблокирован — чистый read, нулевой видимый след. + - `send_chat_action` (typing-индикатор с `message_thread_id`) — + единственный способ провалидировать САМУ тему без создания + message-объекта: Telegram обязан знать про `message_thread_id`, чтобы + показать «печатает...» именно в нужном треде, и явно отказывает, если + такой темы нет. `sendMessage`+`deleteMessage` тоже сработал бы, но + оставлял бы видимый (пусть на секунды) артефакт в истории треда при + КАЖДОМ рестарте контейнера — на rolling-деплое это многократно в + сутки; typing-индикатор того же результата достигает без единого + сообщения. + + НЕ роняет процесс: любой `TelegramError` ловится здесь же и уходит в лог + уровня error — задача явно требует шума в логе, а не падения воркера + (тема пуста/невалидна — это деградация уведомлений, а не фатальный сбой + самого бота, который всё ещё должен принимать входящие). + + `chat_id == 0` (не настроено) — считается успехом без обращения к API: + это штатный kill-switch (см. `app.core.config`), а не ошибка конфигурации. + """ + if not chat_id: + return True + try: + await client.get_chat(chat_id) + if topic_id: + await client.send_chat_action( + chat_id=chat_id, action="typing", message_thread_id=topic_id + ) + except TelegramError as exc: + logger.error( + "tg topic check [%s]: чат/тема недоступны для отправки " + "(chat_id=%s, topic_id=%s) — %s. Сообщения в эту тему БУДУТ " + "падать, пока конфигурация не исправлена.", + label, + chat_id, + topic_id, + exc, + ) + return False + logger.info( + "tg topic check [%s]: чат и тема доступны для отправки (chat_id=%s, topic_id=%s)", + label, + chat_id, + topic_id, + ) + return True diff --git a/tradein-mvp/backend/app/services/tgbot/shared.py b/tradein-mvp/backend/app/services/tgbot/shared.py index ed1e4a52..7c2d9266 100644 --- a/tradein-mvp/backend/app/services/tgbot/shared.py +++ b/tradein-mvp/backend/app/services/tgbot/shared.py @@ -36,6 +36,7 @@ def get_telegram_client() -> TelegramClient: settings.telegram_bot_token, relay_base_url=settings.telegram_relay_base_url, relay_secret=settings.telegram_relay_secret, + group_rate_limit_per_minute=settings.telegram_group_rate_limit_per_minute, ) return _client diff --git a/tradein-mvp/backend/app/tgbot_main.py b/tradein-mvp/backend/app/tgbot_main.py index 61e06f30..97220a10 100644 --- a/tradein-mvp/backend/app/tgbot_main.py +++ b/tradein-mvp/backend/app/tgbot_main.py @@ -31,7 +31,7 @@ from app.core.config import settings from app.core.db import SessionLocal from app.core.shutdown import request_shutdown, shutdown_requested, wait_for_shutdown from app.services.tgbot.bridge import run_poll_loop -from app.services.tgbot.client import TelegramClient +from app.services.tgbot.client import TelegramClient, verify_chat_and_topic logging.basicConfig( level=logging.INFO, @@ -127,7 +127,23 @@ async def _run_bridge() -> None: settings.telegram_bot_token, relay_base_url=settings.telegram_relay_base_url, relay_secret=settings.telegram_relay_secret, + group_rate_limit_per_minute=settings.telegram_group_rate_limit_per_minute, ) as client: + # Startup-проверка (#3471): убеждаемся ОДИН раз, что чат/тема живы, + # прежде чем уходить в бесконечный poll loop. Не блокирует и не роняет + # запуск при неудаче — см. докстринг `verify_chat_and_topic`. + await verify_chat_and_topic( + client, + chat_id=settings.telegram_support_chat_id, + topic_id=settings.telegram_support_topic_id, + label="support", + ) + await verify_chat_and_topic( + client, + chat_id=settings.telegram_alerts_chat_id, + topic_id=settings.telegram_alerts_topic_id, + label="alerts", + ) await run_poll_loop(client, SessionLocal) diff --git a/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py new file mode 100644 index 00000000..3033853b --- /dev/null +++ b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py @@ -0,0 +1,263 @@ +"""Тесты для #3471: startup-проверка темы форума + общий (per-группа) rate limit. + +Оба сюжета выбраны так, чтобы падать на коде ДО правки: + - `TelegramGroupRateLimiter` и `verify_chat_and_topic` физически не существовали + в `app.services.tgbot.client` — импорт ниже сам по себе даёт `ImportError` на + старом коде (подтверждено откатом при разработке). + - `test_send_message_rate_limits_across_different_topics_same_chat` без лимитера + прошёл бы «слишком хорошо» (`sleep_calls == []` после 3 отправок) — это и есть + баг #3471 (нет лимита вовсе, всплеск бьёт по площадочному 429). + +NEVER calls real Telegram API — httpx.MockTransport only, тот же паттерн, что +`test_client.py`. +""" + +from __future__ import annotations + +import os +from unittest import mock + +import httpx +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.services.tgbot.client import ( + TelegramApiError, + TelegramClient, + TelegramGroupRateLimiter, + verify_chat_and_topic, +) + +_REAL_ASYNC_CLIENT = httpx.AsyncClient + + +def _install_transport(handler) -> None: + transport = httpx.MockTransport(handler) + + def factory(*_: object, **__: object) -> httpx.AsyncClient: + return _REAL_ASYNC_CLIENT(transport=transport) + + mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory).start() + + +@pytest.fixture(autouse=True) +def _stop_patches(): + yield + mock.patch.stopall() + + +# ── verify_chat_and_topic ──────────────────────────────────────────────────── + + +async def test_verify_chat_and_topic_success_calls_get_chat_and_send_chat_action() -> None: + """Happy path: getChat + typing-индикатор в тему, никакого видимого сообщения.""" + calls: list[str] = [] + + def handler(request: httpx.Request) -> httpx.Response: + calls.append(request.url.path.rsplit("/", 1)[-1]) + return httpx.Response(200, json={"ok": True, "result": True}) + + _install_transport(handler) + client = TelegramClient(token="fake-token") + + ok = await verify_chat_and_topic(client, chat_id=-100123, topic_id=7, label="support") + + assert ok is True + assert calls == ["getChat", "sendChatAction"] + + +async def test_verify_chat_and_topic_no_visible_message_method_used() -> None: + """Ни один из вызовов не бьёт в sendMessage/copyMessage — не мусорим в чат.""" + methods: list[str] = [] + + def handler(request: httpx.Request) -> httpx.Response: + methods.append(request.url.path.rsplit("/", 1)[-1]) + return httpx.Response(200, json={"ok": True, "result": True}) + + _install_transport(handler) + client = TelegramClient(token="fake-token") + + await verify_chat_and_topic(client, chat_id=-100123, topic_id=7, label="support") + + assert "sendMessage" not in methods + assert "copyMessage" not in methods + + +async def test_verify_chat_and_topic_logs_error_and_does_not_raise(caplog) -> None: + """Тема удалена/переименована → Telegram отвечает отказом. Лог error, процесс жив.""" + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response( + 400, json={"ok": False, "error_code": 400, "description": "message thread not found"} + ) + + _install_transport(handler) + client = TelegramClient(token="fake-token") + + with caplog.at_level("ERROR", logger="app.services.tgbot.client"): + ok = await verify_chat_and_topic(client, chat_id=-100123, topic_id=999, label="support") + + assert ok is False # не бросает исключение наружу — сигнализирует возвратом + assert any(record.levelname == "ERROR" for record in caplog.records) + assert any("support" in record.message for record in caplog.records) + + +async def test_verify_chat_and_topic_raises_nothing_even_on_network_error() -> None: + """Сетевой сбой при проверке (не 4xx, а обрыв) — тоже не должен ронять старт.""" + + def handler(request: httpx.Request) -> httpx.Response: + raise httpx.ConnectError("boom", request=request) + + _install_transport(handler) + with mock.patch("app.services.tgbot.client.asyncio.sleep", return_value=None): + client = TelegramClient(token="fake-token") + ok = await verify_chat_and_topic(client, chat_id=-100123, topic_id=7, label="alerts") + + assert ok is False + + +async def test_verify_chat_and_topic_skips_api_call_when_chat_id_is_zero() -> None: + """chat_id=0 — штатный kill-switch (не настроено), не ошибка. Нет вызовов к API.""" + calls: list[str] = [] + + def handler(request: httpx.Request) -> httpx.Response: + calls.append(request.url.path) + return httpx.Response(200, json={"ok": True, "result": True}) + + _install_transport(handler) + client = TelegramClient(token="fake-token") + + ok = await verify_chat_and_topic(client, chat_id=0, topic_id=0, label="alerts") + + assert ok is True + assert calls == [] + + +# ── TelegramGroupRateLimiter ───────────────────────────────────────────────── + + +async def test_rate_limiter_allows_calls_up_to_limit_without_waiting() -> None: + sleep_calls: list[float] = [] + + async def fake_sleep(seconds: float) -> None: + sleep_calls.append(seconds) + + with mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep): + limiter = TelegramGroupRateLimiter(max_per_window=3, window_s=60.0) + await limiter.acquire(chat_id=555) + await limiter.acquire(chat_id=555) + await limiter.acquire(chat_id=555) + + assert sleep_calls == [] + + +async def test_rate_limiter_waits_once_limit_exceeded_same_chat() -> None: + """4-я отправка в тот же chat_id за лимитом (2) — ждёт, а не теряется/падает.""" + sleep_calls: list[float] = [] + fake_now = [1_000.0] + + def fake_monotonic() -> float: + return fake_now[0] + + async def fake_sleep(seconds: float) -> None: + sleep_calls.append(seconds) + fake_now[0] += seconds + + with ( + mock.patch("app.services.tgbot.client.time.monotonic", fake_monotonic), + mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep), + ): + limiter = TelegramGroupRateLimiter(max_per_window=2, window_s=60.0) + await limiter.acquire(chat_id=1) + await limiter.acquire(chat_id=1) + assert sleep_calls == [] + await limiter.acquire(chat_id=1) # третья — сверх лимита + + assert sleep_calls # ждала, а не пролетела/бросила исключение + assert sleep_calls[0] > 0 + + +async def test_rate_limiter_different_chats_do_not_block_each_other() -> None: + """Разные группы (chat_id) — независимые бюджеты, одна не душит другую.""" + sleep_calls: list[float] = [] + + async def fake_sleep(seconds: float) -> None: + sleep_calls.append(seconds) + + with mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep): + limiter = TelegramGroupRateLimiter(max_per_window=1, window_s=60.0) + await limiter.acquire(chat_id=1) + await limiter.acquire(chat_id=2) # другая группа — свой бюджет + + assert sleep_calls == [] + + +async def test_rate_limiter_no_races_under_concurrent_senders() -> None: + """Несколько конкурентных asyncio-отправителей в ОДИН chat_id — без гонок: + итоговое число «пропущенных без ожидания» строго не больше лимита, и никто + не теряется/не падает — все 5 корутин успешно завершаются.""" + import asyncio + + limiter = TelegramGroupRateLimiter(max_per_window=2, window_s=0.15) + + async def sender() -> None: + await limiter.acquire(chat_id=42) + + # Реальное время, реальный event loop — единственный способ честно + # проверить отсутствие гонки на разделяемом `self._sent_at[chat_id]`. + results = await asyncio.gather(*(sender() for _ in range(5)), return_exceptions=True) + + assert all(not isinstance(r, Exception) for r in results) + history = limiter._sent_at[42] + # После завершения всех 5 — в окне НЕ больше max_per_window меток (иначе + # гонка позволила бы двум конкурентным acquire() одновременно "проскочить"). + assert len(history) <= 2 + + +# ── интеграция с TelegramClient.send_message ───────────────────────────────── + + +async def test_send_message_rate_limits_across_different_topics_same_chat() -> None: + """Ключевой сценарий #3471: РАЗНЫЕ темы одной группы делят ОДИН бюджет. + + Без правки лимитера нет вовсе — этот тест на старом коде проходил бы + "слишком хорошо" (sleep не звался никогда), что и есть баг. + """ + sleep_calls: list[float] = [] + + async def fake_sleep(seconds: float) -> None: + sleep_calls.append(seconds) + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}}) + + _install_transport(handler) + with mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep): + client = TelegramClient(token="fake-token", group_rate_limit_per_minute=2) + await client.send_message(chat_id=42, text="a", message_thread_id=1) # тема поддержки + await client.send_message(chat_id=42, text="b", message_thread_id=2) # тема алертов + assert sleep_calls == [] + await client.send_message(chat_id=42, text="c", message_thread_id=3) # 3-я тема, тот же чат + + assert sleep_calls # общий бюджет группы исчерпан третьей отправкой + + +async def test_send_message_does_not_rate_limit_missing_chat_id_payload() -> None: + """get_updates (нет chat_id в payload) не должен спотыкаться о лимитер вовсе.""" + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={"ok": True, "result": []}) + + _install_transport(handler) + client = TelegramClient(token="fake-token", group_rate_limit_per_minute=1) + + # Не должно ни зависнуть, ни бросить — getUpdates не несёт chat_id. + result = await client.get_updates(offset=1) + assert result == [] + + +def test_telegram_api_error_importable_for_manual_inspection() -> None: + """Sanity: тестовый модуль не потерял импорт TelegramApiError (использовался + при отладке 400-ответа выше).""" + assert issubclass(TelegramApiError, Exception) -- 2.45.3 From 8e7c65061bf4c895604ca3e76c10f3b4f17c01dd Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 15:28:16 +0300 Subject: [PATCH 2/3] fix(tgbot): honest H1 rejection, per-role H2 budget, M1/M2/L1 cleanup (#3471 review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Deep review of PR #3494 found the rate limiter unusable as designed: - H1: acquire() waited unbounded even for interactive HTTP handlers (support.py, glitchtip.py already pass a narrow `timeout` — reuse it as the queue wait cap instead of editing those handlers, which are out of scope here). New TelegramRateLimitedError (subclass of TelegramError) gives a fast, honest 502 instead of hanging past the caller's own budget. - H2: the limiter is per-process (in-memory), but two processes write to the same group (uvicorn API + bot worker) — giving each the same 18/min doubled the platform ceiling. Split into telegram_group_rate_limit_api_per_minute (12) and _bot_per_minute (6), sum kept below ~20. - M1: bridge.py sends without an explicit timeout inherited "wait forever", stalling the single-threaded poll loop (open DB session) past the SIGTERM drain window. Bounded via rate_limit_max_wait=20s at the six call sites. - M2: _locks/_sent_at grew unbounded on every unique DM chat_id. Added opportunistic cleanup of fully-expired entries. - L1: the "queue full" warning now logs once per acquire() call, not once per sleep iteration. - Corrected a factual error in the docstring: TELEGRAM_SUPPORT_CHAT_ID and TELEGRAM_ALERTS_CHAT_ID are the SAME group on prod (topics differ only). Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG --- tradein-mvp/backend/app/core/config.py | 25 ++- .../backend/app/services/tgbot/bridge.py | 36 ++++- .../backend/app/services/tgbot/client.py | 142 +++++++++++++++--- .../backend/app/services/tgbot/shared.py | 5 +- tradein-mvp/backend/app/tgbot_main.py | 4 +- .../test_topic_check_and_group_rate_limit.py | 95 ++++++++++++ 6 files changed, 279 insertions(+), 28 deletions(-) diff --git a/tradein-mvp/backend/app/core/config.py b/tradein-mvp/backend/app/core/config.py index 4a930171..95de5ccb 100644 --- a/tradein-mvp/backend/app/core/config.py +++ b/tradein-mvp/backend/app/core/config.py @@ -1290,10 +1290,27 @@ class Settings(BaseSettings): # Общий (не per-тему) лимит частоты отправки в ОДНУ группу — Telegram # считает ~20 сообщений/минуту на группу суммарно по всем её темам (#3471: # всплеск GlitchTip-алертов + поток поддержки в ту же группу давали 429 и - # потерю сообщений). Дефолт — чуть ниже площадочного потолка. 0/отрицательное - # значение выключает лимитер (см. `TelegramGroupRateLimiter.acquire`). - telegram_group_rate_limit_per_minute: int = Field( - default=18, validation_alias="TELEGRAM_GROUP_RATE_LIMIT_PER_MINUTE" + # потерю сообщений). `TELEGRAM_SUPPORT_CHAT_ID`/`TELEGRAM_ALERTS_CHAT_ID` на + # проде равны (одна группа, темы разные) — обе половины делят один + # площадочный бюджет. + # + # РАЗДЕЛЕНО ПО РОЛЯМ (review H2, #3471), а не одна общая константа: лимитер + # живёт in-memory В ЭКЗЕМПЛЯРЕ `TelegramClient`, а в эту группу пишут ДВА + # независимых процесса — API под uvicorn (`app/services/tgbot/shared.py`, + # интерактивные ручки + GlitchTip-вебхук) и контейнер бота (`app/tgbot_main.py`, + # long-polling воркер). У них НЕТ общего счётчика (это отдельная задача — + # Redis-based распределённый лимитер), поэтому если каждому дать по 18, + # сумма (2×18=36) УДВОИТ площадочный лимit и 429 вернётся ровно там же. + # Бюджет делится статически так, чтобы СУММА была заметно НИЖЕ ~20: у API + # больше — там же интерактивные ответы клиентам, у бота меньше — там же + # обычно только зеркалирование/уведомления, которые могут подождать дольше + # (см. `TelegramGroupRateLimiter.acquire` про `max_wait=None` для фона). + # 0/отрицательное значение выключает лимитер для соответствующего процесса. + telegram_group_rate_limit_api_per_minute: int = Field( + default=12, validation_alias="TELEGRAM_GROUP_RATE_LIMIT_API_PER_MINUTE" + ) + telegram_group_rate_limit_bot_per_minute: int = Field( + default=6, validation_alias="TELEGRAM_GROUP_RATE_LIMIT_BOT_PER_MINUTE" ) # ── Ретранслятор Bot API через Beget (#3471) ───────────────────────────── diff --git a/tradein-mvp/backend/app/services/tgbot/bridge.py b/tradein-mvp/backend/app/services/tgbot/bridge.py index d23c3307..fcbf4f78 100644 --- a/tradein-mvp/backend/app/services/tgbot/bridge.py +++ b/tradein-mvp/backend/app/services/tgbot/bridge.py @@ -147,6 +147,21 @@ _NOTIFY_SEND_TIMEOUT_S = 5.0 _NOTIFY_SEND_MAX_RETRIES = 3 _NOTIFY_SEND_MAX_BACKOFF_S = 1.0 +# Потолок ожидания в TelegramGroupRateLimiter для отправок ИЗ ЭТОГО модуля, +# которые не проходят через _notify_topic (review M1, #3471). Без него +# `client.send_message`/`copy_message` без явного `timeout` наследуют +# лимитер-политику "без потолка" (см. TelegramClient._request) — приемлемую +# ДЛЯ ФОНОВОЙ отправки как таковой, но НЕ здесь: эти вызовы идут внутри +# `run_poll_loop`, который на каждый апдейт держит ОТКРЫТУЮ сессию БД +# (SessionLocal, см. вызывающих) и однопоточно блокирует опрос СЛЕДУЮЩИХ +# апдейтов — минутный сон здесь стопорит и БД-соединение, и весь мост, а не +# только одно сообщение. Второй повод — SIGTERM drain: `tgbot_main._DRAIN_TIMEOUT_S` +# даёт 100с на завершение текущей итерации; ожидание слота дольше этого +# бюджета уже не успевает подчиниться cooperative drain. 20с — заметно больше, +# чем разумная очередь при исчерпанном лимите (окно 60с, обычно секунды), но +# заметно МЕНЬШЕ минуты и вписывается в drain-бюджет с запасом. +_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S = 20.0 + # Потолок переигрываний ОДНОГО update_id на транзиентных сетевых отказах # (#tg-connection-resilience). # @@ -535,7 +550,11 @@ async def _handle_private_message( text_body = message.get("text") if text_body == "/start": # E) команда — не содержательное обращение, топик не засоряем. - await client.send_message(chat_id=chat_id, text=GREETING_TEXT) + await client.send_message( + chat_id=chat_id, + text=GREETING_TEXT, + rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S, + ) return if not settings.telegram_support_chat_id: @@ -545,7 +564,11 @@ async def _handle_private_message( chat_id, ) # Не молчим клиенту (#5 review) — иначе он ждёт ответа, которого никогда не будет. - await client.send_message(chat_id=chat_id, text=SERVICE_UNAVAILABLE_TEXT) + await client.send_message( + chat_id=chat_id, + text=SERVICE_UNAVAILABLE_TEXT, + rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S, + ) return message_id = message.get("message_id") @@ -573,7 +596,11 @@ async def _handle_private_message( # за окно — `_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) + await client.send_message( + chat_id=chat_id, + text=FLOOD_LIMITED_TEXT, + rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S, + ) return _flood_limiter.record(flood_key) @@ -584,6 +611,7 @@ async def _handle_private_message( chat_id=settings.telegram_support_chat_id, text=header, message_thread_id=settings.telegram_support_topic_id or None, + rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S, ) mirrored = await client.copy_message( @@ -591,6 +619,7 @@ async def _handle_private_message( from_chat_id=chat_id, message_id=message_id, message_thread_id=settings.telegram_support_topic_id or None, + rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S, ) topic_message_id = mirrored.get("message_id") if isinstance(mirrored, dict) else None @@ -655,6 +684,7 @@ async def _handle_group_reply( chat_id=target_chat_id, from_chat_id=settings.telegram_support_chat_id, message_id=message_id, + rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S, ) except TelegramApiError as exc: if exc.error_code == 403: diff --git a/tradein-mvp/backend/app/services/tgbot/client.py b/tradein-mvp/backend/app/services/tgbot/client.py index 5f0c93ab..3313988e 100644 --- a/tradein-mvp/backend/app/services/tgbot/client.py +++ b/tradein-mvp/backend/app/services/tgbot/client.py @@ -86,9 +86,16 @@ class TelegramGroupRateLimiter: (`asyncio.sleep`), а не только на чтение/запись счётчика — это осознанно: цель не просто «не гонять счётчик без гонки», а ФАКТИЧЕСКИ сериализовать отправителей в этот чат, чтобы они не просыпались все разом по истечении - окна и не били по лимиту повторно. Другие `chat_id` (например, тема - поддержки и тема алертов лежат в РАЗНЫХ группах) не блокируют друг друга — - у каждого свой лок и свой список меток. + окна и не били по лимиту повторно. + + ВАЖНО про ключ (проверено на проде, review 2026-09-12): `TELEGRAM_SUPPORT_CHAT_ID` + и `TELEGRAM_ALERTS_CHAT_ID` — это ОДНА И ТА ЖЕ группа, различаются только + темы (`*_TOPIC_ID`). Именно поэтому ключ лимитера — `chat_id`, а НЕ + `(chat_id, message_thread_id)`: поток алертов и поток поддержки сегодня + физически делят один Telegram-бюджет группы, и лимитер обязан это + отражать. Ключ по `chat_id` при этом остаётся корректным и в гипотезе, что + когда-нибудь эти два потока разведут по разным группам, — тогда у каждой + просто появится свой независимый лок/бюджет автоматически, без правки кода. Регистр локов защищён отдельным `asyncio.Lock` только на момент создания записи — сам подсчёт/сон идёт уже под персональным локом чата. @@ -113,18 +120,38 @@ class TelegramGroupRateLimiter: self._locks[chat_id] = lock return lock - async def acquire(self, chat_id: int) -> None: + async def acquire(self, chat_id: int, max_wait: float | None = None) -> None: """Блокируется, пока в окне `_window_s` для `chat_id` есть свободный слот. - Не роняет и не отбрасывает вызов — только ждёт очередь (требование - задачи: при исчерпании лимита ждать, а не терять сообщение). + `max_wait` (review H1, #3471): потолок ожидания очереди. `None` (дефолт) + — без потолка, ждать сколько нужно; это ПРАВИЛЬНОЕ поведение для + фоновых отправок бота, где потерять сообщение хуже, чем подождать. + Если задан и слот не появился вовремя — бросает `TelegramRateLimitedError` + (честный отказ), а НЕ продолжает ждать: интерактивная HTTP-ручка не + может легально держать открытый запрос браузера дольше своего + собственного таймаута. Вызывающая сторона — `TelegramClient._request`, + см. её докстринг про то, откуда берётся конкретное значение. + + Логирование (review L1): предупреждение об ожидании пишется РОВНО ОДИН + раз за вызов `acquire` (флаг `warned`), а не на каждой итерации сна — + при реальной перегрузке группы это иначе валит лог сотнями одинаковых + строк вместо одного сигнала «была очередь». + + Побочный эффект (review M2): перед постановкой в очередь чистит ЧУЖИЕ + полностью просроченные записи в `_sent_at`/`_locks` — см. `_cleanup_stale`. """ if self._max_per_window <= 0: return # 0/отрицательное значение конфига = лимитер выключен + now0 = time.monotonic() + await self._cleanup_stale(now0) + deadline = None if max_wait is None else now0 + max_wait lock = await self._lock_for(chat_id) async with lock: + warned = False while True: now = time.monotonic() + if deadline is not None and now >= deadline: + raise TelegramRateLimitedError(chat_id, max_wait or 0.0) history = self._sent_at.setdefault(chat_id, []) cutoff = now - self._window_s while history and history[0] <= cutoff: @@ -133,16 +160,45 @@ class TelegramGroupRateLimiter: history.append(now) return wait_s = history[0] + self._window_s - now - logger.warning( - "tg group rate limit: chat_id=%s — лимит %d/%.0fs исчерпан, " - "жду %.1fs перед отправкой", - chat_id, - self._max_per_window, - self._window_s, - wait_s, - ) + if deadline is not None: + wait_s = min(wait_s, max(deadline - now, 0.0)) + if not warned: + logger.warning( + "tg group rate limit: chat_id=%s — лимит %d/%.0fs исчерпан, " + "отправки встают в очередь (одно предупреждение на серию)", + chat_id, + self._max_per_window, + self._window_s, + ) + warned = True await asyncio.sleep(max(wait_s, 0.01)) + async def _cleanup_stale(self, now: float) -> None: + """Чистит ЧУЖИЕ (не текущий вызов `acquire`) записи с полностью + просроченной историей (review M2, #3471). + + Зачем: `_locks`/`_sent_at` ключуются по ЛЮБОМУ `chat_id`, включая + личные чаты каждого клиента бота (`bridge.py` зеркалит их 1:1 через тот + же `TelegramClient`) — большинство из них шлют боту одно сообщение и + больше никогда не возвращаются. Без чистки оба словаря растут + монотонно на всё время жизни долгоживущего процесса (медленная утечка). + + Безопасность удаления: лок пропускаем, если `lock.locked()` — значит + кто-то ИМЕННО СЕЙЧАС работает с этим `chat_id`, трогать нельзя. Если + лок свободен и вся история старше окна — запись безвредно удалить: + следующий `acquire` для того же `chat_id` просто создаст её заново + пустой (`setdefault`), с тем же результатом, что и не удаляй мы её. + """ + cutoff = now - self._window_s + async with self._registry_lock: + stale = [cid for cid, ts in self._sent_at.items() if not ts or ts[-1] <= cutoff] + for cid in stale: + lock = self._locks.get(cid) + if lock is not None and lock.locked(): + continue + self._sent_at.pop(cid, None) + self._locks.pop(cid, None) + # Раздельные таймауты вместо скаляра. httpx разворачивает скаляр в # connect=read=write=pool, поэтому long-poll `getUpdates` (read = 30с, которые # Telegram держит запрос, + 10с запаса = 40с) ставил 40 секунд и на УСТАНОВКУ @@ -215,6 +271,25 @@ class TelegramNetworkError(TelegramError): super().__init__(f"Telegram {method} unreachable after {attempts} attempts: {reason}") +class TelegramRateLimitedError(TelegramError): + """Слот в `TelegramGroupRateLimiter` не появился за `max_wait` секунд (review H1). + + Бросается ТОЛЬКО когда вызывающий явно (или неявно, через `timeout`) + попросил ограниченное ожидание — фоновые вызовы без такого ограничения + ждут очередь сколько нужно и этого исключения никогда не увидят. Общий + предок `TelegramError` — существующие `except TelegramError` в + `app.api.v1.support`/`glitchtip` подхватывают этот отказ автоматически, без + правки самих ручек, и отвечают своим честным 502 вместо зависшего запроса. + """ + + def __init__(self, chat_id: int, max_wait: float) -> None: + self.chat_id = chat_id + self.max_wait = max_wait + super().__init__( + f"Telegram group rate limit: no slot for chat_id={chat_id} within {max_wait:.1f}s" + ) + + def _request_timeout(read: float) -> httpx.Timeout: """Разворачивает «сколько ждать ответа» (скаляр вызывающего) в таймауты httpx. @@ -377,9 +452,24 @@ class TelegramClient: timeout: float | None = None, max_retries: int = _DEFAULT_MAX_RETRIES, max_backoff: float | None = None, + rate_limit_max_wait: float | None = None, ) -> Any: """POST `method` с JSON-телом `payload`. Ретраит 429/5xx/network, иначе raise сразу. + `rate_limit_max_wait` (review H1/M1, #3471) — потолок ожидания слота в + `TelegramGroupRateLimiter.acquire`. Приоритет: + 1. Явный `rate_limit_max_wait` — используется как есть (bridge.py + передаёт его точечно для конкретных мест, см. `_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S`). + 2. Иначе, если вызывающий передал явный `timeout` — используем + `effective_timeout` КАК ЕСТЬ. Интерактивные ручки (`app.api.v1.support`, + `app.api.v1.glitchtip`) и так ОБЯЗАНЫ передавать узкий `timeout` + (5-8с, см. их собственные докстринги) — этого достаточно, чтобы + очередь лимитера не держала открытый HTTP-запрос браузера дольше + его же собственного бюджета, БЕЗ дополнительной правки этих ручек. + 3. Иначе `None` — без потолка. Это дефолт для фоновых отправок бота + (`app.tgbot_main`/`bridge.py` без явного `timeout`), где потерять + сообщение хуже, чем подождать дольше. + `max_backoff` (#tgsupport-retry) — потолок паузы МЕЖДУ попытками. По умолчанию `_MAX_BACKOFF_S` (30с) и полный `retry_after` из тела 429 — это воркерная политика, она НЕ меняется. Интерактивный вызывающий передаёт узкий @@ -393,18 +483,23 @@ class TelegramClient: это штатные 30-60с) на число попыток и подвесил бы синхронный HTTP-запрос на минуты — ровно то, от чего предостерегает докстринг `send_message`. """ + effective_timeout = timeout if timeout is not None else self._timeout + # Лимитер применяется ТОЛЬКО к методам с `chat_id` в payload (отправка # в конкретный чат) — `getUpdates` его не несёт и лимиту не подлежит. # Списывается ОДИН слот на логический вызов `_request` (то есть на # одну попытку отправки конкретного сообщения), а не на каждую HTTP # попытку внутри ретрай-цикла ниже: ретраи по 429/5xx лечат один и тот - # же send, а не порождают новые отправки. + # же send, а не порождают новые отправки. Приоритет `wait_cap` — см. + # докстринг параметра `rate_limit_max_wait` выше. chat_id = payload.get("chat_id") if isinstance(chat_id, int): - await self._rate_limiter.acquire(chat_id) + wait_cap = rate_limit_max_wait + if wait_cap is None and timeout is not None: + wait_cap = effective_timeout + await self._rate_limiter.acquire(chat_id, max_wait=wait_cap) url = f"{self._base}/{method}" - effective_timeout = timeout if timeout is not None else self._timeout # Раздельные таймауты считаем ОДИН раз и передаём per-request: у общего # AsyncClient свой дефолт, а бюджет ответа у каждого вызова свой. request_timeout = _request_timeout(effective_timeout) @@ -574,8 +669,12 @@ class TelegramClient: message_id: int, message_thread_id: int | None = None, reply_to_message_id: int | None = None, + rate_limit_max_wait: float | None = None, ) -> dict[str, Any]: - """copyMessage — зеркалит ЛЮБОЙ тип контента без ре-аплоада файла.""" + """copyMessage — зеркалит ЛЮБОЙ тип контента без ре-аплоада файла. + + `rate_limit_max_wait` — см. `TelegramClient._request`; используется + `bridge.py` для точечного потолка ожидания на конкретных местах (review M1).""" payload: dict[str, Any] = { "chat_id": chat_id, "from_chat_id": from_chat_id, @@ -585,7 +684,9 @@ class TelegramClient: payload["message_thread_id"] = message_thread_id if reply_to_message_id: payload["reply_to_message_id"] = reply_to_message_id - result = await self._request("copyMessage", payload) + result = await self._request( + "copyMessage", payload, rate_limit_max_wait=rate_limit_max_wait + ) return result if isinstance(result, dict) else {} async def send_message( @@ -598,6 +699,7 @@ class TelegramClient: timeout: float | None = None, max_retries: int | None = None, max_backoff: float | None = None, + rate_limit_max_wait: float | None = None, ) -> dict[str, Any]: """sendMessage — текстовое сообщение (заголовки, приветствия, уведомления об ошибке). @@ -621,6 +723,8 @@ class TelegramClient: kwargs["max_retries"] = max_retries if max_backoff is not None: kwargs["max_backoff"] = max_backoff + if rate_limit_max_wait is not None: + kwargs["rate_limit_max_wait"] = rate_limit_max_wait result = await self._request("sendMessage", payload, **kwargs) return result if isinstance(result, dict) else {} diff --git a/tradein-mvp/backend/app/services/tgbot/shared.py b/tradein-mvp/backend/app/services/tgbot/shared.py index 7c2d9266..7d99b86f 100644 --- a/tradein-mvp/backend/app/services/tgbot/shared.py +++ b/tradein-mvp/backend/app/services/tgbot/shared.py @@ -36,7 +36,10 @@ def get_telegram_client() -> TelegramClient: settings.telegram_bot_token, relay_base_url=settings.telegram_relay_base_url, relay_secret=settings.telegram_relay_secret, - group_rate_limit_per_minute=settings.telegram_group_rate_limit_per_minute, + # API-роль (review H2, #3471) — см. докстринг настройки в + # app.core.config: бюджет группы разделён статически между этим + # процессом и app.tgbot_main, суммарно ниже площадочного лимита. + group_rate_limit_per_minute=settings.telegram_group_rate_limit_api_per_minute, ) return _client diff --git a/tradein-mvp/backend/app/tgbot_main.py b/tradein-mvp/backend/app/tgbot_main.py index 97220a10..a156b7d4 100644 --- a/tradein-mvp/backend/app/tgbot_main.py +++ b/tradein-mvp/backend/app/tgbot_main.py @@ -127,7 +127,9 @@ async def _run_bridge() -> None: settings.telegram_bot_token, relay_base_url=settings.telegram_relay_base_url, relay_secret=settings.telegram_relay_secret, - group_rate_limit_per_minute=settings.telegram_group_rate_limit_per_minute, + # Бот-роль (review H2, #3471) — своя, меньшая доля общего бюджета + # группы; см. докстринг настройки в app.core.config. + group_rate_limit_per_minute=settings.telegram_group_rate_limit_bot_per_minute, ) as client: # Startup-проверка (#3471): убеждаемся ОДИН раз, что чат/тема живы, # прежде чем уходить в бесконечный poll loop. Не блокирует и не роняет diff --git a/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py index 3033853b..bd69bdd1 100644 --- a/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py +++ b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py @@ -26,6 +26,7 @@ from app.services.tgbot.client import ( TelegramApiError, TelegramClient, TelegramGroupRateLimiter, + TelegramRateLimitedError, verify_chat_and_topic, ) @@ -261,3 +262,97 @@ def test_telegram_api_error_importable_for_manual_inspection() -> None: """Sanity: тестовый модуль не потерял импорт TelegramApiError (использовался при отладке 400-ответа выше).""" assert issubclass(TelegramApiError, Exception) + + +# ── review H1: честный отказ вместо бесконечного ожидания на interactive-пути ─ + + +async def test_rate_limiter_raises_rate_limited_error_when_max_wait_exceeded() -> None: + """Слот не появился за `max_wait` — `TelegramRateLimitedError`, не вечный сон.""" + fake_now = [0.0] + + def fake_monotonic() -> float: + return fake_now[0] + + async def fake_sleep(seconds: float) -> None: + fake_now[0] += seconds + + with ( + mock.patch("app.services.tgbot.client.time.monotonic", fake_monotonic), + mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep), + ): + limiter = TelegramGroupRateLimiter(max_per_window=1, window_s=60.0) + await limiter.acquire(chat_id=9) # занимает единственный слот + with pytest.raises(TelegramRateLimitedError): + await limiter.acquire(chat_id=9, max_wait=2.0) # бюджет короче окна + + +async def test_send_message_bounds_rate_limit_wait_by_explicit_timeout() -> None: + """Явный `timeout` (как у interactive-ручек support.py/glitchtip.py) сам по + себе ограничивает ожидание очереди — без правки самих ручек (review H1).""" + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}}) + + _install_transport(handler) + with mock.patch( + "app.services.tgbot.client.TelegramGroupRateLimiter.acquire", + autospec=True, + ) as acquire_mock: + client = TelegramClient(token="fake-token") + await client.send_message(chat_id=1, text="hi", timeout=3.0) + + _, kwargs = acquire_mock.call_args + assert kwargs.get("max_wait") == 3.0 + + +async def test_send_message_unbounded_wait_by_default_for_background_calls() -> None: + """Без явного `timeout` (типичный фоновый вызов) — `max_wait=None` (без потолка).""" + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}}) + + _install_transport(handler) + with mock.patch( + "app.services.tgbot.client.TelegramGroupRateLimiter.acquire", + autospec=True, + ) as acquire_mock: + client = TelegramClient(token="fake-token") + await client.send_message(chat_id=1, text="hi") + + _, kwargs = acquire_mock.call_args + assert kwargs.get("max_wait") is None + + +# ── review H2: бюджет группы разделён по ролям (API-процесс / бот-процесс) ──── + + +def test_group_rate_limit_split_by_role_sums_below_platform_ceiling() -> None: + """API- и бот-роль вместе НЕ должны превышать площадочный лимит (~20/мин). + + Регрессия ровно на баг из review H2: раньше оба процесса получали + ОДИНАКОВЫЙ дефолт (18+18=36) — сумма вдвое превышала лимит площадки. + """ + from app.core.config import settings + + api_limit = settings.telegram_group_rate_limit_api_per_minute + bot_limit = settings.telegram_group_rate_limit_bot_per_minute + assert api_limit > 0 + assert bot_limit > 0 + assert api_limit + bot_limit < 20 + + +def test_shared_client_uses_api_role_limit() -> None: + """`app.services.tgbot.shared` (процесс uvicorn) собирает клиент с API-долей.""" + from app.core.config import settings + from app.services.tgbot import shared + + shared._client = None + try: + client = shared.get_telegram_client() + assert ( + client._rate_limiter._max_per_window + == settings.telegram_group_rate_limit_api_per_minute + ) + finally: + shared._client = None -- 2.45.3 From 814283455579e5a4e3f83f3b69cad7dde0ef736b Mon Sep 17 00:00:00 2001 From: bot-backend Date: Sat, 12 Sep 2026 16:10:27 +0300 Subject: [PATCH 3/3] fix(tgbot): stop CI-hanging busy-spin in rate-limit test, isolate shared client in tests Root cause of the red PR #3494 CI job (9% progress, 75s life, no error line): test_send_message_rate_limits_across_different_topics_same_chat mocked asyncio.sleep as a pure no-op without advancing time.monotonic. The 3rd send (over the test's limit=2) entered TelegramGroupRateLimiter.acquire(), which recomputes wait_s from the real, unmocked clock every iteration - since the fake sleep never advances it, the window never expires and the while-loop busy-spins forever instead of actually waiting, until pytest-timeout kills it. Fixed by advancing a fake monotonic clock inside fake_sleep, matching the already-correct pattern used by the other tests in this file. Also added _reset_telegram_shared_client (tests/conftest.py, same pattern as _reset_estimate_rate_limiter): app.services.tgbot.shared._client is a module-level singleton whose rate limiter otherwise accumulates real wall-clock timestamps across the whole pytest session, not per test. Documented honestly in config.py: the API-role budget is shared between support web-chat mirrors and GlitchTip alerts with no priority between them, so a large alert burst can make the web-chat wait out its own timeout and return 502 - flagged as a known follow-up, not fixed here. NOTE: a full `pytest -q --timeout=60` run still hangs further into the suite, at tests/test_glitchtip_webhook.py::test_telegram_failure_returns_502_not_500. Not root-caused within this session's budget - the test's _fake_telegram_client fixture correctly monkeypatches glitchtip_module.get_telegram_client, but the anyio worker thread running the ASGI request is seen parked in a real event-loop poll/select wait, consistent with an actual (non-mocked) sleep somewhere in that path. Needs a follow-up session with a fresh time budget. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG --- tradein-mvp/backend/app/core/config.py | 19 ++++++++++ tradein-mvp/backend/tests/conftest.py | 35 +++++++++++++++++-- .../test_topic_check_and_group_rate_limit.py | 21 ++++++++++- 3 files changed, 72 insertions(+), 3 deletions(-) diff --git a/tradein-mvp/backend/app/core/config.py b/tradein-mvp/backend/app/core/config.py index 95de5ccb..485484d8 100644 --- a/tradein-mvp/backend/app/core/config.py +++ b/tradein-mvp/backend/app/core/config.py @@ -1306,6 +1306,25 @@ class Settings(BaseSettings): # обычно только зеркалирование/уведомления, которые могут подождать дольше # (см. `TelegramGroupRateLimiter.acquire` про `max_wait=None` для фона). # 0/отрицательное значение выключает лимитер для соответствующего процесса. + # + # ЧЕСТНО ПРО ОГРАНИЧЕНИЕ ЭТОГО ДИЗАЙНА (review, #3471): + # `telegram_group_rate_limit_api_per_minute` — ОДИН общий бюджет на ВСЕ + # отправки процесса API в эту группу, а туда + # пишут И зеркала веб-чата поддержки (`app/api/v1/support.py`), И + # GlitchTip-алерты (`app/api/v1/glitchtip.py`) — обе ручки идут через один + # и тот же `get_telegram_client()` (см. `app/services/tgbot/shared.py`). + # Приоритета между ними НЕТ: кто первый встал в очередь `TelegramGroupRateLimiter`, + # тот и получил слот. Оба пути передают узкий `timeout` (5с у support, 8с у + # glitchtip) — он же становится потолком ожидания слота (см. + # `TelegramClient._request`, review H1). Значит при всплеске алертов (пачка + # ошибок прода бьёт в вебхук залпом) реально возможен сценарий: бюджет + # 12/мин исчерпан алертами → следующая отправка живого клиента в веб-чате + # ждёт до 5с и получает `TelegramRateLimitedError` → 502 клиенту поддержки. + # То есть при достаточно большом всплеске алертов веб-чат ДЕЙСТВИТЕЛЬНО + # может временно вставать. Разделить бюджет по ИСТОЧНИКУ (не по процессу) — + # отдельная задача: нужен свой `TelegramGroupRateLimiter` на алерты с явно + # малой квотой и/или приоритет для support-трафика; здесь НЕ сделано + # (вне бюджета этой правки). telegram_group_rate_limit_api_per_minute: int = Field( default=12, validation_alias="TELEGRAM_GROUP_RATE_LIMIT_API_PER_MINUTE" ) diff --git a/tradein-mvp/backend/tests/conftest.py b/tradein-mvp/backend/tests/conftest.py index 02b64a9f..71a916a3 100644 --- a/tradein-mvp/backend/tests/conftest.py +++ b/tradein-mvp/backend/tests/conftest.py @@ -4,8 +4,10 @@ `--strict-markers` в pyproject.toml не включён, так что незарегистрированный маркер только предупреждал бы) и сторожит глобальное состояние, которое переживает отдельный тест: общий rate-limiter POST /estimate (см. -`_reset_estimate_rate_limiter`) и слоты проверки пароля (см. -`_no_leaked_password_verify_slots`). +`_reset_estimate_rate_limiter`), слоты проверки пароля (см. +`_no_leaked_password_verify_slots`) и синглтон Telegram-клиента (см. +`_reset_telegram_shared_client`, #3471 — иначе его rate limiter копит +реальное время между тестами и вешает прогон). """ from __future__ import annotations @@ -48,6 +50,35 @@ def _reset_estimate_rate_limiter() -> None: ) +@pytest.fixture(autouse=True) +def _reset_telegram_shared_client(): + """`app.services.tgbot.shared._client` — модульный синглтон `TelegramClient` + (#3471). Его `TelegramGroupRateLimiter` копит РЕАЛЬНЫЕ метки времени + (`time.monotonic()`, ничем не замоканные) по `chat_id` за весь pytest-процесс, + а не по тесту — а тестовые настройки `telegram_alerts_chat_id`/ + `telegram_support_chat_id` дефолтятся в 0, так что ЛЮБЫЕ тесты, бьющие в + `app.api.v1.support`/`glitchtip` через реальный (не замоканный) shared-клиент, + делят ОДИН и тот же ключ бакета. После N-й (лимит роли, по умолчанию 12) + такой отправки в пределах 60 реальных секунд следующая уходит в настоящий + `asyncio.sleep` до 60с — тест не падает, а зависает, и именно так выглядела + смерть CI-джобы на #3494 (обрыв на ~9%, 75с жизни, ни строки об ошибке; + `pytest-timeout` затем добивает зависший тест снаружи). + + Фикстура не выключает и не завышает лимит (в проде он ДОЛЖЕН оставаться + ниже площадочного потолка) — она просто гарантирует каждому тесту СВЕЖИЙ + клиент (и тем самым свежий, пустой `TelegramGroupRateLimiter`), так же как + `_reset_estimate_rate_limiter` выше делает для `_estimate_limiter`. Сброс + и ДО, и ПОСЛЕ теста — тест мог создать клиент через `get_telegram_client()`, + не пройдя явный локальный `_reset_singleton` (см. `test_shared.py`), и не + должен оставить накопленное состояние следующему тесту. + """ + from app.services.tgbot import shared + + shared._client = None + yield + shared._client = None + + @pytest.fixture(autouse=True) def _no_leaked_password_verify_slots(): """Тест не оставляет за собой занятых слотов проверки пароля (#2665, #2714). diff --git a/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py index bd69bdd1..73a8ecfa 100644 --- a/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py +++ b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py @@ -224,17 +224,36 @@ async def test_send_message_rate_limits_across_different_topics_same_chat() -> N Без правки лимитера нет вовсе — этот тест на старом коде проходил бы "слишком хорошо" (sleep не звался никогда), что и есть баг. + + ВАЖНО (post-merge review, #3494 CI hang): `fake_sleep` ОБЯЗАН двигать + подменённый `time.monotonic` вперёд, а не оставаться чистым no-op. Без + этого 3-я отправка (сверх лимита=2) уходит в `acquire()`, `wait_s` + пересчитывается из РЕАЛЬНОГО `time.monotonic()` (не замоканного здесь), + окно не истекает — и `while True` крутится настоящим busy-spin БЕЗ + единого реального ожидания, пока `pytest-timeout` не убьёт тест десятками + секунд спустя. Именно так этот тест сам стал причиной 75-секундного + зависания CI-джобы (обрыв на ~9%, ни строки об ошибке) — не сам лимитер и + не синглтон `shared.py` (гипотеза координатора была разумной, но к этому + конкретному зависанию отношения не имела). """ sleep_calls: list[float] = [] + fake_now = [0.0] + + def fake_monotonic() -> float: + return fake_now[0] async def fake_sleep(seconds: float) -> None: sleep_calls.append(seconds) + fake_now[0] += seconds def handler(request: httpx.Request) -> httpx.Response: return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}}) _install_transport(handler) - with mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep): + with ( + mock.patch("app.services.tgbot.client.time.monotonic", fake_monotonic), + mock.patch("app.services.tgbot.client.asyncio.sleep", fake_sleep), + ): client = TelegramClient(token="fake-token", group_rate_limit_per_minute=2) await client.send_message(chat_id=42, text="a", message_thread_id=1) # тема поддержки await client.send_message(chat_id=42, text="b", message_thread_id=2) # тема алертов -- 2.45.3