3 Commits

Author SHA1 Message Date
7549b53ec9 release: версия 0.0.12
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
2026-07-28 19:00:09 +03:00
6d65b620fe perf(backend): явный пул соединений БД и несколько воркеров uvicorn
Пул создавался с дефолтом SQLAlchemy (5 + 10) и под нагрузкой выгребался
за секунды. Теперь параметры заданы явно и вынесены в настройки:
DB_POOL_SIZE=10, DB_MAX_OVERFLOW=10, DB_POOL_TIMEOUT=10. Таймаут снижен с
дефолтных 30 секунд намеренно — пусть запрос падает быстро и показывает
проблему, а не висит полминуты.

Backend запускался одним процессом uvicorn: любой блокирующий вызов
останавливал и параллельные запросы, и WS-чат всех участников. Добавлен
UVICORN_WORKERS с дефолтом 2 — не по числу ядер, потому что на
четырёхъядерном сервере ядра делятся с LiveKit, а медиа важнее API.

Бюджет соединений считается на весь инстанс: каждый воркер держит свой
пул, поэтому UVICORN_WORKERS × (DB_POOL_SIZE + DB_MAX_OVERFLOW) должно
оставаться заметно ниже max_connections у Postgres.

Многопроцессность безопасна: бутстрап настроек в lifespan идемпотентен
(INSERT ... ON CONFLICT DO NOTHING), а WS-чат разносит сообщения через
Redis pub/sub и состояния в памяти процесса не держит.
2026-07-28 19:00:05 +03:00
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
11 changed files with 454 additions and 52 deletions

View File

@@ -81,6 +81,20 @@ LIVEKIT_NODE_IP=127.0.0.1
# с точкой монтирования тома в обоих сервисах. # с точкой монтирования тома в обоих сервисах.
RECORDINGS_DIR=/recordings RECORDINGS_DIR=/recordings
# --- Производительность backend ---
# Число процессов uvicorn. Один процесс означает, что любой блокирующий вызов
# в обработчике останавливает весь event loop: параллельные входы в конференцию
# и WS-чат всех участников встают в очередь. Дефолт 2 рассчитан на 4-ядерный
# сервер, где ядра делятся с LiveKit (медиа важнее API).
UVICORN_WORKERS=2
# Пул соединений с БД НА КАЖДЫЙ воркер. Общий расход инстанса —
# UVICORN_WORKERS × (DB_POOL_SIZE + DB_MAX_OVERFLOW), плюс соединения Celery
# и alembic. Держите сумму заметно ниже max_connections у Postgres (по
# умолчанию 100), иначе вместо понятной ошибки приложения получите отказ БД.
DB_POOL_SIZE=10
DB_MAX_OVERFLOW=10
DB_POOL_TIMEOUT=10
# --- Email (рассылка саммари + .ics-приглашения) --- # --- Email (рассылка саммари + .ics-приглашения) ---
# `console` — дефолт для dev (письмо только логируется, ссылка подтверждения # `console` — дефолт для dev (письмо только логируется, ссылка подтверждения
# email берётся из логов); `smtp` — реальная отправка через aiosmtplib. # email берётся из логов); `smtp` — реальная отправка через aiosmtplib.
@@ -98,7 +112,7 @@ SMTP_TIMEOUT_S=30
# --- Версия инстанса (релиз v0.0.1) --- # --- Версия инстанса (релиз v0.0.1) ---
# install.sh копирует значение из корневого файла VERSION при каждой # install.sh копирует значение из корневого файла VERSION при каждой
# установке/обновлении — руками менять не нужно. # установке/обновлении — руками менять не нужно.
VIDCONF_VERSION=0.0.11 VIDCONF_VERSION=0.0.12
# --- Профили compose. Дефолт ниже (`media,monitoring`) — только для ручного # --- Профили compose. Дефолт ниже (`media,monitoring`) — только для ручного
# `docker compose up` БЕЗ install.sh: медиа (LiveKit+coturn) + мониторинг, # `docker compose up` БЕЗ install.sh: медиа (LiveKit+coturn) + мониторинг,

View File

@@ -3,6 +3,35 @@
Формат основан на [Keep a Changelog](https://keepachangelog.com/ru/1.1.0/), Формат основан на [Keep a Changelog](https://keepachangelog.com/ru/1.1.0/),
проект придерживается [семантического версионирования](https://semver.org/lang/ru/). проект придерживается [семантического версионирования](https://semver.org/lang/ru/).
## [0.0.12] — 2026-07-28
Разблокировка backend под нагрузкой: вход в конференцию перестаёт отваливаться,
когда участники включают микрофоны.
### Исправлено
- Обработчик webhook `track_published` больше не запускает Track Egress внутри
своей транзакции. Раньше каждый опубликованный микрофон уходил в сетевой
вызов, а на инстансе без профиля `transcribe` (egress-сервиса в деплое нет)
этот вызов висел 2025 секунд, всё это время удерживая соединение с БД. На
нагрузочном тесте с 19 участниками пул соединений выгребался за секунды, и
вход в конференцию начинал отвечать 500. Теперь запуск уходит в фоновую
задачу со своей сессией, а webhook отвечает сразу — LiveKit перестаёт копить
очередь доставки и терять события.
- Track Egress не запускается вовсе, если транскрибация выключена в настройках
инстанса — тот же guard, что уже был в обработчике `room_finished`.
- Запуск Track Egress ограничен таймаутом (по умолчанию 3 секунды) вместо
ожидания собственного таймаута LiveKit.
### Изменено
- Пул соединений с БД задаётся явно (`DB_POOL_SIZE`, `DB_MAX_OVERFLOW`,
`DB_POOL_TIMEOUT`) вместо дефолта SQLAlchemy 5 + 10. Считайте бюджет на весь
инстанс: каждый воркер держит свой пул, и сумма должна оставаться заметно
ниже `max_connections` у Postgres.
- Backend запускается с несколькими процессами uvicorn (`UVICORN_WORKERS`,
по умолчанию 2). Одиночный процесс означал, что любой блокирующий вызов
останавливает и параллельные запросы, и WS-чат всех участников. Дефолт 2, а
не по числу ядер: на четырёхъядерном сервере ядра делятся с LiveKit.
## [0.0.11] — 2026-07-28 ## [0.0.11] — 2026-07-28
Выбор режима показа участников в конференции, скрытие остальных и круглая Выбор режима показа участников в конференции, скрытие остальных и круглая

View File

@@ -1 +1 @@
0.0.11 0.0.12

View File

@@ -37,4 +37,13 @@ EXPOSE 8000
HEALTHCHECK --interval=10s --timeout=5s --retries=10 --start-period=15s \ HEALTHCHECK --interval=10s --timeout=5s --retries=10 --start-period=15s \
CMD python -c "import urllib.request; urllib.request.urlopen('http://localhost:8000/api/health')" || exit 1 CMD python -c "import urllib.request; urllib.request.urlopen('http://localhost:8000/api/health')" || exit 1
CMD ["uv", "run", "uvicorn", "main:create_app", "--factory", "--host", "0.0.0.0", "--port", "8000"] # Число воркеров — из окружения (`UVICORN_WORKERS`, см. docker-compose.yml).
# Один процесс означает, что любой блокирующий вызов в обработчике
# останавливает весь event loop: параллельные запросы и WS-чат всех
# участников встают в очередь. Дефолт 2, а не «по числу ядер»: медиа
# важнее API, и на 4-ядерном сервере LiveKit в пике забирает 1.6 ядра.
#
# `sh -c` нужен ради подстановки переменной (exec-форма её не делает),
# `exec` — чтобы uvicorn получил PID 1 и корректно принимал SIGTERM.
CMD ["sh", "-c", \
"exec uv run uvicorn main:create_app --factory --host 0.0.0.0 --port 8000 --workers ${UVICORN_WORKERS:-2}"]

View File

@@ -3,7 +3,7 @@
import logging import logging
from typing import Annotated, Any, cast 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 livekit import api
from sqlalchemy import CursorResult from sqlalchemy import CursorResult
from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.dialects.postgresql import insert as pg_insert
@@ -22,6 +22,7 @@ router = APIRouter(prefix="/api/v1/livekit", tags=["livekit"])
@router.post("/webhook") @router.post("/webhook")
async def receive_webhook( async def receive_webhook(
request: Request, request: Request,
background: BackgroundTasks,
session: Annotated[AsyncSession, Depends(get_session)], session: Annotated[AsyncSession, Depends(get_session)],
authorization: Annotated[str | None, Header()] = None, authorization: Annotated[str | None, Header()] = None,
) -> dict[str, str]: ) -> dict[str, str]:
@@ -30,6 +31,11 @@ async def receive_webhook(
Дедупликация по `event.id`: `INSERT ... ON CONFLICT DO NOTHING` в Дедупликация по `event.id`: `INSERT ... ON CONFLICT DO NOTHING` в
`livekit_webhook_events` в одной транзакции с эффектами обработчика — `livekit_webhook_events` в одной транзакции с эффектами обработчика —
при конфликте (дубль) эффекты пропускаются, но ответ всё равно 200. при конфликте (дубль) эффекты пропускаются, но ответ всё равно 200.
Всё, что требует сети (запуск Track Egress), уходит в `background` и
выполняется уже после ответа: обработчик держит соединение с БД и
открытую транзакцию, а LiveKit при медленном ответе копит очередь
доставки и в итоге дропает события (см. `services.egress.run_track_egress`).
""" """
settings = get_settings() settings = get_settings()
raw_body = await request.body() raw_body = await request.body()
@@ -57,7 +63,7 @@ async def receive_webhook(
await session.commit() await session.commit()
return {"status": "duplicate"} return {"status": "duplicate"}
dispatcher = WebhookDispatcher(session) dispatcher = WebhookDispatcher(session, schedule=background.add_task)
await dispatcher.dispatch(event) await dispatcher.dispatch(event)
await session.commit() await session.commit()
return {"status": "ok"} return {"status": "ok"}

View File

@@ -20,6 +20,24 @@ class Settings(BaseSettings):
redis_url: str = "redis://localhost:6379/0" redis_url: str = "redis://localhost:6379/0"
plugins_config_path: str = "../config/plugins.yaml" plugins_config_path: str = "../config/plugins.yaml"
# --- Пул соединений с БД ---
# Дефолт SQLAlchemy (5 + 10) на нагрузочном тесте 28.07.2026 выгребался
# за секунды: 226 ошибок `QueuePool limit of size 5 overflow 10 reached`
# и 37 ответов 500 на путях входа в конференцию.
#
# ⚠️ Бюджет соединений считается на ВЕСЬ инстанс, а не на процесс: каждый
# воркер uvicorn (`UVICORN_WORKERS`) держит собственный пул, плюс
# соединения нужны Celery-воркерам и alembic при миграциях. При
# `max_connections=100` у Postgres и двух воркерах 2 × (10 + 10) = 40
# оставляет запас. Поднимая значения на более крупном сервере, поднимайте
# и `max_connections` — иначе вместо понятной ошибки приложения получите
# отказ Postgres, который диагностируется куда хуже.
db_pool_size: int = 10
db_max_overflow: int = 10
# 10 секунд вместо дефолтных 30 — сознательно: пусть запрос падает быстро
# и показывает проблему, а не висит полминуты, делая вид, что всё живо.
db_pool_timeout: int = 10
# --- Версия инстанса (релиз v0.0.1) --- # --- Версия инстанса (релиз v0.0.1) ---
# install.sh копирует значение из файла `VERSION` (корень репозитория) в # install.sh копирует значение из файла `VERSION` (корень репозитория) в
# `.env` при каждой установке/обновлении — здесь только чтение готового # `.env` при каждой установке/обновлении — здесь только чтение готового
@@ -59,6 +77,13 @@ class Settings(BaseSettings):
# (см. `deploy/docker-compose.yml`); в тестах переопределяется на `tmp_path`. # (см. `deploy/docker-compose.yml`); в тестах переопределяется на `tmp_path`.
recordings_dir: str = "/recordings" recordings_dir: str = "/recordings"
# Таймаут запуска Track Egress. Когда egress-сервиса в деплое нет (профиль
# `transcribe` не поднят), LiveKit ждёт ответа воркера через Redis до
# собственного таймаута psrpc — на тесте 28.07.2026 это давало по 2025
# секунд на каждый вызов. Ждать столько бессмысленно: если egress жив, он
# отвечает за доли секунды.
egress_start_timeout_s: float = 3.0
# --- Email (SMTP-бэкенд) --- # --- Email (SMTP-бэкенд) ---
# `console` — дефолт для dev (письмо только логируется); `smtp` — реальная # `console` — дефолт для dev (письмо только логируется); `smtp` — реальная
# отправка через aiosmtplib. Секреты SMTP — только в `.env` (инвариант №6), # отправка через aiosmtplib. Секреты SMTP — только в `.env` (инвариант №6),

View File

@@ -13,7 +13,16 @@ from core.config import get_settings
settings = get_settings() settings = get_settings()
engine: AsyncEngine = create_async_engine(settings.database_url, pool_pre_ping=True) engine: AsyncEngine = create_async_engine(
settings.database_url,
pool_pre_ping=True,
# Параметры пула — в настройках (`core/config.py`, там же расчёт бюджета
# соединений на инстанс). Дефолт SQLAlchemy 5 + 10 под нагрузкой
# выгребался за секунды.
pool_size=settings.db_pool_size,
max_overflow=settings.db_max_overflow,
pool_timeout=settings.db_pool_timeout,
)
async_session_maker = async_sessionmaker(engine, expire_on_commit=False) async_session_maker = async_sessionmaker(engine, expire_on_commit=False)

View File

@@ -5,14 +5,24 @@
транскодирования в `.ogg` на общий volume `recordings_dir`. Финализация транскодирования в `.ogg` на общий volume `recordings_dir`. Финализация
результата (итоговый `location`/ошибка) приходит асинхронно через webhook результата (итоговый `location`/ошибка) приходит асинхронно через webhook
`egress_ended` — здесь только сам запуск и `egress_id`/`started_at` из ответа. `egress_ended` — здесь только сам запуск и `egress_id`/`started_at` из ответа.
Запуск вынесен из тела webhook-обработчика в фоновую задачу
(`run_track_egress`) — см. докстринг этой функции.
""" """
import asyncio
import logging
import uuid
from dataclasses import dataclass from dataclasses import dataclass
from datetime import UTC, datetime from datetime import UTC, datetime
from livekit import api from livekit import api
from core.config import get_settings 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) @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, api_secret=settings.livekit_api_secret,
) )
try: try:
info = await lkapi.egress.start_track_egress( # Без таймаута вызов висит до собственного таймаута psrpc LiveKit
api.TrackEgressRequest( # (2025 с, когда egress-воркера в деплое нет). `TimeoutError`
room_name=room_name, # ловит вызывающая сторона наравне с прочими ошибками запуска.
track_id=track_sid, async with asyncio.timeout(settings.egress_start_timeout_s):
file=api.DirectFileOutput(filepath=filepath), info = await lkapi.egress.start_track_egress(
api.TrackEgressRequest(
room_name=room_name,
track_id=track_sid,
file=api.DirectFileOutput(filepath=filepath),
)
) )
)
finally: finally:
await lkapi.aclose() 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) else datetime.now(UTC)
) )
return EgressStartResult(egress_id=info.egress_id, started_at=started_at) 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,
)

View File

@@ -13,6 +13,7 @@
import logging import logging
import uuid import uuid
from collections.abc import Callable
from datetime import UTC, datetime from datetime import UTC, datetime
from livekit.protocol.egress import EgressStatus from livekit.protocol.egress import EgressStatus
@@ -26,7 +27,7 @@ from repositories.conferences import (
ConferenceRepository, ConferenceRepository,
ConferenceSessionRepository, ConferenceSessionRepository,
) )
from services.egress import start_track_egress from services.egress import run_track_egress
from services.instance_settings import InstanceSettingsService from services.instance_settings import InstanceSettingsService
from services.pipeline_producer import enqueue_pipeline from services.pipeline_producer import enqueue_pipeline
@@ -53,13 +54,27 @@ def _egress_ns_to_datetime(nanoseconds: int) -> datetime | None:
class WebhookDispatcher: 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._conferences = ConferenceRepository(session)
self._sessions = ConferenceSessionRepository(session) self._sessions = ConferenceSessionRepository(session)
self._audio_tracks = AudioTrackRepository(session) self._audio_tracks = AudioTrackRepository(session)
self._instance_settings = InstanceSettingsService(session) self._instance_settings = InstanceSettingsService(session)
self._schedule = schedule
async def dispatch(self, event: WebhookEvent) -> None: async def dispatch(self, event: WebhookEvent) -> None:
"""Обработать одно webhook-событие; неизвестный тип события — no-op.""" """Обработать одно webhook-событие; неизвестный тип события — no-op."""
@@ -161,15 +176,32 @@ class WebhookDispatcher:
) )
async def _on_track_published(self, event: WebhookEvent) -> None: async def _on_track_published(self, event: WebhookEvent) -> None:
"""Запустить Track Egress для опубликованного аудиотрека микрофона (ADR-002). """Запланировать Track Egress для опубликованного аудиотрека микрофона (ADR-002).
Видео/скриншеринг и т.п. — no-op (диаризация не нужна: транскрибируем Видео/скриншеринг и т.п. — no-op (диаризация не нужна: транскрибируем
только речь, трек = спикер). Идемпотентно: если строка трека уже только речь, трек = спикер). Идемпотентно: если строка трека уже
существует (гонка повторной доставки), egress повторно не запускается. существует (гонка повторной доставки), egress повторно не запускается.
Сам запуск уходит в фоновую задачу (`services.egress.run_track_egress`):
здесь остаются только быстрые проверки по БД, потому что обработчик
выполняется внутри открытой транзакции вебхука. Обоснование с цифрами —
в докстринге `run_track_egress`.
""" """
if event.track.type != TrackType.AUDIO or event.track.source != TrackSource.MICROPHONE: if event.track.type != TrackType.AUDIO or event.track.source != TrackSource.MICROPHONE:
return 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) conference = await self._conferences.get_by_slug(event.room.name)
if conference is None: if conference is None:
logger.warning( logger.warning(
@@ -216,37 +248,13 @@ class WebhookDispatcher:
filepath = ( filepath = (
f"{settings.recordings_dir}/{session_record.id}/{participant.id}_{event.track.sid}.ogg" f"{settings.recordings_dir}/{session_record.id}/{participant.id}_{event.track.sid}.ogg"
) )
try: self._schedule(
result = await start_track_egress(event.room.name, event.track.sid, filepath) run_track_egress,
except Exception as exc: # noqa: BLE001 — недоступность egress не должна ронять webhook room_name=event.room.name,
# Деплой-профиль (блок D): egress — необязательный сервис профиля track_sid=event.track.sid,
# `transcribe`; без него запись просто не стартует для этого трека filepath=filepath,
# (риск «Потерян 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(
session_id=session_record.id, session_id=session_record.id,
participant_id=participant.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: async def _on_egress_ended(self, event: WebhookEvent) -> None:

View File

@@ -3,9 +3,13 @@
цикл закреплённой/незакреплённой конференции — на фикстурах payload'ов LiveKit. цикл закреплённой/незакреплённой конференции — на фикстурах payload'ов LiveKit.
""" """
import asyncio
import base64 import base64
import hashlib import hashlib
import json
import uuid import uuid
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from datetime import UTC, datetime, timedelta from datetime import UTC, datetime, timedelta
from pathlib import Path from pathlib import Path
from unittest.mock import AsyncMock, Mock from unittest.mock import AsyncMock, Mock
@@ -13,10 +17,13 @@ from unittest.mock import AsyncMock, Mock
import httpx import httpx
import jwt import jwt
import pytest import pytest
from google.protobuf.json_format import ParseDict
from livekit.protocol.webhook import WebhookEvent
from sqlalchemy import select from sqlalchemy import select
from sqlalchemy.dialects.postgresql import insert as pg_insert from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
import services.egress as egress_module
import services.webhook_handlers as webhook_handlers_module import services.webhook_handlers as webhook_handlers_module
from core.config import get_settings from core.config import get_settings
from core.security import hash_password 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") 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: async def _post_webhook(client: httpx.AsyncClient, body: bytes) -> httpx.Response:
return await client.post( return await client.post(
WEBHOOK_URL, WEBHOOK_URL,
@@ -319,8 +372,10 @@ async def test_track_published_by_guest_starts_egress_and_creates_track_row(
mock_start = AsyncMock( mock_start = AsyncMock(
return_value=EgressStartResult(egress_id="EG_guest_track", started_at=started_at) 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()) conference = await _make_conference(db_session, generate_slug())
guest = await _make_guest(db_session, conference) guest = await _make_guest(db_session, conference)
await db_session.commit() await db_session.commit()
@@ -382,8 +437,10 @@ async def test_track_published_survives_egress_unavailable(
) -> None: ) -> None:
"""Недоступность egress не должна ронять webhook (блок D): 200 + warning, без строки трека.""" """Недоступность egress не должна ронять webhook (блок D): 200 + warning, без строки трека."""
mock_start = AsyncMock(side_effect=RuntimeError("egress service unavailable")) 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()) conference = await _make_conference(db_session, generate_slug())
user = await _make_user(db_session, "webhook-track-egress-down@example.com") user = await _make_user(db_session, "webhook-track-egress-down@example.com")
await db_session.commit() await db_session.commit()
@@ -417,13 +474,169 @@ async def test_track_published_survives_egress_unavailable(
assert tracks == [] 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 —
2025 секунд внутри открытой транзакции вебхука. На нагрузочном тесте
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-транспорт дожидается фоновых
задач), а сама суть: пока открыта транзакция вебхука, сетевого вызова не
происходит — обработчик только планирует задачу. Именно это разгружает пул
соединений: до правки вызов жил внутри транзакции и держал соединение
2025 секунд, когда 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 это было 2025 секунд на каждый микрофонный
трек. Живой 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( async def test_track_published_video_is_noop(
client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None: ) -> None:
"""№9 плана (часть 1): video-трек — no-op, egress не запускается.""" """№9 плана (часть 1): video-трек — no-op, egress не запускается."""
mock_start = AsyncMock() 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()) conference = await _make_conference(db_session, generate_slug())
user = await _make_user(db_session, "webhook-track-video@example.com") user = await _make_user(db_session, "webhook-track-video@example.com")
await db_session.commit() await db_session.commit()
@@ -462,8 +675,10 @@ async def test_track_published_repeated_webhook_creates_single_row(
mock_start = AsyncMock( mock_start = AsyncMock(
return_value=EgressStartResult(egress_id="EG_repeat", started_at=datetime.now(UTC)) 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()) conference = await _make_conference(db_session, generate_slug())
user = await _make_user(db_session, "webhook-track-repeat@example.com") user = await _make_user(db_session, "webhook-track-repeat@example.com")
await db_session.commit() 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 client: httpx.AsyncClient, db_session: AsyncSession, monkeypatch: pytest.MonkeyPatch
) -> None: ) -> None:
"""№10 плана (часть 1): `egress_ended` — 'recorded' на успехе, 'failed' на ошибке.""" """№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()) conference = await _make_conference(db_session, generate_slug())
user = await _make_user(db_session, "webhook-egress-ended@example.com") user = await _make_user(db_session, "webhook-egress-ended@example.com")
await db_session.commit() await db_session.commit()
@@ -527,7 +744,7 @@ async def test_egress_ended_finalizes_track_success_and_failure(
# Успешная запись. # Успешная запись.
ok_started = datetime.now(UTC) ok_started = datetime.now(UTC)
monkeypatch.setattr( monkeypatch.setattr(
webhook_handlers_module, egress_module,
"start_track_egress", "start_track_egress",
AsyncMock(return_value=EgressStartResult(egress_id="EG_ok", started_at=ok_started)), 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( monkeypatch.setattr(
webhook_handlers_module, egress_module,
"start_track_egress", "start_track_egress",
AsyncMock( AsyncMock(
return_value=EgressStartResult(egress_id="EG_fail", started_at=datetime.now(UTC)) 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() mock_enqueue = Mock()
monkeypatch.setattr(webhook_handlers_module, "enqueue_pipeline", mock_enqueue) 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()) conference = await _make_conference(db_session, generate_slug())
await db_session.commit() await db_session.commit()

View File

@@ -82,7 +82,12 @@ services:
MEDIA_ROOT: ${MEDIA_ROOT:-/app/media} MEDIA_ROOT: ${MEDIA_ROOT:-/app/media}
# Версия инстанса (релиз v0.0.1) — install.sh копирует значение # Версия инстанса (релиз v0.0.1) — install.sh копирует значение
# из файла VERSION (корень репозитория) в .env; отдаётся в GET /api/health. # из файла VERSION (корень репозитория) в .env; отдаётся в GET /api/health.
VIDCONF_VERSION: ${VIDCONF_VERSION:-0.0.11} VIDCONF_VERSION: ${VIDCONF_VERSION:-0.0.12}
# Число процессов uvicorn (см. backend/Dockerfile). Дефолт 2 рассчитан
# на 4-ядерный сервер, где ядра делятся с LiveKit. Поднимая значение,
# проверьте бюджет соединений с БД: каждый воркер держит свой пул
# (DB_POOL_SIZE + DB_MAX_OVERFLOW), а у Postgres есть max_connections.
UVICORN_WORKERS: ${UVICORN_WORKERS:-2}
# config/ лежит в корне репозитория и не попадает в образ (контекст сборки — # config/ лежит в корне репозитория и не попадает в образ (контекст сборки —
# только backend/), поэтому plugins.yaml монтируется отдельно. # только backend/), поэтому plugins.yaml монтируется отдельно.
volumes: volumes: