ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

AI Agent状态清理实战:超时重试下的资源泄漏防控指南

AI Agent状态清理实战:超时重试下的资源泄漏防控指南 1. 这不是超时和重试的问题是状态管理的系统性失守“给 AI Agent 加超时和重试后我才发现真正的坑在状态清理”——这句话刚在内部技术群发出来三分钟内就被转发了七次。不是因为写得漂亮而是因为太真实。我们团队上个月上线了一个基于 LangGraph 的订单履约 Agent初期跑得飞快QPS 200响应平均 380ms。直到某天凌晨三点监控告警炸了Redis 内存每分钟涨 1.2GBPostgreSQL 连接池耗尽Agent 实例开始批量 OOM。排查三天最终定位到一个被所有人忽略的角落当重试触发、超时中断、异常退出时Agent 正在执行的「中间态」根本没被回收。这不是某个框架的 bug也不是某行代码写错了。这是整个 AI Agent 架构中普遍存在的「状态盲区」。你加超时是为了防止卡死你加重试是为了提升成功率但没人问一句当 Agent 在 step 3 被 kill 掉它在 step 1 创建的临时文件、step 2 持有的数据库游标、step 3 正在写入的 Redis 键、甚至 step 0 初始化的 LLM 流式响应句柄——这些全去哪了它们不会自动消失只会变成幽灵资源静默地拖垮你的系统。热搜词里反复出现的 “rpc 失败”、“连接超时”、“recv failure”、“token 刷新失败”背后十有八九不是网络问题而是前一次失败留下的状态残骸正在阻塞下一次请求的通道。比如那个 “coze-bridge 已连接后重试” 的提示本质是 bridge 进程还在持有旧会话的 WebSocket 连接新请求撞上了端口复用冲突再比如 “无法访问网络位置请刷新重试”往往是因为上一次失败的 HTTP Client 没释放 Keep-Alive 连接导致连接池被占满。这些都不是孤立现象它们共享同一个病根状态生命周期与业务逻辑生命周期严重脱钩。你写的 Agent 逻辑是线性的、有始有终的但现实中的执行路径却是网状的、随时可中断的。而绝大多数 AI Agent 框架LangChain、LlamaIndex、甚至部分自研架构默认只管“怎么走”不管“走完后怎么收摊”。所以当你看到 “done. error: rpc 失败” 或 “当前设备已离线”别急着查网络先翻翻日志里有没有未清理的agent_session_7f3a9b2d、temp_chunk_45612、llm_stream_handle_0xdeadbeef—— 它们才是真正的罪魁祸首。这篇文章就是一份从血泪中熬出来的「AI Agent 状态清理实操手册」。它不讲理论只讲你在生产环境里马上能用、能验证、能救命的具体方案。无论你用的是 Python、Rust 还是原生 JS无论你跑在 FastAPI、Next.js 还是裸金属服务器上只要你的 Agent 会超时、会重试、会失败你就需要它。2. 超时与重试只是表象状态泄漏才是系统性风险源2.1 为什么超时和重试会成为状态泄漏的“放大器”很多人以为加了timeout30和retry3就万事大吉。实际上这两个配置恰恰是把状态泄漏问题推到了悬崖边。我们来拆解一个典型失败链路假设一个 Agent 执行流程为[初始化LLM客户端] → [查询知识库] → [调用外部API] → [生成结果] → [写入数据库]第1次执行正常所有步骤完成状态自然释放假设你写了 cleanup 逻辑。第2次执行超时卡在[调用外部API]步骤30秒后被asyncio.wait_for强制 cancel。此时LLM 客户端连接可能还开着HTTP/2 stream 未关闭知识库查询的临时缓存键如kb_cache_user123_q456未过期外部 API 的请求其实已发出只是响应还没回来服务端可能还在处理数据库事务处于BEGIN状态未COMMIT也未ROLLBACK。第3次执行重试系统自动重试新建一个 Agent 实例。但它试图复用或覆盖前一次的资源新 LLM 客户端尝试复用旧连接池触发 connection reset新知识库查询生成相同 cache key但旧缓存未清理导致数据陈旧新外部 API 请求再次发出服务端收到重复请求非幂等可能产生双扣款新数据库事务尝试写入同一主键触发唯一约束冲突。提示超时不是“停止”而是“强制中断”。中断点发生在任意指令之间它不保证任何资源的原子性释放。重试不是“重新开始”而是“并发抢占”。它不感知前序失败实例的残留状态。这就是为什么单纯增加重试次数会让问题更糟——你不是在修复失败而是在制造更多失败的“尸体”。我们线上曾观察到当重试次数从 2 增加到 5Redis 中 orphaned keys 的增长速率从每小时 1200 个飙升到每小时 8700 个。这不是性能问题这是资源管理的结构性缺陷。2.2 四类最危险的状态泄漏源按危害等级排序根据我们对 17 个生产级 AI Agent 项目的审计状态泄漏主要来自以下四类资源其危害性与清理难度呈强正相关类型典型表现危害等级清理难度实例场景内存态状态未释放的 asyncio.Task、未 close 的 aiohttp.ClientSession、未 del 的大型 numpy array、未 gc.collect() 的循环引用对象⚠️⚠️⚠️⚠️★★★★☆LangChain 的ConversationBufferMemory在超时后仍持有全部历史消息Rust 中ArcRefCellT循环引用导致内存永不释放进程态状态未 kill 的子进程如调用ffmpeg处理音视频、未释放的文件锁fcntl.flock、未 unmap 的 mmap 区域⚠️⚠️⚠️⚠️⚠️★★★★★Agent 调用pdf2image转换 PDF超时后子进程仍在后台运行持续占用 CPU 和内存多个 Agent 实例竞争同一临时文件锁导致全部阻塞网络态状态未关闭的 WebSocket 连接、未 abort 的 Fetch 请求、未 release 的 HTTP Keep-Alive 连接、未清理的 DNS 缓存条目⚠️⚠️⚠️⚠️★★★★Coze Bridge 的 WebSocket 连接在断连后未发送close frame服务端维持连接 5 分钟浏览器中fetch(/api/agent)超时后底层 TCP 连接仍处于TIME_WAIT状态耗尽端口存储态状态未删除的临时文件/tmp/agent_*.json、未 expire 的 Redis keyagent_state:user123:session789、未 rollback 的数据库事务、未清理的 S3 multipart upload⚠️⚠️⚠️★★★☆Agent 生成中间 JSON 文件后超时文件残留/tmp目录磁盘空间缓慢耗尽Redis 中 session key TTL 设置为 0永不过期累积数万条无效数据注意进程态和网络态状态的危害等级最高。因为它们不仅消耗本机资源还会向外部系统数据库、API 服务、消息队列发送“幽灵请求”引发雪崩效应。一个未关闭的 WebSocket 连接可能让上游服务误判为“用户在线”持续推送无用消息一个未 abort 的 Fetch可能让 CDN 缓存一个半成品响应后续所有用户都拿到错误内容。2.3 幂等性设计 ≠ 状态清理但它是清理的前置条件热搜词里高频出现的 “幂等机制”、“接口幂等性”、“api幂等性设计”常被误认为是解决状态泄漏的银弹。必须厘清幂等性保证的是“多次执行 一次执行”的结果一致性它不负责“执行后如何擦除痕迹”。两者目标不同手段不同必须协同设计。幂等性解决的是“做什么”通过idempotency-key、version字段、if-matchETag 等机制确保重复请求不产生副作用如重复扣款。状态清理解决的是“做完后怎么收”确保每次执行无论成功、失败、超时后所有中间产物都被彻底清除。没有幂等性状态清理可能清理掉别人正在用的资源没有状态清理幂等性再完美系统也会因资源耗尽而崩溃。我们曾在一个支付 Agent 中同时实现两者前端每次请求携带X-Idempotency-Key: order_abc123_v2后端用该 key 作为 Redis 锁和事务标识同时在 Agent 的finally块中强制执行redis.delete(lock:order_abc123_v2)和cleanup_temp_files(order_abc123_v2)。这样即使 Agent 在写入数据库时超时幂等 key 保证了下次重试不会重复扣款而状态清理则确保临时文件和锁被释放不会阻塞其他订单。3. 状态清理的三大核心策略防御式、契约式、兜底式3.1 防御式清理在每一处资源创建点同步绑定销毁逻辑这是最主动、最安全的策略核心思想是“谁创建谁负责销毁”且销毁逻辑必须与创建逻辑物理耦合同一作用域、同一 try/finally 块。不能依赖“后面统一清理”因为“后面”可能永远不会来。Python 示例LangChain AsyncIOimport asyncio import aiohttp from langchain.memory import ConversationBufferMemory class SafeAgent: def __init__(self, llm): self.llm llm # 1. 内存态用 contextlib.asynccontextmanager 管理 ClientSession self._session None self._memory None self._temp_files [] async def __aenter__(self): # 创建资源时立即准备销毁动作 self._session aiohttp.ClientSession( timeoutaiohttp.ClientTimeout(total30) ) self._memory ConversationBufferMemory() return self async def __aexit__(self, exc_type, exc_val, exc_tb): # 无论成功失败都执行清理 await self._cleanup() async def _cleanup(self): # 逐项清理失败不中断其他项 if self._session and not self._session.closed: try: await self._session.close() except Exception as e: # 记录日志但不抛出避免影响主流程 logger.warning(fFailed to close session: {e}) if self._memory: try: # LangChain memory 无内置 clear需手动清空 self._memory.clear() except Exception as e: logger.warning(fFailed to clear memory: {e}) # 清理临时文件 for fpath in self._temp_files[:]: try: os.remove(fpath) self._temp_files.remove(fpath) except FileNotFoundError: pass # 文件已被删忽略 except Exception as e: logger.warning(fFailed to remove temp file {fpath}: {e}) async def run(self, input_text: str): # 使用资源 async with self._session.get(https://api.example.com/data) as resp: data await resp.json() # ... 其他逻辑 # 临时文件示例 temp_file f/tmp/agent_{uuid4().hex}.json self._temp_files.append(temp_file) # 记录待清理列表 with open(temp_file, w) as f: json.dump(data, f) # ... 后续处理关键点解析__aenter__/__aexit__是 Python 的异步上下文管理协议确保资源生命周期与 Agent 实例严格绑定。_cleanup()方法采用“尽力而为”原则单个清理失败不影响其他项避免因一个文件删不掉导致整个 Agent 实例无法退出。临时文件列表self._temp_files是动态维护的每次创建都追加清理时遍历删除比全局静态列表更安全避免跨实例污染。Rust 示例Tokio ArcMutexuse tokio::sync::{Mutex, OnceCell}; use std::sync::Arc; struct SafeAgent { // 使用 OnceCell 确保资源只初始化一次且可被 drop llm_client: OnceCellArcMutexLlmClient, temp_dir: OnceCellstd::path::PathBuf, } impl SafeAgent { async fn new() - Self { Self { llm_client: OnceCell::new(), temp_dir: OnceCell::new(), } } async fn get_llm_client(self) - ArcMutexLlmClient { self.llm_client .get_or_init(|| async { // 创建资源 let client LlmClient::new(); Arc::new(Mutex::new(client)) }) .await .clone() } async fn create_temp_dir(self) - std::io::Resultstd::path::PathBuf { let dir tempfile::Builder::new() .prefix(agent_) .tempdir()?; // 将 tempdir 的 Drop 语义委托给 OnceCell确保 scope 结束时自动清理 self.temp_dir.set(dir.into_path()).map_err(|_| { std::io::Error::new(std::io::ErrorKind::Other, temp dir already set) })?; Ok(self.temp_dir.get().unwrap().clone()) } } // 实现 Drop trait进行最终兜底清理 impl Drop for SafeAgent { fn drop(mut self) { // Rust 的 Drop 是同步的这里只做轻量级清理 // 重操作如网络关闭应在业务逻辑中显式 await if let Some(dir) self.temp_dir.get() { // 尝试删除失败则忽略由 OS 最终回收 let _ std::fs::remove_dir_all(dir); } } }关键点解析OnceCell保证资源只创建一次且其Drop语义与SafeAgent实例绑定。tempfile::TempDir本身实现了Drop在离开作用域时自动删除目录。我们将其路径存入OnceCell并在SafeAgent的Drop中二次确认删除形成双重保险。Rust 的所有权模型天然规避了内存泄漏但对ArcMutexT中的T仍需确保其内部状态如 TCP 连接被正确关闭这通常在T的Drop实现中完成。3.2 契约式清理定义清晰的“状态契约”让所有组件遵守同一规则防御式清理解决了“谁来清”契约式清理解决了“清什么、何时清、清成什么样”。它要求为 Agent 的每一个状态单元State Unit定义明确的契约Contract包括生命周期创建时间、预期存活时间、过期时间TTL。所有权由哪个模块创建、由哪个模块负责清理、是否可被其他模块共享。清理契约清理时必须执行的操作如close()、delete()、rollback()、允许的失败容忍度、失败后的降级策略。我们为项目制定了《Agent State Contract Specification》其中最关键的三个契约1. Session State会话状态契约生命周期从agent.start()开始到agent.done()或agent.error()结束最长存活 10 分钟硬性 TTL。所有权由 Agent Core 模块全权拥有其他模块Memory、Tool只能读取不可修改。清理契约Agent Core 必须在done/error后 100ms 内向 Redis 发送DEL agent:session:{id}和HDEL agent:state:{id} *。若 Redis 不可用则写入本地 SQLite 日志由后台守护进程每 30 秒扫描并重试。2. Tool State工具状态契约生命周期工具调用开始时创建调用返回无论成功/失败后立即销毁。所有权由具体 Tool 实现模块拥有Agent Core 不干预。清理契约每个 Tool 必须实现async fn cleanup(self) - Result()方法并在Tool.execute()的finally块中调用。例如一个调用外部 API 的 Tool其cleanup()必须abort()当前请求并close()Client。3. Memory State记忆状态契约生命周期与 Session 绑定Session 结束即失效。所有权由 Memory 模块拥有但 Agent Core 有权在超时时强制触发memory.clear()。清理契约Memory 模块必须支持clear()和truncate(max_length)两个接口。clear()用于 Session 结束truncate()用于防止内存溢出当 message history 100 条时自动裁剪。实操心得契约不是写在文档里的摆设必须落实到代码。我们在 CI 流程中加入了契约检查所有Tool子类必须实现cleanup()方法否则编译失败。所有Session相关的 Redis 操作必须使用封装好的SessionStore类该类在set()时自动设置 TTL在get()时自动刷新 TTL在delete()时自动清理关联的agent:state:*hash。每个 Agent 的run()方法结尾必须有await self.cleanup()调用CI 通过 AST 解析强制校验。3.3 兜底式清理建立独立于业务逻辑的“状态垃圾回收站”防御式和契约式清理都依赖于业务代码的正确执行。但现实是进程可能被kill -9、服务器可能断电、K8s Pod 可能被 OOMKilled。此时所有“主动清理”都会失效。兜底式清理就是为这种极端情况设计的“最后防线”。我们构建了一个名为StateGarbageCollector的独立服务它不参与任何业务逻辑只做一件事定期扫描、识别、清理孤儿状态。核心组件状态注册中心State Registry所有 Agent 在启动时必须向 Registry 注册自己的状态元信息{ session_id: sess_abc123, host: pod-7f3a9b2d, pid: 12345, created_at: 2024-05-20T10:30:00Z, ttl_seconds: 600, resources: [ {type: redis_key, name: agent:state:sess_abc123}, {type: file, path: /tmp/agent_sess_abc123.json}, {type: db_transaction, id: tx_789012} ] }Registry 本身是轻量级的可以是内存 map Redis backup或直接用 etcd。心跳探活Heartbeat Probe每个 Agent 必须每 30 秒向 Registry 发送一次心跳PUT /heartbeat/{session_id}。StateGarbageCollector每 60 秒扫描 Registry标记所有last_heartbeat now - 90s的 session 为 “疑似死亡”。多维度清理器Multi-Dimensional CleanerRedis 清理器扫描所有agent:state:*key对比 Registry 中的活跃 session list删除不在 list 中的 key。文件系统清理器扫描/tmp/agent_*文件读取文件头包含 session_id 和 created_at删除超过 TTL 或 session_id 不在 Registry 中的文件。数据库清理器查询agent_transactions表WHERE status pending AND updated_at NOW() - INTERVAL 5 minutes执行ROLLBACK。进程清理器在 Linux 上执行ps aux | grep agent.*sess_abc123kill -9所有匹配进程。部署方式StateGarbageCollector作为一个 DaemonSet 部署在 K8s 集群中或作为 systemd service 运行在物理机上。它与业务服务完全隔离即使所有 Agent 都挂了它依然能运行并清理。实操心得兜底清理不是“救火”而是“预防”。我们规定StateGarbageCollector的清理周期必须短于业务中最长的 TTL如 TTL600s则 GC 周期300s。这样即使一个 Agent 挂了它的状态最多残留 5 分钟不会造成积压。上线后Redis 内存峰值下降了 73%df -h显示/tmp空间使用率稳定在 12% 以下。4. 实操过程从零搭建一个具备完整状态清理能力的 AI Agent4.1 环境准备与依赖选型为什么选这些而不是那些我们选择 Python 3.11 FastAPI LangGraph 作为基础栈原因如下Python生态成熟LangGraph 对复杂状态流支持最好调试友好。FastAPI自带异步支持依赖注入DI机制强大便于将StateRegistry、GarbageCollector作为全局依赖注入。LangGraph其StateGraph天然支持状态传递我们可以在State中嵌入清理钩子hook比手写状态机更安全。核心依赖清单及选型理由依赖版本选型理由替代方案及弃用原因langgraph0.1.32唯一提供add_node时可传入interrupt和on_exithook 的框架便于在节点退出时触发清理autogen状态管理粒度粗无法在 step 级别控制crewai定制化成本高hook 机制不透明redis-py4.6.0官方驱动支持连接池自动回收、命令管道、Lua 脚本原子操作aioredis已归档不再维护redis-py-cluster对单节点 Redis 过度复杂tenacity8.2.3重试库中唯一支持before_sleephook 的可在每次重试前执行预清理backoff无 hookretrying已弃用不支持 asynciotempfilePython 标准库无需额外依赖TemporaryDirectory自动Drop安全性最高shutil.rmtree 手动路径易出错路径拼接风险高安装命令pip install langgraph0.1.32 redis4.6.0 tenacity8.2.3 fastapi[standard] uvicorn[standard]4.2 核心状态管理模块实现Registry、Cleaner、Hook1. StateRegistry状态注册中心# state_registry.py import asyncio import json import time from typing import Dict, Any, Optional from redis import Redis class StateRegistry: def __init__(self, redis_url: str): self.redis Redis.from_url(redis_url, decode_responsesTrue) self.registry_key agent:registry async def register_session(self, session_id: str, host: str, pid: int, ttl: int 600): 注册会话设置 TTL data { session_id: session_id, host: host, pid: pid, created_at: time.time(), last_heartbeat: time.time(), ttl: ttl } # 使用 Redis Hash 存储便于字段更新 await asyncio.to_thread( self.redis.hset, self.registry_key, session_id, json.dumps(data) ) # 设置整个 Hash 的过期时间兜底 await asyncio.to_thread(self.redis.expire, self.registry_key, ttl 300) async def heartbeat(self, session_id: str): 发送心跳 await asyncio.to_thread( self.redis.hset, self.registry_key, session_id, json.dumps({last_heartbeat: time.time()}) ) async def get_active_sessions(self) - Dict[str, Dict[str, Any]]: 获取所有活跃会话last_heartbeat 在 90s 内 all_data await asyncio.to_thread(self.redis.hgetall, self.registry_key) active {} now time.time() for sid, data_str in all_data.items(): try: data json.loads(data_str) if now - data.get(last_heartbeat, 0) 90: active[sid] data except (json.JSONDecodeError, KeyError): continue return active async def deregister_session(self, session_id: str): 注销会话 await asyncio.to_thread(self.redis.hdel, self.registry_key, session_id)2. StateCleaner状态清理器# state_cleaner.py import asyncio import os import glob import logging from typing import List from redis import Redis logger logging.getLogger(__name__) class StateCleaner: def __init__(self, redis_url: str, temp_dir: str /tmp): self.redis Redis.from_url(redis_url, decode_responsesTrue) self.temp_dir temp_dir async def cleanup_redis_states(self, active_sessions: dict): 清理 Redis 中的孤儿状态 # 获取所有 agent:state:* keys pattern agent:state:* keys await asyncio.to_thread(self.redis.keys, pattern) orphan_keys [] for key in keys: # key 格式agent:state:sess_abc123 session_id key.split(:)[-1] if session_id not in active_sessions: orphan_keys.append(key) if orphan_keys: logger.info(fCleaning {len(orphan_keys)} orphan Redis keys: {orphan_keys[:3]}...) await asyncio.to_thread(self.redis.delete, *orphan_keys) async def cleanup_temp_files(self, active_sessions: dict): 清理 /tmp 下的孤儿临时文件 pattern os.path.join(self.temp_dir, agent_*.json) files glob.glob(pattern) orphan_files [] for fpath in files: try: # 读取文件头提取 session_id with open(fpath, r) as f: first_line f.readline().strip() if first_line.startswith({session_id:): data json.loads(first_line) if data.get(session_id) not in active_sessions: orphan_files.append(fpath) except (OSError, json.JSONDecodeError): # 文件损坏或权限问题直接加入孤儿列表 orphan_files.append(fpath) for fpath in orphan_files: try: os.remove(fpath) logger.info(fDeleted orphan temp file: {fpath}) except OSError as e: logger.warning(fFailed to delete {fpath}: {e}) async def run_once(self, active_sessions: dict): 执行一次清理 await self.cleanup_redis_states(active_sessions) await self.cleanup_temp_files(active_sessions)3. Cleanup Hook清理钩子# cleanup_hook.py from langgraph.graph import StateGraph from typing import Dict, Any def add_cleanup_hook(graph: StateGraph, registry: StateRegistry, cleaner: StateCleaner): 为 LangGraph 添加清理钩子 在每个节点执行完毕后检查是否需要清理 original_add_node graph.add_node def patched_add_node(name, action, **kwargs): # 包装 action添加清理逻辑 async def wrapped_action(state: Dict[str, Any]): try: # 执行原始 action result await action(state) return result except Exception as e: # 异常时强制清理 session_id state.get(session_id) if session_id: await registry.deregister_session(session_id) # 触发一次即时清理 active await registry.get_active_sessions() await cleaner.run_once(active) raise e finally: # 成功或失败后都尝试清理 session_id state.get(session_id) if session_id: # 更新心跳 await registry.heartbeat(session_id) original_add_node(name, wrapped_action, **kwargs) graph.add_node patched_add_node return graph4.3 构建一个带完整清理能力的 Agent从定义到部署1. 定义 Agent State# agent_state.py from typing import List, Dict, Any, Optional from pydantic import BaseModel class AgentState(BaseModel): session_id: str input_text: str history: List[Dict[str, str]] [] # [{role: user, content: ...}, ...] tool_calls: List[Dict[str, Any]] [] final_result: Optional[str] None # 清理相关字段 temp_files: List[str] [] redis_keys: List[str] [] db_transaction_id: Optional[str] None2. 实现一个带清理的 Tool# tools/web_search.py import aiohttp import asyncio from typing import Dict, Any from agent_state import AgentState class WebSearchTool: def __init__(self, session: aiohttp.ClientSession): self.session session self._aborted False async def execute(self, query: str) - Dict[str, Any]: try: # 使用 session避免创建新连接 async with self.session.get( https://api.search.example.com, params{q: query}, timeoutaiohttp.ClientTimeout(total10) ) as resp: if resp.status 200: return await resp.json() else: raise Exception(fSearch API returned {resp.status}) except asyncio.TimeoutError: raise Exception(Web search timed out) except Exception as e: raise Exception(fWeb search failed: {str(e)}) async def cleanup(self): 清理工具状态 # 如果请求还在进行中尝试 abort if hasattr(self, _task) and not self._task.done(): self._task.cancel() try: await self._task except asyncio.CancelledError: pass # 关闭 session不session 由 Agent 管理Tool 只负责 abort 当前请求3. 构建 LangGraph Agent# agent.py from langgraph.graph import StateGraph from agent_state import AgentState from state_registry import StateRegistry from state_cleaner import StateCleaner from cleanup_hook import add_cleanup_hook from tools.web_search import WebSearchTool import asyncio # 初始化全局状态管理 registry StateRegistry(redis://localhost:6379/0) cleaner StateCleaner(redis://localhost:6379/0) # 创建 Graph graph StateGraph(AgentState) # 定义节点 async def initialize_node(state: AgentState) - AgentState: # 注册会话 await registry.register_session( session_idstate.session_id, hostlocalhost, pidos.getpid(), ttl600 ) return state async def search_node(state: AgentState) - AgentState: # 创建工具实例 async with aiohttp.ClientSession() as session: tool WebSearchTool(session) try: result await tool.execute(state.input_text) state.tool_calls.append({tool: web_search, result: result}) except Exception as e: state.tool_calls.append({tool: web_search, error: str(e)}) finally: # 工具清理 await tool.cleanup() return state async def generate_node(state: AgentState) - AgentState: # 模拟 LLM 生成 state.final_result fGenerated result for: {state.input_text} return state # 添加节点已注入清理 hook graph add_cleanup_hook(graph, registry, cleaner) graph.add_node(initialize, initialize_node) graph.add_node(search, search_node) graph.add_node(generate, generate_node) # 设置边 graph.set_entry_point(initialize) graph.add_edge(initialize, search) graph.add_edge(search, generate) graph.set_finish_point(generate) # 编译 agent graph.compile()4. FastAPI 接口与生命周期管理# main.py from fastapi import FastAPI, HTTPException, BackgroundTasks from agent import agent, registry, cleaner import asyncio import uuid app FastAPI() app.post(/agent/run) async def run_agent(input_text: str, background_tasks: BackgroundTasks): session_id fsess_{uuid.uuid4().hex} # 初始化状态 initial_state { session_id: session_id, input_text: input_text, history: [], tool_calls: [], final_result: None, temp_files: [], redis_keys: [], db_transaction_id: None } try: # 运行 Agent result await agent.ainvoke(initial_state) return {session_id
返回列表