"""Веб-чат поддержки (#tgsupport-web) — поверх уже существующего Telegram support-моста (`app.services.tgbot.bridge`, data/sql/186_tg_support.sql). Источник обращения — сайт (не Telegram-личка клиента): пользователь пишет через это API, сообщение зеркалится `sendMessage`-ом в тот же support-топик, оператор отвечает РЕПЛАЕМ ровно так же, как на Telegram-клиента — маршрутизация ответа обратно реализована в `bridge._handle_group_reply` (ветка добавлена там же, без изменения существующего Telegram-пути). Изоляция тредов: все 4 ручки резолвят тред ИСКЛЮЧИТЕЛЬНО по `X-Authenticated-User` (rbac_guard в app/main.py гарантирует его наличие и валидность для non-public путей). thread_id НИКОГДА не принимается снаружи (ни в query, ни в body) — чужой тред прочитать/отметить нельзя ни при каких параметрах запроса, потому что параметра, которым можно было бы адресовать чужой тред, попросту не существует. Копия зеркала в топике всегда помечена "[С САЙТА] : ..." — оператор не должен путать веб-обращение с 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 import logging from typing import Annotated, Literal from fastapi import APIRouter, Depends, HTTPException, Query, Request from pydantic import BaseModel, Field, field_validator from sqlalchemy.orm import Session from app.core.config import settings from app.core.db import get_db from app.core.ratelimit import SlidingWindowLimiter from app.services.tgbot import web_support_storage as storage from app.services.tgbot.bridge import SERVICE_UNAVAILABLE_TEXT from app.services.tgbot.client import TelegramApiError, TelegramClient logger = logging.getLogger(__name__) router = APIRouter() # Лимит Telegram sendMessage (4096) с запасом — см. #tgsupport-web AC ("~4000"). MAX_MESSAGE_LENGTH = 4000 # Жёстче общего RateLimitMiddleware (300 req/60с на пользователя, app/main.py): # бот-токен общий на ВСЕХ клиентов веб-чата, флуд одного клиента иначе может # упереться в 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). `.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 def _bot_configured() -> bool: """TELEGRAM_BOT_TOKEN и TELEGRAM_SUPPORT_CHAT_ID оба обязательны — без них зеркалировать в топик некуда (см. app/tgbot_main.py._should_run для токена, bridge.py для chat_id).""" return bool(settings.telegram_bot_token) and bool(settings.telegram_support_chat_id) class SupportMessageInput(BaseModel): text: str = Field(min_length=1, max_length=MAX_MESSAGE_LENGTH) @field_validator("text") @classmethod def _not_blank(cls, value: str) -> str: stripped = value.strip() if not stripped: raise ValueError("text must not be blank") return stripped class SupportMessageOut(BaseModel): id: int direction: Literal["in", "out"] text_body: str operator_tg_id: int | None = None created_at: str @field_validator("created_at", mode="before") @classmethod def _isoformat(cls, value: object) -> str: if hasattr(value, "isoformat"): return value.isoformat() # type: ignore[no-any-return] return str(value) class UnreadOut(BaseModel): unread: int class StatusOut(BaseModel): status: Literal["ok"] = "ok" def _format_mirror_text(username: str, message_text: str) -> str: """Помечает зеркало как пришедшее С САЙТА, от какого пользователя — оператор иначе не отличит веб-обращение от Telegram-клиента (#tgsupport-web AC). Использует ТОЛЬКО username — thread_id здесь не нужен (см. H1 в docstring модуля), это то, что делает возможным отложить БД-запись до после отправки.""" return f"[С САЙТА] {username}:\n{message_text}" @router.post("/support/messages", response_model=SupportMessageOut) async def send_support_message( payload: SupportMessageInput, username: Annotated[str, Depends(_require_username)], db: Annotated[Session, Depends(get_db)], ) -> SupportMessageOut: """Отправляет сообщение от лица *username* в support-топик (`sendMessage` — не `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) # review L3: peek без расхода бюджета — неудачная отправка НЕ должна стоить # пользователю попытки (расходуем `.record()` только на успех, ниже). retry_after = _send_limiter.retry_after(username) if retry_after is not None: raise HTTPException( status_code=429, detail="Слишком много сообщений. Попробуйте через минуту.", headers={"Retry-After": str(int(retry_after) + 1)}, ) 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 (переписка — ПДн) и НЕ логируем токен (его в # TelegramApiError и не бывает — см. client.py docstring про redaction). logger.exception( "web support: не удалось отправить зеркало в топик (username=%s)", username ) raise HTTPException(status_code=502, detail=SERVICE_UNAVAILABLE_TEXT) from 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() logger.info("web support: message sent username=%s thread_id=%d", username, thread_id) return SupportMessageOut(**row) @router.get("/support/messages", response_model=list[SupportMessageOut]) 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, limit=_LIST_MESSAGES_LIMIT ) return [SupportMessageOut(**r) for r in rows] @router.get("/support/unread", response_model=UnreadOut) 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) return UnreadOut(unread=storage.count_unread(db, thread_id=thread_id)) @router.post("/support/read", response_model=StatusOut) 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) db.commit() return StatusOut()