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
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
995 lines
60 KiB
Python
995 lines
60 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`),
|
||
после КАЖДОГО апдейта — рестарт воркера не переигрывает уже обработанные
|
||
апдейты и не подвисает вечно на «ядовитом» апдейте. Исключение —
|
||
ТРАНЗИЕНТНЫЙ сетевой отказ (`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)")
|