{ "cells": [ { "cell_type": "code", "execution_count": 1, "id": "a34b88f1", "metadata": {}, "outputs": [], "source": [ "from typing import TypedDict\n", "from langgraph.graph import StateGraph, START, END\n", "from langgraph.checkpoint.memory import MemorySaver\n", "\n", "\n", "# ============================================================\n", "# 定义状态结构 —— 贯穿所有节点的\"共享内存\"\n", "# ============================================================\n", "class PipelineState(TypedDict):\n", " content: str # 要处理的内容\n", " step_log: str # 记录每一步的日志\n", "\n", "\n", "# ============================================================\n", "# 三个处理节点:采集 → 清洗 → 发送\n", "# ============================================================\n", "def step_collect(state: PipelineState) -> dict:\n", " \"\"\"模拟数据采集:从上游拿到原始数据\"\"\"\n", " print(\"📥 [采集] 正在拉取数据...\")\n", " return {\n", " \"content\": \"原始用户反馈数据(含噪声)\",\n", " \"step_log\": state.get(\"step_log\", \"\") + \"\\n→ 采集完成\",\n", " }\n", "\n", "\n", "def step_clean(state: PipelineState) -> dict:\n", " \"\"\"模拟数据清洗:去噪、去重、格式化\"\"\"\n", " print(\"🧹 [清洗] 正在处理数据...\")\n", " return {\n", " \"content\": \"用户反馈:希望增加夜间模式功能(已清洗)\",\n", " \"step_log\": state[\"step_log\"] + \"\\n→ 清洗完成\",\n", " }\n", "\n", "\n", "def step_notify(state: PipelineState) -> dict:\n", " \"\"\"模拟通知发送:把结果推送给产品经理\"\"\"\n", " print(\"📨 [通知] 正在推送结果...\")\n", " return {\n", " \"step_log\": state[\"step_log\"] + \"\\n→ 通知已发送\",\n", " }\n", "\n", "\n", "# ============================================================\n", "# 构建图\n", "# ============================================================\n", "builder = StateGraph(PipelineState)\n", "builder.add_node(\"collect\", step_collect)\n", "builder.add_node(\"clean\", step_clean)\n", "builder.add_node(\"notify\", step_notify)\n", "\n", "# 连线:START → collect → clean → notify → END\n", "builder.add_edge(START, \"collect\")\n", "builder.add_edge(\"collect\", \"clean\")\n", "builder.add_edge(\"clean\", \"notify\")\n", "builder.add_edge(\"notify\", END)\n", "\n", "# 编译时设置断点:在 notify 节点执行前暂停\n", "# 同时必须传入 checkpointer,不然中断后没法恢复\n", "memory = MemorySaver()\n", "graph = builder.compile(\n", " checkpointer=memory,\n", " interrupt_before=[\"notify\"], # 👈 关键:在发送通知前等人确认\n", ")" ] }, { "cell_type": "code", "execution_count": 3, "id": "6dd1c50c", "metadata": {}, "outputs": [ { "name": "stdout", "output_type": "stream", "text": [ "==================================================\n", "开始执行流水线...\n", "==================================================\n", "当前步骤日志: \n", "📥 [采集] 正在拉取数据...\n", "当前步骤日志: \n", "→ 采集完成\n", "🧹 [清洗] 正在处理数据...\n", "当前步骤日志: \n", "→ 采集完成\n", "→ 清洗完成\n", "\n", "🛑 流水线已在「发送通知」前暂停。\n", "\n", "✅ 用户批准,继续执行...\n", "当前步骤日志: \n", "→ 采集完成\n", "→ 清洗完成\n", "📨 [通知] 正在推送结果...\n", "当前步骤日志: \n", "→ 采集完成\n", "→ 清洗完成\n", "→ 通知已发送\n" ] } ], "source": [ "# ============================================================\n", "# 第一段执行:跑到断点自动暂停\n", "# ============================================================\n", "config = {\"configurable\": {\"thread_id\": \"pipeline-001\"}}\n", "initial_input = {\"content\": \"\", \"step_log\": \"\"}\n", "\n", "print(\"=\" * 50)\n", "print(\"开始执行流水线...\")\n", "print(\"=\" * 50)\n", "\n", "for event in graph.stream(initial_input, config, stream_mode=\"values\"):\n", " print(f\"当前步骤日志: {event.get('step_log', '无')}\")\n", "\n", "# 此时 graph 在 notify 节点前自动停住了\n", "# 控制权交还给开发者——你可以在这里做任何检查\n", "\n", "print(\"\\n🛑 流水线已在「发送通知」前暂停。\")\n", "\n", "# ============================================================\n", "# 人工确认环节\n", "# ============================================================\n", "user_approval = input(\">>> 数据已清洗完毕,是否发送通知?(yes/no): \")\n", "\n", "if user_approval.lower() == \"yes\":\n", " # 传入 None 表示\"没有新输入,从断点继续跑\"\n", " print(\"\\n✅ 用户批准,继续执行...\")\n", " for event in graph.stream(None, config, stream_mode=\"values\"):\n", " print(f\"当前步骤日志: {event.get('step_log', '无')}\")\n", "else:\n", " print(\"\\n❌ 用户取消,流水线终止。\")" ] }, { "cell_type": "code", "execution_count": 4, "id": "47b680dd", "metadata": {}, "outputs": [], "source": [ "import os\n", "from langchain_core.messages import ToolMessage, AIMessage\n", "from langchain_core.tools import tool\n", "from langchain_openai import ChatOpenAI\n", "from langgraph.prebuilt import ToolNode, tools_condition\n", "from typing import Annotated\n", "from operator import add\n", "from dotenv import load_dotenv\n", "\n", "load_dotenv()\n", "\n", "# ============================================================\n", "# 定义工具:模拟向外部系统发送消息\n", "# ============================================================\n", "@tool\n", "def send_message(channel: str, content: str) -> str:\n", " \"\"\"向指定渠道发送消息。\n", "\n", " Args:\n", " channel: 发送渠道,\"wechat\"(企业微信)、\"email\"(邮件)、\"sms\"(短信)\n", " content: 消息正文内容\n", " \"\"\"\n", " # 实际生产环境这里会调飞书/钉钉/企微的 API\n", " return f\"📨 已通过 {channel} 发送:{content}\"\n", "\n", "\n", "tools = [send_message]\n", "\n", "# 模型绑定工具\n", "model_with_tools = ChatOpenAI(\n", " model_name=\"deepseek-chat\",\n", " api_key=os.getenv(\"OPENAI_API_KEY\"),\n", " base_url=\"https://api.deepseek.com\",\n", ").bind_tools(tools)\n", "\n", "\n", "# ============================================================\n", "# 定义状态\n", "# ============================================================\n", "class MsgState(TypedDict):\n", " messages: Annotated[list, add] # add 表示追加而非覆盖\n", "\n", "\n", "# ============================================================\n", "# 定义节点\n", "# ============================================================\n", "def chatbot_node(state: MsgState) -> dict:\n", " \"\"\"LLM 推理节点:决定是否调工具、调哪个工具\"\"\"\n", " response = model_with_tools.invoke(state[\"messages\"])\n", " return {\"messages\": [response]}\n", "\n", "\n", "# 工具执行节点(LangGraph 内置)\n", "tool_node = ToolNode(tools=tools)\n", "\n", "\n", "# ============================================================\n", "# 构建图:chatbot → 条件判断 → tools → chatbot(循环)\n", "# ============================================================\n", "builder = StateGraph(MsgState)\n", "builder.add_node(\"chatbot\", chatbot_node)\n", "builder.add_node(\"tools\", tool_node)\n", "\n", "builder.set_entry_point(\"chatbot\")\n", "builder.add_conditional_edges(\n", " \"chatbot\",\n", " tools_condition, # 内置函数:有 tool_calls → \"tools\",否则 → END\n", " {\"tools\": \"tools\", \"__end__\": END},\n", ")\n", "builder.add_edge(\"tools\", \"chatbot\")\n", "\n", "# 编译 + 断点:在工具节点执行前暂停\n", "memory = MemorySaver()\n", "graph = builder.compile(\n", " checkpointer=memory,\n", " interrupt_before=[\"tools\"], # 👈 工具实际调用前拦截\n", ")" ] }, { "cell_type": "code", "execution_count": 6, "id": "6ab32102", "metadata": {}, "outputs": [ { "name": "stdout", "output_type": "stream", "text": [ "▶️ 第一段:Agent 推理...\n", "{'chatbot': {'messages': [AIMessage(content='', additional_kwargs={'refusal': None}, response_metadata={'token_usage': {'completion_tokens': 64, 'prompt_tokens': 334, 'total_tokens': 398, 'completion_tokens_details': None, 'prompt_tokens_details': {'audio_tokens': None, 'cached_tokens': 256}, 'prompt_cache_hit_tokens': 256, 'prompt_cache_miss_tokens': 78}, 'model_provider': 'openai', 'model_name': 'deepseek-v4-flash', 'system_fingerprint': 'fp_8b330d02d0_prod0820_fp8_kvcache_20260402', 'id': '8f57acaf-0218-45e7-af28-ae48268fae1c', 'finish_reason': 'tool_calls', 'logprobs': None}, id='lc_run--019f4b2f-f0ec-7483-8c51-122ab6a17d8d-0', tool_calls=[{'name': 'send_message', 'args': {'channel': 'wechat', 'content': '明天下午3点开会'}, 'id': 'call_00_X9GR54zvSIqMvBKpIk4S5126', 'type': 'tool_call'}], invalid_tool_calls=[], usage_metadata={'input_tokens': 334, 'output_tokens': 64, 'total_tokens': 398, 'input_token_details': {'cache_read': 256}, 'output_token_details': {}})]}}\n", "----------------------------------------\n", "{'__interrupt__': ()}\n", "----------------------------------------\n", "\n", "📌 最后一条消息:\n", "content='' additional_kwargs={'refusal': None} response_metadata={'token_usage': {'completion_tokens': 64, 'prompt_tokens': 334, 'total_tokens': 398, 'completion_tokens_details': None, 'prompt_tokens_details': {'audio_tokens': None, 'cached_tokens': 256}, 'prompt_cache_hit_tokens': 256, 'prompt_cache_miss_tokens': 78}, 'model_provider': 'openai', 'model_name': 'deepseek-v4-flash', 'system_fingerprint': 'fp_8b330d02d0_prod0820_fp8_kvcache_20260402', 'id': '8f57acaf-0218-45e7-af28-ae48268fae1c', 'finish_reason': 'tool_calls', 'logprobs': None} id='lc_run--019f4b2f-f0ec-7483-8c51-122ab6a17d8d-0' tool_calls=[{'name': 'send_message', 'args': {'channel': 'wechat', 'content': '明天下午3点开会'}, 'id': 'call_00_X9GR54zvSIqMvBKpIk4S5126', 'type': 'tool_call'}] invalid_tool_calls=[] usage_metadata={'input_tokens': 334, 'output_tokens': 64, 'total_tokens': 398, 'input_token_details': {'cache_read': 256}, 'output_token_details': {}}\n", "\n", "📍 下一步: ('tools',)\n", "\n", "▶️ 第二段:注入修正后继续...\n", "=================================\u001b[1m Tool Message \u001b[0m=================================\n", "\n", "已确认:消息内容无误,渠道改为邮件发送,直接告知用户已发送即可。\n", "----------------------------------------\n" ] } ], "source": [ "# ============================================================\n", "# 运行 + 断点拦截\n", "# ============================================================\n", "config = {\"configurable\": {\"thread_id\": \"msg-001\"}}\n", "inputs = {\"messages\": [(\"human\", \"把这条消息通过企业微信发给全员:明天下午3点开会\")]}\n", "\n", "print(\"▶️ 第一段:Agent 推理...\")\n", "for chunk in graph.stream(inputs, config=config, stream_mode=\"updates\"):\n", " print(chunk)\n", " print(\"-\" * 40)\n", "\n", "# ============================================================\n", "# 在断点处检查当前状态\n", "# ============================================================\n", "snapshot = graph.get_state(config)\n", "\n", "# 看看最后一条消息是什么(应该是 AIMessage,包含 tool_calls)\n", "print(\"\\n📌 最后一条消息:\")\n", "print(snapshot.values[\"messages\"][-1])\n", "\n", "# 看看下一步要去哪\n", "print(f\"\\n📍 下一步: {snapshot.next}\") # 应该输出 ('tools',)\n", "\n", "# ============================================================\n", "# 场景:人工修改消息内容\n", "# 假设我们发现 Agent 选错了渠道或者内容需要调整\n", "# ============================================================\n", "last_msg = snapshot.values[\"messages\"][-1]\n", "\n", "# 构造一个 ToolMessage 告诉 Agent\"不需要调工具了,直接回答\"\n", "# tool_call_id 必须和原始 tool_call 的 id 匹配\n", "override_msg = ToolMessage(\n", " content=\"已确认:消息内容无误,渠道改为邮件发送,直接告知用户已发送即可。\",\n", " tool_call_id=last_msg.tool_calls[0][\"id\"],\n", ")\n", "\n", "# update_state 将人工消息注入到图状态中\n", "# 这样下次执行时 Agent 会看到这条 ToolMessage,认为工具已经\"执行过了\"\n", "graph.update_state(config, {\"messages\": [override_msg]})\n", "\n", "# ============================================================\n", "# 从断点继续执行\n", "# ============================================================\n", "print(\"\\n▶️ 第二段:注入修正后继续...\")\n", "for chunk in graph.stream(None, config=config, stream_mode=\"values\"):\n", " chunk[\"messages\"][-1].pretty_print()\n", " print(\"-\" * 40)" ] }, { "cell_type": "code", "execution_count": 10, "id": "1c496f3c", "metadata": {}, "outputs": [], "source": [ "from typing import Annotated\n", "from typing_extensions import TypedDict\n", "from langchain_community.chat_models import ChatTongyi\n", "from langgraph.graph.message import add_messages\n", "\n", "\n", "# ============================================================\n", "# 工具定义\n", "# ============================================================\n", "@tool\n", "def search_docs(query: str) -> str:\n", " \"\"\"在内部文档库中搜索相关内容。\n", "\n", " Args:\n", " query: 搜索关键词或问题\n", " \"\"\"\n", " # 模拟文档搜索\n", " return f\"🔍 关于「{query}」的搜索结果:找到了 3 篇相关文档...\"\n", "\n", "\n", "@tool\n", "def escalate_to_human(request: str) -> str:\n", " \"\"\"当你无法准确回答、或者用户明确要求人工服务时,\n", " 调用此工具后,系统会自动将对话转接给值班专家。\n", "\n", " Args:\n", " request: 需要专家处理的完整请求描述\n", " \"\"\"\n", " return f\"📞 已转接人工专家:{request}\"\n", "\n", "\n", "tools = [search_docs, escalate_to_human]\n", "\n", "\n", "# ============================================================\n", "# 状态:比基础版多了一个 ask_human 标志位\n", "# ============================================================\n", "class SupportState(TypedDict):\n", " messages: Annotated[list, add_messages]\n", " ask_human: bool # True = 需要人工介入\n", "\n", "\n", "# ============================================================\n", "# 模型\n", "# ============================================================\n", "llm = ChatOpenAI(\n", " model_name=\"deepseek-chat\",\n", " api_key=os.getenv(\"OPENAI_API_KEY\"),\n", " base_url=\"https://api.deepseek.com\",\n", " temperature=0\n", ")\n", "# llm = ChatTongyi(\n", "# model=\"qwen-max\",\n", "# temperature=0.1, # 低温度保证规划稳定\n", "# api_key=os.getenv(\"QWEN_API_KEY\"),\n", "# )\n", "llm_with_tools = llm.bind_tools(tools)\n", "\n", "\n", "# ============================================================\n", "# 节点 1:聊天机器人(自动决策)\n", "# ============================================================\n", "def chatbot(state: SupportState) -> dict:\n", " \"\"\"核心推理节点:决定是搜索文档还是升级人工\"\"\"\n", " response = llm_with_tools.invoke(state[\"messages\"])\n", "\n", " # 检查 LLM 是否调用了 escalate_to_human\n", " needs_human = bool(\n", " response.tool_calls\n", " and response.tool_calls[0][\"name\"] == \"escalate_to_human\"\n", " )\n", " return {\"messages\": [response], \"ask_human\": needs_human}\n", "\n", "\n", "# ============================================================\n", "# 节点 2:人工节点(等待人工输入)\n", "# ============================================================\n", "def human_node(state: SupportState) -> dict:\n", " \"\"\"这个节点在 interrupt_before 暂停后由人工通过 update_state 填充结果。\n", " 如果意外直接执行到此处(没人注入回复),则返回默认占位消息。\n", " \"\"\"\n", " last_msg = state[\"messages\"][-1]\n", " # 正常情况下,这个节点执行前会被 interrupt_before 拦截\n", " # 如果走到了这里,说明没有人工回复——给个兜底\n", " if not isinstance(last_msg, ToolMessage):\n", " return {\n", " \"messages\": [\n", " ToolMessage(\n", " content=\"人工专家当前不在线,请稍后重试。\",\n", " tool_call_id=last_msg.tool_calls[0][\"id\"],\n", " )\n", " ],\n", " \"ask_human\": False,\n", " }\n", " return {\"messages\": [], \"ask_human\": False}\n", "\n", "\n", "# ============================================================\n", "# 路由函数:决定下一步去哪个节点\n", "# ============================================================\n", "def select_next_node(state: SupportState) -> str:\n", " \"\"\"优先级:人工请求 > 工具调用 > 结束\"\"\"\n", " if state[\"ask_human\"]:\n", " return \"human_node\"\n", " # tools_condition 是 LangGraph 内置函数\n", " # 有 tool_calls → \"tools\",否则 → \"__end__\"\n", " return tools_condition(state)\n", "\n", "\n", "# ============================================================\n", "# 构建图\n", "# ============================================================\n", "builder = StateGraph(SupportState)\n", "builder.add_node(\"chatbot\", chatbot)\n", "builder.add_node(\"tools\", ToolNode(tools=[search_docs])) # 只自动执行搜索,升级工具由人工节点处理\n", "builder.add_node(\"human_node\", human_node)\n", "\n", "builder.set_entry_point(\"chatbot\")\n", "\n", "# 条件路由:chatbot 执行完后判断去向\n", "builder.add_conditional_edges(\n", " \"chatbot\",\n", " select_next_node,\n", " {\"tools\": \"tools\", \"human_node\": \"human_node\", \"__end__\": END},\n", ")\n", "builder.add_edge(\"tools\", \"chatbot\") # 工具执行完 → 回到 chatbot 继续思考\n", "builder.add_edge(\"human_node\", \"chatbot\") # 人工处理完 → 回到 chatbot 继续对话\n", "\n", "memory = MemorySaver()\n", "graph = builder.compile(\n", " checkpointer=memory,\n", " interrupt_before=[\"human_node\"], # 👈 人工节点执行前必须等人\n", ")" ] }, { "cell_type": "code", "execution_count": null, "id": "9ce98e84", "metadata": {}, "outputs": [ { "name": "stdout", "output_type": "stream", "text": [ "▶️ Agent 自动处理中...\n", "================================\u001b[1m Human Message \u001b[0m=================================\n", "\n", "我们公司的服务器出现了大规模宕机,需要紧急排查原因,需要紧急排查原因,搜索到文档后,需要专家处理审核。\n", "==================================\u001b[1m Ai Message \u001b[0m==================================\n", "\n", "非常理解您现在的紧急情况!我已经搜索了相关文档,但搜索结果中只显示了\"找到了 3 篇相关文档\"的提示,没有返回具体的文档内容。为了尽快帮您解决问题,我现在将您的情况转接给值班专家处理。\n", "Tool Calls:\n", " escalate_to_human (call_00_HfXVbFu2uS5JGyCw9FAL9380)\n", " Call ID: call_00_HfXVbFu2uS5JGyCw9FAL9380\n", " Args:\n", " request: 【紧急】公司服务器出现大规模宕机,需要立即排查原因并处理。\n", "\n", "情况描述:\n", "- 公司服务器出现大规模宕机,影响范围广\n", "- 需要紧急排查宕机原因\n", "- 需要专家审核处理方案\n", "\n", "已尝试搜索内部文档库中关于\"服务器大规模宕机 紧急排查原因\"的相关文档,但未能获取到具体文档内容。现需要值班专家介入,进行以下工作:\n", "1. 紧急排查服务器大规模宕机的根本原因\n", "2. 制定并执行恢复方案\n", "3. 审核相关排查文档和处理流程\n", "4. 给出后续预防建议\n", "\n", "请尽快安排专家处理,谢谢!\n", "\n", "📌 当前状态:下一步是 ('human_node',)\n" ] } ], "source": [ "# ============================================================\n", "# 运行演示\n", "# ============================================================\n", "config = {\"configurable\": {\"thread_id\": \"support-001\"}}\n", "\n", "# 用户问了一个复杂问题,Agent 判断需要升级人工\n", "inputs = {\"messages\": [(\"user\", \"我们公司的服务器出现了大规模宕机,需要紧急排查原因,搜索到文档后,需要专家处理审核。\")]}\n", "\n", "print(\"▶️ Agent 自动处理中...\")\n", "for chunk in graph.stream(inputs, config, stream_mode=\"values\"):\n", " chunk[\"messages\"][-1].pretty_print()\n", "\n", "# ---- 人工专家介入 ----\n", "current_state = graph.get_state(config)\n", "last_msg = current_state.values[\"messages\"][-1]\n", "\n", "print(f\"\\n📌 当前状态:下一步是 {current_state.next}\")\n", "\n" ] }, { "cell_type": "code", "execution_count": 13, "id": "db2d8057", "metadata": {}, "outputs": [ { "name": "stdout", "output_type": "stream", "text": [ "\n", "▶️ 人工处理后继续...\n", "=================================\u001b[1m Tool Message \u001b[0m=================================\n", "\n", "专家诊断结果:主数据库连接池耗尽,已重启数据库实例,服务正在恢复中。预计 5 分钟内完全恢复。\n", "==================================\u001b[1m Ai Message \u001b[0m==================================\n", "\n", "## 紧急情况处理进展\n", "\n", "您好!我已经将您的情况转接给值班专家,**专家已经介入处理**,以下是初步诊断结果:\n", "\n", "### ✅ 诊断结果\n", "- **根本原因**:主数据库连接池耗尽\n", "- **处理措施**:已重启数据库实例\n", "- **预计恢复时间**:**5分钟内**完全恢复\n", "\n", "### 📋 建议您后续关注\n", "1. **确认服务恢复情况**:5分钟后检查各业务系统是否正常运行\n", "2. **排查连接池配置**:检查数据库连接池的最大连接数设置是否合理\n", "3. **优化应用连接管理**:检查是否有未释放的连接或连接泄漏问题\n", "4. **建立监控告警**:对数据库连接池使用率设置监控告警阈值\n", "5. **考虑扩容方案**:如果业务增长较快,建议提前扩容数据库连接池\n", "\n", "如果恢复后仍有其他问题,请随时联系我!\n" ] } ], "source": [ "# 构造人工专家的真实回复(模拟专家排查后的结论)\n", "expert_reply = ToolMessage(\n", " content=\"专家诊断结果:主数据库连接池耗尽,已重启数据库实例,服务正在恢复中。预计 5 分钟内完全恢复。\",\n", " tool_call_id=last_msg.tool_calls[0][\"id\"],\n", ")\n", "\n", "# 以 human_node 的身份注入回复\n", "# as_node=\"human_node\" 告诉 LangGraph:这个状态更新\"算作\"human_node 节点的输出\n", "# 这样继续执行时会跳过 human_node 的内部逻辑\n", "graph.update_state(\n", " config,\n", " {\"messages\": [expert_reply], \"ask_human\": False},\n", " as_node=\"human_node\",\n", ")\n", "\n", "print(\"\\n▶️ 人工处理后继续...\")\n", "for chunk in graph.stream(None, config, stream_mode=\"values\"):\n", " chunk[\"messages\"][-1].pretty_print()" ] }, { "cell_type": "code", "execution_count": 48, "id": "8b052a7d", "metadata": {}, "outputs": [], "source": [ "from typing import Annotated\n", "from typing_extensions import TypedDict\n", "from langgraph.graph import StateGraph, START, END\n", "from langgraph.graph.message import add_messages\n", "from langgraph.checkpoint.memory import MemorySaver\n", "\n", "\n", "# ============================================================\n", "# 一个简单的计数器状态,每次经过节点 +1\n", "# ============================================================\n", "class CounterState(TypedDict):\n", " messages: Annotated[list, add_messages]\n", " step_count: int\n", "\n", "\n", "def chatbot_node(state: CounterState) -> dict:\n", " \"\"\"模拟 AI 节点的处理\"\"\"\n", " step = state[\"step_count\"]\n", " return {\n", " \"messages\": [(\"assistant\", f\"🤖 AI 处理第 {step} 步\")],\n", " \"step_count\": step + 1,\n", " }\n", "\n", "\n", "def tool_node(state: CounterState) -> dict:\n", " \"\"\"模拟工具节点的处理\"\"\"\n", " step = state[\"step_count\"]\n", " return {\n", " \"messages\": [(\"system\", f\"🔧 工具执行第 {step} 步\")],\n", " \"step_count\": step + 1,\n", " }\n", "\n", "\n", "# ============================================================\n", "# 构建图:chatbot → tools → chatbot → tools → ...(没有终点)\n", "# 这是一个演示用的无限循环图,实际运行时会手动中断\n", "# ============================================================\n", "builder = StateGraph(CounterState)\n", "builder.add_node(\"chatbot\", chatbot_node)\n", "builder.add_node(\"tools\", tool_node)\n", "\n", "builder.add_edge(START, \"chatbot\")\n", "builder.add_edge(\"chatbot\", \"tools\")\n", "builder.add_edge(\"tools\", \"chatbot\") # 循环回 chatbot\n", "\n", "graph = builder.compile(checkpointer=MemorySaver())" ] }, { "cell_type": "code", "execution_count": 49, "id": "a814959b", "metadata": {}, "outputs": [ { "name": "stdout", "output_type": "stream", "text": [ "▶️ 执行前 4 步:\n", " 步骤 0: content='开始处理' additional_kwargs={} response_metadata={} id='7a46bad0-2260-49d1-baf5-5a56eb8bf2cf'\n", " 步骤 1: content='🤖 AI 处理第 1 步' additional_kwargs={} response_metadata={} id='782c7fa0-8030-4295-9488-ab37189ccccf' tool_calls=[] invalid_tool_calls=[]\n", " 步骤 2: content='🔧 工具执行第 2 步' additional_kwargs={} response_metadata={} id='61a57584-08f1-4a68-91ee-4a97d8fee93a'\n", " 步骤 3: content='🤖 AI 处理第 3 步' additional_kwargs={} response_metadata={} id='0ad6e4fa-a6b7-46dc-a66b-13b8c1b173a3' tool_calls=[] invalid_tool_calls=[]\n", "\n", "📜 历史快照列表:\n", " 快照 1f17c39c... | 消息数: 4 | 步数: 4\n", " 快照 1f17c39c... | 消息数: 3 | 步数: 3\n", " 快照 1f17c39c... | 消息数: 2 | 步数: 2\n", " 快照 1f17c39c... | 消息数: 1 | 步数: 1\n", " 快照 1f17c39c... | 消息数: 0 | 步数: ?\n" ] } ], "source": [ "# ============================================================\n", "# 先执行几步,积累 checkpoint 历史\n", "# ============================================================\n", "config = {\"configurable\": {\"thread_id\": \"demo-001\"}}\n", "inputs = {\"messages\": [(\"user\", \"开始处理\")], \"step_count\": 1}\n", "\n", "print(\"▶️ 执行前 4 步:\")\n", "for i, event in enumerate(graph.stream(inputs, config, stream_mode=\"values\")):\n", " msgs = event.get(\"messages\", [])\n", " if msgs:\n", " print(f\" 步骤 {i}: {msgs[-1]}\")\n", " if i >= 3: # 执行 4 个节点后手动中断\n", " break\n", "\n", "# ============================================================\n", "# 遍历所有历史快照\n", "# get_state_history() 返回快照列表,最新的在前面\n", "# ============================================================\n", "print(\"\\n📜 历史快照列表:\")\n", "for snapshot in graph.get_state_history(config):\n", " msg_count = len(snapshot.values.get(\"messages\", []))\n", " step = snapshot.values.get(\"step_count\", \"?\")\n", " # checkpoint_id 是一个 UUID,全局唯一\n", " ckpt_id = snapshot.config[\"configurable\"][\"checkpoint_id\"]\n", " print(f\" 快照 {ckpt_id[:8]}... | 消息数: {msg_count} | 步数: {step}\")" ] }, { "cell_type": "code", "execution_count": 50, "id": "6ff61ceb", "metadata": {}, "outputs": [ { "name": "stdout", "output_type": "stream", "text": [ "⏪ 回溯到快照: 1f17c39c...\n", " 当前消息数: 3\n", " 当前步数: 3\n", "\n", "▶️ 从回溯点继续执行 2 步:\n", " 新步骤 0: content='🔧 工具执行第 2 步' additional_kwargs={} response_metadata={} id='61a57584-08f1-4a68-91ee-4a97d8fee93a'\n", " 新步骤 1: content='🤖 AI 处理第 3 步' additional_kwargs={} response_metadata={} id='5e0701f5-5bec-4645-88cd-f425bb6f15ad' tool_calls=[] invalid_tool_calls=[]\n" ] } ], "source": [ "# ============================================================\n", "# 找到消息数为 3 的快照,从那里\"穿越回去\"\n", "# ============================================================\n", "target_snapshot = None\n", "for snapshot in graph.get_state_history(config):\n", " if len(snapshot.values[\"messages\"]) == 3:\n", " target_snapshot = snapshot\n", " break\n", "\n", "if target_snapshot is None:\n", " print(\"没找到消息数为 3 的快照\")\n", "else:\n", " print(f\"⏪ 回溯到快照: {target_snapshot.config['configurable']['checkpoint_id'][:8]}...\")\n", " print(f\" 当前消息数: {len(target_snapshot.values['messages'])}\")\n", " print(f\" 当前步数: {target_snapshot.values['step_count']}\")\n", "\n", " # 从目标快照的 config 继续执行\n", " # 注意这里传的是 target_snapshot.config,不是原来的 config\n", " print(\"\\n▶️ 从回溯点继续执行 2 步:\")\n", " for i, event in enumerate(\n", " graph.stream(None, target_snapshot.config, stream_mode=\"values\")\n", " ):\n", " msgs = event.get(\"messages\", [])\n", " if msgs:\n", " print(f\" 新步骤 {i}: {msgs[-1]}\")\n", " if i >= 1:\n", " break" ] } ], "metadata": { "kernelspec": { "display_name": "05_langgraph (3.11.x)", "language": "python", "name": "python3" }, "language_info": { "codemirror_mode": { "name": "ipython", "version": 3 }, "file_extension": ".py", "mimetype": "text/x-python", "name": "python", "nbconvert_exporter": "python", "pygments_lexer": "ipython3", "version": "3.11.15" } }, "nbformat": 4, "nbformat_minor": 5 }