一、条件路由与循环

1、条件路由

基于预定的规则的路径选择

2、LLM驱动路由

将决策权交给LLM,处理模糊和复杂的输入

3、循环模式

构建“粗筛后细分”的高效漏斗的决策流程

二、Streaming输出

  • 提升用户体验,长耗时的任务时提供即时反馈
  • 实时过程监控,观察agent的思考过程
  • token-by-token输出

1、values

输出每一步执行后的完整图状态
触发时机:每个Node执行结束后
数据格式:包含state中定义的所有字段

2、update

仅输出状态发生变化,返回刚刚结束的那个节点的输出
触发时机:每个节点执行结束后
数据格式:一个字典, key是节点名称,values是节点的输出结果
场景:高效更新UI

3、messages

直接流式输出LLM生成的Token
场景:实时聊天机器人

三、Interrupt人工干预

  • 人工审核,在关键点暂停
  • 动态修正,允许用户在流程中修改数据或重试失败步骤
  • 信息注入,在AI需要额外上下文时,由人类提供

四、Time Travel状态回溯

  • 查看历史,完整复现每一次执行的完整历史
  • 回溯状态
  • 从历史状态创建一个新的分支

五、 LangGraph实战—交互式内容审核

# =============================================================================
# 交互式内容审核工作流
# =============================================================================
# 功能:
#   1. 自动内容分析 - LLM 驱动的内容检测与风险评级
#   2. 流式输出分析进度 - 实时展示各节点执行进度
#   3. 多级人工审核 - 一级审核 + 高级审核,支持升级机制
#   4. 历史追溯 - SqliteSaver 持久化所有检查点,支持完整审计追踪
#   5. 可撤销决策 - 通过 update_state 回退到任意历史检查点
#
# 图流程:
#   START → analyze → risk_assess → [条件路由]
#       低风险 → finalize → END
#       中风险 → first_review → [通过/驳回/升级]
#               通过 → finalize → END
#               驳回 → revise → analyze (循环)
#               升级 → senior_review → [通过/驳回]
#                      通过 → finalize → END
#                      驳回 → revise → analyze (循环)
#       高风险 → senior_review → [通过/驳回]
#               通过 → finalize → END
#               驳回 → revise → analyze (循环)
# =============================================================================

import json
import os
import sqlite3
import urllib.request
from datetime import datetime
from typing import TypedDict, Annotated

from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.sqlite import SqliteSaver
from langgraph.types import interrupt, Command


# ========================= 配置区域 =========================
# DeepSeek API 配置(从环境变量读取密钥)
API_KEY = os.getenv("DEEPSEEK_API_KEY")
API_URL = "https://api.deepseek.com/chat/completions"
MODEL = "deepseek-v4-pro"


# ========================= SQLite 持久化配置 =========================
# 数据库文件路径:当前脚本所在目录下的 moderation.db
DB_PATH = os.path.join(os.path.dirname(os.path.abspath(__file__)), "moderation.db")
# 创建 SQLite 连接(check_same_thread=False 允许跨线程访问)
conn = sqlite3.connect(DB_PATH, check_same_thread=False)
# SqliteSaver 实例:用于持久化保存每个步骤的检查点状态
# 支持历史追溯(get_state_history)和撤销决策(update_state)
sqlite_saver = SqliteSaver(conn)


# ========================= 状态定义 =========================
class ModerationState(TypedDict):
    """
    内容审核工作流状态定义:
        - content: 待审核的原始内容
        - analysis_result: AI 内容分析结果(JSON 字符串)
        - sensitive_words: 检测到的敏感词(逗号分隔)
        - risk_level: 风险评级(低/中/高)
        - risk_reason: 风险评级原因
        - review_decision: 审核决策(approved/rejected/escalate)
        - review_comments: 审核意见或驳回原因
        - review_level: 当前审核级别(first/senior)
        - final_result: 最终处理结果
        - audit_trail: 审计追踪日志,使用 Annotated + reducer 实现追加(而非覆盖)
    """
    content: str
    analysis_result: str
    sensitive_words: str
    risk_level: str
    risk_reason: str
    review_decision: str
    review_comments: str
    review_level: str
    final_result: str
    audit_trail: Annotated[list, lambda old, new: old + new]


# ========================= LLM 工具函数 =========================
def call_deepseek(messages: list) -> str:
    """
    调用 DeepSeek API 发送聊天请求
    使用 urllib 发送 POST 请求,无需额外安装 requests 库

    :param messages: 消息列表,格式为 [{"role": "system/user", "content": "..."}]
    :return: LLM 返回的文本内容
    """
    data = json.dumps({
        "model": MODEL,
        "messages": messages,
        "temperature": 0.3
    }).encode("utf-8")

    req = urllib.request.Request(
        API_URL,
        data=data,
        headers={
            "Authorization": f"Bearer {API_KEY}",
            "Content-Type": "application/json"
        }
    )

    with urllib.request.urlopen(req) as resp:
        result = json.loads(resp.read().decode("utf-8"))

    return result["choices"][0]["message"]["content"]


# ========================= 节点函数定义 =========================
# 节点1 - 内容分析节点
def analyze_node(state: ModerationState) -> dict:
    """
    职责:调用 LLM 对内容进行全方位安全性检测
    检测维度:敏感词、政治内容、脏话、暴力、歧视性言论等
    返回:检测结果、敏感词列表、风险等级、评级原因
    """
    content = state["content"]
    print(f"  🔍 正在分析内容安全性...")

    messages = [
        {"role": "system",
         "content": """你是一个专业的内容安全检测专家。请对以下内容进行全面安全性检测,检测维度包括:
            1. 敏感词汇(政治敏感、宗教敏感等)
            2. 脏话/侮辱性语言
            3. 暴力/血腥内容
            4. 歧视性言论(种族、性别、地域等)
            5. 违法违规信息

            请严格以JSON格式返回检测结果,不要包含任何其他文字:
            {
                "sensitive_words": ["检测到的敏感词/问题词列表,没有则为空数组"],
                "risk_level": "风险评级:低/中/高",
                "risk_reason": "评级原因的详细说明",
                "analysis_detail": "详细分析说明"
            }"""
         },
        {"role": "user", "content": f"请检测以下内容的安全性:\n{content}"}
    ]

    result = call_deepseek(messages)
    print(f"  📊 分析完成")

    # 解析 JSON,若 LLM 返回格式异常则使用安全默认值
    try:
        parsed = json.loads(result)
    except json.JSONDecodeError:
        parsed = {
            "sensitive_words": [],
            "risk_level": "中",
            "risk_reason": "LLM返回格式异常,建议人工审核",
            "analysis_detail": result
        }

    return {
        "analysis_result": result,
        "sensitive_words": ", ".join(parsed.get("sensitive_words", [])) or "无",
        "risk_level": parsed.get("risk_level", "中"),
        "risk_reason": parsed.get("risk_reason", ""),
        "audit_trail": [{
            "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
            "action": "AI内容分析",
            "detail": f"风险等级: {parsed.get('risk_level', '中')}"
        }]
    }

# 节点2 - 风险评估节点
def risk_assess_node(state: ModerationState) -> dict:
    """
    职责:根据 AI 分析结果确定最终风险等级,决定后续审核路径
    路由规则:
      - 低风险 → 自动通过(finalize)
      - 中风险 → 一级人工审核(first_review)
      - 高风险 → 高级人工审核(senior_review)
    """
    risk_level = state.get("risk_level", "中")
    sensitive_words = state.get("sensitive_words", "无")

    print(f"  📋 风险评估: {risk_level}风险")

    return {
        "audit_trail": [{
            "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
            "action": "风险评估",
            "detail": f"风险等级={risk_level}, 敏感词={sensitive_words}"
        }]
    }

# 节点3a - 一级审核节点(中断点)
def first_review_node(state: ModerationState) -> dict:
    """
    职责:
      - 中风险内容在此暂停,等待一级审核人员决策
      - 审核人员可选择:通过 / 驳回 / 升级到高级审核
    """
    print("  👤 等待一级审核...")

    # interrupt() 暂停图执行,等待审核人员决策
    # 恢复时,Command(resume=value) 中的 value 成为 interrupt() 的返回值
    decision = interrupt({
        "type": "first_review",
        "message": f"📝 一级审核 - 内容需要您的审核\n"
                   f"  内容摘要: {state['content'][:100]}...\n"
                   f"  敏感词: {state['sensitive_words']}\n"
                   f"  风险等级: {state['risk_level']}\n"
                   f"  分析原因: {state['risk_reason']}\n\n"
                   f"请选择:\n"
                   f"  1. 通过 (approve)\n"
                   f"  2. 驳回 (reject)\n"
                   f"  3. 升级到高级审核 (escalate)"
    })

    # ===== 以下代码在 resume 后执行 =====
    print(f"  📩 一级审核决策: {decision}")

    if decision in ["approve", "approved", "1"]:
        return {
            "review_decision": "approved",
            "review_level": "first",
            "audit_trail": [{
                "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
                "action": "一级审核-通过",
                "detail": "一级审核人员批准发布"
            }]
        }
    elif decision in ["escalate", "3"]:
        return {
            "review_decision": "escalate",
            "review_level": "first",
            "audit_trail": [{
                "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
                "action": "一级审核-升级",
                "detail": "一级审核人员升级到高级审核"
            }]
        }
    else:
        return {
            "review_decision": "rejected",
            "review_level": "first",
            "audit_trail": [{
                "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
                "action": "一级审核-驳回",
                "detail": "一级审核人员驳回"
            }]
        }

# 节点3b - 高级审核节点(中断点)
def senior_review_node(state: ModerationState) -> dict:
    """
        职责:
      - 高风险内容或被一级审核升级的内容在此暂停
      - 高级审核人员可选择:通过 / 驳回
    """
    print("  👤 等待高级审核...")

    # interrupt() 暂停图执行,等待高级审核人员决策
    decision = interrupt({
        "type": "senior_review",
        "message": f"📝 高级审核 - 内容需要您的审核\n"
                   f"  内容摘要: {state['content'][:100]}...\n"
                   f"  敏感词: {state['sensitive_words']}\n"
                   f"  风险等级: {state['risk_level']}\n"
                   f"  分析原因: {state['risk_reason']}\n\n"
                   f"请选择:\n"
                   f"  1. 通过 (approve)\n"
                   f"  2. 驳回 (reject)"
    })

    # ===== 以下代码在 resume 后执行 =====
    print(f"  📩 高级审核决策: {decision}")

    if decision in ["approve", "approved", "1"]:
        return {
            "review_decision": "approved",
            "review_level": "senior",
            "audit_trail": [{
                "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
                "action": "高级审核-通过",
                "detail": "高级审核人员批准发布"
            }]
        }
    else:
        return {
            "review_decision": "rejected",
            "review_level": "senior",
            "audit_trail": [{
                "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
                "action": "高级审核-驳回",
                "detail": "高级审核人员驳回"
            }]
        }

# 节点4 - 最终决策节点
def finalize_node(state: ModerationState) -> dict:
    """
        职责:根据审核结果生成最终处理报告
    """
    decision = state.get("review_decision", "approved")
    risk_level = state.get("risk_level", "低")
    review_level = state.get("review_level", "auto")

    if decision == "approved":
        if risk_level == "低":
            result_text = (
                f"\n{'=' * 60}\n"
                f"  ✅ 内容自动通过(低风险)\n"
                f"{'=' * 60}\n"
                f"  📄 内容: {state['content'][:200]}\n"
                f"  🔍 敏感词: {state['sensitive_words']}\n"
                f"  📊 风险等级: {risk_level}\n"
                f"  📋 处理方式: 自动发布\n"
                f"{'=' * 60}"
            )
        else:
            result_text = (
                f"\n{'=' * 60}\n"
                f"  ✅ 内容审核通过\n"
                f"{'=' * 60}\n"
                f"  📄 内容: {state['content'][:200]}\n"
                f"  🔍 敏感词: {state['sensitive_words']}\n"
                f"  📊 风险等级: {risk_level}\n"
                f"  👤 审核级别: {review_level}\n"
                f"  📋 处理方式: 审核通过,可以发布\n"
                f"{'=' * 60}"
            )
    else:
        result_text = (
            f"\n{'=' * 60}\n"
            f"  ❌ 内容审核未通过\n"
            f"{'=' * 60}\n"
            f"  📄 内容: {state['content'][:200]}\n"
            f"  📊 风险等级: {risk_level}\n"
            f"  💬 审核意见: {state.get('review_comments', '无')}\n"
            f"  📋 处理方式: 需要修改后重新提交\n"
            f"{'=' * 60}"
        )

    print(result_text)

    return {
        "final_result": result_text,
        "audit_trail": [{
            "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
            "action": "最终决策",
            "detail": f"决策={decision}, 审核级别={review_level}"
        }]
    }

# 节点5 - 内容修改节点
def revise_node(state: ModerationState) -> dict:
    """
    职责:当内容被驳回时,记录修改意见,准备重新提交分析
    流程:revise → analyze(循环回内容分析节点)
    """
    comments = state.get("review_comments", "请根据审核意见修改内容")
    print(f"  ✏️ 内容被驳回,修改意见: {comments}")
    print(f"  🔄 将重新进入内容分析流程...")

    return {
        "review_decision": "",
        "review_comments": "",
        "audit_trail": [{
            "timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
            "action": "内容修改",
            "detail": f"驳回原因: {comments},准备重新分析"
        }]
    }


# ========================= 路由函数 =========================
# 条件路由:根据风险等级决定审核路径
def route_by_risk(state: ModerationState) -> str:
    """
    - 低风险 → finalize(自动通过)
    - 中风险 → first_review(一级人工审核)
    - 高风险 → senior_review(高级人工审核)
    """
    risk = state.get("risk_level", "中")
    if risk == "低":
        return "finalize"
    elif risk == "高":
        return "senior_review"
    else:
        return "first_review"

# 条件路由:审核后的流程走向
def route_after_review(state: ModerationState) -> str:
    """
    - approved → finalize(最终决策)
    - rejected → revise(内容修改,然后重新分析)
    - escalate → senior_review(升级到高级审核,仅一级审核可能返回此值)
    """
    decision = state.get("review_decision", "")
    if decision == "approved":
        return "finalize"
    elif decision == "escalate":
        return "senior_review"
    else:
        return "revise"


# ========================= 构建 LangGraph 工作流 =========================

# 1. 创建状态图,绑定 ModerationState 作为状态结构
graph = StateGraph(ModerationState)

# 2. 注册六个功能节点
graph.add_node("analyze", analyze_node)               # 内容分析:AI 检测
graph.add_node("risk_assess", risk_assess_node)       # 风险评估:确定审核路径
graph.add_node("first_review", first_review_node)     # 一级审核:中断点
graph.add_node("senior_review", senior_review_node)   # 高级审核:中断点
graph.add_node("finalize", finalize_node)             # 最终决策:生成报告
graph.add_node("revise", revise_node)                 # 内容修改:驳回后修改

# 3. 添加固定边
graph.add_edge(START, "analyze")                              # 开始 → 内容分析
graph.add_edge("analyze", "risk_assess")              # 分析 → 风险评估
graph.add_edge("revise", "analyze")                   # 修改 → 重新分析(循环)
graph.add_edge("finalize", END)                               # 最终决策 → 结束

# 4. 条件边:风险评估之后根据风险等级动态路由
#    低风险 → finalize | 中风险 → first_review | 高风险 → senior_review
graph.add_conditional_edges(
    "risk_assess",
    route_by_risk,
    {
        "finalize": "finalize",
        "first_review": "first_review",
        "senior_review": "senior_review",
    }
)

# 5. 条件边:审核之后根据审核决策动态路由
#    通过 → finalize | 驳回 → revise | 升级 → senior_review
graph.add_conditional_edges(
    "first_review",
    route_after_review,
    {
        "finalize": "finalize",
        "revise": "revise",
        "senior_review": "senior_review",
    }
)

# 6. 条件边:高级审核后只能通过与驳回
graph.add_conditional_edges(
    "senior_review",
    route_after_review,
    {
        "finalize": "finalize",
        "revise": "revise",
    }
)

# 7. 编译图,绑定 SqliteSaver 作为检查点存储器
#    所有状态变更都会持久化到 SQLite,支持历史追溯和撤销决策
app = graph.compile(checkpointer=sqlite_saver)


# ========================= 生成并导出流程图 =========================
def generate_flowchart():
    """
    使用 Mermaid 语法生成可视化流程图,封装为 HTML 文件
    代码运行完毕后调用,在浏览器中打开即可查看工作流结构
    """
    mermaid_str = app.get_graph().draw_mermaid()
    html_content = f"""<!DOCTYPE html>
        <html>
        <head>
        <meta charset="utf-8">
        <title>内容审核工作流 - Interactive Content Moderation</title>
        </head>
        <body>
        <h2 style="text-align:center; margin:20px;">交互式内容审核工作流流程图</h2>
        <pre class="mermaid">{mermaid_str}</pre>
        <script src="https://cdn.jsdelivr.net/npm/mermaid/dist/mermaid.min.js"></script>
        <script>mermaid.initialize({{startOnLoad:true}});</script>
        </body>
        </html>"""

    output_path = os.path.join(
        os.path.dirname(os.path.abspath(__file__)),
        "flowchart_moderation.html"
    )
    with open(output_path, "w", encoding="utf-8") as f:
        f.write(html_content)
    print(f"\n📊 流程图已保存到: {output_path}")
    print("   请在浏览器中打开该文件查看流程图")


# ========================= 辅助函数 =========================
def display_analysis_result(state: dict):
    """
    格式化展示 AI 分析结果
    :param state: 当前工作流状态
    """
    print(f"\n{'─' * 60}")
    print(f"  📄 待审核内容: {state.get('content', '')[:200]}")
    print(f"  🔍 敏感词:     {state.get('sensitive_words', '无')}")
    print(f"  📊 风险等级:   {state.get('risk_level', '未知')}")
    print(f"  📋 风险原因:   {state.get('risk_reason', '未知')}")
    print(f"{'─' * 60}")


def display_audit_trail(audit_trail: list):
    """
    格式化展示审计追踪日志
    :param audit_trail: 审计日志列表
    """
    if not audit_trail:
        print("  (暂无审核记录)")
        return
    for i, entry in enumerate(audit_trail, 1):
        print(f"  [{i}] {entry.get('timestamp', '')} | "
              f"{entry.get('action', '')} | {entry.get('detail', '')}")


# ========================= 主交互循环 =========================

def main():
    """
    主交互循环
    支持三种操作:
      1. 提交内容进行审核(含流式进度展示 + 多级人工审核)
      2. 查看历史审核记录(审计追溯)
      3. 撤销决策(回退到任意历史检查点)
    """
    print("=" * 65)
    print("  🛡️  交互式内容审核工作流 v1.0")
    print("=" * 65)
    print(f"  📌 LLM 模型: {MODEL}")
    print(f"  📌 数据库: {DB_PATH}")
    print(f"  📌 功能: 自动分析 | 流式进度 | 多级审核 | 历史追溯 | 撤销决策")
    print("=" * 65)

    thread_counter = 0

    while True:
        print(f"\n{'─' * 40}")
        print("  请选择操作:")
        print("  1. 📝 提交内容进行审核")
        print("  2. 📜 查看历史审核记录")
        print("  3. ↩️  撤销决策(回退到历史检查点)")
        print("  0. 🚪 退出系统")
        print(f"{'─' * 40}")

        choice = input("  请输入选项 (0/1/2/3): ").strip()

        if choice == "0":
            print("\n  👋 感谢使用内容审核工作流,再见!")
            break

        elif choice == "1":
            # ===== 提交内容进行审核 =====
            content = input("\n  请输入待审核内容(输入 'quit' 返回主菜单):\n  > ").strip()
            if content.lower() == "quit":
                continue
            if not content:
                print("  ⚠️ 内容不能为空!")
                continue

            thread_counter += 1
            thread_id = f"moderation-{thread_counter:03d}"
            config = {"configurable": {"thread_id": thread_id}}

            # 初始化状态
            initial_state = {
                "content": content,
                "analysis_result": "",
                "sensitive_words": "",
                "risk_level": "",
                "risk_reason": "",
                "review_decision": "",
                "review_comments": "",
                "review_level": "",
                "final_result": "",
                "audit_trail": []
            }

            print(f"\n{'=' * 65}")
            print(f"  🔄 开始审核流程 [线程: {thread_id}]")
            print(f"{'=' * 65}")

            # ---- 阶段1:流式执行内容分析 + 风险评估 ----
            # graph.stream() 逐节点产出结果,实时展示执行进度
            print(f"\n  ▶ 阶段1: AI 自动内容分析(流式输出)")
            result = None
            try:
                for chunk in app.stream(initial_state, config):
                    # stream() 每次产出一个节点的执行结果(dict: {node_name: state_update})
                    if isinstance(chunk, dict):
                        for node_name, node_output in chunk.items():
                            print(f"    ✅ [{node_name}] 执行完成")
                    result = chunk
            except Exception as e:
                print(f"    ❌ 执行出错: {e}")
                continue

            # 获取完整的当前状态
            current_state = app.get_state(config)
            state_values = current_state.values if current_state else {}

            # 展示 AI 分析结果
            display_analysis_result(state_values)

            risk_level = state_values.get("risk_level", "中")

            # ---- 阶段2:根据风险等级判断是否需要人工审核 ----
            if risk_level == "低":
                # 低风险自动通过,流程已完成(无中断)
                print(f"\n  ✅ 低风险内容,自动通过审核!")
                print(f"  📋 最终结果: {state_values.get('final_result', '')[:200]}")

            else:
                # 中/高风险需要人工审核
                # 图已在 first_review 或 senior_review 节点处暂停(interrupt)
                print(f"\n  ⏸️  图已在审核节点暂停,等待人工审核...")

                # ---- 阶段3:一级审核 ----
                if risk_level in ["中", "高"]:
                    display_analysis_result(state_values)

                    print(f"\n  👤 一级审核决策:")
                    print(f"    1. ✅ 通过 (approve)")
                    print(f"    2. ❌ 驳回 (reject)")
                    if risk_level == "中":
                        print(f"    3. ⬆️  升级到高级审核 (escalate)")

                    review_input = input("  请输入决策 (1/2/3): ").strip()

                    # 映射输入到决策值
                    decision_map = {"1": "approve", "2": "reject", "3": "escalate"}
                    decision = decision_map.get(review_input, "approve" if risk_level == "高" else "escalate")

                    # 恢复图执行,将决策通过 Command(resume=...) 传入 interrupt() 的返回值
                    print(f"\n  ▶ 恢复执行...")
                    try:
                        result = app.invoke(Command(resume=decision), config)
                    except Exception as e:
                        print(f"    ❌ 恢复执行出错: {e}")
                        continue

                    # 获取恢复后的状态
                    current_state = app.get_state(config)
                    state_values = current_state.values if current_state else {}

                    # 检查是否有后续中断(一级审核升级到高级审核时)
                    has_interrupt = hasattr(current_state, 'next') and current_state.next

                    if has_interrupt:
                        # ---- 阶段4:高级审核(升级场景) ----
                        print(f"\n  ⬆️  内容已升级到高级审核")
                        print(f"  ⏸️  图已在高级审核节点暂停...")

                        display_analysis_result(state_values)

                        print(f"\n  👤 高级审核决策:")
                        print(f"    1. ✅ 通过 (approve)")
                        print(f"    2. ❌ 驳回 (reject)")

                        senior_input = input("  请输入决策 (1/2): ").strip()
                        senior_decision = "approve" if senior_input == "1" else "reject"

                        # 恢复图执行
                        print(f"\n  ▶ 恢复执行...")
                        try:
                            result = app.invoke(Command(resume=senior_decision), config)
                        except Exception as e:
                            print(f"    ❌ 恢复执行出错: {e}")
                            continue

                        # 获取最终状态
                        current_state = app.get_state(config)
                        state_values = current_state.values if current_state else {}

                    # 展示最终结果
                    print(f"\n{'=' * 65}")
                    print(f"  📋 审核完成!最终结果:")
                    print(f"{'=' * 65}")
                    print(f"  {state_values.get('final_result', '处理完成')}")

            # 展示完整审计追踪
            print(f"\n  📜 本次审核完整审计追踪:")
            display_audit_trail(state_values.get("audit_trail", []))

        elif choice == "2":
            # ===== 查看历史审核记录 =====
            print(f"\n{'=' * 65}")
            print(f"  📜 历史审核记录")
            print(f"{'=' * 65}")

            # 遍历所有已知的 thread_id,获取每个线程的状态
            found_records = False
            for i in range(1, thread_counter + 1):
                tid = f"moderation-{i:03d}"
                cfg = {"configurable": {"thread_id": tid}}
                try:
                    state_snap = app.get_state(cfg)
                    if state_snap and state_snap.values:
                        vals = state_snap.values
                        found_records = True
                        print(f"\n  ── 记录 #{i} (线程: {tid}) ──")
                        print(f"    📄 内容: {vals.get('content', '')[:80]}")
                        print(f"    📊 风险等级: {vals.get('risk_level', '未知')}")
                        print(f"    🔍 敏感词: {vals.get('sensitive_words', '无')}")
                        print(f"    📋 审核决策: {vals.get('review_decision', '自动处理')}")
                        print(f"    👤 审核级别: {vals.get('review_level', '自动')}")
                        print(f"    📜 审计追踪:")
                        display_audit_trail(vals.get("audit_trail", []))
                except Exception:
                    continue

            if not found_records:
                print("  (暂无审核记录,请先提交内容进行审核)")

        elif choice == "3":
            # ===== 撤销决策(回退到历史检查点) =====
            print(f"\n{'=' * 65}")
            print(f"  ↩️  撤销决策 - 查看所有检查点历史")
            print(f"{'=' * 65}")

            if thread_counter == 0:
                print("  (暂无可撤销的记录)")
                continue

            # 列出所有线程的检查点
            all_checkpoints = []
            for i in range(1, thread_counter + 1):
                tid = f"moderation-{i:03d}"
                cfg = {"configurable": {"thread_id": tid}}
                try:
                    history = app.get_state_history(cfg)
                    for cp in history:
                        all_checkpoints.append((tid, cp))
                except Exception:
                    continue

            if not all_checkpoints:
                print("  (暂无检查点记录)")
                continue

            # 展示所有检查点
            print(f"\n  共找到 {len(all_checkpoints)} 个检查点:\n")
            for idx, (tid, cp) in enumerate(all_checkpoints):
                vals = cp.values or {}
                ts = getattr(cp, 'timestamp', '未知时间')
                content_preview = vals.get('content', '')[:50]
                risk = vals.get('risk_level', '未知')
                decision = vals.get('review_decision', '处理中')
                next_nodes = cp.next if cp.next else ('已完成',)
                print(f"  [{idx + 1}] 线程: {tid}")
                print(f"      时间: {ts}")
                print(f"      内容: {content_preview}...")
                print(f"      风险: {risk} | 决策: {decision}")
                print(f"      下一步: {next_nodes}")
                print()

            # 用户选择要回退到的检查点
            try:
                cp_choice = int(input("  请输入要回退到的检查点编号 (0 取消): ").strip())
            except ValueError:
                print("  ⚠️ 无效输入")
                continue

            if cp_choice == 0:
                continue

            if 1 <= cp_choice <= len(all_checkpoints):
                selected_tid, selected_cp = all_checkpoints[cp_choice - 1]

                # 使用 update_state 回退到选定的检查点
                # 这会将指定线程的状态恢复到该检查点时刻的值
                try:
                    app.update_state(
                        {"configurable": {"thread_id": selected_tid}},
                        selected_cp.values,
                    )
                    print(f"\n  ✅ 已成功回退到检查点 #{cp_choice}!")
                    print(f"    线程: {selected_tid}")
                    print(f"    内容: {selected_cp.values.get('content', '')[:80]}")
                    print(f"    风险等级: {selected_cp.values.get('risk_level', '未知')}")
                    print(f"  💡 提示:您可以重新提交内容或继续审核流程")
                except Exception as e:
                    print(f"\n  ❌ 回退失败: {e}")
            else:
                print("  ⚠️ 无效的编号!")

    # 关闭数据库连接
    conn.close()
    print("\n  📦 数据库连接已关闭")


# ========================= 启动演示 =========================
if __name__ == "__main__":
    try:
        main()
    finally:
        # 无论是否正常退出,都生成流程图
        generate_flowchart()

状态图:

在这里插入图片描述

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐