153 lines
7.9 KiB
Python
153 lines
7.9 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 отправителю тоже).
|
||
"""
|
||
|
||
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)
|