544 lines
26 KiB
Python
544 lines
26 KiB
Python
"""Репозитории доступа к `conferences` (сущность) и `conference_sessions` (сеанс, ADR-001).
|
||
|
||
`ConferenceSessionRepository` используется webhook-обработчиками LiveKit
|
||
(`services/webhook_handlers.py`) и beat-задачей обслуживания
|
||
(`workers/tasks/maintenance.py`) — методы get-or-create/guard-стиля, чтобы
|
||
быть безопасными при пропущенных или дублирующихся событиях (шаги
|
||
пост-обработки идемпотентны).
|
||
"""
|
||
|
||
import uuid
|
||
from datetime import datetime
|
||
from typing import cast
|
||
|
||
from sqlalchemy import delete, func, or_, select
|
||
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from models.audio_track import SessionAudioTrack
|
||
from models.conference import Conference
|
||
from models.invitee import ConferenceInvitee
|
||
from models.participant import ConferenceParticipant
|
||
from models.session import ConferenceSession
|
||
from models.user import User
|
||
|
||
|
||
class ConferenceRepository:
|
||
"""Инкапсулирует SQL-запросы к конференциям (`conferences`)."""
|
||
|
||
def __init__(self, session: AsyncSession) -> None:
|
||
self._session = session
|
||
|
||
async def get_by_id(self, conference_id: uuid.UUID) -> Conference | None:
|
||
"""Найти конференцию по id."""
|
||
return await self._session.get(Conference, conference_id)
|
||
|
||
async def get_by_slug(self, slug: str) -> Conference | None:
|
||
"""Найти конференцию по постоянной ссылке (= имени LiveKit-комнаты)."""
|
||
result = await self._session.execute(select(Conference).where(Conference.slug == slug))
|
||
return result.scalar_one_or_none()
|
||
|
||
async def get_by_number(self, number: str) -> Conference | None:
|
||
"""Найти конференцию по человеко-диктуемому номеру."""
|
||
result = await self._session.execute(select(Conference).where(Conference.number == number))
|
||
return result.scalar_one_or_none()
|
||
|
||
async def get_status_by_id(self, conference_id: uuid.UUID) -> str | None:
|
||
"""Прочитать АКТУАЛЬНЫЙ статус конференции, минуя identity map сессии.
|
||
|
||
В отличие от `get_by_id`/`session.get(...)`, SELECT одной колонки не
|
||
возвращает уже загруженный в эту сессию ORM-объект `Conference` из
|
||
кэша identity map — а при `expire_on_commit=False` (`core/db.py`)
|
||
такой объект, однажды загруженный долгоживущей WS-сессией чата,
|
||
никогда сам не увидит статус, изменённый вебхуком в ДРУГОЙ сессии/
|
||
процессе (например, `room_finished` -> `ended`). Используется
|
||
`ChatService.persist_and_publish` перед созданием новой сессии
|
||
пайплайна.
|
||
"""
|
||
result = await self._session.execute(
|
||
select(Conference.status).where(Conference.id == conference_id)
|
||
)
|
||
return result.scalar_one_or_none()
|
||
|
||
async def add(self, conference: Conference) -> Conference:
|
||
"""Добавить конференцию в сессию и сделать flush (нарушение unique — здесь же)."""
|
||
self._session.add(conference)
|
||
await self._session.flush()
|
||
return conference
|
||
|
||
async def delete(self, conference: Conference) -> None:
|
||
"""Удалить конференцию (сеансы/гости удаляются каскадом на уровне БД)."""
|
||
await self._session.delete(conference)
|
||
await self._session.flush()
|
||
|
||
async def list_owned(
|
||
self, user_id: uuid.UUID, *, email: str, now: datetime
|
||
) -> list[Conference]:
|
||
"""Конференции владельца ИЛИ приглашённого для «Моих конференций» (решение от
|
||
2026-07-20 поверх ADR-003: приглашённый видит конференцию в своих списках) —
|
||
закреплённые + предстоящие разовые.
|
||
|
||
«Приглашённый» — есть строка `conference_invitees` с `user_id == user_id`
|
||
ИЛИ с `lower(email) == lower(этого email)` (внешнее приглашение на адрес,
|
||
под которым человек впоследствии зарегистрировался). `EXISTS`-подзапрос
|
||
не размножает строки `Conference` — `DISTINCT` не требуется.
|
||
"""
|
||
is_invitee = (
|
||
select(ConferenceInvitee.id)
|
||
.where(
|
||
ConferenceInvitee.conference_id == Conference.id,
|
||
or_(
|
||
ConferenceInvitee.user_id == user_id,
|
||
func.lower(ConferenceInvitee.email) == email.lower(),
|
||
),
|
||
)
|
||
.exists()
|
||
)
|
||
result = await self._session.execute(
|
||
select(Conference)
|
||
.where(
|
||
or_(Conference.owner_id == user_id, is_invitee),
|
||
(Conference.is_pinned.is_(True))
|
||
| (
|
||
(Conference.scheduled_at.isnot(None))
|
||
& (Conference.scheduled_at >= now)
|
||
& (Conference.status == "scheduled")
|
||
),
|
||
)
|
||
.order_by(Conference.is_pinned.desc(), Conference.scheduled_at.asc().nulls_last())
|
||
)
|
||
return list(result.scalars().all())
|
||
|
||
async def list_calendar_candidates(self, user_id: uuid.UUID, *, email: str) -> list[Conference]:
|
||
"""Конференции владельца ИЛИ приглашённого, потенциально дающие вхождения в календаре.
|
||
|
||
Тот же принцип видимости приглашённого, что и `list_owned` (решение от
|
||
2026-07-20). Развёртка диапазона — в сервисном слое
|
||
(`services/conferences.py::list_calendar`), здесь только грубая выборка
|
||
кандидатов (есть recurrence или указано scheduled_at).
|
||
"""
|
||
is_invitee = (
|
||
select(ConferenceInvitee.id)
|
||
.where(
|
||
ConferenceInvitee.conference_id == Conference.id,
|
||
or_(
|
||
ConferenceInvitee.user_id == user_id,
|
||
func.lower(ConferenceInvitee.email) == email.lower(),
|
||
),
|
||
)
|
||
.exists()
|
||
)
|
||
result = await self._session.execute(
|
||
select(Conference).where(
|
||
or_(Conference.owner_id == user_id, is_invitee),
|
||
(Conference.recurrence.isnot(None)) | (Conference.scheduled_at.isnot(None)),
|
||
)
|
||
)
|
||
return list(result.scalars().all())
|
||
|
||
async def list_expired_unpinned_scheduled(self, *, now: datetime) -> list[Conference]:
|
||
"""Незакреплённые плановые конференции без единого сеанса (кандидаты на `ended`).
|
||
|
||
Точная проверка истечения времени (с учётом `duration_minutes` и
|
||
запаса) — в `workers/tasks/maintenance.py`; здесь — грубая выборка по
|
||
`NOT EXISTS` сеанса, чтобы не тянуть в память всё лишнее.
|
||
"""
|
||
has_session = (
|
||
select(ConferenceSession.id)
|
||
.where(ConferenceSession.conference_id == Conference.id)
|
||
.exists()
|
||
)
|
||
result = await self._session.execute(
|
||
select(Conference).where(
|
||
Conference.is_pinned.is_(False),
|
||
Conference.status == "scheduled",
|
||
Conference.scheduled_at.isnot(None),
|
||
Conference.scheduled_at < now,
|
||
~has_session,
|
||
)
|
||
)
|
||
return list(result.scalars().all())
|
||
|
||
|
||
class ConferenceInviteeRepository:
|
||
"""Инкапсулирует SQL-запросы к приглашённым на конференцию (`conference_invitees`, ADR-003)."""
|
||
|
||
def __init__(self, session: AsyncSession) -> None:
|
||
self._session = session
|
||
|
||
async def list_with_user(
|
||
self, conference_id: uuid.UUID
|
||
) -> list[tuple[ConferenceInvitee, str | None, str | None, str | None]]:
|
||
"""Приглашённые конференции + имя/email/путь аватара их пользователя (LEFT JOIN).
|
||
|
||
Для внешних приглашённых (`user_id IS NULL`) три последних элемента
|
||
кортежа — `None` (нет привязанного `User`).
|
||
"""
|
||
result = await self._session.execute(
|
||
select(ConferenceInvitee, User.name_user, User.email, User.avatar_path)
|
||
.outerjoin(User, ConferenceInvitee.user_id == User.id)
|
||
.where(ConferenceInvitee.conference_id == conference_id)
|
||
)
|
||
# `User.name_user`/`User.email` типизированы как non-optional (NOT NULL
|
||
# в модели) — но при LEFT JOIN без совпадения (внешний приглашённый,
|
||
# `user_id IS NULL`) значения реально приходят `NULL`; mypy не видит
|
||
# nullability, вносимую `outerjoin`, отсюда явный `cast`.
|
||
return cast(
|
||
"list[tuple[ConferenceInvitee, str | None, str | None, str | None]]",
|
||
list(result.all()),
|
||
)
|
||
|
||
async def existing_user_ids(self, user_ids: set[uuid.UUID]) -> set[uuid.UUID]:
|
||
"""Подмножество `user_ids`, реально существующее в `users` (проверка перед вставкой)."""
|
||
if not user_ids:
|
||
return set()
|
||
result = await self._session.execute(select(User.id).where(User.id.in_(user_ids)))
|
||
return set(result.scalars().all())
|
||
|
||
async def exists_for_user(
|
||
self, conference_id: uuid.UUID, *, user_id: uuid.UUID, email: str
|
||
) -> bool:
|
||
"""Приглашён ли `user_id` на конференцию — по `user_id` или по `lower(email)`.
|
||
|
||
Используется проверкой доступа к детальной карточке (`GET /conferences/{id}`,
|
||
решение от 2026-07-20 поверх ADR-003) — приглашённый должен её видеть
|
||
наравне с владельцем/администратором.
|
||
"""
|
||
result = await self._session.execute(
|
||
select(ConferenceInvitee.id)
|
||
.where(
|
||
ConferenceInvitee.conference_id == conference_id,
|
||
or_(
|
||
ConferenceInvitee.user_id == user_id,
|
||
func.lower(ConferenceInvitee.email) == email.lower(),
|
||
),
|
||
)
|
||
.limit(1)
|
||
)
|
||
return result.scalar_one_or_none() is not None
|
||
|
||
async def replace_all(
|
||
self, conference_id: uuid.UUID, invitees: list[ConferenceInvitee]
|
||
) -> None:
|
||
"""Полностью заменить состав приглашённых конференции (ADR-003, п.3 — PUT-семантика)."""
|
||
await self._session.execute(
|
||
delete(ConferenceInvitee).where(ConferenceInvitee.conference_id == conference_id)
|
||
)
|
||
for invitee in invitees:
|
||
self._session.add(invitee)
|
||
await self._session.flush()
|
||
|
||
|
||
class ConferenceSessionRepository:
|
||
"""Инкапсулирует SQL-запросы к сеансам конференций и их участникам."""
|
||
|
||
def __init__(self, session: AsyncSession) -> None:
|
||
self._session = session
|
||
|
||
async def get_open_by_conference(self, conference_id: uuid.UUID) -> ConferenceSession | None:
|
||
"""Вернуть открытый (`t_end IS NULL`) сеанс конференции, если есть."""
|
||
result = await self._session.execute(
|
||
select(ConferenceSession)
|
||
.where(
|
||
ConferenceSession.conference_id == conference_id,
|
||
ConferenceSession.t_end.is_(None),
|
||
)
|
||
.order_by(ConferenceSession.t_start.desc())
|
||
.limit(1)
|
||
)
|
||
return result.scalar_one_or_none()
|
||
|
||
async def create(
|
||
self, *, conference_id: uuid.UUID, title: str | None, t_start: datetime
|
||
) -> ConferenceSession:
|
||
"""Создать новый (открытый) сеанс конференции."""
|
||
record = ConferenceSession(conference_id=conference_id, title=title, t_start=t_start)
|
||
self._session.add(record)
|
||
await self._session.flush()
|
||
return record
|
||
|
||
async def get_or_create_open(
|
||
self, *, conference_id: uuid.UUID, title: str | None, t_start: datetime
|
||
) -> ConferenceSession:
|
||
"""Get-or-create открытого сеанса конференции.
|
||
|
||
Гвард на случай, если `room_started` было пропущено и первым пришло
|
||
`participant_joined`.
|
||
"""
|
||
existing = await self.get_open_by_conference(conference_id)
|
||
if existing is not None:
|
||
return existing
|
||
return await self.create(conference_id=conference_id, title=title, t_start=t_start)
|
||
|
||
async def close(self, session_record: ConferenceSession, *, t_end: datetime) -> None:
|
||
"""Закрыть сеанс, проставив `t_end`."""
|
||
session_record.t_end = t_end
|
||
|
||
async def list_open(self) -> list[ConferenceSession]:
|
||
"""Список всех открытых (`t_end IS NULL`) сеансов.
|
||
|
||
Используется maintenance-задачей (`workers/tasks/maintenance.py`)
|
||
для обхода всех "зависших" сеансов разом.
|
||
"""
|
||
result = await self._session.execute(
|
||
select(ConferenceSession).where(ConferenceSession.t_end.is_(None))
|
||
)
|
||
return list(result.scalars().all())
|
||
|
||
async def list_stuck_summarizing(self, *, older_than: datetime) -> list[ConferenceSession]:
|
||
"""Сеансы, зависшие на шаге суммаризации: `pipeline_status='summarizing'`,
|
||
`summary_data` ещё не заполнен, а сеанс завершился раньше `older_than`.
|
||
|
||
Кандидаты на повторную постановку `summarize_session` —
|
||
`workers/tasks/maintenance.py::recover_stuck_summaries` (уровень 2
|
||
защиты от потери постановки задачи при сбое брокера в `run_pipeline`).
|
||
"""
|
||
result = await self._session.execute(
|
||
select(ConferenceSession).where(
|
||
ConferenceSession.pipeline_status == "summarizing",
|
||
ConferenceSession.summary_data.is_(None),
|
||
ConferenceSession.t_end.isnot(None),
|
||
ConferenceSession.t_end < older_than,
|
||
)
|
||
)
|
||
return list(result.scalars().all())
|
||
|
||
async def list_stuck_notifying(self, *, older_than: datetime) -> list[ConferenceSession]:
|
||
"""Сеансы, зависшие на шаге уведомления: саммари готово, но `notify_session`
|
||
так и не перевела пайплайн в `notified`.
|
||
|
||
`pipeline_status='summarizing'` + `summary_data IS NOT NULL` +
|
||
`t_end < older_than` — кандидаты на повторную постановку
|
||
`notify_session` (`workers/tasks/maintenance.py::recover_stuck_notifications`,
|
||
уровень 2 защиты от потери постановки задачи при сбое брокера в
|
||
`summarize_session`, аналог `list_stuck_summarizing`).
|
||
"""
|
||
result = await self._session.execute(
|
||
select(ConferenceSession).where(
|
||
ConferenceSession.pipeline_status == "summarizing",
|
||
ConferenceSession.summary_data.isnot(None),
|
||
ConferenceSession.t_end.isnot(None),
|
||
ConferenceSession.t_end < older_than,
|
||
)
|
||
)
|
||
return list(result.scalars().all())
|
||
|
||
async def count_by_pipeline_status(self) -> dict[str, int]:
|
||
"""Число сеансов в каждом статусе пайплайна (`pipeline_status`).
|
||
|
||
Используется метриками Prometheus (`backend/api/metrics.py`,
|
||
блок 3) для gauge `vidconf_pipeline_sessions{status=...}` — считается
|
||
заново при каждом scrape, не кешируется. Статусы без единого сеанса в
|
||
результат не попадают (пустая группа), это ожидаемо: вызывающая
|
||
сторона сама проставляет 0 для отсутствующих в словаре статусов
|
||
(полный перечень — `models.session.PipelineStatus`), чтобы метрика не
|
||
"пропадала" из `/metrics` между сборами.
|
||
"""
|
||
result = await self._session.execute(
|
||
select(ConferenceSession.pipeline_status, func.count()).group_by(
|
||
ConferenceSession.pipeline_status
|
||
)
|
||
)
|
||
return {status: count for status, count in result.all()}
|
||
|
||
async def has_active_participants(self, session_id: uuid.UUID) -> bool:
|
||
"""Есть ли у сеанса хотя бы один участник без `left_at` (кто-то ещё внутри)."""
|
||
result = await self._session.execute(
|
||
select(ConferenceParticipant.id)
|
||
.where(
|
||
ConferenceParticipant.session_id == session_id,
|
||
ConferenceParticipant.left_at.is_(None),
|
||
)
|
||
.limit(1)
|
||
)
|
||
return result.scalar_one_or_none() is not None
|
||
|
||
async def add_participant(
|
||
self,
|
||
*,
|
||
session_id: uuid.UUID,
|
||
user_id: uuid.UUID | None,
|
||
guest_id: uuid.UUID | None,
|
||
joined_at: datetime,
|
||
) -> ConferenceParticipant:
|
||
"""Get-or-create активной (не покинувшей) записи участия — пользователя ИЛИ гостя.
|
||
|
||
Защищает от дублей при повторной доставке `participant_joined` для
|
||
уже присутствующего участника. Ровно один из `user_id`/`guest_id`
|
||
должен быть передан (см. CHECK-constraint модели).
|
||
"""
|
||
stmt = select(ConferenceParticipant).where(
|
||
ConferenceParticipant.session_id == session_id,
|
||
ConferenceParticipant.left_at.is_(None),
|
||
)
|
||
stmt = stmt.where(
|
||
ConferenceParticipant.user_id == user_id
|
||
if user_id is not None
|
||
else ConferenceParticipant.guest_id == guest_id
|
||
)
|
||
existing = (await self._session.execute(stmt)).scalar_one_or_none()
|
||
if existing is not None:
|
||
return existing
|
||
|
||
participant = ConferenceParticipant(
|
||
session_id=session_id, user_id=user_id, guest_id=guest_id, joined_at=joined_at
|
||
)
|
||
self._session.add(participant)
|
||
await self._session.flush()
|
||
return participant
|
||
|
||
async def get_active_participant(
|
||
self,
|
||
session_id: uuid.UUID,
|
||
*,
|
||
user_id: uuid.UUID | None = None,
|
||
guest_id: uuid.UUID | None = None,
|
||
) -> ConferenceParticipant | None:
|
||
"""Найти активную (`left_at IS NULL`) запись присутствия пользователя/гостя в сеансе.
|
||
|
||
Используется атрибуцией аудиотрека к участнику (`track_published`,
|
||
ADR-002): ровно один из `user_id`/`guest_id` должен быть передан.
|
||
"""
|
||
stmt = select(ConferenceParticipant).where(
|
||
ConferenceParticipant.session_id == session_id,
|
||
ConferenceParticipant.left_at.is_(None),
|
||
)
|
||
stmt = stmt.where(
|
||
ConferenceParticipant.user_id == user_id
|
||
if user_id is not None
|
||
else ConferenceParticipant.guest_id == guest_id
|
||
)
|
||
stmt = stmt.order_by(ConferenceParticipant.joined_at.desc()).limit(1)
|
||
return (await self._session.execute(stmt)).scalar_one_or_none()
|
||
|
||
async def mark_participant_left(
|
||
self,
|
||
*,
|
||
session_id: uuid.UUID,
|
||
user_id: uuid.UUID | None,
|
||
guest_id: uuid.UUID | None,
|
||
left_at: datetime,
|
||
) -> None:
|
||
"""Проставить `left_at` последней открытой записи участника (пользователя или гостя)."""
|
||
stmt = select(ConferenceParticipant).where(
|
||
ConferenceParticipant.session_id == session_id,
|
||
ConferenceParticipant.left_at.is_(None),
|
||
)
|
||
stmt = stmt.where(
|
||
ConferenceParticipant.user_id == user_id
|
||
if user_id is not None
|
||
else ConferenceParticipant.guest_id == guest_id
|
||
)
|
||
stmt = stmt.order_by(ConferenceParticipant.joined_at.desc()).limit(1)
|
||
result = await self._session.execute(stmt)
|
||
participant = result.scalar_one_or_none()
|
||
if participant is not None:
|
||
participant.left_at = left_at
|
||
|
||
async def close_all_open_participants(
|
||
self, *, session_id: uuid.UUID, left_at: datetime
|
||
) -> None:
|
||
"""Закрыть все записи участников сеанса с `left_at IS NULL` (`room_finished`)."""
|
||
open_participants = await self._session.scalars(
|
||
select(ConferenceParticipant).where(
|
||
ConferenceParticipant.session_id == session_id,
|
||
ConferenceParticipant.left_at.is_(None),
|
||
)
|
||
)
|
||
for participant in open_participants:
|
||
participant.left_at = left_at
|
||
|
||
|
||
class AudioTrackRepository:
|
||
"""Инкапсулирует SQL-запросы к записанным аудиотрекам сеансов (`session_audio_tracks`)."""
|
||
|
||
def __init__(self, session: AsyncSession) -> None:
|
||
self._session = session
|
||
|
||
async def get_by_session_and_track(
|
||
self, session_id: uuid.UUID, track_sid: str
|
||
) -> SessionAudioTrack | None:
|
||
"""Найти строку трека по (`session_id`, `track_sid`) — ключ идемпотентности `create`."""
|
||
result = await self._session.execute(
|
||
select(SessionAudioTrack).where(
|
||
SessionAudioTrack.session_id == session_id,
|
||
SessionAudioTrack.track_sid == track_sid,
|
||
)
|
||
)
|
||
return result.scalar_one_or_none()
|
||
|
||
async def create(
|
||
self,
|
||
*,
|
||
session_id: uuid.UUID,
|
||
participant_id: uuid.UUID,
|
||
track_sid: str,
|
||
egress_id: str | None,
|
||
file_path: str | None,
|
||
started_at: datetime,
|
||
) -> SessionAudioTrack:
|
||
"""Идемпотентно создать строку трека по (`session_id`, `track_sid`).
|
||
|
||
Повторная доставка `track_published` (например, гонка между двумя
|
||
одновременными вебхуками) — `INSERT ... ON CONFLICT DO NOTHING` по
|
||
уникальному индексу `uq_session_track`, затем возврат уже
|
||
существующей строки. `status` по умолчанию `recording`.
|
||
"""
|
||
insert_stmt = (
|
||
pg_insert(SessionAudioTrack)
|
||
.values(
|
||
session_id=session_id,
|
||
participant_id=participant_id,
|
||
track_sid=track_sid,
|
||
egress_id=egress_id,
|
||
file_path=file_path,
|
||
started_at=started_at,
|
||
)
|
||
.on_conflict_do_nothing(constraint="uq_session_track")
|
||
.returning(SessionAudioTrack.id)
|
||
)
|
||
inserted_id = (await self._session.execute(insert_stmt)).scalar_one_or_none()
|
||
if inserted_id is None:
|
||
existing = await self.get_by_session_and_track(session_id, track_sid)
|
||
assert existing is not None # конфликт гарантирует существование строки
|
||
return existing
|
||
|
||
await self._session.flush()
|
||
record = await self._session.get(SessionAudioTrack, inserted_id)
|
||
assert record is not None
|
||
return record
|
||
|
||
async def finalize(
|
||
self,
|
||
*,
|
||
egress_id: str,
|
||
status: str,
|
||
ended_at: datetime,
|
||
file_path: str | None = None,
|
||
) -> SessionAudioTrack | None:
|
||
"""Финализировать строку трека по `egress_id` (обработка webhook `egress_ended`).
|
||
|
||
`status` — `'recorded'` при успехе, `'failed'` при ошибке egress.
|
||
Отсутствие строки (например, `track_published` был потерян) — не
|
||
ошибка, лог оставляет вызывающая сторона.
|
||
"""
|
||
result = await self._session.execute(
|
||
select(SessionAudioTrack).where(SessionAudioTrack.egress_id == egress_id)
|
||
)
|
||
record = result.scalar_one_or_none()
|
||
if record is None:
|
||
return None
|
||
|
||
record.status = status
|
||
record.ended_at = ended_at
|
||
if file_path is not None:
|
||
record.file_path = file_path
|
||
return record
|
||
|
||
async def list_by_session(self, session_id: uuid.UUID) -> list[SessionAudioTrack]:
|
||
"""Список всех треков сеанса (используется оркестрацией пайплайна)."""
|
||
result = await self._session.execute(
|
||
select(SessionAudioTrack).where(SessionAudioTrack.session_id == session_id)
|
||
)
|
||
return list(result.scalars().all())
|