003.Agent 智能体中间件应用实战

一:Agent 开发的可观测性基石

对于开发者来说,还需要进一步了解和掌握 LangChain Agent 必备的开发者套件。分别是 LangChain Agent 运行监控框架 LangSmith、底层 LangGraph 图结构可视化与调试框架 LangGraph Studio 和 LangGraph 服务部署工具 LangGraph Cli。可以说这些开发工具套件,是真正推动 LangGraph 的企业级应用开发效率大幅提升的关键。同时监控、调试和部署工具,也是全新一代企业级 Agent 开发框架的必备工具,也是开发者必须要掌握的基础工具。

1.1. LangGraph 图结构可视化与调试框架:LangGraph Studio

langGraph Studio 对于 LangChian Agent 来说,则是比 LangSmith 更加方便和高效的可视化调试工具平台。

LangGraph Studio 在本地可视化运行时会自动把调用过程上传到 LangSmith;而在 LangSmith 网页端查看任何 Trace 时,又能一键 Run in Studio 回放整条执行链,所以它是通过统一 Trace SDK 与 LangSmith 紧密集成。而 LangGraph CLI 则是构建这个项目的关键

1.2. LangGraph 服务部署工具:LangGraph Cli

LangGraph CLI 是用于本地启动、调试、测试和托管 LangGraph 智能体图的开发者命令行工具。

功能类别 命令示例 说明
✅ 启动 Graph 服务 langgraph dev 启动 Graph 的开发服务器,供前端(如 Agent Chat UI)调用
🧪 测试 Graph 输入 langgraph run graph:graph --input '{"input": "你好"}' 本地 CLI 输入测试,输出结果
🧭 管理项目结构 langgraph init 初始化一个标准 Graph 项目目录结构
📦 部署 Graph(未来) langgraph deploy(预留) 发布 graph 至 LangGraph 云端(已对接 Studio)
🧱 显示 Assistant 列表 langgraph list 显示当前 graph 中有哪些 assistant(即 entrypoint)
🔄 重载运行时 自动热重载 修改 graph.py 时,dev 模式自动重启生效

而一旦应用成功部署上线,LangGraph Cli 还会非常贴心的提供后端接口说明文档:

Pasted image 20260707153630.png

1.3. 创建完整 LangGraph 智能体项目流程

import os
from dotenv import load_dotenv
from langchain_deepseek import ChatDeepSeek
from langchain.agents import create_agent
from typing_extensions import TypedDict
from langchain_tavily import TavilySearch
from langchain_core.tools import tool
from pydantic import BaseModel, Field
import requests,json

load_dotenv(override=True)

# 内置搜索工具
search_tool = TavilySearch(max_results=5, topic="general")

class WeatherQuery(BaseModel):
    loc: str = Field(description="The location name of the city")

@tool(args_schema = WeatherQuery)
def get_weather(loc):
    """
    查询即时天气函数
    :param loc: 必要参数,字符串类型,用于表示查询天气的具体城市名称,\
    注意,中国的城市需要用对应城市的英文名称代替,例如如果需要查询北京市天气,则 loc 参数需要输入 'Beijing';
    :return:OpenWeather API 查询即时天气的结果,具体 URL 请求地址为:https://api.openweathermap.org/data/2.5/weather\
    返回结果对象类型为解析之后的 JSON 格式对象,并用字符串形式进行表示,其中包含了全部重要的天气信息
    """
    # Step 1. 构建请求
    url = "https://api.openweathermap.org/data/2.5/weather"

    # Step 2. 设置查询参数
    params = {
        "q": loc,
        "appid": os.getenv("OPENWEATHER_API_KEY"),    # 输入 API key
        "units": "metric",            # 使用摄氏度而不是华氏度
        "lang":"zh_cn"                # 输出语言为简体中文
    }

    # Step 3. 发送 GET 请求
    response = requests.get(url, params=params)

    # Step 4. 解析响应
    data = response.json()
    return json.dumps(data)

tools = [search_tool, get_weather]

# 创建模型
model = ChatDeepSeek(model="deepseek-chat")

prompt = """
你是一名乐于助人的智能助手,擅长根据用户的问题选择合适的工具来查询信息并回答。
当用户的问题涉及**天气信息**时,你应优先调用 `get_weather` 工具,查询用户指定城市的实时天气,并在回答中总结查询结果。
当用户的问题涉及**新闻、事件、实时动态**时,你应优先调用 `search_tool` 工具,检索相关的最新信息,并在回答中简要概述。
如果问题既包含天气又包含新闻,请先使用 `get_weather` 查询天气,再使用 `search_tool` 查询新闻,最后将结果合并后回复用户。

所有回答应使用**简体中文**,条理清晰、简洁友好。
"""

# 创建图
graph = create_agent(
    model=model,
    tools=tools,
    system_prompt=prompt)

代码编写时需要注意,如果需要使用 langgraph studio 进行可视化调试,则需要注意下面两点:

Pasted image 20260707154205.png

Pasted image 20260707154249.png


二:LangChain 1.0 中间件 (Middleware) 概览

2.1. 中间件架构原理

中间件(Middleware)系统,旨在解决早期版本中 Agent 抽象无法灵活定制的痛点。中间件(Middleware)被正式定义为一种用于拦截、修改、控制和增强 Agent 执行流程的机制。允许开发者在模型调用前后、Agent 启动前后或工具调用前后插入自定义逻辑,从而实现日志记录、权限控制、上下文压缩等功能,而无需修改 Agent 的核心业务逻辑。

中间件的工作原理,我们可以将其比作一个洋葱。每一层中间件都包裹着核心的 Agent 功能,就像洋葱的每一层都有其特定的作用。当用户请求到达时,它必须逐层通过这些中间件,每一层都会对请求进行特定的处理,然后将其传递给下一层。简而言之,借助中间件,一个 React Agent 的运行模型,就可以由这种:

Pasted image 20260707163321.png

变为这种:

Pasted image 20260707163330.png

2.2. 中间件的作用

Pasted image 20260707172556.png

2.3. 核心设计理念的深度解析

中间件的设计严格遵循了 SOLID 原则,这些原则不仅仅是理论概念,而是实际指导开发的行动准则。

单一职责原则(Single Responsibility Principle:SRP)在中间件中得到了完美的体现。每一个中间件都专注于一个特定的横切关注点,比如 SummarizationMiddleware 只负责处理上下文压缩,CostTrackingMiddleware 只负责成本统计。这种设计使得每个中间件的职责边界清晰,便于测试、维护和重用。

开闭原则(Open/Closed Principle:OCP)的应用使得系统具备了良好的扩展性。ModelSelectorMiddleware 可以通过配置来支持不同的模型选择策略,而无需修改代码本身。这意味着当有新的模型或者新的选择策略出现时,我们可以通过配置文件来适应,而不是修改核心代码。

里氏替换原则(Liskov Substitution Principle:LSP)是子类必须能够替换其父类而不破坏程序的正确性,DatabaseAuthMiddleware 和 JWTAuthMiddleware 可互相替换,Agent 无差别运行。

装饰器模式(Interface Segregation Principle:ISP)是中间件架构中的核心模式。通过 wrap_model_call,我们可以在不修改原始 Agent 代码的情况下,为其添加额外的功能。这种无侵入性的扩展方式使得 Agent 的核心逻辑保持简洁,而功能的增强通过中间件来实现。

责任链模式(Dependency Inversion Principle:DIP)则体现在多个中间件的顺序处理中。每个中间件都有机会处理请求,然后将请求传递给下一个中间件。这种模式提供了极大的灵活性,我们可以根据具体需求调整中间件的执行顺序,甚至动态地添加或删除中间件。

2.4. 性能考量与优化策略

不同的中间件类型对系统性能的影响是不同的。监控类中间件通常对性能影响最小,因为它们主要是记录信息而不进行复杂的处理。修改类中间件的影响相对较大,因为它们需要对数据进行处理和转换。控制类中间件可能会引入一些延迟,因为它们需要等待外部决策(如人工审批)。强制类中间件则可能涉及复杂的验证逻辑,对性能有一定影响。

优化策略需要从多个层面进行考虑。延迟优化可以通过异步处理来实现,对于那些不直接影响主流程的操作(如日志记录、性能统计),我们可以采用异步方式处理。批处理是另一个重要的优化手段,特别是对于那些需要对多个请求进行相似处理的中间件。

Pasted image 20260707173958.png

中间件类型 典型场景 响应时间增加 吞吐量衰减 资源消耗 性能瓶颈点 优化策略
日志监控中间件
@after_model
请求日志、指标统计 < 1 ms < 5% CPU 1-3%
内存 10-50MB
磁盘 I/O
(异步后极小)
异步写入、采样率 10%
安全脱敏中间件
@before_model
PII 脱敏、权限检查 1-5 ms 5-10% CPU 5-10%
内存 20-100MB
正则表达式
字符串拷贝
编译缓存、原地修改
参数校验中间件
@before_model
输入验证、格式检查 1-3 ms 3-8% CPU 3-8%
内存 10-30MB
复杂校验规则 缓存校验结果
对话总结中间件
@before_model
长文本压缩 50-200 ms 15-30% CPU 20-40%
内存 100-500MB
模型调用
token 计算
仅在消息数 >10 触发
缓存中间件
@wrap_model_call
结果缓存 0.5-2 ms
(缓存命中)
-50% ~ +10%
(命中时提升)
内存 50-200MB 缓存命中率 LRU 策略、TTL 设置
模型降级中间件
@wrap_model_call
动态切换模型 5-10 ms 10-15% CPU 2-5%
内存 10-20MB
模型初始化
API 调用
模型池化、预热
限流中间件
@wrap_model_call
QPS 限制、熔断 0.1-1 ms < 5% CPU 1-2%
内存 5-10MB
原子计数器
锁竞争
滑动窗口算法
工具审计中间件
@wrap_tool_call
调用记录、熔断 2-8 ms 8-15% CPU 5-12%
内存 30-80MB
数据库写入
网络请求
异步批量写入
工具重试中间件
@wrap_tool_call
失败重试 10-50 ms
(视重试次数)
20-50% CPU 10-20%
内存 20-50MB
重试延迟
指数退避
限制重试次数 ≤3
外部 API 调用中间件
@wrap_tool_call
调用第三方服务 100-500 ms
(网络延迟主导)
30-70% CPU 5-10%
内存 10-30MB
网络 I/O
超时设置
连接池、超时 3s
消息队列中间件
@before_model
异步任务解耦 5-20 ms
(仅入队)
10-25% CPU 10-15%
内存 50-150MB
消息序列化
队列持久化
批量发送、压缩
注册中心中间件
@before_agent
服务发现 10-30 ms
(首次查询)
5-15% CPU 5-10%
内存 20-60MB
DNS 查询
缓存过期
本地缓存 60s

三:中间件的分类与应用场景

3.1. 中间件四大分类

中间件的四大分类——监控类、修改类、控制类和强制类——每一种都有其独特的定位和应用场景。这种分类不是人为的划分,而是基于实际需求和功能特性的自然归类。

分类 核心功能 解决的问题 典型应用场景
Monitor (监控类) 观察执行状态、日志记录 调试困难、缺乏可观测性 记录所有的 Prompt 和 Response、性能分析、成本核算。
Modify (修改类) 修改输入 / 输出、上下文管理 上下文窗口溢出、Prompt 优化 SummarizationMiddleware(自动压缩历史对话)、动态注入 System Prompt。
Control (控制类) 流程阻断、人工介入 AI 幻觉、高风险操作失控 HumanInTheLoopMiddleware(敏感操作需人工审批)、重试机制。
Enforce (强制类) 安全过滤、限流、合规检查 数据泄露、API 滥用 PIIMiddleware(敏感信息脱敏)、ModelCallLimit(防止死循环)。

3.2. 深入生命周期 - 中间件的 6 个切入点 (Hooks)

1. Hook 执行顺序的深层逻辑

中间件的 6 个 Hook 点——before_agent、before_model、wrap_model_call、wrap_tool_call、after_model 和 after_agent——不仅仅是一个执行顺序,更是一个精心设计的数据流处理流程。这个流程体现了从输入处理到输出生成的完整生命周期。

在 Agent 执行的最开始,before_agent Hook 提供了全局初始化的机会。这个阶段通常用于设置全局状态、检查环境配置、初始化资源等。

before_model 阶段是输入预处理的关键节点。在这个阶段,中间件可以对输入数据进行预处理、验证、清洗等操作。这是确保数据质量的第一道防线,任何在这个阶段发现的问题都可以避免后续的无谓计算。

wrap_model_call 是最核心的 Hook,它包装了实际的模型调用过程。这个阶段的处理逻辑决定了如何与底层模型交互,是实现高级功能(如缓存、重试、熔断等)的关键位置。

wrap_tool_call 是用于拦截和控制工具的实际执行过程的 Hook,它包装了每次工具调用。这个阶段的处理逻辑可以是权限、重试、日志、审批。

after_model 阶段则处理模型返回的原始结果。这个阶段的任务是验证输出质量、进行格式转换、提取关键信息等。由于模型输出往往包含大量无用信息,这个阶段的处理对于提高整体效率至关重要。after_agent 阶段是整个生命周期的收尾工作。在这个阶段,系统需要清理资源、记录最终状态、生成报告等。这是确保系统处于良好状态,为下一次请求做好准备的关键步骤。

2. 数据传递机制的复杂性

每个 Hook 都能访问和修改共享的上下文对象,这个对象就像是整个处理过程的"记忆"。上下文对象不仅包含请求信息和响应数据,还包含运行时状态、中间计算结果、元数据等。

元数据传递机制允许中间件在不直接共享状态的情况下传递信息。例如,一个中间件可以在元数据中标记某个输入是 VIP 用户的请求,后续的中间件可以根据这个标记来调整处理策略。

中间件需要维护状态,但这些状态可能会因为各种原因而发生变化。错误处理机制、回滚机制、事务性操作等都需要在状态管理中得到体现。

3. 高级 Hook 使用模式的深度解析

条件 Hook 执行是高级使用模式中最常见的一种。它通过智能的条件判断来决定是否执行特定的 Hook,这样可以提高处理效率,避免不必要的计算。

错误恢复 Hook 的设计体现了系统的韧性。不同层级的错误需要不同的恢复策略。有些错误可以通过简单的重试来解决,有些错误需要切换到备用模型,还有些错误需要人工介入。通过在不同 Hook 层级实现恢复机制,系统可以在不同层次上处理错误,提高整体的可靠性。

性能优化 Hook 则关注系统的运行效率。

缓存 Hook 通过在 wrap_model_call 阶段检查缓存来决定是否直接返回缓存结果,避免重复的模型调用。

预加载 Hook 则通过预测用户需求,提前加载可能需要的数据和资源。


四:中间件集成工具使用

4.1. before_model 模型调用前

中间件类型:before_model - 模型调用前中间件

Pasted image 20260707175902.png

1. SummarizationMiddleware 上下文压缩

核心特性

  1. 官方中间件集成:使用 from langchain.agents.middleware import SummarizationMiddleware
  2. 自动压缩:在 create_agent 中通过 middleware 参数集成
  3. 智能保留:自动压缩历史消息,保留最近的对话
  4. 无需手动管理:中间件自动处理压缩逻辑

工作原理

当历史消息的 token 数量超过阈值(500)且消息数量超过保留数量(5 条)时,中间件会自动:

  1. 将旧消息发送给摘要模型进行压缩
  2. 保留最近的 N 条消息
  3. 将摘要结果作为上下文传递给 Agent

预期结果

# ==================== 配置中间件 ====================
summarization_middleware = SummarizationMiddleware(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.1),
    # max_tokens_before_summary=200,          # 历史消息 token 数量超过 200 时触发压缩
    messages_to_keep=5,                     # 保留最近 5 条消息
    summary_prompt="请将以下对话历史进行摘要,保留关键决策点和技术细节:\n\n{messages}\n\n摘要:"   # 摘要提示词
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.2),
    tools=tools,
    middleware=[summarization_middleware],
    context_schema=UserContext,
    debug=True,
)

2. PIIMiddleware PII 信息脱敏

核心特性

  1. 自动 PII 检测:使用 from langchain.agents.middleware import PIIMiddleware
  2. 智能脱敏:自动识别并处理敏感信息
  3. 多种策略:支持 block、redact、mask、hash 四种处理策略
  4. 无缝集成:在模型调用前自动处理,对业务逻辑透明

工作原理

在模型调用前,中间件会自动:

  1. 扫描消息内容,识别指定类型的 PII 信息
  2. 根据策略处理敏感信息(阻止 / 脱敏 / 遮蔽 / 哈希)
  3. 将处理后的消息传递给模型

支持的 PII 类型

处理策略

预期结果

# ==================== 定义用户上下文 ====================
class UserContext(BaseModel):
    """用户上下文 Schema"""
    user_id: str = Field(..., description="用户唯一标识")
    department: str = Field(..., description="所属部门")
    security_level: str = Field(default="normal", description="安全级别")

# ==================== 配置 PIIMiddleware ====================
# 核心配置:信用卡掩码中间件
piim_credit_card = PIIMiddleware(
    "credit_card",
    detector=r"\b(?:\d{4}[-\s]?){3}\d{4}\b",  # 匹配格式: 1234-5678-9012-3456
    strategy="mask",       # 掩码策略
    apply_to_input=True,   # 对输入消息进行掩码
    apply_to_output=False,  # 不对工具输出进行掩码(工具返回的是业务结果)
)

# ==================== 创建智能体 ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.2),
    tools=tools,
    # 中间件:只启用 PII 掩码(生产环境可添加日志等)
    middleware=[
        piim_credit_card,  # 信用卡掩码中间件
    ],
    # 启用上下文
    context_schema=UserContext,
    debug=True,
)

3. ModelCallLimitMiddleware 模型调用限制

核心特性

  1. 安全防护:防止 Agent 陷入无限循环
  2. 简单配置:通过 max_calls 参数设置最大调用次数
  3. 自动熔断:达到限制后自动停止并返回错误或特定消息

工作原理

中间件会跟踪当前会话中的模型调用次数。当调用次数达到设定的阈值时,中间件会阻止后续的模型调用,并引发异常或返回预设的响应。

预期结果

# ==================== 配置中间件 ====================
limit_middleware = ModelCallLimitMiddleware(
    run_limit=3,  # 每次运行最多调用模型 3 次
    exit_behavior='error'  # 超限时抛出异常
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.1),
    tools=tools,
    middleware=[limit_middleware],
    context_schema=UserContext,
    debug=False,  # 关闭调试模式以减少输出
)

4.2. wrap_model_call (包裹模型调用)

wrap_model_call - 模型调用包装中间件

1. ContextEditingMiddleware 管理上下文大小

核心特性

  1. 自动上下文管理:当 token 数量超过阈值时自动清理旧的工具结果
  2. 灵活配置:支持自定义触发阈值、保留数量、排除工具等
  3. 智能清理:保留最近的 N 个工具结果,清理较旧的内容
  4. 无缝集成:在模型调用前自动处理,对业务逻辑透明

工作原理

当消息历史的 token 数量超过配置的阈值时,中间件会自动:

  1. 统计当前消息的 token 数量
  2. 如果超过阈值,清理旧的工具调用结果
  3. 保留最近的 N 个工具结果
  4. 将清理后的消息传递给模型

ClearToolUsesEdit 配置参数

预期结果

# 关键:设置较低的触发阈值,确保能够触发清理
custom_context_middleware = ContextEditingMiddleware(
    edits=[
        ClearToolUsesEdit(
            trigger=800,  # 当 token 数超过 800 时触发清理(约 3-4 次工具调用后)
            keep=1,  # 只保留最近的 1 个工具结果
            clear_at_least=0,  # 清理所有超出 keep 数量的内容
            clear_tool_inputs=False,  # 不清理工具输入参数
            exclude_tools=["generate_report"],  # 不清理 generate_report 的结果
            placeholder="[已清理以节省空间]",  # 自定义占位符
        )
    ],
    token_count_method="approximate"  # 使用近似计数(更快)
)

# ==================== 创建 Agent(使用 checkpointer 来累积消息)====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.1),
    tools=tools,
    middleware=[
        custom_context_middleware,  # 使用自定义配置
    ],
    context_schema=UserContext,  # 定义上下文参数,这里是 UserContext
    checkpointer=MemorySaver(),  # 关键:使用 checkpointer 来保存消息历史
    debug=True,  # 开启调试模式以观察中间件行为
)

ContextEditingMiddleware 工作原理说明

  1. 使用 checkpointer 在同一线程中累积消息历史
  2. 当消息历史超过 800 tokens 时触发清理
  3. 只保留最近的 1 个工具调用结果
  4. 'generate_report' 工具的结果不会被清理(exclude_tools)
  5. 被清理的内容会被替换为 '[已清理以节省空间]'
  6. 每个工具返回约 250 tokens,3-4 次调用后应触发清理

2. ModelFallbackMiddleware 模型故障自动切换

核心特性

  1. 自动故障转移:主模型失败时自动切换到备用模型
  2. 多级备份:支持配置多个备用模型,按顺序尝试
  3. 无缝切换:对业务逻辑透明,自动处理重试逻辑
  4. 提高可用性:显著提升系统的稳定性和可靠性

工作原理

当模型调用失败时,中间件会自动:

  1. 捕获主模型的异常
  2. 按顺序尝试备用模型
  3. 返回第一个成功的模型响应
  4. 如果所有模型都失败,抛出最后一个异常

配置参数

预期结果

# ==================== 配置中间件 ====================
# 配置模型故障转移:主模型 -> 备用模型 1 -> 备用模型 2
# 注意:这里使用相同的模型作为演示,实际应用中应使用不同的模型
fallback_middleware = ModelFallbackMiddleware(
    ChatDeepSeek(model="deepseek-chat", temperature=0.3),  # 第一个备用模型
    ChatDeepSeek(model="deepseek-reasoner", temperature=0.5),  # 第二个备用模型
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.1),  # 主模型
    tools=tools,
    middleware=[
        fallback_middleware,  # 添加故障转移中间件
    ],
    context_schema=UserContext,  # 定义上下文参数,这里是 UserContext
    debug=False,    # 关闭调试模式,避免在测试中输出详细信息
)

💡 提示:
在生产环境中,建议配置不同提供商的模型,例如:
主模型: openai:gpt-4 o
备用 1: anthropic:claude-sonnet-4-5-20250929
备用 2: deepseek:deepseek-chat
这样可以在某个提供商服务中断时,自动切换到其他提供商。

3. LLMToolSelectorMiddleware 智能工具选择

核心特性

  1. 智能工具筛选:使用 LLM 分析查询并选择最相关的工具
  2. 减少 Token 消耗:只将相关工具传递给主模型,降低成本
  3. 提高准确性:帮助主模型聚焦于正确的工具,提升响应质量
  4. 灵活配置:支持限制工具数量、指定必选工具、自定义选择模型

工作原理

在主模型调用前,中间件会自动:

  1. 使用选择模型分析用户查询
  2. 从所有可用工具中选择最相关的 N 个工具
  3. 将筛选后的工具列表传递给主模型
  4. 主模型只能看到和使用被选中的工具

配置参数

预期结果

# 所有工具列表(模拟拥有大量工具的场景)
all_tools = [
    search_weather,
    search_news,
    calculate_math,
    translate_text,
    search_database,
    send_email,
    get_stock_price,
    book_meeting,
]

# ==================== 配置中间件 ====================
# 配置工具选择中间件:使用 LLM 智能选择最相关的工具
tool_selector_middleware = LLMToolSelectorMiddleware(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.1),  # 使用较小的模型进行工具选择
    max_tools=3,  # 最多选择 3 个工具
    always_include=["calculate_math"],  # 始终包含数学计算工具
    system_prompt="分析用户查询,选择最相关的工具。优先选择直接相关的工具。"
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.2),  # 主模型
    tools=all_tools,  # 提供所有 8 个工具
    middleware=[
        tool_selector_middleware,  # 添加工具选择中间件
    ],
    context_schema=UserContext,
    debug=True,  # 开启调试模式以观察工具选择过程
)

💡 优势:

4.3 wrap_tool_call (包裹工具调用)

wrap_tool_call - 工具调用包装中间件

1. ToolRetryMiddleware 自动重试工具调用

核心特性

  1. 自动重试:工具调用失败时自动重试,无需手动处理
  2. 指数退避:支持指数退避策略,避免过度请求
  3. 灵活配置:可配置重试次数、退避因子、延迟时间等
  4. 异常过滤:支持只重试特定类型的异常
  5. 工具级控制:可以针对特定工具配置重试策略

工作原理

当工具调用失败时,中间件会自动:

  1. 捕获工具调用异常
  2. 检查是否应该重试(基于异常类型和重试次数)
  3. 等待一段时间(使用指数退避策略)
  4. 重新执行工具调用
  5. 返回成功结果或最终失败消息

配置参数

预期结果

# ==================== 配置中间件 ====================
# 配置工具重试中间件:自动重试失败的工具调用
retry_middleware = ToolRetryMiddleware(
    max_retries=3,  # 最多重试 3 次
    tools=["unreliable_api_call", "random_failure_tool"],  # 只对这两个工具启用重试
    retry_on=(ConnectionError, RuntimeError),  # 只重试这些异常
    on_failure="return_message",  # 失败时返回友好消息而不是抛出异常
    backoff_factor=1.5,  # 退避因子,每次重试延迟增加 1.5 倍
    initial_delay=0.5,  # 初始延迟 0.5 秒
    max_delay=5.0,  # 最大延迟 5 秒
    jitter=True,  # 添加随机抖动,避免同时重试
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.2),
    tools=tools,
    middleware=[
        retry_middleware,  # 添加重试中间件
    ],
    context_schema=UserContext,
    debug=False,
)

ToolRetryMiddleware 工作原理说明

  1. unreliable_api_call 工具前 2 次调用失败
  2. 中间件自动捕获 ConnectionError 异常
  3. 使用指数退避策略等待后重试
  4. 第 3 次调用成功,返回结果
  5. stable_tool 工具始终成功,不需要重试
  6. 重试机制对业务逻辑完全透明

💡 重试策略:

🎯 适用场景:

2. LLMToolEmulator 模拟工具执行

核心特性

  1. LLM 模拟执行:使用 LLM 生成模拟的工具执行结果
  2. 选择性模拟:可以选择模拟特定工具或所有工具
  3. 安全测试:在不执行真实操作的情况下测试 Agent 逻辑
  4. 快速原型:无需实现真实工具即可测试 Agent 流程

工作原理

当工具被调用时,中间件会自动:

  1. 拦截工具调用请求
  2. 检查该工具是否在模拟列表中
  3. 使用 LLM 根据工具描述和参数生成模拟结果
  4. 返回模拟结果而不是执行真实工具

配置参数

预期结果

# ==================== 配置中间件 ====================
# 配置工具模拟中间件:使用 LLM 模拟危险操作,避免真实执行
emulator_middleware = LLMToolEmulator(
    tools=["send_real_email", "charge_credit_card", "delete_database_record"],  # 只模拟这些危险工具
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.7),  # 使用 DeepSeek 进行模拟
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.2),
    tools=tools,
    middleware=[
        emulator_middleware,  # 添加工具模拟中间件
    ],
    context_schema=UserContext,
    debug=True,  # 开启调试模式以观察模拟过程
)

LLMToolEmulator 工作原理说明

  1. send_real_email, charge_credit_card, delete_database_record 被 LLM 模拟
  2. 这些工具的代码不会被真实执行
  3. LLM 根据工具描述和参数生成合理的模拟结果
  4. safe_query_tool 不在模拟列表中,会真实执行
  5. 日志中可以看到哪些工具被真实调用(⚠️)或模拟(无标记)

🎯 使用场景:

💡 最佳实践:

4.2 after_model 模型调用后

Pasted image 20260708112556.png

after_model - 模型调用后中间件

1. HumanInTheLoopMiddleware 人工干预中间件

核心特性

  1. 官方中间件集成:使用 from langchain.agents.middleware import HumanInTheLoopMiddleware
  2. 工具调用拦截:在 create_agent 中通过 middleware 参数集成
  3. 灵活审批策略:支持 approve(批准)、edit(编辑)、reject(拒绝)三种决策
  4. 无缝集成:中间件自动处理中断和恢复逻辑

工作原理

当 AI 决定调用需要审批的工具时,中间件会自动:

  1. 拦截工具调用请求
  2. 触发中断(interrupt),等待人工决策
  3. 根据人工决策执行相应操作(批准 / 编辑 / 拒绝)
  4. 继续或终止执行流程

审批决策类型

预期结果

model = ChatDeepSeek(model="deepseek-chat")

# ---------------------------------------------------------------------------
# 创建带 HumanInTheLoopMiddleware 的图
# ---------------------------------------------------------------------------
system_prompt = """
你是一个专业的行政助手。
当用户请求发送邮件时,你必须直接调用 `send_email` 工具。
不要问任何后续问题,不要要求确认,直接生成工具调用。
"""

# 定义中间件:指定 'send_email' 工具需要中断审批
# interrupt_on 字典中的 True 表示允许批准、编辑和拒绝
hitl_middleware = HumanInTheLoopMiddleware(
    interrupt_on={"send_email": True},
    description_prefix="需要人工批准才能发送邮件"
)

# 使用 create_agent 创建图,并注入中间件
# LangGraph Studio 会自动处理持久化,不需要传入 checkpointer
graph = create_agent(
    model=model,
    tools=tools,
    system_prompt=system_prompt,
    middleware=[hitl_middleware]
)

2. ToolCallLimitMiddleware 工具调用限制

核心特性

  1. 资源保护:防止特定工具被频繁调用
  2. 灵活配置:支持全局限制或针对特定工具的限制
  3. 自动熔断:达到限制后阻止工具执行并返回错误

工作原理

中间件会跟踪当前会话中的工具调用次数。当特定工具或总工具调用次数达到设定的阈值时,中间件会阻止后续的工具调用,并引发异常或返回预设的响应。

预期结果

# ==================== 配置中间件 ====================
# 方式1: 限制所有工具的调用次数(全局限制)
global_tool_limiter = ToolCallLimitMiddleware(
    tool_name=None,  # None = 限制所有工具
    run_limit=3,     # 每次运行最多调用 3 次工具
    exit_behavior="continue"  # 超限后阻止工具调用,但继续执行
)

# 方式2: 限制特定工具的调用次数
specific_tool_limiter = ToolCallLimitMiddleware(
    tool_name="check_server_status",  # 只限制 check_server_status 工具
    thread_limit=5,   # 整个线程最多调用 5 次
    run_limit=2,      # 每次运行最多调用 2 次
    exit_behavior="error"  # 超限后返回错误消息
)

# ==================== 创建 Agent ====================
agent = create_agent(
    model=ChatDeepSeek(model="deepseek-chat", temperature=0.1),
    tools=tools,
    middleware=[
        specific_tool_limiter,  # 使用特定工具限制器
    ],
    context_schema=UserContext,
    debug=False,
)

五:自定义中间件应用

5.1. 中间件的参数传递机制

1. ModelRequest / Response

# ModelRequest 结构示例
request = ModelRequest(
    model=ChatOpenAI(...),  # 可动态替换的模型实例
    messages=[...],         # 本次调用的消息列表(可被修改)
    tools=[...],            # 本次可用的工具列表(可被过滤)
    runtime=context         # 包含 AgentState 的引用
)
# 典型结构
request.model      # 当前要调用的模型实例(可修改)
request.tools      # 可用工具列表(可修改)
request.messages   # 消息历史
request.state      # Agent 的当前状态
request.runtime    # 运行时上下文,包含 request.runtime.context 等

关键特性:

2. handler: Callable[[ModelRequest], ModelResponse] —— 执行器函数

handler 是实际执行模型调用的函数,它是 LangChain 内部封装的"下一个"执行环节。你可以理解为:

核心作用:

# handler 本质上等价于:
def handler(request: ModelRequest) -> ModelResponse:
    # 实际调用 LLM 模型并返回响应
    return actual_model_call(request)

3. AgentState:整个 Agent 生命周期的"全局状态容器"

总结:在 @wrap_model_call 和 @wrap_tool_call 等包裹类 hook 中,所有 request 字段都可直接赋值修改(框架特批),但 config 除外;在 @before_model、@after_model 等钩子中,只能读取 state,不能访问 request。

参数类型 本质 所属模块 生命周期
ModelRequest 数据类(dataclass) langchain.agents.middleware.base 单次模型调用
ModelResponse 数据类(dataclass) langchain.agents.middleware.base 单次模型调用
AgentState 字典结构(TypedDict) langchain.agents.agent 整个 Agent 会话
handler 函数对象(可调用) 动态传递 单次包装调用

Pasted image 20260708113714.png

5.2. 装饰器 dynamic_prompt 动态提示词(wrap_model_call)

from langchain.agents.middleware import dynamic_prompt

# 动态提示函数
@dynamic_prompt
def role_based_prompt(request:ModelRequest):
    """根据用户角色生成不同提示词"""
    user_role = request.runtime.context.get("user_role", "user")

    if user_role == "expert":
        return "你是一个专业气象分析师,提供详细数据"
    else:
        return "你是一个简洁的天气助手"

# 创建动态 Agent
agent_dynamic = create_agent(
    model="openai:gpt-4o-mini",
    tools=[get_weather],
    middleware=[role_based_prompt],  # 注入动态提示
    context_schema=Context
)

5.3. 使用装饰器实现模型动态切换(wrap_model_call)

# 定义两个模型实例
small_model = ChatOpenAI(model="gpt-4o-mini")
large_model = ChatOpenAI(model="gpt-4o")

hard_keywords = ("证明", "推导", "严谨", "规划","复杂", "多步骤", "chain of thought", "step-by-step", "reason step by step", "数学", "逻辑证明", "约束求解")

# 定义动态模型切换的中间件
@wrap_model_call
async def dynamic_model_router(
    request: ModelRequest,
    handler: Callable[[ModelRequest], ModelResponse]
) -> ModelResponse:
    """
    根据对话上下文动态切换模型
    """
    # 获取当前对话的状态(例如消息列表)
    state = request.state
    messages = state.get("messages", [])

    # 获取上下文中的用户角色
    print(f"打印运行时上下文: {request.runtime.context}")

    # === 逻辑判断示例 ===
    # 场景 A: 如果对话轮数超过 5 轮,切换到大模型处理复杂上下文
    if len(messages) > 5:
        # 使用 .override() 方法替换本次调用的模型
        request = request.override(model=large_model)

    # 场景 B: 如果用户输入包含特定关键词 (仅作演示,实际可用分类器)
    elif messages and "复杂分析" in messages[-1].content or any(kw.lower() in messages[-1].content for kw in hard_keywords):
         request = request.override(model=large_model)

    else:
        print("--- [Middleware] 使用默认小模型 GPT-4o-mini ---")
        # 默认使用 create_agent 初始化时传入的模型(即 small_model)

    # 继续执行调用
    return await handler(request)

# 创建 Agent 并注入中间件
agent = create_agent(
    model=small_model,  # 默认模型
    tools=[get_weather],
    middleware=[dynamic_model_router],  # <--- 关键:注入动态路由中间件
    context_schema=Context    # 上下文类型
)

5.4. persist_session (会话持久化)

中间件类型

after_model - 会话持久化包装器

核心特性

  1. 自动保存:在 Agent 执行过程中自动保存会话状态
  2. 简单易用:作为中间件直接集成

预期结果

每次对话后,会话状态会被保存到指定目录。

from langchain.agents.middleware import AgentMiddleware, AgentState
import os
import json
import time
from typing import Dict, Any

class PersistSessionMiddleware(AgentMiddleware):
    def __init__(self, path: str):
        super().__init__()
        self.path = path
        if not os.path.exists(path):
            os.makedirs(path)

    def after_model(self, state: AgentState, runtime) -> None:
        """在模型调用后保存会话"""
        # 使用当前时间戳作为简单的会话标识或检查点
        timestamp = int(time.time())
        messages = state.get("messages", [])
        if not messages:
            return

        # 仅保存最后一条 AI 消息作为演示,或者保存整个历史
        filename = os.path.join(self.path, f"state_{timestamp}.json")

        # 简化的序列化
        serialized_msgs = []
        for msg in messages:
            msg_data = {
                "role": msg.type,
                "content": msg.content
            }
            if hasattr(msg, 'tool_calls') and msg.tool_calls:
                 msg_data['tool_calls'] = msg.tool_calls
            serialized_msgs.append(msg_data)

        try:
            with open(filename, "w", encoding="utf-8") as f:
                json.dump(serialized_msgs, f, ensure_ascii=False, indent=2)
            print(f" [persist_session] 会话状态已保存: {filename}")
        except Exception as e:
            print(f" [persist_session] 保存失败: {e}")

def persist_session(path: str = "./sessions"):
    """
    persist_session 轻量包装器
    Args:
        path: 会话持久化保存的目录路径
    """
    return PersistSessionMiddleware(path)

from langchain.agents import create_agent
from langchain_openai import ChatOpenAI

# 创建 Agent
# 这里我们使用 persist_session 包装器
agent = create_agent(
    model=ChatOpenAI(model="gpt-4o-mini"),
    tools=[],
    middleware=[persist_session(path="./my_agent_sessions")]
)

5.5. 多中间件组合应用 - IT 运维 Agent

"""
IT 运维 Agent 中间件演示 - 具备 RBAC 和审计功能

本示例展示了一个完整的 IT 运维 Agent,集成了:
- RBAC(基于角色的权限控制)
- 操作审计和日志记录
- 安全检查和验证
- 上下文管理(使用 ContextEditingMiddleware)
- 动态系统提示词注入(使用 @wrap_model_call)
- 智能模型切换(使用 @wrap_model_call)

【中间件执行顺序】
1. before_agent: SecurityGuardrail (安全检查)
2. before_agent: RBACMiddleware (权限验证)
3. before_model: RBACMiddleware (注入用户信息到 state)
4. wrap_model_call: dynamic_system_prompt (动态提示词注入)
5. wrap_model_call: dynamic_model_router (智能模型切换)
6. wrap_model_call: ContextEditingMiddleware (上下文管理和清理)
7. Model Execution (模型执行)
注意:after_model 钩子按注册顺序相反执行(先注册后执行)
8. after_model: AuditLogger (审计日志)
9. after_model: ResponseValidator (响应验证)

【支持的运维操作】
- 服务器状态查询
- 服务重启
- 日志查看
- 系统资源监控
"""

from typing import Any, Dict, Optional, List, Callable,TypedDict

from datetime import datetime
from enum import Enum
from dotenv import load_dotenv
import json

from langchain_deepseek import ChatDeepSeek
from langchain.agents import create_agent
from langchain.agents.middleware import (
    AgentMiddleware,
    AgentState,
    hook_config,
    ContextEditingMiddleware,
    ModelRequest,
    ModelResponse,
    wrap_model_call
)
from langchain.agents.middleware.context_editing import ClearToolUsesEdit
from langgraph.checkpoint.memory import MemorySaver
from langchain_core.tools import tool
from langchain_core.messages import AIMessage, SystemMessage
from langchain_core.language_models import BaseChatModel
from pydantic import BaseModel, Field

load_dotenv(override=True)

# 用户角色和权限定义
class UserRole(str, Enum):
    """用户角色枚举"""
    ADMIN = "admin"           # 管理员:所有权限
    OPERATOR = "operator"     # 运维人员:查询和重启权限
    VIEWER = "viewer"         # 查看者:仅查询权限
    GUEST = "guest"           # 访客:无权限

class Permission(str, Enum):
    """权限枚举"""
    VIEW_STATUS = "view_status"           # 查看状态
    VIEW_LOGS = "view_logs"               # 查看日志
    RESTART_SERVICE = "restart_service"   # 重启服务
    MODIFY_CONFIG = "modify_config"       # 修改配置
    VIEW_METRICS = "view_metrics"         # 查看监控指标

# 角色权限映射
ROLE_PERMISSIONS: Dict[UserRole, List[Permission]] = {
    UserRole.ADMIN: [
        Permission.VIEW_STATUS,
        Permission.VIEW_LOGS,
        Permission.RESTART_SERVICE,
        Permission.MODIFY_CONFIG,
        Permission.VIEW_METRICS
    ],
    UserRole.OPERATOR: [
        Permission.VIEW_STATUS,
        Permission.VIEW_LOGS,
        Permission.RESTART_SERVICE,
        Permission.VIEW_METRICS
    ],
    UserRole.VIEWER: [
        Permission.VIEW_STATUS,
        Permission.VIEW_LOGS,
        Permission.VIEW_METRICS
    ],
    UserRole.GUEST: []
}

# 工具权限映射
TOOL_PERMISSIONS: Dict[str, Permission] = {
    "check_server_status": Permission.VIEW_STATUS,
    "view_service_logs": Permission.VIEW_LOGS,
    "restart_service": Permission.RESTART_SERVICE,
    "get_system_metrics": Permission.VIEW_METRICS
}
# 定义中间件
class RBACMiddleware(AgentMiddleware):
    """
    [阶段 1: before_agent & before_model] RBAC 权限控制中间件
    在执行任何操作前验证用户权限
    """
    def __init__(self):
        super().__init__()

    def _get_user_from_runtime(self, runtime) -> Dict[str, Any]:
        """从 runtime.context 获取用户信息"""
        try:
            # 尝试从 runtime.context 获取用户信息
            if hasattr(runtime, 'context') and runtime.context:
                context_data = runtime.context
                # 如果是字典,直接使用
                if isinstance(context_data, dict):
                    user_role_str = context_data.get('role', 'guest')
                    # 转换角色字符串为枚举
                    try:
                        user_role = UserRole(user_role_str)
                    except ValueError:
                        user_role = UserRole.GUEST

                    return {
                        'user_id': context_data.get('user_id', 'unknown'),
                        'username': context_data.get('username', 'unknown'),
                        'role': user_role,
                        'department': context_data.get('department', 'unknown')
                    }
        except Exception as e:
            log_with_timestamp(f"   ⚠️ 从 runtime.context 获取用户信息失败: {str(e)}", "WARN")

        # 如果无法从 runtime 获取,使用默认用户
        return get_current_user()

    @hook_config(can_jump_to=["end"])  # 允许在 before_agent 阶段跳转到 end
    def before_agent(self, state: AgentState, runtime) -> Optional[Dict[str, Any]]:
        try:
            # 从 runtime 获取用户信息
            current_user = self._get_user_from_runtime(runtime)

            # 获取用户角色的权限列表
            user_permissions = ROLE_PERMISSIONS.get(current_user['role'], [])

            return None

        except Exception as e:
            log_with_timestamp(f"   ❌ 权限验证异常: {str(e)}", "ERROR")
            return {
                "messages": [AIMessage(
                    content="权限验证失败,请联系管理员。"
                )],
                "jump_to": "end"
            }

    def before_model(self, state: AgentState, runtime) -> Optional[Dict[str, Any]]:
        """在 before_model 阶段注入用户信息到 state"""
        try:
            # 从 runtime 获取用户信息
            current_user = self._get_user_from_runtime(runtime)

            # 获取用户角色的权限列表
            user_permissions = ROLE_PERMISSIONS.get(current_user['role'], [])

            # 将用户信息注入到 state
            return {
                "user_info": current_user,
                "user_permissions": [p.value for p in user_permissions]
            }
        except Exception as e:
            log_with_timestamp(f"   ❌ 用户信息注入异常: {str(e)}", "ERROR")
            return None
class SecurityGuardrail(AgentMiddleware):
    """
    [阶段 2: before_agent] 安全护栏
    检查请求中的危险操作和敏感关键词
    """
    # 危险操作关键词
    DANGEROUS_KEYWORDS = ["删除数据库", "drop database", "rm -rf", "format","删除所有", "清空", "hack", "攻击", "入侵"]

    def __init__(self):
        super().__init__()

    @hook_config(can_jump_to=["end"])
    def before_agent(self, state: AgentState, runtime) -> Optional[Dict[str, Any]]:
        try:
            messages = state.get("messages", [])
            if not messages:
                return None

            last_msg = messages[-1]
            if last_msg.type == "human":
                content_lower = last_msg.content.lower()

                # 检查危险关键词
                for keyword in self.DANGEROUS_KEYWORDS:
                    if keyword in content_lower:
                        return {
                            "messages": [AIMessage(
                                content=f"⚠️ 安全警告:检测到危险操作关键词 '{keyword}',"
                                       f"该操作已被拦截。如需执行此类操作,请联系管理员。"
                            )],
                            "jump_to": "end"
                        }

            return None

        except Exception as e:
            log_with_timestamp(f"   ❌ 安全检查异常: {str(e)}", "ERROR")
            return {
                "messages": [AIMessage(
                    content="安全检查失败,操作已被拦截。"
                )],
                "jump_to": "end"
            }
# 动态系统提示词中间件(使用 @wrap_model_call)
def create_dynamic_system_prompt_middleware():
    """
    创建动态系统提示词中间件
    根据用户角色动态注入不同的系统提示词

    使用 @wrap_model_call 装饰器,从 ModelRequest 获取用户信息
    """

    # 角色专属提示词
    ROLE_PROMPTS = {
        UserRole.ADMIN: """
    【管理员模式】
    你拥有完整的系统权限,可以执行所有操作。
    - 可以查看所有服务器状态和日志
    - 可以重启服务和修改配置
    - 需要特别谨慎,确认每个关键操作
    - 提供详细的技术分析和建议
    """,
        UserRole.OPERATOR: """
    【运维人员模式】
    你是运维团队成员,拥有日常运维权限。
    - 可以查看服务器状态和日志
    - 可以重启服务(需确认)
    - 不能修改系统配置
    - 提供实用的运维建议
    """,
        UserRole.VIEWER: """
    【查看者模式】
    你只有查看权限,不能执行任何操作。
    - 可以查看服务器状态和日志
    - 可以查看系统监控指标
    - 不能执行任何修改操作
    - 提供信息查询和数据分析
    """,
        UserRole.GUEST: """
    【访客模式】
    你的权限受到严格限制。
    - 只能进行基本的信息查询
    - 不能访问敏感数据
    - 不能执行任何操作
    - 提供有限的帮助信息
    """
    }

    @wrap_model_call
    def dynamic_system_prompt(
        request: ModelRequest,
        handler: Callable[[ModelRequest], ModelResponse]
    ) -> ModelResponse:
        """
        根据用户角色动态注入系统提示词
        从 ModelRequest 中获取:
        1. request.runtime.context - 包含原始用户上下文
        """
        try:
            user_role = UserRole.GUEST  # 默认角色

            # 尝试从 runtime.context 获取
            if hasattr(request, 'runtime') and hasattr(request.runtime, 'context'):
                context_data = request.runtime.context
                if isinstance(context_data, dict):
                    user_role_str = context_data.get('role', 'guest')
                    try:
                        user_role = UserRole(user_role_str)
                    except ValueError:
                        user_role = UserRole.GUEST
                    log_with_timestamp(f"   ℹ️ 从 runtime.context 获取用户角色: {user_role.value}")

            # 确保 user_role 是 UserRole 枚举类型
            if isinstance(user_role, str):
                try:
                    user_role = UserRole(user_role)
                except ValueError:
                    log_with_timestamp(f"   ⚠️ 无效的角色字符串: {user_role},使用 GUEST", "WARN")
                    user_role = UserRole.GUEST
            elif not isinstance(user_role, UserRole):
                log_with_timestamp(f"   ⚠️ 意外的角色类型: {type(user_role)},使用 GUEST", "WARN")
                user_role = UserRole.GUEST

            log_with_timestamp(f"   🔍 检测到用户角色: {user_role.value}")

            # 获取角色专属提示词
            role_prompt = ROLE_PROMPTS.get(user_role, ROLE_PROMPTS[UserRole.GUEST])

            # 获取当前消息列表
            messages = list(request.messages)

            # 检查是否已有系统消息
            has_system_msg = any(msg.type == "system" for msg in messages)

            if not has_system_msg:
                # 创建增强的系统消息
                enhanced_system_msg = SystemMessage(content=role_prompt)
                messages.insert(0, enhanced_system_msg)

                # 使用修改后的消息列表覆盖请求
                request = request.override(messages=messages)
            else:
                log_with_timestamp("   ℹ️ 系统消息已存在,跳过注入")

            # 继续执行调用
            return handler(request)

        except Exception as e:
            log_with_timestamp(f"   ❌ 提示词注入异常: {str(e)}", "ERROR")
            import traceback
            log_with_timestamp(f"   详细错误: {traceback.format_exc()}", "ERROR")
            # 发生异常时,继续执行原始请求
            return handler(request)

    @wrap_model_call
    async def dynamic_system_prompt_async(
        request: ModelRequest,
        handler: Callable[[ModelRequest], ModelResponse]
    ) -> ModelResponse:
        """异步版本:与同步版本逻辑相同"""
        # 直接调用同步版本的逻辑
        return dynamic_system_prompt.invoke(request, handler)

    return dynamic_system_prompt
# 智能模型切换中间件(使用 @wrap_model_call)
# 定义复杂任务关键词
COMPLEXITY_KEYWORDS = ["分析", "建议", "优化", "故障排查", "诊断", "复杂", "详细", "深入", "为什么", "如何","证明", "推导", "严谨", "规划", "多步骤"]

def create_dynamic_model_router(fast_model: BaseChatModel, smart_model: BaseChatModel):
    """
    创建动态模型路由中间件
    Args:
        fast_model: 快速模型(用于简单查询)
        smart_model: 智能模型(用于复杂任务)
    Returns:
        使用 @wrap_model_call 装饰的中间件函数
    """

    def _analyze_complexity(messages: List) -> tuple[bool, str]:
        """分析请求复杂度"""
        should_use_smart_model = False
        reason = ""

        # 场景 1: 长对话(超过 5 轮)
        if len(messages) > 5:
            should_use_smart_model = True
            reason = f"长对话 ({len(messages)} 条消息)"

        # 场景 2: 检查最后一条用户消息
        elif messages:
            last_human_msg = None
            for msg in reversed(messages):
                if msg.type == "human":
                    last_human_msg = msg
                    break

            if last_human_msg:
                content = last_human_msg.content
                content_lower = content.lower()

                # 检查复杂任务关键词
                for keyword in COMPLEXITY_KEYWORDS:
                    if keyword in content_lower:
                        should_use_smart_model = True
                        reason = f"包含复杂关键词 '{keyword}'"
                        break

                # 检查消息长度
                if not should_use_smart_model and len(content) > 100:
                    should_use_smart_model = True
                    reason = f"长消息 ({len(content)} 字符)"

        return should_use_smart_model, reason

    @wrap_model_call
    def dynamic_model_router(
        request: ModelRequest,
        handler: Callable[[ModelRequest], ModelResponse]
    ) -> ModelResponse:
        """
        同步版本:根据对话上下文和请求复杂度动态切换模型
        切换逻辑:
        1. 如果对话轮数超过 5 轮 → 使用智能模型(处理复杂上下文)
        2. 如果包含复杂任务关键词 → 使用智能模型
        3. 如果消息长度超过 100 字符 → 使用智能模型
        4. 其他情况 → 使用快速模型
        """
        # 获取当前对话状态
        state = request.state
        messages = state.get("messages", [])

        should_use_smart_model, reason = _analyze_complexity(messages)

        # 根据判断结果切换模型
        if should_use_smart_model:
            request = request.override(model=smart_model)
        else:
            log_with_timestamp(f"   ⚡ 使用快速模型 - 简单请求")

        # 继续执行调用
        return handler(request)

    @wrap_model_call
    async def dynamic_model_router_async(
        request: ModelRequest,
        handler: Callable[[ModelRequest], ModelResponse]
    ) -> ModelResponse:
        """
        异步版本:根据对话上下文和请求复杂度动态切换模型
        """
        # 获取当前对话状态
        state = request.state
        messages = state.get("messages", [])

        should_use_smart_model, reason = _analyze_complexity(messages)

        # 根据判断结果切换模型
        if should_use_smart_model:
            request = request.override(model=smart_model)
        else:
            log_with_timestamp(f"   ⚡ 使用快速模型 - 简单请求")

        # 继续执行调用
        return await handler(request)

    return dynamic_model_router
class ResponseValidator(AgentMiddleware):
    """
    [阶段 6: after_model] 响应验证器
    验证模型响应的格式和内容
    """
    def __init__(self):
        super().__init__()

    def after_model(self, state: AgentState, runtime) -> Optional[Dict[str, Any]]:
        try:
            messages = state.get("messages", [])
            if not messages:
                return None

            last_msg = messages[-1]
            if last_msg.type == "ai":
                has_tool_calls = hasattr(last_msg, 'tool_calls') and last_msg.tool_calls

                if has_tool_calls:
                    # 验证工具调用权限
                    user_permissions = state.get('user_permissions', [])
                    for tool_call in last_msg.tool_calls:
                        tool_name = tool_call.get('name', '')
                        required_permission = TOOL_PERMISSIONS.get(tool_name)

                        if required_permission and required_permission.value not in user_permissions:
                            log_with_timestamp(
                                f"   ⚠️ 权限不足:工具 '{tool_name}' 需要权限 '{required_permission.value}'",
                                "WARN"
                            )

                elif last_msg.content:
                    log_with_timestamp(f"   ✅ 响应有效 (长度: {len(last_msg.content)} 字符)")

            return None

        except Exception as e:
            log_with_timestamp(f"   ❌ 响应验证异常: {str(e)}", "ERROR")
            return None
class AuditLogger(AgentMiddleware):
    """
    [阶段 7: after_model] 审计日志中间件
    记录所有操作的审计日志
    """
    def __init__(self):
        super().__init__()
        self.audit_records = []

    def after_model(self, state: AgentState, runtime) -> Optional[Dict[str, Any]]:
        try:
            messages = state.get("messages", [])
            if not messages:
                return None

            last_msg = messages[-1]
            user_info = state.get("user_info", {})

            # 记录工具调用
            if last_msg.type == "ai" and hasattr(last_msg, 'tool_calls') and last_msg.tool_calls:
                for tool_call in last_msg.tool_calls:
                    audit_record = {
                        "timestamp": datetime.now().isoformat(),
                        "user_id": user_info.get("user_id", "unknown"),
                        "username": user_info.get("username", "unknown"),
                        "role": user_info.get("role", "unknown"),
                        "action": "tool_call",
                        "tool_name": tool_call.get('name', ''),
                        "tool_args": tool_call.get('args', {}),
                        "status": "initiated"
                    }
                    self.audit_records.append(audit_record)

            # 记录最终响应
            elif last_msg.type == "ai" and last_msg.content:
                audit_record = {
                    "timestamp": datetime.now().isoformat(),
                    "user_id": user_info.get("user_id", "unknown"),
                    "username": user_info.get("username", "unknown"),
                    "role": user_info.get("role", "unknown"),
                    "action": "response",
                    "response_length": len(last_msg.content),
                    "status": "completed"
                }
                self.audit_records.append(audit_record)

            return None

        except Exception as e:
            log_with_timestamp(f"   ❌ 审计日志异常: {str(e)}", "ERROR")
            return None
# 定义 IT 运维工具
class ServerStatusSchema(BaseModel):
    server_name: str = Field(description="服务器名称,例如: web-server-01")

@tool(args_schema=ServerStatusSchema)
def check_server_status(server_name: str) -> str:
    """
    查询服务器状态
    需要权限: VIEW_STATUS
    """
    # 模拟查询服务器状态
    status_data = {
        "server_name": server_name,
        "status": "running",
        "cpu_usage": "45%",
        "memory_usage": "62%",
        "uptime": "15 days 3 hours",
        "last_check": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
    }

    return json.dumps(status_data, ensure_ascii=False, indent=2)

class ServiceLogsSchema(BaseModel):
    service_name: str = Field(description="服务名称,例如: nginx, mysql")
    lines: int = Field(default=50, description="显示的日志行数")

@tool(args_schema=ServiceLogsSchema)
def view_service_logs(service_name: str, lines: int = 50) -> str:
    """
    查看服务日志
    需要权限: VIEW_LOGS
    """
    result = f"=== {service_name} 服务日志 (最近 {lines} 行) ===\n"
    result += "\n".join(log_entries[:lines])

    return result

class RestartServiceSchema(BaseModel):
    service_name: str = Field(description="要重启的服务名称")
    force: bool = Field(default=False, description="是否强制重启")

@tool(args_schema=RestartServiceSchema)
def restart_service(service_name: str, force: bool = False) -> str:
    """
    重启服务
    需要权限: RESTART_SERVICE
    """
    # 模拟重启服务
    restart_type = "强制重启" if force else "正常重启"

    result = f"""
    🔄 服务重启操作
    服务名称: {service_name}
    重启类型: {restart_type}
    操作时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}
    状态: 成功
    预计恢复时间: 30秒
    """

    return result.strip()

class SystemMetricsSchema(BaseModel):
    metric_type: str = Field(
        description="指标类型: cpu, memory, disk, network"
    )

@tool(args_schema=SystemMetricsSchema)
def get_system_metrics(metric_type: str) -> str:
    """
    获取系统监控指标
    需要权限: VIEW_METRICS
    """
    # 模拟返回系统指标
    metrics = {
        "cpu": {
            "usage": "45.2%",
            "cores": 8,
            "load_average": [2.1, 2.3, 2.5]
        },
        "memory": {
            "total": "16 GB",
            "used": "10 GB",
            "free": "6 GB",
            "usage": "62.5%"
        },
        "disk": {
            "total": "500 GB",
            "used": "320 GB",
            "free": "180 GB",
            "usage": "64%"
        },
        "network": {
            "rx_bytes": "1.2 TB",
            "tx_bytes": "890 GB",
            "connections": 156
        }
    }

    metric_data = metrics.get(metric_type.lower(), {})
    return json.dumps(metric_data, ensure_ascii=False, indent=2)

# 工具列表
tools = [
    check_server_status,
    view_service_logs,
    restart_service,
    get_system_metrics
]
# 创建 IT 运维 Agent
# 创建模型
model = ChatDeepSeek(model="deepseek-chat", temperature=0.1)

# 创建快速模型和智能模型(用于动态切换)
fast_model = ChatDeepSeek(model="deepseek-chat", temperature=0, max_tokens=500)
smart_model = ChatDeepSeek(model="deepseek-chat", temperature=0.3, max_tokens=2000)

# ==================== 配置 ContextEditingMiddleware ====================
# 关键:设置较低的触发阈值,确保能够触发清理
custom_context_middleware = ContextEditingMiddleware(
    edits=[
        ClearToolUsesEdit(
            trigger=800,  # 当 token 数超过 800 时触发清理(约 3-4 次工具调用后)
            keep=1,  # 只保留最近的 1 个工具结果
            clear_at_least=0,  # 清理所有超出 keep 数量的内容
            clear_tool_inputs=False,  # 不清理工具输入参数
            exclude_tools=["restart_service"],  # 不清理 restart_service 的结果(重要操作)
            placeholder="[已清理以节省空间]",  # 自定义占位符
        )
    ],
    token_count_method="approximate"  # 使用近似计数(更快)
)

dynamic_model_router = create_dynamic_model_router(fast_model, smart_model)

dynamic_system_prompt = create_dynamic_system_prompt_middleware()

# 按顺序注册中间件
# 注意:
# - before_agent/before_model 
# - @wrap_model_call 中间件按注册顺序执行(先注册先执行)
# - after_model 钩子按注册顺序相反执行(先注册后执行)
middlewares = [
    SecurityGuardrail(),                                 # 1. before_agent: 安全检查
    RBACMiddleware(),                                    # 2. before_agent & before_model: 权限验证(注入 user_info)
    dynamic_system_prompt,                               # 3. @wrap_model_call: 动态提示词注入
    dynamic_model_router,                                # 4. @wrap_model_call: 智能模型切换
    custom_context_middleware,                           # 5. @wrap_model_call: 上下文管理和清理
    ResponseValidator(),                                 # 6. after_model: 响应验证
    AuditLogger(),                                       # 7. after_model: 审计日志
]

# 系统提示词
system_prompt = """你是一个专业的 IT 运维助手,负责帮助运维人员管理服务器和服务。

你的职责包括:
1. 查询服务器状态和系统指标
2. 查看服务日志
3. 在获得授权后重启服务
4. 提供运维建议和故障排查

注意事项:
- 始终确认用户的操作意图
- 对于重启等关键操作,需要明确确认
- 提供清晰、专业的响应
- 遵守权限控制规则
"""

# 导出 graph 变量供 LangGraph Studio 使用
graph = create_agent(
    model=model,
    tools=tools,
    system_prompt=system_prompt,
    middleware=middlewares,
    context_schema=UserContext,  # 添加上下文 schema
    # checkpointer=MemorySaver(),  # 关键:使用 checkpointer 来保存消息历史
    debug=True  # 开启调试模式以观察中间件行为
)

📊 中间件功能总结:

  1. ✅ SecurityGuardrail - 危险操作拦截
  2. ✅ RBACMiddleware - 基于角色的权限控制
  3. ✅ dynamic_system_prompt (@wrap_model_call) - 动态系统提示词注入
  4. ✅ dynamic_model_router (@wrap_model_call) - 智能模型切换
  5. ✅ ContextEditingMiddleware - 上下文管理和清理
  6. ✅ ResponseValidator - 响应验证
  7. ✅ AuditLogger - 审计日志记录

六:中间件编排黄金法则

6.1. 执行顺序设计的哲学思考

middlewares = [
    SecurityGuardrail(),                                 # 1. before_agent: 安全检查
    RBACMiddleware(),                                    # 2. before_agent & before_model: 权限验证(注入 user_info)
    dynamic_system_prompt,                               # 3. @wrap_model_call: 动态提示词注入
    dynamic_model_router,                                # 4. @wrap_model_call: 智能模型切换
    custom_context_middleware,                           # 5. @wrap_model_call: 上下文管理和清理
    ResponseValidator(),                                 # 6. after_model: 响应验证
    AuditLogger(),                                       # 7. after_model: 审计日志
    ToolMiddleware,                                      # 8. wrap_tool_call: 工具调用拦截和控制
]
钩子类型 执行顺序 控制粒度 最佳实践
before_agent 正序 调用级 安全/初始化优先
before_model 正序 单次调用 预算/合规前置
wrap_model_call 嵌套 调用全过程 路由/重试/缓存核心
after_model 逆序 单次调用 验证/日志后置
wrap_tool_call 嵌套 工具执行 参数校验/重试保护

洋葱模型的应用体现了"分层治理"的思想:外层负责安全防护,中层处理业务逻辑,内层保证核心功能,中心提供基础服务。

这种分层的价值在于每一层都有其特定的职责和优先级。安全检查必须在最外层进行,因为任何安全问题都不应该进入到业务处理层面。业务逻辑应该在中间层处理,这样可以确保核心功能的纯粹性。监控和日志记录应该在最内层进行,这样可以不干扰主要的业务逻辑。

提前失败原则是执行顺序设计中的重要原则。通过在最早的可执行点检测问题,我们可以避免后续的无谓计算。这不仅提高了效率,更重要的是提高了系统的可靠性。当系统发现输入有问题时,应该立即拒绝处理,而不是继续进行复杂的计算后发现问题。

缓存优先原则体现了效率优化的思想。对于相同的请求,系统应该尽量返回缓存的结果,而不是重新计算。这不仅减少了计算资源的消耗,更重要的是提高了响应速度。

6.2. 安全优先原则的深度实施

编排黄金法则

  1. 安全类永远优先
    • 安全 > 业务 > 成本 > 监控
  2. 写操作按依赖排序
    • 如果 B 依赖 A 的写入结果,则 A 必须在 B 之前
  3. 中断操作置后
    • 中断类中间件应放在同类型最后
  4. 性能开销降序排列
    • 轻量级 → 重量级(快速失败)

"先安全,后业务;先读再写,轻量在前;同类型按声明,不同类型按层级"

分层安全检查是实施安全优先原则的有效方法。身份认证是最外层的安全检查,它确保只有合法用户才能访问系统。输入验证是第二层安全检查,它确保用户输入的内容不包含恶意代码或不当内容。业务授权是第三层安全检查,它确保用户只能执行其权限范围内的操作。合规检查是最内层的安全检查,它确保整个操作符合相关法规要求。

6.3. 依赖关系管理的复杂性

依赖类型分析帮助我们理解不同类型依赖的特点和处理方式。顺序依赖是最基本的依赖类型,它要求特定的执行顺序。条件依赖则更加复杂,它要求根据中间结果来决定后续的执行路径。数据依赖确保下游中间件能够获得上游中间件的输出作为输入。

依赖排序策略则提供了处理复杂依赖关系的指导原则。无依赖的中间件可以首先执行,单向依赖按照依赖关系排序,条件依赖需要特殊的处理逻辑,聚合依赖则需要等待所有相关的中间件完成。

6.4. 冲突解决机制的设计

资源冲突是最常见的冲突类型。当多个中间件需要同时访问某个资源时,系统需要确保资源的正确分配和使用。锁机制可以保证资源的独占访问,但可能会降低系统的并发性能。队列管理可以提供更公平的资源分配,但可能会增加系统的延迟。

结果冲突则涉及多个中间件产生不同结果的情况。权重决策机制可以根据预定义的规则来选择最终结果。投票机制则通过多个中间件的投票来决定最终结果。在某些情况下,可能需要指定最终决策者来处理冲突。


七:实际应用场景与最佳实践

7.1. 客服机器人中间件方案的深度实践

在典型的客服机器人架构中,用户识别中间件扮演着"门卫"的角色。当客户首次接触系统时,这个中间件需要快速识别客户的重要程度、历史问题、个人偏好等信息。VIP 客户可能需要转接到专业客服或使用更高级的服务,而普通客户则可以通过标准的自动化流程得到服务。这种差异化的处理不仅提高了服务效率,也提升了客户体验。

意图理解中间件是整个系统的"大脑"。它需要从客户的自然语言表达中提取真实的意图和情感。这不仅仅是简单的关键词匹配,而是需要理解语言的深层含义。例如,当客户说"这个产品太贵了"时,系统需要理解这可能是价格敏感的信号,或者是在寻求折扣信息,或者是在比较不同产品。

知识检索中间件就像是一个智能图书管理员。它需要从庞大的知识库中快速找到最相关的信息。知识库可能包含产品信息、政策条款、使用指南、故障排除等多个方面的内容。中间件需要根据客户的意图和问题的上下文,从知识库中选择最合适的信息进行回答。

7.2. 企业知识助手中间件方案

文档解析中间件负责处理各种格式的企业文档,包括 Word、PDF、Excel 等。这个中间件需要能够提取文档的结构化信息和非结构化内容。对于技术文档,它需要提取技术细节和实施步骤;对于政策文档,它需要提取关键条款和适用范围;对于流程文档,它需要提取操作步骤和责任分工。

知识抽取中间件则负责从解析后的文档中提取有价值的知识点。这个过程涉及到自然语言处理、知识图谱构建、语义分析等高级技术。系统需要理解文档的内容,识别重要的概念、关系、事实等,并将其组织成结构化的知识表示。

语义索引中间件建立企业知识的语义化检索索引。这不仅仅是为了提高检索的准确性,更是为了支持复杂的知识推理和发现。系统需要理解知识之间的关联关系,支持基于语义相似度的检索,以及基于知识图谱的推理查询。

7.3. 代码助手中间件方案

代码理解中间件需要具备强大的代码分析能力。它需要理解代码的结构、逻辑、意图和功能。这涉及到语法分析、语义分析、控制流分析、数据流分析等多个层面。系统不仅要理解代码的表面结构,还要理解代码的深层逻辑和设计意图。

语义分析中间件则专注于理解代码的语义含义。它需要识别函数的用途、变量的含义、类的关系、模块的依赖等。这种分析不仅要准确,还要深入。系统需要能够理解设计模式、算法原理、最佳实践等软件工程概念。

最佳实践中间件基于软件工程的最佳实践,为开发者提供改进建议。这包括代码重构建议、性能优化建议、安全性改进建议、测试覆盖率改进建议等。这些建议不是简单的规则匹配,而是基于深度分析和专业知识的智能推荐。

7.4. 最佳实践总结与思考

模块化设计是基础。每个中间件都应该有明确的职责边界和功能范围,避免功能重叠和相互干扰。同时,中间件之间应该有清晰的接口定义,保证系统的可维护性。

性能优先是要求。中间件不应该成为系统的性能瓶颈,需要通过各种优化技术来提高处理效率和减少资源消耗。

安全可靠是底线。特别是对于生产环境中的中间件,必须具备完善的错误处理、安全防护和监控机制。

易于扩展是目标。系统设计应该支持新中间件的动态添加和现有中间件的功能升级,保持系统的长期生命力。