Files
vidconf/backend/services/egress.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

135 lines
6.5 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.
"""Тонкая обёртка над LiveKit `EgressService.start_track_egress` — единственная точка для
мокирования в тестах (по образцу `workers/livekit_client.py::delete_livekit_room`).
Запускает Track Egress для одного аудиотрека: пишет исходный opus без
транскодирования в `.ogg` на общий volume `recordings_dir`. Финализация
результата (итоговый `location`/ошибка) приходит асинхронно через webhook
`egress_ended` — здесь только сам запуск и `egress_id`/`started_at` из ответа.
Запуск вынесен из тела webhook-обработчика в фоновую задачу
(`run_track_egress`) — см. докстринг этой функции.
"""
import asyncio
import logging
import uuid
from dataclasses import dataclass
from datetime import UTC, datetime
from livekit import api
from core.config import get_settings
from core.db import async_session_maker
from repositories.conferences import AudioTrackRepository
logger = logging.getLogger(__name__)
@dataclass(frozen=True, slots=True)
class EgressStartResult:
"""Результат запуска Track Egress: идентификатор задания и время старта (UTC)."""
egress_id: str
started_at: datetime
async def start_track_egress(room_name: str, track_sid: str, filepath: str) -> EgressStartResult:
"""Запустить запись одного аудиотрека комнаты в файл `filepath` на общем volume.
`DirectFileOutput` без указания облачного хранилища (s3/gcp/azure) пишет
файл напрямую на диск egress-контейнера — тот же volume `recordings_dir`,
что и у воркера транскрибации (см. `deploy/docker-compose.yml`).
"""
settings = get_settings()
lkapi = api.LiveKitAPI(
settings.livekit_url,
api_key=settings.livekit_api_key,
api_secret=settings.livekit_api_secret,
)
try:
# Без таймаута вызов висит до собственного таймаута psrpc LiveKit
# (2025 с, когда egress-воркера в деплое нет). `TimeoutError`
# ловит вызывающая сторона наравне с прочими ошибками запуска.
async with asyncio.timeout(settings.egress_start_timeout_s):
info = await lkapi.egress.start_track_egress(
api.TrackEgressRequest(
room_name=room_name,
track_id=track_sid,
file=api.DirectFileOutput(filepath=filepath),
)
)
finally:
await lkapi.aclose()
# `EgressInfo.started_at` — unix-наносекунды (см. документацию livekit/egress,
# `pkg/config/manifest.go`); при отсутствии (ещё не проставлен на момент
# ответа STARTING) считаем стартом текущий момент.
started_at = (
datetime.fromtimestamp(info.started_at / 1_000_000_000, tz=UTC)
if info.started_at
else datetime.now(UTC)
)
return EgressStartResult(egress_id=info.egress_id, started_at=started_at)
async def run_track_egress(
*,
room_name: str,
track_sid: str,
filepath: str,
session_id: uuid.UUID,
participant_id: uuid.UUID,
) -> None:
"""Запустить Track Egress и записать строку трека — ФОНОВАЯ задача вебхука.
Почему не внутри обработчика. `POST /api/v1/livekit/webhook` держит
соединение с БД и открытую транзакцию всё время своей работы (INSERT в
`livekit_webhook_events` сделан, commit — после обработчика). Пока здесь
жил сетевой вызов к egress, каждое событие `track_published` занимало
соединение на 2025 секунд, и на нагрузочном тесте 28.07.2026 пул
выгребался за секунды: 226 ошибок `QueuePool limit`, 37 ответов 500 на
путях входа в конференцию, 33 `dropped webhook` со стороны LiveKit
(очередь доставки росла до 56 секунд).
Поэтому обработчик отвечает 200 сразу, а сюда попадает только то, что
требует сети. Своя сессия БД (`async_session_maker`) обязательна: сессия
запроса к этому моменту уже закрыта вместе с ответом.
Идемпотентность сохраняется: `AudioTrackRepository.create` — это
`INSERT ... ON CONFLICT DO NOTHING` по `uq_session_track`, а проверка
«трек уже пишется» осталась в обработчике.
"""
try:
result = await start_track_egress(room_name, track_sid, filepath)
except Exception as exc: # noqa: BLE001 — недоступность egress не должна ронять фон
# Деплой-профиль (блок D): egress — необязательный сервис профиля
# `transcribe`; без него запись просто не стартует для этого трека.
# Строку `session_audio_tracks` не создаём — у нас нет `egress_id`,
# по которому её мог бы финализировать `egress_ended`.
logger.warning(
"track_published: не удалось запустить egress для трека %s сеанса %s: %s",
track_sid,
session_id,
exc,
)
return
async with async_session_maker() as session:
await AudioTrackRepository(session).create(
session_id=session_id,
participant_id=participant_id,
track_sid=track_sid,
egress_id=result.egress_id,
file_path=filepath,
started_at=result.started_at,
)
await session.commit()
logger.info(
"track_published: сеанс=%s участник=%s трек=%s egress=%s",
session_id,
participant_id,
track_sid,
result.egress_id,
)