# Workers — Celery Tasks Асинхронная пост-обработка сеансов конференций: транскрибация аудиотреков, суммаризация, email-уведомления и рассылка .ics-приглашений. Работает через Celery + Redis, использует общее окружение backend (модели, репозитории, плагины). ## Структура ``` workers/ ├── celery_app.py Конфигурация Celery: broker, beat schedule, task_routes ├── db.py Утилиты для работы с БД в воркерах (async-сессии) ├── livekit_client.py Обёртка для LiveKit API (получение записанных треков) ├── tasks/ │ ├── __init__.py │ ├── dispatch.py send_task_with_retry — безопасная постановка следующей │ │ задачи пайплайна при временной недоступности брокера │ ├── maintenance.py Периодические задачи (beat): очистка зависших сеансов/ │ │ конференций, восстановление зависших саммари/уведомлений │ ├── pipeline.py run_pipeline — диспетчер: транскрибация треков + реконструкция фраз │ ├── summarize.py summarize_session — суммаризация (map-reduce через LLM) │ ├── notify.py notify_session — email с резюме сеанса │ └── invitations.py send_invitations — рассылка .ics-приглашений ├── transcription/ │ ├── __init__.py │ ├── phrases.py Реконструкция фраз из сегментов треков (чистая функция) │ └── README.md Документация модуля ├── summarizer/ │ ├── __init__.py │ ├── transcript.py Сборка транскрипта из фраз (чистая функция) │ ├── prompts/ │ │ ├── summary_map_ru.txt Промпт map-шага map-reduce суммаризации │ │ └── summary_reduce_ru.txt Промпт reduce-шага │ └── eval/ Eval-корпус для сравнения качества по уровням AI │ ├── run_tiers.py │ ├── corpus/ │ └── results/ └── README.md Этот файл ``` ## Архитектура пайплайна Celery + Redis для асинхронной пост-обработки сеансов. ### Состав Celery-задач | Задача | Входные данные | Выходные данные | |--------|----------------|-----------------| | `run_pipeline(session_id)` | session_id | `session_audio_tracks.segments`, `phrases` | | `summarize_session(session_id)` | session_id | `conference_sessions.summary_data` | | `notify_session(session_id)` | session_id | email с резюме отправлен, запись в `email_deliveries` | | `send_invitations(conference_id, emails)` | conference_id, список email | .ics-приглашение отправлено | ### State machine pipeline_status ``` recording ↓ (webhook room_finished → enqueue run_pipeline) transcribing ↓ (транскрибация + реконструкция фраз успешны) summarizing ↓ (summarize_session записывает summary_data; статус остаётся summarizing) ↓ (notify_session отправляет email) notified ↓ SUCCESS: история сохранена в БД или при ошибке на любом шаге: failed ↓ Лог ошибки, повторный запуск вручную ``` **Идемпотентность:** каждый шаг пайплайна проверяет текущий `pipeline_status`: - Если уже прошёл этот шаг → skip (guard от дублирования) - Если failed → воркер логирует и пропускает - Если recording/transcribing → может переходить по ступеням ### Защита от потери постановки следующей задачи Каждый шаг пайплайна ставит следующую задачу в очередь через `workers.tasks.dispatch.send_task_with_retry` (несколько попыток при временной недоступности Redis-брокера). Если все попытки исчерпаны, шаг не откатывается и не проваливается — восстановление берёт на себя beat-задача `workers.tasks.maintenance`: - `recover_stuck_summaries` — переставляет `summarize_session` для сеансов, зависших в `pipeline_status='summarizing'` без `summary_data` дольше порога. - `recover_stuck_notifications` — переставляет `notify_session` для сеансов с готовым `summary_data`, но без отправленного уведомления. Оба идемпотентны: повторная постановка безопасна благодаря собственным guard'ам задач. ### Транскрибация (workers/tasks/pipeline.py) **Задача `run_pipeline(session_id)`** — диспетчер пайплайна пост-обработки сеанса. 1. Проверка статуса сеанса (`pipeline_status`) и `config/plugins.yaml` (enabled/disabled) 2. Ожидание завершения egress: retry-цикл с backoff (egress завершается позже room_finished) 3. Переход в `pipeline_status='transcribing'` 4. **Транскрибация каждого трека** (per-track идемпотентность): - Загрузка плагина (`FasterWhisperCPU`/`FasterWhisperGPU` или другого из конфига) - `transcriber.transcribe(audio_path, language='ru')` - Сохранение сегментов в `session_audio_tracks.segments` (JSONB) - Коммит после каждого трека (точка возобновления) 5. **Реконструкция фраз** (если хотя бы один трек успешен): - Сборка сегментов по участникам (с учётом смещений треков) - `build_phrases(segments_by_participant, track_offsets)` — чистая функция - DELETE phrases WHERE session_id + INSERT новые фразы в одной транзакции - Переход в `pipeline_status='summarizing'` 6. **Обработка ошибок:** - Если все треки failed → `pipeline_status='failed'` - Зависшие треки (исчерпаны retry) → помечаются `failed`, продолжаем с остальными **Плагин FasterWhisperCPU/GPU** (`backend/core/plugins/faster_whisper.py`): - Встроенный Silero VAD (подавление молчания) - Отбрасывание сегментов < 0.3 сек (галлюцинации Whisper) - Ленивый импорт (не грузит зависимость в API-процесс) - Синглтон модели на процесс воркера **Реконструкция фраз** (`workers/transcription/phrases.py`, чистая функция): - Слияние сегментов всех участников в таймлайн - Группировка подряд идущих сегментов одного спикера → фраза - Короткие вставки другого спикера (<1.5 сек) не рвут фразу - Перекрытия речи → обе фразы - Тишина не рвёт фразу Типы данных: ```python segments_by_participant: dict[uuid.UUID, list[Segment]] # Выход faster-whisper track_offsets: dict[uuid.UUID, float] # Смещения по времени → list[PhraseDraft] # Фразы (без БД) ``` Детали: [workers/transcription/README.md](transcription/README.md). ### Суммаризация (workers/tasks/summarize.py) **Задача `summarize_session(session_id)`** — собрать транскрипт сеанса, разбить на чанки, создать резюме: 1. Guard'ы: сеанс найден, не завершён, статус `summarizing`, фраз достаточно 2. Собрать транскрипт: фразы → `build_transcript()` → строки `[Имя MM:SS] текст` 3. Инстанциирование плагина `Summarizer` (из эффективной конфигурации уровня AI, `backend/services/instance_settings.py::load_effective_config`) 4. Вызвать `summarizer.summarize(transcript)`: - `chunk_transcript()` — разбить на чанки (20 минут, max_chunk_tokens зависит от уровня) - Map: LLM по каждому чанку (`summary_map_ru.txt`, max_tokens_map per-tier) - Reduce: объединить резюме (иерархически при переполнении, `summary_reduce_ru.txt`, max_tokens_reduce per-tier) 5. `session.summary_data = результат`, commit 6. **Статус ОСТАЁТСЯ `summarizing`** — готовность определяется парой: `pipeline_status='summarizing'` И `summary_data IS NOT NULL`; переход в `notified` — задача `notify_session` **Уровни AI** (см. ADR-004 `docs/architecture/adr/004-ai-tier-matrix.md`): - **min:** Qwen3.5-4B, max_tokens_map=1024, max_tokens_reduce=1536, CPU llama.cpp - **medium:** Qwen3.5-9B, max_tokens_map=1024, max_tokens_reduce=2048, CPU/GPU опционально - **max:** Qwen3.5-35B-A3B (MoE), max_tokens_map=1536, max_tokens_reduce=2560, GPU llama.cpp обязателен **Надёжность:** - Retry + circuit breaker в HTTP-клиенте (3 попытки, экспоненциальный backoff, breaker после 5 ошибок) - Celery-retry задачи (bind + acks_late) - Идемпотентные guard'ы: `summary_data IS NOT NULL` → no-op **Конфигурация:** - [config/plugins.yaml](../config/plugins.yaml) — дефолт `provider: "null"`, для пресета 3+ = `"qwen_local"` - [backend/services/ai_tiers.py](../backend/services/ai_tiers.py) — матрица уровней min/medium/max **Установка модели:** [docs/deploy/llm-setup.md](../docs/deploy/llm-setup.md) **Документация:** [docs/plugins/summarizer.md](../docs/plugins/summarizer.md) **Качество:** [docs/deploy/quality-tiers.md](../docs/deploy/quality-tiers.md) (параметры per-tier, eval-корпус) ### Уведомления (workers/tasks/notify.py) и приглашения (workers/tasks/invitations.py) `notify_session(session_id)` собирает получателей (владелец, участники по `summary_recipients`), формирует письмо (`backend/services/email_templates.py`) и отправляет через `backend/services/email.py` (backend `console`/`smtp`), записывая доставку в `email_deliveries` (идемпотентность: уже отправленным получателям письмо повторно не шлётся). `send_invitations(conference_id, emails)` формирует .ics-приглашение (`backend/services/ics.py`) и рассылает его приглашённым и организатору. ## Промпты суммаризации Русскоязычные промпты для Qwen3.5, единые для всех уровней AI: - `summarizer/prompts/summary_map_ru.txt` — map-шаг (суммаризация чанка) - `summarizer/prompts/summary_reduce_ru.txt` — reduce-шаг (объединение сумм) Качество наращивается размером модели (4B/9B/35B-A3B), а не правкой промптов. Per-tier параметры генерации (max_tokens, temperature) задаются в [backend/services/ai_tiers.py](../backend/services/ai_tiers.py) (ADR-004). ## Инфраструктура Celery ### celery_app.py ```python from celery import Celery from core.config import get_settings settings = get_settings() app = Celery("vidconf", broker=settings.redis_url, backend=None) 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"}, } 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, }, } ``` ### tasks/maintenance.py **cleanup_conferences:** Закрытие зависших сеансов и завершение просроченных плановых конференций: 1. Открытые сеансы без активных участников дольше `IDLE_THRESHOLD` (10 мин) — закрываются напрямую, статус родительской конференции — `scheduled` (закреплённая) или `ended` (незакреплённая) 2. Незакреплённые плановые конференции, чьё плановое окно истекло без единого сеанса — переводятся в `ended` **recover_stuck_summaries / recover_stuck_notifications** — уровень 2 защиты от потери постановки задачи при недоступности Redis-брокера (см. «Защита от потери постановки следующей задачи» выше). **Идемпотентность:** каждый шаг проверяет текущее состояние перед действием, поэтому повторные запуски безопасны. ### Запуск **Docker Compose:** ```bash docker compose -f deploy/docker-compose.yml up -d # Включает сервисы worker/worker-transcriber docker compose -f deploy/docker-compose.yml logs -f worker ``` **Локально:** ```bash cd backend PYTHONPATH=/path/to/backend uv run celery -A workers.celery_app worker -B --loglevel=info ``` ## Раздельные очереди Celery Воркеры слушают выделенные очереди для изоляции нагрузки (см. `app.conf.task_routes` выше): | Очередь | Задачи | Назначение | |---------|--------|-----------| | `transcription` | `run_pipeline` | Интенсивная обработка аудио (CPU/GPU) | | `summarize` | `summarize_session` | Инференс LLM | | `notify` | `notify_session`, `send_invitations` | Email, низкий приоритет | | `celery` (default) | остальные задачи (beat, maintenance) | Периодические + сервисные | Подробнее о раздельных воркер-контейнерах и профилях compose — [docs/deploy/scaling.md](../docs/deploy/scaling.md). ## Метрики **Endpoint:** `GET /metrics` (Prometheus, реализован в `backend/api/metrics.py`) - `vidconf_http_request_duration_seconds` — гистограмма латентности HTTP-запросов backend (по маршруту/методу/статусу) - `vidconf_pipeline_sessions` — gauge числа сеансов в каждом статусе `pipeline_status`, пересчитывается при каждом scrape - `vidconf_celery_queue_depth` — gauge длины очередей Celery в Redis (`transcription`/`summarize`/`notify`/`celery`), `LLEN` при каждом scrape **Визуализация:** Grafana (compose-профиль `monitoring`, см. [docs/deploy/monitoring.md](../docs/deploy/monitoring.md)) ## Eval-корпус для качества **Директория:** `workers/summarizer/eval/` **Назначение:** верификация качества суммаризации на разных уровнях AI (min/medium/max). **Структура:** - `run_tiers.py` — скрипт для запуска eval на всех уровнях - `corpus/` — тестовые транскрипты - `results/` — результаты прогонов по уровням **Использование:** ```bash cd workers python summarizer/eval/run_tiers.py --level min # eval уровня min (Qwen3.5-4B) python summarizer/eval/run_tiers.py --level medium python summarizer/eval/run_tiers.py --level max ``` Методика и результаты прогонов — [docs/deploy/quality-tiers.md](../docs/deploy/quality-tiers.md). ## Конфигурация ### .env переменные Полный справочник — [docs/deploy/env.md](../docs/deploy/env.md). Ключевые для воркеров: ```bash REDIS_URL=redis://redis:6379/0 DATABASE_URL=postgresql+asyncpg://vidconf:vidconf@postgres:5432/vidconf PLUGINS_CONFIG_PATH=config/plugins.yaml # Email EMAIL_BACKEND=console # или smtp SMTP_HOST=localhost SMTP_PORT=587 # LiveKit (получение записанных аудиотреков) LIVEKIT_API_KEY=devkey LIVEKIT_API_SECRET=change-me-livekit-secret LIVEKIT_PUBLIC_URL=ws://localhost:7880 # Обнаруженное железо (заполняет install.sh, используется для детекта уровней AI) HW_CPUS= HW_RAM_MB= HW_GPU_NAME= HW_VRAM_MB= ``` ### Docker Compose ```bash # Запустить весь стек с воркерами docker compose -f deploy/docker-compose.yml up -d # Только воркер (при условии postgres/redis уже запущены) docker compose -f deploy/docker-compose.yml up -d worker # Логи docker compose -f deploy/docker-compose.yml logs -f worker # Статус docker compose -f deploy/docker-compose.yml ps ``` ## Тестирование ```bash cd backend # Workers используют окружение backend uv run pytest tests/ -v # С режимом Celery eager (синхронно, без реальной очереди) CELERY_ALWAYS_EAGER=True uv run pytest tests/ -v ``` ## Ссылки ### Backend & плагины - **Backend README:** `backend/README.md` — FastAPI, контракты плагинов - **Плагины:** - Контракты: `docs/plugins/contracts.md` - Транскрибатор: `docs/plugins/transcriber.md` - Суммаризатор: `docs/plugins/summarizer.md` - **Уровни AI:** - Матрица уровней (min/medium/max): `backend/services/ai_tiers.py` - ADR-004 (архитектурное решение): `docs/architecture/adr/004-ai-tier-matrix.md` ### Модули воркеров - **Пайплайн:** `workers/tasks/pipeline.py` - **Реконструкция фраз:** `workers/transcription/phrases.py` + `workers/transcription/README.md` - **Суммаризация:** `workers/tasks/summarize.py` - **Уведомления и рассылка:** `workers/tasks/notify.py`, `workers/tasks/invitations.py` ### Развёртывание & мониторинг - **Инсталлятор:** `docs/deploy/install.md` (5 пресетов, автодетект) - **Профили оборудования:** `docs/deploy/hardware-profiles.md` (требования ресурсов) - **Качество (eval-корпус):** `docs/deploy/quality-tiers.md` - **Мониторинг:** `docs/deploy/monitoring.md` (Prometheus + Grafana) - **Масштабирование:** `docs/deploy/scaling.md` ### Модели и миграции - **Схема БД:** `docs/db/schema.md` - **Модели:** `backend/models/audio_track.py`, `backend/models/phrase.py`, `backend/models/session.py` ### Архитектурные решения - **ADR-001:** `docs/architecture/adr/001-dynamic-conferences-pivot.md` (динамические конференции) - **ADR-002:** `docs/architecture/adr/002-phrase-attribution-session-participant.md` (атрибуция фраз к участнику сеанса) - **ADR-004:** `docs/architecture/adr/004-ai-tier-matrix.md` (матрица уровней AI) - **Общая архитектура:** `docs/architecture/README.md`