Связь с Telegram не встаёт колом, ответ оператора не теряется #3458
12 changed files with 776 additions and 46 deletions
|
|
@ -53,7 +53,8 @@ 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 TelegramClient, TelegramError
|
||||
from app.services.tgbot.client import TelegramError
|
||||
from app.services.tgbot.shared import get_telegram_client
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
|
@ -213,7 +214,10 @@ async def glitchtip_webhook(
|
|||
received_at = datetime.now(UTC)
|
||||
text = _build_message(raw_body, received_at)
|
||||
|
||||
client = TelegramClient(settings.telegram_bot_token)
|
||||
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
|
||||
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
|
||||
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
|
||||
client = get_telegram_client()
|
||||
try:
|
||||
await client.send_message(
|
||||
chat_id=settings.telegram_alerts_chat_id,
|
||||
|
|
|
|||
|
|
@ -73,7 +73,8 @@ 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 TelegramClient, TelegramError
|
||||
from app.services.tgbot.client import TelegramError
|
||||
from app.services.tgbot.shared import get_telegram_client
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
|
@ -219,7 +220,10 @@ async def send_support_message(
|
|||
headers={"Retry-After": str(int(retry_after) + 1)},
|
||||
)
|
||||
|
||||
client = TelegramClient(settings.telegram_bot_token)
|
||||
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
|
||||
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
|
||||
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
|
||||
client = get_telegram_client()
|
||||
try:
|
||||
mirrored = await client.send_message(
|
||||
chat_id=settings.telegram_support_chat_id,
|
||||
|
|
@ -411,7 +415,10 @@ async def send_anon_support_message(
|
|||
)
|
||||
|
||||
display_id = _anon_display_id(token)
|
||||
client = TelegramClient(settings.telegram_bot_token)
|
||||
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
|
||||
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
|
||||
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
|
||||
client = get_telegram_client()
|
||||
try:
|
||||
mirrored = await client.send_message(
|
||||
chat_id=settings.telegram_support_chat_id,
|
||||
|
|
|
|||
|
|
@ -50,6 +50,7 @@ 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__)
|
||||
|
||||
|
|
@ -216,7 +217,19 @@ 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.
|
||||
yield
|
||||
|
||||
# Общий на приложение 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()
|
||||
|
||||
|
||||
app = FastAPI(
|
||||
|
|
|
|||
|
|
@ -29,9 +29,13 @@
|
|||
Telegram 403 (клиент заблокировал бота) → is_blocked=true + уведомление в
|
||||
топике (только для Telegram-ветки — у веб-клиента нет "заблокировал бота").
|
||||
C) Дедуп: update_id <= сохранённого offset — skip. Offset сохраняется И
|
||||
коммитится в той же транзакции, что и запись сообщения (см. `process_update`
|
||||
`finally`), после КАЖДОГО апдейта — рестарт воркера не переигрывает уже
|
||||
обработанные апдейты и не подвисает вечно на «ядовитом» апдейте.
|
||||
коммитится в той же транзакции, что и запись сообщения (см. `process_update`),
|
||||
после КАЖДОГО апдейта — рестарт воркера не переигрывает уже обработанные
|
||||
апдейты и не подвисает вечно на «ядовитом» апдейте. Исключение —
|
||||
ТРАНЗИЕНТНЫЙ сетевой отказ (`TelegramNetworkError`, т.е. исчерпанный бюджет
|
||||
ретраев клиента): такой апдейт СОЗНАТЕЛЬНО остаётся неподтверждённым, чтобы
|
||||
Telegram отдал его снова, иначе ответ оператора пропадал бы навсегда
|
||||
(#tg-connection-resilience). Потолок переигрываний — `_MAX_NETWORK_REPLAYS`.
|
||||
D) TELEGRAM_BOT_TOKEN пуст → бот выключен — проверяется в `app.tgbot_main`
|
||||
(entrypoint), не здесь.
|
||||
E) /start клиенту → короткое приветствие МЕРЫ, без зеркалирования в топик
|
||||
|
|
@ -76,7 +80,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
|
||||
from app.services.tgbot.client import TelegramApiError, TelegramClient, TelegramNetworkError
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
|
@ -131,6 +135,44 @@ 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):
|
||||
|
|
@ -400,6 +442,43 @@ 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
|
||||
|
|
@ -553,14 +632,14 @@ async def _handle_group_reply(
|
|||
if exc.error_code == 403:
|
||||
# Клиент заблокировал бота — фиксируем и уведомляем оператора в топике.
|
||||
storage.mark_blocked(target_chat_id)
|
||||
await client.send_message(
|
||||
chat_id=settings.telegram_support_chat_id,
|
||||
await _notify_topic(
|
||||
client,
|
||||
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
|
||||
|
|
@ -592,11 +671,11 @@ async def _handle_group_reply(
|
|||
web_thread_id,
|
||||
kind,
|
||||
)
|
||||
await client.send_message(
|
||||
chat_id=settings.telegram_support_chat_id,
|
||||
await _notify_topic(
|
||||
client,
|
||||
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
|
||||
|
||||
|
|
@ -636,25 +715,41 @@ async def _handle_group_reply(
|
|||
|
||||
async def process_update(
|
||||
update: dict[str, Any], client: TelegramClient, storage: BridgeStorage
|
||||
) -> None:
|
||||
) -> bool:
|
||||
"""Маршрутизирует один Telegram update. Дедуп (C) + атомарный offset-commit.
|
||||
|
||||
Дедуп: update_id <= сохранённого offset — skip без side-effects. Offset
|
||||
сохраняется и коммитится ПОСЛЕ обработки (в т.ч. если обработка упала —
|
||||
иначе «ядовитый» апдейт блокировал бы весь поток навсегда).
|
||||
Возвращает True, если offset сдвинут (апдейт подтверждён, Telegram его больше
|
||||
не отдаст), и False, если апдейт СОЗНАТЕЛЬНО оставлен неподтверждённым ради
|
||||
переигрывания. На False вызывающий (`run_poll_loop`) ОБЯЗАН прервать разбор
|
||||
пачки: offset у Telegram — единая «высшая отметка», подтверждение любого
|
||||
СЛЕДУЮЩЕГО апдейта неявно подтвердило бы и этот, и переигрывания не было бы.
|
||||
|
||||
Различаем сбой БД (`SQLAlchemyError`) от прочих (Telegram API и т.п.):
|
||||
сбой БД оставляет сессию в failed-transaction state — `rollback()` ОБЯЗАН
|
||||
отработать ПЕРЕД `save_offset`, иначе тот сам кинет `PendingRollbackError`,
|
||||
`process_update` вылетит без сохранения offset'а, следующая итерация
|
||||
`run_poll_loop` получит СТАРЫЙ offset от `get_offset()` и переиграет тот же
|
||||
апдейт заново — copyMessage задублирует зеркало клиента в топике на
|
||||
каждый повтор поллинга (#3 review, воспроизведено).
|
||||
Дедуп: 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 двигаем, «ядовитый» апдейт
|
||||
не блокирует поток.
|
||||
"""
|
||||
update_id = update.get("update_id")
|
||||
if not isinstance(update_id, int):
|
||||
logger.warning("tgbot bridge: update без валидного update_id — игнор")
|
||||
return
|
||||
return True
|
||||
|
||||
current_offset = storage.get_offset()
|
||||
if update_id <= current_offset:
|
||||
|
|
@ -663,7 +758,7 @@ async def process_update(
|
|||
update_id,
|
||||
current_offset,
|
||||
)
|
||||
return
|
||||
return True
|
||||
|
||||
message = update.get("message")
|
||||
try:
|
||||
|
|
@ -677,6 +772,35 @@ 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 перед сохранением "
|
||||
|
|
@ -690,9 +814,12 @@ async def process_update(
|
|||
"(не блокируем поток на 'ядовитом' апдейте)",
|
||||
update_id,
|
||||
)
|
||||
finally:
|
||||
storage.save_offset(update_id)
|
||||
storage.commit()
|
||||
|
||||
# Апдейт подтверждён — счётчик переигрываний больше не нужен (словарь не растёт).
|
||||
_network_replay_attempts.pop(update_id, None)
|
||||
storage.save_offset(update_id)
|
||||
storage.commit()
|
||||
return True
|
||||
|
||||
|
||||
# ── Long-polling loop ─────────────────────────────────────────────────────────
|
||||
|
|
@ -717,8 +844,20 @@ async def run_poll_loop(
|
|||
updates = await client.get_updates(
|
||||
offset=offset + 1, timeout=poll_timeout_s, allowed_updates=["message"]
|
||||
)
|
||||
for update in updates:
|
||||
await process_update(update, client, storage)
|
||||
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
|
||||
consecutive_errors = 0
|
||||
except Exception:
|
||||
consecutive_errors += 1
|
||||
|
|
|
|||
|
|
@ -54,6 +54,38 @@ _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», и «не ответил вовсе».
|
||||
|
|
@ -94,6 +126,21 @@ 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:
|
||||
|
|
@ -127,7 +174,24 @@ def _error_from_body(response: httpx.Response) -> tuple[int, str]:
|
|||
|
||||
|
||||
class TelegramClient:
|
||||
"""Bot API клиент на httpx.AsyncClient. Каждый вызов — отдельное короткоживущее соединение."""
|
||||
"""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,
|
||||
|
|
@ -137,6 +201,32 @@ 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,
|
||||
|
|
@ -164,14 +254,19 @@ 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:
|
||||
async with httpx.AsyncClient(timeout=effective_timeout) as client:
|
||||
response = await client.post(url, json=payload)
|
||||
response = await client.post(url, json=payload, timeout=request_timeout)
|
||||
except httpx.TransportError as exc:
|
||||
# Ловим ВЕСЬ `TransportError`, а не узкий кортеж
|
||||
# `(TimeoutException, NetworkError)`: `RemoteProtocolError`
|
||||
|
|
@ -306,8 +401,11 @@ class TelegramClient:
|
|||
) -> list[dict[str, Any]]:
|
||||
"""Long-polling getUpdates. `timeout` — сколько Telegram держит запрос открытым (сек).
|
||||
|
||||
HTTP-таймаут запроса берётся с запасом (`timeout + 10s`), чтобы не обрывать
|
||||
соединение раньше, чем ответит сам Telegram long-poll.
|
||||
Запас `+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:
|
||||
|
|
|
|||
50
tradein-mvp/backend/app/services/tgbot/shared.py
Normal file
50
tradein-mvp/backend/app/services/tgbot/shared.py
Normal file
|
|
@ -0,0 +1,50 @@
|
|||
"""Общий на приложение `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: пул соединений закрыт")
|
||||
|
|
@ -114,8 +114,12 @@ def _should_run() -> bool:
|
|||
|
||||
|
||||
async def _run_bridge() -> None:
|
||||
client = TelegramClient(settings.telegram_bot_token)
|
||||
await run_poll_loop(client, SessionLocal)
|
||||
# `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)
|
||||
|
||||
|
||||
async def _await_bridge(task: asyncio.Task[None]) -> None:
|
||||
|
|
|
|||
|
|
@ -22,6 +22,12 @@ 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 не задан → клиенту уходит "сервис недоступен"
|
||||
|
|
@ -45,7 +51,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
|
||||
from app.services.tgbot.client import TelegramClient, TelegramNetworkError
|
||||
|
||||
SUPPORT_CHAT_ID = -100123456789
|
||||
SUPPORT_TOPIC_ID = 42
|
||||
|
|
@ -973,3 +979,219 @@ 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
|
||||
|
|
|
|||
|
|
@ -17,6 +17,9 @@ 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,
|
||||
|
|
@ -25,13 +28,19 @@ from app.services.tgbot.client import (
|
|||
_REAL_ASYNC_CLIENT = httpx.AsyncClient
|
||||
|
||||
|
||||
def _install_transport(handler) -> None:
|
||||
def _install_transport(handler) -> list[httpx.AsyncClient]:
|
||||
"""Подменяет транспорт. Возвращает список СОЗДАННЫХ AsyncClient — по нему
|
||||
видно, переиспользуется ли один клиент или он плодится на каждый запрос."""
|
||||
transport = httpx.MockTransport(handler)
|
||||
created: list[httpx.AsyncClient] = []
|
||||
|
||||
def factory(*_: object, **__: object) -> httpx.AsyncClient:
|
||||
return _REAL_ASYNC_CLIENT(transport=transport)
|
||||
client = _REAL_ASYNC_CLIENT(transport=transport)
|
||||
created.append(client)
|
||||
return client
|
||||
|
||||
mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory).start()
|
||||
return created
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
|
|
@ -304,3 +313,99 @@ 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
|
||||
|
|
|
|||
80
tradein-mvp/backend/tests/services/tgbot/test_shared.py
Normal file
80
tradein-mvp/backend/tests/services/tgbot/test_shared.py
Normal file
|
|
@ -0,0 +1,80 @@
|
|||
"""Жизненный цикл общего на приложение `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
|
||||
|
|
@ -59,7 +59,11 @@ class _FakeTelegramClient:
|
|||
def _fake_telegram_client(monkeypatch: pytest.MonkeyPatch) -> Any:
|
||||
_FakeTelegramClient.calls = []
|
||||
_FakeTelegramClient._response = {"message_id": 1}
|
||||
monkeypatch.setattr(glitchtip_module, "TelegramClient", _FakeTelegramClient)
|
||||
# Ручка берёт ОБЩИЙ клиент приложения (#tg-connection-resilience), а не
|
||||
# создаёт свой на запрос — подменяем аксессор, а не класс.
|
||||
monkeypatch.setattr(
|
||||
glitchtip_module, "get_telegram_client", lambda: _FakeTelegramClient("fake-token")
|
||||
)
|
||||
return _FakeTelegramClient
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -69,7 +69,11 @@ class _FakeTelegramClient:
|
|||
def _fake_telegram_client(monkeypatch: pytest.MonkeyPatch) -> Any:
|
||||
_FakeTelegramClient.calls = []
|
||||
_FakeTelegramClient._response = {"message_id": 555}
|
||||
monkeypatch.setattr(support_module, "TelegramClient", _FakeTelegramClient)
|
||||
# Ручка берёт ОБЩИЙ клиент приложения (#tg-connection-resilience), а не
|
||||
# создаёт свой на запрос — подменяем аксессор, а не класс.
|
||||
monkeypatch.setattr(
|
||||
support_module, "get_telegram_client", lambda: _FakeTelegramClient("fake-token")
|
||||
)
|
||||
return _FakeTelegramClient
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue