Files
vidconf/workers/README.md
Max Ronzhin 8757bec8ac
Some checks failed
CI / backend (push) Has been cancelled
CI / frontend (push) Has been cancelled
first commit
2026-07-23 02:38:05 +03:00

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