819 lines
34 KiB
Python
819 lines
34 KiB
Python
#!/usr/bin/env python3
|
||
"""正文候选的单次收敛编排:机械审查、语义检查、授权终态与 CAS。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import hashlib
|
||
import os
|
||
import pathlib
|
||
import re
|
||
import sys
|
||
import tempfile
|
||
from dataclasses import dataclass
|
||
from typing import Any, Callable, Mapping, Protocol, Sequence
|
||
|
||
SCRIPT_DIR = pathlib.Path(__file__).resolve().parent
|
||
SKILLS_DIR = SCRIPT_DIR.parents[1]
|
||
DETECT_DIR = SKILLS_DIR / "check-content-consistency" / "scripts"
|
||
READ_CONTEXT_DIR = SKILLS_DIR / "assemble-context" / "scripts"
|
||
for path in (DETECT_DIR, READ_CONTEXT_DIR):
|
||
if str(path) not in sys.path:
|
||
sys.path.insert(0, str(path))
|
||
|
||
from check_writer_candidate import check_writer_candidate # noqa: E402
|
||
from writer_contract import ( # noqa: E402
|
||
ContractError,
|
||
canonical_json,
|
||
build_writer_creative_input,
|
||
validate_writer_context,
|
||
validate_writer_output,
|
||
)
|
||
|
||
|
||
_HASH_PATTERN = re.compile(r"^sha256:[0-9a-f]{64}$")
|
||
|
||
|
||
class PipelineError(RuntimeError):
|
||
"""携带稳定失败码和终态结果的编排错误。"""
|
||
|
||
def __init__(
|
||
self,
|
||
code: str,
|
||
message: str,
|
||
*,
|
||
details: Mapping[str, Any] | None = None,
|
||
result: Mapping[str, Any] | None = None,
|
||
):
|
||
super().__init__(message)
|
||
self.code = code
|
||
self.details = dict(details or {})
|
||
self.result = dict(result or {})
|
||
self.acceptance_eligible = False
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class CasToken:
|
||
"""一次候选状态的不可变 CAS 身份。"""
|
||
|
||
run_id: str
|
||
attempt: int
|
||
candidate_version: int
|
||
state: str
|
||
revision: int
|
||
|
||
|
||
class CasStateStore(Protocol):
|
||
"""候选存储必须实现的最小 CAS 接口。"""
|
||
|
||
def create(self, run_id: str, *, attempt: int, candidate_version: int) -> CasToken: ...
|
||
|
||
def transition(self, expected: CasToken, target_state: str) -> CasToken | None: ...
|
||
|
||
def start_next(
|
||
self, expected: CasToken, *, attempt: int, candidate_version: int
|
||
) -> CasToken | None: ...
|
||
|
||
def latest(self, run_id: str) -> CasToken | None: ...
|
||
|
||
|
||
class SemanticDetector(Protocol):
|
||
"""语义 detector 的明确边界;实现方不得读写状态或候选文件。"""
|
||
|
||
def __call__(
|
||
self,
|
||
context: Mapping[str, Any],
|
||
candidate: Mapping[str, Any],
|
||
mechanical_report: Mapping[str, Any],
|
||
) -> Mapping[str, Any]: ...
|
||
|
||
|
||
class InMemoryCasStateStore:
|
||
"""仅供实验和测试使用的线程外内存 CAS fake。"""
|
||
|
||
_ALLOWED = {
|
||
"DRAFT": frozenset({"CHECKING"}),
|
||
"CHECKING": frozenset({"PASSED", "REJECTED"}),
|
||
"PASSED": frozenset(),
|
||
"REJECTED": frozenset(),
|
||
}
|
||
|
||
def __init__(self) -> None:
|
||
self._latest: dict[str, CasToken] = {}
|
||
|
||
def create(self, run_id: str, *, attempt: int, candidate_version: int) -> CasToken:
|
||
"""只允许为未存在的 run 创建第一条 DRAFT。"""
|
||
|
||
if run_id in self._latest:
|
||
raise PipelineError("CAS_CONFLICT", "run 已存在,不能重复创建初始状态")
|
||
token = CasToken(run_id, attempt, candidate_version, "DRAFT", 1)
|
||
self._latest[run_id] = token
|
||
return token
|
||
|
||
def transition(self, expected: CasToken, target_state: str) -> CasToken | None:
|
||
"""仅当最新 token 完全匹配且转换合法时更新。"""
|
||
|
||
current = self._latest.get(expected.run_id)
|
||
if current != expected or target_state not in self._ALLOWED.get(expected.state, frozenset()):
|
||
return None
|
||
updated = CasToken(
|
||
expected.run_id,
|
||
expected.attempt,
|
||
expected.candidate_version,
|
||
target_state,
|
||
expected.revision + 1,
|
||
)
|
||
self._latest[expected.run_id] = updated
|
||
return updated
|
||
|
||
def start_next(
|
||
self, expected: CasToken, *, attempt: int, candidate_version: int
|
||
) -> CasToken | None:
|
||
"""只允许从最新 REJECTED 创建单调递增的新 DRAFT。"""
|
||
|
||
current = self._latest.get(expected.run_id)
|
||
if (
|
||
current != expected
|
||
or expected.state != "REJECTED"
|
||
or attempt <= expected.attempt
|
||
or candidate_version <= expected.candidate_version
|
||
):
|
||
return None
|
||
updated = CasToken(
|
||
expected.run_id,
|
||
attempt,
|
||
candidate_version,
|
||
"DRAFT",
|
||
expected.revision + 1,
|
||
)
|
||
self._latest[expected.run_id] = updated
|
||
return updated
|
||
|
||
def latest(self, run_id: str) -> CasToken | None:
|
||
"""返回 run 的最新不可变状态。"""
|
||
|
||
return self._latest.get(run_id)
|
||
|
||
|
||
def atomic_write_json(path: str | pathlib.Path, value: Mapping[str, Any]) -> None:
|
||
"""先写同目录临时文件并 fsync,再以原子替换发布正式结果。"""
|
||
|
||
target = pathlib.Path(path)
|
||
target.parent.mkdir(parents=True, exist_ok=True)
|
||
descriptor, temporary_name = tempfile.mkstemp(
|
||
prefix=f".{target.name}.", suffix=".tmp", dir=target.parent
|
||
)
|
||
temporary = pathlib.Path(temporary_name)
|
||
try:
|
||
with os.fdopen(descriptor, "w", encoding="utf-8", newline="\n") as handle:
|
||
handle.write(canonical_json(value))
|
||
handle.write("\n")
|
||
handle.flush()
|
||
os.fsync(handle.fileno())
|
||
os.replace(str(temporary), str(target))
|
||
except BaseException:
|
||
# replace 成功后临时路径已不存在;失败时只清理本次临时文件。
|
||
try:
|
||
temporary.unlink(missing_ok=True)
|
||
finally:
|
||
raise
|
||
|
||
|
||
def _cas_or_fail(token: CasToken | None, action: str) -> CasToken:
|
||
"""把所有 CAS 竞争统一为稳定失败码。"""
|
||
|
||
if token is None:
|
||
raise PipelineError("CAS_CONFLICT", f"状态 CAS 失败: {action}")
|
||
return token
|
||
|
||
|
||
def _rebind_semantic_report_hash(report: dict[str, Any]) -> dict[str, Any]:
|
||
"""报告内容变更后重算 hash,保持落库自洽。"""
|
||
|
||
body = {key: value for key, value in report.items() if key != "reportSha256"}
|
||
report["reportSha256"] = "sha256:" + hashlib.sha256(
|
||
canonical_json(body).encode("utf-8")
|
||
).hexdigest()
|
||
return report
|
||
|
||
|
||
def _promote_unfillable_gaps_to_settings(report: Mapping[str, Any]) -> dict[str, Any]:
|
||
"""检索零命中的缺口升格为新设定提案,status 回到 passed,交给人闸。
|
||
|
||
不改候选正文。新设定是否进入正典由人决定;检测只拦与既有正典的冲突。
|
||
"""
|
||
|
||
released = dict(report)
|
||
gaps = [dict(item) for item in (released.get("evidenceGaps") or []) if isinstance(item, Mapping)]
|
||
settings = [
|
||
dict(item) for item in (released.get("newSettingCandidates") or []) if isinstance(item, Mapping)
|
||
]
|
||
known_quotes = {str(item.get("candidateQuote") or "") for item in settings}
|
||
for gap in gaps:
|
||
quote = str(gap.get("candidateQuote") or "")
|
||
if quote and quote in known_quotes:
|
||
continue
|
||
settings.append(
|
||
{
|
||
"settingId": f"promoted-{gap.get('gapId') or len(settings) + 1}",
|
||
"factType": "unspecified",
|
||
"text": str(gap.get("reason") or gap.get("query") or "正典检索未命中的新设定提案"),
|
||
"candidateSha256": gap.get("candidateSha256"),
|
||
"candidateQuote": quote,
|
||
"startCodePoint": gap.get("startCodePoint"),
|
||
"endCodePoint": gap.get("endCodePoint"),
|
||
}
|
||
)
|
||
known_quotes.add(quote)
|
||
released["newSettingCandidates"] = settings
|
||
released["evidenceGaps"] = []
|
||
if released.get("status") == "needs_evidence":
|
||
released["status"] = "passed"
|
||
return _rebind_semantic_report_hash(released)
|
||
|
||
|
||
def _validate_semantic_report(
|
||
report: Mapping[str, Any], context: Mapping[str, Any], candidate: Mapping[str, Any]
|
||
) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]:
|
||
"""校验 SemanticDetection v3 的绑定、状态和补证缺口。"""
|
||
|
||
if not isinstance(report, Mapping):
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "语义 detector 必须返回对象")
|
||
if report.get("schemaVersion") != "semantic-detection-v3":
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "语义 detector 报告版本非法")
|
||
if report.get("status") not in {"passed", "failed", "needs_evidence"}:
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "语义 detector status 非法")
|
||
expected_fields = {
|
||
"schemaVersion", "runId", "sampleId", "opaqueArmId", "inputSha256",
|
||
"candidateVersion", "candidateSha256", "contextSnapshotSha256",
|
||
"modelReceiptSha256", "claims", "findings", "assertionVerdicts",
|
||
"hardConstraintVerdicts", "newSettingCandidates", "evidenceGaps", "status",
|
||
"reportSha256",
|
||
}
|
||
if set(report) != expected_fields:
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "语义 detector 报告字段不完整或含额外字段")
|
||
for field in ("sampleId", "opaqueArmId"):
|
||
if not isinstance(report.get(field), str) or not report[field]:
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", f"语义 detector {field} 非法")
|
||
for field in ("inputSha256", "modelReceiptSha256"):
|
||
value = report.get(field)
|
||
if not isinstance(value, str) or not _HASH_PATTERN.fullmatch(value):
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", f"语义 detector {field} 非法")
|
||
if (
|
||
report.get("runId") != context["runId"]
|
||
or report.get("candidateVersion") != candidate["candidateVersion"]
|
||
or report.get("candidateSha256") != candidate["candidateSha256"]
|
||
or report.get("contextSnapshotSha256") != context["contextSnapshot"]["contextSha256"]
|
||
):
|
||
raise PipelineError("SEMANTIC_DETECTOR_STALE", "语义 detector 结果未绑定当前候选")
|
||
required_arrays = (
|
||
"claims", "findings", "assertionVerdicts", "hardConstraintVerdicts",
|
||
"newSettingCandidates", "evidenceGaps",
|
||
)
|
||
if any(
|
||
not isinstance(report.get(field), list)
|
||
or any(not isinstance(item, Mapping) for item in report[field])
|
||
for field in required_arrays
|
||
):
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "语义 detector 数组字段非法")
|
||
report_hash = report.get("reportSha256")
|
||
expected_hash = "sha256:" + hashlib.sha256(
|
||
canonical_json({key: value for key, value in report.items() if key != "reportSha256"}).encode("utf-8")
|
||
).hexdigest()
|
||
if report_hash != expected_hash:
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "语义 detector 报告哈希非法")
|
||
gaps = [dict(item) for item in report["evidenceGaps"]]
|
||
if any(
|
||
not isinstance(item.get("gapId"), str)
|
||
or not isinstance(item.get("query"), str)
|
||
or not isinstance(item.get("reason"), str)
|
||
or item.get("priority") not in {"high", "medium", "low"}
|
||
for item in gaps
|
||
):
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "evidenceGaps 结构非法")
|
||
failures: list[dict[str, Any]] = []
|
||
for item in report["findings"]:
|
||
if item.get("severity") == "high":
|
||
failures.append({"code": f"SEMANTIC_{str(item.get('category') or 'FINDING').upper()}", "message": str(item.get("message") or "语义冲突")})
|
||
for field in ("assertionVerdicts", "hardConstraintVerdicts"):
|
||
for item in report[field]:
|
||
if item.get("verdict") == "fail":
|
||
failures.append({"code": "SEMANTIC_VERDICT_FAILED", "message": f"{field} 未通过"})
|
||
for item in report["claims"]:
|
||
if item.get("coverageState") == "conflict":
|
||
failures.append({"code": "SEMANTIC_CLAIM_CONFLICT", "message": "detector 断言与冻结事实冲突"})
|
||
if report["status"] == "passed" and (failures or gaps):
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "passed 报告不得含阻塞项或证据缺口")
|
||
if report["status"] == "needs_evidence" and not gaps:
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "needs_evidence 必须携带 evidenceGaps")
|
||
if report["status"] == "failed" and not failures:
|
||
raise PipelineError("SEMANTIC_DETECTOR_INVALID", "failed 报告必须含阻塞语义项")
|
||
return failures, gaps
|
||
|
||
|
||
def _validate_mechanical_report(
|
||
report: Any, candidate: Mapping[str, Any]
|
||
) -> list[dict[str, Any]]:
|
||
"""严格校验机械 detector 报告版本、候选绑定和通过状态。"""
|
||
|
||
if not isinstance(report, Mapping):
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector 必须返回对象")
|
||
if report.get("schemaVersion") != "writer-detector-report-v1":
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector 报告版本非法")
|
||
expected = {
|
||
"runId": candidate.get("runId"),
|
||
"attempt": candidate.get("attempt"),
|
||
"candidateVersion": candidate.get("candidateVersion"),
|
||
"candidateSha256": candidate.get("candidateSha256"),
|
||
}
|
||
if any(report.get(field) != value for field, value in expected.items()):
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector 报告未绑定当前候选")
|
||
passed = report.get("passed")
|
||
failures = report.get("blockingFailures")
|
||
if not isinstance(passed, bool):
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector passed 必须是布尔值")
|
||
if not isinstance(failures, list) or any(not isinstance(item, Mapping) for item in failures):
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector blockingFailures 非法")
|
||
if any(not isinstance(item.get("code"), str) or not item["code"] for item in failures):
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector 阻塞项缺少稳定 code")
|
||
if passed != (not failures):
|
||
raise PipelineError("MECHANICAL_DETECTOR_INVALID", "机械 detector 通过状态与阻塞项矛盾")
|
||
return [dict(item) for item in failures]
|
||
|
||
|
||
def _terminal_result(
|
||
*,
|
||
run_id: str,
|
||
status: str,
|
||
token: CasToken,
|
||
evidence_request_count: int,
|
||
rewrite_count: int,
|
||
trace: Sequence[Mapping[str, Any]],
|
||
candidate: Mapping[str, Any] | None,
|
||
failure_code: str | None = None,
|
||
) -> dict[str, Any]:
|
||
"""生成可原子落盘的稳定终态结果。"""
|
||
|
||
result = {
|
||
"schemaVersion": "writer-pipeline-result-v1",
|
||
"runId": run_id,
|
||
"status": status,
|
||
"attempt": token.attempt,
|
||
"candidateVersion": token.candidate_version,
|
||
"candidateSha256": candidate.get("candidateSha256") if candidate else None,
|
||
"failureCode": failure_code,
|
||
"evidenceRequestCount": evidence_request_count,
|
||
"rewriteCount": rewrite_count,
|
||
"trace": [dict(item) for item in trace],
|
||
}
|
||
# 只有通过全部合同和审查门的终态才暴露可进入 Shadow 的最终候选。
|
||
if status == "PASSED" and candidate is not None:
|
||
result["candidateArtifact"] = dict(candidate)
|
||
return result
|
||
|
||
|
||
def _publish_if_requested(
|
||
result_path: str | pathlib.Path | None, result: Mapping[str, Any]
|
||
) -> None:
|
||
"""仅在调用方显式给出路径时发布终态文件。"""
|
||
|
||
if result_path is not None:
|
||
atomic_write_json(result_path, result)
|
||
|
||
|
||
def _terminal_failure(
|
||
*,
|
||
code: str,
|
||
message: str,
|
||
token: CasToken,
|
||
state_store: CasStateStore,
|
||
run_id: str,
|
||
evidence_request_count: int,
|
||
rewrite_count: int,
|
||
trace: Sequence[Mapping[str, Any]],
|
||
candidate: Mapping[str, Any] | None,
|
||
result_path: str | pathlib.Path | None,
|
||
details: Mapping[str, Any] | None = None,
|
||
) -> PipelineError:
|
||
"""将已进入检查的失败统一收敛为可发布的 REJECTED 终态。"""
|
||
|
||
if token.state != "CHECKING":
|
||
raise PipelineError("CAS_CONFLICT", f"失败收敛不接受状态: {token.state}")
|
||
rejected = _cas_or_fail(
|
||
state_store.transition(token, "REJECTED"), "CHECKING -> REJECTED"
|
||
)
|
||
terminal_trace = [dict(item) for item in trace]
|
||
terminal_trace.append(
|
||
{
|
||
"attempt": rejected.attempt,
|
||
"candidateVersion": rejected.candidate_version,
|
||
"candidateSha256": candidate.get("candidateSha256") if candidate else None,
|
||
"status": "pipeline_failed",
|
||
"failureCodes": [code],
|
||
}
|
||
)
|
||
result = _terminal_result(
|
||
run_id=run_id,
|
||
status="REJECTED",
|
||
token=rejected,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=terminal_trace,
|
||
candidate=candidate,
|
||
failure_code=code,
|
||
)
|
||
_publish_if_requested(result_path, result)
|
||
return PipelineError(code, message, details=details, result=result)
|
||
|
||
|
||
def run_writer_pipeline(
|
||
*,
|
||
context: Mapping[str, Any],
|
||
requirements: Mapping[str, Any],
|
||
writer: Callable[[Mapping[str, Any], int], Mapping[str, Any]],
|
||
evidence_provider: Callable[
|
||
[Mapping[str, Any], list[dict[str, Any]], int], Mapping[str, Any]
|
||
],
|
||
semantic_detector: SemanticDetector,
|
||
state_store: CasStateStore,
|
||
result_path: str | pathlib.Path | None = None,
|
||
initial_candidate_version: int = 1,
|
||
) -> dict[str, Any]:
|
||
"""正文候选单次收敛:写手 -> 机械门 -> 语义检查,全分支终态化。
|
||
|
||
语义检查发现证据缺口时先做确定性检索(evidence_provider,只读):
|
||
零命中升格为新设定提案进人闸;命中则收敛为 AUTHORIZATION_REQUIRED 终态,
|
||
携带缺口报告与重组上下文哈希,由编排方在人授权后发起新运行继续重写。
|
||
补证/重写业务次数不设硬上限,管控归人;授权门只挡模型重调用,不挡检索。
|
||
|
||
initial_candidate_version 给生产编排接续既有版本号:候选表对
|
||
(作品, 章, candidate_version) 唯一,同一章重跑必须从「已有最大版本+1」起,
|
||
否则与上一轮留库的被拒候选撞版本。
|
||
"""
|
||
|
||
if isinstance(initial_candidate_version, bool) or not isinstance(initial_candidate_version, int) or initial_candidate_version < 1:
|
||
raise PipelineError("PIPELINE_INPUT_INVALID", "initial_candidate_version 必须是 >=1 的整数")
|
||
try:
|
||
current_context = validate_writer_context(context)
|
||
except ContractError as exc:
|
||
raise PipelineError("WRITER_CONTEXT_INVALID", str(exc)) from exc
|
||
run_id = current_context["runId"]
|
||
candidate_version = initial_candidate_version
|
||
evidence_request_count = 0
|
||
rewrite_count = 0
|
||
trace: list[dict[str, Any]] = []
|
||
draft = state_store.create(
|
||
run_id,
|
||
attempt=current_context["attempt"],
|
||
candidate_version=candidate_version,
|
||
)
|
||
|
||
while True:
|
||
checking = _cas_or_fail(state_store.transition(draft, "CHECKING"), "DRAFT -> CHECKING")
|
||
try:
|
||
raw_candidate = writer(current_context, candidate_version)
|
||
except PipelineError as exc:
|
||
raise _terminal_failure(
|
||
code=exc.code,
|
||
message=str(exc),
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=None,
|
||
result_path=result_path,
|
||
details=exc.details,
|
||
) from exc
|
||
except Exception as exc:
|
||
raise _terminal_failure(
|
||
code="WRITER_FAILED",
|
||
message="写手 adapter 调用失败",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=None,
|
||
result_path=result_path,
|
||
) from exc
|
||
if not isinstance(raw_candidate, Mapping):
|
||
raise _terminal_failure(
|
||
code="WRITER_OUTPUT_INVALID",
|
||
message="写手必须返回对象",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=None,
|
||
result_path=result_path,
|
||
)
|
||
candidate = dict(raw_candidate)
|
||
if (
|
||
candidate.get("runId") != run_id
|
||
or candidate.get("attempt") != current_context["attempt"]
|
||
or candidate.get("candidateVersion") != candidate_version
|
||
):
|
||
raise _terminal_failure(
|
||
code="WRITER_OUTPUT_STALE",
|
||
message="写手结果未绑定当前 attempt 或 candidateVersion",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
)
|
||
try:
|
||
candidate = validate_writer_output(candidate)
|
||
except ContractError as exc:
|
||
# 合同错误作为可定位审查失败进入机械门定位,而不是绕过状态机。
|
||
try:
|
||
mechanical_report = check_writer_candidate(current_context, candidate, requirements)
|
||
candidate_failures = _validate_mechanical_report(mechanical_report, candidate)
|
||
except PipelineError as detector_exc:
|
||
raise _terminal_failure(
|
||
code=detector_exc.code,
|
||
message=str(detector_exc),
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
details=detector_exc.details,
|
||
) from detector_exc
|
||
except Exception as detector_exc:
|
||
raise _terminal_failure(
|
||
code="MECHANICAL_DETECTOR_FAILED",
|
||
message="机械 detector 调用失败",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
) from detector_exc
|
||
if not candidate_failures:
|
||
candidate_failures = [{"code": "OUTPUT_CONTRACT_INVALID", "message": str(exc)}]
|
||
else:
|
||
candidate_failures = []
|
||
mechanical_report = {}
|
||
|
||
if not mechanical_report:
|
||
try:
|
||
mechanical_report = check_writer_candidate(current_context, candidate, requirements)
|
||
candidate_failures = _validate_mechanical_report(mechanical_report, candidate)
|
||
except PipelineError as detector_exc:
|
||
raise _terminal_failure(
|
||
code=detector_exc.code,
|
||
message=str(detector_exc),
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
details=detector_exc.details,
|
||
) from detector_exc
|
||
except Exception as exc:
|
||
raise _terminal_failure(
|
||
code="MECHANICAL_DETECTOR_FAILED",
|
||
message="机械 detector 调用失败",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
) from exc
|
||
semantic_report: Mapping[str, Any] | None = None
|
||
evidence_gaps: list[dict[str, Any]] = []
|
||
if not candidate_failures:
|
||
try:
|
||
semantic_report = semantic_detector(current_context, candidate, mechanical_report)
|
||
semantic_failures, evidence_gaps = _validate_semantic_report(
|
||
semantic_report, current_context, candidate
|
||
)
|
||
candidate_failures.extend(semantic_failures)
|
||
except PipelineError as exc:
|
||
# 语义阶段终态失败也保留机械门审查痕迹,审计链不在 detector 失败处断裂。
|
||
trace.append(
|
||
{
|
||
"attempt": checking.attempt,
|
||
"candidateVersion": checking.candidate_version,
|
||
"candidateSha256": candidate.get("candidateSha256"),
|
||
"mechanicalPassed": mechanical_report.get("passed", False),
|
||
"semanticStatus": "not_run",
|
||
"failureCodes": [exc.code],
|
||
"mechanicalReport": dict(mechanical_report),
|
||
"semanticReport": None,
|
||
}
|
||
)
|
||
raise _terminal_failure(
|
||
code=exc.code,
|
||
message=str(exc),
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
details=exc.details,
|
||
) from exc
|
||
except Exception as exc:
|
||
trace.append(
|
||
{
|
||
"attempt": checking.attempt,
|
||
"candidateVersion": checking.candidate_version,
|
||
"candidateSha256": candidate.get("candidateSha256"),
|
||
"mechanicalPassed": mechanical_report.get("passed", False),
|
||
"semanticStatus": "not_run",
|
||
"failureCodes": ["SEMANTIC_DETECTOR_FAILED"],
|
||
"mechanicalReport": dict(mechanical_report),
|
||
"semanticReport": None,
|
||
}
|
||
)
|
||
raise _terminal_failure(
|
||
code="SEMANTIC_DETECTOR_FAILED",
|
||
message="语义 detector 调用失败",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
) from exc
|
||
if evidence_gaps:
|
||
# 缺口先做确定性检索(只读、无模型开销);结果决定收敛方向:
|
||
# 零命中 = 新设定提案,升格 passed 进人闸;命中 = 需要模型重写,
|
||
# 停在授权终态等人授权,由编排方以新运行继续,不自动重写。
|
||
next_attempt = current_context["attempt"] + 1
|
||
previous_snapshot = current_context["contextSnapshot"]["contextSha256"]
|
||
previous_fact_ids = {
|
||
str(item.get("evidenceId"))
|
||
for item in (current_context.get("factEvidence") or [])
|
||
if isinstance(item, Mapping)
|
||
}
|
||
try:
|
||
provided = evidence_provider(
|
||
current_context,
|
||
[dict(item) for item in evidence_gaps],
|
||
next_attempt,
|
||
)
|
||
next_context = validate_writer_context(provided)
|
||
except Exception as exc:
|
||
raise _terminal_failure(
|
||
code="REASSEMBLED_CONTEXT_INVALID",
|
||
message=str(exc),
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
) from exc
|
||
next_fact_ids = {
|
||
str(item.get("evidenceId"))
|
||
for item in (next_context.get("factEvidence") or [])
|
||
if isinstance(item, Mapping)
|
||
}
|
||
if not (next_fact_ids - previous_fact_ids):
|
||
# 正典检索零命中 = 缺口是新设定提案,不是可补的检索债。
|
||
# 无语义冲突时升格为 passed 交人闸;有冲突则不升格、走失败关闭。
|
||
if not candidate_failures and semantic_report is not None:
|
||
semantic_report = _promote_unfillable_gaps_to_settings(semantic_report)
|
||
evidence_gaps = []
|
||
else:
|
||
if (
|
||
next_context["runId"] != run_id
|
||
or next_context["attempt"] != next_attempt
|
||
or next_context["contextSnapshot"]["contextSha256"] == previous_snapshot
|
||
):
|
||
raise _terminal_failure(
|
||
code="REASSEMBLED_CONTEXT_INVALID",
|
||
message="补证必须形成同 run 的新 attempt 与新上下文快照",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
)
|
||
trace.append({
|
||
"attempt": checking.attempt,
|
||
"candidateVersion": checking.candidate_version,
|
||
"candidateSha256": candidate["candidateSha256"],
|
||
"status": "needs_authorization",
|
||
"gapIds": [item["gapId"] for item in evidence_gaps],
|
||
"mechanicalPassed": mechanical_report.get("passed", False),
|
||
"mechanicalReport": dict(mechanical_report),
|
||
"semanticReport": dict(semantic_report) if semantic_report else None,
|
||
})
|
||
raise _terminal_failure(
|
||
code="AUTHORIZATION_REQUIRED",
|
||
message="补证检索命中,重写需要人授权(以新运行继续)",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
details={
|
||
"evidenceGaps": [dict(item) for item in evidence_gaps],
|
||
"mechanicalPassed": mechanical_report.get("passed", False),
|
||
"semanticStatus": semantic_report.get("status") if semantic_report else None,
|
||
"nextAttempt": next_attempt,
|
||
"reassembledContextSha256": next_context["contextSnapshot"]["contextSha256"],
|
||
},
|
||
)
|
||
trace.append(
|
||
{
|
||
"attempt": checking.attempt,
|
||
"candidateVersion": checking.candidate_version,
|
||
"candidateSha256": candidate.get("candidateSha256"),
|
||
"mechanicalPassed": mechanical_report.get("passed", False),
|
||
"semanticStatus": semantic_report.get("status") if semantic_report else "not_run",
|
||
"failureCodes": [item.get("code") for item in candidate_failures],
|
||
"mechanicalReport": dict(mechanical_report),
|
||
"semanticReport": dict(semantic_report) if semantic_report else None,
|
||
}
|
||
)
|
||
if not candidate_failures:
|
||
passed = _cas_or_fail(
|
||
state_store.transition(checking, "PASSED"), "CHECKING -> PASSED"
|
||
)
|
||
result = _terminal_result(
|
||
run_id=run_id,
|
||
status="PASSED",
|
||
token=passed,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
)
|
||
_publish_if_requested(result_path, result)
|
||
return result
|
||
|
||
if semantic_report is not None:
|
||
raise _terminal_failure(
|
||
code="SEMANTIC_DETECTION_REJECTED",
|
||
message="语义检测报告阻断当前候选且没有可补证缺口",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
details={"failureCodes": [item.get("code") for item in candidate_failures]},
|
||
)
|
||
raise _terminal_failure(
|
||
code="MECHANICAL_DETECTION_REJECTED",
|
||
message="机械检测阻断当前候选且没有可形成新创作输入的证据缺口",
|
||
token=checking,
|
||
state_store=state_store,
|
||
run_id=run_id,
|
||
evidence_request_count=evidence_request_count,
|
||
rewrite_count=rewrite_count,
|
||
trace=trace,
|
||
candidate=candidate,
|
||
result_path=result_path,
|
||
details={"failureCodes": [item.get("code") for item in candidate_failures]},
|
||
)
|
||
|
||
|
||
__all__ = [
|
||
"CasToken",
|
||
"CasStateStore",
|
||
"InMemoryCasStateStore",
|
||
"PipelineError",
|
||
"SemanticDetector",
|
||
"atomic_write_json",
|
||
"run_writer_pipeline",
|
||
]
|