"""Тонкая обёртка над 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, )