"""Бизнес-логика WS-чата конференции: аутентификация LiveKit-токеном, история, publish. Единая аутентификация для пользователей и гостей — LiveKit access-токен (`livekit.api.TokenVerifier`), а не backend-JWT: у гостя backend-JWT нет вовсе (ADR-001, п.6), только LiveKit-токен, выданный при входе (`services/livekit_tokens.py`). Grant `video.room` доказывает допуск именно в эту конференцию — совпадение с `conference.slug` (имя LiveKit-комнаты, ADR-001, п.4). """ import logging import uuid from dataclasses import dataclass from datetime import UTC, datetime from livekit import api from sqlalchemy.ext.asyncio import AsyncSession from core.config import get_settings from core.redis import redis_client from models.chat import ChatMessage from models.conference import Conference from repositories.chat import ChatMessageRepository from repositories.conferences import ConferenceRepository, ConferenceSessionRepository from schemas.chat import ChatMessageOut from services.instance_settings import InstanceSettingsService logger = logging.getLogger(__name__) # Последние N сообщений открытой сессии, отправляемых новому подключению. CHAT_HISTORY_LIMIT = 50 # Префикс identity гостя в LiveKit-токене (см. `services/webhook_handlers.py`). GUEST_IDENTITY_PREFIX = "guest:" def chat_channel(conference_id: uuid.UUID) -> str: """Имя Redis pub/sub канала чата конкретной конференции.""" return f"chat:{conference_id}" class ChatAuthError(Exception): """Базовая ошибка допуска WS-подключения к чату — несёт WS close-код.""" def __init__(self, close_code: int) -> None: self.close_code = close_code super().__init__(close_code) class InvalidTokenError(ChatAuthError): """Нет auth-сообщения, таймаут, невалидный/нераспознанный LiveKit-токен (close 4401).""" def __init__(self) -> None: super().__init__(4401) class WrongRoomError(ChatAuthError): """Токен валиден, но выдан не для этой конференции (close 4403).""" def __init__(self) -> None: super().__init__(4403) class ChatUnavailableError(ChatAuthError): """Чат выключен настройкой инстанса ИЛИ конференция не найдена/завершена. Единый close-код 4404 для обоих случаев (номер/ ссылку конференции не перебираем, наличие конкретной конференции не палим отдельным кодом). """ def __init__(self) -> None: super().__init__(4404) @dataclass(frozen=True) class ChatIdentity: """Идентичность автора сообщения, восстановленная из LiveKit access-токена.""" user_id: uuid.UUID | None guest_access_id: uuid.UUID | None author_name: str room: str class ChatService: """Инкапсулирует сценарии WS-чата: auth, допуск, история, persist+publish.""" def __init__(self, session: AsyncSession) -> None: self._session = session self._conferences = ConferenceRepository(session) self._sessions = ConferenceSessionRepository(session) self._messages = ChatMessageRepository(session) async def authenticate(self, token: str) -> ChatIdentity: """Проверить LiveKit access-токен и восстановить identity автора сообщений. Любая ошибка формата/подписи токена — единообразный `InvalidTokenError` (детали причины не раскрываются клиенту). """ settings = get_settings() verifier = api.TokenVerifier(settings.livekit_api_key, settings.livekit_api_secret) try: claims = verifier.verify(token) except Exception as exc: # noqa: BLE001 — любая ошибка JWT/формата токена = 4401 raise InvalidTokenError from exc if claims.video is None or not claims.video.room or not claims.identity: raise InvalidTokenError identity = claims.identity room = claims.video.room if identity.startswith(GUEST_IDENTITY_PREFIX): try: guest_id = uuid.UUID(identity.removeprefix(GUEST_IDENTITY_PREFIX)) except ValueError as exc: raise InvalidTokenError from exc return ChatIdentity( user_id=None, guest_access_id=guest_id, author_name=claims.name, room=room ) try: user_id = uuid.UUID(identity) except ValueError as exc: raise InvalidTokenError from exc return ChatIdentity( user_id=user_id, guest_access_id=None, author_name=claims.name, room=room ) async def ensure_chat_open( self, conference_id: uuid.UUID, *, identity: ChatIdentity ) -> Conference: """Проверить допуск identity к чату конкретной конференции; вернуть конференцию. Порядок проверок важен: тоггл и существование/статус конференции — единый `ChatUnavailableError` (4404, инвариант №2), несовпадение комнаты токена — отдельный `WrongRoomError` (4403), но только после того, как убедились, что сама конференция легитимна. """ cfg = await InstanceSettingsService(self._session).get() if not cfg.chat.enabled: raise ChatUnavailableError conference = await self._conferences.get_by_id(conference_id) if conference is None or conference.status == "ended": raise ChatUnavailableError if identity.room != conference.slug: raise WrongRoomError return conference async def history(self, conference: Conference) -> list[ChatMessageOut]: """Последние сообщения открытой сессии конференции (пусто, если сессии ещё нет).""" session_record = await self._sessions.get_open_by_conference(conference.id) if session_record is None: return [] rows = await self._messages.last_for_session(session_record.id, limit=CHAT_HISTORY_LIMIT) return [_to_out(row) for row in rows] async def persist_and_publish( self, conference: Conference, *, identity: ChatIdentity, text: str ) -> None: """Сохранить сообщение в открытой сессии конференции и опубликовать его в Redis. Допуск проверяется только при коннекте (`ensure_chat_open`), а WS (с LiveKit-токеном TTL 6 часов, `services/livekit_tokens.py`) может жить намного дольше одной конференции — клиент способен слать сообщения уже ПОСЛЕ `room_finished`. Если открытой сессии нет, брать для решения "можно ли создать новую" статус из уже загруженного объекта `conference` нельзя (`expire_on_commit=False`, объект мог устареть за время жизни WS-сессии) — статус перечитывается свежим SELECT (`ConferenceRepository.get_status_by_id`, минует identity map). Для `ended`-конференции — `ChatUnavailableError` (close 4404), сообщение отклоняется, новая "фантомная" сессия НЕ создаётся (идемпотентность пайплайна). Если открытая сессия уже есть (обычный случай) — пишем в неё без пересчёта статуса. Порядок обязателен: сначала INSERT+commit в БД, потом publish — отправитель получает своё сообщение обратно через pub/sub-echo, порядок доставки единый у всех подписчиков канала. """ session_record = await self._sessions.get_open_by_conference(conference.id) if session_record is None: fresh_status = await self._conferences.get_status_by_id(conference.id) if fresh_status is None or fresh_status == "ended": raise ChatUnavailableError session_record = await self._sessions.create( conference_id=conference.id, title=conference.title, t_start=datetime.now(UTC) ) message = await self._messages.add( session_id=session_record.id, user_id=identity.user_id, guest_access_id=identity.guest_access_id, author_name=identity.author_name, text=text, ) await self._session.commit() payload = _to_out(message) await redis_client.publish(chat_channel(conference.id), payload.model_dump_json()) def _to_out(message: ChatMessage) -> ChatMessageOut: """Собрать `ChatMessageOut` из ORM-строки сообщения.""" if message.user_id is not None: author_id: str | None = str(message.user_id) elif message.guest_access_id is not None: author_id = str(message.guest_access_id) else: author_id = None return ChatMessageOut( id=message.id, author_id=author_id, author_name=message.author_name, is_guest=message.guest_access_id is not None, text=message.text, created_at=message.created_at, )