358 lines
16 KiB
Python
358 lines
16 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 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.egress import start_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` на обработчик по типу события."""
|
||
|
||
def __init__(self, session: AsyncSession) -> None:
|
||
self._conferences = ConferenceRepository(session)
|
||
self._sessions = ConferenceSessionRepository(session)
|
||
self._audio_tracks = AudioTrackRepository(session)
|
||
self._instance_settings = InstanceSettingsService(session)
|
||
|
||
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
|
||
|
||
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 повторно не запускается.
|
||
"""
|
||
if event.track.type != TrackType.AUDIO or event.track.source != TrackSource.MICROPHONE:
|
||
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"
|
||
)
|
||
try:
|
||
result = await start_track_egress(event.room.name, event.track.sid, filepath)
|
||
except Exception as exc: # noqa: BLE001 — недоступность egress не должна ронять webhook
|
||
# Деплой-профиль (блок D): egress — необязательный сервис профиля
|
||
# `transcribe`; без него запись просто не стартует для этого трека
|
||
# (риск «Потерян webhook track_published»).
|
||
# Строку `session_audio_tracks` не создаём — у нас нет `egress_id`,
|
||
# по которому её мог бы финализировать `egress_ended`.
|
||
logger.warning(
|
||
"livekit webhook track_published: не удалось запустить egress для трека %s "
|
||
"сеанса %s: %s",
|
||
event.track.sid,
|
||
session_record.id,
|
||
exc,
|
||
)
|
||
return
|
||
|
||
await self._audio_tracks.create(
|
||
session_id=session_record.id,
|
||
participant_id=participant.id,
|
||
track_sid=event.track.sid,
|
||
egress_id=result.egress_id,
|
||
file_path=filepath,
|
||
started_at=result.started_at,
|
||
)
|
||
logger.info(
|
||
"livekit webhook track_published: сеанс=%s участник=%s трек=%s egress=%s",
|
||
session_record.id,
|
||
participant.id,
|
||
event.track.sid,
|
||
result.egress_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)
|
||
|
||
# Незакреплённая умирает по завершении (история/саммари остаются);
|
||
# закреплённая возвращается в ожидание следующего вхождения (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
|