三事: (1)遥测防御性加固(INSERT OR IGNORE+WAL+busy_timeout+try/except不冒泡) 根治call_id重试复用导致的主键冲突污染重试与熔断器; (2)判据双轨 (detect重扫empty_field+progress文件)视频级断点续跑,首次用detect 自动接历史无需手补状态; (3)asyncio.Semaphore并发默认16,熔断阈值 max(5,concurrency*2),--concurrency走CLI避免与tree.concurrency串台
9.2 KiB
建树修复管线:熔断根因修复 + 断点续跑 + 并发改造
设计日期:2026-07-08 状态:待批准 关联:
app/tree/repair/、tools/repair_trees.py、adapters/telemetry.py、adapters/llm.py、adapters/breaker.py
1. 背景与触发
当前修复管线串行跑 300 个视频,耗时 34 小时。上次运行日志(logs/repair_trees.log)与遥测库(logs/repair_telemetry.db)交叉定位出:VLM 熔断并非服务挂了,而是遥测写入的主键冲突污染了重试与熔断器,导致 102 个视频完全未修复(L3=0)即被跳过,而程序最终汇报"失败数: 0"(假象)。
本设计解决三件事:(1) 根治熔断误触发;(2) 视频级断点续跑;(3) 并发执行,默认 16 路。
2. 根因分析
VLM 偶发的真实瞬时错误(500/超时)按设计走 transient → record_failure + 写遥测 + 退避重试。但 GovernedLLMClient.chat()(adapters/llm.py:299)的 call_id 在重试循环外只生成一次,重试时复用同一 call_id:
flowchart TD
A["attempt 1: VLM 真实失败"] --> B["transient 分支<br/>record_failure + 写遥测 call_id=X ✓<br/>sleep 退避, continue"]
B --> C["attempt 2: 又失败"]
C --> D["transient 分支<br/>写遥测 call_id=X"]
D --> E["💥 UNIQUE constraint<br/>SQLiteRecorder 无冲突容忍"]
E --> F["IntegrityError 冒泡<br/>_is_transient_error 不认识它"]
F --> G["非瞬时非致命分支<br/>再写遥测 call_id=X → 又冲突<br/>raise"]
G --> H["regenerator except<br/>吞掉异常,跳过节点"]
三个独立缺陷叠加放大:
| # | 缺陷 | 位置 | 后果 |
|---|---|---|---|
| 1 | call_id 重试时不重新生成 |
llm.py:299 在 for attempt 外 |
同一 call_id 多次 INSERT,必然主键冲突 |
| 2 | SQLiteTelemetryRecorder 无冲突容忍 |
telemetry.py:_write 裸 INSERT,无 except |
IntegrityError 冒泡,污染调用方 |
| 3 | _is_transient_error 不识别 IntegrityError |
llm.py:424 只认 httpx/StreamLiveness |
IntegrityError 当"非瞬时非致命"抛出,重试失效 |
并发改造还会引爆第四个隐患:SQLiteTelemetryRecorder 每次写新建 sqlite3.connect(),并发 16 路同时写会撞 SQLite 表锁(database is locked)。
3. 三项改造设计
3.1 遥测防御性加固(根治根因)
核心原则:遥测是观测侧信道,绝不能拖垮主调用链(CLAUDE.md P5)。
与 P5 的权衡说明:P5 反对的是"掩盖数据正确性错误"(如 LLM 返回解析失败却用默认值继续,污染业务数据)。遥测写失败不损害任何业务数据——它只丢一条观测记录。这是错误隔离(isolation),不是掩盖错误。平衡点:捕获后必须 logger.warning,错误可见可追溯,但不冒泡到 LLM 重试链。静默 pass 才违反 P5,warning 不违反。
SQLiteTelemetryRecorder._write 三层加固:
| 层 | 做法 | 解决 |
|---|---|---|
| SQL 层 | INSERT → INSERT OR IGNORE |
主键冲突静默 |
| 连接层 | PRAGMA journal_mode=WAL + PRAGMA busy_timeout=5000 |
并发写锁降级为排队 |
| 异常层 | 整个 _write 包 try/except sqlite3.Error,仅 logger.warning |
DB 任何错误不冒泡到 LLM 重试链 |
WAL 模式允许"1 写 + 多读"并发,写之间靠 busy_timeout 自动排队等待(毫秒级,不报错),不引入新瓶颈。
GovernedLLMClient.chat() 的 call_id 生成移入重试循环内(每次 attempt 重新 uuid4()),消除根因——虽 OR IGNORE 后冲突不再致命,但 call_id 唯一性本身是对的。
call_id 移入循环的遥测语义:当前一次 chat() 调用在成功/各失败分支共用一个 call_id,语义是"一次逻辑调用 = 一条最终记录(最后一次 attempt 的结果)"。移入循环后语义变为"一次逻辑调用 = N 条记录(每 attempt 一条,按 created_at 可追溯重试轨迹)"。这更利于事后诊断重试行为。parent_call_id 是 chat() 入参(llm.py:276),在循环外固定,不受影响——每条 attempt 记录都正确关联到父 agent step。
3.2 视频级断点续跑
判据双轨:
| 轨 | 作用 | 内容 |
|---|---|---|
| 数据驱动判据 | 跳过已修干净的视频 | detect_issues 重扫,有 empty_field / L2 event_description 空 / L1 scene_summary 空才进队 |
| progress 文件 | 加速 + 审计 | logs/repair_progress.json 记 finished_video_ids,已修干净的视频按 ID 直接跳过,省 detect 开销 |
关键判据边界:跳过判据只用 empty_field + L2/L1 空字段,不用 missing_frame——修复根本不处理缺帧,用它判跳过会让缺帧视频永远进队死循环。
其他 issue_type 的处理:no_children(L2/L1 无子节点)是结构性缺陷,修复不处理(regenerator 只重生成 card,不改树结构)——排除出跳过判据,避免误判。time_gap(相邻 L2 时间间隙 >1s)是可接受的时间分布特征,非缺陷——排除。即只有 empty_field(L3 四必填字段 + 新增 L2 event_description + L1 scene_summary)参与跳过判定。
首次续跑零成本接历史:progress 文件不存在时,用 detect_issues 扫一遍初始化它,自动识别上次的 102 个未修视频,无需手动补状态文件。
重修语义:队列里的视频跑完整级联(修空 L3 后重生成所有 L2 和所有 L1),刷新临界区"陈旧但非空"的上层卡片。LLM 不可用时逐节点 try/except 自动优雅降级为"只修剩余"。
临界区漏修处理(A 方案):临界区视频(L3 好、L2/L1 非空但过时)的卡片非空,detect_issues 扫不出 → 进不了队列 → 首次漏修。提供 --reaggregate-all 标志兜底(强制全量重聚合)。后续有 progress 文件即再无此问题。
3.3 并发执行
编排结构:照搬建树 video_builder.py 的 asyncio.Semaphore(concurrency) 范式。
| 点 | 做法 |
|---|---|
| 并发粒度 | Semaphore 限视频数 = concurrency,视频内四步(detect→repair→verify→supplement→save)串行 |
| 失败隔离 | 单视频失败只记该视频 error,不影响其他路继续 |
| progress 写入 | 并发下多协程向同一 finished_video_ids 追加,是读-改-写场景。用 asyncio.Lock 保护读改写 + 临时文件 os.replace 原子替换:锁内读旧 json → 追加 ID → 写 .tmp → os.replace 原子替换。锁保证不丢更新,rename 保证崩溃不留半写文件 |
| 完成汇报 | 主循环每 N 个视频汇总进度 |
熔断阈值适配并发:LLM_CIRCUIT_BREAKER_THRESHOLD(.env)默认从 5 改为 max(5, concurrency*2),保持单实例共享——上游 API 配额是全局的,熔断本就该全局生效。
阈值覆盖关系(D7 配置优先级):CLI --concurrency > .env 的 LLM_CIRCUIT_BREAKER_THRESHOLD。实际阈值为 max(.env 显式值, concurrency*2)——用户在 .env 显式设的阈值是下限保护(绝不低于它),concurrency*2 是并发自适应下限,两者取大。若用户未在 .env 设(用默认 5),则按 concurrency*2 生效。这样既尊重用户的显式运维配置,又保证并发下不会过激熔断。
配置归属:并发数用 CLI 参数 --concurrency(默认 16),不进 config/default.yaml。理由:config/default.yaml 的 tree.concurrency 是建树扫动参数(科研对比),修复并发是运维调度参数(本机 CPU/网络),混进同一 YAML 会串台(D7 规则)。
4. 改动范围
| 文件 | 改动 | 性质 |
|---|---|---|
adapters/telemetry.py |
_write 加 WAL + busy_timeout + INSERT OR IGNORE + try/except sqlite3.Error |
防御加固 |
adapters/llm.py |
call_id 生成移入重试循环内 |
根因修复 |
adapters/breaker.py |
无改动 | — |
tools/repair_trees.py |
并发编排(Semaphore)+ 断点续跑(progress 文件 + detect 判据)+ --concurrency/--reaggregate-all CLI |
新增能力 |
app/tree/repair/detector.py |
空字段检测扩展到 L2 event_description / L1 scene_summary |
增强(零 LLM 成本) |
.env.example |
熔断阈值说明更新 | 文档 |
5. 测试策略
| 场景 | 验证 |
|---|---|
| 遥测主键冲突静默 | 重复 call_id 写入不抛异常 |
| 遥测 DB 错误不冒泡 | 模拟 sqlite3.OperationalError,record_llm_call 不影响主调用 |
| 并发写不报锁错 | 16 路 to_thread 并发写,无 database is locked |
call_id 重试唯一 |
transient 重试后 DB 中各 attempt 独立记录 |
| 断点续跑幂等 | 已修视频重跑直接跳过;progress 丢失靠 detect 恢复 |
| 临界区完整级联 | 队列视频跑完后 L2/L1 全部重生成 |
| 并发编排 | Semaphore 限流生效,单视频失败不阻断其他 |
6. 待确认风险
- 熔断后 progress 仍写入:熔断期视频虽未修复但会跑完四步(repair 跳过→verify→supplement 失败→save),需要判定这种"跑完但没修"是否计为
finished。建议:不计入finished,只记 detect 抓到的问题数为 0 且实际未调 LLM 的视频为skipped,保证 progress 语义=真正修复完成。