项目准备
- Java的ERP系统
- 模拟的供应商网站
- 魔塔社区免费的可视化分析工具
- DeepAgents ==0.4.12版本/ 先不要更新到最新的 DeepAgents
- 后台长周期复杂任务的执行(0.5后可以)AsyncSubAgent(子Agent)—-ASGI协议——> 主Agent(langgraph dev或者langsmith的服务器中托管)
- 5月12日-0.6 (上下文优化策略:内置代码解释器:支持TypeScript,下一版本会支持Python语言)
项目目录架构
src/
├── agent/ # Agent 核心
│ ├── main_agent.py # 主 Agent 入口(11阶段初始化流水线)
│ ├── config.py # 全局配置(模型/沙箱/MongoDB/路径/Store)
│ ├── schema.py # 数据结构(ProcurementContext / SSE事件模型 / API模型)
│ ├── env_utils.py # 环境变量读取
│ ├── middleware_config.py # 子 Agent 中间件工厂
│ │
│ ├── memory/
│ │ ├── AGENTS.md # 主 Agent 通用行为准则(沙箱 /AGENTS.md)
│ │ └── prompts.py # 主 Agent system_prompt(精简版协调者角色)
│ │
│ ├── backends/
│ │ ├── sandbox_setup.py # 沙箱创建 + 基础文件播种(seed)
│ │ └── custom_opensandbox.py # OpenSandboxBackend 自定义封装
│ │
│ ├── middlewares/ # 中间件栈(7个,按主 Agent 注入顺序)
│ │ ├── context_injection.py # 1. 注入用户上下文到 system prompt
│ │ ├── skills_sync.py # 2. 本地技能 → 沙箱同步(增量)
│ │ ├── user_skills_restore.py # 3. StoreBackend → 沙箱恢复持久化技能
│ │ ├── tools_summarization.py # 4. SummarizationToolMiddleware 封装
│ │ └── memory_update.py # 5. 自动更新 recent_suppliers / recent_queries
│ │
│ ├── subagents/
│ │ ├── loader.py # YAML 配置加载 + 工具名称匹配解析
│ │ └── configs/
│ │ ├── procurement_analyst.yaml # 分析子 Agent 配置
│ │ └── procurement_order.yaml # 订单子 Agent 配置
│ │
│ └── tools/
│ ├── mcp_client.py # MCP 多服务器连接 + 工具按前缀分组
│ ├── chart_generator.py # 26种图→1个 generate_visualization 入口
│ ├── web_search.py # 智谱搜狗 Web 搜索
│ ├── hitl_tools.py # request_order_info(数据补充中断)
│ ├── assign_skill.py # 技能分配(复制+StoreBackend持久化+清理)
│ └── download_sandbox_file.py # 沙箱文件下载到本地
│
├── skills/ # 技能资源(本地 → 沙箱同步)
│ ├── main/skill-management/ # 技能生命周期管理
│ └── procurement/
│ ├── chart_params.md # 26种图表完整参数参考
│ ├── procurement-analysis/SKILL.md # 分析流程操作手册(5步)
│ ├── supplier-price-urls/ # 供应商报价URL映射表
│ └── web-scraper/ # 网页抓取(直接HTTP+HTML→MD)
│
├── api_view/ # FastAPI Web层
│ ├── web_main.py # FastAPI app 入口(lifespan/路由/CORS)
│ ├── agent_loader.py # Agent 懒加载 + MongoDB 管理
│ ├── web_config.py # API 元信息
│ └── api/
│ ├── chat.py # SSE 流式对话 + 中断检测 + resume 端点
│ └── history.py # 历史会话管理(MongoDB CRUD)
│
├── mcp_server/ # MCP Server(ERP 业务接口代理)
│ ├── server_main.py # MCP 服务入口(端口8000)
│ └── tools/
│ ├── suppliers_tools.py # supplier_query / supplier_*
│ ├── parts_tools.py # part_query / part_search / part_by_supplier
│ ├── order_tools.py # order_create / order_update / order_search_details
│ └── inventory_tools.py # inventory_warning
│
└── test/
└── agent_test.py # 终端测试脚本
Multi-Agent的描述
主 Agent(Orchestrator)
| 项目 |
内容 |
| 模型 |
deepseek-v4-pro (temperature=1.1) |
| 角色 |
协调者 — 理解需求、委派子Agent、管理记忆,不直接调用ERP业务工具 |
| 直接工具 |
用途 |
| web_search |
通用知识问答(行业资讯、技术概念、市场行情) |
| assign_skill |
将已测试技能分配给指定子 Agent |
| read_file / write_file / edit_file |
读写 /memories/{user_id}/preferences.md(由框架内置 FilesystemMiddleware 提供) |
| task |
委派子Agent(SubagentsMiddleware 内置) |
| compact_conversation |
主动压缩上下文(SummarizationMiddleware 内置) |
| 预置技能 |
路径 |
用途 |
| skill-management |
/skills/main/skill-management/ |
技能下载/创建/测试/分配/持久化 全生命周期管理 |
| 中间件 |
功能 |
| ContextInjectionMiddleware |
注入 user_id/username 到 system prompt |
| SkillsSyncMiddleware |
本地 src/skills/ → 沙箱同步 |
| UserSkillsRestoreMiddleware |
StoreBackend 持久化技能 → 沙箱恢复 |
| SummarizationToolMiddleware |
提供 compact_conversation 工具 + 自动摘要 |
| ModelCallLimitMiddleware |
最多 50 次模型调用 |
| ToolCallLimitMiddleware |
最多 200 次工具调用 |
委派规则:
- 触发”分析/对比/报告/建议/评估/行情/比价/筛选” →
procurement-analyst
- 触发”下单/创建订单/修改/更新/取消/订单状态” →
procurement-order
- 问候/知识问答/技能管理 → 主Agent自行处理
procurement-analyst(采购分析子Agent)
| 项目 |
内容 |
| 模型 |
继承主Agent (deepseek-v4-pro) |
| 角色 |
数据收集 → 执行分析 → 生成图表 → 输出报告 |
| MCP 业务工具(ERP查询) |
来源 |
功能 |
| supplier_query |
erp-api |
按名称模糊搜索供应商 → GET /suppliers/search |
| part_query |
erp-api |
分页查询采购零部件(可按名称/分类/供应商ID筛选) |
| part_search |
erp-api |
按名称搜索零部件 → GET /parts/search |
| part_by_supplier |
erp-api |
按供应商ID查其所有零部件 → GET /parts/supplier/{id} |
| order_search_details |
erp-api |
搜索订单明细(含零部件详情+供应商信息) |
| inventory_warning |
erp-api |
查询库存预警列表 |
| 图表工具 |
来源 |
功能 |
| generate_visualization |
魔塔社区 MCP |
26合1可视化入口,通过 chart_type 参数路由到 19种标准图表 + 3种地图 + 3种关系图 + 1种思维导图 |
| 通用工具 |
功能 |
| web_search |
外部市场价格搜索、供应商背景调查、行业资讯 |
| execute / read_file / write_file / ls / glob / grep |
沙箱内 Python 分析、文件操作(框架内置) |
| 预置技能(/skills/procurement/) |
用途 |
| supplier-price-urls |
零部件+供应商名 → 报价URL映射表,外部数据采集第一步 |
| procurement-analysis |
采购分析主流程:5阶段(数据收集→外部数据→分析→图表→报告) |
| web-content-fetcher |
反爬虫备用:jina.ai / markdown.new / defuddle.md 三级方案 |
| 中间件 |
功能 |
| SummarizationToolMiddleware |
阶段完成后主动压缩上下文 |
| ModelCallLimitMiddleware |
最多 50 次模型调用 |
| ToolCallLimitMiddleware |
最多 200 次工具调用 |
分析工作流(5阶段):
1. ls /skills/procurement/ → 扫描可用技能
2. MCP工具获取ERP数据 + supplier-price-urls找外部URL + web-content-fetcher爬取报价
3. Python脚本分析 → /data/analysis_result.json
4. generate_visualization → 2-4张关键图表
5. write_file /analysis/report_*.md → 返回结构化摘要+结论+建议
procurement-order(采购订单子Agent)
| 项目 |
内容 |
| 模型 |
继承主Agent (deepseek-v4-pro) |
| 角色 |
理解订单操作需求 → 调用MCP工具 → 验证结果 → 返回确认 |
| MCP 业务工具 |
来源 |
功能 |
| order_create |
erp-api |
创建采购订单 → POST /orders/create |
| order_update |
erp-api |
更新采购订单 → PUT /orders/update/{order_id} |
| order_search_details |
erp-api |
搜索订单明细(含零部件详情+供应商信息) |
| 通用工具 |
功能 |
| web_search |
验证供应商信息、查询外部参考 |
| read_file / write_file(框架内置) |
辅助文件操作 |
| 预置技能(/skills/order/) |
用途 |
| (暂无预置,留空扩展) |
未来可添加订单模板、审批流程等技能 |
| 中间件 |
功能 |
| ModelCallLimitMiddleware |
最多 20 次模型调用(订单操作短链,限制更严) |
| ToolCallLimitMiddleware |
最多 50 次工具调用 |
核心数据流程
1. 用户输入 "帮我对摩托车火花塞做供应商比价分析"
2. ContextInjectionMiddleware
→ runtime.context 提取 user_id="laoxiao", username="laoxiao"
→ 注入 SystemMessage: "当前用户 user_id: laoxiao, 偏好文件: /memories/laoxiao/preferences.md"
3. MemoryMiddleware
→ 从沙箱 /AGENTS.md 加载全局准则 → 注入 system prompt
4. Agent 推理
→ 读取 /memories/laoxiao/preferences.md (StoreBackend)
→ 判断: "比价分析" → 委派 procurement-analyst
5. task(procurement-analyst) 委派格式:
【任务目标】【用户偏好】【分析需求正文】【输出要求】
6. 子 Agent (procurement-analyst) 执行:
a. ls /skills/procurement/ → 扫描可用技能
b. supplier_query("火花塞") → MCP → Java后端 → ERP数据
c. part_by_supplier(supplier_id) → 按供应商查零部件
d. 技能: 读取 supplier-price-urls/data/url_mapping.yaml
→ 匹配 URL → web-content-fetcher 爬取外部报价
e. Python 分析脚本 → /data/analysis_result.json
f. generate_visualization("bar", chart_config={...}) → 图表URL
g. write_file("/analysis/report_20260509.md") → 最终报告
7. 子 Agent → compact_conversation → 返回结构化结果给主 Agent
8. 主 Agent → 组织用户友好回复 → 更新 /memories/laoxiao/preferences.md (如有新偏好)
存储架构(CompositeBackend 分流)
| 路径 |
后端 |
持久化 |
内容 |
| /AGENTS.md |
OpenSandbox(default路由) |
沙箱生命周期 |
全局行为准则(Phase 1.4 上传) |
| /skills/ |
OpenSandbox(default路由) |
沙箱生命周期 + StoreBackend恢复 |
技能文件(预置 + 持久化混合) |
| /memories/{user_id}/ |
StoreBackend → InMemoryStore |
跨会话 |
用户偏好文件 |
| /persisted-skills/ |
StoreBackend → InMemoryStore |
跨会话 |
持久化技能 |
| 其余路径 |
OpenSandbox(default路由) |
沙箱生命周期 |
临时文件、代码执行、分析产物 |
项目中的自定义中间件
设计原则总结
- 读写分离:ContextInjection(读记忆)↔ MemoryUpdate(写记忆),注入时明确告知 Agent “系统自动维护,你无需手动更新”
- 三层 skills 同步:SkillsSync(本地→沙箱)+ UserSkillsRestore(Store→沙箱)+ SkillsMiddleware(沙箱→Agent 发现),各管一段
- 主/子差异化管理:主 Agent 有全量中间件,analyst 保留摘要,order 只要限制——按需减配
- 静默失败 > 崩溃:ContextInjection 三级降级、SkillsSync 异常捕获、MemoryUpdate 全 try/except——任何中间件挂了都不能拖垮 Agent
before_agent ──────────────────────────────────────► after_agent
│ │
▼ ▼
[1] ContextInjection → 注入用户身份 SystemMessage
[2] SkillsSync → 同步本地 skills → 沙箱
[3] UserSkillsRestore → 恢复持久化 skills → 沙箱
[4] SummarizationTool → 上下文压缩 + compact_conversation 工具
[5] MemoryUpdate → 自动更新用户记忆
[6] ModelCallLimit → 限制 LLM 调用 ≤ 50 次
[7] ToolCallLimit → 限制工具调用 ≤ 200 次
LangChain AgentMiddleware 提供 6 个钩子,项目使用了其中 4 个:
| 钩子 |
使用者 |
语义 |
| before_agent |
ContextInjection, SkillsSync, UserSkillsRestore |
Agent run 开始前,注入信息/准备环境 |
| abefore_agent |
同上(异步版) |
同上,异步流 |
| after_agent |
MemoryUpdate(仅占位) |
Agent run 结束后 |
| aafter_agent |
MemoryUpdate(核心逻辑) |
Agent run 结束后,写记忆 |
| before_model |
SummarizationTool(内部) |
每次 LLM 调用前,检查是否需要压缩 |
| wrap_tool_call |
未使用 |
包装每次工具调用 |
chat.py:282 agent_loader.agent.astream(...)
│ ← astream() 被调用
│ ┌─────────────────────────────────┐
│ │ LangGraph 内部启动: │
│ │ 1. 根据 config 恢复 checkpoint │
│ │ 2. 初始化 state(含历史消息) │
│ │ 3. 构建 runtime.context │
│ │ │
│ │ 4. before_agent 钩子触发 ← 这里! │ ← 在 astream 内部
│ │ ├─ ContextInjection │
│ │ ├─ SkillsSync │
│ │ └─ UserSkillsRestore │
│ │ │
│ │ 5. before_model 钩子触发 │
│ │ 6. LLM 推理 │
│ │ 7. 生成第一个 token │
│ └─────────────────────────────────┘
│ ← 第一条 chunk yield
chat.py:282 async for chunk in ...: ← chat.py 收到第一条数据
after_agent 钩子 ← "结束后"
│ └─ MemoryUpdate 在这里写 Store
│ astream() 生成器退出
chat.py:xxx # async for 循环结束,代码继续往下走 ← chat.py 继续执行
每个中间件的说明
[1] ContextInjectionMiddleware before_agent
| 项目 |
说明 |
| 文件 |
middlewares/context_injection.py |
| 钩子 |
before_agent + abefore_agent |
| 触发时机 |
每次 Agent run 开始时(仅一次) |
| 数据来源 |
runtime.context(由 context_schema=ProcurementContext 从 configurable 映射) |
| 做了什么 |
注入一条 SystemMessage:当前用户 user_id: xxx,偏好文件路径: /memories/{user_id}/preferences.md,请首先 read_file 读取偏好 |
| 返回值 |
{“messages”: [SystemMessage(…)]} 或 None(跳过) |
| 容错 |
三级降级:context 为 None → 跳过;user_id 为空 → 跳过;username 为空 → 用 user_id 兜底 |
[2] SkillsSyncMiddleware before_agent
| 项目 |
说明 |
| 文件 |
middlewares/skills_sync.py |
| 钩子 |
before_agent + abefore_agent(异步版用 run_in_executor 执行同步 I/O) |
| 触发时机 |
每次 Agent run 开始时 |
| 数据来源 |
本地 src/skills/ 目录 vs 沙箱 /skills/ 目录 |
| 做了什么 |
① 遍历本地 skills 目录 ② MD5 比对本地 vs 沙箱文件 ③ 使用 test -f 避免 404 ERROR 日志 ④ 增量上传变更文件 ⑤ 有变更时注入 SystemMessage 通知 Agent |
| 返回值 |
{“messages”: [SystemMessage(“以下技能包已更新: …”)]} 或 None(无变更) |
| 性能优化 |
_last_hashes 字典缓存本地文件 MD5,避免重复计算和上传 |
[3] UserSkillsRestoreMiddleware before_agent
| 项目 |
说明 |
| 文件 |
middlewares/user_skills_restore.py |
| 钩子 |
abefore_agent(仅异步版;同步版直接返回 None) |
| 触发时机 |
每次 Agent run 开始时 |
| 数据来源 |
runtime.store(StoreBackend,内存存储) |
| 做了什么 |
① store.asearch(namespace) 搜索所有持久化 skill ② 映射 key /{scope}/{skill_name}/… → 沙箱路径 /skills/{scope}/{skill_name}/… ③ backend.aupload_files() 恢复到沙箱 |
| 返回值 |
始终 None(不注入消息) |
| 与 [2] 的分工 |
SkillsSync = 预置技能同步(本地→沙箱);UserSkillsRestore = 持久化技能恢复(Store→沙箱) |
| 项目 |
说明 |
| 文件 |
middlewares/tools_summarization.py |
| 来源 |
DeepAgents 内置 create_summarization_tool_middleware() 工厂 |
| 模型 |
SUMMARY_MODEL(deepseek-v4-flash,轻量模型) |
| 做了什么 |
① 内部含一个 SummarizationMiddleware:上下文达到 ~85% 上限时自动压缩历史消息 ② 额外注入 compact_conversation 工具:Agent 可在关键节点(如收到子 Agent 报告后)主动压缩 ③ 被压缩的完整对话持久化到 CompositeBackend |
| 返回值 |
内部处理,不显式返回 |
[5] MemoryUpdateMiddleware after_agent
| 项目 |
说明 |
| 文件 |
middlewares/memory_update.py |
| 钩子 |
aafter_agent(仅异步版;同步版直接返回 None) |
| 触发时机 |
Agent run 完成后 |
| 数据来源 |
state[“messages”] + runtime.context + runtime.store |
| 做了什么 |
详见下方流程图 |
| 返回值 |
始终 None(副作用写入 Store,不修改消息) |
MemoryUpdate 内部流程:
after_agent 触发
│
├─ ① 从 runtime.context 获取 user_id
│
├─ ② _is_meaningful_erp_exchange(messages)
│ ├─ 找最后一条 HumanMessage
│ ├─ 跳过闲聊(_SKIP_PATTERNS:你好/在吗/你是谁...)
│ ├─ 检查是否包含 ERP 关键词(供应商/采购/订单/报价...)
│ └─ 兜底:检查是否有 task(子 Agent)工具调用
│
├─ ③ _extract_ai_summary(messages)
│ 取最后一条 AIMessage 前 300 字符
│
├─ ④ _extract_entities(model, user_msg, ai_summary)
│ LLM 提取 JSON: {"suppliers": [...], "query": "..."}
│
├─ ⑤ store.aget(namespace, key)
│ 读取 /memories/{user_id}/preferences.md 当前内容
│
├─ ⑥ _merge_preferences(current_lines, new_suppliers, new_query)
│ 解析旧的 recent_suppliers / recent_queries 区块
│ → 移除旧区块 → 合并新旧 → 去重 → 截断(suppliers ≤ 10, queries ≤ 5)
│ → 追加到文件末尾
│
└─ ⑦ store.aput(namespace, key, file_value)
写回 StoreBackend
[6] ModelCallLimitMiddleware
| 项目 |
说明 |
| 来源 |
langchain.agents.middleware.ModelCallLimitMiddleware |
| 限制 |
run_limit=50(最多 50 次 LLM 调用) |
| 作用 |
防止 Agent 陷入无限推理循环,超限后抛出异常终止 |
| 项目 |
说明 |
| 来源 |
langchain.agents.middleware.ToolCallLimitMiddleware |
| 限制 |
run_limit=200(最多 200 次工具调用) |
| 作用 |
防止工具调用爆炸(如循环调用 web_search),超限后终止 |
Agent懒加载的设计
为什么需要懒加载
create_main_agent() 内部是 11 个阶段的初始化流水线:
沙箱创建 → MCP 工具加载 → 26→1 工具合并 → 子Agent YAML加载 →
中间件构建 → create_deep_agent() 编译 StateGraph
这个流程涉及网络 I/O(MCP 连接、沙箱 API)和计算(LangGraph 编译),可能耗时 5-15 秒。如果放在模块 import 阶段执行,会导致:
import agent.main_agent 卡死 10 秒
- 模块级别错误无法被 try/except 捕获
- 测试/CLI 每次启动都要等,即使不调用 Agent
懒加载的核心思想:模块导入时只创建一个空壳,真正使用时才初始化。
实际流程时序
FastAPI 启动
│
├─ lifespan 事件触发
│ agent_loader.initialize()
│ from agent.main_agent import get_agent_async ← 导入模块
│ │ ← _AgentProxy() 创建(空壳)
│ │ ← 无事件循环,什么都不发生
│ │
│ agent = await get_agent_async() ← 显式异步初始化
│ ├─ _is_initialized? → False
│ ├─ await _create_agent() ← 11 个阶段,5-15 秒
│ │ ├─ Phase 1: setup_sandbox()
│ │ ├─ Phase 2: load_mcp_tools() ← 连接 MCP 服务器
│ │ ├─ Phase 3: create_generate_chart_tool()
│ │ ├─ ...
│ │ └─ Phase 9: create_deep_agent() ← 编译 StateGraph
│ └─ agent._agent = graph
│
├─ HTTP 请求到达
│ agent.astream(...) ← _AgentProxy._agent 已有值
│ ← __getattr__ → 直接返回 graph
│ ← 不经过 _ensure_initialized 检测
_AgentProxy 的三层设计
flowchart TD
%% 定义组件样式
classDef proxy fill:#f5f5f5,stroke:#333,stroke-width:3px,color:#333;
classDef branch fill:#e1f5fe,stroke:#0277bd,stroke-width:2px,color:#0277bd;
classDef internal fill:#fff3e0,stroke:#ef6c00,stroke-width:2px,stroke-dasharray: 5 5;
classDef result fill:#e8f5e9,stroke:#2e7d32,stroke-width:2px;
subgraph EntryPoint ["代理类结构"]
Proxy["<b>_AgentProxy()</b><br/>self._agent = None<br/><i>(导入时:空壳,0 开销)</i>"]:::proxy
end
subgraph AccessModes ["接入模式"]
direction LR
getattr["<b>__getattr__('astream')</b><br/>(透明属性访问)"]:::branch
async_init["<b>get_agent_async()</b><br/>(显式异步初始化)"]:::branch
end
subgraph ExecutionLogic ["初始化逻辑"]
ensure["_ensure_initialized()"]:::internal
create["_create_agent()"]:::internal
run_sync["asyncio.run()"]:::internal
run_async["await _create_agent()"]:::internal
end
%% 连接关系
Proxy --> getattr
Proxy --> async_init
%% 路径映射
getattr --> ensure --> run_sync --> create
async_init --> run_async --> create
%% 最终结果
create --> Result["<b>self._agent = graph</b>"]:::result
%% 布局优化
Proxy --- ExecutionLogic
getattr 透明代理(agent.astream() 直接用)
# main_agent.py:297
def __getattr__(self, name):
return getattr(self._ensure_initialized(), name)
任何对 agent.xxx 的访问(.astream(), .ainvoke(), .aget_state() 等)都先触发 _ensure_initialized(),然后委托给真实的编译后 graph 对象。调用方不需要感知代理的存在:
# agent_test.py:7
from agent.main_agent import agent # ← 空壳,不阻塞
# agent_test.py:74 —— 首次访问触发初始化
async for chunk in agent.astream(...): # ← __getattr__ → 初始化 → 委托
_ensure_initialized() 事件循环检测
# main_agent.py:272-295
def _ensure_initialized(self):
if self._agent is not None:
return self._agent # 已初始化,直接返回
try:
loop = asyncio.get_event_loop()
if loop.is_running(): # ← 检测:当前在事件循环中?
raise RuntimeError( # 是 → 无法用 asyncio.run(),报错
"Agent 尚未初始化且当前在事件循环中,"
"请使用 await get_agent_async() 获取 agent"
)
except RuntimeError as e:
if "Agent 尚未初始化" in str(e):
raise # ← 我们自己的错,往上抛
# asyncio.get_event_loop() 抛的错(无事件循环)→ 可以安全用 asyncio.run()
self._agent = asyncio.run(_create_agent()) # ← 同步包装异步创建
return self._agent
这里用了一个双重 RuntimeError 技巧:
| 场景 |
asyncio.get_event_loop() 行为 |
结果 |
| 不在事件循环中(模块顶层、普通脚本) |
抛出 RuntimeError(“no current event loop”) |
被 except 捕获 → 放行 → asyncio.run() |
| 在运行中的事件循环(FastAPI 内部) |
返回运行中的 loop → loop.is_running() = True |
抛出我们的 RuntimeError → 被重新抛出 |
用代码路径表示:
场景 A: agent_test.py 控制台
─────────────────────────────
import agent.main_agent → _AgentProxy() 创建(agent._agent = None)
...其他代码...
asyncio.run(main()) → 事件循环启动
main() 内:
agent.astream() → __getattr__ → _ensure_initialized()
├─ _agent 不是 None?→ 不是,继续
├─ get_event_loop() → RuntimeError("no current event loop") ← 已在事件循环中!
│ 等,这不对...
等一下——这里有一个实际使用中的矛盾。 agent_test.py 在 asyncio.run(main()) 内调用 agent.astream(),此时事件循环正在运行。_ensure_initialized() 会检测到运行中的循环并抛出 RuntimeError。所以 agent_test.py 不能直接用 getattr 路径——它必须在进入事件循环前调用 get_agent() 完成初始化。
get_agent() / get_agent_async() 显式初始化
# main_agent.py:310-325 —— 同步路径
def get_agent():
global agent
if isinstance(agent, _AgentProxy):
if agent._is_initialized:
return agent._agent
return agent._ensure_initialized() # 必须在事件循环外调用
return agent
# main_agent.py:328-344 —— 异步路径
async def get_agent_async():
global agent
if isinstance(agent, _AgentProxy):
if agent._is_initialized:
return agent._agent
agent._agent = await _create_agent() # 在事件循环内 await
return agent._agent
return agent
设计优势
- 模块导入零开销:
import agent.main_agent 瞬间完成,不阻塞
- 透明代理:使用者写
agent.astream() 和直接持有 graph 对象写 graph.astream() 完全一样
- 双路径兼容:同步测试脚本用
get_agent(),异步服务用 get_agent_async(),同一个 agent 全局变量
- 只初始化一次:
self._agent 一旦赋值,后续调用全部跳过
MCP服务器中26+个工具如何合并
它解决什么问题
魔塔社区 MCP Server 暴露了 26 个独立的可视化工具:
generate_bar_chart, generate_line_chart, generate_pie_chart,
generate_scatter_chart, generate_radar_chart, generate_sankey_chart,
generate_area_chart, generate_column_chart, generate_boxplot_chart,
generate_violin_chart, generate_histogram_chart, generate_funnel_chart,
generate_treemap_chart, generate_word_cloud_chart, generate_waterfall_chart,
generate_dual_axes_chart, generate_venn_chart, generate_liquid_chart,
generate_organization_chart, generate_mind_map, generate_fishbone_diagram,
generate_flow_diagram, generate_network_graph,
generate_district_map, generate_pin_map, generate_path_map
如果全部直接注入 Agent,工具列表描述会占用大量上下文,每次 LLM 调用都要遍历 26 个工具的 schema。chart_generator.py 把这些合并成 1 个入口:generate_visualization(chart_type, chart_config),通过 chart_type 字符串路由到实际的 MCP 工具。
工具描述的 token 预算策略
合并后 generate_visualization 的工具描述分两部分:
第一部分:工具本体(~50 tokens)
# 第 267-277 行
"Generate a data visualization and return the image URL.\n"
"Use this tool to create charts, maps, diagrams, or mind maps for analysis reports.\n"
f"Available chart_type values: {chart_list}\n"
+ _COMPACT_TABLE # ~800 tokens
+ "\nWhen unsure about the exact fields for chart_config, read "
+ "`/skills/procurement/chart_params.md` first, then call this tool."
第二部分:_COMPACT_TABLE 速查表(~800 tokens)
按 6 种数据模式分组,每组列出 chart_type → data 格式 → 特有参数:
category-value 模式 ─→ bar, column, pie, funnel, treemap, word_cloud
time-value 模式 ─→ line, area
分布/统计模式 ─→ boxplot, violin, histogram, scatter
多维度/流向/集合 ─→ radar, sankey, venn, waterfall, dual_axes
特殊图表 ─→ liquid, organization, mind_map, fishbone_diagram, flow_diagram, network_graph
地图 ─→ district_map, pin_map, path_map
速查表只写 data 格式和特有参数(如 liquid 的 percent 必填、pie 的 innerRadius)。通用参数(width, height, title, theme, style)只在表头一句话概括,不逐图重复。
第三部分:完整参数文档(沙箱文件)
速查表装不下的细节(如 district_map 的 subdistricts/showAllSubdistricts/dataType 等字段含义、dual_axes 的嵌套 series 结构示例代码)放入沙箱 /skills/procurement/chart_params.md,Agent 不确定时主动 read_file 查看。
设计意图:工具描述是每次 LLM 调用都加载的,必须短(控制在 ~850 tokens)。完整参数参考只在需要时才加载进上下文。这就是”渐进式披露”在工具参数层面的应用。
def create_generate_chart_tool(chart_mcp_tools: list):
"""
Returns:
(generate_visualization, other_tools) 二元组
- generate_visualization: 合并后的可视化入口工具
- other_tools: 未被合并且应保留的独立工具列表
"""
调用方从 load_mcp_tools() 拿到 chart_mcp_tools 后传入,拿到两个产物:
# main_agent.py:158-160
generate_visualization, extra_mcp_tools = create_generate_chart_tool(chart_mcp_tools)
generate_visualization 进入工具池,供主 Agent 和子 Agent 使用
extra_mcp_tools(如 generate_spreadsheet)也加入工具池,作为独立工具
运行时路由:generate_visualization 内部
# 第 325-335 行
@tool
async def generate_visualization(chart_type: str, chart_config: dict) -> str:
chart_tool = tool_map.get(chart_type)
if not chart_tool:
available = ", ".join(sorted(tool_map.keys()))
return f"Error: Unknown chart type '{chart_type}'. Available types: {available}"
result = await chart_tool.ainvoke(chart_config)
return result
逻辑极其简单——就是一个字符串查字典 + 委托调用:
Agent 调用: generate_visualization(chart_type="pie", chart_config={"data": [...], "title": "..."})
│
▼
tool_map["pie"] → generate_pie_chart (原始 MCP StructuredTool)
│
▼
chart_tool.ainvoke({"data": [...], "title": "..."})
│
▼
魔塔社区 MCP Server → 返回图片 URL
一个容易被忽略的细节
函数签名写的是 chart_config: dict,不是一个 Pydantic 模型。这意味着:
- Agent 传入的参数直接透传给 MCP 工具,不做任何校验
- 如果传错了字段名或格式,错误由魔塔社区 MCP 返回,信息不可控
- 这也是为什么
chart_params.md 存在——在没有结构化校验的情况下,只能靠文档引导 Agent 传对参数
当前设计的核心机制
graph TD
%% --- 定义样式 ---
classDef toolStyle fill:#e3f2fd,stroke:#1565c0,stroke-width:1px,rx:5,ry:5;
classDef mergerStyle fill:#fff3e0,stroke:#ef6c00,stroke-width:2px,stroke-dasharray: 5 5;
classDef gatewayStyle fill:#fbe9e7,stroke:#bf360c,stroke-width:2px,rx:10,ry:10;
classDef resourceStyle fill:#e8f5e9,stroke:#2e7d32,stroke-width:1px;
%% --- 区域划分 ---
subgraph Server [魔塔社区 MCP Server]
Tools["<b>26 个工具集合</b><br/>(generate_bar, line, pie, scatter, radar, sankey, ... )"]:::toolStyle
end
%% 合并逻辑
MergeLogic("create_generate_chart_tool()<br/>[转换与合并逻辑]"):::mergerStyle
subgraph EntryPoint [统一路由入口]
Gateway["<b>generate_visualization</b><br/>(chart_type, chart_config)"]:::gatewayStyle
Routing["路由逻辑:<br/>tool_map[chart_type].ainvoke(chart_config)"]:::gatewayStyle
end
subgraph Dependency [运行时执行依赖]
Desc("工具描述 (~800t)<br/>(每次 LLM 调用)"):::resourceStyle
Doc("chart_params.md<br/>(沙箱参考文件)"):::resourceStyle
Invoke("运行时动态路由<br/>(ainvoke)"):::resourceStyle
end
%% --- 逻辑连线 ---
Server -->|注册| MergeLogic
MergeLogic --> Gateway
Gateway --> Routing
Routing -.-> Desc
Routing -.-> Doc
Routing -.-> Invoke
%% 设置统一的背景颜色,可选
style Server fill:#fcfcfc
style EntryPoint fill:#fafafa
style Dependency fill:#fcfcfc
三个关键设计决策:
| 决策 |
做法 |
目的 |
| 工具描述只放速查表 |
_COMPACT_TABLE ~800 tokens |
每次 LLM 调用都加载工具描述,必须短 |
| 完整参数放沙箱文件 |
chart_params.md 手动维护 |
Agent 需要时主动 read_file 查看 |
| 两段式调用 |
先 read_file 参考 → 再 generate_visualization |
需要时才加载详细 schema,不浪费上下文 |
五、Agent 使用时的实际调用流程
Agent 收到任务: "对比博世、电装、德尔福的价格,用柱状图展示"
Step 1: 分析需求 → 柱状图适合类别对比 → chart_type = "bar"
Step 2: 看工具描述中的速查表 → bar 属于 category-value 模式
→ data 格式: [{category, value}]
→ 没有特殊必填参数 → 直接调用,不需要 read_file
Step 3: generate_visualization(
chart_type="bar",
chart_config={
"data": [
{"category": "博世", "value": 45.5},
{"category": "电装", "value": 38.2},
{"category": "德尔福", "value": 52.0}
],
"title": "火花塞供应商价格对比",
"axisYTitle": "单价 (元)"
}
)
Step 4: 返回图片 URL → Agent 展示给用户
Agent 收到任务: "用思维导图梳理采购策略"
Step 1: 分析 → chart_type = "mind_map"
Step 2: 看速查表 → mind_map 属于"特殊图表"分组
→ data: {name, children: [{name, children: [...]}]},最大深度3
→ 格式不直观 → 不确定具体字段 → 先 read_file
Step 3: read_file("/skills/procurement/chart_params.md")
→ 看到完整的 data 结构示例,确认格式
Step 4: generate_visualization("mind_map", {...})
两种路径:简单图表 Agent 可以凭速查表直接调,复杂图表 Agent 自己判断需要参考文件。这个”自己判断”的能力取决于模型对数据结构的推理——速查表的 {name, children: [...]} 语法是否足够清晰。