diff --git a/tradein-mvp/backend/app/core/config.py b/tradein-mvp/backend/app/core/config.py index 46a15b08..485484d8 100644 --- a/tradein-mvp/backend/app/core/config.py +++ b/tradein-mvp/backend/app/core/config.py @@ -1287,6 +1287,51 @@ 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 и + # потерю сообщений). `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/отрицательное значение выключает лимитер для соответствующего процесса. + # + # ЧЕСТНО ПРО ОГРАНИЧЕНИЕ ЭТОГО ДИЗАЙНА (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" + ) + telegram_group_rate_limit_bot_per_minute: int = Field( + default=6, validation_alias="TELEGRAM_GROUP_RATE_LIMIT_BOT_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/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 eb8c2777..3313988e 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,150 @@ _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`), а не только на чтение/запись счётчика — это осознанно: + цель не просто «не гонять счётчик без гонки», а ФАКТИЧЕСКИ сериализовать + отправителей в этот чат, чтобы они не просыпались все разом по истечении + окна и не били по лимиту повторно. + + ВАЖНО про ключ (проверено на проде, 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` только на момент + создания записи — сам подсчёт/сон идёт уже под персональным локом чата. + """ + + 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, 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: + history.pop(0) + if len(history) < self._max_per_window: + history.append(now) + return + wait_s = history[0] + self._window_s - now + 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 секунд и на УСТАНОВКУ @@ -126,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. @@ -200,6 +364,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 +377,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 (вне цикла ретраев).""" @@ -282,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 — это воркерная политика, она НЕ меняется. Интерактивный вызывающий передаёт узкий @@ -298,8 +483,23 @@ class TelegramClient: это штатные 30-60с) на число попыток и подвесил бы синхронный HTTP-запрос на минуты — ровно то, от чего предостерегает докстринг `send_message`. """ - url = f"{self._base}/{method}" effective_timeout = timeout if timeout is not None else self._timeout + + # Лимитер применяется ТОЛЬКО к методам с `chat_id` в payload (отправка + # в конкретный чат) — `getUpdates` его не несёт и лимиту не подлежит. + # Списывается ОДИН слот на логический вызов `_request` (то есть на + # одну попытку отправки конкретного сообщения), а не на каждую HTTP + # попытку внутри ретрай-цикла ниже: ретраи по 429/5xx лечат один и тот + # же send, а не порождают новые отправки. Приоритет `wait_cap` — см. + # докстринг параметра `rate_limit_max_wait` выше. + chat_id = payload.get("chat_id") + if isinstance(chat_id, int): + 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}" # Раздельные таймауты считаем ОДИН раз и передаём per-request: у общего # AsyncClient свой дефолт, а бюджет ответа у каждого вызова свой. request_timeout = _request_timeout(effective_timeout) @@ -469,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, @@ -480,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( @@ -493,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 — текстовое сообщение (заголовки, приветствия, уведомления об ошибке). @@ -516,5 +723,103 @@ 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 {} + + 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..7d99b86f 100644 --- a/tradein-mvp/backend/app/services/tgbot/shared.py +++ b/tradein-mvp/backend/app/services/tgbot/shared.py @@ -36,6 +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, + # 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 61e06f30..a156b7d4 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,25 @@ async def _run_bridge() -> None: settings.telegram_bot_token, relay_base_url=settings.telegram_relay_base_url, relay_secret=settings.telegram_relay_secret, + # Бот-роль (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. Не блокирует и не роняет + # запуск при неудаче — см. докстринг `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/conftest.py b/tradein-mvp/backend/tests/conftest.py index 2e9f9d5e..0838150c 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 new file mode 100644 index 00000000..73a8ecfa --- /dev/null +++ b/tradein-mvp/backend/tests/services/tgbot/test_topic_check_and_group_rate_limit.py @@ -0,0 +1,377 @@ +"""Тесты для #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, + TelegramRateLimitedError, + 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 не звался никогда), что и есть баг. + + ВАЖНО (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.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) # тема алертов + 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) + + +# ── 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