mirror of
https://github.com/zhouxiaoka/autoclip.git
synced 2026-10-02 02:34:34 +08:00
perf(studio): trim the fast-output chain and record model token usage
- Fast output runs the content pipeline without clustering (one unused model call) and without re-encoding every clip and collection (studio renders from the source); clip rows still sync from step 4 for the candidate list. - Packaging no longer asks the model to restate the transcript when no translation is needed; its rewrite was discarded for same-language output. - One face-detection pass per clip serves both the 4:3 interview window and the 9:16 podcast crop. - LLM inputs are compact JSON (indentation cost ~10 tokens per subtitle row). - Every text and vision call records tokens per project and stage into metadata/llm_usage.jsonl (DashScope native usage is captured now), so the cost of one video can be measured. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -428,6 +428,9 @@ class LLMManager:
|
||||
try:
|
||||
response = self.current_provider.call(prompt, input_data, **kwargs)
|
||||
content = response.content
|
||||
from backend.core import llm_usage
|
||||
llm_usage.record(response.model or getattr(self.current_provider, 'model_name', None), response.usage,
|
||||
prompt_chars=len(self.current_provider._build_full_input(prompt, input_data)), completion_chars=len(content or ''))
|
||||
if cache_path is not None:
|
||||
cache_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
cache_path.write_text(content, encoding="utf-8")
|
||||
|
||||
@@ -94,7 +94,8 @@ class LLMProvider(ABC):
|
||||
"""构建完整的输入"""
|
||||
if input_data:
|
||||
if isinstance(input_data, (dict, list, tuple)):
|
||||
return f"{prompt}\n\n输入内容:\n{json.dumps(input_data, ensure_ascii=False, indent=2, default=str)}"
|
||||
# Compact JSON: indentation adds ~10 tokens per subtitle row for no gain.
|
||||
return f"{prompt}\n\n输入内容:\n{json.dumps(input_data, ensure_ascii=False, default=str)}"
|
||||
else:
|
||||
return f"{prompt}\n\n输入内容:\n{input_data}"
|
||||
return prompt
|
||||
@@ -148,8 +149,10 @@ class DashScopeProvider(LLMProvider):
|
||||
del os.environ["DASHSCOPE_API_KEY"]
|
||||
if resp and getattr(resp, 'status_code', 200) == 200:
|
||||
if getattr(resp, 'output', None) and getattr(resp.output, 'text', None) is not None:
|
||||
usage = getattr(resp, 'usage', None)
|
||||
return LLMResponse(
|
||||
content=resp.output.text,
|
||||
usage={'input_tokens': usage.get('input_tokens'), 'output_tokens': usage.get('output_tokens')} if usage else None,
|
||||
model=self.model_name,
|
||||
finish_reason=getattr(resp.output, 'finish_reason', None)
|
||||
)
|
||||
|
||||
@@ -0,0 +1,103 @@
|
||||
"""Model token usage per project and stage, so the cost of one video can be measured.
|
||||
|
||||
`tracking(project_id)` binds a sink for the current context (one studio run); `stage(name)` names
|
||||
the step. `llm_manager.call` and the studio vision call record every response into the bound
|
||||
sink, appended to <project>/metadata/llm_usage.jsonl. Without a sink nothing is recorded.
|
||||
Providers that do not report tokens get an estimate from characters (`estimated: true`).
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import contextvars
|
||||
import json
|
||||
import threading
|
||||
import time
|
||||
from contextlib import contextmanager
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
_sink: contextvars.ContextVar[Path | None] = contextvars.ContextVar('llm_usage_sink', default=None)
|
||||
_stage: contextvars.ContextVar[str] = contextvars.ContextVar('llm_usage_stage', default='other')
|
||||
_lock = threading.Lock()
|
||||
FILE = 'llm_usage.jsonl'
|
||||
|
||||
|
||||
def usage_path(project_id: str) -> Path:
|
||||
from backend.core.path_utils import get_project_directory
|
||||
return get_project_directory(project_id) / 'metadata' / FILE
|
||||
|
||||
|
||||
@contextmanager
|
||||
def tracking(project_id: str):
|
||||
token = _sink.set(usage_path(project_id))
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
_sink.reset(token)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def stage(name: str):
|
||||
token = _stage.set(name)
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
_stage.reset(token)
|
||||
|
||||
|
||||
def set_stage(name: str) -> None:
|
||||
"""Name the step for the rest of a sequential run (the pipeline steps)."""
|
||||
_stage.set(name)
|
||||
|
||||
|
||||
def _estimate(chars: int) -> int:
|
||||
# Mixed Chinese/English transcripts: about 1.5 characters per token on Qwen/GPT tokenizers.
|
||||
return round(chars / 1.5)
|
||||
|
||||
|
||||
def record(model: str | None, usage: dict[str, Any] | None, *, prompt_chars: int, completion_chars: int, kind: str = 'text',
|
||||
images: int = 0) -> None:
|
||||
path = _sink.get()
|
||||
if path is None:
|
||||
return
|
||||
usage = usage or {}
|
||||
prompt = usage.get('prompt_tokens', usage.get('input_tokens'))
|
||||
completion = usage.get('completion_tokens', usage.get('output_tokens'))
|
||||
estimated = prompt is None or completion is None
|
||||
row = {'at': round(time.time(), 1), 'stage': _stage.get(), 'kind': kind, 'model': model or '',
|
||||
'prompt_tokens': int(prompt) if prompt is not None else _estimate(prompt_chars),
|
||||
'completion_tokens': int(completion) if completion is not None else _estimate(completion_chars),
|
||||
'prompt_chars': prompt_chars, 'completion_chars': completion_chars, 'images': images, 'estimated': estimated}
|
||||
try:
|
||||
with _lock:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
with path.open('a', encoding='utf-8') as handle:
|
||||
handle.write(json.dumps(row, ensure_ascii=False) + '\n')
|
||||
except OSError:
|
||||
pass # accounting never blocks output
|
||||
|
||||
|
||||
def summary(project_id: str) -> dict[str, Any]:
|
||||
"""Calls and tokens per stage plus totals for one project."""
|
||||
stages: dict[str, dict[str, Any]] = {}
|
||||
try:
|
||||
lines = usage_path(project_id).read_text(encoding='utf-8').splitlines()
|
||||
except OSError:
|
||||
lines = []
|
||||
for line in lines:
|
||||
try:
|
||||
row = json.loads(line)
|
||||
except ValueError:
|
||||
continue
|
||||
item = stages.setdefault(row.get('stage', 'other'), {'calls': 0, 'prompt_tokens': 0, 'completion_tokens': 0, 'estimated_calls': 0})
|
||||
item['calls'] += 1
|
||||
item['prompt_tokens'] += row.get('prompt_tokens') or 0
|
||||
item['completion_tokens'] += row.get('completion_tokens') or 0
|
||||
item['estimated_calls'] += bool(row.get('estimated'))
|
||||
total = {key: sum(item[key] for item in stages.values()) for key in ('calls', 'prompt_tokens', 'completion_tokens', 'estimated_calls')}
|
||||
return {'stages': stages, 'total': total}
|
||||
|
||||
|
||||
def run_in_context(fn):
|
||||
"""Wrap `fn` so a worker thread records into the caller's sink and stage."""
|
||||
context = contextvars.copy_context()
|
||||
return lambda *args, **kwargs: context.copy().run(fn, *args, **kwargs)
|
||||
@@ -153,7 +153,8 @@ class DataSyncService:
|
||||
project_dir / "step4_title" / "step4_title.json",
|
||||
project_dir / "step4_titles.json",
|
||||
project_dir / "clips_metadata.json",
|
||||
project_dir / "metadata" / "clips_metadata.json"
|
||||
project_dir / "metadata" / "clips_metadata.json",
|
||||
project_dir / "metadata" / "step4_titles.json", # fast output skips step 6
|
||||
]
|
||||
|
||||
clips_data = None
|
||||
|
||||
@@ -17,6 +17,7 @@ from backend.pipeline.step1_outline import run_step1_outline
|
||||
from backend.pipeline.step2_timeline import run_step2_timeline
|
||||
from backend.pipeline.step3_scoring import run_step3_scoring
|
||||
from backend.pipeline.step4_title import run_step4_title
|
||||
from backend.core import llm_usage
|
||||
from backend.pipeline.step5_clustering import run_step5_clustering
|
||||
from backend.pipeline.step6_video import run_step6_video
|
||||
|
||||
@@ -136,13 +137,15 @@ class SimplePipelineAdapter:
|
||||
f"没有可用的 LLM 提供商(当前选择:{name} · {model}),缺少 API Key 或本地服务地址。",
|
||||
)
|
||||
|
||||
async def process_project_sync(self, input_video_path: str, input_srt_path: str) -> Dict[str, Any]:
|
||||
async def process_project_sync(self, input_video_path: str, input_srt_path: str, clips_only: bool = False) -> Dict[str, Any]:
|
||||
"""
|
||||
同步处理项目 - 使用简化的进度系统
|
||||
|
||||
Args:
|
||||
input_video_path: 输入视频路径
|
||||
input_srt_path: 输入SRT路径
|
||||
clips_only: 快速出片只需要带标题的片段时间:跳过主题聚类(一次模型调用)与逐片段
|
||||
重新编码(Studio 直接从原片渲染),省下模型费用和大量 CPU。
|
||||
|
||||
Returns:
|
||||
处理结果
|
||||
@@ -191,6 +194,7 @@ class SimplePipelineAdapter:
|
||||
|
||||
# Step 1: 大纲提取(字幕为空 / 模型全部失败 / 无法解析时由 step1 自己抛 PipelineFailure)
|
||||
logger.info("执行Step 1: 大纲提取")
|
||||
llm_usage.set_stage("outline")
|
||||
outlines = run_step1_outline(srt_path, metadata_dir=metadata_dir, prompt_files=prompt_files)
|
||||
emit_progress(self.project_id, "SUBTITLE", "字幕处理完成", subpercent=50)
|
||||
|
||||
@@ -199,6 +203,7 @@ class SimplePipelineAdapter:
|
||||
|
||||
# Step 2: 时间线提取
|
||||
logger.info("执行Step 2: 时间线提取")
|
||||
llm_usage.set_stage("timeline")
|
||||
timeline_data = run_step2_timeline(
|
||||
metadata_dir / "step1_outline.json",
|
||||
metadata_dir=metadata_dir,
|
||||
@@ -212,6 +217,7 @@ class SimplePipelineAdapter:
|
||||
|
||||
# Step 3: 内容评分
|
||||
logger.info("执行Step 3: 内容评分")
|
||||
llm_usage.set_stage("scoring")
|
||||
scored_clips = run_step3_scoring(
|
||||
metadata_dir / "step2_timeline.json",
|
||||
metadata_dir=metadata_dir,
|
||||
@@ -231,6 +237,7 @@ class SimplePipelineAdapter:
|
||||
|
||||
# Step 4: 标题生成
|
||||
logger.info("执行Step 4: 标题生成")
|
||||
llm_usage.set_stage("titles")
|
||||
titled_clips = run_step4_title(
|
||||
metadata_dir / "step3_high_score_clips.json",
|
||||
metadata_dir=str(metadata_dir),
|
||||
@@ -238,36 +245,41 @@ class SimplePipelineAdapter:
|
||||
)
|
||||
emit_progress(self.project_id, "HIGHLIGHT", "标题生成完成", subpercent=40)
|
||||
|
||||
# Step 5: 主题聚类
|
||||
logger.info("执行Step 5: 主题聚类")
|
||||
collections = run_step5_clustering(
|
||||
metadata_dir / "step4_titles.json",
|
||||
metadata_dir=str(metadata_dir),
|
||||
prompt_files=prompt_files,
|
||||
)
|
||||
emit_progress(self.project_id, "HIGHLIGHT", "片段定位完成", subpercent=100)
|
||||
|
||||
# 阶段5: 视频导出
|
||||
emit_progress(self.project_id, "EXPORT", "开始视频导出")
|
||||
|
||||
# Step 6: 视频切割
|
||||
logger.info("执行Step 6: 视频切割")
|
||||
video_result = run_step6_video(
|
||||
metadata_dir / "step4_titles.json",
|
||||
metadata_dir / "step5_collections.json",
|
||||
input_video_path,
|
||||
output_dir=output_dir,
|
||||
clips_dir=str(clips_output_dir),
|
||||
collections_dir=str(collections_output_dir),
|
||||
metadata_dir=str(metadata_dir)
|
||||
)
|
||||
if titled_clips and not video_result.get("clips_generated"):
|
||||
raise PipelineFailure(
|
||||
"EXPORT",
|
||||
f"视频切割没有产出任何文件({len(titled_clips)} 个片段待切)。",
|
||||
HINT_CHECK_FFMPEG,
|
||||
if clips_only:
|
||||
collections, video_result = [], {}
|
||||
emit_progress(self.project_id, "HIGHLIGHT", "片段定位完成", subpercent=100)
|
||||
else:
|
||||
# Step 5: 主题聚类
|
||||
logger.info("执行Step 5: 主题聚类")
|
||||
llm_usage.set_stage("clustering")
|
||||
collections = run_step5_clustering(
|
||||
metadata_dir / "step4_titles.json",
|
||||
metadata_dir=str(metadata_dir),
|
||||
prompt_files=prompt_files,
|
||||
)
|
||||
emit_progress(self.project_id, "EXPORT", "视频导出完成", subpercent=100)
|
||||
emit_progress(self.project_id, "HIGHLIGHT", "片段定位完成", subpercent=100)
|
||||
|
||||
# 阶段5: 视频导出
|
||||
emit_progress(self.project_id, "EXPORT", "开始视频导出")
|
||||
|
||||
# Step 6: 视频切割
|
||||
logger.info("执行Step 6: 视频切割")
|
||||
video_result = run_step6_video(
|
||||
metadata_dir / "step4_titles.json",
|
||||
metadata_dir / "step5_collections.json",
|
||||
input_video_path,
|
||||
output_dir=output_dir,
|
||||
clips_dir=str(clips_output_dir),
|
||||
collections_dir=str(collections_output_dir),
|
||||
metadata_dir=str(metadata_dir)
|
||||
)
|
||||
if titled_clips and not video_result.get("clips_generated"):
|
||||
raise PipelineFailure(
|
||||
"EXPORT",
|
||||
f"视频切割没有产出任何文件({len(titled_clips)} 个片段待切)。",
|
||||
HINT_CHECK_FFMPEG,
|
||||
)
|
||||
emit_progress(self.project_id, "EXPORT", "视频导出完成", subpercent=100)
|
||||
|
||||
# 阶段6: 处理完成
|
||||
emit_progress(self.project_id, "DONE", "处理完成")
|
||||
|
||||
@@ -346,5 +346,6 @@ def refine_clips(rows: list[Row], clips: list[tuple[float, float]], call: Callab
|
||||
logger.warning('Boundary refinement fell back: %s', type(error).__name__)
|
||||
return fallback
|
||||
|
||||
from backend.core.llm_usage import run_in_context
|
||||
with ThreadPoolExecutor(max_workers=4, thread_name_prefix='studio-bounds') as pool:
|
||||
return list(pool.map(one, zip(clips, snapped)))
|
||||
return list(pool.map(run_in_context(one), zip(clips, snapped)))
|
||||
|
||||
@@ -346,26 +346,18 @@ def merge_points(points: list[dict[str, Any]], min_jump: float = MIN_JUMP) -> li
|
||||
return merged
|
||||
|
||||
|
||||
def auto_frame(video: Path, draft: Draft, source_w: int, source_h: int, *, window: tuple[int, int] | None = None) -> dict[str, Any]:
|
||||
"""Shot-aligned crop tracks per scene.
|
||||
|
||||
Scenes get a track only when somebody is visible somewhere in the draft; a clip with no
|
||||
people at all (gameplay, screen recording) keeps the layout the user chose untouched.
|
||||
`window` overrides the output size, e.g. the 4:3 window of the interview template.
|
||||
"""
|
||||
def scan_speakers(video: Path, scenes) -> list[dict[str, Any]]:
|
||||
"""Shots and sampled speaker centres per scene: the expensive part, independent of the output window."""
|
||||
if not is_installed():
|
||||
raise RuntimeError("人物识别组件未安装")
|
||||
out_w, out_h = window or {"portrait": (1080, 1920), "landscape": (1920, 1080)}.get(draft.aspect, (source_w, source_h))
|
||||
fraction = window_fraction(source_w, source_h, out_w, out_h)
|
||||
scenes = []
|
||||
scans = []
|
||||
with tempfile.TemporaryDirectory(prefix="ac-framing-") as temp:
|
||||
folder = Path(temp)
|
||||
for scene in draft.scenes:
|
||||
for scene in scenes:
|
||||
length = scene.end - scene.start
|
||||
shots = split_shots(length, detect_cuts(video, scene.start, length))
|
||||
points: list[dict[str, Any]] = []
|
||||
shots = []
|
||||
faces = grabbed = 0
|
||||
for index, (shot_start, shot_end) in enumerate(shots):
|
||||
for index, (shot_start, shot_end) in enumerate(split_shots(length, detect_cuts(video, scene.start, length))):
|
||||
samples: list[tuple[float, float | None]] = []
|
||||
for j, at in enumerate(sample_offsets(shot_start, shot_end)):
|
||||
pair = _grab_pair(video, scene.start + at, folder, f"{scene.id}-{index}-{j}")
|
||||
@@ -375,10 +367,32 @@ def auto_frame(video: Path, draft: Draft, source_w: int, source_h: int, *, windo
|
||||
center = _speaker_center(pair)
|
||||
faces += center is not None
|
||||
samples.append((at - shot_start, center))
|
||||
points += frame_shot(samples, shot_start, shot_end - shot_start, fraction)
|
||||
track = merge_points(points)
|
||||
scenes.append({"id": scene.id, "crop_x": track[0]["crop_x"], "crop_track": track, "faces": faces, "samples": grabbed,
|
||||
"shots": len(shots), "fit_shots": sum(p["mode"] == "fit" for p in track), "switches": len(track) - 1})
|
||||
shots.append((shot_start, shot_end, samples))
|
||||
scans.append({"shots": shots, "faces": faces, "samples": grabbed})
|
||||
return scans
|
||||
|
||||
|
||||
def auto_frame(video: Path, draft: Draft, source_w: int, source_h: int, *, window: tuple[int, int] | None = None,
|
||||
scans: list[dict[str, Any]] | None = None) -> dict[str, Any]:
|
||||
"""Shot-aligned crop tracks per scene.
|
||||
|
||||
Scenes get a track only when somebody is visible somewhere in the draft; a clip with no
|
||||
people at all (gameplay, screen recording) keeps the layout the user chose untouched.
|
||||
`window` overrides the output size, e.g. the 4:3 window of the interview template.
|
||||
`scans` reuses a `scan_speakers` pass, so several windows cost one detection.
|
||||
"""
|
||||
if scans is None:
|
||||
scans = scan_speakers(video, draft.scenes)
|
||||
out_w, out_h = window or {"portrait": (1080, 1920), "landscape": (1920, 1080)}.get(draft.aspect, (source_w, source_h))
|
||||
fraction = window_fraction(source_w, source_h, out_w, out_h)
|
||||
scenes = []
|
||||
for scene, scan in zip(draft.scenes, scans):
|
||||
points: list[dict[str, Any]] = []
|
||||
for shot_start, shot_end, samples in scan["shots"]:
|
||||
points += frame_shot(samples, shot_start, shot_end - shot_start, fraction)
|
||||
track = merge_points(points)
|
||||
scenes.append({"id": scene.id, "crop_x": track[0]["crop_x"], "crop_track": track, "faces": scan["faces"], "samples": scan["samples"],
|
||||
"shots": len(scan["shots"]), "fit_shots": sum(p["mode"] == "fit" for p in track), "switches": len(track) - 1})
|
||||
if not any(s["faces"] for s in scenes):
|
||||
for s in scenes:
|
||||
s.update(crop_x=None, crop_track=None, fit_shots=0, switches=0)
|
||||
|
||||
@@ -99,6 +99,11 @@ def vision_call(content, config=None):
|
||||
raise failure('connection', '视觉模型连接中断,请稍后重试;原素材已保留') from None
|
||||
except (ValueError, UnicodeError):
|
||||
raise failure('invalid_response', '视觉模型返回了无法解析的响应,请检查接口兼容性后重试') from None
|
||||
from backend.core import llm_usage
|
||||
llm_usage.record(model, result.get('usage') if isinstance(result, dict) else None, kind='vision',
|
||||
prompt_chars=sum(len(part.get('text', '')) for part in content if isinstance(part, dict)),
|
||||
completion_chars=len(str(((result.get('choices') or [{}])[0].get('message') or {}).get('content') or '')) if isinstance(result, dict) else 0,
|
||||
images=sum(1 for part in content if isinstance(part, dict) and part.get('type') == 'image_url'))
|
||||
try:
|
||||
choice = result['choices'][0]
|
||||
if choice.get('finish_reason') == 'length':
|
||||
|
||||
@@ -6,6 +6,7 @@ from copy import deepcopy
|
||||
import uuid
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from pathlib import Path
|
||||
from backend.core import llm_usage
|
||||
from backend.services.studio import audio, intelligence, store
|
||||
from backend.services.studio.models import Draft, Preferences, Scene
|
||||
from backend.services.studio.intelligence import analyze, make_drafts, VisionRequestError
|
||||
@@ -50,6 +51,20 @@ def export(project_id, draft, *, brand_outro=False):
|
||||
raise ValueError(message) from None
|
||||
return {k: v for k, v in added.items() if k not in ('instance', 'snapshot')}
|
||||
|
||||
def _tracked(stage):
|
||||
"""Record model token usage of a studio job (first argument: project id) under `stage`."""
|
||||
def wrap(fn):
|
||||
from functools import wraps
|
||||
|
||||
@wraps(fn)
|
||||
def run(project_id, *args, **kwargs):
|
||||
with llm_usage.tracking(project_id), llm_usage.stage(stage):
|
||||
return fn(project_id, *args, **kwargs)
|
||||
return run
|
||||
return wrap
|
||||
|
||||
|
||||
@_tracked('render')
|
||||
def _render(project_id, draft, job_id, *, brand_outro=False):
|
||||
started = monotonic()
|
||||
def update(**values):
|
||||
@@ -92,6 +107,7 @@ def mark_project(project_id, status, **config):
|
||||
p.completed_at = datetime.now(timezone.utc)
|
||||
db.commit()
|
||||
|
||||
@_tracked('visual_analysis')
|
||||
def _analyze(project_id, prefs, url, browser):
|
||||
try:
|
||||
mark_project(project_id, 'processing')
|
||||
@@ -215,6 +231,7 @@ def run_content(project_id, video):
|
||||
result = process_video_pipeline.apply(kwargs={
|
||||
'project_id': project_id, 'input_video_path': str(video),
|
||||
'input_srt_path': str(srt) if srt.exists() else None,
|
||||
'clips_only': True, # studio renders from the source: no clustering, no per-clip re-encode
|
||||
}, throw=True).get()
|
||||
if not result or not result.get('success'):
|
||||
message = (result or {}).get('error') or '内容切片未完成,请检查语音与文字模型设置后重试'
|
||||
@@ -283,7 +300,8 @@ def _complete_thought_bounds(project_id, clips):
|
||||
call = intelligence.text_json
|
||||
except Exception: # noqa: BLE001
|
||||
call = None
|
||||
return boundaries.refine_clips(rows, clips, call, boundaries.audio_silences(source(project_id)) if audio.has_audio(source(project_id)) else None)
|
||||
with llm_usage.stage('boundaries'):
|
||||
return boundaries.refine_clips(rows, clips, call, boundaries.audio_silences(source(project_id)) if audio.has_audio(source(project_id)) else None)
|
||||
|
||||
|
||||
def _content_drafts(project_id, plan, video):
|
||||
@@ -363,7 +381,7 @@ def _prepare_speaker_framing(platforms):
|
||||
logger.warning('Framing install could not start: %s', type(error).__name__)
|
||||
|
||||
|
||||
def _speaker_framing(value, video, *, window=None, wait_sec=120):
|
||||
def _speaker_framing(value, video, *, window=None, wait_sec=120, scans=None):
|
||||
"""Speaker-following crop tracks for a vertical draft (or the given output window).
|
||||
|
||||
Returns (scenes, framing): `speaker` when faces drive the crop, `full_frame` when the clip has
|
||||
@@ -379,7 +397,15 @@ def _speaker_framing(value, video, *, window=None, wait_sec=120):
|
||||
info = _probe(video)
|
||||
if not info.get('width') or not info.get('height'):
|
||||
return None, 'full_frame'
|
||||
result = framing.auto_frame(video, Draft.model_validate({**value, 'aspect': 'portrait', 'layout': 'crop'}), int(info['width']), int(info['height']), window=window)
|
||||
draft = Draft.model_validate({**value, 'aspect': 'portrait', 'layout': 'crop'})
|
||||
# Detection does not depend on the window: one pass serves the interview and podcast crops.
|
||||
key = tuple((scene.start, scene.end) for scene in draft.scenes)
|
||||
if scans is None or key not in scans:
|
||||
detected = framing.scan_speakers(video, draft.scenes)
|
||||
if scans is None:
|
||||
scans = {}
|
||||
scans[key] = detected
|
||||
result = framing.auto_frame(video, draft, int(info['width']), int(info['height']), window=window, scans=scans[key])
|
||||
tracks = {scene['id']: scene for scene in result['scenes']}
|
||||
if not any(scene.get('faces') for scene in result['scenes']):
|
||||
return None, 'full_frame'
|
||||
@@ -411,7 +437,7 @@ def _apply_framing(project_id, value, strategy_id, video, burned, cache):
|
||||
key = (tuple((scene['start'], scene['end']) for scene in value['scenes']), window)
|
||||
if key not in cache:
|
||||
try:
|
||||
cache[key] = _speaker_framing(value, video, window=window)
|
||||
cache[key] = _speaker_framing(value, video, window=window, scans=cache.setdefault('scans', {}))
|
||||
except Exception as error: # noqa: BLE001 - fall back to the full frame rather than fail output
|
||||
logger.warning('Speaker framing failed: %s', type(error).__name__)
|
||||
capture_studio_exception(error, 'auto_frame')
|
||||
@@ -447,8 +473,9 @@ def _apply_packaging(project_id, value, strategy_id, burned, cache):
|
||||
if key not in cache:
|
||||
lines = packaging.draft_lines(cache['entries'], value['scenes'])
|
||||
used = cache.setdefault('palettes', [])
|
||||
cache[key] = packaging.build_packaging(value, lines, strategy, burned=burned, known_names=cache['names'],
|
||||
avoid_palettes=tuple(used[-2:]))
|
||||
with llm_usage.stage('packaging'):
|
||||
cache[key] = packaging.build_packaging(value, lines, strategy, burned=burned, known_names=cache['names'],
|
||||
avoid_palettes=tuple(used[-2:]))
|
||||
if cache[key].get('palette'):
|
||||
used.append(cache[key]['palette'])
|
||||
return {**value, 'packaging': cache[key]}
|
||||
@@ -483,6 +510,7 @@ def _fit_platform_limit(project_id, value, strategy):
|
||||
return {**value, 'scenes': kept}, limit
|
||||
|
||||
|
||||
@_tracked('production')
|
||||
def _auto_generate(project_id, plan):
|
||||
"""Produce and render platform variants after the cheap screening pass."""
|
||||
started = monotonic()
|
||||
@@ -695,6 +723,7 @@ def inspect_project(project_id, options, url=None, browser=None):
|
||||
return state['analysis']['run_id']
|
||||
|
||||
|
||||
@_tracked('screening')
|
||||
def _inspect(project_id, options, url, browser):
|
||||
started = monotonic()
|
||||
try:
|
||||
@@ -783,6 +812,7 @@ def confirm_project(project_id, body):
|
||||
return state['analysis']['run_id']
|
||||
|
||||
|
||||
@_tracked('production')
|
||||
def _produce_selected(project_id, plan):
|
||||
from backend.services.studio import intelligence
|
||||
started = monotonic()
|
||||
|
||||
@@ -34,7 +34,8 @@ PROMPT = (
|
||||
'accent_line 是需要强调的那一行下标(0 或 1)。'
|
||||
'segments 把 lines 按完整句子重新分段:from/to 是连续的行 id 区间,按顺序首尾相接、覆盖全部行、不重叠;'
|
||||
'每段只含 1–2 句话、最多覆盖 4 行,不要把大段内容合成一段;'
|
||||
'translate 为 true 时 text 是该段翻译成 audience_language 的口语化译文(去掉口头禅,不添加事实),为 false 时 text 是整理后的原文。'
|
||||
'translate 为 true 时 text 是该段翻译成 audience_language 的口语化译文(去掉口头禅,不添加事实);'
|
||||
'translate 为 false 时字幕直接用原文,segments 返回空数组 [],不要复述原文。'
|
||||
'speakers 只填写在 lines 或 known_names 中明确出现过的人名,role 写其公开身份(不确定就留空),line 是此人第一次说话的行;'
|
||||
'不确定就返回空数组,绝不猜测身份。'
|
||||
'tags 仅当 template 为 interview_zh 时给出 2–4 个编辑点评(中文,每个不超过 10 个字),必须具体点出这一句最有冲击力的内容,'
|
||||
|
||||
@@ -67,6 +67,7 @@ def process_video_pipeline(
|
||||
project_id: str,
|
||||
input_video_path: str,
|
||||
input_srt_path: Optional[str] = None,
|
||||
clips_only: bool = False,
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
处理视频流水线任务 - 使用Pipeline适配器
|
||||
@@ -127,7 +128,7 @@ def process_video_pipeline(
|
||||
|
||||
# 执行Pipeline处理 - 使用异步包装器
|
||||
import asyncio
|
||||
result = asyncio.run(pipeline_adapter.process_project_sync(input_video_path, input_srt_path))
|
||||
result = asyncio.run(pipeline_adapter.process_project_sync(input_video_path, input_srt_path, clips_only=clips_only))
|
||||
|
||||
with session_scope() as db:
|
||||
task = db.query(Task).filter(Task.id == task_row_id).first()
|
||||
|
||||
@@ -24,7 +24,7 @@ def test_landscape_outputs_are_not_reframed(monkeypatch):
|
||||
|
||||
def test_faces_give_speaker_crop_and_one_detection_per_window_shape(monkeypatch):
|
||||
calls = []
|
||||
def detect(value, video, *, window=None):
|
||||
def detect(value, video, *, window=None, scans=None):
|
||||
calls.append(window)
|
||||
return [{**scene, 'crop_x': .3, 'crop_track': [{'start': 0, 'crop_x': .3, 'mode': 'crop'}]} for scene in value['scenes']], 'speaker'
|
||||
monkeypatch.setattr(jobs, '_speaker_framing', detect)
|
||||
@@ -37,6 +37,20 @@ def test_faces_give_speaker_crop_and_one_detection_per_window_shape(monkeypatch)
|
||||
assert calls == [jobs.INTERVIEW_WINDOW, None] # 4:3 interview window once, 9:16 once
|
||||
|
||||
|
||||
def test_both_window_shapes_share_one_face_detection_pass(monkeypatch):
|
||||
from backend.services import publish_export
|
||||
from backend.services.studio import framing
|
||||
scans = []
|
||||
monkeypatch.setattr(framing, 'is_installed', lambda: True)
|
||||
monkeypatch.setattr(publish_export, '_probe', lambda _v: {'width': 1920, 'height': 1080})
|
||||
monkeypatch.setattr(framing, 'scan_speakers', lambda _v, scenes: scans.append(1) or [{'shots': [(0.0, 10.0, [(1.0, .3), (5.0, .3)])], 'faces': 2, 'samples': 2} for _ in scenes])
|
||||
cache = {}
|
||||
for strategy_id in ('douyin', 'tiktok'):
|
||||
_, framed = jobs._apply_framing('p1', _value(), strategy_id, 'video.mp4', False, cache)
|
||||
assert framed == 'speaker'
|
||||
assert len(scans) == 1
|
||||
|
||||
|
||||
@pytest.mark.parametrize('strategy_id', ['douyin', 'tiktok', 'instagram_reels', 'youtube_shorts', 'youtube_long', 'bilibili', 'xiaohongshu', 'original'])
|
||||
def test_redirecting_a_derived_draft_to_any_platform_keeps_a_valid_title_version(strategy_id):
|
||||
derived = {**_value(), 'title_style': 'comic', 'title_template_version': 6} # e.g. an earlier Douyin variant
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
"""Token usage is recorded per project and stage, including worker threads (no model calls)."""
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
|
||||
from backend.core import llm_usage
|
||||
|
||||
|
||||
def test_usage_is_recorded_per_stage_and_summarised(monkeypatch, tmp_path):
|
||||
monkeypatch.setattr(llm_usage, 'usage_path', lambda project_id: tmp_path / project_id / 'metadata' / llm_usage.FILE)
|
||||
llm_usage.record('qwen-plus', {'input_tokens': 10}, prompt_chars=30, completion_chars=3) # no sink: ignored
|
||||
with llm_usage.tracking('p1'):
|
||||
with llm_usage.stage('packaging'):
|
||||
llm_usage.record('qwen-plus', {'input_tokens': 1200, 'output_tokens': 300}, prompt_chars=2000, completion_chars=500)
|
||||
llm_usage.set_stage('outline')
|
||||
llm_usage.record('qwen-plus', None, prompt_chars=1500, completion_chars=150) # provider without usage
|
||||
with ThreadPoolExecutor(max_workers=2) as pool:
|
||||
with llm_usage.stage('boundaries'):
|
||||
record = llm_usage.run_in_context(lambda: llm_usage.record('qwen-plus', {'prompt_tokens': 50, 'completion_tokens': 5},
|
||||
prompt_chars=80, completion_chars=8))
|
||||
list(pool.map(lambda _: record(), range(2)))
|
||||
result = llm_usage.summary('p1')
|
||||
assert result['stages']['packaging'] == {'calls': 1, 'prompt_tokens': 1200, 'completion_tokens': 300, 'estimated_calls': 0}
|
||||
assert result['stages']['outline'] == {'calls': 1, 'prompt_tokens': 1000, 'completion_tokens': 100, 'estimated_calls': 1}
|
||||
assert result['stages']['boundaries']['calls'] == 2
|
||||
assert result['total']['calls'] == 4
|
||||
|
||||
|
||||
def test_manager_calls_record_provider_usage(monkeypatch, tmp_path):
|
||||
from backend.core import llm_manager
|
||||
from backend.core.llm_providers import LLMResponse
|
||||
|
||||
class Provider:
|
||||
model_name = 'qwen-plus'
|
||||
|
||||
def _build_full_input(self, prompt, data):
|
||||
return prompt + str(data)
|
||||
|
||||
def call(self, prompt, data, **_):
|
||||
return LLMResponse(content='{"ok": true}', usage={'input_tokens': 42, 'output_tokens': 7}, model='qwen-plus')
|
||||
|
||||
manager = llm_manager.LLMManager.__new__(llm_manager.LLMManager)
|
||||
manager.current_provider = Provider()
|
||||
monkeypatch.setattr(manager, '_reload_if_settings_changed', lambda: None, raising=False)
|
||||
monkeypatch.setattr(llm_manager, '_llm_cache_path', lambda *_: None)
|
||||
monkeypatch.setattr(llm_usage, 'usage_path', lambda project_id: tmp_path / llm_usage.FILE)
|
||||
with llm_usage.tracking('p1'), llm_usage.stage('scoring'):
|
||||
assert manager.call('score', {'x': 1}) == '{"ok": true}'
|
||||
assert llm_usage.summary('p1')['stages']['scoring']['prompt_tokens'] == 42
|
||||
@@ -317,6 +317,26 @@ def test_adapter_fails_when_ffmpeg_produced_no_clip(adapter, monkeypatch, tmp_pa
|
||||
assert "ffmpeg" in result["error"]
|
||||
|
||||
|
||||
def test_fast_output_skips_clustering_and_clip_encoding(adapter, monkeypatch, tmp_path):
|
||||
import pytest
|
||||
from backend.services import simple_pipeline_adapter as mod
|
||||
|
||||
_fake_manager(monkeypatch, available=True)
|
||||
srt = tmp_path / "in.srt"
|
||||
srt.write_text(SRT, encoding="utf-8")
|
||||
monkeypatch.setattr(mod, "run_step1_outline", lambda *a, **k: [{"title": "t"}])
|
||||
monkeypatch.setattr(mod, "run_step2_timeline", lambda *a, **k: [{"id": 1}])
|
||||
monkeypatch.setattr(mod, "run_step3_scoring", lambda *a, **k: [{"id": 1, "final_score": 0.9}])
|
||||
monkeypatch.setattr(mod, "run_step4_title", lambda *a, **k: [{"id": 1, "generated_title": "x"}])
|
||||
monkeypatch.setattr(mod, "run_step5_clustering", lambda *a, **k: pytest.fail("studio never reads collections"))
|
||||
monkeypatch.setattr(mod, "run_step6_video", lambda *a, **k: pytest.fail("studio renders from the source"))
|
||||
|
||||
result = asyncio.run(adapter.process_project_sync(str(tmp_path / "in.mp4"), str(srt), clips_only=True))
|
||||
|
||||
assert result["status"] == "succeeded"
|
||||
assert result["result"]["titled_clips"] == [{"id": 1, "generated_title": "x"}]
|
||||
|
||||
|
||||
def test_adapter_happy_path_still_succeeds(adapter, monkeypatch, tmp_path):
|
||||
from backend.services import simple_pipeline_adapter as mod
|
||||
|
||||
|
||||
Reference in New Issue
Block a user