gendesign/tradein-mvp/backend/app/api/v1/support.py
bot-backend d5c876e3d0
All checks were successful
CI Trade-In / backend-tests (pull_request) Successful in 7m35s
CI Trade-In / changes (pull_request) Successful in 9s
CI / changes (pull_request) Successful in 11s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Successful in 1m10s
fix(tradein): включаем idempotency-key на фронте, чиним assert-crash и честность докстрингов (#3471)
Ревью PR #3495 нашло, что механизм был мёртвым кодом: фронт не отправлял
Idempotency-Key ни в одном запросе, весь прод-трафик шёл по ненадёжному
fallback-отпечатку. Плюс два "assert" в record_inbound после ON CONFLICT
давали AssertionError (не ловится except SQLAlchemyError) уже ПОСЛЕ
доставки в Telegram — под `python -O` assert и вовсе исчезает.

- useSupportChat.ts: useSendSupportMessage генерирует Idempotency-Key
  (crypto.randomUUID()) на намерение отправить, переиспользует его при
  повторной отправке ТОГО ЖЕ текста, сбрасывает на успехе.
- web_support_storage.record_inbound: assert -> явные ветки с логом;
  логируем отброшенный topic_message_id проигравшего гонку (не молча).
- Докстринги/комментарии переписаны честно: что именно закрывает
  pre-check (ответ клиенту потерян / двойной клик после успеха), а что
  НЕ закрывает (сетевую потерю на плече Selectel -> Telegram — там
  сообщение просто не доставлено, повтор это законная первая попытка).

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JY6iWDnGDthdvsMWgK1BMG
2026-09-12 15:24:21 +03:00

857 lines
54 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Веб-чат поддержки (#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) — чужой
тред прочитать/отметить нельзя ни при каких параметрах запроса, потому что
параметра, которым можно было бы адресовать чужой тред, попросту не существует.
Копия зеркала в топике всегда помечена "[С САЙТА] <username>: ..." — оператор
не должен путать веб-обращение с 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. Бонус: неудачная
отправка больше не создаёт тред.
Анонимная ветка (`/support/anon/*`, инцидент 2026-07-31)
-------------------------------------------------------
Ровно те же 4 действия, но БЕЗ авторизации — доступны с экрана входа. Причина:
после cutover'а на свою авторизацию (#2558) единственным каналом в поддержку был
чат ЗА логином, а самая частая причина писать в поддержку — как раз «не могу
войти». 2026-07-31 «Практика» весь день билась в форму (5 login_failed, 0
успешных) и достучаться из продукта не могла ничем.
Идентичность анонима — opaque-токен в httpOnly-куке (`_ANON_COOKIE_NAME`),
тред живёт в тех же `web_support_threads` под ключом `anon:<token>`. Двоеточие
делает коллизию с реальным логином структурно невозможной: `tradein_users`
допускает только `^[A-Za-z0-9._-]{3,64}$` (CHECK из миграции 193 + Pydantic),
двоеточия там быть не может — аноним НИКОГДА не попадёт в чужой тред и не
«станет» существующим юзером.
Изоляция тредов та же, что у авторизованной ветки, и по той же причине:
thread_id не принимается снаружи ни в каком виде, тред резолвится
ИСКЛЮЧИТЕЛЬНО из куки. Кука здесь — bearer-токен своего треда, поэтому
httpOnly+Secure+SameSite=Lax (как session-cookie) и `token_urlsafe(18)`
(144 бита) вместо чего-то угадываемого.
В Telegram-топик уходит НЕ сам токен, а `anon-<6 hex от sha256(токен)>`
(`_anon_display_id`): оператору нужен стабильный ярлык треда, а не bearer —
зеркало топика читают люди и пересылают дальше.
"""
from __future__ import annotations
import hashlib
import logging
import re
import secrets
import time
from datetime import UTC, datetime
from typing import Annotated, Literal
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
from pydantic import BaseModel, Field, field_validator
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.db import get_db
from app.core.ratelimit import SlidingWindowLimiter, _client_ip
from app.observability.metrics import SUPPORT_MESSAGES
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 TelegramError
from app.services.tgbot.shared import get_telegram_client
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с — легальный суммарный бюджет минуты). Узкий бюджет здесь: интерактивный
# клиент должен получить ответ (даже если это ошибка) за секунды, а не висеть до
# исчерпания воркерных ретраев.
#
# Числа подобраны по замеру прода 01.09.2026 (#tgsupport-retry), а не на глаз.
# Что измерено на самом хосте, из контейнера бота:
# - канал до api.telegram.org рвётся постоянно: 353 ConnectTimeout за сутки в
# логе long-polling'а; доля отказов на попытку 15-38% всплесками (пять проб:
# 3/8, 15/40, 5/20, 3/20, 1/25). Транспорт ни при чём — httpx и сырой сокет
# отваливаются одинаково (25% против 35% в чередующемся замере);
# - успешный запрос отвечает за 0.13с, максимум из 25 проб — 0.18с;
# - неудачный НИКОГДА не отваливается быстро: все отказы упираются в таймаут
# целиком (10.02с при timeout=10.0), быстрых — ноль;
# - через `SCRAPER_PROXY_URL` не легчает, а хуже: 0 из 20. Прокси не решение.
#
# Отсюда три следствия. Таймаут 10с был чистой платой за неудачу (успех не занимает
# и половины секунды) — снижен до 5с, это ~28-кратный запас к измеренному максимуму.
# Одного повтора мало: при 30% отказов на попытку до пользователя доходило ~9%
# отказов, то есть каждое одиннадцатое сообщение терялось с 502; три повтора уводят
# это к ~0.8%. Экспоненциальная пауза 2→4→8с осмысленна против троттлинга, но здесь
# отказ — неустановленное соединение, паузе нечего пережидать, поэтому потолок 1с:
# он же не даёт `retry_after` из 429 (для группы штатные 30-60с) умножиться на число
# попыток. Худший случай: 4 попытки × 5с + 3 паузы × 1с = 23с, и он требует четырёх
# отказов подряд. Типичный случай не меняется — 0.13с.
_INTERACTIVE_SEND_TIMEOUT_S = 5.0
_INTERACTIVE_SEND_MAX_RETRIES = 3
_INTERACTIVE_SEND_MAX_BACKOFF_S = 1.0
# Счётчик ОТКАЗОВ отправки — отдельный от бюджетов выше (#tgsupport-fail-cooldown).
#
# Зачем вообще второй счётчик. `_send_limiter`/`_anon_ip_limiter` расходуются
# ТОЛЬКО на успехе (review L3, non-destructive peek выше) — это верно для
# «не наказывать за чужую аварию», но имеет обратную сторону: пока Telegram
# недоступен, лимита нет ВООБЩЕ. Каждый повтор пользователя при этом стоит до
# `1 + _INTERACTIVE_SEND_MAX_RETRIES` = 4 попыток к api.telegram.org и не
# расходует ни один бюджет. Двух-трёх вкладок с авто-ретраем хватает, чтобы
# выесть лимиты группы ровно в тот момент, когда канал и так еле жив.
#
# Отсюда — дешёвый gate ПЕРЕД походом в Telegram, на своих ключах (тех же, что
# у основных лимитеров: username / anon-thread-key и client IP).
#
# N=5. Отказ одной отправки — не событие: по замеру #tgsupport-retry доля отказов
# на попытку 15-38%, но после 4 попыток до пользователя доходит ~0.8-2% отказов.
# Пять подряд на живом канале — вероятность порядка 1e-10, то есть cooldown
# физически не может сработать на «просто не повезло»; он срабатывает только на
# настоящей недоступности. Плюс `reset()` на успехе: считаем именно ПОДРЯД, одна
# успешная отправка стирает историю.
#
# 30с окна (оно же длительность cooldown: блок держится, пока самый старый из
# N отказов не выпадет из окна). Верхняя граница осмысленности — реальная
# недоступность Telegram по тому же замеру длится минутами, так что 30с заведомо
# короче и пользователя после восстановления канала не наказывают. Нижняя —
# один отказавший интерактивный запрос сам по себе занимает до 23с
# (4 попытки × 5с + 3 паузы × 1с), cooldown короче этого просто не имел бы смысла.
_SEND_FAILURE_LIMIT = 5
_SEND_FAILURE_WINDOW_S = 30.0
_send_failure_limiter = SlidingWindowLimiter(
limit=_SEND_FAILURE_LIMIT, window_s=_SEND_FAILURE_WINDOW_S
)
# Отдельный экземпляр для IP-ключей — ровно как `_anon_ip_limiter` отдельный от
# `_send_limiter`: ключи разных пространств (IP vs thread-key) в одном словаре
# смешивать нельзя.
_anon_ip_failure_limiter = SlidingWindowLimiter(
limit=_SEND_FAILURE_LIMIT, window_s=_SEND_FAILURE_WINDOW_S
)
_SEND_UNAVAILABLE_DETAIL = "Telegram сейчас недоступен. Попробуйте через полминуты."
# #tgsupport-web review M5: без LIMIT каждое монтирование виджета на старом
# треде отдавало бы ВЕСЬ лог переписки. См. `web_support_storage.list_messages`.
_LIST_MESSAGES_LIMIT = 200
# Идемпотентность отправки (#3471 retry-storm) — см. `_resolve_idempotency_key`.
_IDEMPOTENCY_HEADER = "idempotency-key"
# Форма клиентского ключа — как у anon-токена (`_ANON_TOKEN_RE`): произвольная
# opaque-строка клиента, без пробелов/спецсимволов, которые попали бы в SQL-параметр
# как есть. Не матчится — считаем заголовок отсутствующим и уходим на fallback,
# а не пытаемся его "починить" (тот же принцип, что у `_read_anon_token`).
_CLIENT_IDEMPOTENCY_KEY_RE = re.compile(r"^[A-Za-z0-9_.-]{8,128}\Z")
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
# Дошло ли сообщение до БД. `True` по умолчанию — все существующие пути
# (`GET /support/messages`, успешный POST) строят модель из storage-строки и
# ничего про флаг не знают, контракт для них не меняется.
#
# `False` ставит ТОЛЬКО `_unpersisted_message_out`: доставлено оператору, но
# не записано. Флаг нужен потому, что без него деградация неотличима от
# тишины: фронт выбрасывает тело POST и рендерит транскрипт исключительно из
# `GET /support/messages` (`useSupportChat.ts`), где этого сообщения нет —
# поле ввода очищается, сообщение не появляется, ошибки нет. Пользователь
# решает, что не отправилось, и шлёт снова — ровно тот дубль в топике,
# против которого вся эта ветка и сделана.
persisted: bool = True
@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}"
def _too_many_failures_error(retry_after: float) -> HTTPException:
"""429 вместо похода в Telegram — канал только что отказал N раз подряд."""
return HTTPException(
status_code=429,
detail=_SEND_UNAVAILABLE_DETAIL,
headers={"Retry-After": str(int(retry_after) + 1)},
)
def _unpersisted_message_out(text_body: str) -> SupportMessageOut:
"""Синтетический ответ для случая «в топик доставлено, а в БД не записано».
Почему ответ вообще УСПЕШНЫЙ. Сообщение оператору реально доставлено —
отдать на это ошибку значит соврать: пользователь повторит, и в топике
окажется дубль (плюс второе предупреждение оператору). Успех здесь честнее
отказа, но он неполный, и об этом клиенту надо сказать явно.
Что сообщает `persisted=False`: «доставлено, но ответить тебе через чат не
смогут; повторять не надо». Именно флаг — контракт для клиента, а НЕ `id=0`:
по идентификатору клиент отличить деградацию не обязан и не будет.
`id` при этом всё равно нужен модели, и 0 — сознательный сентинел, а не
выдумка: реальные id в `web_support_messages` начинаются с 1 (serial),
поэтому 0 ни с чем не столкнётся, а `GET /support/messages` принимает
`since>=0`. Даже если фронт когда-нибудь начнёт курсорить по возвращённому
id (сейчас он перечитывает тред целиком с `since=0`, `useSupportChat.ts`),
нулевой курсор не перепрыгнет ни одного реального сообщения: занижение
курсора безопасно, завышение — нет.
"""
return SupportMessageOut(
id=0,
direction="in",
text_body=text_body,
operator_tg_id=None,
created_at=datetime.now(UTC),
persisted=False,
)
async def _warn_operator_message_not_recorded(*, label: str, topic_message_id: int | None) -> None:
"""Реплаем к только что доставленному зеркалу предупреждает оператора, что
ответ на ЭТО сообщение не смаршрутизируется: треда в БД нет, на реплай к
осиротевшему зеркалу `bridge._handle_group_reply` напишет только WARNING, а
клиент не увидит ничего. Без предупреждения оператор отвечает в пустоту и
считает, что помог.
Текст обращения сюда НЕ дублируется (ПДн) — оператор видит его в сообщении,
к которому это реплай.
Формулировка учитывает, что клиенту ответили успехом (`persisted=False`,
«принято, повторять не надо»): рассчитывать на «клиент напишет снова»
оператору нельзя, поднимать тред придётся иначе.
Любой отказ гасится логом: основное сообщение УЖЕ доставлено, превращать
неудачу служебного уведомления в 500 поверх успеха нельзя.
"""
text = (
f"ВНИМАНИЕ · {label}: сообщение доставлено в топик, но НЕ записано в базу "
"(сбой БД). Треда у этого обращения нет — ответить клиенту через бота "
"НЕЛЬЗЯ: реплай на это сообщение никуда не уйдёт. Повтора тоже не ждите, "
"клиенту показано, что сообщение принято и отправлять его снова не нужно. "
"Если в обращении есть контакт — свяжитесь напрямую; иначе передайте "
"дежурному и сообщите об отказе БД."
)
try:
client = get_telegram_client()
await client.send_message(
chat_id=settings.telegram_support_chat_id,
text=text,
message_thread_id=settings.telegram_support_topic_id or None,
# reply_to может отсутствовать (sendMessage не вернул message_id) —
# тогда уведомление уходит отдельным сообщением в топик: хуже, чем
# реплай, но несравнимо лучше тишины.
reply_to_message_id=topic_message_id,
# Тот же узкий интерактивный бюджет (review H1): клиент ждёт ответа
# ручки, а не доставки служебного уведомления.
timeout=_INTERACTIVE_SEND_TIMEOUT_S,
max_retries=_INTERACTIVE_SEND_MAX_RETRIES,
max_backoff=_INTERACTIVE_SEND_MAX_BACKOFF_S,
)
except Exception:
# Шире `TelegramError` намеренно: это best-effort хвост уже успешного
# запроса, любой отказ здесь обязан остаться в логе, а не у клиента.
logger.exception(
"web support: не удалось предупредить оператора о несохранённом сообщении (%s)",
label,
)
def _rollback_quietly(db: Session) -> None:
"""Откат после `SQLAlchemyError`. Сам rollback на мёртвом соединении тоже
может бросить, а мы уже решили отдать клиенту успех — `get_db` в любом
случае закроет сессию в `finally`."""
try:
db.rollback()
except SQLAlchemyError:
logger.warning("web support: rollback после сбоя БД тоже не удался", exc_info=True)
def _resolve_idempotency_key(request: Request, *, identity_key: str, text: str) -> str:
"""Ключ идемпотентности inbound-отправки (#3471, миграция 301).
Почему заголовок + fallback, а не что-то одно. Явный `Idempotency-Key` —
ЕДИНСТВЕННЫЙ надёжный путь: клиент генерирует ключ ОДИН раз на "намерение
отправить" и переиспользует его на любом ретрае (fetch retry / переотправка
после таймаута) независимо от того, что именно менялось в UI между
попытками (см. `useSendSupportMessage` во фронтенде — реальный трафик
обязан идти этим путём). Fallback без заголовка — детерминированный
отпечаток sha256(identity, текст, минутное окно) — существует ТОЛЬКО для
клиентов, которые заголовок не прислали (п.5 требования — старое поведение
не должно сломаться), и у него два честных изъяна, оба снимаются самим
фактом использования заголовка:
1. два РАЗНЫХ по смыслу сообщения с одинаковым текстом от одного и того
же человека в течение одной минуты ("да" и ещё раз "да", "+", "ок")
СХЛОПНУТСЯ в одно — второе будет молча проглочено, клиент получит id
первого, оператор не увидит второе сообщение вообще;
2. окно — не скользящие 60 секунд, а округление `time.time() // 60` вниз:
фактический срок дедупликации случаен от 0 до 60с в зависимости от
момента внутри минуты, и часть настоящих повторов (ретрай ровно на
границе окна) эту защиту не получит.
Оба пункта перестают быть важны, когда клиент реально шлёт заголовок
(см. `useSendSupportMessage`).
Хранит и потенциально логирует (см. вызовы в этом файле) только САМ ключ —
он либо непрозрачный клиентский токен, либо хэш. Текст сообщения сюда
попадает ТОЛЬКО как вход в sha256, в открытом виде не сохраняется и не
возвращается — жёсткое правило проекта "текст не в логах" (ПДн) при этом
не нарушается, даже если бы этот ключ где-то залогировали.
Префиксы `client:`/`auto:` разводят два namespace'а ключей (защита от
случайного совпадения клиентского токена с fallback-хэшем) — дёшево и не
требует отдельной колонки.
"""
header = (request.headers.get(_IDEMPOTENCY_HEADER) or "").strip()
if header and _CLIENT_IDEMPOTENCY_KEY_RE.match(header):
return f"client:{header}"
window = int(time.time() // 60)
fingerprint = f"{identity_key}|{text}|{window}"
return "auto:" + hashlib.sha256(fingerprint.encode("utf-8")).hexdigest()
@router.post("/support/messages", response_model=SupportMessageOut)
async def send_support_message(
request: Request,
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)},
)
# Cooldown по ОТКАЗАМ (см. `_SEND_FAILURE_LIMIT`): канал только что отказал
# N раз подряд — не тратим на этот запрос ещё четыре попытки к Telegram.
cooldown = _send_failure_limiter.retry_after(username)
if cooldown is not None:
raise _too_many_failures_error(cooldown)
# Идемпотентность (#3471). Что именно закрывает этот pre-check — важно не
# переоценить: строка в БД (и её idempotency_key) появляется ТОЛЬКО ПОСЛЕ
# успешной доставки в Telegram (H1 ниже), поэтому потерю самого запроса на
# плече Selectel -> api.telegram.org этот механизм НЕ дедуплицирует — если
# `send_message` не удался, ключ нигде не записан, и повтор клиента после
# неудачи это законная первая попытка. Реально закрывается другой, тоже
# частый случай: доставка УЖЕ состоялась (Telegram принял, строка
# закоммичена), но ответ до клиента не дошёл (обрыв на обратном пути,
# клиентский таймаут) или пользователь кликнул "отправить" второй раз по
# той же ещё не отрисовавшейся отправке — тогда pre-check находит уже
# записанное сообщение по ключу и отдаёт его же, не уходя в Telegram снова.
# Тред может ещё не существовать (это первая попытка этого ключа) — тогда
# сравнивать не с чем, идём в Telegram как обычно; сам факт "нет треда"
# здесь безопасен, т.к. тред создаётся ТОЛЬКО в этой же ручке (см. H1) —
# если бы сообщение под этим ключом уже было записано, тред уже был бы.
idempotency_key = _resolve_idempotency_key(request, identity_key=username, text=payload.text)
existing_thread_id = storage.find_thread_id(db, username)
if existing_thread_id is not None:
existing = storage.find_inbound_by_idempotency_key(
db, thread_id=existing_thread_id, idempotency_key=idempotency_key
)
if existing is not None:
logger.info(
"web support: repeat send (idempotency key already recorded) "
"username=%s thread_id=%d message_id=%s",
username,
existing_thread_id,
existing["id"],
)
return SupportMessageOut(**existing)
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
client = get_telegram_client()
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,
max_backoff=_INTERACTIVE_SEND_MAX_BACKOFF_S,
)
except TelegramError:
# НЕ логируем payload.text (переписка — ПДн) и НЕ логируем токен (его в
# TelegramApiError и не бывает — см. client.py docstring про redaction).
# Предок, а не `TelegramApiError`: при таймауте до Telegram пользователь
# должен увидеть тот же «сервис недоступен», а не 500 (#3456).
logger.exception(
"web support: не удалось отправить зеркало в топик (username=%s)", username
)
_send_failure_limiter.record(username)
raise HTTPException(status_code=502, detail=SERVICE_UNAVAILABLE_TEXT) from None
# Отправка удалась — теперь и только теперь расходуем rate-limit бюджет.
_send_limiter.record(username)
SUPPORT_MESSAGES.labels(channel="web").inc()
# Канал жив — счётчик отказов считает именно ПОДРЯД идущие отказы.
_send_failure_limiter.reset(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 модуля.
#
# Оборотная сторона этого порядка: отказ БД здесь означает, что сообщение
# оператору УЖЕ доставлено. Отдавать на это 500 (как было) — худший из
# вариантов: клиент видит ошибку, шлёт повторно, в топике дубль, а на
# осиротевшее зеркало оператор отвечает в пустоту. Поэтому отвечаем успехом
# (доставка правда состоялась) и отдельным сообщением предупреждаем
# оператора, что отвечать на это зеркало бесполезно.
try:
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,
idempotency_key=idempotency_key,
)
db.commit()
except SQLAlchemyError:
# Ни текста обращения (ПДн), ни токена в логе — только кто и что сломалось.
logger.exception(
"web support: сообщение доставлено в топик, но не записано в БД (username=%s)",
username,
)
_rollback_quietly(db)
await _warn_operator_message_not_recorded(
label=f"[С САЙТА] {username}", topic_message_id=topic_message_id
)
return _unpersisted_message_out(payload.text)
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()
# ---------------------------------------------------------------------------
# Анонимная ветка — поддержка без входа (см. блок в докстринге модуля)
# ---------------------------------------------------------------------------
_ANON_COOKIE_NAME = "tradein_support_anon"
# 30 дней: тред должен пережить «напишу вечером — отвечут утром», но не жить вечно.
_ANON_COOKIE_MAX_AGE_S = 30 * 24 * 3600
# Двоеточие → структурная невозможность коллизии с реальным логином (докстринг).
_ANON_THREAD_PREFIX = "anon:"
# Форма того, что МЫ выдаём (`token_urlsafe(18)` → 24 символа из [A-Za-z0-9_-]).
# Кука клиент-контролируема: без этой проверки в ключ треда (а значит в SQL-параметр
# и в лог) уехала бы произвольная строка из браузера. Не матчится — считаем куку
# отсутствующей и выдаём новую, а не пытаемся «починить» присланное.
_ANON_TOKEN_RE = re.compile(r"^[A-Za-z0-9_-]{16,64}\Z")
# Публичная ручка записи в общий Telegram-топик — поверхность для спама, которой у
# авторизованной ветки нет. Два независимых бюджета:
# 1) per-token (`_send_limiter`, 12/мин — тот же объект, ключи не пересекаются:
# анонимные начинаются с "anon:", что невозможно для username);
# 2) per-IP — именно он ловит обход ротацией куки (сбросил куку → новый токен →
# бюджет (1) снова пуст). Окно широкое и щедрое для живого диалога: реальный
# сценарий — «не могу войти, помогите», несколько сообщений подряд.
_ANON_IP_RATE_LIMIT = 10
_ANON_IP_RATE_WINDOW_S = 600.0
_anon_ip_limiter = SlidingWindowLimiter(limit=_ANON_IP_RATE_LIMIT, window_s=_ANON_IP_RATE_WINDOW_S)
def _anon_display_id(token: str) -> str:
"""Стабильный НЕсекретный ярлык треда для оператора — см. докстринг модуля.
sha256, а не префикс токена: префикс — это часть bearer'а, а зеркало уходит
в Telegram-топик, который читают люди и пересылают дальше.
"""
return f"anon-{hashlib.sha256(token.encode('utf-8')).hexdigest()[:6]}"
def _read_anon_token(request: Request) -> str | None:
"""Токен из куки, если он валидной формы; иначе None (кука считается отсутствующей)."""
raw = request.cookies.get(_ANON_COOKIE_NAME)
if raw is None or not _ANON_TOKEN_RE.match(raw):
return None
return raw
def _anon_thread_key(token: str) -> str:
return f"{_ANON_THREAD_PREFIX}{token}"
def _set_anon_cookie(response: Response, token: str) -> None:
response.set_cookie(
key=_ANON_COOKIE_NAME,
value=token,
max_age=_ANON_COOKIE_MAX_AGE_S,
httponly=True,
secure=True,
samesite="lax",
path="/",
)
@router.post("/support/anon/messages", response_model=SupportMessageOut)
async def send_anon_support_message(
payload: SupportMessageInput,
request: Request,
response: Response,
db: Annotated[Session, Depends(get_db)],
) -> SupportMessageOut:
"""Сообщение в поддержку БЕЗ входа. Порядок операций — как в авторизованной
ветке (H1 в докстринге модуля): БД трогаем только после успешного sendMessage.
Кука выставляется тоже только на успехе — иначе первая же неудачная попытка
(бот не настроен / Telegram лёг) закрепляла бы за посетителем пустой тред.
"""
if not _bot_configured():
raise HTTPException(status_code=503, detail=SERVICE_UNAVAILABLE_TEXT)
token = _read_anon_token(request)
is_new_token = token is None
if token is None:
token = secrets.token_urlsafe(18)
thread_key = _anon_thread_key(token)
ip = _client_ip(request)
# Оба бюджета — non-destructive peek (review L3): неудачная отправка не
# должна стоить посетителю попытки. `.record()` только на успех, ниже.
for retry_after in (_send_limiter.retry_after(thread_key), _anon_ip_limiter.retry_after(ip)):
if retry_after is not None:
raise HTTPException(
status_code=429,
detail="Слишком много сообщений. Попробуйте позже.",
headers={"Retry-After": str(int(retry_after) + 1)},
)
# Cooldown по ОТКАЗАМ (см. `_SEND_FAILURE_LIMIT`) — оба ключа, как и у
# бюджетов выше: per-token ловит одну вкладку с авто-ретраем, per-IP — ту же
# петлю после сброса куки.
for cooldown in (
_send_failure_limiter.retry_after(thread_key),
_anon_ip_failure_limiter.retry_after(ip),
):
if cooldown is not None:
raise _too_many_failures_error(cooldown)
display_id = _anon_display_id(token)
# Идемпотентность (#3471) — см. развёрнутый комментарий в авторизованной
# ветке. Ограничение специфично для анонимной ветки: если предыдущая
# попытка сама провалилась ДО выдачи куки (`is_new_token=True` тогда и
# сейчас), identity_key каждый раз новый (случайный токен) и pre-check
# структурно не может найти прошлую попытку — это тот же класс проблемы,
# что и потеря куки в сети, вне scope этой правки.
idempotency_key = _resolve_idempotency_key(request, identity_key=thread_key, text=payload.text)
existing_thread_id = storage.find_thread_id(db, thread_key)
if existing_thread_id is not None:
existing = storage.find_inbound_by_idempotency_key(
db, thread_id=existing_thread_id, idempotency_key=idempotency_key
)
if existing is not None:
logger.info(
"web support (anon): repeat send (idempotency key already recorded) "
"%s thread_id=%d message_id=%s",
display_id,
existing_thread_id,
existing["id"],
)
if is_new_token:
_set_anon_cookie(response, token)
return SupportMessageOut(**existing)
# Общий клиент приложения (#tg-connection-resilience): на каждый запрос
# свой создавать нельзя — это ноль keep-alive и полный TCP+TLS-хендшейк
# до api.telegram.org перед каждой отправкой. Живёт в lifespan.
client = get_telegram_client()
try:
mirrored = await client.send_message(
chat_id=settings.telegram_support_chat_id,
text=_format_anon_mirror_text(display_id, payload.text),
message_thread_id=settings.telegram_support_topic_id or None,
timeout=_INTERACTIVE_SEND_TIMEOUT_S,
max_retries=_INTERACTIVE_SEND_MAX_RETRIES,
max_backoff=_INTERACTIVE_SEND_MAX_BACKOFF_S,
)
except TelegramError:
# Ни текст сообщения (ПДн), ни токен (bearer треда) в лог не попадают.
# Про предок вместо `TelegramApiError` — см. комментарий в парной ручке.
logger.exception(
"web support (anon): не удалось отправить зеркало в топик (%s)", display_id
)
_send_failure_limiter.record(thread_key)
_anon_ip_failure_limiter.record(ip)
raise HTTPException(status_code=502, detail=SERVICE_UNAVAILABLE_TEXT) from None
_send_limiter.record(thread_key)
_anon_ip_limiter.record(ip)
SUPPORT_MESSAGES.labels(channel="anon").inc()
# Канал жив — счётчики отказов считают именно ПОДРЯД идущие отказы.
_send_failure_limiter.reset(thread_key)
_anon_ip_failure_limiter.reset(ip)
topic_message_id = mirrored.get("message_id") if isinstance(mirrored, dict) else None
if topic_message_id is None:
logger.warning(
"web support (anon): Telegram sendMessage не вернул message_id (%s) — "
"ответ оператора на это сообщение не будет смаршрутизирован",
display_id,
)
# Отказ БД после успешной отправки — см. развёрнутый комментарий в парной
# (авторизованной) ручке: сообщение оператору уже доставлено, 500 тут создаёт
# дубли в топике и «осиротевшее» зеркало, на которое оператор отвечает зря.
try:
thread_id = storage.get_or_create_thread(db, thread_key)
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,
idempotency_key=idempotency_key,
)
db.commit()
except SQLAlchemyError:
# Ни текста обращения (ПДн), ни токена (bearer треда) в логе — только
# несекретный ярлык треда.
logger.exception(
"web support (anon): сообщение доставлено в топик, но не записано в БД (%s)",
display_id,
)
_rollback_quietly(db)
await _warn_operator_message_not_recorded(
label=f"[С САЙТА · БЕЗ ВХОДА] {display_id}", topic_message_id=topic_message_id
)
# Куку ставим ВСЁ РАВНО: тред в БД не создан, но идентичность посетителя
# обязана пережить этот сбой — иначе следующее сообщение (когда база
# поднимется) заведёт ВТОРОЙ тред, и переписка разъедется на два.
if is_new_token:
_set_anon_cookie(response, token)
return _unpersisted_message_out(payload.text)
if is_new_token:
_set_anon_cookie(response, token)
logger.info("web support (anon): message sent %s thread_id=%d", display_id, thread_id)
return SupportMessageOut(**row)
def _format_anon_mirror_text(display_id: str, message_text: str) -> str:
"""Помечает зеркало как пришедшее с сайта ОТ НЕЗАЛОГИНЕННОГО посетителя.
Оператору это ключевой контекст: у такого обращения нет аккаунта, по которому
можно посмотреть историю, и самая вероятная причина написать — как раз
невозможность войти (инцидент 2026-07-31).
"""
return f"[С САЙТА · БЕЗ ВХОДА] {display_id}:\n{message_text}"
@router.get("/support/anon/messages", response_model=list[SupportMessageOut])
def list_anon_support_messages(
request: Request,
db: Annotated[Session, Depends(get_db)],
since: Annotated[int, Query(ge=0)] = 0,
) -> list[SupportMessageOut]:
"""Свой тред по куке. Нет куки / нет треда → пустой список, НЕ 401: виджет
поллит эту ручку и до первого сообщения, 401 там был бы ложной ошибкой.
Sync `def` (review M3) — см. `list_support_messages`.
"""
token = _read_anon_token(request)
if token is None:
return []
thread_id = storage.find_thread_id(db, _anon_thread_key(token))
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/anon/unread", response_model=UnreadOut)
def get_anon_support_unread(
request: Request,
db: Annotated[Session, Depends(get_db)],
) -> UnreadOut:
"""Sync `def` (review M3) — см. `list_support_messages`."""
token = _read_anon_token(request)
if token is None:
return UnreadOut(unread=0)
thread_id = storage.find_thread_id(db, _anon_thread_key(token))
if thread_id is None:
return UnreadOut(unread=0)
return UnreadOut(unread=storage.count_unread(db, thread_id=thread_id))
@router.post("/support/anon/read", response_model=StatusOut)
def mark_anon_support_read(
request: Request,
db: Annotated[Session, Depends(get_db)],
) -> StatusOut:
"""Sync `def` (review M3) — см. `list_support_messages`."""
token = _read_anon_token(request)
if token is None:
return StatusOut()
thread_id = storage.find_thread_id(db, _anon_thread_key(token))
if thread_id is not None:
storage.mark_read(db, thread_id=thread_id)
db.commit()
return StatusOut()