跳到主要内容

数据截至 (上游 commit 5359534c6f00)

第 2 章 · HTTPEnvServer:端点、会话与并发

本章讲服务端。核心类是 HTTPEnvServer(src/openenv/core/env_server/http_server.py:140),1716 行文件里最重的一块。读完你能看懂一条 WebSocket 连接从建立到销毁发生了什么。


2.1 它要解决的小问题

把一个 Python 类包成 HTTP 服务,本身不难。难的是这三件事:

  1. 一次 GRPO 采样要开 32 条并行轨迹,32 条轨迹不能共用同一个棋盘;
  2. 很多环境内部塞着 Playwright、greenlet 这类线程敏感的同步库,不能随便在线程池里乱跑;
  3. 同一份代码,训练时要暴露 reset,上线推理时必须不暴露

HTTPEnvServer 就是这三件事的答案。


2.2 第一个关键决定:服务端不持有环境,持有工厂

构造函数第一件事就是检查 env 是不是 callable,不是就直接 TypeError,错误信息还专门提醒「传类,别传实例」:

if not callable(env):
raise TypeError(
f"env must be a callable (class or factory function), got {type(env)}. "
f"Pass the environment class (e.g., MyEnvironment) not an instance ..."
)

真实源码 http_server.py:208-212。存下来的字段叫 self._env_factory(http_server.py:214)。

这个决定往下传导出整章其余所有设计:因为服务端手里只有工厂,它才能在每条连接来的时候现造一个新环境

对照 envs/echo_env/server/app.py:44-50,环境作者传的确实是类本身:

app = create_app(
EchoEnvironment, # 类,不是 EchoEnvironment()
CallToolAction,
CallToolObservation,
env_name="echo_env",
max_concurrent_envs=max_concurrent,
)

2.3 端点全景

路由注册在 register_routes()(http_server.py:626)。按「哪些模式下存在」分成三类:

端点方法仿真模式生产模式干什么
/wsWebSocket会话通道:reset/step/state/close/mcp 五种消息
/mcpWebSocket纯 MCP JSON-RPC 通道
/mcpPOST无状态或带 session_id 的 MCP JSON-RPC
/healthGET健康检查,容器就绪探测用
/schemaGET一次返回 action/observation/state 三份 JSON Schema
/metadataGET环境名、描述、版本、README
/resetPOST单次重置(无状态)
/stepPOST单次动作(无状态)
/stateGET查内部状态

模式判断在 http_server.py:1259(if mode == ServerMode.SIMULATION:)和 http_server.py:1386(/state 的条件插入)。

一个容易踩的坑:HTTP /reset/step 是无状态的

reset_handler(http_server.py:670)的结构:

_env = self._env_factory() # 每次请求现造
try:
...
finally:
_env.close() # 用完就扔

真实源码 http_server.py:674:669,step_handler 同理(:683:707)。所以你用 HTTP /step 连着发两次动作,第二次不会记得第一次——它们是两个不同的环境实例。要有状态,必须走 WebSocket。

/state/metadata 的处理函数也是同样的「造—用—扔」模式(http_server.py:1346-1358),这意味着 GET /state 在仿真模式下返回的是一个全新环境的初始状态,不是某条会话的状态。这是个不太直观、但从「HTTP 无状态」角度自洽的设计。


2.4 会话生命周期:本章核心

一条 WebSocket 连接换来什么

_create_session()(http_server.py:385)。一次成功的会话创建会产生四样东西,分别存进四个字典:

字典存什么定义处
_sessions环境实例http_server.py:247
_session_executors该会话独占的 ThreadPoolExecutor(max_workers=1)http_server.py:248
_session_stacksAsyncExitStack(持有 MCP 会话)http_server.py:249
_session_infoSessionInfo(建立时间、活跃时间、步数)http_server.py:276

创建流程:两段加锁,中间放开

这是本章最值得学的一段代码。怎么读下图:方框是步骤,[锁] 标记表示在 asyncio.Lock 保护下执行。

[锁] ① 检查容量 → 满了抛 SessionCapacityError
② 生成 session_id、建线程池
③ _sessions[sid] = None ← 占位!先把坑占住
─────────────── 放开锁 ───────────────
④ 在线程池里跑 env_factory() ← 慢操作,不占锁
⑤ 进 MCP 会话(如果环境有)
─────────────── 重新加锁 ──────────────
[锁] ⑥ _sessions[sid] = env、写入 SessionInfo

占位符那一步的注释写得很直白(http_server.py:380-387):

# Create executor and reserve slot so capacity is not exceeded while
# we create the env outside the lock (avoids blocking other sessions)
...
self._sessions[session_id] = None # placeholder until env is ready

妙在哪: 环境构造可能很慢(拉起一个浏览器、初始化一个模拟器)。如果整段都锁着,并发建连就串行化了。用 None 占位就同时拿到两个好处——容量计数立刻生效(不会超卖),慢操作却在锁外跑。

代价是系统里存在一个「已占坑但环境还没好」的中间态。代码里处处要处理它,比如 MCP 的 openenv/session/close 方法遇到 env is None 时,会把占位符重新塞回去再报错,避免 _create_session 后半段写进一个已被删掉的键(http_server.py:823-836)。

失败路径也照顾到了

两处清理值得注意:

  • 工厂抛异常 → 回滚线程池和占位符,再包装成 EnvironmentFactoryError(http_server.py:393-402);
  • MCP 传输启动失败 → 关 stack、删三个字典项、清理环境和线程池,注释明说是为了「不让它们永久地占着 _max_concurrent_envs」(http_server.py:415-425)。

销毁流程

_destroy_session()(http_server.py:441)先在锁内把四个字典项 pop 出来,然后在锁外调 _cleanup_session_resources()(http_server.py:457)。清理顺序是有讲究的:

  1. await stack.aclose() —— 优雅退出 MCP 会话;
  2. run_in_executor(executor, env.close) —— 在创建它的那个线程里关环境;
  3. 最后 executor.shutdown(wait=False)

第 2 步的注释点明了原因(http_server.py:473-474):

# Run close() in the same executor where the env was created
# This is required for thread-sensitive libraries like Playwright/greenlet

2.5 线程模型:三种执行器

服务端同时维护三种执行器,各有分工。

执行器规模谁用定义处
_executormax_workers=32HTTP /reset/step 这类一次性请求http_server.py:255
每会话 executormax_workers=1WebSocket 会话的同步 reset/stephttp_server.py:385
_shared_session_executormax_workers=1,全局唯一声明了 REQUIRES_SINGLE_THREAD_EXECUTOR 的环境http_server.py:258-260

每会话单线程是关键设计。同一个环境实例的所有调用永远落在同一个线程上,这样 Playwright 那种「对象绑定线程」的库才不会炸。

第三种更极端:某些环境要求所有会话共用一个线程(比如底层库全局只允许一个事件循环)。这时 _detect_single_thread_requirement()(http_server.py:307)在启动时读类属性,建一个共享执行器,所有会话都用它——注意这会让并发退化成串行。共享执行器不会被单个会话 shutdown(http_server.py:492),只在应用关闭时统一关(http_server.py:661-662)。

分派逻辑

WebSocket 的 step 分支(http_server.py:1573-1593):

is_async = session_env.step_async.__func__ is not Environment.step_async

if is_async:
observation = await session_env.step_async(action) # 事件循环上直接跑
else:
observation = await self._run_in_session_executor( # 丢进会话独占线程
session_id, session_env.step, action
)

_run_in_session_executor()http_server.py:570,取不到会话执行器时会退回全局的 _executor


2.6 并发闸门与容量

声明式闸门

_validate_concurrency_safety()(http_server.py:276)在构造函数里就跑。逻辑很短:

  1. max_concurrent_envs <= 1 → 直接放行;
  2. 拆开 functools.partial,拿到真正的类;
  3. 如果工厂不是类(是个函数),先造一个临时环境读属性再 close() 掉(http_server.py:296-299);
  4. 没有 SUPPORTS_CONCURRENT_SESSIONS=True → 抛 ConcurrencyConfigurationError

错误信息本身就是使用说明(src/openenv/core/env_server/exceptions.py:32-37):要么把 max_concurrent_envs 调回 1,要么确认环境真的隔离了状态再打开开关。

容量满了怎么办

SessionCapacityError 带着 active_sessionsmax_sessions 两个字段(exceptions.py:42),在两条通道上被翻译成不同格式:

通道表现
/wsWSErrorResponse,code = CAPACITY_REACHED(http_server.py:1666-1675)
/mcpJSON-RPC 错误,code = SERVER_ERROR (-32000)(http_server.py:1227-1236)

两边都把 active/max 数字放进 data,客户端可以据此退避重试。

谁来设这个数

环境作者在 app.py 里从环境变量读。以 echo 为例(envs/echo_env/server/app.py:42):

max_concurrent = int(os.getenv("MAX_CONCURRENT_ENVS", "8"))

这是个约定俗成的写法,仓库里 textarena_envbrowsergym_envtbench2_env 等都用同一个变量名,jupyter_env 默认值是 4。


2.7 空闲会话回收

如果客户端连上就消失(网络断了、进程被杀),会话会永远占着容量。ConcurrencyConfig.session_timeout(types.py:327)开启后,后台任务 _reap_idle_sessions()(http_server.py:512)负责收尸。

它的循环有两个不显然的细节:

其一,检查间隔是自适应的:

interval = max(timeout / 4, 5.0) # check frequently enough

真实源码 http_server.py:517。超时越短查得越勤,但不低于 5 秒。

其二,两阶段确认。 先在锁内快照一份过期 id 列表,放开锁,再逐个重新加锁复查(http_server.py:527-537)。注释解释了为什么不能一把梭:快照之后可能有新活动到达;而且要重新读 now,否则前面几个 _destroy_session 耗时会让后面的条目拿着过期时钟做判断。

生命周期挂在 FastAPI 的 startup/shutdown 事件上,并用 app.router._openenv_reaper_registered 这个自定义标记做幂等,避免同一个 app 被 register_routes 多次时重复注册(http_server.py:664-667)。


2.8 WebSocket 消息协议

/ws 的消息用 type 字段区分,服务端用 match 语句分派(http_server.py:1537)。

客户端发服务端回说明
{"type": "reset", "data": {...}}{"type": "observation", "data": {...}}data 是 reset 的 kwargs
{"type": "step", "data": {...}}{"type": "observation", "data": {...}}data 是动作 dict
{"type": "state"}{"type": "state", "data": {...}}
{"type": "mcp", "data": {...}}{"type": "mcp", "data": {...}}data 是 JSON-RPC 请求/响应
{"type": "close"}(断开)服务端 break 循环

出错时统一回 {"type": "error", "data": {"message": ..., "code": ...}},错误码枚举在 types.py:33-42:

错误码触发条件
INVALID_JSON消息不是合法 JSON
UNKNOWN_TYPEtype 不认识
VALIDATION_ERRORPydantic 校验失败(带 errors 明细)
EXECUTION_ERROR环境执行时抛异常
CAPACITY_REACHED会话数满
FACTORY_ERROR环境工厂建不出来
SESSION_ERROR其他会话级失败

注意错误的粒度差别: 前四种是「单条消息失败」——回一个 error 帧然后继续循环(http_server.py:1532:1483:1491);后三种是「整条连接失败」——回完就走 finally 销毁会话(http_server.py:1666-1692)。

消息类型也是 Pydantic 模型

定义在 types.py:250-287,并用 Field(discriminator="type") 组成判别联合 WSIncomingMessage。有意思的是服务端没有直接用这个联合,而是手写 match 再逐个 WSResetMessage(**message_dict)——因为 MCP 消息类型 WSMCPMessage 定义在另一个模块以避免循环导入,联合里放不下(见 types.py:281-283 的注释)。


2.9 两个应用工厂

函数位置干什么
create_apphttp_server.py:1699ENABLE_WEB_INTERFACE 环境变量,决定要不要带 Gradio Web UI
create_fastapi_apphttp_server.py:1790纯 API,不带 UI,自带完整 OpenAPI 文档配置

开关判断在 http_server.py:1755-1759:

enable_web = os.getenv("ENABLE_WEB_INTERFACE", "false").lower() in ("true", "1", "yes")

默认是的——本地开发保持轻量。但环境的 Dockerfile 会显式打开,比如 envs/echo_env/server/Dockerfile:73ENV ENABLE_WEB_INTERFACE=true。所以「本地裸跑没有 UI,容器里有 UI」是预期行为,不是 bug。

Web UI 本身是 Gradio 实现的,入口在 src/openenv/core/env_server/web_interface.py:425(create_web_interface_app),会根据 Action 类的 JSON Schema 自动生成输入表单(_extract_action_fields,web_interface.py:656)。


2.10 代码地图

主题文件符号
服务端主类src/openenv/core/env_server/http_server.pyHTTPEnvServer
工厂校验src/openenv/core/env_server/http_server.py__init__(callable 检查)
并发闸门src/openenv/core/env_server/http_server.py_validate_concurrency_safety
单线程需求探测src/openenv/core/env_server/http_server.py_detect_single_thread_requirement
会话创建src/openenv/core/env_server/http_server.py_create_session
会话销毁src/openenv/core/env_server/http_server.py_destroy_session_cleanup_session_resources
空闲回收src/openenv/core/env_server/http_server.py_reap_idle_sessions_start_reaper
路由注册src/openenv/core/env_server/http_server.pyregister_routes
WebSocket 主循环src/openenv/core/env_server/http_server.pywebsocket_endpoint
kwargs 过滤src/openenv/core/env_server/http_server.py_get_valid_kwargs
应用工厂src/openenv/core/env_server/http_server.pycreate_appcreate_fastapi_app
GET 端点声明式注册src/openenv/core/env_server/route_config.pyGetEndpointConfigregister_get_endpoints
异常族src/openenv/core/env_server/exceptions.pySessionCapacityErrorConcurrencyConfigurationErrorEnvironmentFactoryError
并发配置类型src/openenv/core/env_server/types.pyConcurrencyConfigServerCapacityStatusSessionInfo
Web UIsrc/openenv/core/env_server/web_interface.pycreate_web_interface_appWebInterfaceManager