Раньше HandQueueMenu.tsx рендерился только организатору — теперь очередь видит любой участник, но опустить чужую руку по-прежнему может только организатор (сервер это уже проверял, менял только фронт). Кнопка «Опустить» показывается у записи, только если это своя рука или пользователь — организатор. Модуль «поднятие руки» (кнопка «Рука» + очередь целиком) — отключаемый в админке (instance_settings.hand_queue, дефолт enabled=true, как у chat_enabled). Настройка едет участнику в JoinOut ещё до входа в комнату; выключенный модуль гасит кнопки и на фронте, и на бэке — raise_hand/lower_hand отклоняются кодом hand_queue_disabled, если модуль выключен, даже если у клиента на руках старый JoinOut.
228 lines
11 KiB
Python
228 lines
11 KiB
Python
"""Бизнес-логика 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 hand_queue_enabled(self) -> bool:
|
||
"""Тоггл инстанса `hand_queue.enabled` — снимается один раз при подключении WS
|
||
(см. `api/chat.py::chat_websocket`), а не на каждое сообщение: та же
|
||
осознанная «застылость» на время жизни соединения, что и у
|
||
`is_organizer` в `_pump_websocket_to_service` — переключение модуля
|
||
администратором применяется со следующего подключения."""
|
||
cfg = await InstanceSettingsService(self._session).get()
|
||
return cfg.hand_queue.enabled
|
||
|
||
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,
|
||
)
|