zizi fa922f8cc5 实现: 装配正文五章真实冻结评测
将卡稳定选择器、冻结历史正文和目标 scaffold 绑定到同一只读快照;固定上下文预算并净化诊断索引,收紧盲评输入边界,防止目标事实和选卡信息泄漏。
2026-07-21 16:09:23 +08:00

1129 lines
47 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""正文 A/B/C 同条件冻结回放编排器。
默认 dry-run 只组装并校验 WriterContext,公开计划、清单和脱敏摘要,
不调用模型。真实 semantic detector 与盲评 adapter 尚未接线,因此 CLI
的 ``--execute`` 明确失败关闭;测试只能通过专用注入接口替换 subprocess
runner,并且候选仍必须经过正式 writer adapter 和完整审查管线。
"""
from __future__ import annotations
import argparse
import copy
import json
import random
import re
import subprocess
import sys
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable, Mapping, Protocol
SCRIPT_DIR = Path(__file__).resolve().parent
SKILLS_DIR = SCRIPT_DIR.parents[1]
CONTINUATION_DIR = SKILLS_DIR / "continuation" / "scripts"
READ_CONTEXT_DIR = SKILLS_DIR / "read-context" / "scripts"
for import_path in (CONTINUATION_DIR, READ_CONTEXT_DIR):
if str(import_path) not in sys.path:
sys.path.insert(0, str(import_path))
from assemble_writer_context import AssemblyError, assemble_context # noqa: E402
from audit_leakage import audit_snapshot # noqa: E402
from build_snapshot import SnapshotError, build_snapshot, normalize_chapter, sha256_value # noqa: E402
from check_snapshot import ( # noqa: E402
STATUS_READY,
check_arm_manifests,
check_authorization,
check_target_sources,
)
from run_writer import WriterAdapterError, run_writer # noqa: E402
from run_writer_pipeline import ( # noqa: E402
InMemoryCasStateStore,
PipelineError,
atomic_write_json,
run_writer_pipeline,
)
from writer_contract import ( # noqa: E402
ContractError,
calculate_target_chars,
retrieval_identity,
validate_writer_context,
)
from writer_rubric import adjudicate_reviews, deblind_reports # noqa: E402
PROFILE = "writer_replay"
REQUIRED_ARMS = ("A", "B", "C")
EVIDENCE_STRATEGIES = {
"A": "historical_prose_only",
"B": "card_index_only",
"C": "card_index_plus_prose",
}
PRIVATE_TMP = Path("/private/tmp").resolve()
SNAPSHOT_COMPAT_ARMS = (
"outline_only",
"outline_plus_cards",
"outline_plus_placebo_cards",
)
CONTENT_MODES = frozenset({"sanitized_contract_fixture", "canonical_frozen_prose"})
# Judge 只能复用写手细纲合同中的四个字段,不能直接信任原始评测配置。
FINE_OUTLINE_JUDGE_FIELDS = (
"sourceRef",
"hardConstraints",
"adjustableBeats",
"declaredNewFacts",
)
class WriterReplayError(ValueError):
"""正文回放配置或可信边界不合法。"""
class AuthorizationBindingError(WriterReplayError):
"""WriterContext 输入试图脱离顶层不可变授权快照。"""
class ReplayJudge(Protocol):
"""测试盲评边界;只接收盲化内容,生产 judge adapter 未接线。"""
def __call__(
self,
blind_input: Mapping[str, Any],
reviewer_id: str,
) -> Mapping[str, Any]: ...
@dataclass(frozen=True)
class WriterReplayTestAdapters:
"""仅供无模型测试注入的三类边界,不能由 CLI 构造。"""
writer_runner: Callable[..., subprocess.CompletedProcess[str]]
semantic_detector: Callable[
[Mapping[str, Any], Mapping[str, Any], Mapping[str, Any]], Mapping[str, Any]
]
judge: ReplayJudge
def _safe_json(value: Any) -> str:
"""生成稳定 JSON,便于哈希与机械比较。"""
return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
def _mapping(value: Any, field: str) -> Mapping[str, Any]:
"""读取必填对象字段并失败关闭。"""
if not isinstance(value, Mapping):
raise WriterReplayError(f"{field} 必须是对象")
return value
def _sequence(value: Any, field: str) -> list[Any]:
"""读取必填数组字段,不猜测字符串或单对象。"""
if not isinstance(value, list):
raise WriterReplayError(f"{field} 必须是数组")
return list(value)
def _run_id(value: str) -> str:
"""把外部 run/sample 标识收敛到 WriterContext 允许的稳定字符集。"""
normalized = re.sub(r"[^A-Za-z0-9._:-]+", "-", value).strip("-._:")
if not normalized:
raise WriterReplayError("runId/sampleId 不能规范化为空")
return normalized[:128]
def _length_bounds(target_chars: int) -> tuple[int, int]:
"""按预注册目标的正负 10% 计算边界,并执行正文合同的总限幅。"""
return max(2000, target_chars * 9 // 10), min(10000, target_chars * 11 // 10)
def _validate_preregistered_length(sample: Mapping[str, Any], as_of: int) -> int:
"""只用冻结点前四章汉字数机械复算篇幅,禁止读取目标章实际长度。"""
counts = _sequence(sample.get("frozenRecentHanCounts"), "frozenRecentHanCounts")
if (
len(counts) != 4
or any(isinstance(value, bool) or not isinstance(value, int) or value < 500 for value in counts)
):
raise WriterReplayError("frozenRecentHanCounts 必须是冻结点前四章的有效汉字数")
basis = _mapping(sample.get("targetLengthBasis"), "targetLengthBasis")
expected_fields = {
"algorithm",
"sourceChapters",
"hardEventCount",
"foreshadowingActionCount",
"requiredSceneCount",
"minChars",
"maxChars",
"usesTargetChapterLength",
}
if set(basis) != expected_fields:
raise WriterReplayError("targetLengthBasis 字段必须严格匹配预注册篇幅合同")
if basis.get("algorithm") != "calculate_target_chars":
raise WriterReplayError("targetLengthBasis.algorithm 非法")
source_chapters = _sequence(basis.get("sourceChapters"), "targetLengthBasis.sourceChapters")
if source_chapters != list(range(max(1, as_of - 3), as_of + 1)):
raise WriterReplayError("篇幅基线必须严格来自冻结点前连续四章")
if basis.get("usesTargetChapterLength") is not False:
raise WriterReplayError("篇幅预注册禁止读取目标章实际长度")
integer_fields = (
"hardEventCount",
"foreshadowingActionCount",
"requiredSceneCount",
"minChars",
"maxChars",
)
if any(isinstance(basis.get(field), bool) or not isinstance(basis.get(field), int) for field in integer_fields):
raise WriterReplayError("targetLengthBasis 计数与边界必须是整数")
try:
target = calculate_target_chars(
recent_chapter_han_counts=counts,
hard_event_count=int(basis["hardEventCount"]),
foreshadowing_action_count=int(basis["foreshadowingActionCount"]),
required_scene_count=int(basis["requiredSceneCount"]),
min_chars=int(basis["minChars"]),
max_chars=int(basis["maxChars"]),
)
except ContractError as error:
raise WriterReplayError(f"篇幅预注册参数非法: {error}") from error
if sample.get("targetChars") != target:
raise WriterReplayError("targetChars 与冻结历史机械复算结果不一致")
expected_min, expected_max = _length_bounds(target)
expected_length = _mapping(sample.get("expectedLength"), "expectedLength")
if dict(expected_length) != {
"targetChars": target,
"minChars": expected_min,
"maxChars": expected_max,
}:
raise WriterReplayError("expectedLength 必须等于机械目标的正负 10% 限幅")
context_input = _mapping(sample.get("writerContextInput"), "writerContextInput")
output_contract = _mapping(context_input.get("outputContract"), "writerContextInput.outputContract")
for field, expected in (
("targetChars", target),
("minChars", expected_min),
("maxChars", expected_max),
):
if output_contract.get(field) != expected:
raise WriterReplayError(f"writerContextInput.outputContract.{field} 未绑定预注册篇幅")
return target
def _validate_new_character_ratio(sample: Mapping[str, Any], as_of: int) -> None:
"""校验具名必需角色在冻结点前无记录的比例;泛称角色必须保持未知。"""
context_input = _mapping(sample.get("writerContextInput"), "writerContextInput")
requirements = _mapping(context_input.get("requirements"), "writerContextInput.requirements")
required = _sequence(
requirements.get("requiredCharacters"),
"writerContextInput.requirements.requiredCharacters",
)
if not required or any(not isinstance(item, str) or not item.strip() for item in required):
raise WriterReplayError("requiredCharacters 必须是非空字符串数组")
status = sample.get("newCharacterRatioStatus")
basis = _mapping(sample.get("newCharacterBasis"), "newCharacterBasis")
expected_fields = {
"definition",
"asOfChapter",
"requiredCharacters",
"knownBeforeAsOf",
"absentBeforeAsOf",
"genericRoles",
}
if set(basis) != expected_fields:
raise WriterReplayError("newCharacterBasis 字段必须严格匹配预注册合同")
if basis.get("definition") != "named_required_characters_absent_before_as_of_ratio":
raise WriterReplayError("newCharacterBasis.definition 非法")
if basis.get("asOfChapter") != as_of:
raise WriterReplayError("newCharacterBasis.asOfChapter 未绑定冻结点")
if _sequence(basis.get("requiredCharacters"), "newCharacterBasis.requiredCharacters") != required:
raise WriterReplayError("newCharacterBasis.requiredCharacters 与细纲要求不一致")
known = _sequence(basis.get("knownBeforeAsOf"), "newCharacterBasis.knownBeforeAsOf")
absent = _sequence(basis.get("absentBeforeAsOf"), "newCharacterBasis.absentBeforeAsOf")
generic = _sequence(basis.get("genericRoles"), "newCharacterBasis.genericRoles")
for field, values in (("knownBeforeAsOf", known), ("absentBeforeAsOf", absent), ("genericRoles", generic)):
if any(not isinstance(item, str) or not item.strip() for item in values) or len(set(values)) != len(values):
raise WriterReplayError(f"newCharacterBasis.{field} 必须是无重复非空字符串数组")
if status == "unresolved_generic_role":
if sample.get("newCharacterRatio") is not None or not generic:
raise WriterReplayError("泛称角色无法判定时 newCharacterRatio 必须为 null")
if set(generic) != set(required) or known or absent:
raise WriterReplayError("泛称角色不得被默认归为已知或新角色")
return
if status != "resolved" or generic:
raise WriterReplayError("newCharacterRatioStatus 非法")
if set(known).intersection(absent) or set(known).union(absent) != set(required):
raise WriterReplayError("具名 requiredCharacters 必须被已知/无记录集合完整且互斥地覆盖")
retrieval_result = _mapping(context_input.get("retrievalResult"), "writerContextInput.retrievalResult")
raw_hints = retrieval_result.get("indexHints")
if raw_hints is None:
raw_hints = retrieval_result.get("unverifiedIndexHints", [])
hints = _sequence(raw_hints, "writerContextInput.retrievalResult.indexHints")
hinted_names = {
str(hint.get("name") or "")
for hint in hints
if isinstance(hint, Mapping)
}
if set(absent).intersection(hinted_names):
raise WriterReplayError("冻结点前无记录角色不得同时出现在卡索引提示中")
expected_ratio = len(absent) / len(required)
ratio = sample.get("newCharacterRatio")
if isinstance(ratio, bool) or not isinstance(ratio, (int, float)) or float(ratio) != expected_ratio:
raise WriterReplayError("newCharacterRatio 与冻结点前记录分区不一致")
def _common_controls(
config: Mapping[str, Any], sample: Mapping[str, Any]
) -> dict[str, Any]:
"""冻结三臂共享的作品、输入、模型、采样和 detector。"""
common = dict(_mapping(config.get("commonControls"), "commonControls"))
as_of = normalize_chapter(sample.get("asOfChapter"))
target = normalize_chapter(sample.get("targetChapter"))
if as_of is None or target != as_of + 1:
raise WriterReplayError("样本 targetChapter 必须等于 asOfChapter+1")
_validate_preregistered_length(sample, as_of)
_validate_new_character_ratio(sample, as_of)
for field in (
"outlineSource",
"fineOutlineSource",
"targetChars",
"modelVersion",
"sampling",
"detectorProfile",
):
value = sample.get(field, common.get(field))
if value is None:
raise WriterReplayError(f"样本缺少公共控制字段: {field}")
common[field] = value
reference = _mapping(config.get("referenceWork"), "referenceWork")
common.update(
{
"workId": reference.get("id"),
"referenceWorkVersion": reference.get("version"),
"asOfChapter": as_of,
"targetChapter": target,
}
)
if isinstance(common["workId"], bool) or not isinstance(common["workId"], int):
raise WriterReplayError("referenceWork.id 必须是正整数")
if common["workId"] <= 0:
raise WriterReplayError("referenceWork.id 必须是正整数")
if sample.get("workId", common["workId"]) != common["workId"]:
raise WriterReplayError("样本 workId 与 referenceWork 不一致")
return common
def _retrieval_plan(
*, run_id: str, common: Mapping[str, Any], context_input: Mapping[str, Any]
) -> dict[str, Any]:
"""为三臂生成同一份确定性检索计划;证据差异只进入 manifest。"""
token_budget = _mapping(context_input.get("tokenBudget"), "writerContextInput.tokenBudget")
plan = {
"planVersion": "writer-retrieval-plan-v1",
"runId": run_id,
"asOf": common["asOfChapter"],
"queries": copy.deepcopy(_sequence(context_input.get("retrievalQueries", []), "retrievalQueries")),
"cardIndexVersion": str(context_input.get("cardIndexVersion") or "writer-replay-cards-v1"),
"proseIndexVersion": str(context_input.get("proseIndexVersion") or "writer-replay-prose-v1"),
"filters": {
"workId": common["workId"],
"asOfChapter": common["asOfChapter"],
"sourceStatus": str(context_input.get("sourceStatus") or ""),
"authorizationRequired": True,
},
"tieBreak": "score DESC, sourceVersion ASC, sourceId ASC, sourceOffset ASC",
"tokenBudget": {"maxContextChars": token_budget.get("maxContextChars")},
}
plan["planId"] = retrieval_identity(plan)
return plan
def _requirements(context_input: Mapping[str, Any]) -> dict[str, Any]:
"""读取机械 detector 的完整输入,dry-run 也先校验结构。"""
raw = dict(_mapping(context_input.get("requirements"), "writerContextInput.requirements"))
for field in (
"requiredEvents",
"requiredCharacters",
"foreshadowingActions",
"detectedNewSettings",
"frozenConflicts",
):
_sequence(raw.get(field), f"writerContextInput.requirements.{field}")
_mapping(raw.get("chapterEndHook"), "writerContextInput.requirements.chapterEndHook")
return raw
def _validate_context_authorization_binding(
config: Mapping[str, Any], context_input: Mapping[str, Any]
) -> None:
"""把 WriterContext 来源、版本、用途和校验时间绑定到顶层授权快照。
顶层 ``offline_evaluation`` 在 WriterContext 授权快照中允许保留原值,
或规范化为合同目的 ``evaluation``。顶层来源状态是授权登记状态,WriterContext 的
``active`` 是其运行期投影;只有下表明确列出的投影允许通过。
"""
reference = _mapping(config.get("referenceWork"), "referenceWork")
authorization = _mapping(config.get("authorization"), "authorization")
top_snapshot = _mapping(
authorization.get("authorizationSnapshot"),
"authorization.authorizationSnapshot",
)
context_snapshot = _mapping(
context_input.get("authorizationSnapshot"),
"writerContextInput.authorizationSnapshot",
)
versions = {
"referenceWork.version": str(reference.get("version") or ""),
"authorization.sourceVersion": str(authorization.get("sourceVersion") or ""),
"authorization.authorizationSnapshot.sourceVersion": str(
top_snapshot.get("sourceVersion") or ""
),
"writerContextInput.sourceVersion": str(context_input.get("sourceVersion") or ""),
}
if any(not value for value in versions.values()) or len(set(versions.values())) != 1:
raise AuthorizationBindingError(
"WriterContext sourceVersion 未绑定 referenceWork 与不可变授权快照"
)
if str(context_snapshot.get("snapshotId") or "") != str(top_snapshot.get("id") or ""):
raise AuthorizationBindingError("WriterContext authorizationSnapshot.snapshotId 不匹配")
context_allowed_purpose = str(context_snapshot.get("allowedPurpose") or "")
if context_allowed_purpose not in {"evaluation", "offline_evaluation"}:
raise AuthorizationBindingError(
"WriterContext 授权用途必须对应 evaluation/offline_evaluation"
)
top_allowed = authorization.get("allowedPurpose")
snapshot_allowed = top_snapshot.get("allowedPurpose")
if (
not isinstance(top_allowed, list)
or "offline_evaluation" not in top_allowed
or not isinstance(snapshot_allowed, list)
or "offline_evaluation" not in snapshot_allowed
):
raise AuthorizationBindingError("顶层授权快照未授权 offline_evaluation")
if str(context_snapshot.get("verifiedAt") or "") != str(top_snapshot.get("checkedAt") or ""):
raise AuthorizationBindingError(
"WriterContext authorizationSnapshot.verifiedAt 未绑定顶层 checkedAt"
)
context_snapshot_source = context_snapshot.get("sourceVersion")
if context_snapshot_source is not None and str(context_snapshot_source) != versions[
"referenceWork.version"
]:
raise AuthorizationBindingError(
"WriterContext authorizationSnapshot.sourceVersion 不匹配"
)
top_status = str(authorization.get("sourceStatus") or "").lower()
context_status = str(context_input.get("sourceStatus") or "").lower()
allowed_context_statuses = {
"active": frozenset({"active"}),
"approved": frozenset({"active"}),
"authorized": frozenset({"active", "authorized", "frozen_authorized"}),
"licensed": frozenset({"active"}),
}
if context_status not in allowed_context_statuses.get(top_status, frozenset()):
raise AuthorizationBindingError(
"WriterContext sourceStatus 与顶层授权状态语义不一致"
)
def _build_arm_contexts(
*,
config: Mapping[str, Any],
sample: Mapping[str, Any],
common: Mapping[str, Any],
replay_run_id: str,
) -> tuple[dict[str, dict[str, Any]], dict[str, Any], str]:
"""经唯一组装入口构造 A/B/C,并再次执行 WriterContext 合同校验。"""
context_input = _mapping(sample.get("writerContextInput"), "writerContextInput")
_validate_context_authorization_binding(config, context_input)
content_mode = str(context_input.get("contentMode") or "")
if content_mode not in CONTENT_MODES:
raise WriterReplayError(
"writerContextInput.contentMode 必须明确为 sanitized_contract_fixture 或 canonical_frozen_prose"
)
requirements = _requirements(context_input)
sample_run_id = _run_id(f"{replay_run_id}:{sample['sampleId']}")
plan = _retrieval_plan(run_id=sample_run_id, common=common, context_input=context_input)
retrieval_result = _mapping(
context_input.get("retrievalResult"), "writerContextInput.retrievalResult"
)
contexts: dict[str, dict[str, Any]] = {}
for arm in REQUIRED_ARMS:
try:
assembled = assemble_context(
run_id=sample_run_id,
attempt=1,
mode="diagnostic_only",
purpose="evaluation",
quality_policy_version="writer-eval-v1",
work_id=int(common["workId"]),
target_chapter=int(common["targetChapter"]),
as_of=int(common["asOfChapter"]),
source_version=str(context_input.get("sourceVersion") or ""),
authorization_snapshot=_mapping(
context_input.get("authorizationSnapshot"),
"writerContextInput.authorizationSnapshot",
),
source_status=str(context_input.get("sourceStatus") or ""),
retrieval_plan=plan,
retrieval_result=retrieval_result,
fine_outline=_mapping(
context_input.get("fineOutline"), "writerContextInput.fineOutline"
),
narrative_state=_mapping(
context_input.get("narrativeState"), "writerContextInput.narrativeState"
),
recent_chapters=_sequence(
context_input.get("recentChapters"), "writerContextInput.recentChapters"
),
output_contract=_mapping(
context_input.get("outputContract"), "writerContextInput.outputContract"
),
token_budget=_mapping(
context_input.get("tokenBudget"), "writerContextInput.tokenBudget"
),
pattern_references=_sequence(
context_input.get("patternReferences", []),
"writerContextInput.patternReferences",
),
generated_at=str(context_input.get("generatedAt") or ""),
evidence_strategy=EVIDENCE_STRATEGIES[arm],
)
# 组装器内部已校验;这里保留显式二次校验,证明每一臂都经过合同入口。
contexts[arm] = validate_writer_context(assembled["context"])
except (AssemblyError, ContractError, TypeError, ValueError) as error:
raise WriterReplayError(f"{sample['sampleId']} {arm} 臂 WriterContext 非法: {error}") from error
return contexts, requirements, content_mode
def _context_summary(context: Mapping[str, Any], *, content_mode: str) -> dict[str, Any]:
"""只公开合同、冻结和来源元数据,不泄露证据文本或提示内容。"""
prose = context["proseEvidence"]
hints = context.get("indexHints", [])
return {
"schemaVersion": context["schemaVersion"],
"mode": context["mode"],
"purpose": context["purpose"],
"evidenceStrategy": context["evidenceStrategy"],
"acceptanceEligible": context["acceptanceEligible"],
"asOf": context["asOf"],
"targetChapter": context["targetChapter"],
"contentMode": content_mode,
"contextSnapshotId": context["contextSnapshot"]["manifestId"],
"contextSnapshotSha256": context["contextSnapshot"]["contextSha256"],
"retrievalPlanId": context["retrievalPlan"]["planId"],
"retrievalManifestId": context["retrievalManifest"]["manifestId"],
"factEvidenceCount": len(context["factEvidence"]),
"proseEvidenceCount": len(prose),
"proseChapters": [item["chapter"] for item in prose],
"recentBaselineChapters": [
item["chapter"] for item in prose if item["isRecentBaseline"] is True
],
"indexHintCount": len(hints),
"indexHintAsOf": sorted({item["asOf"] for item in hints}),
"sourceCount": len(context["retrievalManifest"]["sources"]),
"omittedSourceCount": len(context["retrievalManifest"]["omittedSources"]),
}
def _arm_manifests(
*,
common: Mapping[str, Any],
contexts: Mapping[str, Mapping[str, Any]],
snapshot_sha256: str,
content_mode: str,
) -> dict[str, dict[str, Any]]:
"""从已校验上下文生成三臂清单,避免配置宣称与实际输入漂移。"""
common_hash = sha256_value(common)
manifests: dict[str, dict[str, Any]] = {}
for arm in REQUIRED_ARMS:
context = contexts[arm]
prose = context["proseEvidence"]
manifests[arm] = {
"arm": arm,
"profile": PROFILE,
"evidenceStrategy": EVIDENCE_STRATEGIES[arm],
"commonControlsSha256": common_hash,
"snapshotManifestSha256": snapshot_sha256,
"acceptanceEligible": False,
"retrievalPlan": copy.deepcopy(context["retrievalPlan"]),
"retrievalManifest": copy.deepcopy(context["retrievalManifest"]),
"contextSummary": _context_summary(context, content_mode=content_mode),
"diagnosticContext": {
"proseSourceIds": [item["sourceRef"]["sourceId"] for item in prose],
"indexHintSourceIds": [item["sourceId"] for item in context.get("indexHints", [])],
"cardExpandedProseSourceIds": [
item["sourceRef"]["sourceId"]
for item in prose
if item["isRecentBaseline"] is False
],
},
}
control_manifests = {
arm: {
"arm": arm,
"asOfChapter": common["asOfChapter"],
"targetChapter": common["targetChapter"],
"commonControls": dict(common),
"snapshotManifestSha256": snapshot_sha256,
}
for arm in REQUIRED_ARMS
}
control_check = check_arm_manifests(
control_manifests,
required_arms=REQUIRED_ARMS,
as_of_chapter=int(common["asOfChapter"]),
target_chapter=int(common["targetChapter"]),
smoke=True,
)
if not control_check["ok"]:
raise WriterReplayError("三臂公共控制不一致: " + "; ".join(control_check["errors"]))
return manifests
def _prepare_sample(
config: Mapping[str, Any], sample: Mapping[str, Any], *, replay_run_id: str
) -> dict[str, Any]:
"""复用授权、章界、冻结和泄露逻辑,再组装三臂 WriterContext。"""
sample_id = str(sample.get("sampleId") or "")
if not sample_id:
raise WriterReplayError("样本缺少 sampleId")
common = _common_controls(config, sample)
sources = sample.get("sources")
if not isinstance(sources, list):
raise WriterReplayError(f"{sample_id}.sources 必须是数组")
source_check = check_target_sources(int(common["targetChapter"]), sources)
if not source_check["ok"]:
return {
"sampleId": sample_id,
"ok": False,
"status": source_check["status"],
"errors": source_check["errors"],
}
reference = _mapping(config.get("referenceWork"), "referenceWork")
authorization = _mapping(config.get("authorization"), "authorization")
try:
frozen = build_snapshot(
_mapping(sample.get("snapshotData"), f"{sample_id}.snapshotData"),
int(common["asOfChapter"]),
str(sample.get("snapshotVersion") or ""),
target_chapter=int(common["targetChapter"]),
manifest_metadata={
"referenceWork": {
"id": str(reference.get("id")),
"version": str(reference.get("version")),
},
"evaluationSetVersion": str(config.get("evaluationSetVersion") or ""),
"strategyVersion": str(config.get("strategyVersion") or ""),
"authorizationSnapshot": authorization.get("authorizationSnapshot"),
"runPermissions": {"purpose": "offline_evaluation", "mode": "diagnostic_only"},
"armConfig": {"arms": list(SNAPSHOT_COMPAT_ARMS)},
},
)
except SnapshotError as error:
return {
"sampleId": sample_id,
"ok": False,
"status": "invalid_snapshot",
"errors": [str(error)],
}
future_omissions = [
item
for item in frozen["manifest"].get("omittedSources", [])
if isinstance(item, Mapping) and item.get("reason") == "future_or_crosses_as_of"
]
if future_omissions:
return {
"sampleId": sample_id,
"ok": False,
"status": "invalid_leakage",
"errors": ["冻结输入包含目标章或未来章记录"],
"leakageAudit": {
"status": "invalid_snapshot",
"findingCount": len(future_omissions),
},
}
leakage = _mapping(sample.get("leakageAudit"), f"{sample_id}.leakageAudit")
context_input = _mapping(sample.get("writerContextInput"), f"{sample_id}.writerContextInput")
# 目标章细纲本来就带 targetChapter/sourceRef,不能交给“历史来源章界”扫描器;
# 泄露审计只扫描真正注入写手的冻结历史证据和卡索引提示。
audit = audit_snapshot(
{
"snapshot": frozen["snapshot"],
"writerEvidence": {
"recentChapters": context_input.get("recentChapters", []),
"retrievalResult": context_input.get("retrievalResult", {}),
"narrativeState": context_input.get("narrativeState", {}),
},
},
leakage.get("targetFacts"),
as_of=int(common["asOfChapter"]),
target=int(common["targetChapter"]),
)
if not audit["ok"]:
return {
"sampleId": sample_id,
"ok": False,
"status": audit["status"],
"errors": audit["errors"],
"leakageAudit": audit,
}
try:
contexts, requirements, content_mode = _build_arm_contexts(
config=config,
sample=sample,
common=common,
replay_run_id=replay_run_id,
)
arms = _arm_manifests(
common=common,
contexts=contexts,
snapshot_sha256=frozen["manifest"]["manifestSha256"],
content_mode=content_mode,
)
except WriterReplayError as error:
error_text = str(error)
if isinstance(error, AuthorizationBindingError):
status = "blocked_authorization"
elif "超出冻结线" in error_text or "目标章或未来章" in error_text:
status = "invalid_leakage"
else:
status = "invalid_writer_context"
failure = {
"sampleId": sample_id,
"ok": False,
"status": status,
"errors": [error_text],
}
if status == "invalid_leakage":
failure["leakageAudit"] = {
"status": "invalid_snapshot",
"findingCount": 1,
}
return failure
return {
"sampleId": sample_id,
"scenario": str(sample.get("scenario") or ""),
"ok": True,
"status": STATUS_READY,
"asOfChapter": common["asOfChapter"],
"targetChapter": common["targetChapter"],
"snapshotManifestSha256": frozen["manifest"]["manifestSha256"],
"leakageAudit": {
"status": audit["status"],
"findingCount": audit["findingCount"],
},
"requirementsSha256": sha256_value(requirements),
"arms": arms,
"commonControls": common,
"_contexts": contexts,
"_requirements": requirements,
}
def _public_sample(prepared: Mapping[str, Any]) -> dict[str, Any]:
"""移除仅供执行使用的完整上下文,避免 dry-run 回显脱敏夹具文本。"""
return {key: copy.deepcopy(value) for key, value in prepared.items() if not key.startswith("_")}
def _blind_mapping(evaluation_set_version: str, sample_id: str) -> tuple[dict[str, str], list[str]]:
"""从预注册评测集和样本 ID 派生稳定盲化顺序。"""
seed = int(sha256_value(f"{evaluation_set_version}:{sample_id}")[:16], 16)
shuffled = list(REQUIRED_ARMS)
random.Random(seed).shuffle(shuffled)
order = [f"blind-{index}" for index in range(1, 4)]
return dict(zip(order, shuffled, strict=True)), order
def _candidate_summary(candidate: Mapping[str, Any], arm: str) -> dict[str, Any]:
"""只保留安全候选元数据,正文和完整响应不得进入结构化结果。"""
if candidate.get("acceptanceEligible") is not False:
raise WriterReplayError(f"{arm} 臂候选必须 acceptanceEligible=false")
candidate_hash = candidate.get("candidateSha256")
candidate_version = candidate.get("candidateVersion")
if not isinstance(candidate_hash, str) or not candidate_hash.startswith("sha256:"):
raise WriterReplayError(f"{arm} 臂候选缺少 candidateSha256")
if isinstance(candidate_version, bool) or not isinstance(candidate_version, int):
raise WriterReplayError(f"{arm} 臂候选缺少 candidateVersion")
return {
"candidateId": f"{candidate['runId']}:{arm}:v{candidate_version}",
"candidateVersion": candidate_version,
"candidateSha256": candidate_hash,
"acceptanceEligible": False,
}
def _build_blind_input(
*,
raw_sample: Mapping[str, Any],
writer_fine_outline: Mapping[str, Any],
candidates: Mapping[str, Mapping[str, Any]],
mapping: Mapping[str, str],
blind_id: str,
candidate_order: list[str],
) -> dict[str, Any]:
"""构造 judge 唯一可见输入;只含共同参考与盲化候选。"""
if blind_id not in mapping or set(mapping.values()) != set(REQUIRED_ARMS):
raise WriterReplayError("盲评映射非法")
if set(candidate_order) != set(mapping):
raise WriterReplayError("盲评候选顺序与预注册盲 ID 不一致")
blind_candidates: list[dict[str, Any]] = []
for current_blind_id in candidate_order:
arm = mapping[current_blind_id]
candidate = candidates[arm]
body = candidate.get("candidateBody")
candidate_hash = candidate.get("candidateSha256")
if not isinstance(body, str) or not body:
raise WriterReplayError("盲评候选缺少正文内容")
if not isinstance(candidate_hash, str) or not candidate_hash.startswith("sha256:"):
raise WriterReplayError("盲评候选缺少正文哈希")
blind_candidates.append(
{
"blindCandidateId": current_blind_id,
"candidateSha256": candidate_hash,
"candidateBody": body,
}
)
context_input = _mapping(raw_sample.get("writerContextInput"), "writerContextInput")
requirements = _mapping(context_input.get("requirements"), "writerContextInput.requirements")
if set(writer_fine_outline) != set(FINE_OUTLINE_JUDGE_FIELDS):
raise WriterReplayError("写手细纲字段必须严格匹配 judge 共同参考合同")
shared_reference = {
# 只复制写手已通过合同校验的四字段视图,禁止从原始 fineOutline 倾倒评委不可见字段。
"fineOutline": {
field: copy.deepcopy(writer_fine_outline[field])
for field in FINE_OUTLINE_JUDGE_FIELDS
},
"requirements": {
"requiredEvents": copy.deepcopy(
_sequence(requirements.get("requiredEvents"), "requirements.requiredEvents")
),
"requiredCharacters": copy.deepcopy(
_sequence(requirements.get("requiredCharacters"), "requirements.requiredCharacters")
),
"foreshadowingActions": copy.deepcopy(
_sequence(
requirements.get("foreshadowingActions"),
"requirements.foreshadowingActions",
)
),
"chapterEndHook": copy.deepcopy(
_mapping(requirements.get("chapterEndHook"), "requirements.chapterEndHook")
),
},
"historicalProseBaseline": copy.deepcopy(
_sequence(context_input.get("recentChapters"), "writerContextInput.recentChapters")
),
}
return {
"schemaVersion": "writer-blind-input-v1",
"sample": {
"sampleId": str(raw_sample.get("sampleId") or ""),
"scenario": str(raw_sample.get("scenario") or ""),
"targetChapter": raw_sample.get("targetChapter"),
"targetTitle": str(raw_sample.get("targetTitle") or ""),
"targetDescription": str(raw_sample.get("targetDescription") or ""),
},
"blindCandidateId": blind_id,
"candidateOrder": list(candidate_order),
"sharedEvaluationReference": shared_reference,
"candidates": blind_candidates,
}
def _execute_sample_for_test(
*,
raw_sample: Mapping[str, Any],
prepared: Mapping[str, Any],
evaluation_set_version: str,
adapters: WriterReplayTestAdapters,
raw_dir: Path,
) -> dict[str, Any]:
"""仅供测试验证正式 writer/pipeline/detector 编排,不代表生产已接线。"""
raw_candidates: dict[str, dict[str, Any]] = {}
candidate_summaries: dict[str, dict[str, Any]] = {}
detector: dict[str, dict[str, Any]] = {}
for arm in REQUIRED_ARMS:
context = prepared["_contexts"][arm]
captured: dict[str, Mapping[str, Any]] = {}
def writer(
current_context: Mapping[str, Any],
candidate_version: int,
_repair_failures: list[dict[str, Any]],
) -> Mapping[str, Any]:
"""通过正式 adapter 调 subprocess runner,禁止直接返回候选。"""
try:
candidate = run_writer(current_context, runner=adapters.writer_runner)
except WriterAdapterError as error:
raise PipelineError(error.code, str(error), details=error.details) from error
if candidate["candidateVersion"] != candidate_version:
raise PipelineError("WRITER_OUTPUT_STALE", "候选版本未绑定本次管线状态")
captured["candidate"] = candidate
return candidate
def no_evidence_provider(*_args: Any, **_kwargs: Any) -> Mapping[str, Any]:
"""固定三臂实验不允许生成中途改变证据集合。"""
raise PipelineError("REPLAY_EVIDENCE_REQUEST_FORBIDDEN", "固定回放禁止补证")
try:
pipeline_result = run_writer_pipeline(
context=context,
requirements=prepared["_requirements"],
writer=writer,
evidence_provider=no_evidence_provider,
semantic_detector=adapters.semantic_detector,
state_store=InMemoryCasStateStore(),
result_path=raw_dir / f"pipeline-{arm}.json",
)
except PipelineError as error:
raise WriterReplayError(f"{prepared['sampleId']} {arm} 臂管线失败: {error.code}") from error
candidate = captured.get("candidate")
if candidate is None:
raise WriterReplayError(f"{arm} 臂管线未留下已校验候选")
atomic_write_json(raw_dir / f"candidate-{arm}.json", candidate)
raw_candidates[arm] = candidate
candidate_summaries[arm] = _candidate_summary(candidate, arm)
last_trace = pipeline_result["trace"][-1]
detector[arm] = {
"passed": pipeline_result["status"] == "PASSED",
"highSeverityCount": 0,
"hardConstraintCoverage": 1.0,
"mechanicalPassed": last_trace.get("mechanicalPassed") is True,
"semanticStatus": last_trace.get("semanticStatus"),
}
mapping, blind_order = _blind_mapping(evaluation_set_version, str(prepared["sampleId"]))
# 共同参考取自实际写手上下文,并在进入 judge 前确认三臂视图完全一致。
writer_fine_outline = _mapping(
prepared["_contexts"][REQUIRED_ARMS[0]].get("fineOutline"),
"writerContext.fineOutline",
)
for arm in REQUIRED_ARMS[1:]:
current_fine_outline = _mapping(
prepared["_contexts"][arm].get("fineOutline"),
f"writerContext.{arm}.fineOutline",
)
if current_fine_outline != writer_fine_outline:
raise WriterReplayError("三臂写手细纲不一致,不能构造共同评测参考")
reviews: dict[str, Any] = {}
first_reports: list[Mapping[str, Any]] = []
for blind_id in blind_order:
first_input = _build_blind_input(
raw_sample=raw_sample,
writer_fine_outline=writer_fine_outline,
candidates=raw_candidates,
mapping=mapping,
blind_id=blind_id,
candidate_order=blind_order,
)
first = adapters.judge(first_input, "judge-1")
first_reports.append(first)
second_input = _build_blind_input(
raw_sample=raw_sample,
writer_fine_outline=writer_fine_outline,
candidates=raw_candidates,
mapping=mapping,
blind_id=blind_id,
candidate_order=list(reversed(blind_order)),
)
second = adapters.judge(second_input, "judge-2")
adjudication = adjudicate_reviews(first, second)
if adjudication.get("status") == "needs_third_reviewer":
third_input = _build_blind_input(
raw_sample=raw_sample,
writer_fine_outline=writer_fine_outline,
candidates=raw_candidates,
mapping=mapping,
blind_id=blind_id,
candidate_order=blind_order,
)
third = adapters.judge(third_input, "judge-3")
adjudication = adjudicate_reviews(first, second, third)
reviews[blind_id] = adjudication
deblind_reports(first_reports, mapping)
reviews_by_arm = {mapping[blind_id]: reviews[blind_id] for blind_id in blind_order}
return {
**_public_sample(prepared),
"status": "completed_test_pipeline",
"candidates": candidate_summaries,
"detector": detector,
"blindOrderSha256": sha256_value(blind_order),
"reviews": reviews_by_arm,
}
def run_writer_replay(
config: Mapping[str, Any],
*,
run_id: str | None = None,
output_dir: Path | None = None,
execute: bool = False,
test_adapters: WriterReplayTestAdapters | None = None,
) -> dict[str, Any]:
"""准备正文回放;测试可验证执行管线,生产 execute 当前失败关闭。"""
if not isinstance(config, Mapping) or config.get("profile") != PROFILE:
raise WriterReplayError(f"profile 必须是 {PROFILE}")
evaluation_set_version = str(config.get("evaluationSetVersion") or "")
strategy_version = str(config.get("strategyVersion") or "")
replay_run_id = _run_id(str(run_id or "dry-run"))
if not evaluation_set_version or not strategy_version:
raise WriterReplayError("evaluationSetVersion/strategyVersion 不能为空")
authorization = check_authorization(config.get("authorization"))
if not authorization["ok"]:
return {
"runId": replay_run_id,
"mode": "execute" if execute else "dry_run",
"status": authorization["status"],
"ok": False,
"errors": authorization["errors"],
"samples": [],
}
samples = config.get("samples")
if not isinstance(samples, list) or not samples or any(
not isinstance(item, Mapping) for item in samples
):
raise WriterReplayError("samples 必须是非空对象数组")
sample_ids = [str(item.get("sampleId") or "") for item in samples]
if any(not item for item in sample_ids) or len(set(sample_ids)) != len(sample_ids):
raise WriterReplayError("sampleId 必须非空且不能重复")
prepared = [
_prepare_sample(config, sample, replay_run_id=replay_run_id) for sample in samples
]
invalid = [item for item in prepared if not item["ok"]]
result: dict[str, Any] = {
"runId": replay_run_id,
"mode": "execute" if execute else "dry_run",
"profile": PROFILE,
"evaluationSetVersion": evaluation_set_version,
"strategyVersion": strategy_version,
"status": "ready",
"ok": not invalid,
"executeAdapterStatus": "semantic_and_judge_not_wired",
"samples": [_public_sample(item) for item in prepared],
}
if invalid:
statuses = {item["status"] for item in invalid}
result["status"] = sorted(statuses)[0]
if any(item.get("leakageAudit") for item in invalid):
result["status"] = "invalid_leakage"
result["ok"] = False
return result
if not execute:
return result
if output_dir is None:
raise WriterReplayError("execute 模式必须显式提供 output_dir")
resolved_output = output_dir.resolve()
if resolved_output == PRIVATE_TMP or not resolved_output.is_relative_to(PRIVATE_TMP):
raise WriterReplayError("execute 输出目录必须位于 /private/tmp 的独立子目录")
if test_adapters is None:
raise WriterReplayError(
"真实 execute 未接线:缺少生产 semantic detector 与 judge adapter,已失败关闭"
)
# 这一路径仅证明边界可测试;CLI 无法注入 test_adapters,不能宣称 real-run 完成。
raw_dir = resolved_output / "raw"
raw_dir.mkdir(parents=True, exist_ok=True)
executed: list[dict[str, Any]] = []
for raw_sample, prepared_sample in zip(samples, prepared, strict=True):
sample_raw_dir = raw_dir / prepared_sample["sampleId"]
sample_raw_dir.mkdir(parents=True, exist_ok=True)
executed.append(
_execute_sample_for_test(
raw_sample=raw_sample,
prepared=prepared_sample,
evaluation_set_version=evaluation_set_version,
adapters=test_adapters,
raw_dir=sample_raw_dir,
)
)
result.update(
{
"status": "completed_test_pipeline",
"ok": True,
"executeAdapterStatus": "test_injection_only",
"samples": executed,
}
)
resolved_output.mkdir(parents=True, exist_ok=True)
atomic_write_json(resolved_output / "manifest.json", result)
return result
def _parse_args() -> argparse.Namespace:
"""解析固定 CLI;默认 dry-run,显式 execute 仍受未接线门阻断。"""
parser = argparse.ArgumentParser(description="运行 writer_replay A/B/C 回放")
parser.add_argument("--config", type=Path, required=True)
parser.add_argument("--run-id")
parser.add_argument("--output-dir", type=Path)
mode = parser.add_mutually_exclusive_group()
mode.add_argument("--dry-run", action="store_true")
mode.add_argument("--execute", action="store_true")
return parser.parse_args()
def main() -> int:
"""CLI 只真正开放零模型 dry-run;execute 未接线时返回稳定阻塞。"""
args = _parse_args()
config = json.loads(args.config.read_text(encoding="utf-8"))
result = run_writer_replay(
config,
run_id=args.run_id,
output_dir=args.output_dir,
execute=args.execute,
)
print(json.dumps(result, ensure_ascii=False, sort_keys=True, indent=2))
return 0 if result["ok"] else 2
if __name__ == "__main__":
try:
raise SystemExit(main())
except (WriterReplayError, OSError, json.JSONDecodeError) as error:
print(
json.dumps(
{"status": "blocked_writer_replay", "error": str(error)},
ensure_ascii=False,
)
)
raise SystemExit(2)
__all__ = [
"CONTENT_MODES",
"EVIDENCE_STRATEGIES",
"PROFILE",
"REQUIRED_ARMS",
"WriterReplayError",
"WriterReplayTestAdapters",
"run_writer_replay",
]