Обработчик `track_published` вызывал `start_track_egress` внутри своей транзакции. На инстансе без профиля `transcribe` egress-сервиса нет, и LiveKit ждал ответа воркера через Redis до собственного таймаута psrpc — 20–25 секунд на каждый микрофонный трек. Всё это время webhook удерживал соединение с БД и открытую транзакцию. На нагрузочном тесте с 19 участниками (28.07.2026) это дало 226 ошибок `QueuePool limit of size 5 overflow 10 reached` и 37 ответов 500 на путях входа в конференцию, а со стороны LiveKit — 33 дропнутых webhook при очереди доставки до 56 секунд. Что изменилось: - запуск ушёл в фоновую задачу `run_track_egress` со своей сессией БД; обработчик только планирует её и отвечает 200 сразу; - добавлен ранний выход по `transcriber.enabled` — симметрично guard'у, который уже был в `room_finished`; - запуск ограничен таймаутом `egress_start_timeout_s` (по умолчанию 3 с). Идемпотентность сохранена: проверка «трек уже пишется» осталась в обработчике, а `AudioTrackRepository.create` — это INSERT ... ON CONFLICT DO NOTHING. Попутно: `test_room_finished_enqueues_pipeline` падал в зависимости от того, что осталось в локальной БД, — теперь выставляет `transcriber` явно, как и остальные тесты этой группы.
366 lines
17 KiB
Python
366 lines
17 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.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
|
||
|
||
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)
|
||
|
||
# Незакреплённая умирает по завершении (история/саммари остаются);
|
||
# закреплённая возвращается в ожидание следующего вхождения (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
|