1. 项目概述与核心价值
聊到AI Agent,现在大家都不陌生了,各种框架和Demo满天飞。但说实话,很多朋友跟我交流时都提到一个痛点:看教程跑通一个简单的Agent例子不难,但一旦想把它做成一个能稳定运行、能处理真实业务、能上线给用户用的“产品级”应用,立刻就感觉无从下手。这中间的鸿沟,远比想象的要大。我自己在从实验原型到生产系统的过程中,踩过无数的坑,也正是在这些实战里,我逐渐总结出了一套构建“产品级Agent”的方法论和工具链,我称之为Agent Harness。
这个系列,我们就来彻底拆解这个“Harness”。它不是某个特定的框架,而是一套工程化的约束、规范和最佳实践的集合。你可以把它理解为一套为Agent量身定制的“安全带”和“缰绳”,目的是让这个能力强大但行为可能难以预测的“智能体”,能够在可控、可靠、可观测的轨道上运行,最终交付稳定的业务价值。今天这第三篇,我们将深入到最核心的部分:状态管理、工具调用与编排、以及异常处理与自愈。这是决定你的Agent是“玩具”还是“工具”的关键分水岭。
2. 架构设计思路:为什么需要“Harness”?
在动手写代码之前,我们必须先想清楚为什么。一个在Jupyter Notebook里跑得欢的Agent,直接搬到生产环境为什么大概率会“翻车”?核心原因在于生产环境对系统的要求是根本性的不同。
2.1 实验环境与生产环境的本质差异
在实验阶段,我们的目标是“验证可能性”。我们关心的是:这个Agent能不能理解我的指令?能不能调用正确的工具?输出的结果看起来对不对?这个过程往往是单次、交互式的,环境是纯净的,数据是准备好的,我们作为开发者全程监控。
而生产环境要求的是“保障确定性”。系统必须满足:
- 可靠性:7x24小时稳定运行,处理高并发请求,不能轻易崩溃。
- 可观测性:任何时候都能知道Agent内部在“想”什么、做了什么、为什么出错。
- 可控性:能够限制Agent的行为边界(比如不能执行危险操作),能够设置超时、重试等策略。
- 可维护性:代码结构清晰,模块解耦,便于迭代、调试和团队协作。
- 成本可控:每一次LLM API调用、每一次工具执行都有成本,需要精细化管理。
“Harness”就是为了弥合这中间的差距而生的。它通过一系列设计模式和技术选型,在Agent强大的认知能力之上,叠加一层工程化的保障。
2.2 Agent Harness的核心组件模型
基于上述目标,一个完整的Agent Harness通常包含以下几个核心组件,它们共同构成了Agent的“运行时环境”:
- 状态管理引擎:负责维护Agent在一次会话或一次任务执行周期内的所有上下文信息。这不仅仅是聊天历史,还包括工具调用结果、中间决策、用户会话数据等。它必须支持持久化、版本化和并发安全。
- 工具编排与执行层:负责管理Agent可用的所有工具(Tools)。包括工具的注册、发现、描述生成、参数验证、安全执行、结果格式化等。这是Agent与外部世界交互的桥梁。
- 工作流与决策控制器:控制Agent的执行逻辑。是简单的“思考-行动”循环(ReAct模式),还是更复杂的多步骤规划(Plan-and-Execute)?是否需要子任务分解?这部分定义了Agent的“行为模式”。
- 可观测性与监控套件:贯穿始终的日志、指标(Metrics)和追踪(Tracing)系统。必须能记录每一次LLM调用(输入/输出/Token消耗)、每一次工具调用(参数/结果/耗时)、每一次状态变更。
- 异常处理与自愈机制:预设各种故障场景(如网络超时、工具错误、LLM返回格式异常、内容安全审核失败等)的应对策略,如重试、降级、转人工或安全终止。
接下来的内容,我们将聚焦于前三个核心组件的实现细节。
3. 核心实现一:持久化与并发安全的状态管理
状态管理是Agent的“记忆”系统。一个糟糕的状态管理设计,会导致上下文丢失、会话混乱、难以调试。
3.1 状态数据模型设计
首先,我们需要定义状态里到底存什么。一个丰富的状态对象可能包含以下字段:
from pydantic import BaseModel, Field from datetime import datetime from typing import Dict, Any, List, Optional from enum import Enum class TaskStatus(Enum): PENDING = "pending" RUNNING = "running" SUCCESS = "success" FAILED = "failed" CANCELLED = "cancelled" class AgentState(BaseModel): """Agent核心状态模型""" # 会话标识 session_id: str task_id: str user_id: Optional[str] = None # 核心上下文 conversation_history: List[Dict[str, Any]] = Field(default_factory=list) # 消息历史 current_goal: Optional[str] = None # 当前任务目标 extracted_entities: Dict[str, Any] = Field(default_factory=dict) # 从对话中提取的实体信息 context_variables: Dict[str, Any] = Field(default_factory=dict) # 自定义上下文变量 # 执行轨迹 execution_stack: List[str] = Field(default_factory=list) # 执行步骤栈(用于复杂任务分解) tool_calls_history: List[Dict[str, Any]] = Field(default_factory=list) # 工具调用历史 # 元数据 status: TaskStatus = TaskStatus.PENDING created_at: datetime = Field(default_factory=datetime.utcnow) updated_at: datetime = Field(default_factory=datetime.utcnow) metadata: Dict[str, Any] = Field(default_factory=dict) # 扩展元数据 class Config: use_enum_values = True # 序列化时使用枚举值设计要点解析:
- 使用Pydantic:利用其数据验证和序列化能力,确保状态数据的结构一致性。
- 区分历史与当前上下文:
conversation_history存储原始对话,extracted_entities和context_variables存储结构化信息,便于工具使用。 - 执行轨迹记录:
execution_stack和tool_calls_history对于调试和实现复杂逻辑(如回退、继续)至关重要。 - 状态枚举:明确定义任务生命周期,便于监控和管理。
3.2 状态存储后端选型与实现
状态存储需要根据数据量、并发量和持久化要求来选择。
场景一:单实例/轻量级应用——内存 + 文件备份适用于原型或低并发场景。使用内存字典存储活跃会话,定期序列化到文件(如JSON)做持久化。
import json import asyncio from pathlib import Path from typing import Dict import aiofiles class FileBackedStateManager: def __init__(self, storage_path: Path = Path("./agent_states")): self.storage_path = storage_path self.storage_path.mkdir(exist_ok=True) self._in_memory_cache: Dict[str, AgentState] = {} self._lock = asyncio.Lock() # 简易锁,处理并发写入 async def get_state(self, session_id: str) -> Optional[AgentState]: """获取状态:内存优先,文件回退""" # 1. 检查内存缓存 if session_id in self._in_memory_cache: return self._in_memory_cache[session_id].copy(deep=True) # 2. 从文件加载 file_path = self.storage_path / f"{session_id}.json" if file_path.exists(): async with aiofiles.open(file_path, 'r', encoding='utf-8') as f: data = json.loads(await f.read()) state = AgentState(**data) async with self._lock: self._in_memory_cache[session_id] = state return state.copy(deep=True) return None async def save_state(self, state: AgentState): """保存状态:更新内存,异步写入文件""" async with self._lock: self._in_memory_cache[state.session_id] = state.copy(deep=True) # 异步写入文件,避免阻塞主流程 file_path = self.storage_path / f"{state.session_id}.json" state.updated_at = datetime.utcnow() async with aiofiles.open(file_path, 'w', encoding='utf-8') as f: await f.write(state.json(indent=2, ensure_ascii=False))注意:这种方案在服务器重启时会丢失内存中的状态,但可以从文件恢复。对于生产环境,仅适用于可容忍短暂状态丢失或会话无关紧要的场景。
场景二:生产环境——RedisRedis是生产环境中最常见的选择,它提供了高性能、持久化、数据结构丰富和分布式支持。
import redis.asyncio as redis from redis.commands.json.path import Path import pickle # 或使用msgpack, orjson class RedisStateManager: def __init__(self, redis_url: str, ttl: int = 3600): self.client = redis.from_url(redis_url, decode_responses=False) self.ttl = ttl # 状态过期时间,避免内存泄漏 async def get_state(self, session_id: str) -> Optional[AgentState]: # 使用pickle序列化复杂对象,或使用RedisJSON模块 data = await self.client.get(f"agent:state:{session_id}") if data: # 使用pickle反序列化 state_dict = pickle.loads(data) return AgentState(**state_dict) return None async def save_state(self, state: AgentState): state.updated_at = datetime.utcnow() state_dict = state.dict() # 使用pickle序列化 data = pickle.dumps(state_dict) await self.client.setex( name=f"agent:state:{session_id}", time=self.ttl, value=data )场景三:高要求生产环境——数据库(PostgreSQL/MySQL)当状态数据非常庞大,需要复杂查询(如按用户、时间、状态筛选)、强一致性或与其他业务数据关联时,需要使用关系型数据库或文档数据库。
# 以SQLAlchemy异步ORM为例(简化) from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine from sqlalchemy.orm import declarative_base, sessionmaker from sqlalchemy import Column, String, JSON, DateTime, Enum Base = declarative_base() class AgentStateORM(Base): __tablename__ = 'agent_states' session_id = Column(String, primary_key=True) state_data = Column(JSON) # 存储序列化的状态字典 status = Column(String) created_at = Column(DateTime) updated_at = Column(DateTime) class DBStateManager: def __init__(self, database_url: str): self.engine = create_async_engine(database_url) self.async_session = sessionmaker(self.engine, class_=AsyncSession, expire_on_commit=False) async def get_state(self, session_id: str) -> Optional[AgentState]: async with self.async_session() as session: result = await session.get(AgentStateORM, session_id) if result: return AgentState(**result.state_data) return None async def save_state(self, state: AgentState): async with self.async_session() as session: state_orm = AgentStateORM( session_id=state.session_id, state_data=state.dict(), status=state.status.value, updated_at=datetime.utcnow() ) await session.merge(state_orm) # 使用merge处理upsert await session.commit()选型心得:
- 开发/测试阶段:用文件备份或内存存储,最简单快捷。
- 中小型生产应用:Redis是首选。性能极高,支持丰富数据结构(Hash, List, Sorted Set可用于优先级队列),设置TTL自动清理。记得配置RDB/AOF持久化。
- 大型、状态复杂、需关联查询的应用:用关系型数据库。虽然性能不如Redis,但保证了数据的可靠性和查询灵活性。可以采用缓存+数据库的组合,热数据放Redis,冷数据或需要分析的数据落库。
3.3 状态管理的并发与锁
在生产环境中,同一个会话可能同时收到多个请求(比如用户快速发送消息)。如果不加控制,可能导致状态覆盖,出现“丢失中间步骤”的诡异问题。
解决方案:乐观锁或分布式锁。
乐观锁实现(基于版本号): 在状态模型中增加一个version字段。每次更新时,检查当前版本号是否与读取时一致。
class AgentState(BaseModel): # ... 其他字段同上 version: int = 0 class OptimisticLockingStateManager(RedisStateManager): async def save_state(self, state: AgentState, expected_version: int) -> bool: """保存状态,使用乐观锁。返回是否成功""" # 使用Redis的WATCH/MULTI/EXEC实现乐观锁 async with self.client.pipeline(transaction=True) as pipe: try: await pipe.watch(f"agent:state:{state.session_id}") current_data = await pipe.get(f"agent:state:{state.session_id}") if current_data: current_state = AgentState(**pickle.loads(current_data)) if current_state.version != expected_version: await pipe.unwatch() return False # 版本冲突,保存失败 state.version = expected_version + 1 state.updated_at = datetime.utcnow() data = pickle.dumps(state.dict()) pipe.multi() pipe.setex(f"agent:state:{state.session_id}", self.ttl, data) await pipe.execute() return True except redis.WatchError: # 在WATCH期间键被其他客户端修改 return False在实际调用时,流程如下:
state = await state_manager.get_state(session_id) # ... 基于state进行一系列处理,生成新的state_new ... success = await state_manager.save_state(state_new, expected_version=state.version) if not success: # 处理冲突:重试或返回错误给用户 raise StateConflictError("会话状态已过期,请重试")注意事项:
- 对于简单的Agent,如果处理逻辑是线性的(一个请求处理完才接收下一个),可以在业务逻辑层用队列或锁来序列化请求,避免并发写。但乐观锁是更通用、更 scalable 的方案。
- 锁的粒度要仔细设计。太粗(如锁整个管理器)会严重影响性能;太细(如每个字段)会增加复杂度。通常以会话(session_id)为粒度是合理的。
4. 核心实现二:健壮的工具调用与编排层
工具(Tools)是Agent能力的延伸。一个健壮的工具层需要解决:如何让Agent知道有哪些工具可用?如何安全地执行工具?如何处理工具的错误?
4.1 工具的定义与注册中心
首先,我们需要一个统一的方式来定义工具。一个工具至少包含:名称、描述、参数模式、执行函数。
from typing import Callable, Any, Dict, List, Optional, get_type_hints from pydantic import BaseModel, create_model import inspect class Tool(BaseModel): """工具定义""" name: str description: str args_schema: Optional[BaseModel] = None # Pydantic模型,用于参数验证 func: Callable class Config: arbitrary_types_allowed = True async def execute(self, **kwargs) -> Any: """执行工具,并做基础验证""" # 1. 参数验证 if self.args_schema: validated_args = self.args_schema(**kwargs).dict() else: validated_args = kwargs # 2. 执行(支持同步和异步函数) if inspect.iscoroutinefunction(self.func): result = await self.func(**validated_args) else: result = self.func(**validated_args) return result def get_openai_function_schema(self) -> Dict[str, Any]: """生成OpenAI Function Calling格式的schema""" schema = { "name": self.name, "description": self.description, } if self.args_schema: # 将Pydantic模型转换为JSON Schema schema["parameters"] = self.args_schema.schema() else: schema["parameters"] = {"type": "object", "properties": {}} return schema工具注册中心管理所有可用工具:
class ToolRegistry: def __init__(self): self._tools: Dict[str, Tool] = {} def register(self, tool: Tool): if tool.name in self._tools: raise ValueError(f"Tool '{tool.name}' already registered.") self._tools[tool.name] = tool def register_from_function(self, func: Callable, name: str = None, description: str = None): """从普通函数自动创建并注册工具""" tool_name = name or func.__name__ tool_desc = description or func.__doc__ or "" # 从函数签名推断参数schema sig = inspect.signature(func) fields = {} for param_name, param in sig.parameters.items(): if param_name == 'self': continue # 简化处理:这里需要根据实际类型映射到Pydantic字段,此处省略复杂逻辑 # 实际项目中可以使用更完善的类型推断库 fields[param_name] = (Optional[Any], ...) # 简化示例 args_model = create_model(f"{tool_name}Args", **fields) if fields else None tool = Tool(name=tool_name, description=tool_desc, args_schema=args_model, func=func) self.register(tool) def get_tool(self, name: str) -> Optional[Tool]: return self._tools.get(name) def list_tools(self) -> List[Tool]: return list(self._tools.values()) def get_openai_functions(self) -> List[Dict[str, Any]]: return [tool.get_openai_function_schema() for tool in self._tools.values()]实操示例:定义几个常用工具
from datetime import datetime # 1. 使用装饰器注册(更优雅) registry = ToolRegistry() def register_tool(name: str = None, description: str = None): def decorator(func): registry.register_from_function(func, name=name, description=description) return func return decorator @register_tool( name="get_current_time", description="获取当前的日期和时间。当用户询问时间或日期时使用此工具。" ) async def get_current_time(timezone: str = "UTC") -> str: """获取指定时区的当前时间""" # 这里简化处理,实际应使用pytz等库 now = datetime.utcnow() return f"当前时间({timezone})是:{now.isoformat()}" @register_tool( name="search_web", description="在互联网上搜索信息。当你需要获取最新、未知的或特定网站的信息时使用。" ) async def search_web(query: str, max_results: int = 5) -> List[Dict[str, str]]: """模拟网络搜索""" # 实际应接入Serper API、Google Search API等 # 此处返回模拟数据 return [ {"title": f"关于 {query} 的结果1", "snippet": "这是摘要1...", "url": "https://example.com/1"}, {"title": f"关于 {query} 的结果2", "snippet": "这是摘要2...", "url": "https://example.com/2"}, ] @register_tool( name="calculate", description="执行数学计算。支持加(+)、减(-)、乘(*)、除(/)、幂(**)等基本运算。" ) async def calculate(expression: str) -> float: """计算数学表达式""" # 警告:直接使用eval有安全风险!生产环境应用用ast.literal_eval或安全计算库 # 此处仅为示例,务必进行严格的输入验证和沙箱隔离 try: # 非常简单的安全过滤示例(不完善) allowed_chars = set("0123456789+-*/(). ") if not all(c in allowed_chars for c in expression): raise ValueError("表达式包含不安全字符") result = eval(expression, {"__builtins__": {}}, {}) return float(result) except Exception as e: raise ValueError(f"计算失败: {e}")重要安全警告:
calculate工具中的eval用法是极其危险的,仅用于演示。在生产环境中,绝对禁止直接eval用户输入的字符串。必须使用安全的表达式求值库(如asteval),或在严格沙箱环境中执行。这是构建可靠Agent的底线之一。
4.2 工具执行器:安全、超时与隔离
工具执行不能是“裸奔”的。我们需要一个执行器来包裹所有工具调用,提供统一的保障。
import asyncio from concurrent.futures import ThreadPoolExecutor from contextlib import asynccontextmanager import traceback from typing import Tuple class ToolExecutor: def __init__(self, registry: ToolRegistry, timeout: int = 30, max_workers: int = 10): self.registry = registry self.timeout = timeout # 用于执行同步的、可能阻塞的工具 self.thread_pool = ThreadPoolExecutor(max_workers=max_workers) async def execute( self, tool_name: str, arguments: Dict[str, Any], state: AgentState ) -> Tuple[bool, Any, str]: """ 执行工具。 返回: (是否成功, 执行结果或错误信息, 可读的日志) """ tool = self.registry.get_tool(tool_name) if not tool: return False, None, f"错误:未找到工具 '{tool_name}'" log_parts = [f"调用工具: {tool_name}"] if arguments: log_parts.append(f"参数: {arguments}") try: # 1. 参数验证(已在Tool.execute中处理) # 2. 带超时执行 if inspect.iscoroutinefunction(tool.func): # 异步函数 task = asyncio.create_task(tool.execute(**arguments)) result = await asyncio.wait_for(task, timeout=self.timeout) else: # 同步函数,放到线程池执行,避免阻塞事件循环 loop = asyncio.get_event_loop() func = tool.execute result = await loop.run_in_executor( self.thread_pool, lambda: func(**arguments) ) log_parts.append(f"结果: {str(result)[:200]}...") # 截断长结果 # 记录到状态 state.tool_calls_history.append({ "tool": tool_name, "arguments": arguments, "result": result, "timestamp": datetime.utcnow().isoformat(), "success": True }) return True, result, " | ".join(log_parts) except asyncio.TimeoutError: error_msg = f"工具 '{tool_name}' 执行超时(>{self.timeout}秒)" log_parts.append(error_msg) state.tool_calls_history.append({ "tool": tool_name, "arguments": arguments, "error": error_msg, "timestamp": datetime.utcnow().isoformat(), "success": False }) return False, None, " | ".join(log_parts) except Exception as e: error_msg = f"工具 '{tool_name}' 执行出错: {str(e)}" log_parts.append(error_msg) # 记录详细堆栈到日志系统,但返回给用户的信息要简化 state.tool_calls_history.append({ "tool": tool_name, "arguments": arguments, "error": error_msg, "traceback": traceback.format_exc(), "timestamp": datetime.utcnow().isoformat(), "success": False }) return False, None, " | ".join(log_parts)设计要点:
- 统一错误处理:所有工具异常都在这里捕获,避免单个工具崩溃导致整个Agent崩溃。
- 超时控制:防止某些工具(如网络请求)无限期挂起,拖垮整个系统。
- 线程池执行同步代码:避免同步的CPU密集型或阻塞IO操作阻塞异步事件循环。
- 执行日志记录:将每次工具调用的详情记录到Agent状态中,便于后续调试和审计。
4.3 工具编排与Agent核心循环
有了状态管理和工具执行器,我们可以构建Agent的核心决策与执行循环了。这里以经典的ReAct (Reasoning + Acting)模式为例。
from openai import AsyncOpenAI # 或其他LLM客户端 class ReActAgent: def __init__( self, llm_client: AsyncOpenAI, tool_executor: ToolExecutor, state_manager: StateManager, max_steps: int = 10 # 防止无限循环 ): self.llm = llm_client self.tool_executor = tool_executor self.state_manager = state_manager self.max_steps = max_steps async def run(self, session_id: str, user_input: str) -> str: """运行一次Agent循环""" # 1. 加载或创建状态 state = await self.state_manager.get_state(session_id) if not state: state = AgentState(session_id=session_id, task_id=f"task_{int(datetime.utcnow().timestamp())}") state.conversation_history.append({"role": "user", "content": user_input}) state.current_goal = user_input state.status = TaskStatus.RUNNING step_count = 0 final_answer = None while step_count < self.max_steps and state.status == TaskStatus.RUNNING: step_count += 1 # 2. 准备LLM的上下文(包含对话历史、工具schema、之前的工具调用结果) messages = self._prepare_messages(state) tools = self.tool_executor.registry.get_openai_functions() # 3. 调用LLM,获取决策(思考+行动) llm_response = await self.llm.chat.completions.create( model="gpt-4", # 或你使用的模型 messages=messages, tools=tools, tool_choice="auto", # 让模型决定是否调用工具 temperature=0.1, # 低温度,让输出更确定 ) message = llm_response.choices[0].message state.conversation_history.append(message.model_dump()) # 4. 处理LLM响应 if message.tool_calls: # LLM决定调用工具 for tool_call in message.tool_calls: tool_name = tool_call.function.name try: import json arguments = json.loads(tool_call.function.arguments) except json.JSONDecodeError: arguments = {} # 执行工具 success, result, log = await self.tool_executor.execute( tool_name, arguments, state ) # 将工具执行结果作为新的消息追加到历史 tool_result_msg = { "role": "tool", "tool_call_id": tool_call.id, "content": str(result) if success else f"Error: {result}", "name": tool_name, } state.conversation_history.append(tool_result_msg) # 保存状态(每次工具调用后都保存,保证状态持久化) await self.state_manager.save_state(state) if not success: # 工具执行失败,可以决定让Agent继续尝试或终止 # 这里简单处理:将错误信息反馈给LLM,让它决定下一步 pass else: # LLM给出了最终答案 final_answer = message.content state.status = TaskStatus.SUCCESS state.conversation_history.append({"role": "assistant", "content": final_answer}) break # 循环结束 if not final_answer and step_count >= self.max_steps: final_answer = "抱歉,经过多次尝试仍未能完成任务。可能是问题太复杂或工具暂时不可用。" state.status = TaskStatus.FAILED state.updated_at = datetime.utcnow() await self.state_manager.save_state(state) return final_answer or "未生成回答。" def _prepare_messages(self, state: AgentState) -> List[Dict[str, Any]]: """构建LLM的对话上下文""" messages = [] # 可以添加系统提示词,定义Agent的角色和行为约束 system_prompt = """你是一个有帮助的AI助手,可以调用工具来获取信息或执行操作。 请逐步思考,如果需要,就调用合适的工具。工具调用结果会提供给你。 请用中文回复用户。""" messages.append({"role": "system", "content": system_prompt}) # 添加上下文历史(可以截断或总结,避免超出Token限制) # 这里简单添加全部历史,生产环境需要做Token管理和历史总结 messages.extend(state.conversation_history[-20:]) # 限制最近20轮 return messages这个核心循环的要点:
- 状态驱动:每一步都依赖和更新状态。
- 工具调用集成:LLM通过Function Calling格式决定调用哪个工具。
- 循环与终止:通过
max_steps防止Agent陷入死循环。 - 持久化点:在关键步骤(如工具调用后、最终回答后)保存状态,保证中断后可恢复。
5. 核心实现三:异常处理、自愈与监控
一个健壮的系统必须能妥善处理失败。Agent的异常来源多样:LLM API错误、工具执行异常、网络问题、无效输入等。
5.1 分层异常处理策略
我们需要一个分层的异常处理框架:
class AgentError(Exception): """Agent基础异常""" pass class LLMError(AgentError): """LLM服务相关错误""" pass class ToolExecutionError(AgentError): """工具执行错误""" pass class StateError(AgentError): """状态管理错误""" pass class AgentRuntime: def __init__(self, agent: ReActAgent, retry_policy: Dict[str, Any]): self.agent = agent self.retry_policy = retry_policy # 配置重试策略 async def process_request(self, session_id: str, user_input: str) -> Dict[str, Any]: """处理用户请求,包含完整的异常处理""" start_time = datetime.utcnow() result = {"success": False, "response": None, "error": None, "session_id": session_id} try: # 输入验证与清理 cleaned_input = self._sanitize_input(user_input) # 带重试的Agent执行 response = await self._execute_with_retry(session_id, cleaned_input) result["success"] = True result["response"] = response result["processing_time"] = (datetime.utcnow() - start_time).total_seconds() except LLMError as e: # LLM错误:可能是额度不足、模型过载、内容过滤 result["error"] = f"智能服务暂时不可用: {e}" # 可以触发降级策略,如切换到备用模型或返回缓存答案 await self._trigger_fallback(session_id, user_input, result) except ToolExecutionError as e: # 工具错误:可能是外部API失败、参数错误 result["error"] = f"执行操作时出错: {e}" # 可以尝试替代工具或提示用户提供更多信息 except StateError as e: # 状态错误:并发冲突、存储失败 result["error"] = "会话状态异常,请稍后重试。" # 可能需要清理或重置该会话的状态 except asyncio.TimeoutError: result["error"] = "请求处理超时,请简化您的问题或稍后再试。" except Exception as e: # 未知异常 result["error"] = "系统内部错误,请稍后再试。" # 记录详细日志到监控系统 self._log_critical_error(session_id, e, traceback.format_exc()) # 无论成功失败,记录本次请求的指标 await self._record_metrics(result, start_time) return result async def _execute_with_retry(self, session_id: str, input_text: str, max_retries: int = 2) -> str: """带重试的Agent执行""" last_exception = None for attempt in range(max_retries + 1): try: return await self.agent.run(session_id, input_text) except (LLMError, ToolExecutionError) as e: last_exception = e if attempt == max_retries: raise # 根据错误类型决定等待时间(指数退避) wait_time = (2 ** attempt) + (random.random() * 0.5) await asyncio.sleep(wait_time) # 可以在这里根据异常类型进行一些恢复操作,如重置部分状态 raise last_exception def _sanitize_input(self, text: str) -> str: """简单的输入清理,防止注入攻击""" # 移除过长的输入 if len(text) > 2000: text = text[:2000] + "...[已截断]" # 这里可以添加更多安全检查,如敏感词过滤、特殊字符检查等 return text.strip()5.2 可观测性:日志、指标与追踪
没有可观测性,线上问题就是“黑盒”。我们需要三个维度的数据:
- 日志(Logging):记录离散事件。使用结构化日志(如JSON格式),便于检索和分析。
import structlog logger = structlog.get_logger() # 在关键位置记录 await logger.info("agent_tool_called", session_id=session_id, tool_name=tool_name, arguments=arguments, duration=duration_ms, success=success )- 指标(Metrics):聚合性能数据。使用Prometheus等工具。
from prometheus_client import Counter, Histogram, Gauge AGENT_REQUESTS_TOTAL = Counter('agent_requests_total', 'Total agent requests', ['status']) AGENT_PROCESSING_TIME = Histogram('agent_processing_seconds', 'Request processing time') LLM_TOKEN_USAGE = Counter('llm_token_usage_total', 'Total tokens used', ['type']) # prompt, completion # 在请求处理中记录 AGENT_REQUESTS_TOTAL.labels(status='success').inc() AGENT_PROCESSING_TIME.observe(processing_time)- 分布式追踪(Tracing):跟踪一个请求在微服务或复杂调用链中的完整路径。使用OpenTelemetry。
from opentelemetry import trace tracer = trace.get_tracer(__name__) async def run_agent(session_id, input_text): with tracer.start_as_current_span("agent.run") as span: span.set_attribute("session_id", session_id) span.set_attribute("input.length", len(input_text)) # ... 在LLM调用、工具调用处创建子span监控看板应包含的关键指标:
- 请求量 & 成功率:总请求数、成功/失败率、按错误类型分类。
- 延迟:P50、P95、P99处理时间。
- LLM相关:Token消耗(分prompt/completion)、API调用次数与错误率、成本估算。
- 工具相关:各工具调用次数、平均耗时、错误率。
- 业务相关:会话平均轮次、任务完成率、用户满意度(如有评分)。
5.3 自愈与降级策略
当某些组件故障时,系统应能优雅降级,而不是完全崩溃。
- LLM降级:当主LLM(如GPT-4)不可用或响应慢时,自动切换到备用模型(如GPT-3.5-Turbo、或本地部署的模型)。可以在配置中定义降级链。
- 工具降级:当某个关键工具(如搜索)失败时,可以尝试使用缓存的结果,或者用其他工具组合来近似实现功能,甚至提示用户“该功能暂不可用,但您可以...”。
- 限流与熔断:对LLM API和关键外部工具接口实施限流(rate limiting)和熔断(circuit breaker),防止雪崩效应。例如,使用
pybreaker库。 - 会话恢复:当检测到状态异常(如版本冲突)时,可以尝试从最近的检查点恢复,或者引导用户开始一个新的会话。
6. 部署与运维考量
将上述所有组件组合起来,我们就得到了一个具备产品级雏形的Agent系统。最后,谈谈部署和运维。
6.1 配置管理
所有可变参数(如API密钥、模型名称、超时时间、重试次数)必须外部化配置。推荐使用环境变量或配置文件(如YAML),并区分开发、测试、生产环境。
# config/production.yaml agent: max_steps: 15 default_model: "gpt-4" fallback_model: "gpt-3.5-turbo" temperature: 0.1 tools: search_web: api_key: ${SEARCH_API_KEY} timeout: 10 state: backend: "redis" redis_url: ${REDIS_URL} ttl_hours: 24 monitoring: metrics_port: 9090 log_level: "INFO"6.2 容器化与编排
使用Docker将Agent服务及其依赖(如Python环境)打包。使用Docker Compose(开发)或Kubernetes(生产)进行编排。
# Dockerfile FROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY . . CMD ["python", "-m", "uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]6.3 健康检查与就绪探针
在K8s中,必须配置健康检查端点。
# app/health.py from fastapi import APIRouter, Depends from redis import Redis router = APIRouter() @router.get("/health") async def health_check(redis: Redis = Depends(get_redis)): """健康检查:检查核心依赖(如Redis、数据库)""" try: # 检查Redis连接 await redis.ping() return {"status": "healthy", "timestamp": datetime.utcnow().isoformat()} except Exception as e: raise HTTPException(status_code=503, detail=f"Service unhealthy: {e}")6.4 持续集成与持续部署(CI/CD)
- 代码检查:使用 black, isort, mypy, flake8 确保代码质量。
- 单元测试与集成测试:对工具、状态管理器、Agent核心逻辑进行测试。模拟LLM响应可以使用
unittest.mock。 - 安全扫描:在CI流水线中加入依赖漏洞扫描(如
safety,trivy)。 - 自动化部署:使用GitLab CI/CD、GitHub Actions或Jenkins,实现测试通过后自动部署到相应环境。
构建产品级Agent是一个系统工程,远不止是调通一个API。它要求我们在追求智能的同时,用工程化的思维去约束和保障这份智能。从状态管理、工具编排到异常处理与监控,每一层设计都在为系统的稳定性、可维护性和可扩展性添砖加瓦。这套“Harness”可能初期会带来一些开发复杂度,但它能让你在凌晨三点被报警电话叫醒时,能快速定位问题;在业务量翻十倍时,系统依然坚挺;在需要增加一个新工具或修改决策逻辑时,能够从容不迫。这才是将AI能力转化为实际生产力的关键。