Files
vidconf/backend/services/pipeline_producer.py

71 lines
4.2 KiB
Python
Raw 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.
"""Постановка задачи `run_pipeline` в очередь Celery `transcription`.
Отправляет задачу по имени через голый Celery-клиент (`broker_url` без
`backend`) — backend не импортирует пакет `workers` (инвариант разделения
слоёв: API-процесс не должен тянуть зависимости воркеров, в т.ч. `faster-whisper`
через транзитивный импорт `workers.tasks.pipeline`). Задача регистрируется
и обрабатывается в `workers/tasks/pipeline.py` (блок C).
"""
import uuid
from celery import Celery
from kombu.exceptions import OperationalError
from core.config import get_settings
RUN_PIPELINE_TASK_NAME = "workers.tasks.pipeline.run_pipeline"
TRANSCRIPTION_QUEUE = "transcription"
def enqueue_pipeline(session_id: uuid.UUID) -> None:
"""Поставить в очередь `transcription` запуск AI-пайплайна для сеанса `session_id`.
Идемпотентно на стороне задачи (`run_pipeline` продолжает с последнего
успешного шага) — повторная постановка (например,
из `_on_room_finished` при повторной обработке) безопасна.
"""
settings = get_settings()
client = Celery("vidconf-producer", broker=settings.redis_url)
client.send_task(RUN_PIPELINE_TASK_NAME, args=[str(session_id)], queue=TRANSCRIPTION_QUEUE)
def transcription_queue_served(timeout: float = 1.0) -> bool:
"""Проверить, обслуживается ли очередь `transcription` хотя бы одним воркером Celery.
Детект доступности уровня
AI (`services.ai_levels.detect_ai_levels`) смотрит только на железо и
файлы моделей на диске, но не видит, запущен ли вообще воркер
транскрибации, — админка (`GET /admin/settings`) использует эту функцию,
чтобы предупредить «AI включён, но обработка недоступна».
`app.control.inspect(timeout=...).active_queues()` — блокирующий вызов
(ждёт ответа брокера/воркеров); возвращает `{hostname: [{"name": ...}, ...]}`
для ответивших воркеров либо `None`, если за `timeout` не ответил НИ ОДИН
(нет запущенных воркеров либо брокер Redis недоступен) — оба случая здесь
трактуются как «очередь не обслуживается». Вызывающая сторона (API-хендлер)
должна оборачивать в `anyio.to_thread.run_sync`, чтобы не блокировать event loop.
Полностью недоступный брокер (Redis лежит/не резолвится) — отдельный
случай: `kombu`/`redis-py` не возвращают `None`, а бросают исключение,
которое Celery оборачивает в `kombu.exceptions.OperationalError`
(«Recoverable message transport connection error» — проверено
эмпирически: недоступный/несуществующий хост даёт именно этот тип).
По контракту «нет воркеров ИЛИ брокер недоступен → `False`» это тоже
трактуется как «очередь не обслуживается», а не пробрасывается 500-кой
наружу в `GET`/`PUT /admin/settings`.
"""
settings = get_settings()
client = Celery("vidconf-producer", broker=settings.redis_url)
try:
active_queues = client.control.inspect(timeout=timeout).active_queues()
except OperationalError:
return False
if not active_queues:
return False
return any(
queue.get("name") == TRANSCRIPTION_QUEUE
for queues in active_queues.values()
for queue in queues
)