Files
vidconf/backend/services/webhook_handlers.py
Max Ronzhin 32949ebc66 fix(webhook): запуск egress не блокирует транзакцию track_published
Обработчик `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`
явно, как и остальные тесты этой группы.
2026-07-28 18:59:53 +03:00

366 lines
17 KiB
Python
Raw 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.
"""Бизнес-логика обработки 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-контейнера в деплое нет) каждый микрофон превращался в
# заведомо безнадёжный сетевой вызов длиной в 2025 секунд.
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