feat(tradein/tgbot): Telegram support-мост @MERAsupport_bot (#2526)
All checks were successful
Deploy Trade-In / changes (push) Successful in 13s
Deploy Trade-In / build-frontend (push) Successful in 44s
Deploy Trade-In / build-browser (push) Successful in 3m18s
Deploy Trade-In / test (push) Successful in 5m2s
Deploy Trade-In / build-backend (push) Successful in 2m16s
Deploy Trade-In / deploy (push) Successful in 1m59s
All checks were successful
Deploy Trade-In / changes (push) Successful in 13s
Deploy Trade-In / build-frontend (push) Successful in 44s
Deploy Trade-In / build-browser (push) Successful in 3m18s
Deploy Trade-In / test (push) Successful in 5m2s
Deploy Trade-In / build-backend (push) Successful in 2m16s
Deploy Trade-In / deploy (push) Successful in 1m59s
This commit is contained in:
commit
fcfc777baa
15 changed files with 2233 additions and 7 deletions
|
|
@ -483,7 +483,11 @@ jobs:
|
||||||
# рискует ложно отменить НЕ относящийся к этому recreate run (напр.
|
# рискует ложно отменить НЕ относящийся к этому recreate run (напр.
|
||||||
# admin-triggered scrape внутри backend, если backend в этом деплое
|
# admin-triggered scrape внутри backend, если backend в этом деплое
|
||||||
# не пересоздавался — его heartbeat продолжит расти после checkpoint'а).
|
# не пересоздавался — его heartbeat продолжит расти после checkpoint'а).
|
||||||
SERVICES="browser backend frontend"
|
# tgbot: тот же backend-образ (rebuild уже покрыт filters.backend —
|
||||||
|
# tradein-mvp/backend/** включает app/tgbot_main.py), никакого
|
||||||
|
# in-flight state вроде scrape_runs → пересоздаётся безусловно вместе
|
||||||
|
# с browser/backend/frontend, отдельного graceful-drain не требует.
|
||||||
|
SERVICES="browser backend frontend tgbot"
|
||||||
SCRAPER_STOP_TS=""
|
SCRAPER_STOP_TS=""
|
||||||
if [ "${SCRAPER_CHANGED:-true}" = "true" ]; then
|
if [ "${SCRAPER_CHANGED:-true}" = "true" ]; then
|
||||||
echo "→ scraper paths changed — waiting for in-flight scrape_runs to drain (up to 5 min)"
|
echo "→ scraper paths changed — waiting for in-flight scrape_runs to drain (up to 5 min)"
|
||||||
|
|
|
||||||
|
|
@ -40,3 +40,14 @@ DADATA_API_SECRET=
|
||||||
POSTGRES_USER=tradein
|
POSTGRES_USER=tradein
|
||||||
POSTGRES_PASSWORD=tradein
|
POSTGRES_PASSWORD=tradein
|
||||||
POSTGRES_DB=tradein
|
POSTGRES_DB=tradein
|
||||||
|
|
||||||
|
# === Telegram support-bot bridge (tgbot service, docker-compose.prod.yml) ===
|
||||||
|
# Long-polling worker: пересылает support-обращения в Telegram-топик. Не FastAPI,
|
||||||
|
# отдельный процесс (app/tgbot_main.py), env читается из backend/.env.runtime на VPS.
|
||||||
|
#
|
||||||
|
# BotFather token. Пусто = бот не стартует (выключен).
|
||||||
|
TELEGRAM_BOT_TOKEN=
|
||||||
|
# ID супергруппы-форума с включёнными топиками (вида -100XXXXXXXXXX).
|
||||||
|
TELEGRAM_SUPPORT_CHAT_ID=
|
||||||
|
# ID топика (thread) внутри супергруппы, куда падают support-сообщения.
|
||||||
|
TELEGRAM_SUPPORT_TOPIC_ID=
|
||||||
|
|
|
||||||
|
|
@ -89,10 +89,47 @@ GLITCHTIP_DSN=<dsn или пусто>
|
||||||
# Регистрация ключей: https://dadata.ru/api/clean/
|
# Регистрация ключей: https://dadata.ru/api/clean/
|
||||||
DADATA_API_TOKEN=<token или пусто>
|
DADATA_API_TOKEN=<token или пусто>
|
||||||
DADATA_API_SECRET=<secret или пусто>
|
DADATA_API_SECRET=<secret или пусто>
|
||||||
|
|
||||||
|
# Telegram support-bot bridge (сервис tgbot, docker-compose.prod.yml).
|
||||||
|
# Long-polling воркер (app/tgbot_main.py), тот же образ что backend/scraper,
|
||||||
|
# отдельный контейнер tradein-tgbot. Пусто TELEGRAM_BOT_TOKEN = бот НЕ падает
|
||||||
|
# и НЕ рестарт-лупится — процесс стартует, уходит в idle-блокировку и просто
|
||||||
|
# висит (это норма для окружений без токена, не сбой; см. tgbot_main.py).
|
||||||
|
#
|
||||||
|
# 1. TELEGRAM_BOT_TOKEN — токен от @BotFather (/newbot). Пусто = бот выключен
|
||||||
|
# (idle, не polling).
|
||||||
|
TELEGRAM_BOT_TOKEN=<token или пусто>
|
||||||
|
# 2. TELEGRAM_SUPPORT_CHAT_ID — id супергруппы-форума (Topics включены в
|
||||||
|
# настройках группы), вида -100XXXXXXXXXX. Получить: добавить бота в группу,
|
||||||
|
# отправить любое сообщение в любой топик, дернуть
|
||||||
|
# https://api.telegram.org/bot<token>/getUpdates — в ответе
|
||||||
|
# message.chat.id (для супергруппы всегда отрицательный, начинается с -100).
|
||||||
|
TELEGRAM_SUPPORT_CHAT_ID=<-100... или пусто>
|
||||||
|
# 3. TELEGRAM_SUPPORT_TOPIC_ID — id конкретного топика (thread) внутри группы,
|
||||||
|
# куда падают support-обращения. Открыть нужный топик в Telegram Desktop/Web →
|
||||||
|
# в URL топика (t.me/c/<chat>/<topic_id>) последнее число — это topic_id.
|
||||||
|
# Либо взять message_thread_id из того же getUpdates-ответа (п.2), отправив
|
||||||
|
# тестовое сообщение именно в целевой топик.
|
||||||
|
TELEGRAM_SUPPORT_TOPIC_ID=<topic_id или пусто>
|
||||||
```
|
```
|
||||||
|
|
||||||
Оба файла создаются вручную при первом деплое.
|
Оба файла создаются вручную при первом деплое.
|
||||||
|
|
||||||
|
### Деплой / рестарт `tgbot`
|
||||||
|
|
||||||
|
```bash
|
||||||
|
# .env.runtime читается на старте container — `compose restart` НЕ перечитывает.
|
||||||
|
docker compose -p gendesign-tradein -f docker-compose.prod.yml \
|
||||||
|
up -d --force-recreate --no-deps tgbot
|
||||||
|
```
|
||||||
|
|
||||||
|
`restart: unless-stopped` + `stop_grace_period: 120s` в compose (см. `docker-compose.prod.yml`)
|
||||||
|
— автозапуск после ребута VPS гарантирован (`on-failure` сюда не годится: код выхода
|
||||||
|
контейнера при ребуте — гонка с long-poll таймаутом 30с, `unless-stopped`/`always`
|
||||||
|
не зависят от exit-кода). 120s grace даёт time докрутить long-poll + отработать
|
||||||
|
кооперативный drain (`_DRAIN_TIMEOUT_S=100s` в `tgbot_main.py`) до docker SIGKILL —
|
||||||
|
паттерн скопирован с `scraper` (см. комментарий там же).
|
||||||
|
|
||||||
### После изменения `backend/.env.runtime`
|
### После изменения `backend/.env.runtime`
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
|
|
|
||||||
|
|
@ -664,5 +664,19 @@ class Settings(BaseSettings):
|
||||||
# честная маркировка.
|
# честная маркировка.
|
||||||
sell_time_sensitivity_min_n_lots: int = 10
|
sell_time_sensitivity_min_n_lots: int = 10
|
||||||
|
|
||||||
|
# ── Telegram support bridge (@MERAsupport_bot) ───────────────────────────
|
||||||
|
# Клиент пишет боту в личку → зеркалится в топик support-группы → оператор
|
||||||
|
# отвечает реплаем в топике → бот доставляет ответ клиенту. Standalone
|
||||||
|
# long-polling воркер (app.tgbot_main), НЕ webhook — см. app/services/tgbot/.
|
||||||
|
# Пусто/0 = бот выключен: tgbot_main логирует «disabled» и выходит с кодом 0
|
||||||
|
# (чтобы контейнер без секрета не крутил рестарт-луп). ENV: TELEGRAM_BOT_TOKEN,
|
||||||
|
# TELEGRAM_SUPPORT_CHAT_ID, TELEGRAM_SUPPORT_TOPIC_ID.
|
||||||
|
telegram_bot_token: str = Field(default="", validation_alias="TELEGRAM_BOT_TOKEN")
|
||||||
|
# Telegram id форум-группы (супергруппы с включёнными топиками), куда
|
||||||
|
# зеркалятся обращения клиентов. Отрицательный для supergroup id (напр. -100...).
|
||||||
|
telegram_support_chat_id: int = Field(default=0, validation_alias="TELEGRAM_SUPPORT_CHAT_ID")
|
||||||
|
# message_thread_id топика внутри support-группы, в который идут зеркала.
|
||||||
|
telegram_support_topic_id: int = Field(default=0, validation_alias="TELEGRAM_SUPPORT_TOPIC_ID")
|
||||||
|
|
||||||
|
|
||||||
settings = Settings()
|
settings = Settings()
|
||||||
|
|
|
||||||
|
|
@ -1,22 +1,51 @@
|
||||||
"""Хук before_send для GlitchTip/Sentry SDK (tradein-local, #396).
|
"""Хуки before_send для GlitchTip/Sentry SDK (tradein-local, #396, #tgsupport).
|
||||||
|
|
||||||
Redact-ит consumer-PII (client_name / client_phone / client_email и пр.)
|
Redact-ит consumer-PII (client_name / client_phone / client_email и пр.)
|
||||||
из error events до отправки в GlitchTip — estimator/trade-in flow таскает
|
из error events до отправки в GlitchTip — estimator/trade-in flow таскает
|
||||||
эти поля, а send_default_pii=False их не покрывает (это user-data в
|
эти поля, а send_default_pii=False их не покрывает (это user-data в
|
||||||
request.data / extra / contexts, не PII-заголовки).
|
request.data / extra / contexts, не PII-заголовки).
|
||||||
|
|
||||||
|
`redact_telegram_bot_token` — отдельный хук (#tgsupport review): Telegram Bot
|
||||||
|
API токен живёт в URL-пути (`https://api.telegram.org/bot<id>:<secret>/...`),
|
||||||
|
а не в query/userinfo, поэтому НЕ покрывается sentry_sdk `sanitize_url` (тот
|
||||||
|
режет только `user:pass@` и query-параметры). Токен утекает ДВУМЯ путями,
|
||||||
|
которые `_scrub`/`scrub_pii_event` (ключ-based, PII-словарь) не ловят:
|
||||||
|
1. `include_local_variables=True` (sentry_sdk default) кладёт locals
|
||||||
|
stack-фрейма (`self._base`, `url` в `TelegramClient._request`) в
|
||||||
|
traceback → полный токен открытым текстом.
|
||||||
|
2. `HttpxIntegration` кладёт полный request URL в span `data` (виден при
|
||||||
|
любом ненулевом `traces_sample_rate`), а не только в traceback.
|
||||||
|
Поэтому редактор — НЕ ключ-based, а regex full-text по КАЖДОЙ строке во всём
|
||||||
|
event (глубокий обход dict/list/tuple) — токен может всплыть в любом поле.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import re
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from sentry_sdk.types import Event
|
from sentry_sdk.types import Event
|
||||||
|
|
||||||
_REDACTED = "[REDACTED]"
|
_REDACTED = "[REDACTED]"
|
||||||
# Ключи consumer-PII (нижний регистр; сверка case-insensitive).
|
# Ключи consumer-PII (нижний регистр; сверка case-insensitive).
|
||||||
_PII_KEYS = frozenset(
|
_PII_KEYS = frozenset({"client_name", "client_phone", "client_email", "phone", "email", "name"})
|
||||||
{"client_name", "client_phone", "client_email", "phone", "email", "name"}
|
|
||||||
)
|
# Telegram Bot API токен в пути URL: /bot<numeric_id>:<secret-part>/<method>.
|
||||||
|
# Матчим ровно этот сегмент (не весь URL) — сохраняет остальной путь/query
|
||||||
|
# читаемым для диагностики (метод API, error code и т.п.).
|
||||||
|
_TG_BOT_TOKEN_RE = re.compile(r"/bot\d+:[A-Za-z0-9_-]+")
|
||||||
|
_TG_BOT_TOKEN_REPLACEMENT = "/bot[REDACTED]"
|
||||||
|
|
||||||
|
# Тот же токен БЕЗ префикса `/bot` — форма `<numeric_id>:<secret>` сама по себе
|
||||||
|
# (напр. локаль `token` в конструкторе TelegramClient, или если его кто-то
|
||||||
|
# засунет в log-сообщение). Сейчас единственный путь такой формы в event —
|
||||||
|
# locals стек-фрейма, а они выключены через include_local_variables=False в
|
||||||
|
# tgbot_main. Но именно на отказ того флага этот редактор и страхует: без этой
|
||||||
|
# ветки рубеж был бы один, а не два. Формат токена BotFather: 8-12 цифр `:` 35
|
||||||
|
# символов base64url — нижние границы взяты с запасом, чтобы не промахнуться
|
||||||
|
# на нестандартных id, но остаться уже, чем `\d+:\S+` (тот бил бы по любым
|
||||||
|
# `id:value` в логах, напр. `chat_id:12345`).
|
||||||
|
_TG_BOT_TOKEN_BARE_RE = re.compile(r"\b\d{6,12}:[A-Za-z0-9_-]{30,}\b")
|
||||||
|
|
||||||
|
|
||||||
def _scrub(obj: Any) -> None:
|
def _scrub(obj: Any) -> None:
|
||||||
|
|
@ -42,3 +71,34 @@ def scrub_pii_event(event: Event, _hint: dict[str, Any]) -> Event | None:
|
||||||
_scrub(event.get("extra"))
|
_scrub(event.get("extra"))
|
||||||
_scrub(event.get("contexts"))
|
_scrub(event.get("contexts"))
|
||||||
return event
|
return event
|
||||||
|
|
||||||
|
|
||||||
|
def _redact_strings(obj: Any) -> Any:
|
||||||
|
"""Рекурсивно проходит dict/list/tuple и прогоняет обе токен-регулярки по КАЖДОЙ
|
||||||
|
строке (не только по конкретным ключам) — токен может оказаться в locals
|
||||||
|
stack-фрейма, span data, breadcrumb message, request.url и т.д. Возвращает
|
||||||
|
НОВУЮ структуру (не мутирует `obj` — в отличие от `_scrub`, чтобы не зависеть
|
||||||
|
от того, какие контейнеры sentry_sdk считает mutable в своём event dict)."""
|
||||||
|
if isinstance(obj, str):
|
||||||
|
redacted = _TG_BOT_TOKEN_RE.sub(_TG_BOT_TOKEN_REPLACEMENT, obj)
|
||||||
|
return _TG_BOT_TOKEN_BARE_RE.sub(_REDACTED, redacted)
|
||||||
|
if isinstance(obj, dict):
|
||||||
|
return {k: _redact_strings(v) for k, v in obj.items()}
|
||||||
|
if isinstance(obj, list):
|
||||||
|
return [_redact_strings(v) for v in obj]
|
||||||
|
if isinstance(obj, tuple):
|
||||||
|
return tuple(_redact_strings(v) for v in obj)
|
||||||
|
return obj
|
||||||
|
|
||||||
|
|
||||||
|
def redact_telegram_bot_token(event: Event, _hint: dict[str, Any]) -> Event | None:
|
||||||
|
"""Full-text regex redaction Telegram Bot API токена по ВСЕМУ event (#tgsupport).
|
||||||
|
|
||||||
|
Ловит оба вектора утечки токена в GlitchTip, которые ключ-based `scrub_pii_event`
|
||||||
|
не покрывает: locals stack-фреймов (`include_local_variables=True`) и httpx-span
|
||||||
|
`data` (полный request URL). Композировать с `scrub_pii_event`, не вместо него —
|
||||||
|
разные классы секретов (PII полей формы vs bot-токен в URL).
|
||||||
|
"""
|
||||||
|
if not isinstance(event, dict):
|
||||||
|
return event
|
||||||
|
return _redact_strings(event) # type: ignore[return-value]
|
||||||
|
|
|
||||||
9
tradein-mvp/backend/app/services/tgbot/__init__.py
Normal file
9
tradein-mvp/backend/app/services/tgbot/__init__.py
Normal file
|
|
@ -0,0 +1,9 @@
|
||||||
|
"""Telegram support bridge (@MERAsupport_bot) — long-polling мост клиент↔оператор.
|
||||||
|
|
||||||
|
Клиент пишет боту в личку → зеркалится в топик support-группы (`client.py` —
|
||||||
|
тонкая HTTP-обёртка над Bot API; `bridge.py` — маршрутизация апдейтов и
|
||||||
|
персистентность через tg_support_* таблицы). Standalone entrypoint —
|
||||||
|
`app.tgbot_main` (long-polling воркер, НЕ webhook, отдельный контейнер/процесс).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
537
tradein-mvp/backend/app/services/tgbot/bridge.py
Normal file
537
tradein-mvp/backend/app/services/tgbot/bridge.py
Normal file
|
|
@ -0,0 +1,537 @@
|
||||||
|
"""Маршрутизация Telegram-апдейтов для support-моста (#tgsupport).
|
||||||
|
|
||||||
|
Поток:
|
||||||
|
A) Клиент пишет боту в личку (chat.type == 'private') →
|
||||||
|
upsert tg_support_users → (если первое сообщение за последний час — шапка
|
||||||
|
с идентификацией клиента в топик) → copyMessage контента в support-топик →
|
||||||
|
запись в tg_support_messages (direction='in', topic_message_id — ключ
|
||||||
|
маршрутизации ответа).
|
||||||
|
B) Оператор отвечает РЕПЛАЕМ в support-группе на зеркало клиента →
|
||||||
|
находим chat_id по topic_message_id → copyMessage ответа в личку клиента →
|
||||||
|
запись (direction='out'). Реплай не на зеркало (или не реплай вообще) —
|
||||||
|
обычная болтовня в топике, тихий игнор. Telegram 403 (клиент заблокировал
|
||||||
|
бота) → is_blocked=true + уведомление в топике.
|
||||||
|
C) Дедуп: update_id <= сохранённого offset — skip. Offset сохраняется И
|
||||||
|
коммитится в той же транзакции, что и запись сообщения (см. `process_update`
|
||||||
|
`finally`), после КАЖДОГО апдейта — рестарт воркера не переигрывает уже
|
||||||
|
обработанные апдейты и не подвисает вечно на «ядовитом» апдейте.
|
||||||
|
D) TELEGRAM_BOT_TOKEN пуст → бот выключен — проверяется в `app.tgbot_main`
|
||||||
|
(entrypoint), не здесь.
|
||||||
|
E) /start клиенту → короткое приветствие МЕРЫ, без зеркалирования в топик
|
||||||
|
(команда — не содержательное обращение, не должна засорять топик).
|
||||||
|
|
||||||
|
Персистентность вынесена за `BridgeStorage`-протокол — маршрутизирующая логика
|
||||||
|
(`process_update` и приватные `_handle_*`) не завязана на реальную БД, тестируется
|
||||||
|
на in-memory fake storage + mock httpx (см. tests/services/tgbot/).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
from collections.abc import Callable
|
||||||
|
from typing import Any, Protocol
|
||||||
|
|
||||||
|
from sqlalchemy import text
|
||||||
|
from sqlalchemy.exc import SQLAlchemyError
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
from app.core.config import settings
|
||||||
|
from app.core.shutdown import shutdown_requested
|
||||||
|
from app.services.tgbot.client import TelegramApiError, TelegramClient
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# Ключ в tg_support_state под который сохраняется last processed update_id
|
||||||
|
# (см. data/sql/186_tg_support.sql — комментарий на колонке .key).
|
||||||
|
_OFFSET_KEY = "last_update_id"
|
||||||
|
|
||||||
|
# Окно, за которое повторное сообщение клиента НЕ дублирует шапку-идентификацию
|
||||||
|
# в топике (одна шапка на "сессию" обращения).
|
||||||
|
_HEADER_THROTTLE_WINDOW_S = 3600
|
||||||
|
|
||||||
|
GREETING_TEXT = (
|
||||||
|
"Здравствуйте! Это служба поддержки МЕРА (сервис trade-in квартир). "
|
||||||
|
"Опишите ваш вопрос — оператор ответит вам в этом чате в ближайшее время."
|
||||||
|
)
|
||||||
|
|
||||||
|
# Отправляется клиенту вместо тихой потери сообщения, если TELEGRAM_SUPPORT_CHAT_ID
|
||||||
|
# не сконфигурирован (иначе клиент ждёт ответа, которого никогда не будет — #5 review).
|
||||||
|
SERVICE_UNAVAILABLE_TEXT = (
|
||||||
|
"Служба поддержки временно недоступна. Пожалуйста, попробуйте написать позже."
|
||||||
|
)
|
||||||
|
|
||||||
|
# Kinds, задокументированные в data/sql/186_tg_support.sql (COMMENT ON COLUMN
|
||||||
|
# tg_support_messages.kind): "text | photo | document | video | voice | other".
|
||||||
|
_KNOWN_KINDS = ("text", "photo", "document", "video", "voice")
|
||||||
|
|
||||||
|
|
||||||
|
# ── Storage abstraction (testable без реальной БД) ──────────────────────────
|
||||||
|
class BridgeStorage(Protocol):
|
||||||
|
"""Persistence-контракт моста. `SqlBridgeStorage` — прод-реализация поверх
|
||||||
|
tg_support_* (см. data/sql/186_tg_support.sql). Тесты используют in-memory fake."""
|
||||||
|
|
||||||
|
def get_offset(self) -> int: ...
|
||||||
|
|
||||||
|
def save_offset(self, update_id: int) -> None: ...
|
||||||
|
|
||||||
|
def commit(self) -> None: ...
|
||||||
|
|
||||||
|
def rollback(self) -> None: ...
|
||||||
|
|
||||||
|
def upsert_user(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
username: str | None,
|
||||||
|
first_name: str | None,
|
||||||
|
last_name: str | None,
|
||||||
|
language_code: str | None,
|
||||||
|
) -> None: ...
|
||||||
|
|
||||||
|
def had_recent_inbound(self, chat_id: int, window_seconds: int) -> bool: ...
|
||||||
|
|
||||||
|
def record_message(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
direction: str,
|
||||||
|
tg_message_id: int | None,
|
||||||
|
topic_message_id: int | None,
|
||||||
|
kind: str,
|
||||||
|
text_body: str | None,
|
||||||
|
operator_tg_id: int | None,
|
||||||
|
) -> int | None: ...
|
||||||
|
|
||||||
|
def find_chat_by_topic_message(self, topic_message_id: int) -> int | None: ...
|
||||||
|
|
||||||
|
def mark_blocked(self, chat_id: int) -> None: ...
|
||||||
|
|
||||||
|
|
||||||
|
class SqlBridgeStorage:
|
||||||
|
"""`BridgeStorage` поверх SQLAlchemy Session (psycopg v3), tg_support_* таблицы.
|
||||||
|
|
||||||
|
Методы исполняют SQL немедленно, но НЕ коммитят по отдельности — коммит
|
||||||
|
один раз в конце `process_update` (после записи сообщения И offset'а), чтобы
|
||||||
|
оба изменения фиксировались атомарно в одной транзакции (требование C).
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, db: Session) -> None:
|
||||||
|
self._db = db
|
||||||
|
|
||||||
|
def get_offset(self) -> int:
|
||||||
|
row = self._db.execute(
|
||||||
|
text("SELECT value FROM tg_support_state WHERE key = CAST(:key AS text)"),
|
||||||
|
{"key": _OFFSET_KEY},
|
||||||
|
).fetchone()
|
||||||
|
if row is None or row[0] is None:
|
||||||
|
return 0
|
||||||
|
try:
|
||||||
|
return int(row[0])
|
||||||
|
except (TypeError, ValueError):
|
||||||
|
logger.warning("tgbot storage: невалидный offset в БД (%r) — считаем 0", row[0])
|
||||||
|
return 0
|
||||||
|
|
||||||
|
def save_offset(self, update_id: int) -> None:
|
||||||
|
self._db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
INSERT INTO tg_support_state (key, value, updated_at)
|
||||||
|
VALUES (CAST(:key AS text), CAST(:value AS text), NOW())
|
||||||
|
ON CONFLICT (key) DO UPDATE
|
||||||
|
SET value = EXCLUDED.value, updated_at = NOW()
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"key": _OFFSET_KEY, "value": str(update_id)},
|
||||||
|
)
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
self._db.commit()
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
"""Откатывает текущую (возможно failed-transaction) сессию перед save_offset.
|
||||||
|
|
||||||
|
Нужно, когда исключение пришло от самой БД (напр. обрыв коннекта к
|
||||||
|
postgres при деплое) — SQLAlchemy Session после такого исключения
|
||||||
|
переходит в failed-transaction state, и ЛЮБОЙ следующий `execute()`
|
||||||
|
(включая `save_offset`) кидает `PendingRollbackError` без явного
|
||||||
|
rollback() (#3 review — иначе update_id не сдвигается, апдейт
|
||||||
|
переигрывается на следующей итерации, copyMessage дублирует зеркало
|
||||||
|
клиента в топик на каждый повтор).
|
||||||
|
"""
|
||||||
|
self._db.rollback()
|
||||||
|
|
||||||
|
def upsert_user(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
username: str | None,
|
||||||
|
first_name: str | None,
|
||||||
|
last_name: str | None,
|
||||||
|
language_code: str | None,
|
||||||
|
) -> None:
|
||||||
|
self._db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
INSERT INTO tg_support_users
|
||||||
|
(chat_id, username, first_name, last_name, language_code,
|
||||||
|
created_at, last_seen_at, is_blocked)
|
||||||
|
VALUES
|
||||||
|
(CAST(:chat_id AS bigint), :username, :first_name, :last_name,
|
||||||
|
:language_code, NOW(), NOW(), FALSE)
|
||||||
|
ON CONFLICT (chat_id) DO UPDATE
|
||||||
|
SET username = EXCLUDED.username,
|
||||||
|
first_name = EXCLUDED.first_name,
|
||||||
|
last_name = EXCLUDED.last_name,
|
||||||
|
language_code = EXCLUDED.language_code,
|
||||||
|
last_seen_at = NOW(),
|
||||||
|
is_blocked = FALSE
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"username": username,
|
||||||
|
"first_name": first_name,
|
||||||
|
"last_name": last_name,
|
||||||
|
"language_code": language_code,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
def had_recent_inbound(self, chat_id: int, window_seconds: int) -> bool:
|
||||||
|
row = self._db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
SELECT 1
|
||||||
|
FROM tg_support_messages
|
||||||
|
WHERE chat_id = CAST(:chat_id AS bigint)
|
||||||
|
AND direction = 'in'
|
||||||
|
AND created_at > NOW() - make_interval(secs => CAST(:window_seconds AS integer))
|
||||||
|
LIMIT 1
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"chat_id": chat_id, "window_seconds": window_seconds},
|
||||||
|
).fetchone()
|
||||||
|
return row is not None
|
||||||
|
|
||||||
|
def record_message(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
direction: str,
|
||||||
|
tg_message_id: int | None,
|
||||||
|
topic_message_id: int | None,
|
||||||
|
kind: str,
|
||||||
|
text_body: str | None,
|
||||||
|
operator_tg_id: int | None,
|
||||||
|
) -> int | None:
|
||||||
|
row = self._db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
INSERT INTO tg_support_messages
|
||||||
|
(chat_id, direction, tg_message_id, topic_message_id, kind,
|
||||||
|
text_body, operator_tg_id, created_at)
|
||||||
|
VALUES
|
||||||
|
(CAST(:chat_id AS bigint), CAST(:direction AS text),
|
||||||
|
CAST(:tg_message_id AS bigint), CAST(:topic_message_id AS bigint),
|
||||||
|
CAST(:kind AS text), :text_body, CAST(:operator_tg_id AS bigint), NOW())
|
||||||
|
RETURNING id
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"direction": direction,
|
||||||
|
"tg_message_id": tg_message_id,
|
||||||
|
"topic_message_id": topic_message_id,
|
||||||
|
"kind": kind,
|
||||||
|
"text_body": text_body,
|
||||||
|
"operator_tg_id": operator_tg_id,
|
||||||
|
},
|
||||||
|
).fetchone()
|
||||||
|
return int(row[0]) if row is not None else None
|
||||||
|
|
||||||
|
def find_chat_by_topic_message(self, topic_message_id: int) -> int | None:
|
||||||
|
row = self._db.execute(
|
||||||
|
text(
|
||||||
|
"""
|
||||||
|
SELECT chat_id
|
||||||
|
FROM tg_support_messages
|
||||||
|
WHERE topic_message_id = CAST(:topic_message_id AS bigint)
|
||||||
|
AND direction = 'in'
|
||||||
|
ORDER BY created_at DESC
|
||||||
|
LIMIT 1
|
||||||
|
"""
|
||||||
|
),
|
||||||
|
{"topic_message_id": topic_message_id},
|
||||||
|
).fetchone()
|
||||||
|
return int(row[0]) if row is not None else None
|
||||||
|
|
||||||
|
def mark_blocked(self, chat_id: int) -> None:
|
||||||
|
self._db.execute(
|
||||||
|
text(
|
||||||
|
"UPDATE tg_support_users SET is_blocked = TRUE "
|
||||||
|
"WHERE chat_id = CAST(:chat_id AS bigint)"
|
||||||
|
),
|
||||||
|
{"chat_id": chat_id},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Pure helpers ──────────────────────────────────────────────────────────────
|
||||||
|
def _infer_kind(message: dict[str, Any]) -> str:
|
||||||
|
"""Content-type сообщения → kind-строка. Неизвестные типы (voice/sticker/location/
|
||||||
|
etc.) сворачиваются в 'other' — см. документированный набор в COMMENT ON COLUMN."""
|
||||||
|
for field in _KNOWN_KINDS:
|
||||||
|
if field in message:
|
||||||
|
return field
|
||||||
|
return "other"
|
||||||
|
|
||||||
|
|
||||||
|
def _format_topic_header(
|
||||||
|
chat_id: int, username: str | None, first_name: str | None, last_name: str | None
|
||||||
|
) -> str:
|
||||||
|
"""Короткая шапка-идентификация клиента для support-топика."""
|
||||||
|
display_name = " ".join(p for p in (first_name, last_name) if p) or "Без имени"
|
||||||
|
username_part = f", @{username}" if username else ""
|
||||||
|
return f"Новое обращение от {display_name}{username_part} (chat_id={chat_id})"
|
||||||
|
|
||||||
|
|
||||||
|
# ── Update routing ────────────────────────────────────────────────────────────
|
||||||
|
async def _handle_private_message(
|
||||||
|
message: dict[str, Any], client: TelegramClient, storage: BridgeStorage
|
||||||
|
) -> None:
|
||||||
|
"""A) Личка клиента → бот. upsert user → (опц. шапка) → зеркало в топик."""
|
||||||
|
chat = message.get("chat") or {}
|
||||||
|
chat_id = chat.get("id")
|
||||||
|
if not isinstance(chat_id, int):
|
||||||
|
logger.warning("tgbot bridge: приватное сообщение без валидного chat.id — игнор")
|
||||||
|
return
|
||||||
|
|
||||||
|
from_user = message.get("from") or {}
|
||||||
|
username = from_user.get("username")
|
||||||
|
first_name = from_user.get("first_name")
|
||||||
|
last_name = from_user.get("last_name")
|
||||||
|
language_code = from_user.get("language_code")
|
||||||
|
|
||||||
|
storage.upsert_user(
|
||||||
|
chat_id=chat_id,
|
||||||
|
username=username,
|
||||||
|
first_name=first_name,
|
||||||
|
last_name=last_name,
|
||||||
|
language_code=language_code,
|
||||||
|
)
|
||||||
|
|
||||||
|
text_body = message.get("text")
|
||||||
|
if text_body == "/start":
|
||||||
|
# E) команда — не содержательное обращение, топик не засоряем.
|
||||||
|
await client.send_message(chat_id=chat_id, text=GREETING_TEXT)
|
||||||
|
return
|
||||||
|
|
||||||
|
if not settings.telegram_support_chat_id:
|
||||||
|
logger.warning(
|
||||||
|
"tgbot bridge: TELEGRAM_SUPPORT_CHAT_ID не задан — сообщение от chat_id=%d "
|
||||||
|
"не может быть зеркалировано; отвечаем клиенту вместо тихой потери (#5 review)",
|
||||||
|
chat_id,
|
||||||
|
)
|
||||||
|
# Не молчим клиенту (#5 review) — иначе он ждёт ответа, которого никогда не будет.
|
||||||
|
await client.send_message(chat_id=chat_id, text=SERVICE_UNAVAILABLE_TEXT)
|
||||||
|
return
|
||||||
|
|
||||||
|
message_id = message.get("message_id")
|
||||||
|
if not isinstance(message_id, int):
|
||||||
|
logger.warning("tgbot bridge: приватное сообщение без message_id — игнор")
|
||||||
|
return
|
||||||
|
|
||||||
|
# Шапка — только на первое сообщение клиента за окно, иначе топик засоряется.
|
||||||
|
if not storage.had_recent_inbound(chat_id, window_seconds=_HEADER_THROTTLE_WINDOW_S):
|
||||||
|
header = _format_topic_header(chat_id, username, first_name, last_name)
|
||||||
|
await client.send_message(
|
||||||
|
chat_id=settings.telegram_support_chat_id,
|
||||||
|
text=header,
|
||||||
|
message_thread_id=settings.telegram_support_topic_id or None,
|
||||||
|
)
|
||||||
|
|
||||||
|
mirrored = await client.copy_message(
|
||||||
|
chat_id=settings.telegram_support_chat_id,
|
||||||
|
from_chat_id=chat_id,
|
||||||
|
message_id=message_id,
|
||||||
|
message_thread_id=settings.telegram_support_topic_id or None,
|
||||||
|
)
|
||||||
|
topic_message_id = mirrored.get("message_id") if isinstance(mirrored, dict) else None
|
||||||
|
|
||||||
|
storage.record_message(
|
||||||
|
chat_id=chat_id,
|
||||||
|
direction="in",
|
||||||
|
tg_message_id=message_id,
|
||||||
|
topic_message_id=topic_message_id,
|
||||||
|
kind=_infer_kind(message),
|
||||||
|
text_body=text_body or message.get("caption"),
|
||||||
|
operator_tg_id=None,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def _handle_group_reply(
|
||||||
|
message: dict[str, Any], client: TelegramClient, storage: BridgeStorage
|
||||||
|
) -> None:
|
||||||
|
"""B) Реплай оператора в support-группе → доставка ответа клиенту."""
|
||||||
|
reply_to = message.get("reply_to_message")
|
||||||
|
if not isinstance(reply_to, dict):
|
||||||
|
return # не реплай вообще — обычная болтовня в топике, тихий игнор
|
||||||
|
|
||||||
|
mirror_message_id = reply_to.get("message_id")
|
||||||
|
if not isinstance(mirror_message_id, int):
|
||||||
|
return
|
||||||
|
|
||||||
|
target_chat_id = storage.find_chat_by_topic_message(mirror_message_id)
|
||||||
|
if target_chat_id is None:
|
||||||
|
# Обычная болтовня в топике (реплай на чьё-то ещё сообщение) — не логируем,
|
||||||
|
# это ожидаемый шум. НО реплай на сообщение, отправленное САМИМ БОТОМ
|
||||||
|
# (is_bot=True) и при этом отсутствующее в tg_support_messages — подозрительно:
|
||||||
|
# вероятная причина — осиротевшее зеркало (воркер упал МЕЖДУ copyMessage и
|
||||||
|
# storage.commit() в `_handle_private_message` — зеркало ушло в Telegram, а
|
||||||
|
# запись в БД потерялась). Дискриминатор неидеальный (шапка-идентификация
|
||||||
|
# тоже от бота, но не routing-ключ — тоже даст этот WARNING), но лучше редкий
|
||||||
|
# ложный WARNING, чем оператор молча решает, что ответ клиенту доставлен,
|
||||||
|
# хотя реплай тихо утонул (#4 review — двухфазный протокол НЕ делаем, overkill).
|
||||||
|
reply_from = reply_to.get("from") or {}
|
||||||
|
if reply_from.get("is_bot"):
|
||||||
|
logger.warning(
|
||||||
|
"tgbot bridge: реплай на сообщение бота (message_id=%d) не найден в "
|
||||||
|
"tg_support_messages как зеркало клиента — возможно, осиротевшее "
|
||||||
|
"зеркало (крах между copyMessage и commit) или шапка-идентификация; "
|
||||||
|
"ответ оператора НЕ доставлен клиенту",
|
||||||
|
mirror_message_id,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
|
message_id = message.get("message_id")
|
||||||
|
if not isinstance(message_id, int):
|
||||||
|
return
|
||||||
|
|
||||||
|
operator = message.get("from") or {}
|
||||||
|
operator_id = operator.get("id")
|
||||||
|
|
||||||
|
try:
|
||||||
|
delivered = await client.copy_message(
|
||||||
|
chat_id=target_chat_id,
|
||||||
|
from_chat_id=settings.telegram_support_chat_id,
|
||||||
|
message_id=message_id,
|
||||||
|
)
|
||||||
|
except TelegramApiError as exc:
|
||||||
|
if exc.error_code == 403:
|
||||||
|
# Клиент заблокировал бота — фиксируем и уведомляем оператора в топике.
|
||||||
|
storage.mark_blocked(target_chat_id)
|
||||||
|
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,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
raise
|
||||||
|
|
||||||
|
tg_message_id = delivered.get("message_id") if isinstance(delivered, dict) else None
|
||||||
|
storage.record_message(
|
||||||
|
chat_id=target_chat_id,
|
||||||
|
direction="out",
|
||||||
|
tg_message_id=tg_message_id,
|
||||||
|
topic_message_id=None,
|
||||||
|
kind=_infer_kind(message),
|
||||||
|
text_body=message.get("text") or message.get("caption"),
|
||||||
|
operator_tg_id=operator_id,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def process_update(
|
||||||
|
update: dict[str, Any], client: TelegramClient, storage: BridgeStorage
|
||||||
|
) -> None:
|
||||||
|
"""Маршрутизирует один Telegram update. Дедуп (C) + атомарный offset-commit.
|
||||||
|
|
||||||
|
Дедуп: update_id <= сохранённого offset — skip без side-effects. 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
|
||||||
|
|
||||||
|
current_offset = storage.get_offset()
|
||||||
|
if update_id <= current_offset:
|
||||||
|
logger.debug(
|
||||||
|
"tgbot bridge: update_id=%d уже обработан (offset=%d) — skip",
|
||||||
|
update_id,
|
||||||
|
current_offset,
|
||||||
|
)
|
||||||
|
return
|
||||||
|
|
||||||
|
message = update.get("message")
|
||||||
|
try:
|
||||||
|
if isinstance(message, dict):
|
||||||
|
chat = message.get("chat") or {}
|
||||||
|
chat_type = chat.get("type")
|
||||||
|
chat_id = chat.get("id")
|
||||||
|
if chat_type == "private":
|
||||||
|
await _handle_private_message(message, client, storage)
|
||||||
|
elif settings.telegram_support_chat_id and chat_id == settings.telegram_support_chat_id:
|
||||||
|
await _handle_group_reply(message, client, storage)
|
||||||
|
# иначе — необрабатываемый тип чата/апдейта (edited_message, канал и
|
||||||
|
# т.п.) — тихий игнор, но offset всё равно сдвигаем ниже.
|
||||||
|
except SQLAlchemyError:
|
||||||
|
logger.exception(
|
||||||
|
"tgbot bridge: DB-ошибка на update_id=%d — rollback перед сохранением "
|
||||||
|
"offset (иначе save_offset сам упадёт на failed-transaction state)",
|
||||||
|
update_id,
|
||||||
|
)
|
||||||
|
storage.rollback()
|
||||||
|
except Exception:
|
||||||
|
logger.exception(
|
||||||
|
"tgbot bridge: обработка update_id=%d упала — offset всё равно сдвигаем "
|
||||||
|
"(не блокируем поток на 'ядовитом' апдейте)",
|
||||||
|
update_id,
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
storage.save_offset(update_id)
|
||||||
|
storage.commit()
|
||||||
|
|
||||||
|
|
||||||
|
# ── Long-polling loop ─────────────────────────────────────────────────────────
|
||||||
|
async def run_poll_loop(
|
||||||
|
client: TelegramClient,
|
||||||
|
session_factory: Callable[[], Session],
|
||||||
|
poll_timeout_s: int = 30,
|
||||||
|
) -> None:
|
||||||
|
"""Бесконечный long-polling цикл до `shutdown_requested()`.
|
||||||
|
|
||||||
|
Свежая DB-сессия на каждую итерацию (одна итерация = один getUpdates-вызов +
|
||||||
|
обработка полученной пачки апдейтов) — не держим соединение открытым на
|
||||||
|
неопределённый срок между итерациями.
|
||||||
|
"""
|
||||||
|
logger.info("tgbot bridge: старт poll loop (timeout=%ds)", poll_timeout_s)
|
||||||
|
consecutive_errors = 0
|
||||||
|
while not shutdown_requested():
|
||||||
|
try:
|
||||||
|
with session_factory() as db:
|
||||||
|
storage = SqlBridgeStorage(db)
|
||||||
|
offset = storage.get_offset()
|
||||||
|
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)
|
||||||
|
consecutive_errors = 0
|
||||||
|
except Exception:
|
||||||
|
consecutive_errors += 1
|
||||||
|
backoff = min(5 * consecutive_errors, 60)
|
||||||
|
logger.exception("tgbot bridge: итерация poll loop упала — retry через %ds", backoff)
|
||||||
|
await asyncio.sleep(backoff)
|
||||||
|
logger.info("tgbot bridge: poll loop остановлен (shutdown)")
|
||||||
249
tradein-mvp/backend/app/services/tgbot/client.py
Normal file
249
tradein-mvp/backend/app/services/tgbot/client.py
Normal file
|
|
@ -0,0 +1,249 @@
|
||||||
|
"""Тонкая httpx-обёртка над Telegram Bot API (#tgsupport).
|
||||||
|
|
||||||
|
Зачем свой клиент, а не aiogram: единственные нужные методы — `getUpdates`
|
||||||
|
(long-polling), `copyMessage` (зеркалирование ЛЮБОГО типа контента без ре-аплоада)
|
||||||
|
и `sendMessage` (заголовки/приветствия/уведомления). aiogram — избыточная
|
||||||
|
зависимость (webhook-framework, dispatcher, FSM) ради трёх HTTP-вызовов; в стеке
|
||||||
|
уже есть httpx (см. `app.services.dadata`, `app.services.geocoder` — тот же паттерн
|
||||||
|
retry/timeout).
|
||||||
|
|
||||||
|
Docs: https://core.telegram.org/bots/api
|
||||||
|
|
||||||
|
Ретраи:
|
||||||
|
- HTTP 429 (Too Many Requests) — уважаем `parameters.retry_after` из тела ответа
|
||||||
|
(Telegram сам говорит сколько ждать), fallback на `_DEFAULT_RETRY_AFTER_S`.
|
||||||
|
- HTTP 5xx / сетевые ошибки (timeout/connect) — экспоненциальный backoff,
|
||||||
|
`capped` на `_MAX_BACKOFF_S`.
|
||||||
|
- Любая другая 4xx (400/401/403/404) — НЕ ретраится, сразу `TelegramApiError`
|
||||||
|
(запрос некорректен или прав нет — повтор не поможет).
|
||||||
|
|
||||||
|
БЕЗОПАСНОСТЬ: наши `logger.*`-вызовы здесь содержат только имя метода API,
|
||||||
|
HTTP-статус и `description` из ответа Telegram — токен туда не пишем.
|
||||||
|
Это НЕ гарантирует, что токен не утечёт по другим стокам: он живёт в
|
||||||
|
`self._base`/`url` (локальные переменные stack-фрейма `_request`), а GlitchTip
|
||||||
|
(sentry_sdk) по умолчанию прикладывает locals к traceback и Httpx-интеграция
|
||||||
|
кладёт полный URL в span data. Эти стоки закрываются НЕ здесь, а в
|
||||||
|
`app.tgbot_main` (`include_local_variables=False`, `before_send`-редактор,
|
||||||
|
`traces_sample_rate=0.0`) и подавлением INFO-логов самого `httpx`-логгера
|
||||||
|
(который печатает полный request URL, включая токен, на уровне INFO).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
_DEFAULT_TIMEOUT_S = 15.0
|
||||||
|
_DEFAULT_RETRY_AFTER_S = 5.0
|
||||||
|
_MAX_BACKOFF_S = 30.0
|
||||||
|
_DEFAULT_MAX_RETRIES = 5
|
||||||
|
|
||||||
|
|
||||||
|
class TelegramApiError(Exception):
|
||||||
|
"""Telegram Bot API ответил `ok: false` (после исчерпания ретраев, если применимо)."""
|
||||||
|
|
||||||
|
def __init__(self, method: str, error_code: int, description: str) -> None:
|
||||||
|
self.method = method
|
||||||
|
self.error_code = error_code
|
||||||
|
self.description = description
|
||||||
|
super().__init__(f"Telegram API {method} failed: {error_code} {description}")
|
||||||
|
|
||||||
|
|
||||||
|
def _extract_retry_after(
|
||||||
|
response: httpx.Response, default: float = _DEFAULT_RETRY_AFTER_S
|
||||||
|
) -> float:
|
||||||
|
"""Достаёт `parameters.retry_after` из тела 429-ответа. Fallback — `default`."""
|
||||||
|
try:
|
||||||
|
data = response.json()
|
||||||
|
except ValueError:
|
||||||
|
return default
|
||||||
|
if not isinstance(data, dict):
|
||||||
|
return default
|
||||||
|
params = data.get("parameters")
|
||||||
|
if isinstance(params, dict):
|
||||||
|
retry_after = params.get("retry_after")
|
||||||
|
if isinstance(retry_after, int | float):
|
||||||
|
return float(retry_after)
|
||||||
|
return default
|
||||||
|
|
||||||
|
|
||||||
|
def _error_from_body(response: httpx.Response) -> tuple[int, str]:
|
||||||
|
"""Парсит (error_code, description) из тела ответа Telegram; fallback на HTTP-статус."""
|
||||||
|
try:
|
||||||
|
data = response.json()
|
||||||
|
except ValueError:
|
||||||
|
return response.status_code, (response.text or "")[:200]
|
||||||
|
if not isinstance(data, dict):
|
||||||
|
return response.status_code, str(data)[:200]
|
||||||
|
error_code = data.get("error_code", response.status_code)
|
||||||
|
description = data.get("description", "")
|
||||||
|
code = int(error_code) if isinstance(error_code, int | float) else response.status_code
|
||||||
|
return code, str(description)
|
||||||
|
|
||||||
|
|
||||||
|
class TelegramClient:
|
||||||
|
"""Bot API клиент на httpx.AsyncClient. Каждый вызов — отдельное короткоживущее соединение."""
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
token: str,
|
||||||
|
base_url: str = "https://api.telegram.org",
|
||||||
|
timeout: float = _DEFAULT_TIMEOUT_S,
|
||||||
|
) -> None:
|
||||||
|
self._base = f"{base_url}/bot{token}"
|
||||||
|
self._timeout = timeout
|
||||||
|
|
||||||
|
async def _request(
|
||||||
|
self,
|
||||||
|
method: str,
|
||||||
|
payload: dict[str, Any],
|
||||||
|
*,
|
||||||
|
timeout: float | None = None,
|
||||||
|
max_retries: int = _DEFAULT_MAX_RETRIES,
|
||||||
|
) -> Any:
|
||||||
|
"""POST `method` с JSON-телом `payload`. Ретраит 429/5xx/network, иначе raise сразу."""
|
||||||
|
url = f"{self._base}/{method}"
|
||||||
|
effective_timeout = timeout if timeout is not None else self._timeout
|
||||||
|
attempt = 0
|
||||||
|
|
||||||
|
while True:
|
||||||
|
attempt += 1
|
||||||
|
try:
|
||||||
|
async with httpx.AsyncClient(timeout=effective_timeout) as client:
|
||||||
|
response = await client.post(url, json=payload)
|
||||||
|
except (httpx.TimeoutException, httpx.NetworkError) as exc:
|
||||||
|
if attempt > max_retries:
|
||||||
|
logger.error(
|
||||||
|
"tg client: %s — network error после %d попыток: %s", method, attempt, exc
|
||||||
|
)
|
||||||
|
raise
|
||||||
|
backoff = min(2.0**attempt, _MAX_BACKOFF_S)
|
||||||
|
logger.warning(
|
||||||
|
"tg client: %s — network error (попытка %d/%d): %s — retry через %.0fs",
|
||||||
|
method,
|
||||||
|
attempt,
|
||||||
|
max_retries,
|
||||||
|
exc,
|
||||||
|
backoff,
|
||||||
|
)
|
||||||
|
await asyncio.sleep(backoff)
|
||||||
|
continue
|
||||||
|
|
||||||
|
if response.status_code == 429:
|
||||||
|
retry_after = _extract_retry_after(response)
|
||||||
|
if attempt > max_retries:
|
||||||
|
error_code, description = _error_from_body(response)
|
||||||
|
logger.error("tg client: %s — 429 после %d попыток, сдаёмся", method, attempt)
|
||||||
|
raise TelegramApiError(method, error_code, description)
|
||||||
|
logger.warning(
|
||||||
|
"tg client: %s — HTTP 429 (попытка %d/%d), retry_after=%.0fs",
|
||||||
|
method,
|
||||||
|
attempt,
|
||||||
|
max_retries,
|
||||||
|
retry_after,
|
||||||
|
)
|
||||||
|
await asyncio.sleep(retry_after)
|
||||||
|
continue
|
||||||
|
|
||||||
|
if response.status_code >= 500:
|
||||||
|
if attempt > max_retries:
|
||||||
|
error_code, description = _error_from_body(response)
|
||||||
|
logger.error(
|
||||||
|
"tg client: %s — HTTP %d после %d попыток, сдаёмся",
|
||||||
|
method,
|
||||||
|
response.status_code,
|
||||||
|
attempt,
|
||||||
|
)
|
||||||
|
raise TelegramApiError(method, error_code, description)
|
||||||
|
backoff = min(2.0**attempt, _MAX_BACKOFF_S)
|
||||||
|
logger.warning(
|
||||||
|
"tg client: %s — HTTP %d (попытка %d/%d) — retry через %.0fs",
|
||||||
|
method,
|
||||||
|
response.status_code,
|
||||||
|
attempt,
|
||||||
|
max_retries,
|
||||||
|
backoff,
|
||||||
|
)
|
||||||
|
await asyncio.sleep(backoff)
|
||||||
|
continue
|
||||||
|
|
||||||
|
if response.status_code >= 400:
|
||||||
|
# 4xx кроме 429 — запрос некорректен/прав нет, повтор не поможет.
|
||||||
|
error_code, description = _error_from_body(response)
|
||||||
|
raise TelegramApiError(method, error_code, description)
|
||||||
|
|
||||||
|
try:
|
||||||
|
data = response.json()
|
||||||
|
except ValueError as exc:
|
||||||
|
raise TelegramApiError(
|
||||||
|
method, response.status_code, f"invalid json: {exc}"
|
||||||
|
) from exc
|
||||||
|
|
||||||
|
if not isinstance(data, dict) or not data.get("ok"):
|
||||||
|
error_code, description = _error_from_body(response)
|
||||||
|
raise TelegramApiError(method, error_code, description)
|
||||||
|
|
||||||
|
return data.get("result")
|
||||||
|
|
||||||
|
async def get_updates(
|
||||||
|
self,
|
||||||
|
offset: int,
|
||||||
|
timeout: int = 30,
|
||||||
|
allowed_updates: list[str] | None = None,
|
||||||
|
) -> list[dict[str, Any]]:
|
||||||
|
"""Long-polling getUpdates. `timeout` — сколько Telegram держит запрос открытым (сек).
|
||||||
|
|
||||||
|
HTTP-таймаут запроса берётся с запасом (`timeout + 10s`), чтобы не обрывать
|
||||||
|
соединение раньше, чем ответит сам Telegram long-poll.
|
||||||
|
"""
|
||||||
|
payload: dict[str, Any] = {"offset": offset, "timeout": timeout}
|
||||||
|
if allowed_updates is not None:
|
||||||
|
payload["allowed_updates"] = allowed_updates
|
||||||
|
result = await self._request(
|
||||||
|
"getUpdates", payload, timeout=float(timeout) + 10.0, max_retries=3
|
||||||
|
)
|
||||||
|
return result if isinstance(result, list) else []
|
||||||
|
|
||||||
|
async def copy_message(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
from_chat_id: int,
|
||||||
|
message_id: int,
|
||||||
|
message_thread_id: int | None = None,
|
||||||
|
reply_to_message_id: int | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
"""copyMessage — зеркалит ЛЮБОЙ тип контента без ре-аплоада файла."""
|
||||||
|
payload: dict[str, Any] = {
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"from_chat_id": from_chat_id,
|
||||||
|
"message_id": message_id,
|
||||||
|
}
|
||||||
|
if message_thread_id:
|
||||||
|
payload["message_thread_id"] = message_thread_id
|
||||||
|
if reply_to_message_id:
|
||||||
|
payload["reply_to_message_id"] = reply_to_message_id
|
||||||
|
result = await self._request("copyMessage", payload)
|
||||||
|
return result if isinstance(result, dict) else {}
|
||||||
|
|
||||||
|
async def send_message(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
text: str,
|
||||||
|
message_thread_id: int | None = None,
|
||||||
|
reply_to_message_id: int | None = None,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
"""sendMessage — текстовое сообщение (заголовки, приветствия, уведомления об ошибке)."""
|
||||||
|
payload: dict[str, Any] = {"chat_id": chat_id, "text": text}
|
||||||
|
if message_thread_id:
|
||||||
|
payload["message_thread_id"] = message_thread_id
|
||||||
|
if reply_to_message_id:
|
||||||
|
payload["reply_to_message_id"] = reply_to_message_id
|
||||||
|
result = await self._request("sendMessage", payload)
|
||||||
|
return result if isinstance(result, dict) else {}
|
||||||
180
tradein-mvp/backend/app/tgbot_main.py
Normal file
180
tradein-mvp/backend/app/tgbot_main.py
Normal file
|
|
@ -0,0 +1,180 @@
|
||||||
|
"""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_pii_event
|
||||||
|
|
||||||
|
def _before_send(event: Any, hint: dict[str, Any]) -> Any:
|
||||||
|
"""Композиция 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.
|
||||||
|
"""
|
||||||
|
scrubbed = scrub_pii_event(event, 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,
|
||||||
|
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:
|
||||||
|
client = TelegramClient(settings.telegram_bot_token)
|
||||||
|
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())
|
||||||
93
tradein-mvp/backend/data/sql/186_tg_support.sql
Normal file
93
tradein-mvp/backend/data/sql/186_tg_support.sql
Normal file
|
|
@ -0,0 +1,93 @@
|
||||||
|
-- 186_tg_support.sql
|
||||||
|
-- Telegram support bridge: @MERAsupport_bot mirrors client DMs into a
|
||||||
|
-- support-group topic; an operator replies in-thread; the bot relays the
|
||||||
|
-- reply back to the client's private chat.
|
||||||
|
--
|
||||||
|
-- WHY:
|
||||||
|
-- No durable state existed for this flow. Two things are required to make
|
||||||
|
-- it work reliably:
|
||||||
|
-- 1. A mapping from "message mirrored into the support topic" back to
|
||||||
|
-- "which client chat_id it came from" — this is how an operator's
|
||||||
|
-- reply (a Telegram reply-to a topic message) gets routed to the
|
||||||
|
-- right client. `topic_message_id` on tg_support_messages is that
|
||||||
|
-- routing key.
|
||||||
|
-- 2. A durable long-polling offset (`tg_support_state`) so a worker
|
||||||
|
-- restart does not replay already-processed Telegram updates.
|
||||||
|
--
|
||||||
|
-- WHAT:
|
||||||
|
-- - tg_support_users — one row per client Telegram private chat
|
||||||
|
-- (chat_id is the Telegram chat id, stable per
|
||||||
|
-- client, used directly as PK — no surrogate key
|
||||||
|
-- needed).
|
||||||
|
-- - tg_support_messages — full conversation log, both directions.
|
||||||
|
-- - tg_support_state — singleton key/value store for worker offsets
|
||||||
|
-- (e.g. key='last_update_id').
|
||||||
|
--
|
||||||
|
-- 152-FZ:
|
||||||
|
-- tg_support_users / tg_support_messages hold personal data (Telegram
|
||||||
|
-- username/name + free-text conversation content). ON DELETE CASCADE from
|
||||||
|
-- tg_support_users -> tg_support_messages makes client erasure a single
|
||||||
|
-- `DELETE FROM tg_support_users WHERE chat_id = :chat_id` statement, no
|
||||||
|
-- separate cleanup pass needed.
|
||||||
|
--
|
||||||
|
-- IDEMPOTENCY / SAFETY:
|
||||||
|
-- CREATE TABLE IF NOT EXISTS + CREATE INDEX IF NOT EXISTS throughout —
|
||||||
|
-- safe re-run. Purely additive: no existing table/view/column touched.
|
||||||
|
--
|
||||||
|
-- Dependencies: none (new standalone tables). Auto-applied on deploy via
|
||||||
|
-- _schema_migrations tracking (tradein-mvp/backend/data/sql convention).
|
||||||
|
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS tg_support_users (
|
||||||
|
chat_id bigint PRIMARY KEY,
|
||||||
|
username text,
|
||||||
|
first_name text,
|
||||||
|
last_name text,
|
||||||
|
language_code text,
|
||||||
|
created_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
last_seen_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
is_blocked boolean NOT NULL DEFAULT false
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE tg_support_users IS '152-ФЗ: ПДн клиентов Telegram-поддержки (@MERAsupport_bot). Удаление клиента — DELETE FROM tg_support_users WHERE chat_id=...; ON DELETE CASCADE в tg_support_messages подчищает переписку одной операцией.';
|
||||||
|
COMMENT ON COLUMN tg_support_users.chat_id IS 'Telegram private chat id клиента (стабильный, используется как PK напрямую).';
|
||||||
|
COMMENT ON COLUMN tg_support_users.is_blocked IS 'true, если клиент заблокировал бота (Telegram 403 на отправку) — бот перестаёт пытаться слать сообщения.';
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS tg_support_messages (
|
||||||
|
id bigserial PRIMARY KEY,
|
||||||
|
chat_id bigint NOT NULL REFERENCES tg_support_users (chat_id) ON DELETE CASCADE,
|
||||||
|
direction text NOT NULL CHECK (direction IN ('in', 'out')),
|
||||||
|
tg_message_id bigint,
|
||||||
|
topic_message_id bigint,
|
||||||
|
kind text NOT NULL,
|
||||||
|
text_body text,
|
||||||
|
operator_tg_id bigint,
|
||||||
|
created_at timestamptz NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE tg_support_messages IS '152-ФЗ: полный лог переписки Telegram-поддержки (ПДн, содержимое сообщений). Каскадно удаляется вместе с tg_support_users по chat_id.';
|
||||||
|
COMMENT ON COLUMN tg_support_messages.direction IS '''in'' — сообщение от клиента боту; ''out'' — ответ бота/оператора клиенту.';
|
||||||
|
COMMENT ON COLUMN tg_support_messages.tg_message_id IS 'id сообщения в личном чате с клиентом (Telegram message_id в chat_id).';
|
||||||
|
COMMENT ON COLUMN tg_support_messages.topic_message_id IS 'id зеркала сообщения в support-топике группы — ключ маршрутизации: реплай оператора на это сообщение адресуется данному chat_id.';
|
||||||
|
COMMENT ON COLUMN tg_support_messages.kind IS 'text | photo | document | video | voice | other.';
|
||||||
|
COMMENT ON COLUMN tg_support_messages.text_body IS 'Текст сообщения или caption медиа; NULL для медиа без подписи.';
|
||||||
|
COMMENT ON COLUMN tg_support_messages.operator_tg_id IS 'Telegram user id оператора, ответившего в топике; заполняется только для direction=''out''.';
|
||||||
|
|
||||||
|
CREATE UNIQUE INDEX IF NOT EXISTS tg_support_messages_topic_message_id_uq
|
||||||
|
ON tg_support_messages (topic_message_id)
|
||||||
|
WHERE topic_message_id IS NOT NULL;
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS tg_support_messages_chat_id_created_at_idx
|
||||||
|
ON tg_support_messages (chat_id, created_at DESC);
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS tg_support_state (
|
||||||
|
key text PRIMARY KEY,
|
||||||
|
value text NOT NULL,
|
||||||
|
updated_at timestamptz NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE tg_support_state IS 'Singleton key/value store для состояния Telegram-поддержки (например last_update_id для long-polling), переживает рестарт воркера.';
|
||||||
|
COMMENT ON COLUMN tg_support_state.key IS 'e.g. ''last_update_id''.';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
0
tradein-mvp/backend/tests/services/tgbot/__init__.py
Normal file
0
tradein-mvp/backend/tests/services/tgbot/__init__.py
Normal file
654
tradein-mvp/backend/tests/services/tgbot/test_bridge.py
Normal file
654
tradein-mvp/backend/tests/services/tgbot/test_bridge.py
Normal file
|
|
@ -0,0 +1,654 @@
|
||||||
|
"""Unit tests for `app.services.tgbot.bridge` — чистая логика роутинга.
|
||||||
|
|
||||||
|
Coverage (per task spec + review follow-up):
|
||||||
|
- user → topic (личка клиента зеркалится в support-топик, с шапкой на первое
|
||||||
|
сообщение за throttle-окно, без шапки на повторное В окне, и снова с шапкой
|
||||||
|
после истечения окна — #6 review)
|
||||||
|
- реплай оператора → user (доставка ответа клиенту + запись direction='out')
|
||||||
|
- реплай не на зеркало (или не реплай вообще) — тихий игнор, не мусорим в чат;
|
||||||
|
реплай на СООБЩЕНИЕ БОТА без записи в БД — WARNING про осиротевшее зеркало
|
||||||
|
(#4 review)
|
||||||
|
- дедуп update_id (<=offset — skip без side-effects; poison-pill апдейт всё
|
||||||
|
равно сдвигает offset, чтобы не подвесить весь поток)
|
||||||
|
- сбой БД (SQLAlchemyError) во время обработки → rollback() ПЕРЕД save_offset,
|
||||||
|
offset всё равно сдвигается — без этого следующий поллинг переиграл бы тот
|
||||||
|
же апдейт и задублировал зеркало клиента в топике (#3 review)
|
||||||
|
- /start → приветствие без зеркалирования
|
||||||
|
- Telegram 403 на доставку оператору → is_blocked + уведомление в топике
|
||||||
|
- TELEGRAM_SUPPORT_CHAT_ID не задан → клиенту уходит "сервис недоступен"
|
||||||
|
вместо тихой потери сообщения (#5 review)
|
||||||
|
|
||||||
|
NEVER calls real Telegram API — все HTTP-запросы mock'аются через
|
||||||
|
httpx.MockTransport (consistent с tests/services/test_dadata.py).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
import pytest
|
||||||
|
from sqlalchemy.exc import SQLAlchemyError
|
||||||
|
|
||||||
|
# DATABASE_URL required by app.core.config before any app import (см. test_dadata.py).
|
||||||
|
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
|
||||||
|
|
||||||
|
SUPPORT_CHAT_ID = -100123456789
|
||||||
|
SUPPORT_TOPIC_ID = 42
|
||||||
|
|
||||||
|
|
||||||
|
# ── Fake in-memory storage (БД не нужна) ─────────────────────────────────────
|
||||||
|
class FakeBridgeStorage:
|
||||||
|
"""In-memory `BridgeStorage` — никакой реальной БД, чистая логика роутинга.
|
||||||
|
|
||||||
|
`clock_s` — управляемые тестом "фейковые часы" (просто float, продвигается
|
||||||
|
вручную через `storage.clock_s += ...`), чтобы честно проверить throttle-окно
|
||||||
|
в `had_recent_inbound` (#6 review) без реального `time.sleep`/datetime-моков.
|
||||||
|
|
||||||
|
`fail_next_record_message` — если True, следующий вызов `record_message`
|
||||||
|
кидает `SQLAlchemyError` (симулирует обрыв коннекта к БД) и сбрасывается в
|
||||||
|
False — для теста rollback-пути в `process_update` (#3 review).
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(self, offset: int = 0) -> None:
|
||||||
|
self._offset = offset
|
||||||
|
self.users: dict[int, dict[str, Any]] = {}
|
||||||
|
self.messages: list[dict[str, Any]] = []
|
||||||
|
self.blocked: set[int] = set()
|
||||||
|
self.commits = 0
|
||||||
|
self.rollbacks = 0
|
||||||
|
self._next_id = 1
|
||||||
|
self.clock_s: float = 0.0
|
||||||
|
self.fail_next_record_message = False
|
||||||
|
|
||||||
|
def get_offset(self) -> int:
|
||||||
|
return self._offset
|
||||||
|
|
||||||
|
def save_offset(self, update_id: int) -> None:
|
||||||
|
self._offset = update_id
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
self.commits += 1
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
self.rollbacks += 1
|
||||||
|
|
||||||
|
def upsert_user(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
username: str | None,
|
||||||
|
first_name: str | None,
|
||||||
|
last_name: str | None,
|
||||||
|
language_code: str | None,
|
||||||
|
) -> None:
|
||||||
|
self.users[chat_id] = {
|
||||||
|
"username": username,
|
||||||
|
"first_name": first_name,
|
||||||
|
"last_name": last_name,
|
||||||
|
"language_code": language_code,
|
||||||
|
}
|
||||||
|
|
||||||
|
def had_recent_inbound(self, chat_id: int, window_seconds: int) -> bool:
|
||||||
|
return any(
|
||||||
|
m["chat_id"] == chat_id
|
||||||
|
and m["direction"] == "in"
|
||||||
|
and (self.clock_s - m["recorded_at_s"]) < window_seconds
|
||||||
|
for m in self.messages
|
||||||
|
)
|
||||||
|
|
||||||
|
def record_message(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
chat_id: int,
|
||||||
|
direction: str,
|
||||||
|
tg_message_id: int | None,
|
||||||
|
topic_message_id: int | None,
|
||||||
|
kind: str,
|
||||||
|
text_body: str | None,
|
||||||
|
operator_tg_id: int | None,
|
||||||
|
) -> int:
|
||||||
|
if self.fail_next_record_message:
|
||||||
|
self.fail_next_record_message = False
|
||||||
|
raise SQLAlchemyError("simulated DB failure (deploy connection reset)")
|
||||||
|
row_id = self._next_id
|
||||||
|
self._next_id += 1
|
||||||
|
self.messages.append(
|
||||||
|
{
|
||||||
|
"id": row_id,
|
||||||
|
"chat_id": chat_id,
|
||||||
|
"direction": direction,
|
||||||
|
"tg_message_id": tg_message_id,
|
||||||
|
"topic_message_id": topic_message_id,
|
||||||
|
"kind": kind,
|
||||||
|
"text_body": text_body,
|
||||||
|
"operator_tg_id": operator_tg_id,
|
||||||
|
"recorded_at_s": self.clock_s,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
return row_id
|
||||||
|
|
||||||
|
def find_chat_by_topic_message(self, topic_message_id: int) -> int | None:
|
||||||
|
for m in reversed(self.messages):
|
||||||
|
if m["direction"] == "in" and m["topic_message_id"] == topic_message_id:
|
||||||
|
return m["chat_id"]
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_blocked(self, chat_id: int) -> None:
|
||||||
|
self.blocked.add(chat_id)
|
||||||
|
|
||||||
|
|
||||||
|
# ── httpx mocking helpers (mirrors tests/services/test_dadata.py) ───────────
|
||||||
|
_REAL_ASYNC_CLIENT = httpx.AsyncClient
|
||||||
|
|
||||||
|
|
||||||
|
def _method_from_url(url: httpx.URL) -> str:
|
||||||
|
return str(url).rsplit("/", 1)[-1]
|
||||||
|
|
||||||
|
|
||||||
|
def _make_client(
|
||||||
|
responses: dict[str, Any], calls: list[tuple[str, dict[str, Any]]]
|
||||||
|
) -> TelegramClient:
|
||||||
|
"""TelegramClient wired to a MockTransport. `responses[method]` may be a dict
|
||||||
|
(returned as Bot API `result`), an int (HTTP error status), or a callable
|
||||||
|
`(payload) -> dict`. Every request is recorded into `calls`."""
|
||||||
|
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
import json as _json
|
||||||
|
|
||||||
|
method = _method_from_url(request.url)
|
||||||
|
payload = _json.loads(request.content.decode("utf-8")) if request.content else {}
|
||||||
|
calls.append((method, payload))
|
||||||
|
|
||||||
|
canned = responses.get(method)
|
||||||
|
if isinstance(canned, int):
|
||||||
|
return httpx.Response(
|
||||||
|
canned, json={"ok": False, "error_code": canned, "description": "mocked error"}
|
||||||
|
)
|
||||||
|
if callable(canned):
|
||||||
|
canned = canned(payload)
|
||||||
|
result = canned if canned is not None else {"message_id": 999}
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": result})
|
||||||
|
|
||||||
|
transport = httpx.MockTransport(handler)
|
||||||
|
|
||||||
|
def factory(*_: object, **__: object) -> httpx.AsyncClient:
|
||||||
|
return _REAL_ASYNC_CLIENT(transport=transport)
|
||||||
|
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
import unittest.mock as mock
|
||||||
|
|
||||||
|
# Патчим httpx.AsyncClient ТОЛЬКО внутри client-модуля — не трогаем глобальный httpx.
|
||||||
|
patcher = mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory)
|
||||||
|
patcher.start()
|
||||||
|
return client
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _support_chat_settings(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
"""Все тесты по умолчанию считают support-группу/топик настроенными."""
|
||||||
|
monkeypatch.setattr(bridge.settings, "telegram_support_chat_id", SUPPORT_CHAT_ID)
|
||||||
|
monkeypatch.setattr(bridge.settings, "telegram_support_topic_id", SUPPORT_TOPIC_ID)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _stop_patches():
|
||||||
|
"""Останавливает httpx.AsyncClient monkeypatch после каждого теста (unittest.mock.patch.start()
|
||||||
|
без контекст-менеджера требует явного stop, чтобы не утекать в соседние тесты)."""
|
||||||
|
import unittest.mock as mock
|
||||||
|
|
||||||
|
yield
|
||||||
|
mock.patch.stopall()
|
||||||
|
|
||||||
|
|
||||||
|
def _private_message(
|
||||||
|
*,
|
||||||
|
message_id: int = 1,
|
||||||
|
chat_id: int = 555,
|
||||||
|
text: str | None = "Здравствуйте, вопрос по trade-in",
|
||||||
|
username: str | None = "client_ivan",
|
||||||
|
first_name: str | None = "Иван",
|
||||||
|
last_name: str | None = "Петров",
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
msg: dict[str, Any] = {
|
||||||
|
"message_id": message_id,
|
||||||
|
"chat": {"id": chat_id, "type": "private"},
|
||||||
|
"from": {
|
||||||
|
"id": chat_id,
|
||||||
|
"username": username,
|
||||||
|
"first_name": first_name,
|
||||||
|
"last_name": last_name,
|
||||||
|
"language_code": "ru",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
if text is not None:
|
||||||
|
msg["text"] = text
|
||||||
|
return msg
|
||||||
|
|
||||||
|
|
||||||
|
def _group_reply_message(
|
||||||
|
*,
|
||||||
|
message_id: int = 200,
|
||||||
|
reply_to_message_id: int | None = 100,
|
||||||
|
text: str = "Ответ оператора",
|
||||||
|
operator_id: int = 777,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
msg: dict[str, Any] = {
|
||||||
|
"message_id": message_id,
|
||||||
|
"chat": {"id": SUPPORT_CHAT_ID, "type": "supergroup"},
|
||||||
|
"from": {"id": operator_id, "username": "operator1"},
|
||||||
|
"text": text,
|
||||||
|
}
|
||||||
|
if reply_to_message_id is not None:
|
||||||
|
msg["reply_to_message"] = {"message_id": reply_to_message_id}
|
||||||
|
return msg
|
||||||
|
|
||||||
|
|
||||||
|
# ── A) user → topic ──────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def test_private_message_mirrors_to_topic_with_header_on_first_contact() -> None:
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({"copyMessage": {"message_id": 555}}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
update = {"update_id": 10, "message": _private_message()}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
# Первое сообщение за окно → шапка ПЕРЕД зеркалом контента.
|
||||||
|
assert methods == ["sendMessage", "copyMessage"]
|
||||||
|
|
||||||
|
header_call = calls[0][1]
|
||||||
|
assert header_call["chat_id"] == SUPPORT_CHAT_ID
|
||||||
|
assert header_call["message_thread_id"] == SUPPORT_TOPIC_ID
|
||||||
|
assert "Иван Петров" in header_call["text"]
|
||||||
|
assert "@client_ivan" in header_call["text"]
|
||||||
|
|
||||||
|
mirror_call = calls[1][1]
|
||||||
|
assert mirror_call["from_chat_id"] == 555
|
||||||
|
assert mirror_call["chat_id"] == SUPPORT_CHAT_ID
|
||||||
|
assert mirror_call["message_id"] == 1
|
||||||
|
assert mirror_call["message_thread_id"] == SUPPORT_TOPIC_ID
|
||||||
|
|
||||||
|
assert len(storage.messages) == 1
|
||||||
|
rec = storage.messages[0]
|
||||||
|
assert rec["direction"] == "in"
|
||||||
|
assert rec["chat_id"] == 555
|
||||||
|
assert rec["topic_message_id"] == 555 # copyMessage result.message_id
|
||||||
|
assert rec["kind"] == "text"
|
||||||
|
assert rec["text_body"] == "Здравствуйте, вопрос по trade-in"
|
||||||
|
|
||||||
|
assert storage.get_offset() == 10
|
||||||
|
assert storage.commits == 1
|
||||||
|
assert 555 in storage.users
|
||||||
|
|
||||||
|
|
||||||
|
async def test_private_message_second_message_within_window_skips_header() -> None:
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({"copyMessage": {"message_id": 556}}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
# Симулируем уже существующее inbound-сообщение за последний час.
|
||||||
|
storage.record_message(
|
||||||
|
chat_id=555,
|
||||||
|
direction="in",
|
||||||
|
tg_message_id=0,
|
||||||
|
topic_message_id=100,
|
||||||
|
kind="text",
|
||||||
|
text_body="первое сообщение",
|
||||||
|
operator_tg_id=None,
|
||||||
|
)
|
||||||
|
|
||||||
|
update = {"update_id": 11, "message": _private_message(message_id=2)}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
# Шапка НЕ отправляется повторно — только зеркало.
|
||||||
|
assert methods == ["copyMessage"]
|
||||||
|
assert len(storage.messages) == 2
|
||||||
|
|
||||||
|
|
||||||
|
async def test_private_message_header_resent_after_window_expires() -> None:
|
||||||
|
"""#6 review: throttle-окно (3600с) реально проверяется по времени — после
|
||||||
|
истечения окна шапка отправляется заново (не одна на весь чат навсегда)."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({"copyMessage": {"message_id": 557}}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
storage.record_message(
|
||||||
|
chat_id=555,
|
||||||
|
direction="in",
|
||||||
|
tg_message_id=0,
|
||||||
|
topic_message_id=100,
|
||||||
|
kind="text",
|
||||||
|
text_body="первое сообщение (час назад)",
|
||||||
|
operator_tg_id=None,
|
||||||
|
)
|
||||||
|
# Продвигаем фейковые часы за throttle-окно (3600с).
|
||||||
|
storage.clock_s += bridge._HEADER_THROTTLE_WINDOW_S + 1
|
||||||
|
|
||||||
|
update = {"update_id": 13, "message": _private_message(message_id=3)}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
assert methods == ["sendMessage", "copyMessage"] # шапка снова отправлена
|
||||||
|
assert len(storage.messages) == 2
|
||||||
|
|
||||||
|
|
||||||
|
async def test_private_message_start_sends_greeting_without_mirroring() -> None:
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
update = {"update_id": 12, "message": _private_message(text="/start")}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
assert methods == ["sendMessage"]
|
||||||
|
greeting_call = calls[0][1]
|
||||||
|
assert greeting_call["chat_id"] == 555
|
||||||
|
assert "МЕРА" in greeting_call["text"]
|
||||||
|
# /start не зеркалируется и не попадает в лог переписки.
|
||||||
|
assert storage.messages == []
|
||||||
|
assert storage.get_offset() == 12
|
||||||
|
|
||||||
|
|
||||||
|
async def test_private_message_notifies_client_when_support_chat_unset(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
"""#5 review: TELEGRAM_SUPPORT_CHAT_ID не задан → клиент получает "сервис
|
||||||
|
недоступен" вместо того, чтобы молча ждать ответа, который никогда не придёт."""
|
||||||
|
monkeypatch.setattr(bridge.settings, "telegram_support_chat_id", 0)
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
update = {"update_id": 14, "message": _private_message()}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
assert methods == ["sendMessage"]
|
||||||
|
notify_call = calls[0][1]
|
||||||
|
assert notify_call["chat_id"] == 555
|
||||||
|
assert notify_call["text"] == bridge.SERVICE_UNAVAILABLE_TEXT
|
||||||
|
# Ничего не зеркалируется и не пишется в лог переписки — support-группа не настроена.
|
||||||
|
assert storage.messages == []
|
||||||
|
assert storage.get_offset() == 14
|
||||||
|
|
||||||
|
|
||||||
|
# ── B) реплай оператора → user ──────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def test_group_reply_delivers_to_client_and_records_outbound() -> None:
|
||||||
|
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": 20,
|
||||||
|
"message": _group_reply_message(reply_to_message_id=100),
|
||||||
|
}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
assert methods == ["copyMessage"]
|
||||||
|
delivery = calls[0][1]
|
||||||
|
assert delivery["chat_id"] == 555
|
||||||
|
assert delivery["from_chat_id"] == SUPPORT_CHAT_ID
|
||||||
|
|
||||||
|
assert len(storage.messages) == 2
|
||||||
|
out_rec = storage.messages[-1]
|
||||||
|
assert out_rec["direction"] == "out"
|
||||||
|
assert out_rec["chat_id"] == 555
|
||||||
|
assert out_rec["operator_tg_id"] == 777
|
||||||
|
assert out_rec["text_body"] == "Ответ оператора"
|
||||||
|
assert storage.get_offset() == 20
|
||||||
|
|
||||||
|
|
||||||
|
async def test_group_reply_not_a_reply_is_ignored() -> None:
|
||||||
|
"""Обычное сообщение в топике (не реплай) — тихий игнор, никаких Telegram-вызовов."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
update = {
|
||||||
|
"update_id": 21,
|
||||||
|
"message": _group_reply_message(reply_to_message_id=None),
|
||||||
|
}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert calls == []
|
||||||
|
assert storage.messages == []
|
||||||
|
# offset всё равно сдвигается — апдейт "обработан" (даже если ничего не сделано).
|
||||||
|
assert storage.get_offset() == 21
|
||||||
|
|
||||||
|
|
||||||
|
async def test_group_reply_to_unknown_message_is_ignored() -> None:
|
||||||
|
"""Реплай на сообщение, которого нет в tg_support_messages как зеркало клиента —
|
||||||
|
тихий игнор (обычная болтовня в топике на постороннее сообщение, не от бота)."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
update = {
|
||||||
|
"update_id": 22,
|
||||||
|
"message": _group_reply_message(reply_to_message_id=999999),
|
||||||
|
}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert calls == []
|
||||||
|
assert storage.messages == []
|
||||||
|
assert storage.get_offset() == 22
|
||||||
|
|
||||||
|
|
||||||
|
async def test_group_reply_to_bot_message_without_record_logs_orphaned_mirror_warning(
|
||||||
|
caplog: pytest.LogCaptureFixture,
|
||||||
|
) -> None:
|
||||||
|
"""#4 review: реплай на сообщение БОТА, которого нет в tg_support_messages, —
|
||||||
|
вероятное осиротевшее зеркало (крах между copyMessage и commit). WARNING, не
|
||||||
|
тихий игнор — оператор иначе решит, что ответ клиенту доставлен."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
message = _group_reply_message(reply_to_message_id=100)
|
||||||
|
message["reply_to_message"]["from"] = {"id": 999, "is_bot": True, "username": "MERAsupport_bot"}
|
||||||
|
update = {"update_id": 24, "message": message}
|
||||||
|
|
||||||
|
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert calls == [] # ответ НЕ доставлен — routing-ключ потерян
|
||||||
|
assert "осиротевшее" in caplog.text
|
||||||
|
assert storage.get_offset() == 24
|
||||||
|
|
||||||
|
|
||||||
|
async def test_group_reply_to_non_bot_message_without_record_stays_silent(
|
||||||
|
caplog: pytest.LogCaptureFixture,
|
||||||
|
) -> None:
|
||||||
|
"""Обычный реплай на сообщение ДРУГОГО ЧЕЛОВЕКА (не бота) в топике — реальная
|
||||||
|
болтовня, никакого WARNING (дискриминатор `is_bot` работает в обе стороны)."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
message = _group_reply_message(reply_to_message_id=101)
|
||||||
|
message["reply_to_message"]["from"] = {"id": 42, "is_bot": False, "username": "colleague"}
|
||||||
|
update = {"update_id": 25, "message": message}
|
||||||
|
|
||||||
|
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert calls == []
|
||||||
|
assert caplog.text == ""
|
||||||
|
|
||||||
|
|
||||||
|
async def test_group_reply_403_marks_blocked_and_notifies_topic() -> None:
|
||||||
|
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,
|
||||||
|
)
|
||||||
|
|
||||||
|
update = {
|
||||||
|
"update_id": 23,
|
||||||
|
"message": _group_reply_message(reply_to_message_id=100),
|
||||||
|
}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
methods = [m for m, _ in calls]
|
||||||
|
assert methods == ["copyMessage", "sendMessage"]
|
||||||
|
notify_call = calls[1][1]
|
||||||
|
assert notify_call["chat_id"] == SUPPORT_CHAT_ID
|
||||||
|
assert "заблокирован" in notify_call["text"]
|
||||||
|
|
||||||
|
assert 555 in storage.blocked
|
||||||
|
# Неудачная доставка НЕ должна создавать фейковую запись 'out'.
|
||||||
|
assert len(storage.messages) == 1
|
||||||
|
assert storage.get_offset() == 23
|
||||||
|
|
||||||
|
|
||||||
|
# ── C) дедуп ──────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def test_dedup_update_id_leq_offset_is_skipped_without_side_effects() -> None:
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage(offset=50)
|
||||||
|
|
||||||
|
update = {"update_id": 50, "message": _private_message()}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert calls == []
|
||||||
|
assert storage.messages == []
|
||||||
|
assert storage.commits == 0 # ранний return — offset уже актуален, коммитить нечего
|
||||||
|
assert storage.get_offset() == 50
|
||||||
|
|
||||||
|
update_older = {"update_id": 10, "message": _private_message()}
|
||||||
|
await bridge.process_update(update_older, client, storage)
|
||||||
|
assert calls == []
|
||||||
|
assert storage.get_offset() == 50
|
||||||
|
|
||||||
|
|
||||||
|
async def test_poison_pill_update_still_advances_offset() -> None:
|
||||||
|
"""Апдейт, на котором обработчик упал (например, message без chat), не должен
|
||||||
|
подвесить весь поток — offset сдвигается даже при исключении внутри handler'а."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
malformed_message = {"message_id": 1, "chat": {"type": "private"}} # нет chat.id
|
||||||
|
update = {"update_id": 30, "message": malformed_message}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert storage.get_offset() == 30
|
||||||
|
assert storage.commits == 1
|
||||||
|
assert storage.messages == []
|
||||||
|
|
||||||
|
|
||||||
|
async def test_db_error_during_processing_rolls_back_and_still_advances_offset() -> None:
|
||||||
|
"""#3 review: SQLAlchemyError (напр. обрыв коннекта к БД при деплое) во время
|
||||||
|
`record_message` → storage.rollback() ПЕРЕД save_offset, offset всё равно
|
||||||
|
сдвигается. Без rollback() save_offset сам кинул бы PendingRollbackError →
|
||||||
|
process_update вылетел бы без сохранения offset'а → следующая итерация
|
||||||
|
переиграла бы тот же апдейт → copyMessage задублировал бы зеркало в топике
|
||||||
|
на каждый повтор поллинга (воспроизведено ревьюером)."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({"copyMessage": {"message_id": 558}}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
storage.fail_next_record_message = True
|
||||||
|
|
||||||
|
update = {"update_id": 15, "message": _private_message()}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
# copyMessage успел уйти в Telegram (реальная утечка мирроринга при DB-сбое —
|
||||||
|
# известное ограничение атомарности между внешним API и БД, вне scope этого фикса),
|
||||||
|
# но rollback() отработал, offset сдвинут, commit вызван РОВНО один раз (в finally).
|
||||||
|
assert storage.rollbacks == 1
|
||||||
|
assert storage.commits == 1
|
||||||
|
assert storage.get_offset() == 15
|
||||||
|
# Запись сообщения НЕ попала в storage (record_message упал до append).
|
||||||
|
assert storage.messages == []
|
||||||
|
|
||||||
|
# Повторный вызов с тем же update_id теперь корректно дедупится — НЕ переигрывается.
|
||||||
|
calls.clear()
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
assert calls == []
|
||||||
|
assert storage.get_offset() == 15
|
||||||
|
|
||||||
|
|
||||||
|
async def test_update_without_update_id_is_ignored() -> None:
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
await bridge.process_update({"message": _private_message()}, client, storage)
|
||||||
|
|
||||||
|
assert calls == []
|
||||||
|
assert storage.commits == 0
|
||||||
|
assert storage.get_offset() == 0
|
||||||
|
|
||||||
|
|
||||||
|
# ── kind inference ────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("message_extra", "expected_kind"),
|
||||||
|
[
|
||||||
|
({"text": "hi"}, "text"),
|
||||||
|
({"photo": [{"file_id": "x"}]}, "photo"),
|
||||||
|
({"document": {"file_id": "x"}}, "document"),
|
||||||
|
({"video": {"file_id": "x"}}, "video"),
|
||||||
|
({"voice": {"file_id": "x"}}, "voice"),
|
||||||
|
({"sticker": {"file_id": "x"}}, "other"),
|
||||||
|
({"location": {"latitude": 1, "longitude": 2}}, "other"),
|
||||||
|
({}, "other"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_infer_kind(message_extra: dict[str, Any], expected_kind: str) -> None:
|
||||||
|
message = {"message_id": 1, "chat": {"id": 1, "type": "private"}, **message_extra}
|
||||||
|
assert bridge._infer_kind(message) == expected_kind
|
||||||
|
|
||||||
|
|
||||||
|
# ── unrelated chat types ─────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
async def test_update_from_unrelated_chat_is_ignored_but_offset_advances() -> None:
|
||||||
|
"""Апдейт не из личного чата и не из support-группы (например, другой чат/канал)
|
||||||
|
— молча игнорируется, offset всё равно сдвигается."""
|
||||||
|
calls: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
client = _make_client({}, calls)
|
||||||
|
storage = FakeBridgeStorage()
|
||||||
|
|
||||||
|
message = {
|
||||||
|
"message_id": 1,
|
||||||
|
"chat": {"id": -999, "type": "group"},
|
||||||
|
"from": {"id": 1},
|
||||||
|
"text": "болтовня в постороннем чате",
|
||||||
|
}
|
||||||
|
update = {"update_id": 40, "message": message}
|
||||||
|
await bridge.process_update(update, client, storage)
|
||||||
|
|
||||||
|
assert calls == []
|
||||||
|
assert storage.messages == []
|
||||||
|
assert storage.get_offset() == 40
|
||||||
164
tradein-mvp/backend/tests/services/tgbot/test_client.py
Normal file
164
tradein-mvp/backend/tests/services/tgbot/test_client.py
Normal file
|
|
@ -0,0 +1,164 @@
|
||||||
|
"""Unit tests for `app.services.tgbot.client.TelegramClient` retry/backoff logic.
|
||||||
|
|
||||||
|
NEVER calls real Telegram API — httpx.MockTransport only (consistent с
|
||||||
|
tests/services/test_dadata.py). `asyncio.sleep` is patched to a no-op so retry
|
||||||
|
tests run instantly regardless of configured backoff/retry_after durations.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from typing import Any
|
||||||
|
from unittest import mock
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
from app.services.tgbot.client import TelegramApiError, TelegramClient
|
||||||
|
|
||||||
|
_REAL_ASYNC_CLIENT = httpx.AsyncClient
|
||||||
|
|
||||||
|
|
||||||
|
def _install_transport(handler) -> None:
|
||||||
|
transport = httpx.MockTransport(handler)
|
||||||
|
|
||||||
|
def factory(*_: object, **__: object) -> httpx.AsyncClient:
|
||||||
|
return _REAL_ASYNC_CLIENT(transport=transport)
|
||||||
|
|
||||||
|
mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory).start()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture(autouse=True)
|
||||||
|
def _stop_patches_and_noop_sleep():
|
||||||
|
sleep_patcher = mock.patch("app.services.tgbot.client.asyncio.sleep", return_value=None)
|
||||||
|
sleep_patcher.start()
|
||||||
|
yield
|
||||||
|
mock.patch.stopall()
|
||||||
|
|
||||||
|
|
||||||
|
async def test_get_updates_happy_path_returns_list() -> None:
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
assert request.url.path.endswith("/getUpdates")
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": [{"update_id": 1}]})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
updates = await client.get_updates(offset=1)
|
||||||
|
assert updates == [{"update_id": 1}]
|
||||||
|
|
||||||
|
|
||||||
|
async def test_never_logs_or_leaks_token_in_request_url_host() -> None:
|
||||||
|
"""Sanity: token lives only in the path, base host stays api.telegram.org."""
|
||||||
|
captured: dict[str, str] = {}
|
||||||
|
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
captured["url"] = str(request.url)
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": {}})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="super-secret-token")
|
||||||
|
await client.send_message(chat_id=1, text="hi")
|
||||||
|
assert "bot" + "super-secret-token" in captured["url"] # goes over the wire, not logged
|
||||||
|
|
||||||
|
|
||||||
|
async def test_copy_message_retries_on_429_then_succeeds() -> None:
|
||||||
|
calls = {"n": 0}
|
||||||
|
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
calls["n"] += 1
|
||||||
|
if calls["n"] == 1:
|
||||||
|
return httpx.Response(
|
||||||
|
429,
|
||||||
|
json={
|
||||||
|
"ok": False,
|
||||||
|
"error_code": 429,
|
||||||
|
"description": "Too Many Requests",
|
||||||
|
"parameters": {"retry_after": 3},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": {"message_id": 5}})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
result = await client.copy_message(chat_id=1, from_chat_id=2, message_id=3)
|
||||||
|
|
||||||
|
assert result == {"message_id": 5}
|
||||||
|
assert calls["n"] == 2
|
||||||
|
|
||||||
|
|
||||||
|
async def test_send_message_retries_on_5xx_then_succeeds() -> 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": "bad gw"}
|
||||||
|
)
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": {"message_id": 9}})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
result = await client.send_message(chat_id=1, text="retrying")
|
||||||
|
|
||||||
|
assert result == {"message_id": 9}
|
||||||
|
assert calls["n"] == 3
|
||||||
|
|
||||||
|
|
||||||
|
async def test_send_message_raises_immediately_on_non_retryable_4xx() -> None:
|
||||||
|
calls = {"n": 0}
|
||||||
|
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
calls["n"] += 1
|
||||||
|
return httpx.Response(
|
||||||
|
403, json={"ok": False, "error_code": 403, "description": "Forbidden: bot blocked"}
|
||||||
|
)
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
|
||||||
|
with pytest.raises(TelegramApiError) as exc_info:
|
||||||
|
await client.send_message(chat_id=1, text="hi")
|
||||||
|
|
||||||
|
assert exc_info.value.error_code == 403
|
||||||
|
assert calls["n"] == 1 # НЕ ретраится
|
||||||
|
|
||||||
|
|
||||||
|
async def test_copy_message_gives_up_after_max_retries_on_persistent_5xx() -> None:
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
return httpx.Response(500, json={"ok": False, "error_code": 500, "description": "boom"})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
|
||||||
|
with pytest.raises(TelegramApiError) as exc_info:
|
||||||
|
await client.copy_message(chat_id=1, from_chat_id=2, message_id=3)
|
||||||
|
|
||||||
|
assert exc_info.value.error_code == 500
|
||||||
|
|
||||||
|
|
||||||
|
async def test_get_updates_returns_empty_list_on_malformed_result() -> None:
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": "not-a-list"})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
assert await client.get_updates(offset=1) == []
|
||||||
|
|
||||||
|
|
||||||
|
async def test_optional_thread_and_reply_params_omitted_when_falsy() -> None:
|
||||||
|
captured: dict[str, Any] = {}
|
||||||
|
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
import json as _json
|
||||||
|
|
||||||
|
captured["body"] = _json.loads(request.content.decode("utf-8"))
|
||||||
|
return httpx.Response(200, json={"ok": True, "result": {"message_id": 1}})
|
||||||
|
|
||||||
|
_install_transport(handler)
|
||||||
|
client = TelegramClient(token="fake-token")
|
||||||
|
await client.copy_message(chat_id=1, from_chat_id=2, message_id=3, message_thread_id=0)
|
||||||
|
|
||||||
|
assert "message_thread_id" not in captured["body"]
|
||||||
|
|
@ -1,9 +1,21 @@
|
||||||
"""Unit-тесты consumer-PII scrubber для GlitchTip (#396)."""
|
"""Unit-тесты consumer-PII scrubber для GlitchTip (#396) + Telegram bot-токен
|
||||||
|
redaction (#tgsupport review — CRITICAL: токен утекал в GlitchTip двумя
|
||||||
|
независимыми векторами, воспроизведёнными ревьюером на реальном событии
|
||||||
|
sentry_sdk: stack-frame locals (`include_local_variables=True`) и httpx-span
|
||||||
|
`data` (`HttpxIntegration`). Тесты ниже конструируют event-словари ровно в той
|
||||||
|
форме, в которой их реально отдаёт sentry_sdk (frames[].vars, spans[].data),
|
||||||
|
чтобы regression на этот конкретный CRITICAL не проскочил молча."""
|
||||||
|
|
||||||
import os
|
import os
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
from app.observability.sentry_scrub import scrub_pii_event
|
from app.observability.sentry_scrub import (
|
||||||
|
redact_telegram_bot_token,
|
||||||
|
scrub_pii_event,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_redacts_pii_in_request_data() -> None:
|
def test_redacts_pii_in_request_data() -> None:
|
||||||
|
|
@ -88,3 +100,162 @@ def test_returns_event_not_none() -> None:
|
||||||
out = scrub_pii_event(event, {})
|
out = scrub_pii_event(event, {})
|
||||||
assert out is not None
|
assert out is not None
|
||||||
assert out is event
|
assert out is event
|
||||||
|
|
||||||
|
|
||||||
|
# ── Telegram bot-token redaction (#tgsupport review, CRITICAL) ──────────────
|
||||||
|
|
||||||
|
_LEAKED_TOKEN_URL = "https://api.telegram.org/bot8663867262:AAExampleSecretPartAbCdEf123/getMe"
|
||||||
|
|
||||||
|
|
||||||
|
def test_redacts_token_in_stack_frame_locals() -> None:
|
||||||
|
"""Вектор #1 (ревьюер): `include_local_variables=True` кладёт locals
|
||||||
|
`TelegramClient._request` (`self`, `url`) в traceback frame `vars`."""
|
||||||
|
event = {
|
||||||
|
"exception": {
|
||||||
|
"values": [
|
||||||
|
{
|
||||||
|
"type": "NetworkError",
|
||||||
|
"stacktrace": {
|
||||||
|
"frames": [
|
||||||
|
{
|
||||||
|
"function": "_request",
|
||||||
|
"vars": {
|
||||||
|
"url": _LEAKED_TOKEN_URL,
|
||||||
|
"self": {"_base": _LEAKED_TOKEN_URL.rsplit("/", 1)[0]},
|
||||||
|
"method": "getMe",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
out = redact_telegram_bot_token(event, {})
|
||||||
|
assert out is not None
|
||||||
|
frame_vars = out["exception"]["values"][0]["stacktrace"]["frames"][0]["vars"]
|
||||||
|
assert "8663867262:AAExampleSecretPartAbCdEf123" not in frame_vars["url"]
|
||||||
|
assert "8663867262:AAExampleSecretPartAbCdEf123" not in frame_vars["self"]["_base"]
|
||||||
|
assert frame_vars["url"] == "https://api.telegram.org/bot[REDACTED]/getMe"
|
||||||
|
# Не PII/секрет — остаётся как есть.
|
||||||
|
assert frame_vars["method"] == "getMe"
|
||||||
|
|
||||||
|
|
||||||
|
def test_redacts_token_in_httpx_span_data() -> None:
|
||||||
|
"""Вектор #2 (ревьюер): `HttpxIntegration` кладёт полный request URL в span
|
||||||
|
`data` независимо от traceback — `traces_sample_rate=0.0` спасает сейчас, но
|
||||||
|
редактор — belt-and-suspenders на случай если трейсинг когда-нибудь включат."""
|
||||||
|
event = {
|
||||||
|
"spans": [
|
||||||
|
{
|
||||||
|
"op": "http.client",
|
||||||
|
"description": "POST " + _LEAKED_TOKEN_URL,
|
||||||
|
"data": {"url": _LEAKED_TOKEN_URL, "http.method": "POST"},
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
out = redact_telegram_bot_token(event, {})
|
||||||
|
assert out is not None
|
||||||
|
span = out["spans"][0]
|
||||||
|
assert "8663867262:AAExampleSecretPartAbCdEf123" not in span["data"]["url"]
|
||||||
|
assert "8663867262:AAExampleSecretPartAbCdEf123" not in span["description"]
|
||||||
|
assert span["data"]["http.method"] == "POST"
|
||||||
|
|
||||||
|
|
||||||
|
def test_redacts_token_in_breadcrumb_message() -> None:
|
||||||
|
"""httpx-логгер (INFO, до нашего getLogger('httpx').setLevel(WARNING) в
|
||||||
|
tgbot_main) может всплыть breadcrumb'ом с полным URL через LoggingIntegration."""
|
||||||
|
event = {
|
||||||
|
"breadcrumbs": {
|
||||||
|
"values": [
|
||||||
|
{
|
||||||
|
"category": "httpx",
|
||||||
|
"message": f'HTTP Request: POST {_LEAKED_TOKEN_URL} "HTTP/1.1 401"',
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
out = redact_telegram_bot_token(event, {})
|
||||||
|
assert out is not None
|
||||||
|
msg = out["breadcrumbs"]["values"][0]["message"]
|
||||||
|
assert "8663867262:AAExampleSecretPartAbCdEf123" not in msg
|
||||||
|
assert "/bot[REDACTED]/getMe" in msg
|
||||||
|
|
||||||
|
|
||||||
|
def test_redact_token_leaves_unrelated_urls_and_strings_untouched() -> None:
|
||||||
|
event = {"extra": {"public_url": "https://gendsgn.ru/trade-in", "region": "66"}}
|
||||||
|
out = redact_telegram_bot_token(event, {})
|
||||||
|
assert out is not None
|
||||||
|
assert out["extra"]["public_url"] == "https://gendsgn.ru/trade-in"
|
||||||
|
assert out["extra"]["region"] == "66"
|
||||||
|
|
||||||
|
|
||||||
|
def test_redact_token_handles_non_dict_event() -> None:
|
||||||
|
assert redact_telegram_bot_token(None, {}) is None # type: ignore[arg-type]
|
||||||
|
|
||||||
|
|
||||||
|
# Голая форма `<id>:<secret>` без префикса `/bot` — так токен выглядит в локали
|
||||||
|
# `token` конструктора TelegramClient. Штатно в event не попадает
|
||||||
|
# (include_local_variables=False в tgbot_main), но редактор обязан покрывать и
|
||||||
|
# этот случай: иначе рубеж против утечки токена всего один — флаг SDK.
|
||||||
|
_LEAKED_TOKEN_BARE = "8663867262:AAFakeSecretPartForTestsOnly1234567"
|
||||||
|
|
||||||
|
|
||||||
|
def test_redacts_bare_token_without_bot_prefix() -> None:
|
||||||
|
event = {
|
||||||
|
"logentry": {"message": f"init failed for {_LEAKED_TOKEN_BARE}"},
|
||||||
|
"exception": {
|
||||||
|
"values": [{"stacktrace": {"frames": [{"vars": {"token": _LEAKED_TOKEN_BARE}}]}}]
|
||||||
|
},
|
||||||
|
}
|
||||||
|
out = redact_telegram_bot_token(event, {})
|
||||||
|
assert out is not None
|
||||||
|
assert "AAFakeSecretPartForTestsOnly1234567" not in repr(out)
|
||||||
|
frame = out["exception"]["values"][0]["stacktrace"]["frames"][0]
|
||||||
|
assert frame["vars"]["token"] == "[REDACTED]"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"benign",
|
||||||
|
[
|
||||||
|
"chat_id:12345",
|
||||||
|
"ratio 3:4",
|
||||||
|
"2026-07-16T10:00:00",
|
||||||
|
"postgresql://user:pass@postgres:5432/tradein",
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_bare_token_redaction_leaves_benign_colon_strings_untouched(benign: str) -> None:
|
||||||
|
"""Голая регулярка не должна бить по любым `x:y` — иначе диагностика ослепнет."""
|
||||||
|
out = redact_telegram_bot_token({"logentry": {"message": benign}}, {})
|
||||||
|
assert out is not None
|
||||||
|
assert out["logentry"]["message"] == benign
|
||||||
|
|
||||||
|
|
||||||
|
def test_composed_before_send_scrubs_pii_and_token_together() -> None:
|
||||||
|
"""Композиция, реально используемая в `app.tgbot_main._before_send`: PII-scrub
|
||||||
|
(ключ-based) И token-redaction (regex full-text) применяются оба, не заменяя
|
||||||
|
друг друга — разные классы секретов, разные механизмы обнаружения."""
|
||||||
|
event = {
|
||||||
|
"request": {"data": {"client_phone": "+79991234567"}},
|
||||||
|
"exception": {
|
||||||
|
"values": [
|
||||||
|
{
|
||||||
|
"stacktrace": {
|
||||||
|
"frames": [{"vars": {"url": _LEAKED_TOKEN_URL}}],
|
||||||
|
}
|
||||||
|
}
|
||||||
|
]
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
def composed_before_send(evt, hint):
|
||||||
|
scrubbed = scrub_pii_event(evt, hint)
|
||||||
|
if scrubbed is None:
|
||||||
|
return None
|
||||||
|
return redact_telegram_bot_token(scrubbed, hint)
|
||||||
|
|
||||||
|
out = composed_before_send(event, {})
|
||||||
|
assert out is not None
|
||||||
|
assert out["request"]["data"]["client_phone"] == "[REDACTED]"
|
||||||
|
frame_url = out["exception"]["values"][0]["stacktrace"]["frames"][0]["vars"]["url"]
|
||||||
|
assert "8663867262:AAExampleSecretPartAbCdEf123" not in frame_url
|
||||||
|
|
|
||||||
|
|
@ -201,6 +201,49 @@ services:
|
||||||
networks:
|
networks:
|
||||||
- tradein-net
|
- tradein-net
|
||||||
|
|
||||||
|
# Telegram support-bot bridge: long-polling worker (asyncio/httpx), тот же
|
||||||
|
# backend-образ что и scraper/backend (один Dockerfile, другая команда).
|
||||||
|
# Ничего не слушает (long-polling исходящий к Telegram API) → без expose/ports.
|
||||||
|
# НЕ подписан на gendesign_shared — не проксируется Caddy, только tradein-net
|
||||||
|
# (нужен доступ к postgres). TELEGRAM_* переменные — ТОЛЬКО из backend/.env.runtime
|
||||||
|
# (env_file ниже); пусто/нет токена = бот молча не стартует (см. tgbot_main.py).
|
||||||
|
# ⚠️ НЕ дублировать TELEGRAM_* в блоке environment: (см. предупреждение у backend
|
||||||
|
# выше про environment: перекрывающий env_file при пустой host-env).
|
||||||
|
tgbot:
|
||||||
|
image: ghcr.io/lekss361/gendesign-tradein-backend:${IMAGE_TAG:-latest}
|
||||||
|
container_name: tradein-tgbot
|
||||||
|
# mem: нет прод-замера (новый сервис, 2026-07-16) — консервативный старт по
|
||||||
|
# аналогии с профилем scraper (лёгкая asyncio-оркестрация, без headless-браузера
|
||||||
|
# и без больших payload'ов). Пересмотреть после первого docker stats на проде.
|
||||||
|
mem_limit: 256m
|
||||||
|
memswap_limit: 256m
|
||||||
|
logging: *default-logging
|
||||||
|
command: ["python", "-m", "app.tgbot_main"]
|
||||||
|
env_file:
|
||||||
|
- path: ./backend/.env.runtime
|
||||||
|
required: false
|
||||||
|
environment:
|
||||||
|
DATABASE_URL: "postgresql+psycopg://${TRADEIN_POSTGRES_USER:-tradein}:${TRADEIN_POSTGRES_PASSWORD}@postgres:5432/tradein"
|
||||||
|
ENVIRONMENT: "production"
|
||||||
|
GLITCHTIP_DSN: "${GLITCHTIP_DSN:-}"
|
||||||
|
# Этот процесс — не scheduler_main; false на всякий случай, если общий
|
||||||
|
# Settings()-объект образа где-то читает флаг при импорте (defense-in-depth,
|
||||||
|
# аналогично backend). tgbot_main.py не должен зависеть от этого значения.
|
||||||
|
SCHEDULER_ENABLE: "false"
|
||||||
|
depends_on:
|
||||||
|
postgres:
|
||||||
|
condition: service_healthy
|
||||||
|
# unless-stopped (как и весь остальной стек): только always/unless-stopped
|
||||||
|
# гарантируют автозапуск после ребута VPS. Kill-switch при пустом
|
||||||
|
# TELEGRAM_BOT_TOKEN реализован idle-блокировкой в tgbot_main.py, а НЕ exit(0) —
|
||||||
|
# иначе эта политика дала бы рестарт-луп (перезапускает независимо от кода).
|
||||||
|
restart: unless-stopped
|
||||||
|
# >_DRAIN_TIMEOUT_S (100s) в tgbot_main.py, иначе docker убьёт по дефолтным 10с
|
||||||
|
# раньше, чем докрутится long-poll (до 30с), и кооперативный drain не сработает.
|
||||||
|
stop_grace_period: 120s
|
||||||
|
networks:
|
||||||
|
- tradein-net
|
||||||
|
|
||||||
frontend:
|
frontend:
|
||||||
image: ghcr.io/lekss361/gendesign-tradein-frontend:${IMAGE_TAG:-latest}
|
image: ghcr.io/lekss361/gendesign-tradein-frontend:${IMAGE_TAG:-latest}
|
||||||
container_name: tradein-frontend
|
container_name: tradein-frontend
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue