H1 (critical): reorder send_support_message — get_or_create_thread now runs AFTER a successful Telegram send, not before. Prod runs a single uvicorn process with no --workers on a sync SQLAlchemy engine sharing one event loop; holding a row-lock across the Telegram call (which can legally take minutes on 429/5xx worker-grade retries) risked freezing the entire API on a second concurrent request from the same user. send_message now also accepts explicit timeout/max_retries so the interactive endpoint uses a bounded budget instead of inheriting the long-polling worker's retry policy. M1: web_support_messages and tg_support_messages both gain a support_chat_id column (188 migration for the already-applied tg_support_messages table; 187 edited in place since it hasn't shipped yet). bridge._handle_group_reply now resolves BOTH the Telegram and web candidate under the CURRENT support_chat_id and refuses delivery loudly (logger.error) if both match, instead of silently preferring the Telegram path — closing a cross-table misroute risk that would surface if the support group is ever recreated. M2: a media reply to a web-mirrored message (including photos with a caption) is now refused in full with an explicit notice back in the topic, instead of silently delivering just the caption text or logging a WARNING nobody sees. M3: list_support_messages/get_support_unread/mark_support_read switched from async def to def — their bodies are pure sync psycopg calls; running them as async def executed blocking DB work directly on the event loop. M5: list_messages now takes a bounded LIMIT (last N, chronological) so a long-lived thread doesn't return its entire history on every poll. L1-L4: warn when Telegram doesn't return a message_id (routing dead-end), rate limit is peeked before send and only recorded on success (failed attempts no longer burn the budget), SlidingWindowLimiter gained the same empty-bucket cleanup as RateLimitMiddleware, and _require_username now strips whitespace so a proxy-injected space can't fork a second thread. L5: removed 187 from _manifest_applied.txt — the deploy pipeline tracks applied migrations via _schema_migrations, not this file, and keeping an unmerged migration name out of it preserves the option to rename before merge without tripping the "can't rename applied migrations" test.
665 lines
33 KiB
Python
665 lines
33 KiB
Python
"""Маршрутизация 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`
|
||
`finally`), после КАЖДОГО апдейта — рестарт воркера не переигрывает уже
|
||
обработанные апдейты и не подвисает вечно на «ядовитом» апдейте.
|
||
D) TELEGRAM_BOT_TOKEN пуст → бот выключен — проверяется в `app.tgbot_main`
|
||
(entrypoint), не здесь.
|
||
E) /start клиенту → короткое приветствие МЕРЫ, без зеркалирования в топик
|
||
(команда — не содержательное обращение, не должна засорять топик).
|
||
|
||
Персистентность вынесена за `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.shutdown import shutdown_requested
|
||
from app.services.tgbot import web_support_storage
|
||
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")
|
||
|
||
# #tgsupport-web review M2: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ)
|
||
# на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор
|
||
# получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил).
|
||
_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT = "Веб-чат поддерживает только текст, сообщение не доставлено."
|
||
|
||
|
||
# ── 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
|
||
) -> 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-матч (единственный
|
||
действовавший чат на тот момент)."""
|
||
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'
|
||
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
|
||
) -> None:
|
||
web_support_storage.record_outbound(
|
||
self._db,
|
||
thread_id=thread_id,
|
||
text_body=text_body,
|
||
operator_tg_id=operator_tg_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,
|
||
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,
|
||
)
|
||
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,
|
||
)
|
||
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 client.send_message(
|
||
chat_id=settings.telegram_support_chat_id,
|
||
text=_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT,
|
||
message_thread_id=settings.telegram_support_topic_id or None,
|
||
reply_to_message_id=message_id if isinstance(message_id, int) else None,
|
||
)
|
||
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")
|
||
storage.record_web_out_message(
|
||
thread_id=web_thread_id,
|
||
text_body=text_body,
|
||
operator_tg_id=operator_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
|
||
) -> 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)")
|