"""Тонкая 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