All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m19s
Замер прода за сутки 12.09.2026: 576 строк `network error` в логе `tradein-tgbot` и 7 полных исчерпаний бюджета ретраев, после которых падала итерация poll loop. Три причины, все подтверждены на коде и в рантайме. ## Ответ оператора мог пропасть навсегда `process_update` заканчивался безусловным `finally: save_offset(update_id)`. Замысел верный — «ядовитый» апдейт не должен блокировать поток, — но он не отличал неисправимый апдейт от транзиентного сетевого отказа. Оператор отвечает клиенту в топике, `copy_message` падает по сети, `TelegramNetworkError` улетает в общий `except Exception`, offset сдвигается. Telegram этот апдейт больше не отдаст, `record_message` не выполнился, оператор уверен, что ответил. Следа нет нигде, кроме строчки в логе. Теперь `process_update` возвращает `bool`. На `TelegramNetworkError` делается `rollback()`, offset НЕ сохраняется, возвращается `False`, и `run_poll_loop` прерывает разбор пачки — offset у Telegram единая «высшая отметка», подтверждение любого следующего апдейта неявно подтвердило бы и этот. Остаток пачки Telegram отдаст заново. Переигрывания ограничены сверху `_MAX_NETWORK_REPLAYS = 3`: без потолка «вечно недоставляемый» апдейт заклинил бы очередь навсегда, а это хуже потери одного сообщения. На потолке offset всё-таки двигается, но с `logger.error` и с `chat_id`/`message_id`, по которым человек найдёт ответ в топике и перешлёт руками. Текст переписки в лог по-прежнему не идёт. Дубли: `TelegramNetworkError` означает исчерпанный бюджет ретраев, при этом запрос мог дойти до Telegram, а ответ потеряться. Переигрывание тогда доставит сообщение второй раз. Это осознанный at-least-once компромисс — дубль видят и клиент, и оператор, а тихая потеря не видна никому. Полная идемпотентность по паре (update_id, target_chat_id) потребовала бы новой персистентной таблицы ради редкого случая; вместо неё число дублей жёстко ограничено сверху. Ветка `except TelegramApiError` с разбором `error_code == 403` («бот заблокирован») не тронута — там повтор действительно ничего не изменит. ## Таймаут задавался скаляром, поэтому connect ждал сорок секунд `httpx.AsyncClient(timeout=effective_timeout)` разворачивается в connect=read=write=pool. Для `getUpdates` бюджет ответа 40 секунд (30 держит Telegram плюс запас), и те же 40 секунд уходили на установку соединения — при живом connect в 0.036 секунды. Худший цикл: четыре попытки по 40 секунд плюс backoff, около трёх минут, в течение которых бот не видит ответов оператора. В логе это ровно те разрывы: 06:40:10, 06:42:22, 06:43:35. Теперь `httpx.Timeout(connect=5, read=<бюджет вызывающего>, write=10, pool=5)`, значения в именованных константах. Запас `+10s` у `get_updates` относится к read, докстринг поправлен. ## Клиент создавался заново на каждую попытку `httpx.AsyncClient` стоял ВНУТРИ цикла ретраев — keep-alive не было вовсе: полный TCP+TLS-хендшейк на каждый запрос и на каждый повтор, и заново кидался кубик «встанет ли коннект». Для long-polling это была основная статья сетевых отказов. Плюс три HTTP-ручки создавали `TelegramClient` на каждый входящий запрос. Теперь один ленивый переиспользуемый `AsyncClient` на экземпляр, с `aclose()` и `async with`. Общий клиент приложения живёт в новом `app/services/tgbot/shared.py`, создаётся и закрывается в lifespan; воркер бота держит свой на время поллинга. `keepalive_expiry` задан явно: дефолт httpx — 5 секунд, и с ним пул не давал бы ничего там, где нужнее всего. Poll loop переиспользует соединение и так, а вот веб-поддержка шлёт раз в минуты и за 5 секунд теряла бы его каждый раз. Плата за длинный keep-alive — шанс взять из пула закрытое той стороной соединение; httpx отдаёт это как `RemoteProtocolError`, который ретраится с #3457. ## Уведомления оператору шли с воркерным бюджетом внутри poll loop Обе отправки в топик («бот заблокирован», «веб-чат не поддерживает медиа») звались без своего бюджета, то есть с дефолтом в 5 ретраев и backoff до 30 секунд. Одна такая отправка стопорила весь цикл на минуты, а её отказ решал судьбу апдейта. Вынесены в `_notify_topic` с узким бюджетом и собственным `except`: провал вторичного действия больше не отменяет основную ветку. ## Тесты `tests/services/tgbot/test_shared.py` — новый, на жизненный цикл общего клиента. В `test_bridge.py` — сетевой отказ оставляет offset нетронутым и апдейт переигрывается, потолок разблокирует поток, отказ уведомления не отменяет основную ветку, прежнее поведение на 403 не изменилось. В `test_client.py` — раздельные таймауты доезжают до httpx per-request, два вызова используют один `AsyncClient`, `aclose()` его закрывает. Прогон по затронутым файлам: 127 passed. Ruff check и format чистые. Прокси намеренно не добавлялся: замер был на восьми запросах, это не статистика, и решение инфраструктурное. Если обрывы останутся — мерить сотней попыток отдельно.
198 lines
9.8 KiB
Python
198 lines
9.8 KiB
Python
"""Standalone entrypoint для Telegram support-bridge воркера (#tgsupport).
|
||
|
||
Зачем отдельный процесс/контейнер: `getUpdates` long-polling держит открытый
|
||
HTTP-запрос к Telegram до 30с за раз в бесконечном цикле — деплой основного API
|
||
(docker restart tradein-backend) не должен обрывать эту петлю на середине, как и
|
||
API не должен блокироваться долгим poll'ом. Тот же паттерн, что и
|
||
`scheduler_main.py` (#1182) для scraper'ов — отдельный контейнер с тем же образом,
|
||
другая команда.
|
||
|
||
Запуск: python -m app.tgbot_main
|
||
|
||
Kill-switch: TELEGRAM_BOT_TOKEN пуст (дефолт) → воркер логирует «disabled» и
|
||
блокируется на `wait_for_shutdown()` (idle, ~0 CPU) — НЕ `sys.exit(0)`. Сервис в
|
||
compose поднят с `restart: unless-stopped`, который рестартует контейнер
|
||
независимо от кода выхода — чистый exit(0) без токена дал бы бесконечный
|
||
рестарт-луп. Idle-блокировка держит процесс живым (автозапуск после ребута VPS
|
||
работает штатно через restart-policy) без CPU-луп и без спама рестартов;
|
||
SIGTERM просто убивает процесс — восстанавливать здесь нечего (bridge-задача
|
||
не запущена).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import os
|
||
import signal
|
||
from contextlib import suppress
|
||
from typing import Any
|
||
|
||
from app.core.config import settings
|
||
from app.core.db import SessionLocal
|
||
from app.core.shutdown import request_shutdown, shutdown_requested, wait_for_shutdown
|
||
from app.services.tgbot.bridge import run_poll_loop
|
||
from app.services.tgbot.client import TelegramClient
|
||
|
||
logging.basicConfig(
|
||
level=logging.INFO,
|
||
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||
)
|
||
|
||
# httpx INFO-логи печатают ПОЛНЫЙ request URL, включая Telegram Bot API токен
|
||
# в пути (https://api.telegram.org/bot<id>:<secret>/...) — `httpx: HTTP Request:
|
||
# POST https://api.telegram.org/bot<TOKEN>/getMe "HTTP/1.1 401 Unauthorized"`.
|
||
# При бесконечном long-polling'e это боевой токен в `docker logs` каждые ~30с.
|
||
# WARNING+ у httpx не логирует URL запроса (#tgsupport review, воспроизведено).
|
||
logging.getLogger("httpx").setLevel(logging.WARNING)
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Тот же safety-net паттерн, что scheduler_main.py — ниже docker stop_grace_period.
|
||
_DRAIN_TIMEOUT_S = 100.0
|
||
|
||
# Мониторинг ошибок — GlitchTip (Sentry-совместимый, #396). Только integrations
|
||
# без Starlette/FastAPI — здесь нет ASGI-приложения (тот же выбор что scheduler_main).
|
||
if settings.glitchtip_dsn:
|
||
import sentry_sdk
|
||
from sentry_sdk.integrations.httpx import HttpxIntegration
|
||
from sentry_sdk.integrations.logging import LoggingIntegration
|
||
|
||
from app.observability.sentry_scrub import (
|
||
redact_telegram_bot_token,
|
||
scrub_payment_request_body,
|
||
scrub_pii_event,
|
||
)
|
||
|
||
def _before_send(event: Any, hint: dict[str, Any]) -> Any:
|
||
"""Композиция платёжный body-wipe (PR-D2) + PII-scrub (form-данные) +
|
||
Telegram bot-токен redaction (#tgsupport review). Токен утекает ДВУМЯ
|
||
независимыми векторами, которые `include_local_variables=False` ниже и
|
||
этот хук закрывают вместе:
|
||
1. `include_local_variables=True` (sentry_sdk default) кладёт stack-frame
|
||
locals (`self._base`/`url` в `TelegramClient._request`) в traceback —
|
||
закрыто через `include_local_variables=False` в `sentry_sdk.init`.
|
||
2. `HttpxIntegration` кладёт полный request URL в span `data` (не только
|
||
traceback) — `traces_sample_rate=0.0` спасает СЕЙЧАС, но молча
|
||
перестанет спасать, если трейсинг когда-нибудь включат. Regex-редактор
|
||
— belt-and-suspenders на случай #1 (если include_local_variables
|
||
случайно вернут) И на span data.
|
||
|
||
Платёжный body-wipe — belt-and-suspenders: этот процесс не держит ASGI-
|
||
приложения (нет `request` в event сегодня), но тот же обработчик передан
|
||
ОБОИМ каналам ниже (before_send/before_send_transaction) ради единообразия
|
||
со всеми точками инициализации sentry_sdk в проекте (см. app/main.py).
|
||
"""
|
||
scrubbed = scrub_payment_request_body(event, hint)
|
||
if scrubbed is None:
|
||
return None
|
||
scrubbed = scrub_pii_event(scrubbed, hint)
|
||
if scrubbed is None:
|
||
return None
|
||
return redact_telegram_bot_token(scrubbed, hint)
|
||
|
||
sentry_sdk.init(
|
||
dsn=settings.glitchtip_dsn,
|
||
environment=settings.environment,
|
||
release=os.getenv("GIT_SHA") or os.getenv("SENTRY_RELEASE") or "unknown",
|
||
traces_sample_rate=0.0,
|
||
send_default_pii=False,
|
||
include_local_variables=False,
|
||
before_send=_before_send,
|
||
before_send_transaction=_before_send,
|
||
integrations=[
|
||
HttpxIntegration(),
|
||
LoggingIntegration(level=logging.INFO, event_level=logging.ERROR),
|
||
],
|
||
)
|
||
logger.info("GlitchTip monitoring enabled (tgbot_main)")
|
||
|
||
|
||
def _should_run() -> bool:
|
||
"""Kill-switch: TELEGRAM_BOT_TOKEN не задан → бот выключен (dev/staging без секрета)."""
|
||
return bool(settings.telegram_bot_token)
|
||
|
||
|
||
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)
|
||
|
||
|
||
async def _await_bridge(task: asyncio.Task[None]) -> None:
|
||
"""Кооперативный SIGTERM-drain — идентичная семантика scheduler_main._await_scheduler.
|
||
|
||
`run_poll_loop` сам проверяет `shutdown_requested()` между итерациями (между
|
||
getUpdates-вызовами) — long-polling запрос к Telegram (до 30с) докручивается,
|
||
затем цикл выходит сам. Safety-net здесь на случай зависшего HTTP-вызова.
|
||
"""
|
||
shutdown_waiter = asyncio.create_task(wait_for_shutdown())
|
||
try:
|
||
await asyncio.wait({task, shutdown_waiter}, return_when=asyncio.FIRST_COMPLETED)
|
||
finally:
|
||
shutdown_waiter.cancel()
|
||
with suppress(asyncio.CancelledError):
|
||
await shutdown_waiter
|
||
|
||
if task.done():
|
||
task.result()
|
||
logger.info("tgbot_main: bridge task exited cleanly")
|
||
return
|
||
|
||
logger.info(
|
||
"tgbot_main: SIGTERM-drain — waiting up to %.0fs for current poll iteration to finish",
|
||
_DRAIN_TIMEOUT_S,
|
||
)
|
||
try:
|
||
await asyncio.wait_for(task, timeout=_DRAIN_TIMEOUT_S)
|
||
logger.info("tgbot_main: bridge drained and exited cleanly")
|
||
except TimeoutError:
|
||
logger.warning(
|
||
"tgbot_main: drain exceeded %.0fs grace — hard-cancelling bridge task",
|
||
_DRAIN_TIMEOUT_S,
|
||
)
|
||
task.cancel()
|
||
with suppress(asyncio.CancelledError):
|
||
await task
|
||
|
||
|
||
async def _run() -> None:
|
||
task = asyncio.create_task(_run_bridge())
|
||
|
||
loop = asyncio.get_running_loop()
|
||
|
||
def _on_signal(signum: int) -> None:
|
||
logger.info("tgbot_main: signal %d received — requesting cooperative drain", signum)
|
||
request_shutdown()
|
||
|
||
try:
|
||
loop.add_signal_handler(signal.SIGTERM, lambda: _on_signal(signal.SIGTERM))
|
||
loop.add_signal_handler(signal.SIGINT, lambda: _on_signal(signal.SIGINT))
|
||
except NotImplementedError:
|
||
# Windows dev: signal handlers через loop не поддерживаются
|
||
logger.warning("tgbot_main: loop.add_signal_handler not supported (Windows dev)")
|
||
|
||
await _await_bridge(task)
|
||
|
||
if shutdown_requested():
|
||
logger.info("tgbot_main: bridge drained cleanly (SIGTERM)")
|
||
else:
|
||
logger.info("tgbot_main: bridge task exited")
|
||
|
||
|
||
if __name__ == "__main__":
|
||
if not _should_run():
|
||
# NOT sys.exit(0): compose service has `restart: unless-stopped`, который
|
||
# рестартует контейнер независимо от кода выхода — чистый exit(0) без
|
||
# токена дал бы бесконечный рестарт-луп. Idle-блокировка вместо этого:
|
||
# ~0 CPU, SIGTERM просто убивает процесс (нечего дренировать).
|
||
logger.warning(
|
||
"tgbot_main: TELEGRAM_BOT_TOKEN не задан — бот выключен, "
|
||
"блокируемся на idle (не exit, чтобы не было рестарт-лупа с restart:unless-stopped)"
|
||
)
|
||
asyncio.run(wait_for_shutdown())
|
||
else:
|
||
asyncio.run(_run())
|