Files
LLM-infa/api/lib/processing.py
jze9 e87f6cd215 Фоновые задачи: плейлисты, длинные лекции, вкладка настроек
- таблица jobs + очередь с воркером: POST /jobs сразу отвечает, элементы
  плейлиста обрабатываются последовательно, есть отмена и живой прогресс
  (скачивание %, минуты распознанного Vosk-аудио)
- expand_url: плейлист раскрывается в список видео
- keep_video=false теперь качает только аудио (для 4-5ч лекций)
- настройки (LLM-ключ, модель, длина выжимки, хранение видео, путь Vosk)
  хранятся в Postgres: переживают перезапуск и перезагрузку страницы
- длина выжимки прокидывается в LLM-промпт ({n} предложений)
- UI: вкладки «Главная»/«Настройки», карточка задач с автообновлением
  каждые 3 сек, прогресс-бар, отмена, библиотека обновляется по мере
  готовности элементов

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-14 15:01:11 +05:00

230 lines
8.2 KiB
Python
Raw Permalink 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.
"""Ядро пайплайна: одно видео от URL до записи в БД.
Используется и синхронным роутом /pipeline/process, и фоновыми задачами /jobs.
Все дефолты (LLM, длина выжимки, хранение видео, Vosk-модель) берутся из
настроек в БД (таблица settings) с фолбэком на переменные окружения.
"""
from __future__ import annotations
import asyncio
import logging
import uuid as _uuid
from pathlib import Path
from typing import Callable, Optional
from api.config import settings
from api.db.setting import get_all_settings
from api.db.video import create_video
from api.lib.llm import LLMConfig
from api.lib.summarize import summarize
from api.lib.transcribe import transcribe_via_ffmpeg
from api.lib.youtube import download_video_with_subs, vtt_to_text
logger = logging.getLogger("processing")
_PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent
_VIDEO_DIR = _PROJECT_ROOT / "data" / "videos"
_TEXT_DIR = _PROJECT_ROOT / "data" / "text"
_MODELS_DIR = _PROJECT_ROOT / "models"
class PipelineError(RuntimeError):
def __init__(self, message: str, status: int = 500):
super().__init__(message)
self.status = status
async def effective_defaults() -> dict:
"""Настройки из БД поверх переменных окружения."""
db = await get_all_settings()
base_url = db.get("llm_base_url") or settings.LLM_BASE_URL
model = db.get("llm_model") or settings.LLM_MODEL
llm_cfg = None
if base_url and model:
llm_cfg = LLMConfig(
base_url=base_url,
api_key=db.get("llm_api_key") or settings.LLM_API_KEY,
model=model,
timeout=settings.LLM_TIMEOUT,
max_tokens=settings.LLM_MAX_TOKENS,
temperature=settings.LLM_TEMPERATURE,
)
try:
max_sentences = int(db.get("summary_max_sentences") or 15)
except ValueError:
max_sentences = 15
return {
"llm_cfg": llm_cfg,
"summary_max_sentences": max_sentences,
"keep_video": db.get("keep_video", "1") != "0",
"model_path": db.get("vosk_model_path") or settings.DEFAULT_VOSK_MODEL,
}
def resolve_model_path(raw: str) -> Path:
p = Path(raw)
if not p.is_absolute():
p = _PROJECT_ROOT / raw
p = p.resolve()
try:
p.relative_to(_MODELS_DIR.resolve())
except ValueError:
raise PipelineError(
"model_path must point inside the project's models directory", status=400
)
if not p.exists():
raise PipelineError(f"Model not found: {p}", status=400)
return p
async def process_single_video(
url: str,
*,
model_path: str | None = None,
summary_max_sentences: int | None = None,
keep_video: bool | None = None,
section_id: int | None = None,
job_id: str | None = None,
llm_cfg: Optional[LLMConfig] = None,
system_prompt: str | None = None,
user_template: str | None = None,
progress: Optional[Callable[[str], None]] = None,
) -> dict:
"""Скачивает, расшифровывает, сокращает и сохраняет одно видео.
Параметры со значением None подтягиваются из настроек (БД/env).
progress(msg) — живой статус для UI задач.
"""
notify = progress or (lambda msg: None)
defaults = await effective_defaults()
if llm_cfg is None:
llm_cfg = defaults["llm_cfg"]
if summary_max_sentences is None:
summary_max_sentences = defaults["summary_max_sentences"]
if keep_video is None:
keep_video = defaults["keep_video"]
raw_model_path = model_path or defaults["model_path"]
video_uuid = _uuid.uuid4().hex
_VIDEO_DIR.mkdir(parents=True, exist_ok=True)
_TEXT_DIR.mkdir(parents=True, exist_ok=True)
loop = asyncio.get_running_loop()
logger.info("[%s] downloading %s (keep_video=%s)", video_uuid, url, keep_video)
notify("скачивание…")
try:
info = await loop.run_in_executor(
None,
lambda: download_video_with_subs(
url,
video_uuid,
_VIDEO_DIR,
settings.subtitle_langs_list,
# видео не сохраняем — достаточно аудио (для лекций в разы быстрее)
audio_only=not keep_video,
on_progress=lambda pct: notify(f"скачивание {pct}%"),
),
)
except Exception as e:
raise PipelineError(f"Download failed: {e}")
title: str = info["title"]
video_path: Path = info["video_path"]
sub_path: Path | None = info["subtitle_path"]
if not video_path.exists():
raise PipelineError("Video file missing after download")
full_text = ""
method = ""
if sub_path is not None and sub_path.exists():
try:
content = sub_path.read_text(encoding="utf-8", errors="ignore")
full_text = vtt_to_text(content)
if full_text.strip():
method = f"subtitles:{info.get('subtitle_kind') or 'auto'}:{info.get('subtitle_lang') or 'unknown'}"
logger.info("[%s] using subtitles (%s)", video_uuid, method)
except Exception as e:
logger.warning("[%s] subtitle parse failed: %s", video_uuid, e)
full_text = ""
if not full_text.strip():
logger.info("[%s] no usable subtitles, falling back to Vosk", video_uuid)
notify("распознавание речи (Vosk)…")
resolved_model = resolve_model_path(raw_model_path)
try:
full_text = await loop.run_in_executor(
None,
lambda: transcribe_via_ffmpeg(
video_path,
resolved_model,
on_progress=lambda m: notify(f"распознано {m} мин аудио"),
),
)
method = f"vosk:{resolved_model.name}"
except Exception as e:
raise PipelineError(f"Transcription failed: {e}")
if not full_text.strip():
raise PipelineError("No subtitles and transcription returned empty text")
text_full_path = _TEXT_DIR / f"{video_uuid}.txt"
text_summary_path = _TEXT_DIR / f"{video_uuid}_summary.txt"
text_full_path.write_text(full_text, encoding="utf-8")
notify("выжимка (LLM)…" if llm_cfg else "выжимка…")
try:
summary = await loop.run_in_executor(
None,
lambda: summarize(
full_text,
llm_cfg=llm_cfg,
system_prompt=system_prompt,
user_template=user_template,
max_sentences=summary_max_sentences,
),
)
except Exception as e:
raise PipelineError(f"Summarization failed: {e}")
text_summary_path.write_text(summary, encoding="utf-8")
if keep_video:
rel_video = str(video_path.relative_to(_PROJECT_ROOT))
else:
# пользователь просил не хранить видео: текст уже извлечён, файлы не нужны
video_path.unlink(missing_ok=True)
for p in _VIDEO_DIR.glob(f"{video_uuid}*.vtt"):
p.unlink(missing_ok=True)
rel_video = ""
rel_full = str(text_full_path.relative_to(_PROJECT_ROOT))
rel_summary = str(text_summary_path.relative_to(_PROJECT_ROOT))
try:
await create_video(
uuid=video_uuid,
source_url=url,
title=title,
video_path=rel_video,
text_full_path=rel_full,
text_summary_path=rel_summary,
transcription_method=method,
section_id=section_id,
job_id=job_id,
)
except Exception as e:
raise PipelineError(f"DB insert failed: {e}")
logger.info("[%s] done (method=%s)", video_uuid, method)
return {
"uuid": video_uuid,
"title": title,
"source_url": url,
"video_path": rel_video,
"text_full_path": rel_full,
"text_summary_path": rel_summary,
"transcription_method": method,
"summary_preview": summary[:300],
}