fix(tgbot): проверка темы при старте + общий rate limit на группу (#3471)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 8m55s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 8m55s
Бот падал молча на каждом сообщении, если тему форума удалили/переименовали: узнавали об этом только по отсутствию сообщений у людей. 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 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG
This commit is contained in:
parent
994eb79323
commit
6433477f7c
5 changed files with 491 additions and 1 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
Loading…
Add table
Reference in a new issue