#!/usr/bin/env python3
"""真写一章 · 阶段二:走完整生产链写下一章(默认写目标作品的下一章)。
生产链(meta/chains continuation 登记的保护节点序列落地):
读已确认细纲 + 前章正文基线
→ build_retrieval_plan + retrieve_writer_sources(生产仓储;新书无卡诚实返空)
→ assemble_context 冻结 WriterContext v1
→ run_writer_pipeline(持久 CAS 状态链 + 机械门 + 语义 detector,单次收敛 + 授权终态合同)
writer 经两阶段框架派发(探索取材 → 单次成稿):模型证据记派发运行,生产候选账本以
writer_raw_ref 显式关联,不双套记账;语义 detector 走冻结 detector profile 真调;
语义 needs_evidence 时走 production_evidence_reassemble(检索已有正典摘录);零命中当新设定交人闸,不禁写不重写
→ persist_writer_execution(冻结 + 运行注册 + 候选[含语义状态] + 回执 + 机械/语义质量证据一次落库)
→ accept_preflight(check_writer_acceptance 纯函数 + acceptance_state 实时重读)
→ 停止并展示候选,等待用户明确选择改 / 丢弃 / 采纳
用法:.venv/bin/python .agent/skills/write-next-chapter/scripts/produce_next_chapter.py [目标章号]
--provider
--model [--thinking T] [--instruction "本轮人指令原文"]
写手唯一执行形态是两阶段框架派发(2026-08-23 对照裁决:直调链退出创作生成);
--provider/--model 必须显式给出,不从环境变量推断。
缺省写下一章(库内最大章序 +1)。前置:目标章已建立、存在 confirmed 细纲,且门锚合同
GATE_ANCHORS 已登记该章。
--instruction:写入本轮输入第 7 项(人指令),经 styleConstraints 注入 writer;
与 05 §2.3 同序拼装合同对齐。正式采纳只由用户明确决定后调用 decide-candidate。
"""
import hashlib
import json
import sys
import uuid
from datetime import datetime, timezone
from decimal import Decimal
from pathlib import Path
from typing import Any, Mapping
SCRIPT_DIR = Path(__file__).resolve().parent
REPO_ROOT = SCRIPT_DIR.parents[3]
AGENT_ROOT = SCRIPT_DIR.parents[2]
SKILLS = AGENT_ROOT / "skills"
for sub in (
"assemble-context/scripts",
"write-next-chapter/scripts",
"record-run-evidence/scripts",
"check-content-consistency/scripts",
"decide-candidate/scripts",
"prevent-ai-flavor/scripts",
"diagnose-ai-flavor/scripts",
"replay-writer-gate/scripts",
):
p = str(SKILLS / sub)
if p not in sys.path:
sys.path.insert(0, p)
from muse_db import connect, DSN # noqa: E402
from assemble_writer_context import assemble_context # noqa: E402
from prevent_ai_flavor import ( # noqa: E402
PreventionContractError, build_prevention_contract, persist_prevention, render_writer_constraints,
)
from diagnose_ai_flavor import run_diagnosis, persist_diagnosis # noqa: E402
from retrieve_writer_sources import ( # noqa: E402
ProductionCardIndexRepository, FrozenProseRepository, RetrievalError,
build_retrieval_plan, retrieve_writer_sources, load_confirmed_fine_outline,
load_confirmed_pattern_bindings, load_confirmed_style,
)
from run_writer import ( # noqa: E402
build_production_length_contracts, build_writer_execution_profile,
calculate_dynamic_output_contract,
)
from run_writer_pipeline import PipelineError, run_writer_pipeline # noqa: E402
from production_evidence_reassemble import ( # noqa: E402
EvidenceReassembleError,
reassemble_writer_context_for_gaps,
)
from candidate_cas import PostgresCasStateStore # noqa: E402
from run_writer_semantic_detector import ( # noqa: E402
SEMANTIC_DETECTOR_REPORT_JSON_SCHEMA, build_safe_semantic_diagnostic,
build_semantic_input_v3, is_semantic_schema_specialization,
run_writer_semantic_detector,
)
from muse_role import ( # noqa: E402
RoleExecutionProfile,
run_role,
sha256_json,
)
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 two_phase_writer import ExplorationError, run_two_phase_writer # 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
WORK_ID = 12
GATE_A_CONFIG = json.loads((SKILLS / "replay-writer-gate" / "configs" /
"writer-gate-a-deep-space-v1.json").read_text(encoding="utf-8"))
GATE_A_WRITER = GATE_A_CONFIG["executionProfiles"]["writer"]
GATE_A_DETECTOR = GATE_A_CONFIG["executionProfiles"]["semantic_detector"]
PRODUCTION_LENGTH_PROMPT = (
"\n9. 篇幅是硬门:按输入 lengthContract 的汉字数口径,成稿必须达到 minChars;"
"输出前自行估算,若不足就继续展开场景。宁可接近 maxChars,也不得低于 minChars。"
"不要把标点、数字或拉丁字母计入汉字数。"
"\n10. 写手可以设计本章新出现的地名、能力、器物、感知或宇宙规则;"
"它们只是候选正文里的新设定,不是已确认正典,不要写成设定文档里早已成立的事实。"
"与已给出的事实约束冲突的内容不要写。新设定是否入库由人决定。"
)
SYSTEM_PROMPT = GATE_A_WRITER["systemPrompt"] + PRODUCTION_LENGTH_PROMPT
SYSTEM_PROMPT_ID = "writer-production-system-v4-new-settings"
SYSTEM_PROMPT_SHA256 = "sha256:" + hashlib.sha256(SYSTEM_PROMPT.encode("utf-8")).hexdigest()
ARTIFACTS = REPO_ROOT / "docs" / "write-chapter" / "artifacts"
# 门锚合同按章登记:锚点是章级创作判断,any-hit 子串匹配。新章必须先登记再跑。
# 机械门(check_writer_candidate)只认这些子串,不认语义等价;必须投影给写手,
# 否则写手只见细纲「完成第一次升级」却因禁抄长句而避开「升级」二字 → HARD_EVENT_MISSING。
GATE_ANCHORS = {
2: {
"requiredEvents": [
{"requirementId": "event-1-isolation",
"anchors": ["隔离", "收押", "禁闭", "关押", "封锁"]},
{"requirementId": "event-2-interrogation",
"anchors": ["审讯", "审问", "询问", "盘问", "讯问"]},
{"requirementId": "event-3-conceal",
"anchors": ["隐瞒", "没有告诉", "没说", "没有说", "咽了回去", "沉默", "闭上嘴"]},
{"requirementId": "event-4-hunger",
"anchors": ["饥饿", "渴望", "吞噬", "进食", "吃", "贪"]},
],
"requiredCharacters": ["林深", "何岚"],
"foreshadowingActions": [
{"requirementId": "foreshadow-upgrade",
"anchors": ["异种核心", "融合", "升级"]},
],
"chapterEndHook": {
"requirementId": "hook-ch2",
"anchors": ["调令", "实战", "出击", "部署", "任务", "出征", "离不开", "不愿离开"],
"maxDistanceFromEnd": 900,
},
},
3: {
# 接第2章结尾硬钩子(茧撕开舱门出击、要吃掉更强核心)。锚点 any-hit 子串匹配;
# requiredCharacters 只硬约束主角(避免过度约束触发 costly 重抽),其余靠事件锚点。
"requiredEvents": [
{"requirementId": "event-1-sortie",
"anchors": ["出击", "实战", "战斗", "交火", "搏杀", "拦截", "扑向", "战场"]},
{"requirementId": "event-2-devour",
"anchors": ["吞噬", "吞食", "吃掉", "进食", "撕碎", "吸收", "吞下", "吞"]},
{"requirementId": "event-3-upgrade",
"anchors": ["升级", "蜕变", "进化", "变强", "增强", "新的力量", "蜕变"]},
{"requirementId": "event-4-pollution",
"anchors": ["黑纹", "污染", "扩散", "蔓延", "加深", "恶化"]},
],
"requiredCharacters": ["林深"],
"foreshadowingActions": [
{"requirementId": "foreshadow-voice-merge",
"anchors": ["分不清", "像他自己", "脑内的声音", "低语", "渴望", "哪个念头", "另一个"]},
],
"chapterEndHook": {
"requirementId": "hook-ch3",
"anchors": ["深渊", "更深", "回应", "召唤", "更大", "下一", "不止", "饥饿", "注视", "凝视"],
"maxDistanceFromEnd": 900,
},
},
}
def format_mechanical_gate_constraints(requirements: Mapping[str, Any]) -> list[str]:
"""把 GATE_ANCHORS 投影成写手可见约束(机械门规则对写手透明)。
规则(check_writer_candidate._anchors_present):每个 requirement 的 anchors
列表里**任意一个子串**出现在正文即通过;缺则 HARD_EVENT_MISSING 等,且
机械失败不进补证环、直接拒绝。语义「写了升级这件事」不够,必须命中子串。
"""
lines: list[str] = [
"机械门验收标记(确定性子串,非细纲长句):下列每组至少在正文自然出现其中一个词;"
"这是验收标记,允许写进戏剧化句子,禁止整句复读硬约束长句。",
]
for item in requirements.get("requiredEvents") or []:
if not isinstance(item, Mapping):
continue
rid = item.get("requirementId") or "event"
anchors = [a for a in (item.get("anchors") or []) if isinstance(a, str) and a]
if anchors:
lines.append(f"硬事件[{rid}] 须含其一:{' / '.join(anchors)}")
chars = [c for c in (requirements.get("requiredCharacters") or []) if isinstance(c, str) and c]
if chars:
lines.append(f"必须出场角色(全文须出现姓名):{'、'.join(chars)}")
for item in requirements.get("foreshadowingActions") or []:
if not isinstance(item, Mapping):
continue
rid = item.get("requirementId") or "foreshadow"
anchors = [a for a in (item.get("anchors") or []) if isinstance(a, str) and a]
if anchors:
lines.append(f"伏笔动作[{rid}] 须含其一:{' / '.join(anchors)}")
hook = requirements.get("chapterEndHook")
if isinstance(hook, Mapping):
rid = hook.get("requirementId") or "hook"
anchors = [a for a in (hook.get("anchors") or []) if isinstance(a, str) and a]
dist = hook.get("maxDistanceFromEnd")
if anchors:
lines.append(
f"章末钩子[{rid}] 须在结尾约 {dist} 字内含其一:{' / '.join(anchors)}"
)
return lines
def build_semantic_detector_profile() -> RoleExecutionProfile:
"""按正式配置重建语义 detector 的冻结角色 profile。"""
return profile_from_mapping(GATE_A_DETECTOR, role="semantic_detector")
class ProductionSemanticRunner:
"""语义 detector 生产 runner:角色调用随 run_id 统一落库。"""
def __init__(self, profile: RoleExecutionProfile, *, run_id: str) -> None:
self.profile = profile
self.run_id = run_id
def run(self, *, adapter_role: str, model_input: Mapping[str, Any],
output_schema: Mapping[str, Any]) -> Mapping[str, Any]:
if self.profile.adapter_role != adapter_role:
raise PipelineError("SEMANTIC_DETECTOR_FAILED", "RoleExecutionProfile 与语义 detector 适配不一致")
if not is_semantic_schema_specialization(self.profile.json_schema, output_schema):
raise PipelineError(
"SEMANTIC_DETECTOR_FAILED",
"output_schema 必须是基座 schema 的闭集特化",
)
from dataclasses import replace
call_profile = replace(
self.profile,
json_schema=dict(output_schema),
json_schema_sha256=sha256_json(output_schema),
)
def persist_detector_event(event):
# runtime 默认把非 writer 调用标为 evaluation;生产链语义审查改标 production_detection。
event = dict(event)
event["purpose"] = "production_detection"
return persist_llm_event(event)
result = run_role(
call_profile,
model_input,
run_id=self.run_id,
caller="semantic_detector",
persist_call=persist_detector_event,
)
receipt = result.receipt
receipt_dict = receipt.as_dict() if hasattr(receipt, "as_dict") else receipt
return {"structuredOutput": dict(result.structured_output),
"modelReceiptSha256": sha256_json(receipt_dict)}
def resolve_target_chapter(requested: int | None) -> int:
with connect(readonly=True) as conn:
if requested is not None:
return requested
row = conn.execute(
"SELECT COALESCE(MAX(order_no),0)+1 FROM muse_content_chapter "
"WHERE work_id=%s AND deleted=false", (WORK_ID,)).fetchone()
return int(row[0])
def _dump(path: Path, value: Any) -> None:
path.write_text(json.dumps(value, ensure_ascii=False, indent=1), encoding="utf-8")
def main():
argv = sys.argv[1:]
if "--dry-run" in argv:
raise SystemExit("--dry-run 已移除:生成入口只落 Shadow,不执行 accept")
human_instruction = ""
if "--instruction" in argv:
idx = argv.index("--instruction")
if idx + 1 >= len(argv) or argv[idx + 1].startswith("--"):
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
# 写手唯一执行形态是两阶段框架派发(2026-08-23 对照裁决);模型参数必须显式传入。
dispatch_provider = _take("--provider")
dispatch_model = _take("--model")
dispatch_thinking = _take("--thinking")
continue_from = _take("--continue-from")
if not dispatch_provider or not dispatch_model:
raise SystemExit("生产写手必须显式给出 --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
run_id = f"run-prod-work12-ch{target}-{uuid.uuid4().hex[:8]}"
generated_at = datetime.now(timezone.utc).isoformat()
ARTIFACTS.mkdir(exist_ok=True)
if target not in GATE_ANCHORS:
raise SystemExit(f"第{target}章门锚合同未登记(GATE_ANCHORS),先登记锚点再跑。")
# 1) 已确认细纲(read-context 统一消费点)+ 前章全文基线(asOf 起连续四章;不足四章从第1章起)
with connect(readonly=True) as conn:
try:
fine_outline = load_confirmed_fine_outline(conn, work_id=WORK_ID, target_chapter=target)
except RetrievalError as exc:
raise SystemExit(f"{exc}(目标章必须先建立并确认细纲)")
first = max(1, as_of - 3)
recent_rows = conn.execute(
"SELECT c.order_no, b.id, b.revision, b.content_text FROM muse_content_chapter c "
"JOIN muse_content_block b ON b.chapter_id=c.id AND b.deleted=false "
"WHERE c.work_id=%s AND c.deleted=false AND c.order_no BETWEEN %s AND %s "
"ORDER BY c.order_no", (WORK_ID, first, as_of)).fetchall()
style_constraints = load_confirmed_style(conn, work_id=WORK_ID)
pattern_references = load_confirmed_pattern_bindings(conn, work_id=WORK_ID)
# 第 7 项人指令:与文风并列注入 styleConstraints,看板/raw 可回看
if human_instruction:
style_constraints = list(style_constraints) + [f"本轮人指令:{human_instruction}"]
print(f"本轮人指令已注入({len(human_instruction)} 字)")
# 机械门锚点必须对写手可见(根因修复:#120 写了升级语义但未命中子串)
gate_lines = format_mechanical_gate_constraints(GATE_ANCHORS[target])
style_constraints = list(style_constraints) + gate_lines
print(f"机械门锚点已注入写手约束({len(gate_lines)} 条)")
recent_chapters = [{
"chapter": order_no,
"sourceRef": {"sourceId": f"content-block:{block_id}",
"sourceVersion": f"rev{revision}",
"blockId": int(block_id), "chapter": order_no,
"startCodePoint": 0, "endCodePoint": len(body)},
"text": body,
} for order_no, block_id, revision, body in recent_rows]
expected = list(range(first, as_of + 1))
got = [item["chapter"] for item in recent_chapters]
if got != expected:
raise SystemExit(f"连续前章基线缺章:期望 {expected},实际 {got}。")
print(f"目标第{target}章(asOf={as_of});细纲已读,基线 {got},"
f"基线总字数 {sum(len(item['text']) for item in recent_chapters)}")
# 1.5) 人感前置预防:规则/声音账先形成合同,再冻结进 WriterContext。
try:
humanization_contract = build_prevention_contract(
f"work:{WORK_ID}", load_database=True
)
humanization_contract["writer_constraints"] = render_writer_constraints(humanization_contract)
prevention_receipt = persist_prevention(humanization_contract)
except PreventionContractError as exc:
raise SystemExit(f"人感前置预防合同失败:{exc}") from exc
print(f"人感前置预防:约束 {len(humanization_contract['writer_constraints'])} 条,"
f"规则库 {humanization_contract['built_from']['rule_library_version']},"
f"run_id={prevention_receipt['run_id']}")
# 2) 检索计划 + 执行(生产仓储;新书无卡诚实返空)
token_budget = {"maxContextChars": 200000}
plan = build_retrieval_plan(
run_id=run_id, work_id=WORK_ID, target_chapter=target, as_of=as_of,
fine_outline=fine_outline, card_index_version="knowledge-index-v1",
prose_index_version="content-block-v1", token_budget=token_budget)
retrieval = retrieve_writer_sources(
plan=plan, card_repository=ProductionCardIndexRepository(),
prose_repository=FrozenProseRepository(dsn=DSN, tenant_id=0))
_dump(ARTIFACTS / f"{run_id}-retrieval.json",
{"plan": plan, "resultCounts": {k: len(v) for k, v in retrieval.items()
if isinstance(v, list)}})
print(f"检索: 卡={len(retrieval['cards'])} 事实={len(retrieval['factEvidence'])} "
f"原文={len(retrieval['proseEvidence'])}")
# 3) 组装冻结 WriterContext v1
dynamic_output_contract = calculate_dynamic_output_contract(
fine_outline=fine_outline,
recent_chapter_bodies=[item["text"] for item in recent_chapters])
output_contract, generation_length_contract = build_production_length_contracts(
dynamic_output_contract)
narrative_state = {
"time": f"第{as_of}章结束后",
"location": "承接上一章结尾的场景",
"characterPositions": {},
"immediateSituation": "按细纲 chapterGoal 展开(上一章结尾状态见基线正文)。",
}
authorization_snapshot = {"snapshotId": "auth-work12-production-v1",
"allowedPurpose": "production_generation",
"verifiedAt": generated_at, "sourceVersion": "v1"}
assembled = assemble_context(
run_id=run_id, attempt=1, mode="production", purpose="production",
quality_policy_version="writer-production-v1", work_id=WORK_ID,
target_chapter=target, as_of=as_of, source_version="outline@v1",
authorization_snapshot=authorization_snapshot, source_status="active",
retrieval_plan=plan, retrieval_result=retrieval, fine_outline=fine_outline,
narrative_state=narrative_state, recent_chapters=recent_chapters,
output_contract=output_contract, token_budget=token_budget,
generation_length_contract=generation_length_contract,
pattern_references=pattern_references, style_constraints=style_constraints,
humanization_contract=humanization_contract, generated_at=generated_at,
evidence_strategy="production_dual_evidence")
writer_context = assembled["context"]
(ARTIFACTS / f"{run_id}-writer-context.json").write_text(
assembled["contextJson"], encoding="utf-8")
print(f"上下文冻结: contextSha256={writer_context['contextSnapshot']['contextSha256'][:24]}... "
f"写手篇幅 {generation_length_contract['minChars']}-{generation_length_contract['maxChars']}"
f"(目标 {generation_length_contract['targetChars']}),"
f"机械接受 >3000(系统上限 {output_contract['maxChars']}),"
f"范式绑定 {len(pattern_references)} 张,"
f"文风约束 {len(style_constraints)} 条")
# 4) 生产 pipeline:持久 CAS + 两阶段写手派发 + 机械门 + 语义 detector(先审后入)
detector_profile = build_semantic_detector_profile()
semantic_runner = ProductionSemanticRunner(detector_profile, run_id=run_id)
state_store = PostgresCasStateStore(
work_id=WORK_ID, target_chapter=target, creator="continuation")
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]] = {}
explorations_by_version: dict[int, dict] = {}
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 适配:两阶段框架派发是唯一形态,探索取材后单次成稿,失败关闭。"""
contexts_by_attempt[current_context["attempt"]] = dict(current_context)
try:
candidate, receipt, raw_ref, exploration = run_two_phase_writer(
current_context,
candidate_version=candidate_version,
repo_root=REPO_ROOT,
provider=dispatch_provider,
model=dispatch_model,
thinking=dispatch_thinking,
human_instruction=human_instruction,
spec_dir=ARTIFACTS,
)
except ExplorationError as exc:
raise PipelineError(exc.code, f"两阶段写手失败: {exc}",
details=exc.details) from exc
explorations_by_version[candidate_version] = exploration
_dump(ARTIFACTS / f"{run_id}-exploration-summary-v{candidate_version}.json",
exploration)
print(f"两阶段写手: exploration_run={exploration['explorationRunId']} "
f"材料={exploration['materialCount']} 生成_run={exploration['generationRunId']}")
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
return candidate
def production_semantic_detector(current_context, candidate, mechanical_report):
"""语义 detector 适配:构造冻结输入、真调模型、留档输入输出。"""
version = candidate["candidateVersion"]
detector_input = build_semantic_input_v3(
run_id=run_id, sample_id=f"writer-ch{target}", opaque_arm_id="production",
writer_context=current_context, candidate=candidate)
_dump(ARTIFACTS / f"{run_id}-semantic-input-v{version}.json", detector_input)
outcome = run_writer_semantic_detector(detector_input, model_runner=semantic_runner)
_dump(ARTIFACTS / f"{run_id}-semantic-output-v{version}.json", outcome)
if outcome.get("ok") is not True or not isinstance(outcome.get("report"), Mapping):
diagnostic = build_safe_semantic_diagnostic(outcome)
raise PipelineError(
"SEMANTIC_DETECTOR_FAILED",
f"语义 detector 未产生有效报告: {diagnostic['primaryCode']}",
details={"safeDiagnostic": diagnostic})
print(f"语义 detector(v{version}): status={outcome['status']} "
f"调用={outcome['attemptCount']}次 纠错={outcome['correctionCount']}次")
return outcome["report"]
def production_evidence_provider(context, gaps, attempt):
"""语义缺口 → 正典检索;命中则注入摘录。零命中由 pipeline 当新设定交人闸。"""
try:
next_ctx = reassemble_writer_context_for_gaps(
context, gaps, attempt, work_id=WORK_ID)
except EvidenceReassembleError as exc:
raise PipelineError(
"PRODUCTION_EVIDENCE_REASSEMBLE_FAILED",
f"补证重组装失败(缺口 {len(gaps)}): {exc}",
) from exc
print(
f"补证重组装: attempt={next_ctx['attempt']} "
f"facts={len(next_ctx.get('factEvidence') or [])} "
f"gaps={len(gaps)}"
)
return next_ctx
# 候选版本接续:候选表对 (作品,章,candidate_version) 唯一,重跑同章必须从已有最大版本+1 起,
# 否则与上一轮留库的被拒候选撞版本。
with connect(readonly=True) as conn:
max_version_row = conn.execute(
"SELECT COALESCE(MAX(CASE WHEN candidate_version ~ '^[0-9]+$' "
"THEN candidate_version::integer END),0) FROM example_candidate "
"WHERE tenant_id=0 AND work_id=%s AND target_chapter=%s AND deleted=false",
(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"],
**({"continueFrom": continue_from} if continue_from else {}),
"writerMode": "two-phase"},
creator="continuation")
try:
pipeline_result = run_writer_pipeline(
context=writer_context,
requirements=GATE_ANCHORS[target],
writer=production_writer,
evidence_provider=production_evidence_provider,
semantic_detector=production_semantic_detector,
state_store=state_store,
result_path=ARTIFACTS / f"{run_id}-pipeline-result.json",
initial_candidate_version=initial_candidate_version,
)
except PipelineError as exc:
# 被拒候选留痕:凡跑过机械门的版本都落 Shadow(state=rejected + 机械/语义证据)
trace = (exc.result or {}).get("trace") or []
audit_entry = next((entry for entry in reversed(trace)
if isinstance(entry.get("mechanicalReport"), Mapping)), None)
if audit_entry is not None:
version = audit_entry.get("candidateVersion")
failed_candidate = candidates_by_version.get(version)
failed_receipt = receipts_by_version.get(version)
if failed_candidate is not None and failed_receipt is not None:
try:
persisted = persist_writer_execution(
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,
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: # 留痕失败不掩盖原始失败码
print(f"[警告] 被拒候选留库失败: {persist_exc}", file=sys.stderr)
finish_run(run_id, "failed", creator="continuation",
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')})")
print("授权后由主代理发起新运行继续补证(本运行已收敛 REJECTED,新运行接续候选版本)。")
print(f"[停止] 生产 pipeline 未通过: code={exc.code};{exc}")
print(f"复核 artifacts/{run_id}-pipeline-result.json 后决定下一步。")
print(f"\nRUN_ID={run_id}")
raise SystemExit(1)
except Exception as exc:
# 非 PipelineError(库连接断、适配层异常等)也要收口运行态,不留 running 悬挂
try:
finish_run(run_id, "failed", creator="continuation",
trigger_detail={"stage": "writer-production-pipeline",
"error": type(exc).__name__})
except Exception:
pass
raise
# 5) pipeline 通过:机械门 + 语义 detector 双证据落库(候选 semantic_status=passed)
candidate = pipeline_result["candidateArtifact"]
final_context = contexts_by_attempt.get(pipeline_result["attempt"], writer_context)
final_trace = pipeline_result["trace"][-1]
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,
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}")
print(f"落库: candidate_id={cand_id}, receipt_id={persisted['receipt_id']}, "
f"raw_content_id={persisted['raw_content_id']}, state={persisted['state']}, "
f"semantic={persisted['semantic_status']}")
# 6) 接受前置检查:实时状态重读 + 纯函数全检(上下文/授权/来源/有效期/detector 终态)
try:
live_state = build_live_acceptance_state(final_context)
preflight = check_writer_acceptance(
decision="accept", confirmed=True, context=final_context, candidate=candidate,
detector_result=pipeline_result, live_state=live_state,
expected_revision=live_state["canonicalRevision"])
except (LiveStateError, AcceptanceError) as exc:
finish_run(run_id, "failed", creator="continuation",
trigger_detail={"stage": "accept-preflight",
"error": getattr(exc, "code", type(exc).__name__)})
print(f"[停止] 接受前置检查未通过: {getattr(exc, 'code', '')} {exc}")
print(f"候选 {cand_id} 已留库(state=passed, semantic=passed),人工复核后决定。")
print(f"\nCANDIDATE_ID={cand_id}\nRUN_ID={run_id}")
raise SystemExit(1)
print(f"接受前置检查: {preflight['status']} canonicalRevision={live_state['canonicalRevision']}")
# 7) 人闸:生成入口到此停止,正式正文只能由用户明确决定后走 decide-candidate。
print(f"候选 {cand_id} 已通过机械门、语义门与接受前置检查,尚未写入正式正文。")
print(f"候选详情: http://127.0.0.1:8765/candidates/{cand_id}")
print("请决定:改:<具体要求> / 丢弃 / 采纳")
print(f"\nCANDIDATE_ID={cand_id}\nRUN_ID={run_id}")
if __name__ == "__main__":
main()