diff --git a/lib/air_runtime/modes/arc_mode.py b/lib/air_runtime/modes/arc_mode.py index 1fe716e..0651334 100644 --- a/lib/air_runtime/modes/arc_mode.py +++ b/lib/air_runtime/modes/arc_mode.py @@ -78,10 +78,12 @@ def _ensure_dirs(paths: dict[str, Path]) -> None: def _export_task_graph_json(graph: TaskGraph, path: Path) -> None: data = { + "dispatchFrozen": graph.dispatch_frozen, # P1-21 "nodes": {nid: {"id": n.id, "status": n.status, "task": n.task, "filesDirs": n.files_dirs, "doneWhen": n.done_when, "inDegree": n.in_degree, "outEdges": n.out_edges, - "writeSet": n.write_set, "testRequired": n.test_required} + "writeSet": n.write_set, "testRequired": n.test_required, + "adrRefs": n.adr_refs} # P1-21 for nid, n in graph.nodes.items()}, "edges": [{"source": e.source, "target": e.target, "kind": e.kind} for e in graph.edges], } @@ -96,9 +98,14 @@ def _build_graph_from_todo(todo_path: Path) -> tuple[TaskGraph, list[str]]: violations = [] # P1-19.1: 记录 Done When 不含"测试通过"的任务 for t in tasks: + # P1-21: 从 todo.md ADR 列提取 adr_refs + adr_refs = [] + if hasattr(t, "adr") and t.adr: + adr_refs = [a.strip() for a in t.adr.split(",") if a.strip()] node = TaskNode( id=t.task_id, status=t.status, task=t.task, files_dirs=t.files_dirs, done_when=t.done_when, + adr_refs=adr_refs, ) graph.add_node(node) diff --git a/lib/air_runtime/modes/eng_mode.py b/lib/air_runtime/modes/eng_mode.py index 839182e..f40d495 100644 --- a/lib/air_runtime/modes/eng_mode.py +++ b/lib/air_runtime/modes/eng_mode.py @@ -26,7 +26,7 @@ from air_runtime.modes.merge_pipeline import ( sync_engine_managed_docs, update_todo_after_merge, ) -from air_runtime.task_graph import TaskGraph +from air_runtime.task_graph import TaskGraph, CascadeReport, PlanDelta from air_runtime.todo_parser import parse_tasks from air_runtime.utils import now_iso, session_stamp, truncate_history @@ -192,6 +192,23 @@ def dispatch_worker_group(project_root: Path, group_name: str = "") -> dict: atomic_json_write(paths["state"], state) state = safe_json_load(paths["state"]) or _init_state(project_root) + # P1-21: 检查调度冻结(ADR 级联失效期间) + tg_json = airplan_root(project_root) / "state" / "airarc" / "reviews" / "task-graph.json" + if tg_json.exists(): + try: + graph = TaskGraph.load(tg_json) + if graph.dispatch_frozen: + log = EventLog(event_log_path(project_root)) + log.emit("eng.blocked", {"reason": "dispatch frozen — ADR cascade invalidation in progress"}) + return { + "blocked": True, + "reason": "dispatch frozen — ADR cascade invalidation in progress", + "waveId": "", + "taskIds": [], + } + except Exception: + pass + # P1-19.3: 检查是否有 block-release verdict,阻止所有后续派发 from air_runtime.review_runtime import ReviewRuntime rvr = ReviewRuntime(project_root) @@ -546,6 +563,166 @@ def _select_ready_tasks(project_root: Path, max_count: int) -> list[str]: return [t.task_id for t in tasks if t.status == "TODO"][:max_count] +def handle_adr_invalidation(project_root: Path, adr_id: str) -> dict: + """P1-21: ADR 变更级联失效处理。 + + 10步流程: + 1. 加载 task-graph.json + 2. 调用 invalidate_by_adr() 级联失效 + 3. 冻结调度 + 4. 中止进行中的相关 Worker + 5. 创建回滚快照(git tag) + 6. git revert 已合并的旧代码 + 7. 写回更新后的 task-graph.json + 8. 等待 Arc 重新生成受影响部分的任务 + 9. apply_delta() 吸收新任务 + 10. 解冻调度 + """ + tg_json = airplan_root(project_root) / "state" / "airarc" / "reviews" / "task-graph.json" + if not tg_json.exists(): + return {"error": "task-graph.json not found", "adrId": adr_id} + + graph = TaskGraph.load(tg_json) + delta = PlanDelta() + + # 2-4: 级联失效 + report = graph.invalidate_by_adr(adr_id, delta) + log = EventLog(event_log_path(project_root)) + log.emit("adr.invalidation", { + "adrId": adr_id, + "invalidatedCompleted": report.invalidated_completed, + "terminatedInProgress": report.terminated_in_progress, + "cascadedDownstream": report.cascaded_downstream, + }) + + # 4: 中止进行中的相关 Worker + paths = _paths(project_root) + _ensure_dirs(paths) + state = safe_json_load(paths["state"]) or _init_state(project_root) + terminated_workers = [] + for worker in list(state.get("activeWorkers", [])): + if worker.get("taskId") in report.invalidated_task_ids: + terminated_workers.append(worker["taskId"]) + state["activeWorkers"] = [ + w for w in state.get("activeWorkers", []) + if w.get("taskId") not in report.invalidated_task_ids + ] + + # 5: 创建回滚快照 + rollback_ref = _create_rollback_snapshot(project_root, report.invalidated_task_ids) + report.rollback_ref = rollback_ref + delta.rollback_ref = rollback_ref + + # 6: git revert 已合并的旧代码(按 task_id 查找对应 commit) + revert_results = _git_revert_invalidated(project_root, report.invalidated_task_ids) + + # 7: 写回更新后的 task-graph.json + from air_runtime.modes.arc_mode import _export_task_graph_json + _export_task_graph_json(graph, tg_json) + + # 更新引擎状态 + state["dispatchFrozen"] = True + state["adrInvalidationInProgress"] = { + "adrId": adr_id, + "startedAt": now_iso(), + "invalidatedTaskIds": report.invalidated_task_ids, + "rollbackRef": rollback_ref, + } + atomic_json_write(paths["state"], state) + + return { + "adrId": adr_id, + "cascadeReport": { + "invalidatedCompleted": report.invalidated_completed, + "terminatedInProgress": report.terminated_in_progress, + "cascadedDownstream": report.cascaded_downstream, + "rollbackRef": rollback_ref, + "invalidatedTaskIds": report.invalidated_task_ids, + }, + "terminatedWorkers": terminated_workers, + "revertResults": revert_results, + "nextStep": "arc-replan-then-unfreeze", + } + + +def _create_rollback_snapshot(project_root: Path, invalidated_task_ids: list[str]) -> str: + """P1-21: 为失效任务创建 git tag 回滚点。""" + import subprocess + ref = f"airplan/adr-invalidate-{session_stamp()}" + try: + subprocess.run( + ["git", "tag", ref], + cwd=project_root, capture_output=True, timeout=30, + ) + except Exception: + pass + return ref + + +def _git_revert_invalidated(project_root: Path, invalidated_task_ids: list[str]) -> list[dict]: + """P1-21: 尝试 git revert 已合并的失效任务对应的 commit。""" + import subprocess + results = [] + for tid in invalidated_task_ids: + try: + # 查找包含 task_id 的 commit + r = subprocess.run( + ["git", "log", "--oneline", "--all", "--grep", tid, "-1"], + cwd=project_root, capture_output=True, text=True, timeout=10, + ) + if r.returncode == 0 and r.stdout.strip(): + commit_hash = r.stdout.strip().split()[0] + rv = subprocess.run( + ["git", "revert", "--no-commit", commit_hash], + cwd=project_root, capture_output=True, text=True, timeout=30, + ) + results.append({"taskId": tid, "commit": commit_hash, "reverted": rv.returncode == 0}) + if rv.returncode == 0: + subprocess.run( + ["git", "commit", "-m", f"AirPlan: revert invalidated task {tid}"], + cwd=project_root, capture_output=True, timeout=10, + ) + else: + results.append({"taskId": tid, "commit": None, "reverted": False, "reason": "no commit found"}) + except Exception as e: + results.append({"taskId": tid, "commit": None, "reverted": False, "reason": str(e)}) + return results + + +def unfreeze_after_replan(project_root: Path, new_task_graph_path: Path | None = None) -> dict: + """P1-21: Arc 重新生成受影响部分后,apply_delta + 解冻调度。""" + paths = _paths(project_root) + _ensure_dirs(paths) + + tg_json = airplan_root(project_root) / "state" / "airarc" / "reviews" / "task-graph.json" + if not tg_json.exists(): + return {"error": "task-graph.json not found"} + + graph = TaskGraph.load(tg_json) + + # 如果 Arc 生成了新的任务图,增量合并 + if new_task_graph_path and new_task_graph_path.exists(): + new_graph = TaskGraph.load(new_task_graph_path) + delta = new_graph.diff(graph) + graph.apply_delta(delta) + + # 解冻 + graph.unfreeze_dispatch() + from air_runtime.modes.arc_mode import _export_task_graph_json + _export_task_graph_json(graph, tg_json) + + # 更新引擎状态 + state = safe_json_load(paths["state"]) or _init_state(project_root) + state["dispatchFrozen"] = False + adr_info = state.pop("adrInvalidationInProgress", {}) + atomic_json_write(paths["state"], state) + + log = EventLog(event_log_path(project_root)) + log.emit("adr.unfreezed", {"previousAdrInvalidation": adr_info}) + + return {"frozen": False, "readyTasks": graph.ready_tasks()} + + def maybe_replan(project_root: Path, todo_path: Path | None = None) -> dict | None: """检查 todo.md mtime vs task_graph.json mtime,若 todo 更新则触发 replan。""" from air_runtime.modes.arc_mode import incremental_replan_mode diff --git a/lib/air_runtime/review_runtime.py b/lib/air_runtime/review_runtime.py index 4f8a3c0..1195234 100644 --- a/lib/air_runtime/review_runtime.py +++ b/lib/air_runtime/review_runtime.py @@ -136,6 +136,25 @@ class ReviewRuntime: "deliveryVerdict": report.get("highRiskAudit", {}).get("deliveryVerdict", "safe-to-ship"), # P1-19.2 } + def check_invalidated_cleanup(self, invalidated_task_ids: list[str]) -> dict: + """P1-21: 检查 INVALIDATED 任务的代码是否已清理(无残留)。""" + from air_runtime.io import safe_json_load + residual = [] + for tid in invalidated_task_ids: + # 检查是否有残留的 result 文件(说明旧代码未被 revert) + result_dir = self._state_dir.parent / "airdo" / "tasks" / tid + if result_dir.exists(): + result_file = result_dir / "result.json" + if result_file.exists(): + data = safe_json_load(result_file) + if data and data.get("status") == "done": + residual.append({"taskId": tid, "reason": "done result still exists — code may not be reverted"}) + return { + "cleaned": len(residual) == 0, + "residualCount": len(residual), + "residualDetails": residual, + } + @staticmethod def _report_to_dict(report: ReviewReport) -> dict: # P1-19.2: highRiskAudit 序列化 diff --git a/lib/air_runtime/task_graph.py b/lib/air_runtime/task_graph.py index 0172488..e955c9f 100644 --- a/lib/air_runtime/task_graph.py +++ b/lib/air_runtime/task_graph.py @@ -12,7 +12,7 @@ from typing import Any @dataclass class TaskNode: id: str - status: str = "TODO" # TODO | DISPATCHED | DONE | BLOCKED + status: str = "TODO" # TODO | DISPATCHED | DONE | BLOCKED | INVALIDATED task: str = "" files_dirs: str = "" done_when: str = "" @@ -21,6 +21,7 @@ class TaskNode: write_set: list[str] = field(default_factory=list) meta: dict[str, Any] = field(default_factory=dict) test_required: bool = False # P1-19.1: 边界测试强制标记 + adr_refs: list[str] = field(default_factory=list) # P1-21: ADR→任务溯源链 @dataclass @@ -36,6 +37,16 @@ class EdgeChange: removed: list[Edge] = field(default_factory=list) +@dataclass +class CascadeReport: + """P1-21: ADR 变更级联失效报告。""" + invalidated_completed: int = 0 + terminated_in_progress: int = 0 + cascaded_downstream: int = 0 + rollback_ref: str = "" + invalidated_task_ids: list[str] = field(default_factory=list) + + @dataclass class PlanDelta: """Arc 重规划产出的增量差异,替代全量覆盖 todo.md。""" @@ -43,6 +54,7 @@ class PlanDelta: added_tasks: list[TaskNode] = field(default_factory=list) modified_tasks: list[TaskNode] = field(default_factory=list) edge_changes: EdgeChange = field(default_factory=EdgeChange) + rollback_ref: str = "" # P1-21: 回滚快照引用 class TaskGraph: @@ -51,6 +63,7 @@ class TaskGraph: def __init__(self): self.nodes: dict[str, TaskNode] = {} self.edges: list[Edge] = [] + self.dispatch_frozen: bool = False # P1-21: 调度冻结 def add_node(self, node: TaskNode) -> None: self.nodes[node.id] = node @@ -90,7 +103,9 @@ class TaskGraph: self.nodes[edge.source].out_edges.append(edge.target) def ready_tasks(self) -> list[str]: - """返回当前入度为 0 且状态为 TODO 的任务。""" + """返回当前入度为 0 且状态为 TODO 的任务。调度冻结时返回空。""" + if self.dispatch_frozen: + return [] return [nid for nid, n in self.nodes.items() if n.in_degree == 0 and n.status == "TODO"] def diff(self, other: TaskGraph) -> PlanDelta: @@ -138,6 +153,7 @@ class TaskGraph: graph = cls() if not data or not isinstance(data, dict): return graph + graph.dispatch_frozen = data.get("dispatchFrozen", False) for nid, nd in data.get("nodes", {}).items(): graph.nodes[nid] = TaskNode( id=nd.get("id", nid), @@ -149,6 +165,7 @@ class TaskGraph: out_edges=list(nd.get("outEdges", [])), write_set=list(nd.get("writeSet", [])), test_required=nd.get("testRequired", False), + adr_refs=list(nd.get("adrRefs", [])), ) for ed in data.get("edges", []): graph.edges.append(Edge( @@ -205,7 +222,8 @@ class TaskGraph: if node.id in self.nodes: existing_status = self.nodes[node.id].status self.nodes[node.id] = node - if existing_status in ("DISPATCHED", "DONE"): + # INVALIDATED 可覆盖 DONE/DISPATCHED(P1-21: ADR 级联失效) + if existing_status in ("DISPATCHED", "DONE") and node.status != "INVALIDATED": self.nodes[node.id].status = existing_status def _remove_edge(self, edge: Edge) -> None: @@ -220,3 +238,70 @@ class TaskGraph: self.nodes[edge.target].in_degree += 1 if edge.source in self.nodes: self.nodes[edge.source].out_edges.append(edge.target) + + # P1-21: ADR 级联失效 + + def tasks_by_adr(self, adr_id: str) -> list[TaskNode]: + """查找所有引用指定 ADR 的任务(含已完成)。""" + return [n for n in self.nodes.values() if adr_id in n.adr_refs] + + def _find_downstream(self, task_ids: list[str]) -> list[str]: + """BFS 遍历下游依赖任务。""" + visited: set[str] = set() + queue = list(task_ids) + while queue: + current = queue.pop(0) + if current in visited: + continue + visited.add(current) + node = self.nodes.get(current) + if node: + for target in node.out_edges: + if target not in visited: + queue.append(target) + # 排除起点自身 + return [tid for tid in visited if tid not in set(task_ids)] + + def invalidate_by_adr(self, adr_id: str, delta: PlanDelta) -> CascadeReport: + """P1-21: ADR 变更时级联失效所有相关任务。""" + affected = self.tasks_by_adr(adr_id) + completed = [t for t in affected if t.status == "DONE"] + in_progress = [t for t in affected if t.status == "DISPATCHED"] + + # 1. 冻结调度 + self.dispatch_frozen = True + + # 2. 标记已完成任务为 INVALIDATED + for t in completed: + t.status = "INVALIDATED" + delta.removed_tasks.append(t.id) + + # 3. 标记进行中任务为 INVALIDATED(调用方负责中止 Worker) + for t in in_progress: + t.status = "INVALIDATED" + delta.removed_tasks.append(t.id) + + # 4. 级联失效下游 + downstream_ids = self._find_downstream([t.id for t in completed + in_progress]) + cascaded = [] + for tid in downstream_ids: + node = self.nodes.get(tid) + if node and node.status in ("TODO", "DISPATCHED"): + node.status = "INVALIDATED" + delta.removed_tasks.append(tid) + cascaded.append(tid) + + # 5. 回滚快照引用(由调用方在 git revert 后填入) + all_invalidated = [t.id for t in completed + in_progress] + cascaded + + return CascadeReport( + invalidated_completed=len(completed), + terminated_in_progress=len(in_progress), + cascaded_downstream=len(cascaded), + rollback_ref=delta.rollback_ref, + invalidated_task_ids=all_invalidated, + ) + + def unfreeze_dispatch(self) -> None: + """P1-21: 解冻调度,在 Arc 重新生成受影响任务后调用。""" + self.dispatch_frozen = False diff --git a/skills/airplan/SKILL.md b/skills/airplan/SKILL.md index c7c7777..86e5ae1 100644 --- a/skills/airplan/SKILL.md +++ b/skills/airplan/SKILL.md @@ -90,6 +90,14 @@ AirDo 处理 UI 任务时必须使用 frontend-design Skill: - frontend-design Skill 不存在则自动安装 - 安装失败时阻止任务执行 +### INV-15 ADR 变更级联失效 +架构方案变更(如 ffmpeg → gstreamer)时,基于旧 ADR 已完成的任务必须自动失效: +- TaskGraph 维护 ADR→任务溯源链(adr_refs) +- ADR 变更时级联失效所有相关任务(含已完成和下游依赖) +- 调度冻结直到 Arc 重新生成受影响部分的任务 +- 自动 git revert 已合并的旧代码 +- AirRvr 终审检查 INVALIDATED 任务代码是否已清理 + ## L1 代码级保障(不依赖 LLM 自觉) 以下功能由引擎代码强制执行,SKILL.md 指令仅作辅助: @@ -108,6 +116,7 @@ AirDo 处理 UI 任务时必须使用 frontend-design Skill: 12. **边界测试强制** — `air_runtime/modes/arc_mode.py:_inject_boundary_tests()` 为每个模块注入测试任务 13. **block-release 阻止派发** — `air_runtime/modes/eng_mode.py:dispatch_worker_group()` 检查 deliveryVerdict 14. **UI Skill 路由** — `air_runtime/modes/do_mode.py:route_ui_task()` 检测并确保 frontend-design Skill +15. **ADR 级联失效** — `air_runtime/task_graph.py:invalidate_by_adr()` ADR 变更时级联失效 + `air_runtime/modes/eng_mode.py:handle_adr_invalidation()` 调度冻结 + git revert ## 使用示例 diff --git a/test_p1_21.py b/test_p1_21.py new file mode 100644 index 0000000..f66a629 --- /dev/null +++ b/test_p1_21.py @@ -0,0 +1,407 @@ +#!/usr/bin/env python3 +"""P1-21 ADR 变更级联失效 功能测试""" + +import sys +import tempfile +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).parent / "lib")) + + +def test_task_node_adr_refs(): + """测试 TaskNode.adr_refs 字段""" + from air_runtime.task_graph import TaskNode + + node = TaskNode(id="T-001", adr_refs=["ADR-0005", "ADR-0012"]) + assert node.adr_refs == ["ADR-0005", "ADR-0012"] + + # 默认为空列表 + node2 = TaskNode(id="T-002") + assert node2.adr_refs == [] + + print("✓ TaskNode.adr_refs 测试通过") + return True + + +def test_invalidated_status(): + """测试 INVALIDATED 状态""" + from air_runtime.task_graph import TaskNode, TaskGraph + + node = TaskNode(id="T-001", status="DONE", task="用 ffmpeg 实现解码器", + adr_refs=["ADR-0005"]) + + # 直接设置 INVALIDATED + node.status = "INVALIDATED" + assert node.status == "INVALIDATED" + + # _update_node 允许 INVALIDATED 覆盖 DONE + graph = TaskGraph() + graph.add_node(TaskNode(id="T-001", status="DONE", task="旧描述", adr_refs=["ADR-0005"])) + graph._update_node(TaskNode(id="T-001", status="INVALIDATED", task="新描述", adr_refs=["ADR-0005"])) + assert graph.nodes["T-001"].status == "INVALIDATED" + + # _update_node 仍然保护 DONE 不被普通修改覆盖 + graph2 = TaskGraph() + graph2.add_node(TaskNode(id="T-002", status="DONE", task="旧描述")) + graph2._update_node(TaskNode(id="T-002", status="TODO", task="新描述")) + assert graph2.nodes["T-002"].status == "DONE" # 被保护,不允许降级为 TODO + + print("✓ INVALIDATED 状态测试通过") + return True + + +def test_tasks_by_adr(): + """测试 ADR→任务溯源链查询""" + from air_runtime.task_graph import TaskGraph, TaskNode + + graph = TaskGraph() + graph.add_node(TaskNode(id="T-001", task="ffmpeg 解码器", adr_refs=["ADR-0005"])) + graph.add_node(TaskNode(id="T-002", task="ffmpeg 编码器", adr_refs=["ADR-0005", "ADR-0010"])) + graph.add_node(TaskNode(id="T-003", task="UI 界面", adr_refs=["ADR-0008"])) + graph.add_node(TaskNode(id="T-004", task="日志模块", adr_refs=[])) + + # ADR-0005 关联 T-001 和 T-002 + result = graph.tasks_by_adr("ADR-0005") + assert len(result) == 2 + assert {n.id for n in result} == {"T-001", "T-002"} + + # ADR-0010 只关联 T-002 + result2 = graph.tasks_by_adr("ADR-0010") + assert len(result2) == 1 + assert result2[0].id == "T-002" + + # 不存在的 ADR + result3 = graph.tasks_by_adr("ADR-9999") + assert len(result3) == 0 + + print("✓ ADR→任务溯源链查询测试通过") + return True + + +def test_invalidate_by_adr(): + """测试 ADR 变更级联失效""" + from air_runtime.task_graph import TaskGraph, TaskNode, Edge, PlanDelta + + graph = TaskGraph() + # 已完成的 ffmpeg 任务 + graph.add_node(TaskNode(id="T-001", status="DONE", task="ffmpeg 解码器", adr_refs=["ADR-0005"])) + graph.add_node(TaskNode(id="T-002", status="DONE", task="ffmpeg 编码器", adr_refs=["ADR-0005"])) + # 进行中的 ffmpeg 任务 + graph.add_node(TaskNode(id="T-003", status="DISPATCHED", task="ffmpeg 流媒体", adr_refs=["ADR-0005"])) + # 下游任务(依赖上面的 ffmpeg 模块) + graph.add_node(TaskNode(id="T-004", status="TODO", task="集成测试")) + graph.add_node(TaskNode(id="T-005", status="TODO", task="部署")) + # 无关任务 + graph.add_node(TaskNode(id="T-006", status="TODO", task="UI 优化", adr_refs=["ADR-0008"])) + + graph.add_edge(Edge(source="T-001", target="T-004")) + graph.add_edge(Edge(source="T-002", target="T-004")) + graph.add_edge(Edge(source="T-004", target="T-005")) + + delta = PlanDelta() + report = graph.invalidate_by_adr("ADR-0005", delta) + + # 已完成: T-001, T-002 + assert report.invalidated_completed == 2 + # 进行中: T-003 + assert report.terminated_in_progress == 1 + # 下游: T-004, T-005 + assert report.cascaded_downstream == 2 + # 总失效 + assert len(report.invalidated_task_ids) == 5 + + # 验证状态已变为 INVALIDATED + assert graph.nodes["T-001"].status == "INVALIDATED" + assert graph.nodes["T-002"].status == "INVALIDATED" + assert graph.nodes["T-003"].status == "INVALIDATED" + assert graph.nodes["T-004"].status == "INVALIDATED" + assert graph.nodes["T-005"].status == "INVALIDATED" + + # 无关任务不受影响 + assert graph.nodes["T-006"].status == "TODO" + + # 调度已冻结 + assert graph.dispatch_frozen == True + + # ready_tasks 为空(冻结状态) + assert graph.ready_tasks() == [] + + print(f"✓ ADR 变更级联失效测试通过 (完成={report.invalidated_completed}, " + f"进行中={report.terminated_in_progress}, 下游={report.cascaded_downstream})") + return True + + +def test_dispatch_freeze_unfreeze(): + """测试调度冻结和解冻""" + from air_runtime.task_graph import TaskGraph, TaskNode + + graph = TaskGraph() + graph.add_node(TaskNode(id="T-001", status="TODO", task="任务1")) + + # 正常状态 + assert graph.ready_tasks() == ["T-001"] + + # 冻结 + graph.dispatch_frozen = True + assert graph.ready_tasks() == [] + + # 解冻 + graph.unfreeze_dispatch() + assert graph.dispatch_frozen == False + assert graph.ready_tasks() == ["T-001"] + + print("✓ 调度冻结/解冻测试通过") + return True + + +def test_cascade_report(): + """测试 CascadeReport 数据结构""" + from air_runtime.task_graph import CascadeReport + + report = CascadeReport( + invalidated_completed=2, + terminated_in_progress=1, + cascaded_downstream=3, + rollback_ref="airplan/adr-invalidate-2026-06-11", + invalidated_task_ids=["T-001", "T-002", "T-003", "T-004", "T-005", "T-006"], + ) + + assert report.invalidated_completed == 2 + assert report.terminated_in_progress == 1 + assert report.cascaded_downstream == 3 + assert report.rollback_ref.startswith("airplan/") + assert len(report.invalidated_task_ids) == 6 + + print("✓ CascadeReport 测试通过") + return True + + +def test_plan_delta_rollback_ref(): + """测试 PlanDelta.rollback_ref""" + from air_runtime.task_graph import PlanDelta + + delta = PlanDelta(rollback_ref="airplan/adr-invalidate-2026-06-11") + assert delta.rollback_ref == "airplan/adr-invalidate-2026-06-11" + + # 默认为空字符串 + delta2 = PlanDelta() + assert delta2.rollback_ref == "" + + print("✓ PlanDelta.rollback_ref 测试通过") + return True + + +def test_rvr_invalidated_cleanup_check(): + """测试 AirRvr INVALIDATED 清理检查""" + from air_runtime.review_runtime import ReviewRuntime + + with tempfile.TemporaryDirectory() as tmpdir: + project_root = Path(tmpdir) + rvr = ReviewRuntime(project_root) + + # 无残留 + result = rvr.check_invalidated_cleanup(["T-001", "T-002"]) + assert result["cleaned"] == True + assert result["residualCount"] == 0 + + # 创建残留的 result 文件 + result_dir = project_root / "AirPlan" / "state" / "airdo" / "tasks" / "T-001" + result_dir.mkdir(parents=True, exist_ok=True) + (result_dir / "result.json").write_text('{"taskId": "T-001", "status": "done"}') + + result2 = rvr.check_invalidated_cleanup(["T-001"]) + assert result2["cleaned"] == False + assert result2["residualCount"] == 1 + assert result2["residualDetails"][0]["taskId"] == "T-001" + + print("✓ AirRvr INVALIDATED 清理检查测试通过") + return True + + +def test_eng_handle_adr_invalidation(): + """测试 Eng ADR 失效处理入口""" + from air_runtime.modes.eng_mode import handle_adr_invalidation + import json + + with tempfile.TemporaryDirectory() as tmpdir: + project_root = Path(tmpdir) + + # 初始化引擎 + from air_runtime.modes.eng_mode import enter_engine + enter_engine(project_root) + + # 创建 task-graph.json(含 ADR 引用) + arc_dir = project_root / "AirPlan" / "state" / "airarc" / "reviews" + arc_dir.mkdir(parents=True, exist_ok=True) + + from air_runtime.task_graph import TaskGraph, TaskNode, Edge + from air_runtime.modes.arc_mode import _export_task_graph_json + + graph = TaskGraph() + graph.add_node(TaskNode(id="T-001", status="DONE", task="ffmpeg 解码器", adr_refs=["ADR-0005"])) + graph.add_node(TaskNode(id="T-002", status="TODO", task="集成测试")) + graph.add_edge(Edge(source="T-001", target="T-002")) + _export_task_graph_json(graph, arc_dir / "task-graph.json") + + # 执行 ADR 失效 + result = handle_adr_invalidation(project_root, "ADR-0005") + + assert "cascadeReport" in result + assert result["cascadeReport"]["invalidatedCompleted"] == 1 + assert result["cascadeReport"]["cascadedDownstream"] == 1 + assert result["nextStep"] == "arc-replan-then-unfreeze" + + # 验证 task-graph.json 已更新 + updated = TaskGraph.load(arc_dir / "task-graph.json") + assert updated.nodes["T-001"].status == "INVALIDATED" + assert updated.dispatch_frozen == True + + print("✓ Eng ADR 失效处理测试通过") + return True + + +def test_eng_unfreeze_after_replan(): + """测试 Eng 解冻调度""" + from air_runtime.modes.eng_mode import handle_adr_invalidation, unfreeze_after_replan + + with tempfile.TemporaryDirectory() as tmpdir: + project_root = Path(tmpdir) + + from air_runtime.modes.eng_mode import enter_engine + enter_engine(project_root) + + arc_dir = project_root / "AirPlan" / "state" / "airarc" / "reviews" + arc_dir.mkdir(parents=True, exist_ok=True) + + from air_runtime.task_graph import TaskGraph, TaskNode + from air_runtime.modes.arc_mode import _export_task_graph_json + + graph = TaskGraph() + graph.add_node(TaskNode(id="T-001", status="DONE", task="ffmpeg 解码器", adr_refs=["ADR-0005"])) + graph.dispatch_frozen = True + _export_task_graph_json(graph, arc_dir / "task-graph.json") + + # 解冻 + result = unfreeze_after_replan(project_root) + + assert result["frozen"] == False + + # 验证 task-graph.json 已解冻 + updated = TaskGraph.load(arc_dir / "task-graph.json") + assert updated.dispatch_frozen == False + + print("✓ Eng 解冻调度测试通过") + return True + + +def test_ffmpeg_to_gstreamer_scenario(): + """完整场景测试:ffmpeg → gstreamer 架构变更""" + from air_runtime.task_graph import TaskGraph, TaskNode, Edge, PlanDelta + + # 初始状态:ffmpeg 架构 + graph = TaskGraph() + graph.add_node(TaskNode(id="T-001", status="DONE", task="用 ffmpeg 实现视频解码器", + files_dirs="src/decoder.cpp", adr_refs=["ADR-0005"])) + graph.add_node(TaskNode(id="T-002", status="DONE", task="用 ffmpeg 实现音频解码器", + files_dirs="src/audio_decoder.cpp", adr_refs=["ADR-0005"])) + graph.add_node(TaskNode(id="T-003", status="DONE", task="ffmpeg 流媒体传输", + files_dirs="src/streaming.cpp", adr_refs=["ADR-0005", "ADR-0006"])) + graph.add_node(TaskNode(id="T-004", status="TODO", task="解码器集成测试", + files_dirs="tests/decoder_test.cpp")) + graph.add_node(TaskNode(id="T-005", status="TODO", task="流媒体集成测试", + files_dirs="tests/streaming_test.cpp")) + graph.add_node(TaskNode(id="T-006", status="TODO", task="部署播放器", + files_dirs="deploy/")) + # 无关任务 + graph.add_node(TaskNode(id="T-007", status="DONE", task="用户界面", + files_dirs="src/ui/", adr_refs=["ADR-0008"])) + graph.add_node(TaskNode(id="T-008", status="TODO", task="UI 测试", + files_dirs="tests/ui_test.cpp")) + + graph.add_edge(Edge(source="T-001", target="T-004")) + graph.add_edge(Edge(source="T-002", target="T-004")) + graph.add_edge(Edge(source="T-003", target="T-005")) + graph.add_edge(Edge(source="T-004", target="T-006")) + graph.add_edge(Edge(source="T-005", target="T-006")) + graph.add_edge(Edge(source="T-007", target="T-008")) + + # ADR-0005 变更:ffmpeg → gstreamer + delta = PlanDelta(rollback_ref="airplan/adr-invalidate-snapshot-001") + report = graph.invalidate_by_adr("ADR-0005", delta) + + # T-001, T-002, T-003 完成 → INVALIDATED + assert report.invalidated_completed == 3 + # T-004, T-005, T-006 下游 → INVALIDATED + assert report.cascaded_downstream == 3 + total_invalidated = len(report.invalidated_task_ids) + assert total_invalidated == 6 # T-001~T-006 全部失效 + + # 无关任务不受影响 + assert graph.nodes["T-007"].status == "DONE" + assert graph.nodes["T-008"].status == "TODO" + + # 调度冻结 + assert graph.dispatch_frozen == True + assert graph.ready_tasks() == [] + + # 模拟 Arc 重新生成 gstreamer 任务后解冻 + graph.add_node(TaskNode(id="T-101", status="TODO", task="用 gstreamer 实现视频解码器", + files_dirs="src/decoder.cpp", adr_refs=["ADR-0005"])) + graph.add_node(TaskNode(id="T-102", status="TODO", task="用 gstreamer 实现音频解码器", + files_dirs="src/audio_decoder.cpp", adr_refs=["ADR-0005"])) + graph.unfreeze_dispatch() + + assert graph.dispatch_frozen == False + # 新任务可以派发 + ready = [tid for tid, n in graph.nodes.items() if n.status == "TODO" and n.in_degree == 0] + assert "T-101" in ready + assert "T-102" in ready + # T-008 有入度(依赖 T-007),不在 ready 中 + assert "T-008" not in ready + + print(f"✓ ffmpeg→gstreamer 完整场景测试通过 (失效={total_invalidated}, 新任务可派发)") + return True + + +def main(): + print("=" * 50) + print("P1-21 ADR 变更级联失效 功能测试") + print("=" * 50) + + tests = [ + ("TaskNode.adr_refs", test_task_node_adr_refs), + ("INVALIDATED 状态", test_invalidated_status), + ("ADR→任务溯源链", test_tasks_by_adr), + ("级联失效", test_invalidate_by_adr), + ("调度冻结/解冻", test_dispatch_freeze_unfreeze), + ("CascadeReport", test_cascade_report), + ("PlanDelta.rollback_ref", test_plan_delta_rollback_ref), + ("AirRvr 清理检查", test_rvr_invalidated_cleanup_check), + ("Eng ADR 失效处理", test_eng_handle_adr_invalidation), + ("Eng 解冻调度", test_eng_unfreeze_after_replan), + ("ffmpeg→gstreamer 场景", test_ffmpeg_to_gstreamer_scenario), + ] + + passed = 0 + failed = 0 + + for name, test_fn in tests: + try: + test_fn() + passed += 1 + except Exception as e: + print(f"✗ {name} 失败: {e}") + import traceback + traceback.print_exc() + failed += 1 + + print("=" * 50) + print(f"测试结果: {passed} 通过, {failed} 失败") + print("=" * 50) + + return failed == 0 + + +if __name__ == "__main__": + success = main() + sys.exit(0 if success else 1)