004.DeepAgents 框架介绍与应用实战

一:DeepAgents 框架定位与核心价值

1.1. 什么是 DeepAgents

DeepAgents 是一个基于 LangChain 和 LangGraph 构建的企业级高级智能体框架。它建立在 LangGraph(底层运行时)和 LangChain(工具 / 模型层)之上,是一个高阶的 Agent Harness(智能体装备 / 套件)**。

DeepAgents 开源地址:https://github.com/langchain-ai/deepagents

Pasted image 20260708141904.png

1.2. 解决的核心问题

传统的 Agent 开发通常运行一个简单的循环:思考 → 调用工具 → 观察 → 重复。这种模式在处理多小时或多天的任务时,容易遇到以下"浅层陷阱"(Shallow Agent Problem)

  1. 规划能力缺失:原生 Agent 倾向于“走一步看一步”,缺乏全局视角的任务拆解,容易在多步任务中迷失方向。
  2. 遗忘与混乱:在执行超过 10-20 步的长任务时,由于 Context Window(上下文窗口)限制,传统 Agent 容易忘记初始目标或陷入死循环。
  3. 环境交互困难:文件系统操作、代码执行环境(沙箱)的配置和安全管理复杂。
  4. 上下文污染:所有工具返回结果都堆积在一个对话历史中,导致噪声过大。
  5. 协作编排复杂:多智能体(Multi-Agent)之间的任务分发和上下文隔离难以实现。

DeepAgents 通过引入"类人"的工作流解决了这些问题:先做计划(Plan),再执行,利用文件系统管理记忆,遇到复杂子任务时"外包"给子 Agent。将规划工具、文件系统访问、子代理和详细提示词等关键机制整合在一起,以支持复杂的深度任务 。

Pasted image 20260708142502.png

1.3. 应用场景

DeepAgents 不适合用来做简单的聊天机器人(Chatbot),它是为重任务设计的,它适用于任务需规划、上下文海量、需多专家协作、要求持久记忆的场景,将 LangChain 生态从单步响应提升至自主完成复杂项目的高度:

DeepAgents 适用场景

场景类型 能力说明 工作逻辑 / 技术特点 代表性案例
深度调研与报告生成 支持长周期、多步骤、多来源信息整合的研究任务 • 自动生成研究计划(Todo)
• 调用搜索工具获取资料
• 将关键信息写入文件系统(长期记忆)
• 使用子代理(Sub-Agents)深入研究子课题
• 主代理统一规划、整合结果
• LangChain Deep Research 示例(Tavily 搜索 + 多子代理拆分研究)
• OpenAI Deep Research(官方生产级深度调研助理)
自动编程与代码助理 理解代码、修改代码、生成新文件、执行工具链 • 代理可读写虚拟文件系统
• 自动分析源码并输出 diff
• 人工审批(Human-in-the-Loop)保证安全写入
• 调用 Shell / 测试工具执行流程
• 可将项目规范写入 / memories 用作长期记忆
• LangChain DeepAgents CLI(终端自动编码)
• Anthropic Claude Code(深度自动重构与编程)
• Manus(多步骤代码智能体)
复杂流程自动化(业务流程 Orchestration) 将多个步骤串联为可控流程,适合企业级自动化任务 • 任务分解 → 多步骤规划 → 调用不同工具
• 搜索、筛选、处理、生成等多环节协作
• 使用文件系统存储中间数据(如列表、分析结果等)
• 支持多工具、多子任务并行处理
• DeepAgents 求职助手(职位搜索 → 筛选排序 → 求职信生成 → 打包结果)
• 企业场景如:自动生成分析报告 / 客服知识库构建 / 数据采集 + 处理流

二:与 LangChain 及 LangGraph 的区别

特性 LangChain LangGraph DeepAgents
层级 基础组件库 (Foundation) 编排引擎 (Orchestration) 应用框架 (Application Framework)
核心抽象 Chain, Runnable, Tool StateGraph, Node, Edge DeepAgent, Middleware, Backend
灵活性 极高 (积木块) 高 (自定义图结构) 中 (Opinionated / 约定优于配置)
开箱即用 低 (需自行组装) 中 (需定义图逻辑) 高 (内置规划、文件系统、子代理)
适用对象 库开发者 / 底层构建 复杂流程序列化开发者 应用开发者 / 企业级解决方案

Pasted image 20260708143207.png


三:DeepAgents 核心功能介绍

3.1. 核心入口:create_deep_agent()

这是整个框架的核心函数,它创建了一个功能完整的深度智能体。

默认配置

关键参数

from deepagents import create_deep_agent
from langchain_tavily import TavilySearch
from langgraph.checkpoint.memory import InMemorySaver

# 1. 初始化 Tavily 搜索工具
tavily = TavilySearch(max_results=3)

# 2. 编写系统提示词
research_instructions = """
您是一位资深的研究人员。您的工作是进行深入的研究,然后撰写一份精美的报告。
您可以通过互联网搜索引擎作为主要的信息收集工具。
## 可用工具
### `互联网搜索`
使用此功能针对给定的查询进行互联网搜索。您可以指定要返回的最大结果数量、主题以及是否包含原始内容。
### `写入本地文件`
使用此功能将研究报告保存到本地文件。当您完成研究并生成报告后,请使用此工具将完整的报告内容保存到文
件中。
- 文件路径建议使用 .md 格式(Markdown),例如 "research_report.md" 或 "./reports/报告名
称.md"
- 请确保报告内容完整、结构清晰,包含所有章节和引用来源
## 工作流程
在进行研究时:
1. 首先将研究任务分解为清晰的步骤
2. 使用互联网搜索来收集全面的信息
3. 将信息整合成一份结构清晰的报告
4. **重要**:完成报告后,务必使用 `写入本地文件` 工具将完整报告保存到本地文件
5. 务必引用你的资料来源
**注意**:请确保在完成研究后,将完整的报告内容保存到文件中,这样用户可以方便地查看和保存报告。
"""

# 2. 创建 DeepAgents 智能体
agent = create_deep_agent(
    name="DeepAgents_Agent",       # 智能体名称
    tools=[tavily],                # 可调用工具 Tool
    model=model,                   # 模型 Model
    system_prompt=research_instructions,  # 系统提示词
    checkpointer=InMemorySaver(),  # 检查点 Checkpointer,内存检查点
)

# 3. 配置线程 ID
config = {"configurable": {"thread_id": "1"}}

result = agent.invoke({"messages": [{"role": "user", "content": "帮我查询一下有关deepagents框架的最新动态"}]}, config=config)

3.2. create_deep_agent 内部结构

源码参数截图

Pasted image 20260708144151.png

除了自己定义的工具(如 Tavily 搜索),DeepAgents 还默认添加了一些其他工具:

Pasted image 20260708145115.png

这些都是 DeepAgents 特有的功能,用于支持智能体在实际应用中的各种场景。那么,接下来我们来看看这些功能的具体应用。

Pasted image 20260708145031.png


四:四大核心内置工具与组件详解

DeepAgents 通过中间件 (Middleware) 的形式,为智能体注入了四项核心能力,构成了框架的四大支柱(Four Pillars)

Pasted image 20260708145209.png

其核心价值在于将原本需要手动编排的规划 - 存储 - 委托 - 执行流程,固化为中心化、可复用、可观测的中间件体系,标志着 AI Agent 从"脚本化"向"产品化"的关键演进。

维度 系统提示词 (System Prompt) 规划工具 (Planning Tool) 文件系统 (File System) 子代理 (Sub Agents)
角色定位 行为总导演:定义 Agent 的"世界观"与工具使用范式,确保三大中间件协同不偏离目标 任务架构师:将模糊需求转化为可执行、可追踪、可动态调整的结构化任务蓝图 上下文仓库:虚拟化存储引擎,解决长任务中的信息溢出与状态持久化难题 执行特派员:实现上下文隔离与专业分工,防止主 Agent 因深层递归导致状态混乱
核心功能 内置 Claude Code 风格指令,涵盖规划逻辑、文件操作规范、子代理调用协议;支持场景化自定义覆盖 write_todos: 生成带优先级 / 依赖关系的 JSON 任务列表
read_todos: 实时查询任务执行进度与状态
ls / glob: 文件浏览与模式匹配
read / write / edit: CRUD 操作
grep: 内容检索
execute: 沙箱命令执行
task: 动态生成同构或异构子 Agent
支持独立上下文窗口与工具集配置
结果通过文件系统回传
技术实现 字符串模板,在 create_deep_agent 时注入;默认提示词约 2000 tokens,包含 ReAct 循环与三大中间件调用示例 TodoListMiddleware:拦截 LLM 输出中的 todo_list 字段,解析为 agent_state.todos 字典,状态变更触发图节点重计算 FilesystemMiddleware:基于 LangGraph State 的 files 字段实现内存级虚拟文件系统,大工具结果(> 2KB)自动触发 write_file 落盘 SubAgentMiddleware:将 task 调用编译为独立的 StateGraph 子图,通过命名空间隔离状态,父图通过 files 读取子图输出
状态管理 静态配置,单次会话内不可变;可通过configurable_agent实现热更新 动态状态机:每个 todo 含 id / description / status / priority / dependencies 字段,执行后状态从 pendingcompleted,支持 update_todos 动态调整 持久化存储:默认存储在 LangGraph State,支持切换 StateBackend(内存 / Redis / Postgres)实现跨会话文件共享 完全隔离:子 Agent 拥有独立的 messagesfiles 命名空间,异常不会污染父 Agent 状态;支持 max_iterations 限制防止无限递归
使用场景 ① 垂直领域定制:金融研究 / 医疗诊断等需强化专业约束的场景
② 多 Agent 协作:统一多个子 Agent 的行为规范
③ 安全合规:注入数据脱敏、权限检查等硬性规则
① 长周期研究:自动拆解为文献检索 → 数据收集 → 分析 → 撰写的阶段性任务
② 故障恢复:崩溃后通过 read_todos 快速定位断点续跑
③ 动态重规划:执行中发现信息不足时新增补充任务
① 大结果处理:搜索返回 100KB 内容自动落盘,避免上下文溢出
② 知识沉淀:中间分析结果写入文件供后续步骤复用
③ 多 Agent 数据共享:父 Agent 与子 Agent 通过文件交换数据,无需序列化传递
① 高风险操作隔离:网页抓取 / 代码执行等易失败任务委托给子 Agent
② 专业化分工:主 Agent 负责任务编排,子 Agent 专注领域执行(如专门的数据分析 Agent)
③ 资源优化:子 Agent 可使用轻量化模型,降低整体 token 成本

4.1. 系统提示词 (System Prompts)

Pasted image 20260708145933.png

Pasted image 20260708145944.png

Pasted image 20260708145958.png

Pasted image 20260708150006.png

4.2. 规划工具 (Planning System / Todo List)

Pasted image 20260708153151.png

4.3. 子代理 (Sub-Agent Delegation)

1. 显示传入 subAgent 参数

pip install langchain-mcp-adapters mcp
import asyncio
import os
from dotenv import load_dotenv
from rich.console import Console
from rich.panel import Panel
from rich.tree import Tree
from deepagents import create_deep_agent
from langchain_deepseek import ChatDeepSeek
from langchain_mcp_adapters.client import MultiServerMCPClient
from langchain_community.tools import TavilySearchResults

from langchain_core.messages import ToolMessage, BaseMessage

load_dotenv(override=True)

async def setup_mcp_tools():
    try:
        client = MultiServerMCPClient({
            "context7": {
                "transport": "stdio",
                "command": "npx",
                "args": ["-y", "@upstash/context7-mcp@latest"],
            }
        })
        # 获取工具
        tools = await client.get_tools()
        return client, tools
    except Exception as e:
        return None, []

# 定义子 Agent 配置
def get_subagents_config(mcp_tools):
    # 子 Agent 1: 官方文档专家
    doc_tools = mcp_tools if mcp_tools else [TavilySearchResults(max_results=3)]

    docs_researcher = {
        "name": "DocsResearcher",
        "description": "负责查阅官方文档和技术规范的专家 Agent。",
        "system_prompt": "你是一名专门查阅官方文档的技术专家。请使用工具获取准确的技术细节。不要猜测。",
        "tools": doc_tools,
        "model": "deepseek-chat"
    }

    # 子 Agent 2: 社区生态专家
    community_researcher = {
        "name": "CommunityResearcher",
        "description": "负责搜索社区博客、教程和最佳实践的专家 Agent。",
        "system_prompt": "你是一名关注社区动态的开发者。请搜索博客、论坛和 GitHub 讨论。",
        "tools": [TavilySearchResults(max_results=3)],
        "model": "deepseek-chat"
    }

    return [docs_researcher, community_researcher]

# 主运行逻辑
async def run_parallel_demo():
    # 初始化 MCP
    mcp_client, mcp_tools = await setup_mcp_tools()

    # 获取子 Agent 配置
    subagents = get_subagents_config(mcp_tools)

    # 创建主 Agent
    llm = ChatDeepSeek(model="deepseek-chat", temperature=0)

    agent = create_deep_agent(
        model=llm,
        tools=[],
        subagents=subagents,        # 传入子 Agent 配置
        system_prompt="""你是一名技术总监。你的任务是协调 DocsResearcher 和 CommunityResearcher 完成调研任务。
                        请根据用户需求,将任务拆解并分发给这两个子 Agent。
                        如果任务允许,请务必并行调用它们以提高效率。
                        最后汇总它们的报告。"""
    )

    task = "请详细调研 'LangChain DeepAgents' 框架。我需要官方的技术架构说明(来自文档)以及社区的最佳实践案例。请对比两者。"

五:文件系统集成 (Filesystem & Sandbox)

5.1. 核心工具集

# 使用本地文件系统后端,根目录设为 ./workspace
# 这样我们可以看到真实文件的创建
# 注意:设置 virtual_mode=True 以支持绝对路径 (如 /hello_world.py) 映射到 ./workspace
backend = FilesystemBackend(root_dir="./workspace", virtual_mode=True)

# 创建 Agent
# 注意:create_deep_agent 默认会自动包含 FilesystemMiddleware
agent = create_deep_agent(
    model=model,
    backend=backend,
    system_prompt="你是一个文件系统操作助手。请根据用户指令使用相应的工具。"
)

5.2. Backend 后端应用

1. 默认模式(内存沙箱)

内存沙箱,支持在内存中执行代码,防止代码注入攻击。代码执行结果会被存储在内存中,不会对本地环境造成影响。执行结束后,内存中的数据会被清除。

from deepagents import create_deep_agent
from dotenv import load_dotenv
from langchain_openai import ChatOpenAI

load_dotenv(override=True)

model = ChatOpenAI(model="gpt-4o", temperature=0)

# 默认情况下,DeepAgent 会自动加载 FilesystemMiddleware 并使用内存后端
agent = create_deep_agent(
    model=model,
    system_prompt="你是一个数据处理助手。"
)

# Agent 可以自由创建文件、读取文件,但这些文件只存在于内存中
# 任务结束后,这些文件会自动消失
agent.invoke({"messages": [("user", "创建一个名为 test.txt 的文件并写入 'Hello, World!'")]})

2. 持久化模式(操作真实文件)

持久化沙箱,支持将代码执行结果持久化到本地磁盘,防止代码执行过程中数据丢失。代码执行结束后,内存中的数据会被清除,但是磁盘上的文件会被保留。

from deepagents import create_deep_agent
from deepagents.backends import FilesystemBackend
from langchain_openai import ChatOpenAI

model = ChatOpenAI(model="gpt-4o", temperature=0)

# 将 Agent 的根目录映射到本地的 "./workspace" 文件夹,virtual_mode=True 表示启用虚拟文件系统,作用是将 Agent 执行的所有文件操作都映射到本地的 "./workspace" 文件夹,而不是直接操作本地文件系统。
backend = FilesystemBackend(root_dir="./workspace", virtual_mode=True)

agent = create_deep_agent(
    model=model,
    backend=backend  # 传入后端,中间件会自动使用它
)

# 此时 Agent 执行 write_file("/readme.md", ...) 会在本地 ./workspace/readme.md 创建文件
agent.invoke({"messages": [("user", "创建一个名为 readme.md 的文件并写入 'Hello, World!'")]})

3. 演示大文件读取分页

在 FilesystemBackend 分页的参数是固定死的,默认每次读取 500 行,我们可以在系统提示词中指定每次读取的行数。

import asyncio
import os
import shutil
from pathlib import Path
from dotenv import load_dotenv
from langchain_openai import ChatOpenAI
from langchain_core.messages import BaseMessage, ToolMessage
from deepagents import create_deep_agent
from deepagents.backends import FilesystemBackend

load_dotenv(override=True)

# 定义工作目录
WORK_DIR = Path("workspace/pagination_demo").resolve()
LARGE_FILE_NAME = "server_logs.txt"
LARGE_FILE_PATH = WORK_DIR / LARGE_FILE_NAME
TARGET_SECRET = "CRITICAL_ERROR_CODE_998877"

async def run_pagination_demo():
    # 初始化 FilesystemBackend
    # 将 backend 指向我们的测试目录
    # virtual_mode=True 确保 Agent 只能访问该目录下的文件,不能访问宿主机其他目录
    backend = FilesystemBackend(root_dir=WORK_DIR, virtual_mode=True)

    # 创建 Agent
    # 我们明确指示 Agent 使用分页读取,每次读取 300 行,Agent 默认每次读取 500 行(DEFAULT_READ_LIMIT)
    system_prompt = """
    你是一个专业的系统管理员。
    你的任务是从日志文件中查找特定的错误代码。
    注意:日志文件可能非常大,为了避免上下文溢出,你必须使用 `read_file` 工具的分页功能。
    每次读取请限制在 300 行以内 (limit=300),并使用 offset 参数向后滚动。
    直到找到目标信息为止。
    """

    llm = ChatOpenAI(model="gpt-4o", temperature=0)

    agent = create_deep_agent(
        model=llm,
        backend=backend,
        system_prompt=system_prompt
    )

4. E2BBackend 使用 E2B 云沙箱

E2B(Environment To Be) 是一个专为 AI 智能体设计的云端安全执行环境 。当 AI 需要写代码、运行脚本或操作文件时,它不会在你的本地机器上操作,而是连接到 E2B 的云端环境中进行。

核心特性

  1. 安全性与隔离性 (Security & Isolation) :
    • AI 生成的代码(可能包含错误或恶意逻辑)完全运行在云端沙箱中, 绝不会破坏你本地的电脑环境 。
    • 演示代码中,Agent 即使执行了 rm -rf / ,也只是删除了云端临时的沙箱,对宿主机毫发无损。
  2. 持久化会话 (Long-running Sessions) :
    • 沙箱可以保持运行状态。Agent 可以先创建一个文件(如演示中的 /home/user/hello.py ),然后在后续步骤中运行它。环境状态在会话期间是保持的。
  3. 标准 Linux 环境 :
    • 它提供标准的 Linux Shell。演示中 Agent 执行了 uname -a 和 python --version ,就像在真实的服务器上一样。
# 定义 E2B 沙箱后端类
# 该类继承自 BaseSandbox,实现了在 E2B 沙箱中执行命令的功能

import base64
from typing import Any, Optional

# 导入 DeepAgents 后端协议,定义了沙箱后端的接口
from deepagents.backends.protocol import (
    ExecuteResponse,
    FileDownloadResponse,
    FileUploadResponse,
    SandboxBackendProtocol,
)
# 导入基础沙箱类,用于实现沙箱后端的基本功能
from deepagents.backends.sandbox import BaseSandbox

try:
    from e2b import Sandbox
except ImportError:
    Sandbox = None

class E2BBackend(BaseSandbox):
    """E2B 沙箱后端实现,用于 DeepAgents。
    该后端使用 E2B(https://e2b.dev)提供安全、隔离的执行环境。
    """
    def __init__(
        self,
        template: str = "base",
        api_key: Optional[str] = None,
        timeout: Optional[int] = None,
        metadata: Optional[dict[str, str]] = None,
    ) -> None:
        """初始化 E2B 沙箱。

        参数:
            template: E2B 沙箱模板 ID(默认:"base")
            api_key: E2B API 密钥(可选,默认使用 E2B_API_KEY 环境变量)
            timeout: 沙箱超时时间(秒)
            metadata: 自定义沙箱元数据
        """
        if Sandbox is None:
            raise ImportError(
                "e2b package is not installed. "
                "Please install it with `pip install e2b`."
            )

        self.sandbox = Sandbox.create(
            template=template,
            api_key=api_key,
            timeout=timeout,
            metadata=metadata,
        )

    @property
    def id(self) -> str:
        """Unique identifier for the sandbox backend."""
        return self.sandbox.sandbox_id

    def execute(self, command: str) -> ExecuteResponse:
        """Execute a command in the sandbox."""
        try:
            # E2B commands.run returns CommandResult with stdout, stderr, exit_code
            result = self.sandbox.commands.run(command)

            # 返回执行结果,包含 stdout 是标准输出,stderr 是标准错误输出,exit_code 是退出码
            return ExecuteResponse(
                output=result.stdout + result.stderr,
                exit_code=result.exit_code,
                truncated=False,
            )
        except Exception as e:
            return ExecuteResponse(
                output=f"Error executing command: {str(e)}",
                exit_code=1,
                truncated=False,
            )

    def upload_files(self, files: list[tuple[str, bytes]]) -> list[FileUploadResponse]:
        """Upload multiple files to the sandbox."""
        responses = []
        for path, content in files:
            try:
                # Ensure directory exists before writing
                # We can use execute to mkdir -p
                parent_dir = path.rsplit("/", 1)[0]
                if parent_dir:
                    self.sandbox.commands.run(f"mkdir -p {parent_dir}")

                # Write file
                self.sandbox.files.write(path, content)
                responses.append(FileUploadResponse(path=path, error=None))
            except Exception as e:
                error_msg = str(e).lower()
                error = "invalid_path"
                if "permission" in error_msg:
                    error = "permission_denied"

                responses.append(FileUploadResponse(path=path, error=error))
        return responses

    def download_files(self, paths: list[str]) -> list[FileDownloadResponse]:
        """Download multiple files from the sandbox."""
        responses = []
        for path in paths:
            try:
                content = self.sandbox.files.read(path)
                # Ensure content is bytes
                if isinstance(content, str):
                    content = content.encode("utf-8")

                responses.append(FileDownloadResponse(path=path, content=content, error=None))
            except Exception as e:
                error_msg = str(e).lower()
                error = "invalid_path"
                if "not found" in error_msg:
                    error = "file_not_found"
                elif "directory" in error_msg:
                    error = "is_directory"
                elif "permission" in error_msg:
                    error = "permission_denied"

                responses.append(FileDownloadResponse(path=path, content=None, error=error))
        return responses

    def close(self):
        """Close the sandbox session."""
        self.sandbox.kill()

将自定义好的沙箱后端类注册到 DeepAgents 中

import asyncio
import os
import sys
from dotenv import load_dotenv
from deepagents import create_deep_agent
from langchain_openai import ChatOpenAI
from langchain_mcp_adapters.client import MultiServerMCPClient
from langchain_core.messages import BaseMessage, ToolMessage

# 假设 e2b_backend.py 在同一目录下
try:
    from e2b_backend import E2BBackend
except ImportError:
    # 尝试从当前路径导入
    sys.path.append(os.path.dirname(os.path.abspath(__file__)))
    from e2b_backend import E2BBackend

load_dotenv(override=True)

async def setup_mcp_tools():
    """连接 Context7 MCP 服务器并获取工具"""
    print("正在连接 Context7 MCP 服务器...")
    try:
        client = MultiServerMCPClient({
            "context7": {
                "transport": "stdio",
                "command": "npx",
                "args": ["-y", "@upstash/context7-mcp@latest"],
            }
        })
        # 获取工具
        tools = await client.get_tools()
        return client, tools
    except Exception as e:
        return None, []

async def run_e2b_demo():
    # 1. 检查 API Key
    if not os.getenv("E2B_API_KEY"):
        print("❌ 错误: 未找到 E2B_API_KEY 环境变量。")
        return

    # 2. 初始化 MCP (用于查询文档等辅助任务)
    mcp_client, mcp_tools = await setup_mcp_tools()

    # 3. 初始化 E2B Backend
    try:
        backend = E2BBackend(template="base")
    except Exception as e:
        return

    try:
        # 4. 创建 DeepAgent
        llm = ChatOpenAI(model="gpt-4o", temperature=0)

        agent = create_deep_agent(
            model=llm,
            tools=mcp_tools,     # 赋予 MCP 工具能力
            backend=backend,     # 赋予 E2B 沙箱能力
            system_prompt="""你是一个拥有云端沙箱环境的高级技术助手。
            你的任务是演示如何在沙箱中进行操作。
            请执行以下步骤:
            1. 使用 'execute_command' 运行 'uname -a' 和 'python --version' 来展示环境信息。
            2. 创建一个 Python 脚本 '/home/user/hello.py',内容是打印 'Hello from E2B Sandbox!'。
            3. 运行这个 Python 脚本并显示输出。
            """
        )

5. DockerBackend 使用 Docker 容器

作用与优势

# 这个 Backend 会在本地启动一个 Docker 容器,并将会话隔离在容器内部。

# - 核心功能 :
#   - 自动生命周期管理:初始化时启动容器,结束时自动销毁 (auto_remove=True)。
#   - 高效文件传输:使用 tar 流在宿主机和容器之间传输文件,支持批量操作。
#   - 资源限制:支持设置 CPU (cpu_quota) 和内存 (memory_limit) 限制,防止 Agent 耗尽本机资源。
#   - 网络控制:可选禁用网络 (network_disabled=True) 以增强安全性。

import io
import tarfile
import time
import uuid
from typing import Optional

# 加载 DeepAgents 后端协议
from deepagents.backends.protocol import (
    ExecuteResponse,
    FileDownloadResponse,
    FileUploadResponse,
    SandboxBackendProtocol,
)

# 加载 DeepAgents 基础沙箱类
from deepagents.backends.sandbox import BaseSandbox

try:
    import docker
    from docker.errors import NotFound, APIError
except ImportError:
    docker = None

class DockerBackend(BaseSandbox):
    """Docker 沙箱后端实现,用于 DeepAgents。

    该后端使用本地 Docker 守护进程提供隔离的执行环境。
    需要安装 `docker` Python 包,并确保 Docker 守护进程正在运行。
    """
    def __init__(
        self,
        image: str = "python:3.11-slim",
        auto_remove: bool = True,
        cpu_quota: int = 50000,  # 50% CPU
        memory_limit: str = "512m",
        network_disabled: bool = False,
        working_dir: str = "/workspace",
        volumes: dict[str, dict[str, str]] | None = None,
    ) -> None:
        """初始化 Docker 沙箱。

        参数:
            image: 使用的 Docker 镜像(默认:"python:3.11-slim")
            auto_remove: 是否在关闭时移除容器(默认:True)
            cpu_quota: CPU 配额,单位为微秒(默认:50000)
            memory_limit: 内存限制(默认:"512m")
            network_disabled: 是否禁用网络访问(默认:False)
            working_dir: 容器内的工作目录(默认:"/workspace")
            volumes: Docker 卷配置,例如 {'/宿主机路径': {'bind': '/容器路径', 'mode': 'rw'}}
        """
        if docker is None:
            raise ImportError(
                "docker package is not installed. "
                "Please install it with `pip install docker`."
            )
        # 初始化 Docker 客户端,from_env() 会自动从环境变量中读取 Docker 配置
        self.client = docker.from_env()
        self.image = image
        self.auto_remove = auto_remove
        self.working_dir = working_dir
        self.volumes = volumes or {}
        self._container = None

        # Start container
        try:
            # Ensure image exists
            try:
                self.client.images.get(image)
            except NotFound:
                print(f"Pulling image {image}...")
                self.client.images.pull(image)

            self._container = self.client.containers.run(
                image,
                command="tail -f /dev/null",  # Keep container running
                detach=True,
                tty=True,
                cpu_quota=cpu_quota,
                mem_limit=memory_limit,
                network_disabled=network_disabled,
                working_dir=working_dir,
                volumes=self.volumes,
            )

            # Ensure working directory exists
            self.execute(f"mkdir -p {working_dir}")

        except Exception as e:
            raise RuntimeError(f"Failed to start Docker container: {e}")

    @property
    def id(self) -> str:
        """Unique identifier for the sandbox backend."""
        return self._container.id if self._container else "unknown"

    def execute(self, command: str) -> ExecuteResponse:
        """Execute a command in the sandbox."""
        if not self._container:
            return ExecuteResponse(
                output="Container not running",
                exit_code=1,
                truncated=False
            )

        try:
            # Docker exec_run 返回 (exit_code, output)
            # output 是字节类型
            # 使用列表形式的 cmd 以避免 shell 转义问题
            exit_code, output = self._container.exec_run(
                cmd=["bash", "-c", command],
                workdir=self.working_dir,
                demux=False # Combine stdout and stderr
            )

            return ExecuteResponse(
                output=output.decode("utf-8", errors="replace"),
                exit_code=exit_code,
                truncated=False,
            )
        except Exception as e:
            return ExecuteResponse(
                output=f"Error executing command: {str(e)}",
                exit_code=1,
                truncated=False,
            )

    def upload_files(self, files: list[tuple[str, bytes]]) -> list[FileUploadResponse]:
        """Upload multiple files to the sandbox using tar archive."""
        if not self._container:
            return [FileUploadResponse(path=p, error="permission_denied") for p, _ in files]

        responses = []

        # Create a tar archive in memory
        tar_stream = io.BytesIO()
        with tarfile.open(fileobj=tar_stream, mode='w') as tar:
            for path, content in files:
                # Docker put_archive expects relative paths inside the tar to be relative to the destination
                # But here we want absolute paths to be respected.
                # Actually put_archive extracts to a directory.
                # To support absolute paths, we should probably upload to root /?
                # Or handle relative paths relative to working_dir.

                # Let's handle paths:
                # If path is absolute, we strip leading / and upload to root.
                # If path is relative, we upload to working_dir.

                # Simplification: We will create a tar with full structure and extract to /

                # Normalize path
                if path.startswith("/"):
                    arcname = path.lstrip("/")
                    dest_path = "/"
                else:
                    arcname = path
                    dest_path = self.working_dir

                info = tarfile.TarInfo(name=arcname)
                info.size = len(content)
                info.mtime = time.time()
                tar.addfile(info, io.BytesIO(content))

                responses.append(FileUploadResponse(path=path, error=None))

        tar_stream.seek(0)

        try:
            # We extract to / to support absolute paths in the tar
            # Note: This assumes all files in the batch can be extracted to the same root.
            # If mixed absolute/relative, this might be tricky.
            # For robustness, we might need to upload one by one if paths are mixed,
            # or group them.
            # Strategy: Always extract to / (root), and ensure arcnames are full paths (without leading /)

            self._container.put_archive(
                path="/",
                data=tar_stream
            )
        except Exception as e:
            # Mark all as failed if batch fails
            return [FileUploadResponse(path=p, error="permission_denied") for p, _ in files]

        return responses

    def download_files(self, paths: list[str]) -> list[FileDownloadResponse]:
        """Download multiple files from the sandbox."""
        if not self._container:
            return [FileDownloadResponse(path=p, error="permission_denied") for p in paths]

        responses = []
        for path in paths:
            try:
                # get_archive returns a tuple (generator, stat)
                bits, stat = self._container.get_archive(path)

                # Reconstruct tar from bits
                file_content = io.BytesIO()
                for chunk in bits:
                    file_content.write(chunk)
                file_content.seek(0)

                # Extract file from tar
                with tarfile.open(fileobj=file_content, mode='r') as tar:
                    # There should be only one file/dir
                    member = tar.next()
                    if member.isdir():
                        responses.append(FileDownloadResponse(path=path, error="is_directory"))
                        continue

                    f = tar.extractfile(member)
                    if f:
                        content = f.read()
                        responses.append(FileDownloadResponse(path=path, content=content, error=None))
                    else:
                        responses.append(FileDownloadResponse(path=path, error="file_not_found"))

            except NotFound:
                responses.append(FileDownloadResponse(path=path, error="file_not_found"))
            except Exception as e:
                error_msg = str(e).lower()
                error = "invalid_path"
                if "permission" in error_msg:
                    error = "permission_denied"

                responses.append(FileDownloadResponse(path=path, content=None, error=error))
        return responses

    def close(self):
        """Close the sandbox session."""
        if self._container:
            try:
                if self.auto_remove:
                    self._container.remove(force=True)
                else:
                    self._container.stop()
            except Exception:
                pass
            self._container = None

使用自定义的 DockerBackend 实现沙箱环境

import asyncio
import os
import sys
from dotenv import load_dotenv
from deepagents import create_deep_agent
from langchain_openai import ChatOpenAI
from langchain_core.messages import BaseMessage

# 尝试导入 DockerBackend
try:
    from docker_backend import DockerBackend
except ImportError:
    sys.path.append(os.path.dirname(os.path.abspath(__file__)))
    from docker_backend import DockerBackend

load_dotenv(override=True)

async def run_docker_demo():
    # 1. 检查 Docker 环境
    try:
        import docker
        docker.from_env().ping()
    except Exception as e:
        return

    # 2. 初始化 Docker Backend
    try:
        backend = DockerBackend(
            image="python:3.11-slim",
            auto_remove=True
        )
    except Exception as e:
        return

    try:
        # 3. 创建 Agent
        llm = ChatOpenAI(model="gpt-4o-mini", temperature=0)

        agent = create_deep_agent(
            model=llm,
            backend=backend,
            system_prompt="""你是一个运行在 Docker 容器中的 AI 助手。
            你的任务是演示环境隔离性。

            请执行以下步骤:
            1. 运行 'cat /etc/os-release' 查看容器操作系统。
            2. 运行 'python --version' 确认 Python 环境。
            3. 创建文件 '/workspace/hello_docker.py',内容为打印 'Hello from Docker Container!'。
            4. 运行该脚本。
            """
        )

DockerBackend 和 E2BBackend 的区别

特性 DockerBackend E2BBackend
核心定位 本地轻量级容器化沙箱 云端安全沙箱环境 (SaaS)
部署位置 运行在本地机器 (Localhost) 运行在 E2B 云端集群 (Remote Cloud)
依赖环境 需要本地安装并运行 Docker Desktop / Daemon 仅需安装 e2b Python SDK,无需本地 Docker
资源消耗 消耗本地 CPU / 内存资源 消耗 E2B 云端资源 (不占用本地算力)
启动速度 快 (本地镜像启动,毫秒 - 秒级) 较快 (云端冷启动约 1 - 3秒)
网络隔离 可配置 (支持完全离线 network_disabled=True) 默认联网 (支持访问公网 API)
持久化 支持挂载本地卷 (Volumes) 实现数据持久化 临时环境 (会话结束即销毁),数据需手动导出
适用场景 • 本地开发 / 调试
• 数据隐私敏感 (不想数据出本地)
• 离线环境使用
• 生产环境部署 (无需维护 Docker)
• 多租户隔离 (每个用户一个云沙箱)
• 本地资源受限设备
成本 免费 (使用自有硬件) 付费 (按使用时长 / 资源计费)
配置复杂度 中 (需管理镜像、卷挂载、Docker 进程) 低 (API Key 开箱即用)

6. StoreBackend 使用数据库存储

import asyncio
import os
import uuid
import traceback
from dotenv import load_dotenv

# LangChain / LangGraph Imports
from langchain_openai import ChatOpenAI
from langchain_mcp_adapters.client import MultiServerMCPClient
from langchain_core.messages import BaseMessage, ToolMessage
from langgraph.store.postgres import PostgresStore
from langgraph.checkpoint.memory import MemorySaver
from psycopg_pool import ConnectionPool

# DeepAgents Imports
from deepagents import create_deep_agent
from deepagents.backends import StoreBackend

load_dotenv(override=True)

DB_URI = "postgresql://myuser:123456@localhost:5432/mydatabase"

async def setup_mcp_tools():
    """
    连接 Context7 MCP 服务器并获取工具。
    """
    try:
        client = MultiServerMCPClient({
            "context7": {
                "transport": "stdio",
                "command": "npx",
                "args": ["-y", "@upstash/context7-mcp@latest"],
            }
        })

        # 获取工具列表
        tools = await client.get_tools()

        return client, tools
    except Exception as e:
        return None, []

async def run_store_backend_demo():
    """
    运行 StoreBackend 演示:
    展示如何使用 PostgreSQL 作为 DeepAgents 的文件系统后端。
    """
    print_header()

    # Step 1. 初始化 MCP 工具 (Connect to Context7 MCP Server)
    mcp_client, mcp_tools = await setup_mcp_tools()

    # Step 2. 初始化 PostgreSQL 连接池和存储组件 (Init DB & Store)
    # 使用 ConnectionPool 管理数据库连接
    # StoreBackend 需要同步的 ConnectionPool (因为它在线程中运行同步操作)
    try:
        with ConnectionPool(conninfo=DB_URI, kwargs={"autocommit": True}) as pool:
            # 2.1 初始化 Checkpointer (用于保存会话状态 / 聊天记录)
            # 注意: 这里使用 MemorySaver 是为了避开 langgraph-checkpoint-postgres 在某些环境下的异步兼容性问题
            checkpointer = MemorySaver()

            # 2.2 初始化 Store (用于保存文件 / 长期记忆)
            # DeepAgents 的 StoreBackend 会将文件系统操作映射到这个 PostgresStore
            store = PostgresStore(pool)

            # 2.3 执行数据库迁移 (首次运行需要初始化表结构)
            # 这会在数据库中创建 'store' 表,用于存储 JSON 数据 (即我们的文件)
            with pool.connection() as conn:
                with conn.cursor() as cur:
                    for migration in store.MIGRATIONS:
                        cur.execute(migration)

            # Step 3. 创建基于 StoreBackend 的 Agent (Create Agent)
            # StoreBackend 将文件系统操作映射到 LangGraph 的 Store
            # 这意味着文件将持久化存储在 PostgreSQL 数据库中,跨重启可用
            llm = ChatOpenAI(model="gpt-4o", temperature=0)

            # 定义 Backend Factory
            # 这是一个工厂函数,用于在运行时(runtime)把 StoreBackend 实例化。
            # create_deep_agent 会在内部把 LangGraph 的 runtime 对象作为参数 rt 传进来,
            # 这样 StoreBackend 就能拿到 store / checkpointer 等资源,实现文件系统到 PostgreSQL 的映射。
            backend_factory = lambda rt: StoreBackend(rt)

            agent = create_deep_agent(
                model=llm,
                tools=mcp_tools,
                backend=backend_factory,
                store=store,              # 关键: 必须传入 store 实例
                checkpointer=checkpointer, # 关键: 传入 checkpointer 实例
                system_prompt="""你是一个高级技术助手。
                你的任务是使用 Context7 工具查询关于 'DeepAgents StoreBackend' 的用法。
                查询后,创建一个总结文件 '/knowledge/store_backend_notes.md',并写入关键信息。
                由于你使用的是 StoreBackend,这个文件将直接存储在 PostgreSQL 数据库中。
                最后,请读取该文件以验证存储成功。"""
            )
            
            # Step 4. 执行任务 (Execute Task)
            thread_id = str(uuid.uuid4())
            config = {"configurable": {"thread_id": thread_id}}

            task = "请查询 StoreBackend 的用法,并将总结写入 /knowledge/store_backend_notes.md,最后读取它验证。"
参数组件 Backend (后端) Checkpointer (检查点) Store (存储)
核心定义 环境层 (Environment) 状态层 (State / Short-term Memory) 记忆层 (Memory / Long-term Memory)
负责什么? "外部世界" 的交互能力。
即:文件存在哪?代码在哪跑?
"当前对话" 的上下文。
即:刚才说了什么?现在运行到哪一步了?
"跨会话" 的知识积累。
即:用户叫什么名字?上次任务学到了什么?
数据类型 非结构化文件 (.py.md.txt)
运行时环境 (Shell, Process)
BaseMessage 列表 (User / AI / Tool Message)
Graph 节点状态
结构化 JSON 数据 (Key-Value)
用户偏好、长期笔记
典型实现 DockerBackend (容器)
FilesystemBackend (磁盘)
E2BBackend (云沙箱)
MemorySaver (内存)
PostgresSaver (数据库)
SqliteSaver (本地DB)
InMemoryStore (内存)
PostgresStore (数据库)
生命周期 任务级
(任务结束容器可能销毁)
线程级 (Thread)
(换个 thread_id 就没了)
全局级 (Global)
(所有 thread 都能查到)
形象比喻 工作台 / 电脑 大脑的工作记忆 (只会死记硬背当前对话) 日记本 / 知识库 (记录永久信息)

7. CompositeBackend 使用混合模式

在单一后端模式下,Agent 所有的文件操作(读、写、列出目录)都只能去往同一个地方(要么全是本地磁盘,要么全是 Docker 容器内)。而 CompositeBackend 允许你根据 文件路径前缀 ,将请求分发给不同的后端。

核心优势

  1. 性能与开销优化 (Performance) 这是最关键的技术优势。
    • DockerBackend 的局限 : 向 Docker 容器内读写文件(尤其是大文件)需要经过 Docker Daemon 的 API (如 put_archive / get_archive ),涉及网络通信和打包解包,开销较大。
    • CompositeBackend 的解法 : 对于数据文件( /data ),直接通过 FilesystemBackend 进行本地 I/O 操作, 完全绕过了 Docker API ,读写速度是操作系统原生的速度。
  2. 计算与存储分离 (Decoupling)
    • 计算是临时的 : processor.py 脚本可能只需要运行一次,运行环境(Python 依赖)可能很复杂且容易冲突。放在 Docker 里最合适。
    • 数据是永恒的 : raw_metrics.txt 和 health_report.txt 是业务资产。通过路由直接落盘到宿主机,即使 Docker 容器崩溃、被删除或重启, 数据毫发无损且立即可在宿主机访问
  3. 给予 Agent "混合云" 的能力,Agent 可以像人类工程师一样工作:
    • "我在临时的沙箱里写代码测试(Docker)。"
    • "测试好了,我把结果保存到公司的共享网盘里(Filesystem / Mount)。"
import asyncio
import shutil
import os
import time
from pathlib import Path
from dotenv import load_dotenv

# DeepAgents 导入
from deepagents import create_deep_agent
from deepagents.backends.composite import CompositeBackend
from deepagents.backends.filesystem import FilesystemBackend
from langchain_openai import ChatOpenAI
from langchain_mcp_adapters.client import MultiServerMCPClient
from langchain_core.messages import BaseMessage, ToolMessage

# 导入 DockerBackend
try:
    from docker_backend import DockerBackend
except ImportError:
    try:
        from deepagents.backends.docker import DockerBackend
    except ImportError:
        DockerBackend = None

async def setup_mcp_tools():
    try:
        client = MultiServerMCPClient({
            "context7": {
                "transport": "stdio",
                "command": "npx",
                "args": ["-y", "@upstash/context7-mcp@latest"],
            }
        })
        tools = await client.get_tools()
        return client, tools
    except Exception as e:
        return None, []

async def run_composite_demo():
    load_dotenv(override=True)
    print_header()

    if DockerBackend is None:
        return

    # Step 1
    host_work_dir = Path("workspace/data_analysis_project").resolve()
    if host_work_dir.exists():
        shutil.rmtree(host_work_dir)
    host_work_dir.mkdir(parents=True, exist_ok=True)

    container_mount_path = "/data"
    docker_volumes = {
        str(host_work_dir): {'bind': container_mount_path, 'mode': 'rw'}
    }

    # Step 2
    fs_backend = FilesystemBackend(root_dir=host_work_dir, virtual_mode=True)

    docker_backend = DockerBackend(
        image="python:3.11-slim",
        auto_remove=True,
        volumes=docker_volumes
    )

    routes = {
        container_mount_path: fs_backend
    }
    backend = CompositeBackend(default=docker_backend, routes=routes)

    # Step 3
    mcp_client, mcp_tools = await setup_mcp_tools()

    system_prompt = f"""你是一名在混合环境中工作的高级数据工程师。

    环境地图:
    1. 执行层 (根目录 `/`):
       - 临时的 Docker 容器。
       - 用于创建脚本 (`.py`) 和运行命令。
       - 这里的文会在会话结束后消失。

    2. 存储层 (`{container_mount_path}`):
       - 从宿主机挂载的持久化存储。
       - 用于存放 输入 数据和 输出 报告。
       - 这里的文件会永久保存。

    你的任务:
    1. **摄入**: 创建一个文件 `{container_mount_path}/raw_metrics.txt`,内容为 "CPU: 45%, Mem: 60%"。
       (注意: 这使用了 'write_file' 工具,该工具通过路由直接写入宿主机文件系统)。

    2. **处理**: 创建一个 Python 脚本 `/processor.py` (在根目录),该脚本:
       - 读取 `{container_mount_path}/raw_metrics.txt`。
       - 计算 "健康分数" (模拟一下即可)。
       - 将报告写入 `{container_mount_path}/health_report.txt`。
       - 打印 "Analysis Complete"。

    3. **执行**: 使用 `python /processor.py` 运行脚本。
       (注意: 这在 Docker 内部运行。Docker 因为卷挂载能看到这些文件)。

    4. **验证**: 读取 `{container_mount_path}/health_report.txt` 并显示它。
    """

    agent = create_deep_agent(
        model=ChatOpenAI(model="gpt-4o", temperature=0),
        tools=mcp_tools,
        backend=backend,
        system_prompt=system_prompt
    )

5.3. interrupt_on 参数

1. Human-in-the-loop(HITL) 人工干预

  1. Human-in-the-loop (HITL):对于 write_file 或 delegate_task 等关键操作,利用 LangGraph 的中断机制加入人工审批。
  2. 异步运行:DeepAgent 的任务通常耗时较长,务必使用异步 Webhooks 接收结果。
  3. 监控与调试:强烈建议结合 LangSmith 使用。由于 DeepAgent 内部有复杂的子 Agent 递归调用,使用 LangSmith 的 Tracing 功能是排查问题的有效手段。
  4. 后端挂载:在生产环境中,建议将 VFS 挂载到云端存储(如 S3),以防止容器重启导致 Agent 的"记忆"丢失。
  5. interrupt_on:这个参数其实是一个 HITL 的开关,就是把 HITL 的中间件插入到 DeepAgent 的执行流程中,当 DeepAgent 执行到需要人工审批的操作时,就会中断执行,等待人工审批。
    • 类型: dict[str, bool | InterruptOnConfig]
    • 作用: 映射“工具名称”到“中断配置”。
    • 示例: interrupt_on={"write_file": True} 表示当 Agent 试图调用 write_file 工具时,程序会暂停(Suspend),等待人工(Human)在 LangGraph 层面进行 Approve 、 Reject 或 Edit 操作后才能继续。
import os
import asyncio
from typing import Optional, Set
from langchain_openai import ChatOpenAI
from deepagents import create_deep_agent
from deepagents.backends import FilesystemBackend
from langgraph.checkpoint.memory import InMemorySaver
from dotenv import load_dotenv
from langchain_tavily import TavilySearch
from langchain.agents.middleware.human_in_the_loop import (
    HITLResponse,
    ApproveDecision,
    EditDecision,
    RejectDecision
)
from langgraph.types import Command
from langchain_core.messages import BaseMessage, ToolMessage, AIMessage

load_dotenv(override=True)

async def run_interrupt_test():
    """
    示例 1: 基础中断功能 (封装版)
    在工具调用前中断,让用户确认是否继续执行
    """
    llm = ChatOpenAI(model="gpt-4o", temperature=0)
    search_tool = TavilySearch(max_results=2)

    # 创建 Agent,设置在 "tools" 节点中断
    agent = create_deep_agent(
        model=llm,
        tools=[search_tool],
        backend=FilesystemBackend(root_dir="./workspace", virtual_mode=True),
        checkpointer=InMemorySaver(),  # 必需!用于支持中断和恢复
        interrupt_on={"tavily_search": True},  # 在特定工具调用时中断
    )

    task = "搜索 'Python 异步编程' 的最新信息,并创建一个总结文件"

    # 为了避免之前的状态干扰,我们使用一个新的 thread_id
    config = {"configurable": {"thread_id": "demo_basic_refactored_v1"}}

    # 追踪已打印的消息数量,避免重复打印
    message_history_len = 0

    # --- 第一次执行 ---
    async for event in agent.astream({"messages": [("user", task)]}, config=config):
        if "messages" in event:
            current_messages = event["messages"]
            if len(current_messages) > message_history_len:
                # 打印新增的消息
                for i in range(message_history_len, len(current_messages)):
                    msg = current_messages[i]
                    if msg.type == "ai":
                        if hasattr(msg, 'tool_calls') and msg.tool_calls:
                            print(f"🔧 AI 决定调用工具: {msg.tool_calls[0]['name']}")
                            print(f"   参数: {msg.tool_calls[0]['args']}")
                        elif msg.content:
                            print(f"💬 AI: {msg.content}")
                    elif msg.type == "tool":
                        print(f"✅ 工具输出: {msg.content[:100]}..." if len(msg.content) > 100 else f"✅ 工具输出: {msg.content}")

                message_history_len = len(current_messages)

    # 检查是否中断
    # 使用 aget_state (async) 获取状态
    state = await agent.aget_state(config)
    print(f"\n⏸️  执行状态: {state.next}")

    if state.tasks:
        print(f"\n--- 🛑 执行已暂停 (HITL Middleware) ---")
        last_message = state.values["messages"][-1]

        if hasattr(last_message, "tool_calls") and last_message.tool_calls:
            tool_call = last_message.tool_calls[0]

            # === 人工介入 ===
            approval = input("\n[管理员]: 是否批准执行此操作? (y/n/e[编辑]): ")

            if approval.lower() == 'y':
                print("\n[系统]: 操作已批准,继续执行...")

                hitl_response = HITLResponse(
                    decisions=[ApproveDecision(type="approve")]
                )

                # === 恢复执行 ===
                # 使用 Command(resume=...)
                async for event in agent.astream(
                    Command(resume=hitl_response),
                    config=config,
                    stream_mode="values"
                ):
                    if "messages" in event:
                        current_messages = event["messages"]
                        if len(current_messages) > message_history_len:
                            for i in range(message_history_len, len(current_messages)):
                                msg = current_messages[i]

                                # 优化打印逻辑,清晰展示 AI 回复
                                if msg.type == "tool":
                                    print(f"\n[工具输出]:\n{msg.content[:300]}..." if len(msg.content) > 300 else f"\n[工具输出]:\n{msg.content}")
                                elif msg.type == "ai":
                                    if msg.content:
                                        print(f"\n[AI 回复]:\n{msg.content}\n")
                                    elif msg.tool_calls:
                                        print(f"\n🔧 AI 决定调用工具: {msg.tool_calls[0]['name']}")
                                        print(f"   参数: {msg.tool_calls[0]['args']}")

                            message_history_len = len(current_messages)

            else:
                print("\n[系统]: 操作被拒绝或您选择了其他选项 (本演示仅处理 'y')。")
    else:
        print("流程已完成,没有触发中断。")
        if state.values.get("messages"):
            last_msg = state.values["messages"][-1]
            if last_msg.type == "ai" and last_msg.content:
                print(f"\n[最终回复]: {last_msg.content}")

if __name__ == "__main__":
    try:
        # asyncio.run(run_interrupt_test())
        await run_interrupt_test()
    except KeyboardInterrupt:
        print("\n程序已停止")