1. 项目概述:为什么我们需要在Agent执行流中“插一脚”?
最近和几个做AI应用落地的朋友聊天,大家不约而同地提到了同一个痛点:Agent(智能体)好用,但不好管。无论是基于LangChain、AutoGen还是其他框架构建的Agent,一旦把复杂的业务逻辑交给它去自主执行,就像把车钥匙交给了一个驾驶技术一流但对本地路况一无所知的新司机。它能从A点开到B点,但中途会不会超速?会不会误入单行道?遇到突发路况(比如业务规则临时变更)能不能灵活处理?我们心里完全没底。
这就是“中间件系统”要解决的核心问题。它不是一个独立运行的服务,而是一套精巧的“拦截器”或“插件”机制,允许我们在Agent既定的执行流水线中,像设置检查站和加工站一样,插入我们自定义的逻辑。想象一下,你给那位新司机配了一个实时在线的副驾,这个副驾不干涉方向盘,但能在司机做出每个动作(起步、转弯、加速)前后,进行安全提醒、路况播报,甚至根据你的指令临时修改目的地。这个“副驾系统”,就是Agent执行流中的中间件。
它的价值远不止于“监控”。对于企业级应用而言,它关乎可控性、可观测性和可扩展性。你可以用它来:
- 注入上下文:在执行前,自动从数据库或缓存中拉取最新的用户会话历史、产品信息,丰富Agent的“记忆”。
- 实施合规审查:在Agent调用工具(如发送邮件、查询数据库)前,校验操作是否符合安全策略,比如检查邮件内容是否包含敏感词,数据库查询是否越权。
- 实现业务熔断:当检测到Agent连续多次尝试失败或陷入循环时,主动中断流程,转由人工或备用流程处理,避免资源浪费和不良体验。
- 增强可观测性:以非侵入的方式,收集每个步骤的输入输出、耗时、Token消耗,为性能分析和成本核算提供精准数据。
- 进行后处理:对Agent生成的结果进行格式化、翻译、敏感信息脱敏,再返回给用户。
简单说,中间件系统让你在不重写Agent核心大脑的前提下,为它的每一次“思考”和“行动”戴上符合你业务需求的“手套”与“眼镜”。下面,我就结合具体的实现思路和踩过的坑,来拆解如何构建这样一个系统。
2. 核心设计思路:洋葱模型与责任链模式
在设计中间件系统时,最忌讳的就是写成一堆散乱、耦合的if-else判断硬编码在Agent的主循环里。那样会导致代码迅速腐化,难以维护和扩展。业界对此有非常成熟的设计模式可以借鉴,其中最贴切的就是洋葱模型(Onion Model)和责任链模式(Chain of Responsibility)。
2.1 理解执行流的“洋葱”结构
把Agent的一次完整调用(例如,处理用户问题“帮我查一下上季度的销售额”)想象成一个洋葱。最核心的“芯”是Agent真正的执行逻辑:理解问题、规划步骤、调用工具、合成答案。而中间件,就是包裹在这个芯外面的一层层“皮”。
每一层皮(中间件)都有机会在请求向内传递到下一层/核心前,和响应向外返回到上一层/客户端前,执行自己的逻辑。这个顺序至关重要:
- 请求阶段(Inbound):从最外层中间件开始,依次向内执行。适合做参数校验、上下文注入、身份认证。
- 核心处理:Agent执行核心任务。
- 响应阶段(Outbound):从最内层中间件开始,依次向外执行。适合做结果格式化、日志记录、错误包装。
这种结构保证了中间件之间的解耦。一个负责日志的中间件完全不需要知道前面有没有权限校验中间件。
2.2 实现责任链的两种常见方式
在代码层面,我们常用责任链模式来实现这个洋葱模型。假设我们有一个基础的Agent执行函数async def execute_agent(input: str) -> str:。
方式一:装饰器(Decorator)链这是最Pythonic、最直观的方式。每个中间件都是一个装饰器函数,它接收一个执行函数(可能是核心函数,也可能是已被其他装饰器包装过的函数),返回一个新的包装函数。
def logging_middleware(next_fn): async def wrapper(input_text): start_time = time.time() print(f"[LOG] 开始执行,输入: {input_text}") try: result = await next_fn(input_text) # 调用链中的下一个函数(可能是核心,也可能是其他中间件) print(f"[LOG] 执行成功,耗时: {time.time() - start_time:.2f}s") return result except Exception as e: print(f"[LOG] 执行失败,错误: {e}") raise return wrapper def validation_middleware(next_fn): async def wrapper(input_text): if not input_text or len(input_text.strip()) == 0: raise ValueError("输入不能为空") # 可以在这里加入更复杂的业务规则校验 return await next_fn(input_text) return wrapper # 组合中间件:注意顺序,先应用的装饰器在外层 @logging_middleware @validation_middleware async def execute_agent(input: str) -> str: # 这里是Agent的核心逻辑 return f"处理结果: {input}" # 调用 result = await execute_agent("查询销售额")方式二:显式的中间件处理器列表这种方式更结构化,适合中间件数量多、需要动态配置的场景。定义一个中间件基类或协议,然后在一个执行器中顺序调用。
from abc import ABC, abstractmethod from typing import Any, Awaitable, Callable class Middleware(ABC): @abstractmethod async def handle(self, request: Any, next_fn: Callable[[Any], Awaitable[Any]]) -> Any: pass class LoggingMiddleware(Middleware): async def handle(self, request: str, next_fn): print(f"[MW-LOG] 请求入参: {request}") response = await next_fn(request) # 将控制权交给下一个中间件或核心逻辑 print(f"[MW-LOG] 响应结果: {response}") return response class AgentCore: async def execute(self, input_text: str) -> str: return f"核心处理结果: {input_text}" class MiddlewareExecutor: def __init__(self, core: AgentCore, middlewares: list[Middleware]): self.core = core self.middlewares = middlewares async def execute(self, input_text: str) -> str: # 构建一个最终的“下一个函数”链 def create_handler(index: int) -> Callable[[str], Awaitable[str]]: async def handler(current_input: str) -> str: if index < len(self.middlewares): # 调用当前中间件,并传入下一个handler return await self.middlewares[index].handle(current_input, create_handler(index + 1)) else: # 所有中间件执行完毕,调用核心逻辑 return await self.core.execute(current_input) return handler # 从第一个中间件开始执行 return await create_handler(0)(input_text) # 使用 core = AgentCore() executor = MiddlewareExecutor(core, [LoggingMiddleware()]) result = await executor.execute("测试输入")实操心得:对于中小型项目或中间件逻辑相对简单的场景,装饰器链更轻量、易读。而对于需要从配置文件加载中间件、中间件有复杂状态或依赖注入需求的企业级应用,显式的处理器列表模式更具优势,结构清晰,易于测试和动态编排。
3. 五大核心中间件场景与实现详解
理解了设计模式,我们来看几个必须实现的、具有代表性的中间件场景。我会给出关键代码和配置要点。
3.1 上下文注入中间件:让Agent拥有“记忆”
Agent的每次调用通常是独立的,但它需要历史对话、用户画像、业务数据作为上下文才能做出精准判断。硬编码在Prompt里不灵活,每次都去数据库查又低效。
实现要点:
- 在Inbound阶段工作:在执行核心Agent前,从外部源(Redis、数据库、向量库)获取数据。
- 数据拼接:将获取的数据以特定格式(如
## 用户历史:\n{history}\n## 当前问题:)拼接到原始用户输入之前,形成新的、增强版的Prompt。 - 缓存策略:为减轻外部源压力,引入本地缓存(如
functools.lru_cache)或分布式缓存(Redis),并设置合理的过期时间。
import redis.asyncio as redis from functools import lru_cache class ContextEnrichmentMiddleware: def __init__(self, redis_client: redis.Redis): self.redis = redis_client async def handle(self, request: dict, next_fn) -> dict: """ request 结构: {'user_id': '123', 'input_text': '...', 'session_id': '...'} """ user_id = request.get('user_id') session_id = request.get('session_id') # 1. 获取用户历史对话(带缓存) history = await self._get_conversation_history(session_id) # 2. 获取用户画像(如偏好、权限) profile = await self._get_user_profile(user_id) # 3. 获取实时业务数据(如产品库存) biz_data = await self._get_business_data(request['input_text']) # 构建增强的Prompt enriched_prompt = f""" 用户信息:{profile} 最近对话历史: {history} 当前问题:{request['input_text']} 相关业务数据:{biz_data} 请根据以上信息回答。 """ # 更新请求 request['enriched_prompt'] = enriched_prompt # 将增强后的请求传递给下一个环节 return await next_fn(request) @lru_cache(maxsize=128) async def _get_conversation_history(self, session_id: str) -> str: # 这里可以使用缓存,避免频繁查询数据库 cache_key = f"chat_history:{session_id}" history = await self.redis.get(cache_key) if history: return history.decode('utf-8') # 缓存未命中,查询数据库... # db_history = await db.query(...) # await self.redis.setex(cache_key, 300, db_history) # 缓存5分钟 # return db_history return "暂无历史对话"注意事项:上下文注入要避免“信息过载”。给Agent太多无关信息会干扰其判断,增加Token消耗和成本。最佳实践是动态相关性检索,例如使用向量数据库,只检索与当前问题最相关的几条历史记录或文档片段,而不是全部加载。
3.2 合规与安全审查中间件:守住红线
这是企业应用的生死线。必须在Agent调用外部工具(Tool Call)或最终输出答案前,进行审查。
实现要点:
- 双检查点:
- 工具调用前(Pre-tool-call):检查工具名、参数是否被允许。例如,是否允许该用户调用“send_email”工具?参数中的收件人是否在白名单内?
- 最终输出后(Post-output):对Agent生成的文本进行敏感词过滤、个人隐私信息(PII)脱敏、内容安全审核。
- 异步与同步审查:对于简单的关键词过滤,可以同步进行;对于复杂的AI内容审核(如调用内容安全API),应采用异步非阻塞方式,避免影响主流程响应速度。
- 熔断机制:当连续多次审查不通过时,应触发熔断,暂时禁用某些功能或转人工。
class SecurityAndComplianceMiddleware: def __init__(self, forbidden_tools: set, pii_detector): self.forbidden_tools = forbidden_tools # 例如 {'delete_database', 'format_system'} self.pii_detector = pii_detector # 一个用于检测手机号、邮箱等的正则或模型 async def handle(self, request: dict, next_fn) -> dict: # 假设request中包含了Agent规划的行动步骤 action_plan = request.get('action_plan', []) for step in action_plan: if step['type'] == 'tool_call': tool_name = step['name'] # 1. 工具权限校验 if tool_name in self.forbidden_tools: raise PermissionError(f"无权调用工具: {tool_name}") # 2. 参数安全检查(示例:检查SQL注入风险) if tool_name == 'query_database': sql = step['args'].get('query', '') if self._has_sql_injection_risk(sql): raise SecurityError("检测到潜在的SQL注入风险") # 安全校验通过,执行后续逻辑 response = await next_fn(request) # 3. 对最终输出进行PII脱敏 if 'final_output' in response: response['final_output'] = self.pii_detector.redact(response['final_output']) return response def _has_sql_injection_risk(self, sql: str) -> bool: # 简单的关键词检查,生产环境应使用更专业的SQL解析库或ORM的参数化查询来根本性避免 dangerous_patterns = ["DROP", "DELETE FROM", "INSERT INTO", "--", ";", "UNION SELECT"] import re for pattern in dangerous_patterns: if re.search(rf'\b{pattern}\b', sql, re.IGNORECASE): return True return False3.3 可观测性与监控中间件:看清每一个环节
没有监控的系统就是“盲开”。这个中间件负责无侵入地收集全链路数据。
实现要点:
- 结构化日志:不仅打印文本,更要以JSON等结构化格式记录,方便后续接入ELK(Elasticsearch, Logstash, Kibana)或时序数据库。记录内容应包括:
request_id,user_id,stage(阶段),input,output,latency(延迟),token_usage,error(如果有)。 - 指标埋点:使用像Prometheus这样的工具,记录计数器(如请求总数、错误数)、直方图(如响应时间分布)。
- 分布式追踪:如果系统是分布式的(例如Agent、工具、模型服务分开部署),集成OpenTelemetry来追踪一个请求在所有服务间的流转路径。
import time import json import logging from contextvars import ContextVar from prometheus_client import Counter, Histogram # 定义指标 REQUEST_COUNT = Counter('agent_requests_total', 'Total agent requests', ['endpoint', 'status']) REQUEST_LATENCY = Histogram('agent_request_duration_seconds', 'Request latency in seconds', ['endpoint']) # 用于传递请求ID的上下文变量 request_id_ctx: ContextVar[str] = ContextVar('request_id', default='unknown') class ObservabilityMiddleware: def __init__(self, logger: logging.Logger): self.logger = logger async def handle(self, request: dict, next_fn) -> dict: request_id = request.get('request_id', 'default_id') request_id_ctx.set(request_id) start_time = time.time() endpoint = request.get('path', 'execute') # 记录请求开始日志 self.logger.info(json.dumps({ "request_id": request_id, "event": "request_start", "input": request.get('input_text', '')[:200], # 截断避免日志过大 "user": request.get('user_id') })) try: response = await next_fn(request) status = "success" # 记录成功日志和指标 self.logger.info(json.dumps({ "request_id": request_id, "event": "request_end", "status": status, "latency_sec": time.time() - start_time, "token_used": response.get('usage', {}).get('total_tokens', 0) })) REQUEST_COUNT.labels(endpoint=endpoint, status=status).inc() REQUEST_LATENCY.labels(endpoint=endpoint).observe(time.time() - start_time) return response except Exception as e: status = "error" # 记录错误日志和指标 self.logger.error(json.dumps({ "request_id": request_id, "event": "request_error", "error_type": e.__class__.__name__, "error_msg": str(e), "latency_sec": time.time() - start_time }), exc_info=True) REQUEST_COUNT.labels(endpoint=endpoint, status=status).inc() raise # 重新抛出异常,让上层处理 # 在其他地方可以通过 request_id_ctx.get() 获取当前请求ID3.4 流式输出中间件:提升用户体验
对于生成时间较长的内容,让用户干等着是一种糟糕的体验。中间件可以拦截Agent的流式生成结果(如果底层LLM支持),并实时推送给前端。
实现要点:
- 拦截生成器:核心Agent的
execute方法需要改为异步生成器(async generator),逐步yieldtokens或chunks。 - 中间件也需支持生成器:中间件的
handle方法需要能够接收并包装一个生成器,在yield每个片段前后执行逻辑(例如,计算已生成字数、进行部分结果的安全扫描)。 - 与Web框架集成:对于FastAPI或类似框架,可以将这个生成器作为
StreamingResponse的内容返回。
from typing import AsyncGenerator class StreamingMiddleware: """处理流式输出的中间件示例""" async def handle_stream(self, request: dict, next_fn) -> AsyncGenerator[str, None]: """ next_fn 现在返回一个异步生成器 (AsyncGenerator) """ # Inbound 逻辑:可以在这里做一些初始化 buffer = [] print(f"开始流式处理请求: {request['input_text'][:50]}...") # 调用下一个环节,获取流式生成器 stream_generator = await next_fn(request) # 注意这里,next_fn可能返回一个生成器 # 处理流式响应 async for chunk in stream_generator: # 这里可以加入后处理逻辑,例如: # 1. 缓存片段(用于后续完整内容分析) buffer.append(chunk) # 2. 简单的实时敏感词检测(生产环境可能需要更复杂的缓冲窗口检测) if "敏感词" in chunk: chunk = chunk.replace("敏感词", "***") # 3. 添加SSE (Server-Sent Events) 格式包装 # formatted_chunk = f"data: {chunk}\n\n" yield chunk # 将处理后的片段返回给上游 # 流结束后,可以做一些清理或最终分析 full_content = ''.join(buffer) print(f"流式处理结束,总长度: {len(full_content)}") # 例如,将完整内容存入审计日志 # await audit_log(request, full_content) # 假设核心Agent的流式执行函数 async def agent_execute_stream(request: dict) -> AsyncGenerator[str, None]: # 模拟LLM流式生成 simulated_response = "这是一个流式生成的回答,它被分成多个部分返回。" for word in simulated_response.split(): await asyncio.sleep(0.1) # 模拟生成延迟 yield word + " "3.5 故障熔断与降级中间件:保障系统韧性
当依赖的下游服务(如LLM API、向量数据库、工具服务)不稳定时,不能让Agent无限重试或长时间等待,需要有快速失败和降级策略。
实现要点:
- 集成熔断器模式:使用如
pybreaker库,为每个关键依赖设置熔断器。当失败次数超过阈值,熔断器“打开”,短时间内直接拒绝请求,避免雪崩;经过一段时间后进入“半开”状态试探性放行少量请求。 - 超时控制:为每个中间件环节或工具调用设置独立的超时时间,使用
asyncio.wait_for防止单个慢请求阻塞整个链路。 - 优雅降级:当核心LLM服务不可用时,可以降级到规则引擎返回一个预设的兜底答案,或者返回一个友好的错误提示,引导用户稍后重试。
import pybreaker import asyncio from datetime import datetime # 为LLM调用定义一个熔断器 llm_breaker = pybreaker.CircuitBreaker(fail_max=5, reset_timeout=60) # 5次失败后熔断60秒 class CircuitBreakerMiddleware: def __init__(self): self.breaker = llm_breaker async def handle(self, request: dict, next_fn) -> dict: # 尝试通过熔断器执行 try: # 使用熔断器包装调用 response = await self.breaker.call_async(self._call_with_timeout, request, next_fn) return response except pybreaker.CircuitBreakerError as e: # 熔断器已打开,直接拒绝请求,执行降级逻辑 self.logger.warning(f"熔断器已打开,请求被快速失败: {e}") return await self._fallback_response(request) except asyncio.TimeoutError: self.logger.error("请求超时") # 超时也算作一次失败,熔断器会计数 raise except Exception as e: # 其他异常也会被熔断器捕获并计数 self.logger.error(f"请求执行失败: {e}") raise async def _call_with_timeout(self, request: dict, next_fn) -> dict: """为关键调用添加超时控制""" try: # 设置10秒超时 return await asyncio.wait_for(next_fn(request), timeout=10.0) except asyncio.TimeoutError: self.logger.error("下游服务响应超时") raise async def _fallback_response(self, request: dict) -> dict: """优雅降级:返回预设的兜底答案""" # 这里可以根据请求类型返回不同的兜底内容 # 例如,从缓存中获取旧答案,或返回一个静态提示 return { 'final_output': '当前服务暂时繁忙,您可以稍后重试,或尝试描述其他问题。', 'source': 'fallback', 'is_fallback': True }4. 中间件系统的配置、测试与部署实践
设计好了中间件,如何管理它们,并确保它们稳定可靠地工作?
4.1 动态配置与编排
我们不应该把中间件的启用顺序和参数硬编码在代码里。理想的方式是通过配置文件(如YAML)或配置中心来管理。
# middleware_config.yaml middlewares: - name: request_logging class: "app.middleware.ObservabilityMiddleware" enabled: true params: log_level: "INFO" - name: context_enrichment class: "app.middleware.ContextEnrichmentMiddleware" enabled: true params: redis_url: "redis://localhost:6379/0" cache_ttl: 300 - name: security_check class: "app.middleware.SecurityAndComplianceMiddleware" enabled: true params: forbidden_tools: ["send_email", "execute_code"] - name: circuit_breaker class: "app.middleware.CircuitBreakerMiddleware" enabled: true # 注意:顺序很重要!先执行的中间件在外层。然后在应用启动时,根据这个配置动态加载和实例化中间件,构建执行链。这样,在不停机的情况下,通过修改配置就能启用、禁用或调整某个中间件。
4.2 中间件的单元测试与集成测试
中间件逻辑必须被充分测试,尤其是安全、合规相关的。
- 单元测试:针对每个中间件的
handle方法,模拟输入request和next_fn(通常是一个返回固定值的简单异步函数),断言其输出或副作用(如日志记录、指标增长)是否符合预期。import pytest from unittest.mock import AsyncMock, MagicMock @pytest.mark.asyncio async def test_security_middleware_blocks_forbidden_tool(): middleware = SecurityAndComplianceMiddleware(forbidden_tools={"dangerous_tool"}) mock_next = AsyncMock(return_value={"output": "ok"}) malicious_request = { 'action_plan': [{'type': 'tool_call', 'name': 'dangerous_tool'}] } with pytest.raises(PermissionError): await middleware.handle(malicious_request, mock_next) # 确保下一个函数没有被调用 mock_next.assert_not_called() - 集成测试:将多个中间件与核心Agent串联起来,进行端到端测试。使用测试专用的LLM Mock(如
unittest.mock.patch掉OpenAI客户端)和工具Mock,验证整个链路的输入输出。
4.3 性能考量与最佳实践
中间件会增加开销,设计不当会成为性能瓶颈。
- 异步(Async)是所有I/O操作的必须:任何涉及网络调用(数据库、Redis、API)的中间件逻辑,都必须使用异步库(如
aioredis,aiohttp)和async/await语法,避免阻塞事件循环。 - 避免在中间件中进行重型同步计算:如复杂的字符串处理、大文件解析。如果必须,考虑将其放入线程池执行(
asyncio.to_thread),防止阻塞其他并发请求。 - 设置合理的超时和超时传递:为每个可能耗时的中间件环节设置独立的超时,并且超时时间应从外向内逐层递减,确保整体响应时间可控。
- 中间件应保持无状态(Stateless):尽可能让中间件本身不保存请求相关的状态。状态应存储在请求上下文(
ContextVar)或外部存储中。这有利于水平扩展和避免并发问题。 - 谨慎使用全局变量:在多线程/协程环境下,修改全局变量是危险的。如果需要共享配置,使用依赖注入或单例模式。
5. 常见问题排查与实战技巧
在实际部署和运行中,你肯定会遇到各种问题。这里记录了几个最典型的坑和解决办法。
5.1 中间件执行顺序错乱或未生效
问题现象:日志中间件没记录到错误,或者安全校验在工具调用后才执行。排查步骤:
- 检查注册顺序:确认中间件加载到执行链(
MiddlewareExecutor)的顺序是否正确。最外层的中间件最先执行handle的“入站”逻辑,但最后执行“出站”逻辑。画一个简单的执行顺序图来验证。 - 验证中间件是否被正确包装:在装饰器模式下,检查装饰器应用的顺序。
@a @b def fn()等价于fn = a(b(fn)),所以a是最外层。 - 检查中间件的
next_fn调用:确保每个中间件的handle方法内部都正确调用了await next_fn(request),这是责任链向下传递的关键。如果某个中间件因为条件判断提前返回了响应,链就会在此中断。
5.2 内存泄漏或性能下降
问题现象:服务运行一段时间后,内存持续增长,响应变慢。可能原因与解决:
- 中间件中未释放资源:例如,在中间件中打开了数据库连接或文件句柄,但在异常情况下没有正确关闭。使用
try...finally块或异步上下文管理器(async with)来确保资源释放。 - 缓存失控:上下文注入中间件如果使用了无过期时间或过大的缓存(如
lru_cache缓存了过大的对象),会导致内存堆积。为缓存设置大小限制和合理的TTL。 - 日志中间件记录全量数据:如果日志中间件无差别地记录完整的请求和响应体(特别是包含大文件Base64编码),会迅速耗尽磁盘和内存。务必进行截断或采样记录。
- 使用
asyncio.create_task后未管理:如果在中间件内部分发了后台任务(如异步发送审计日志),但没有妥善收集或等待它们,可能导致任务堆积。考虑使用后台任务管理器。
5.3 错误处理与异常传播
问题现象:某个中间件抛出异常后,整个请求无声无息地失败了,或者返回了不友好的错误。最佳实践:
- 中间件应捕获自身逻辑的异常:如果中间件的辅助功能(如发送监控指标)失败,不应中断主业务流程。应使用
try...except记录错误后,继续调用next_fn。async def handle(self, request, next_fn): try: # 非核心的监控上报 await self._report_metrics(request) except Exception as e: self.logger.error(f"监控上报失败,但不中断流程: {e}") # 核心流程继续 return await next_fn(request) - 区分业务异常和系统异常:安全中间件校验不通过,应抛出明确的业务异常(如
PermissionError),并由全局异常处理器转换为对用户友好的错误信息。而网络超时等系统异常,则应触发熔断或重试机制。 - 设计统一的错误响应格式:在所有中间件和核心逻辑的最外层,有一个全局异常捕获中间件,负责将任何未处理的异常转换为结构化的错误响应(如
{"code": 500, "msg": "系统内部错误", "request_id": "xxx"}),并确保记录到日志。
5.4 调试困难:请求在哪个中间件卡住了?
技巧:为每个请求生成唯一的request_id,并在每个中间件的入口和出口都打印带有该ID的日志。这样,在分布式日志系统中,你可以通过request_id轻松串联起一个请求在所有中间件和核心服务中的完整生命周期,快速定位延迟或错误的环节。这就是前面ObservabilityMiddleware和ContextVar所做的事情。
构建一个健壮的Agent中间件系统,初期会花费一些设计精力,但一旦建成,它将成为你AI应用架构中不可或缺的“神经系统”,让你对Agent的每一次运行都了如指掌、掌控自如。从简单的日志记录开始,逐步引入安全、监控、流式处理等中间件,你的Agent系统会在这个过程中变得越来越强大和可靠。