"""Периодические задачи обслуживания конференций (beat), ADR-001. Комнаты как отдельная сущность отсутствуют: закрываются зависшие сеансы (пропущенный `room_finished`) и завершаются просроченные незакреплённые плановые конференции, за которые так никто и не подключился. Также содержит `recover_stuck_summaries` — уровень 2 защиты от «зависания» сеанса в `pipeline_status='summarizing'` без `summary_data`: если постановка `summarize_session` в очередь из `workers.tasks.pipeline.run_pipeline` не удалась даже после её собственных retry (брокер был недоступен дольше, чем длится backoff), эта задача находит такие сеансы по расписанию и переставляет `summarize_session` повторно — безопасно благодаря идемпотентным guard'ам самой задачи суммаризации. И `recover_stuck_notifications` — тот же уровень 2 защиты, но для следующего шага пайплайна: саммари уже готово (`summary_data IS NOT NULL`), а `notify_session` из `workers.tasks.summarize.summarize_session_async` не была поставлена (сбой брокера дольше её собственного retry в `send_task_with_retry`) — эта задача переставляет `notify_session` повторно. """ import logging from datetime import UTC, datetime, timedelta from core.plugins.config import InstanceConfig from repositories.conferences import ConferenceRepository, ConferenceSessionRepository from services.instance_settings import load_effective_config from workers.celery_app import app as app from workers.db import open_session, run_async from workers.livekit_client import delete_livekit_room logger = logging.getLogger(__name__) # Страховочный порог простоя сеанса: если он открыт (`t_end IS NULL`) без # активных участников дольше этого времени — считаем его зависшим (например, # если webhook `room_finished` потерялся) и закрываем напрямую. IDLE_THRESHOLD = timedelta(minutes=10) # Запас после конца планового окна (`scheduled_at` + `duration_minutes`), # прежде чем считать незакреплённую плановую конференцию без единого сеанса # просроченной и переводить её в `ended`. SCHEDULED_GRACE = timedelta(minutes=60) # Сколько сеанс может провисеть в `pipeline_status='summarizing'` без # `summary_data`, прежде чем считать постановку `summarize_session` потерянной # и переставить задачу повторно. 30 минут — заведомо больше суммарного # backoff'а retry в `run_pipeline._send_summarize_task` (секунды) и # retry самой `summarize_session` при обрыве LLM (минуты), так что recovery # не пересекается с их штатной работой. STUCK_SUMMARIZING_THRESHOLD = timedelta(minutes=30) # Сколько сеанс может провисеть с готовым summary_data без перехода в # 'notified', прежде чем считать постановку `notify_session` потерянной и # переставить задачу повторно. Тот же порог, что у суммаризации (30 минут) — # заведомо больше суммарного backoff'а `send_task_with_retry` и retry самой # `notify_session` при временном сбое SMTP. STUCK_NOTIFYING_THRESHOLD = timedelta(minutes=30) @app.task(name="workers.tasks.maintenance.cleanup_conferences") def cleanup_conferences() -> None: """Точка входа Celery beat — синхронная обёртка над асинхронной логикой.""" run_async(lambda: cleanup_conferences_async(datetime.now(UTC))) async def cleanup_conferences_async(now: datetime) -> None: """Закрыть зависшие сеансы и завершить просроченные незакреплённые плановые конференции. Идемпотентна: повторный запуск с теми же (или более поздними) данными в БД — no-op, т.к. каждый шаг проверяет текущее состояние перед действием (`t_end IS NULL`, наличие активных участников, `status='scheduled'`). Два независимых шага: 1. Открытые сеансы (`t_end IS NULL`) без активных участников, идущие дольше `IDLE_THRESHOLD` — закрываются напрямую (`t_end = now`), не дожидаясь LiveKit; статус родительской конференции переводится в `scheduled` (закреплённая) или `ended` (незакреплённая) — так же, как это сделал бы webhook `room_finished`. Дополнительно запрашивается удаление LiveKit-комнаты (идемпотентно) — страховка на случай, если комната в LiveKit всё ещё существует. 2. Незакреплённые плановые конференции (`status='scheduled'`), чьё плановое окно (`scheduled_at` + `duration_minutes` + запас) истекло, а сеанс так и не появился — переводятся в `ended` напрямую. """ async with open_session() as session: conferences = ConferenceRepository(session) sessions = ConferenceSessionRepository(session) for open_session_record in await sessions.list_open(): if open_session_record.t_start > now - IDLE_THRESHOLD: continue if await sessions.has_active_participants(open_session_record.id): continue conference = await conferences.get_by_id(open_session_record.conference_id) await sessions.close(open_session_record, t_end=now) await sessions.close_all_open_participants( session_id=open_session_record.id, left_at=now ) if conference is not None: if conference.is_pinned: conference.status = "scheduled" else: conference.status = "ended" conference.ended_at = now await delete_livekit_room(conference.slug) logger.info( "cleanup_conferences: сеанс %s конференции %s закрыт как простаивающий " "(t_start=%s), новый статус=%s", open_session_record.id, conference.id, open_session_record.t_start, conference.status, ) for conference in await conferences.list_expired_unpinned_scheduled(now=now): assert conference.scheduled_at is not None # гарантировано запросом репозитория grace_minutes = conference.duration_minutes or 0 expiry = conference.scheduled_at + timedelta(minutes=grace_minutes) + SCHEDULED_GRACE if expiry >= now: continue conference.status = "ended" conference.ended_at = now logger.info( "cleanup_conferences: незакреплённая плановая конференция %s просрочена без " "сеансов (scheduled_at=%s) — переведена в ended", conference.id, conference.scheduled_at, ) await session.commit() @app.task(name="workers.tasks.maintenance.recover_stuck_summaries") def recover_stuck_summaries() -> None: """Точка входа Celery beat — синхронная обёртка над асинхронной логикой.""" run_async(lambda: recover_stuck_summaries_async(datetime.now(UTC))) async def recover_stuck_summaries_async( now: datetime, *, plugins_config: InstanceConfig | None = None ) -> None: """Переставить `summarize_session` для сеансов, зависших в `summarizing` без summary. Уровень 2 защиты от потери постановки задачи суммаризации при недоступности брокера в `run_pipeline` (см. докстринг `workers.tasks.dispatch.send_task_with_retry`). Идемпотентна: повторная отправка `summarize_session` безопасна — задача сама завершится no-op, если `summary_data` уже заполнен либо `pipeline_status` уже другой (её собственные guard'ы). Разница с `summarizer.enabled=false`: в этом случае `summarize_session` тоже отвечает no-op, но КАЖДЫЙ запуск этой задачи заново слал бы её впустую каждые `beat`-интервал — дёшево отличить заранее (эффективная конфигурация, которую читает и сама `summarize_session`) и выйти сразу, не выполняя запрос списка зависших сеансов. """ async with open_session() as session: cfg = plugins_config or await load_effective_config(session) if not cfg.summarizer.enabled: logger.info("recover_stuck_summaries: summarizer отключён (enabled=false) — no-op") return sessions = ConferenceSessionRepository(session) stuck = await sessions.list_stuck_summarizing(older_than=now - STUCK_SUMMARIZING_THRESHOLD) for session_record in stuck: app.send_task( "workers.tasks.summarize.summarize_session", args=[str(session_record.id)] ) logger.warning( "recover_stuck_summaries: сеанс %s завис в pipeline_status=summarizing " "(t_end=%s) без summary_data — summarize_session переставлена в очередь", session_record.id, session_record.t_end, ) @app.task(name="workers.tasks.maintenance.recover_stuck_notifications") def recover_stuck_notifications() -> None: """Точка входа Celery beat — синхронная обёртка над асинхронной логикой.""" run_async(lambda: recover_stuck_notifications_async(datetime.now(UTC))) async def recover_stuck_notifications_async(now: datetime) -> None: """Переставить `notify_session` для сеансов, зависших с готовым summary без уведомления. Уровень 2 защиты от потери постановки `notify_session` при недоступности брокера в `summarize_session_async` (см. докстринг `workers.tasks.dispatch.send_task_with_retry`). Идемпотентна: повторная отправка `notify_session` безопасна — задача сама завершится no-op, если `pipeline_status` уже `notified`/`failed` либо получатели уже все уведомлены (её собственные guard'ы и таблица `email_deliveries`). Разница с `summarizer.enabled=false`: в этом случае сеансов с `summary_data IS NOT NULL` в статусе `summarizing` просто не появится (summarize_session сама выходит раньше) — отдельная проверка конфигурации не нужна, в отличие от `recover_stuck_summaries_async`. """ async with open_session() as session: sessions = ConferenceSessionRepository(session) stuck = await sessions.list_stuck_notifying(older_than=now - STUCK_NOTIFYING_THRESHOLD) for session_record in stuck: app.send_task("workers.tasks.notify.notify_session", args=[str(session_record.id)]) logger.warning( "recover_stuck_notifications: сеанс %s завис в pipeline_status=summarizing " "(t_end=%s) с готовым summary_data — notify_session переставлена в очередь", session_record.id, session_record.t_end, )