"""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 отправителю тоже). """ import asyncio import logging import uuid from typing import Annotated from fastapi import APIRouter, Depends, WebSocket, WebSocketDisconnect from pydantic import 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 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 @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) # Подписка ДО чтения истории: сообщение, # опубликованное другим клиентом в окне между SELECT истории и # subscribe, иначе теряется для подключающегося клиента — Redis начинает # буферизовать входящие publish для этого соединения сразу после # subscribe, до первого вызова `get_message`. На стыке возможен дубликат # (то же сообщение и в history, и в первом pub/sub-сообщении) — безопаснее # дедуплицировать по `id`, чем потерять сообщение. await pubsub.subscribe(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} async with asyncio.TaskGroup() as tg: tg.create_task(_pump_pubsub_to_websocket(websocket, pubsub, 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) # `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, 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 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")) async def _pump_websocket_to_service( websocket: WebSocket, service: ChatService, conference: Conference, identity: ChatIdentity ) -> None: """Читать текстовые сообщения клиента, валидировать и сохранять+публиковать их.""" while True: raw = await websocket.receive_text() try: envelope = ChatMessageIn.model_validate_json(raw) except ValidationError: await websocket.send_json(ChatErrorOut(code="invalid_message").model_dump(mode="json")) continue await service.persist_and_publish(conference, identity=identity, text=envelope.text) 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)