"""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 hand_queue_enabled = await service.hand_queue_enabled() 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, hand_queue_enabled ) ) 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, hand_queue_enabled: bool, ) -> None: """Читать сообщения клиента (текст чата / поднять-опустить руку), валидировать и обработать. `hand_queue_enabled` — снятый один раз при подключении тоггл модуля «поднятие руки» (см. `ChatService.hand_queue_enabled`): при `False` `raise_hand`/`lower_hand` отклоняются кодом `hand_queue_disabled` — вторая линия защиты сверх того, что фронт при выключенном модуле вообще не рисует кнопки (см. `RoomToolbar`/`HandQueueMenu`). """ 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): if not hand_queue_enabled: await websocket.send_json( ChatErrorOut(code="hand_queue_disabled").model_dump(mode="json") ) continue await hand_queue.raise_hand( conference.id, identity=_identity_key(identity), name=identity.author_name ) await hand_queue.publish_snapshot(conference.id) else: if not hand_queue_enabled: await websocket.send_json( ChatErrorOut(code="hand_queue_disabled").model_dump(mode="json") ) continue 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)