Files
vidconf/backend/api/chat.py
Max Ronzhin 8e5eda88a2 feat(room): поднятие руки и очередь для организатора
Транспорт — существующий аутентифицированный WS чата (api/chat.py), а не
отдельный эндпоинт: сервер уже держит это соединение на каждого участника
(обоснование — докстринг chat_websocket и useChat.ts). Состояние очереди —
Redis (services/hand_queue.py), не Postgres: это эфемерное состояние звонка,
а не история, и два процесса uvicorn делают наивную память одного процесса
недостаточной. HSETNX даёт идемпотентное «поднять» (повторный клик не
переставляет в конец очереди), снапшот шлётся всем участникам при любом
изменении — организатор, зашедший позже, сразу видит актуальную картину.

Опустить чужую руку может организатор (решение оператора) — проверка через
conference.owner_id, не через identity клиента. Участник, вышедший из
комнаты LiveKit (webhook participant_left), теряет место в очереди
автоматически; переподключение WS чата место не сбрасывает (Redis не привязан
к жизни соединения). room_finished чистит очередь целиком — она не должна
пережить завершение звонка.

Побочный эффект транспортного решения: поднять руку нельзя, если чат выключен
настройкой инстанса (WS вообще не открывается) — принятый компромисс ради
переиспользования уже готового канала.

UI: кнопка «Рука» в тулбаре (у всех, бейдж — общий счётчик), бейдж на плитке
говорящего (видно всем), панель «Очередь» организатору (HandQueuePanel).
Кнопка «Рука» и панель «Очередь» намеренно НЕ прячутся в мобильную шторку
настроек, в отличие от «Вида», — поднятие руки посреди разговора требует
кнопки под рукой, а не в два клика вглубь настроек.

Этим же коммитом (файлы разделяемые с задачей B2, RoomParticipantTile.tsx/
useChat.ts/RoomStage.tsx/RoomPage.tsx/room.css) — проброс conferenceId и
каркас forced_mute-обработки, без которых кнопки принудительного мьюта не
скомпилировались бы; сама реализация мьюта — следующим коммитом.
2026-08-01 22:06:43 +03:00

209 lines
11 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.
"""WS-роутер комнаты конференции: `WS /api/v1/conferences/{id}/chat`.
Протокол: `connect` -> `accept()` -> клиент шлёт `{"type":"auth","token":...}`
первым сообщением (таймаут 10 с; токен не query-параметр — не палим его в
логах nginx) -> сервер проверяет тоггл `chat.enabled` и LiveKit-токен ->
история последних 50 сообщений чата + текущая очередь поднятых рук ->
двунаправленный обмен: `{"type":"message","text":...}` (чат, Redis pub/sub,
echo отправителю тоже), `{"type":"raise_hand"}`/`{"type":"lower_hand"}`
(очередь рук, задача B1 — состояние в Redis, см. `services/hand_queue.py`,
НЕ в БД: это эфемерное состояние звонка, а не история). Название файла и
эндпоинта («чат») оставлено как есть — эндпоинт исторически первый и
единственный аутентифицированный WS комнаты, поэтому очередь рук едет по
нему же, а не заводит отдельное соединение (дешевле: сервер уже держит
это соединение на каждого участника).
"""
import asyncio
import logging
import uuid
from typing import Annotated
from fastapi import APIRouter, Depends, WebSocket, WebSocketDisconnect
from pydantic import Field, TypeAdapter, ValidationError
from redis.asyncio.client import PubSub
from sqlalchemy.ext.asyncio import AsyncSession
from core.db import get_session
from core.redis import redis_client
from models.conference import Conference
from schemas.chat import (
ChatAuthIn,
ChatErrorOut,
ChatHistoryOut,
ChatMessageEventOut,
ChatMessageIn,
ChatMessageOut,
)
from schemas.room_events import LowerHandIn, RaiseHandIn
from services import hand_queue
from services.chat import ChatAuthError, ChatIdentity, ChatService, InvalidTokenError, chat_channel
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/v1/conferences", tags=["chat"])
# Таймаут ожидания первого (auth) сообщения клиента.
AUTH_TIMEOUT_SECONDS = 10.0
# Дискриминированное объединение сообщений клиента ПОСЛЕ auth — по полю `type`.
_ClientEnvelope = Annotated[
ChatMessageIn | RaiseHandIn | LowerHandIn, Field(discriminator="type")
]
_client_envelope_adapter: TypeAdapter[ChatMessageIn | RaiseHandIn | LowerHandIn] = TypeAdapter(
_ClientEnvelope
)
@router.websocket("/{conference_id}/chat")
async def chat_websocket(
websocket: WebSocket,
conference_id: uuid.UUID,
session: Annotated[AsyncSession, Depends(get_session)],
) -> None:
"""WS-эндпоинт комнаты конференции — единая аутентификация LiveKit-токеном."""
await websocket.accept()
service = ChatService(session)
try:
identity = await _authenticate(websocket, service)
conference = await service.ensure_chat_open(conference_id, identity=identity)
except ChatAuthError as exc:
await _close_quietly(websocket, exc.close_code)
return
pubsub = redis_client.pubsub()
channel = chat_channel(conference.id)
room_channel = hand_queue.hand_queue_channel(conference.id)
# Подписка ДО чтения истории/снапшота очереди: событие,
# опубликованное другим клиентом в окне между SELECT/HGETALL и subscribe,
# иначе теряется для подключающегося клиента — Redis начинает
# буферизовать входящие publish для этого соединения сразу после
# subscribe, до первого вызова `get_message`. На стыке возможен дубликат
# (то же сообщение чата и в history, и в первом pub/sub-сообщении) —
# безопаснее дедуплицировать по `id`, чем потерять сообщение; снапшот
# очереди дублировать безвредно (полная замена состояния на клиенте).
await pubsub.subscribe(channel, room_channel)
try:
history = await service.history(conference)
await websocket.send_json(ChatHistoryOut(messages=history).model_dump(mode="json"))
seen_ids = {item.id for item in history}
queue_out = await hand_queue.get_snapshot_out(conference.id)
await websocket.send_json(queue_out.model_dump(mode="json"))
async with asyncio.TaskGroup() as tg:
tg.create_task(_pump_pubsub_to_websocket(websocket, pubsub, channel, seen_ids))
tg.create_task(_pump_websocket_to_service(websocket, service, conference, identity))
except* WebSocketDisconnect:
# Штатное закрытие соединения клиентом — не ошибка.
pass
except* ChatAuthError as eg:
# Допуск был проверен только при коннекте — за время жизни
# долгоживущего WS (LiveKit-токен TTL 6 часов) конференция могла
# завершиться; `persist_and_publish` бросает `ChatUnavailableError`
# при попытке создать сессию пайплайна для уже мёртвой
# конференции — закрываем с тем же кодом, что и при отказе
# на коннекте.
# `except*` всегда связывает `ExceptionGroup` (PEP 654) — на рантайме
# `eg.exceptions[0]` гарантированно `ChatAuthError`; mypy после
# нескольких подряд идущих `except*` моделирует тип `eg` неточно
# (union с "голым" `ChatAuthError`, не имеющим `.exceptions`).
await _close_quietly(websocket, eg.exceptions[0].close_code) # type: ignore[union-attr]
finally:
# Всегда отписываемся и закрываем pubsub-соединение, иначе при частых
# обрывах соединений копятся забытые подписки на стороне Redis.
await pubsub.unsubscribe(channel, room_channel)
# `PubSub.aclose` в redis-py не аннотирован (untyped def) несмотря на
# `py.typed` пакета — узкий игнор именно этого вызова.
await pubsub.aclose() # type: ignore[no-untyped-call]
async def _authenticate(websocket: WebSocket, service: ChatService) -> ChatIdentity:
"""Дождаться первого (auth) сообщения клиента с таймаутом и проверить LiveKit-токен."""
try:
raw = await asyncio.wait_for(websocket.receive_text(), timeout=AUTH_TIMEOUT_SECONDS)
except (TimeoutError, WebSocketDisconnect) as exc:
raise InvalidTokenError from exc
try:
envelope = ChatAuthIn.model_validate_json(raw)
except ValidationError as exc:
raise InvalidTokenError from exc
return await service.authenticate(envelope.token)
async def _pump_pubsub_to_websocket(
websocket: WebSocket, pubsub: PubSub, chat_channel_name: str, seen_ids: set[int]
) -> None:
"""Читать оба Redis pub/sub канала комнаты (чат + очередь рук) и пересылать клиенту.
`seen_ids` — id сообщений чата, уже отправленных клиенту в `history` (на
стыке подписки и SELECT истории возможен дубликат, см. докстринг
`chat_websocket`) — такие сообщения не пересылаются повторно. Снапшоты
очереди рук такой дедупликации не требуют (полная замена состояния).
"""
while True:
raw = await pubsub.get_message(ignore_subscribe_messages=True, timeout=None)
if raw is None:
continue
if raw["channel"] == chat_channel_name:
message = ChatMessageOut.model_validate_json(raw["data"])
if message.id in seen_ids:
continue
seen_ids.add(message.id)
await websocket.send_json(
ChatMessageEventOut(message=message).model_dump(mode="json")
)
else:
# Канал комнаты (`hand_queue.hand_queue_channel`) — уже готовый
# JSON исходящего конверта (`HandQueueOut`/`ForcedMuteOut`,
# см. `services/hand_queue.py::publish_snapshot` и эндпоинт мьюта
# в `api/conferences.py`), пересылаем как есть без пересборки.
await websocket.send_text(raw["data"])
async def _pump_websocket_to_service(
websocket: WebSocket, service: ChatService, conference: Conference, identity: ChatIdentity
) -> None:
"""Читать сообщения клиента (текст чата / поднять-опустить руку), валидировать и обработать."""
is_organizer = conference.owner_id is not None and conference.owner_id == identity.user_id
while True:
raw = await websocket.receive_text()
try:
envelope = _client_envelope_adapter.validate_json(raw)
except ValidationError:
await websocket.send_json(ChatErrorOut(code="invalid_message").model_dump(mode="json"))
continue
if isinstance(envelope, ChatMessageIn):
await service.persist_and_publish(conference, identity=identity, text=envelope.text)
elif isinstance(envelope, RaiseHandIn):
await hand_queue.raise_hand(
conference.id, identity=_identity_key(identity), name=identity.author_name
)
await hand_queue.publish_snapshot(conference.id)
else:
target = envelope.identity or _identity_key(identity)
if target != _identity_key(identity) and not is_organizer:
await websocket.send_json(
ChatErrorOut(code="forbidden").model_dump(mode="json")
)
continue
await hand_queue.lower_hand(conference.id, identity=target)
await hand_queue.publish_snapshot(conference.id)
def _identity_key(identity: ChatIdentity) -> str:
"""Identity участника в формате LiveKit/очереди рук — `str(user_id)` либо `guest:{id}`."""
if identity.user_id is not None:
return str(identity.user_id)
return f"guest:{identity.guest_access_id}"
async def _close_quietly(websocket: WebSocket, code: int) -> None:
"""Закрыть WS с заданным кодом, не роняя обработчик, если клиент уже отвалился."""
try:
await websocket.close(code=code)
except Exception: # noqa: BLE001 — соединение уже могло быть разорвано клиентом
logger.debug("chat websocket: close(%s) на уже разорванном соединении", code)