Files
vidconf/backend/services/hand_queue.py
Max Ronzhin 8e5eda88a2 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-обработки, без которых кнопки принудительного мьюта не
скомпилировались бы; сама реализация мьюта — следующим коммитом.
2026-08-01 22:06:43 +03:00

145 lines
7.9 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Очередь поднятых рук конференции — состояние в Redis, не в Postgres (задача B1).
Транспорт для клиентов — тот же аутентифицированный WS чата (`api/chat.py`):
переиспользуем уже открытые и держащиеся сервером соединения вместо отдельного
эндпоинта. Хранение — Redis, а не БД: очередь существует ровно во время звонка
и не должна переживать его завершение (в отличие от истории чата), а два
процесса uvicorn (`UVICORN_WORKERS`) делают наивное состояние в памяти одного
процесса недостаточным — организатор и участник могут оказаться на разных
воркерах.
Один Redis-ключ (HASH) на конференцию: поле — identity участника (тот же
формат, что в LiveKit-токене и вебхуках — `str(user_id)` или
`guest:{guest_access.id}`), значение — JSON `{"name": ..., "raised_at": <unix
epoch>}`. `HSETNX` даёт атомарное «добавить, только если ещё нет» — повторное
поднятие уже поднятой руки НЕ сбрасывает её место в очереди (идемпотентно).
Порядок — сортировкой по `raised_at` при чтении снапшота (участников в одной
конференции — единицы-десятки, сортировка в Python здесь дешевле, чем держать
вторую структуру (ZSET) синхронно с первой).
"""
import json
import time
import uuid
from dataclasses import dataclass
from datetime import UTC, datetime
from core.redis import redis_client
from schemas.room_events import ForcedMuteOut, ForcedMuteSource, HandQueueEntryOut, HandQueueOut
# TTL ключа очереди — подстраховка на случай пропущенного webhook
# `room_finished` (см. `services/webhook_handlers.py::_on_room_finished`,
# который чистит очередь явно при штатном завершении). Сама конференция
# столько не длится ни при каких сценариях.
HAND_QUEUE_TTL_SECONDS = 24 * 60 * 60
def hand_queue_key(conference_id: uuid.UUID) -> str:
"""Redis-ключ HASH очереди поднятых рук конкретной конференции."""
return f"hand_queue:{conference_id}"
def hand_queue_channel(conference_id: uuid.UUID) -> str:
"""Redis pub/sub канал событий комнаты (очередь рук + принудительный мьют, задача B2)."""
return f"room_events:{conference_id}"
@dataclass(frozen=True, slots=True)
class HandQueueEntry:
"""Один участник в очереди поднятых рук."""
identity: str
name: str
raised_at: float
async def raise_hand(conference_id: uuid.UUID, *, identity: str, name: str) -> bool:
"""Поднять руку участника; `True` — рука реально поднялась (не была поднята раньше).
`HSETNX` — атомарная проверка-и-запись: если участник уже в очереди,
ничего не меняет (в т.ч. НЕ обновляет `raised_at`) — переподключение и
повторный клик не переставляют его в конец очереди.
"""
key = hand_queue_key(conference_id)
payload = json.dumps({"name": name, "raised_at": time.time()})
added = await redis_client.hsetnx(key, identity, payload)
await redis_client.expire(key, HAND_QUEUE_TTL_SECONDS)
return bool(added)
async def lower_hand(conference_id: uuid.UUID, *, identity: str) -> bool:
"""Опустить руку участника; `True` — рука была поднята и теперь снята."""
removed = await redis_client.hdel(hand_queue_key(conference_id), identity)
return bool(removed)
async def snapshot(conference_id: uuid.UUID) -> list[HandQueueEntry]:
"""Текущая очередь, упорядоченная по времени поднятия (раньше — раньше в списке)."""
raw = await redis_client.hgetall(hand_queue_key(conference_id))
entries = []
for identity, payload in raw.items():
try:
data = json.loads(payload)
entries.append(
HandQueueEntry(
identity=str(identity), name=data["name"], raised_at=data["raised_at"]
)
)
except (ValueError, KeyError, TypeError):
# Побитый/устаревшего формата элемент — пропускаем, а не роняем всю очередь.
continue
entries.sort(key=lambda entry: entry.raised_at)
return entries
async def clear(conference_id: uuid.UUID) -> None:
"""Полностью снести очередь конференции (штатное завершение — `room_finished`)."""
await redis_client.delete(hand_queue_key(conference_id))
def _to_out(entries: list[HandQueueEntry]) -> HandQueueOut:
"""Собрать исходящий снапшот из внутренних записей очереди."""
return HandQueueOut(
queue=[
HandQueueEntryOut(
identity=entry.identity,
name=entry.name,
raised_at=datetime.fromtimestamp(entry.raised_at, tz=UTC),
)
for entry in entries
]
)
async def get_snapshot_out(conference_id: uuid.UUID) -> HandQueueOut:
"""Текущая очередь в исходящем формате — для отправки сразу после подключения к WS."""
return _to_out(await snapshot(conference_id))
async def publish_snapshot(conference_id: uuid.UUID) -> None:
"""Опубликовать текущий снапшот очереди всем подписчикам канала комнаты.
Вызывается после любого изменения очереди (`raise_hand`/`lower_hand` —
из `api/chat.py`, а также `participant_left`/`room_finished` — из
`services/webhook_handlers.py`), чтобы у всех участников (и особенно у
организатора, зашедшего позже) была всегда актуальная картина.
"""
payload = _to_out(await snapshot(conference_id))
await redis_client.publish(hand_queue_channel(conference_id), payload.model_dump_json())
async def publish_forced_mute(
conference_id: uuid.UUID, *, identity: str, source: ForcedMuteSource
) -> None:
"""Оповестить всех участников комнаты о принудительном мьюте (задача B2).
Тот же канал, что и у очереди рук (`hand_queue_channel`) — `api/chat.py`
пересылает с него ЛЮБОЙ JSON как есть, различая события по полю `type`
(см. `_pump_pubsub_to_websocket`). Рассылается ВСЕМ, а не адресно
затронутому участнику: канал общий на конференцию, адресной доставки
одному соединению тут нет, поэтому клиент сам сверяет `identity` со
своей (см. `ForcedMuteOut` в `schemas/room_events.py`).
"""
payload = ForcedMuteOut(identity=identity, source=source)
await redis_client.publish(hand_queue_channel(conference_id), payload.model_dump_json())