gendesign/tradein-mvp/backend/app/services/tgbot/client.py
bot-backend d708f15019
All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 11s
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 4m50s
fix(tradein/support): один повтор терял каждое одиннадцатое сообщение в поддержку
Замер прода 01.09.2026 из контейнера бота: канал до api.telegram.org рвётся
всплесками, доля отказов на попытку 15-38% (пять проб: 3/8, 15/40, 5/20, 3/20,
1/25), в логе long-polling'а 353 ConnectTimeout за сутки. Транспорт ни при чём —
httpx и сырой сокет отваливаются одинаково (25% против 35% в чередующемся
замере), и прокси не помогает, а мешает: через SCRAPER_PROXY_URL 0 из 20.

Ручка веб-поддержки ходила с max_retries=1, то есть двумя попытками. При 30%
отказов на попытку до пользователя доходило ~9% отказов — каждое одиннадцатое
сообщение возвращало 502 «сервис недоступен».

Два других числа из того же замера задают конструкцию. Успешный запрос отвечает
за 0.13с (максимум из 25 проб — 0.18с), а неудачный НИКОГДА не отваливается
быстро: все отказы упираются в таймаут целиком (10.02с при timeout=10.0). Значит
десятисекундный таймаут не покупал ничего, кроме цены за неудачу, — снижен до 5с,
это ~28-кратный запас к измеренному максимуму. И экспоненциальная пауза 2→4→8с
здесь бессмысленна: отказ — неустановленное соединение, а не троттлинг, пережидать
нечего; она лишь добавляла 14с к ожиданию.

Правка: бюджет ручки — 3 повтора, таймаут 5с, потолок паузы 1с. Худший случай
4 попытки × 5с + 3 паузы × 1с = 23с и требует четырёх отказов подряд; типичный
случай не меняется (0.13с). Расчётная потеря падает с ~9% до ~0.8%.

В TelegramClient добавлен необязательный max_backoff. Воркерная политика НЕ
меняется: без явного потолка откат прежний экспоненциальный до 30с, а retry_after
из 429 уважается целиком — эту границу держит отдельный тест, потому что первая
версия правки её сломала (капала 60с до 30с и для воркера тоже). Потолок на
retry_after применяется только когда его передали явно: интерактивному пути
нельзя ждать Telegram-овские 30-60с, за ним стоит открытый запрос от браузера.

Тесты: 8 новых (потолок на network/429/5xx, неизменность воркерного пути,
арифметика «max_retries=N → N+1 попыток», границы бюджета ручки).
2026-09-01 09:53:21 +03:00

302 lines
15 KiB
Python
Raw Permalink 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 / сетевые ошибки (timeout/connect) — экспоненциальный backoff,
`capped` на `_MAX_BACKOFF_S`.
- Любая другая 4xx (400/401/403/404) — НЕ ретраится, сразу `TelegramApiError`
(запрос некорректен или прав нет — повтор не поможет).
БЕЗОПАСНОСТЬ: наши `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
class TelegramApiError(Exception):
"""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}")
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. Каждый вызов — отдельное короткоживущее соединение."""
def __init__(
self,
token: str,
base_url: str = "https://api.telegram.org",
timeout: float = _DEFAULT_TIMEOUT_S,
) -> None:
self._base = f"{base_url}/bot{token}"
self._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
backoff_cap = _MAX_BACKOFF_S if max_backoff is None else max_backoff
attempt = 0
while True:
attempt += 1
try:
async with httpx.AsyncClient(timeout=effective_timeout) as client:
response = await client.post(url, json=payload)
except (httpx.TimeoutException, httpx.NetworkError) as exc:
# Тип исключения обязан попасть в строку (#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
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
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 держит запрос открытым (сек).
HTTP-таймаут запроса берётся с запасом (`timeout + 10s`), чтобы не обрывать
соединение раньше, чем ответит сам Telegram long-poll.
"""
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 {}