From 25f5537974b1d555bb0fe0b40327658b7e4be173 Mon Sep 17 00:00:00 2001 From: iomgaa Date: Sat, 11 Jul 2026 11:11:02 -0400 Subject: [PATCH] docs(plans): add Spec-2 batch tree build implementation plan --- research-wiki/graph/edges.json | 12 + research-wiki/index.md | 6 +- research-wiki/log.md | 3 + .../plans/2026-07-11-batch-tree-build.md | 880 ++++++++++++++++++ research-wiki/plans/batch-tree-build-plan.md | 9 + 5 files changed, 908 insertions(+), 2 deletions(-) create mode 100644 research-wiki/plans/2026-07-11-batch-tree-build.md create mode 100644 research-wiki/plans/batch-tree-build-plan.md diff --git a/research-wiki/graph/edges.json b/research-wiki/graph/edges.json index 488ecf1..39b3084 100644 --- a/research-wiki/graph/edges.json +++ b/research-wiki/graph/edges.json @@ -105,6 +105,11 @@ "id": "finding:eval-spec1check", "label": "Harness 评估: Spec-1 修复验证 (infer_spec1check)", "type": "finding" + }, + { + "id": "plan:batch-tree-build-plan", + "label": "Spec-2 建树批量并行实现计划", + "type": "plan" } ], "links": [ @@ -170,6 +175,13 @@ "relation": "implements", "evidence": "实现 Spec-1 三处改动 A1/A2/B", "added": "2026-07-11T12:06:25.879643+00:00" + }, + { + "source": "plan:batch-tree-build-plan", + "target": "design:batch-tree-build", + "relation": "implements", + "evidence": "实现 Spec-2 四任务", + "added": "2026-07-11T15:07:06.178664+00:00" } ] } \ No newline at end of file diff --git a/research-wiki/index.md b/research-wiki/index.md index 7cbcc10..05bf1db 100644 --- a/research-wiki/index.md +++ b/research-wiki/index.md @@ -1,6 +1,6 @@ # Research Wiki 索引 -> 自动生成,更新时间:2026-07-11 13:33 UTC +> 自动生成,更新时间:2026-07-11 15:07 UTC ## design (21) - [2026-07-06-core-agent-adapters-llm-design](designs/2026-07-06-core-agent-adapters-llm-design.md) `design:2026-07-06-core-agent-adapters-llm-design` @@ -30,7 +30,7 @@ - [2026-07-11-question-gen-calibration-analysis](findings/2026-07-11-question-gen-calibration-analysis.md) `finding:2026-07-11-question-gen-calibration-analysis` - [Harness 评估: Spec-1 修复验证 (infer_spec1check)](findings/eval-spec1check.md) `finding:eval-spec1check` -## plan (19) +## plan (21) - [2026-07-06-core-agent-adapters-llm](plans/2026-07-06-core-agent-adapters-llm.md) `plan:2026-07-06-core-agent-adapters-llm` - [2026-07-07-app-harness](plans/2026-07-07-app-harness.md) `plan:2026-07-07-app-harness` - [2026-07-07-core-evolution](plans/2026-07-07-core-evolution.md) `plan:2026-07-07-core-evolution` @@ -40,12 +40,14 @@ - [2026-07-09-question-gen-synth](plans/2026-07-09-question-gen-synth.md) `plan:2026-07-09-question-gen-synth` - [2026-07-09-tree-repair-resilience](plans/2026-07-09-tree-repair-resilience.md) `plan:2026-07-09-tree-repair-resilience` - [2026-07-11-agent-runtime-fixes](plans/2026-07-11-agent-runtime-fixes.md) `plan:2026-07-11-agent-runtime-fixes` +- [2026-07-11-batch-tree-build](plans/2026-07-11-batch-tree-build.md) `plan:2026-07-11-batch-tree-build` - [app/harness/ 训练循环编排层实现计划](plans/app-harness.md) `plan:app-harness` - [app/search/ 搜索 Agent 装配层实现计划](plans/2026-07-07-search-module.md) `plan:2026-07-07-search-module` - [core/agent/ + adapters/llm 基础设施实现计划](plans/core-agent-adapters-llm.md) `plan:core-agent-adapters-llm` - [main.py 推理入口 + 初始 Prompt 集实现计划](plans/main-inference-entry.md) `plan:main-inference-entry` - [question_gen 模块实现计划](plans/question-gen.md) `plan:question-gen` - [Spec-1 Agent 执行环境修复实现计划](plans/agent-runtime-fixes-plan.md) `plan:agent-runtime-fixes-plan` +- [Spec-2 建树批量并行实现计划](plans/batch-tree-build-plan.md) `plan:batch-tree-build-plan` - [建树修复管线三项改造实现计划](plans/tree-repair-resilience.md) `plan:tree-repair-resilience` - [建树模块竖切实现计划](plans/tree-module-vertical-slice.md) `plan:tree-module-vertical-slice` - [赛题生成工具实现计划](plans/question-gen-synth.md) `plan:question-gen-synth` diff --git a/research-wiki/log.md b/research-wiki/log.md index eb07d1e..a874aa6 100644 --- a/research-wiki/log.md +++ b/research-wiki/log.md @@ -47,3 +47,6 @@ - [2026-07-11 12:06 UTC] 重建索引: 42 篇页面 - [2026-07-11 13:33 UTC] 新增 finding: Harness 评估: Spec-1 修复验证 (infer_spec1check) (finding:eval-spec1check) - [2026-07-11 13:33 UTC] 重建索引: 43 篇页面 +- [2026-07-11 15:07 UTC] 新增 plan: Spec-2 建树批量并行实现计划 (plan:batch-tree-build-plan) +- [2026-07-11 15:07 UTC] 新增边: plan:batch-tree-build-plan --implements--> design:batch-tree-build +- [2026-07-11 15:07 UTC] 重建索引: 45 篇页面 diff --git a/research-wiki/plans/2026-07-11-batch-tree-build.md b/research-wiki/plans/2026-07-11-batch-tree-build.md new file mode 100644 index 0000000..44473bb --- /dev/null +++ b/research-wiki/plans/2026-07-11-batch-tree-build.md @@ -0,0 +1,880 @@ +# 建树批量并行入口(Spec-2)Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use subagent-driven-development to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** 新增多视频批量并行建树入口(视频级 16 并发 + 全局共享 API 信号量 16),单视频建树行为零变化。 + +**Architecture:** `VideoTreeBuilder` 公开 `build_async()`(消除事件循环嵌套)并支持外部注入 API Semaphore;`tools/build_trees.py` 复刻 `repair_trees.py` 的编排惯例(视频级 Semaphore + gather + progress.json 断点续跑 + 熔断阈值缩放);`scripts/build_trees.sh` 自包含实验入口。设计文档:`research-wiki/designs/2026-07-11-batch-tree-build-design.md`。 + +**Tech Stack:** Python 3.11 asyncio / pytest + pytest-asyncio / loguru / PyYAML。 + +**验证策略说明**:仓库 data/ 下当前无视频文件(`find data -name "*.mp4"` 为空),设计 §6 的"3-4 短视频小批量集成"改为**桩编排集成测试**(monkeypatch VideoTreeBuilder,验证并发上限、进度跳过、信号量共享、中断恢复),真实视频烟测流程写入 sh 脚本头部注释,首次实际建树时执行。 + +--- + +### Task 1: VideoTreeBuilder 异步入口公开 + API Semaphore 注入 + +**Files:** +- Modify: `app/tree/video_builder.py`(`__init__` 172-192 行;`build()` 284-299 行;`_build_async` 改名 `build_async` 305 行起;`vlm_sem = asyncio.Semaphore(...)` 357 行) +- Test: `tests/unit/test_video_builder.py`(追加测试类) + +- [ ] **Step 1: 写失败测试** + +`tests/unit/test_video_builder.py` 末尾追加(fixture `mock_vlm`/`mock_llm`/`tree_config` 已存在于该文件 159-191 行): + +```python +# ── Semaphore 注入与异步入口(Spec-2)───────────────────────── + + +class TestApiSemaphoreInjection: + """API Semaphore 注入与 build_async 公开入口。""" + + def test_injected_semaphore_stored( + self, + mock_vlm: MockVLMProvider, + mock_llm: MockLLMProvider, + tree_config: TreeConfig, + ) -> None: + """构造器注入的 Semaphore 应被保存供 build_async 使用。""" + sem = asyncio.Semaphore(3) + builder = VideoTreeBuilder( + vlm=mock_vlm, llm=mock_llm, config=tree_config, api_semaphore=sem + ) + assert builder._api_semaphore is sem + + def test_default_no_injection( + self, + mock_vlm: MockVLMProvider, + mock_llm: MockLLMProvider, + tree_config: TreeConfig, + ) -> None: + """未注入时属性为 None(build_async 内部自建,单视频行为零变化)。""" + builder = VideoTreeBuilder(vlm=mock_vlm, llm=mock_llm, config=tree_config) + assert builder._api_semaphore is None + + def test_build_async_is_public( + self, + mock_vlm: MockVLMProvider, + mock_llm: MockLLMProvider, + tree_config: TreeConfig, + ) -> None: + """build_async 必须是公开协程方法(供批量编排在事件循环内调用)。""" + builder = VideoTreeBuilder(vlm=mock_vlm, llm=mock_llm, config=tree_config) + assert hasattr(builder, "build_async") + assert asyncio.iscoroutinefunction(builder.build_async) +``` + +(文件顶部如无 `import asyncio` 则补上。) + +- [ ] **Step 2: 运行测试确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_video_builder.py::TestApiSemaphoreInjection -v` +Expected: FAIL——`api_semaphore` 是未知构造参数(TypeError),`build_async` 不存在。 + +- [ ] **Step 3: 实现** + +三处修改: + +其一,`__init__` 增加 keyword-only 可选参数(docstring 同步补充): + +```python + def __init__( + self, + vlm: VLMProvider, + llm: LLMProvider, + config: TreeConfig, + *, + api_semaphore: asyncio.Semaphore | None = None, + ) -> None: +``` + +赋值区追加: + +```python + self._api_semaphore = api_semaphore +``` + +docstring 参数段追加一行:`api_semaphore: 外部注入的全局 VLM/LLM 并发信号量(批量建树时跨视频共享);None 时 build_async 内部按 config.concurrency 自建,单视频行为零变化。` + +其二,`_build_async` 改名为公开 `build_async`(方法体不动,docstring 首行相应更新),`build()` 同步壳改为: + +```python + return asyncio.run(self.build_async(video_path, srt_entries)) +``` + +同时更新模块 docstring(16-21 行)与类 docstring(158-163 行)中的 `_build_async` 字样为 `build_async`。 + +其三,`build_async` 内 357 行的信号量创建改为: + +```python + # 创建 VLM/LLM 并发控制信号量(外部注入时跨视频全局共享,Spec-2) + vlm_sem = ( + self._api_semaphore + if self._api_semaphore is not None + else asyncio.Semaphore(self._config.concurrency) + ) +``` + +- [ ] **Step 4: 运行测试确认通过** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_video_builder.py -v` +Expected: 全部 PASS(既有测试 + 新增 3 个,确认改名无回归——若有既有测试直接引用 `_build_async`,同步改为 `build_async`)。 + +- [ ] **Step 5: Commit** + +```bash +git add app/tree/video_builder.py tests/unit/test_video_builder.py +git commit -m "feat(tree): expose build_async and accept injected API semaphore" +``` + +(commit body 标注:核心算法 #1/#2/#3 不变——仅入口封装与信号量来源切换,构建逻辑零改动。) + +--- + +### Task 2: tools/build_trees.py 辅助纯函数 + +**Files:** +- Create: `tools/build_trees.py`(本任务只写模块骨架 + 纯函数;编排在 Task 3) +- Test: `tests/unit/test_build_trees.py`(新建) + +- [ ] **Step 1: 写失败测试** + +新建 `tests/unit/test_build_trees.py`: + +```python +"""tools/build_trees.py 单元测试。 + +覆盖纯函数(完整性校验、待建清单发现、SRT 查找)与编排集成(Task 3 追加)。 +""" + +from __future__ import annotations + +import json +import sys +from pathlib import Path + +import pytest + +PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent +if str(PROJECT_ROOT) not in sys.path: + sys.path.insert(0, str(PROJECT_ROOT)) + +from app.tree.index import IndexMeta, L1Card, L1Node, TreeIndex +from tools.build_trees import _discover_pending, _find_srt_entries, _tree_is_complete + + +def _write_valid_tree(tree_path: Path) -> None: + """写入一棵最小合法树。""" + l1 = L1Node( + id="vid_L1_000", + card=L1Card("场景", "室内", ["实体"], ["动作"], ["关键词"], [], "线性"), + time_range=(0.0, 10.0), + children=[], + ) + index = TreeIndex(metadata=IndexMeta("/v.mp4", "video"), roots=[l1]) + tree_path.parent.mkdir(parents=True, exist_ok=True) + index.save_json(str(tree_path)) + + +class TestTreeIsComplete: + """_tree_is_complete 测试。""" + + def test_valid_tree(self, tmp_path: Path) -> None: + """合法 tree.json 判定完整。""" + tree_path = tmp_path / "vid" / "tree.json" + _write_valid_tree(tree_path) + assert _tree_is_complete(tree_path) is True + + def test_missing_file(self, tmp_path: Path) -> None: + """文件不存在判定不完整。""" + assert _tree_is_complete(tmp_path / "nope" / "tree.json") is False + + def test_corrupt_json(self, tmp_path: Path) -> None: + """损坏 JSON 判定不完整(不抛异常)。""" + p = tmp_path / "vid" / "tree.json" + p.parent.mkdir(parents=True) + p.write_text("{broken", encoding="utf-8") + assert _tree_is_complete(p) is False + + +class TestDiscoverPending: + """_discover_pending 测试。""" + + def _touch_videos(self, videos_dir: Path, names: list[str]) -> None: + videos_dir.mkdir(parents=True, exist_ok=True) + for n in names: + (videos_dir / n).write_bytes(b"") + + def test_all_pending_when_fresh(self, tmp_path: Path) -> None: + """无进度无产物时全部待建,按名排序。""" + videos = tmp_path / "videos" + self._touch_videos(videos, ["b.mp4", "a.mkv", "c.txt"]) + pending = _discover_pending(videos, tmp_path / "out", set()) + assert [p.name for p in pending] == ["a.mkv", "b.mp4"] # 非视频扩展名被忽略 + + def test_skips_finished_and_complete(self, tmp_path: Path) -> None: + """progress 已记录或 tree.json 完整的视频被跳过。""" + videos = tmp_path / "videos" + out = tmp_path / "out" + self._touch_videos(videos, ["a.mp4", "b.mp4", "c.mp4"]) + _write_valid_tree(out / "b" / "tree.json") # b 已有完整树 + pending = _discover_pending(videos, out, {"a"}) # a 在 progress 中 + assert [p.name for p in pending] == ["c.mp4"] + + def test_incomplete_tree_not_skipped(self, tmp_path: Path) -> None: + """tree.json 损坏的视频仍待建(重建覆盖)。""" + videos = tmp_path / "videos" + out = tmp_path / "out" + self._touch_videos(videos, ["a.mp4"]) + (out / "a").mkdir(parents=True) + (out / "a" / "tree.json").write_text("{broken", encoding="utf-8") + pending = _discover_pending(videos, out, set()) + assert [p.name for p in pending] == ["a.mp4"] + + +class TestFindSrtEntries: + """_find_srt_entries 测试。""" + + def test_found(self, tmp_path: Path) -> None: + """同名 .srt 存在时解析返回条目。""" + srt = tmp_path / "vid.srt" + srt.write_text( + "1\n00:00:01,000 --> 00:00:03,000\nhello world\n\n", + encoding="utf-8", + ) + entries = _find_srt_entries(tmp_path / "vid.mp4", tmp_path) + assert entries is not None + assert len(entries) == 1 + + def test_missing_returns_none(self, tmp_path: Path) -> None: + """无同名 .srt 返回 None。""" + assert _find_srt_entries(tmp_path / "vid.mp4", tmp_path) is None +``` + +- [ ] **Step 2: 运行确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_build_trees.py -v` +Expected: FAIL(ModuleNotFoundError: tools.build_trees)。 + +- [ ] **Step 3: 实现模块骨架 + 纯函数** + +新建 `tools/build_trees.py`: + +```python +#!/usr/bin/env python3 +"""批量并行建树入口:多视频并发构建三层 TreeIndex。 + +并发模型(Spec-2): + 视频级 Semaphore(video_concurrency) + gather —— 复刻 repair_trees.py 惯例; + 全局共享一个 API Semaphore(api_concurrency) 注入所有 VideoTreeBuilder, + 端点压力与单视频建树完全一致,吞吐提升来自非 API 阶段跨视频重叠。 + +用法: + conda activate Video-Tree-TRM + python tools/build_trees.py --videos-dir [--out-dir store/videos] + [--srt-dir ] [--video-concurrency 16] [--limit 0] + +api_concurrency 为工程配置,从 .env 读取 TREE_BUILD_API_CONCURRENCY(默认 16)。 +app/core/adapters 不 import 此脚本。 +""" + +from __future__ import annotations + +import argparse +import asyncio +import json +import os +import sys +import time +from pathlib import Path + +PROJECT_ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(PROJECT_ROOT)) + +from dotenv import load_dotenv +from loguru import logger + +load_dotenv(PROJECT_ROOT / ".env") + +import yaml + +from app.tree.config import TreeConfig +from app.tree.index import TreeIndex +from app.tree.subtitle import parse_srt +from app.tree.video_builder import VideoTreeBuilder + +# --------------------------------------------------------------------------- +# 日志配置:不缓存,立即输出 +# --------------------------------------------------------------------------- + +logger.remove() +logger.add( + sys.stderr, + format="{time:HH:mm:ss} | {level:<7} | {message}", + level="DEBUG", + colorize=True, +) +logger.add( + PROJECT_ROOT / "logs" / "build_trees.log", + format="{time:YYYY-MM-DD HH:mm:ss} | {level:<7} | {message}", + level="DEBUG", + rotation="50 MB", +) + +# --------------------------------------------------------------------------- +# 断点续跑 — progress 文件管理(复刻 repair_trees.py 惯例) +# --------------------------------------------------------------------------- + +PROGRESS_FILE = "build_progress.json" + +_VIDEO_SUFFIXES = frozenset({".mp4", ".mkv", ".avi", ".webm"}) + + +def load_progress(path: Path) -> set[str]: + """读取 progress 文件,返回已完成视频 ID 集合。 + + 参数: + path: progress JSON 文件路径。 + + 返回: + 已完成视频 ID 集合。文件不存在或损坏时返回空集。 + """ + if not path.exists(): + return set() + try: + data = json.loads(path.read_text(encoding="utf-8")) + return set(data.get("finished_video_ids", [])) + except (json.JSONDecodeError, KeyError, TypeError, AttributeError): + logger.warning("progress 文件损坏,忽略: {}", path) + return set() + + +async def save_progress(path: Path, lock: asyncio.Lock, vid: str) -> None: + """原子追加一个视频 ID 到 progress 文件。 + + 参数: + path: progress JSON 文件路径。 + lock: asyncio.Lock,防并发读改写丢更新。 + vid: 要追加的视频 ID。 + """ + async with lock: + finished = load_progress(path) + finished.add(vid) + tmp = path.with_suffix(".tmp") + tmp.write_text( + json.dumps({"finished_video_ids": sorted(finished)}, ensure_ascii=False, indent=2), + encoding="utf-8", + ) + os.replace(str(tmp), str(path)) + + +# --------------------------------------------------------------------------- +# 待建发现与完整性校验 +# --------------------------------------------------------------------------- + + +def _tree_is_complete(tree_path: Path) -> bool: + """判断 tree.json 是否存在且可加载为非空树。 + + 参数: + tree_path: tree.json 路径。 + + 返回: + True 表示完整(跳过重建);文件缺失/损坏/空树返回 False。 + """ + if not tree_path.exists(): + return False + try: + index = TreeIndex.load_json(str(tree_path)) + except (json.JSONDecodeError, KeyError, TypeError, ValueError, AssertionError) as exc: + logger.warning("tree.json 无法加载,视为不完整: {} ({})", tree_path, exc) + return False + return len(index.roots) > 0 + + +def _discover_pending( + videos_dir: Path, + out_dir: Path, + finished: set[str], +) -> list[Path]: + """扫描视频目录,返回待建视频文件列表(按文件名排序)。 + + 跳过条件:video_id 在 progress 中,或 out_dir//tree.json 完整。 + + 参数: + videos_dir: 视频文件目录。 + out_dir: 树输出根目录。 + finished: progress 中已完成的视频 ID 集合。 + + 返回: + 待建视频文件路径列表。 + """ + pending: list[Path] = [] + for f in sorted(videos_dir.iterdir()): + if not f.is_file() or f.suffix.lower() not in _VIDEO_SUFFIXES: + continue + vid = f.stem + if vid in finished: + continue + if _tree_is_complete(out_dir / vid / "tree.json"): + continue + pending.append(f) + return pending + + +def _find_srt_entries(video_path: Path, srt_dir: Path): + """按视频同名规则查找并解析 SRT 字幕。 + + 参数: + video_path: 视频文件路径。 + srt_dir: SRT 目录。 + + 返回: + SRTEntry 列表;无同名 .srt 时返回 None。 + """ + srt_path = srt_dir / f"{video_path.stem}.srt" + if not srt_path.exists(): + return None + return parse_srt(str(srt_path)) +``` + +(编排函数与 main 在 Task 3 追加;本步模块以纯函数收尾即可运行测试。) + +- [ ] **Step 4: 运行确认通过** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_build_trees.py -v` +Expected: 全部 PASS。 + +- [ ] **Step 5: Commit** + +```bash +git add tools/build_trees.py tests/unit/test_build_trees.py +git commit -m "feat(tools): add build_trees skeleton with discovery helpers" +``` + +--- + +### Task 3: 编排 main_async + CLI + 桩编排集成测试 + +**Files:** +- Modify: `tools/build_trees.py`(追加客户端构建、编排、CLI) +- Test: `tests/unit/test_build_trees.py`(追加编排测试) + +- [ ] **Step 1: 写失败测试(桩 builder 编排集成)** + +`tests/unit/test_build_trees.py` 追加: + +```python +# ── 编排集成(桩 builder,无真实视频/LLM)────────────────────── + +import asyncio + +from tools.build_trees import main_async + + +class _StubBuilder: + """记录并发与注入信号量的桩 builder。""" + + instances: list["_StubBuilder"] = [] + inflight = 0 + max_inflight = 0 + + def __init__(self, vlm, llm, config, *, api_semaphore=None) -> None: + self.api_semaphore = api_semaphore + _StubBuilder.instances.append(self) + + async def build_async(self, video_path: str, srt_entries=None): + _StubBuilder.inflight += 1 + _StubBuilder.max_inflight = max(_StubBuilder.max_inflight, _StubBuilder.inflight) + await asyncio.sleep(0.02) + _StubBuilder.inflight -= 1 + from app.tree.index import IndexMeta, L1Card, L1Node, TreeIndex + + l1 = L1Node( + id="x_L1_000", + card=L1Card("s", "室内", [], [], [], [], "线性"), + time_range=(0.0, 1.0), + children=[], + ) + return TreeIndex(metadata=IndexMeta(video_path, "video"), roots=[l1]) + + +@pytest.fixture() +def batch_env(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + """5 个假视频 + 桩 builder + 隔离的 progress 路径。""" + import tools.build_trees as bt + + _StubBuilder.instances = [] + _StubBuilder.inflight = 0 + _StubBuilder.max_inflight = 0 + monkeypatch.setattr(bt, "VideoTreeBuilder", _StubBuilder) + monkeypatch.setattr(bt, "_build_clients", lambda api_concurrency: (None, None)) + + videos = tmp_path / "videos" + videos.mkdir() + for i in range(5): + (videos / f"v{i}.mp4").write_bytes(b"") + + return { + "videos": videos, + "out": tmp_path / "out", + "progress": tmp_path / "build_progress.json", + } + + +def _make_args(env: dict, video_concurrency: int = 2, limit: int = 0) -> "argparse.Namespace": + import argparse + + return argparse.Namespace( + videos_dir=str(env["videos"]), + out_dir=str(env["out"]), + srt_dir=str(env["videos"]), + video_concurrency=video_concurrency, + limit=limit, + progress_path=str(env["progress"]), + ) + + +class TestOrchestration: + """main_async 编排行为(桩 builder)。""" + + @pytest.mark.asyncio + async def test_video_concurrency_capped(self, batch_env: dict) -> None: + """同时在建视频数不得超过 video_concurrency。""" + await main_async(_make_args(batch_env, video_concurrency=2)) + assert _StubBuilder.max_inflight <= 2 + assert len(_StubBuilder.instances) == 5 + + @pytest.mark.asyncio + async def test_shared_api_semaphore(self, batch_env: dict) -> None: + """全部 builder 实例共享同一个 API Semaphore 对象。""" + await main_async(_make_args(batch_env)) + sems = {id(b.api_semaphore) for b in _StubBuilder.instances} + assert len(sems) == 1 + assert _StubBuilder.instances[0].api_semaphore is not None + + @pytest.mark.asyncio + async def test_trees_saved_and_progress_recorded(self, batch_env: dict) -> None: + """每个视频产出 tree.json 且 progress 记录全部完成。""" + await main_async(_make_args(batch_env)) + for i in range(5): + assert (batch_env["out"] / f"v{i}" / "tree.json").exists() + from tools.build_trees import load_progress + + assert load_progress(batch_env["progress"]) == {f"v{i}" for i in range(5)} + + @pytest.mark.asyncio + async def test_resume_skips_finished(self, batch_env: dict) -> None: + """第二次运行跳过全部已完成视频。""" + await main_async(_make_args(batch_env)) + n_first = len(_StubBuilder.instances) + await main_async(_make_args(batch_env)) + assert len(_StubBuilder.instances) == n_first # 无新建 + + @pytest.mark.asyncio + async def test_partial_completion_resume(self, batch_env: dict) -> None: + """部分完成后重跑只建剩余视频(模拟中断后恢复)。""" + batch_env["progress"].write_text( + json.dumps({"finished_video_ids": ["v0", "v1"]}), + encoding="utf-8", + ) + await main_async(_make_args(batch_env)) + assert len(_StubBuilder.instances) == 3 # 仅 v2/v3/v4 + + @pytest.mark.asyncio + async def test_limit(self, batch_env: dict) -> None: + """--limit 2 只建前两个(烟测入口)。""" + await main_async(_make_args(batch_env, limit=2)) + assert len(_StubBuilder.instances) == 2 +``` + +- [ ] **Step 2: 运行确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_build_trees.py::TestOrchestration -v` +Expected: FAIL(ImportError: main_async 不存在)。 + +- [ ] **Step 3: 实现客户端构建 + 编排 + CLI** + +`tools/build_trees.py` 追加: + +```python +# --------------------------------------------------------------------------- +# LLM/VLM 客户端构建(复刻 repair_trees.py 惯例) +# --------------------------------------------------------------------------- + + +def _build_clients(api_concurrency: int): + """构建 GovernedLLMClient(LLM + VLM),熔断阈值随 API 并发缩放。 + + 参数: + api_concurrency: 全局 API 并发上限(熔断阈值取 max(cfg, api_concurrency*2))。 + + 返回: + (llm_client, vlm_client) 元组。 + """ + from adapters.breaker import CircuitBreaker + from adapters.llm import GovernedLLMClient + from adapters.telemetry import SQLiteTelemetryRecorder + from adapters.vlm import GovernedVLMClient + + (PROJECT_ROOT / "logs").mkdir(exist_ok=True) + telemetry = SQLiteTelemetryRecorder(str(PROJECT_ROOT / "logs" / "build_trees_telemetry.db")) + + breaker_threshold = int(os.getenv("LLM_CIRCUIT_BREAKER_THRESHOLD", "5")) + breaker_threshold = max(breaker_threshold, api_concurrency * 2) + breaker_cooldown = int(os.getenv("LLM_CIRCUIT_BREAKER_COOLDOWN", "60")) + timeout_s = float(os.getenv("LLM_TIMEOUT", "120")) + max_retries = int(os.getenv("LLM_MAX_RETRIES", "3")) + base_delay = float(os.getenv("LLM_RETRY_BASE_DELAY", "2.0")) + max_delay = float(os.getenv("LLM_RETRY_MAX_DELAY", "30.0")) + ttft = float(os.getenv("LLM_TTFT_TIMEOUT", "30")) + inter_token = float(os.getenv("LLM_INTER_TOKEN_TIMEOUT", "15")) + + llm = GovernedLLMClient( + model=os.environ["SEARCH_LLM_MODEL"], + base_url=os.environ["SEARCH_LLM_BASE_URL"], + api_key=os.environ["SEARCH_LLM_API_KEY"], + provider="deepseek", + thinking=False, + breaker=CircuitBreaker(fail_threshold=breaker_threshold, cooldown_s=breaker_cooldown), + cache=None, + telemetry=telemetry, + timeout_s=timeout_s, + ttft_timeout_s=ttft, + inter_token_timeout_s=inter_token, + max_retries=max_retries, + retry_base_delay_s=base_delay, + retry_max_delay_s=max_delay, + ) + vlm_base = GovernedLLMClient( + model=os.environ["VL_LLM_MODEL"], + base_url=os.environ["VL_LLM_BASE_URL"], + api_key=os.environ["VL_LLM_API_KEY"], + provider="qwen", + thinking=False, + breaker=CircuitBreaker(fail_threshold=breaker_threshold, cooldown_s=breaker_cooldown), + cache=None, + telemetry=telemetry, + timeout_s=timeout_s, + ttft_timeout_s=ttft, + inter_token_timeout_s=inter_token, + max_retries=max_retries, + retry_base_delay_s=base_delay, + retry_max_delay_s=max_delay, + ) + return llm, GovernedVLMClient(vlm_base) + + +# --------------------------------------------------------------------------- +# 主编排 +# --------------------------------------------------------------------------- + + +async def main_async(args: argparse.Namespace) -> None: + """异步主流程:视频级并发建树 + 全局共享 API 信号量。 + + 参数: + args: CLI 参数(videos_dir/out_dir/srt_dir/video_concurrency/limit/progress_path)。 + """ + videos_dir = Path(args.videos_dir) + out_dir = Path(args.out_dir) + srt_dir = Path(args.srt_dir) + assert videos_dir.is_dir(), f"视频目录不存在: {videos_dir}" + + api_concurrency = int(os.getenv("TREE_BUILD_API_CONCURRENCY", "16")) + + # Phase 1: 待建发现(progress + 完整性双重跳过) + progress_path = Path(args.progress_path) + finished = load_progress(progress_path) + if finished: + logger.info("progress 已记录 {} 个完成视频", len(finished)) + pending = _discover_pending(videos_dir, out_dir, finished) + if args.limit > 0: + pending = pending[: args.limit] + logger.info( + "待建 {} 个视频, video_concurrency={}, api_concurrency={}", + len(pending), args.video_concurrency, api_concurrency, + ) + if not pending: + return + + # Phase 2: 客户端与共享信号量 + llm, vlm = _build_clients(api_concurrency) + with open(PROJECT_ROOT / "config" / "default.yaml", encoding="utf-8") as f: + tree_cfg = TreeConfig.from_dict(yaml.safe_load(f)["tree"]) + api_sem = asyncio.Semaphore(api_concurrency) + video_sem = asyncio.Semaphore(args.video_concurrency) + progress_lock = asyncio.Lock() + start_time = time.time() + completed = 0 + failed: list[str] = [] + + # Phase 3: 视频级并发编排(复刻 repair_trees 模式) + async def _build_one(video_path: Path) -> None: + nonlocal completed + async with video_sem: + vid = video_path.stem + logger.info("开始建树 {}", vid) + builder = VideoTreeBuilder(vlm=vlm, llm=llm, config=tree_cfg, api_semaphore=api_sem) + srt_entries = _find_srt_entries(video_path, srt_dir) + try: + index = await builder.build_async(str(video_path), srt_entries) + except Exception as exc: + logger.error("建树失败 {} ({}): {}", vid, type(exc).__name__, exc) + failed.append(vid) + return + tree_path = out_dir / vid / "tree.json" + tree_path.parent.mkdir(parents=True, exist_ok=True) + index.save_json(str(tree_path)) + await save_progress(progress_path, progress_lock, vid) + completed += 1 + if completed % 5 == 0: + elapsed = time.time() - start_time + rate = completed / elapsed * 60 if elapsed > 0 else 0 + logger.info( + "进度: {}/{}, 已用 {:.0f}s, 速率 {:.2f} 视频/分钟", + completed, len(pending), elapsed, rate, + ) + + await asyncio.gather(*[asyncio.create_task(_build_one(p)) for p in pending]) + + # Phase 4: 汇总 + elapsed = time.time() - start_time + logger.info( + "批量建树完成: 成功 {}, 失败 {}, 总耗时 {:.0f}s{}", + completed, len(failed), elapsed, + f", 失败清单: {failed}" if failed else "", + ) + + +def main() -> None: + """同步入口。""" + parser = argparse.ArgumentParser(description="批量并行建树") + parser.add_argument("--videos-dir", type=str, required=True, help="视频文件目录") + parser.add_argument("--out-dir", type=str, default="store/videos", help="树输出根目录") + parser.add_argument("--srt-dir", type=str, default="", help="SRT 目录(默认同 videos-dir)") + parser.add_argument("--video-concurrency", type=int, default=16, help="同时在建视频数") + parser.add_argument("--limit", type=int, default=0, help="只建前 N 个(0=全部,烟测用)") + parser.add_argument( + "--progress-path", + type=str, + default=str(PROJECT_ROOT / "logs" / PROGRESS_FILE), + dest="progress_path", + help="progress 文件路径", + ) + args = parser.parse_args() + if not args.srt_dir: + args.srt_dir = args.videos_dir + asyncio.run(main_async(args)) + + +if __name__ == "__main__": + main() +``` + +注意:`_build_one` 中 `except Exception` 用于**单视频错误隔离**(一个视频失败不拖垮整批,失败清单汇总上报,vid 不进 progress 可重跑)——这是编排层的错误隔离语义而非吞错,日志含异常类型与内容。 + +- [ ] **Step 4: 运行确认通过** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_build_trees.py -v` +Expected: 全部 PASS(纯函数 8 个 + 编排 6 个,共 14 个)。 + +- [ ] **Step 5: Commit** + +```bash +git add tools/build_trees.py tests/unit/test_build_trees.py +git commit -m "feat(tools): batch tree build orchestration with shared API semaphore" +``` + +--- + +### Task 4: scripts/build_trees.sh + 工程配置 + 回归收尾 + +**Files:** +- Create: `scripts/build_trees.sh` +- Modify: `.env.example`(追加 TREE_BUILD_API_CONCURRENCY) + +- [ ] **Step 1: 写 sh 脚本** + +新建 `scripts/build_trees.sh`: + +```bash +#!/usr/bin/env bash +# 批量并行建树(Spec-2) +# 职责:对目录下所有视频并发构建三层 TreeIndex,断点续跑。 +# +# 用法: +# bash scripts/build_trees.sh --videos-dir # 全量 +# bash scripts/build_trees.sh --videos-dir --limit 2 # 真实视频烟测 +# VIDEO_CONCURRENCY=8 bash scripts/build_trees.sh --videos-dir +# +# 真实视频烟测流程(首次使用时执行): +# 1. --limit 2 建两个视频,观察日志速率与 API 并发 +# 2. 中途 Ctrl+C 后重跑,确认已完成视频被跳过、未完成视频续跑 +# 3. 检查 store/videos//tree.json 可被 TreeIndex.load_json 加载 +# +# 日志输出: +# stderr → 终端实时显示 +# logs/build_trees.log → 完整日志(自动 rotation 50MB) +# logs/build_trees_telemetry.db → LLM/VLM 调用遥测 +# logs/build_progress.json → 断点续跑进度 +# +# api_concurrency 走 .env 的 TREE_BUILD_API_CONCURRENCY(工程配置,默认 16) + +set -euo pipefail + +cd "$(dirname "$0")/.." + +VIDEO_CONCURRENCY="${VIDEO_CONCURRENCY:-16}" + +export PYTHONUNBUFFERED=1 + +# shellcheck source=/dev/null +source "$(conda info --base)/etc/profile.d/conda.sh" +conda activate Video-Tree-TRM + +python tools/build_trees.py \ + --out-dir store/videos \ + --video-concurrency "$VIDEO_CONCURRENCY" \ + "$@" +``` + +```bash +chmod +x scripts/build_trees.sh +``` + +- [ ] **Step 2: 补 .env.example** + +`.env.example` 追加(若文件不存在则跳过并在 NOTES 说明): + +``` +# 建树批量并行:全局 VLM/LLM 在途调用上限(Spec-2 工程配置) +TREE_BUILD_API_CONCURRENCY=16 +``` + +- [ ] **Step 3: 全量回归 + lint** + +Run: `conda run -n Video-Tree-TRM ruff format --check tools/build_trees.py tests/unit/test_build_trees.py app/tree/video_builder.py` +Run: `conda run -n Video-Tree-TRM ruff check tools/build_trees.py tests/unit/test_build_trees.py app/tree/video_builder.py` +Run: `conda run -n Video-Tree-TRM radon cc tools/build_trees.py app/tree/video_builder.py -n C -s`(新增函数不得出现 C 级;video_builder 既有 C 级为存量) +Run: `conda run -n Video-Tree-TRM pytest tests/ -q` +Expected: 全绿。 + +- [ ] **Step 4: Commit** + +```bash +git add scripts/build_trees.sh .env.example +git commit -m "feat(scripts): add batch tree build entry with smoke-test guide" +``` + +--- + +## Self-Review 记录 + +1. **Spec 覆盖**:设计 §3 全部落地——build_async 公开 + Semaphore 注入(Task 1)、两层并发参数与全局共享信号量(Task 3)、熔断阈值缩放(Task 3 `_build_clients`)、progress 断点续跑与完整性跳过(Task 2/3)、sh 入口与 D7 配置归属(Task 4);§4 速率日志(Task 3 Phase 3);§6 验证改为桩编排测试(头部"验证策略说明"已声明理由:仓库无视频文件),自动化覆盖"已完成跳过"(`test_resume_skips_finished`)与"部分完成后恢复"(`test_partial_completion_resume`);**段级断点续跑**由 builder 既有 `_load_resume_state` 机制承担(算法 #3,本计划零改动)+ 真实视频烟测流程(sh 头注释步骤 2)在首次实际建树时验证。 +2. **占位符扫描**:无 TBD/TODO;所有代码步骤含完整代码。 +3. **类型一致性**:`_discover_pending(videos_dir, out_dir, finished)`、`_find_srt_entries(video_path, srt_dir)`、`main_async(args)`、`_build_clients(api_concurrency)` 在测试与实现间签名一致;`_make_args` 构造的 Namespace 字段与 `main_async` 读取字段一一对应(含 progress_path)。 + +## 核心算法保真校验 + +本计划涉及**算法 #1(L2 轴心建树)/#2(VLM 批量帧描述)/#3(断点续跑)**所在文件 `app/tree/video_builder.py`:改动仅为 (a) `_build_async` 改名公开(方法体零改动)、(b) 信号量来源三元切换(默认路径与现状逐字等价)、(c) 构造器加可选参数。三项算法的核心逻辑(L2→L3 链式触发、批量 VLM+fallback、progress 段级恢复)均未触碰。执行时质量门第 7 项应 diff 确认 `build_async` 方法体与原 `_build_async` 逐行一致。其余算法不涉及。 diff --git a/research-wiki/plans/batch-tree-build-plan.md b/research-wiki/plans/batch-tree-build-plan.md new file mode 100644 index 0000000..a7c5f92 --- /dev/null +++ b/research-wiki/plans/batch-tree-build-plan.md @@ -0,0 +1,9 @@ +--- +type: plan +node_id: plan:batch-tree-build-plan +title: "Spec-2 建树批量并行实现计划" +date: 2026-07-11 +--- + +# Spec-2 建树批量并行实现计划 +