All checks were successful
CI Trade-In / changes (pull_request) Successful in 11s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 16s
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 6m17s
Deep review PR #3479 нашёл дефект в предыдущем фиксе (#3471 пункт 3): новые direction='out' строки стали видимы резолверам (find_chat_by_topic_message, find_thread_by_topic_message), но писались без support_chat_id. Резолверы матчат support_chat_id IS NULL как лениентный wildcard "любой текущий чат" (легаси-строки до 187/188) — то есть КАЖДАЯ out-строка становилась таким wildcard. При ротации support-группы новый message_id мог бы случайно совпасть со старой out-строкой: TG-путь увёл бы ответ ЧУЖОМУ клиенту через copyMessage, веб-путь записал бы ответ в чужой тред. Ровно от этого защищали миграции 187/188 (review M1). - bridge.py: TG- и веб-ветка `_handle_group_reply` теперь передают support_chat_id=settings.telegram_support_chat_id в record_message / record_web_out_message (симметрично уже существующей in-ветке). - web_support_storage.record_outbound: добавлен параметр support_chat_id, пишется в INSERT (колонка уже существовала, DDL не нужен). - Тест test_group_reply_to_own_previous_tg_reply_resolves_target_chat сидел предыдущую out-строку с уже заполненным support_chat_id вручную, хотя код писал NULL — маскировал дефект. Добавлены прямые проверки на записанное support_chat_id (TG и веб), обе падают на прежней реализации (проверено локальным откатом изменения — 2 failed, restore — 41 passed). - Комментарий про "апдейт частично применён в Telegram" в except-ветке веб-ответа был неверен для этого случая (на веб-пути ничего не уходит в Telegram до сбоя БД) — переписан на настоящую причину: сбой БД не переигрывается по общей политике process_update, а не из-за частичной доставки. Refs #3471
1360 lines
63 KiB
Python
1360 lines
63 KiB
Python
"""Unit tests for `app.services.tgbot.bridge` — чистая логика роутинга.
|
||
|
||
Coverage (per task spec + review follow-up):
|
||
- user → topic (личка клиента зеркалится в support-топик, с шапкой на первое
|
||
сообщение за throttle-окно, без шапки на повторное В окне, и снова с шапкой
|
||
после истечения окна — #6 review)
|
||
- реплай оператора → user (доставка ответа клиенту + запись direction='out')
|
||
- реплай оператора → веб-чат (#tgsupport-web): зеркало веб-сообщения резолвится
|
||
в web-тред (скоуп по (topic_message_id, support_chat_id) — review M1), ответ
|
||
пишется direction='out' БЕЗ Telegram-доставки; медиа-реплай (в т.ч. фото С
|
||
ПОДПИСЬЮ) на веб-зеркало — отказ ЦЕЛИКОМ + уведомление оператору в топике
|
||
(review M2, никакой частичной доставки одной подписи); зеркало от ЧУЖОГО/
|
||
устаревшего support_chat_id — не матчится (ротация группы); NULL
|
||
support_chat_id (легаси) — wildcard-матч; совпадение ОБЕИХ сторон
|
||
одновременно (tg И web) — громкий отказ (logger.error), а не молчаливый
|
||
выбор tg-пути (152-ФЗ misroute risk)
|
||
- реплай не на зеркало (или не реплай вообще) — тихий игнор, не мусорим в чат;
|
||
реплай на СООБЩЕНИЕ БОТА без записи в БД — WARNING про осиротевшее зеркало
|
||
(#4 review)
|
||
- дедуп update_id (<=offset — skip без side-effects; poison-pill апдейт всё
|
||
равно сдвигает offset, чтобы не подвесить весь поток)
|
||
- сбой БД (SQLAlchemyError) во время обработки → rollback() ПЕРЕД save_offset,
|
||
offset всё равно сдвигается — без этого следующий поллинг переиграл бы тот
|
||
же апдейт и задублировал зеркало клиента в топике (#3 review)
|
||
- ТРАНЗИЕНТНЫЙ сетевой отказ (TelegramNetworkError) на доставке ответа оператора
|
||
→ offset НЕ сдвигается, апдейт переигрывается и доходит до клиента; на потолке
|
||
`_MAX_NETWORK_REPLAYS` offset всё-таки сдвигается (поток не заклинен); poll loop
|
||
прерывает разбор пачки на неподтверждённом апдейте (#tg-connection-resilience)
|
||
- провал ВТОРИЧНОГО уведомления оператору в топик не отменяет основную ветку
|
||
(is_blocked остаётся, offset сдвигается) и логируется отдельной строкой
|
||
- /start → приветствие без зеркалирования
|
||
- Telegram 403 на доставку оператору → is_blocked + уведомление в топике
|
||
- TELEGRAM_SUPPORT_CHAT_ID не задан → клиенту уходит "сервис недоступен"
|
||
вместо тихой потери сообщения (#5 review)
|
||
|
||
NEVER calls real Telegram API — все HTTP-запросы mock'аются через
|
||
httpx.MockTransport (consistent с tests/services/test_dadata.py).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
import os
|
||
from typing import Any
|
||
|
||
import httpx
|
||
import pytest
|
||
from sqlalchemy.exc import SQLAlchemyError
|
||
|
||
# DATABASE_URL required by app.core.config before any app import (см. test_dadata.py).
|
||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||
|
||
from app.services.tgbot import bridge
|
||
from app.services.tgbot.client import TelegramClient, TelegramNetworkError
|
||
|
||
SUPPORT_CHAT_ID = -100123456789
|
||
SUPPORT_TOPIC_ID = 42
|
||
|
||
|
||
# ── Fake in-memory storage (БД не нужна) ─────────────────────────────────────
|
||
class FakeBridgeStorage:
|
||
"""In-memory `BridgeStorage` — никакой реальной БД, чистая логика роутинга.
|
||
|
||
`clock_s` — управляемые тестом "фейковые часы" (просто float, продвигается
|
||
вручную через `storage.clock_s += ...`), чтобы честно проверить throttle-окно
|
||
в `had_recent_inbound` (#6 review) без реального `time.sleep`/datetime-моков.
|
||
|
||
`fail_next_record_message` — если True, следующий вызов `record_message`
|
||
кидает `SQLAlchemyError` (симулирует обрыв коннекта к БД) и сбрасывается в
|
||
False — для теста rollback-пути в `process_update` (#3 review).
|
||
"""
|
||
|
||
def __init__(self, offset: int = 0) -> None:
|
||
self._offset = offset
|
||
self.users: dict[int, dict[str, Any]] = {}
|
||
self.messages: list[dict[str, Any]] = []
|
||
self.blocked: set[int] = set()
|
||
self.commits = 0
|
||
self.rollbacks = 0
|
||
self._next_id = 1
|
||
self.clock_s: float = 0.0
|
||
self.fail_next_record_message = False
|
||
# #3471 P0: следующий вызов record_web_out_message кидает SQLAlchemyError
|
||
# (симулирует обрыв коннекта к БД на веб-ветке — для веб-треда эта запись
|
||
# И ЕСТЬ доставка клиенту) и сбрасывается в False.
|
||
self.fail_next_record_web_out_message = False
|
||
# #tgsupport-web: web_support_messages-эквивалент, topic_message_id ->
|
||
# (thread_id, support_chat_id) — второй элемент моделирует колонку
|
||
# web_support_messages.support_chat_id (review M1); None = легаси wildcard.
|
||
# + журнал outbound-записей, записанных через реплай оператора.
|
||
self.web_topic_to_thread: dict[int, tuple[int, int | None]] = {}
|
||
self.web_out_messages: list[dict[str, Any]] = []
|
||
|
||
def get_offset(self) -> int:
|
||
return self._offset
|
||
|
||
def save_offset(self, update_id: int) -> None:
|
||
self._offset = update_id
|
||
|
||
def commit(self) -> None:
|
||
self.commits += 1
|
||
|
||
def rollback(self) -> None:
|
||
self.rollbacks += 1
|
||
|
||
def upsert_user(
|
||
self,
|
||
*,
|
||
chat_id: int,
|
||
username: str | None,
|
||
first_name: str | None,
|
||
last_name: str | None,
|
||
language_code: str | None,
|
||
) -> None:
|
||
self.users[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:
|
||
return any(
|
||
m["chat_id"] == chat_id
|
||
and m["direction"] == "in"
|
||
and (self.clock_s - m["recorded_at_s"]) < window_seconds
|
||
for m in self.messages
|
||
)
|
||
|
||
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:
|
||
if self.fail_next_record_message:
|
||
self.fail_next_record_message = False
|
||
raise SQLAlchemyError("simulated DB failure (deploy connection reset)")
|
||
row_id = self._next_id
|
||
self._next_id += 1
|
||
self.messages.append(
|
||
{
|
||
"id": row_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,
|
||
"recorded_at_s": self.clock_s,
|
||
}
|
||
)
|
||
return row_id
|
||
|
||
def find_chat_by_topic_message(self, topic_message_id: int, support_chat_id: int) -> int | None:
|
||
"""support_chat_id-скоуп (review M1): запись со ЧУЖИМ (не None, не текущим)
|
||
support_chat_id не матчится — None (легаси/дефолт) матчится всегда.
|
||
БЕЗ фильтра по direction (#3471 P0) — реплай на СВОЙ предыдущий ответ
|
||
(direction='out') резолвится так же, как реплай на зеркало клиента."""
|
||
for m in reversed(self.messages):
|
||
if m["topic_message_id"] != topic_message_id:
|
||
continue
|
||
entry_chat_id = m.get("support_chat_id")
|
||
if entry_chat_id is not None and entry_chat_id != support_chat_id:
|
||
continue
|
||
return m["chat_id"]
|
||
return None
|
||
|
||
def mark_blocked(self, chat_id: int) -> None:
|
||
self.blocked.add(chat_id)
|
||
|
||
# ── #tgsupport-web ────────────────────────────────────────────────────
|
||
def find_web_thread_by_topic_message(
|
||
self, topic_message_id: int, support_chat_id: int
|
||
) -> int | None:
|
||
"""Тот же support_chat_id-скоуп, что и `find_chat_by_topic_message` (review M1)."""
|
||
entry = self.web_topic_to_thread.get(topic_message_id)
|
||
if entry is None:
|
||
return None
|
||
thread_id, entry_chat_id = entry
|
||
if entry_chat_id is not None and entry_chat_id != support_chat_id:
|
||
return None
|
||
return thread_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:
|
||
if self.fail_next_record_web_out_message:
|
||
self.fail_next_record_web_out_message = False
|
||
raise SQLAlchemyError("simulated DB failure (deploy connection reset)")
|
||
self.web_out_messages.append(
|
||
{
|
||
"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,
|
||
}
|
||
)
|
||
# Зеркалим в `web_topic_to_thread` (deep review PR #3479) — реальный
|
||
# `record_outbound` пишет ту же строку в web_support_messages, которую
|
||
# потом читает `find_thread_by_topic_message`; без этого фейк не мог бы
|
||
# поймать баг "out-строка с support_chat_id=NULL — вечный wildcard".
|
||
if topic_message_id is not None:
|
||
self.web_topic_to_thread[topic_message_id] = (thread_id, support_chat_id)
|
||
|
||
|
||
# ── httpx mocking helpers (mirrors tests/services/test_dadata.py) ───────────
|
||
_REAL_ASYNC_CLIENT = httpx.AsyncClient
|
||
|
||
|
||
def _method_from_url(url: httpx.URL) -> str:
|
||
return str(url).rsplit("/", 1)[-1]
|
||
|
||
|
||
def _make_client(
|
||
responses: dict[str, Any], calls: list[tuple[str, dict[str, Any]]]
|
||
) -> TelegramClient:
|
||
"""TelegramClient wired to a MockTransport. `responses[method]` may be a dict
|
||
(returned as Bot API `result`), an int (HTTP error status), or a callable
|
||
`(payload) -> dict`. Every request is recorded into `calls`."""
|
||
|
||
def handler(request: httpx.Request) -> httpx.Response:
|
||
import json as _json
|
||
|
||
method = _method_from_url(request.url)
|
||
payload = _json.loads(request.content.decode("utf-8")) if request.content else {}
|
||
calls.append((method, payload))
|
||
|
||
canned = responses.get(method)
|
||
if isinstance(canned, int):
|
||
return httpx.Response(
|
||
canned, json={"ok": False, "error_code": canned, "description": "mocked error"}
|
||
)
|
||
if callable(canned):
|
||
canned = canned(payload)
|
||
result = canned if canned is not None else {"message_id": 999}
|
||
return httpx.Response(200, json={"ok": True, "result": result})
|
||
|
||
transport = httpx.MockTransport(handler)
|
||
|
||
def factory(*_: object, **__: object) -> httpx.AsyncClient:
|
||
return _REAL_ASYNC_CLIENT(transport=transport)
|
||
|
||
client = TelegramClient(token="fake-token")
|
||
import unittest.mock as mock
|
||
|
||
# Патчим httpx.AsyncClient ТОЛЬКО внутри client-модуля — не трогаем глобальный httpx.
|
||
patcher = mock.patch("app.services.tgbot.client.httpx.AsyncClient", factory)
|
||
patcher.start()
|
||
return client
|
||
|
||
|
||
@pytest.fixture(autouse=True)
|
||
def _support_chat_settings(monkeypatch: pytest.MonkeyPatch) -> None:
|
||
"""Все тесты по умолчанию считают support-группу/топик настроенными."""
|
||
monkeypatch.setattr(bridge.settings, "telegram_support_chat_id", SUPPORT_CHAT_ID)
|
||
monkeypatch.setattr(bridge.settings, "telegram_support_topic_id", SUPPORT_TOPIC_ID)
|
||
|
||
|
||
@pytest.fixture(autouse=True)
|
||
def _reset_flood_limiters(monkeypatch: pytest.MonkeyPatch) -> None:
|
||
"""F) `_flood_limiter`/`_flood_notify_limiter` — module-level singletons (тот же
|
||
паттерн, что `_send_limiter` в app/api/v1/support.py); большинство тестов в
|
||
этом файле шлют сообщения от одного и того же chat_id=555, поэтому без сброса
|
||
накопленные хиты одного теста бы протекали в следующий и ломали его
|
||
предположения (тест флуда должен видеть ЧИСТЫЙ бюджет)."""
|
||
monkeypatch.setattr(
|
||
bridge,
|
||
"_flood_limiter",
|
||
bridge.SlidingWindowLimiter(limit=bridge._FLOOD_LIMIT, window_s=bridge._FLOOD_WINDOW_S),
|
||
)
|
||
monkeypatch.setattr(
|
||
bridge,
|
||
"_flood_notify_limiter",
|
||
bridge.SlidingWindowLimiter(limit=1, window_s=bridge._FLOOD_WINDOW_S),
|
||
)
|
||
|
||
|
||
@pytest.fixture(autouse=True)
|
||
def _stop_patches():
|
||
"""Останавливает httpx.AsyncClient monkeypatch после каждого теста (unittest.mock.patch.start()
|
||
без контекст-менеджера требует явного stop, чтобы не утекать в соседние тесты)."""
|
||
import unittest.mock as mock
|
||
|
||
yield
|
||
mock.patch.stopall()
|
||
|
||
|
||
def _private_message(
|
||
*,
|
||
message_id: int = 1,
|
||
chat_id: int = 555,
|
||
text: str | None = "Здравствуйте, вопрос по trade-in",
|
||
username: str | None = "client_ivan",
|
||
first_name: str | None = "Иван",
|
||
last_name: str | None = "Петров",
|
||
) -> dict[str, Any]:
|
||
msg: dict[str, Any] = {
|
||
"message_id": message_id,
|
||
"chat": {"id": chat_id, "type": "private"},
|
||
"from": {
|
||
"id": chat_id,
|
||
"username": username,
|
||
"first_name": first_name,
|
||
"last_name": last_name,
|
||
"language_code": "ru",
|
||
},
|
||
}
|
||
if text is not None:
|
||
msg["text"] = text
|
||
return msg
|
||
|
||
|
||
def _group_reply_message(
|
||
*,
|
||
message_id: int = 200,
|
||
reply_to_message_id: int | None = 100,
|
||
text: str = "Ответ оператора",
|
||
operator_id: int = 777,
|
||
) -> dict[str, Any]:
|
||
msg: dict[str, Any] = {
|
||
"message_id": message_id,
|
||
"chat": {"id": SUPPORT_CHAT_ID, "type": "supergroup"},
|
||
"from": {"id": operator_id, "username": "operator1"},
|
||
"text": text,
|
||
}
|
||
if reply_to_message_id is not None:
|
||
msg["reply_to_message"] = {"message_id": reply_to_message_id}
|
||
return msg
|
||
|
||
|
||
# ── A) user → topic ──────────────────────────────────────────────────────────
|
||
|
||
|
||
async def test_private_message_mirrors_to_topic_with_header_on_first_contact() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 555}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
update = {"update_id": 10, "message": _private_message()}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
# Первое сообщение за окно → шапка ПЕРЕД зеркалом контента.
|
||
assert methods == ["sendMessage", "copyMessage"]
|
||
|
||
header_call = calls[0][1]
|
||
assert header_call["chat_id"] == SUPPORT_CHAT_ID
|
||
assert header_call["message_thread_id"] == SUPPORT_TOPIC_ID
|
||
assert "Иван Петров" in header_call["text"]
|
||
assert "@client_ivan" in header_call["text"]
|
||
|
||
mirror_call = calls[1][1]
|
||
assert mirror_call["from_chat_id"] == 555
|
||
assert mirror_call["chat_id"] == SUPPORT_CHAT_ID
|
||
assert mirror_call["message_id"] == 1
|
||
assert mirror_call["message_thread_id"] == SUPPORT_TOPIC_ID
|
||
|
||
assert len(storage.messages) == 1
|
||
rec = storage.messages[0]
|
||
assert rec["direction"] == "in"
|
||
assert rec["chat_id"] == 555
|
||
assert rec["topic_message_id"] == 555 # copyMessage result.message_id
|
||
assert rec["kind"] == "text"
|
||
assert rec["text_body"] == "Здравствуйте, вопрос по trade-in"
|
||
|
||
assert storage.get_offset() == 10
|
||
assert storage.commits == 1
|
||
assert 555 in storage.users
|
||
|
||
|
||
async def test_private_message_second_message_within_window_skips_header() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 556}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
# Симулируем уже существующее inbound-сообщение за последний час.
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=0,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="первое сообщение",
|
||
operator_tg_id=None,
|
||
)
|
||
|
||
update = {"update_id": 11, "message": _private_message(message_id=2)}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
# Шапка НЕ отправляется повторно — только зеркало.
|
||
assert methods == ["copyMessage"]
|
||
assert len(storage.messages) == 2
|
||
|
||
|
||
async def test_private_message_header_resent_after_window_expires() -> None:
|
||
"""#6 review: throttle-окно (3600с) реально проверяется по времени — после
|
||
истечения окна шапка отправляется заново (не одна на весь чат навсегда)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 557}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=0,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="первое сообщение (час назад)",
|
||
operator_tg_id=None,
|
||
)
|
||
# Продвигаем фейковые часы за throttle-окно (3600с).
|
||
storage.clock_s += bridge._HEADER_THROTTLE_WINDOW_S + 1
|
||
|
||
update = {"update_id": 13, "message": _private_message(message_id=3)}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["sendMessage", "copyMessage"] # шапка снова отправлена
|
||
assert len(storage.messages) == 2
|
||
|
||
|
||
async def test_private_message_start_sends_greeting_without_mirroring() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
update = {"update_id": 12, "message": _private_message(text="/start")}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["sendMessage"]
|
||
greeting_call = calls[0][1]
|
||
assert greeting_call["chat_id"] == 555
|
||
assert "МЕРА" in greeting_call["text"]
|
||
# /start не зеркалируется и не попадает в лог переписки.
|
||
assert storage.messages == []
|
||
assert storage.get_offset() == 12
|
||
|
||
|
||
async def test_private_message_notifies_client_when_support_chat_unset(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""#5 review: TELEGRAM_SUPPORT_CHAT_ID не задан → клиент получает "сервис
|
||
недоступен" вместо того, чтобы молча ждать ответа, который никогда не придёт."""
|
||
monkeypatch.setattr(bridge.settings, "telegram_support_chat_id", 0)
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
update = {"update_id": 14, "message": _private_message()}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["sendMessage"]
|
||
notify_call = calls[0][1]
|
||
assert notify_call["chat_id"] == 555
|
||
assert notify_call["text"] == bridge.SERVICE_UNAVAILABLE_TEXT
|
||
# Ничего не зеркалируется и не пишется в лог переписки — support-группа не настроена.
|
||
assert storage.messages == []
|
||
assert storage.get_offset() == 14
|
||
|
||
|
||
# ── F) флуд-лимит на отправителя (низкий приоритет) ─────────────────────────
|
||
|
||
|
||
async def test_private_message_flood_limit_blocks_excess_and_notifies_once() -> None:
|
||
"""Больше `_FLOOD_LIMIT` сообщений от ОДНОГО chat_id за окно — зеркалирование
|
||
сверх лимита отключается (никакого copyMessage, никакой записи в
|
||
tg_support_messages — маршрутизировать ответ всё равно нечего без
|
||
topic_message_id). Клиент получает уведомление о недоставке РОВНО один раз
|
||
за окно, а не на каждое следующее превышение — иначе само уведомление стало
|
||
бы вторым источником флуда."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 900}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
update_id = 100
|
||
for i in range(bridge._FLOOD_LIMIT):
|
||
update = {"update_id": update_id, "message": _private_message(message_id=i + 1)}
|
||
await bridge.process_update(update, client, storage)
|
||
update_id += 1
|
||
|
||
# Ровно _FLOOD_LIMIT сообщений прошли мирроринг: первое — шапка + зеркало,
|
||
# остальные — только зеркало.
|
||
mirrored_calls = [m for m, _ in calls if m == "copyMessage"]
|
||
assert len(mirrored_calls) == bridge._FLOOD_LIMIT
|
||
assert len(storage.messages) == bridge._FLOOD_LIMIT
|
||
|
||
calls.clear()
|
||
over_limit_update = {
|
||
"update_id": update_id,
|
||
"message": _private_message(message_id=bridge._FLOOD_LIMIT + 1),
|
||
}
|
||
await bridge.process_update(over_limit_update, client, storage)
|
||
update_id += 1
|
||
|
||
# Сверх лимита — НЕ зеркалируется, НЕ пишется в лог переписки, клиент
|
||
# получает уведомление о недоставке (не тихий игнор — клиент не должен
|
||
# решить, что оператор получил сообщение).
|
||
assert len(calls) == 1
|
||
method, payload = calls[0]
|
||
assert method == "sendMessage"
|
||
assert payload["chat_id"] == 555
|
||
assert payload["text"] == bridge.FLOOD_LIMITED_TEXT
|
||
assert len(storage.messages) == bridge._FLOOD_LIMIT
|
||
|
||
calls.clear()
|
||
second_over_limit_update = {
|
||
"update_id": update_id,
|
||
"message": _private_message(message_id=bridge._FLOOD_LIMIT + 2),
|
||
}
|
||
await bridge.process_update(second_over_limit_update, client, storage)
|
||
|
||
# Повторное превышение в ТОМ ЖЕ окне — уведомление подавлено (не второй
|
||
# источник флуда), никаких Telegram-вызовов вообще.
|
||
assert calls == []
|
||
assert len(storage.messages) == bridge._FLOOD_LIMIT
|
||
|
||
|
||
async def test_private_message_flood_limit_does_not_block_other_client() -> None:
|
||
"""Флуд-лимит — per-chat_id: клиент А исчерпал свой бюджет, но клиент Б
|
||
(другой chat_id) продолжает получать зеркалирование как обычно — один
|
||
флудящий клиент не блокирует доставку сообщений остальным (сама суть
|
||
задачи — воркер однопоточный, но лимит не даёт флудеру монополизировать
|
||
его через Telegram 429)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 901}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
flooding_chat_id = 555
|
||
update_id = 300
|
||
for i in range(bridge._FLOOD_LIMIT + 2):
|
||
update = {
|
||
"update_id": update_id,
|
||
"message": _private_message(chat_id=flooding_chat_id, message_id=i + 1),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
update_id += 1
|
||
|
||
calls.clear()
|
||
|
||
other_chat_id = 777001
|
||
other_update = {
|
||
"update_id": update_id,
|
||
"message": _private_message(chat_id=other_chat_id, message_id=1, username="another_client"),
|
||
}
|
||
await bridge.process_update(other_update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
# Другой клиент получает шапку (первое обращение) + зеркало как обычно —
|
||
# флуд первого клиента на него не влияет.
|
||
assert methods == ["sendMessage", "copyMessage"]
|
||
mirror_call = calls[1][1]
|
||
assert mirror_call["from_chat_id"] == other_chat_id
|
||
|
||
|
||
# ── B) реплай оператора → user ──────────────────────────────────────────────
|
||
|
||
|
||
async def test_group_reply_delivers_to_client_and_records_outbound() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 42}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
)
|
||
|
||
update = {
|
||
"update_id": 20,
|
||
"message": _group_reply_message(reply_to_message_id=100),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["copyMessage"]
|
||
delivery = calls[0][1]
|
||
assert delivery["chat_id"] == 555
|
||
assert delivery["from_chat_id"] == SUPPORT_CHAT_ID
|
||
|
||
assert len(storage.messages) == 2
|
||
out_rec = storage.messages[-1]
|
||
assert out_rec["direction"] == "out"
|
||
assert out_rec["chat_id"] == 555
|
||
assert out_rec["operator_tg_id"] == 777
|
||
assert out_rec["text_body"] == "Ответ оператора"
|
||
assert storage.get_offset() == 20
|
||
|
||
|
||
async def test_group_reply_not_a_reply_is_ignored() -> None:
|
||
"""Обычное сообщение в топике (не реплай) — тихий игнор, никаких Telegram-вызовов."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
update = {
|
||
"update_id": 21,
|
||
"message": _group_reply_message(reply_to_message_id=None),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == []
|
||
assert storage.messages == []
|
||
# offset всё равно сдвигается — апдейт "обработан" (даже если ничего не сделано).
|
||
assert storage.get_offset() == 21
|
||
|
||
|
||
async def test_group_reply_to_unknown_message_is_ignored() -> None:
|
||
"""Реплай на сообщение, которого нет в tg_support_messages как зеркало клиента —
|
||
тихий игнор (обычная болтовня в топике на постороннее сообщение, не от бота)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
update = {
|
||
"update_id": 22,
|
||
"message": _group_reply_message(reply_to_message_id=999999),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == []
|
||
assert storage.messages == []
|
||
assert storage.get_offset() == 22
|
||
|
||
|
||
async def test_group_reply_to_bot_message_without_record_logs_orphaned_mirror_warning(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
"""#4 review: реплай на сообщение БОТА, которого нет в tg_support_messages, —
|
||
вероятное осиротевшее зеркало (крах между copyMessage и commit). WARNING, не
|
||
тихий игнор — оператор иначе решит, что ответ клиенту доставлен."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
message = _group_reply_message(reply_to_message_id=100)
|
||
message["reply_to_message"]["from"] = {"id": 999, "is_bot": True, "username": "MERAsupport_bot"}
|
||
update = {"update_id": 24, "message": message}
|
||
|
||
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == [] # ответ НЕ доставлен — routing-ключ потерян
|
||
assert "осиротевшее" in caplog.text
|
||
assert storage.get_offset() == 24
|
||
|
||
|
||
async def test_group_reply_to_non_bot_message_without_record_stays_silent(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
"""Обычный реплай на сообщение ДРУГОГО ЧЕЛОВЕКА (не бота) в топике — реальная
|
||
болтовня, никакого WARNING (дискриминатор `is_bot` работает в обе стороны)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
message = _group_reply_message(reply_to_message_id=101)
|
||
message["reply_to_message"]["from"] = {"id": 42, "is_bot": False, "username": "colleague"}
|
||
update = {"update_id": 25, "message": message}
|
||
|
||
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == []
|
||
assert caplog.text == ""
|
||
|
||
|
||
async def test_group_reply_403_marks_blocked_and_notifies_topic() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": 403}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
)
|
||
|
||
update = {
|
||
"update_id": 23,
|
||
"message": _group_reply_message(reply_to_message_id=100),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["copyMessage", "sendMessage"]
|
||
notify_call = calls[1][1]
|
||
assert notify_call["chat_id"] == SUPPORT_CHAT_ID
|
||
assert "заблокирован" in notify_call["text"]
|
||
|
||
assert 555 in storage.blocked
|
||
# Неудачная доставка НЕ должна создавать фейковую запись 'out'.
|
||
assert len(storage.messages) == 1
|
||
assert storage.get_offset() == 23
|
||
|
||
|
||
# ── B') реплай оператора → веб-чат (#tgsupport-web) ──────────────────────────
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_records_outbound_web_message() -> None:
|
||
"""Реплай на зеркало веб-сообщения (не найдено в tg_support_messages, найдено
|
||
среди web_support_messages) → записывается в веб-тред, БЕЗ Telegram-доставки
|
||
(у веб-клиента нет личного чата с ботом)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
# topic_message_id=300 -> thread_id=42, под ТЕКУЩИМ support_chat_id.
|
||
storage.web_topic_to_thread[300] = (42, SUPPORT_CHAT_ID)
|
||
|
||
update = {
|
||
"update_id": 60,
|
||
"message": _group_reply_message(reply_to_message_id=300, text="Ответ по веб-чату"),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
# Никакого Telegram API вызова — веб-клиент не имеет личного чата с ботом.
|
||
assert calls == []
|
||
assert len(storage.web_out_messages) == 1
|
||
rec = storage.web_out_messages[0]
|
||
assert rec["thread_id"] == 42
|
||
assert rec["text_body"] == "Ответ по веб-чату"
|
||
assert rec["operator_tg_id"] == 777
|
||
# tg-путь тоже не тронут — ни одной записи в tg_support_messages.
|
||
assert storage.messages == []
|
||
assert storage.get_offset() == 60
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_with_null_support_chat_id_matches_current_chat() -> None:
|
||
"""Легаси-строка (до 187/188, support_chat_id=None) — лениентный wildcard,
|
||
матчится под ЛЮБЫМ текущим support_chat_id (review M1)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.web_topic_to_thread[305] = (46, None)
|
||
|
||
update = {
|
||
"update_id": 63,
|
||
"message": _group_reply_message(reply_to_message_id=305, text="Ответ по легаси-зеркалу"),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert len(storage.web_out_messages) == 1
|
||
assert storage.web_out_messages[0]["thread_id"] == 46
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_from_stale_support_chat_is_not_matched() -> None:
|
||
"""#tgsupport-web review M1: зеркало, записанное под ДРУГИМ (не текущим,
|
||
не None) support_chat_id — исторический артефакт ротации группы, НЕ валидный
|
||
маршрут сегодня. Не матчится → падает в orphan-check (не-bot реплай — тихий
|
||
игнор, никакой доставки в чужой/устаревший тред)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
stale_chat_id = -999999999999
|
||
storage.web_topic_to_thread[306] = (47, stale_chat_id)
|
||
|
||
update = {
|
||
"update_id": 64,
|
||
"message": _group_reply_message(reply_to_message_id=306),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == []
|
||
assert storage.web_out_messages == [] # НЕ доставлено в устаревший тред
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_without_text_is_refused_with_operator_notice(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
"""Веб-чат — текстовый MVP: реплай медиа-типом (нет text/caption) на веб-зеркало
|
||
не может быть доставлен — WARNING в лог И явное уведомление оператору в топике
|
||
(review M2: раньше был тихий игнор, оператор был уверен что ответил)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.web_topic_to_thread[301] = (43, SUPPORT_CHAT_ID)
|
||
|
||
message = _group_reply_message(reply_to_message_id=301, message_id=201)
|
||
del message["text"] # медиа-реплай без текста/caption
|
||
message["voice"] = {"file_id": "x"}
|
||
update = {"update_id": 61, "message": message}
|
||
|
||
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert storage.web_out_messages == []
|
||
assert "не текст" in caplog.text
|
||
assert storage.get_offset() == 61
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["sendMessage"]
|
||
notice_call = calls[0][1]
|
||
assert notice_call["chat_id"] == SUPPORT_CHAT_ID
|
||
assert notice_call["text"] == bridge._WEB_UNSUPPORTED_MEDIA_REPLY_TEXT
|
||
assert notice_call["reply_to_message_id"] == 201
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_with_photo_and_caption_is_refused_not_partial() -> None:
|
||
"""Фото С ПОДПИСЬЮ на веб-зеркало — НЕ доставляем только подпись молча
|
||
(клиент решил бы, что подпись — весь ответ): отказ целиком, как и без caption
|
||
(review M2)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.web_topic_to_thread[302] = (44, SUPPORT_CHAT_ID)
|
||
|
||
message = _group_reply_message(reply_to_message_id=302, message_id=202)
|
||
del message["text"]
|
||
message["photo"] = [{"file_id": "x"}]
|
||
message["caption"] = "Смотрите скриншот"
|
||
update = {"update_id": 65, "message": message}
|
||
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert storage.web_out_messages == [] # подпись НЕ доставлена как "весь ответ"
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["sendMessage"]
|
||
assert calls[0][1]["text"] == bridge._WEB_UNSUPPORTED_MEDIA_REPLY_TEXT
|
||
|
||
|
||
async def test_group_reply_refuses_delivery_when_both_tg_and_web_match(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
"""#tgsupport-web review M1: если topic_message_id одновременно резолвится и в
|
||
tg_support_messages, И в web_support_messages (под ОДНИМ и тем же
|
||
support_chat_id — целостность нарушена) — ГРОМКИЙ отказ (logger.error), НИКАКОЙ
|
||
доставки ни в Telegram-личку, ни в веб-тред. Раньше tg-путь выбирался молча —
|
||
misroute постороннему Telegram-клиенту (152-ФЗ risk)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 999}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=400,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
support_chat_id=SUPPORT_CHAT_ID,
|
||
)
|
||
storage.web_topic_to_thread[400] = (99, SUPPORT_CHAT_ID)
|
||
|
||
update = {
|
||
"update_id": 62,
|
||
"message": _group_reply_message(reply_to_message_id=400),
|
||
}
|
||
with caplog.at_level(logging.ERROR, logger="app.services.tgbot.bridge"):
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == [] # ничего не доставлено НИ В ОДНУ сторону
|
||
assert storage.web_out_messages == []
|
||
assert len(storage.messages) == 1 # только исходное 'in', никакого 'out'
|
||
assert "ОДНОВРЕМЕННО" in caplog.text
|
||
assert storage.get_offset() == 62
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_db_failure_notifies_operator_and_advances_offset() -> None:
|
||
"""#3471 P0: сбой БД на `record_web_out_message` — для веб-треда эта запись И
|
||
ЕСТЬ доставка клиенту (веб-фронт читает её polling'ом), поэтому тихий откат
|
||
означал бы навсегда потерянный ответ оператора (воспроизведено на проде
|
||
31.08.2026 — клиент kopylov). Теперь: rollback → уведомление оператору
|
||
реплаем в топик, что ответ НЕ доставлен → offset всё равно сдвигается (та же
|
||
политика, что у любого другого `SQLAlchemyError` в `process_update` — сбой
|
||
БД не переигрывается, human-in-the-loop retry заменяет технический) →
|
||
исключение наружу НЕ улетает."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.web_topic_to_thread[310] = (50, SUPPORT_CHAT_ID)
|
||
storage.fail_next_record_web_out_message = True
|
||
|
||
update = {
|
||
"update_id": 70,
|
||
"message": _group_reply_message(reply_to_message_id=310, message_id=210),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
# Ответ НЕ попал в web_out_messages (запись упала), но и не потерян молча.
|
||
assert storage.web_out_messages == []
|
||
assert storage.rollbacks == 1
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["sendMessage"]
|
||
notice = calls[0][1]
|
||
assert notice["chat_id"] == SUPPORT_CHAT_ID
|
||
assert notice["reply_to_message_id"] == 210
|
||
assert "НЕ доставлен" in notice["text"]
|
||
# Сбой БД не переигрывается (общая политика SQLAlchemyError) — offset сдвинут и закоммичен.
|
||
assert storage.get_offset() == 70
|
||
assert storage.commits == 1
|
||
|
||
|
||
async def test_group_reply_to_web_mirror_db_failure_and_notify_failure_logs_error(
|
||
caplog: pytest.LogCaptureFixture,
|
||
) -> None:
|
||
"""#3471 P0: сбой БД НА веб-ответе, а следом ещё и уведомление оператору не
|
||
ушло (Telegram недоступен) — полная тишина по обоим каналам. `process_update`
|
||
всё равно не падает и offset сдвигает, но остаётся `logger.error` с
|
||
идентификаторами (thread_id/message_id), НЕ текстом — по нему человек найдёт
|
||
ответ оператора в топике вручную."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
# 403 (не 5xx) — не ретраится клиентом, `_notify_topic` падает быстро и
|
||
# детерминированно (без реальных retry-пауз).
|
||
client = _make_client({"sendMessage": 403}, calls)
|
||
|
||
storage = FakeBridgeStorage()
|
||
storage.web_topic_to_thread[311] = (51, SUPPORT_CHAT_ID)
|
||
storage.fail_next_record_web_out_message = True
|
||
|
||
update = {
|
||
"update_id": 71,
|
||
"message": _group_reply_message(reply_to_message_id=311, message_id=211),
|
||
}
|
||
with caplog.at_level(logging.ERROR, logger="app.services.tgbot.bridge"):
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert storage.web_out_messages == []
|
||
assert "потерян молча" in caplog.text
|
||
assert "51" in caplog.text # thread_id узнаваем в логе
|
||
assert storage.get_offset() == 71 # сбой БД не переигрывается — offset сдвинут
|
||
assert storage.commits == 1
|
||
|
||
|
||
async def test_group_reply_to_own_previous_web_reply_resolves_thread() -> None:
|
||
"""#3471 P0 (пункт 3): реплай оператора на СВОЙ предыдущий веб-ответ (не на
|
||
исходное зеркало клиента) теперь тоже резолвится — `record_outbound`
|
||
сохраняет topic_message_id исходящей записи, `find_thread_by_topic_message`
|
||
больше не фильтрует по direction.
|
||
|
||
Deep review PR #3479: новая out-строка ОБЯЗАНА писаться с ТЕКУЩИМ
|
||
`support_chat_id`, а не NULL — NULL матчится `find_thread_by_topic_message`
|
||
как лениентный wildcard "любой чат" (легаси до 187/188), т.е. NULL сделал бы
|
||
КАЖДУЮ out-строку вечным wildcard-совпадением при ротации support-группы.
|
||
Эта проверка падает на дефектной реализации (support_chat_id не передавался
|
||
в `record_web_out_message`), даже когда resolve выше внешне "работает"."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
# Симулируем уже сохранённый предыдущий ответ оператора (topic_message_id=320)
|
||
# так, как это сделал бы реальный record_outbound после фикса.
|
||
storage.web_topic_to_thread[320] = (52, SUPPORT_CHAT_ID)
|
||
|
||
update = {
|
||
"update_id": 72,
|
||
"message": _group_reply_message(
|
||
reply_to_message_id=320, message_id=220, text="Продолжение"
|
||
),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert len(storage.web_out_messages) == 1
|
||
assert storage.web_out_messages[0]["thread_id"] == 52
|
||
assert storage.web_out_messages[0]["topic_message_id"] == 220
|
||
# Deep review PR #3479: НЕ NULL — иначе эта строка стала бы вечным wildcard.
|
||
assert storage.web_out_messages[0]["support_chat_id"] == SUPPORT_CHAT_ID
|
||
|
||
|
||
async def test_group_reply_to_own_previous_tg_reply_resolves_target_chat() -> None:
|
||
"""#3471 P0 (пункт 3): та же история для Telegram-пути — реплай оператора на
|
||
СВОЙ предыдущий ответ клиенту (direction='out', topic_message_id теперь
|
||
заполнен) резолвится в chat_id, а не проваливается в orphan-check.
|
||
|
||
Deep review PR #3479: новая out-строка ОБЯЗАНА писаться с ТЕКУЩИМ
|
||
`support_chat_id` (симметрично in-ветке) — иначе она стала бы вечным
|
||
wildcard в `find_chat_by_topic_message` при ротации support-группы, и
|
||
reply мог бы увести ответ ЧУЖОМУ клиенту."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 601}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
# Предыдущий ответ оператора клиенту, зафиксированный с topic_message_id
|
||
# (id этого ответа В ТОПИКЕ) — то, что теперь пишет TG-путь `_handle_group_reply`.
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="out",
|
||
tg_message_id=500,
|
||
topic_message_id=330,
|
||
kind="text",
|
||
text_body="Первый ответ оператора",
|
||
operator_tg_id=777,
|
||
support_chat_id=SUPPORT_CHAT_ID,
|
||
)
|
||
|
||
update = {
|
||
"update_id": 73,
|
||
"message": _group_reply_message(reply_to_message_id=330, message_id=230, text="Уточнение"),
|
||
}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
methods = [m for m, _ in calls]
|
||
assert methods == ["copyMessage"]
|
||
assert calls[0][1]["chat_id"] == 555
|
||
assert len(storage.messages) == 2 # исходный 'out' + новый 'out'
|
||
assert storage.messages[-1]["topic_message_id"] == 230
|
||
# Deep review PR #3479: НЕ NULL — иначе эта строка стала бы вечным wildcard.
|
||
assert storage.messages[-1]["support_chat_id"] == SUPPORT_CHAT_ID
|
||
|
||
|
||
# ── C) дедуп ──────────────────────────────────────────────────────────────────
|
||
|
||
|
||
async def test_dedup_update_id_leq_offset_is_skipped_without_side_effects() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage(offset=50)
|
||
|
||
update = {"update_id": 50, "message": _private_message()}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == []
|
||
assert storage.messages == []
|
||
assert storage.commits == 0 # ранний return — offset уже актуален, коммитить нечего
|
||
assert storage.get_offset() == 50
|
||
|
||
update_older = {"update_id": 10, "message": _private_message()}
|
||
await bridge.process_update(update_older, client, storage)
|
||
assert calls == []
|
||
assert storage.get_offset() == 50
|
||
|
||
|
||
async def test_poison_pill_update_still_advances_offset() -> None:
|
||
"""Апдейт, на котором обработчик упал (например, message без chat), не должен
|
||
подвесить весь поток — offset сдвигается даже при исключении внутри handler'а."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
malformed_message = {"message_id": 1, "chat": {"type": "private"}} # нет chat.id
|
||
update = {"update_id": 30, "message": malformed_message}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert storage.get_offset() == 30
|
||
assert storage.commits == 1
|
||
assert storage.messages == []
|
||
|
||
|
||
async def test_db_error_during_processing_rolls_back_and_still_advances_offset() -> None:
|
||
"""#3 review: SQLAlchemyError (напр. обрыв коннекта к БД при деплое) во время
|
||
`record_message` → storage.rollback() ПЕРЕД save_offset, offset всё равно
|
||
сдвигается. Без rollback() save_offset сам кинул бы PendingRollbackError →
|
||
process_update вылетел бы без сохранения offset'а → следующая итерация
|
||
переиграла бы тот же апдейт → copyMessage задублировал бы зеркало в топике
|
||
на каждый повтор поллинга (воспроизведено ревьюером)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 558}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.fail_next_record_message = True
|
||
|
||
update = {"update_id": 15, "message": _private_message()}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
# copyMessage успел уйти в Telegram (реальная утечка мирроринга при DB-сбое —
|
||
# известное ограничение атомарности между внешним API и БД, вне scope этого фикса),
|
||
# но rollback() отработал, offset сдвинут, commit вызван РОВНО один раз (в finally).
|
||
assert storage.rollbacks == 1
|
||
assert storage.commits == 1
|
||
assert storage.get_offset() == 15
|
||
# Запись сообщения НЕ попала в storage (record_message упал до append).
|
||
assert storage.messages == []
|
||
|
||
# Повторный вызов с тем же update_id теперь корректно дедупится — НЕ переигрывается.
|
||
calls.clear()
|
||
await bridge.process_update(update, client, storage)
|
||
assert calls == []
|
||
assert storage.get_offset() == 15
|
||
|
||
|
||
async def test_update_without_update_id_is_ignored() -> None:
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
await bridge.process_update({"message": _private_message()}, client, storage)
|
||
|
||
assert calls == []
|
||
assert storage.commits == 0
|
||
assert storage.get_offset() == 0
|
||
|
||
|
||
# ── kind inference ────────────────────────────────────────────────────────────
|
||
|
||
|
||
@pytest.mark.parametrize(
|
||
("message_extra", "expected_kind"),
|
||
[
|
||
({"text": "hi"}, "text"),
|
||
({"photo": [{"file_id": "x"}]}, "photo"),
|
||
({"document": {"file_id": "x"}}, "document"),
|
||
({"video": {"file_id": "x"}}, "video"),
|
||
({"voice": {"file_id": "x"}}, "voice"),
|
||
({"sticker": {"file_id": "x"}}, "other"),
|
||
({"location": {"latitude": 1, "longitude": 2}}, "other"),
|
||
({}, "other"),
|
||
],
|
||
)
|
||
def test_infer_kind(message_extra: dict[str, Any], expected_kind: str) -> None:
|
||
message = {"message_id": 1, "chat": {"id": 1, "type": "private"}, **message_extra}
|
||
assert bridge._infer_kind(message) == expected_kind
|
||
|
||
|
||
# ── unrelated chat types ─────────────────────────────────────────────────────
|
||
|
||
|
||
async def test_update_from_unrelated_chat_is_ignored_but_offset_advances() -> None:
|
||
"""Апдейт не из личного чата и не из support-группы (например, другой чат/канал)
|
||
— молча игнорируется, offset всё равно сдвигается."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
|
||
message = {
|
||
"message_id": 1,
|
||
"chat": {"id": -999, "type": "group"},
|
||
"from": {"id": 1},
|
||
"text": "болтовня в постороннем чате",
|
||
}
|
||
update = {"update_id": 40, "message": message}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert calls == []
|
||
assert storage.messages == []
|
||
assert storage.get_offset() == 40
|
||
|
||
|
||
# ── сетевая устойчивость (#tg-connection-resilience) ─────────────────────────
|
||
|
||
|
||
@pytest.fixture(autouse=True)
|
||
def _reset_network_replay_attempts() -> None:
|
||
"""`bridge._network_replay_attempts` — module-level словарь, его состояние
|
||
иначе протекало бы между тестами (потолок переигрываний виден глобально)."""
|
||
bridge._network_replay_attempts.clear()
|
||
|
||
|
||
def _network_boom(method: str = "copyMessage"):
|
||
"""Асинхронная заглушка метода клиента, изображающая исчерпанный бюджет ретраев."""
|
||
|
||
async def _raise(*_args: object, **_kwargs: object) -> dict[str, Any]:
|
||
raise TelegramNetworkError(method, "ConnectTimeout", 4)
|
||
|
||
return _raise
|
||
|
||
|
||
async def test_network_failure_on_operator_reply_keeps_offset_for_replay(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Главный сценарий потери: оператор ответил, Telegram в этот момент недоступен.
|
||
|
||
Раньше `finally: save_offset` подтверждал апдейт — ответ не доходил до клиента
|
||
НИКОГДА (Telegram апдейт больше не отдаёт, записи нет, оператор уверен, что
|
||
ответил). Теперь offset остаётся прежним, апдейт переигрывается и доходит.
|
||
"""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": {"message_id": 42}}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
)
|
||
|
||
update = {"update_id": 40, "message": _group_reply_message(reply_to_message_id=100)}
|
||
|
||
# Падает РОВНО первый copyMessage (monkeypatch.undo() тут не годится: он снял бы
|
||
# и autouse-патчи настроек support-группы, и ветка просто перестала бы работать).
|
||
real_copy = client.copy_message
|
||
failures = {"left": 1}
|
||
|
||
async def flaky_copy(**kwargs: Any) -> dict[str, Any]:
|
||
if failures["left"] > 0:
|
||
failures["left"] -= 1
|
||
raise TelegramNetworkError("copyMessage", "ConnectTimeout", 4)
|
||
return await real_copy(**kwargs)
|
||
|
||
monkeypatch.setattr(client, "copy_message", flaky_copy)
|
||
advanced = await bridge.process_update(update, client, storage)
|
||
|
||
assert advanced is False
|
||
assert storage.get_offset() == 0 # апдейт НЕ подтверждён — Telegram отдаст его снова
|
||
assert storage.commits == 0
|
||
assert storage.rollbacks == 1 # частичные записи не уедут чужим commit'ом
|
||
assert len(storage.messages) == 1 # фейковой записи 'out' не появилось
|
||
|
||
# Переигрывание: сеть починилась — тот же апдейт доставляется и подтверждается.
|
||
advanced = await bridge.process_update(update, client, storage)
|
||
|
||
assert advanced is True
|
||
assert storage.get_offset() == 40
|
||
assert [m for m, _ in calls] == ["copyMessage"]
|
||
out_rec = storage.messages[-1]
|
||
assert out_rec["direction"] == "out"
|
||
assert out_rec["chat_id"] == 555
|
||
|
||
|
||
async def test_network_failure_stops_advancing_only_until_replay_cap(
|
||
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
|
||
) -> None:
|
||
"""Потолок переигрываний: «вечно недоставляемый» апдейт не должен заклинить
|
||
поток навсегда (ровно то, от чего защищал прежний безусловный `finally`)."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
)
|
||
monkeypatch.setattr(client, "copy_message", _network_boom())
|
||
|
||
update = {"update_id": 41, "message": _group_reply_message(reply_to_message_id=100)}
|
||
|
||
for _ in range(bridge._MAX_NETWORK_REPLAYS - 1):
|
||
assert await bridge.process_update(update, client, storage) is False
|
||
assert storage.get_offset() == 0
|
||
|
||
with caplog.at_level(logging.ERROR, logger="app.services.tgbot.bridge"):
|
||
advanced = await bridge.process_update(update, client, storage)
|
||
|
||
assert advanced is True
|
||
assert storage.get_offset() == 41 # поток разблокирован
|
||
assert "потолок переигрываний" in caplog.text
|
||
# Счётчик снят — словарь не растёт от апдейта к апдейту.
|
||
assert bridge._network_replay_attempts == {}
|
||
|
||
|
||
async def test_poll_loop_stops_batch_on_unconfirmed_update(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Offset у Telegram — единая «высшая отметка»: подтвердив СЛЕДУЮЩИЙ апдейт
|
||
пачки, мы неявно подтвердили бы неудавшийся, и переигрывания не случилось бы.
|
||
Поэтому разбор пачки обрывается на первом неподтверждённом апдейте."""
|
||
processed: list[int] = []
|
||
|
||
async def fake_process(update: dict[str, Any], client: object, storage: object) -> bool:
|
||
processed.append(update["update_id"])
|
||
return update["update_id"] != 51 # 51 — сетевой отказ
|
||
|
||
monkeypatch.setattr(bridge, "process_update", fake_process)
|
||
|
||
iteration = {"n": 0}
|
||
|
||
def fake_shutdown() -> bool:
|
||
iteration["n"] += 1
|
||
return iteration["n"] > 1 # ровно одна итерация poll loop
|
||
|
||
monkeypatch.setattr(bridge, "shutdown_requested", fake_shutdown)
|
||
monkeypatch.setattr(bridge, "SqlBridgeStorage", lambda _db: FakeBridgeStorage())
|
||
|
||
class _FakeSession:
|
||
def __enter__(self) -> _FakeSession:
|
||
return self
|
||
|
||
def __exit__(self, *_exc: object) -> bool:
|
||
return False
|
||
|
||
class _FakeClient:
|
||
async def get_updates(self, **_kwargs: object) -> list[dict[str, Any]]:
|
||
return [{"update_id": 51}, {"update_id": 52}, {"update_id": 53}]
|
||
|
||
await bridge.run_poll_loop(_FakeClient(), _FakeSession, poll_timeout_s=1)
|
||
|
||
assert processed == [51] # 52/53 придут заново следующим getUpdates
|
||
|
||
|
||
async def test_topic_notification_failure_does_not_cancel_main_branch(
|
||
monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture
|
||
) -> None:
|
||
"""Уведомление в топик — вторичное действие: его сетевой отказ не отменяет
|
||
пометку is_blocked и не превращает апдейт в переигрываемый."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": 403}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
)
|
||
monkeypatch.setattr(client, "send_message", _network_boom("sendMessage"))
|
||
|
||
update = {"update_id": 42, "message": _group_reply_message(reply_to_message_id=100)}
|
||
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
|
||
advanced = await bridge.process_update(update, client, storage)
|
||
|
||
assert advanced is True
|
||
assert 555 in storage.blocked # основная ветка отработала
|
||
assert storage.get_offset() == 42
|
||
assert storage.commits == 1
|
||
assert len(storage.messages) == 1
|
||
assert "уведомление оператору в топик" in caplog.text
|
||
assert "chat_id=555" in caplog.text # только идентификатор, без текста переписки
|
||
|
||
|
||
async def test_topic_notification_uses_interactive_send_budget(
|
||
monkeypatch: pytest.MonkeyPatch,
|
||
) -> None:
|
||
"""Уведомление шлётся с УЗКИМ бюджетом: воркерный дефолт (5 ретраев, backoff
|
||
до 30с, полный retry_after на 429) застопорил бы весь poll loop на минуты."""
|
||
calls: list[tuple[str, dict[str, Any]]] = []
|
||
client = _make_client({"copyMessage": 403}, calls)
|
||
storage = FakeBridgeStorage()
|
||
storage.record_message(
|
||
chat_id=555,
|
||
direction="in",
|
||
tg_message_id=1,
|
||
topic_message_id=100,
|
||
kind="text",
|
||
text_body="вопрос клиента",
|
||
operator_tg_id=None,
|
||
)
|
||
|
||
sent: list[dict[str, Any]] = []
|
||
|
||
async def recording_send(**kwargs: Any) -> dict[str, Any]:
|
||
sent.append(kwargs)
|
||
return {"message_id": 1}
|
||
|
||
monkeypatch.setattr(client, "send_message", recording_send)
|
||
|
||
update = {"update_id": 43, "message": _group_reply_message(reply_to_message_id=100)}
|
||
await bridge.process_update(update, client, storage)
|
||
|
||
assert len(sent) == 1
|
||
assert sent[0]["timeout"] == bridge._NOTIFY_SEND_TIMEOUT_S
|
||
assert sent[0]["max_retries"] == bridge._NOTIFY_SEND_MAX_RETRIES
|
||
assert sent[0]["max_backoff"] == bridge._NOTIFY_SEND_MAX_BACKOFF_S
|
||
assert sent[0]["chat_id"] == SUPPORT_CHAT_ID
|