Files
vidconf/workers/tasks/maintenance.py

213 lines
13 KiB
Python
Raw Permalink 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.
"""Периодические задачи обслуживания конференций (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,
)