gendesign/tradein-mvp/backend/app/services/tgbot/bridge.py
bot-backend 8e7c65061b
Some checks failed
CI Trade-In / backend-tests (pull_request) Failing after 2m29s
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 11s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
fix(tgbot): honest H1 rejection, per-role H2 budget, M1/M2/L1 cleanup (#3471 review)
Deep review of PR #3494 found the rate limiter unusable as designed:
- H1: acquire() waited unbounded even for interactive HTTP handlers (support.py,
  glitchtip.py already pass a narrow `timeout` — reuse it as the queue wait cap
  instead of editing those handlers, which are out of scope here). New
  TelegramRateLimitedError (subclass of TelegramError) gives a fast, honest
  502 instead of hanging past the caller's own budget.
- H2: the limiter is per-process (in-memory), but two processes write to the
  same group (uvicorn API + bot worker) — giving each the same 18/min doubled
  the platform ceiling. Split into telegram_group_rate_limit_api_per_minute
  (12) and _bot_per_minute (6), sum kept below ~20.
- M1: bridge.py sends without an explicit timeout inherited "wait forever",
  stalling the single-threaded poll loop (open DB session) past the SIGTERM
  drain window. Bounded via rate_limit_max_wait=20s at the six call sites.
- M2: _locks/_sent_at grew unbounded on every unique DM chat_id. Added
  opportunistic cleanup of fully-expired entries.
- L1: the "queue full" warning now logs once per acquire() call, not once
  per sleep iteration.
- Corrected a factual error in the docstring: TELEGRAM_SUPPORT_CHAT_ID and
  TELEGRAM_ALERTS_CHAT_ID are the SAME group on prod (topics differ only).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG
2026-09-12 15:28:16 +03:00

995 lines
60 KiB
Python
Raw Permalink 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, #tgsupport-web).
Поток:
A) Клиент пишет боту в личку (chat.type == 'private') →
upsert tg_support_users → (если первое сообщение за последний час — шапка
с идентификацией клиента в топик) → copyMessage контента в support-топик →
запись в tg_support_messages (direction='in', topic_message_id — ключ
маршрутизации ответа).
A') Пользователь сайта пишет через `app.api.v1.support` (веб-чат поддержки,
#tgsupport-web) → тот эндпоинт САМ зеркалит sendMessage'ом в топик и пишет
web_support_messages (direction='in') — этот модуль в этой ветке не участвует,
только в разборе ответа (B ниже).
B) Оператор отвечает РЕПЛАЕМ в support-группе на зеркало клиента →
резолвим topic_message_id ОБЕ стороны (tg_support_messages И
web_support_messages), скоупя к ТЕКУЩЕМУ TELEGRAM_SUPPORT_CHAT_ID
(#tgsupport-web review M1 — если группу когда-нибудь сменят/пересоздадут,
Telegram message_id стартует заново и может совпасть со старым числом из
другой таблицы; без скоупинга это была бы ТИХАЯ доставка постороннему
клиенту). Совпадение НА ОБЕИХ сторонах одновременно — громкий отказ
(`logger.error`, ничего не доставляем) вместо произвольного выбора одной из
них. Иначе: chat_id найден → copyMessage ответа в личку клиента → запись
(direction='out'); thread_id найден (веб-зеркало) → доставка идёт НЕ в
Telegram (у веб-клиента нет личного чата с ботом), а записью direction='out'
в web_support_messages (веб-фронт вычитывает её обычным polling'ом); реплай
медиа-типом на веб-зеркало — веб-чат текстовый MVP, доставка целиком
отклоняется (не частично — фото с подписью НЕ превращается в "ответ = только
подпись"), оператор получает уведомление в топике (review M2). Реплай не на
зеркало (или не реплай вообще) — обычная болтовня в топике, тихий игнор.
Telegram 403 (клиент заблокировал бота) → is_blocked=true + уведомление в
топике (только для Telegram-ветки — у веб-клиента нет "заблокировал бота").
C) Дедуп: update_id <= сохранённого offset — skip. Offset сохраняется И
коммитится в той же транзакции, что и запись сообщения (см. `process_update`),
после КАЖДОГО апдейта — рестарт воркера не переигрывает уже обработанные
апдейты и не подвисает вечно на «ядовитом» апдейте. Исключение —
ТРАНЗИЕНТНЫЙ сетевой отказ (`TelegramNetworkError`, т.е. исчерпанный бюджет
ретраев клиента): такой апдейт СОЗНАТЕЛЬНО остаётся неподтверждённым, чтобы
Telegram отдал его снова, иначе ответ оператора пропадал бы навсегда
(#tg-connection-resilience). Потолок переигрываний — `_MAX_NETWORK_REPLAYS`.
D) TELEGRAM_BOT_TOKEN пуст → бот выключен — проверяется в `app.tgbot_main`
(entrypoint), не здесь.
E) /start клиенту → короткое приветствие МЕРЫ, без зеркалирования в топик
(команда — не содержательное обращение, не должна засорять топик).
F) Флуд-лимит на отправителя (низкий приоритет, per-chat_id): воркер
long-polling однопоточный и обрабатывает апдейты СТРОГО последовательно, а
Telegram ограничивает саму support-группу ~20 сообщениями/минуту — ОДНИМ
бюджетом на ВСЕХ клиентов разом (зеркала + шапки + ответы оператора).
Превышение — 429 с ожиданием 30-60с, на которые воркер не может обработать
НИ ОДНОГО следующего апдейта — один флудящий клиент подвешивает доставку
всем остальным. `_flood_limiter` (тот же `SlidingWindowLimiter`, что и
веб-чат поддержки, ключ — TELEGRAM chat_id) режет per-sender поток заметно
ниже группового лимита; сообщения сверх бюджета НЕ зеркалируются (иначе
сам факт мирроринга уже съедает групповой бюджет, который мы и защищаем) и
НЕ пишутся в tg_support_messages (нечего маршрутизировать без
topic_message_id). Клиент получает уведомление, что сообщение НЕ
доставлено (молчать нельзя — иначе клиент решит, что оператор его получил),
но не чаще ОДНОГО РАЗА за то же окно (`_flood_notify_limiter`, limit=1) —
иначе само уведомление стало бы вторым источником флуда.
Персистентность вынесена за `BridgeStorage`-протокол — маршрутизирующая логика
(`process_update` и приватные `_handle_*`) не завязана на реальную БД, тестируется
на in-memory fake storage + mock httpx (см. tests/services/tgbot/). Веб-чат
таблицы (web_support_threads/web_support_messages) сознательно ОТДЕЛЬНЫ от
tg_support_* — обоснование в data/sql/187_web_support_chat.sql; здесь `BridgeStorage`
несёт два дополнительных метода (`find_web_thread_by_topic_message`,
`record_web_out_message`), делегирующих в `web_support_storage`.
"""
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.ratelimit import SlidingWindowLimiter
from app.core.shutdown import shutdown_requested
from app.services.tgbot import web_support_storage
from app.services.tgbot.client import TelegramApiError, TelegramClient, TelegramNetworkError
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")
# (низкий приоритет, флуд-защита) — см. пункт F) в докстринге модуля. Порог
# НАМЕРЕННО заметно ниже группового лимита Telegram (~20 msg/min): бюджет
# делится с шапками-идентификациями и ответами оператора, и с другими
# одновременными клиентами — щедрый лимит одного отправителя всё равно упёрся
# бы в общий групповой 429. Тот же примитив, что и веб-чат поддержки
# (app/api/v1/support.py `_send_limiter`), ключ здесь — TELEGRAM chat_id
# отправителя (не username — у Telegram-клиента username может отсутствовать).
_FLOOD_LIMIT = 5
_FLOOD_WINDOW_S = 60.0
_flood_limiter = SlidingWindowLimiter(limit=_FLOOD_LIMIT, window_s=_FLOOD_WINDOW_S)
# Уведомление о флуде — не чаще ОДНОГО раза за то же окно, иначе само
# уведомление стало бы вторым источником флуда. Отдельный лимитер с limit=1 на
# то же окно: `check()` возвращает None (и фиксирует попытку) ровно один раз за
# окно, дальше молчит до его истечения — без отдельной структуры "когда в
# последний раз уведомляли".
_flood_notify_limiter = SlidingWindowLimiter(limit=1, window_s=_FLOOD_WINDOW_S)
FLOOD_LIMITED_TEXT = (
"Сообщение не доставлено — вы отправляете сообщения слишком часто. "
"Пожалуйста, подождите немного и напишите ещё раз."
)
# #tgsupport-web review M2: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ)
# на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор
# получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил).
_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT = "Веб-чат поддерживает только текст, сообщение не доставлено."
# Уведомления оператору в топике отправляются ИЗНУТРИ poll loop, который
# однопоточный: пока висит одна отправка, не обрабатывается НИ ОДИН следующий
# апдейт. Поэтому им нужен интерактивный бюджет, а не воркерный дефолт клиента
# (5 ретраев, backoff до 30с, полный retry_after на 429 — минуты стопа на
# ВТОРИЧНОМ действии). Числа — те же, что у интерактивных отправок веб-чата
# поддержки (`_INTERACTIVE_SEND_*` в app/api/v1/support.py, подобраны замером
# прода #tgsupport-retry); сознательно ДУБЛИРУЕМ, а не импортируем из слоя API —
# воркер не должен зависеть от роутера.
_NOTIFY_SEND_TIMEOUT_S = 5.0
_NOTIFY_SEND_MAX_RETRIES = 3
_NOTIFY_SEND_MAX_BACKOFF_S = 1.0
# Потолок ожидания в TelegramGroupRateLimiter для отправок ИЗ ЭТОГО модуля,
# которые не проходят через _notify_topic (review M1, #3471). Без него
# `client.send_message`/`copy_message` без явного `timeout` наследуют
# лимитер-политику "без потолка" (см. TelegramClient._request) — приемлемую
# ДЛЯ ФОНОВОЙ отправки как таковой, но НЕ здесь: эти вызовы идут внутри
# `run_poll_loop`, который на каждый апдейт держит ОТКРЫТУЮ сессию БД
# (SessionLocal, см. вызывающих) и однопоточно блокирует опрос СЛЕДУЮЩИХ
# апдейтов — минутный сон здесь стопорит и БД-соединение, и весь мост, а не
# только одно сообщение. Второй повод — SIGTERM drain: `tgbot_main._DRAIN_TIMEOUT_S`
# даёт 100с на завершение текущей итерации; ожидание слота дольше этого
# бюджета уже не успевает подчиниться cooperative drain. 20с — заметно больше,
# чем разумная очередь при исчерпанном лимите (окно 60с, обычно секунды), но
# заметно МЕНЬШЕ минуты и вписывается в drain-бюджет с запасом.
_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S = 20.0
# Потолок переигрываний ОДНОГО update_id на транзиентных сетевых отказах
# (#tg-connection-resilience).
#
# Зачем потолок: без него «вечно недоставляемый» апдейт заклинил бы очередь
# НАВСЕГДА — ровно то, от чего защищал прежний безусловный `finally: save_offset`.
# Потерять одно сообщение плохо, потерять весь поток — хуже, поэтому на потолке
# offset всё-таки сдвигается, но ГРОМКО (`logger.error`), а не молча.
#
# Про дубли: `TelegramNetworkError` означает исчерпанный бюджет ретраев клиента,
# при этом запрос МОГ дойти до Telegram (потерялся ответ) — переигрывание тогда
# доставит клиенту то же сообщение второй раз. Это осознанный at-least-once
# компромисс: дубль и клиент, и оператор ВИДЯТ и могут поправить, а тихая потеря
# ответа не оставляет следа нигде, кроме строчки в логе. Полноценная
# идемпотентность по (update_id, target_chat_id) потребовала бы нового
# персистентного состояния (колонка/таблица + миграция) ради редкого случая;
# вместо этого число возможных дублей жёстко ограничено сверху — не больше
# (_MAX_NETWORK_REPLAYS - 1) повторов на апдейт.
_MAX_NETWORK_REPLAYS = 3
# update_id -> сколько раз мы уже отказались подтверждать этот апдейт.
# In-memory осознанно: воркер long-polling однопоточный, запись живёт ровно до
# подтверждения апдейта (`pop` в `process_update`), так что словарь не растёт.
# Рестарт воркера обнуляет счётчик — это допустимо (новый процесс = новая сеть),
# потолок всё равно действует в пределах каждой жизни процесса.
_network_replay_attempts: dict[int, int] = {}
# ── Storage abstraction (testable без реальной БД) ──────────────────────────
class BridgeStorage(Protocol):
"""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,
support_chat_id: int | None = None,
) -> int | None: ...
def find_chat_by_topic_message(
self, topic_message_id: int, support_chat_id: int
) -> int | None: ...
def mark_blocked(self, chat_id: int) -> None: ...
def find_web_thread_by_topic_message(
self, topic_message_id: int, support_chat_id: int
) -> int | None: ...
def record_web_out_message(
self,
*,
thread_id: int,
text_body: str,
operator_tg_id: int | None,
topic_message_id: int | None = None,
support_chat_id: int | None = None,
) -> 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,
support_chat_id: int | None = 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, support_chat_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),
CAST(:support_chat_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,
"support_chat_id": support_chat_id,
},
).fetchone()
return int(row[0]) if row is not None else None
def find_chat_by_topic_message(self, topic_message_id: int, support_chat_id: int) -> int | None:
"""Скоупим к ТЕКУЩЕМУ support_chat_id (#tgsupport-web review M1) — строка
со ЧУЖИМ (не NULL, не текущим) support_chat_id — исторический артефакт
ротации support-группы, не валидный маршрут сегодня. NULL (строки до
миграции 188, если есть) — лениентный wildcard-матч (единственный
действовавший чат на тот момент).
БЕЗ фильтра по direction (#3471 P0, было `AND direction = 'in'`): с тех
пор как `record_message` на исходящем ответе тоже сохраняет
`topic_message_id` (id сообщения оператора В ТОПИКЕ), реплай оператора
на СВОЙ предыдущий ответ обязан резолвиться так же, как реплай на
зеркало клиента — иначе продолжение диалога без повторного цитирования
клиента тихо проваливалось в orphan-check."""
row = self._db.execute(
text(
"""
SELECT chat_id
FROM tg_support_messages
WHERE topic_message_id = CAST(:topic_message_id AS bigint)
AND (support_chat_id = CAST(:support_chat_id AS bigint)
OR support_chat_id IS NULL)
ORDER BY created_at DESC
LIMIT 1
"""
),
{"topic_message_id": topic_message_id, "support_chat_id": support_chat_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},
)
def find_web_thread_by_topic_message(
self, topic_message_id: int, support_chat_id: int
) -> int | None:
"""Делегирует в `web_support_storage` (#tgsupport-web) — то же соединение/
транзакцию, что и tg-путь, коммитится вместе offset'ом в `process_update`."""
return web_support_storage.find_thread_by_topic_message(
self._db, topic_message_id, support_chat_id
)
def record_web_out_message(
self,
*,
thread_id: int,
text_body: str,
operator_tg_id: int | None,
topic_message_id: int | None = None,
support_chat_id: int | None = None,
) -> None:
web_support_storage.record_outbound(
self._db,
thread_id=thread_id,
text_body=text_body,
operator_tg_id=operator_tg_id,
topic_message_id=topic_message_id,
support_chat_id=support_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})"
async def _notify_topic(
client: TelegramClient,
*,
text: str,
reply_to_message_id: int | None,
context: str,
) -> bool:
"""Служебное уведомление оператору в support-топик (вторичное действие).
Два свойства, которых не было у прямых `client.send_message` вызовов:
1) узкий интерактивный бюджет (`_NOTIFY_SEND_*`) — иначе одна такая отправка
стопорит весь однопоточный poll loop на минуты;
2) собственный `except` — провал УВЕДОМЛЕНИЯ не отменяет основную ветку
обработки (клиент уже помечен заблокированным / медиа-реплай уже отклонён)
и не решает судьбу апдейта.
`context` — только технические идентификаторы (chat_id/thread_id), НЕ текст
переписки: логи моста принципиально не содержат ПДн.
Возвращает True, если уведомление реально ушло, False — если само
уведомление тоже упало (напр. Telegram недоступен). Вызывающий, для
которого проваленное уведомление означает ПОЛНУЮ тишину (ни клиенту, ни
оператору), обязан на False залогировать `logger.error` с идентификаторами
(#3471 P0) — иначе единственный след остаётся только в этом WARNING.
"""
try:
await client.send_message(
chat_id=settings.telegram_support_chat_id,
text=text,
message_thread_id=settings.telegram_support_topic_id or None,
reply_to_message_id=reply_to_message_id,
timeout=_NOTIFY_SEND_TIMEOUT_S,
max_retries=_NOTIFY_SEND_MAX_RETRIES,
max_backoff=_NOTIFY_SEND_MAX_BACKOFF_S,
)
return True
except Exception:
logger.warning(
"tgbot bridge: не удалось отправить уведомление оператору в топик (%s) — "
"основная ветка обработки не отменяется",
context,
exc_info=True,
)
return False
# ── 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,
rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S,
)
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,
rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S,
)
return
message_id = message.get("message_id")
if not isinstance(message_id, int):
logger.warning("tgbot bridge: приватное сообщение без message_id — игнор")
return
# F) Флуд-лимит на отправителя — peek БЕЗ расхода бюджета (тот же паттерн,
# что `_send_limiter` в app/api/v1/support.py: под лимитом ниже сразу
# `.record()`-им попытку). Над лимитом — НЕ зеркалируем (иначе сам мирроринг
# уже съедает групповой Telegram-бюджет, который лимит и защищает) и НЕ
# пишем в tg_support_messages (без topic_message_id маршрутизировать ответ
# всё равно нечего).
flood_key = str(chat_id) # SlidingWindowLimiter — ключ str (см. app/core/ratelimit.py)
if _flood_limiter.retry_after(flood_key) is not None:
logger.warning(
"tgbot bridge: chat_id=%d превысил флуд-лимит (%d msg/%.0fs) — "
"сообщение НЕ зеркалируется в топик (защита группового Telegram-лимита)",
chat_id,
_FLOOD_LIMIT,
_FLOOD_WINDOW_S,
)
# Уведомляем клиента, что сообщение НЕ доставлено (молчать нельзя —
# иначе клиент решит, что оператор его получил), но не чаще одного раза
# за окно — `_flood_notify_limiter.check()` возвращает None (и сам
# фиксирует попытку) ровно один раз за окно.
if _flood_notify_limiter.check(flood_key) is None:
await client.send_message(
chat_id=chat_id,
text=FLOOD_LIMITED_TEXT,
rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S,
)
return
_flood_limiter.record(flood_key)
# Шапка — только на первое сообщение клиента за окно, иначе топик засоряется.
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,
rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S,
)
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,
rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S,
)
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,
support_chat_id=settings.telegram_support_chat_id,
)
async def _handle_group_reply(
message: dict[str, Any], client: TelegramClient, storage: BridgeStorage
) -> None:
"""B) Реплай оператора в support-группе → доставка ответа клиенту (Telegram
ЛИБО веб-чат, #tgsupport-web — см. модульный docstring)."""
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
# #tgsupport-web review M1: резолвим ОБЕ стороны с текущим support_chat_id
# (НЕ short-circuit на первом найденном) — если topic_message_id совпал в
# ОБЕИХ таблицах одновременно, это значит инвариант "уникален в пределах
# текущей support-группы" нарушен (баг/ручная правка данных) — отказываем в
# доставке ГРОМКО, вместо того чтобы молча выбрать tg-путь и отправить ответ
# постороннему Telegram-клиенту (152-ФЗ misroute risk).
current_chat_id = settings.telegram_support_chat_id
target_chat_id = storage.find_chat_by_topic_message(mirror_message_id, current_chat_id)
web_thread_id = storage.find_web_thread_by_topic_message(mirror_message_id, current_chat_id)
if target_chat_id is not None and web_thread_id is not None:
logger.error(
"tgbot bridge: topic_message_id=%d резолвится ОДНОВРЕМЕННО в Telegram "
"(chat_id=%d) и веб-чат (thread_id=%d) под support_chat_id=%d — отказ в "
"доставке, требуется ручной разбор tg_support_messages/web_support_messages",
mirror_message_id,
target_chat_id,
web_thread_id,
current_chat_id,
)
return
if target_chat_id is not None:
# Существующий Telegram-путь — НЕ ТРОНУТ.
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,
rate_limit_max_wait=_BRIDGE_SEND_RATE_LIMIT_MAX_WAIT_S,
)
except TelegramApiError as exc:
if exc.error_code == 403:
# Клиент заблокировал бота — фиксируем и уведомляем оператора в топике.
storage.mark_blocked(target_chat_id)
await _notify_topic(
client,
text=(
f"Не удалось доставить сообщение клиенту (chat_id={target_chat_id}) — "
"бот заблокирован."
),
reply_to_message_id=message_id,
context=f"403 на доставке клиенту chat_id={target_chat_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,
# #3471 P0: id ЭТОГО сообщения оператора В ТОПИКЕ (было безусловно
# None) — без него реплай оператора на СВОЙ предыдущий ответ не
# резолвился (искать было нечего), маршрут держался только на
# зеркале клиента. Совпадение с `in`-записью структурно исключено:
# `message_id` — id реплая оператора, а зеркало клиента уже занимает
# другой message_id в том же чате.
topic_message_id=message_id,
kind=_infer_kind(message),
text_body=message.get("text") or message.get("caption"),
operator_tg_id=operator_id,
# Deep review PR #3479: без этого out-строка была бы вечным
# wildcard для `find_chat_by_topic_message` (матчит support_chat_id
# IS NULL под ЛЮБЫМ текущим чатом) — при ротации support-группы
# (188) новый message_id мог бы совпасть со старой out-строкой и
# увести ответ ЧУЖОМУ клиенту. Симметрично in-ветке выше (строка ~601).
support_chat_id=settings.telegram_support_chat_id,
)
return
if web_thread_id is not None:
message_id = message.get("message_id")
kind = _infer_kind(message)
if kind != "text":
# #tgsupport-web review M2: НЕ доставляем частично (фото С ПОДПИСЬЮ
# молча превратилось бы в "ответ = только текст подписи", клиент решил
# бы что это весь ответ) — отказ целиком + явное уведомление оператору
# в топике (тот же паттерн, что 403-уведомление выше), иначе оператор
# уверен, что ответ доставлен, хотя веб-чат не поддерживает медиа.
logger.warning(
"tgbot bridge: реплай на веб-зеркало (thread_id=%d) содержит %s, "
"не текст — веб-чат поддерживает только текст, доставка отклонена",
web_thread_id,
kind,
)
await _notify_topic(
client,
text=_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT,
reply_to_message_id=message_id if isinstance(message_id, int) else None,
context=f"медиа-реплай на веб-зеркало thread_id={web_thread_id}",
)
return
text_body = message.get("text")
if not text_body:
# Текстовый kind, но пустой text (защитный edge case) — нечего доставлять.
return
operator = message.get("from") or {}
operator_id = operator.get("id")
try:
storage.record_web_out_message(
thread_id=web_thread_id,
text_body=text_body,
operator_tg_id=operator_id,
# #3471 P0: id ЭТОГО сообщения оператора в топике — без него
# реплай оператора на СВОЙ предыдущий веб-ответ не резолвится
# (см. `find_thread_by_topic_message`, direction-фильтр снят).
topic_message_id=message_id if isinstance(message_id, int) else None,
# Deep review PR #3479: БЕЗ этого out-строка писалась бы с
# support_chat_id=NULL — `find_thread_by_topic_message` матчит
# NULL под ЛЮБЫМ текущим чатом (лениентный wildcard для легаси
# строк до 187/188), т.е. каждая out-строка стала бы вечным
# wildcard. При ротации support-группы новый message_id мог бы
# совпасть со старой out-строкой и увести ответ в ЧУЖОЙ тред —
# ровно то, от чего защищала скоупинг-миграция 187/188.
support_chat_id=settings.telegram_support_chat_id,
)
except SQLAlchemyError:
# #3471 P0: для веб-треда ЭТА запись — и есть доставка клиенту (веб-
# фронт вычитывает ответ обычным polling'ом web_support_messages).
# Откат без уведомления означал бы: оператор уверен, что ответил,
# клиент ждёт молча. Offset ниже всё равно сдвигается — НЕ потому,
# что апдейт "частично применён в Telegram" (в этой ветке до сбоя в
# Telegram ничего не уходило вообще: сам реплай оператора Telegram
# уже полностью доставил ДО того, как мы начали его разбирать,
# ретраить на стороне площадки нечего), а потому что действует общая
# политика `process_update` для `SQLAlchemyError` — сбой БД не
# переигрывается (в отличие от `TelegramNetworkError`), а
# сигнализируется громко; здесь это explicit-просьба оператору
# прислать ответ заново — human-in-the-loop retry вместо
# технического. rollback() ОБЯЗАН отработать ДО уведомления —
# сессия в failed-transaction state, а `_notify_topic` шлёт через
# `client`, не через `storage`, поэтому сам rollback тут не нужен для
# отправки, но нужен, чтобы process_update дальше не упал на
# save_offset/commit тем же PendingRollbackError (см. #3 review).
storage.rollback()
notified = await _notify_topic(
client,
text=(
"Не удалось сохранить ваш ответ из-за сбоя базы данных — клиенту "
"он НЕ доставлен. Пожалуйста, отправьте ответ ещё раз."
),
reply_to_message_id=message_id if isinstance(message_id, int) else None,
context=f"сбой БД на доставке веб-ответа thread_id={web_thread_id}",
)
if not notified:
# Оба канала молчат (БД и уведомление) — единственный след,
# который останется, это эта строка. Идентификаторы, НЕ текст
# (ПДн в лог не идёт) — по ним человек найдёт ответ оператора в
# топике и перешлёт его руками (#3471 P0).
logger.error(
"tgbot bridge: сбой БД на веб-ответе И не удалось уведомить "
"оператора (thread_id=%d, message_id=%s) — ответ клиенту "
"потерян молча, требуется ручной разбор support-топика",
web_thread_id,
message_id,
)
return
# Обычная болтовня в топике (реплай на чьё-то ещё сообщение) — не логируем,
# это ожидаемый шум. НО реплай на сообщение, отправленное САМИМ БОТОМ
# (is_bot=True) и при этом отсутствующее ни в tg_support_messages, ни в
# web_support_messages — подозрительно: вероятная причина — осиротевшее
# зеркало (воркер/API упал МЕЖДУ отправкой зеркала и commit'ом записи в БД).
# Дискриминатор неидеальный (шапка-идентификация тоже от бота, но не
# 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, ни в web_support_messages как зеркало — возможно, "
"осиротевшее зеркало (крах между отправкой и commit'ом) или "
"шапка-идентификация; ответ оператора НЕ доставлен",
mirror_message_id,
)
async def process_update(
update: dict[str, Any], client: TelegramClient, storage: BridgeStorage
) -> bool:
"""Маршрутизирует один Telegram update. Дедуп (C) + атомарный offset-commit.
Возвращает True, если offset сдвинут (апдейт подтверждён, Telegram его больше
не отдаст), и False, если апдейт СОЗНАТЕЛЬНО оставлен неподтверждённым ради
переигрывания. На False вызывающий (`run_poll_loop`) ОБЯЗАН прервать разбор
пачки: offset у Telegram — единая «высшая отметка», подтверждение любого
СЛЕДУЮЩЕГО апдейта неявно подтвердило бы и этот, и переигрывания не было бы.
Дедуп: update_id <= сохранённого offset — skip без side-effects.
Судьба offset'а по классам отказа:
- `TelegramNetworkError` (транзиентный: Telegram не ответил, бюджет ретраев
клиента исчерпан) — offset НЕ двигаем, `rollback()` частичных записей,
апдейт переигрывается на следующей итерации. Иначе ответ оператора
терялся НАВСЕГДА: copyMessage не дошёл, `record_message` не выполнился,
Telegram апдейт больше не отдаст, а оператор уверен, что ответил
(#tg-connection-resilience). Ограничено `_MAX_NETWORK_REPLAYS` — на
потолке offset всё-таки сдвигается с `logger.error`, иначе «вечно
недоставляемый» апдейт заклинил бы поток навсегда.
- `SQLAlchemyError` (сбой БД) — offset двигаем, но `rollback()` ОБЯЗАН
отработать ПЕРЕД `save_offset`: сбой БД оставляет сессию в
failed-transaction state, иначе `save_offset` сам кинет
`PendingRollbackError`, `process_update` вылетит без сохранения offset'а,
следующая итерация получит СТАРЫЙ offset от `get_offset()` и переиграет
тот же апдейт — copyMessage задублирует зеркало клиента в топике на
каждый повтор поллинга (#3 review, воспроизведено). Для веб-ветки
(`_handle_group_reply` → `record_web_out_message`) этот `SQLAlchemyError`
перехватывается ЛОКАЛЬНО, до этого места: там запись в БД И ЕСТЬ
доставка клиенту, поэтому rollback сопровождается уведомлением оператору
в топике, что ответ НЕ доставлен (#3471 P0) — сюда, на верхний уровень,
это исключение уже не долетает.
- любое прочее исключение (в т.ч. `TelegramApiError` — площадка ОТВЕТИЛА
отказом, повтор ничего не изменит) — offset двигаем, «ядовитый» апдейт
не блокирует поток.
"""
update_id = update.get("update_id")
if not isinstance(update_id, int):
logger.warning("tgbot bridge: update без валидного update_id — игнор")
return True
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 True
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 TelegramNetworkError:
attempts = _network_replay_attempts.get(update_id, 0) + 1
# Частичные записи этого апдейта не должны уехать в БД чужим commit'ом
# (сессия одна на всю пачку) — переигрывание начинается с чистого листа.
storage.rollback()
if attempts < _MAX_NETWORK_REPLAYS:
_network_replay_attempts[update_id] = attempts
logger.warning(
"tgbot bridge: Telegram недоступен на update_id=%d (отказ %d из %d) — "
"offset НЕ сдвигаем, апдейт переиграется на следующей итерации",
update_id,
attempts,
_MAX_NETWORK_REPLAYS,
)
return False
logger.error(
"tgbot bridge: update_id=%d исчерпал потолок переигрываний (%d сетевых "
"отказов подряд) — сдвигаем offset, содержимое апдейта ПОТЕРЯНО; поток не "
"блокируем, требуется ручной разбор support-топика "
"(chat_id=%s, message_id=%s)",
update_id,
_MAX_NETWORK_REPLAYS,
# Идентификаторы, а НЕ текст: это единственная строка, по которой
# человек найдёт потерянный ответ оператора в топике и перешлёт его
# руками. Без них в логе остаётся только update_id, которого в
# интерфейсе Telegram не видно. Текст сообщения — ПДн, в лог не идёт.
(message or {}).get("chat", {}).get("id") if isinstance(message, dict) else None,
(message or {}).get("message_id") if isinstance(message, dict) else None,
)
except SQLAlchemyError:
logger.exception(
"tgbot bridge: DB-ошибка на update_id=%d — rollback перед сохранением "
"offset (иначе save_offset сам упадёт на failed-transaction state)",
update_id,
)
storage.rollback()
except Exception:
logger.exception(
"tgbot bridge: обработка update_id=%d упала — offset всё равно сдвигаем "
"(не блокируем поток на 'ядовитом' апдейте)",
update_id,
)
# Апдейт подтверждён — счётчик переигрываний больше не нужен (словарь не растёт).
_network_replay_attempts.pop(update_id, None)
storage.save_offset(update_id)
storage.commit()
return True
# ── 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 idx, update in enumerate(updates):
if not await process_update(update, client, storage):
# Апдейт намеренно не подтверждён (транзиентный сетевой
# отказ). Обрабатывать остаток пачки НЕЛЬЗЯ: offset —
# единая «высшая отметка», подтверждение следующего
# апдейта неявно подтвердило бы и этот. Остаток Telegram
# отдаст заново на следующей итерации.
logger.warning(
"tgbot bridge: update_id=%s не подтверждён — остаток пачки "
"(%d апдейтов) разберём на следующей итерации",
update.get("update_id"),
len(updates) - idx - 1,
)
break
consecutive_errors = 0
except Exception:
consecutive_errors += 1
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)")