"""单机教学 Runtime。任务事务与模拟支付事务分开；不提供生产认证或分布式调度。"""
from contextlib import contextmanager
import json
import sqlite3
import time
from pathlib import Path
from payment_simulator import PaymentSimulator, fingerprint


class RuntimeStore:
    def __init__(self, directory: Path):
        directory.mkdir(parents=True, exist_ok=True)
        self.path = directory / "tasks.sqlite"
        self.payment = PaymentSimulator(directory / "payment.sqlite")
        with self.connect() as db:
            db.executescript("""
                CREATE TABLE IF NOT EXISTS tasks (
                    id TEXT PRIMARY KEY, body TEXT NOT NULL);
                CREATE TABLE IF NOT EXISTS events (
                    task_id TEXT NOT NULL, seq INTEGER NOT NULL,
                    kind TEXT NOT NULL, payload TEXT NOT NULL,
                    PRIMARY KEY (task_id, seq));
                CREATE TABLE IF NOT EXISTS outbox (
                    id TEXT PRIMARY KEY, task_id TEXT NOT NULL,
                    payload TEXT NOT NULL, sent INTEGER NOT NULL DEFAULT 0);
            """)

    @contextmanager
    def connect(self):
        db = sqlite3.connect(self.path, timeout=10)
        try:
            # 连接上下文提交或回滚事务；finally 另外关闭连接，不能混淆这两件事。
            with db:
                yield db
        finally:
            db.close()

    @staticmethod
    def _get(db, task_id):
        row = db.execute("SELECT body FROM tasks WHERE id=?", (task_id,)).fetchone()
        if row is None:
            raise ValueError("任务不存在")
        return json.loads(row[0])

    @staticmethod
    def _save(db, task, kind, payload=None):
        # 调用者已经持有 BEGIN IMMEDIATE 写事务。状态、事件、Outbox 在同一事务中提交。
        db.execute("INSERT INTO tasks VALUES (?, ?) ON CONFLICT(id) DO UPDATE SET body=excluded.body",
                   (task["id"], json.dumps(task)))
        seq = db.execute("SELECT COALESCE(MAX(seq), 0)+1 FROM events WHERE task_id=?",
                         (task["id"],)).fetchone()[0]
        data = json.dumps(payload or {"status": task["status"]})
        db.execute("INSERT INTO events VALUES (?, ?, ?, ?)", (task["id"], seq, kind, data))
        db.execute("INSERT INTO outbox(id, task_id, payload) VALUES (?, ?, ?)",
                   (f'{task["id"]}:{seq}', task["id"], json.dumps({"seq": seq, "kind": kind})))

    def get(self, task_id):
        with self.connect() as db:
            return self._get(db, task_id)

    def submit(self, task_id: str, amount: int = 29900):
        if not task_id.strip() or type(amount) is not int or not 0 < amount <= 29900:
            raise ValueError("任务 ID 不能为空；金额必须为 1—29900 整数分")
        # task ID 是本例的业务申请身份；同一次申请的所有重试复用它。
        # 生产中还需将客户端幂等键与经过认证的租户绑定，不接受任意人指定身份。
        with self.connect() as db:
            db.execute("BEGIN IMMEDIATE")
            existing = db.execute("SELECT body FROM tasks WHERE id=?", (task_id,)).fetchone()
            if existing:
                task = json.loads(existing[0])
                if task["amount_minor"] != amount:
                    raise ValueError("同一任务 ID 的金额发生变化")
                return task
            task = {"id": task_id, "tenant": "shop-demo", "order_id": "A1042",
                    "order_revision": 13, "quality_verified": True,
                    "amount_minor": amount, "currency": "CNY",
                    "operation_id": "refund:" + task_id, "status": "awaiting_approval",
                    "schema_version": 1, "approval": None}
            self._save(db, task, "task.created")
            return task

    def decide(self, task_id: str, approved: bool, decision_id: str):
        if type(approved) is not bool or not decision_id:
            raise ValueError("审批必须为真正的布尔值，并携带 decision ID")
        with self.connect() as db:
            db.execute("BEGIN IMMEDIATE")
            task = self._get(db, task_id)
            prior = task["approval"]
            if prior and prior["decision_id"] == decision_id:
                if prior["approved"] != approved:
                    raise ValueError("同一决定 ID 的内容冲突")
                return task  # 重复投递的同一决定不产生第二条审批记录。
            if task["status"] != "awaiting_approval":
                raise ValueError("任务当前不接受审批")
            task["approval"] = {"decision_id": decision_id, "approved": approved,
                                "actor": "demo-reviewer", "expires_at": time.time() + 3600,
                                "order_revision": task["order_revision"],
                                "fingerprint": fingerprint(task["tenant"], task["order_id"],
                                                           task["amount_minor"], task["currency"])}
            # demo-reviewer 是虚构身份；CLI 没有登录，不能用于真实审批。
            task["status"] = "ready" if approved else "rejected"
            self._save(db, task, "approval.accepted" if approved else "approval.rejected")
            return task

    def advance(self, task_id: str, *, lose_response: bool = False):
        task = self.get(task_id)
        if task["status"] in {"completed", "rejected", "blocked"}:
            return task
        if task["status"] not in {"ready", "executing", "reconciling"}:
            raise ValueError("任务尚未获得审批")
        fp = fingerprint(task["tenant"], task["order_id"], task["amount_minor"], task["currency"])
        # 先核对旧回执。即使审批此刻过期，已发生的支付事实仍需被正确记录。
        receipt = self.payment.lookup(task["operation_id"], fp)
        if receipt is None:
            with self.connect() as db:
                db.execute("BEGIN IMMEDIATE")
                task = self._get(db, task_id)
                if task["status"] in {"completed", "rejected", "blocked"}:
                    return task
                approval = task["approval"]
                # 本例订单事实不可编辑。真实系统这里必须重读订单、权限和规则版本。
                if (not approval or approval["approved"] is not True
                        or approval["expires_at"] <= time.time()
                        or approval["fingerprint"] != fp
                        or approval["order_revision"] != task["order_revision"]):
                    task["status"] = "blocked"
                    self._save(db, task, "approval.invalid")
                    return task
                task["status"] = "executing"
                self._save(db, task, "payment.intent")
            # 远端调用不能放进本地事务并假装两者原子提交。
            try:
                receipt = self.payment.refund(task["operation_id"], task["tenant"],
                                              task["order_id"], task["amount_minor"],
                                              task["currency"], lose_response=lose_response)
            except (TimeoutError, ValueError) as error:
                with self.connect() as db:
                    db.execute("BEGIN IMMEDIATE")
                    task = self._get(db, task_id)
                    if task["status"] == "completed":
                        return task  # 迟到的失败不能覆盖另一执行器已记录的成功。
                    task["status"] = "reconciling" if isinstance(error, TimeoutError) else "blocked"
                    task["error"] = str(error)
                    self._save(db, task, "payment." + task["status"])
                    return task
        with self.connect() as db:
            db.execute("BEGIN IMMEDIATE")
            task = self._get(db, task_id)
            if task["status"] != "completed":
                task.update(status="completed", receipt=receipt)
                task.pop("error", None)
                self._save(db, task, "task.completed", receipt)
            return task

    def events(self, task_id: str, after: int = 0):
        if type(after) is not int or after < 0:
            raise ValueError("游标必须为非负整数")
        with self.connect() as db:
            self._get(db, task_id)  # 不将不存在的任务静默返回成空历史。
            rows = db.execute("SELECT seq, kind, payload FROM events WHERE task_id=? AND seq>? ORDER BY seq",
                              (task_id, after)).fetchall()
        return [{"seq": seq, "kind": kind, "payload": json.loads(payload)} for seq, kind, payload in rows]
