"""Конфигурация Celery-приложения: broker Redis и расписание периодических задач. Пакет `workers` запускается в окружении backend (тот же venv, `models`/`core` доступны как top-level модули — см. `deploy/docker-compose.yml`, сервис `worker`, и зависимости Celery в `backend/pyproject.toml`). """ from celery import Celery from core.config import get_settings settings = get_settings() app = Celery("vidconf", broker=settings.redis_url, backend=None) # Периодические задачи (beat). cleanup_conferences (ADR-001): # закрытие зависших сеансов (>10 мин без участников) и завершение # просроченных незакреплённых плановых конференций без единого сеанса. # recover_stuck_summaries — уровень 2 защиты от потери # постановки `summarize_session` при сбое брокера в `run_pipeline`: раз в # 5 минут переставляет задачу для сеансов, зависших в `summarizing` без # `summary_data` дольше `STUCK_SUMMARIZING_THRESHOLD` (30 мин) — интервал # планировщика намного короче порога, чтобы «зависание» было устранено # в течение нескольких минут после порога, а не одним запросом на грани. # recover_stuck_notifications — тот же уровень 2 защиты для # следующего шага пайплайна: переставляет `notify_session` для сеансов с уже # готовым `summary_data`, для которых её постановка из `summarize_session` # не удалась (тот же интервал/порог 5мин/30мин, симметрично recover-stuck-summaries). app.conf.beat_schedule = { "cleanup-conferences": { "task": "workers.tasks.maintenance.cleanup_conferences", "schedule": 60.0, }, "recover-stuck-summaries": { "task": "workers.tasks.maintenance.recover_stuck_summaries", "schedule": 300.0, }, "recover-stuck-notifications": { "task": "workers.tasks.maintenance.recover_stuck_notifications", "schedule": 300.0, }, } # Маршрутизация задач по очередям. # # `run_pipeline` — отдельная очередь `transcription`: слушает только # `worker-transcriber` (`--pool=solo`, faster-whisper/ctranslate2 несовместимы # с prefork-пулом Celery), базовый `worker` эту очередь не слушает. # `pipeline_producer.enqueue_pipeline` уже передаёт `queue=` явно при # отправке — этот маршрут страхует остальные пути постановки задачи # (beat/ретраи/ручной вызов по имени без явной очереди). # # `summarize_session` — очередь `summarize`, `notify_session`/`send_invitations` # (рассылка .ics-приглашений) — очередь `notify`: обе ставятся через # `app.send_task` (общий клиент `app` из этого модуля, `workers.tasks.dispatch`) # либо через отдельный клиент `services/invitations_producer.py` (backend), у # которого нет собственного `task_routes` — там очередь передаётся явным # `queue="notify"` при отправке (см. этот файл), а маршрут ниже страхует # остальные пути постановки (ретраи `send_invitations` через `task.retry` # используют исходную очередь сообщения, а не эту конфигурацию). # # Задачи обслуживания (`workers.tasks.maintenance.*` — `cleanup_conferences`, # `recover_stuck_summaries`, `recover_stuck_notifications`) явного маршрута не # получают и остаются на дефолтной очереди Celery `celery` — базовый `worker` # слушает `celery,summarize,notify` (`docs/deploy/scaling.md`), поэтому они # обрабатываются тем же контейнером, что и summarize/notify на малых пресетах. # # Приоритет между очередями реализуем ИЗОЛЯЦИЕЙ (отдельные очереди слушаются # отдельными воркерами/репликами при масштабировании), а НЕ Redis-priorities # Celery (`Kombu`/`redis` транспорт эмулирует приоритеты через несколько # внутренних списков и не даёт строгих гарантий порядка — по сути ненадёжны # на брокере Redis, в отличие от RabbitMQ; см. `docs/deploy/scaling.md`). app.conf.task_routes = { "workers.tasks.pipeline.run_pipeline": {"queue": "transcription"}, "workers.tasks.summarize.*": {"queue": "summarize"}, "workers.tasks.notify.*": {"queue": "notify"}, "workers.tasks.invitations.*": {"queue": "notify"}, } # Явный импорт модулей с задачами: наши задачи лежат в `workers/tasks/*.py`, # а не в `workers/tasks.py`, поэтому `autodiscover_tasks` со стандартным # суффиксом `.tasks` их бы не нашёл. `workers.tasks.pipeline` и # `workers.tasks.summarize` регистрируются в обоих процессах (базовый # `worker` и `worker-transcriber` запускают один и тот же `-A # workers.celery_app`) — очередь, а не импорт, разграничивает, где задача # реально выполняется (`task_routes` ниже); `pipeline.py` ссылается на # `summarize_session`, а `summarize.py` — на `notify_session` # по имени задачи (`app.send_task`/`workers.tasks.dispatch.send_task_with_retry`), # не импортируя модули друг друга напрямую. import workers.tasks.invitations # noqa: E402,F401 import workers.tasks.maintenance # noqa: E402,F401 import workers.tasks.notify # noqa: E402,F401 import workers.tasks.pipeline # noqa: E402,F401 import workers.tasks.summarize # noqa: E402,F401