1. 案例标题

电商智能客服多 Agent 系统:基于 LangGraph Checkpoint + Interrupt + MCP + Neo4j 构建可中断、可恢复、可升级的人机协同客服 Workflow


2. 一句话说明

通过 LangGraph Checkpoint + Interrupt + Multi-Agent 技术,解决电商客服场景下 AI 处理复杂工单时无法暂停等待人工审批、进程崩溃后无法恢复、多轮对话状态丢失 的核心问题,并借助 Neo4j 知识图谱、Redis 缓存、MCP 工具协议、PostgreSQL 持久化 构建生产级智能客服系统。


3. 为什么需要这个技术

3.1 传统客服系统的困境

用户发起咨询:"我买的手机屏幕碎了,订单号 20260901-88721,我要退款"

传统方案(无状态函数调用):
┌─────────────────────────────────────────────────────┐
│  Step 1: 意图识别 → "退款"                          │
│  Step 2: 查询订单 → 找到订单,金额 ¥5,999            │
│  Step 3: 判断金额 > ¥500 → 需要主管审批              │
│  Step 4: 等待审批...                                 │
│                                                     │
│  ❌ 此时:                                           │
│     - 用户关闭了页面 → 状态丢失,一切重来              │
│     - 服务器重启 → 所有进行中的对话全部丢失             │
│     - 审批需要 2 小时 → 进程一直阻塞                  │
│     - 1000 个用户同时等待 → 1000 个线程阻塞           │
│     - 审批被拒绝 → 无法优雅回退到上一步               │
└─────────────────────────────────────────────────────┘

3.2 核心痛点

痛点传统方案后果
用户中途离开内存中的对话状态丢失用户回来后必须重新描述问题
等待人工审批线程阻塞 / 轮询资源浪费,无法扩展
进程崩溃所有进行中的工单丢失用户投诉,需要人工重新处理
多步骤流程用 if-else 硬编码无法暂停、恢复、回退
多轮对话手动管理上下文代码复杂,容易出 bug
工具调用每个 Agent 各自实现工具无法复用,维护成本高

3.3 为什么需要这套技术栈

LangGraph Checkpoint  →  解决"状态持久化 + 崩溃恢复"
LangGraph Interrupt   →  解决"暂停等待人工 + 无阻塞"
Multi-Agent           →  解决"不同问题需要不同专业能力"
Neo4j                 →  解决"产品知识关联查询"
Redis                 →  解决"会话缓存 + 限流 + 实时通知"
PostgreSQL            →  解决"持久化 Checkpoint + 业务数据"
MCP                   →  解决"工具标准化 + 跨 Agent 复用"
FastAPI               →  解决"异步 API + SSE 流式输出"

没有 Checkpoint,Agent 中断后必须从头执行。没有 Interrupt,人工审批只能阻塞等待。没有 Multi-Agent,一个巨型 Agent 无法胜任所有任务。


4. 场景设计

4.1 业务流程

用户发送消息:"我买的手机屏幕碎了,订单 20260901-88721,我要退款"
        │
        ▼
┌──────────────────┐
│  FastAPI 接入层   │  ← SSE 流式响应 / Redis 限流
└────────┬─────────┘
         ▼
┌──────────────────┐
│  Router Agent    │  ← 意图识别:退款 / 查单 / 咨询 / 投诉
└────────┬─────────┘
         ▼
    ┌────┴────┬──────────┬──────────────┐
    ▼         ▼          ▼              ▼
┌───────┐ ┌───────┐ ┌────────┐  ┌──────────┐
│ Order │ │Refund │ │Know-   │  │Escalation│
│ Agent │ │ Agent │ │ledge   │  │  Agent   │
│       │ │       │ │Agent   │  │          │
│查订单 │ │处理退款│ │产品问答│  │投诉升级  │
│查物流 │ │       │ │(Neo4j) │  │          │
└───┬───┘ └───┬───┘ └───┬────┘  └────┬─────┘
    │         │         │            │
    │    ┌────┴────┐    │       始终触发
    │    │金额>500?│    │       Interrupt
    │    └────┬────┘    │            │
    │    Yes  │  No     │            │
    │    ┌────┘  └───┐  │            │
    │    ▼           ▼  ▼            ▼
    │ ┌────────┐  自动  自动     ┌────────┐
    │ │INTERRUPT│  处理  回复     │INTERRUPT│
    │ │暂停等待 │       │       │人工介入 │
    │ └───┬────┘       │       └───┬────┘
    │     ▼            │           ▼
    │ ┌────────┐       │      ┌────────┐
    │ │主管审批 │       │      │人工客服 │
    │ │(人工)  │       │      │处理    │
    │ └───┬────┘       │      └───┬────┘
    │     ▼            │          ▼
    │  Resume          │       Resume
    │  继续执行         │       继续执行
    │     ▼            ▼          ▼
    └────►┌────────────────────────┐◄────┘
          │   生成最终回复          │
          │   更新工单状态          │
          │   Redis 通知用户        │
          └────────────────────────┘

4.2 关键场景示例

场景触发条件行为
普通查单用户问"我的快递到哪了"Order Agent 直接回复,无需中断
小额退款退款金额 ≤ ¥500Refund Agent 自动处理
大额退款退款金额 > ¥500Interrupt → 主管审批 → Resume
产品咨询用户问"这款手机支持5G吗"Knowledge Agent 查 Neo4j 回复
投诉升级用户表达强烈不满始终 Interrupt → 人工客服介入
进程崩溃任意时刻从 Checkpoint 恢复
用户离开用户关闭页面状态保存,下次回来继续

5. 技术映射表

业务需求技术能力具体实现
保存对话/工单状态LangGraph CheckpointAsyncPostgresSaver / MemorySaver
暂停等待人工审批LangGraph Interruptinterrupt() 函数
人工审批/介入Human-in-the-loopFastAPI 审批接口 + Command(resume=...)
恢复暂停的 WorkflowLangGraph CommandCommand(resume={"decision": ...})
多用户会话隔离Threadthread_id = f"cs-{user_id}-{ticket_id}"
意图路由Multi-Agent Routeradd_conditional_edges
专业领域处理Multi-Agent 子图每个 Agent 独立子图
产品知识查询Neo4j 知识图谱Cypher 查询产品/配件/兼容性
会话缓存 + 限流Redisredis-py + Rate Limiter
持久化存储PostgreSQLlanggraph-checkpoint-postgres
工具标准化MCP (Model Context Protocol)MCP Tool Server
API 接入 + 流式输出FastAPI + SSEStreamingResponse
状态管理LangGraph StateTypedDict + Annotated

6. 最小系统架构

                        ┌─────────────────────────────────────────┐
                        │            FastAPI Gateway              │
                        │  ┌─────────┐  ┌────────┐  ┌─────────┐  │
                        │  │ POST    │  │ POST   │  │ GET     │  │
                        │  │ /chat   │  │/approve│  │ /events │  │
                        │  │(SSE)    │  │(人工)  │  │(SSE推送)│  │
                        │  └────┬────┘  └───┬────┘  └────┬────┘  │
                        └───────┼───────────┼─────────────┼──────┘
                                │           │             │
                    ┌───────────▼───────────▼─────────────▼──────┐
                    │              Redis Layer                    │
                    │  • 会话缓存 (session:{thread_id})           │
                    │  • 限流 (rate:{user_id})                    │
                    │  • 审批队列 (pending_approvals)             │
                    │  • Pub/Sub 实时通知                          │
                    └───────────────────┬────────────────────────┘
                                        │
                    ┌───────────────────▼────────────────────────┐
                    │           LangGraph Workflow                │
                    │                                            │
                    │   START                                    │
                    │     │                                      │
                    │     ▼                                      │
                    │  [Router Agent] ─── 意图分类                │
                    │     │                                      │
                    │     ├──► [Order Agent] ──── 查单/查物流     │
                    │     │         │                            │
                    │     │         └──► MCP: order_tools        │
                    │     │                                      │
                    │     ├──► [Refund Agent] ── 退款处理         │
                    │     │         │                            │
                    │     │         ├──► MCP: refund_tools       │
                    │     │         │                            │
                    │     │         └──► interrupt() ──► 人工审批 │
                    │     │                    │                  │
                    │     │              Command(resume)          │
                    │     │                    │                  │
                    │     │                    ▼                  │
                    │     │              继续执行                  │
                    │     │                                      │
                    │     ├──► [Knowledge Agent] ─ 产品问答       │
                    │     │         │                            │
                    │     │         └──► Neo4j 知识图谱           │
                    │     │                                      │
                    │     └──► [Escalation Agent] ─ 投诉升级      │
                    │               │                            │
                    │               └──► interrupt() ──► 人工客服 │
                    │                                            │
                    │  [Checkpoint: PostgreSQL]                   │
                    │  • 每个 Node 执行后自动保存                  │
                    │  • Interrupt 时保存完整状态                  │
                    │  • 按 thread_id 隔离                       │
                    │                                            │
                    │     END                                    │
                    └────────────────────────────────────────────┘
                                        │
                    ┌───────────────────▼────────────────────────┐
                    │           MCP Tool Server                   │
                    │  ┌──────────────┐  ┌───────────────┐       │
                    │  │ order_tools  │  │ refund_tools  │       │
                    │  │ • get_order  │  │ • calc_refund │       │
                    │  │ • get_logis  │  │ • exec_refund │       │
                    │  └──────┬───────┘  └──────┬────────┘       │
                    └─────────┼─────────────────┼────────────────┘
                              │                 │
                    ┌─────────▼─────────────────▼────────────────┐
                    │           PostgreSQL                        │
                    │  • orders 表 (订单数据)                      │
                    │  • checkpoints 表 (LangGraph 状态)          │
                    │  • approval_logs 表 (审批记录)               │
                    └────────────────────────────────────────────┘

每个组件存在的原因

组件为什么需要
FastAPI异步 API 网关,支持 SSE 流式输出,不阻塞等待
Redis会话缓存避免重复查询;限流防刷;Pub/Sub 实现审批通知实时推送
LangGraph核心编排引擎,原生支持 Checkpoint/Interrupt/Multi-Agent
Router Agent将用户意图分发到正确的专业 Agent,避免单一 Agent 过载
Order/Refund/Knowledge/Escalation Agent各司其职,独立可测试,可独立升级
Neo4j产品之间存在复杂关系(配件、兼容性、替代品),关系型数据库难以表达
PostgreSQL持久化 Checkpoint(进程重启不丢状态)+ 业务数据(订单、审批记录)
MCP工具标准化,Order Agent 和 Refund Agent 可共享同一个 MCP Server 的订单工具

7. 项目目录

ecommerce-cs-agent/
│
├── app/
│   ├── __init__.py
│   ├── config.py                      # 全局配置
│   ├── state.py                       # LangGraph State 定义
│   ├── graph.py                       # 主 Graph 构建(核心)
│   │
│   ├── agents/
│   │   ├── __init__.py
│   │   ├── router.py                  # 意图路由节点
│   │   ├── order_agent.py             # 订单查询 Agent
│   │   ├── refund_agent.py            # 退款处理 Agent(含 Interrupt)
│   │   ├── knowledge_agent.py         # 产品知识 Agent(Neo4j)
│   │   └── escalation_agent.py        # 投诉升级 Agent(含 Interrupt)
│   │
│   ├── tools/
│   │   ├── __init__.py
│   │   ├── mcp_server.py             # MCP Tool Server
│   │   ├── order_tools.py            # 订单工具(MCP 注册)
│   │   ├── refund_tools.py           # 退款工具(MCP 注册)
│   │   └── knowledge_tools.py        # Neo4j 查询工具
│   │
│   ├── infra/
│   │   ├── __init__.py
│   │   ├── redis_client.py           # Redis 连接 + 缓存 + 限流
│   │   ├── neo4j_client.py           # Neo4j 连接 + 查询
│   │   ├── postgres.py               # PostgreSQL 连接
│   │   └── checkpointer.py           # Checkpointer 工厂
│   │
│   ├── api/
│   │   ├── __init__.py
│   │   ├── routes.py                 # FastAPI 路由
│   │   └── schemas.py               # Pydantic 请求/响应模型
│   │
│   └── models/
│       ├── __init__.py
│       └── enums.py                  # 枚举定义
│
├── tests/
│   ├── __init__.py
│   ├── test_checkpoint.py
│   ├── test_interrupt.py
│   ├── test_multi_agent.py
│   └── test_resume.py
│
├── scripts/
│   ├── seed_neo4j.py                 # 初始化知识图谱
│   └── seed_postgres.py             # 初始化订单数据
│
├── requirements.txt
├── docker-compose.yml
├── README.md
├── main.py                           # 生产启动入口
└── run_demo.py                       # 本地演示(无需外部依赖)

8. 核心代码

8.1 依赖版本

# requirements.txt
langgraph>=0.4.10
langgraph-checkpoint-postgres>=2.0.0
langchain-core>=0.3.50
langchain-openai>=0.3.20
fastapi>=0.115.0
uvicorn[standard]>=0.34.0
redis>=5.2.0
neo4j>=5.26.0
psycopg[binary]>=3.2.0
pydantic>=2.10.0
mcp>=1.0.0
python-dotenv>=1.0.0
sse-starlette>=2.0.0

⚠️ 版本敏感提示:本案例基于 langgraph>=0.4.10 编写。interrupt()Command(resume=...) API 在 0.4.0+ 版本中引入。如果你使用更早版本,请查阅对应文档。


8.2 状态定义 — app/state.py

"""LangGraph 全局状态定义"""

from __future__ import annotations

from typing import Annotated, Any, Literal, TypedDict
from operator import add


class CustomerMessage(TypedDict):
    """用户消息"""
    role: Literal["user", "assistant", "system", "tool"]
    content: str


class AgentState(TypedDict):
    """
    整个 Multi-Agent Workflow 的共享状态。
    每个节点读取并返回部分更新,LangGraph 自动合并。
    """
    # 对话历史(追加式)
    messages: Annotated[list[CustomerMessage], add]

    # 路由结果
    intent: str                          # "order" | "refund" | "knowledge" | "escalation"
    confidence: float                    # 意图识别置信度

    # 业务上下文
    user_id: str
    ticket_id: str
    order_id: str | None
    order_info: dict[str, Any] | None    # 订单详情(Order Agent 填充)

    # 退款相关
    refund_amount: float
    refund_reason: str
    refund_decision: str | None          # "approved" | "rejected" | "pending"
    refund_decision_by: str | None       # 审批人

    # 投诉/升级相关
    escalation_reason: str | None
    human_agent_id: str | None           # 介入的人工客服

    # 最终回复
    final_response: str | None

    # 执行状态
    status: str                          # "running" | "interrupted" | "completed" | "failed"
    error: str | None


def create_initial_state(
    user_id: str,
    ticket_id: str,
    user_message: str,
) -> AgentState:
    """创建初始状态"""
    return AgentState(
        messages=[{"role": "user", "content": user_message}],
        intent="",
        confidence=0.0,
        user_id=user_id,
        ticket_id=ticket_id,
        order_id=None,
        order_info=None,
        refund_amount=0.0,
        refund_reason="",
        refund_decision=None,
        refund_decision_by=None,
        escalation_reason=None,
        human_agent_id=None,
        final_response=None,
        status="running",
        error=None,
    )

8.3 基础设施 — app/infra/

app/infra/checkpointer.py
"""Checkpointer 工厂:开发用内存,生产用 PostgreSQL"""

from __future__ import annotations

from langgraph.checkpoint.memory import MemorySaver

try:
    from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
    HAS_PG_CHECKPOINT = True
except ImportError:
    HAS_PG_CHECKPOINT = False


def get_checkpointer(mode: str = "memory", pg_dsn: str = ""):
    """
    获取 Checkpointer 实例。

    Args:
        mode: "memory" 用于开发/测试,"postgres" 用于生产
        pg_dsn: PostgreSQL 连接字符串,如
                "postgresql://user:pass@localhost:5432/cs_agent"
    """
    if mode == "postgres" and HAS_PG_CHECKPOINT:
        # 生产环境:持久化到 PostgreSQL
        # LangGraph 会自动创建 checkpoints / checkpoint_writes 表
        return AsyncPostgresSaver.from_conn_string(pg_dsn)
    else:
        # 开发/演示:内存存储(进程重启后丢失)
        return MemorySaver()
app/infra/redis_client.py
"""Redis 客户端:会话缓存 + 限流 + 审批通知"""

from __future__ import annotations

import json
from typing import Any

import redis.asyncio as redis


class RedisManager:
    def __init__(self, url: str = "redis://localhost:6379/0"):
        self._url = url
        self._client: redis.Redis | None = None

    async def connect(self) -> None:
        self._client = redis.from_url(self._url, decode_responses=True)

    async def close(self) -> None:
        if self._client:
            await self._client.aclose()

    # ── 会话缓存 ──────────────────────────────────────────
    async def cache_session(self, thread_id: str, data: dict[str, Any], ttl: int = 3600) -> None:
        """缓存当前会话状态(快速读取,避免每次查 PG)"""
        await self._client.set(
            f"session:{thread_id}",
            json.dumps(data, ensure_ascii=False),
            ex=ttl,
        )

    async def get_session(self, thread_id: str) -> dict | None:
        raw = await self._client.get(f"session:{thread_id}")
        return json.loads(raw) if raw else None

    # ── 限流 ─────────────────────────────────────────────
    async def check_rate_limit(self, user_id: str, max_calls: int = 10, window: int = 60) -> bool:
        """滑动窗口限流:每个用户每分钟最多 max_calls 次请求"""
        key = f"rate:{user_id}"
        current = await self._client.incr(key)
        if current == 1:
            await self._client.expire(key, window)
        return current <= max_calls

    # ── 审批队列 ──────────────────────────────────────────
    async def push_pending_approval(self, ticket_id: str, info: dict) -> None:
        """将待审批工单推入队列,供人工审批面板消费"""
        await self._client.lpush(
            "pending_approvals",
            json.dumps({"ticket_id": ticket_id, **info}, ensure_ascii=False),
        )

    async def publish_event(self, channel: str, data: dict) -> None:
        """Pub/Sub:实时通知前端"""
        await self._client.publish(channel, json.dumps(data, ensure_ascii=False))


# 全局单例
redis_manager = RedisManager()
app/infra/neo4j_client.py
"""Neo4j 客户端:产品知识图谱查询"""

from __future__ import annotations

from neo4j import AsyncGraphDatabase


class Neo4jManager:
    def __init__(self, uri: str = "bolt://localhost:7687", user: str = "neo4j", password: str = "password"):
        self._driver = AsyncGraphDatabase.driver(uri, auth=(user, password))

    async def close(self) -> None:
        await self._driver.close()

    async def query_product(self, product_name: str) -> list[dict]:
        """查询产品信息及其关联(配件、替代品、兼容性)"""
        query = """
        MATCH (p:Product)
        WHERE p.name CONTAINS $name
        OPTIONAL MATCH (p)-[:HAS_ACCESSORY]->(acc:Product)
        OPTIONAL MATCH (p)-[:COMPATIBLE_WITH]->(comp:Product)
        OPTIONAL MATCH (p)-[:ALTERNATIVE_TO]->(alt:Product)
        RETURN p,
               collect(DISTINCT acc) AS accessories,
               collect(DISTINCT comp) AS compatibles,
               collect(DISTINCT alt) AS alternatives
        LIMIT 5
        """
        async with self._driver.session() as session:
            result = await session.run(query, name=product_name)
            records = await result.data()
            return [
                {
                    "product": dict(r["p"]),
                    "accessories": [dict(a) for a in r["accessories"]],
                    "compatibles": [dict(c) for c in r["compatibles"]],
                    "alternatives": [dict(a) for a in r["alternatives"]],
                }
                for r in records
            ]

    async def query_by_category(self, category: str) -> list[dict]:
        """按品类查询产品"""
        query = """
        MATCH (p:Product {category: $category})
        RETURN p ORDER BY p.price DESC LIMIT 10
        """
        async with self._driver.session() as session:
            result = await session.run(query, category=category)
            return [dict(r["p"]) for r in await result.data()]


neo4j_manager = Neo4jManager()

8.4 Agent 节点实现 — app/agents/

app/agents/router.py — 意图路由
"""Router Agent:意图识别 + 路由分发"""

from __future__ import annotations

import re

from app.state import AgentState


# 【示例业务规则】实际生产中应使用 LLM 或分类模型
INTENT_KEYWORDS: dict[str, list[str]] = {
    "refund": ["退款", "退货", "退钱", "不想要了", "申请退"],
    "order": ["订单", "快递", "物流", "发货", "到货", "到哪了"],
    "knowledge": ["参数", "规格", "支持", "兼容", "配件", "推荐", "怎么样"],
    "escalation": ["投诉", "差评", "举报", "经理", "工商", "315", "太差了"],
}


def router_node(state: AgentState) -> dict:
    """
    意图识别节点。
    根据用户最新消息分类意图。

    【示例业务规则】这里使用关键词匹配演示。
    生产环境应替换为 LLM 调用或专用分类模型。
    """
    last_message = state["messages"][-1]["content"].lower()

    best_intent = "order"  # 默认走订单查询
    best_score = 0.0

    for intent, keywords in INTENT_KEYWORDS.items():
        score = sum(1 for kw in keywords if kw in last_message)
        if score > best_score:
            best_score = score
            best_intent = intent

    confidence = min(best_score / 3.0, 1.0)  # 简单归一化

    # 提取订单号(如果有)
    order_match = re.search(r"(\d{8}-\d{5})", last_message)
    order_id = order_match.group(1) if order_match else state.get("order_id")

    return {
        "intent": best_intent,
        "confidence": confidence,
        "order_id": order_id,
    }


def route_by_intent(state: AgentState) -> str:
    """条件边:根据意图路由到不同 Agent"""
    intent = state.get("intent", "order")
    mapping = {
        "refund": "refund_agent",
        "order": "order_agent",
        "knowledge": "knowledge_agent",
        "escalation": "escalation_agent",
    }
    return mapping.get(intent, "order_agent")
app/agents/order_agent.py — 订单查询
"""Order Agent:查询订单和物流信息"""

from __future__ import annotations

from app.state import AgentState

# 【示例业务规则】模拟订单数据库
# 生产环境通过 MCP 调用真实订单服务
MOCK_ORDERS: dict[str, dict] = {
    "20260901-88721": {
        "product": "iPhone 17 Pro",
        "amount": 5999.0,
        "status": "已签收",
        "logistics": "顺丰 SF1234567890 - 已签收(2026-09-03)",
        "buyer": "user_001",
    },
    "20260902-11034": {
        "product": "AirPods Pro 3",
        "amount": 1899.0,
        "status": "运输中",
        "logistics": "中通 ZT9876543210 - 运输中(预计明天到达)",
        "buyer": "user_002",
    },
}


def order_agent_node(state: AgentState) -> dict:
    """
    订单查询节点。
    根据订单号查询订单详情和物流信息。
    """
    order_id = state.get("order_id")

    if not order_id:
        return {
            "final_response": "请提供您的订单号,我帮您查询。",
            "status": "completed",
            "messages": [{"role": "assistant", "content": "请提供您的订单号,我帮您查询。"}],
        }

    order = MOCK_ORDERS.get(order_id)

    if not order:
        return {
            "final_response": f"未找到订单 {order_id},请核实订单号。",
            "status": "completed",
            "messages": [{"role": "assistant", "content": f"未找到订单 {order_id},请核实订单号。"}],
        }

    response = (
        f"📦 订单 {order_id}\n"
        f"商品:{order['product']}\n"
        f"金额:¥{order['amount']:.2f}\n"
        f"状态:{order['status']}\n"
        f"物流:{order['logistics']}"
    )

    return {
        "order_info": order,
        "final_response": response,
        "status": "completed",
        "messages": [{"role": "assistant", "content": response}],
    }
app/agents/refund_agent.py — 退款处理(核心:Interrupt)
"""Refund Agent:退款处理 + 大额退款人工审批(Interrupt)"""

from __future__ import annotations

from langgraph.types import interrupt

from app.state import AgentState

# 【示例业务规则】退款审批阈值
REFUND_AUTO_APPROVE_THRESHOLD = 500.0  # ¥500 以下自动审批

# 复用 Order Agent 的模拟数据
from app.agents.order_agent import MOCK_ORDERS


def refund_agent_node(state: AgentState) -> dict:
    """
    退款处理节点。

    流程:
    1. 查询订单信息
    2. 计算退款金额
    3. 判断是否需要人工审批
       - 金额 <= 500:自动通过
       - 金额 >  500:调用 interrupt() 暂停,等待人工审批
    4. 执行退款
    """
    order_id = state.get("order_id")
    user_message = state["messages"][-1]["content"]

    # ── Step 1: 查询订单 ──
    if not order_id or order_id not in MOCK_ORDERS:
        return {
            "final_response": "无法找到相关订单,请提供订单号。",
            "status": "completed",
            "messages": [{"role": "assistant", "content": "无法找到相关订单,请提供订单号。"}],
        }

    order = MOCK_ORDERS[order_id]
    refund_amount = order["amount"]  # 【示例业务规则】全额退款

    # ── Step 2: 判断是否需要人工审批 ──
    if refund_amount > REFUND_AUTO_APPROVE_THRESHOLD:
        # ★★★ 核心:调用 interrupt() 暂停 Workflow ★★★
        #
        # 调用 interrupt() 后:
        #   1. 当前节点执行立即停止
        #   2. LangGraph 将当前 State 保存到 Checkpoint
        #   3. Graph 执行暂停,控制权返回给调用方
        #   4. 调用方(FastAPI)收到中断信息,可以通知人工审批
        #   5. 人工审批后,通过 Command(resume=...) 恢复执行
        #   6. 恢复后,interrupt() 返回 resume 传入的值,继续往下执行
        #
        human_decision = interrupt({
            "type": "refund_approval",
            "ticket_id": state["ticket_id"],
            "order_id": order_id,
            "product": order["product"],
            "refund_amount": refund_amount,
            "reason": user_message,
            "question": f"退款金额 ¥{refund_amount:.2f} 超过自动审批阈值 ¥{REFUND_AUTO_APPROVE_THRESHOLD:.2f},需要主管审批。是否批准?",
        })

        # ── 以下代码在 Resume 后执行 ──
        # human_decision 是 Command(resume=...) 传入的值
        # 例如:{"decision": "approved", "approver": "张主管", "comment": "同意退款"}

        decision = human_decision.get("decision", "rejected")
        approver = human_decision.get("approver", "unknown")
        comment = human_decision.get("comment", "")

        if decision == "approved":
            response = (
                f"✅ 退款已批准!\n"
                f"订单:{order_id}\n"
                f"商品:{order['product']}\n"
                f"退款金额:¥{refund_amount:.2f}\n"
                f"审批人:{approver}\n"
                f"备注:{comment}\n"
                f"预计 3-5 个工作日到账。"
            )
            return {
                "refund_amount": refund_amount,
                "refund_decision": "approved",
                "refund_decision_by": approver,
                "final_response": response,
                "status": "completed",
                "messages": [{"role": "assistant", "content": response}],
            }
        else:
            response = (
                f"❌ 退款申请未通过。\n"
                f"订单:{order_id}\n"
                f"审批人:{approver}\n"
                f"原因:{comment}\n"
                f"如有疑问,您可以拨打客服热线 400-000-0000。"
            )
            return {
                "refund_amount": refund_amount,
                "refund_decision": "rejected",
                "refund_decision_by": approver,
                "final_response": response,
                "status": "completed",
                "messages": [{"role": "assistant", "content": response}],
            }
    else:
        # ── 小额退款:自动通过 ──
        response = (
            f"✅ 退款已自动处理。\n"
            f"订单:{order_id}\n"
            f"退款金额:¥{refund_amount:.2f}\n"
            f"预计 1-3 个工作日到账。"
        )
        return {
            "refund_amount": refund_amount,
            "refund_decision": "approved",
            "refund_decision_by": "system",
            "final_response": response,
            "status": "completed",
            "messages": [{"role": "assistant", "content": response}],
        }
app/agents/knowledge_agent.py — 产品知识(Neo4j)
"""Knowledge Agent:基于 Neo4j 知识图谱的产品问答"""

from __future__ import annotations

from app.state import AgentState

# 【示例业务规则】模拟知识图谱查询结果
# 生产环境通过 neo4j_manager.query_product() 查询
MOCK_KNOWLEDGE: dict[str, dict] = {
    "iphone": {
        "product": {"name": "iPhone 17 Pro", "category": "手机", "price": 5999, "specs": "A19 Pro / 8GB / 256GB"},
        "accessories": [
            {"name": "MagSafe 保护壳", "price": 299},
            {"name": "USB-C 快充线", "price": 149},
        ],
        "compatibles": [
            {"name": "AirPods Pro 3", "type": "蓝牙耳机"},
            {"name": "Apple Watch Ultra 3", "type": "智能手表"},
        ],
        "alternatives": [
            {"name": "iPhone 17", "price": 4999, "diff": "无 Pro 摄像头"},
        ],
    },
    "airpods": {
        "product": {"name": "AirPods Pro 3", "category": "耳机", "price": 1899, "specs": "H3 芯片 / ANC / USB-C"},
        "accessories": [{"name": "硅胶耳塞套装", "price": 49}],
        "compatibles": [{"name": "iPhone 17 Pro", "type": "手机"}],
        "alternatives": [{"name": "AirPods 4", "price": 999, "diff": "无主动降噪"}],
    },
}


def knowledge_agent_node(state: AgentState) -> dict:
    """
    产品知识查询节点。
    从知识图谱中检索产品信息、配件、兼容性、替代品。

    【示例业务规则】这里使用 Mock 数据。
    生产环境替换为:
        result = await neo4j_manager.query_product(keyword)
    """
    user_message = state["messages"][-1]["content"].lower()

    # 简单关键词匹配找到相关产品
    matched = None
    for key, data in MOCK_KNOWLEDGE.items():
        if key in user_message:
            matched = data
            break

    if not matched:
        return {
            "final_response": "抱歉,暂未找到相关产品信息。您可以描述得更具体一些吗?",
            "status": "completed",
            "messages": [{"role": "assistant", "content": "抱歉,暂未找到相关产品信息。"}],
        }

    p = matched["product"]
    parts = [f"🔍 {p['name']}{p['category']})"]
    parts.append(f"价格:¥{p['price']}")
    if "specs" in p:
        parts.append(f"规格:{p['specs']}")

    if matched["accessories"]:
        parts.append("\n📎 推荐配件:")
        for acc in matched["accessories"]:
            parts.append(f"  • {acc['name']} - ¥{acc['price']}")

    if matched["compatibles"]:
        parts.append("\n🔗 兼容设备:")
        for c in matched["compatibles"]:
            parts.append(f"  • {c['name']}{c['type']})")

    if matched["alternatives"]:
        parts.append("\n🔄 替代选择:")
        for alt in matched["alternatives"]:
            diff = f"({alt['diff']})" if "diff" in alt else ""
            parts.append(f"  • {alt['name']} - ¥{alt['price']} {diff}")

    response = "\n".join(parts)

    return {
        "final_response": response,
        "status": "completed",
        "messages": [{"role": "assistant", "content": response}],
    }
app/agents/escalation_agent.py — 投诉升级(始终 Interrupt)
"""Escalation Agent:投诉/升级处理,始终需要人工介入"""

from __future__ import annotations

from langgraph.types import interrupt

from app.state import AgentState


def escalation_agent_node(state: AgentState) -> dict:
    """
    投诉升级节点。

    【示例业务规则】所有投诉类请求必须转人工处理。
    无论投诉内容是什么,都触发 Interrupt,等待人工客服介入。
    """
    user_message = state["messages"][-1]["content"]

    # ★★★ 始终触发 Interrupt:投诉必须人工处理 ★★★
    human_result = interrupt({
        "type": "escalation",
        "ticket_id": state["ticket_id"],
        "user_id": state["user_id"],
        "reason": user_message,
        "question": "用户发起投诉/升级请求,需要人工客服介入处理。",
    })

    # ── Resume 后执行 ──
    # human_result 示例:
    # {
    #     "agent_id": "cs_agent_007",
    #     "resolution": "已为用户安排补发 + 50元优惠券",
    #     "internal_note": "用户情绪激动,已安抚"
    # }

    agent_id = human_result.get("agent_id", "unknown")
    resolution = human_result.get("resolution", "问题已记录,我们将尽快处理。")

    response = (
        f"📋 您的反馈已由客服 {agent_id} 处理。\n"
        f"处理结果:{resolution}\n"
        f"工单号:{state['ticket_id']}\n"
        f"如有其他问题,随时联系我们。"
    )

    return {
        "human_agent_id": agent_id,
        "escalation_reason": user_message,
        "final_response": response,
        "status": "completed",
        "messages": [{"role": "assistant", "content": response}],
    }

8.5 主 Graph 构建 — app/graph.py(核心)

"""
主 Graph 构建:Multi-Agent Workflow + Checkpoint + Interrupt

这是整个系统的核心文件。
"""

from __future__ import annotations

from langgraph.graph import StateGraph, START, END

from app.state import AgentState
from app.agents.router import router_node, route_by_intent
from app.agents.order_agent import order_agent_node
from app.agents.refund_agent import refund_agent_node
from app.agents.knowledge_agent import knowledge_agent_node
from app.agents.escalation_agent import escalation_agent_node


def build_graph(checkpointer):
    """
    构建完整的 Multi-Agent 客服 Graph。

    拓扑:
        START → router → (conditional) → order / refund / knowledge / escalation → END

    关键设计:
    - router 节点做意图分类
    - conditional_edges 根据意图路由到不同 Agent 节点
    - refund_agent 和 escalation_agent 内部调用 interrupt()
    - checkpointer 在每个节点执行后自动保存状态
    """
    graph = StateGraph(AgentState)

    # ── 注册节点 ──
    graph.add_node("router", router_node)
    graph.add_node("order_agent", order_agent_node)
    graph.add_node("refund_agent", refund_agent_node)
    graph.add_node("knowledge_agent", knowledge_agent_node)
    graph.add_node("escalation_agent", escalation_agent_node)

    # ── 边 ──
    # START → Router
    graph.add_edge(START, "router")

    # Router → 条件路由
    graph.add_conditional_edges(
        "router",
        route_by_intent,
        {
            "order_agent": "order_agent",
            "refund_agent": "refund_agent",
            "knowledge_agent": "knowledge_agent",
            "escalation_agent": "escalation_agent",
        },
    )

    # 每个 Agent → END
    graph.add_edge("order_agent", END)
    graph.add_edge("refund_agent", END)
    graph.add_edge("knowledge_agent", END)
    graph.add_edge("escalation_agent", END)

    # ── 编译(绑定 Checkpointer)──
    # checkpointer 使得每个节点执行后自动保存状态
    # 当遇到 interrupt() 时,状态被持久化,等待 resume
    app = graph.compile(checkpointer=checkpointer)

    return app

8.6 MCP Tool Server — app/tools/mcp_server.py

"""
MCP (Model Context Protocol) Tool Server

将订单查询、退款计算等工具通过 MCP 协议暴露,
使得任何 Agent(包括外部 Agent)都可以调用这些工具。

生产环境中,MCP Server 独立部署,通过 stdio 或 HTTP 与 Agent 通信。
这里展示核心注册逻辑。
"""

from __future__ import annotations

import json
from typing import Any

# MCP SDK(mcp >= 1.0.0)
from mcp.server import Server
from mcp.server.stdio import stdio_server
from mcp.types import Tool, TextContent

from app.agents.order_agent import MOCK_ORDERS

server = Server("ecommerce-cs-tools")


@server.list_tools()
async def list_tools() -> list[Tool]:
    """注册所有可用工具"""
    return [
        Tool(
            name="get_order",
            description="根据订单号查询订单详情(商品、金额、状态、物流)",
            inputSchema={
                "type": "object",
                "properties": {
                    "order_id": {"type": "string", "description": "订单号,如 20260901-88721"},
                },
                "required": ["order_id"],
            },
        ),
        Tool(
            name="calculate_refund",
            description="计算退款金额(含运费、优惠券抵扣等)",
            inputSchema={
                "type": "object",
                "properties": {
                    "order_id": {"type": "string"},
                    "refund_type": {
                        "type": "string",
                        "enum": ["full", "partial"],
                        "description": "全额退款或部分退款",
                    },
                },
                "required": ["order_id", "refund_type"],
            },
        ),
    ]


@server.call_tool()
async def call_tool(name: str, arguments: dict[str, Any]) -> list[TextContent]:
    """执行工具调用"""
    if name == "get_order":
        order_id = arguments["order_id"]
        order = MOCK_ORDERS.get(order_id)
        if order:
            return [TextContent(type="text", text=json.dumps(order, ensure_ascii=False))]
        return [TextContent(type="text", text=f"订单 {order_id} 不存在")]

    elif name == "calculate_refund":
        order_id = arguments["order_id"]
        refund_type = arguments.get("refund_type", "full")
        order = MOCK_ORDERS.get(order_id)
        if not order:
            return [TextContent(type="text", text=f"订单 {order_id} 不存在")]

        amount = order["amount"]
        if refund_type == "partial":
            amount *= 0.8  # 【示例业务规则】部分退款扣 20%

        return [TextContent(type="text", text=json.dumps({
            "order_id": order_id,
            "refund_amount": amount,
            "refund_type": refund_type,
        }, ensure_ascii=False))]

    return [TextContent(type="text", text=f"未知工具: {name}")]


async def run_mcp_server():
    """启动 MCP Server(stdio 模式)"""
    async with stdio_server() as (read, write):
        await server.run(read, write, server.create_initialization_options())


# 启动方式:python -m app.tools.mcp_server
if __name__ == "__main__":
    import asyncio
    asyncio.run(run_mcp_server())

8.7 FastAPI 接入层 — app/api/routes.py

"""FastAPI 路由:对话、审批、事件推送"""

from __future__ import annotations

import json
import uuid
from typing import AsyncGenerator

from fastapi import APIRouter, HTTPException
from fastapi.responses import StreamingResponse
from langgraph.types import Command

from app.api.schemas import ChatRequest, ApprovalRequest
from app.state import create_initial_state

router = APIRouter()

# graph 实例在 main.py 中初始化后注入
_graph = None
_redis = None


def set_dependencies(graph, redis_mgr):
    global _graph, _redis
    _graph = graph
    _redis = redis_mgr


# ── 1. 用户对话(SSE 流式) ─────────────────────────────────
@router.post("/chat")
async def chat(req: ChatRequest):
    """
    用户发送消息,触发 Multi-Agent Workflow。

    如果 Workflow 遇到 interrupt(),返回中断信息(包含待审批详情)。
    如果正常完成,返回最终回复。
    """
    # 限流检查
    if _redis and not await _redis.check_rate_limit(req.user_id):
        raise HTTPException(status_code=429, detail="请求过于频繁,请稍后再试")

    thread_id = req.thread_id or f"cs-{req.user_id}-{uuid.uuid4().hex[:8]}"
    config = {"configurable": {"thread_id": thread_id}}

    initial_state = create_initial_state(
        user_id=req.user_id,
        ticket_id=thread_id,
        user_message=req.message,
    )

    try:
        # invoke 会执行 Graph 直到完成或遇到 interrupt
        result = _graph.invoke(initial_state, config)

        # 检查是否被 interrupt
        # LangGraph 在 interrupt 时,结果中会包含 __interrupt__ 信息
        snapshot = _graph.get_state(config)

        if snapshot.next:
            # Graph 被中断,还有未执行的节点
            # 从 checkpoint 中提取 interrupt 信息
            pending_tasks = snapshot.tasks
            interrupt_info = None
            for task in pending_tasks:
                if task.interrupts:
                    interrupt_info = task.interrupts[0].value
                    break

            # 将待审批信息推入 Redis 队列
            if _redis and interrupt_info:
                await _redis.push_pending_approval(thread_id, interrupt_info)

            return {
                "thread_id": thread_id,
                "status": "interrupted",
                "message": "您的请求需要人工审核,已提交审批,请稍候。",
                "interrupt_detail": interrupt_info,
            }

        # 正常完成
        return {
            "thread_id": thread_id,
            "status": "completed",
            "message": result.get("final_response", "处理完成"),
        }

    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


# ── 2. 人工审批(Resume) ──────────────────────────────────
@router.post("/approve")
async def approve(req: ApprovalRequest):
    """
    人工审批接口。
    主管/客服审核后调用此接口,通过 Command(resume=...) 恢复 Workflow。
    """
    thread_id = req.thread_id
    config = {"configurable": {"thread_id": thread_id}}

    # 检查是否存在待处理的 interrupt
    snapshot = _graph.get_state(config)
    if not snapshot.next:
        raise HTTPException(
            status_code=400,
            detail=f"工单 {thread_id} 没有待处理的审批请求(可能已处理或不存在)。",
        )

    # 构建 resume 数据
    resume_data = {
        "decision": req.decision,       # "approved" | "rejected"
        "approver": req.approver,
        "comment": req.comment or "",
    }

    # 如果是投诉升级类型
    if req.agent_id:
        resume_data["agent_id"] = req.agent_id
        resume_data["resolution"] = req.resolution or ""

    # ★★★ 核心:通过 Command(resume=...) 恢复被中断的 Workflow ★★★
    result = _graph.invoke(Command(resume=resume_data), config)

    # 发布通知(Pub/Sub)
    if _redis:
        await _redis.publish_event(
            f"notify:{thread_id}",
            {"status": "completed", "message": result.get("final_response", "")},
        )

    return {
        "thread_id": thread_id,
        "status": "completed",
        "message": result.get("final_response", "审批完成"),
    }


# ── 3. 查询工单状态 ────────────────────────────────────────
@router.get("/status/{thread_id}")
async def get_status(thread_id: str):
    """查询指定工单的当前状态"""
    config = {"configurable": {"thread_id": thread_id}}
    snapshot = _graph.get_state(config)

    if not snapshot.values:
        raise HTTPException(status_code=404, detail="工单不存在")

    return {
        "thread_id": thread_id,
        "status": "interrupted" if snapshot.next else "completed",
        "intent": snapshot.values.get("intent"),
        "final_response": snapshot.values.get("final_response"),
        "pending_approval": bool(snapshot.next),
    }
app/api/schemas.py
"""Pydantic 请求/响应模型"""

from pydantic import BaseModel, Field
from typing import Literal


class ChatRequest(BaseModel):
    user_id: str = Field(..., description="用户ID")
    message: str = Field(..., description="用户消息")
    thread_id: str | None = Field(None, description="会话ID(续聊时传入)")


class ApprovalRequest(BaseModel):
    thread_id: str = Field(..., description="工单/会话ID")
    decision: Literal["approved", "rejected"] = Field(..., description="审批决定")
    approver: str = Field(..., description="审批人")
    comment: str | None = Field(None, description="审批备注")
    # 投诉升级时额外字段
    agent_id: str | None = Field(None, description="处理客服ID")
    resolution: str | None = Field(None, description="处理结果")

8.8 启动入口 — main.py

"""生产启动入口"""

from __future__ import annotations

import asyncio
from contextlib import asynccontextmanager

from fastapi import FastAPI

from app.config import settings
from app.graph import build_graph
from app.infra.checkpointer import get_checkpointer
from app.infra.redis_client import redis_manager
from app.api.routes import router as api_router, set_dependencies


@asynccontextmanager
async def lifespan(app: FastAPI):
    """应用生命周期管理"""
    # ── 启动 ──
    # 1. 初始化 Checkpointer
    checkpointer = get_checkpointer(
        mode=settings.CHECKPOINT_MODE,       # "memory" 或 "postgres"
        pg_dsn=settings.POSTGRES_DSN,
    )
    # 如果是 AsyncPostgresSaver,需要异步初始化
    if hasattr(checkpointer, "__aenter__"):
        checkpointer = await checkpointer.__aenter__()

    # 2. 构建 Graph
    graph = build_graph(checkpointer)

    # 3. 连接 Redis
    await redis_manager.connect()

    # 4. 注入依赖
    set_dependencies(graph, redis_manager)

    yield

    # ── 关闭 ──
    await redis_manager.close()
    if hasattr(checkpointer, "__aexit__"):
        await checkpointer.__aexit__(None, None, None)


app = FastAPI(title="电商智能客服多 Agent 系统", lifespan=lifespan)
app.include_router(api_router)


if __name__ == "__main__":
    import uvicorn
    uvicorn.run("main:app", host="0.0.0.0", port=8000, reload=True)
app/config.py
"""全局配置"""

from pydantic_settings import BaseSettings


class Settings(BaseSettings):
    CHECKPOINT_MODE: str = "memory"  # "memory" | "postgres"
    POSTGRES_DSN: str = "postgresql://postgres:postgres@localhost:5432/cs_agent"
    REDIS_URL: str = "redis://localhost:6379/0"
    NEO4J_URI: str = "bolt://localhost:7687"
    NEO4J_USER: str = "neo4j"
    NEO4J_PASSWORD: str = "password"

    class Config:
        env_file = ".env"


settings = Settings()

9. 核心 API 深入解释

9.1 interrupt(value) — 暂停 Workflow

from langgraph.types import interrupt

human_decision = interrupt({
    "type": "refund_approval",
    "refund_amount": 5999.0,
    "question": "是否批准退款?",
})
维度说明
为什么调用当前节点需要外部输入(人工审批),无法自动继续
调用后发生什么① 当前节点立即停止执行 ② interrupt() 抛出 GraphInterrupt 异常 ③ LangGraph 捕获该异常 ④ 将当前 State 写入 Checkpoint ⑤ 将 value 参数作为中断信息返回给调用方 ⑥ invoke() 返回,但 snapshot.next 不为空
Graph 是否继续执行❌ 不会。Graph 完全暂停,不占用任何计算资源
状态是否保存✅ 完整保存到 Checkpointer(包括所有 State 字段、执行位置、待执行节点)
如何恢复调用 graph.invoke(Command(resume=data), config)
resume 时传递什么Command(resume=...) 中的值会作为 interrupt() 的返回值
恢复后从哪里继续interrupt() 调用的下一行代码继续执行

9.2 Command(resume=value) — 恢复 Workflow

from langgraph.types import Command

result = graph.invoke(
    Command(resume={"decision": "approved", "approver": "张主管"}),
    config={"configurable": {"thread_id": "cs-user001-ticket001"}},
)
维度说明
做了什么① 从 Checkpointer 加载指定 thread_id 的最新 Checkpoint ② 恢复 State ③ 将 resume 值注入到之前 interrupt() 的位置 ④ 从中断点继续执行后续代码
前提条件必须存在有效的 thread_id 且该 thread 处于中断状态
重复调用如果 Graph 已经不在中断状态(snapshot.next 为空),会报错
数据流向Command(resume=X) 中的 Xinterrupt() 的返回值

9.3 Checkpointer — 状态持久化

from langgraph.checkpoint.memory import MemorySaver
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver

# 开发
checkpointer = MemorySaver()

# 生产
checkpointer = AsyncPostgresSaver.from_conn_string(
    "postgresql://user:pass@host:5432/db"
)
维度说明
保存时机每个节点执行完成后自动保存(不需要手动调用)
保存内容完整 State、当前执行位置、待执行节点列表、中间结果
存储位置MemorySaver → 内存;AsyncPostgresSaver → PostgreSQL 表
恢复方式通过 thread_id 查找最新 Checkpoint,恢复 State 继续执行
隔离机制每个 thread_id 独立的 Checkpoint 链,互不干扰

9.4 thread_id — 任务隔离

config = {"configurable": {"thread_id": "cs-user001-ticket001"}}
graph.invoke(state, config)
  • 每个 thread_id 对应一条独立的执行链
  • 不同用户的对话完全隔离
  • 同一用户的多轮对话使用相同 thread_id
  • Checkpoint 按 thread_id 索引,恢复时精确定位

10. 完整运行流程

场景:用户申请 ¥5,999 退款 → 触发人工审批 → 主管批准 → 完成

═══════════════════════════════════════════════════════════
Thread ID: cs-user001-a3f2b1
═══════════════════════════════════════════════════════════

[Step 1] 用户发送消息
─────────────────────────────────────────────────────────
  POST /chat
  {
    "user_id": "user_001",
    "message": "我买的手机屏幕碎了,订单 20260901-88721,我要退款"
  }

  → 创建初始 State
  → thread_id = "cs-user001-a3f2b1"

[Step 2] Router Agent 执行
─────────────────────────────────────────────────────────
  输入: "我买的手机屏幕碎了,订单 20260901-88721,我要退款"
  意图识别: "退款" → intent = "refund"
  提取订单号: order_id = "20260901-88721"
  置信度: 0.67

  → 条件路由 → refund_agent
  → ✅ Checkpoint #1 保存(router 执行完毕)

[Step 3] Refund Agent 执行
─────────────────────────────────────────────────────────
  查询订单: iPhone 17 Pro, ¥5,999
  退款金额: ¥5,999
  判断: ¥5,999 > ¥500(自动审批阈值)

  ★★★ 调用 interrupt() ★★★

  interrupt 信息:
  {
    "type": "refund_approval",
    "ticket_id": "cs-user001-a3f2b1",
    "order_id": "20260901-88721",
    "product": "iPhone 17 Pro",
    "refund_amount": 5999.0,
    "question": "退款金额 ¥5999.00 超过自动审批阈值 ¥500.00,需要主管审批。是否批准?"
  }

  → ✅ Checkpoint #2 保存(interrupt 状态)
  → Graph 状态: INTERRUPTED
  → 进程不再占用任何资源
  → 用户收到: "您的请求需要人工审核,已提交审批,请稍候。"

[Step 4] 等待人工审批(可能等 5 分钟,也可能等 2 天)
─────────────────────────────────────────────────────────
  审批信息已推入 Redis 队列: pending_approvals
  审批面板显示:
    工单: cs-user001-a3f2b1
    商品: iPhone 17 Pro
    退款: ¥5,999
    原因: 屏幕碎了

  ⏳ 此时:
    - 无进程阻塞
    - 无内存占用(状态在 PostgreSQL)
    - 服务器可以处理其他 1000 个请求
    - 即使服务器重启,状态也不丢失

[Step 5] 主管审批
─────────────────────────────────────────────────────────
  POST /approve
  {
    "thread_id": "cs-user001-a3f2b1",
    "decision": "approved",
    "approver": "张主管",
    "comment": "屏幕损坏属实,同意全额退款"
  }

  → 验证: snapshot.next 不为空 → 确认处于中断状态
  → 构建: Command(resume={"decision": "approved", "approver": "张主管", ...})

[Step 6] Resume 执行
─────────────────────────────────────────────────────────
  → 从 Checkpoint #2 恢复 State
  → interrupt() 返回 {"decision": "approved", "approver": "张主管", ...}
  → 继续执行 interrupt() 之后的代码
  → 生成最终回复

  → ✅ Checkpoint #3 保存(完成状态)
  → Graph 状态: COMPLETED

[Step 7] 最终响应
─────────────────────────────────────────────────────────
  {
    "thread_id": "cs-user001-a3f2b1",
    "status": "completed",
    "message": "✅ 退款已批准!\n订单:20260901-88721\n商品:iPhone 17 Pro\n退款金额:¥5999.00\n审批人:张主管\n备注:屏幕损坏属实,同意全额退款\n预计 3-5 个工作日到账。"
  }

  → Redis Pub/Sub 通知用户
  → 审批记录写入 PostgreSQL

═══════════════════════════════════════════════════════════
  总耗时: 用户等待时间 = 人工审批时间(而非进程阻塞时间)
  资源占用: 仅在节点实际执行时占用,等待期间为零
═══════════════════════════════════════════════════════════

11. Checkpoint 深入分析

11.1 Checkpoint 保存了什么?

AsyncPostgresSaver 为例,Checkpoint 包含:

┌─────────────────────────────────────────────────────┐
│  Checkpoint (PostgreSQL: checkpoints 表)             │
│                                                     │
│  1. checkpoint_id     — 唯一标识(UUID)              │
│  2. thread_id         — 所属会话/工单                 │
│  3. parent_id         — 上一个 Checkpoint(形成链)    │
│  4. checkpoint        — 序列化的完整 State            │
│     ├── messages      — 完整对话历史                  │
│     ├── intent        — 当前意图                     │
│     ├── order_info    — 订单数据                     │
│     ├── refund_amount — 退款金额                     │
│     ├── status        — 执行状态                     │
│     └── ...所有 State 字段                           │
│  5. metadata          — 元数据                       │
│     ├── step          — 当前执行到第几步              │
│     ├── source        — 触发来源(input/loop)        │
│     └── writes        — 本次写入的内容                │
│  6. pending_sends     — 待发送的消息                  │
│  7. versions          — 版本信息(用于并发控制)       │
│                                                     │
│  额外表: checkpoint_writes                           │
│  ├── 每个节点写入的具体字段和值                       │
│  └
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐