"""DBOS 本地持久工作流示例。数据库文件落盘；不提供真实认证或支付。"""
import argparse
import json
from pathlib import Path
from dbos import DBOS, SetWorkflowID
from runtime_store import RuntimeStore


@DBOS.step()
def create_task(task_id: str, directory: str):
    # 业务写入在 step 中。重放已有 step 时 DBOS 可返回已记录结果。
    return RuntimeStore(Path(directory)).submit(task_id)


@DBOS.step()
def apply_decision(task_id: str, decision: dict, directory: str):
    store = RuntimeStore(Path(directory))
    store.decide(task_id, decision["approved"], decision["decision_id"])
    # 即使本 step 在业务成功后、DBOS 保存结果前崩溃，稳定操作 ID 仍可核对原回执。
    return store.advance(task_id) if decision["approved"] else store.get(task_id)


@DBOS.workflow()
def refund(task_id: str, directory: str):
    create_task(task_id, directory)
    decision = DBOS.recv(topic="approval", timeout_seconds=3600)
    if decision is None:
        # 等待超时是明确结果；示例不自动重新开审批，也不宣称业务已经完成。
        return {"status": "approval_wait_expired", "task_id": task_id}
    if (not isinstance(decision, dict) or type(decision.get("approved")) is not bool
            or not decision.get("decision_id")):
        raise ValueError("审批消息格式不正确")
    return apply_decision(task_id, decision, directory)


def main():
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("command", choices=["start", "approve", "reject"])
    parser.add_argument("task_id")
    parser.add_argument("--data-dir", type=Path, default=Path(".dbos-refund-demo"))
    args = parser.parse_args()
    directory = args.data_dir.resolve()
    directory.mkdir(parents=True, exist_ok=True)
    # 同一路径、应用名及兼容代码是恢复的前提；不把进程内对象当作持久状态。
    DBOS(config={"name": "refund-demo", "system_database_url":
                 "sqlite:///" + str(directory / "dbos.sqlite")})
    DBOS.launch()
    try:
        if args.command == "start":
            with SetWorkflowID(args.task_id):
                handle = DBOS.start_workflow(refund, args.task_id, str(directory))
            # 命令保持运行等待结果；另一终端 approve/reject。重启后复用工作流 ID。
            print(json.dumps(handle.get_result(), ensure_ascii=False, indent=2))
        else:
            # 消息去重键与业务审批 ID 分别明确；重复相同命令仍发送同一决定。
            DBOS.send(args.task_id, {"approved": args.command == "approve",
                                    "decision_id": "decision-1"}, topic="approval",
                      idempotency_key="approval:" + args.task_id)
            print("教学审批消息已发送；执行结果由工作流与业务状态共同确认。")
    finally:
        DBOS.destroy()


if __name__ == "__main__":
    main()
