95 lines
6.4 KiB
Python
95 lines
6.4 KiB
Python
"""Конфигурация 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
|