diff --git a/research-wiki/graph/edges.json b/research-wiki/graph/edges.json index 2292939..a3368fb 100644 --- a/research-wiki/graph/edges.json +++ b/research-wiki/graph/edges.json @@ -280,6 +280,11 @@ "id": "design:gate-speedup", "label": "gate 验证提速:预灌 BaselineCache + 双臂并行", "type": "design" + }, + { + "id": "plan:gate-speedup", + "label": "连续并发 gate + Redis 复用实现计划", + "type": "plan" } ], "links": [ @@ -562,6 +567,13 @@ "relation": "informs", "evidence": "跨 task 终审确认设计 21 项缺陷全部落地 + 4 项集成发现", "added": "2026-07-16T11:06:43.709183+00:00" + }, + { + "source": "plan:gate-speedup", + "target": "design:gate-speedup", + "relation": "implements", + "evidence": "实现设计 v3 的 ②″+③", + "added": "2026-07-17T03:12:53.530801+00:00" } ] } \ No newline at end of file diff --git a/research-wiki/index.md b/research-wiki/index.md index 96f9764..d2e0a64 100644 --- a/research-wiki/index.md +++ b/research-wiki/index.md @@ -1,6 +1,6 @@ # Research Wiki 索引 -> 自动生成,更新时间:2026-07-16 18:23 UTC +> 自动生成,更新时间:2026-07-17 03:12 UTC ## design (40) - [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` @@ -56,7 +56,7 @@ - [v3 §9 帧感知抽取机制 — 小样本实证验证结果](findings/2026-07-15-v3-frame-perception-spike-validation.md) `finding:2026-07-15-v3-frame-perception-spike-validation` - [出题范式转变 — 从"生成-打分-过滤"转向"构造优先 + 两正交独立信号"(六篇原文深读)](findings/2026-07-15-question-gen-paradigm-shift-construction-over-filtering.md) `finding:2026-07-15-question-gen-paradigm-shift-construction-over-filtering` -## plan (48) +## plan (50) - [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` @@ -76,6 +76,7 @@ - [2026-07-14-task-type-strategy-framework](plans/2026-07-14-task-type-strategy-framework.md) `plan:2026-07-14-task-type-strategy-framework` - [2026-07-15-question-gen-v3-phase1-contract](plans/2026-07-15-question-gen-v3-phase1-contract.md) `plan:2026-07-15-question-gen-v3-phase1-contract` - [2026-07-15-results-driven-video-split](plans/2026-07-15-results-driven-video-split.md) `plan:2026-07-15-results-driven-video-split` +- [2026-07-16-gate-speedup](plans/2026-07-16-gate-speedup.md) `plan:2026-07-16-gate-speedup` - [2026-07-16-preflight-wp1-asset-migration](plans/2026-07-16-preflight-wp1-asset-migration.md) `plan:2026-07-16-preflight-wp1-asset-migration` - [2026-07-16-preflight-wp2-split-wiring](plans/2026-07-16-preflight-wp2-split-wiring.md) `plan:2026-07-16-preflight-wp2-split-wiring` - [2026-07-16-preflight-wp3-train-loop](plans/2026-07-16-preflight-wp3-train-loop.md) `plan:2026-07-16-preflight-wp3-train-loop` @@ -104,6 +105,7 @@ - [建树模块竖切实现计划](plans/tree-module-vertical-slice.md) `plan:tree-module-vertical-slice` - [结果驱动的视频级切分实现计划](plans/results-driven-video-split-plan.md) `plan:results-driven-video-split-plan` - [赛题生成工具实现计划](plans/question-gen-synth.md) `plan:question-gen-synth` +- [连续并发 gate + Redis 复用实现计划](plans/gate-speedup.md) `plan:gate-speedup` - [项目基础设施初始化计划](plans/infrastructure-setup.md) `plan:infrastructure-setup` ## review (4) diff --git a/research-wiki/log.md b/research-wiki/log.md index d9f9329..6e04723 100644 --- a/research-wiki/log.md +++ b/research-wiki/log.md @@ -139,3 +139,6 @@ - [2026-07-16 12:47 UTC] 重建索引: 113 篇页面 - [2026-07-16 18:23 UTC] 新增 design: gate 验证提速:预灌 BaselineCache + 双臂并行 (design:gate-speedup) - [2026-07-16 18:23 UTC] 重建索引: 115 篇页面 +- [2026-07-17 03:12 UTC] 新增 plan: 连续并发 gate + Redis 复用实现计划 (plan:gate-speedup) +- [2026-07-17 03:12 UTC] 新增边: plan:gate-speedup --implements--> design:gate-speedup +- [2026-07-17 03:12 UTC] 重建索引: 117 篇页面 diff --git a/research-wiki/plans/2026-07-16-gate-speedup.md b/research-wiki/plans/2026-07-16-gate-speedup.md new file mode 100644 index 0000000..fd3cffa --- /dev/null +++ b/research-wiki/plans/2026-07-16-gate-speedup.md @@ -0,0 +1,1349 @@ +# 连续并发 gate + Redis 复用 实现计划 + +> **For agentic workers:** REQUIRED SUB-SKILL: Use subagent-driven-development to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** 把 gate 验证从"题型串行 × 块串行 × 双臂串行"重构为"推理乱序并发 × 统计阶梯序前缀有序"的连续调度器,单 step gate 墙钟 ~9-18h → ~1.5-2.5h;顺手修复 step 重跑不清 gate 行的幂等 bug;Redis TTL 提到 7 天。 + +**Architecture:** 设计文档 `research-wiki/designs/2026-07-16-gate-speedup-design.md`(v3,已批准)。统计内核(`core/evolution/gate.py` 的 e-值/四出口、`pair_block` 配对、信息量阶梯、试用期)**一行不动**;只重构 `app/harness/validate.py` 的推理编排层与 `app/harness/runner.py` 的 `_gate_batch_skills` 装配层。核心不变量:**(W,L) 只按预声明阶梯序的"已配齐前缀"推进,严禁按完成到达序消费**(到达序消费会因两臂延迟不对称产生假接受偏差,Codex 审查 C1)。 + +**Tech Stack:** Python 3.11 asyncio(Semaphore/Lock/gather),pytest-asyncio,SQLite(HarnessLog 单连接+锁模式)。 + +**保真声明:** 本计划涉及核心算法 #6(块顺序验证)的**已批准语义修订**:块序贯 → 阶梯序前缀逐对序贯,判据不变。#4(CE-Gate e-process)、#5(信息阶梯)不动。每个触及 #6 的 commit message 须标注 `(algo #6)`。 + +--- + +## 全局约束(每个 Task 都适用) + +- 所有命令在 `Video-Tree-TRM` conda 环境:`conda run -n Video-Tree-TRM pytest ...` +- 禁止 `except Exception: pass`;中文 docstring;loguru 日志 +- 现有纯函数复用,不重写:`_ladder_units` / `_load_run_rows` / `_infra_question_ids_from_db` / `_candidate_correctness_from_db` / `_count_infra_units` / `_check_infra_guard` / `_finalize_outcome` / `materialize_candidate_skill` / `unit_correctness_view` / `pair_block` / `gate_decision` +- 每个 Task 结束跑 `conda run -n Video-Tree-TRM pytest tests/unit -x -q` 保持全绿后才 commit + +--- + +### Task 1: 前缀消费纯逻辑 + 数据结构(validate.py 新增,不删旧码) + +**Files:** +- Modify: `app/harness/validate.py`(文件末尾追加新 section) +- Test: `tests/unit/test_gate_prefix.py`(新建) + +- [ ] **Step 1: 写失败测试**(核心不变量:乱序到达下统计严格按阶梯序) + +```python +"""连续并发 gate 的前缀消费纯逻辑测试。""" +from __future__ import annotations + +import pytest + +from app.harness.validate import GateSpec, _advance_prefix, _GateRun, _UnitSlot +from core.evolution import GateParams +from core.types import GeneratedQuestion, QuestionUnit + + +def _mk_question(qid: str, task_type: str = "Action Reasoning") -> GeneratedQuestion: + """构造最小可用的 single 题(字段名以 core.types 实际定义为准,缺省值从简)。""" + return GeneratedQuestion( + question_id=qid, video_id="v1", task_type=task_type, + question=f"q-{qid}", options=["A. x", "B. y"], answer="A", + ) + + +def _mk_unit(qid: str, task_type: str = "Action Reasoning") -> QuestionUnit: + return QuestionUnit(unit_id=qid, questions=[_mk_question(qid, task_type)]) + + +def _mk_run(n_units: int) -> _GateRun: + spec = GateSpec( + task_type="Action Reasoning", target_file="action-reasoning.md", + candidate_content="cand", base_skill_content="base", + units=[_mk_unit(f"q{i}") for i in range(n_units)], + gate_run_prefix="r_e1_s0_gate_action-reasoning", + ) + return _GateRun.from_spec(spec) + + +_PARAMS = GateParams( + e_confirm=20.0, e_provisional=3.0, w_net_min=2, + delta_min=0.02, lambda_dir=-0.642, e_rollback=10.0, +) + + +def test_prefix_blocks_on_unresolved_head() -> None: + """阶梯头部单元未配齐时,即使尾部全部配齐也一个都不消费。""" + run = _mk_run(4) + for i in (1, 2, 3): # 尾部三个先到 + run.slots[i].base = False + run.slots[i].cand_per_q = {f"q{i}": True} + _advance_prefix(run, _PARAMS) + assert run.n_used == 0 and run.w == 0 and run.verdict is None + + +def test_prefix_consumes_in_ladder_order_after_head_arrives() -> None: + """头部补齐后一次性顺序消费到最长已配齐前缀。""" + run = _mk_run(4) + for i in (0, 1, 2): + run.slots[i].base = False + run.slots[i].cand_per_q = {f"q{i}": True} + _advance_prefix(run, _PARAMS) + assert run.n_used == 3 and run.w == 3 and run.l == 0 + assert [r["ladder_rank"] for r in run.evidence_rows] == [0, 1, 2] + + +def test_freeze_on_terminal_verdict_stops_consumption() -> None: + """过线(此处 8 连胜 E≥20 confirmed)即冻结,后续已配齐单元不再消费。""" + run = _mk_run(12) + for i in range(12): + run.slots[i].base = False + run.slots[i].cand_per_q = {f"q{i}": True} + _advance_prefix(run, _PARAMS) + assert run.frozen and run.verdict is not None + assert run.verdict.decision == "accept_confirmed" + assert run.n_used < 12 # 早停,没吃满 + + +def test_infra_unit_skipped_not_counted() -> None: + """INFRA 单元(任一臂)剔除:不入 (W,L)、计入 n_excluded、前缀继续推进。 + + 注意 futility 出口在小 n_remaining 下很敏感,用 6 单元(首个 INFRA、其余 5 个 + W 翻转)保证消费全程不提前触发 futility:题尽走 accept_provisional 终态。 + """ + run = _mk_run(6) + run.slots[0].base_infra = True + run.slots[0].cand_per_q = {"q0": True} + for i in range(1, 6): + run.slots[i].base = False + run.slots[i].cand_per_q = {f"q{i}": True} + _advance_prefix(run, _PARAMS) + assert run.n_excluded == 1 and run.n_used == 5 and run.w == 5 and run.l == 0 + assert run.verdict is not None and run.verdict.decision == "accept_provisional" + + +def test_ties_hit_futility_early_and_freeze() -> None: + """全打平(无翻转对)时 futility 出口尽早触发并冻结——早停语义(数值:W=L=0 + 时乐观 E = E(n_remaining, 0),n 小易 bool: + """双臂均已出结果(含 INFRA 判定)。""" + base_done = self.base is not None or self.base_infra + cand_done = self.cand_per_q is not None or self.cand_infra + return base_done and cand_done + + def excluded(self) -> bool: + """任一臂 INFRA 即整单元剔除(不入配对)。""" + return self.base_infra or self.cand_infra + + +@dataclass +class _GateRun: + """单题型 gate 的运行时状态(计数器 + 前缀指针 + 证据)。""" + + spec: GateSpec + slots: list[_UnitSlot] + s_hash: str + prefix_ptr: int = 0 + w: int = 0 + l: int = 0 # noqa: E741 + n_used: int = 0 + n_excluded: int = 0 + errors: int = 0 + infra_denom: int = 0 + frozen: bool = False + verdict: GateVerdict | None = None + base_obs: dict[str, bool] = field(default_factory=dict) + cand_obs: dict[str, bool] = field(default_factory=dict) + candidate_per_q: dict[str, bool] = field(default_factory=dict) + evidence_rows: list[dict] = field(default_factory=list) + + @classmethod + def from_spec(cls, spec: GateSpec) -> _GateRun: + """由规格构造初始状态(slots 与阶梯序一一对应)。""" + return cls( + spec=spec, + slots=[_UnitSlot(unit=u) for u in spec.units], + s_hash=skill_hash(spec.base_skill_content), + ) + + +def _advance_prefix(run: _GateRun, params: GateParams) -> None: + """沿阶梯序消费"已配齐前缀",逐单元更新 (W,L) 并判定,过线即冻结。 + + 统计合法性关键(设计 v3 §1 / Codex C1):严禁按到达序消费——base 臂缓存命中 + 瞬间返回、cand 臂必新鲜跑,两臂延迟不对称,若 cand 延迟与对错相关,早到翻转 + 对系统性偏向 W 型 → e-值虚高假接受。前缀消费把判定顺序钉回预声明阶梯序, + anytime-valid 无条件成立;INFRA 单元视为"已解决(剔除)"不阻塞前缀。 + """ + while not run.frozen and run.prefix_ptr < len(run.slots): + slot = run.slots[run.prefix_ptr] + if not slot.resolved(): + return + rank = run.prefix_ptr + run.prefix_ptr += 1 + if slot.excluded(): + run.n_excluded += 1 + continue + uid = slot.unit.unit_id + assert slot.base is not None and slot.cand_per_q is not None + c_units = unit_correctness_view([slot.unit], slot.cand_per_q) + pair_result = pair_block({uid: slot.base}, c_units, [uid]) + run.candidate_per_q.update(slot.cand_per_q) + for u, (b, c) in pair_result.observed.items(): + run.base_obs[u] = b + run.cand_obs[u] = c + run.w += pair_result.w + run.l += pair_result.l + run.n_used += 1 + n_remaining = (len(run.slots) - run.n_excluded) - run.n_used + run.verdict = gate_decision(run.w, run.l, run.n_used, n_remaining, params=params) + run.evidence_rows.append( + { + "question_id": uid, + "task_type": run.spec.task_type, + "ladder_rank": rank, + "baseline_correct": slot.base, + "candidate_correct": c_units[uid], + "e_value": run.verdict.e_value, + "stop_reason": "", + } + ) + if run.verdict.decision != "continue": + run.frozen = True +``` + +- [ ] **Step 4: 跑测试确认通过** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_prefix.py -x -q` +Expected: 5 passed + +- [ ] **Step 5: Commit** + +```bash +git add app/harness/validate.py tests/unit/test_gate_prefix.py +git commit -m "feat: gate prefix-ordered consumption core (algo #6)" +``` + +--- + +### Task 2: 题槽并发闸 + 单元臂执行任务 + +**Files:** +- Modify: `app/harness/validate.py`(继续追加) +- Test: `tests/unit/test_gate_unit_arm.py`(新建) + +- [ ] **Step 1: 写失败测试** + +```python +"""单元臂执行任务测试:缓存命中/新鲜跑/INFRA/冻结跳过/题槽并发上限。""" +from __future__ import annotations + +import asyncio +from pathlib import Path + +import pytest + +from app.harness.gate_ladder import BaselineCache +from app.harness.validate import ( + GateSpec, _GateRun, _QuestionSlots, _run_unit_arm, +) +from core.evolution import GateParams + +from tests.unit.test_gate_prefix import _mk_unit, _PARAMS # 复用 fixture + + +class _FakeLog: + """假 HarnessLog:query 返回预置 predictions 行。""" + + def __init__(self) -> None: + self.rows: list[dict] = [] + + def query(self, sql: str, params: tuple = ()) -> list[dict]: + run_id = params[0] + return [r for r in self.rows if r["run_id"] == run_id] + + +def _mk_gate_run(n: int, tmp_path: Path) -> tuple[_GateRun, BaselineCache]: + spec = GateSpec( + task_type="Action Reasoning", target_file="action-reasoning.md", + candidate_content="cand", base_skill_content="base", + units=[_mk_unit(f"q{i}") for i in range(n)], + gate_run_prefix="r_e1_s0_gate_action-reasoning", + ) + return _GateRun.from_spec(spec), BaselineCache(tmp_path / "bc.json") + + +def _fake_run_inference(log: _FakeLog, correct: bool, stop_reason: str = "finished"): + """构造假推理:把每题结果写进 _FakeLog 并返回带 total 的结果对象。""" + + class _R: + def __init__(self, run_id: str, total: int) -> None: + self.run_id = run_id + self.total = total + + async def _run(questions, *, run_id: str, skills_dir: Path): + for q in questions: + log.rows.append({ + "run_id": run_id, "question_id": q.question_id, + "prediction": "A" if correct else "B", "answer": "A", + "stop_reason": stop_reason, "steps_json": "[]", + }) + return _R(run_id, len(questions)) + + return _run + + +@pytest.mark.asyncio +async def test_base_arm_cache_hit_skips_inference(tmp_path) -> None: + """base 臂缓存命中:不调推理,slot.base 直接就位,infra_denom 不增。""" + run, cache = _mk_gate_run(1, tmp_path) + cache.put("Action Reasoning", run.s_hash, "v1", "q0", True) + called = {"n": 0} + + async def _boom(questions, *, run_id, skills_dir): + called["n"] += 1 + raise AssertionError("缓存命中不应触发推理") + + slots = _QuestionSlots(4) + await _run_unit_arm( + run, 0, "base", slots, _boom, _FakeLog(), cache, "v1", + Path("/nonexistent"), Path("/nonexistent"), _PARAMS, 0.10, + ) + assert called["n"] == 0 and run.slots[0].base is True and run.infra_denom == 0 + + +@pytest.mark.asyncio +async def test_base_arm_fresh_run_writes_cache(tmp_path) -> None: + """base 臂 miss 新鲜跑:结果折叠入 slot 并回写缓存。""" + run, cache = _mk_gate_run(1, tmp_path) + log = _FakeLog() + slots = _QuestionSlots(4) + await _run_unit_arm( + run, 0, "base", slots, _fake_run_inference(log, correct=True), log, + cache, "v1", tmp_path, tmp_path, _PARAMS, 0.10, + ) + assert run.slots[0].base is True + assert cache.get("Action Reasoning", run.s_hash, "v1", "q0") is True + assert run.infra_denom == 1 + + +@pytest.mark.asyncio +async def test_infra_arm_marks_excluded_and_no_cache(tmp_path) -> None: + """INFRA 臂:标记 infra、errors+1、不写缓存。""" + run, cache = _mk_gate_run(1, tmp_path) + log = _FakeLog() + slots = _QuestionSlots(4) + await _run_unit_arm( + run, 0, "base", slots, _fake_run_inference(log, correct=False, stop_reason="error"), + log, cache, "v1", tmp_path, tmp_path, _PARAMS, 0.10, + ) + assert run.slots[0].base_infra and run.errors == 1 + assert cache.get("Action Reasoning", run.s_hash, "v1", "q0") is None + + +@pytest.mark.asyncio +async def test_frozen_run_skips_launch(tmp_path) -> None: + """已冻结题型的排队臂:直接返回,不占槽不推理。""" + run, cache = _mk_gate_run(1, tmp_path) + run.frozen = True + called = {"n": 0} + + async def _boom(questions, *, run_id, skills_dir): + called["n"] += 1 + + await _run_unit_arm( + run, 0, "cand", _QuestionSlots(4), _boom, _FakeLog(), cache, "v1", + tmp_path, tmp_path, _PARAMS, 0.10, + ) + assert called["n"] == 0 + + +@pytest.mark.asyncio +async def test_question_slots_caps_inflight() -> None: + """题槽闸:峰值在飞数严格 ≤ 宽度(多槽获取不交错死锁)。""" + slots = _QuestionSlots(2) + peak = {"cur": 0, "max": 0} + + async def _job(n: int) -> None: + await slots.acquire(n) + peak["cur"] += n + peak["max"] = max(peak["max"], peak["cur"]) + await asyncio.sleep(0.01) + peak["cur"] -= n + slots.release(n) + + await asyncio.gather(*[_job(1) for _ in range(6)], *[_job(2) for _ in range(3)]) + assert peak["max"] <= 2 +``` + +- [ ] **Step 2: 跑测试确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_unit_arm.py -x -q` +Expected: FAIL(ImportError: cannot import name '_QuestionSlots') + +- [ ] **Step 3: 实现 `_QuestionSlots` 与 `_run_unit_arm`**(validate.py 追加) + +```python +class _QuestionSlots: + """按题数计数的共享并发闸:峰值在飞请求恒 ≤ width(设计 v3 §2.4)。 + + 多槽获取(AR pair 一单元两题)经内部锁串行化,防多任务半持有交错死锁。 + asyncio.Semaphore 等待队列 FIFO,任务按创建序(题型 round-robin)获得槽, + 即公平调度的实现载体(Codex I2)。 + """ + + def __init__(self, width: int) -> None: + assert width > 0, f"并发宽度必须为正: {width}" + self._sem = asyncio.Semaphore(width) + self._acquire_lock = asyncio.Lock() + + async def acquire(self, n: int) -> None: + """原子获取 n 个题槽。""" + async with self._acquire_lock: + for _ in range(n): + await self._sem.acquire() + + def release(self, n: int) -> None: + """归还 n 个题槽。""" + for _ in range(n): + self._sem.release() + + +async def _run_unit_arm( + run: _GateRun, + slot_idx: int, + arm: str, + slots: _QuestionSlots, + run_inference: RunInferenceFn, + log: HarnessLog, + baseline_cache: BaselineCache, + prompts_version: str, + base_skills_dir: Path, + cand_dir: Path, + gate_params: GateParams, + gate_guard_err: float, +) -> None: + """执行一个 (单元, 臂) 任务:缓存/推理 → 到达登记 → 前缀消费推进。 + + 冻结检查两次:启动时(排队任务撤销点)与获得题槽后(获槽期间被冻结)。 + base 臂缓存命中不占题槽(零推理);INFRA 单元不写缓存(不永久污染基线快照)。 + 护栏在每次臂完成时检查(等价迁移自跨块累计,设计 v3 §2.3),超阈值 raise + 中止整轮(与现行行为一致)。 + """ + assert arm in ("base", "cand") + if run.frozen: + return + slot = run.slots[slot_idx] + spec = run.spec + + if arm == "base": + cached = baseline_cache.get( + spec.task_type, run.s_hash, prompts_version, slot.unit.unit_id + ) + if cached is not None: + slot.base = cached + _advance_prefix(run, gate_params) + return + + questions = list(slot.unit.questions) + await slots.acquire(len(questions)) + try: + if run.frozen: + return + run_id = f"{spec.gate_run_prefix}_{arm}" + skills_dir = base_skills_dir if arm == "base" else cand_dir + r = await run_inference(questions, run_id=run_id, skills_dir=skills_dir) + infra_qids = _infra_question_ids_from_db(log, r.run_id, questions) + run.infra_denom += r.total + if infra_qids: + run.errors += 1 + if arm == "base": + slot.base_infra = True + else: + slot.cand_infra = True + else: + per_q = _candidate_correctness_from_db(log, r.run_id, questions) + if arm == "base": + folded = unit_correctness_view([slot.unit], per_q) + slot.base = folded[slot.unit.unit_id] + baseline_cache.put( + spec.task_type, run.s_hash, prompts_version, + slot.unit.unit_id, slot.base, + ) + else: + slot.cand_per_q = per_q + _check_infra_guard(run.errors, run.infra_denom, gate_guard_err) + finally: + slots.release(len(questions)) + _advance_prefix(run, gate_params) +``` + +- [ ] **Step 4: 跑测试确认通过** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_unit_arm.py tests/unit/test_gate_prefix.py -x -q` +Expected: 10 passed + +- [ ] **Step 5: Commit** + +```bash +git add app/harness/validate.py tests/unit/test_gate_unit_arm.py +git commit -m "feat: gate unit-arm tasks with question-slot gate (algo #6)" +``` + +--- + +### Task 3: 调度编排 `validate_skills_concurrent` + evidence 表 ladder_rank + +**Files:** +- Modify: `app/harness/validate.py`(追加编排函数) +- Modify: `app/harness/observation.py`(`_GATE_EVIDENCE_COLS` 的 `block_idx` → `ladder_rank`;同步改 docstring) +- Test: `tests/unit/test_gate_concurrent.py`(新建) + +- [ ] **Step 1: 写失败测试** + +```python +"""连续并发 gate 编排测试:乱序到达/多题型隔离/round-robin/终态组装。""" +from __future__ import annotations + +import asyncio +from pathlib import Path + +import pytest + +from app.harness.gate_ladder import BaselineCache +from app.harness.validate import GateSpec, validate_skills_concurrent +from tests.unit.test_gate_prefix import _mk_unit, _PARAMS +from tests.unit.test_gate_unit_arm import _FakeLog + + +def _mk_spec(task_type: str, slug: str, n: int) -> GateSpec: + return GateSpec( + task_type=task_type, target_file=f"{slug}.md", + candidate_content=f"cand-{slug}", base_skill_content=f"base-{slug}", + units=[_mk_unit(f"{slug}-q{i}", task_type) for i in range(n)], + gate_run_prefix=f"r_e1_s0_gate_{slug}", + ) + + +def _scripted_inference(log: _FakeLog, script: dict[str, tuple[bool, float]]): + """脚本化假推理:按 question_id 决定 (对错, 延迟秒),制造乱序到达。""" + + class _R: + def __init__(self, run_id: str, total: int) -> None: + self.run_id = run_id + self.total = total + + async def _run(questions, *, run_id: str, skills_dir: Path): + arm = "cand" if run_id.endswith("_cand") else "base" + correct, delay = script[f"{questions[0].question_id}|{arm}"] + await asyncio.sleep(delay) + for q in questions: + log.rows.append({ + "run_id": run_id, "question_id": q.question_id, + "prediction": "A" if correct else "B", "answer": "A", + "stop_reason": "finished", "steps_json": "[]", + }) + return _R(run_id, len(questions)) + + return _run + + +@pytest.mark.asyncio +async def test_out_of_order_arrival_still_ladder_order(tmp_path, monkeypatch) -> None: + """尾部先到、头部后到:判定结果与顺序到达完全相同(前缀有序性端到端)。""" + spec = _mk_spec("Action Reasoning", "action-reasoning", 4) + log = _FakeLog() + # 头部 q0 最慢;全部翻转为 W(base 错 cand 对) + script = {} + for i in range(4): + script[f"action-reasoning-q{i}|base"] = (False, 0.05 if i == 0 else 0.0) + script[f"action-reasoning-q{i}|cand"] = (True, 0.05 if i == 0 else 0.0) + monkeypatch.setattr( + "app.harness.validate.materialize_candidate_skill", + lambda *a, **k: tmp_path / "cand", + ) + outcomes = await validate_skills_concurrent( + workspace_dir=tmp_path, base_skills_version="v1", specs=[spec], + gate_params=_PARAMS, gate_guard_err=0.10, + baseline_cache=BaselineCache(tmp_path / "bc.json"), + prompts_version="v1", run_inference=_scripted_inference(log, script), + log=log, concurrency=8, + ) + o = outcomes["Action Reasoning"] + assert o.w == 4 and o.l == 0 + assert [r["ladder_rank"] for r in o.evidence_rows] == [0, 1, 2, 3] + + +@pytest.mark.asyncio +async def test_two_types_isolated(tmp_path, monkeypatch) -> None: + """两题型并行:计数互不污染,各自独立判定。 + + A 型 4 单元全 W(题尽 accept_provisional);B 型 2 单元全平 + (futility 早停,W=L=0)——两型结果都不受对方污染。 + """ + spec_a = _mk_spec("Action Reasoning", "action-reasoning", 4) + spec_b = _mk_spec("Counting Problem", "counting-problem", 2) + log = _FakeLog() + script = {} + for i in range(4): + script[f"action-reasoning-q{i}|base"] = (False, 0.0) + script[f"action-reasoning-q{i}|cand"] = (True, 0.0) + for i in range(2): + script[f"counting-problem-q{i}|base"] = (True, 0.0) + script[f"counting-problem-q{i}|cand"] = (True, 0.0) + monkeypatch.setattr( + "app.harness.validate.materialize_candidate_skill", + lambda *a, **k: tmp_path / "cand", + ) + outcomes = await validate_skills_concurrent( + workspace_dir=tmp_path, base_skills_version="v1", + specs=[spec_a, spec_b], gate_params=_PARAMS, gate_guard_err=0.10, + baseline_cache=BaselineCache(tmp_path / "bc.json"), + prompts_version="v1", run_inference=_scripted_inference(log, script), + log=log, concurrency=8, + ) + assert outcomes["Action Reasoning"].w == 4 + assert outcomes["Counting Problem"].w == 0 + assert outcomes["Counting Problem"].l == 0 + + +@pytest.mark.asyncio +async def test_all_infra_raises(tmp_path, monkeypatch) -> None: + """全单元 INFRA:保留现行 RuntimeError 语义(检查推理基础设施)。""" + spec = _mk_spec("Action Reasoning", "action-reasoning", 2) + log = _FakeLog() + + class _R: + def __init__(self, run_id, total): + self.run_id, self.total = run_id, total + + async def _infra_run(questions, *, run_id, skills_dir): + for q in questions: + log.rows.append({ + "run_id": run_id, "question_id": q.question_id, + "prediction": "", "answer": "A", + "stop_reason": "error", "steps_json": "[]", + }) + return _R(run_id, len(questions)) + + monkeypatch.setattr( + "app.harness.validate.materialize_candidate_skill", + lambda *a, **k: tmp_path / "cand", + ) + with pytest.raises(RuntimeError): + await validate_skills_concurrent( + workspace_dir=tmp_path, base_skills_version="v1", specs=[spec], + gate_params=_PARAMS, gate_guard_err=0.99, # 护栏放宽,逼出全 INFRA 分支 + baseline_cache=BaselineCache(tmp_path / "bc.json"), + prompts_version="v1", run_inference=_infra_run, + log=log, concurrency=8, + ) +``` + +- [ ] **Step 2: 跑测试确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_concurrent.py -x -q` +Expected: FAIL(ImportError: cannot import name 'validate_skills_concurrent') + +- [ ] **Step 3: 实现编排函数**(validate.py 追加) + +```python +async def validate_skills_concurrent( + workspace_dir: Path, + base_skills_version: str, + specs: list[GateSpec], + gate_params: GateParams, + gate_guard_err: float, + baseline_cache: BaselineCache, + prompts_version: str, + run_inference: RunInferenceFn, + log: HarnessLog, + concurrency: int, +) -> dict[str, ValidationOutcome]: + """连续并发 gate:多题型全部臂共享题槽并发,统计按阶梯序前缀有序推进。 + + 发射顺序 = 题型 round-robin × 题型内阶梯序(base 先 cand 后),asyncio.Semaphore + FIFO 等待队列保证获槽顺序即发射顺序 → 公平调度,大题型不饿死小题型。 + 题型过线即冻结:其排队任务启动时自查冻结标志撤销;in-flight 跑完落库但 + 不计入(τ 之后样本,合法丢弃)。全部题型判定后统一组装 ValidationOutcome。 + + 参数: + workspace_dir: workspace 根目录(候选物化用)。 + base_skills_version: 基线 skills 版本名。 + specs: 各题型 gate 规格(units 已阶梯序 + 截断 n_max)。 + gate_params / gate_guard_err: e-process 判据与 INFRA 护栏阈值。 + baseline_cache / prompts_version: 基线缓存及其键成分。 + run_inference: 注入推理函数(调用方须绑定共享 HarnessLog,见 runner)。 + log: HarnessLog 共享实例(推理后读预测,与 run_inference 同库)。 + concurrency: 题槽宽度(峰值在飞题数上限)。 + + 返回: + {task_type: ValidationOutcome}。 + + 异常: + RuntimeError: INFRA 护栏超阈值,或某题型全部单元被 INFRA 排除。 + ValueError: spec 校验失败(空阶梯 / run_prefix 缺 "_gate_")。 + """ + for spec in specs: + if "_gate_" not in spec.gate_run_prefix: + raise ValueError(f"gate_run_prefix 必须含 '_gate_': {spec.gate_run_prefix!r}") + if not spec.units: + raise ValueError(f"task_type={spec.task_type} 阶梯为空,无法验证") + + base_skills_dir = workspace_dir / "skills" / base_skills_version + runs = [_GateRun.from_spec(s) for s in specs] + cand_dirs = { + r.spec.task_type: materialize_candidate_skill( + workspace_dir, base_skills_version, r.spec.target_file, r.spec.candidate_content + ) + for r in runs + } + slots_gate = _QuestionSlots(concurrency) + try: + # 发射队列:round-robin 交错(rank 0 各题型 → rank 1 各题型 → ...) + coros = [] + max_rank = max(len(r.slots) for r in runs) + for rank in range(max_rank): + for r in runs: + if rank < len(r.slots): + for arm in ("base", "cand"): + coros.append( + _run_unit_arm( + r, rank, arm, slots_gate, run_inference, log, + baseline_cache, prompts_version, + base_skills_dir, cand_dirs[r.spec.task_type], + gate_params, gate_guard_err, + ) + ) + await asyncio.gather(*coros) + finally: + for d in cand_dirs.values(): + try: + shutil.rmtree(d) + except OSError as e: + logger.warning("候选临时目录清理失败 {}: {}", d, e) + + outcomes: dict[str, ValidationOutcome] = {} + for r in runs: + if r.verdict is None: + raise RuntimeError( + f"gate[{r.spec.task_type}] 全部 unit 被判为 INFRA 排除,无法验证(检查推理基础设施)" + ) + outcomes[r.spec.task_type] = _finalize_outcome( + verdict=r.verdict, w=r.w, l=r.l, n_used=r.n_used, n_plan=len(r.slots), + base_obs=r.base_obs, cand_obs=r.cand_obs, + candidate_per_q=r.candidate_per_q, evidence_rows=r.evidence_rows, + task_type=r.spec.task_type, + ) + return outcomes +``` + +顶部 import 需补 `import asyncio`。 + +- [ ] **Step 4: `observation.py` 改 evidence 列** + +找到 `_GATE_EVIDENCE_COLS` 定义(`grep -n "_GATE_EVIDENCE_COLS" app/harness/observation.py`),把 `"block_idx"` 列改名 `"ladder_rank"`(类型不变 INTEGER),`write_gate_evidence` docstring 中 `block_idx` 描述同步改为 `ladder_rank(阶梯序号)`。 + +- [ ] **Step 5: 跑测试确认通过** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_concurrent.py tests/unit/test_gate_prefix.py tests/unit/test_gate_unit_arm.py -x -q` +Expected: 13 passed +注意:`test_out_of_order_arrival_still_ladder_order` 里 cand 臂识别依赖 `run_id.endswith("_cand")`,与实现的 run_id 命名 `{prefix}_{arm}` 一致;若失败先核对 run_id 拼接。 + +- [ ] **Step 6: Commit** + +```bash +git add app/harness/validate.py app/harness/observation.py tests/unit/test_gate_concurrent.py +git commit -m "feat: continuous concurrent gate orchestrator (algo #6)" +``` + +--- + +### Task 4: step 重跑幂等修复(现行潜伏 bug) + +**Files:** +- Modify: `app/harness/runner.py:1105-1109`(`_run_step` 开头的清行块) +- Test: `tests/unit/test_step_rerun_idempotent.py`(新建) + +- [ ] **Step 1: 写失败测试** + +```python +"""step 重跑幂等:gate 派生行必须随 step 清理,否则崩溃重跑累积重复。""" +from __future__ import annotations + +import sqlite3 +from pathlib import Path + +from app.harness.runner import _clear_step_rows + + +def _mk_db(tmp_path: Path) -> Path: + db = tmp_path / "harness.db" + conn = sqlite3.connect(db) + conn.execute("CREATE TABLE predictions (run_id TEXT, question_id TEXT)") + conn.execute("CREATE TABLE traces (run_id TEXT, question_id TEXT)") + conn.execute("CREATE TABLE gate_evidence (run_id TEXT, epoch INTEGER, step INTEGER)") + conn.execute("CREATE TABLE quadrant_pair (run_id TEXT, epoch INTEGER, step INTEGER)") + rows = [ + ("infer_adhoc_e1_s0", "q1"), # rollout 行 + ("infer_adhoc_e1_s0_gate_action-reasoning_base", "q2"), # gate base 臂 + ("infer_adhoc_e1_s0_gate_action-reasoning_cand", "q3"), # gate cand 臂 + ("infer_adhoc_e1_s1", "q4"), # 其他 step,不许误删 + ("infer_adhoc_e1_s10_gate_x_base", "q5"), # s10 前缀陷阱,不许误删 + ] + conn.executemany("INSERT INTO predictions VALUES (?, ?)", rows) + conn.executemany("INSERT INTO traces VALUES (?, ?)", rows) + conn.execute("INSERT INTO gate_evidence VALUES ('infer_adhoc', 1, 0)") + conn.execute("INSERT INTO gate_evidence VALUES ('infer_adhoc', 1, 1)") + conn.execute("INSERT INTO quadrant_pair VALUES ('infer_adhoc', 1, 0)") + conn.commit() + conn.close() + return db + + +def test_clear_step_rows_removes_rollout_and_gate_rows(tmp_path) -> None: + db = _mk_db(tmp_path) + _clear_step_rows(str(db), baseline_run_id="infer_adhoc", epoch=1, step=0) + conn = sqlite3.connect(db) + left = {r[0] for r in conn.execute("SELECT run_id FROM predictions")} + assert left == {"infer_adhoc_e1_s1", "infer_adhoc_e1_s10_gate_x_base"} + left_t = {r[0] for r in conn.execute("SELECT run_id FROM traces")} + assert left_t == left + ge = list(conn.execute("SELECT step FROM gate_evidence")) + assert ge == [(1,)] # 只剩 step=1 的行 + assert list(conn.execute("SELECT COUNT(*) FROM quadrant_pair"))[0][0] == 0 + conn.close() + + +def test_clear_step_rows_missing_tables_is_noop(tmp_path) -> None: + """gate_evidence/quadrant_pair 表尚未建(首个 step)时不报错。""" + db = tmp_path / "harness.db" + conn = sqlite3.connect(db) + conn.execute("CREATE TABLE predictions (run_id TEXT)") + conn.execute("CREATE TABLE traces (run_id TEXT)") + conn.commit() + conn.close() + _clear_step_rows(str(db), baseline_run_id="infer_adhoc", epoch=1, step=0) +``` + +- [ ] **Step 2: 跑测试确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_step_rerun_idempotent.py -x -q` +Expected: FAIL(ImportError: cannot import name '_clear_step_rows') + +- [ ] **Step 3: 实现 `_clear_step_rows`(runner.py 模块级函数)并替换 `_run_step` 内清行块** + +新函数(放在 `_write_skip_report` 附近的模块级函数区): + +```python +def _clear_step_rows(db_path: str, *, baseline_run_id: str, epoch: int, step: int) -> None: + """清空一个 step 的全部旧行(rollout + gate 派生),保证崩溃重跑幂等。 + + 修复前序潜伏 bug:旧实现只清 rollout run_id,gate 派生 run_id + (`{step_run_id}_gate_%`)从不清理,重跑会累积重复 predictions(HarnessLog + 无主键去重),_load_run_rows 的 dict 覆盖使结果依赖 SELECT 顺序。 + gate_evidence / quadrant_pair 以 (run_id, epoch, step) 过滤删除; + 表不存在(首个 step)时跳过。step_report 为按文件名覆盖写的 JSON,天然幂等。 + """ + from app.harness.inference import PREDICTIONS_SCHEMA, TRACES_SCHEMA + from app.harness.log import HarnessLog + + step_run_id = f"{baseline_run_id}_e{epoch}_s{step}" + with HarnessLog(db_path, step_run_id, register_run=False) as log: + log.create_table("predictions", PREDICTIONS_SCHEMA) + log.create_table("traces", TRACES_SCHEMA) + for table in ("predictions", "traces"): + log.execute(f"DELETE FROM {table} WHERE run_id=?", (step_run_id,)) + # ESCAPE 显式声明,防 run_id 中出现 '_' 通配歧义(_gate_ 前缀含字面下划线, + # LIKE 的 '_' 单字符通配在此无害但语义须钉死为字面匹配 + '%' 后缀) + log.execute( + f"DELETE FROM {table} WHERE run_id LIKE ? ESCAPE '\\'", + (step_run_id.replace("_", r"\_") + r"\_gate\_%",), + ) + for table in ("gate_evidence", "quadrant_pair"): + exists = log.query( + "SELECT name FROM sqlite_master WHERE type='table' AND name=?", (table,) + ) + if exists: + log.execute( + f"DELETE FROM {table} WHERE run_id=? AND epoch=? AND step=?", + (baseline_run_id, epoch, step), + ) +``` + +`_run_step` 内原 `with HarnessLog(...) as log: ... DELETE ...` 五行(runner.py:1105-1109)替换为一行: + +```python + _clear_step_rows( + str(self._paths.db_path), + baseline_run_id=pools.baseline_run_id, epoch=epoch, step=step, + ) +``` + +- [ ] **Step 4: 跑测试确认通过 + 回归** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_step_rerun_idempotent.py tests/unit -x -q` +Expected: 全绿(重点:runner 相关既有测试不破) + +- [ ] **Step 5: Commit** + +```bash +git add app/harness/runner.py tests/unit/test_step_rerun_idempotent.py +git commit -m "fix: clear gate-derived rows on step rerun (idempotency)" +``` + +--- + +### Task 5: runner 重接线 —— 并行进化 + 装配 + 连续 gate + 字母序落账 + +**Files:** +- Modify: `app/harness/runner.py`(`_gate_batch_skills`:1143-1252 重写;`_run_gate_validation`:1254-1334 删除;`_make_validate_run_inference_fn` 改共享 log) +- Test: `tests/unit/test_gate_batch_parallel.py`(新建) + +- [ ] **Step 1: 写失败测试**(target_file 冲突 fail-fast 是纯逻辑,可直接测) + +```python +"""_gate_batch_skills 并行装配的纯逻辑护栏测试。""" +from __future__ import annotations + +import pytest + +from app.harness.runner import _assert_disjoint_target_files + + +def test_disjoint_target_files_pass() -> None: + _assert_disjoint_target_files( + {"Action Reasoning": "action-reasoning.md", "Counting Problem": "counting-problem.md"} + ) + + +def test_shared_target_file_fails_fast() -> None: + """两题型 fallback 到同一文件:并行进化会互相覆盖,必须 fail-fast。""" + with pytest.raises(RuntimeError, match="default-strategy.md"): + _assert_disjoint_target_files( + {"OCR Problems": "default-strategy.md", "Spatial Reasoning": "default-strategy.md"} + ) +``` + +- [ ] **Step 2: 跑测试确认失败** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_batch_parallel.py -x -q` +Expected: FAIL(ImportError) + +- [ ] **Step 3: 实现 runner 侧改造** + +3a. 模块级函数(放 `_write_skip_report` 附近): + +```python +def _assert_disjoint_target_files(targets_by_type: dict[str, str]) -> None: + """断言本 step 各题型进化目标文件互不相同(设计 v3 §1 fail-fast)。 + + 题型并行进化 + 并行 gate 的前提是 skill 文件不相交;两题型 fallback 到 + 同一 default-strategy.md 时并行会互相覆盖候选与 accept,必须显式中止 + 而非静默串行(当前 12 题型均有专属文件,此断言防未来配置漂移)。 + """ + seen: dict[str, str] = {} + for task_type, target in targets_by_type.items(): + if target in seen: + raise RuntimeError( + f"题型 {seen[target]!r} 与 {task_type!r} 映射同一 skill 文件 {target!r}," + "并行进化/gate 不支持共享目标文件(设计 v3 §1)" + ) + seen[target] = task_type +``` + +3b. `_make_validate_run_inference_fn` 改为共享单一 HarnessLog(消除 per-unit 新建连接的锁竞争,对齐单连接+锁模式): + +```python + def _make_validate_run_inference_fn(self, gate_log: HarnessLog): + """构造 validate 用的 RunInferenceFn(绑定共享依赖与共享 HarnessLog)。 + + 连续并发 gate 下本函数被逐单元高频并发调用:每次调用新建 HarnessLog + 连接会重现多连接争 SQLite 写锁(遥测同款教训),故复用调用方传入的 + 单一 gate_log(单连接 + threading.Lock 串行化)。_record_run 按 run_id + 去重,避免逐单元重复 upsert。 + """ + from app.harness.inference import run_inference + + recorded: set[str] = set() + + async def _run( + questions: list[GeneratedQuestion], + *, + run_id: str, + skills_dir: Path, + ) -> InferenceResult: + if run_id not in recorded: + recorded.add(run_id) + self._record_run(run_id) + return await run_inference( + questions=questions, + llm=self._llm, + tool_dispatch_fn=self._make_tool_dispatch_fn(skills_dir=skills_dir), + prompt_builder=self._make_prompt_builder( + skills_dir=skills_dir, prompts_dir=self._paths.prompts_dir + ), + log=gate_log, + run_id=run_id, + concurrency=self._config.concurrency, + max_steps=self._config.max_steps, + skill_mode=self._config.skill_mode, + ) + + return _run +``` + +3c. `_gate_batch_skills` 重写为四阶段(替换 1143-1252 的串行 for;`_run_gate_validation` 整个删除,其阶梯装配逻辑收编进 Phase B): + +```python + async def _gate_batch_skills( + self, + epoch: int, + step: int, + diagnosis: DiagnosisResult, + total_steps: int, + pools: Pools, + state: _TrainState, + ) -> None: + """按 task_type 并行 evolve → 连续并发 gate → 字母序统一落账。 + + 四阶段(设计 v3 §2.1):Phase A 并行进化(cooldown/无改动照旧跳过); + Phase B 装配 GateSpec(阶梯出题 + 案例单元排除 + n_max 截断); + Phase C validate_skills_concurrent(共享题槽,统计按阶梯序前缀); + Phase D 唯一写 state 阶段——按字母序 accept/reject 落账,与原串行 + 语义等价(题型 skill 文件不相交,合并顺序仅为确定性)。 + """ + from app.harness.workspace import VersionedSkillStore + from core.evolution import evolve_single_skill + + budget = edit_budget_at( + global_step=state.global_step, + total_steps=total_steps, + start=self._config.edit_budget_start, + end=self._config.edit_budget_end, + ) + + # ---- Phase A: 并行进化(冷却/跳过路径先出清) ---- + active_types: list[str] = [] + for task_type in sorted(diagnosis.skill_case_packs): + if state.gate_cooldown.get(task_type, 0) > 0: + _write_skip_report( + self._config.workspace_dir, epoch, step, state.global_step, + task_type, action="cooldown", + baseline_acc=self._class_baseline_acc( + task_type, pools.validation, state.correctness + ), + budget=budget, + ) + continue + active_types.append(task_type) + if not active_types: + return + + evolve_prompts = self._load_evolve_prompts() + skills_version = self._current_version("skills") + + async def _evolve_one(task_type: str) -> EvolutionRecord: + pack = diagnosis.skill_case_packs[task_type] + skill_store = VersionedSkillStore(self._paths.skills_dir) + return await evolve_single_skill( + self._evolve_llm, pack, skill_store, evolve_prompts, + skills_version, budget, self._config.appendix_consolidate_threshold, + skill_update_mode=self._config.skill_update_mode, + rejected=state.rejected_buffer.get(task_type, []), + ) + + records = dict( + zip( + active_types, + await asyncio.gather(*[_evolve_one(t) for t in active_types]), + strict=True, + ) + ) + + # 无真实改动的题型照旧写 skipped 后出队 + gated_types: list[str] = [] + for task_type in active_types: + record = records[task_type] + if record.status in ("rejected", "skipped") or ( + record.evolved_content == record.original_content + ): + _write_skip_report( + self._config.workspace_dir, epoch, step, state.global_step, + task_type, action="skipped", + baseline_acc=self._class_baseline_acc( + task_type, pools.validation, state.correctness + ), + budget=budget, + rank_clip_triggered=bool(record.clip_info.get("triggered", False)), + ) + continue + gated_types.append(task_type) + if not gated_types: + return + + _assert_disjoint_target_files({t: records[t].target_file for t in gated_types}) + + # ---- Phase B: 装配 GateSpec(阶梯出题,收编原 _run_gate_validation 前半) ---- + specs: list[GateSpec] = [] + for task_type in gated_types: + record = records[task_type] + pack = diagnosis.skill_case_packs[task_type] + exclude_units = { + self._gate_questions_by_id[c.question_id].unit_id + for c in pack.failure_cases + pack.success_cases + if c.question_id in self._gate_questions_by_id + } + ladder_unit_ids = state.gate_pools.ladder_for( + task_type, + exclude_units, + p_low=self._config.gate_p_low, + p_high=self._config.gate_p_high, + cold=not state.gate_epoch_observed, + ) + missing = [uid for uid in ladder_unit_ids if uid not in self._gate_units_by_id] + if missing: + raise RuntimeError( + f"gate 阶梯引用未知 unit: {missing[:5]}(gate_pools.json 与题库失配)" + ) + ladder_items = [ + q for uid in ladder_unit_ids for q in self._gate_units_by_id[uid].questions + ] + slug = task_type.lower().replace(" ", "-") + specs.append( + GateSpec( + task_type=task_type, + target_file=record.target_file, + candidate_content=record.evolved_content, + base_skill_content=( + self._paths.skills_dir / record.target_file + ).read_text(encoding="utf-8"), + units=_ladder_units(ladder_items)[: self._config.gate_n_max], + gate_run_prefix=f"{pools.baseline_run_id}_e{epoch}_s{step}_gate_{slug}", + ) + ) + + # ---- Phase C: 连续并发 gate(只读 state) ---- + with HarnessLog(str(self._paths.db_path), f"gate_e{epoch}_s{step}") as gate_log: + outcomes = await validate_skills_concurrent( + workspace_dir=self._config.workspace_dir, + base_skills_version=self._current_version("skills"), + specs=specs, + gate_params=GateParams( + e_confirm=self._config.gate_e_confirm, + e_provisional=self._config.gate_e_provisional, + w_net_min=self._config.gate_w_net_min, + delta_min=self._config.gate_delta_min, + lambda_dir=self._config.gate_lambda_dir, + e_rollback=self._config.gate_e_rollback, + ), + gate_guard_err=self._config.gate_guard_err, + baseline_cache=state.baseline_cache, + prompts_version=self._current_version("prompts"), + run_inference=self._make_validate_run_inference_fn(gate_log), + log=gate_log, + concurrency=self._config.concurrency, + ) + + # ---- Phase D: 唯一写 state 阶段(字母序确定性落账) ---- + for task_type in sorted(outcomes): + record = records[task_type] + outcome = outcomes[task_type] + write_gate_evidence( + str(self._paths.db_path), + run_id=pools.baseline_run_id, epoch=epoch, step=step, + rows=outcome.evidence_rows, + ) + write_step_report( + self._config.workspace_dir, + epoch=epoch, step=step, global_step=state.global_step, + task_type=task_type, gate_action=outcome.action, + candidate_acc=outcome.candidate_acc, + class_baseline_acc=outcome.baseline_acc, + edit_budget=budget, + rank_clip_triggered=bool(record.clip_info.get("triggered", False)), + gate_w=outcome.w, gate_l=outcome.l, gate_e_value=outcome.e_value, + gate_n_used=outcome.n_used, gate_stop_reason=outcome.stop_reason, + ) + write_quadrant_pairs( + str(self._paths.db_path), + run_id=pools.baseline_run_id, epoch=epoch, step=step, + pairs=_outcome_to_quadrant_pairs(task_type, outcome), + ) + if outcome.accepted: + self._accept_skill(task_type, record, outcome, state, pools) + else: + self._record_rejected_skill( + state.rejected_buffer, task_type, record, outcome, state.global_step + ) +``` + +导入区补:`from app.harness.validate import GateSpec, _ladder_units, validate_skills_concurrent`(与既有 validate 导入合并;`_ladder_units` 若被视为私有,可在 validate.py 里改名 `ladder_units` 导出并同步旧调用点)。 + +**语义等价性注意(实现者必读):** +- 原串行版中,题型 B 的 `base_skill_content` / `base_skills_version` 读的是"A accept 之后"的版本;新版全部读 step 起点版本。设计 v3 §1 已论证内容等价(文件不相交)+ 用户批准。**但** `_accept_skill` 内部按当前 manifest 版本推进——Phase D 串行调用时每次 accept 后版本前移,第二个 accept 基于第一个 accept 后的版本追加自己的 target_file 内容,文件不相交所以互不覆盖。逐一确认 `_accept_skill`(runner.py:1340 起)满足此性质:它以 `record.evolved_content` 写 `target_file`,复制其余文件自当前版本 → 满足。 +- `gate_epoch_observed` / cooldown 递减等 step 级状态,原逻辑不在本函数内改动的,保持不动。 + +- [ ] **Step 4: 跑测试 + 全量回归** + +Run: `conda run -n Video-Tree-TRM pytest tests/unit/test_gate_batch_parallel.py tests/unit -x -q` +Expected: 新测试过;既有 runner/gate 测试中直接调用 `_run_gate_validation` 或 `validate_skill_local` 的会失败——**属预期**,Task 6 处理;本步允许用 `--deselect` 记录清单但不得改产品代码迁就旧测试。 + +- [ ] **Step 5: Commit** + +```bash +git add app/harness/runner.py tests/unit/test_gate_batch_parallel.py +git commit -m "feat: parallel evolve + continuous gate wiring in runner (algo #6)" +``` + +--- + +### Task 6: 删除旧块路径 + gate_block 全线移除 + 旧测试迁移 + +**Files:** +- Modify: `app/harness/validate.py`(删 `_run_local_validation`、`validate_skill_local`、`_resolve_baseline_block`、`_run_candidate_block`、`_build_evidence_rows`) +- Modify: `app/harness/config.py:73,128,372-375`(删 gate_block 字段/文档/校验;校验改为 `gate_n_max > 0`) +- Modify: `app/harness/checkpoint.py:61`(从字段清单删 `"gate_block"`) +- Modify: `config/*.yaml` 全部含 `gate_block:` 的文件(train_videomme / default / question_gen_180_补 / question_gen_360 / train_action_recognition / train_ar30)删该行 +- Modify/Delete: `tests/unit/test_gate_block_unit.py`、`tests/unit/test_harness_validate.py` 等引用被删符号的测试 + +- [ ] **Step 1: 全局定位被删符号的引用** + +Run: `grep -rn "validate_skill_local\|_run_local_validation\|_resolve_baseline_block\|_run_candidate_block\|_build_evidence_rows\|gate_block" app/ core/ config/ tests/ main.py --include="*.py" --include="*.yaml"` +逐一处置:产品代码引用应已在 Task 5 清零(若有残留即 Task 5 遗漏,回去补);测试引用见 Step 2。 + +- [ ] **Step 2: 迁移旧测试** + +原则:**测的语义保留,测的载体更新**。 +- 纯配对/四象限/单元折叠测试(`pair_block`/`classify_quadrants`/`unit_correctness_view`,在 test_validate.py 等):不动,它们测 core 纯函数。 +- 块循环编排测试(`test_gate_block_unit.py`、`test_harness_validate.py` 中调 `validate_skill_local` 的用例):改写为经 `validate_skills_concurrent` 的等价断言(单 spec 输入,断言 action/W/L/E/evidence 与原期望一致;块边界断言删除,替换为 ladder_rank 连续性断言)。改写时保持原测试意图注释。 +- INFRA 护栏测试:改为构造假推理触发 `_check_infra_guard` 阈值,断言 RuntimeError(语义同前)。 + +- [ ] **Step 3: 删除产品代码旧路径与 gate_block** + +按 Files 清单逐个删除;`config.py` 校验块 372-375 改为: + +```python + if config.gate_n_max <= 0: + raise ValueError(f"需 gate_n_max > 0,实际: n_max={config.gate_n_max}") +``` + +- [ ] **Step 4: 全量回归 + lint** + +Run: `conda run -n Video-Tree-TRM pytest tests/ -x -q && conda run -n Video-Tree-TRM ruff check app/ core/ adapters/ && conda run -n Video-Tree-TRM ruff format --check app/ core/ adapters/` +Expected: 全绿、零 lint 错误。覆盖率:`conda run -n Video-Tree-TRM pytest tests/ --cov=app --cov=core -q | tail -3` ≥ 80%。 + +- [ ] **Step 5: Commit** + +```bash +git add -u app/ config/ tests/ +git commit -m "refactor: remove block-sequential gate path and gate_block knob (algo #6)" +``` + +--- + +### Task 7: Redis TTL + 收尾验证 + +**Files:** +- Modify: `.env`(仅 `REDIS_CACHE_TTL` 一行,**严禁**触碰其他行) +- Modify: `.env.example`(同键同值,保持模板同步) + +- [ ] **Step 1: 改 TTL** + +```bash +sed -i 's/^REDIS_CACHE_TTL=86400$/REDIS_CACHE_TTL=604800/' .env +sed -i 's/^REDIS_CACHE_TTL=.*/REDIS_CACHE_TTL=604800/' .env.example +grep -n "^REDIS_CACHE_TTL=" .env .env.example +``` + +Expected: 两文件均为 604800。 + +- [ ] **Step 2: 全量测试终验** + +Run: `conda run -n Video-Tree-TRM pytest tests/ -q --cov=app --cov=core --cov-report=term-missing | tail -15` +Expected: 全绿,覆盖率 ≥80%。 + +- [ ] **Step 3: Commit**(.env 不入库,仅 .env.example) + +```bash +git add .env.example +git commit -m "chore: raise Redis cache TTL to 7 days" +``` + +--- + +## 运维 Runbook(实现完成后、重启训练前,控制器执行,不属实现任务) + +1. **一次性续期今日 Redis 键**(scratchpad 脚本,用 dotenv 读 REDIS_URL): + `for k in scan_iter('llm_cache:*'): 0 < ttl(k) < 604800 → expire(k, 604800)` +2. **重启训练**: + `tmux new-session -d -s train_videomme "CUDA_VISIBLE_DEVICES=0 bash scripts/train_videomme.sh 2>&1 | tee logs/train_videomme.log"` + (config 的 run_id 保持 `train_videomme_v2`;rollout/诊断的 Redis 盐派生自 `infer_adhoc` 不变 → 已烧调用命中复用) +3. **验证里程碑**:启动检查点(预检剔除 4 题型 / gate 阶梯 8 题型 / Best=0.7167)→ 首 batch rollout(应大量缓存命中,显著快于 27min)→ 首个 step 的连续 gate(观察多题型交错的 `_gate_` run 记录与判定日志)。 + +## 保真校验记录 + +- **#4 CE-Gate e-process**:`core/evolution/gate.py` 零改动;`gate_decision` 调用参数语义不变(n_remaining 等价迁移)。✅ +- **#5 信息阶梯**:`ladder_for` / `_ladder_units` 排序与截断逻辑零改动;案例单元排除保留。✅ +- **#6 块顺序验证**:**已批准语义修订**(设计 v3):块序贯 → 阶梯序前缀逐对序贯;基线缓存、INFRA 护栏、配对翻转等价迁移;本计划 Task 1-3/5-6 涉及,commit 均标注 `(algo #6)`。⚠️ 已声明 +- 其余 #1-#3、#7-#12:不涉及。 + +## 风险与回滚 + +- 每 Task 独立 commit,任一集成点失败可 `git revert` 到上一 Task。 +- Phase C 若在真实训练暴露调度死锁(理论上 `_QuestionSlots` 的锁序已防),现象为 gate 阶段静默无推理调用 → 立即 `tmux kill` + 回滚 Task 5/6 的 commit,退回块序贯版本重启(旧路径在 Task 6 前仍在)。 diff --git a/research-wiki/plans/gate-speedup.md b/research-wiki/plans/gate-speedup.md new file mode 100644 index 0000000..b2c5352 --- /dev/null +++ b/research-wiki/plans/gate-speedup.md @@ -0,0 +1,9 @@ +--- +type: plan +node_id: plan:gate-speedup +title: "连续并发 gate + Redis 复用实现计划" +date: 2026-07-17 +--- + +# 连续并发 gate + Redis 复用实现计划 +