gendesign/tradein-mvp/backend/app/services/tgbot/client.py
bot-backend 24c2052057
All checks were successful
CI Trade-In / changes (pull_request) Successful in 13s
CI / changes (pull_request) Successful in 15s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (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 / backend-tests (pull_request) Successful in 5m58s
style(tg): перенос длинной строки заголовков ретранслятора
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG
2026-09-12 14:27:10 +03:00

520 lines
31 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Тонкая 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
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
# Раздельные таймауты вместо скаляра. 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}")
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 = "",
) -> 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
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,
) -> Any:
"""POST `method` с JSON-телом `payload`. Ретраит 429/5xx/network, иначе raise сразу.
`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`.
"""
url = f"{self._base}/{method}"
effective_timeout = timeout if timeout is not None else self._timeout
# Раздельные таймауты считаем ОДИН раз и передаём per-request: у общего
# AsyncClient свой дефолт, а бюджет ответа у каждого вызова свой.
request_timeout = _request_timeout(effective_timeout)
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,
) -> dict[str, Any]:
"""copyMessage — зеркалит ЛЮБОЙ тип контента без ре-аплоада файла."""
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)
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,
) -> 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
result = await self._request("sendMessage", payload, **kwargs)
return result if isinstance(result, dict) else {}