"""可信单进程教学 Harness：真实文件操作、有限调用、事实校验和恢复。

安全前提：目录由当前用户独占，没有不可信并发进程修改文件。
工具不接受任意路径，也没有 shell；这仍不是 OS 沙箱或生产身份认证。
状态与产物分开替换，没有跨文件事务；恢复时重新核对候选。
"""
import json
import os
from pathlib import Path
from contracts import ORDER, POLICY, REPORT_SCHEMA, digest
from verifier import verify

TOOLS = [
    {"name": "read_source", "description": "读取本任务获准的订单和政策快照。",
     "input_schema": {"type": "object", "properties": {}, "additionalProperties": False}},
    {"name": "write_report", "description": "保存候选报告；写入成功不代表事实校验通过。",
     "input_schema": {"type": "object", "properties": {"report": REPORT_SCHEMA},
                      "required": ["report"], "additionalProperties": False}},
    {"name": "verify_report", "description": "使用可信输入重新验证磁盘中的候选。",
     "input_schema": {"type": "object", "properties": {}, "additionalProperties": False}},
    {"name": "finish", "description": "再次核验候选，并导出同一份内容；失败时不能完成。",
     "input_schema": {"type": "object", "properties": {}, "additionalProperties": False}},
]

class ToolRejected(ValueError):
    """可交给模型理解的受控拒绝；磁盘异常等由上层终止任务。"""


def read_json(path):
    # 限制载入大小，避免意外把超大文件读入内存；不防御敌对路径竞态。
    if path.is_symlink():
        raise ToolRejected("symbolic_link_not_allowed")
    with path.open("rb") as stream:
        raw = stream.read(100_001)
    if len(raw) > 100_000:
        raise ToolRejected("file_too_large")
    return json.loads(raw)


def atomic_json(path, value):
    """同目录写入再替换；单写入者约定下避免读取半份 JSON。"""
    text = json.dumps(value, ensure_ascii=False, indent=2)
    if len(text.encode("utf-8")) > 100_000:
        raise ToolRejected("output_too_large")
    if path.is_symlink():
        raise ToolRejected("symbolic_link_not_allowed")
    temporary = path.with_suffix(".pending")
    if temporary.is_symlink():
        raise ToolRejected("symbolic_link_not_allowed")
    with temporary.open("w", encoding="utf-8") as stream:
        stream.write(text)
        stream.flush()
        os.fsync(stream.fileno())
    os.replace(temporary, path)


class Runtime:
    def __init__(self, root):
        self.root = Path(root).absolute()
        # 新任务只能使用不存在的目录；不覆盖已有用户目录。
        if not self.root.exists():
            self.root.mkdir(parents=True)
            (self.root / "inputs").mkdir()
            (self.root / "outputs").mkdir()
            (self.root / "export").mkdir()
            source = {"order": ORDER, "policy": POLICY}
            atomic_json(self.root / "inputs/source.json", source)
            atomic_json(self.root / "manifest.json", {
                "source_hash": digest(source), "harness_version": "teaching-v1",
                "output": "outputs/report.json", "identity": "shop-a/user-417",
            })
            atomic_json(self.root / "state.json", {
                "phase": "ready", "calls": 0, "attempts": 0, "source_hash": digest(source),
            })
        # 已有目录必须是完整任务目录；部分初始化失败要求人工检查，不能偷偷重置。
        for path in (self.root, self.root/'inputs', self.root/'outputs', self.root/'export'):
            if path.is_symlink() or not path.is_dir():
                raise ToolRejected("invalid_workspace")
        self.source = read_json(self.root / "inputs/source.json")
        self.state = read_json(self.root / "state.json")
        manifest = read_json(self.root / "manifest.json")
        if manifest["harness_version"] != "teaching-v1":
            raise ToolRejected("harness_version_changed")
        if digest(self.source) != manifest["source_hash"] or digest(self.source) != self.state["source_hash"]:
            raise ToolRejected("source_changed_requires_review")
        # 固定身份只用于说明授权位置；实际系统从认证会话注入，并在动作前重新验证。
        order = self.source["order"]
        if (order["tenant_id"], order["user_id"], order["order_id"]) != ("shop-a", "user-417", "A1042"):
            raise ToolRejected("resource_not_authorized")
        self.allowed = {item["name"] for item in TOOLS}

    @property
    def candidate(self):
        return self.root / "outputs/report.json"

    def save(self):
        atomic_json(self.root / "state.json", self.state)

    def require_allowed(self, name):
        if type(name) is not str or name not in self.allowed:
            raise ToolRejected("tool_not_allowed")

    def consume_call(self):
        if self.state["calls"] >= 12:
            raise ToolRejected("task_call_budget_exhausted")
        # 执行前保存额度；崩溃可能消耗一次未完成调用，但不会借恢复重置预算。
        self.state["calls"] += 1
        self.save()

    def check(self):
        if not self.candidate.exists():
            return {"passed": False, "errors": [{"check": "candidate_missing"}]}
        return verify(read_json(self.candidate), self.source)

    def dispatch(self, name, arguments):
        if type(arguments) is not dict:
            raise ToolRejected("arguments_must_be_object")
        expected = {"report"} if name == "write_report" else set()
        if set(arguments) != expected:
            raise ToolRejected("unexpected_arguments")
        if name == "read_source":
            return self.source
        if name == "write_report":
            if self.state["attempts"] >= 2:
                raise ToolRejected("repair_budget_exhausted")
            # 保存尝试次数后写文件；不把事实检查放在这里，便于展示候选与验证的区别。
            self.state.update(phase="writing", attempts=self.state["attempts"] + 1)
            self.save()
            atomic_json(self.candidate, arguments["report"])
            self.state["phase"] = "candidate"
            self.save()
            return {"candidate": "outputs/report.json", "verified": False}
        if name == "verify_report":
            result = self.check()
            self.state.update(phase="verified" if result["passed"] else "needs_repair", verification=result)
            self.save()
            return result
        if name == "finish":
            # 读取一次候选并核验，导出相同对象；不信任之前保存的 passed 标志。
            report = read_json(self.candidate) if self.candidate.exists() else None
            result = verify(report, self.source)
            if not result["passed"]:
                return result
            atomic_json(self.root / "export/report.json", report)
            self.state.update(phase="completed", artifact_hash=digest(report), verification=result)
            self.save()
            return {"completed": True, "artifact": "export/report.json", "hash": digest(report)}
        raise ToolRejected("tool_not_implemented")

    def call(self, name, arguments):
        self.require_allowed(name)
        if self.state["phase"] == "completed":
            raise ToolRejected("completed_task_is_immutable")
        self.consume_call()
        return self.dispatch(name, arguments)

    def inspect_completed(self):
        """恢复已完成任务时核对导出内容；不允许只凭状态字符串报成功。"""
        artifact = read_json(self.root / "export/report.json")
        if digest(artifact) != self.state["artifact_hash"] or not verify(artifact, self.source)["passed"]:
            raise ToolRejected("completed_artifact_changed")
        return {"completed": True, "artifact": str(self.root / "export/report.json")}
