gendesign/tradein-mvp/backend/app/services/tgbot/bridge.py
bot-backend b579fa4ced
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m4s
feat(tradein/tgbot): Telegram support-мост @MERAsupport_bot
Клиент пишет боту в личку → воркер зеркалит сообщение через copyMessage
в топик супергруппы-форума → оператор отвечает реплаем на зеркало → бот
доставляет ответ клиенту. Полный лог переписки в Postgres.

Отдельный контейнер на long-polling, а не webhook в tradein-backend:
не нужно пробивать дырку в auth-middleware (_PUBLIC_PATHS, #2213) и
маршрут в Caddy, нулевая внешняя поверхность, падение бота не задевает API.
Без aiogram — httpx уже в зависимостях, нужны только getUpdates/copyMessage.

Маршрутизация ответа — по topic_message_id: message_id в Telegram уникален
в пределах чата сквозь все топики, а все зеркала лежат в одном support-чате,
поэтому спутать адресата нельзя. Реплай на шапку/на ответ другого оператора
не резолвится (у direction='out' topic_message_id IS NULL) → тихий игнор.

Безопасность (найдено ревью, воспроизведено эмпирически):
- токен Telegram живёт в PATH URL, поэтому sanitize_url его не режет;
  утекал в GlitchTip через locals стек-фреймов (include_local_variables
  по умолчанию True) и через span data HttpxIntegration. Закрыто
  include_local_variables=False + regex-редактор в before_send (обе формы:
  /bot<id>:<secret> и голая <id>:<secret>), поверх существующего PII-scrub.
- httpx-логгер печатает полный URL на INFO → боевой токен уходил бы в
  docker logs каждые 30с. Приглушён до WARNING.

Надёжность:
- kill-switch при пустом токене — idle-блокировка, не exit(0): при
  restart: unless-stopped выход с любым кодом даёт рестарт-луп.
  unless-stopped выбран сознательно — только он гарантирует автозапуск
  после ребута VPS.
- stop_grace_period: 120s — дефолтные 10с убивали бы контейнер раньше,
  чем докрутится long-poll (30с) и отработает drain (100с).
- сбой SQL теперь ловится отдельно и делает rollback перед сдвигом offset:
  иначе сессия в failed-transaction не давала сохранить offset, апдейт
  переигрывался и зеркалился в топик по кругу.

152-ФЗ: переписка — ПДн, ON DELETE CASCADE по chat_id, удаление клиента
одним DELETE. Ретенция — follow-up.

Бот не включается автоматически: TELEGRAM_* задаются в runtime-env на VPS,
без них воркер штатно висит в idle. Порядок — в DEPLOY.md.

Тесты: 51 passed (маршрутизация обоих направлений, дедуп, 403→is_blocked,
throttle-окно шапки, redaction токена во всех формах event).
2026-07-16 16:58:53 +03:00

537 lines
25 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Маршрутизация 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)")