Files
vidconf/backend/api/livekit_webhook.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

70 lines
3.0 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.
"""Приёмник webhook-событий LiveKit (без JWT — верификация подписью LiveKit)."""
import logging
from typing import Annotated, Any, cast
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
from sqlalchemy.ext.asyncio import AsyncSession
from core.config import get_settings
from core.db import get_session
from models.webhook_event import LivekitWebhookEvent
from services.webhook_handlers import WebhookDispatcher
logger = logging.getLogger(__name__)
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]:
"""Принять, верифицировать и обработать webhook-событие LiveKit.
Дедупликация по `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()
receiver = api.WebhookReceiver(
api.TokenVerifier(settings.livekit_api_key, settings.livekit_api_secret)
)
try:
event = receiver.receive(raw_body.decode(), authorization or "")
except Exception as exc: # noqa: BLE001 — SDK кидает generic Exception на невалидную подпись
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED, detail="invalid_signature"
) from exc
insert_result = cast(
CursorResult[Any],
await session.execute(
pg_insert(LivekitWebhookEvent)
.values(event_id=event.id, event_type=event.event)
.on_conflict_do_nothing(index_elements=["event_id"])
),
)
if insert_result.rowcount == 0:
# Дубль уже обработанного события — пропускаем эффекты, но отвечаем 200.
await session.commit()
return {"status": "duplicate"}
dispatcher = WebhookDispatcher(session, schedule=background.add_task)
await dispatcher.dispatch(event)
await session.commit()
return {"status": "ok"}