"""LangGraph 同步 SQLite Checkpoint 示例；每次命令都是独立进程，不调用模型。"""
import argparse
import json
from pathlib import Path
from typing import TypedDict
from langgraph.checkpoint.sqlite import SqliteSaver
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
from runtime_store import RuntimeStore


class RefundState(TypedDict):
    task_id: str
    status: str


def build_graph(store: RuntimeStore, saver: SqliteSaver):
    def approval(state: RefundState):
        task = store.get(state["task_id"])
        # 恢复会从本节点开头重执行，所以 interrupt 前只读状态，不创建支付。
        # 输入包含动作参数，帮助审批人判断其决定具体绑定什么。
        decision = interrupt({"task_id": task["id"], "order_id": task["order_id"],
                              "amount_minor": task["amount_minor"], "currency": "CNY"})
        if not isinstance(decision, dict) or type(decision.get("approved")) is not bool:
            raise ValueError("resume 必须包含布尔 approved 与 decision_id")
        result = store.decide(task["id"], decision["approved"], decision["decision_id"])
        return {"status": result["status"]}

    def execute(state: RefundState):
        # 图节点可能重执行。稳定业务 operation ID 与支付回执负责避免重复退款。
        return {"status": store.advance(state["task_id"])["status"]}

    graph = StateGraph(RefundState)
    graph.add_node("approval", approval)
    graph.add_node("execute", execute)
    graph.add_edge(START, "approval")
    graph.add_conditional_edges("approval", lambda s: "execute" if s["status"] == "ready" else END)
    graph.add_edge("execute", END)
    return graph.compile(checkpointer=saver)


def main():
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("command", choices=["submit", "approve", "reject", "retry", "status"])
    parser.add_argument("task_id")
    parser.add_argument("--data-dir", type=Path, default=Path(".runtime-demo"))
    parser.add_argument("--decision-id", default="decision-1")
    args = parser.parse_args()
    store = RuntimeStore(args.data_dir)
    config = {"configurable": {"thread_id": args.task_id}}
    # 使用文件路径而非 :memory:。下一次命令必须使用同一路径与 thread_id。
    with SqliteSaver.from_conn_string(str(args.data_dir / "checkpoints.sqlite")) as saver:
        graph = build_graph(store, saver)
        if args.command == "submit":
            task = store.submit(args.task_id)
            snapshot = graph.get_state(config)
            # 重复 submit 不重新创建一轮图执行；已有暂停点应走 approve/reject。
            if not snapshot.values:
                graph.invoke({"task_id": task["id"], "status": task["status"]}, config,
                             durability="sync")
        elif args.command in {"approve", "reject"}:
            graph.invoke(Command(resume={"approved": args.command == "approve",
                                         "decision_id": args.decision_id}), config,
                         durability="sync")
        elif args.command == "retry":
            # 处理进程异常留下的待执行节点；不会自动替用户做审批决定。
            graph.invoke(None, config, durability="sync")
        snapshot = graph.get_state(config)
        print(json.dumps({"business": store.get(args.task_id),
                          "graph_values": snapshot.values, "next": snapshot.next},
                         ensure_ascii=False, indent=2))
        # 图 END 只表示本轮图停止。若业务为 reconciling，使用 refund_cli advance 对账。
        # 对账后图里的旧 status 可能滞后；展示业务结果应读取上面的 business。


if __name__ == "__main__":
    main()
