Транспорт — существующий аутентифицированный 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-обработки, без которых кнопки принудительного мьюта не скомпилировались бы; сама реализация мьюта — следующим коммитом.
382 lines
18 KiB
Python
382 lines
18 KiB
Python
"""Бизнес-логика обработки webhook-событий LiveKit (ADR-001).
|
||
|
||
Дедупликация по `event.id` выполняется на уровне API-роутера
|
||
(`api/livekit_webhook.py`) в одной транзакции с эффектами обработчика —
|
||
этим обеспечивается идемпотентность пайплайна. Обработчики здесь
|
||
дополнительно используют get-or-create/guard-паттерны в репозитории
|
||
конференций, чтобы корректно восстанавливаться после пропущенных событий
|
||
(например, потерянного `room_started`).
|
||
|
||
Комната LiveKit больше не отдельная сущность — её имя всегда равно
|
||
`conferences.slug` (ADR-001, п.4), поэтому lookup идёт напрямую по slug.
|
||
"""
|
||
|
||
import logging
|
||
import uuid
|
||
from collections.abc import Callable
|
||
from datetime import UTC, datetime
|
||
|
||
from livekit.protocol.egress import EgressStatus
|
||
from livekit.protocol.models import TrackSource, TrackType
|
||
from livekit.protocol.webhook import WebhookEvent
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from core.config import get_settings
|
||
from repositories.conferences import (
|
||
AudioTrackRepository,
|
||
ConferenceRepository,
|
||
ConferenceSessionRepository,
|
||
)
|
||
from services import hand_queue
|
||
from services.egress import run_track_egress
|
||
from services.instance_settings import InstanceSettingsService
|
||
from services.pipeline_producer import enqueue_pipeline
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Префикс identity гостя в LiveKit-токене (ADR-001, п.6): `guest:{guest_access.id}`.
|
||
GUEST_IDENTITY_PREFIX = "guest:"
|
||
|
||
# Статусы EgressInfo, означающие неуспешное завершение записи (webhook `egress_ended`).
|
||
_EGRESS_FAILURE_STATUSES = frozenset(
|
||
{EgressStatus.EGRESS_FAILED, EgressStatus.EGRESS_ABORTED, EgressStatus.EGRESS_LIMIT_REACHED}
|
||
)
|
||
|
||
|
||
def _egress_ns_to_datetime(nanoseconds: int) -> datetime | None:
|
||
"""Перевести unix-наносекунды `EgressInfo.started_at`/`ended_at` в UTC datetime.
|
||
|
||
`0` (поле не проставлено) — валидное protobuf-значение по умолчанию, не
|
||
временная метка.
|
||
"""
|
||
if not nanoseconds:
|
||
return None
|
||
return datetime.fromtimestamp(nanoseconds / 1_000_000_000, tz=UTC)
|
||
|
||
|
||
class WebhookDispatcher:
|
||
"""Диспатчит `WebhookEvent` на обработчик по типу события.
|
||
|
||
`schedule` — планировщик фоновых задач: вызывается как
|
||
`schedule(coro_func, **kwargs)` и обязан вернуть управление немедленно,
|
||
не дожидаясь выполнения. В приложении это `BackgroundTasks.add_task`
|
||
FastAPI (задача стартует после отправки ответа), в тестах — вызовы
|
||
просто записываются. Через него уходит запуск Track Egress: сетевому
|
||
вызову не место внутри транзакции вебхука (см. `_on_track_published`).
|
||
"""
|
||
|
||
def __init__(
|
||
self,
|
||
session: AsyncSession,
|
||
*,
|
||
schedule: Callable[..., object],
|
||
) -> None:
|
||
self._conferences = ConferenceRepository(session)
|
||
self._sessions = ConferenceSessionRepository(session)
|
||
self._audio_tracks = AudioTrackRepository(session)
|
||
self._instance_settings = InstanceSettingsService(session)
|
||
self._schedule = schedule
|
||
|
||
async def dispatch(self, event: WebhookEvent) -> None:
|
||
"""Обработать одно webhook-событие; неизвестный тип события — no-op."""
|
||
handlers = {
|
||
"room_started": self._on_room_started,
|
||
"participant_joined": self._on_participant_joined,
|
||
"participant_left": self._on_participant_left,
|
||
"track_published": self._on_track_published,
|
||
"egress_ended": self._on_egress_ended,
|
||
"room_finished": self._on_room_finished,
|
||
}
|
||
handler = handlers.get(event.event)
|
||
if handler is None:
|
||
# Штатный шум: остальные типы событий (track_unpublished и т.п.) не нужны.
|
||
logger.debug("livekit webhook: необрабатываемый тип события %s", event.event)
|
||
return
|
||
await handler(event)
|
||
|
||
async def _on_room_started(self, event: WebhookEvent) -> None:
|
||
conference = await self._conferences.get_by_slug(event.room.name)
|
||
if conference is None:
|
||
logger.warning(
|
||
"livekit webhook room_started: конференция %s не найдена", event.room.name
|
||
)
|
||
return
|
||
|
||
conference.status = "active"
|
||
session_record = await self._sessions.get_or_create_open(
|
||
conference_id=conference.id, title=conference.title, t_start=datetime.now(UTC)
|
||
)
|
||
logger.info(
|
||
"livekit webhook room_started: конференция=%s сеанс=%s",
|
||
conference.id,
|
||
session_record.id,
|
||
)
|
||
|
||
async def _on_participant_joined(self, event: WebhookEvent) -> None:
|
||
conference = await self._conferences.get_by_slug(event.room.name)
|
||
if conference is None:
|
||
logger.warning(
|
||
"livekit webhook participant_joined: конференция %s не найдена", event.room.name
|
||
)
|
||
return
|
||
|
||
identity = _parse_identity(event.participant.identity)
|
||
if identity is None:
|
||
return
|
||
user_id, guest_id = identity
|
||
|
||
# Fallback на случай, если событие room_started было пропущено.
|
||
session_record = await self._sessions.get_or_create_open(
|
||
conference_id=conference.id, title=conference.title, t_start=datetime.now(UTC)
|
||
)
|
||
await self._sessions.add_participant(
|
||
session_id=session_record.id,
|
||
user_id=user_id,
|
||
guest_id=guest_id,
|
||
joined_at=datetime.now(UTC),
|
||
)
|
||
logger.info(
|
||
"livekit webhook participant_joined: конференция=%s identity=%s сеанс=%s",
|
||
conference.id,
|
||
event.participant.identity,
|
||
session_record.id,
|
||
)
|
||
|
||
async def _on_participant_left(self, event: WebhookEvent) -> None:
|
||
conference = await self._conferences.get_by_slug(event.room.name)
|
||
if conference is None:
|
||
logger.warning(
|
||
"livekit webhook participant_left: конференция %s не найдена", event.room.name
|
||
)
|
||
return
|
||
|
||
identity = _parse_identity(event.participant.identity)
|
||
if identity is None:
|
||
return
|
||
user_id, guest_id = identity
|
||
|
||
# Очередь поднятых рук живёт в Redis по `conference.id`, независимо
|
||
# от `ConferenceSession` (задача B1) — снимаем руку СРАЗУ, до guard'а
|
||
# на отсутствующий открытый сеанс ниже: пропущенный/задержанный
|
||
# `room_started` не должен оставлять фантомную запись в очереди у
|
||
# реально вышедшего участника. Не путать с обрывом WS-соединения
|
||
# самой очереди рук — то живёт своей жизнью и переживается без
|
||
# потери места (см. `services/hand_queue.py`).
|
||
removed = await hand_queue.lower_hand(conference.id, identity=event.participant.identity)
|
||
if removed:
|
||
await hand_queue.publish_snapshot(conference.id)
|
||
|
||
session_record = await self._sessions.get_open_by_conference(conference.id)
|
||
if session_record is None:
|
||
logger.warning(
|
||
"livekit webhook participant_left: нет открытого сеанса для конференции %s",
|
||
conference.id,
|
||
)
|
||
return
|
||
|
||
await self._sessions.mark_participant_left(
|
||
session_id=session_record.id,
|
||
user_id=user_id,
|
||
guest_id=guest_id,
|
||
left_at=datetime.now(UTC),
|
||
)
|
||
logger.info(
|
||
"livekit webhook participant_left: конференция=%s identity=%s сеанс=%s",
|
||
conference.id,
|
||
event.participant.identity,
|
||
session_record.id,
|
||
)
|
||
|
||
async def _on_track_published(self, event: WebhookEvent) -> None:
|
||
"""Запланировать Track Egress для опубликованного аудиотрека микрофона (ADR-002).
|
||
|
||
Видео/скриншеринг и т.п. — no-op (диаризация не нужна: транскрибируем
|
||
только речь, трек = спикер). Идемпотентно: если строка трека уже
|
||
существует (гонка повторной доставки), egress повторно не запускается.
|
||
|
||
Сам запуск уходит в фоновую задачу (`services.egress.run_track_egress`):
|
||
здесь остаются только быстрые проверки по БД, потому что обработчик
|
||
выполняется внутри открытой транзакции вебхука. Обоснование с цифрами —
|
||
в докстринге `run_track_egress`.
|
||
"""
|
||
if event.track.type != TrackType.AUDIO or event.track.source != TrackSource.MICROPHONE:
|
||
return
|
||
|
||
# Транскрибация выключена — записывать нечего. Тот же guard, что и в
|
||
# `_on_room_finished`: без него на инстансе без профиля `transcribe`
|
||
# (egress-контейнера в деплое нет) каждый микрофон превращался в
|
||
# заведомо безнадёжный сетевой вызов длиной в 20–25 секунд.
|
||
cfg = await self._instance_settings.get()
|
||
if not cfg.transcriber.enabled:
|
||
logger.debug(
|
||
"livekit webhook track_published: транскрибация выключена — трек %s пропущен",
|
||
event.track.sid,
|
||
)
|
||
return
|
||
|
||
conference = await self._conferences.get_by_slug(event.room.name)
|
||
if conference is None:
|
||
logger.warning(
|
||
"livekit webhook track_published: конференция %s не найдена", event.room.name
|
||
)
|
||
return
|
||
|
||
session_record = await self._sessions.get_open_by_conference(conference.id)
|
||
if session_record is None:
|
||
logger.warning(
|
||
"livekit webhook track_published: нет открытого сеанса для конференции %s",
|
||
conference.id,
|
||
)
|
||
return
|
||
|
||
existing_track = await self._audio_tracks.get_by_session_and_track(
|
||
session_record.id, event.track.sid
|
||
)
|
||
if existing_track is not None:
|
||
logger.debug(
|
||
"livekit webhook track_published: трек %s уже записывается (сеанс=%s)",
|
||
event.track.sid,
|
||
session_record.id,
|
||
)
|
||
return
|
||
|
||
identity = _parse_identity(event.participant.identity)
|
||
if identity is None:
|
||
return
|
||
user_id, guest_id = identity
|
||
|
||
participant = await self._sessions.get_active_participant(
|
||
session_record.id, user_id=user_id, guest_id=guest_id
|
||
)
|
||
if participant is None:
|
||
logger.warning(
|
||
"livekit webhook track_published: нет активного участника identity=%s сеанса %s",
|
||
event.participant.identity,
|
||
session_record.id,
|
||
)
|
||
return
|
||
|
||
settings = get_settings()
|
||
filepath = (
|
||
f"{settings.recordings_dir}/{session_record.id}/{participant.id}_{event.track.sid}.ogg"
|
||
)
|
||
self._schedule(
|
||
run_track_egress,
|
||
room_name=event.room.name,
|
||
track_sid=event.track.sid,
|
||
filepath=filepath,
|
||
session_id=session_record.id,
|
||
participant_id=participant.id,
|
||
)
|
||
|
||
async def _on_egress_ended(self, event: WebhookEvent) -> None:
|
||
"""Финализировать строку аудиотрека по результату Track Egress.
|
||
|
||
Успех (`EGRESS_COMPLETE`) -> `status='recorded'`; ошибка (failed/
|
||
aborted/limit_reached) -> `status='failed'`. Отсутствие строки трека
|
||
(например, потерянный `track_published`) — предупреждение, не ошибка.
|
||
"""
|
||
egress_info = event.egress_info
|
||
status = "failed" if egress_info.status in _EGRESS_FAILURE_STATUSES else "recorded"
|
||
ended_at = _egress_ns_to_datetime(egress_info.ended_at) or datetime.now(UTC)
|
||
file_path = egress_info.file.filename or egress_info.file.location or None
|
||
|
||
record = await self._audio_tracks.finalize(
|
||
egress_id=egress_info.egress_id,
|
||
status=status,
|
||
ended_at=ended_at,
|
||
file_path=file_path,
|
||
)
|
||
if record is None:
|
||
logger.warning(
|
||
"livekit webhook egress_ended: строка трека для egress %s не найдена",
|
||
egress_info.egress_id,
|
||
)
|
||
return
|
||
|
||
logger.info(
|
||
"livekit webhook egress_ended: трек=%s egress=%s статус=%s",
|
||
record.id,
|
||
egress_info.egress_id,
|
||
status,
|
||
)
|
||
|
||
async def _on_room_finished(self, event: WebhookEvent) -> None:
|
||
conference = await self._conferences.get_by_slug(event.room.name)
|
||
if conference is None:
|
||
logger.warning(
|
||
"livekit webhook room_finished: конференция %s не найдена", event.room.name
|
||
)
|
||
return
|
||
|
||
session_record = await self._sessions.get_open_by_conference(conference.id)
|
||
if session_record is None:
|
||
logger.warning(
|
||
"livekit webhook room_finished: нет открытого сеанса для конференции %s",
|
||
conference.id,
|
||
)
|
||
return
|
||
|
||
now = datetime.now(UTC)
|
||
await self._sessions.close(session_record, t_end=now)
|
||
await self._sessions.close_all_open_participants(session_id=session_record.id, left_at=now)
|
||
# Очередь поднятых рук — состояние звонка, не история; следующий
|
||
# заход (в т.ч. у закреплённой конференции) должен начинать с чистой
|
||
# очереди, а не наследовать поднятые руки из прошлого раза.
|
||
await hand_queue.clear(conference.id)
|
||
|
||
# Незакреплённая умирает по завершении (история/саммари остаются);
|
||
# закреплённая возвращается в ожидание следующего вхождения (ADR-001, п.2).
|
||
if conference.is_pinned:
|
||
conference.status = "scheduled"
|
||
else:
|
||
conference.status = "ended"
|
||
conference.ended_at = now
|
||
|
||
logger.info(
|
||
"livekit webhook room_finished: конференция=%s сеанс=%s новый статус=%s",
|
||
conference.id,
|
||
session_record.id,
|
||
conference.status,
|
||
)
|
||
|
||
# Постановка AI-пайплайна: запись треков
|
||
# завершена (egress ещё может дописывать файлы — это ждёт шаг 3
|
||
# `run_pipeline`), сеанс закрыт — можно ставить задачу в очередь.
|
||
#
|
||
# Guard от зависания сеанса:
|
||
# если транскрибация выключена в настройках инстанса (пресеты 1/2
|
||
# инсталлятора — без AI, воркеров/LLM в деплое нет), задача `transcribe`
|
||
# уйдёт в очередь `transcription`, которую некому обслужить, и сеанс
|
||
# навсегда застрянет в `pipeline_status='recording'`. Вместо постановки
|
||
# в очередь сразу проставляем терминальный статус без AI-шагов —
|
||
# 'notified' (тот же статус, которым штатно завершается полный
|
||
# пайплайн; отдельный enum-статус/миграция не нужны).
|
||
cfg = await self._instance_settings.get()
|
||
if cfg.transcriber.enabled:
|
||
enqueue_pipeline(session_record.id)
|
||
else:
|
||
session_record.pipeline_status = "notified"
|
||
logger.info(
|
||
"livekit webhook room_finished: транскрибация выключена — "
|
||
"сеанс=%s сразу переведён в pipeline_status='notified'",
|
||
session_record.id,
|
||
)
|
||
|
||
|
||
def _parse_identity(identity: str) -> tuple[uuid.UUID | None, uuid.UUID | None] | None:
|
||
"""Распарсить identity участника в пару (`user_id`, `guest_id`) — ровно один заполнен.
|
||
|
||
`guest:{guest_access.id}` — гость; иначе — `str(user.id)` зарегистрированного
|
||
пользователя. Невалидный/пустой identity — предупреждение в лог и пропуск
|
||
события (не должно приводить к 500).
|
||
"""
|
||
try:
|
||
if identity.startswith(GUEST_IDENTITY_PREFIX):
|
||
guest_id = uuid.UUID(identity.removeprefix(GUEST_IDENTITY_PREFIX))
|
||
return None, guest_id
|
||
return uuid.UUID(identity), None
|
||
except (ValueError, AttributeError):
|
||
logger.warning("livekit webhook: невалидный identity участника %r", identity)
|
||
return None
|