跳到主要内容

数据截至 (上游 commit 5359534c6f00)

第 4 章 · MCP 层:工具即动作

本章讲 OpenEnv 怎么把 MCP(Model Context Protocol,模型上下文协议——一套让模型发现并调用外部工具的标准)接进 Gym 式的 step()。设计依据是仓库里的 rfcs/003-mcp-support.md


4.1 它要解决的小问题

传统 RL 环境的动作空间是环境作者拍脑袋定的:Discrete(4){"code": str}{"move": "e2e4"}。每个环境一套,模型每换个环境就要重新学「这里的动作长什么样」。

而 LLM 世界已经有了现成的标准答案:MCP 工具。工具有名字、有描述、有 JSON Schema 参数表,模型天生会用。

RFC 003 的判断很直接:列出所有可能动作 = tools/list;执行一个动作 = tools/call 这两件事和 MCP 的语义完全重合,那就别再发明第三套。

于是:

传统 RL OpenEnv + MCP

动作空间是什么? ────────▶ tools/list
执行这个动作 ────────▶ tools/call

4.2 思路:把 MCP 方法包成 Action 类型

映射只用了两个 Pydantic 类(src/openenv/core/env_server/mcp_types.py):

MCP 方法Action 类位置对应 Observation
tools/listListToolsActionmcp_types.py:225ListToolsObservation(:259)
tools/callCallToolActionmcp_types.py:239CallToolObservation(:269)

两个 Action 都带一个 type 字面量字段做判别器("list_tools" / "call_tool")。CallToolAction 另有 tool_name: strarguments: Dict[str, Any]

关键是它们是普通的 Action 子类。所以模型「调工具」这件事,在协议层就是一次普通的 step()

一个细节:反序列化会优先认 MCP 类型

deserialize_action()(src/openenv/core/env_server/serialization.py:42)在交给环境自己的 action_cls 之前,先看 type 字段是不是 list_tools / call_tool:

_MCP_ACTION_TYPES: Dict[str, Type[Action]] = {
"list_tools": ListToolsAction,
"call_tool": CallToolAction,
}

真实源码 serialization.py:20-23。但拦截是有条件的——只有当 action_cls 本身是 Action 基类或某个 MCP 类型时才生效(serialization.py:36)。注释说明了理由:「这让环境自己的动作校验保持权威」。

换句话说:如果你的环境声明了 MyGameAction,那 type: "call_tool" 的载荷不会被偷偷改道,还是走 MyGameAction.model_validate(),该报错就报错。


4.3 MCPEnvironment:把 FastMCP 服务器接进 step()

定义在 src/openenv/core/env_server/mcp_environment.py:133,继承 Environment

环境作者要写什么

envs/echo_env/server/echo_environment.py:71-104 这段真实代码的结构:

  1. __init__ 里建一个 FastMCP("echo_env");
  2. @mcp.tool 装饰几个普通函数;
  3. super().__init__(mcp) 把服务器交给基类;
  4. 实现 reset()state_step_impl()

注意第 4 步是 _step_impl 而不是 step——因为 step 已经被基类占用来做路由了。

step() 的路由

MCPEnvironment.step()(mcp_environment.py:418)只做三岔分流:

if isinstance(action, ListToolsAction):
return self._handle_list_tools()
elif isinstance(action, CallToolAction):
return self._handle_call_tool(action, timeout_s=timeout_s)
else:
return self._step_impl(action, timeout_s=timeout_s, **kwargs)

真实源码 mcp_environment.py:447-452MCP 动作被基类吃掉,其余的落到子类。 纯 MCP 环境(比如 echo)的 _step_impl 就只回一句「不认识这个动作类型」(echo_environment.py:156-163)。

同步与异步:两条镜像路径

这是本模块最容易看晕的地方。基类为 list_tools 和 call_tool 各准备了两个版本:

同步入口异步实现关系
_handle_list_tools(:454)_async_handle_list_tools(:492)同步版用 run_async_safely 包异步版
_handle_call_tool(:468)_async_handle_call_tool(:512)同上

真正的逻辑只在异步版里,同步版是一行转调。而 step_async()(mcp_environment.py:597)直接调异步版,跳过 run_async_safely

为什么要这么绕?step_async 的文档字符串给了答案(mcp_environment.py:606-608):

WebSocket 处理器在外层事件循环上直接调用它,而 MCP 会话已经在那个循环上打开了——这避免了走同步 step() 经由 run_in_executor 时出现的线程/事件循环死锁。

串起来看就是:run_async_safely()(src/openenv/core/utils.py:9)在已有事件循环时会另起一个线程跑 asyncio.run。而 MCP 客户端会话绑在原来那个循环上,跨线程访问就死锁。所以异步路径必须一路 await 到底,不能中途落回同步。

✅ WebSocket 路径(推荐)
websocket_endpoint → step_async → _async_handle_call_tool → await client.call_tool
(全程同一个事件循环,MCP 会话一直开着)

⚠️ 同步路径(仍支持,但要小心)
step → _handle_call_tool → run_async_safely → 新线程 + 新循环 → ...

4.4 mcp_session:一个被反复解释的上下文管理器

mcp_session()(mcp_environment.py:211)只有三行实体代码:

client = self._require_mcp_client()
async with client:
yield client

但它有 30 行文档字符串,讲了两个作用(mcp_environment.py:213-239):

作用一:空值守卫。 close() 之后 mcp_client 会被置 None(mcp_environment.py:653),这里给出清晰报错而不是 AttributeError

作用二(有意思的那个):AsyncExitStack 适配器。 FastMCP 的 Client.__aenter__ 会创建一个后台 asyncio.Task 管理会话。如果在 HTTP 会话路径里直接把它塞进 AsyncExitStack,某些 ASGI 测试工具(比如 Starlette 的 TestClient)会在请求之间取消这个孤儿任务,会话状态就损坏了。

包一层 @asynccontextmanager 生成器就能隔离:生成器帧把 async with client: 挂在 yield,只有当 stack 显式关闭这个生成器时清理才会跑,事件循环取消孤儿任务时不会误伤。

这段注释还顺带说明了两件事:FastMCP 的 Client 上下文管理器是可重入的(内部引用计数,最外层退出才真关);它内部已经有 anyio.Lock 串行化连接状态变更,所以不需要额外加锁。

会话为什么要一直开着

服务端在三处把 MCP 会话钉在连接生命周期上:

位置场景
http_server.py:409-414HTTP MCP 会话创建时进 stack
http_server.py:1179-1185/mcp WebSocket 连接期间
http_server.py:1510-1516/ws WebSocket 连接期间

目的一致:避免每条消息都重建 MCP 传输,同时保住 FastMCP 的 ctx.set_state / ctx.get_state 跨调用有效。


4.5 两条通道:训练用哪条,推理用哪条

这是本章的重点。同一批工具,有两条到达路径。

训练侧 推理侧
│ │
CallToolAction JSON-RPC tools/call
│ │
▼ ▼
/ws {"type":"step"} /mcp POST 或 WebSocket
│ │
▼ ▼
env.step_async() mcp_handler()
│ │
└──────────┬────────────────────┘

同一个 FastMCP 工具函数

差别一览

step 通道/mcp 通道
消息格式OpenEnv 帧(type + data)JSON-RPC 2.0
有没有 reward(Observation 带 reward/done)没有——直接返回工具结果
走不走环境的 step 计数会话有 session_id 时也计数
典型用户RL 训练循环推理服务、agent 框架

mcp_client.py 顶部的架构图注释把这个分工画得很清楚(src/openenv/core/mcp_client.py:13-30)。

/mcp 的会话管理:两个自定义方法

纯 HTTP 的 MCP 是无状态的,但环境有状态。OpenEnv 加了两个非标准的 JSON-RPC 方法补这个洞(http_server.py:760:764):

方法干什么
openenv/session/create建一个持久会话,返回 session_id
openenv/session/close关掉它

之后每次 tools/callparams 里带上 session_id,服务端就能找到那个专属环境实例(http_server.py:852-871)。不带 session_id 的话,服务端现造一个临时环境、用完就扔(http_server.py:872-874:979-980)。

客户端侧的对应实现是 _ensure_production_session()(mcp_client.py:201),它加锁缓存 session_id,只创建一次。

openenv/session/close 有一条保护:不允许关掉当前 WebSocket 正在用的会话(http_server.py:799-804),否则连接会莫名其妙失去环境。


4.6 MCPToolClient:客户端侧的两副面孔

定义在 mcp_client.py:372,继承 MCPClientBase(:82)。

强制 production 模式

构造函数里有个硬校验(mcp_client.py:136-145):mode 缺省是 "production",传别的直接 ValueError,并指路 GenericEnvClient

use_production_mode 开关切换通道

注意这是另一个标志,默认 False(mcp_client.py:157)。它决定 list_toolscall_tool 走哪条路:

if getattr(self, "use_production_mode", False):
# 走 HTTP /mcp,JSON-RPC
...
else:
# 走 step 通道
action = CallToolAction(tool_name=name, arguments=kwargs)
result = await self.step(action)

真实源码见 call_tool(mcp_client.py:447-468)。用 getattr(self, ..., False) 而不是 self.use_production_mode,注释解释是因为有些测试用 __new__ 绕过 __init__(mcp_client.py:242)。

结果解包

call_tool 最后要处理 FastMCP 的 CallToolResult,它可能是对象也可能已经被 JSON 化成 dict,两种都取 data(mcp_client.py:481-488):

if hasattr(result, "data"):
return result.data
if isinstance(result, dict) and "data" in result:
return result["data"]
return result

传输层错误(obs.error 非空)会被抛成 RuntimeError,消息里带上错误类型(mcp_client.py:471-476)。

注意区分两种「错误」——这是 mcp_types.py 里刻意的设计:

错误来源放在哪例子
工具自己返回的失败result 字段里「文件不存在」这种业务结果
传输/框架层失败error 字段(ToolError)工具找不到、超时、参数非法

ToolErrorType 枚举有 5 个值(mcp_types.py:198-205):execution_errorinvalid_argstransport_errortool_not_foundtimeout


4.7 两个进阶特性

模式感知工具

MCPEnvironment.tool() 装饰器(mcp_environment.py:343)接受 mode 参数:

@self.tool(mode="production")
def fetch(url: str) -> str: ... # 真的去抓网页

@self.tool(mode="simulation")
def fetch(url: str) -> str: ... # 返回缓存的固定响应

实现上,mode=None 的工具正常注册进 FastMCP;带 mode 的工具不注册进 FastMCP,而是存进 _mode_tools[name][mode] 自管(mcp_environment.py:378-380)。schema 靠 inspect.signature 手工推导,只认 int/float/bool 三种,其余一律 "string"(mcp_environment.py:390-401)——这是个明显的简化实现。

list_tools 时按当前模式过滤(_schema_for_mode,mcp_environment.py:124),调用时也按模式挑函数,挑不到就返回 TOOL_NOT_FOUND(mcp_environment.py:522-536)。

Code mode(CodeAct)

RFC 003 说要同时支持「一次一个工具调用」和「写一段代码调多个工具」。后者的实现是 execute_code()(mcp_environment.py:290):

namespace = self.get_callables() # 工具函数直接作为 Python 可调用对象
exec(code, namespace, result_dict)
result = result_dict.get("result")

get_callables()(mcp_environment.py:259)从 FastMCP 服务器上把每个工具的 .fn 掏出来,再合并模式感知工具。语法错误和运行时错误分别捕获,包成带 metadata["error"]Observation

这里要诚实说明边界: execute_code 用的是裸 exec(),没有任何沙箱。安全性完全依赖「整个环境跑在 Docker 容器里」这一层——容器就是沙箱边界,exec 本身不是。


4.8 关键细节与坑

  • FastMCP 2.x 和 3.x 的兼容层。 get_server_tools()(mcp_environment.py:87)先试 get_tools()(2.x,返回 dict),再试 list_tools()(3.x,返回 list)。pyproject.toml:32 要求 fastmcp>=3.0.0,但兼容代码保留着。
  • list_tools 失败不抛异常。 _async_handle_list_tools 捕获所有异常,返回空工具列表 + metadata["error"](mcp_environment.py:506-510)。生产环境里要记得检查这个 metadata。
  • 超时默认 30 秒。 MCP_TOOL_CALL_TIMEOUT = 30.0(mcp_environment.py:81),可以用 step(action, timeout_s=...) 覆盖。
  • 错误分类靠字符串匹配。 _async_handle_call_tool"not found" in error_message.lower() 这类判断给异常归类(mcp_environment.py:578-590)。上游 FastMCP 改文案就会失效——这是个脆弱点。
  • JsonRpcResponse.model_dump 被重写了。 为了符合 JSON-RPC「result 和 error 二选一」的规范,它手工构造 dict 而不是走 Pydantic 默认序列化(mcp_types.py:131-144)。
  • close() 只清引用。 MCPEnvironment.close()mcp_clientmcp_serverNone,依赖上下文管理器已经关掉传输(mcp_environment.py:644-654)。

4.9 代码地图

主题文件符号
MCP 动作/观察类型src/openenv/core/env_server/mcp_types.pyListToolsActionCallToolActionCallToolObservation
工具描述与错误src/openenv/core/env_server/mcp_types.pyToolToolErrorToolErrorType
JSON-RPC 类型src/openenv/core/env_server/mcp_types.pyJsonRpcRequestJsonRpcResponseJsonRpcErrorCodeMcpMethod
保留字src/openenv/core/env_server/mcp_types.pyRESERVED_TOOL_NAMES
MCP 环境基类src/openenv/core/env_server/mcp_environment.pyMCPEnvironment
动作路由src/openenv/core/env_server/mcp_environment.pystepstep_async_step_impl
会话保持src/openenv/core/env_server/mcp_environment.pymcp_session
模式感知工具src/openenv/core/env_server/mcp_environment.pytool_schema_for_mode_mode_tools
Code modesrc/openenv/core/env_server/mcp_environment.pyget_callablesexecute_code
FastMCP 版本兼容src/openenv/core/env_server/mcp_environment.pyget_server_tools
服务端 MCP 处理src/openenv/core/env_server/http_server.pymcp_handler_run_mcp_client_operation
MCP 客户端src/openenv/core/mcp_client.pyMCPClientBaseMCPToolClientcall_tool
生产会话src/openenv/core/mcp_client.py_ensure_production_session_production_mcp_url
参考实现envs/echo_env/server/echo_environment.pyEchoEnvironment
设计提案rfcs/003-mcp-support.md