docker compose определяет .env для подстановки ${VAR} по каталогу
compose-файла (deploy/), а не по текущей директории — repo-root .env,
который использует install.sh и вся документация, молча не подхватывался.
Это и была причина "WARN: LIVEKIT_API_KEY not set" на боевом сервере:
секреты были в .env, но compose их не видел и подставлял небезопасные
дефолты (см. таблицу разведки сервера).
Теперь ставшие обязательными ${VAR:?} в docker-compose.yml (см. предыдущий
коммит) без этого немедленно проваливали бы конфиг на любой из
документированных команд. Добавлен `--env-file .env`/"$ENV_FILE" ко всем
вызовам docker compose в install.sh и в командах из README/docs.
install.sh дополнительно: ensure_default (аналог ensure_secret без генерации
секрета) для новых не-секретных параметров nginx/coturn/livekit
(NGINX_SERVER_NAMES, NGINX_CERT_NAME, LIVEKIT_USE_EXTERNAL_IP,
LIVEKIT_NODE_IP, TURN_EXTERNAL_IP) — дефолты только для локальной
разработки, не перезаписывают значения, заданные вручную на боевом
сервере. Плюс вызов deploy/render-templates.sh перед сборкой/подъёмом
стека.
492 lines
21 KiB
Markdown
492 lines
21 KiB
Markdown
# Плагин Summarizer: QwenLocal
|
||
|
||
## Обзор
|
||
|
||
Плагин `QwenLocal` реализует суммаризацию сеансов конференций с использованием локальной модели **Qwen3.5** (3 уровня качества: 4B/9B/35B-A3B) через OpenAI-совместимый сервер **llama.cpp** (матрица уровней — ADR-004).
|
||
|
||
**Провайдер:** `qwen_local`
|
||
|
||
**Назначение:** создание структурированного резюме текстового транскрипта сеанса, разбитого на логические части (что решено, ответственные, сроки и т.д.) согласно утверждённым промптам `workers/summarizer/prompts/summary_map_ru.txt` и `summary_reduce_ru.txt`.
|
||
|
||
**Архитектурные решения:**
|
||
- [docs/architecture/adr/004-ai-tier-matrix.md](../../docs/architecture/adr/004-ai-tier-matrix.md) — матрица уровней min/medium/max с моделями Qwen3.5
|
||
|
||
## Архитектура
|
||
|
||
### Конвейер обработки
|
||
|
||
```
|
||
Сеанс (phrases в БД)
|
||
↓
|
||
build_transcript(phrases) — собрать строки [Имя MM:SS] текст
|
||
↓
|
||
chunk_transcript() — разбить на чанки (20 мин, ≤8000 токенов)
|
||
↓
|
||
MAP (параллельно по чанкам)
|
||
• summary_map_ru.txt.format(transcript_chunk=чанк)
|
||
• OpenAI API → LLM-сервер (llama.cpp)
|
||
• Результат: одно резюме-чанка
|
||
↓
|
||
REDUCE (иерархически при переполнении)
|
||
• summary_reduce_ru.txt.format(partial_summaries=...)
|
||
• Группировка по бюджету 6000 токенов
|
||
• Итерация до одного финального резюме
|
||
↓
|
||
ConferenceSession.summary_data = JSON
|
||
```
|
||
|
||
### Компоненты
|
||
|
||
#### 1. Сборка транскрипта (`workers/summarizer/transcript.py`)
|
||
|
||
**Функция:** `build_transcript(lines: Sequence[TranscriptLine]) -> str`
|
||
|
||
Входные данные:
|
||
- Список фраз из `phrases` таблицы, преобразованные в `TranscriptLine`:
|
||
- `speaker: str` — имя говорящего (из `users.name_user` ИЛИ `guest_access.display_name`)
|
||
- `offset_s: float` — смещение в секундах от `t_start` сеанса
|
||
- `text: str` — распознанный и реконструированный текст
|
||
|
||
Выходные данные:
|
||
- Строка с фразами, отсортированными по времени:
|
||
```
|
||
[Иван 00:15] Здравствуйте, начнём встречу
|
||
[Мария 01:30] Спасибо, вот моя презентация
|
||
[Иван 15:42] Подводим итоги
|
||
```
|
||
|
||
Формат времени: `MM:SS` (до часа), `ЧЧ:MM:SS` (от часа).
|
||
|
||
#### 2. Чанкинг (`backend/core/summarization/chunking.py`)
|
||
|
||
**Функция:** `chunk_transcript(transcript: str, count_tokens: Callable[[str], int], *, max_chunk_tokens: int = 8000, target_chunk_minutes: int = 20) -> list[str]`
|
||
|
||
Правила закрытия чанка:
|
||
- Каждая строка (фраза) атомарна — не режется пополам
|
||
- Чанк закрывается, когда:
|
||
1. Охваченное время ≥ `target_chunk_minutes` (20 минут дефолт), **ИЛИ**
|
||
2. Добавление следующей фразы превысит `max_chunk_tokens`
|
||
|
||
**Аварийный случай — монолог:**
|
||
- Единственная фраза сама по себе > `max_chunk_tokens` (например, 40-минутный монолог)
|
||
- Фраза делится по границам предложений (`.`, `!`, `?`)
|
||
- Каждая часть получает повторённую исходную метку спикера
|
||
- Части добавляются как отдельные чанки
|
||
|
||
Пример с монологом:
|
||
```
|
||
Входная строка (40 мин): [Директор 00:00] Первое предложение. Второе предложение. ...
|
||
Выход (2 чанка):
|
||
[Директор 00:00] Первое предложение.
|
||
[Директор 00:00] Второе предложение.
|
||
...
|
||
```
|
||
|
||
#### 3. Подсчёт токенов (`backend/core/summarization/tokens.py`)
|
||
|
||
**Класс:** `QwenTokenCounter`
|
||
|
||
```python
|
||
counter = QwenTokenCounter(tokenizer_path="/models/qwen/tokenizer.json")
|
||
token_count = counter("Какой-нибудь текст") # int
|
||
```
|
||
|
||
Поведение:
|
||
- **Штатный режим:** загрузить `tokenizers.Tokenizer` из файла `tokenizer.json` (Qwen2.5-3B)
|
||
- **Ленивая загрузка:** первый вызов — попытка загрузить, результат кэшируется
|
||
- **Фолбэк-эвристика:** если файл отсутствует или повреждён → `len(text) // 3` с предупреждением в логе
|
||
- Пайплайн не падает даже без файла токенизатора
|
||
|
||
#### 4. HTTP-клиент LLM (`backend/core/summarization/llm_client.py`)
|
||
|
||
**Класс:** `OpenAICompatClient`
|
||
|
||
```python
|
||
client = OpenAICompatClient(
|
||
base_url="http://llm:8080/v1",
|
||
model="qwen2.5-3b-instruct-q4_k_m",
|
||
temperature=0.2,
|
||
max_tokens=1024,
|
||
)
|
||
result = client.complete("Ваш промпт здесь") # str
|
||
```
|
||
|
||
**Надёжность:**
|
||
- **Retry с экспоненциальным backoff:** базовая пауза 0.5 сек, растёт как `2^attempt`
|
||
- **Retryable статусы:** 5xx (включая 503 — модель ещё грузится)
|
||
- **Circuit breaker:** после 5 подряд неудач окно отказа 60 сек (все запросы падают сразу без HTTP)
|
||
- **Исключение:** `LlmUnavailableError` при открытом breaker'е или исчерпании retry
|
||
|
||
Пример обработки ошибки:
|
||
```python
|
||
try:
|
||
summary = client.complete(prompt)
|
||
except LlmUnavailableError:
|
||
# Задача пауза и повтор через Celery (countdown зависит от номера попытки)
|
||
task.retry(countdown=60 * (attempt + 1))
|
||
```
|
||
|
||
#### 5. Плагин QwenLocal (`backend/core/plugins/qwen_local.py`)
|
||
|
||
**Класс:** `QwenLocal(Summarizer)`
|
||
|
||
Инкапсулирует весь цикл map-reduce:
|
||
|
||
```python
|
||
@register_summarizer
|
||
class QwenLocal(Summarizer):
|
||
provider: ClassVar[str] = "qwen_local"
|
||
|
||
def __init__(
|
||
self,
|
||
model: str | None = None, # GGUF-модель
|
||
chunk_minutes: int = 20, # Размер чанка
|
||
base_url: str = "http://llm:8080/v1", # LLM-сервер
|
||
tokenizer_path: str = "/models/qwen/tokenizer.json", # Токенизатор
|
||
prompts_dir: str = "workers/summarizer/prompts", # Промпты
|
||
temperature: float = 0.2, # Творчество LLM
|
||
max_tokens: int = 1024, # Макс вывод
|
||
**options: Any, # Доп. параметры
|
||
) -> None: ...
|
||
|
||
def summarize(self, transcript: str) -> str:
|
||
"""Полный цикл: chunk → map → reduce → результат."""
|
||
...
|
||
```
|
||
|
||
**Логика `summarize()`:**
|
||
|
||
1. **Чанкирование:** `chunk_transcript(transcript, self._count_tokens, target_chunk_minutes=...)`
|
||
- Пустой список чанков → вернуть пустую строку (no-op)
|
||
|
||
2. **Map-фаза:** для каждого чанка:
|
||
```python
|
||
prompt = summary_map_ru.txt.replace("{transcript_chunk}", chunk)
|
||
summary = client.complete(prompt)
|
||
```
|
||
|
||
3. **Один чанк?** Вернуть его map-результат напрямую (формат совпадает с reduce)
|
||
|
||
4. **Reduce-фаза:** объединить частичные резюме:
|
||
```python
|
||
if len(partial_summaries) == 1:
|
||
return partial_summaries[0]
|
||
|
||
return self._reduce(partial_summaries) # иерархический reduce
|
||
```
|
||
|
||
**Иерархический reduce:**
|
||
|
||
```python
|
||
def _reduce(self, summaries: list[str]) -> str:
|
||
"""Рекурсивное сведение: группировка по бюджету → reduce каждой группы."""
|
||
while len(summaries) > 1:
|
||
# Все резюме влезают в бюджет 6000 токенов?
|
||
if count_tokens("\n\n".join(summaries)) <= 6000:
|
||
return self._reduce_once(summaries)
|
||
|
||
# Нет → группируем жадно по бюджету
|
||
groups = self._group_by_token_budget(summaries, 6000)
|
||
|
||
# Все ещё одна группа (крупные резюме)? Сводим как есть
|
||
if len(groups) == 1:
|
||
return self._reduce_once(summaries)
|
||
|
||
# Reduce каждую группу отдельно → новый уровень
|
||
summaries = [self._reduce_once(group) for group in groups]
|
||
|
||
return summaries[0]
|
||
```
|
||
|
||
Пример:
|
||
- 10 чанков → 10 map-результатов
|
||
- 10 резюме не влезают в 6000 токенов → группируем на 3 группы
|
||
- 3 reduce-вызова → 3 результата
|
||
- 3 результата влезают → финальный reduce → 1 резюме
|
||
|
||
## Конфигурация
|
||
|
||
### Пример (профиль B — с LLM)
|
||
|
||
**`config/plugins.yaml`:**
|
||
```yaml
|
||
summarizer:
|
||
enabled: true
|
||
provider: qwen_local
|
||
model: qwen2.5-3b-instruct-q4_k_m
|
||
chunk_minutes: 20
|
||
options:
|
||
base_url: http://llm:8080/v1
|
||
tokenizer_path: /models/qwen/tokenizer.json
|
||
temperature: 0.2
|
||
max_tokens: 1024
|
||
```
|
||
|
||
### Параметры
|
||
|
||
| Параметр | Тип | Дефолт | Описание |
|
||
|---|---|---|---|
|
||
| `enabled` | bool | `true` | Включена ли суммаризация |
|
||
| `provider` | str | `"null"` | Провайдер (в репо дефолт `"null"`, для профиля B = `"qwen_local"`) |
|
||
| `model` | str | `"qwen2.5-3b-instruct-q4_k_m"` | GGUF-файл модели (без расширения) |
|
||
| `chunk_minutes` | int | `20` | Целевой размер чанка (минуты) |
|
||
| **options:** | | | |
|
||
| `base_url` | str | `"http://llm:8080/v1"` | OpenAI-совместимый URL сервера llama.cpp |
|
||
| `tokenizer_path` | str | `"/models/qwen/tokenizer.json"` | Путь внутри контейнера worker к файлу токенизатора |
|
||
| `temperature` | float | `0.2` | Творчество генерации (0 = точно, 1 = вариативно) |
|
||
| `max_tokens` | int | `1024` | Максимальный размер одного чанка-резюме |
|
||
|
||
### Дефолт репозитория
|
||
|
||
```yaml
|
||
summarizer:
|
||
enabled: true
|
||
provider: "null" # ← ДЕФОЛТ без LLM-сервера
|
||
model: null
|
||
chunk_minutes: 20
|
||
```
|
||
|
||
Это безопасно для dev-окружений без Docker Compose профиля `llm`. `NullSummarizer` просто возвращает пустую строку.
|
||
|
||
## Надёжность
|
||
|
||
### Retry и circuit breaker
|
||
|
||
**HTTP-клиент** (`OpenAICompatClient`):
|
||
- 3 попытки (max_attempts) на каждый запрос
|
||
- Экспоненциальный backoff: 0.5 сек → 1 сек → 2 сек
|
||
- Circuit breaker открывается после 5 подряд неудач (окно 60 сек)
|
||
|
||
**Celery-задача** (`workers.tasks.summarize.summarize_session`):
|
||
- `max_retries=5` на уровне задачи
|
||
- `acks_late=True` — подтверждение доставки после успеха
|
||
- `countdown=60 * (attempt + 1)` — нарастающая пауза между retry
|
||
|
||
Пример: если на 2-м attempt LLM упадёт → пауза 180 сек (3 минуты) перед 3-й попыткой.
|
||
|
||
### Двухуровневая защита от потери постановки задачи
|
||
|
||
**Уровень 1 — retry при сбое брокера** (`workers.tasks.pipeline._send_summarize_task`):
|
||
- Если `app.send_task` бросит `kombu.exceptions.OperationalError`/`ConnectionError` (Redis недоступен), фразы уже сохранены
|
||
- Делается 3 попытки с линейно растущим backoff (2 сек, 4 сек, 6 сек)
|
||
- Если всё исчерпано — `pipeline_status` остаётся `summarizing` (не откатывается), ошибка логируется
|
||
|
||
**Уровень 2 — периодическое восстановление** (`workers.tasks.maintenance.recover_stuck_summaries`):
|
||
- Beat-задача запускается каждые 5 минут
|
||
- Находит сеансы, зависшие в `pipeline_status='summarizing'` без `summary_data` дольше 30 минут
|
||
- Переставляет `summarize_session` в очередь повторно
|
||
- Безопасна благодаря guard'ам задачи суммаризации (если `summary_data` уже есть — no-op)
|
||
|
||
Таким образом, временная недоступность Redis при постановке задачи не ведёт к зависанию сеанса.
|
||
|
||
### Идемпотентность
|
||
|
||
**Guard'ы в `summarize_session_async()`** (в порядке исполнения):
|
||
|
||
1. **Сеанс не найден ИЛИ не завершён** (`t_end IS NULL`) → логировать warning, выход
|
||
2. **Статус уже не `summarizing`** → no-op (готов к notify или failed)
|
||
3. **`summary_data IS NOT NULL`** → no-op (уже суммаризировано)
|
||
4. **`summarizer.enabled=false`** → логировать info, выход
|
||
5. **Нет фраз** → логировать warning, выход (`summary_data` остаётся NULL)
|
||
|
||
**Пример повторной доставки задачи:**
|
||
```
|
||
Попытка 1: Celery отправляет summarize_session(session_id=abc)
|
||
↓ LLM-сервер упадёт посреди map → LlmUnavailableError
|
||
↓ task.retry(countdown=60)
|
||
↓ Очередь переотправляет через 60 сек
|
||
|
||
Попытка 2: summarize_session(session_id=abc) вызовется снова
|
||
↓ Guard №3 проверит: summary_data IS NOT NULL?
|
||
↓ Если да → no-op (уже готово)
|
||
↓ Если нет → повторить map-reduce
|
||
```
|
||
|
||
## Поведение при `summarizer.enabled=false`
|
||
|
||
**Конфиг:**
|
||
```yaml
|
||
summarizer:
|
||
enabled: false
|
||
```
|
||
|
||
**Поведение:**
|
||
|
||
1. **Пайплайн транскрибации** завершается нормально, сеанс переходит в статус `summarizing`
|
||
2. **Guard №4 в `summarize_session_async()`** проверяет `enabled` → логирует info-уровень и выходит
|
||
3. **`summary_data` остаётся `NULL`** — это нормально (документированное состояние)
|
||
4. **Пайплайн НЕ переходит в `notified`** — остаётся в `summarizing` (уведомление может пропустить сеансы без резюме)
|
||
|
||
Полезно для:
|
||
- dev-окружений без LLM-модели
|
||
- тестирования пайплайна без AI-обработки
|
||
- отключения суммаризации на боевом сервере (дорого по ресурсам)
|
||
|
||
## Фолбэк при отсутствии токенизатора
|
||
|
||
**Сценарий:** dev-окружение без смонтированного volume `llm-models` (файла `tokenizer.json` нет).
|
||
|
||
**QwenTokenCounter** реагирует:
|
||
1. При первом вызове пытается загрузить `tokenizers.Tokenizer` из `tokenizer_path`
|
||
2. Если файл не найден ИЛИ повреждён → логирует warning-уровень
|
||
3. Переходит на эвристику: `count_tokens(text) = len(text) // 3`
|
||
4. **Пайплайн продолжает работать**, но подсчёт токенов менее точен
|
||
|
||
**Последствия:**
|
||
- Чанки могут быть немного больше/меньше целевого размера
|
||
- Reduce может потребовать дополнительную итерацию
|
||
- Общее время обработки увеличится, но не критично
|
||
|
||
**Рекомендация:** в продакшене всегда монтировать volume с токенизатором.
|
||
|
||
## Установка модели
|
||
|
||
Модель и токенизатор устанавливаются в Docker Compose профилем `llm`.
|
||
|
||
**Подробно:** [docs/deploy/llm-setup.md](../deploy/llm-setup.md)
|
||
|
||
**Краткие шаги:**
|
||
|
||
1. Поднять профиль `llm`:
|
||
```bash
|
||
docker compose -f deploy/docker-compose.yml --env-file .env --profile llm up -d llm-model-init llm
|
||
```
|
||
|
||
2. Дождаться готовности:
|
||
```bash
|
||
curl http://localhost:8080/health
|
||
# {"status":"ok"}
|
||
```
|
||
|
||
3. Включить `qwen_local` в конфиге:
|
||
```yaml
|
||
summarizer:
|
||
provider: qwen_local
|
||
```
|
||
|
||
4. Перезапустить worker:
|
||
```bash
|
||
docker compose -f deploy/docker-compose.yml --env-file .env up -d --force-recreate worker
|
||
```
|
||
|
||
## Примеры использования
|
||
|
||
### Прямое использование плагина
|
||
|
||
```python
|
||
from core.plugins.config import load_plugins_config
|
||
from core.plugins.factory import create_summarizer
|
||
|
||
# Загрузить конфиг
|
||
cfg = load_plugins_config("config/plugins.yaml")
|
||
|
||
# Инстанцировать плагин QwenLocal
|
||
summarizer = create_summarizer(cfg.summarizer)
|
||
|
||
# Использовать
|
||
transcript = "[Иван 00:15] Привет...\n[Мария 02:30] Привет..."
|
||
result = summarizer.summarize(transcript)
|
||
print(result)
|
||
# Результат: структурированное резюме
|
||
```
|
||
|
||
### В пайплайне (workers)
|
||
|
||
```python
|
||
# workers/tasks/summarize.py
|
||
|
||
async def summarize_session_async(
|
||
task: RetryableTask,
|
||
session_id: uuid.UUID,
|
||
*,
|
||
plugins_config: PluginsConfig | None = None,
|
||
) -> None:
|
||
"""Суммаризировать сеанс и сохранить в БД."""
|
||
|
||
# ... guard'ы пропущены для краткости ...
|
||
|
||
# Собрать транскрипт
|
||
lines = [
|
||
TranscriptLine(
|
||
speaker=speaker_name,
|
||
offset_s=phrase.offset_s,
|
||
text=phrase.text,
|
||
)
|
||
for phrase in phrases
|
||
]
|
||
transcript = build_transcript(lines)
|
||
|
||
# Инстанцировать плагин
|
||
cfg = plugins_config or load_plugins_config()
|
||
summarizer = create_summarizer(cfg.summarizer)
|
||
|
||
# Получить резюме (может вызвать LlmUnavailableError)
|
||
try:
|
||
summary_text = summarizer.summarize(transcript)
|
||
except LlmUnavailableError:
|
||
# Retry с нарастающей паузой
|
||
countdown = RETRY_COUNTDOWN_BASE_S * (task.request.retries + 1)
|
||
raise task.retry(countdown=countdown)
|
||
|
||
# Сохранить в БД
|
||
session.summary_data = summary_text
|
||
await db_session.commit()
|
||
```
|
||
|
||
## Тестирование
|
||
|
||
### Модульные тесты компонентов
|
||
|
||
```bash
|
||
cd backend
|
||
|
||
# Чанкинг (без LLM)
|
||
uv run pytest tests/test_chunking.py -v
|
||
|
||
# LLM-клиент (на мок-транспорте)
|
||
uv run pytest tests/test_llm_client.py -v
|
||
|
||
# Плагин QwenLocal (на мок-LLM)
|
||
uv run pytest tests/test_qwen_local.py -v
|
||
```
|
||
|
||
### Интеграционный тест пайплайна
|
||
|
||
```bash
|
||
# С мок-LLM (реальной БД и Celery в eager-режиме)
|
||
CELERY_ALWAYS_EAGER=True uv run pytest tests/test_summarize_task.py -v
|
||
```
|
||
|
||
### Ручная проверка с реальной моделью
|
||
|
||
1. Убедиться, что `llm` здоров:
|
||
```bash
|
||
curl http://localhost:8080/health
|
||
```
|
||
|
||
2. Включить `qwen_local` в конфиге, перезапустить worker
|
||
|
||
3. Запустить пайплайн на тестовой конференции:
|
||
```bash
|
||
# В тестовой конференции завершить запись
|
||
# Webhook room_finished → очередь run_pipeline
|
||
# → обработка пайплайна → видеть логи worker
|
||
|
||
docker compose -f deploy/docker-compose.yml --env-file .env logs -f worker
|
||
# Смотреть: transcribing → summarizing → ready to notify
|
||
```
|
||
|
||
4. Проверить БД:
|
||
```sql
|
||
SELECT summary_data FROM conference_sessions WHERE id = 'test-session-id';
|
||
-- Должно содержать структурированное резюме (JSON или plain text)
|
||
```
|
||
|
||
## Ссылки
|
||
|
||
- **Матрица уровней AI:** [ADR-004](../architecture/adr/004-ai-tier-matrix.md)
|
||
- **Установка LLM:** [docs/deploy/llm-setup.md](../deploy/llm-setup.md)
|
||
- **Контракты плагинов:** [docs/plugins/contracts.md](./contracts.md)
|
||
- **Пайплайн:** `workers/tasks/pipeline.py`, `workers/tasks/summarize.py`
|
||
- **Модель Qwen:** https://huggingface.co/Qwen/Qwen2.5-3B-Instruct
|