"""Тонкая 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, ) -> None: self._base = f"{base_url}/bot{token}" 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 _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 client.post(url, json=payload, timeout=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 {}