All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m19s
Замер прода за сутки 12.09.2026: 576 строк `network error` в логе `tradein-tgbot` и 7 полных исчерпаний бюджета ретраев, после которых падала итерация poll loop. Три причины, все подтверждены на коде и в рантайме. ## Ответ оператора мог пропасть навсегда `process_update` заканчивался безусловным `finally: save_offset(update_id)`. Замысел верный — «ядовитый» апдейт не должен блокировать поток, — но он не отличал неисправимый апдейт от транзиентного сетевого отказа. Оператор отвечает клиенту в топике, `copy_message` падает по сети, `TelegramNetworkError` улетает в общий `except Exception`, offset сдвигается. Telegram этот апдейт больше не отдаст, `record_message` не выполнился, оператор уверен, что ответил. Следа нет нигде, кроме строчки в логе. Теперь `process_update` возвращает `bool`. На `TelegramNetworkError` делается `rollback()`, offset НЕ сохраняется, возвращается `False`, и `run_poll_loop` прерывает разбор пачки — offset у Telegram единая «высшая отметка», подтверждение любого следующего апдейта неявно подтвердило бы и этот. Остаток пачки Telegram отдаст заново. Переигрывания ограничены сверху `_MAX_NETWORK_REPLAYS = 3`: без потолка «вечно недоставляемый» апдейт заклинил бы очередь навсегда, а это хуже потери одного сообщения. На потолке offset всё-таки двигается, но с `logger.error` и с `chat_id`/`message_id`, по которым человек найдёт ответ в топике и перешлёт руками. Текст переписки в лог по-прежнему не идёт. Дубли: `TelegramNetworkError` означает исчерпанный бюджет ретраев, при этом запрос мог дойти до Telegram, а ответ потеряться. Переигрывание тогда доставит сообщение второй раз. Это осознанный at-least-once компромисс — дубль видят и клиент, и оператор, а тихая потеря не видна никому. Полная идемпотентность по паре (update_id, target_chat_id) потребовала бы новой персистентной таблицы ради редкого случая; вместо неё число дублей жёстко ограничено сверху. Ветка `except TelegramApiError` с разбором `error_code == 403` («бот заблокирован») не тронута — там повтор действительно ничего не изменит. ## Таймаут задавался скаляром, поэтому connect ждал сорок секунд `httpx.AsyncClient(timeout=effective_timeout)` разворачивается в connect=read=write=pool. Для `getUpdates` бюджет ответа 40 секунд (30 держит Telegram плюс запас), и те же 40 секунд уходили на установку соединения — при живом connect в 0.036 секунды. Худший цикл: четыре попытки по 40 секунд плюс backoff, около трёх минут, в течение которых бот не видит ответов оператора. В логе это ровно те разрывы: 06:40:10, 06:42:22, 06:43:35. Теперь `httpx.Timeout(connect=5, read=<бюджет вызывающего>, write=10, pool=5)`, значения в именованных константах. Запас `+10s` у `get_updates` относится к read, докстринг поправлен. ## Клиент создавался заново на каждую попытку `httpx.AsyncClient` стоял ВНУТРИ цикла ретраев — keep-alive не было вовсе: полный TCP+TLS-хендшейк на каждый запрос и на каждый повтор, и заново кидался кубик «встанет ли коннект». Для long-polling это была основная статья сетевых отказов. Плюс три HTTP-ручки создавали `TelegramClient` на каждый входящий запрос. Теперь один ленивый переиспользуемый `AsyncClient` на экземпляр, с `aclose()` и `async with`. Общий клиент приложения живёт в новом `app/services/tgbot/shared.py`, создаётся и закрывается в lifespan; воркер бота держит свой на время поллинга. `keepalive_expiry` задан явно: дефолт httpx — 5 секунд, и с ним пул не давал бы ничего там, где нужнее всего. Poll loop переиспользует соединение и так, а вот веб-поддержка шлёт раз в минуты и за 5 секунд теряла бы его каждый раз. Плата за длинный keep-alive — шанс взять из пула закрытое той стороной соединение; httpx отдаёт это как `RemoteProtocolError`, который ретраится с #3457. ## Уведомления оператору шли с воркерным бюджетом внутри poll loop Обе отправки в топик («бот заблокирован», «веб-чат не поддерживает медиа») звались без своего бюджета, то есть с дефолтом в 5 ретраев и backoff до 30 секунд. Одна такая отправка стопорила весь цикл на минуты, а её отказ решал судьбу апдейта. Вынесены в `_notify_topic` с узким бюджетом и собственным `except`: провал вторичного действия больше не отменяет основную ветку. ## Тесты `tests/services/tgbot/test_shared.py` — новый, на жизненный цикл общего клиента. В `test_bridge.py` — сетевой отказ оставляет offset нетронутым и апдейт переигрывается, потолок разблокирует поток, отказ уведомления не отменяет основную ветку, прежнее поведение на 403 не изменилось. В `test_client.py` — раздельные таймауты доезжают до httpx per-request, два вызова используют один `AsyncClient`, `aclose()` его закрывает. Прогон по затронутым файлам: 127 passed. Ruff check и format чистые. Прокси намеренно не добавлялся: замер был на восьми запросах, это не статистика, и решение инфраструктурное. Если обрывы останутся — мерить сотней попыток отдельно.
867 lines
49 KiB
Python
867 lines
49 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
|
||
|
||
# Потолок переигрываний ОДНОГО 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
|
||
) -> 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})"
|
||
|
||
|
||
async def _notify_topic(
|
||
client: TelegramClient,
|
||
*,
|
||
text: str,
|
||
reply_to_message_id: int | None,
|
||
context: str,
|
||
) -> None:
|
||
"""Служебное уведомление оператору в support-топик (вторичное действие).
|
||
|
||
Два свойства, которых не было у прямых `client.send_message` вызовов:
|
||
1) узкий интерактивный бюджет (`_NOTIFY_SEND_*`) — иначе одна такая отправка
|
||
стопорит весь однопоточный poll loop на минуты;
|
||
2) собственный `except` — провал УВЕДОМЛЕНИЯ не отменяет основную ветку
|
||
обработки (клиент уже помечен заблокированным / медиа-реплай уже отклонён)
|
||
и не решает судьбу апдейта.
|
||
`context` — только технические идентификаторы (chat_id/thread_id), НЕ текст
|
||
переписки: логи моста принципиально не содержат ПДн.
|
||
"""
|
||
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,
|
||
)
|
||
except Exception:
|
||
logger.warning(
|
||
"tgbot bridge: не удалось отправить уведомление оператору в топик (%s) — "
|
||
"основная ветка обработки не отменяется",
|
||
context,
|
||
exc_info=True,
|
||
)
|
||
|
||
|
||
# ── 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
|
||
|
||
# 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)
|
||
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,
|
||
)
|
||
|
||
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 _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,
|
||
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 _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")
|
||
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
|
||
) -> 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, воспроизведено).
|
||
- любое прочее исключение (в т.ч. `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)")
|