Merge pull request 'fix(tgbot): проверка темы при старте + общий rate limit на группу (#3471)' (#3494) from feat/3471-tg-topic-check-and-group-rate-limit into main
Some checks failed
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / changes (push) Successful in 13s
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / deploy (push) Has been cancelled
Deploy Trade-In / test (push) Successful in 4m28s
Deploy Trade-In / build-backend (push) Successful in 2m0s
Some checks failed
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / changes (push) Successful in 13s
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / deploy (push) Has been cancelled
Deploy Trade-In / test (push) Successful in 4m28s
Deploy Trade-In / build-backend (push) Successful in 2m0s
This commit is contained in:
commit
ba2eb3b149
7 changed files with 819 additions and 9 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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).
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Add table
Reference in a new issue