Some checks failed
CI Trade-In / backend-tests (pull_request) Failing after 2m29s
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 11s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-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 / frontend-tests (pull_request) Has been skipped
Deep review of PR #3494 found the rate limiter unusable as designed: - H1: acquire() waited unbounded even for interactive HTTP handlers (support.py, glitchtip.py already pass a narrow `timeout` — reuse it as the queue wait cap instead of editing those handlers, which are out of scope here). New TelegramRateLimitedError (subclass of TelegramError) gives a fast, honest 502 instead of hanging past the caller's own budget. - H2: the limiter is per-process (in-memory), but two processes write to the same group (uvicorn API + bot worker) — giving each the same 18/min doubled the platform ceiling. Split into telegram_group_rate_limit_api_per_minute (12) and _bot_per_minute (6), sum kept below ~20. - M1: bridge.py sends without an explicit timeout inherited "wait forever", stalling the single-threaded poll loop (open DB session) past the SIGTERM drain window. Bounded via rate_limit_max_wait=20s at the six call sites. - M2: _locks/_sent_at grew unbounded on every unique DM chat_id. Added opportunistic cleanup of fully-expired entries. - L1: the "queue full" warning now logs once per acquire() call, not once per sleep iteration. - Corrected a factual error in the docstring: TELEGRAM_SUPPORT_CHAT_ID and TELEGRAM_ALERTS_CHAT_ID are the SAME group on prod (topics differ only). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG
825 lines
51 KiB
Python
825 lines
51 KiB
Python
"""Тонкая httpx-обёртка над Telegram Bot API (#tgsupport).
|
||
|
||
Зачем свой клиент, а не aiogram: единственные нужные методы — `getUpdates`
|
||
(long-polling), `copyMessage` (зеркалирование ЛЮБОГО типа контента без ре-аплоада)
|
||
и `sendMessage` (заголовки/приветствия/уведомления). aiogram — избыточная
|
||
зависимость (webhook-framework, dispatcher, FSM) ради трёх HTTP-вызовов; в стеке
|
||
уже есть httpx (см. `app.services.dadata`, `app.services.geocoder` — тот же паттерн
|
||
retry/timeout).
|
||
|
||
Docs: https://core.telegram.org/bots/api
|
||
|
||
Ретраи:
|
||
- HTTP 429 (Too Many Requests) — уважаем `parameters.retry_after` из тела ответа
|
||
(Telegram сам говорит сколько ждать), fallback на `_DEFAULT_RETRY_AFTER_S`.
|
||
- HTTP 5xx / транспортные ошибки (`httpx.TransportError`: timeout, connect,
|
||
обрыв протокола, прокси) — экспоненциальный backoff, `capped` на
|
||
`_MAX_BACKOFF_S`.
|
||
- Прочие отказы запроса (`httpx.RequestError`: битый ответ) — НЕ ретряются,
|
||
сразу `TelegramNetworkError`: повтор не чинит ни испорченный ответ, ни
|
||
кривую конфигурацию.
|
||
- Любая другая 4xx (400/401/403/404) — НЕ ретраится, сразу `TelegramApiError`
|
||
(запрос некорректен или прав нет — повтор не поможет).
|
||
|
||
Наружу летит только свой тип: `TelegramApiError` (площадка ответила отказом) или
|
||
`TelegramNetworkError` (не ответила), общий предок — `TelegramError`. Сырые
|
||
httpx-исключения из клиента не выходят: инвариант держат ДВА `except` в
|
||
`_request` — `httpx.TransportError` (ретраится) и страховочный
|
||
`httpx.RequestError` (не ретраится), вместе покрывающие всё дерево отказов
|
||
запроса, включая те, что появятся в httpx позже.
|
||
|
||
БЕЗОПАСНОСТЬ: наши `logger.*`-вызовы здесь содержат только имя метода API,
|
||
HTTP-статус и `description` из ответа Telegram — токен туда не пишем.
|
||
Это НЕ гарантирует, что токен не утечёт по другим стокам: он живёт в
|
||
`self._base`/`url` (локальные переменные stack-фрейма `_request`), а GlitchTip
|
||
(sentry_sdk) по умолчанию прикладывает locals к traceback и Httpx-интеграция
|
||
кладёт полный URL в span data. Эти стоки закрываются НЕ здесь, а в
|
||
`app.tgbot_main` (`include_local_variables=False`, `before_send`-редактор,
|
||
`traces_sample_rate=0.0`) и подавлением INFO-логов самого `httpx`-логгера
|
||
(который печатает полный request URL, включая токен, на уровне INFO).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import time
|
||
from typing import Any
|
||
|
||
import httpx
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
_DEFAULT_TIMEOUT_S = 15.0
|
||
_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 секунд и на УСТАНОВКУ
|
||
# соединения. Живой connect до api.telegram.org из прод-контейнера занимает
|
||
# 0.036с — 40-секундное ожидание коннекта было чистой слепотой: худший цикл
|
||
# 4 попытки × 40с + backoff ≈ 174с, и всё это время бот не видит ответов
|
||
# оператора (замер 12.09.2026: разрывы в логе 06:40:10 → 06:42:22 → 06:43:35,
|
||
# 576 строк `network error` и 7 полных исчерпаний бюджета ретраев за сутки).
|
||
# connect/write/pool к ожиданию ОТВЕТА Telegram отношения не имеют и коротки.
|
||
_CONNECT_TIMEOUT_S = 5.0
|
||
_WRITE_TIMEOUT_S = 10.0
|
||
_POOL_TIMEOUT_S = 5.0
|
||
|
||
# Пул keep-alive соединений на ОДИН экземпляр клиента. Параллелизма тут почти
|
||
# нет (long-polling — один запрос за раз, интерактивные ручки — единицы в
|
||
# минуту), так что смысл пула не в ширине, а в том, чтобы TCP+TLS-хендшейк не
|
||
# повторялся на каждый запрос и каждый ретрай.
|
||
_MAX_KEEPALIVE_CONNECTIONS = 5
|
||
_MAX_CONNECTIONS = 10
|
||
|
||
# Сколько держать простаивающее соединение. Задаём ЯВНО, потому что дефолт
|
||
# httpx — 5 секунд, и с ним пул не давал бы ничего там, где он нужнее всего:
|
||
# poll loop переиспользует соединение (следующий getUpdates уходит сразу), а
|
||
# вот веб-поддержка шлёт сообщения раз в минуты — за 5с соединение протухает и
|
||
# каждое зеркало снова платит полный TCP+TLS.
|
||
#
|
||
# Плата за длинный keep-alive — возросший шанс взять из пула соединение, которое
|
||
# уже закрыла та сторона; httpx отдаёт это как `RemoteProtocolError` («Server
|
||
# disconnected without sending a response»). Он ретраится с #3457, так что
|
||
# сценарий закрыт: попытка на протухшем соединении стоит один повтор, а не отказ.
|
||
_KEEPALIVE_EXPIRY_S = 90.0
|
||
|
||
|
||
class TelegramError(Exception):
|
||
"""Общий предок отказов клиента: и «ответил ok: false», и «не ответил вовсе».
|
||
|
||
Нужен ровно затем, чтобы вызывающий мог одной строкой сказать «Telegram не
|
||
сработал» и отдать свой 502. До #3456 сетевой отказ прилетал наружу сырым
|
||
`httpx.ConnectTimeout`, мимо `except TelegramApiError`, и FastAPI отдавал
|
||
500 — см. `TelegramNetworkError`.
|
||
"""
|
||
|
||
|
||
class TelegramApiError(TelegramError):
|
||
"""Telegram Bot API ответил `ok: false` (после исчерпания ретраев, если применимо)."""
|
||
|
||
def __init__(self, method: str, error_code: int, description: str) -> None:
|
||
self.method = method
|
||
self.error_code = error_code
|
||
self.description = description
|
||
super().__init__(f"Telegram API {method} failed: {error_code} {description}")
|
||
|
||
|
||
class TelegramNetworkError(TelegramError):
|
||
"""Ответа от Telegram не было: таймаут/обрыв, ретраи исчерпаны.
|
||
|
||
Отдельный тип, а не `TelegramApiError`, потому что `error_code`/`description`
|
||
брать неоткуда — Telegram ничего не сказал. Вызывающие, которым важна ТОЛЬКО
|
||
реакция площадки (`bridge`, разбирающий 403 «бот заблокирован»), продолжают
|
||
ловить `TelegramApiError` и этот отказ не перехватывают.
|
||
|
||
Причина сохраняется в `__cause__`: в GlitchTip виден исходный httpx-класс,
|
||
по которому и отличают таймаут соединения от сброса TLS (#3156).
|
||
"""
|
||
|
||
def __init__(self, method: str, reason: str, attempts: int) -> None:
|
||
self.method = method
|
||
self.reason = reason
|
||
self.attempts = attempts
|
||
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.
|
||
|
||
`read` — запрошенный бюджет ОТВЕТА (для long-poll это `poll_timeout + 10s`);
|
||
connect/write/pool фиксированы модульными константами и коротки: ждать
|
||
ответа Telegram — не то же самое, что ждать установки соединения.
|
||
"""
|
||
return httpx.Timeout(
|
||
connect=_CONNECT_TIMEOUT_S,
|
||
read=read,
|
||
write=_WRITE_TIMEOUT_S,
|
||
pool=_POOL_TIMEOUT_S,
|
||
)
|
||
|
||
|
||
def _extract_retry_after(
|
||
response: httpx.Response, default: float = _DEFAULT_RETRY_AFTER_S
|
||
) -> float:
|
||
"""Достаёт `parameters.retry_after` из тела 429-ответа. Fallback — `default`."""
|
||
try:
|
||
data = response.json()
|
||
except ValueError:
|
||
return default
|
||
if not isinstance(data, dict):
|
||
return default
|
||
params = data.get("parameters")
|
||
if isinstance(params, dict):
|
||
retry_after = params.get("retry_after")
|
||
if isinstance(retry_after, int | float):
|
||
return float(retry_after)
|
||
return default
|
||
|
||
|
||
def _error_from_body(response: httpx.Response) -> tuple[int, str]:
|
||
"""Парсит (error_code, description) из тела ответа Telegram; fallback на HTTP-статус."""
|
||
try:
|
||
data = response.json()
|
||
except ValueError:
|
||
return response.status_code, (response.text or "")[:200]
|
||
if not isinstance(data, dict):
|
||
return response.status_code, str(data)[:200]
|
||
error_code = data.get("error_code", response.status_code)
|
||
description = data.get("description", "")
|
||
code = int(error_code) if isinstance(error_code, int | float) else response.status_code
|
||
return code, str(description)
|
||
|
||
|
||
class TelegramClient:
|
||
"""Bot API клиент поверх ОДНОГО долгоживущего `httpx.AsyncClient`.
|
||
|
||
Соединение переиспользуется всё время жизни экземпляра: `AsyncClient`
|
||
создаётся лениво при первом запросе и хранится в `self._http`. Раньше он
|
||
создавался ВНУТРИ цикла ретраев — то есть keep-alive не было вовсе: полный
|
||
TCP+TLS-хендшейк на каждый запрос и на каждую повторную попытку, и заново
|
||
кидался кубик «встанет ли коннект». Для long-polling'а, ходящего каждые
|
||
~30с в бесконечном цикле, это была основная статья сетевых отказов.
|
||
|
||
Отсюда — требование к вызывающим: экземпляр НАДО переиспользовать (один на
|
||
процесс воркера, один на FastAPI-приложение — см.
|
||
`app.services.tgbot.shared`), а не создавать на каждый запрос, и закрывать
|
||
через `aclose()` или `async with`.
|
||
|
||
Таймаут теперь per-request: у `AsyncClient` он стоит дефолтом, а каждый
|
||
вызов `_request` передаёт свой `httpx.Timeout` (long-poll — свои 40с на
|
||
read, интерактивные ручки — свой узкий бюджет).
|
||
"""
|
||
|
||
def __init__(
|
||
self,
|
||
token: str,
|
||
base_url: str = "https://api.telegram.org",
|
||
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`,
|
||
# и `_post` ниже не делает второй попытки (фолбэчить с прямого пути
|
||
# НА прямой же путь бессмысленно). Заданный адрес переключает основной
|
||
# путь на ретранслятор, прямой остаётся ЗАПАСНЫМ на случай его отказа.
|
||
self._direct_base = f"{base_url}/bot{token}"
|
||
self._relay_base = f"{relay_base_url}/bot{token}" if relay_base_url else self._direct_base
|
||
self._relay_secret = relay_secret
|
||
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 (вне цикла ретраев)."""
|
||
if self._http is None:
|
||
self._http = httpx.AsyncClient(
|
||
timeout=_request_timeout(self._timeout),
|
||
limits=httpx.Limits(
|
||
max_keepalive_connections=_MAX_KEEPALIVE_CONNECTIONS,
|
||
max_connections=_MAX_CONNECTIONS,
|
||
keepalive_expiry=_KEEPALIVE_EXPIRY_S,
|
||
),
|
||
)
|
||
return self._http
|
||
|
||
async def aclose(self) -> None:
|
||
"""Закрывает пул соединений. Идемпотентно; после — клиент снова ленив."""
|
||
http, self._http = self._http, None
|
||
if http is not None:
|
||
await http.aclose()
|
||
|
||
async def __aenter__(self) -> TelegramClient:
|
||
return self
|
||
|
||
async def __aexit__(self, *_exc: object) -> None:
|
||
await self.aclose()
|
||
|
||
async def _post(
|
||
self,
|
||
client: httpx.AsyncClient,
|
||
url: str,
|
||
payload: dict[str, Any],
|
||
timeout: httpx.Timeout,
|
||
) -> httpx.Response:
|
||
"""POST с фолбэком на прямой путь при отказе РЕТРАНСЛЯТОРА (#3471).
|
||
|
||
Когда ретранслятор не настроен, `_relay_base is _direct_base`, `via_relay`
|
||
ниже всегда `False`, и метод ведёт себя как простой `client.post` — этот
|
||
путь ничем не отличается от поведения до #3471 (механизм отката).
|
||
|
||
Когда настроен: `url` бьёт в `self._relay_base`. Транспортный отказ (не
|
||
ответ Telegram ЧЕРЕЗ ретранслятор, а отказ ДО него — TCP/TLS до самого
|
||
relay-хоста) даёт РОВНО ОДНУ попытку напрямую к api.telegram.org — не
|
||
рекурсивно: если недоступен и прямой путь, исключение поднимается как
|
||
обычно и подхватывается retry-циклом `_request` на общих основаниях (со
|
||
следующей попытки цикл снова пробует ретранслятор first — временный
|
||
сбой relay не должен постоянно понижать клиента до прямого пути).
|
||
"""
|
||
via_relay = self._relay_base != self._direct_base and url.startswith(self._relay_base)
|
||
headers = (
|
||
{"X-Relay-Secret": self._relay_secret} if via_relay and self._relay_secret else None
|
||
)
|
||
try:
|
||
return await client.post(url, json=payload, timeout=timeout, headers=headers)
|
||
except httpx.TransportError:
|
||
if not via_relay:
|
||
raise
|
||
direct_url = self._direct_base + url[len(self._relay_base) :]
|
||
logger.warning(
|
||
"tg client: ретранслятор недоступен, одна попытка напрямую к Telegram"
|
||
)
|
||
return await client.post(direct_url, json=payload, timeout=timeout)
|
||
|
||
async def _request(
|
||
self,
|
||
method: str,
|
||
payload: dict[str, Any],
|
||
*,
|
||
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 — это
|
||
воркерная политика, она НЕ меняется. Интерактивный вызывающий передаёт узкий
|
||
потолок, потому что у него другой характер отказа: замер прода 01.09.2026 —
|
||
успешный запрос к api.telegram.org отвечает за 0.13с (максимум из 25 проб —
|
||
0.18с), а неудачный ВСЕГДА упирается в таймаут целиком (ConnectTimeout, ни
|
||
одного быстрого отказа). Это не троттлинг, а неустановленное соединение:
|
||
экспоненциальная пауза 2→4→8с не даёт удалённой стороне «остыть», она просто
|
||
добавляет 14 секунд к ожиданию пользователя. Потолок применяется и к 429:
|
||
иначе рост `max_retries` умножил бы Telegram-овский `retry_after` (для группы
|
||
это штатные 30-60с) на число попыток и подвесил бы синхронный HTTP-запрос на
|
||
минуты — ровно то, от чего предостерегает докстринг `send_message`.
|
||
"""
|
||
effective_timeout = timeout if timeout is not None else self._timeout
|
||
|
||
# Лимитер применяется ТОЛЬКО к методам с `chat_id` в payload (отправка
|
||
# в конкретный чат) — `getUpdates` его не несёт и лимиту не подлежит.
|
||
# Списывается ОДИН слот на логический вызов `_request` (то есть на
|
||
# одну попытку отправки конкретного сообщения), а не на каждую HTTP
|
||
# попытку внутри ретрай-цикла ниже: ретраи по 429/5xx лечат один и тот
|
||
# же send, а не порождают новые отправки. Приоритет `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)
|
||
backoff_cap = _MAX_BACKOFF_S if max_backoff is None else max_backoff
|
||
attempt = 0
|
||
# Клиент берём ДО цикла: пересоздавать его на каждую попытку значило бы
|
||
# заново платить за TCP+TLS ровно там, где сеть уже показала себя плохо.
|
||
client = self._http_client()
|
||
|
||
while True:
|
||
attempt += 1
|
||
try:
|
||
response = await self._post(client, url, payload, request_timeout)
|
||
except httpx.TransportError as exc:
|
||
# Ловим ВЕСЬ `TransportError`, а не узкий кортеж
|
||
# `(TimeoutException, NetworkError)`: `RemoteProtocolError`
|
||
# («Server disconnected without sending a response» — бытовой
|
||
# ответ api.telegram.org из РФ), `ProxyError`,
|
||
# `LocalProtocolError` и `UnsupportedProtocol` — СЁСТРЫ
|
||
# `NetworkError` по `TransportError`, а не наследники. Кортеж
|
||
# оставлял дыру ровно того класса, который чинил #3456: отказ
|
||
# вылетал сырым httpx мимо `except TelegramError` в ручках и
|
||
# снова давал 500 вместо 502 — и вдобавок не ретраился ни разу.
|
||
# Расширение ретраев на `RemoteProtocolError` наследует уже
|
||
# принятый здесь риск at-least-once (запрос мог дойти до
|
||
# Telegram, потерялся ответ) — он тот же, что у давно
|
||
# ретраящегося `ReadTimeout`; политика не меняется.
|
||
#
|
||
# Тип исключения обязан попасть в строку (#3156). У
|
||
# httpx.ReadError и httpx.ConnectError `str(exc)` пуст, и лог
|
||
# выглядел так: «network error (попытка 1/3): — retry через 2s»
|
||
# — после двоеточия пустота. По такой строке не отличить таймаут
|
||
# от обрыва соединения от сброса TLS, то есть диагностировать
|
||
# нечего. Замер на проде 27.08: 23 срабатывания за сутки, ни
|
||
# одной строки, по которой можно было бы что-то сказать.
|
||
reason = f"{type(exc).__name__}: {exc}" if str(exc) else type(exc).__name__
|
||
if attempt > max_retries:
|
||
logger.error(
|
||
"tg client: %s — network error после %d попыток: %s",
|
||
method,
|
||
attempt,
|
||
reason,
|
||
)
|
||
raise TelegramNetworkError(method, reason, attempt) from exc
|
||
backoff = min(2.0**attempt, backoff_cap)
|
||
logger.warning(
|
||
"tg client: %s — network error (попытка %d/%d): %s — retry через %.1fs",
|
||
method,
|
||
attempt,
|
||
max_retries,
|
||
reason,
|
||
backoff,
|
||
)
|
||
await asyncio.sleep(backoff)
|
||
continue
|
||
except httpx.RequestError as exc:
|
||
# Страховка на остаток дерева отказов запроса: сегодня это
|
||
# `DecodingError` (битая компрессия в ответе), завтра — всё, что
|
||
# httpx заведёт под `RequestError`. `TooManyRedirects` сюда НЕ
|
||
# относится: клиент создаётся с дефолтным `follow_redirects=False`
|
||
# и редиректы не ходит. Порядок `except`-ов
|
||
# значим: `TransportError` — наследник `RequestError`, и стоять
|
||
# обязан ВЫШЕ, иначе сетевые отказы перестали бы ретраиться.
|
||
#
|
||
# Без ретраев намеренно: это не «площадка недоступна», а
|
||
# испорченный ответ или кривая конфигурация — повтор не лечит
|
||
# ни то, ни другое, а пять попыток с backoff подвесили бы
|
||
# интерактивную ручку почти на минуту впустую. Свой тип тут
|
||
# нужен ровно за тем же, за чем и выше: чтобы ручка увидела
|
||
# `TelegramError` и отдала 502, а не 500.
|
||
reason = f"{type(exc).__name__}: {exc}" if str(exc) else type(exc).__name__
|
||
logger.error(
|
||
"tg client: %s — запрос не состоялся (без ретраев): %s", method, reason
|
||
)
|
||
raise TelegramNetworkError(method, reason, attempt) from exc
|
||
|
||
if response.status_code == 429:
|
||
retry_after = _extract_retry_after(response)
|
||
if max_backoff is not None:
|
||
# Потолок на retry_after — ТОЛЬКО когда его попросили явно.
|
||
# Воркеру Telegram-овские 30-60с надо уважать целиком, иначе
|
||
# мы долбимся в 429 и заводим лимит жёстче; интерактивному
|
||
# пути столько ждать нельзя ни при каких обстоятельствах —
|
||
# за ним стоит открытый HTTP-запрос от браузера.
|
||
retry_after = min(retry_after, max_backoff)
|
||
if attempt > max_retries:
|
||
error_code, description = _error_from_body(response)
|
||
logger.error("tg client: %s — 429 после %d попыток, сдаёмся", method, attempt)
|
||
raise TelegramApiError(method, error_code, description)
|
||
logger.warning(
|
||
"tg client: %s — HTTP 429 (попытка %d/%d), retry_after=%.0fs",
|
||
method,
|
||
attempt,
|
||
max_retries,
|
||
retry_after,
|
||
)
|
||
await asyncio.sleep(retry_after)
|
||
continue
|
||
|
||
if response.status_code >= 500:
|
||
if attempt > max_retries:
|
||
error_code, description = _error_from_body(response)
|
||
logger.error(
|
||
"tg client: %s — HTTP %d после %d попыток, сдаёмся",
|
||
method,
|
||
response.status_code,
|
||
attempt,
|
||
)
|
||
raise TelegramApiError(method, error_code, description)
|
||
backoff = min(2.0**attempt, backoff_cap)
|
||
logger.warning(
|
||
"tg client: %s — HTTP %d (попытка %d/%d) — retry через %.1fs",
|
||
method,
|
||
response.status_code,
|
||
attempt,
|
||
max_retries,
|
||
backoff,
|
||
)
|
||
await asyncio.sleep(backoff)
|
||
continue
|
||
|
||
if response.status_code >= 400:
|
||
# 4xx кроме 429 — запрос некорректен/прав нет, повтор не поможет.
|
||
error_code, description = _error_from_body(response)
|
||
raise TelegramApiError(method, error_code, description)
|
||
|
||
try:
|
||
data = response.json()
|
||
except ValueError as exc:
|
||
raise TelegramApiError(
|
||
method, response.status_code, f"invalid json: {exc}"
|
||
) from exc
|
||
|
||
if not isinstance(data, dict) or not data.get("ok"):
|
||
error_code, description = _error_from_body(response)
|
||
raise TelegramApiError(method, error_code, description)
|
||
|
||
return data.get("result")
|
||
|
||
async def get_updates(
|
||
self,
|
||
offset: int,
|
||
timeout: int = 30,
|
||
allowed_updates: list[str] | None = None,
|
||
) -> list[dict[str, Any]]:
|
||
"""Long-polling getUpdates. `timeout` — сколько Telegram держит запрос открытым (сек).
|
||
|
||
Запас `+10s` относится к READ-таймауту (сколько ждём ответа), чтобы не
|
||
обрывать соединение раньше, чем ответит сам Telegram long-poll. На
|
||
connect/write/pool он НЕ распространяется — они короткие и фиксированы
|
||
(`_CONNECT_TIMEOUT_S` и соседи): установка соединения либо занимает
|
||
десятки миллисекунд, либо не состоится вовсе.
|
||
"""
|
||
payload: dict[str, Any] = {"offset": offset, "timeout": timeout}
|
||
if allowed_updates is not None:
|
||
payload["allowed_updates"] = allowed_updates
|
||
result = await self._request(
|
||
"getUpdates", payload, timeout=float(timeout) + 10.0, max_retries=3
|
||
)
|
||
return result if isinstance(result, list) else []
|
||
|
||
async def copy_message(
|
||
self,
|
||
*,
|
||
chat_id: int,
|
||
from_chat_id: int,
|
||
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 — зеркалит ЛЮБОЙ тип контента без ре-аплоада файла.
|
||
|
||
`rate_limit_max_wait` — см. `TelegramClient._request`; используется
|
||
`bridge.py` для точечного потолка ожидания на конкретных местах (review M1)."""
|
||
payload: dict[str, Any] = {
|
||
"chat_id": chat_id,
|
||
"from_chat_id": from_chat_id,
|
||
"message_id": message_id,
|
||
}
|
||
if message_thread_id:
|
||
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, rate_limit_max_wait=rate_limit_max_wait
|
||
)
|
||
return result if isinstance(result, dict) else {}
|
||
|
||
async def send_message(
|
||
self,
|
||
*,
|
||
chat_id: int,
|
||
text: str,
|
||
message_thread_id: int | None = None,
|
||
reply_to_message_id: int | None = None,
|
||
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 — текстовое сообщение (заголовки, приветствия, уведомления об ошибке).
|
||
|
||
`timeout`/`max_retries` — по умолчанию наследуют воркерную политику
|
||
(`_DEFAULT_TIMEOUT_S`/`_DEFAULT_MAX_RETRIES`: на 429 спим Telegram-овский
|
||
`retry_after` — для группы это штатные 30-60с, на 5xx backoff до 30с).
|
||
Это ПРИЕМЛЕМО для `tgbot_main.py` (изолированный long-polling воркер), но
|
||
ФАТАЛЬНО для интерактивного HTTP-запроса (#tgsupport-web review H1) —
|
||
синхронный request/response путь не может легально висеть минуты. Вызывающая
|
||
сторона на interactive-пути ОБЯЗАНА передать узкий бюджет явно (см.
|
||
`app.api.v1.support.send_support_message`)."""
|
||
payload: dict[str, Any] = {"chat_id": chat_id, "text": text}
|
||
if message_thread_id:
|
||
payload["message_thread_id"] = message_thread_id
|
||
if reply_to_message_id:
|
||
payload["reply_to_message_id"] = reply_to_message_id
|
||
kwargs: dict[str, Any] = {}
|
||
if timeout is not None:
|
||
kwargs["timeout"] = timeout
|
||
if max_retries is not None:
|
||
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
|