From 628b01b77e4e6a2225e280c9ee85a2014f7ed0fc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=91=A8=E5=B0=8F=E8=88=9F?= Date: Thu, 1 Oct 2026 11:47:50 +0800 Subject: [PATCH] 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) --- backend/core/llm_manager.py | 3 + backend/core/llm_providers.py | 5 +- backend/core/llm_usage.py | 103 ++++++++++++++++++++ backend/services/data_sync_service.py | 3 +- backend/services/simple_pipeline_adapter.py | 72 ++++++++------ backend/services/studio/boundaries.py | 3 +- backend/services/studio/framing.py | 50 ++++++---- backend/services/studio/intelligence.py | 5 + backend/services/studio/jobs.py | 42 ++++++-- backend/services/studio/packaging.py | 3 +- backend/tasks/processing.py | 3 +- backend/tests/test_auto_framing.py | 16 ++- backend/tests/test_llm_usage.py | 47 +++++++++ backend/tests/test_pipeline_failures.py | 20 ++++ 14 files changed, 315 insertions(+), 60 deletions(-) create mode 100644 backend/core/llm_usage.py create mode 100644 backend/tests/test_llm_usage.py diff --git a/backend/core/llm_manager.py b/backend/core/llm_manager.py index 111b78ee..bf2b05d9 100644 --- a/backend/core/llm_manager.py +++ b/backend/core/llm_manager.py @@ -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") diff --git a/backend/core/llm_providers.py b/backend/core/llm_providers.py index b09dee9c..571ec537 100644 --- a/backend/core/llm_providers.py +++ b/backend/core/llm_providers.py @@ -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) ) diff --git a/backend/core/llm_usage.py b/backend/core/llm_usage.py new file mode 100644 index 00000000..e0cd2188 --- /dev/null +++ b/backend/core/llm_usage.py @@ -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 /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) diff --git a/backend/services/data_sync_service.py b/backend/services/data_sync_service.py index b3b55a6b..4df26593 100644 --- a/backend/services/data_sync_service.py +++ b/backend/services/data_sync_service.py @@ -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 diff --git a/backend/services/simple_pipeline_adapter.py b/backend/services/simple_pipeline_adapter.py index 12a79410..b95d983b 100644 --- a/backend/services/simple_pipeline_adapter.py +++ b/backend/services/simple_pipeline_adapter.py @@ -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", "处理完成") diff --git a/backend/services/studio/boundaries.py b/backend/services/studio/boundaries.py index 1517725b..1450d463 100644 --- a/backend/services/studio/boundaries.py +++ b/backend/services/studio/boundaries.py @@ -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))) diff --git a/backend/services/studio/framing.py b/backend/services/studio/framing.py index fad8f65e..3df94d07 100644 --- a/backend/services/studio/framing.py +++ b/backend/services/studio/framing.py @@ -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) diff --git a/backend/services/studio/intelligence.py b/backend/services/studio/intelligence.py index 25a108ae..b599beaa 100644 --- a/backend/services/studio/intelligence.py +++ b/backend/services/studio/intelligence.py @@ -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': diff --git a/backend/services/studio/jobs.py b/backend/services/studio/jobs.py index cebf5d6b..84469106 100644 --- a/backend/services/studio/jobs.py +++ b/backend/services/studio/jobs.py @@ -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() diff --git a/backend/services/studio/packaging.py b/backend/services/studio/packaging.py index 6e2aede1..100c74e1 100644 --- a/backend/services/studio/packaging.py +++ b/backend/services/studio/packaging.py @@ -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 个字),必须具体点出这一句最有冲击力的内容,' diff --git a/backend/tasks/processing.py b/backend/tasks/processing.py index fb8368cf..c6ee6483 100644 --- a/backend/tasks/processing.py +++ b/backend/tasks/processing.py @@ -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() diff --git a/backend/tests/test_auto_framing.py b/backend/tests/test_auto_framing.py index 71ff2456..3af30fed 100644 --- a/backend/tests/test_auto_framing.py +++ b/backend/tests/test_auto_framing.py @@ -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 diff --git a/backend/tests/test_llm_usage.py b/backend/tests/test_llm_usage.py new file mode 100644 index 00000000..b14b4faf --- /dev/null +++ b/backend/tests/test_llm_usage.py @@ -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 diff --git a/backend/tests/test_pipeline_failures.py b/backend/tests/test_pipeline_failures.py index 546c6498..73d44f3d 100644 --- a/backend/tests/test_pipeline_failures.py +++ b/backend/tests/test_pipeline_failures.py @@ -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