Compare commits

..

No commits in common. "a3d4fcf0b3f7165faf620c0a8b2369e18cbf56be" and "8994e041cf82f4edde788fe62830a583c1854061" have entirely different histories.

12 changed files with 46 additions and 776 deletions

View file

@ -53,8 +53,7 @@ from fastapi import APIRouter, Header, HTTPException, Query, Request
from pydantic import BaseModel, ConfigDict, ValidationError
from app.core.config import settings
from app.services.tgbot.client import TelegramError
from app.services.tgbot.shared import get_telegram_client
from app.services.tgbot.client import TelegramClient, TelegramError
logger = logging.getLogger(__name__)
@ -214,10 +213,7 @@ async def glitchtip_webhook(
received_at = datetime.now(UTC)
text = _build_message(raw_body, received_at)
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
client = get_telegram_client()
client = TelegramClient(settings.telegram_bot_token)
try:
await client.send_message(
chat_id=settings.telegram_alerts_chat_id,

View file

@ -73,8 +73,7 @@ from app.core.db import get_db
from app.core.ratelimit import SlidingWindowLimiter, _client_ip
from app.services.tgbot import web_support_storage as storage
from app.services.tgbot.bridge import SERVICE_UNAVAILABLE_TEXT
from app.services.tgbot.client import TelegramError
from app.services.tgbot.shared import get_telegram_client
from app.services.tgbot.client import TelegramClient, TelegramError
logger = logging.getLogger(__name__)
@ -220,10 +219,7 @@ async def send_support_message(
headers={"Retry-After": str(int(retry_after) + 1)},
)
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
client = get_telegram_client()
client = TelegramClient(settings.telegram_bot_token)
try:
mirrored = await client.send_message(
chat_id=settings.telegram_support_chat_id,
@ -415,10 +411,7 @@ async def send_anon_support_message(
)
display_id = _anon_display_id(token)
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
client = get_telegram_client()
client = TelegramClient(settings.telegram_bot_token)
try:
mirrored = await client.send_message(
chat_id=settings.telegram_support_chat_id,

View file

@ -50,7 +50,6 @@ from app.core.rbac import rbac_guard
from app.core.request_audit import RequestAuditMiddleware
from app.observability import metrics as app_metrics
from app.observability.sentry_scrub import scrub_pii_event
from app.services.tgbot.shared import close_telegram_client, init_telegram_client
logger = logging.getLogger(__name__)
@ -217,19 +216,7 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
# in the tradein-scraper container (`python -m app.scheduler_main`, kit scheduler).
# Prod backend has always run with SCHEDULER_ENABLE=false (see docker-compose.prod.yml);
# this API process never actually launched scheduler_loop() in production.
# Общий на приложение Telegram-клиент (#tg-connection-resilience): ручки
# support/glitchtip раньше создавали его на КАЖДЫЙ запрос, то есть каждое
# зеркалирование сообщения начиналось с полного TCP+TLS-хендшейка до
# api.telegram.org. Один пул keep-alive на процесс, закрываем на shutdown.
# Без токена не создаём: ручки в этом случае и так отвечают 503.
if settings.telegram_bot_token:
init_telegram_client()
try:
yield
finally:
await close_telegram_client()
yield
app = FastAPI(

View file

@ -29,13 +29,9 @@
Telegram 403 (клиент заблокировал бота) is_blocked=true + уведомление в
топике (только для Telegram-ветки у веб-клиента нет "заблокировал бота").
C) Дедуп: update_id <= сохранённого offset skip. Offset сохраняется И
коммитится в той же транзакции, что и запись сообщения (см. `process_update`),
после КАЖДОГО апдейта рестарт воркера не переигрывает уже обработанные
апдейты и не подвисает вечно на «ядовитом» апдейте. Исключение
ТРАНЗИЕНТНЫЙ сетевой отказ (`TelegramNetworkError`, т.е. исчерпанный бюджет
ретраев клиента): такой апдейт СОЗНАТЕЛЬНО остаётся неподтверждённым, чтобы
Telegram отдал его снова, иначе ответ оператора пропадал бы навсегда
(#tg-connection-resilience). Потолок переигрываний — `_MAX_NETWORK_REPLAYS`.
коммитится в той же транзакции, что и запись сообщения (см. `process_update`
`finally`), после КАЖДОГО апдейта рестарт воркера не переигрывает уже
обработанные апдейты и не подвисает вечно на «ядовитом» апдейте.
D) TELEGRAM_BOT_TOKEN пуст бот выключен проверяется в `app.tgbot_main`
(entrypoint), не здесь.
E) /start клиенту короткое приветствие МЕРЫ, без зеркалирования в топик
@ -80,7 +76,7 @@ from app.core.config import settings
from app.core.ratelimit import SlidingWindowLimiter
from app.core.shutdown import shutdown_requested
from app.services.tgbot import web_support_storage
from app.services.tgbot.client import TelegramApiError, TelegramClient, TelegramNetworkError
from app.services.tgbot.client import TelegramApiError, TelegramClient
logger = logging.getLogger(__name__)
@ -135,44 +131,6 @@ FLOOD_LIMITED_TEXT = (
# получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил).
_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT = "Веб-чат поддерживает только текст, сообщение не доставлено."
# Уведомления оператору в топике отправляются ИЗНУТРИ poll loop, который
# однопоточный: пока висит одна отправка, не обрабатывается НИ ОДИН следующий
# апдейт. Поэтому им нужен интерактивный бюджет, а не воркерный дефолт клиента
# (5 ретраев, backoff до 30с, полный retry_after на 429 — минуты стопа на
# ВТОРИЧНОМ действии). Числа — те же, что у интерактивных отправок веб-чата
# поддержки (`_INTERACTIVE_SEND_*` в app/api/v1/support.py, подобраны замером
# прода #tgsupport-retry); сознательно ДУБЛИРУЕМ, а не импортируем из слоя API —
# воркер не должен зависеть от роутера.
_NOTIFY_SEND_TIMEOUT_S = 5.0
_NOTIFY_SEND_MAX_RETRIES = 3
_NOTIFY_SEND_MAX_BACKOFF_S = 1.0
# Потолок переигрываний ОДНОГО update_id на транзиентных сетевых отказах
# (#tg-connection-resilience).
#
# Зачем потолок: без него «вечно недоставляемый» апдейт заклинил бы очередь
# НАВСЕГДА — ровно то, от чего защищал прежний безусловный `finally: save_offset`.
# Потерять одно сообщение плохо, потерять весь поток — хуже, поэтому на потолке
# offset всё-таки сдвигается, но ГРОМКО (`logger.error`), а не молча.
#
# Про дубли: `TelegramNetworkError` означает исчерпанный бюджет ретраев клиента,
# при этом запрос МОГ дойти до Telegram (потерялся ответ) — переигрывание тогда
# доставит клиенту то же сообщение второй раз. Это осознанный at-least-once
# компромисс: дубль и клиент, и оператор ВИДЯТ и могут поправить, а тихая потеря
# ответа не оставляет следа нигде, кроме строчки в логе. Полноценная
# идемпотентность по (update_id, target_chat_id) потребовала бы нового
# персистентного состояния (колонка/таблица + миграция) ради редкого случая;
# вместо этого число возможных дублей жёстко ограничено сверху — не больше
# (_MAX_NETWORK_REPLAYS - 1) повторов на апдейт.
_MAX_NETWORK_REPLAYS = 3
# update_id -> сколько раз мы уже отказались подтверждать этот апдейт.
# In-memory осознанно: воркер long-polling однопоточный, запись живёт ровно до
# подтверждения апдейта (`pop` в `process_update`), так что словарь не растёт.
# Рестарт воркера обнуляет счётчик — это допустимо (новый процесс = новая сеть),
# потолок всё равно действует в пределах каждой жизни процесса.
_network_replay_attempts: dict[int, int] = {}
# ── Storage abstraction (testable без реальной БД) ──────────────────────────
class BridgeStorage(Protocol):
@ -442,43 +400,6 @@ def _format_topic_header(
return f"Новое обращение от {display_name}{username_part} (chat_id={chat_id})"
async def _notify_topic(
client: TelegramClient,
*,
text: str,
reply_to_message_id: int | None,
context: str,
) -> None:
"""Служебное уведомление оператору в support-топик (вторичное действие).
Два свойства, которых не было у прямых `client.send_message` вызовов:
1) узкий интерактивный бюджет (`_NOTIFY_SEND_*`) иначе одна такая отправка
стопорит весь однопоточный poll loop на минуты;
2) собственный `except` провал УВЕДОМЛЕНИЯ не отменяет основную ветку
обработки (клиент уже помечен заблокированным / медиа-реплай уже отклонён)
и не решает судьбу апдейта.
`context` только технические идентификаторы (chat_id/thread_id), НЕ текст
переписки: логи моста принципиально не содержат ПДн.
"""
try:
await client.send_message(
chat_id=settings.telegram_support_chat_id,
text=text,
message_thread_id=settings.telegram_support_topic_id or None,
reply_to_message_id=reply_to_message_id,
timeout=_NOTIFY_SEND_TIMEOUT_S,
max_retries=_NOTIFY_SEND_MAX_RETRIES,
max_backoff=_NOTIFY_SEND_MAX_BACKOFF_S,
)
except Exception:
logger.warning(
"tgbot bridge: не удалось отправить уведомление оператору в топик (%s) — "
"основная ветка обработки не отменяется",
context,
exc_info=True,
)
# ── Update routing ────────────────────────────────────────────────────────────
async def _handle_private_message(
message: dict[str, Any], client: TelegramClient, storage: BridgeStorage
@ -632,14 +553,14 @@ async def _handle_group_reply(
if exc.error_code == 403:
# Клиент заблокировал бота — фиксируем и уведомляем оператора в топике.
storage.mark_blocked(target_chat_id)
await _notify_topic(
client,
await client.send_message(
chat_id=settings.telegram_support_chat_id,
text=(
f"Не удалось доставить сообщение клиенту (chat_id={target_chat_id}) — "
"бот заблокирован."
),
message_thread_id=settings.telegram_support_topic_id or None,
reply_to_message_id=message_id,
context=f"403 на доставке клиенту chat_id={target_chat_id}",
)
return
raise
@ -671,11 +592,11 @@ async def _handle_group_reply(
web_thread_id,
kind,
)
await _notify_topic(
client,
await client.send_message(
chat_id=settings.telegram_support_chat_id,
text=_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT,
message_thread_id=settings.telegram_support_topic_id or None,
reply_to_message_id=message_id if isinstance(message_id, int) else None,
context=f"медиа-реплай на веб-зеркало thread_id={web_thread_id}",
)
return
@ -715,41 +636,25 @@ async def _handle_group_reply(
async def process_update(
update: dict[str, Any], client: TelegramClient, storage: BridgeStorage
) -> bool:
) -> None:
"""Маршрутизирует один Telegram update. Дедуп (C) + атомарный offset-commit.
Возвращает True, если offset сдвинут (апдейт подтверждён, Telegram его больше
не отдаст), и False, если апдейт СОЗНАТЕЛЬНО оставлен неподтверждённым ради
переигрывания. На False вызывающий (`run_poll_loop`) ОБЯЗАН прервать разбор
пачки: offset у Telegram единая «высшая отметка», подтверждение любого
СЛЕДУЮЩЕГО апдейта неявно подтвердило бы и этот, и переигрывания не было бы.
Дедуп: update_id <= сохранённого offset skip без side-effects. Offset
сохраняется и коммитится ПОСЛЕ обработки (в т.ч. если обработка упала
иначе «ядовитый» апдейт блокировал бы весь поток навсегда).
Дедуп: update_id <= сохранённого offset skip без side-effects.
Судьба offset'а по классам отказа:
- `TelegramNetworkError` (транзиентный: Telegram не ответил, бюджет ретраев
клиента исчерпан) offset НЕ двигаем, `rollback()` частичных записей,
апдейт переигрывается на следующей итерации. Иначе ответ оператора
терялся НАВСЕГДА: copyMessage не дошёл, `record_message` не выполнился,
Telegram апдейт больше не отдаст, а оператор уверен, что ответил
(#tg-connection-resilience). Ограничено `_MAX_NETWORK_REPLAYS` — на
потолке offset всё-таки сдвигается с `logger.error`, иначе «вечно
недоставляемый» апдейт заклинил бы поток навсегда.
- `SQLAlchemyError` (сбой БД) offset двигаем, но `rollback()` ОБЯЗАН
отработать ПЕРЕД `save_offset`: сбой БД оставляет сессию в
failed-transaction state, иначе `save_offset` сам кинет
`PendingRollbackError`, `process_update` вылетит без сохранения offset'а,
следующая итерация получит СТАРЫЙ offset от `get_offset()` и переиграет
тот же апдейт copyMessage задублирует зеркало клиента в топике на
каждый повтор поллинга (#3 review, воспроизведено).
- любое прочее исключение (в т.ч. `TelegramApiError` площадка ОТВЕТИЛА
отказом, повтор ничего не изменит) offset двигаем, «ядовитый» апдейт
не блокирует поток.
Различаем сбой БД (`SQLAlchemyError`) от прочих (Telegram API и т.п.):
сбой БД оставляет сессию в failed-transaction state `rollback()` ОБЯЗАН
отработать ПЕРЕД `save_offset`, иначе тот сам кинет `PendingRollbackError`,
`process_update` вылетит без сохранения offset'а, следующая итерация
`run_poll_loop` получит СТАРЫЙ offset от `get_offset()` и переиграет тот же
апдейт заново copyMessage задублирует зеркало клиента в топике на
каждый повтор поллинга (#3 review, воспроизведено).
"""
update_id = update.get("update_id")
if not isinstance(update_id, int):
logger.warning("tgbot bridge: update без валидного update_id — игнор")
return True
return
current_offset = storage.get_offset()
if update_id <= current_offset:
@ -758,7 +663,7 @@ async def process_update(
update_id,
current_offset,
)
return True
return
message = update.get("message")
try:
@ -772,35 +677,6 @@ async def process_update(
await _handle_group_reply(message, client, storage)
# иначе — необрабатываемый тип чата/апдейта (edited_message, канал и
# т.п.) — тихий игнор, но offset всё равно сдвигаем ниже.
except TelegramNetworkError:
attempts = _network_replay_attempts.get(update_id, 0) + 1
# Частичные записи этого апдейта не должны уехать в БД чужим commit'ом
# (сессия одна на всю пачку) — переигрывание начинается с чистого листа.
storage.rollback()
if attempts < _MAX_NETWORK_REPLAYS:
_network_replay_attempts[update_id] = attempts
logger.warning(
"tgbot bridge: Telegram недоступен на update_id=%d (отказ %d из %d) — "
"offset НЕ сдвигаем, апдейт переиграется на следующей итерации",
update_id,
attempts,
_MAX_NETWORK_REPLAYS,
)
return False
logger.error(
"tgbot bridge: update_id=%d исчерпал потолок переигрываний (%d сетевых "
"отказов подряд) — сдвигаем offset, содержимое апдейта ПОТЕРЯНО; поток не "
"блокируем, требуется ручной разбор support-топика "
"(chat_id=%s, message_id=%s)",
update_id,
_MAX_NETWORK_REPLAYS,
# Идентификаторы, а НЕ текст: это единственная строка, по которой
# человек найдёт потерянный ответ оператора в топике и перешлёт его
# руками. Без них в логе остаётся только update_id, которого в
# интерфейсе Telegram не видно. Текст сообщения — ПДн, в лог не идёт.
(message or {}).get("chat", {}).get("id") if isinstance(message, dict) else None,
(message or {}).get("message_id") if isinstance(message, dict) else None,
)
except SQLAlchemyError:
logger.exception(
"tgbot bridge: DB-ошибка на update_id=%d — rollback перед сохранением "
@ -814,12 +690,9 @@ async def process_update(
"(не блокируем поток на 'ядовитом' апдейте)",
update_id,
)
# Апдейт подтверждён — счётчик переигрываний больше не нужен (словарь не растёт).
_network_replay_attempts.pop(update_id, None)
storage.save_offset(update_id)
storage.commit()
return True
finally:
storage.save_offset(update_id)
storage.commit()
# ── Long-polling loop ─────────────────────────────────────────────────────────
@ -844,20 +717,8 @@ async def run_poll_loop(
updates = await client.get_updates(
offset=offset + 1, timeout=poll_timeout_s, allowed_updates=["message"]
)
for idx, update in enumerate(updates):
if not await process_update(update, client, storage):
# Апдейт намеренно не подтверждён (транзиентный сетевой
# отказ). Обрабатывать остаток пачки НЕЛЬЗЯ: offset —
# единая «высшая отметка», подтверждение следующего
# апдейта неявно подтвердило бы и этот. Остаток Telegram
# отдаст заново на следующей итерации.
logger.warning(
"tgbot bridge: update_id=%s не подтверждён — остаток пачки "
"(%d апдейтов) разберём на следующей итерации",
update.get("update_id"),
len(updates) - idx - 1,
)
break
for update in updates:
await process_update(update, client, storage)
consecutive_errors = 0
except Exception:
consecutive_errors += 1

View file

@ -54,38 +54,6 @@ _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», и «не ответил вовсе».
@ -126,21 +94,6 @@ class TelegramNetworkError(TelegramError):
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:
@ -174,24 +127,7 @@ def _error_from_body(response: httpx.Response) -> tuple[int, str]:
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, интерактивные ручки свой узкий бюджет).
"""
"""Bot API клиент на httpx.AsyncClient. Каждый вызов — отдельное короткоживущее соединение."""
def __init__(
self,
@ -201,32 +137,6 @@ class TelegramClient:
) -> 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,
@ -254,19 +164,14 @@ class TelegramClient:
"""
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)
async with httpx.AsyncClient(timeout=effective_timeout) as client:
response = await client.post(url, json=payload)
except httpx.TransportError as exc:
# Ловим ВЕСЬ `TransportError`, а не узкий кортеж
# `(TimeoutException, NetworkError)`: `RemoteProtocolError`
@ -401,11 +306,8 @@ class TelegramClient:
) -> list[dict[str, Any]]:
"""Long-polling getUpdates. `timeout` — сколько Telegram держит запрос открытым (сек).
Запас `+10s` относится к READ-таймауту (сколько ждём ответа), чтобы не
обрывать соединение раньше, чем ответит сам Telegram long-poll. На
connect/write/pool он НЕ распространяется они короткие и фиксированы
(`_CONNECT_TIMEOUT_S` и соседи): установка соединения либо занимает
десятки миллисекунд, либо не состоится вовсе.
HTTP-таймаут запроса берётся с запасом (`timeout + 10s`), чтобы не обрывать
соединение раньше, чем ответит сам Telegram long-poll.
"""
payload: dict[str, Any] = {"offset": offset, "timeout": timeout}
if allowed_updates is not None:

View file

@ -1,50 +0,0 @@
"""Общий на приложение `TelegramClient` (#tg-connection-resilience).
Зачем: до этого три HTTP-ручки (`api.v1.support` ×2, `api.v1.glitchtip`) делали
`TelegramClient(token)` на КАЖДЫЙ входящий запрос, а клиент внутри пересоздавал
`httpx.AsyncClient` на каждую попытку то есть keep-alive не было ни на каком
уровне и каждый запрос начинался с полного TCP+TLS-хендшейка до
api.telegram.org. Здесь живёт один экземпляр на процесс: создаётся в lifespan
(`app.main`), закрывается на shutdown там же.
Воркер бота (`app.tgbot_main`) сюда НЕ ходит у него свой процесс без ASGI и
свой экземпляр на всё время жизни поллинга.
"""
from __future__ import annotations
import logging
from app.core.config import settings
from app.services.tgbot.client import TelegramClient
logger = logging.getLogger(__name__)
_client: TelegramClient | None = None
def get_telegram_client() -> TelegramClient:
"""Общий клиент приложения. Ленив: создаётся при первом обращении.
Ленивость (а не «только из lifespan») нужна из-за kill-switch: при пустом
`TELEGRAM_BOT_TOKEN` в lifespan создавать нечего, а тесты ручек поднимают
приложение без прохода через startup.
"""
global _client
if _client is None:
_client = TelegramClient(settings.telegram_bot_token)
return _client
def init_telegram_client() -> TelegramClient:
"""Явное создание на старте приложения (lifespan)."""
return get_telegram_client()
async def close_telegram_client() -> None:
"""Закрывает общий клиент на shutdown. Идемпотентно."""
global _client
client, _client = _client, None
if client is not None:
await client.aclose()
logger.info("tg shared client: пул соединений закрыт")

View file

@ -114,12 +114,8 @@ def _should_run() -> bool:
async def _run_bridge() -> None:
# `async with` — чтобы пул keep-alive соединений закрывался при любом выходе
# из поллинга (кооперативный drain по SIGTERM, hard-cancel, исключение).
# Клиент один на весь процесс: пересоздание на запрос убивало keep-alive и
# заставляло каждый long-poll начинаться с TCP+TLS-хендшейка.
async with TelegramClient(settings.telegram_bot_token) as client:
await run_poll_loop(client, SessionLocal)
client = TelegramClient(settings.telegram_bot_token)
await run_poll_loop(client, SessionLocal)
async def _await_bridge(task: asyncio.Task[None]) -> None:

View file

@ -22,12 +22,6 @@ Coverage (per task spec + review follow-up):
- сбой БД (SQLAlchemyError) во время обработки rollback() ПЕРЕД save_offset,
offset всё равно сдвигается без этого следующий поллинг переиграл бы тот
же апдейт и задублировал зеркало клиента в топике (#3 review)
- ТРАНЗИЕНТНЫЙ сетевой отказ (TelegramNetworkError) на доставке ответа оператора
offset НЕ сдвигается, апдейт переигрывается и доходит до клиента; на потолке
`_MAX_NETWORK_REPLAYS` offset всё-таки сдвигается (поток не заклинен); poll loop
прерывает разбор пачки на неподтверждённом апдейте (#tg-connection-resilience)
- провал ВТОРИЧНОГО уведомления оператору в топик не отменяет основную ветку
(is_blocked остаётся, offset сдвигается) и логируется отдельной строкой
- /start приветствие без зеркалирования
- Telegram 403 на доставку оператору is_blocked + уведомление в топике
- TELEGRAM_SUPPORT_CHAT_ID не задан клиенту уходит "сервис недоступен"
@ -51,7 +45,7 @@ from sqlalchemy.exc import SQLAlchemyError
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from app.services.tgbot import bridge
from app.services.tgbot.client import TelegramClient, TelegramNetworkError
from app.services.tgbot.client import TelegramClient
SUPPORT_CHAT_ID = -100123456789
SUPPORT_TOPIC_ID = 42
@ -979,219 +973,3 @@ async def test_update_from_unrelated_chat_is_ignored_but_offset_advances() -> No
assert calls == []
assert storage.messages == []
assert storage.get_offset() == 40
# ── сетевая устойчивость (#tg-connection-resilience) ─────────────────────────
@pytest.fixture(autouse=True)
def _reset_network_replay_attempts() -> None:
"""`bridge._network_replay_attempts` — module-level словарь, его состояние
иначе протекало бы между тестами (потолок переигрываний виден глобально)."""
bridge._network_replay_attempts.clear()
def _network_boom(method: str = "copyMessage"):
"""Асинхронная заглушка метода клиента, изображающая исчерпанный бюджет ретраев."""
async def _raise(*_args: object, **_kwargs: object) -> dict[str, Any]:
raise TelegramNetworkError(method, "ConnectTimeout", 4)
return _raise
async def test_network_failure_on_operator_reply_keeps_offset_for_replay(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Главный сценарий потери: оператор ответил, Telegram в этот момент недоступен.
Раньше `finally: save_offset` подтверждал апдейт ответ не доходил до клиента
НИКОГДА (Telegram апдейт больше не отдаёт, записи нет, оператор уверен, что
ответил). Теперь offset остаётся прежним, апдейт переигрывается и доходит.
"""
calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({"copyMessage": {"message_id": 42}}, calls)
storage = FakeBridgeStorage()
storage.record_message(
chat_id=555,
direction="in",
tg_message_id=1,
topic_message_id=100,
kind="text",
text_body="вопрос клиента",
operator_tg_id=None,
)
update = {"update_id": 40, "message": _group_reply_message(reply_to_message_id=100)}
# Падает РОВНО первый copyMessage (monkeypatch.undo() тут не годится: он снял бы
# и autouse-патчи настроек support-группы, и ветка просто перестала бы работать).
real_copy = client.copy_message
failures = {"left": 1}
async def flaky_copy(**kwargs: Any) -> dict[str, Any]:
if failures["left"] > 0:
failures["left"] -= 1
raise TelegramNetworkError("copyMessage", "ConnectTimeout", 4)
return await real_copy(**kwargs)
monkeypatch.setattr(client, "copy_message", flaky_copy)
advanced = await bridge.process_update(update, client, storage)
assert advanced is False
assert storage.get_offset() == 0 # апдейт НЕ подтверждён — Telegram отдаст его снова
assert storage.commits == 0
assert storage.rollbacks == 1 # частичные записи не уедут чужим commit'ом
assert len(storage.messages) == 1 # фейковой записи 'out' не появилось
# Переигрывание: сеть починилась — тот же апдейт доставляется и подтверждается.
advanced = await bridge.process_update(update, client, storage)
assert advanced is True
assert storage.get_offset() == 40
assert [m for m, _ in calls] == ["copyMessage"]
out_rec = storage.messages[-1]
assert out_rec["direction"] == "out"
assert out_rec["chat_id"] == 555
async def test_network_failure_stops_advancing_only_until_replay_cap(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""Потолок переигрываний: «вечно недоставляемый» апдейт не должен заклинить
поток навсегда (ровно то, от чего защищал прежний безусловный `finally`)."""
calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({}, calls)
storage = FakeBridgeStorage()
storage.record_message(
chat_id=555,
direction="in",
tg_message_id=1,
topic_message_id=100,
kind="text",
text_body="вопрос клиента",
operator_tg_id=None,
)
monkeypatch.setattr(client, "copy_message", _network_boom())
update = {"update_id": 41, "message": _group_reply_message(reply_to_message_id=100)}
for _ in range(bridge._MAX_NETWORK_REPLAYS - 1):
assert await bridge.process_update(update, client, storage) is False
assert storage.get_offset() == 0
with caplog.at_level(logging.ERROR, logger="app.services.tgbot.bridge"):
advanced = await bridge.process_update(update, client, storage)
assert advanced is True
assert storage.get_offset() == 41 # поток разблокирован
assert "потолок переигрываний" in caplog.text
# Счётчик снят — словарь не растёт от апдейта к апдейту.
assert bridge._network_replay_attempts == {}
async def test_poll_loop_stops_batch_on_unconfirmed_update(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Offset у Telegram — единая «высшая отметка»: подтвердив СЛЕДУЮЩИЙ апдейт
пачки, мы неявно подтвердили бы неудавшийся, и переигрывания не случилось бы.
Поэтому разбор пачки обрывается на первом неподтверждённом апдейте."""
processed: list[int] = []
async def fake_process(update: dict[str, Any], client: object, storage: object) -> bool:
processed.append(update["update_id"])
return update["update_id"] != 51 # 51 — сетевой отказ
monkeypatch.setattr(bridge, "process_update", fake_process)
iteration = {"n": 0}
def fake_shutdown() -> bool:
iteration["n"] += 1
return iteration["n"] > 1 # ровно одна итерация poll loop
monkeypatch.setattr(bridge, "shutdown_requested", fake_shutdown)
monkeypatch.setattr(bridge, "SqlBridgeStorage", lambda _db: FakeBridgeStorage())
class _FakeSession:
def __enter__(self) -> _FakeSession:
return self
def __exit__(self, *_exc: object) -> bool:
return False
class _FakeClient:
async def get_updates(self, **_kwargs: object) -> list[dict[str, Any]]:
return [{"update_id": 51}, {"update_id": 52}, {"update_id": 53}]
await bridge.run_poll_loop(_FakeClient(), _FakeSession, poll_timeout_s=1)
assert processed == [51] # 52/53 придут заново следующим getUpdates
async def test_topic_notification_failure_does_not_cancel_main_branch(
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
) -> None:
"""Уведомление в топик — вторичное действие: его сетевой отказ не отменяет
пометку is_blocked и не превращает апдейт в переигрываемый."""
calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({"copyMessage": 403}, calls)
storage = FakeBridgeStorage()
storage.record_message(
chat_id=555,
direction="in",
tg_message_id=1,
topic_message_id=100,
kind="text",
text_body="вопрос клиента",
operator_tg_id=None,
)
monkeypatch.setattr(client, "send_message", _network_boom("sendMessage"))
update = {"update_id": 42, "message": _group_reply_message(reply_to_message_id=100)}
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
advanced = await bridge.process_update(update, client, storage)
assert advanced is True
assert 555 in storage.blocked # основная ветка отработала
assert storage.get_offset() == 42
assert storage.commits == 1
assert len(storage.messages) == 1
assert "уведомление оператору в топик" in caplog.text
assert "chat_id=555" in caplog.text # только идентификатор, без текста переписки
async def test_topic_notification_uses_interactive_send_budget(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Уведомление шлётся с УЗКИМ бюджетом: воркерный дефолт (5 ретраев, backoff
до 30с, полный retry_after на 429) застопорил бы весь poll loop на минуты."""
calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({"copyMessage": 403}, calls)
storage = FakeBridgeStorage()
storage.record_message(
chat_id=555,
direction="in",
tg_message_id=1,
topic_message_id=100,
kind="text",
text_body="вопрос клиента",
operator_tg_id=None,
)
sent: list[dict[str, Any]] = []
async def recording_send(**kwargs: Any) -> dict[str, Any]:
sent.append(kwargs)
return {"message_id": 1}
monkeypatch.setattr(client, "send_message", recording_send)
update = {"update_id": 43, "message": _group_reply_message(reply_to_message_id=100)}
await bridge.process_update(update, client, storage)
assert len(sent) == 1
assert sent[0]["timeout"] == bridge._NOTIFY_SEND_TIMEOUT_S
assert sent[0]["max_retries"] == bridge._NOTIFY_SEND_MAX_RETRIES
assert sent[0]["max_backoff"] == bridge._NOTIFY_SEND_MAX_BACKOFF_S
assert sent[0]["chat_id"] == SUPPORT_CHAT_ID

View file

@ -17,9 +17,6 @@ import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from app.services.tgbot.client import (
_CONNECT_TIMEOUT_S,
_POOL_TIMEOUT_S,
_WRITE_TIMEOUT_S,
TelegramApiError,
TelegramClient,
TelegramNetworkError,
@ -28,19 +25,13 @@ from app.services.tgbot.client import (
_REAL_ASYNC_CLIENT = httpx.AsyncClient
def _install_transport(handler) -> list[httpx.AsyncClient]:
"""Подменяет транспорт. Возвращает список СОЗДАННЫХ AsyncClient — по нему
видно, переиспользуется ли один клиент или он плодится на каждый запрос."""
def _install_transport(handler) -> None:
transport = httpx.MockTransport(handler)
created: list[httpx.AsyncClient] = []
def factory(*_: object, **__: object) -> httpx.AsyncClient:
client = _REAL_ASYNC_CLIENT(transport=transport)
created.append(client)
return client
return _REAL_ASYNC_CLIENT(transport=transport)
mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory).start()
return created
@pytest.fixture(autouse=True)
@ -313,99 +304,3 @@ async def test_network_error_is_not_api_error() -> None:
await TelegramClient(token="t").send_message(chat_id=-1, text="x", max_retries=0)
assert not isinstance(caught.value, TelegramApiError)
# --- раздельные таймауты (#tg-connection-resilience) --------------------------
async def _capture_request_timeout(call) -> dict[str, float]:
captured: dict[str, Any] = {}
def handler(request: httpx.Request) -> httpx.Response:
captured["timeout"] = request.extensions.get("timeout")
return httpx.Response(200, json={"ok": True, "result": []})
_install_transport(handler)
await call(TelegramClient(token="t"))
assert captured["timeout"] is not None, "таймаут не доехал до запроса"
return captured["timeout"]
async def test_long_poll_timeout_applies_to_read_only_not_to_connect() -> None:
"""40 секунд запаса long-poll'а — это бюджет ОТВЕТА, а не установки соединения.
Скаляр в httpx разворачивается в connect=read=write=pool, поэтому
`getUpdates` ждал 40с и коннекта тоже. Живой connect до api.telegram.org из
прод-контейнера 0.036с; худший цикл из-за этого растягивался на ~174с
(4 попытки × 40с + backoff), и всё это время бот не видел ответов оператора.
"""
timeout = await _capture_request_timeout(lambda tg: tg.get_updates(offset=0, timeout=30))
assert timeout["read"] == 40.0, "long-poll обязан сохранить свои 30+10с на ответ"
assert timeout["connect"] == _CONNECT_TIMEOUT_S, "connect не должен наследовать long-poll"
assert timeout["write"] == _WRITE_TIMEOUT_S
assert timeout["pool"] == _POOL_TIMEOUT_S
async def test_interactive_call_keeps_its_own_narrow_read_budget() -> None:
"""Узкий интерактивный бюджет ручки — тоже read, и он не подменяется дефолтом."""
timeout = await _capture_request_timeout(
lambda tg: tg.send_message(chat_id=1, text="x", timeout=6.0)
)
assert timeout["read"] == 6.0
assert timeout["connect"] == _CONNECT_TIMEOUT_S
# --- переиспользование соединения --------------------------------------------
async def test_two_calls_share_one_httpx_client_and_aclose_closes_it() -> None:
"""Один `httpx.AsyncClient` на жизнь `TelegramClient`, а не на запрос.
Раньше клиент создавался ВНУТРИ цикла ретраев: нулевой keep-alive, полный
TCP+TLS-хендшейк на каждый запрос и на каждую попытку.
"""
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}})
created = _install_transport(handler)
tg = TelegramClient(token="t")
await tg.send_message(chat_id=1, text="раз")
await tg.send_message(chat_id=1, text="два")
assert len(created) == 1, f"клиент пересоздаётся на запрос: {len(created)} штук"
assert not created[0].is_closed
await tg.aclose()
assert created[0].is_closed, "aclose() обязан закрыть пул соединений"
async def test_retries_reuse_the_same_http_client() -> None:
"""Ретраи не пересоздают клиент — иначе повтор платит за хендшейк заново."""
calls = {"n": 0}
def handler(request: httpx.Request) -> httpx.Response:
calls["n"] += 1
if calls["n"] < 3:
return httpx.Response(502, json={"ok": False, "error_code": 502, "description": "gw"})
return httpx.Response(200, json={"ok": True, "result": {"message_id": 7}})
created = _install_transport(handler)
await TelegramClient(token="t").send_message(chat_id=1, text="x")
assert calls["n"] == 3, "бюджет ретраев изменился незаметно"
assert len(created) == 1, "на каждую попытку создаётся новый клиент"
async def test_async_context_manager_closes_client_on_exit() -> None:
def handler(request: httpx.Request) -> httpx.Response:
return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}})
created = _install_transport(handler)
async with TelegramClient(token="t") as tg:
await tg.send_message(chat_id=1, text="x")
assert len(created) == 1
assert created[0].is_closed

View file

@ -1,80 +0,0 @@
"""Жизненный цикл общего на приложение `TelegramClient` (`app.services.tgbot.shared`).
Смысл модуля ровно в том, что экземпляр ОДИН: до правки три HTTP-ручки создавали
клиент на каждый входящий запрос, и keep-alive не было ни на каком уровне. Если
синглтон однажды перестанет быть синглтоном, тесты ручек этого не заметят (они
подменяют аксессор целиком), а на проде вернутся TLS-хендшейки на каждое сообщение.
Поэтому проверяем идентичность и закрытие здесь, отдельно.
Реальных сетевых вызовов нет: `TelegramClient` создаёт `httpx.AsyncClient` лениво,
при первом запросе, а мы до запросов не доходим.
"""
from __future__ import annotations
import os
import pytest
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
from app.services.tgbot import shared
from app.services.tgbot.client import TelegramClient
@pytest.fixture(autouse=True)
def _reset_singleton():
"""Глобал модуля не должен течь между тестами — иначе порядок решает исход."""
shared._client = None
yield
shared._client = None
def test_get_returns_same_instance() -> None:
"""Два обращения — один объект. Это и есть весь смысл модуля."""
first = shared.get_telegram_client()
second = shared.get_telegram_client()
assert isinstance(first, TelegramClient)
assert first is second, "аксессор создал второй клиент — пул перестал переиспользоваться"
def test_init_returns_the_shared_instance() -> None:
"""`init_telegram_client` в lifespan и `get_telegram_client` в ручке — один и тот же объект."""
created = shared.init_telegram_client()
assert created is shared.get_telegram_client()
async def test_close_resets_and_next_get_builds_a_fresh_one() -> None:
"""После shutdown глобал обнулён; повторный startup обязан получить рабочий клиент.
Ленивость после закрытия намеренна: держать «остановленный» флаг значило бы,
что повторный `init_telegram_client()` в том же процессе отдаёт закрытый пул.
"""
first = shared.get_telegram_client()
await shared.close_telegram_client()
assert shared._client is None, "глобал не обнулён — следующий startup взял бы закрытый пул"
second = shared.get_telegram_client()
assert second is not first
async def test_close_is_idempotent() -> None:
"""Второй `close` не должен падать: lifespan зовёт его из `finally`."""
shared.get_telegram_client()
await shared.close_telegram_client()
await shared.close_telegram_client()
assert shared._client is None
async def test_close_without_client_is_a_noop() -> None:
"""Пустой токен — клиента в lifespan не создавали, закрывать нечего."""
assert shared._client is None
await shared.close_telegram_client()
assert shared._client is None

View file

@ -59,11 +59,7 @@ class _FakeTelegramClient:
def _fake_telegram_client(monkeypatch: pytest.MonkeyPatch) -> Any:
_FakeTelegramClient.calls = []
_FakeTelegramClient._response = {"message_id": 1}
# Ручка берёт ОБЩИЙ клиент приложения (#tg-connection-resilience), а не
# создаёт свой на запрос — подменяем аксессор, а не класс.
monkeypatch.setattr(
glitchtip_module, "get_telegram_client", lambda: _FakeTelegramClient("fake-token")
)
monkeypatch.setattr(glitchtip_module, "TelegramClient", _FakeTelegramClient)
return _FakeTelegramClient

View file

@ -69,11 +69,7 @@ class _FakeTelegramClient:
def _fake_telegram_client(monkeypatch: pytest.MonkeyPatch) -> Any:
_FakeTelegramClient.calls = []
_FakeTelegramClient._response = {"message_id": 555}
# Ручка берёт ОБЩИЙ клиент приложения (#tg-connection-resilience), а не
# создаёт свой на запрос — подменяем аксессор, а не класс.
monkeypatch.setattr(
support_module, "get_telegram_client", lambda: _FakeTelegramClient("fake-token")
)
monkeypatch.setattr(support_module, "TelegramClient", _FakeTelegramClient)
return _FakeTelegramClient