diff --git a/tradein-mvp/backend/app/api/v1/support.py b/tradein-mvp/backend/app/api/v1/support.py index 8adfdf8e..0a8c1914 100644 --- a/tradein-mvp/backend/app/api/v1/support.py +++ b/tradein-mvp/backend/app/api/v1/support.py @@ -15,6 +15,20 @@ support-моста (`app.services.tgbot.bridge`, data/sql/186_tg_support.sql). Копия зеркала в топике всегда помечена "[С САЙТА] : ..." — оператор не должен путать веб-обращение с Telegram-клиентом (#tgsupport-web AC). + +КРИТИЧНО (review H1) — порядок операций в `send_support_message`: + БД-запись (`get_or_create_thread`) идёт ПОСЛЕ успешного `send_message`, не до. + Прод — один uvicorn-процесс БЕЗ `--workers` (docker-compose.prod.yml) с + синхронным SQLAlchemy engine (пул 5+10 overflow) на ОДНОМ event loop. Если бы + `INSERT ... ON CONFLICT DO UPDATE` уходил ДО Telegram-вызова, строка/row-lock + держались бы всё время, пока `send_message` ждёт Telegram (секунды-минуты при + 429/5xx на воркерных ретраях) — второй параллельный запрос ТОГО ЖЕ юзера + (двойной клик, вторая вкладка) упёрся бы в этот lock ВНУТРИ синхронного + psycopg-вызова внутри `async def`, останавливая event loop целиком (весь API + встаёт, не только этот эндпоинт). `_format_mirror_text` использует только + `username` — thread_id для отправки не нужен вообще, поэтому эту БД-операцию + можно безопасно отложить до после успешного sendMessage. Бонус: неудачная + отправка больше не создаёт тред. """ from __future__ import annotations @@ -42,20 +56,35 @@ MAX_MESSAGE_LENGTH = 4000 # Жёстче общего RateLimitMiddleware (300 req/60с на пользователя, app/main.py): # бот-токен общий на ВСЕХ клиентов веб-чата, флуд одного клиента иначе может -# упереться в Telegram-лимиты (`sendMessage` 429) и застопорить доставку всем -# остальным (см. задачу, п.6). 12 сообщений/минуту — щедро для живого диалога, -# но режет скрипт-флуд на порядок раньше общего API-лимита. +# упереться в Telegram-лимиты (`sendMessage` 429, лимит группы ~20 msg/min) и +# застопорить доставку всем остальным (см. задачу, п.6). 12 сообщений/минуту — +# щедро для живого диалога, но режет скрипт-флуд на порядок раньше общего API-лимита. _SEND_RATE_LIMIT = 12 _SEND_RATE_WINDOW_S = 60.0 _send_limiter = SlidingWindowLimiter(limit=_SEND_RATE_LIMIT, window_s=_SEND_RATE_WINDOW_S) +# #tgsupport-web review H1: интерактивный HTTP-запрос НЕ МОЖЕТ наследовать +# воркерную политику ретраев `TelegramClient` (по умолчанию — до 5 попыток, на +# 429 спит `retry_after` Telegram'а — для группы штатно 30-60с, на 5xx backoff до +# 30с — легальный суммарный бюджет минуты). Узкий бюджет здесь: 1 повтор, короткий +# timeout — интерактивный клиент должен получить ответ (даже если это ошибка) +# за секунды, а не висеть до исчерпания воркерных ретраев. +_INTERACTIVE_SEND_TIMEOUT_S = 10.0 +_INTERACTIVE_SEND_MAX_RETRIES = 1 + +# #tgsupport-web review M5: без LIMIT каждое монтирование виджета на старом +# треде отдавало бы ВЕСЬ лог переписки. См. `web_support_storage.list_messages`. +_LIST_MESSAGES_LIMIT = 200 + def _require_username(request: Request) -> str: """Достаёт X-Authenticated-User. rbac_guard (app/main.py) уже гарантирует его наличие в проде для non-public путей — этот guard здесь defence-in-depth и делает роутер тестируемым без поднятия всего app.main (см. tests/test_support.py, - как test_trade_in_lead.py для /lead).""" - username = request.headers.get("x-authenticated-user") + как test_trade_in_lead.py для /lead). `.strip()` (review L4) — лишний пробел + от прокси иначе завёл бы ВТОРОЙ тред на, по сути, того же пользователя + (username — UNIQUE ключ треда, "alice" != "alice ").""" + username = (request.headers.get("x-authenticated-user") or "").strip() if not username: raise HTTPException(status_code=401, detail="no authenticated user") return username @@ -105,7 +134,9 @@ class StatusOut(BaseModel): def _format_mirror_text(username: str, message_text: str) -> str: """Помечает зеркало как пришедшее С САЙТА, от какого пользователя — оператор - иначе не отличит веб-обращение от Telegram-клиента (#tgsupport-web AC).""" + иначе не отличит веб-обращение от Telegram-клиента (#tgsupport-web AC). + Использует ТОЛЬКО username — thread_id здесь не нужен (см. H1 в docstring + модуля), это то, что делает возможным отложить БД-запись до после отправки.""" return f"[С САЙТА] {username}:\n{message_text}" @@ -116,14 +147,21 @@ async def send_support_message( db: Annotated[Session, Depends(get_db)], ) -> SupportMessageOut: """Отправляет сообщение от лица *username* в support-топик (`sendMessage` — - не `copyMessage`: у веб-сообщения нет исходного Telegram-сообщения для копии).""" + не `copyMessage`: у веб-сообщения нет исходного Telegram-сообщения для копии). + + Порядок операций см. H1 в docstring модуля: rate-limit проверяется (но НЕ + расходуется, review L3) до отправки, thread создаётся ТОЛЬКО после успешного + `send_message` — до этого момента с БД не происходит ничего. + """ if not _bot_configured(): # Предсказуемое поведение вместо 500 (#tgsupport-web AC): бот не настроен # (пустой TELEGRAM_BOT_TOKEN, dev/staging) или support-топик не задан — # мирроринг невозможен физически, ничего не пишем в БД. raise HTTPException(status_code=503, detail=SERVICE_UNAVAILABLE_TEXT) - retry_after = _send_limiter.check(username) + # review L3: peek без расхода бюджета — неудачная отправка НЕ должна стоить + # пользователю попытки (расходуем `.record()` только на успех, ниже). + retry_after = _send_limiter.retry_after(username) if retry_after is not None: raise HTTPException( status_code=429, @@ -131,14 +169,15 @@ async def send_support_message( headers={"Retry-After": str(int(retry_after) + 1)}, ) - thread_id = storage.get_or_create_thread(db, username) - client = TelegramClient(settings.telegram_bot_token) try: mirrored = await client.send_message( chat_id=settings.telegram_support_chat_id, text=_format_mirror_text(username, payload.text), message_thread_id=settings.telegram_support_topic_id or None, + # review H1: узкий интерактивный бюджет — НЕ воркерные 5 ретраев/минуты. + timeout=_INTERACTIVE_SEND_TIMEOUT_S, + max_retries=_INTERACTIVE_SEND_MAX_RETRIES, ) except TelegramApiError: # НЕ логируем payload.text (переписка — ПДн) и НЕ логируем токен (его в @@ -148,13 +187,28 @@ async def send_support_message( ) raise HTTPException(status_code=502, detail=SERVICE_UNAVAILABLE_TEXT) from None - topic_message_id = mirrored.get("message_id") if isinstance(mirrored, dict) else None + # Отправка удалась — теперь и только теперь расходуем rate-limit бюджет. + _send_limiter.record(username) + topic_message_id = mirrored.get("message_id") if isinstance(mirrored, dict) else None + if topic_message_id is None: + # review L1: без topic_message_id реплай оператора на это сообщение + # НИКОГДА не смаршрутизируется обратно (find_thread_by_topic_message ищет + # именно по этому полю) — тихая, но зафиксированная в логе деградация. + logger.warning( + "web support: Telegram sendMessage не вернул message_id (username=%s) — " + "ответ оператора на это сообщение не будет смаршрутизирован", + username, + ) + + # review H1: БД-операция ПОСЛЕ успешной отправки — см. docstring модуля. + thread_id = storage.get_or_create_thread(db, username) row = storage.record_inbound( db, thread_id=thread_id, text_body=payload.text, topic_message_id=topic_message_id, + support_chat_id=settings.telegram_support_chat_id, ) db.commit() @@ -163,25 +217,34 @@ async def send_support_message( @router.get("/support/messages", response_model=list[SupportMessageOut]) -async def list_support_messages( +def list_support_messages( username: Annotated[str, Depends(_require_username)], db: Annotated[Session, Depends(get_db)], since: Annotated[int, Query(ge=0)] = 0, ) -> list[SupportMessageOut]: """Сообщения СВОЕГО треда с id > since. Тред резолвится по username — чужой - тред недостижим (нет параметра, которым его можно адресовать).""" + тред недостижим (нет параметра, которым его можно адресовать). + + Обычный (sync) `def`, не `async def` (review M3): тело — только синхронные + psycopg-вызовы, ни одного `await`; как `async def` это исполнялось бы прямо в + event loop (а фронт поллит эту ручку постоянно). Starlette гонит sync-handlers + в threadpool автоматически — тот же паттерн, что `trade_in.py:get_estimate`. + """ thread_id = storage.find_thread_id(db, username) if thread_id is None: return [] - rows = storage.list_messages(db, thread_id=thread_id, since_id=since) + rows = storage.list_messages( + db, thread_id=thread_id, since_id=since, limit=_LIST_MESSAGES_LIMIT + ) return [SupportMessageOut(**r) for r in rows] @router.get("/support/unread", response_model=UnreadOut) -async def get_support_unread( +def get_support_unread( username: Annotated[str, Depends(_require_username)], db: Annotated[Session, Depends(get_db)], ) -> UnreadOut: + """Sync `def` (review M3) — см. `list_support_messages`.""" thread_id = storage.find_thread_id(db, username) if thread_id is None: return UnreadOut(unread=0) @@ -189,10 +252,11 @@ async def get_support_unread( @router.post("/support/read", response_model=StatusOut) -async def mark_support_read( +def mark_support_read( username: Annotated[str, Depends(_require_username)], db: Annotated[Session, Depends(get_db)], ) -> StatusOut: + """Sync `def` (review M3) — см. `list_support_messages`.""" thread_id = storage.find_thread_id(db, username) if thread_id is not None: storage.mark_read(db, thread_id=thread_id) diff --git a/tradein-mvp/backend/app/core/ratelimit.py b/tradein-mvp/backend/app/core/ratelimit.py index e1fb4352..2b809ddb 100644 --- a/tradein-mvp/backend/app/core/ratelimit.py +++ b/tradein-mvp/backend/app/core/ratelimit.py @@ -97,19 +97,43 @@ class SlidingWindowLimiter: self._window_s = window_s self._hits: dict[str, deque[float]] = defaultdict(deque) - def check(self, key: str) -> float | None: - """Регистрирует попытку под *key*. Возвращает None, если она уложилась в - лимит (и учтена), иначе — сколько секунд ждать до следующей попытки.""" - now = time.monotonic() - bucket = self._hits[key] + def _prune(self, bucket: deque[float], now: float) -> None: cutoff = now - self._window_s while bucket and bucket[0] < cutoff: bucket.popleft() + def retry_after(self, key: str) -> float | None: + """Non-destructive проверка: сколько секунд ждать, если *key* СЕЙЧАС за + лимитом, иначе None. НЕ регистрирует попытку — вызывающая сторона решает + сама, когда звать `record()` (обычно — только на успех действия, #tgsupport-web + review L3: неудачная попытка не должна съедать бюджет).""" + now = time.monotonic() + bucket = self._hits[key] + self._prune(bucket, now) if len(bucket) >= self._limit: return self._window_s - (now - bucket[0]) + return None + def record(self, key: str) -> None: + """Регистрирует одну успешную попытку под *key*.""" + now = time.monotonic() + bucket = self._hits[key] + self._prune(bucket, now) bucket.append(now) + # Лёгкая защита от утечки памяти — чистим пустые корзины изредка (тот же + # паттерн, что RateLimitMiddleware.dispatch). + if len(self._hits) > 10000: + for k in [k for k, v in self._hits.items() if not v]: + del self._hits[k] + + def check(self, key: str) -> float | None: + """Комбинированная проверка+регистрация (peek+record за один вызов) — + для вызывающих, которым не нужно различать "попытка"/"успех" (см. + `retry_after`/`record` для раздельного варианта).""" + retry_after = self.retry_after(key) + if retry_after is not None: + return retry_after + self.record(key) return None diff --git a/tradein-mvp/backend/app/services/tgbot/bridge.py b/tradein-mvp/backend/app/services/tgbot/bridge.py index 51809338..6cbc8f3a 100644 --- a/tradein-mvp/backend/app/services/tgbot/bridge.py +++ b/tradein-mvp/backend/app/services/tgbot/bridge.py @@ -11,15 +11,23 @@ web_support_messages (direction='in') — этот модуль в этой ветке не участвует, только в разборе ответа (B ниже). B) Оператор отвечает РЕПЛАЕМ в support-группе на зеркало клиента → - находим chat_id по topic_message_id → copyMessage ответа в личку клиента → - запись (direction='out'). Если зеркало НЕ Telegram-клиента, а веб-чата - (data/sql/187_web_support_chat.sql) — доставка идёт НЕ в Telegram (у веб- - клиента нет личного чата с ботом), а записью direction='out' в - web_support_messages (веб-фронт вычитывает её обычным polling'ом). Реплай - не на зеркало (или не реплай вообще) — обычная болтовня в топике, тихий - игнор. Telegram 403 (клиент заблокировал бота) → is_blocked=true + - уведомление в топике (только для Telegram-ветки — у веб-клиента нет - "заблокировал бота"). + резолвим topic_message_id ОБЕ стороны (tg_support_messages И + web_support_messages), скоупя к ТЕКУЩЕМУ TELEGRAM_SUPPORT_CHAT_ID + (#tgsupport-web review M1 — если группу когда-нибудь сменят/пересоздадут, + Telegram message_id стартует заново и может совпасть со старым числом из + другой таблицы; без скоупинга это была бы ТИХАЯ доставка постороннему + клиенту). Совпадение НА ОБЕИХ сторонах одновременно — громкий отказ + (`logger.error`, ничего не доставляем) вместо произвольного выбора одной из + них. Иначе: chat_id найден → copyMessage ответа в личку клиента → запись + (direction='out'); thread_id найден (веб-зеркало) → доставка идёт НЕ в + Telegram (у веб-клиента нет личного чата с ботом), а записью direction='out' + в web_support_messages (веб-фронт вычитывает её обычным polling'ом); реплай + медиа-типом на веб-зеркало — веб-чат текстовый MVP, доставка целиком + отклоняется (не частично — фото с подписью НЕ превращается в "ответ = только + подпись"), оператор получает уведомление в топике (review M2). Реплай не на + зеркало (или не реплай вообще) — обычная болтовня в топике, тихий игнор. + Telegram 403 (клиент заблокировал бота) → is_blocked=true + уведомление в + топике (только для Telegram-ветки — у веб-клиента нет "заблокировал бота"). C) Дедуп: update_id <= сохранённого offset — skip. Offset сохраняется И коммитится в той же транзакции, что и запись сообщения (см. `process_update` `finally`), после КАЖДОГО апдейта — рестарт воркера не переигрывает уже @@ -79,6 +87,11 @@ SERVICE_UNAVAILABLE_TEXT = ( # tg_support_messages.kind): "text | photo | document | video | voice | other". _KNOWN_KINDS = ("text", "photo", "document", "video", "voice") +# #tgsupport-web review M2: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ) +# на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор +# получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил). +_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT = "Веб-чат поддерживает только текст, сообщение не доставлено." + # ── Storage abstraction (testable без реальной БД) ────────────────────────── class BridgeStorage(Protocol): @@ -115,13 +128,18 @@ class BridgeStorage(Protocol): 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) -> 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) -> 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 @@ -243,17 +261,19 @@ class SqlBridgeStorage: 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, created_at) + 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), NOW()) + CAST(:kind AS text), :text_body, CAST(:operator_tg_id AS bigint), + CAST(:support_chat_id AS bigint), NOW()) RETURNING id """ ), @@ -265,11 +285,17 @@ class SqlBridgeStorage: "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) -> int | 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( """ @@ -277,11 +303,13 @@ class SqlBridgeStorage: 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}, + {"topic_message_id": topic_message_id, "support_chat_id": support_chat_id}, ).fetchone() return int(row[0]) if row is not None else None @@ -294,10 +322,14 @@ class SqlBridgeStorage: {"chat_id": chat_id}, ) - def find_web_thread_by_topic_message(self, topic_message_id: int) -> int | None: + 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) + 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 @@ -400,6 +432,7 @@ async def _handle_private_message( kind=_infer_kind(message), text_body=text_body or message.get("caption"), operator_tg_id=None, + support_chat_id=settings.telegram_support_chat_id, ) @@ -416,11 +449,30 @@ async def _handle_group_reply( if not isinstance(mirror_message_id, int): return - target_chat_id = storage.find_chat_by_topic_message(mirror_message_id) + # #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-путь — НЕ ТРОНУТ, только обёрнут в explicit if - # (раньше было `if target_chat_id is None: ...; return`, теперь после - # этой ветки идёт ещё веб-резолв, см. ниже). + # Существующий Telegram-путь — НЕ ТРОНУТ. message_id = message.get("message_id") if not isinstance(message_id, int): return @@ -462,19 +514,32 @@ async def _handle_group_reply( ) return - # #tgsupport-web: не найдено среди tg_support_messages — пробуем веб-чат. - web_thread_id = storage.find_web_thread_by_topic_message(mirror_message_id) if web_thread_id is not None: - text_body = message.get("text") or message.get("caption") - if not text_body: - # Веб-чат — текстовый MVP (web_support_messages.text_body NOT NULL, - # нет kind/file_id колонок как у tg_support_messages) — доставить - # фото/документ/voice некуда, фронт это не отрендерит. + 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) без текста " - "(медиа?) — веб-чат текстовый, доставка невозможна, игнор", + "tgbot bridge: реплай на веб-зеркало (thread_id=%d) содержит %s, " + "не текст — веб-чат поддерживает только текст, доставка отклонена", web_thread_id, + kind, ) + await client.send_message( + chat_id=settings.telegram_support_chat_id, + text=_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT, + message_thread_id=settings.telegram_support_topic_id or None, + reply_to_message_id=message_id if isinstance(message_id, int) else None, + ) + return + + text_body = message.get("text") + if not text_body: + # Текстовый kind, но пустой text (защитный edge case) — нечего доставлять. return operator = message.get("from") or {} diff --git a/tradein-mvp/backend/app/services/tgbot/client.py b/tradein-mvp/backend/app/services/tgbot/client.py index 5909d0b3..faeeee66 100644 --- a/tradein-mvp/backend/app/services/tgbot/client.py +++ b/tradein-mvp/backend/app/services/tgbot/client.py @@ -238,12 +238,28 @@ class TelegramClient: text: str, message_thread_id: int | None = None, reply_to_message_id: int | None = None, + timeout: float | None = None, + max_retries: int | None = None, ) -> dict[str, Any]: - """sendMessage — текстовое сообщение (заголовки, приветствия, уведомления об ошибке).""" + """sendMessage — текстовое сообщение (заголовки, приветствия, уведомления об ошибке). + + `timeout`/`max_retries` — по умолчанию наследуют воркерную политику + (`_DEFAULT_TIMEOUT_S`/`_DEFAULT_MAX_RETRIES`: на 429 спим Telegram-овский + `retry_after` — для группы это штатные 30-60с, на 5xx backoff до 30с). + Это ПРИЕМЛЕМО для `tgbot_main.py` (изолированный long-polling воркер), но + ФАТАЛЬНО для интерактивного HTTP-запроса (#tgsupport-web review H1) — + синхронный request/response путь не может легально висеть минуты. Вызывающая + сторона на interactive-пути ОБЯЗАНА передать узкий бюджет явно (см. + `app.api.v1.support.send_support_message`).""" payload: dict[str, Any] = {"chat_id": chat_id, "text": text} if message_thread_id: payload["message_thread_id"] = message_thread_id if reply_to_message_id: payload["reply_to_message_id"] = reply_to_message_id - result = await self._request("sendMessage", payload) + kwargs: dict[str, Any] = {} + if timeout is not None: + kwargs["timeout"] = timeout + if max_retries is not None: + kwargs["max_retries"] = max_retries + result = await self._request("sendMessage", payload, **kwargs) return result if isinstance(result, dict) else {} diff --git a/tradein-mvp/backend/app/services/tgbot/web_support_storage.py b/tradein-mvp/backend/app/services/tgbot/web_support_storage.py index a52ff29d..42334482 100644 --- a/tradein-mvp/backend/app/services/tgbot/web_support_storage.py +++ b/tradein-mvp/backend/app/services/tgbot/web_support_storage.py @@ -62,19 +62,31 @@ def get_or_create_thread(db: Session, username: str) -> int: def record_inbound( - db: Session, *, thread_id: int, text_body: str, topic_message_id: int | None + db: Session, + *, + thread_id: int, + text_body: str, + topic_message_id: int | None, + support_chat_id: int | None, ) -> dict[str, Any]: """Записывает сообщение пользователя сайта (direction='in'). `topic_message_id` — - id зеркала (sendMessage) в support-топике, ключ маршрутизации ответа оператора.""" + id зеркала (sendMessage) в support-топике, ключ маршрутизации ответа оператора. + `support_chat_id` — TELEGRAM_SUPPORT_CHAT_ID В МОМЕНТ отправки (#tgsupport-web + review M1): скоупит будущий резолв `find_thread_by_topic_message` к ТЕКУЩЕЙ + support-группе — если группу когда-нибудь сменят/пересоздадут, Telegram + message_id стартует заново с 1 в новом чате и может совпасть с числом из + старого — без этого поля коллизия была бы ТИХОЙ (см. миграцию 187/188).""" row = ( db.execute( text( """ INSERT INTO web_support_messages - (thread_id, direction, text_body, topic_message_id, operator_tg_id, created_at) + (thread_id, direction, text_body, topic_message_id, + support_chat_id, operator_tg_id, created_at) VALUES (CAST(:thread_id AS bigint), 'in', :text_body, - CAST(:topic_message_id AS bigint), NULL, NOW()) + CAST(:topic_message_id AS bigint), + CAST(:support_chat_id AS bigint), NULL, NOW()) RETURNING id, direction, text_body, operator_tg_id, created_at """ ), @@ -82,6 +94,7 @@ def record_inbound( "thread_id": thread_id, "text_body": text_body, "topic_message_id": topic_message_id, + "support_chat_id": support_chat_id, }, ) .mappings() @@ -90,9 +103,17 @@ def record_inbound( return dict(row) -def find_thread_by_topic_message(db: Session, topic_message_id: int) -> int | None: +def find_thread_by_topic_message( + db: Session, topic_message_id: int, support_chat_id: int +) -> int | None: """Резолвит id зеркала (сообщения оператора reply_to) в thread_id — только - среди direction='in' записей, зеркало-конвенция как в tg_support_messages (186).""" + среди direction='in' записей, зеркало-конвенция как в tg_support_messages (186). + + Скоупим к ТЕКУЩЕМУ `support_chat_id` (#tgsupport-web review M1): строка со + ЧУЖИМ (не NULL и не текущим) support_chat_id — это исторический артефакт + ротации support-группы, НЕ валидный маршрут для сегодняшнего реплая. NULL + (легаси-строки до этой колонки, если такие есть) — лениентно матчатся как + "любой чат", т.к. до введения этого поля был ровно один действующий чат.""" row = db.execute( text( """ @@ -100,11 +121,12 @@ def find_thread_by_topic_message(db: Session, topic_message_id: int) -> int | No FROM web_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}, + {"topic_message_id": topic_message_id, "support_chat_id": support_chat_id}, ).fetchone() return int(row[0]) if row is not None else None @@ -135,8 +157,17 @@ def record_outbound( return int(row[0]) if row is not None else None -def list_messages(db: Session, *, thread_id: int, since_id: int) -> list[dict[str, Any]]: - """Сообщения треда с id > since_id, по возрастанию (обычный polling с фронта).""" +def list_messages( + db: Session, *, thread_id: int, since_id: int, limit: int = 200 +) -> list[dict[str, Any]]: + """Сообщения треда с id > since_id, по возрастанию (обычный polling с фронта). + + `limit` (#tgsupport-web review M5): без него КАЖДОЕ монтирование виджета на + старом треде отдавало бы ВЕСЬ лог переписки. Берём последние `limit` (ORDER + BY id DESC + LIMIT), потом разворачиваем в хронологический порядок — так + incremental-polling (`since_id` = последний известный id, обычно единицы + новых строк) не страдает, а первый холодный load длинного треда получает + последние `limit`, а не самые старые.""" rows = ( db.execute( text( @@ -145,15 +176,16 @@ def list_messages(db: Session, *, thread_id: int, since_id: int) -> list[dict[st FROM web_support_messages WHERE thread_id = CAST(:thread_id AS bigint) AND id > CAST(:since_id AS bigint) - ORDER BY id ASC + ORDER BY id DESC + LIMIT CAST(:limit AS integer) """ ), - {"thread_id": thread_id, "since_id": since_id}, + {"thread_id": thread_id, "since_id": since_id, "limit": limit}, ) .mappings() .all() ) - return [dict(r) for r in rows] + return [dict(r) for r in reversed(rows)] def count_unread(db: Session, *, thread_id: int) -> int: diff --git a/tradein-mvp/backend/data/sql/187_web_support_chat.sql b/tradein-mvp/backend/data/sql/187_web_support_chat.sql index bf4b92cf..1f23cc92 100644 --- a/tradein-mvp/backend/data/sql/187_web_support_chat.sql +++ b/tradein-mvp/backend/data/sql/187_web_support_chat.sql @@ -35,15 +35,30 @@ -- - topic_message_id-маршрутизация (ключевой механизм моста) СОХРАНЕНА -- 1-в-1 по конвенции 186: partial UNIQUE на topic_message_id, -- заполняется только для direction='in', NULL для direction='out'. --- Коллизий между web_support_messages.topic_message_id и --- tg_support_messages.topic_message_id НЕ возникает: оба — Telegram --- message_id ОДНОЙ и той же support-супергруппы, а Telegram message_id --- в пределах одного чата монотонно возрастает и никогда не переиспользуется --- — значит конкретное значение окажется ровно в одной из двух таблиц. -- - bridge.py меняется МИНИМАЛЬНО: _handle_group_reply получает одну --- дополнительную ветку (пробуем tg-резолв, потом web-резолв, потом +-- дополнительную ветку (резолвит tg-путь И web-путь, потом -- existing orphan-warning) — существующий Telegram-путь не трогается. -- +-- ⚠️ CROSS-TABLE КОЛЛИЗИЯ topic_message_id (review M1, зафиксировано ДО +-- первого прод-использования, пока обе таблицы пусты): +-- Инвариант "topic_message_id уникален между tg_support_messages и +-- web_support_messages" на самом деле звучит так: "уникален, ПОКА +-- TELEGRAM_SUPPORT_CHAT_ID не менялся". Это OPS-инвариант, а НЕ DB-инвариант — +-- ничем не гарантирован. Смена/пересоздание support-группы обнуляет счётчик +-- Telegram message_id в новом чате; когда он дорастёт до диапазона, +-- использованного старым чатом, — number, ранее занятый ОДНОЙ таблицей, +-- может совпасть с числом, занятым ДРУГОЙ. Внутри одной таблицы partial +-- UNIQUE превращает такую коллизию в громкий отказ INSERT — это ок. МЕЖДУ +-- таблицами constraint'а нет: без доп. скоупинга бот молча доставил бы ответ +-- оператора НЕ ТОМУ клиенту (152-ФЗ-инцидент, происходящий тихо). +-- Фикс: колонка `support_chat_id` на web_support_messages (симметричная +-- колонка для УЖЕ применённой tg_support_messages — отдельная миграция +-- 188_tg_support_chat_id_scope.sql, эту таблицу нельзя трогать здесь, она +-- уже применена/задеплоена как часть 186). Резолв (bridge.py) матчит ПАРУ +-- (support_chat_id, topic_message_id), а не topic_message_id в одиночку; +-- NULL (легаси-строки без этой колонки) — лениентный wildcard, т.к. на тот +-- момент действовал ровно один чат. +-- -- ЧТО: -- - web_support_threads — один тред на username (сайт = 1 логин = 1 линия -- переписки с поддержкой, без под-тредов). @@ -53,11 +68,16 @@ -- 152-ФЗ: -- Переписка (text_body) — ПДн (может содержать любые данные, которые юзер -- решит написать). ON DELETE CASCADE от web_support_threads делает erasure --- одной операцией: DELETE FROM web_support_threads WHERE username = :u. +-- ОДНОЙ операцией (DELETE FROM web_support_threads WHERE username = :u) ДЛЯ +-- КОПИИ В ЭТОЙ БД. Копия того же текста уже ушла в Telegram-топик (sendMessage +-- зеркало) и живёт ТАМ вне зоны действия этого DELETE — реальное "право на +-- забвение" по всей цепочке требует ОТДЕЛЬНОЙ процедуры (удаление сообщений в +-- Telegram-супергруппе через Bot API deleteMessage, вне scope этой миграции). +-- Не ссылаться на этот комментарий как на доказательство полного erasure. -- -- IDEMPOTENCY: CREATE TABLE/INDEX IF NOT EXISTS — безопасный re-run. -- Зависимости: нет (новые standalone таблицы, никакие существующие --- tg_support_*/иные таблицы не трогаются). +-- tg_support_*/иные таблицы не трогаются — см. 188 для ALTER на tg_support_messages). BEGIN; @@ -80,14 +100,16 @@ CREATE TABLE IF NOT EXISTS web_support_messages ( direction text NOT NULL CHECK (direction IN ('in', 'out')), text_body text NOT NULL CHECK (char_length(btrim(text_body)) > 0), topic_message_id bigint, + support_chat_id bigint, operator_tg_id bigint, created_at timestamptz NOT NULL DEFAULT now() ); -COMMENT ON TABLE web_support_messages IS '152-ФЗ: полный лог веб-чата поддержки (ПДн — содержимое сообщений). Каскадно удаляется вместе с web_support_threads по username.'; +COMMENT ON TABLE web_support_messages IS '152-ФЗ: полный лог веб-чата поддержки (ПДн — содержимое сообщений; удаление подчищает КОПИЮ В ЭТОЙ БД, не Telegram-топик — см. блок 152-ФЗ выше). Каскадно удаляется вместе с web_support_threads по username.'; COMMENT ON COLUMN web_support_messages.direction IS '''in'' — сообщение от пользователя сайта; ''out'' — ответ оператора (доставлен через реплай в Telegram-топике, см. bridge.py _handle_group_reply).'; COMMENT ON COLUMN web_support_messages.text_body IS 'Текст сообщения. Веб-чат — текстовый MVP, медиа не поддерживается (в отличие от tg_support_messages.kind).'; COMMENT ON COLUMN web_support_messages.topic_message_id IS 'id зеркала (sendMessage) в support-топике — ключ маршрутизации ответа, только для direction=''in''. NULL для ''out'' (конвенция 186: маршрутизирующий ключ живёт исключительно на inbound-записи).'; +COMMENT ON COLUMN web_support_messages.support_chat_id IS 'TELEGRAM_SUPPORT_CHAT_ID в момент отправки — скоупит резолв topic_message_id к ТЕКУЩЕЙ support-группе (review M1: без этого поля ротация группы даёт тихую cross-table коллизию, см. блок выше). NULL — лениентный wildcard для строк без этого поля.'; COMMENT ON COLUMN web_support_messages.operator_tg_id IS 'Telegram user id оператора, ответившего в топике; заполняется только для direction=''out''.'; CREATE UNIQUE INDEX IF NOT EXISTS web_support_messages_topic_message_id_uq diff --git a/tradein-mvp/backend/data/sql/188_tg_support_chat_id_scope.sql b/tradein-mvp/backend/data/sql/188_tg_support_chat_id_scope.sql new file mode 100644 index 00000000..753ce2e9 --- /dev/null +++ b/tradein-mvp/backend/data/sql/188_tg_support_chat_id_scope.sql @@ -0,0 +1,40 @@ +-- 188_tg_support_chat_id_scope.sql +-- Симметричная колонка для web_support_messages.support_chat_id (см. +-- data/sql/187_web_support_chat.sql — полный разбор проблемы в блоке "CROSS-TABLE +-- КОЛЛИЗИЯ topic_message_id" там же). +-- +-- ПОЧЕМУ ОТДЕЛЬНАЯ МИГРАЦИЯ, А НЕ ПРАВКА 186: +-- tg_support_messages создана в data/sql/186_tg_support.sql — миграция, которая +-- к моменту написания этого файла уже смержена в main отдельным PR (#2526) и, +-- по конвенции проекта (deploy-tradein.yml применяет каждый data/sql/*.sql РОВНО +-- ОДИН РАЗ по bare filename через _schema_migrations), скорее всего уже +-- применена на проде. Редактирование СОДЕРЖИМОГО уже применённого файла НЕ +-- долетает до прода повторным прогоном — прод просто пропустит файл с тем же +-- именем. Единственный корректный способ добавить колонку в уже существующую +-- таблицу — новый ALTER-файл. +-- +-- ЧТО: tg_support_messages.support_chat_id bigint (nullable) — TELEGRAM_SUPPORT_ +-- CHAT_ID в момент записи 'in'-сообщения. bridge.py.find_chat_by_topic_message +-- матчит (support_chat_id, topic_message_id) вместо topic_message_id в одиночку; +-- NULL (все строки ДО этой миграции) — лениентный wildcard-матч, т.к. до +-- появления этой колонки действовал ровно один support-чат за раз. +-- +-- Бэкфилл существующих строк текущим TELEGRAM_SUPPORT_CHAT_ID НЕ делаем: значение +-- живёт в Python `settings`/env, разное на каждом окружении (dev/staging/prod), а +-- plain-SQL миграция не имеет доступа к процессным env vars — хардкодить +-- конкретный chat_id в SQL-файл было бы хрупко и окружение-специфично. NULL +-- (wildcard) для существующих строк — безопасный дефолт: они писались, когда +-- support-чат был ровно один, коллизии из-за смены чата у НИХ по определению +-- невозможны (см. 187 — только смена чата ПОСЛЕ появления этой колонки создаёт +-- сценарий, который она защищает). +-- +-- IDEMPOTENCY: ADD COLUMN IF NOT EXISTS — безопасный re-run. Не трогает +-- существующие данные/constraints tg_support_messages. + +BEGIN; + +ALTER TABLE tg_support_messages ADD COLUMN IF NOT EXISTS support_chat_id bigint; + +COMMENT ON COLUMN tg_support_messages.support_chat_id IS 'TELEGRAM_SUPPORT_CHAT_ID в момент записи ''in''-сообщения — скоупит резолв topic_message_id к ТЕКУЩЕЙ support-группе (review M1, см. data/sql/187_web_support_chat.sql). NULL — строки до этой колонки (лениентный wildcard-матч).'; + +COMMIT; diff --git a/tradein-mvp/backend/data/sql/_manifest_applied.txt b/tradein-mvp/backend/data/sql/_manifest_applied.txt index fe7aadf7..0d8889dc 100644 --- a/tradein-mvp/backend/data/sql/_manifest_applied.txt +++ b/tradein-mvp/backend/data/sql/_manifest_applied.txt @@ -176,4 +176,3 @@ 169_osm_poi_ekb_local.sql 170_scrape_schedules_seed_osm_poi_ekb_refresh.sql 172_trade_in_leads.sql -187_web_support_chat.sql diff --git a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py index f0a62d0e..fdc82334 100644 --- a/tradein-mvp/backend/tests/services/tgbot/test_bridge.py +++ b/tradein-mvp/backend/tests/services/tgbot/test_bridge.py @@ -6,9 +6,14 @@ Coverage (per task spec + review follow-up): после истечения окна — #6 review) - реплай оператора → user (доставка ответа клиенту + запись direction='out') - реплай оператора → веб-чат (#tgsupport-web): зеркало веб-сообщения резолвится - в web-тред, ответ пишется direction='out' БЕЗ Telegram-доставки; медиа-реплай - на веб-зеркало (нет text/caption) — тихий игнор (веб-чат текстовый MVP); - tg-резолв имеет приоритет над веб-резолвом при (искусственной) коллизии + в 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) @@ -69,9 +74,11 @@ class FakeBridgeStorage: self._next_id = 1 self.clock_s: float = 0.0 self.fail_next_record_message = False - # #tgsupport-web: web_support_messages-эквивалент (topic_message_id -> - # thread_id) + журнал outbound-записей, записанных через реплай оператора. - self.web_topic_to_thread: dict[int, int] = {} + # #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: @@ -120,6 +127,7 @@ class FakeBridgeStorage: 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 @@ -136,23 +144,39 @@ class FakeBridgeStorage: "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) -> int | None: + 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 (легаси/дефолт) матчится всегда.""" for m in reversed(self.messages): - if m["direction"] == "in" and m["topic_message_id"] == topic_message_id: - return m["chat_id"] + if m["direction"] != "in" or 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) -> int | None: - return self.web_topic_to_thread.get(topic_message_id) + 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 @@ -561,7 +585,8 @@ async def test_group_reply_to_web_mirror_records_outbound_web_message() -> None: calls: list[tuple[str, dict[str, Any]]] = [] client = _make_client({}, calls) storage = FakeBridgeStorage() - storage.web_topic_to_thread[300] = 42 # topic_message_id=300 -> thread_id=42 + # topic_message_id=300 -> thread_id=42, под ТЕКУЩИМ support_chat_id. + storage.web_topic_to_thread[300] = (42, SUPPORT_CHAT_ID) update = { "update_id": 60, @@ -581,33 +606,107 @@ async def test_group_reply_to_web_mirror_records_outbound_web_message() -> None: assert storage.get_offset() == 60 -async def test_group_reply_to_web_mirror_without_text_is_ignored( - caplog: pytest.LogCaptureFixture, -) -> None: - """Веб-чат — текстовый MVP: реплай медиа-типом (нет text/caption) на веб-зеркало - не может быть доставлен — тихий (WARNING, не error) игнор, ничего не пишем.""" +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[301] = 43 + storage.web_topic_to_thread[305] = (46, None) - message = _group_reply_message(reply_to_message_id=301) + 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 calls == [] assert storage.web_out_messages == [] - assert "текстовый" in caplog.text + 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_prefers_tg_thread_when_both_would_match() -> None: - """Приоритет резолва — tg СНАЧАЛА: если topic_message_id найден среди - tg_support_messages, веб-резолв даже не вызывается (существующий Telegram-путь - работает как раньше, не деградирует из-за новой ветки).""" + +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() @@ -619,24 +718,22 @@ async def test_group_reply_prefers_tg_thread_when_both_would_match() -> None: kind="text", text_body="вопрос клиента", operator_tg_id=None, + support_chat_id=SUPPORT_CHAT_ID, ) - # Тот же topic_message_id "случайно" тоже был бы в web-мапе — не должен - # переопределять tg-резолв (defensive, в реальности Telegram message_id не - # повторяется в пределах чата). - storage.web_topic_to_thread[400] = 99 + storage.web_topic_to_thread[400] = (99, SUPPORT_CHAT_ID) update = { "update_id": 62, "message": _group_reply_message(reply_to_message_id=400), } - await bridge.process_update(update, client, storage) + with caplog.at_level(logging.ERROR, logger="app.services.tgbot.bridge"): + await bridge.process_update(update, client, storage) - methods = [m for m, _ in calls] - assert methods == ["copyMessage"] + assert calls == [] # ничего не доставлено НИ В ОДНУ сторону assert storage.web_out_messages == [] - out_rec = storage.messages[-1] - assert out_rec["direction"] == "out" - assert out_rec["chat_id"] == 555 + assert len(storage.messages) == 1 # только исходное 'in', никакого 'out' + assert "ОДНОВРЕМЕННО" in caplog.text + assert storage.get_offset() == 62 # ── C) дедуп ────────────────────────────────────────────────────────────────── diff --git a/tradein-mvp/backend/tests/test_ratelimit.py b/tradein-mvp/backend/tests/test_ratelimit.py index bf5fba49..4d1cd8dc 100644 --- a/tradein-mvp/backend/tests/test_ratelimit.py +++ b/tradein-mvp/backend/tests/test_ratelimit.py @@ -19,7 +19,7 @@ from fastapi import FastAPI from fastapi.testclient import TestClient from app.core import config -from app.core.ratelimit import RateLimitMiddleware +from app.core.ratelimit import RateLimitMiddleware, SlidingWindowLimiter @pytest.fixture @@ -112,3 +112,58 @@ def test_anonymous_key_from_rightmost_xff(client): client.get("/api/v1/ping", headers={"X-Forwarded-For": "8.8.8.8, 2.2.2.2"}).status_code == 200 ) + + +# ── SlidingWindowLimiter (#tgsupport-web) — reusable narrower per-feature limit ── + + +def test_sliding_window_limiter_retry_after_does_not_record(): + """`.retry_after()` — non-destructive peek: не расходует бюджет сам по себе + (review L3 — вызывающая сторона решает, когда `.record()`).""" + limiter = SlidingWindowLimiter(limit=1, window_s=60.0) + assert limiter.retry_after("alice") is None + # Повторный peek БЕЗ record() — бюджет не тронут, всё ещё None. + assert limiter.retry_after("alice") is None + + +def test_sliding_window_limiter_record_then_retry_after_blocks(): + limiter = SlidingWindowLimiter(limit=1, window_s=60.0) + assert limiter.retry_after("alice") is None + limiter.record("alice") + retry_after = limiter.retry_after("alice") + assert retry_after is not None + assert retry_after > 0 + + +def test_sliding_window_limiter_check_combines_peek_and_record(): + """`.check()` — комбинированный (обратной совместимости ради) peek+record.""" + limiter = SlidingWindowLimiter(limit=1, window_s=60.0) + assert limiter.check("alice") is None # первый — проходит и сразу учитывается + assert limiter.check("alice") is not None # второй — уже за лимитом + + +def test_sliding_window_limiter_per_key_isolation(): + limiter = SlidingWindowLimiter(limit=1, window_s=60.0) + limiter.record("alice") + assert limiter.retry_after("alice") is not None + assert limiter.retry_after("bob") is None # свой ключ — не задет alice + + +def test_sliding_window_limiter_prunes_empty_buckets_past_threshold(): + """review L2: пустые корзины чистятся при накоплении >10000 ключей (тот же + паттерн, что `RateLimitMiddleware.dispatch`) — не бесконечная утечка памяти. + + `retry_after()`-peek на новом ключе создаёт ПУСТУЮ корзину как побочный эффект + `defaultdict` (даже если вызывающая сторона так и не позвала `.record()` — + напр. запрос отклонён по другой причине выше по стеку). Это главный источник + "мусорных" пустых корзин, которые cleanup обязан подбирать. + """ + limiter = SlidingWindowLimiter(limit=1000, window_s=60.0) + for i in range(10001): + limiter.retry_after(f"user-{i}") + assert len(limiter._hits) == 10001 + + # Следующий record() создаёт СВОЮ (непустую) корзину и заодно подчищает + # все чужие пустые — тот же порог (>10000), что и RateLimitMiddleware. + limiter.record("trigger-cleanup") + assert len(limiter._hits) < 10001 diff --git a/tradein-mvp/backend/tests/test_support.py b/tradein-mvp/backend/tests/test_support.py index fd92484b..4b87d446 100644 --- a/tradein-mvp/backend/tests/test_support.py +++ b/tradein-mvp/backend/tests/test_support.py @@ -107,6 +107,34 @@ def test_list_messages_without_auth_header_401(client: TestClient) -> None: assert r.status_code == 401 +def test_whitespace_only_auth_header_401(client: TestClient) -> None: + """review L4: заголовок из одних пробелов после `.strip()` пуст — не должен + считаться валидным auth (не проваливается тихо в "" как username).""" + r = client.get("/api/v1/trade-in/support/messages", headers={"x-authenticated-user": " "}) + assert r.status_code == 401 + + +def test_auth_header_with_surrounding_whitespace_is_stripped( + client: TestClient, monkeypatch: pytest.MonkeyPatch +) -> None: + """review L4: лишний пробел от прокси не должен заводить ВТОРОЙ тред для, по + сути, того же пользователя — username резолвится в тред по .strip()-нутому + значению.""" + seen_usernames = [] + + def fake_find_thread_id(db, username): + seen_usernames.append(username) + return None + + monkeypatch.setattr(support_module.storage, "find_thread_id", fake_find_thread_id) + + r = client.get( + "/api/v1/trade-in/support/messages", headers={"x-authenticated-user": " alice "} + ) + assert r.status_code == 200 + assert seen_usernames == ["alice"] + + # ── validation ──────────────────────────────────────────────────────────────── @@ -182,11 +210,22 @@ def test_send_message_support_chat_id_unset_returns_503( def test_send_message_happy_path_mirrors_with_website_marker( client: TestClient, db: MagicMock, monkeypatch: pytest.MonkeyPatch, _fake_telegram_client: Any ) -> None: - monkeypatch.setattr(support_module.storage, "get_or_create_thread", lambda db, username: 7) + thread_calls: list[str] = [] + + def fake_get_or_create_thread(db, username): + thread_calls.append(username) + return 7 + + monkeypatch.setattr(support_module.storage, "get_or_create_thread", fake_get_or_create_thread) recorded = {} - def fake_record_inbound(db, *, thread_id, text_body, topic_message_id): - recorded.update(thread_id=thread_id, text_body=text_body, topic_message_id=topic_message_id) + def fake_record_inbound(db, *, thread_id, text_body, topic_message_id, support_chat_id): + recorded.update( + thread_id=thread_id, + text_body=text_body, + topic_message_id=topic_message_id, + support_chat_id=support_chat_id, + ) return { "id": 100, "direction": "in", @@ -216,9 +255,14 @@ def test_send_message_happy_path_mirrors_with_website_marker( assert "У меня вопрос про trade-in" in mirror_call["text"] assert mirror_call["chat_id"] == support_module.settings.telegram_support_chat_id assert mirror_call["message_thread_id"] == 42 + # review H1: интерактивный узкий бюджет ретраев/timeout, не воркерный дефолт. + assert mirror_call["max_retries"] == support_module._INTERACTIVE_SEND_MAX_RETRIES + assert mirror_call["timeout"] == support_module._INTERACTIVE_SEND_TIMEOUT_S assert recorded["thread_id"] == 7 assert recorded["topic_message_id"] == 555 # из FakeTelegramClient.send_message result + assert recorded["support_chat_id"] == support_module.settings.telegram_support_chat_id + assert thread_calls == ["kopylov"] assert db.commit.called @@ -226,7 +270,12 @@ def test_send_message_telegram_failure_returns_502_and_does_not_persist( client: TestClient, monkeypatch: pytest.MonkeyPatch, _fake_telegram_client: Any ) -> None: _fake_telegram_client._response = TelegramApiError("sendMessage", 400, "chat not found") - monkeypatch.setattr(support_module.storage, "get_or_create_thread", lambda db, username: 7) + thread_created = [] + monkeypatch.setattr( + support_module.storage, + "get_or_create_thread", + lambda db, username: thread_created.append(1), + ) record_called = [] monkeypatch.setattr( support_module.storage, @@ -238,6 +287,8 @@ def test_send_message_telegram_failure_returns_502_and_does_not_persist( assert r.status_code == 502 assert record_called == [] # неотправленное сообщение не персистится + # review H1: thread создаётся ПОСЛЕ успешной отправки — на неудаче до БД не доходит вообще. + assert thread_created == [] # ── rate limit ──────────────────────────────────────────────────────────────── @@ -268,6 +319,40 @@ def test_send_message_rate_limited_429(client: TestClient, monkeypatch: pytest.M assert "Retry-After" in second.headers +def test_send_message_failed_attempts_do_not_consume_rate_limit( + client: TestClient, monkeypatch: pytest.MonkeyPatch, _fake_telegram_client: Any +) -> None: + """review L3: неудачная отправка НЕ должна расходовать rate-limit бюджет — + иначе клиент, которому не повезло с транзиентной Telegram-ошибкой, терял бы + попытки, не доставив НИ ОДНОГО сообщения.""" + monkeypatch.setattr( + support_module, "_send_limiter", SlidingWindowLimiter(limit=1, window_s=60.0) + ) + _fake_telegram_client._response = TelegramApiError("sendMessage", 500, "boom") + + for _ in range(3): + r = client.post("/api/v1/trade-in/support/messages", json={"text": "hi"}, headers=_auth()) + assert r.status_code == 502 + + # "Телеграм" снова работает — бюджет (лимит=1) всё ещё цел, ни одна неудачная + # попытка выше его не тронула. + _fake_telegram_client._response = {"message_id": 555} + monkeypatch.setattr(support_module.storage, "get_or_create_thread", lambda db, username: 1) + monkeypatch.setattr( + support_module.storage, + "record_inbound", + lambda *a, **kw: { + "id": 1, + "direction": "in", + "text_body": kw["text_body"], + "operator_tg_id": None, + "created_at": "2026-07-26T00:00:00+00:00", + }, + ) + ok = client.post("/api/v1/trade-in/support/messages", json={"text": "ok"}, headers=_auth()) + assert ok.status_code == 200, ok.text + + def test_send_message_rate_limit_is_per_user( client: TestClient, monkeypatch: pytest.MonkeyPatch ) -> None: @@ -332,7 +417,7 @@ def test_list_messages_returns_thread_scoped_rows( monkeypatch.setattr( support_module.storage, "list_messages", - lambda db, *, thread_id, since_id: [ + lambda db, *, thread_id, since_id, limit: [ { "id": 1, "direction": "in", @@ -351,6 +436,24 @@ def test_list_messages_returns_thread_scoped_rows( assert body[0]["text_body"] == "hi" +def test_list_messages_passes_bounded_limit_to_storage( + client: TestClient, monkeypatch: pytest.MonkeyPatch +) -> None: + """review M5: без LIMIT каждое монтирование виджета отдавало бы весь тред.""" + seen_limits = [] + + def fake_list_messages(db, *, thread_id, since_id, limit): + seen_limits.append(limit) + return [] + + monkeypatch.setattr(support_module.storage, "find_thread_id", lambda db, username: 7) + monkeypatch.setattr(support_module.storage, "list_messages", fake_list_messages) + + r = client.get("/api/v1/trade-in/support/messages", headers=_auth()) + assert r.status_code == 200 + assert seen_limits == [support_module._LIST_MESSAGES_LIMIT] + + def test_two_users_get_independent_threads( client: TestClient, monkeypatch: pytest.MonkeyPatch ) -> None: @@ -361,7 +464,7 @@ def test_two_users_get_independent_threads( monkeypatch.setattr( support_module.storage, "list_messages", - lambda db, *, thread_id, since_id: [ + lambda db, *, thread_id, since_id, limit: [ { "id": 1, "direction": "in",