zizi 1d32a96399 修复: 收紧回放评测失败关闭与盲化边界
复用目录先原子落盘非成功态,并为三类外部 runner 增加可配置超时。detector 改为真正盲检,安全报告与嵌套响应统一执行严格闭集校验。
2026-07-19 21:17:34 +08:00

938 lines
34 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
"""细纲回放编排器。
默认只做 dry-run。真实模式显式调用本地 Claude CLI,并将原始候选限定在运行目录;
最终结果只返回状态、哈希、评分入口和失败类别,不把原始 prompt/response 带出运行目录。
"""
from __future__ import annotations
import argparse
import json
import math
import os
import re
import subprocess
import tempfile
from pathlib import Path
from typing import Any, Mapping, Sequence
from build_snapshot import SnapshotError, build_snapshot, normalize_chapter, sha256_value
from audit_leakage import audit_snapshot
from check_snapshot import (
STATUS_READY,
check_candidate_output,
check_replay,
)
from fine_outline_detector import build_detector_request, validate_detector_report
from fine_outline_rubric import DIMENSIONS, RUBRIC_PROFILE, stability_warning, validate_report
REQUIRED_ARMS = ("outline_only", "outline_plus_cards", "outline_plus_placebo_cards")
REPO_ROOT = Path(__file__).resolve().parents[4]
SKILL_PATH = REPO_ROOT / ".claude/skills/fine-outline/SKILL.md"
PLANNER_PATH = REPO_ROOT / ".claude/agents/planner.md"
JUDGE_IDS = ("judge-primary", "judge-secondary")
DEFAULT_TIMEOUT_SECONDS = 300.0
class ReplayRunError(ValueError):
"""回放配置不符合运行边界。"""
class RunnerInvocationError(ReplayRunError):
"""外部 runner 无法启动或以非零状态退出。"""
class RunnerOutputError(ReplayRunError):
"""外部 runner 返回的内容不符合机器合同。"""
class RunnerTimeoutError(ReplayRunError):
"""外部 runner 超过允许的最长执行时间。"""
def _read_json(path: Path) -> Any:
return json.loads(path.read_text(encoding="utf-8"))
def _safe_json(value: Any) -> str:
return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
def _timeout_stdout(error: subprocess.TimeoutExpired) -> str:
"""规范化超时前捕获的标准输出,供临时目录留痕和哈希审计。"""
output = error.stdout or ""
if isinstance(output, bytes):
return output.decode("utf-8", errors="replace")
return output
def _require_mapping(config: Mapping[str, Any], key: str) -> Mapping[str, Any]:
value = config.get(key)
if not isinstance(value, Mapping):
raise ReplayRunError(f"配置缺少对象字段: {key}")
return value
def _build_arm_manifest(
*,
name: str,
arm: Mapping[str, Any],
common_input: Mapping[str, Any],
as_of: int,
target: int,
snapshot_version: str,
) -> dict[str, Any]:
cards = arm.get("cards", [])
if not isinstance(cards, list):
raise ReplayRunError(f"{name}.cards 必须是数组")
source_ids = arm.get("cardSourceIds", [])
if not isinstance(source_ids, list):
raise ReplayRunError(f"{name}.cardSourceIds 必须是数组")
strategy = str(arm.get("cardStrategy") or ("none" if name == "outline_only" else "correct"))
return {
"arm": name,
"snapshotVersion": snapshot_version,
"asOfChapter": as_of,
"targetChapter": target,
"commonInputSha256": sha256_value(common_input),
"cardInjectionSha256": sha256_value(cards),
"cardInjectionCount": len(cards),
"cardSourceIds": [str(item) for item in source_ids],
"cardStrategy": strategy,
}
def _extract_candidate(output: str) -> Mapping[str, Any]:
"""兼容 Claude JSON 外壳和代码围栏,但不把原始文本返回调用方。"""
outer: Any = output
try:
outer = json.loads(output)
except json.JSONDecodeError:
pass
if isinstance(outer, Mapping) and isinstance(outer.get("result"), str):
outer = outer["result"]
if isinstance(outer, str):
match = re.search(r"```(?:json)?\s*(\{.*\})\s*```", outer, re.DOTALL)
outer = match.group(1) if match else outer.strip()
outer = json.loads(outer)
if not isinstance(outer, Mapping):
raise RunnerOutputError("模型输出不是 JSON 对象")
return outer
def _planner_prompt(
*,
target: int,
as_of: int,
snapshot: Mapping[str, Any],
common_input: Mapping[str, Any],
cards: list[Any],
) -> str:
"""只把冻结后的结构化资料和功能合同送入 planner。"""
skill = SKILL_PATH.read_text(encoding="utf-8")
identity = PLANNER_PATH.read_text(encoding="utf-8")
# 公共快照不能携带知识卡;卡只能通过当前评测臂的独立注入区进入上下文,
# 否则 outline_only 臂会在“无卡”名义下偷看到全局卡片。
public_snapshot = {
str(key): value for key, value in snapshot.items() if str(key) != "cards"
}
context = {
"targetChapter": target,
"asOfChapter": as_of,
"frozenSnapshot": public_snapshot,
"commonContext": common_input,
"cardInjection": cards,
}
return "\n".join(
[
"这是 next_fine_outline_replay_v0 的离线规划任务。只输出一个 JSON 对象,不要 Markdown、正文或解释。",
"你不能调用工具,也不能读取仓库、目标章或任何未列入下面 JSON 的资料。",
"严格执行 fine-outline 合同;目标章号必须保持不变,未知内容写入 unknowns/assumptions。",
"规划上下文冻结到 as_of;卡只是事实索引和补充,不得替代公共大纲与叙事现在时。",
"--- planner identity ---",
identity,
"--- fine-outline skill ---",
skill,
"--- frozen input ---",
_safe_json(context),
"--- output contract ---",
_safe_json(
{
"targetChapter": target,
"requiredFields": [
"chapterGoal",
"keyEvents",
"entities",
"foreshadowing",
"stateChanges",
"hook",
"unknowns",
"assumptions",
],
"eventFields": ["id", "order", "event", "participants", "trigger", "resultDirection"],
"eventOrder": "至少一个事件;order 必须从 1 开始连续严格递增",
"entityFields": ["name", "type", "role"],
"foreshadowingFields": ["action", "subject", "evidence"],
"sourceRefs": "可选;只能引用冻结来源 ID",
}
),
]
)
def _invoke_planner(
*,
prompt: str,
planner_bin: str,
model: str,
output_path: Path,
max_budget_usd: float,
timeout_seconds: float,
) -> Mapping[str, Any]:
"""用无工具、无会话持久化的 Claude print 模式运行 planner。"""
command = [
planner_bin,
"-p",
"--agent",
"planner",
"--model",
model,
"--tools",
"",
"--no-session-persistence",
"--output-format",
"json",
"--max-budget-usd",
str(max_budget_usd),
"--append-system-prompt",
"本次是严格离线回放;不要调用任何工具,不要读取文件,不要输出 JSON 以外内容。",
prompt,
]
try:
completed = subprocess.run(
command,
text=True,
capture_output=True,
check=False,
timeout=timeout_seconds,
)
except subprocess.TimeoutExpired as error:
output_path.write_text(_timeout_stdout(error), encoding="utf-8")
raise RunnerTimeoutError(f"planner 调用超时,限制={timeout_seconds:g}秒") from error
except OSError as error:
raise RunnerInvocationError(f"planner 启动失败: {error}") from error
output_path.write_text(completed.stdout, encoding="utf-8")
if completed.returncode != 0:
raise RunnerInvocationError(f"planner 调用失败,退出码={completed.returncode}")
return _extract_candidate(completed.stdout)
def _invoke_structured_agent(
*,
agent: str,
request: Mapping[str, Any],
runner_bin: str,
model: str,
output_path: Path,
max_budget_usd: float,
identity: str,
timeout_seconds: float,
) -> Mapping[str, Any]:
"""以独立无会话进程调用 detector/judge,并保存仓库外原始响应。"""
command = [
runner_bin,
"-p",
"--agent",
agent,
"--model",
model,
"--tools",
"",
"--no-session-persistence",
"--output-format",
"json",
"--max-budget-usd",
str(max_budget_usd),
"--append-system-prompt",
f"独立身份={identity};只处理给定 JSON;禁止调用工具、读取文件或输出 JSON 以外内容。",
_safe_json(request),
]
try:
completed = subprocess.run(
command,
text=True,
capture_output=True,
check=False,
timeout=timeout_seconds,
)
except subprocess.TimeoutExpired as error:
output_path.write_text(_timeout_stdout(error), encoding="utf-8")
raise RunnerTimeoutError(f"{agent} 调用超时,限制={timeout_seconds:g}秒") from error
except OSError as error:
raise RunnerInvocationError(f"{agent} 启动失败: {error}") from error
output_path.write_text(completed.stdout, encoding="utf-8")
if completed.returncode != 0:
raise RunnerInvocationError(f"{agent} 调用失败,退出码={completed.returncode}")
return _extract_candidate(completed.stdout)
def _public_snapshot(snapshot: Mapping[str, Any]) -> dict[str, Any]:
"""公共评测上下文不携带快照中的卡集合。"""
return {str(key): value for key, value in snapshot.items() if str(key) != "cards"}
def _blind_assignments(
candidates: Mapping[str, Mapping[str, Any]], run_id: str
) -> list[dict[str, Any]]:
"""按运行 ID 稳定打乱三臂,并只向审查模型暴露匿名候选 ID。"""
ordered_arms = sorted(
candidates,
key=lambda arm: sha256_value({"runId": run_id, "purpose": "blind-order", "arm": arm}),
)
return [
{
"arm": arm,
"candidateId": f"candidate-{sha256_value({'runId': run_id, 'arm': arm})[:12]}",
"candidate": candidates[arm],
}
for arm in ordered_arms
]
def _judge_request(
*,
judge_id: str,
assignments: Sequence[Mapping[str, Any]],
as_of: int,
frozen_snapshot: Mapping[str, Any],
common_context: Mapping[str, Any],
reference_proxy: Mapping[str, Any],
) -> dict[str, Any]:
"""构造不含 arm 和卡 manifest 的盲评批输入。"""
return {
"protocol": "fine_outline_judge_v0",
"profile": RUBRIC_PROFILE,
"judgeId": judge_id,
"asOfChapter": as_of,
"candidates": [
{"candidateId": item["candidateId"], "candidate": item["candidate"]}
for item in assignments
],
"frozenContext": {
"snapshot": _public_snapshot(frozen_snapshot),
"commonContext": dict(common_context),
},
"referenceProxy": dict(reference_proxy),
"rules": {
"armIdentityVisible": False,
"cardManifestVisible": False,
"proseDimensionsForbidden": True,
"scoreEvidenceRequired": True,
},
}
def _evaluation_by_candidate(report: Mapping[str, Any]) -> dict[str, Mapping[str, Any]]:
"""把已校验的 judge 批报告按匿名候选 ID 建索引。"""
return {
str(item["candidateId"]): item
for item in report["evaluations"]
if isinstance(item, Mapping)
}
def _aggregate_evaluation(
assignments: Sequence[Mapping[str, Any]], reports: Sequence[Mapping[str, Any]]
) -> dict[str, Any]:
"""去盲汇总双评分;稳定性未通过时不计算卡增量矩阵。"""
indexed = [_evaluation_by_candidate(report) for report in reports]
arm_scores: dict[str, dict[str, float]] = {}
stability_by_arm: dict[str, dict[str, Any]] = {}
for assignment in assignments:
arm = str(assignment["arm"])
candidate_id = str(assignment["candidateId"])
first_scores = {
dimension: float(indexed[0][candidate_id]["scores"][dimension]["score"])
for dimension in DIMENSIONS
}
second_scores = {
dimension: float(indexed[1][candidate_id]["scores"][dimension]["score"])
for dimension in DIMENSIONS
}
stability_by_arm[arm] = stability_warning(first_scores, second_scores)
arm_scores[arm] = {
dimension: round((first_scores[dimension] + second_scores[dimension]) / 2, 3)
for dimension in DIMENSIONS
}
max_gaps = {
dimension: max(
stability_by_arm[arm]["gaps"].get(dimension, float("inf"))
for arm in REQUIRED_ARMS
)
for dimension in DIMENSIONS
}
stable = all(item["stable"] for item in stability_by_arm.values())
evaluation: dict[str, Any] = {
"profile": RUBRIC_PROFILE,
"judgeIds": list(JUDGE_IDS),
"armScores": arm_scores,
"stability": {
"stable": stable,
"threshold": 0.5,
"maxGaps": max_gaps,
"byArm": stability_by_arm,
},
}
if stable:
baseline = arm_scores["outline_only"]
evaluation["deltas"] = {
"B-A": {
dimension: round(arm_scores["outline_plus_cards"][dimension] - baseline[dimension], 3)
for dimension in DIMENSIONS
},
"C-A": {
dimension: round(
arm_scores["outline_plus_placebo_cards"][dimension] - baseline[dimension],
3,
)
for dimension in DIMENSIONS
},
}
return evaluation
def _write_result(output_dir: Path, result: Mapping[str, Any]) -> None:
"""每个阶段都覆盖写入可恢复的结构化运行状态。"""
result_path = output_dir / "run_result.json"
temporary_path: Path | None = None
try:
# 唯一临时文件避免同目录并发写相互覆盖;同目录 replace 保证正式状态原子切换。
with tempfile.NamedTemporaryFile(
mode="w",
encoding="utf-8",
dir=output_dir,
prefix=".run_result.",
suffix=".tmp",
delete=False,
) as handle:
temporary_path = Path(handle.name)
handle.write(_safe_json(result) + "\n")
handle.flush()
os.fsync(handle.fileno())
temporary_path.replace(result_path)
finally:
if temporary_path is not None and temporary_path.exists():
temporary_path.unlink()
def run_replay(
config: Mapping[str, Any],
output_dir: Path,
*,
mode: str = "dry_run",
planner_bin: str = "claude",
detector_bin: str | None = None,
judge_primary_bin: str | None = None,
judge_secondary_bin: str | None = None,
model: str = "opus",
max_budget_usd: float = 1.0,
timeout_seconds: float = DEFAULT_TIMEOUT_SECONDS,
) -> dict[str, Any]:
"""执行一次单目标三臂回放;任何前置门失败都不调用模型。"""
output_dir = output_dir.resolve()
if output_dir.is_relative_to(REPO_ROOT.resolve()):
raise ReplayRunError("原始候选运行目录不得位于仓库内")
output_dir.mkdir(parents=True, exist_ok=True)
result: dict[str, Any] = {
"runId": str(config.get("runId") or "unassigned"),
"mode": mode,
"status": "validating_config",
"ok": False,
"results": {},
}
_write_result(output_dir, result)
try:
if mode not in {"dry_run", "execute"}:
raise ReplayRunError("mode 只能是 dry_run 或 execute")
if (
isinstance(timeout_seconds, bool)
or not isinstance(timeout_seconds, (int, float))
or not math.isfinite(float(timeout_seconds))
or timeout_seconds <= 0
):
raise ReplayRunError("timeout_seconds 必须是正数")
snapshot_config = _require_mapping(config, "snapshot")
as_of = normalize_chapter(snapshot_config.get("asOfChapter"))
target = normalize_chapter(config.get("targetChapter"))
snapshot_version = str(snapshot_config.get("snapshotVersion") or "")
if as_of is None or target is None:
raise ReplayRunError("as_of/target 必须是明确正整数")
if target != as_of + 1:
raise ReplayRunError("targetChapter 必须等于 snapshot.asOfChapter+1")
reference_work = _require_mapping(config, "referenceWork")
authorization = _require_mapping(config, "authorization")
sources = config.get("sources", [])
if not isinstance(sources, list):
raise ReplayRunError("sources 必须是数组")
common_input = _require_mapping(config, "commonContext")
arms = _require_mapping(config, "arms")
if set(arms) != set(REQUIRED_ARMS):
raise ReplayRunError("生产回放必须精确配置三臂")
arm_manifests = {
name: _build_arm_manifest(
name=name,
arm=_require_mapping(arms, name),
common_input=common_input,
as_of=as_of,
target=target,
snapshot_version=snapshot_version,
)
for name in REQUIRED_ARMS
}
preflight = check_replay(
authorization=authorization,
as_of_chapter=as_of,
target_chapter=target,
planner_sources=sources,
arm_manifests=arm_manifests,
)
except ReplayRunError as error:
result["status"] = "config_invalid"
result["errors"] = [str(error)]
_write_result(output_dir, result)
raise
result.update(
{
"referenceWork": str(reference_work.get("id") or ""),
"referenceWorkVersion": str(reference_work.get("version") or ""),
"asOfChapter": as_of,
"targetChapter": target,
"snapshotVersion": snapshot_version,
"preflight": {
"status": preflight["status"],
"errors": preflight["errors"],
"warnings": preflight["warnings"],
},
"arms": arm_manifests,
}
)
if preflight["ok"]:
# 前置门通过不等于快照配置有效;深层冻结完成前始终保持非成功态。
result["status"] = "validating_config"
result["ok"] = False
else:
result["status"] = preflight["status"]
result["ok"] = False
_write_result(output_dir, result)
if not preflight["ok"]:
return result
leakage_audit_config = config.get("leakageAudit")
if not isinstance(leakage_audit_config, Mapping):
# 没有目标事实审计就不能把 dry-run 说成可运行,避免结构冻结掩盖内容泄露。
result["status"] = "blocked_leakage_audit"
result["ok"] = False
result["leakageAudit"] = {
"status": "not_configured",
"ok": False,
"errors": ["缺少 leakageAudit 配置"],
"warnings": [],
"findingCount": 0,
"findings": [],
}
_write_result(output_dir, result)
return result
metadata = {
"targetChapter": target,
"referenceWork": reference_work,
"evaluationSetVersion": config.get("evaluationSetVersion", "unregistered"),
"strategyVersion": config.get("strategyVersion", "unregistered"),
"authorizationSnapshot": authorization["authorizationSnapshot"],
"runPermissions": config.get("runPermissions", {"purpose": "offline_evaluation", "mode": mode}),
"armConfig": {"arms": list(REQUIRED_ARMS)},
}
snapshot_data = snapshot_config.get("data", {})
try:
frozen = build_snapshot(
snapshot_data,
as_of,
snapshot_version,
target_chapter=target,
manifest_metadata=metadata,
)
except SnapshotError as error:
result["status"] = "config_invalid"
result["ok"] = False
result["errors"] = [str(error)]
_write_result(output_dir, result)
raise
# 内容审计同时覆盖公共冻结快照和各臂卡注入区;卡不在公共区,不能因此逃过未来事实检查。
audit_payload = {
"snapshot": frozen["snapshot"],
"armCardInjections": {
name: _require_mapping(arms, name).get("cards", []) for name in REQUIRED_ARMS
},
}
audit_result = audit_snapshot(
audit_payload,
leakage_audit_config.get("targetFacts"),
as_of=as_of,
target=target,
)
result["leakageAudit"] = audit_result
if not audit_result["ok"]:
# 审计失败的快照不生成 manifest,也不允许进入任何 planner 臂。
result["status"] = audit_result["status"]
result["ok"] = False
_write_result(output_dir, result)
return result
(output_dir / "snapshot_manifest.json").write_text(_safe_json(frozen["manifest"]) + "\n", encoding="utf-8")
if mode == "dry_run":
result["status"] = STATUS_READY
result["ok"] = True
result["snapshotManifestSha256"] = frozen["manifest"]["manifestSha256"]
_write_result(output_dir, result)
return result
result["status"] = "running_planner"
result["ok"] = False
result["snapshotManifestSha256"] = frozen["manifest"]["manifestSha256"]
_write_result(output_dir, result)
candidates: dict[str, Mapping[str, Any]] = {}
for name in REQUIRED_ARMS:
arm = _require_mapping(arms, name)
cards = arm.get("cards", [])
raw_path = output_dir / f"planner_{name}.raw.json"
prompt = _planner_prompt(
target=target,
as_of=as_of,
snapshot=frozen["snapshot"],
common_input=common_input,
cards=cards,
)
try:
candidate = _invoke_planner(
prompt=prompt,
planner_bin=planner_bin,
model=model,
output_path=raw_path,
max_budget_usd=max_budget_usd,
timeout_seconds=float(timeout_seconds),
)
candidate_path = output_dir / f"candidate_{name}.json"
candidate_path.write_text(_safe_json(candidate) + "\n", encoding="utf-8")
schema = check_candidate_output(candidate, target, sources)
candidates[name] = candidate
result["results"][name] = {
"status": schema["status"],
"ok": schema["ok"],
"errors": schema["errors"],
"candidateSha256": sha256_value(candidate),
"candidatePath": str(candidate_path),
}
except RunnerTimeoutError as error:
result["results"][name] = {
"status": "planner_timeout",
"ok": False,
"errors": [str(error)],
"rawOutputSha256": sha256_value(raw_path.read_text(encoding="utf-8")) if raw_path.exists() else None,
}
result["status"] = "planner_timeout"
_write_result(output_dir, result)
return result
except RunnerInvocationError as error:
result["results"][name] = {
"status": "planner_failed",
"ok": False,
"errors": [str(error)],
"rawOutputSha256": sha256_value(raw_path.read_text(encoding="utf-8")) if raw_path.exists() else None,
}
result["status"] = "planner_failed"
_write_result(output_dir, result)
return result
except (RunnerOutputError, json.JSONDecodeError) as error:
result["results"][name] = {
"status": "planner_invalid",
"ok": False,
"errors": [str(error)],
"rawOutputSha256": sha256_value(raw_path.read_text(encoding="utf-8")) if raw_path.exists() else None,
}
result["status"] = "planner_invalid"
_write_result(output_dir, result)
return result
if not all(item["ok"] for item in result["results"].values()):
result["status"] = "candidate_blocked"
result["ok"] = False
_write_result(output_dir, result)
return result
run_id = str(result["runId"])
assignments = _blind_assignments(candidates, run_id)
detector_runner = detector_bin or planner_bin
result["status"] = "running_detector"
_write_result(output_dir, result)
detector_invalid = False
detector_blocked = False
for assignment in assignments:
arm = str(assignment["arm"])
candidate_id = str(assignment["candidateId"])
request = build_detector_request(
candidate_id=candidate_id,
candidate=assignment["candidate"],
as_of_chapter=as_of,
frozen_snapshot=_public_snapshot(frozen["snapshot"]),
common_context=common_input,
sources=sources,
)
detector_path = output_dir / f"detector_{candidate_id}.raw.json"
try:
detector_report = _invoke_structured_agent(
agent="detector",
request=request,
runner_bin=detector_runner,
model=model,
output_path=detector_path,
max_budget_usd=max_budget_usd,
identity="blind-detector",
timeout_seconds=float(timeout_seconds),
)
detector_check = validate_detector_report(detector_report, candidate_id)
detector_invalid = detector_invalid or not detector_check["ok"]
detector_blocked = detector_blocked or detector_check["highSeverityCount"] > 0
result["results"][arm]["detector"] = {
"status": (
"invalid"
if not detector_check["ok"]
else "blocked_high"
if detector_check["highSeverityCount"]
else "passed"
),
"findingCount": detector_check["findingCount"],
"coverageFindingCount": detector_check["coverageFindingCount"],
"highSeverityCount": detector_check["highSeverityCount"],
"reportSha256": sha256_value(detector_report),
"errors": detector_check["errors"],
}
if not detector_check["ok"]:
result["status"] = "detector_invalid"
_write_result(output_dir, result)
return result
except RunnerTimeoutError as error:
result["results"][arm]["detector"] = {
"status": "timeout",
"findingCount": 0,
"coverageFindingCount": 0,
"highSeverityCount": 0,
"errors": [str(error)],
"rawOutputSha256": (
sha256_value(detector_path.read_text(encoding="utf-8"))
if detector_path.exists()
else None
),
}
result["status"] = "detector_timeout"
_write_result(output_dir, result)
return result
except RunnerInvocationError as error:
result["results"][arm]["detector"] = {
"status": "failed",
"findingCount": 0,
"coverageFindingCount": 0,
"highSeverityCount": 0,
"errors": [str(error)],
"rawOutputSha256": (
sha256_value(detector_path.read_text(encoding="utf-8"))
if detector_path.exists()
else None
),
}
result["status"] = "detector_failed"
_write_result(output_dir, result)
return result
except (RunnerOutputError, json.JSONDecodeError) as error:
result["results"][arm]["detector"] = {
"status": "invalid",
"findingCount": 0,
"coverageFindingCount": 0,
"highSeverityCount": 0,
"errors": [str(error)],
"rawOutputSha256": (
sha256_value(detector_path.read_text(encoding="utf-8"))
if detector_path.exists()
else None
),
}
result["status"] = "detector_invalid"
_write_result(output_dir, result)
return result
if detector_invalid or detector_blocked:
result["status"] = "detector_invalid" if detector_invalid else "detector_blocked"
result["ok"] = False
_write_result(output_dir, result)
return result
judge_runners = (judge_primary_bin or planner_bin, judge_secondary_bin or planner_bin)
if len(set(JUDGE_IDS)) != 2:
raise ReplayRunError("两个 judge 身份必须不同")
judge_reports: list[Mapping[str, Any]] = []
judge_invalid = False
reference_proxy = leakage_audit_config.get("targetFacts")
if not isinstance(reference_proxy, Mapping):
result["status"] = "judge_invalid"
result["ok"] = False
result["evaluation"] = {"profile": RUBRIC_PROFILE, "errors": ["缺少结构化 reference proxy"]}
_write_result(output_dir, result)
return result
result["status"] = "running_judge"
_write_result(output_dir, result)
for index, judge_id in enumerate(JUDGE_IDS):
ordered = assignments if index == 0 else list(reversed(assignments))
request = _judge_request(
judge_id=judge_id,
assignments=ordered,
as_of=as_of,
frozen_snapshot=frozen["snapshot"],
common_context=common_input,
reference_proxy=reference_proxy,
)
judge_path = output_dir / f"judge_{judge_id}.raw.json"
try:
judge_report = _invoke_structured_agent(
agent="judge",
request=request,
runner_bin=judge_runners[index],
model=model,
output_path=judge_path,
max_budget_usd=max_budget_usd,
identity=judge_id,
timeout_seconds=float(timeout_seconds),
)
candidate_ids = tuple(str(item["candidateId"]) for item in ordered)
judge_errors = validate_report(
judge_report,
expected_judge_id=judge_id,
expected_candidate_ids=candidate_ids,
)
if judge_errors:
result["status"] = "judge_invalid"
result["evaluation"] = {
"profile": RUBRIC_PROFILE,
"judgeIds": list(JUDGE_IDS),
"status": "invalid",
}
_write_result(output_dir, result)
return result
else:
judge_reports.append(judge_report)
except RunnerTimeoutError:
result["status"] = "judge_timeout"
result["evaluation"] = {
"profile": RUBRIC_PROFILE,
"judgeIds": list(JUDGE_IDS),
"status": "timeout",
}
_write_result(output_dir, result)
return result
except RunnerInvocationError:
result["status"] = "judge_failed"
result["evaluation"] = {
"profile": RUBRIC_PROFILE,
"judgeIds": list(JUDGE_IDS),
"status": "failed",
}
_write_result(output_dir, result)
return result
except (RunnerOutputError, json.JSONDecodeError):
result["status"] = "judge_invalid"
result["evaluation"] = {
"profile": RUBRIC_PROFILE,
"judgeIds": list(JUDGE_IDS),
"status": "invalid",
}
_write_result(output_dir, result)
return result
if judge_invalid or len(judge_reports) != 2:
result["status"] = "judge_invalid"
result["ok"] = False
result["evaluation"] = {
"profile": RUBRIC_PROFILE,
"judgeIds": list(JUDGE_IDS),
"status": "invalid",
}
_write_result(output_dir, result)
return result
result["evaluation"] = _aggregate_evaluation(assignments, judge_reports)
if not result["evaluation"]["stability"]["stable"]:
result["status"] = "judge_unstable"
result["ok"] = False
_write_result(output_dir, result)
return result
result["status"] = "completed"
result["ok"] = True
_write_result(output_dir, result)
return result
def _parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="运行 next_fine_outline_replay_v0")
parser.add_argument("--config", type=Path, required=True)
parser.add_argument("--output-dir", type=Path, required=True)
parser.add_argument("--mode", choices=("dry_run", "execute"), default="dry_run")
parser.add_argument("--planner-bin", default="claude")
parser.add_argument("--detector-bin")
parser.add_argument("--judge-primary-bin")
parser.add_argument("--judge-secondary-bin")
parser.add_argument("--model", default="opus")
parser.add_argument("--max-budget-usd", type=float, default=1.0)
parser.add_argument("--timeout-seconds", type=float, default=DEFAULT_TIMEOUT_SECONDS)
return parser.parse_args()
def main() -> int:
args = _parse_args()
result = run_replay(
_read_json(args.config),
args.output_dir,
mode=args.mode,
planner_bin=args.planner_bin,
detector_bin=args.detector_bin,
judge_primary_bin=args.judge_primary_bin,
judge_secondary_bin=args.judge_secondary_bin,
model=args.model,
max_budget_usd=args.max_budget_usd,
timeout_seconds=args.timeout_seconds,
)
return 0 if result["ok"] else 2
if __name__ == "__main__":
raise SystemExit(main())