diff --git a/.agent/skills/write-next-chapter/SKILL.md b/.agent/skills/write-next-chapter/SKILL.md index 990a8b6..2364631 100644 --- a/.agent/skills/write-next-chapter/SKILL.md +++ b/.agent/skills/write-next-chapter/SKILL.md @@ -40,11 +40,14 @@ Writer 不接收 `runId`、权限信息、manifest、hash、候选版本、验 Dashboard 与人工生产入口统一调用 `scripts/produce_next_chapter.py`。该入口只负责编排本 Skill 已登记的冻结、检索、writer、detector、CAS、候选落库和人闸步骤;不提供自动 accept。运行 artifacts 仍由只读看板按 run_id 读取。 +- 默认走直调链(`run_writer_with_receipt`)。`--dispatch-writer --provider P --model M [--thinking T]` 改走框架派发链:写作智能体经 `dispatch-agent-task` 自主探索取材并产正文(`scripts/dispatch_writer_bridge.py`),一章一个稳定会话(复用=继续原写作智能体)。派发链的模型证据记在派发运行下,生产候选账本以显式 `writer_raw_ref` 关联,不双套记账。 +- 授权继续:管线停在 `AUTHORIZATION_REQUIRED` 时留痕 `artifacts/-authorization-request.json`;人授权后以 `--continue-from <前序 run_id>` 发起新运行,编排方按缺口补证重组上下文后继续。 + ## 生产落库 - 生产编排走 `run_writer_pipeline`(机械门→语义 detector→单次收敛;证据缺口收敛为 AUTHORIZATION_REQUIRED 授权终态,补证/重写经人授权后以新运行继续),状态链用 `scripts/candidate_cas.py` 的 `PostgresCasStateStore` 持久化到 `example_candidate_cas`(一次运行一条链,revision 单调,DB 触发器锁方向闭集);内存 `InMemoryCasStateStore` 仅供离线测试。 - `run_writer_with_receipt()` 只负责可信 writer adapter;生产编排在组装后调用 `assemble-context/scripts/persist_context_freeze.py`,审查落库调用 `scripts/persist_writer_run.py`。 -- `persist_writer_run.py` 要求本次 `run_id` 已有成功 writer 调用的 raw 指针,随后登记 `example_candidate`(传入 `semantic_report` 时校验绑定并把 `semantic_status`/`semantic_report_sha256` 固化到候选行)、追加 `example_run_receipt` 和机械/语义两行 `example_quality_result`;它不接受正文,Shadow→Canonical 仍只能由 `decide-candidate/scripts/write_canonical.py` 完成(接受通道 DB 级兜底复检 `state=passed` 且 `semantic_status=passed`)。 +- `persist_writer_run.py` 要求成功 writer 调用的 raw 指针(直调链查本运行调用账;派发链由编排显式传 `writer_raw_ref`),随后登记 `example_candidate`(传入 `semantic_report` 时校验绑定并把 `semantic_status`/`semantic_report_sha256` 固化到候选行)、追加 `example_run_receipt` 和机械/语义两行 `example_quality_result`;它不接受正文,Shadow→Canonical 仍只能由 `decide-candidate/scripts/write_canonical.py` 完成(接受通道 DB 级兜底复检 `state=passed` 且 `semantic_status=passed`)。 ## 输出合同 diff --git a/.agent/skills/write-next-chapter/scripts/dispatch_writer_bridge.py b/.agent/skills/write-next-chapter/scripts/dispatch_writer_bridge.py new file mode 100644 index 0000000..2b72cf4 --- /dev/null +++ b/.agent/skills/write-next-chapter/scripts/dispatch_writer_bridge.py @@ -0,0 +1,239 @@ +#!/usr/bin/env python3 +"""生产写手的框架派发桥(阶段 E 第二部分)。 + +职责分界(边界合同):智能体框架负责模型事件、raw、逐回合调用,记在派发运行下; +生产编排负责候选、CAS、质量结果与业务回执,记在生产运行下。两者以显式 +`writer_raw_ref`(调用账 id + raw 内容 id)关联,禁止同一模型回合两套记账。 + +桥只做三件事: +1. 把 WriterContext 装配成可移植任务包(角色=writer、冻结创作输入、输出 Schema); +2. 经 dispatch-agent-task 公共入口派发(显式 provider/model、会话复用、只读探索工具); +3. 把派发产出绑定为可信候选信封(身份、哈希、版本由本桥绑定,模型只产正文)。 +""" +from __future__ import annotations + +import json +import sys +from dataclasses import dataclass +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" +ASSEMBLE_SCRIPTS = SCRIPT_DIR.parents[1] / "assemble-context" / "scripts" +for _path in (DISPATCH_SCRIPTS, ASSEMBLE_SCRIPTS): + if str(_path) not in sys.path: + sys.path.insert(0, str(_path)) + +from writer_contract import build_candidate_envelope, build_writer_creative_input # 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 + +# writer 探索白名单:工具 server 登记表的全部只读工具(登记表是唯一事实源)。 +READ_TOOL_ALLOWLIST = tuple(sorted(TOOL_REGISTRY)) + +# writer 只产正文;篇幅合同由机械门与可信适配层校验,不在 Schema 里重复定义语义。 +WRITER_DISPATCH_OUTPUT_SCHEMA: dict[str, Any] = { + "$schema": "https://json-schema.org/draft/2020-12/schema", + "type": "object", + "additionalProperties": False, + "required": ["candidateBody"], + "properties": {"candidateBody": {"type": "string", "minLength": 1}}, +} + +DEFAULT_MAX_DURATION_SECONDS = 2400 +SESSION_ROOT = Path("/tmp/muse-agent-runs/writer-sessions") + + +class DispatchWriterError(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 {}) + + +@dataclass(frozen=True) +class DispatchWriterReceipt: + """派发回执的生产账本适配:persist_writer_execution 以属性方式读取。""" + + requested_model_id: str + actual_model_id: str + model_match: bool + usage: Mapping[str, Any] + total_cost_usd: float | None + effort: str | None + stop_reason: str + terminal_reason: str + is_error: bool + dispatch_run_id: str + + +def writer_session_paths(work_id: int, target_chapter: int) -> tuple[str, Path]: + """一章一个写作智能体会话:会话 ID 与目录按作品/章稳定,跨运行复用。""" + + session_id = f"writer-work{work_id}-ch{target_chapter}" + session_dir = SESSION_ROOT / session_id + return session_id, session_dir + + +def build_writer_dispatch_spec( + context: Mapping[str, Any], + *, + target_chapter: int, + human_instruction: str, + candidate_version: int, +) -> dict[str, Any]: + """装配可移植任务包:冻结创作输入原样进 input,不掺框架字段。""" + + instruction = (human_instruction or "").strip() or "按冻结创作输入续写本章完整正文。" + return { + "specVersion": "agent-task-v1", + "role": "writer", + "taskPrompt": ( + f"为第{target_chapter}章写正文候选(候选版本 {candidate_version})。" + f"人的创作指令:{instruction}" + "先用授权只读工具补齐动笔所需的细纲、人物状态与前文衔接,再写整章;" + "正文中的事实必须来自你实际读到的资料。" + ), + "input": { + "workId": context.get("workId"), + "targetChapter": target_chapter, + "candidateVersion": candidate_version, + "creativeInput": build_writer_creative_input(context), + }, + "outputSchema": WRITER_DISPATCH_OUTPUT_SCHEMA, + "outputSchemaId": "writer-candidate-body-v1", + "toolAllowlist": list(READ_TOOL_ALLOWLIST), + "maxDurationSeconds": DEFAULT_MAX_DURATION_SECONDS, + } + + +def _writer_raw_ref(dispatch_run_id: str, connect_factory: Callable[..., Any] | None) -> tuple[Any, Any]: + """从派发运行的模型调用账取最新一条带 raw 的成功调用;缺则失败关闭。""" + + factory = connect_factory + if factory is None: + import muse_db + + factory = lambda: muse_db.connect(readonly=True) # noqa: E731 + with factory() as conn: + row = conn.execute( + "SELECT id, raw_content_id FROM example_llm_call " + "WHERE run_id=%s AND out_tokens>0 AND raw_content_id IS NOT NULL " + "ORDER BY id DESC LIMIT 1", + (dispatch_run_id,), + ).fetchone() + if not row: + raise DispatchWriterError( + "DISPATCH_EVIDENCE_MISSING", + f"派发运行缺少带 raw 的成功模型调用账:{dispatch_run_id}", + ) + return row[0], row[1] + + +def run_writer_via_dispatch( + context: Mapping[str, Any], + *, + candidate_version: int, + repo_root: str | Path, + provider: str, + model: str, + thinking: str | None = None, + human_instruction: str = "", + spec_path: str | Path | None = None, + launcher: Callable[..., Any] | None = None, + connect_factory: Callable[..., Any] | None = None, +) -> tuple[dict[str, Any], DispatchWriterReceipt, tuple[Any, Any]]: + """派发一次写作智能体并绑定候选信封;返回(信封、回执适配、证据引用)。 + + 任何失败抛 DispatchWriterError(失败关闭);不产生半绑定候选。 + """ + + run_id = str(context.get("runId") or "") + work_id = context.get("workId") + target_chapter = context.get("targetChapter") + if not run_id or not isinstance(work_id, int) or not isinstance(target_chapter, int): + raise DispatchWriterError("DISPATCH_CONTEXT_INVALID", "WriterContext 缺 runId/workId/targetChapter") + + spec = build_writer_dispatch_spec( + context, + target_chapter=target_chapter, + human_instruction=human_instruction, + candidate_version=candidate_version, + ) + spec_file = Path(spec_path) if spec_path is not None else SCRIPT_DIR / f"{run_id}-writer-task-v{candidate_version}.json" + spec_file.parent.mkdir(parents=True, exist_ok=True) + spec_file.write_text(json.dumps(spec, ensure_ascii=False, indent=1), encoding="utf-8") + + session_id, session_dir = writer_session_paths(work_id, target_chapter) + session_dir.mkdir(parents=True, mode=0o700, exist_ok=True) + dispatch_run_id = f"{run_id}-writer-v{candidate_version}" + + policy = ExecutionPolicy(provider=provider, model=model, thinking=thinking) + receipt, code = run_dispatch( + spec_file, + repo_root=repo_root, + policy=policy, + run_id=dispatch_run_id, + trigger_source="user", + trigger_detail={"stage": "writer-dispatch", "productionRunId": run_id}, + 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 DispatchWriterError( + str(receipt.get("errorCode") or "DISPATCH_FAILED"), + f"写作智能体派发未成功:{receipt.get('error') or receipt.get('errorCode')}", + details={"dispatchRunId": dispatch_run_id, "exitCode": code}, + ) + + # 结构化输出与事件账本同源:回读派发运行目录的 output.json,不另造权威。 + body = None + output_file = Path(str(receipt.get("runDir") or "")) / "output.json" + try: + loaded = json.loads(output_file.read_text(encoding="utf-8")) + if isinstance(loaded, Mapping): + body = loaded.get("candidateBody") + except (OSError, ValueError): + body = None + if not isinstance(body, str) or not body.strip(): + raise DispatchWriterError( + "DISPATCH_OUTPUT_INVALID", "写作智能体未返回可用正文", + details={"dispatchRunId": dispatch_run_id}, + ) + + envelope = build_candidate_envelope(context, {"candidateBody": body}, candidate_version=candidate_version) + raw_ref = _writer_raw_ref(dispatch_run_id, connect_factory) + + model_ids = receipt.get("actualModelIds") or [] + adapter = DispatchWriterReceipt( + requested_model_id=str(receipt.get("requestedModelId") or f"{provider}/{model}"), + actual_model_id=str(model_ids[-1]) if model_ids else "", + model_match=bool(model_ids) and all(mid == receipt.get("requestedModelId") for mid in model_ids), + usage=dict(receipt.get("usage") or {}), + total_cost_usd=receipt.get("totalCostUsd"), + effort=thinking, + stop_reason="completed", + terminal_reason="writer-dispatch", + is_error=False, + dispatch_run_id=dispatch_run_id, + ) + return envelope, adapter, raw_ref + + +__all__ = [ + "DEFAULT_MAX_DURATION_SECONDS", + "DispatchWriterError", + "DispatchWriterReceipt", + "READ_TOOL_ALLOWLIST", + "WRITER_DISPATCH_OUTPUT_SCHEMA", + "build_writer_dispatch_spec", + "run_writer_via_dispatch", + "writer_session_paths", +] diff --git a/.agent/skills/write-next-chapter/scripts/persist_writer_run.py b/.agent/skills/write-next-chapter/scripts/persist_writer_run.py index 97c7b55..d70e84b 100644 --- a/.agent/skills/write-next-chapter/scripts/persist_writer_run.py +++ b/.agent/skills/write-next-chapter/scripts/persist_writer_run.py @@ -133,12 +133,17 @@ def persist_writer_execution( *, semantic_report: Mapping[str, Any] | None = None, assemble_result: Mapping[str, Any] | None = None, + writer_raw_ref: tuple[Any, Any] | None = None, dry_run: bool = False, ) -> dict[str, Any]: """落库一次 writer Shadow 运行;失败不会创建半套候选账本。 semantic_report 非 None 时校验绑定并把语义状态固化到候选行(接受通道的 DB 兜底依据); None 时候选 semantic_status 留空,write_canonical 一律不得接受(先审后入,失败关闭)。 + + writer_raw_ref:(调用账 id, raw 内容 id)。直调链在本运行下有 caller='writer' + 的成功调用账,缺省即从本运行查;框架派发链的模型证据记在派发运行下, + 由编排方查得后显式传入。两条路径都要求真实证据,缺一失败关闭。 """ run_id = str(context.get("runId") or "") @@ -173,7 +178,13 @@ def persist_writer_execution( try: with connect() as conn: - call_id, raw_content_id = _raw_response(conn, run_id) + if writer_raw_ref is not None: + call_id, raw_content_id = writer_raw_ref + if call_id is None or raw_content_id is None: + raise WriterPersistenceError( + "writer_raw_ref 缺调用账或 raw 内容:禁止写入看似完整的候选") + else: + call_id, raw_content_id = _raw_response(conn, run_id) existing = conn.execute( "SELECT id,candidate_sha256,state FROM example_candidate " "WHERE tenant_id=0 AND work_id=%s AND target_chapter=%s AND candidate_version=%s", diff --git a/.agent/skills/write-next-chapter/scripts/produce_next_chapter.py b/.agent/skills/write-next-chapter/scripts/produce_next_chapter.py index ff9329c..5371a6d 100644 --- a/.agent/skills/write-next-chapter/scripts/produce_next_chapter.py +++ b/.agent/skills/write-next-chapter/scripts/produce_next_chapter.py @@ -82,6 +82,7 @@ from muse_role import ( # noqa: E402 from run_writer_replay import profile_from_mapping # noqa: E402 from persist_llm_call import persist_call as persist_llm_event # noqa: E402 from persist_writer_run import persist_writer_execution # noqa: E402 +from dispatch_writer_bridge import DispatchWriterError, run_writer_via_dispatch # noqa: E402 from run_registry import finish_run, start_run # noqa: E402 from check_writer_acceptance import AcceptanceError, check_writer_acceptance # noqa: E402 from acceptance_state import LiveStateError, build_live_acceptance_state # noqa: E402 @@ -275,6 +276,26 @@ def main(): raise SystemExit("--instruction 后须跟本轮人指令原文") human_instruction = argv[idx + 1].strip() argv = argv[:idx] + argv[idx + 2 :] + def _take(flag: str): + nonlocal argv + if flag not in argv: + return None + idx = argv.index(flag) + if idx + 1 >= len(argv) or argv[idx + 1].startswith("--"): + raise SystemExit(f"{flag} 后须跟值") + value = argv[idx + 1] + argv = argv[:idx] + argv[idx + 2:] + return value + + dispatch_mode = "--dispatch-writer" in argv + if dispatch_mode: + argv = [arg for arg in argv if arg != "--dispatch-writer"] + dispatch_provider = _take("--provider") + dispatch_model = _take("--model") + dispatch_thinking = _take("--thinking") + continue_from = _take("--continue-from") + if dispatch_mode and (not dispatch_provider or not dispatch_model): + raise SystemExit("--dispatch-writer 必须显式给出 --provider 与 --model") args = [arg for arg in argv if not arg.startswith("--")] target = resolve_target_chapter(int(args[0]) if args else None) as_of = target - 1 @@ -401,11 +422,50 @@ def main(): receipts_by_version: dict[int, Any] = {} candidates_by_version: dict[int, dict] = {} contexts_by_attempt: dict[int, dict] = {} + writer_raw_refs: dict[int, tuple[Any, Any]] = {} + + def _diagnose_candidate(candidate_body: str, candidate_version: int) -> None: + """人感技能 3:每个候选先做只读诊断并自动落质量账;不在这里改正文。""" + + deai_artifact = run_diagnosis( + candidate_body, work_ref=f"work:{WORK_ID}", + chapter_ref=f"chapter:{target}", mode="Audit", + ) + _dump(ARTIFACTS / f"{run_id}-ai-flavor-diagnosis-v{candidate_version}.json", deai_artifact) + persist_diagnosis(deai_artifact, text=candidate_body) + print(f"AI 味诊断(v{candidate_version}):发现 {len(deai_artifact['findings'])} 条") def production_writer(current_context: Mapping[str, Any], candidate_version: int): - """writer 适配:篇幅越界自动重抽(≤3 遍),成功回执按候选版本留证。""" + """writer 适配:直调链篇幅越界自动重抽(≤3 遍);派发链失败关闭。""" contexts_by_attempt[current_context["attempt"]] = dict(current_context) + if dispatch_mode: + try: + candidate, receipt, raw_ref = run_writer_via_dispatch( + current_context, + candidate_version=candidate_version, + repo_root=REPO_ROOT, + provider=dispatch_provider, + model=dispatch_model, + thinking=dispatch_thinking, + human_instruction=human_instruction, + spec_path=ARTIFACTS / f"{run_id}-writer-task-v{candidate_version}.json", + ) + except DispatchWriterError as exc: + raise PipelineError(exc.code, f"writer 派发失败: {exc}", + details=exc.details) from exc + receipts_by_version[candidate_version] = receipt + candidates_by_version[candidate_version] = candidate + writer_raw_refs[candidate_version] = raw_ref + try: + _diagnose_candidate(candidate["candidateBody"], candidate_version) + except Exception as exc: + raise PipelineError( + "AI_FLAVOR_DIAGNOSIS_FAILED", "候选 AI 味诊断或落库失败", + details={"errorType": type(exc).__name__, "message": str(exc)}, + ) from exc + print(f"writer 派发模式: dispatch_run={receipt.dispatch_run_id}") + return candidate for retry in range(1, MAX_ATTEMPTS + 1): try: candidate, receipt = run_writer_with_receipt( @@ -418,20 +478,13 @@ def main(): ) receipts_by_version[candidate_version] = receipt candidates_by_version[candidate_version] = candidate - # 人感技能 3:每个候选先做只读诊断并自动落质量账;不在这里改正文。 try: - deai_artifact = run_diagnosis( - candidate["candidateBody"], work_ref=f"work:{WORK_ID}", - chapter_ref=f"chapter:{target}", mode="Audit", - ) - _dump(ARTIFACTS / f"{run_id}-ai-flavor-diagnosis-v{candidate_version}.json", deai_artifact) - persist_diagnosis(deai_artifact, text=candidate["candidateBody"]) + _diagnose_candidate(candidate["candidateBody"], candidate_version) except Exception as exc: raise PipelineError( "AI_FLAVOR_DIAGNOSIS_FAILED", "候选 AI 味诊断或落库失败", details={"errorType": type(exc).__name__, "message": str(exc)}, ) from exc - print(f"AI 味诊断(v{candidate_version}):发现 {len(deai_artifact['findings'])} 条") return candidate except WriterAdapterError as exc: if exc.code != "candidate_length_out_of_range" or retry == MAX_ATTEMPTS: @@ -487,9 +540,26 @@ def main(): (WORK_ID, target)).fetchone() initial_candidate_version = int(max_version_row[0]) + 1 + if continue_from: + request_path = ARTIFACTS / f"{continue_from}-authorization-request.json" + try: + request = json.loads(request_path.read_text(encoding="utf-8")) + except (OSError, ValueError) as exc: + raise SystemExit(f"授权请求读取失败: {request_path}: {exc}") + if request.get("workId") != WORK_ID or request.get("targetChapter") != target: + raise SystemExit("授权请求与本次作品/章不匹配,拒绝继续") + gaps = request.get("evidenceGaps") or [] + if not gaps: + raise SystemExit("授权请求无证据缺口,无需继续") + writer_context = production_evidence_provider(writer_context, gaps, writer_context["attempt"] + 1) + print(f"[授权继续] 前序运行 {continue_from}:命中缺口 {len(gaps)} 项," + f"本次以补证后上下文(attempt={writer_context['attempt']})继续。") + start_run(run_id=run_id, work_id=WORK_ID, target_chapter=target, trigger_detail={"stage": "writer-production-pipeline", - "contextSha256": writer_context["contextSnapshot"]["contextSha256"]}, + "contextSha256": writer_context["contextSnapshot"]["contextSha256"], + **({"continueFrom": continue_from} if continue_from else {}), + **({"writerMode": "dispatch"} if dispatch_mode else {})}, creator="continuation") try: pipeline_result = run_writer_pipeline( @@ -517,7 +587,8 @@ def main(): contexts_by_attempt.get(failed_candidate.get("attempt"), writer_context), failed_candidate, failed_receipt, audit_entry["mechanicalReport"], semantic_report=audit_entry.get("semanticReport"), - assemble_result=assembled) + assemble_result=assembled, + writer_raw_ref=writer_raw_refs.get(failed_candidate.get("candidateVersion"))) print(f"[被拒候选留库] candidate_id={persisted['candidate_id']} " f"state={persisted['state']} semantic={persisted.get('semantic_status')}") except Exception as persist_exc: # 留痕失败不掩盖原始失败码 @@ -526,6 +597,19 @@ def main(): trigger_detail={"stage": "writer-production-pipeline", "failureCode": exc.code}) if exc.code == "AUTHORIZATION_REQUIRED": gaps = (exc.details or {}).get("evidenceGaps") or [] + _dump(ARTIFACTS / f"{run_id}-authorization-request.json", { + "schemaVersion": "authorization-request-v1", + "runId": run_id, + "workId": WORK_ID, + "targetChapter": target, + "humanInstruction": human_instruction, + "evidenceGaps": gaps, + "nextAttempt": (exc.details or {}).get("nextAttempt"), + "reassembledContextSha256": (exc.details or {}).get("reassembledContextSha256"), + "candidateSha256": (exc.result or {}).get("candidateSha256"), + "generatedAt": generated_at, + }) + print(f"授权请求已留痕: artifacts/{run_id}-authorization-request.json") print("[需要授权] 语义检查发现证据缺口,补证或重写需要人授权:") for item in gaps: print(f" - {item.get('gapId')}: {item.get('reason')}(检索:{item.get('query')})") @@ -551,7 +635,8 @@ def main(): receipt = receipts_by_version[pipeline_result["candidateVersion"]] persisted = persist_writer_execution( final_context, candidate, receipt, final_trace["mechanicalReport"], - semantic_report=final_trace.get("semanticReport"), assemble_result=assembled) + semantic_report=final_trace.get("semanticReport"), assemble_result=assembled, + writer_raw_ref=writer_raw_refs.get(candidate.get("candidateVersion"))) cand_id = persisted["candidate_id"] print(f"writer 产出: sha256={candidate['candidateSha256'][:24]}..., " f"实际模型={receipt.actual_model_id}, 成本=${receipt.total_cost_usd}") diff --git a/docs/plans/2026-08-22-阶段E-生产链切换-第二部分-派发接线与授权继续.md b/docs/plans/2026-08-22-阶段E-生产链切换-第二部分-派发接线与授权继续.md new file mode 100644 index 0000000..a6390d8 --- /dev/null +++ b/docs/plans/2026-08-22-阶段E-生产链切换-第二部分-派发接线与授权继续.md @@ -0,0 +1,44 @@ +# 阶段 E:生产链切换(第二部分:派发接线与授权继续) + +日期:2026-08-22 +总 plan:[2026-08-22-agent-example整体收敛总plan.md](2026-08-22-agent-example整体收敛总plan.md) +状态:离线接线完成并提交;真实生产烟测(一章正文经派发链生成、停在人闸)待人工授权。 + +## 1. 意图 + +把写作节点接入框架派发链(可选显式模式),把授权终态接成可继续的流程:授权请求留痕、人授权后以新运行继续补证。证据归属按边界合同分界:智能体框架负责模型事件、raw、逐回合调用(记派发运行);生产编排负责候选、CAS、质量结果与业务回执(记生产运行);两者以显式证据引用关联,禁止同一模型回合两套记账。 + +## 2. 实现落点 + +- `dispatch_writer_bridge.py`(新增):写作智能体派发桥。任务包装配(角色=writer、冻结创作输入、输出 Schema=candidateBody、探索白名单=工具 server 登记表全部只读工具);会话按作品/章稳定(`writer-work{W}-ch{T}`,跨运行复用同一写作智能体);产出经可信信封绑定(身份、哈希、版本由桥绑定,模型只产正文);证据引用从派发运行调用账读取,缺则失败关闭。 +- `produce_next_chapter.py`: + - `--dispatch-writer --provider --model [--thinking]`:写作节点走派发链,失败关闭(不做篇幅自动重抽,篇幅合同由机械门把关)。 + - `AUTHORIZATION_REQUIRED` 时留痕 `artifacts/-authorization-request.json`(缺口、下一 attempt、重组上下文哈希、候选哈希、人指令)。 + - `--continue-from <前序 run_id>`:校验授权请求的作品/章匹配后,按缺口补证重组上下文,以新运行继续;运行登记记录 `continueFrom` 与 `writerMode`。 + - 两处候选落库透传 `writer_raw_ref`(直调链为 None,走原有本运行查账)。 +- `persist_writer_run.py`:`persist_writer_execution` 增加可选 `writer_raw_ref`;给则用显式引用(缺项失败关闭),不给则维持原有本运行查账的严格合同。 + +## 3. 文件台账 + +| 处置 | 文件 | 原因 | +|---|---|---| +| 新增 | `scripts/dispatch_writer_bridge.py` | 生产链接入框架派发链的唯一桥;入口为 `produce_next_chapter --dispatch-writer` | +| 修改 | `scripts/produce_next_chapter.py` | 派发模式、授权请求留痕、授权继续、证据引用透传;诊断逻辑提取为助手函数复用 | +| 修改 | `scripts/persist_writer_run.py` | 显式证据引用路径;直调链合同不变 | +| 新增 | `tests/.../test_dispatch_writer_bridge.py` | 桥合同离线测试:任务包冻结、会话稳定、信封绑定、派发失败与缺证失败关闭(5 项) | +| 修改 | `tests/.../test_persist_writer_run.py` | 显式引用绕过本运行查账、引用缺项失败关闭(2 项) | +| 修改 | `SKILL.md`(write-next-chapter) | 登记派发模式、授权继续与证据分界合同 | +| 修改 | `harness/manifests/test-inventory.json` | 登记新测试条目(111 条) | +| 保留 | 直调链(默认模式) | 对照与回退路径;终态裁决归阶段 F | +| 遗留 | 真实生产烟测 | 需人工授权(一章正文,派发链,停在人闸) | +| 遗留 | 检测智能体派发化 | 语义检测仍走直调链;接入框架派发归阶段 F(含评委) | + +## 4. 验证 + +- 桥离线测试:5 项通过(任务包装配、会话按章稳定、成功绑定信封+证据引用、派发失败关闭、缺证失败关闭)。 +- 落库显式引用测试:2 项通过;原有 2 项不回归。 +- write-next-chapter 全部测试:10 个文件全过(含真实库 CAS 集成)。 +- 架构门禁:导入边界与索引测试通过(跨 Skill 引用遵循既有先例方式)。 +- 全量离线清单:103 通过、8 项外部依赖阻断、0 失败(较第一部分基线 +1)。 +- 技能严格审计:58 个 Skill,阻断 0;索引一致性通过。 +- 真实生产烟测:待授权。建议参数:作品 12、目标章按库内进度、显式 `catproxy-anthropic/claude-opus-5`、思考等级待人定;验收按总 plan §11 的可见性矩阵走查。 diff --git a/harness/manifests/test-inventory.json b/harness/manifests/test-inventory.json index 1a8fc4a..c704ea7 100644 --- a/harness/manifests/test-inventory.json +++ b/harness/manifests/test-inventory.json @@ -1719,6 +1719,23 @@ "classification_confidence": "high", "classification_basis": "Drives the in-memory writer/mechanical/semantic/CAS pipeline and atomic temporary result writes with fake detectors." }, + { + "path": "tests/skills/write-next-chapter/test_dispatch_writer_bridge.py", + "scope": "runtime_skill", + "owner_skill_or_domain": "write-next-chapter", + "kind": "runtime_contract", + "evidence_level": "deterministic_offline", + "requires": [ + "offline", + "filesystem" + ], + "side_effects": [ + "filesystem" + ], + "skill_behavior_eval": false, + "classification_confidence": "high", + "classification_basis": "Exercises writer dispatch bridge contract with fake pi streams and injected DB connections: spec freeze, per-chapter session identity, envelope binding, evidence-reference fail-closed; no network, no real framework binary." + }, { "path": "tests/skills/write-next-chapter/test_semantic_verdict.py", "scope": "runtime_skill", @@ -1803,10 +1820,10 @@ } ], "summary": { - "entry_count": 110, + "entry_count": 111, "by_scope": { "other": 1, - "runtime_skill": 98, + "runtime_skill": 99, "harness": 3, "domain": 8 }, @@ -1819,13 +1836,13 @@ "tool_unit": 32, "fake_pipeline": 22, "runtime_probe": 1, - "runtime_contract": 3 + "runtime_contract": 4 }, "by_evidence_level": { - "deterministic_offline": 94, + "deterministic_offline": 95, "real_dependency_integration": 10, "static_structure": 6 }, - "total": 110 + "total": 111 } } \ No newline at end of file diff --git a/tests/skills/write-next-chapter/test_dispatch_writer_bridge.py b/tests/skills/write-next-chapter/test_dispatch_writer_bridge.py new file mode 100644 index 0000000..9e14203 --- /dev/null +++ b/tests/skills/write-next-chapter/test_dispatch_writer_bridge.py @@ -0,0 +1,173 @@ +#!/usr/bin/env python3 +"""生产写手框架派发桥的离线测试(不连真实框架、不连真实模型)。 + +固定合同:任务包装配(角色/冻结输入/探索白名单)、会话按章稳定、 +成功路径绑定候选信封与证据引用、派发失败与证据缺失一律失败关闭。 +""" +from __future__ import annotations + +import json +import pathlib +import sys +import tempfile +import unittest + +PROJECT_ROOT = pathlib.Path(__file__).resolve().parents[3] +SCRIPT_DIR = PROJECT_ROOT / ".agent" / "skills" / "write-next-chapter" / "scripts" +DISPATCH_TEST_DIR = PROJECT_ROOT / "tests" / "skills" / "dispatch-agent-task" +CHECK_TEST_DIR = PROJECT_ROOT / "tests" / "skills" / "check-content-consistency" +WNC_TEST_DIR = PROJECT_ROOT / "tests" / "skills" / "write-next-chapter" +for path in (SCRIPT_DIR, DISPATCH_TEST_DIR, CHECK_TEST_DIR, WNC_TEST_DIR): + if str(path) not in sys.path: + sys.path.insert(0, str(path)) + +import dispatch_writer_bridge as bridge # noqa: E402 +from dispatch_agent_task import DEFAULT_RUN_DIR_ROOT # noqa: E402 +from dispatch_writer_bridge import ( # noqa: E402 + DispatchWriterError, + build_writer_dispatch_spec, + run_writer_via_dispatch, + writer_session_paths, +) +from read_tools import TOOL_REGISTRY # noqa: E402 +from test_check_writer_candidate import _valid_pair # noqa: E402 +from test_dispatch_agent_task import RecordingConnect, fake_launcher, pi_stream_lines # noqa: E402 + + +class BridgeConnect(RecordingConnect): + """派发桥专用假连接:补证证据引用查询返回 raw_row(可置 None 模拟缺证)。""" + + def __init__(self): + super().__init__() + self.raw_row = (9001, 7777) + + def __call__(self, *args, **kwargs): + outer = self + + class _Cursor: + def execute(self, sql, params=None): + outer.log.append((sql, params)) + return self + + def fetchone(self): + sql = outer.log[-1][0] + if sql.startswith("SELECT run_id, work_id"): + return ("row", None, None, "running") + if sql.startswith("SELECT id, raw_content_id"): + return outer.raw_row + if sql.startswith("INSERT INTO example_run"): + return ("row",) + if sql.startswith("UPDATE example_run"): + state = "failed" if "failed" in (outer.log[-1][1] or ()) else "completed" + return ("row", state, None) + RecordingConnect.next_id += 1 + return (RecordingConnect.next_id,) + + def fetchall(self): + return [] + + def commit(self): + outer.log.append(("COMMIT", None)) + + def rollback(self): + outer.log.append(("ROLLBACK", None)) + + class _Ctx: + def __enter__(self): + return _Cursor() + + def __exit__(self, *exc): + return False + + return _Ctx() + + +def _candidate_body() -> str: + return "茧撕开舱门的一瞬,林深听见了深渊的回声。" * 8 + + +class SpecBuildTest(unittest.TestCase): + def test_spec_freezes_creative_input_and_read_tools(self): + context, _ = _valid_pair() + spec = build_writer_dispatch_spec( + context, target_chapter=3, human_instruction="让恐惧落在身体上", candidate_version=2, + ) + self.assertEqual(spec["role"], "writer") + self.assertEqual(spec["input"]["workId"], context["workId"]) + self.assertEqual(spec["input"]["targetChapter"], 3) + self.assertEqual(spec["input"]["candidateVersion"], 2) + self.assertIn("让恐惧落在身体上", spec["taskPrompt"]) + self.assertIn("creativeInput", spec["input"]) + # 探索白名单 = 工具 server 登记表(单一事实源,防漂移)。 + self.assertEqual(spec["toolAllowlist"], sorted(TOOL_REGISTRY)) + self.assertEqual( + spec["outputSchema"]["required"], ["candidateBody"]) + + def test_session_paths_are_stable_per_chapter(self): + sid_a, dir_a = writer_session_paths(12, 3) + sid_b, dir_b = writer_session_paths(12, 3) + sid_c, dir_c = writer_session_paths(12, 4) + self.assertEqual((sid_a, dir_a), (sid_b, dir_b)) + self.assertNotEqual(sid_a, sid_c) + self.assertNotEqual(dir_a, dir_c) + + +class DispatchRoundTripTest(unittest.TestCase): + def setUp(self): + self.tmp = pathlib.Path(tempfile.mkdtemp()) + self.context, _ = _valid_pair() + # 运行目录按派发运行号固定;测试前清理,避免审计根冲突。 + self.dispatch_run_id = f"{self.context['runId']}-writer-v1" + import shutil + + shutil.rmtree(DEFAULT_RUN_DIR_ROOT / self.dispatch_run_id, ignore_errors=True) + + def _run(self, lines, connect): + return run_writer_via_dispatch( + self.context, + candidate_version=1, + repo_root=PROJECT_ROOT, + provider="p", + model="claude-opus-test", + thinking="low", + human_instruction="", + spec_path=self.tmp / "task.json", + launcher=fake_launcher(lines), + connect_factory=connect, + ) + + def test_success_binds_envelope_and_evidence_ref(self): + final_text = json.dumps({"candidateBody": _candidate_body()}, ensure_ascii=False) + connect = BridgeConnect() + envelope, receipt, raw_ref = self._run(pi_stream_lines(final_text), connect) + # 身份、哈希、版本由桥绑定,不来自模型。 + self.assertEqual(envelope["runId"], self.context["runId"]) + self.assertEqual(envelope["candidateVersion"], 1) + self.assertTrue(envelope["candidateSha256"].startswith("sha256:")) + self.assertEqual(envelope["candidateBody"], _candidate_body()) + # 回执适配供生产账本消费;证据引用来自派发运行的调用账。 + self.assertEqual(receipt.requested_model_id, "p/claude-opus-test") + self.assertTrue(receipt.dispatch_run_id.startswith(self.context["runId"])) + self.assertEqual(raw_ref, (9001, 7777)) + # 任务包落盘可审计。 + saved_spec = json.loads((self.tmp / "task.json").read_text(encoding="utf-8")) + self.assertEqual(saved_spec["role"], "writer") + + def test_dispatch_failure_fails_closed(self): + connect = BridgeConnect() + with self.assertRaises(DispatchWriterError) as caught: + # 空事件流 -> 框架层失败关闭。 + self._run([], connect) + self.assertNotEqual(caught.exception.code, "") + + def test_missing_raw_reference_fails_closed(self): + final_text = json.dumps({"candidateBody": _candidate_body()}, ensure_ascii=False) + connect = BridgeConnect() + connect.raw_row = None + with self.assertRaises(DispatchWriterError) as caught: + self._run(pi_stream_lines(final_text), connect) + self.assertEqual(caught.exception.code, "DISPATCH_EVIDENCE_MISSING") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/skills/write-next-chapter/test_persist_writer_run.py b/tests/skills/write-next-chapter/test_persist_writer_run.py index cd89488..758b48d 100644 --- a/tests/skills/write-next-chapter/test_persist_writer_run.py +++ b/tests/skills/write-next-chapter/test_persist_writer_run.py @@ -24,3 +24,126 @@ class WriterPersistenceContractTest(unittest.TestCase): if __name__ == "__main__": unittest.main() + + +class _FakeCursor: + """按 SQL 关键词回放固定行的假游标(只验证落库合同,不连库)。""" + + def __init__(self, log): + self._log = log + self._last = None + self._ids = iter(range(100, 1000)) + + def execute(self, sql, params=None): + self._log.append((sql, params)) + self._last = sql + return self + + def fetchone(self): + sql = self._last or "" + if sql.startswith("SELECT id,candidate_sha256,state"): + return None + if "COALESCE(MAX(revision)" in sql: + return (0,) + if "RETURNING id" in sql: + return (next(self._ids),) + return (next(self._ids),) + + def commit(self): + self._log.append(("COMMIT", None)) + + def rollback(self): + self._log.append(("ROLLBACK", None)) + + +class _FakeConnCtx: + def __init__(self, log): + self._log = log + + def __enter__(self): + return _FakeCursor(self._log) + + def __exit__(self, *exc): + return False + + +def _minimal_inputs(): + context = { + "runId": "run-persist-dispatch-1", + "workId": 12, + "targetChapter": 3, + "attempt": 1, + "contextSnapshot": {"contextSha256": "sha256:" + "b" * 64}, + } + candidate = { + "runId": "run-persist-dispatch-1", + "attempt": 1, + "candidateVersion": 1, + "candidateSha256": "sha256:" + "c" * 64, + "candidateBody": "正文候选", + } + mechanical = {"passed": True, "requirements": []} + return context, candidate, mechanical + + +class DispatchRawRefTest(unittest.TestCase): + """阶段 E:派发链模型证据在派发运行下,编排显式传引用,禁止两套记账。""" + + def _patch_runtime(self): + import contextlib + + log = [] + + @contextlib.contextmanager + def _patches(): + orig_connect = writer_persist.connect + orig_start = writer_persist.start_run + orig_finish = writer_persist.finish_run + orig_lesson = writer_persist._propose_writer_lesson + writer_persist.connect = lambda *a, **k: _FakeConnCtx(log) + writer_persist.start_run = lambda *a, **k: None + writer_persist.finish_run = lambda *a, **k: None + writer_persist._propose_writer_lesson = lambda *a, **k: None + try: + yield log + finally: + writer_persist.connect = orig_connect + writer_persist.start_run = orig_start + writer_persist.finish_run = orig_finish + writer_persist._propose_writer_lesson = orig_lesson + + return _patches() + + def test显式引用绕过本运行查询并落回执(self): + context, candidate, mechanical = _minimal_inputs() + + def _forbid_raw_lookup(*args, **kwargs): + raise AssertionError("显式引用模式下不得查询本运行调用账") + + with self._patch_runtime() as log: + orig_raw = writer_persist._raw_response + writer_persist._raw_response = _forbid_raw_lookup + try: + result = writer_persist.persist_writer_execution( + context, candidate, receipt=None, mechanical_report=mechanical, + semantic_report=None, assemble_result=None, + writer_raw_ref=(5001, 6002), + ) + finally: + writer_persist._raw_response = orig_raw + self.assertEqual(result["status"], "persisted") + self.assertEqual(result["raw_content_id"], 6002) + # 候选与回执都落了库。 + inserts = [sql for sql, _ in log if sql.startswith("INSERT INTO")] + self.assertTrue(any("example_candidate" in sql for sql in inserts)) + self.assertTrue(any("example_run_receipt" in sql for sql in inserts)) + + def test引用缺项失败关闭(self): + context, candidate, mechanical = _minimal_inputs() + with self._patch_runtime(): + with self.assertRaises(writer_persist.WriterPersistenceError): + writer_persist.persist_writer_execution( + context, candidate, receipt=None, mechanical_report=mechanical, + semantic_report=None, assemble_result=None, + writer_raw_ref=(None, None), + )