数据截至 (上游 commit 5359534c6f00)
第 3 章 · EnvClient:一份代码,两副面孔
本章讲客户端。核心是
EnvClient(src/openenv/core/env_client.py:238)和它的同步影子SyncEnvClient(src/openenv/core/sync_client.py:43)。这一章有 OpenEnv 里最「取巧」的一段代码。
3.1 它要解决的小问题
RL 训练循环有两派人:
- 一派写
asyncio.gather并发跑 64 条 rollout,他们要async; - 一派在 Jupyter 里手搓一局棋、或者在同步的老训练框架里调用,他们要不带
await的普通函数。
通常的做法是维护两套 API(step 和 astep),或者让用户自己包 asyncio.run。OpenEnv 选了第三条路:只写一份异步实现,让调用点自己表现成两副面孔。
3.2 直觉:先看结果
同一个类,两种写法都合法。下面两段直接抄自 EnvClient 的类文档字符串(env_client.py:259-276)。
# 写法 A:async
from envs.coding_env.client import CodingEnv
async with CodingEnv(base_url="ws://localhost:8000") as env:
result = await env.reset(seed=42)
while not result.done:
action = agent.predict(result.observation)
result = await env.step(action)
# 写法 B:sync —— 靠 .sync() 包装器
env = CodingEnv(base_url="ws://localhost:8000").sync()
with env:
result = env.reset(seed=42)
result = env.step(action)
差别只有一处:写法 B 多了个 .sync(),业务调用一个字没改。(base_url 写 http:// 也行,客户端内部会用 convert_to_ws_url 换成 ws://。)
3.3 机关一:_dispatch 的四路分派
所有公开方法都是一行转发。看 reset/step/state/close(env_client.py:828、683、706、720):
def step(self, action: ActT, **kwargs: Any) -> Any:
return self._dispatch(lambda: self._step_async(action, **kwargs))
注意它是 def 不是 async def——调用它不会产生协程,而是立刻执行 _dispatch。
_dispatch(env_client.py:460)是一串守卫子句,不是并列的三选一。怎么读下图:从上往下,命中哪条 就地返回,不再往下走。
_dispatch(coro_factory)
│
├─① 这个 client 已经锁定成 sync 了吗? (371-377)
│ 是 ─┬─ 当前循环就是 sync 客户端自己的后台循环
│ │ → 直接返回协程(避免自己等自己)
│ └─ 否 → _run_sync(...) 阻塞跑完,返回真结果
│
├─② 有没有正在跑的事件循环? (378-381)
│ 有 ─┬─ 已锁定成 async → 直接返回协程
│ └─ 还没锁定 → 返回 _AutoAsyncResult
│
└─③ 都不是(压根没有事件循环) (383)
→ _run_sync(..., allow_async_handoff=True)
跑到底,返回真结果(允许从 async 交接回 sync)
四个出口对应源码 env_client.py:467-479。第 ① 条是最容易被忽略的一条:一旦客户端锁定成 sync,后续调用根本不看有没有事件循环,而是先问「现在这个循环是不是我自己起的那个后台循环」——是的话继续用协程往下传,否则一律走 _run_sync。锁定机制本身见 3.5。
3.4 机关二:_AutoAsyncResult —— 既能 await、又能当结果
定义在 env_client.py:67。它是个包装对象,同时实现了两组协议:
| 方法 | 触发场景 | 行为 |
|---|---|---|
__await__ | await client.step(...) | 锁定 async 模式,await 真协程(env_client.py:80-85) |
__getattr__ | client.step(...).observation | 阻塞跑完,再取属性(env_client.py:93-94) |
__bool__ | if client.step(...) | 阻塞跑完,取真值(env_client.py:96-97) |
换句话说:这个对象在 你「碰它」之前不承诺自己是什么。你 await 它,它就是协程;你读它的字段,它就同步跑完再给你字段。中间结果用 _resolved 缓存,只跑一次(env_client.py:87-91)。
简化示意
# 示意,非源码 —— 只演示核心想法
class MaybeAsync:
def __init__(self, make_coro):
self._make_coro = make_coro
self._value = _UNSET
def __await__(self): # 被 await:走异步路
return self._make_coro().__await__()
def __getattr__(self, name): # 被读字段:走同步路
if self._value is _UNSET:
self._value = run_blocking(self._make_coro())
return getattr(self._value, name)
重点看:异步路和同步路共用同一个 _make_coro,没有第二份业务实现。
3.5 机关三:模式锁,防止两边混用
光有「两副面孔」是危险的——同一个 WebSocket 连接绑定在某个事件循环上,如果一会儿在主循环用、一会儿在后台循环用,就会炸。
所以有 _claim_execution_mode()(env_client.py:435):
def _claim_execution_mode(self, mode: str) -> None:
if self._execution_mode is None:
self._execution_mode = mode
elif self._execution_mode != mode:
raise RuntimeError(
f"EnvClient is already being used in {self._execution_mode} mode. "
"Create a separate client instance when mixing sync and async code."
)
第一次使用决定终身。之后换边直接报错,而且错误信息直接给出解决方案。
唯一的例外是 _run_sync(..., allow_async_handoff=True)(env_client.py:445-455):当没有事件循环在跑时,允许从 async 交接回 sync——也就是 3.3 的第 ③ 条出口。这对应「先 await Client.from_env(...) 建好客户端,后面全用同步调用」这种真实模式。
3.6 SyncEnvClient:一个专属的后台事件循环
.sync()(env_client.py:952)返回 SyncEnvClient(sync_client.py:43),而且是缓存的——同一个 async 客户端只有一个 sync 影子(env_client.py:978-980)。
它的做法很干脆:起一个 daemon 线程,在里面 run_forever 一个专属事件循环(sync_client.py:91-98)。所有调用通过 asyncio.run_coroutine_threadsafe 丢进去再阻塞等结果(sync_client.py:134-141)。
主线程(你的同步代码) 后台线程(openenv-sync-client-loop)
│ │
env.step(a) 事件循环 run_forever
│ │
├── run_coroutine_threadsafe ─▶│ _step_async(a)
│ │ ↕ WebSocket
│◀────── future.result() ──────┤
▼ │
拿到 StepResult
为什么非要固定一个循环? 因为 websockets 的连接对象绑定创建它的循环。如果每次调用都 asyncio.run() 开一个新循环,连接就废了。文档字符串把这点说得很明白(sync_client.py:51-52)。
这个后台循环也正是 3.3 第 ① 条守卫里说的「sync 客户端自己的循环」:调用一旦已经跑在它上面,_dispatch 就不能再往里丢一次 run_coroutine_threadsafe,否则等于自己等自己。
三个工程细节:
- 循环初始化有双重检查锁(
sync_client.py:100-128),防多线程同时首次调用; - 启动等待 5 秒超时,超时抛
RuntimeError(sync_client.py:125-126); __getattr__把未知属性转发给 async 客户端,如果是协程函数就自动包一层同步 wrapper 并缓存(sync_client.py:272-292)。这就是为什么MCPToolClient.call_tool这种子类方法不需要写同步版,.sync()之后自动可用。
__del__ 会尽力停掉后台循环,但文档明说别依赖它——用 with 或显式 close()(sync_client.py:54-56)。
3.7 连接:比想象中讲究
陈旧连接的识别
_connect_async()(env_client.py:484)开头不是简单的 if self._ws is not None: return,而是:
if self._ws is not None:
if self._ws_loop is asyncio.get_running_loop():
return self
# 不同循环 —— 旧连接已经废了
self._ws = None
self._ws_loop = None
真实源码 env_client.py:494-507,那段注释描述的正是「先 asyncio.run 里建连、后走 .sync() 的后台循环」这个真实场景。旧循环通常已经关了,连接既不能复用也不能干净关闭,只能丢掉重连。
配套的 _ensure_connected()(env_client.py:563)干脆无条件转调 _connect_async(),把「能不能复用」的判断权全交给它——注释里明写了这个理由。
本地代理绕行,而且不改全局环境变量
_is_localhost_ws_url()(env_client.py:200)判断目标是不是回环地址。注意它解析 hostname 再判断,不做子串匹配:
hostname = urlsplit(ws_url).hostname
if hostname == "localhost":
return True
try:
return ipaddress.ip_address(hostname).is_loopback
except ValueError:
return False
函数文档特意举了反例:my-localhost-proxy.example.com 和 127.0.0.1.example.com 都是远端主机,不能误判(env_client.py:201-207)。
命中回环时,通过每连接的 proxy=None 参数关代理,而不是改 NO_PROXY 环境变量。注释解释了原因(env_client.py:517-520):asyncio.gather 并发建连时,改 os.environ 会互相竞争、泄漏状态。
超时与保活
构造函数暴露五个旋钮(env_client.py:279-289):
| 参数 | 默认 | 说明 |
|---|---|---|
connect_timeout_s | 10.0 | 建连超时 |
message_timeout_s | 60.0 | 单条消息响应超时 |
max_message_size_mb | 100.0 | 消息上限,文档说明是为了容纳截图、DOM 这类大观察 |
websocket_ping_interval_s | 20.0 | 保活 ping 间隔,None 关闭 |
websocket_ping_timeout_s | 20.0 | pong 超时 |
100MB 这个默认值透露了 OpenEnv 的定位:它预期观察里会有浏览器截图和完整 DOM,不是 CartPole 那种四个浮点数。
3.8 容器生命周期:客户端也能当运维
三种起法
| 入口 | 位置 | 行为 |
|---|---|---|
EnvClient(base_url=...) | env_client.py:279 | 连已有服务,不管容器 |
EnvClient.from_docker_image(image) | env_client.py:506 | LocalDockerProvider 起容器 → 等就绪 → 连 |
EnvClient.from_env(repo_id) | env_client.py:542 | 从 HF Space 拉镜像或用 uv 本地跑 |
EnvClient(provider=...)(不给 URL) | env_client.py:361 | 连接时才让 provider 自己启动 |
最后一种最有意思。_start_provider_if_needed()(env_client.py:361)会先检查 provider 的 start_container() 有没有必填参数——有就报错,并说明「provider 自启动」这条路走不通:
raise ValueError(
f"{type(self._provider).__name__} does not support "
"provider-owned startup because start_container() requires "
f"{required}. Start the provider manually and pass base_url, ..."
)
真实源码 env_client.py:370-375,参数检查在 _required_start_container_parameters(env_client.py:219)。
from_env 的两条路
from_env(repo_id, use_docker=?)
│
┌────┴────┐
True False
│ │
▼ ▼
拉 HF 镜像 UVProvider
registry.hf.space/ git+https://huggingface.co/
{org}-{space}:latest spaces/{repo_id}
│ │
└────┬─────┘
▼
wait_for_ready → connect
镜像名拼接在 env_client.py:699,uv 路径的默认 git URL 在 env_client.py:717。
uv 路径有一处细心的错误处理(env_client.py:741-747):启动失败时,由于客户端还没建出来、调用方没有 close() 可调,这里是唯一能释放子进程和临时 clone 目录的机会,所以显式 provider.stop() 再抛。