diff --git a/backend/api/livekit_webhook.py b/backend/api/livekit_webhook.py index c57eeaf..3eb420a 100644 --- a/backend/api/livekit_webhook.py +++ b/backend/api/livekit_webhook.py @@ -3,7 +3,7 @@ import logging from typing import Annotated, Any, cast -from fastapi import APIRouter, Depends, Header, HTTPException, Request, status +from fastapi import APIRouter, BackgroundTasks, Depends, Header, HTTPException, Request, status from livekit import api from sqlalchemy import CursorResult from sqlalchemy.dialects.postgresql import insert as pg_insert @@ -22,6 +22,7 @@ router = APIRouter(prefix="/api/v1/livekit", tags=["livekit"]) @router.post("/webhook") async def receive_webhook( request: Request, + background: BackgroundTasks, session: Annotated[AsyncSession, Depends(get_session)], authorization: Annotated[str | None, Header()] = None, ) -> dict[str, str]: @@ -30,6 +31,11 @@ async def receive_webhook( Дедупликация по `event.id`: `INSERT ... ON CONFLICT DO NOTHING` в `livekit_webhook_events` в одной транзакции с эффектами обработчика — при конфликте (дубль) эффекты пропускаются, но ответ всё равно 200. + + Всё, что требует сети (запуск Track Egress), уходит в `background` и + выполняется уже после ответа: обработчик держит соединение с БД и + открытую транзакцию, а LiveKit при медленном ответе копит очередь + доставки и в итоге дропает события (см. `services.egress.run_track_egress`). """ settings = get_settings() raw_body = await request.body() @@ -57,7 +63,7 @@ async def receive_webhook( await session.commit() return {"status": "duplicate"} - dispatcher = WebhookDispatcher(session) + dispatcher = WebhookDispatcher(session, schedule=background.add_task) await dispatcher.dispatch(event) await session.commit() return {"status": "ok"} diff --git a/backend/services/egress.py b/backend/services/egress.py index 814e058..bf95ac5 100644 --- a/backend/services/egress.py +++ b/backend/services/egress.py @@ -5,14 +5,24 @@ транскодирования в `.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) @@ -37,13 +47,17 @@ async def start_track_egress(room_name: str, track_sid: str, filepath: str) -> E api_secret=settings.livekit_api_secret, ) try: - info = await lkapi.egress.start_track_egress( - api.TrackEgressRequest( - room_name=room_name, - track_id=track_sid, - file=api.DirectFileOutput(filepath=filepath), + # Без таймаута вызов висит до собственного таймаута 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() @@ -56,3 +70,65 @@ async def start_track_egress(room_name: str, track_sid: str, filepath: str) -> E 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, + ) diff --git a/backend/services/webhook_handlers.py b/backend/services/webhook_handlers.py index 99c8cc6..f21ce83 100644 --- a/backend/services/webhook_handlers.py +++ b/backend/services/webhook_handlers.py @@ -13,6 +13,7 @@ import logging import uuid +from collections.abc import Callable from datetime import UTC, datetime from livekit.protocol.egress import EgressStatus @@ -26,7 +27,7 @@ from repositories.conferences import ( ConferenceRepository, ConferenceSessionRepository, ) -from services.egress import start_track_egress +from services.egress import run_track_egress from services.instance_settings import InstanceSettingsService from services.pipeline_producer import enqueue_pipeline @@ -53,13 +54,27 @@ def _egress_ns_to_datetime(nanoseconds: int) -> datetime | None: class WebhookDispatcher: - """Диспатчит `WebhookEvent` на обработчик по типу события.""" + """Диспатчит `WebhookEvent` на обработчик по типу события. - def __init__(self, session: AsyncSession) -> None: + `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.""" @@ -161,15 +176,32 @@ class WebhookDispatcher: ) async def _on_track_published(self, event: WebhookEvent) -> None: - """Запустить Track Egress для опубликованного аудиотрека микрофона (ADR-002). + """Запланировать 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( @@ -216,37 +248,13 @@ class WebhookDispatcher: filepath = ( f"{settings.recordings_dir}/{session_record.id}/{participant.id}_{event.track.sid}.ogg" ) - try: - result = await start_track_egress(event.room.name, event.track.sid, filepath) - except Exception as exc: # noqa: BLE001 — недоступность egress не должна ронять webhook - # Деплой-профиль (блок D): egress — необязательный сервис профиля - # `transcribe`; без него запись просто не стартует для этого трека - # (риск «Потерян webhook track_published»). - # Строку `session_audio_tracks` не создаём — у нас нет `egress_id`, - # по которому её мог бы финализировать `egress_ended`. - logger.warning( - "livekit webhook track_published: не удалось запустить egress для трека %s " - "сеанса %s: %s", - event.track.sid, - session_record.id, - exc, - ) - return - - await self._audio_tracks.create( + 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, - track_sid=event.track.sid, - egress_id=result.egress_id, - file_path=filepath, - started_at=result.started_at, - ) - logger.info( - "livekit webhook track_published: сеанс=%s участник=%s трек=%s egress=%s", - session_record.id, - participant.id, - event.track.sid, - result.egress_id, ) async def _on_egress_ended(self, event: WebhookEvent) -> None: diff --git a/backend/tests/test_livekit_webhook.py b/backend/tests/test_livekit_webhook.py index 1313423..aee7014 100644 --- a/backend/tests/test_livekit_webhook.py +++ b/backend/tests/test_livekit_webhook.py @@ -3,9 +3,13 @@ цикл закреплённой/незакреплённой конференции — на фикстурах payload'ов LiveKit. """ +import asyncio import base64 import hashlib +import json import uuid +from collections.abc import AsyncGenerator +from contextlib import asynccontextmanager from datetime import UTC, datetime, timedelta from pathlib import Path from unittest.mock import AsyncMock, Mock @@ -13,10 +17,13 @@ from unittest.mock import AsyncMock, Mock import httpx import jwt import pytest +from google.protobuf.json_format import ParseDict +from livekit.protocol.webhook import WebhookEvent from sqlalchemy import select from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.ext.asyncio import AsyncSession +import services.egress as egress_module import services.webhook_handlers as webhook_handlers_module from core.config import get_settings from core.security import hash_password @@ -54,6 +61,52 @@ def _sign(body: bytes) -> str: return jwt.encode(payload, settings.livekit_api_secret, algorithm="HS256") +def _parse_fixture_event(name: str, **placeholders: str) -> WebhookEvent: + """Разобрать фикстуру в `WebhookEvent` — для тестов диспатчера без HTTP-слоя.""" + return ParseDict(json.loads(_load_fixture(name, **placeholders)), WebhookEvent()) + + +async def _set_transcriber_enabled(session: AsyncSession, *, enabled: bool) -> None: + """Выставить `instance_settings.transcriber.enabled`. + + Пишется тем же `db_session` (savepoint), что и обработчик webhook (подмена + `get_session` в фикстуре `app`) — видна обработчику без реального коммита + в dev-БД. Тесты, которым важен `track_published`, обязаны выставлять флаг + ЯВНО: значение в dev-БД непредсказуемо, а обработчик с версии 0.0.12 + выходит на выключенной транскрибации раньше всех остальных проверок. + """ + value: dict[str, object] = { + "enabled": enabled, + "provider": "faster_whisper_cpu" if enabled else "null", + "model": "small" if enabled else None, + "language": "ru", + "options": {}, + } + await session.execute( + pg_insert(InstanceSetting) + .values(key="transcriber", value=value) + .on_conflict_do_update(index_elements=["key"], set_={"value": value}) + ) + + +def _use_test_session_in_background( + monkeypatch: pytest.MonkeyPatch, db_session: AsyncSession +) -> None: + """Заставить фоновую задачу egress работать с тестовой (savepoint) сессией. + + `run_track_egress` намеренно берёт СВОЮ сессию (`async_session_maker`): + в бою сессия запроса к моменту фоновой задачи уже закрыта. В тестах такое + подключение шло бы мимо откатываемой транзакции и не увидело бы ни + конференции, ни участника — поэтому подменяем фабрику на тестовую сессию. + """ + + @asynccontextmanager + async def _maker() -> AsyncGenerator[AsyncSession, None]: + yield db_session + + monkeypatch.setattr(egress_module, "async_session_maker", _maker) + + async def _post_webhook(client: httpx.AsyncClient, body: bytes) -> httpx.Response: return await client.post( WEBHOOK_URL, @@ -319,8 +372,10 @@ async def test_track_published_by_guest_starts_egress_and_creates_track_row( mock_start = AsyncMock( return_value=EgressStartResult(egress_id="EG_guest_track", started_at=started_at) ) - monkeypatch.setattr(webhook_handlers_module, "start_track_egress", mock_start) + monkeypatch.setattr(egress_module, "start_track_egress", mock_start) + _use_test_session_in_background(monkeypatch, db_session) + await _set_transcriber_enabled(db_session, enabled=True) conference = await _make_conference(db_session, generate_slug()) guest = await _make_guest(db_session, conference) await db_session.commit() @@ -382,8 +437,10 @@ async def test_track_published_survives_egress_unavailable( ) -> None: """Недоступность egress не должна ронять webhook (блок D): 200 + warning, без строки трека.""" mock_start = AsyncMock(side_effect=RuntimeError("egress service unavailable")) - monkeypatch.setattr(webhook_handlers_module, "start_track_egress", mock_start) + monkeypatch.setattr(egress_module, "start_track_egress", mock_start) + _use_test_session_in_background(monkeypatch, db_session) + await _set_transcriber_enabled(db_session, enabled=True) conference = await _make_conference(db_session, generate_slug()) user = await _make_user(db_session, "webhook-track-egress-down@example.com") await db_session.commit() @@ -417,13 +474,169 @@ async def test_track_published_survives_egress_unavailable( assert tracks == [] +async def test_track_published_skips_egress_when_transcription_disabled( + client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + """Транскрибация выключена → egress не дёргается вовсе (релиз 0.0.12). + + Симметрично guard'у в `_on_room_finished`. Без него на инстансе без профиля + `transcribe` (egress-контейнера в деплое нет) каждый микрофонный трек + превращался в заведомо безнадёжный вызов длиной в таймаут psrpc LiveKit — + 20–25 секунд внутри открытой транзакции вебхука. На нагрузочном тесте + 28.07.2026 это выгребало пул соединений и роняло вход в конференцию в 500. + """ + mock_start = AsyncMock() + monkeypatch.setattr(egress_module, "start_track_egress", mock_start) + _use_test_session_in_background(monkeypatch, db_session) + + await _set_transcriber_enabled(db_session, enabled=False) + conference = await _make_conference(db_session, generate_slug()) + user = await _make_user(db_session, "webhook-track-transcriber-off@example.com") + await db_session.commit() + + joined = _load_fixture( + "participant_joined.json", + event_id=f"evt-{uuid.uuid4()}", + room_name=conference.slug, + identity=str(user.id), + ) + assert (await _post_webhook(client, joined)).status_code == 200 + + track_sid = "TR_transcriber_off" + published = _load_fixture( + "track_published.json", + event_id=f"evt-{uuid.uuid4()}", + room_name=conference.slug, + identity=str(user.id), + track_sid=track_sid, + ) + assert (await _post_webhook(client, published)).status_code == 200 + + mock_start.assert_not_awaited() + tracks = ( + await db_session.scalars( + select(SessionAudioTrack).where(SessionAudioTrack.track_sid == track_sid) + ) + ).all() + assert tracks == [] + + +async def test_track_published_does_not_call_egress_inside_transaction( + db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch +) -> None: + """Запуск egress уходит в фон, а не выполняется внутри обработчика (релиз 0.0.12). + + Проверяется не время ответа (в тестах ASGI-транспорт дожидается фоновых + задач), а сама суть: пока открыта транзакция вебхука, сетевого вызова не + происходит — обработчик только планирует задачу. Именно это разгружает пул + соединений: до правки вызов жил внутри транзакции и держал соединение + 20–25 секунд, когда egress-сервиса в деплое нет. + """ + mock_start = AsyncMock( + return_value=EgressStartResult(egress_id="EG_bg", started_at=datetime.now(UTC)) + ) + monkeypatch.setattr(egress_module, "start_track_egress", mock_start) + + scheduled: list[tuple[object, dict[str, object]]] = [] + + def _schedule(func: object, **kwargs: object) -> None: + scheduled.append((func, kwargs)) + + await _set_transcriber_enabled(db_session, enabled=True) + conference = await _make_conference(db_session, generate_slug()) + user = await _make_user(db_session, "webhook-track-background@example.com") + await db_session.commit() + + dispatcher = webhook_handlers_module.WebhookDispatcher(db_session, schedule=_schedule) + joined_event = _parse_fixture_event( + "participant_joined.json", + event_id=f"evt-{uuid.uuid4()}", + room_name=conference.slug, + identity=str(user.id), + ) + await dispatcher.dispatch(joined_event) + + track_sid = "TR_background" + published_event = _parse_fixture_event( + "track_published.json", + event_id=f"evt-{uuid.uuid4()}", + room_name=conference.slug, + identity=str(user.id), + track_sid=track_sid, + ) + await dispatcher.dispatch(published_event) + + # Сеть не тронута: обработчик только запланировал задачу. + mock_start.assert_not_awaited() + assert len(scheduled) == 1 + func, kwargs = scheduled[0] + assert func is egress_module.run_track_egress + assert kwargs["room_name"] == conference.slug + assert kwargs["track_sid"] == track_sid + + session_record = await db_session.scalar( + select(ConferenceSession).where(ConferenceSession.conference_id == conference.id) + ) + assert session_record is not None + assert kwargs["session_id"] == session_record.id + + # А вот запущенная задача действительно ходит в egress и пишет строку. + _use_test_session_in_background(monkeypatch, db_session) + await egress_module.run_track_egress(**kwargs) # type: ignore[arg-type] + mock_start.assert_awaited_once() + + track_row = await db_session.scalar( + select(SessionAudioTrack).where(SessionAudioTrack.track_sid == track_sid) + ) + assert track_row is not None + assert track_row.egress_id == "EG_bg" + + +async def test_start_track_egress_gives_up_on_timeout(monkeypatch: pytest.MonkeyPatch) -> None: + """Запуск egress не ждёт дольше `egress_start_timeout_s` (релиз 0.0.12). + + Когда egress-воркера нет, LiveKit держит вызов до собственного таймаута + psrpc — на тесте 28.07.2026 это было 20–25 секунд на каждый микрофонный + трек. Живой egress отвечает за доли секунды, ждать столько незачем. + """ + settings = get_settings() + monkeypatch.setattr(settings, "egress_start_timeout_s", 0.05, raising=False) + + closed = False + + class _HangingEgress: + async def start_track_egress(self, _request: object) -> object: + await asyncio.sleep(5) + raise AssertionError("вызов должен был прерваться по таймауту") + + class _HangingApi: + def __init__(self, *_args: object, **_kwargs: object) -> None: + self.egress = _HangingEgress() + + async def aclose(self) -> None: + nonlocal closed + closed = True + + # Строковая форма: `api` в `services.egress` — реэкспорт из livekit SDK, + # обращение к нему атрибутом mypy считает неявным экспортом. + monkeypatch.setattr("services.egress.api.LiveKitAPI", _HangingApi) + + with pytest.raises(TimeoutError): + await egress_module.start_track_egress("room", "TR_hang", "/recordings/x.ogg") + + # Клиент закрывается и на неуспешном пути — иначе утекали бы соединения. + assert closed is True + + async def test_track_published_video_is_noop( client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch ) -> None: """№9 плана (часть 1): video-трек — no-op, egress не запускается.""" mock_start = AsyncMock() - monkeypatch.setattr(webhook_handlers_module, "start_track_egress", mock_start) + monkeypatch.setattr(egress_module, "start_track_egress", mock_start) + _use_test_session_in_background(monkeypatch, db_session) + await _set_transcriber_enabled(db_session, enabled=True) conference = await _make_conference(db_session, generate_slug()) user = await _make_user(db_session, "webhook-track-video@example.com") await db_session.commit() @@ -462,8 +675,10 @@ async def test_track_published_repeated_webhook_creates_single_row( mock_start = AsyncMock( return_value=EgressStartResult(egress_id="EG_repeat", started_at=datetime.now(UTC)) ) - monkeypatch.setattr(webhook_handlers_module, "start_track_egress", mock_start) + monkeypatch.setattr(egress_module, "start_track_egress", mock_start) + _use_test_session_in_background(monkeypatch, db_session) + await _set_transcriber_enabled(db_session, enabled=True) conference = await _make_conference(db_session, generate_slug()) user = await _make_user(db_session, "webhook-track-repeat@example.com") await db_session.commit() @@ -508,6 +723,8 @@ async def test_egress_ended_finalizes_track_success_and_failure( client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch ) -> None: """№10 плана (часть 1): `egress_ended` — 'recorded' на успехе, 'failed' на ошибке.""" + _use_test_session_in_background(monkeypatch, db_session) + await _set_transcriber_enabled(db_session, enabled=True) conference = await _make_conference(db_session, generate_slug()) user = await _make_user(db_session, "webhook-egress-ended@example.com") await db_session.commit() @@ -527,7 +744,7 @@ async def test_egress_ended_finalizes_track_success_and_failure( # Успешная запись. ok_started = datetime.now(UTC) monkeypatch.setattr( - webhook_handlers_module, + egress_module, "start_track_egress", AsyncMock(return_value=EgressStartResult(egress_id="EG_ok", started_at=ok_started)), ) @@ -559,7 +776,7 @@ async def test_egress_ended_finalizes_track_success_and_failure( # Ошибка записи. monkeypatch.setattr( - webhook_handlers_module, + egress_module, "start_track_egress", AsyncMock( return_value=EgressStartResult(egress_id="EG_fail", started_at=datetime.now(UTC)) @@ -597,6 +814,10 @@ async def test_room_finished_enqueues_pipeline( mock_enqueue = Mock() monkeypatch.setattr(webhook_handlers_module, "enqueue_pipeline", mock_enqueue) + # Флаг выставляется явно: `_on_room_finished` ставит задачу в очередь только + # при включённой транскрибации, а состояние `instance_settings` в dev-БД + # непредсказуемо (тест падал, если в базе оставалось `enabled: false`). + await _set_transcriber_enabled(db_session, enabled=True) conference = await _make_conference(db_session, generate_slug()) await db_session.commit()