Обработчик `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` явно, как и остальные тесты этой группы.
135 lines
6.5 KiB
Python
135 lines
6.5 KiB
Python
"""Тонкая обёртка над 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
|
||
# (20–25 с, когда 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` занимало
|
||
соединение на 20–25 секунд, и на нагрузочном тесте 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,
|
||
)
|