MCP 工具:把别人做好的工具"远程接进来"
前四天的工具都是"本地函数"。但真实世界里,"操作 Notion""读文件系统""查 GitHub"这些能力,往往是别人做成独立服务发布的——这就是 MCP(Model Context Protocol)。今天读 mcp/ 目录(config.py/tool_resolver.py)和 tools/mcp_tool_wrapper.py/mcp_native_tool.py:MCP 服务器怎么用三种传输方式配置、远程工具怎么被"发现"并包装成 Day 27 的 BaseTool、以及连接超时/指数退避重试/并行连接隔离这些"网络编程的脏活"怎么处理。
connect)、问清单(list_tools)、把清单上每项服务做成一张本地"工牌"(包装成 BaseTool),这样 Agent 用它们就跟用本地工具一模一样——完全不知道背后是远程调用。而打电话会遇到占线、断线,所以要有超时和重拨(重试)。痛点:好工具凭什么每个框架都重写一遍
MCPToolResolver 干:连接服务器 → 发现工具清单 → 把每个远程工具包装成本地 BaseTool。包装后,Agent 执行循环(Day 08)根本分不清它是本地还是远程——因为都是 BaseTool 接口(这正是 Day 27 立那个基类的回报)。| 文件 | 职责 |
|---|---|
| mcp/config.py | 三种传输方式的配置模型(Stdio/HTTP/SSE) |
| mcp/tool_resolver.py | 连接、发现、包装的总协调 MCPToolResolver |
| mcp/client.py | 底层 MCP 客户端(连接/列工具/调工具) |
| tools/mcp_native_tool.py | 包装成 BaseTool:每次调用新建连接(并行安全) |
| tools/mcp_tool_wrapper.py | 另一种包装:带超时 + 指数退避重试 |
三种传输:本地进程 / HTTP / SSE
MCP 服务器可以跑在不同地方,config.py 用三个模型描述(config.py:12 / config.py:53 / config.py:90):
# config.py:12 ——本地进程,标准输入输出通信
class MCPServerStdio(BaseModel):
command: str # 'python' / 'node' / 'npx' / 'uvx'
args: list[str] = Field(default_factory=list) # ['server.py'] 等
env: dict[str, str] | None = None # 传给子进程的环境变量
tool_filter: ToolFilter | None = None
cache_tools_list: bool = False # 缓存工具清单,加速后续访问
# config.py:53 ——远程 HTTP(可流式)
class MCPServerHTTP(BaseModel):
url: str # 'https://api.example.com/mcp'
headers: dict[str, str] | None = None # 认证头等
streamable: bool = True
tool_filter: ToolFilter | None = None
cache_tools_list: bool = False
# config.py:90 ——远程 SSE(服务器推送事件)
class MCPServerSSE(BaseModel):
url: str
headers: dict[str, str] | None = None
tool_filter: ToolFilter | None = None
cache_tools_list: bool = False
# config.py:123
MCPServerConfig = MCPServerStdio | MCPServerHTTP | MCPServerSSE
Stdio(本地进程)MCP 服务器是本机的一个子进程,通过标准输入输出通信。像 npx some-mcp-server。适合本地工具(文件系统、本地数据库)。HTTP(远程)服务器在远端,走 HTTP。headers 放认证信息。适合托管的云服务。SSE(远程流式)Server-Sent Events,服务器能持续推送。适合需要实时流的场景。三者共有的字段tool_filter(L07 只暴露部分工具)和 cache_tools_list(缓存"有哪些工具"的清单,省得每次都重新发现)。MCPServerConfig 联合类型三种配置用 | 合成一个类型,后面 resolve 靠 isinstance 分流处理。resolve:一个入口,三种来源分流
MCPToolResolver.resolve(tool_resolver.py:68)是总入口,按引用类型分流:
# tool_resolver.py:68
def resolve(self, mcps: list[str | MCPServerConfig]) -> list[BaseTool]:
all_tools: list[BaseTool] = []
amp_refs: list[tuple[str, str | None]] = []
for mcp_config in mcps:
if isinstance(mcp_config, str) and mcp_config.startswith("https://"):
all_tools.extend(self._resolve_external(mcp_config)) # ① 直接 URL
elif isinstance(mcp_config, str):
amp_refs.append(self._parse_amp_ref(mcp_config)) # ② AMP 短引用("notion")
else:
tools, clients = self._resolve_native(mcp_config) # ③ 配置对象
all_tools.extend(tools)
self._clients.extend(clients)
if amp_refs:
tools, clients = self._resolve_amp(amp_refs) # 批量解析 AMP
all_tools.extend(tools)
self._clients.extend(clients)
return all_tools
三种来源用户可以传:① 完整 https URL(外部服务器);② 短名字如 "notion"(CrewAI AMP 市场的引用);③ 上一节的配置对象。resolve 用 isinstance 和前缀判断分别处理。amp_refs 先收集后批量AMP 短引用先攒进列表,最后 _resolve_amp 一次性批量解析——因为同一个服务器可能被多次引用,批量能只连一次(见 tool_resolver.py:129 的去重)。self._clients.extend把创建的客户端连接登记到 resolver 自己身上。为什么要记?因为用完要断开(L07 的 cleanup)——谁创建谁负责回收。返回 list[BaseTool]★不管哪种来源,最终都返回 BaseTool 列表。Agent 拿到手,用起来和本地工具毫无区别。if is_mcp,越来越乱。源码做法:把 MCP 工具包装成标准 BaseTool,让它对执行系统"伪装"成本地工具。这就是 Day 27 定基类、Day 28 归一执行态的深远回报——新增一种工具来源(今天是 MCP,明天可能是别的协议),执行循环一行都不用改。代价是包装层要"翻译"远程调用的差异(超时、重试、连接管理),但这些被封在包装类里,不外泄。这是"面向接口"带来的扩展性红利。发现工具:问清单,逐个包装
_resolve_native(tool_resolver.py:313)先连服务器列出工具,再把每个包装成 MCPNativeTool:
# tool_resolver.py:333(发现)
async def _setup_client_and_list_tools() -> list[dict]:
...
tools_list = await discovery_client.list_tools() # ★问服务器:你有哪些工具
...
# tool_resolver.py:411(为每次调用准备"造客户端"的闭包)
def _client_factory() -> MCPClient:
return MCPClient(transport=..., cache_tools_list=mcp_config.cache_tools_list, ...)
# tool_resolver.py:419(逐个包装)
for tool_def in tools_list:
tool_name = tool_def.get("name", "")
original_tool_name = tool_def.get("original_name", tool_name)
if not tool_name:
continue
args_schema = None
if tool_def.get("inputSchema"):
try:
args_schema = self._json_schema_to_pydantic( # JSON schema → Pydantic 模型
tool_name, tool_def["inputSchema"])
except Exception as e:
self._logger.log("warning", f"Failed to build args schema ... {e}")
args_schema = None # 建不出 schema 也不放弃
tool_schema = {"description": tool_def.get("description", ""), "args_schema": args_schema}
native_tool = MCPNativeTool( # ★包装成 BaseTool
client_factory=_client_factory, tool_name=tool_name,
tool_schema=tool_schema, server_name=server_name,
original_tool_name=original_tool_name)
tools.append(native_tool)
list_tools()连上服务器后问它"你有哪些工具",返回每个工具的名字、描述、输入 schema(JSON Schema 格式)。_json_schema_to_pydantic★把 MCP 的 JSON Schema 翻译成 CrewAI 认的 Pydantic 模型(args_schema)。这样远程工具的参数也能被 Day 28 的 _parse_args 校验。schema 建失败也继续★健壮性:某个工具的 schema 翻译失败,只记 warning 并把 args_schema 设 None,不影响其他工具。一个坏工具不拖垮整批。client_factory 闭包★不是把一个"连好的客户端"传给工具,而是传一个"能造新客户端的函数"。为什么?见 L05——每次调用要新连接,才能并行安全。MCPNativeTool:每次调用都开一条新连接
MCPNativeTool._run(mcp_native_tool.py:73)要在"可能已有事件循环"的环境里跑异步代码:
# mcp_native_tool.py:73
def _run(self, **kwargs) -> str:
try:
try:
asyncio.get_running_loop() # 当前已在事件循环里?
import concurrent.futures
ctx = contextvars.copy_context() # 复制上下文
with concurrent.futures.ThreadPoolExecutor() as executor:
coro = self._run_async(**kwargs)
future = executor.submit(ctx.run, asyncio.run, coro) # ★另开线程跑
return future.result()
except RuntimeError:
return asyncio.run(self._run_async(**kwargs)) # 没有事件循环 → 直接跑
except Exception as e:
raise RuntimeError(f"Error executing MCP tool {self.original_tool_name}: {e!s}") from e
而 _run_async(mcp_native_tool.py:101)每次都新建一个客户端:
# mcp_native_tool.py:101
async def _run_async(self, **kwargs) -> str:
client = self._client_factory() # ★每次调用造一个全新客户端
await client.connect()
try:
result = await client.call_tool(self.original_tool_name, kwargs)
finally:
await client.disconnect() # 用完必断
if isinstance(result, str):
return result
if hasattr(result, "content") and result.content: # 提取文本内容
if isinstance(result.content, list) and len(result.content) > 0:
content_item = result.content[0]
if hasattr(content_item, "text"):
return str(content_item.text)
return str(content_item)
return str(result.content)
return str(result)
get_running_loop 判断回应 Day 28 的坑:在已有事件循环里不能直接 asyncio.run。所以先探测——有循环就另开线程跑 asyncio.run,没有就直接跑。copy_context()另开线程时复制当前上下文变量,保证子线程也能读到正确的 Agent/Task 上下文(和 Day 08 并行工具同款手法)。client_factory() 每次新建★核心:每次调用造一个全新客户端 + 连接,用完就断。绝不复用连接。try/finally disconnect无论调用成功还是抛错,finally 都保证断开——不泄露连接。提取 content[0].textMCP 返回的是结构化的 content 列表,这里把第一个文本内容抠出来当结果字符串。anyio 的 cancel-scope 错误、状态错乱。源码做法:每次调用新建独立客户端(docstring 明说 "creates a fresh client per invocation … never share mutable connection state")。代价是每次调用多一次建连开销;收益是并行绝对安全——哪怕同一个工具被同时调多次也不会串。对并发场景,正确性压倒性能。(如果确实想复用/带重试,用 L06 的 MCPToolWrapper。)MCPToolWrapper:超时 + 指数退避重试
另一种包装 MCPToolWrapper(mcp_tool_wrapper.py:16)带更完整的容错。先看常量和重试主循环(mcp_tool_wrapper.py:91):
# mcp_tool_wrapper.py:10
MCP_CONNECTION_TIMEOUT = 15
MCP_TOOL_EXECUTION_TIMEOUT = 60
MCP_MAX_RETRIES = 3
# mcp_tool_wrapper.py:91
async def _retry_with_exponential_backoff(self, operation_func, **kwargs) -> str:
last_error = None
for attempt in range(MCP_MAX_RETRIES):
result, error, should_retry = await self._execute_single_attempt(operation_func, **kwargs)
if result is not None:
return result # 成功 → 返回
if not should_retry:
return error # 不可重试的错 → 直接返回错误
last_error = error
if attempt < MCP_MAX_RETRIES - 1:
wait_time = 2**attempt # ★指数退避:1s, 2s, 4s...
await asyncio.sleep(wait_time)
return f"MCP tool execution failed after {MCP_MAX_RETRIES} attempts: {last_error}"
关键在 _execute_single_attempt(mcp_tool_wrapper.py:117)对错误分类——哪些该重试、哪些不该:
# mcp_tool_wrapper.py:117
async def _execute_single_attempt(self, operation_func, **kwargs):
try:
result = await operation_func(**kwargs)
return result, "", False
except ImportError:
return None, "MCP library not available. Please install with: pip install mcp", False
except asyncio.TimeoutError:
return None, f"Connection timed out after {MCP_TOOL_EXECUTION_TIMEOUT} seconds", True # 可重试
except Exception as e:
error_str = str(e).lower()
if "authentication" in error_str or "unauthorized" in error_str:
return None, f"Authentication failed for MCP server: {e!s}", False # 不可重试
if "not found" in error_str:
return None, f"Tool '{self.original_tool_name}' not found on MCP server", False
if "connection" in error_str or "network" in error_str:
return None, f"Network connection failed: {e!s}", True # 可重试
if "json" in error_str or "parsing" in error_str:
return None, f"Server response parsing error: {e!s}", True # 可重试
return None, f"MCP execution error: {e!s}", False
三个超时常量连接 15s、执行 60s、发现 15s。远程调用必须设超时,否则一个卡死的服务器能拖垮整个 Agent。2**attempt 指数退避★重试间隔 1→2→4 秒递增。为什么递增?服务器可能是临时过载,立刻猛重试只会雪上加霜;退避给它喘息时间。错误分类的第三个返回值★should_retry:超时、网络错、解析错 → True(重试有意义);认证失败、工具不存在、缺库 → False(重试一万次也没用)。和 Day 08 "哪些错自愈哪些抛出"是同一个判断标准。为何避免"循环内 try-except"docstring 提到把 try 抽到 _execute_single_attempt 里、循环体保持干净——性能和可读性的小优化。MCPNativeTool(L05)主打并行安全(每次新连接),用于 _resolve_native 路径;MCPToolWrapper(本节)主打网络容错(超时+指数退避+错误分类)。它们体现了同一件事的两个关注点:远程工具既要"并发不打架",又要"网络抖动能扛住"。工具过滤与连接清理
不是每个服务器的工具都想全暴露给 Agent。tool_filter 支持静态和动态两种(filters.py:17 / filters.py:38):
# filters.py:17
class ToolFilterContext(BaseModel):
agent: Any = Field(..., description="The agent requesting tools.")
server_name: str = Field(..., description="Name of the MCP server.")
run_context: dict[str, Any] | None = None
# filters.py:32
ToolFilter = (
Callable[[ToolFilterContext, dict[str, Any]], bool] # 动态:带上下文判断
| Callable[[dict[str, Any]], bool]) # 静态:只看工具本身
# filters.py:38
class StaticToolFilter: # 简单的 allow/block 名单
def __init__(self, allowed_tool_names=None, blocked_tool_names=None): ...
过滤应用在发现循环里(tool_resolver.py:383),而连接用完由 cleanup 统一断开(tool_resolver.py:90):
# tool_resolver.py:90
def cleanup(self) -> None:
if not self._clients:
return
async def _disconnect_all() -> None:
for client in self._clients:
if client and hasattr(client, "connected") and client.connected:
await client.disconnect() # 逐个断开
try:
asyncio.run(_disconnect_all())
except Exception as e:
self._logger.log("error", f"Error during MCP client cleanup: {e}")
finally:
self._clients.clear() # 清空登记表
静态 vs 动态过滤静态:一份"允许/禁止"名单(StaticToolFilter)。动态:一个函数,能根据"哪个 Agent、什么运行上下文"实时决定放不放某工具(拿到 ToolFilterContext)。为什么要过滤安全和聚焦:一个文件系统 MCP 可能同时暴露"读"和"删",你只想让 Agent 读,就 block 掉删除类工具。也减少无关工具占 prompt。cleanup 逐个断开resolver 登记过所有连接(L03),任务结束时 cleanup 遍历断开。finally: clear() 保证即使断开报错,登记表也清空,不留悬空引用。谁创建谁回收连接是 resolver 创建的,就由 resolver 负责回收——资源管理的基本纪律,防连接泄露。边界 + 今日小结
MCPNativeTool 为并行安全每次新连接,高频调用时延迟叠加明显——用 cache_tools_list 至少能省"重复发现工具清单"的开销。② 超时不是万能:执行超时 60s 意味着一个慢工具能卡住一个 Agent 步整整一分钟;MCP 服务器该自己控制单次响应时间。③ schema 可能建不出来:远程工具的 JSON Schema 复杂时翻译成 Pydantic 会失败,此时 args_schema=None,参数就失去校验,模型乱传参也不会被拦——L04 那句 warning 就是提醒你注意这种工具。反模式:把 MCP 工具当本地函数一样无脑高频调用,忽视它背后是网络请求。👶 小白:我在代码里怎么给 Agent 加 MCP 工具?
👨🏫 老师:给 Agent 传 mcps=[...](可以是配置对象、https URL 或 AMP 短名)。框架会自动用 MCPToolResolver.resolve 把它们变成工具、混进 Agent 的 tools 里。执行时 Agent 完全按普通工具用(走 Day 29 的选择、Day 30 的缓存都照样生效,因为它们都是 BaseTool)。用完框架会调 cleanup 断连。你几乎不用手动碰这些底层类。
🧠 今天你应该能回答
- MCP 解决什么问题?为什么说它是"工具界的 USB"?
- 三种传输(Stdio/HTTP/SSE)各适合什么场景?
- 远程工具是怎么被"发现"并变成
BaseTool的?归一成 BaseTool 有什么好处? MCPNativeTool为什么每次调用都新建连接?解决什么并发问题?- 指数退避重试是什么?哪些错该重试、哪些不该?
tool_filter静态和动态的区别?为什么要过滤?- 为什么连接要由 resolver 的
cleanup统一回收?
✋ 10 分钟动手
P=lib/crewai/src/crewai
sed -n '1,123p' $P/mcp/config.py # 三种传输配置
sed -n '68,106p' $P/mcp/tool_resolver.py # resolve + cleanup
sed -n '419,464p' $P/mcp/tool_resolver.py # 发现并包装
sed -n '73,133p' $P/tools/mcp_native_tool.py # 每次新连接
sed -n '91,155p' $P/tools/mcp_tool_wrapper.py # 超时 + 指数退避
# 概念自测:为什么 2**attempt 而不是固定 3 秒?(服务器过载时给它退避喘息)