Транспорт — существующий аутентифицированный 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-обработки, без которых кнопки принудительного мьюта не скомпилировались бы; сама реализация мьюта — следующим коммитом.
209 lines
11 KiB
Python
209 lines
11 KiB
Python
"""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)
|