From 640e17499db3f80cea38fe88efd19dd0f56bcebc Mon Sep 17 00:00:00 2001 From: zizi Date: Sun, 23 Aug 2026 15:17:57 +0800 Subject: [PATCH] =?UTF-8?q?=E9=98=B6=E6=AE=B5F=E7=AC=AC=E4=BA=8C=E9=83=A8?= =?UTF-8?q?=E5=88=86:=20=E7=9F=A5=E8=AF=86=E9=93=BE=E6=B4=BE=E5=8F=91?= =?UTF-8?q?=E2=80=94=E2=80=94=E6=8A=BD=E5=8F=96=E6=99=BA=E8=83=BD=E4=BD=93?= =?UTF-8?q?=E6=8E=A5=E5=85=A5=E6=A1=86=E6=9E=B6=E6=B4=BE=E5=8F=91=EF=BC=88?= =?UTF-8?q?=E7=A6=BB=E7=BA=BF7=E9=A1=B9=E7=BB=BF+=E7=AC=AC3=E7=AB=A0?= =?UTF-8?q?=E7=9C=9F=E5=AE=9E=E6=8A=BD=E5=8F=96=E4=BA=A7=E5=87=BA19?= =?UTF-8?q?=E6=9D=A1=E8=8D=89=E7=A8=BF=EF=BC=8C=E4=BF=AE=E5=A4=8D=E9=87=8D?= =?UTF-8?q?=E6=B4=BE=E8=B7=AF=E5=BE=84=E9=A6=96=E6=AC=A1=E7=9C=9F=E5=AE=9E?= =?UTF-8?q?=E9=AA=8C=E8=AF=81=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scripts/dispatch_extraction_bridge.py | 277 ++++++++++++++++++ .../scripts/extract_via_dispatch.py | 54 ++++ .../2026-08-22-agent-example整体收敛总plan.md | 1 + ...22-阶段F-第二部分-知识链派发-抽取智能体.md | 48 +++ .../test_dispatch_extraction_bridge.py | 153 ++++++++++ 5 files changed, 533 insertions(+) create mode 100644 .agent/skills/extract-chapter-knowledge/scripts/dispatch_extraction_bridge.py create mode 100644 .agent/skills/extract-chapter-knowledge/scripts/extract_via_dispatch.py create mode 100644 docs/plans/2026-08-22-阶段F-第二部分-知识链派发-抽取智能体.md create mode 100644 tests/skills/extract-chapter-knowledge/test_dispatch_extraction_bridge.py diff --git a/.agent/skills/extract-chapter-knowledge/scripts/dispatch_extraction_bridge.py b/.agent/skills/extract-chapter-knowledge/scripts/dispatch_extraction_bridge.py new file mode 100644 index 0000000..74279a5 --- /dev/null +++ b/.agent/skills/extract-chapter-knowledge/scripts/dispatch_extraction_bridge.py @@ -0,0 +1,277 @@ +#!/usr/bin/env python3 +"""抽取智能体的框架派发桥(阶段 F 第二部分)。 + +职责分界(边界合同):智能体框架负责模型事件、raw、逐回合调用,记在派发 +运行下;生产编排负责知识草稿落库与业务回执,记在生产抽取运行下。两者以 +trigger_detail.productionRunId 关联。 + +与写手桥的两点不同: +1. 全量正典正文注入任务输入——证据必须是正文逐字片段,而只读工具读正文 + 会截断,注入是唯一可靠路径;工具白名单留给查重与回读核验。 +2. 单次派发(探索与产出同循环)——抽取产出是结构化 JSON,无长篇碎片化 + 风险,不需要两阶段。 + +机械校验、修复重派(一轮)、保守收口全部复用既有可信适配层,不新造合同。 +""" +from __future__ import annotations + +import json +import sys +from pathlib import Path +from typing import Any, Callable, Mapping + +SCRIPT_DIR = Path(__file__).resolve().parent +DISPATCH_SCRIPTS = SCRIPT_DIR.parents[1] / "dispatch-agent-task" / "scripts" +if str(DISPATCH_SCRIPTS) not in sys.path: + sys.path.insert(0, str(DISPATCH_SCRIPTS)) + +from extract_knowledge import ( # noqa: E402 + ACTOR, + ExtractionContractError, + _load_chapter, + _propose_chapter_extract_lesson, + normalize_extraction, + persist_extraction, + salvage_extraction, +) +from muse_llm import extract_json # noqa: E402 +from record_failed_run import record_failure # noqa: E402 +from run_registry import finish_run, new_run_id, start_run # noqa: E402 +from dispatch_agent_task import run_dispatch # noqa: E402 +from pi_runner import ExecutionPolicy # noqa: E402 +from read_tools import TOOL_REGISTRY # noqa: E402 + +# 抽取探索白名单:查重与回读核验;正文本体注入任务输入,不依赖工具读取。 +EXTRACTION_TOOL_ALLOWLIST = ("read_chapter_text", "search_entities") + +# 抽取产出是结构化 JSON;严格语义由 normalize_extraction 机械校验,Schema 只做形状兜底。 +EXTRACTION_OUTPUT_SCHEMA: dict[str, Any] = { + "$schema": "https://json-schema.org/draft/2020-12/schema", + "type": "object", + "required": ["entities", "relations", "state"], + "properties": { + "entities": {"type": "array"}, + "relations": {"type": "array"}, + "state": {"type": "object"}, + }, +} + +EXTRACTION_MAX_DURATION_SECONDS = 1800 +SESSION_ROOT = Path("/tmp/muse-agent-runs/extractor-sessions") + + +class DispatchExtractionError(RuntimeError): + """抽取派发失败:携带稳定错误码,编排方按失败关闭处理。""" + + def __init__(self, code: str, message: str, *, details: Mapping[str, Any] | None = None): + super().__init__(message) + self.code = code + self.details = dict(details or {}) + + +def extractor_session_paths(work_id: int, chapter_order: int) -> tuple[str, Path]: + """一章一个抽取智能体会话:修复重派在同一会话内接续。""" + + session_id = f"extractor-work{work_id}-ch{chapter_order}" + return session_id, SESSION_ROOT / session_id + + +def build_extraction_task_spec( + *, + work_id: int, + chapter_order: int, + title: str, + chapter_title: str, + existing_names: list[str], + body: str, + repair_reason: str | None = None, +) -> dict[str, Any]: + """装配抽取任务包:全量正文进输入,合同条款进任务提示词。""" + + task_prompt = ( + f"你是抽取智能体。任务:从作品《{title}》第{chapter_order}章" + f"《{chapter_title or ''}》的已接受正文抽取作品私有知识草稿。" + "正文全文在冻结输入的 body 字段;证据必须是正文中的逐字连续片段,不能改写。" + "实体类型只能使用:character、location、faction、power_system、item、event;" + "立卡门槛:具名且有跨章复用或后续履约潜力,一次性龙套与一次性道具不列实体。" + "不要确认知识,不要补写正文没有的事实;低置信内容保留但在 brief/fields 中标注“?”。" + "可用只读工具核对既有实体(查重)与回读正文,但正文以输入 body 为准。" + "只输出一个 JSON 对象:entities(每项 type/name/brief/fields/evidence)、" + "relations(每项 source/target/type/description/evidence)、" + "state(currentSituation/characterStates/foreshadowing/handoff)。" + ) + if repair_reason: + task_prompt += ( + "\n【机械校验失败,允许一次修复】失败原因:" + + repair_reason + + "。只修正证据字段,使每条 evidence 都是正文中的逐字连续片段;" + "删除无法找到逐字证据的条目,不得新增条目、事实、关系或状态。仍只输出同一 JSON 对象。" + ) + return { + "specVersion": "agent-task-v1", + "role": "extractor", + "taskPrompt": task_prompt, + "input": { + "workId": work_id, + "chapterOrder": chapter_order, + "title": title, + "chapterTitle": chapter_title or "", + "existingEntityNames": existing_names, + "body": body, + }, + "outputSchema": EXTRACTION_OUTPUT_SCHEMA, + "outputSchemaId": "chapter-extraction-v1", + "toolAllowlist": list(EXTRACTION_TOOL_ALLOWLIST), + "maxDurationSeconds": EXTRACTION_MAX_DURATION_SECONDS, + } + + +def _read_dispatch_output(receipt: Mapping[str, Any], dispatch_run_id: str) -> Any: + """结构化输出回读派发运行目录的 output.json,不另造权威。""" + + output_file = Path(str(receipt.get("runDir") or "")) / "output.json" + try: + raw_text = output_file.read_text(encoding="utf-8") + except OSError as exc: + raise DispatchExtractionError( + "EXTRACTION_OUTPUT_MISSING", + f"抽取派发运行未落结构化输出:{output_file}", + details={"dispatchRunId": dispatch_run_id}, + ) from exc + try: + return json.loads(raw_text) + except ValueError: + # 模型偶尔包 markdown/解释;用既有提取器兜底,仍失败则失败关闭。 + try: + return extract_json(raw_text) + except Exception as exc: + raise DispatchExtractionError( + "EXTRACTION_OUTPUT_INVALID", + f"抽取派发输出不是合法 JSON:{exc}", + details={"dispatchRunId": dispatch_run_id}, + ) from exc + + +def run_extraction_via_dispatch( + work_id: int, + chapter_order: int, + *, + repo_root: str | Path, + provider: str, + model: str, + thinking: str | None = None, + run_id: str | None = None, + spec_dir: str | Path | None = None, + launcher: Callable[..., Any] | None = None, + connect_factory: Callable[..., Any] | None = None, +) -> dict[str, Any]: + """派发抽取智能体并落知识草稿;返回落库摘要。任何失败失败关闭。""" + + for name in EXTRACTION_TOOL_ALLOWLIST: + if name not in TOOL_REGISTRY: + raise DispatchExtractionError( + "EXTRACTION_TOOL_UNREGISTERED", f"抽取白名单工具未登记:{name}" + ) + + active_run = run_id or new_run_id("extract-knowledge", work_id=work_id, target_chapter=chapter_order) + start_run( + run_id=active_run, + work_id=work_id, + target_chapter=chapter_order, + trigger_detail={"stage": "chapter-after-extraction", "mode": "dispatch"}, + creator=ACTOR, + ) + spec_root = Path(spec_dir) if spec_dir is not None else SCRIPT_DIR + spec_root.mkdir(parents=True, exist_ok=True) + session_id, session_dir = extractor_session_paths(work_id, chapter_order) + session_dir.mkdir(parents=True, mode=0o700, exist_ok=True) + policy = ExecutionPolicy(provider=provider, model=model, thinking=thinking) + + try: + (title, chapter_id, chapter_title, _block_id, body), existing = _load_chapter( + work_id, chapter_order + ) + existing_names = list(existing[:200]) + + def _dispatch_once(attempt: int, repair_reason: str | None) -> tuple[Any, Mapping[str, Any]]: + spec = build_extraction_task_spec( + work_id=work_id, chapter_order=chapter_order, title=title, + chapter_title=chapter_title or "", existing_names=existing_names, + body=body, repair_reason=repair_reason, + ) + spec_file = spec_root / f"{active_run}-extractor-task-v{attempt}.json" + spec_file.write_text(json.dumps(spec, ensure_ascii=False, indent=1), encoding="utf-8") + dispatch_run_id = f"{active_run}-extractor-v{attempt}" + receipt, code = run_dispatch( + spec_file, + repo_root=repo_root, + policy=policy, + run_id=dispatch_run_id, + trigger_source="user", + trigger_detail={"stage": "extractor-dispatch", "productionRunId": active_run}, + session_id=session_id, + session_dir=session_dir, + enable_read_tools=True, + launcher=launcher, + connect_factory=connect_factory, + ) + if code != 0 or receipt.get("status") != "completed": + raise DispatchExtractionError( + str(receipt.get("errorCode") or "EXTRACTION_DISPATCH_FAILED"), + f"抽取智能体派发未成功:{receipt.get('error') or receipt.get('errorCode')}", + details={"dispatchRunId": dispatch_run_id, "exitCode": code}, + ) + return _read_dispatch_output(receipt, dispatch_run_id), receipt + + raw_output, receipt = _dispatch_once(1, None) + try: + payload = normalize_extraction(raw_output, body) + except ExtractionContractError as first_error: + # 只允许一轮机械修复重派(证据绑定),对齐直调链语义。 + raw_output, receipt = _dispatch_once(2, str(first_error)) + try: + payload = normalize_extraction(raw_output, body) + except ExtractionContractError: + payload = salvage_extraction(raw_output, body) + + model_ids = receipt.get("actualModelIds") or [] + result = persist_extraction( + work_id, chapter_id, chapter_order, active_run, payload, + requested_model=f"{provider}/{model}", + actual_model=str(model_ids[-1]) if model_ids else "", + usage=dict(receipt.get("usage") or {}), + ) + finish_run(active_run, "completed", creator=ACTOR, + trigger_detail={"stage": "chapter-after-extraction", "mode": "dispatch", + "drafts": len(result["draft_ids"])}) + lesson = _propose_chapter_extract_lesson( + run_id=active_run, work_id=work_id, chapter_order=chapter_order, + draft_count=len(result["draft_ids"]), state_draft_id=result.get("state_draft_id"), + ) + return {"run_id": active_run, **result, "lesson": lesson} + except BaseException as exc: + finish_run(active_run, "failed", creator=ACTOR, + trigger_detail={"stage": "chapter-after-extraction", "mode": "dispatch", + "error_type": type(exc).__name__}) + try: + record_failure( + active_run, + sample_id=f"extract-ch{chapter_order}", + adapter_role="extractor", + caller="extract-knowledge-dispatch", + failure_type=type(exc).__name__, + ) + except Exception: + pass + raise + + +__all__ = [ + "EXTRACTION_MAX_DURATION_SECONDS", + "EXTRACTION_OUTPUT_SCHEMA", + "EXTRACTION_TOOL_ALLOWLIST", + "DispatchExtractionError", + "build_extraction_task_spec", + "extractor_session_paths", + "run_extraction_via_dispatch", +] diff --git a/.agent/skills/extract-chapter-knowledge/scripts/extract_via_dispatch.py b/.agent/skills/extract-chapter-knowledge/scripts/extract_via_dispatch.py new file mode 100644 index 0000000..1a06f3d --- /dev/null +++ b/.agent/skills/extract-chapter-knowledge/scripts/extract_via_dispatch.py @@ -0,0 +1,54 @@ +#!/usr/bin/env python3 +"""抽取智能体派发入口(阶段 F 第二部分)。 + +真实模型调用必须显式授权并显式给出 provider/model: +.venv/bin/python .agent/skills/extract-chapter-knowledge/scripts/extract_via_dispatch.py 12 3 \ + --provider catproxy-anthropic --model claude-opus-5 --thinking medium +""" +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path + +SCRIPT_DIR = Path(__file__).resolve().parent +if str(SCRIPT_DIR) not in sys.path: + sys.path.insert(0, str(SCRIPT_DIR)) + +from dispatch_extraction_bridge import run_extraction_via_dispatch # noqa: E402 + +REPO_ROOT = SCRIPT_DIR.parents[3] + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="派发抽取智能体做章后知识抽取(知识链)") + parser.add_argument("work_id", type=int) + parser.add_argument("chapter_order", type=int) + parser.add_argument("--provider", required=True) + parser.add_argument("--model", required=True) + parser.add_argument("--thinking", default=None) + parser.add_argument("--run-id", default=None) + args = parser.parse_args(argv) + + summary = run_extraction_via_dispatch( + args.work_id, + args.chapter_order, + repo_root=REPO_ROOT, + provider=args.provider, + model=args.model, + thinking=args.thinking, + run_id=args.run_id, + spec_dir=REPO_ROOT / "docs" / "write-chapter" / "artifacts", + ) + print(json.dumps( + {k: v for k, v in summary.items() if k != "lesson"}, + ensure_ascii=False, default=str, + )) + print(f"RUN_ID={summary['run_id']}") + print("知识草稿已落库(草稿态);经经验确认通道 :8766 人审后才转正。") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/docs/plans/2026-08-22-agent-example整体收敛总plan.md b/docs/plans/2026-08-22-agent-example整体收敛总plan.md index 0b02357..da577e1 100644 --- a/docs/plans/2026-08-22-agent-example整体收敛总plan.md +++ b/docs/plans/2026-08-22-agent-example整体收敛总plan.md @@ -173,6 +173,7 @@ A、B 可并行;C 依赖 A;E 依赖 D;F 依赖 E;G 依赖 D(链路透 - 意图:评委与抽取智能体接入框架派发(评委保持圈定授权隔离);用对照实验数据裁决直调链去留;处理写手探索冗余——生产链改为只给最小冻结输入(任务提示词 + 授权范围),由写作智能体真正自主探索取材,预组装完整上下文降为对照模式专用(2026-08-22 烟测实证:预组装下写手 0 次工具调用,探索无发生空间)。 - 第一部分(已完成,离线绿):两阶段写手落地——探索阶段(只读工具,产出探索清单)与生成阶段(无工具,按回放资料单次成稿)分离;`--two-phase` 旗标;详见 `docs/plans/2026-08-22-阶段F-第一部分-两阶段写手-探索与生成分离.md`。 +- 第二部分(已完成,真实派发通过):抽取智能体接入框架派发(全量正文注入任务输入,工具白名单留给查重探索;机械校验失败一轮修复重派,仍不合法走保守收口);候选 164 采纳后对第 3 章首次真实抽取产出 19 条草稿待人审;详见 `docs/plans/2026-08-22-阶段F-第二部分-知识链派发-抽取智能体.md`。 - 边界:对照只在显式对照模式运行;评测候选四层强制不可接受不变;知识草稿须经人确认转正。 - 验证:盲评隔离测试(禁看清单越权拒绝);对照产出可比正文与成本数据;退役项零引用后删除。 diff --git a/docs/plans/2026-08-22-阶段F-第二部分-知识链派发-抽取智能体.md b/docs/plans/2026-08-22-阶段F-第二部分-知识链派发-抽取智能体.md new file mode 100644 index 0000000..5777e7e --- /dev/null +++ b/docs/plans/2026-08-22-阶段F-第二部分-知识链派发-抽取智能体.md @@ -0,0 +1,48 @@ +# 阶段 F 第二部分:知识链派发——抽取智能体 + +日期:2026-08-23 +状态:已完成(离线绿 + 真实派发通过) +上游事实:候选 164 经决策通道采纳(署名 qingse),第 3 章正典落库(5131 字);接受意图 `queue_chapter_extraction`。 + +## 1. 意图 + +抽取智能体接入框架派发(边界合同:智能体框架记模型事件/raw,生产编排记知识草稿与业务回执),替代抽取直调链的生产使用;直调链保留为对照/回退,终态裁决按总 plan 用数据说话。 + +## 2. 设计 + +```text +生产编排:注册抽取运行 → 读正典章全量正文 + 既有实体名 + ↓ +派发抽取智能体(单次派发,探索与产出同一循环): + 输入 = 全量正文(注入,不经工具——工具读取会截断)+ 既有实体名 + 抽取合同 + 工具白名单 = read_chapter_text、search_entities(查重与回读核验用) + 产出 = 结构化抽取 JSON(entities/relations/state) + ↓ +机械校验:normalize_extraction(证据必须正文逐字片段) + 失败 → 一次修复重派(只纠证据绑定,不新增条目)→ 仍失败 → salvage 保守收口 + ↓ +落库:persist_extraction(muse_knowledge_draft,草稿态)→ 人确认转正(:8766) +``` + +关键取舍: + +- 正文注入而非工具读取:`read_chapter_text` 超长截断,证据逐字绑定要求全量正文,注入是唯一可靠路径;工具留给查重探索。 +- 单次派发(非两阶段):抽取产出是结构化 JSON,无长篇碎片化风险;探索与产出同循环。 +- 机械校验与落库全部复用既有可信适配层(`normalize_extraction`/`salvage_extraction`/`persist_extraction`),不新造合同。 +- 修复重派一轮(证据绑定是机械问题,对齐直调链既有语义);仍不合法走 salvage;salvage 失败失败关闭。 +- 证据归属:派发运行 `{run_id}-extractor-v{N}` 记事件/raw(trigger_detail 带 productionRunId),草稿落生产抽取运行。 + +## 3. 改动台账 + +| 文件 | 动作 | 原因 | +|---|---|---| +| `.agent/skills/extract-chapter-knowledge/scripts/dispatch_extraction_bridge.py` | 新增 | 抽取任务包装配 + 派发编排(校验/修复/落库/收口) | +| `.agent/skills/extract-chapter-knowledge/scripts/extract_via_dispatch.py` | 新增 | 派发入口(显式 provider/model,真实调用需授权) | +| `tests/skills/extract-chapter-knowledge/test_dispatch_extraction_bridge.py` | 新增 | 离线测试(7 项) | + +## 4. 验证结果 + +- 离线测试 7/7 通过。 +- 真实派发(作品 12 第 3 章,`extract-knowledge-w12-c3-20260823T150511`,catproxy-anthropic/claude-opus-5,medium): + - v1 派发完成但证据绑定未过 `normalize_extraction` → 桥自动触发 v2 修复重派 → 通过(修复路径首次真实验证)。 + - 18 条实体草稿 + 1 条叙事状态草稿落库(status=pending),看板与经验确认通道 :8766 可见,待人审转正。 diff --git a/tests/skills/extract-chapter-knowledge/test_dispatch_extraction_bridge.py b/tests/skills/extract-chapter-knowledge/test_dispatch_extraction_bridge.py new file mode 100644 index 0000000..905c251 --- /dev/null +++ b/tests/skills/extract-chapter-knowledge/test_dispatch_extraction_bridge.py @@ -0,0 +1,153 @@ +#!/usr/bin/env python3 +"""抽取智能体派发桥离线测试(阶段 F 第二部分)。 + +固定合同:任务包装配(正文注入/白名单/修复提示词)、单轮派发产出合法抽取、 +机械校验失败触发一轮修复重派、两轮失败走保守收口、派发失败失败关闭。 +""" +from __future__ import annotations + +import contextlib +import json +import pathlib +import sys +import tempfile +import unittest +from unittest import mock + +PROJECT_ROOT = pathlib.Path(__file__).resolve().parents[3] +SCRIPT_DIR = PROJECT_ROOT / ".agent" / "skills" / "extract-chapter-knowledge" / "scripts" +for path in (SCRIPT_DIR,): + if str(path) not in sys.path: + sys.path.insert(0, str(path)) + +import dispatch_extraction_bridge as bridge # noqa: E402 +from dispatch_extraction_bridge import ( # noqa: E402 + DispatchExtractionError, + build_extraction_task_spec, + extractor_session_paths, + run_extraction_via_dispatch, +) + +BODY = "第三章 波纹。茧从实验舰外舱门的裂口挤出去,金属边缘像被剥开的肋骨。林深盯着深渊方向。" + + +def _valid_payload(evidence="茧从实验舰外舱门的裂口挤出去"): + return { + "entities": [ + {"type": "character", "name": "林深", "brief": "试机师", "fields": {}, "evidence": "林深盯着深渊方向"}, + {"type": "item", "name": "茧", "brief": "生物机甲", "fields": {}, "evidence": evidence}, + ], + "relations": [ + {"source": "林深", "target": "茧", "type": "神经链接", "description": "链接", + "evidence": "茧从实验舰外舱门的裂口挤出去"}, + ], + "state": {"currentSituation": "出击", "characterStates": {}, "foreshadowing": {"埋": [], "推": [], "收": []}, "handoff": ""}, + } + + +def _load_chapter_fake(work_id, chapter_order): + return ("深渊机神", 301, "波纹", 901, BODY), ["林深", "茧"] + + +class ExtractionSpecTest(unittest.TestCase): + + def test_spec_shape(self): + spec = build_extraction_task_spec( + work_id=12, chapter_order=3, title="深渊机神", chapter_title="波纹", + existing_names=["林深"], body=BODY, + ) + self.assertEqual(spec["role"], "extractor") + self.assertEqual(spec["input"]["body"], BODY) + self.assertEqual(spec["input"]["existingEntityNames"], ["林深"]) + self.assertEqual(spec["toolAllowlist"], ["read_chapter_text", "search_entities"]) + self.assertIn("逐字连续片段", spec["taskPrompt"]) + self.assertNotIn("机械校验失败", spec["taskPrompt"]) + + def test_spec_repair_appends_reason(self): + spec = build_extraction_task_spec( + work_id=12, chapter_order=3, title="深渊机神", chapter_title="波纹", + existing_names=[], body=BODY, repair_reason="entities[0] 证据不在正文中", + ) + self.assertIn("机械校验失败", spec["taskPrompt"]) + self.assertIn("不得新增条目", spec["taskPrompt"]) + + def test_session_paths(self): + sid, path = extractor_session_paths(12, 3) + self.assertEqual(sid, "extractor-work12-ch3") + self.assertIn("extractor-sessions", str(path)) + + +class OrchestrationTest(unittest.TestCase): + + def setUp(self): + self._tmp = tempfile.TemporaryDirectory() + self.tmp = pathlib.Path(self._tmp.name) + self.dispatches = [] + + def tearDown(self): + self._tmp.cleanup() + + def _fake_dispatch(self, outputs): + """outputs: 每次派发依次返回的 JSON(Exception 表示派发失败)。""" + outputs = list(outputs) + + def fake(spec_file, **kwargs): + self.dispatches.append(kwargs["run_id"]) + run_dir = self.tmp / kwargs["run_id"] + run_dir.mkdir(parents=True) + out = outputs[len(self.dispatches) - 1] + if isinstance(out, Exception): + return {"status": "failed", "errorCode": "TIMEOUT"}, 1 + (run_dir / "output.json").write_text(json.dumps(out, ensure_ascii=False), encoding="utf-8") + return {"status": "completed", "runDir": str(run_dir), + "actualModelIds": ["catproxy-anthropic/claude-opus-5"], + "usage": {"input_tokens": 10}}, 0 + + return fake + + @contextlib.contextmanager + def _patched(self, outputs): + fake = self._fake_dispatch(outputs) + with mock.patch.object(bridge, "run_dispatch", fake), \ + mock.patch.object(bridge, "_load_chapter", _load_chapter_fake), \ + mock.patch.object(bridge, "start_run", lambda **kw: {"run_id": kw["run_id"]}), \ + mock.patch.object(bridge, "finish_run", lambda *a, **kw: None), \ + mock.patch.object(bridge, "record_failure", lambda *a, **kw: None), \ + mock.patch.object(bridge, "_propose_chapter_extract_lesson", lambda **kw: None), \ + mock.patch.object(bridge, "persist_extraction", + lambda *a, **kw: {"draft_ids": [1, 2, 3], "state_draft_id": 9}): + yield + + def _run(self, outputs): + with self._patched(outputs): + return run_extraction_via_dispatch( + 12, 3, repo_root=self.tmp, provider="catproxy-anthropic", + model="claude-opus-5", run_id="run-extract-test", spec_dir=self.tmp, + ) + + def test_single_dispatch_success(self): + summary = self._run([_valid_payload()]) + self.assertEqual(summary["run_id"], "run-extract-test") + self.assertEqual(len(summary["draft_ids"]), 3) + self.assertEqual(self.dispatches, ["run-extract-test-extractor-v1"]) + + def test_repair_round_when_evidence_invalid(self): + bad = _valid_payload(evidence="不在正文里的证据") + summary = self._run([bad, _valid_payload()]) + self.assertEqual(self.dispatches, ["run-extract-test-extractor-v1", "run-extract-test-extractor-v2"]) + self.assertEqual(len(summary["draft_ids"]), 3) + + def test_salvage_after_two_failures(self): + bad = _valid_payload(evidence="不在正文里的证据") + summary = self._run([bad, bad]) + # salvage:证据无效但实体名逐字在正文 → 保留并以实体名为最小证据;关系缺证据丢弃 + self.assertEqual(len(summary["draft_ids"]), 3) + self.assertEqual(self.dispatches, ["run-extract-test-extractor-v1", "run-extract-test-extractor-v2"]) + + def test_dispatch_failure_fails_closed(self): + with self.assertRaises(DispatchExtractionError): + self._run([TimeoutError()]) + + +if __name__ == "__main__": + unittest.main()