fix(tradein/support): address deep-review H1/M1/M2/M3/M5/L1-L5 (web chat backend)

H1 (critical): reorder send_support_message — get_or_create_thread now runs
AFTER a successful Telegram send, not before. Prod runs a single uvicorn
process with no --workers on a sync SQLAlchemy engine sharing one event
loop; holding a row-lock across the Telegram call (which can legally take
minutes on 429/5xx worker-grade retries) risked freezing the entire API on
a second concurrent request from the same user. send_message now also
accepts explicit timeout/max_retries so the interactive endpoint uses a
bounded budget instead of inheriting the long-polling worker's retry policy.

M1: web_support_messages and tg_support_messages both gain a support_chat_id
column (188 migration for the already-applied tg_support_messages table;
187 edited in place since it hasn't shipped yet). bridge._handle_group_reply
now resolves BOTH the Telegram and web candidate under the CURRENT
support_chat_id and refuses delivery loudly (logger.error) if both match,
instead of silently preferring the Telegram path — closing a cross-table
misroute risk that would surface if the support group is ever recreated.

M2: a media reply to a web-mirrored message (including photos with a
caption) is now refused in full with an explicit notice back in the topic,
instead of silently delivering just the caption text or logging a WARNING
nobody sees.

M3: list_support_messages/get_support_unread/mark_support_read switched
from async def to def — their bodies are pure sync psycopg calls; running
them as async def executed blocking DB work directly on the event loop.

M5: list_messages now takes a bounded LIMIT (last N, chronological) so a
long-lived thread doesn't return its entire history on every poll.

L1-L4: warn when Telegram doesn't return a message_id (routing dead-end),
rate limit is peeked before send and only recorded on success (failed
attempts no longer burn the budget), SlidingWindowLimiter gained the same
empty-bucket cleanup as RateLimitMiddleware, and _require_username now
strips whitespace so a proxy-injected space can't fork a second thread.

L5: removed 187 from _manifest_applied.txt — the deploy pipeline tracks
applied migrations via _schema_migrations, not this file, and keeping an
unmerged migration name out of it preserves the option to rename before
merge without tripping the "can't rename applied migrations" test.
This commit is contained in:
bot-backend 2026-07-26 23:33:09 +03:00
parent 7377fb61e5
commit 5eadae1e95
11 changed files with 634 additions and 117 deletions

View file

@ -15,6 +15,20 @@ support-моста (`app.services.tgbot.bridge`, data/sql/186_tg_support.sql).
Копия зеркала в топике всегда помечена "[С САЙТА] <username>: ..." оператор Копия зеркала в топике всегда помечена "[С САЙТА] <username>: ..." оператор
не должен путать веб-обращение с Telegram-клиентом (#tgsupport-web AC). не должен путать веб-обращение с 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 from __future__ import annotations
@ -42,20 +56,35 @@ MAX_MESSAGE_LENGTH = 4000
# Жёстче общего RateLimitMiddleware (300 req/60с на пользователя, app/main.py): # Жёстче общего RateLimitMiddleware (300 req/60с на пользователя, app/main.py):
# бот-токен общий на ВСЕХ клиентов веб-чата, флуд одного клиента иначе может # бот-токен общий на ВСЕХ клиентов веб-чата, флуд одного клиента иначе может
# упереться в Telegram-лимиты (`sendMessage` 429) и застопорить доставку всем # упереться в Telegram-лимиты (`sendMessage` 429, лимит группы ~20 msg/min) и
# остальным (см. задачу, п.6). 12 сообщений/минуту — щедро для живого диалога, # застопорить доставку всем остальным (см. задачу, п.6). 12 сообщений/минуту —
# но режет скрипт-флуд на порядок раньше общего API-лимита. # щедро для живого диалога, но режет скрипт-флуд на порядок раньше общего API-лимита.
_SEND_RATE_LIMIT = 12 _SEND_RATE_LIMIT = 12
_SEND_RATE_WINDOW_S = 60.0 _SEND_RATE_WINDOW_S = 60.0
_send_limiter = SlidingWindowLimiter(limit=_SEND_RATE_LIMIT, window_s=_SEND_RATE_WINDOW_S) _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: def _require_username(request: Request) -> str:
"""Достаёт X-Authenticated-User. rbac_guard (app/main.py) уже гарантирует его """Достаёт X-Authenticated-User. rbac_guard (app/main.py) уже гарантирует его
наличие в проде для non-public путей этот guard здесь defence-in-depth и наличие в проде для non-public путей этот guard здесь defence-in-depth и
делает роутер тестируемым без поднятия всего app.main (см. tests/test_support.py, делает роутер тестируемым без поднятия всего app.main (см. tests/test_support.py,
как test_trade_in_lead.py для /lead).""" как test_trade_in_lead.py для /lead). `.strip()` (review L4) лишний пробел
username = request.headers.get("x-authenticated-user") от прокси иначе завёл бы ВТОРОЙ тред на, по сути, того же пользователя
(username UNIQUE ключ треда, "alice" != "alice ")."""
username = (request.headers.get("x-authenticated-user") or "").strip()
if not username: if not username:
raise HTTPException(status_code=401, detail="no authenticated user") raise HTTPException(status_code=401, detail="no authenticated user")
return username return username
@ -105,7 +134,9 @@ class StatusOut(BaseModel):
def _format_mirror_text(username: str, message_text: str) -> str: 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}" return f"[С САЙТА] {username}:\n{message_text}"
@ -116,14 +147,21 @@ async def send_support_message(
db: Annotated[Session, Depends(get_db)], db: Annotated[Session, Depends(get_db)],
) -> SupportMessageOut: ) -> SupportMessageOut:
"""Отправляет сообщение от лица *username* в support-топик (`sendMessage` — """Отправляет сообщение от лица *username* в support-топик (`sendMessage` —
не `copyMessage`: у веб-сообщения нет исходного Telegram-сообщения для копии).""" не `copyMessage`: у веб-сообщения нет исходного Telegram-сообщения для копии).
Порядок операций см. H1 в docstring модуля: rate-limit проверяется (но НЕ
расходуется, review L3) до отправки, thread создаётся ТОЛЬКО после успешного
`send_message` до этого момента с БД не происходит ничего.
"""
if not _bot_configured(): if not _bot_configured():
# Предсказуемое поведение вместо 500 (#tgsupport-web AC): бот не настроен # Предсказуемое поведение вместо 500 (#tgsupport-web AC): бот не настроен
# (пустой TELEGRAM_BOT_TOKEN, dev/staging) или support-топик не задан — # (пустой TELEGRAM_BOT_TOKEN, dev/staging) или support-топик не задан —
# мирроринг невозможен физически, ничего не пишем в БД. # мирроринг невозможен физически, ничего не пишем в БД.
raise HTTPException(status_code=503, detail=SERVICE_UNAVAILABLE_TEXT) 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: if retry_after is not None:
raise HTTPException( raise HTTPException(
status_code=429, status_code=429,
@ -131,14 +169,15 @@ async def send_support_message(
headers={"Retry-After": str(int(retry_after) + 1)}, headers={"Retry-After": str(int(retry_after) + 1)},
) )
thread_id = storage.get_or_create_thread(db, username)
client = TelegramClient(settings.telegram_bot_token) client = TelegramClient(settings.telegram_bot_token)
try: try:
mirrored = await client.send_message( mirrored = await client.send_message(
chat_id=settings.telegram_support_chat_id, chat_id=settings.telegram_support_chat_id,
text=_format_mirror_text(username, payload.text), text=_format_mirror_text(username, payload.text),
message_thread_id=settings.telegram_support_topic_id or None, 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: except TelegramApiError:
# НЕ логируем payload.text (переписка — ПДн) и НЕ логируем токен (его в # НЕ логируем payload.text (переписка — ПДн) и НЕ логируем токен (его в
@ -148,13 +187,28 @@ async def send_support_message(
) )
raise HTTPException(status_code=502, detail=SERVICE_UNAVAILABLE_TEXT) from None 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( row = storage.record_inbound(
db, db,
thread_id=thread_id, thread_id=thread_id,
text_body=payload.text, text_body=payload.text,
topic_message_id=topic_message_id, topic_message_id=topic_message_id,
support_chat_id=settings.telegram_support_chat_id,
) )
db.commit() db.commit()
@ -163,25 +217,34 @@ async def send_support_message(
@router.get("/support/messages", response_model=list[SupportMessageOut]) @router.get("/support/messages", response_model=list[SupportMessageOut])
async def list_support_messages( def list_support_messages(
username: Annotated[str, Depends(_require_username)], username: Annotated[str, Depends(_require_username)],
db: Annotated[Session, Depends(get_db)], db: Annotated[Session, Depends(get_db)],
since: Annotated[int, Query(ge=0)] = 0, since: Annotated[int, Query(ge=0)] = 0,
) -> list[SupportMessageOut]: ) -> list[SupportMessageOut]:
"""Сообщения СВОЕГО треда с id > since. Тред резолвится по username — чужой """Сообщения СВОЕГО треда с 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) thread_id = storage.find_thread_id(db, username)
if thread_id is None: if thread_id is None:
return [] 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] return [SupportMessageOut(**r) for r in rows]
@router.get("/support/unread", response_model=UnreadOut) @router.get("/support/unread", response_model=UnreadOut)
async def get_support_unread( def get_support_unread(
username: Annotated[str, Depends(_require_username)], username: Annotated[str, Depends(_require_username)],
db: Annotated[Session, Depends(get_db)], db: Annotated[Session, Depends(get_db)],
) -> UnreadOut: ) -> UnreadOut:
"""Sync `def` (review M3) — см. `list_support_messages`."""
thread_id = storage.find_thread_id(db, username) thread_id = storage.find_thread_id(db, username)
if thread_id is None: if thread_id is None:
return UnreadOut(unread=0) return UnreadOut(unread=0)
@ -189,10 +252,11 @@ async def get_support_unread(
@router.post("/support/read", response_model=StatusOut) @router.post("/support/read", response_model=StatusOut)
async def mark_support_read( def mark_support_read(
username: Annotated[str, Depends(_require_username)], username: Annotated[str, Depends(_require_username)],
db: Annotated[Session, Depends(get_db)], db: Annotated[Session, Depends(get_db)],
) -> StatusOut: ) -> StatusOut:
"""Sync `def` (review M3) — см. `list_support_messages`."""
thread_id = storage.find_thread_id(db, username) thread_id = storage.find_thread_id(db, username)
if thread_id is not None: if thread_id is not None:
storage.mark_read(db, thread_id=thread_id) storage.mark_read(db, thread_id=thread_id)

View file

@ -97,19 +97,43 @@ class SlidingWindowLimiter:
self._window_s = window_s self._window_s = window_s
self._hits: dict[str, deque[float]] = defaultdict(deque) self._hits: dict[str, deque[float]] = defaultdict(deque)
def check(self, key: str) -> float | None: def _prune(self, bucket: deque[float], now: float) -> None:
"""Регистрирует попытку под *key*. Возвращает None, если она уложилась в
лимит (и учтена), иначе сколько секунд ждать до следующей попытки."""
now = time.monotonic()
bucket = self._hits[key]
cutoff = now - self._window_s cutoff = now - self._window_s
while bucket and bucket[0] < cutoff: while bucket and bucket[0] < cutoff:
bucket.popleft() 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: if len(bucket) >= self._limit:
return self._window_s - (now - bucket[0]) 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) 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 return None

View file

@ -11,15 +11,23 @@
web_support_messages (direction='in') этот модуль в этой ветке не участвует, web_support_messages (direction='in') этот модуль в этой ветке не участвует,
только в разборе ответа (B ниже). только в разборе ответа (B ниже).
B) Оператор отвечает РЕПЛАЕМ в support-группе на зеркало клиента B) Оператор отвечает РЕПЛАЕМ в support-группе на зеркало клиента
находим chat_id по topic_message_id copyMessage ответа в личку клиента резолвим topic_message_id ОБЕ стороны (tg_support_messages И
запись (direction='out'). Если зеркало НЕ Telegram-клиента, а веб-чата web_support_messages), скоупя к ТЕКУЩЕМУ TELEGRAM_SUPPORT_CHAT_ID
(data/sql/187_web_support_chat.sql) доставка идёт НЕ в Telegram (у веб- (#tgsupport-web review M1 — если группу когда-нибудь сменят/пересоздадут,
клиента нет личного чата с ботом), а записью direction='out' в Telegram message_id стартует заново и может совпасть со старым числом из
web_support_messages (веб-фронт вычитывает её обычным polling'ом). Реплай другой таблицы; без скоупинга это была бы ТИХАЯ доставка постороннему
не на зеркало (или не реплай вообще) обычная болтовня в топике, тихий клиенту). Совпадение НА ОБЕИХ сторонах одновременно громкий отказ
игнор. Telegram 403 (клиент заблокировал бота) is_blocked=true + (`logger.error`, ничего не доставляем) вместо произвольного выбора одной из
уведомление в топике (только для Telegram-ветки у веб-клиента нет них. Иначе: 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 сохраняется И C) Дедуп: update_id <= сохранённого offset skip. Offset сохраняется И
коммитится в той же транзакции, что и запись сообщения (см. `process_update` коммитится в той же транзакции, что и запись сообщения (см. `process_update`
`finally`), после КАЖДОГО апдейта рестарт воркера не переигрывает уже `finally`), после КАЖДОГО апдейта рестарт воркера не переигрывает уже
@ -79,6 +87,11 @@ SERVICE_UNAVAILABLE_TEXT = (
# tg_support_messages.kind): "text | photo | document | video | voice | other". # tg_support_messages.kind): "text | photo | document | video | voice | other".
_KNOWN_KINDS = ("text", "photo", "document", "video", "voice") _KNOWN_KINDS = ("text", "photo", "document", "video", "voice")
# #tgsupport-web review M2: реплай оператора медиа-типом (в т.ч. фото С ПОДПИСЬЮ)
# на веб-зеркало НЕ доставляется частично — веб-чат текстовый MVP, оператор
# получает это уведомление в топике вместо тихого игнора (иначе уверен, что ответил).
_WEB_UNSUPPORTED_MEDIA_REPLY_TEXT = "Веб-чат поддерживает только текст, сообщение не доставлено."
# ── Storage abstraction (testable без реальной БД) ────────────────────────── # ── Storage abstraction (testable без реальной БД) ──────────────────────────
class BridgeStorage(Protocol): class BridgeStorage(Protocol):
@ -115,13 +128,18 @@ class BridgeStorage(Protocol):
kind: str, kind: str,
text_body: str | None, text_body: str | None,
operator_tg_id: int | None, operator_tg_id: int | None,
support_chat_id: int | None = None,
) -> int | 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 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( def record_web_out_message(
self, *, thread_id: int, text_body: str, operator_tg_id: int | None self, *, thread_id: int, text_body: str, operator_tg_id: int | None
@ -243,17 +261,19 @@ class SqlBridgeStorage:
kind: str, kind: str,
text_body: str | None, text_body: str | None,
operator_tg_id: int | None, operator_tg_id: int | None,
support_chat_id: int | None = None,
) -> int | None: ) -> int | None:
row = self._db.execute( row = self._db.execute(
text( text(
""" """
INSERT INTO tg_support_messages INSERT INTO tg_support_messages
(chat_id, direction, tg_message_id, topic_message_id, kind, (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 VALUES
(CAST(:chat_id AS bigint), CAST(:direction AS text), (CAST(:chat_id AS bigint), CAST(:direction AS text),
CAST(:tg_message_id AS bigint), CAST(:topic_message_id AS bigint), 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 RETURNING id
""" """
), ),
@ -265,11 +285,17 @@ class SqlBridgeStorage:
"kind": kind, "kind": kind,
"text_body": text_body, "text_body": text_body,
"operator_tg_id": operator_tg_id, "operator_tg_id": operator_tg_id,
"support_chat_id": support_chat_id,
}, },
).fetchone() ).fetchone()
return int(row[0]) if row is not None else None 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( row = self._db.execute(
text( text(
""" """
@ -277,11 +303,13 @@ class SqlBridgeStorage:
FROM tg_support_messages FROM tg_support_messages
WHERE topic_message_id = CAST(:topic_message_id AS bigint) WHERE topic_message_id = CAST(:topic_message_id AS bigint)
AND direction = 'in' AND direction = 'in'
AND (support_chat_id = CAST(:support_chat_id AS bigint)
OR support_chat_id IS NULL)
ORDER BY created_at DESC ORDER BY created_at DESC
LIMIT 1 LIMIT 1
""" """
), ),
{"topic_message_id": topic_message_id}, {"topic_message_id": topic_message_id, "support_chat_id": support_chat_id},
).fetchone() ).fetchone()
return int(row[0]) if row is not None else None return int(row[0]) if row is not None else None
@ -294,10 +322,14 @@ class SqlBridgeStorage:
{"chat_id": chat_id}, {"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) — то же соединение/ """Делегирует в `web_support_storage` (#tgsupport-web) — то же соединение/
транзакцию, что и tg-путь, коммитится вместе offset'ом в `process_update`.""" транзакцию, что и 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( def record_web_out_message(
self, *, thread_id: int, text_body: str, operator_tg_id: int | None 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), kind=_infer_kind(message),
text_body=text_body or message.get("caption"), text_body=text_body or message.get("caption"),
operator_tg_id=None, 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): if not isinstance(mirror_message_id, int):
return 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: if target_chat_id is not None:
# Существующий Telegram-путь — НЕ ТРОНУТ, только обёрнут в explicit if # Существующий Telegram-путь — НЕ ТРОНУТ.
# (раньше было `if target_chat_id is None: ...; return`, теперь после
# этой ветки идёт ещё веб-резолв, см. ниже).
message_id = message.get("message_id") message_id = message.get("message_id")
if not isinstance(message_id, int): if not isinstance(message_id, int):
return return
@ -462,19 +514,32 @@ async def _handle_group_reply(
) )
return 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: if web_thread_id is not None:
text_body = message.get("text") or message.get("caption") message_id = message.get("message_id")
if not text_body: kind = _infer_kind(message)
# Веб-чат — текстовый MVP (web_support_messages.text_body NOT NULL, if kind != "text":
# нет kind/file_id колонок как у tg_support_messages) — доставить # #tgsupport-web review M2: НЕ доставляем частично (фото С ПОДПИСЬЮ
# фото/документ/voice некуда, фронт это не отрендерит. # молча превратилось бы в "ответ = только текст подписи", клиент решил
# бы что это весь ответ) — отказ целиком + явное уведомление оператору
# в топике (тот же паттерн, что 403-уведомление выше), иначе оператор
# уверен, что ответ доставлен, хотя веб-чат не поддерживает медиа.
logger.warning( logger.warning(
"tgbot bridge: реплай на веб-зеркало (thread_id=%d) без текста " "tgbot bridge: реплай на веб-зеркало (thread_id=%d) содержит %s, "
"(медиа?) — веб-чат текстовый, доставка невозможна, игнор", "не текст — веб-чат поддерживает только текст, доставка отклонена",
web_thread_id, 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 return
operator = message.get("from") or {} operator = message.get("from") or {}

View file

@ -238,12 +238,28 @@ class TelegramClient:
text: str, text: str,
message_thread_id: int | None = None, message_thread_id: int | None = None,
reply_to_message_id: int | None = None, reply_to_message_id: int | None = None,
timeout: float | None = None,
max_retries: int | None = None,
) -> dict[str, Any]: ) -> 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} payload: dict[str, Any] = {"chat_id": chat_id, "text": text}
if message_thread_id: if message_thread_id:
payload["message_thread_id"] = message_thread_id payload["message_thread_id"] = message_thread_id
if reply_to_message_id: if reply_to_message_id:
payload["reply_to_message_id"] = 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 {} return result if isinstance(result, dict) else {}

View file

@ -62,19 +62,31 @@ def get_or_create_thread(db: Session, username: str) -> int:
def record_inbound( 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]: ) -> dict[str, Any]:
"""Записывает сообщение пользователя сайта (direction='in'). `topic_message_id` — """Записывает сообщение пользователя сайта (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 = ( row = (
db.execute( db.execute(
text( text(
""" """
INSERT INTO web_support_messages 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 VALUES
(CAST(:thread_id AS bigint), 'in', :text_body, (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 RETURNING id, direction, text_body, operator_tg_id, created_at
""" """
), ),
@ -82,6 +94,7 @@ def record_inbound(
"thread_id": thread_id, "thread_id": thread_id,
"text_body": text_body, "text_body": text_body,
"topic_message_id": topic_message_id, "topic_message_id": topic_message_id,
"support_chat_id": support_chat_id,
}, },
) )
.mappings() .mappings()
@ -90,9 +103,17 @@ def record_inbound(
return dict(row) 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 — только """Резолвит 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( row = db.execute(
text( text(
""" """
@ -100,11 +121,12 @@ def find_thread_by_topic_message(db: Session, topic_message_id: int) -> int | No
FROM web_support_messages FROM web_support_messages
WHERE topic_message_id = CAST(:topic_message_id AS bigint) WHERE topic_message_id = CAST(:topic_message_id AS bigint)
AND direction = 'in' AND direction = 'in'
AND (support_chat_id = CAST(:support_chat_id AS bigint) OR support_chat_id IS NULL)
ORDER BY created_at DESC ORDER BY created_at DESC
LIMIT 1 LIMIT 1
""" """
), ),
{"topic_message_id": topic_message_id}, {"topic_message_id": topic_message_id, "support_chat_id": support_chat_id},
).fetchone() ).fetchone()
return int(row[0]) if row is not None else None 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 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]]: def list_messages(
"""Сообщения треда с id > since_id, по возрастанию (обычный polling с фронта).""" 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 = ( rows = (
db.execute( db.execute(
text( text(
@ -145,15 +176,16 @@ def list_messages(db: Session, *, thread_id: int, since_id: int) -> list[dict[st
FROM web_support_messages FROM web_support_messages
WHERE thread_id = CAST(:thread_id AS bigint) WHERE thread_id = CAST(:thread_id AS bigint)
AND id > CAST(:since_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() .mappings()
.all() .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: def count_unread(db: Session, *, thread_id: int) -> int:

View file

@ -35,15 +35,30 @@
-- - topic_message_id-маршрутизация (ключевой механизм моста) СОХРАНЕНА -- - topic_message_id-маршрутизация (ключевой механизм моста) СОХРАНЕНА
-- 1-в-1 по конвенции 186: partial UNIQUE на topic_message_id, -- 1-в-1 по конвенции 186: partial UNIQUE на topic_message_id,
-- заполняется только для direction='in', NULL для direction='out'. -- заполняется только для 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 получает одну -- - bridge.py меняется МИНИМАЛЬНО: _handle_group_reply получает одну
-- дополнительную ветку (пробуем tg-резолв, потом web-резолв, потом -- дополнительную ветку (резолвит tg-путь И web-путь, потом
-- existing orphan-warning) — существующий Telegram-путь не трогается. -- 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 линия -- - web_support_threads — один тред на username (сайт = 1 логин = 1 линия
-- переписки с поддержкой, без под-тредов). -- переписки с поддержкой, без под-тредов).
@ -53,11 +68,16 @@
-- 152-ФЗ: -- 152-ФЗ:
-- Переписка (text_body) — ПДн (может содержать любые данные, которые юзер -- Переписка (text_body) — ПДн (может содержать любые данные, которые юзер
-- решит написать). ON DELETE CASCADE от web_support_threads делает erasure -- решит написать). 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. -- IDEMPOTENCY: CREATE TABLE/INDEX IF NOT EXISTS — безопасный re-run.
-- Зависимости: нет (новые standalone таблицы, никакие существующие -- Зависимости: нет (новые standalone таблицы, никакие существующие
-- tg_support_*/иные таблицы не трогаются). -- tg_support_*/иные таблицы не трогаются — см. 188 для ALTER на tg_support_messages).
BEGIN; BEGIN;
@ -80,14 +100,16 @@ CREATE TABLE IF NOT EXISTS web_support_messages (
direction text NOT NULL CHECK (direction IN ('in', 'out')), direction text NOT NULL CHECK (direction IN ('in', 'out')),
text_body text NOT NULL CHECK (char_length(btrim(text_body)) > 0), text_body text NOT NULL CHECK (char_length(btrim(text_body)) > 0),
topic_message_id bigint, topic_message_id bigint,
support_chat_id bigint,
operator_tg_id bigint, operator_tg_id bigint,
created_at timestamptz NOT NULL DEFAULT now() 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.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.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.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''.'; 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 CREATE UNIQUE INDEX IF NOT EXISTS web_support_messages_topic_message_id_uq

View file

@ -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;

View file

@ -176,4 +176,3 @@
169_osm_poi_ekb_local.sql 169_osm_poi_ekb_local.sql
170_scrape_schedules_seed_osm_poi_ekb_refresh.sql 170_scrape_schedules_seed_osm_poi_ekb_refresh.sql
172_trade_in_leads.sql 172_trade_in_leads.sql
187_web_support_chat.sql

View file

@ -6,9 +6,14 @@ Coverage (per task spec + review follow-up):
после истечения окна #6 review) после истечения окна #6 review)
- реплай оператора user (доставка ответа клиенту + запись direction='out') - реплай оператора user (доставка ответа клиенту + запись direction='out')
- реплай оператора веб-чат (#tgsupport-web): зеркало веб-сообщения резолвится - реплай оператора веб-чат (#tgsupport-web): зеркало веб-сообщения резолвится
в web-тред, ответ пишется direction='out' БЕЗ Telegram-доставки; медиа-реплай в web-тред (скоуп по (topic_message_id, support_chat_id) review M1), ответ
на веб-зеркало (нет text/caption) тихий игнор (веб-чат текстовый MVP); пишется direction='out' БЕЗ Telegram-доставки; медиа-реплай (в т.ч. фото С
tg-резолв имеет приоритет над веб-резолвом при (искусственной) коллизии ПОДПИСЬЮ) на веб-зеркало отказ ЦЕЛИКОМ + уведомление оператору в топике
(review M2, никакой частичной доставки одной подписи); зеркало от ЧУЖОГО/
устаревшего support_chat_id не матчится (ротация группы); NULL
support_chat_id (легаси) wildcard-матч; совпадение ОБЕИХ сторон
одновременно (tg И web) громкий отказ (logger.error), а не молчаливый
выбор tg-пути (152-ФЗ misroute risk)
- реплай не на зеркало (или не реплай вообще) тихий игнор, не мусорим в чат; - реплай не на зеркало (или не реплай вообще) тихий игнор, не мусорим в чат;
реплай на СООБЩЕНИЕ БОТА без записи в БД WARNING про осиротевшее зеркало реплай на СООБЩЕНИЕ БОТА без записи в БД WARNING про осиротевшее зеркало
(#4 review) (#4 review)
@ -69,9 +74,11 @@ class FakeBridgeStorage:
self._next_id = 1 self._next_id = 1
self.clock_s: float = 0.0 self.clock_s: float = 0.0
self.fail_next_record_message = False self.fail_next_record_message = False
# #tgsupport-web: web_support_messages-эквивалент (topic_message_id -> # #tgsupport-web: web_support_messages-эквивалент, topic_message_id ->
# thread_id) + журнал outbound-записей, записанных через реплай оператора. # (thread_id, support_chat_id) — второй элемент моделирует колонку
self.web_topic_to_thread: dict[int, int] = {} # 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]] = [] self.web_out_messages: list[dict[str, Any]] = []
def get_offset(self) -> int: def get_offset(self) -> int:
@ -120,6 +127,7 @@ class FakeBridgeStorage:
kind: str, kind: str,
text_body: str | None, text_body: str | None,
operator_tg_id: int | None, operator_tg_id: int | None,
support_chat_id: int | None = None,
) -> int: ) -> int:
if self.fail_next_record_message: if self.fail_next_record_message:
self.fail_next_record_message = False self.fail_next_record_message = False
@ -136,23 +144,39 @@ class FakeBridgeStorage:
"kind": kind, "kind": kind,
"text_body": text_body, "text_body": text_body,
"operator_tg_id": operator_tg_id, "operator_tg_id": operator_tg_id,
"support_chat_id": support_chat_id,
"recorded_at_s": self.clock_s, "recorded_at_s": self.clock_s,
} }
) )
return row_id 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): for m in reversed(self.messages):
if m["direction"] == "in" and m["topic_message_id"] == topic_message_id: if m["direction"] != "in" or m["topic_message_id"] != topic_message_id:
return m["chat_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 return None
def mark_blocked(self, chat_id: int) -> None: def mark_blocked(self, chat_id: int) -> None:
self.blocked.add(chat_id) self.blocked.add(chat_id)
# ── #tgsupport-web ──────────────────────────────────────────────────── # ── #tgsupport-web ────────────────────────────────────────────────────
def find_web_thread_by_topic_message(self, topic_message_id: int) -> int | None: def find_web_thread_by_topic_message(
return self.web_topic_to_thread.get(topic_message_id) 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( def record_web_out_message(
self, *, thread_id: int, text_body: str, operator_tg_id: int | None 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]]] = [] calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({}, calls) client = _make_client({}, calls)
storage = FakeBridgeStorage() 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 = {
"update_id": 60, "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 assert storage.get_offset() == 60
async def test_group_reply_to_web_mirror_without_text_is_ignored( async def test_group_reply_to_web_mirror_with_null_support_chat_id_matches_current_chat() -> None:
caplog: pytest.LogCaptureFixture, """Легаси-строка (до 187/188, support_chat_id=None) — лениентный wildcard,
) -> None: матчится под ЛЮБЫМ текущим support_chat_id (review M1)."""
"""Веб-чат — текстовый MVP: реплай медиа-типом (нет text/caption) на веб-зеркало
не может быть доставлен тихий (WARNING, не error) игнор, ничего не пишем."""
calls: list[tuple[str, dict[str, Any]]] = [] calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({}, calls) client = _make_client({}, calls)
storage = FakeBridgeStorage() 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 del message["text"] # медиа-реплай без текста/caption
message["voice"] = {"file_id": "x"}
update = {"update_id": 61, "message": message} update = {"update_id": 61, "message": message}
with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"): with caplog.at_level(logging.WARNING, logger="app.services.tgbot.bridge"):
await bridge.process_update(update, client, storage) await bridge.process_update(update, client, storage)
assert calls == []
assert storage.web_out_messages == [] assert storage.web_out_messages == []
assert "текстовый" in caplog.text assert "не текст" in caplog.text
assert storage.get_offset() == 61 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 найден среди async def test_group_reply_to_web_mirror_with_photo_and_caption_is_refused_not_partial() -> None:
tg_support_messages, веб-резолв даже не вызывается (существующий Telegram-путь """Фото С ПОДПИСЬЮ на веб-зеркало — НЕ доставляем только подпись молча
работает как раньше, не деградирует из-за новой ветки).""" (клиент решил бы, что подпись весь ответ): отказ целиком, как и без 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]]] = [] calls: list[tuple[str, dict[str, Any]]] = []
client = _make_client({"copyMessage": {"message_id": 999}}, calls) client = _make_client({"copyMessage": {"message_id": 999}}, calls)
storage = FakeBridgeStorage() storage = FakeBridgeStorage()
@ -619,24 +718,22 @@ async def test_group_reply_prefers_tg_thread_when_both_would_match() -> None:
kind="text", kind="text",
text_body="вопрос клиента", text_body="вопрос клиента",
operator_tg_id=None, operator_tg_id=None,
support_chat_id=SUPPORT_CHAT_ID,
) )
# Тот же topic_message_id "случайно" тоже был бы в web-мапе — не должен storage.web_topic_to_thread[400] = (99, SUPPORT_CHAT_ID)
# переопределять tg-резолв (defensive, в реальности Telegram message_id не
# повторяется в пределах чата).
storage.web_topic_to_thread[400] = 99
update = { update = {
"update_id": 62, "update_id": 62,
"message": _group_reply_message(reply_to_message_id=400), "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 calls == [] # ничего не доставлено НИ В ОДНУ сторону
assert methods == ["copyMessage"]
assert storage.web_out_messages == [] assert storage.web_out_messages == []
out_rec = storage.messages[-1] assert len(storage.messages) == 1 # только исходное 'in', никакого 'out'
assert out_rec["direction"] == "out" assert "ОДНОВРЕМЕННО" in caplog.text
assert out_rec["chat_id"] == 555 assert storage.get_offset() == 62
# ── C) дедуп ────────────────────────────────────────────────────────────────── # ── C) дедуп ──────────────────────────────────────────────────────────────────

View file

@ -19,7 +19,7 @@ from fastapi import FastAPI
from fastapi.testclient import TestClient from fastapi.testclient import TestClient
from app.core import config from app.core import config
from app.core.ratelimit import RateLimitMiddleware from app.core.ratelimit import RateLimitMiddleware, SlidingWindowLimiter
@pytest.fixture @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 client.get("/api/v1/ping", headers={"X-Forwarded-For": "8.8.8.8, 2.2.2.2"}).status_code
== 200 == 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

View file

@ -107,6 +107,34 @@ def test_list_messages_without_auth_header_401(client: TestClient) -> None:
assert r.status_code == 401 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 ──────────────────────────────────────────────────────────────── # ── 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( def test_send_message_happy_path_mirrors_with_website_marker(
client: TestClient, db: MagicMock, monkeypatch: pytest.MonkeyPatch, _fake_telegram_client: Any client: TestClient, db: MagicMock, monkeypatch: pytest.MonkeyPatch, _fake_telegram_client: Any
) -> None: ) -> 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 = {} recorded = {}
def fake_record_inbound(db, *, thread_id, text_body, 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) recorded.update(
thread_id=thread_id,
text_body=text_body,
topic_message_id=topic_message_id,
support_chat_id=support_chat_id,
)
return { return {
"id": 100, "id": 100,
"direction": "in", "direction": "in",
@ -216,9 +255,14 @@ def test_send_message_happy_path_mirrors_with_website_marker(
assert "У меня вопрос про trade-in" in mirror_call["text"] assert "У меня вопрос про trade-in" in mirror_call["text"]
assert mirror_call["chat_id"] == support_module.settings.telegram_support_chat_id assert mirror_call["chat_id"] == support_module.settings.telegram_support_chat_id
assert mirror_call["message_thread_id"] == 42 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["thread_id"] == 7
assert recorded["topic_message_id"] == 555 # из FakeTelegramClient.send_message result 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 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 client: TestClient, monkeypatch: pytest.MonkeyPatch, _fake_telegram_client: Any
) -> None: ) -> None:
_fake_telegram_client._response = TelegramApiError("sendMessage", 400, "chat not found") _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 = [] record_called = []
monkeypatch.setattr( monkeypatch.setattr(
support_module.storage, 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 r.status_code == 502
assert record_called == [] # неотправленное сообщение не персистится assert record_called == [] # неотправленное сообщение не персистится
# review H1: thread создаётся ПОСЛЕ успешной отправки — на неудаче до БД не доходит вообще.
assert thread_created == []
# ── rate limit ──────────────────────────────────────────────────────────────── # ── rate limit ────────────────────────────────────────────────────────────────
@ -268,6 +319,40 @@ def test_send_message_rate_limited_429(client: TestClient, monkeypatch: pytest.M
assert "Retry-After" in second.headers 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( def test_send_message_rate_limit_is_per_user(
client: TestClient, monkeypatch: pytest.MonkeyPatch client: TestClient, monkeypatch: pytest.MonkeyPatch
) -> None: ) -> None:
@ -332,7 +417,7 @@ def test_list_messages_returns_thread_scoped_rows(
monkeypatch.setattr( monkeypatch.setattr(
support_module.storage, support_module.storage,
"list_messages", "list_messages",
lambda db, *, thread_id, since_id: [ lambda db, *, thread_id, since_id, limit: [
{ {
"id": 1, "id": 1,
"direction": "in", "direction": "in",
@ -351,6 +436,24 @@ def test_list_messages_returns_thread_scoped_rows(
assert body[0]["text_body"] == "hi" 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( def test_two_users_get_independent_threads(
client: TestClient, monkeypatch: pytest.MonkeyPatch client: TestClient, monkeypatch: pytest.MonkeyPatch
) -> None: ) -> None:
@ -361,7 +464,7 @@ def test_two_users_get_independent_threads(
monkeypatch.setattr( monkeypatch.setattr(
support_module.storage, support_module.storage,
"list_messages", "list_messages",
lambda db, *, thread_id, since_id: [ lambda db, *, thread_id, since_id, limit: [
{ {
"id": 1, "id": 1,
"direction": "in", "direction": "in",