204 lines
11 KiB
Python
204 lines
11 KiB
Python
"""Плагин `Summarizer` на локальной модели семейства Qwen через сервер llama.cpp.
|
||
|
||
Конкретная модель/квант не зашиты в плагине — их задаёт конфигурация
|
||
(`config/plugins.yaml` либо `TierSpec` в `services/ai_tiers.py`, ADR-004);
|
||
дефолт конструктора (`qwen2.5-3b-instruct-q4_k_m`) — только фолбэк на случай
|
||
прямого создания плагина без конфига.
|
||
|
||
Map-reduce целиком инкапсулирован в плагине (контракт `Summarizer.summarize`
|
||
не меняется — ТЗ §1.3): транскрипт делится на чанки
|
||
чистой функцией `chunk_transcript`, каждый чанк резюмируется отдельным
|
||
вызовом LLM (map), частичные резюме объединяются одним reduce-вызовом; если
|
||
частичные резюме суммарно не влезают в бюджет токенов запроса — reduce
|
||
выполняется иерархически, группами, пока не останется одно резюме.
|
||
`max_tokens_map`/`max_tokens_reduce` — раздельные per-tier лимиты генерации
|
||
(ADR-004: reduce всегда ≥ map — 1024 токенов на reduce не хватает).
|
||
|
||
Тексты промптов (`workers/summarizer/prompts/summary_map_ru.txt`,
|
||
`summary_reduce_ru.txt`) утверждены и не меняются в коде — загружаются
|
||
лениво из файлов. Подстановка плейсхолдера — через `str.replace`, а не
|
||
`str.format`: промпты содержат разметку формата вывода (`[что решено] —
|
||
ответственный: [имя]` и т.п.) с квадратными, но потенциально и фигурными
|
||
скобками в будущих правках текста — `str.format` на них падает с
|
||
`KeyError`/`IndexError`, тогда как `str.replace` нечувствителен к остальному
|
||
содержимому файла.
|
||
"""
|
||
|
||
from pathlib import Path
|
||
from typing import Any, ClassVar
|
||
|
||
from core.plugins.factory import register_summarizer
|
||
from core.plugins.summarizer import Summarizer
|
||
from core.summarization.chunking import chunk_transcript
|
||
from core.summarization.llm_client import OpenAICompatClient
|
||
from core.summarization.tokens import QwenTokenCounter
|
||
|
||
_DEFAULT_PROMPTS_DIR = "workers/summarizer/prompts"
|
||
_MAP_PROMPT_FILE = "summary_map_ru.txt"
|
||
_REDUCE_PROMPT_FILE = "summary_reduce_ru.txt"
|
||
_MAP_PLACEHOLDER = "{transcript_chunk}"
|
||
_REDUCE_PLACEHOLDER = "{partial_summaries}"
|
||
|
||
_REDUCE_BUDGET_TOKENS = 6000
|
||
"""Бюджет токенов на один reduce-вызов (частичные резюме + шаблон промпта);
|
||
меньше `max_chunk_tokens` чанкера — запас под текст самого reduce-промпта и
|
||
вывод модели в общем контексте (CTX_SIZE=16384)."""
|
||
|
||
|
||
@register_summarizer
|
||
class QwenLocal(Summarizer):
|
||
"""Summarizer на Qwen2.5-3B-Instruct через OpenAI-совместимый сервер llama.cpp."""
|
||
|
||
provider: ClassVar[str] = "qwen_local"
|
||
|
||
def __init__(
|
||
self,
|
||
model: str | None = None,
|
||
chunk_minutes: int = 20,
|
||
base_url: str = "http://llm:8080/v1",
|
||
tokenizer_path: str = "/models/qwen/tokenizer.json",
|
||
prompts_dir: str = _DEFAULT_PROMPTS_DIR,
|
||
temperature: float = 0.2,
|
||
max_tokens: int = 1024,
|
||
max_tokens_map: int | None = None,
|
||
max_tokens_reduce: int | None = None,
|
||
**options: Any,
|
||
) -> None:
|
||
self.model = model or "qwen2.5-3b-instruct-q4_k_m"
|
||
self.chunk_minutes = chunk_minutes
|
||
self.base_url = base_url
|
||
self.tokenizer_path = tokenizer_path
|
||
self.prompts_dir = prompts_dir
|
||
self.temperature = temperature
|
||
self.max_tokens = max_tokens
|
||
# Раздельные лимиты map/reduce (ADR-004, per-tier параметры генерации);
|
||
# без явного значения оба используют общий `max_tokens` — обратная
|
||
# совместимость со старым форматом конфигурации.
|
||
self.max_tokens_map = max_tokens_map if max_tokens_map is not None else max_tokens
|
||
self.max_tokens_reduce = max_tokens_reduce if max_tokens_reduce is not None else max_tokens
|
||
self.options = options
|
||
|
||
self._count_tokens = QwenTokenCounter(tokenizer_path)
|
||
self._client: OpenAICompatClient | None = None
|
||
self._map_prompt: str | None = None
|
||
self._reduce_prompt: str | None = None
|
||
|
||
def _get_client(self) -> OpenAICompatClient:
|
||
"""Лениво создать HTTP-клиент LLM (переиспользуется в рамках инстанса плагина)."""
|
||
if self._client is None:
|
||
self._client = OpenAICompatClient(
|
||
base_url=self.base_url,
|
||
model=self.model,
|
||
temperature=self.temperature,
|
||
max_tokens=self.max_tokens,
|
||
**self.options,
|
||
)
|
||
return self._client
|
||
|
||
def close(self) -> None:
|
||
"""Закрыть HTTP-клиент LLM, если он был лениво создан (освободить пул соединений).
|
||
|
||
Безопасно вызывать многократно и до первого использования — если
|
||
клиент ни разу не создавался, ничего не делает.
|
||
"""
|
||
if self._client is not None:
|
||
self._client.close()
|
||
self._client = None
|
||
|
||
def _load_prompt(self, filename: str) -> str:
|
||
"""Прочитать текст промпта из `prompts_dir` (без изменений, как есть на диске)."""
|
||
return (Path(self.prompts_dir) / filename).read_text(encoding="utf-8")
|
||
|
||
def _map_prompt_template(self) -> str:
|
||
if self._map_prompt is None:
|
||
self._map_prompt = self._load_prompt(_MAP_PROMPT_FILE)
|
||
return self._map_prompt
|
||
|
||
def _reduce_prompt_template(self) -> str:
|
||
if self._reduce_prompt is None:
|
||
self._reduce_prompt = self._load_prompt(_REDUCE_PROMPT_FILE)
|
||
return self._reduce_prompt
|
||
|
||
def _map_chunk(self, chunk: str) -> str:
|
||
"""Выполнить map-вызов LLM для одного чанка транскрипта."""
|
||
prompt = self._map_prompt_template().replace(_MAP_PLACEHOLDER, chunk)
|
||
return self._get_client().complete(prompt, max_tokens=self.max_tokens_map)
|
||
|
||
def _reduce_once(self, summaries: list[str]) -> str:
|
||
"""Выполнить один reduce-вызов LLM над группой частичных резюме."""
|
||
joined = "\n\n".join(summaries)
|
||
prompt = self._reduce_prompt_template().replace(_REDUCE_PLACEHOLDER, joined)
|
||
return self._get_client().complete(prompt, max_tokens=self.max_tokens_reduce)
|
||
|
||
def _group_by_token_budget(self, summaries: list[str], budget: int) -> list[list[str]]:
|
||
"""Жадно сгруппировать резюме так, чтобы каждая группа влезала в `budget` токенов."""
|
||
groups: list[list[str]] = []
|
||
current: list[str] = []
|
||
current_tokens = 0
|
||
for summary in summaries:
|
||
tokens = self._count_tokens(summary)
|
||
if current and current_tokens + tokens > budget:
|
||
groups.append(current)
|
||
current = []
|
||
current_tokens = 0
|
||
current.append(summary)
|
||
current_tokens += tokens
|
||
if current:
|
||
groups.append(current)
|
||
return groups
|
||
|
||
def _reduce(self, partial_summaries: list[str]) -> str:
|
||
"""Свести частичные резюме к одному, иерархически группами при переполнении бюджета.
|
||
|
||
Группировка по токенам (`_group_by_token_budget`) не гарантирует
|
||
прогресс, если отдельные частичные резюме сами не помещаются в
|
||
`_REDUCE_BUDGET_TOKENS` (например, при неудачно большом `max_tokens`
|
||
в конфиге плагина) — тогда она вырождается в список синглтон-групп,
|
||
и список резюме не сокращается. В этом случае принудительно сводим
|
||
резюме попарно: длина списка минимум делится пополам на каждой
|
||
итерации, что гарантирует завершение цикла за конечное число шагов.
|
||
"""
|
||
summaries = partial_summaries
|
||
while len(summaries) > 1:
|
||
joined_tokens = self._count_tokens("\n\n".join(summaries))
|
||
if joined_tokens <= _REDUCE_BUDGET_TOKENS:
|
||
return self._reduce_once(summaries)
|
||
|
||
groups = self._group_by_token_budget(summaries, _REDUCE_BUDGET_TOKENS)
|
||
if len(groups) >= len(summaries):
|
||
# Группировка по бюджету не уменьшила число групп (каждое
|
||
# резюме — уже отдельная группа) — гарантируем прогресс
|
||
# принудительным объединением попарно.
|
||
groups = [summaries[i : i + 2] for i in range(0, len(summaries), 2)]
|
||
summaries = [self._reduce_once(group) for group in groups]
|
||
return summaries[0]
|
||
|
||
def summarize(self, transcript: str) -> str:
|
||
"""Построить резюме транскрипта: map по чанкам, затем reduce до одного текста.
|
||
|
||
Пустой транскрипт (пустой список чанков) — пустая строка без вызовов
|
||
LLM. Единственный чанк — map-результат уже соответствует формату
|
||
reduce-вывода, дополнительный reduce-вызов не требуется.
|
||
|
||
HTTP-клиент LLM (если он был создан) закрывается по завершении вызова
|
||
независимо от исхода — плагин инстанцируется на одну задачу
|
||
суммаризации (см. `create_summarizer` в фабрике), поэтому держать
|
||
пул соединений открытым дольше одного вызова `summarize` не нужно.
|
||
"""
|
||
try:
|
||
chunks = chunk_transcript(
|
||
transcript,
|
||
self._count_tokens,
|
||
target_chunk_minutes=self.chunk_minutes,
|
||
)
|
||
if not chunks:
|
||
return ""
|
||
|
||
partial_summaries = [self._map_chunk(chunk) for chunk in chunks]
|
||
if len(partial_summaries) == 1:
|
||
return partial_summaries[0]
|
||
|
||
return self._reduce(partial_summaries)
|
||
finally:
|
||
self.close()
|