feat(room): поднятие руки и очередь для организатора
Транспорт — существующий аутентифицированный WS чата (api/chat.py), а не отдельный эндпоинт: сервер уже держит это соединение на каждого участника (обоснование — докстринг chat_websocket и useChat.ts). Состояние очереди — Redis (services/hand_queue.py), не Postgres: это эфемерное состояние звонка, а не история, и два процесса uvicorn делают наивную память одного процесса недостаточной. HSETNX даёт идемпотентное «поднять» (повторный клик не переставляет в конец очереди), снапшот шлётся всем участникам при любом изменении — организатор, зашедший позже, сразу видит актуальную картину. Опустить чужую руку может организатор (решение оператора) — проверка через conference.owner_id, не через identity клиента. Участник, вышедший из комнаты LiveKit (webhook participant_left), теряет место в очереди автоматически; переподключение WS чата место не сбрасывает (Redis не привязан к жизни соединения). room_finished чистит очередь целиком — она не должна пережить завершение звонка. Побочный эффект транспортного решения: поднять руку нельзя, если чат выключен настройкой инстанса (WS вообще не открывается) — принятый компромисс ради переиспользования уже готового канала. UI: кнопка «Рука» в тулбаре (у всех, бейдж — общий счётчик), бейдж на плитке говорящего (видно всем), панель «Очередь» организатору (HandQueuePanel). Кнопка «Рука» и панель «Очередь» намеренно НЕ прячутся в мобильную шторку настроек, в отличие от «Вида», — поднятие руки посреди разговора требует кнопки под рукой, а не в два клика вглубь настроек. Этим же коммитом (файлы разделяемые с задачей B2, RoomParticipantTile.tsx/ useChat.ts/RoomStage.tsx/RoomPage.tsx/room.css) — проброс conferenceId и каркас forced_mute-обработки, без которых кнопки принудительного мьюта не скомпилировались бы; сама реализация мьюта — следующим коммитом.
This commit is contained in:
@@ -1,10 +1,17 @@
|
||||
"""WS-роутер текстового чата конференции: `WS /api/v1/conferences/{id}/chat`.
|
||||
"""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 отправителю тоже).
|
||||
история последних 50 сообщений чата + текущая очередь поднятых рук ->
|
||||
двунаправленный обмен: `{"type":"message","text":...}` (чат, Redis pub/sub,
|
||||
echo отправителю тоже), `{"type":"raise_hand"}`/`{"type":"lower_hand"}`
|
||||
(очередь рук, задача B1 — состояние в Redis, см. `services/hand_queue.py`,
|
||||
НЕ в БД: это эфемерное состояние звонка, а не история). Название файла и
|
||||
эндпоинта («чат») оставлено как есть — эндпоинт исторически первый и
|
||||
единственный аутентифицированный WS комнаты, поэтому очередь рук едет по
|
||||
нему же, а не заводит отдельное соединение (дешевле: сервер уже держит
|
||||
это соединение на каждого участника).
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
@@ -13,7 +20,7 @@ import uuid
|
||||
from typing import Annotated
|
||||
|
||||
from fastapi import APIRouter, Depends, WebSocket, WebSocketDisconnect
|
||||
from pydantic import ValidationError
|
||||
from pydantic import Field, TypeAdapter, ValidationError
|
||||
from redis.asyncio.client import PubSub
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
@@ -28,6 +35,8 @@ from schemas.chat import (
|
||||
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__)
|
||||
@@ -37,6 +46,14 @@ 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(
|
||||
@@ -44,7 +61,7 @@ async def chat_websocket(
|
||||
conference_id: uuid.UUID,
|
||||
session: Annotated[AsyncSession, Depends(get_session)],
|
||||
) -> None:
|
||||
"""WS-эндпоинт текстового чата конференции — единая аутентификация LiveKit-токеном."""
|
||||
"""WS-эндпоинт комнаты конференции — единая аутентификация LiveKit-токеном."""
|
||||
await websocket.accept()
|
||||
service = ChatService(session)
|
||||
|
||||
@@ -57,21 +74,26 @@ async def chat_websocket(
|
||||
|
||||
pubsub = redis_client.pubsub()
|
||||
channel = chat_channel(conference.id)
|
||||
# Подписка ДО чтения истории: сообщение,
|
||||
# опубликованное другим клиентом в окне между SELECT истории и
|
||||
# subscribe, иначе теряется для подключающегося клиента — Redis начинает
|
||||
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)
|
||||
# (то же сообщение чата и в 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, seen_ids))
|
||||
tg.create_task(_pump_pubsub_to_websocket(websocket, pubsub, channel, seen_ids))
|
||||
tg.create_task(_pump_websocket_to_service(websocket, service, conference, identity))
|
||||
except* WebSocketDisconnect:
|
||||
# Штатное закрытие соединения клиентом — не ошибка.
|
||||
@@ -91,7 +113,7 @@ async def chat_websocket(
|
||||
finally:
|
||||
# Всегда отписываемся и закрываем pubsub-соединение, иначе при частых
|
||||
# обрывах соединений копятся забытые подписки на стороне Redis.
|
||||
await pubsub.unsubscribe(channel)
|
||||
await pubsub.unsubscribe(channel, room_channel)
|
||||
# `PubSub.aclose` в redis-py не аннотирован (untyped def) несмотря на
|
||||
# `py.typed` пакета — узкий игнор именно этого вызова.
|
||||
await pubsub.aclose() # type: ignore[no-untyped-call]
|
||||
@@ -111,37 +133,71 @@ async def _authenticate(websocket: WebSocket, service: ChatService) -> ChatIdent
|
||||
|
||||
|
||||
async def _pump_pubsub_to_websocket(
|
||||
websocket: WebSocket, pubsub: PubSub, seen_ids: set[int]
|
||||
websocket: WebSocket, pubsub: PubSub, chat_channel_name: str, seen_ids: set[int]
|
||||
) -> None:
|
||||
"""Читать сообщения Redis pub/sub канала чата и пересылать их подключённому клиенту.
|
||||
"""Читать оба Redis pub/sub канала комнаты (чат + очередь рук) и пересылать клиенту.
|
||||
|
||||
`seen_ids` — id сообщений, уже отправленных клиенту в `history` (на
|
||||
`seen_ids` — id сообщений чата, уже отправленных клиенту в `history` (на
|
||||
стыке подписки и SELECT истории возможен дубликат, см. докстринг
|
||||
`chat_websocket`) — такие сообщения не пересылаются повторно.
|
||||
`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"))
|
||||
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
|
||||
) -> None:
|
||||
"""Читать текстовые сообщения клиента, валидировать и сохранять+публиковать их."""
|
||||
"""Читать сообщения клиента (текст чата / поднять-опустить руку), валидировать и обработать."""
|
||||
is_organizer = conference.owner_id is not None and conference.owner_id == identity.user_id
|
||||
while True:
|
||||
raw = await websocket.receive_text()
|
||||
try:
|
||||
envelope = ChatMessageIn.model_validate_json(raw)
|
||||
envelope = _client_envelope_adapter.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)
|
||||
|
||||
if isinstance(envelope, ChatMessageIn):
|
||||
await service.persist_and_publish(conference, identity=identity, text=envelope.text)
|
||||
elif isinstance(envelope, RaiseHandIn):
|
||||
await hand_queue.raise_hand(
|
||||
conference.id, identity=_identity_key(identity), name=identity.author_name
|
||||
)
|
||||
await hand_queue.publish_snapshot(conference.id)
|
||||
else:
|
||||
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:
|
||||
|
||||
Reference in New Issue
Block a user