387 lines
21 KiB
Markdown
387 lines
21 KiB
Markdown
# 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`
|