数据截至 (上游 commit 5359534c6f00)
第 2 章 · HTTPEnvServer:端点、会话与并发
本章讲服务端。核心类是
HTTPEnvServer(src/openenv/core/env_server/http_server.py:140),1716 行文件里最重的一块。读完你能看懂一条 WebSocket 连接从建立到销毁发生了什么。
2.1 它要解决的小问题
把一个 Python 类包成 HTTP 服务,本身不难。难的是这三件事:
- 一次 GRPO 采样要开 32 条并行轨迹,32 条轨迹不能共用同一个棋盘;
- 很多环境内部塞着 Playwright、greenlet 这类线程敏感的同步库,不能随便在线程池里乱跑;
- 同一份代码,训练时要暴露
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)。按「哪些模式下存在」分成三类:
| 端点 | 方法 | 仿真模式 | 生产模式 | 干什么 |
|---|---|---|---|---|
/ws | WebSocket | ✅ | ✅ | 会话通道:reset/step/state/close/mcp 五种消息 |
/mcp | WebSocket | ✅ | ✅ | 纯 MCP JSON-RPC 通道 |
/mcp | POST | ✅ | ✅ | 无状态或带 session_id 的 MCP JSON-RPC |
/health | GET | ✅ | ✅ | 健康检查,容器就绪探测用 |
/schema | GET | ✅ | ✅ | 一次返回 action/observation/state 三份 JSON Schema |
/metadata | GET | ✅ | ✅ | 环境名、描述、版本、README |
/reset | POST | ✅ | ❌ | 单次重置(无状态) |
/step | POST | ✅ | ❌ | 单次动作(无状态) |
/state | GET | ✅ | ❌ | 查内部状态 |
模式判断在 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_stacks | AsyncExitStack(持有 MCP 会话) | http_server.py:249 |
_session_info | SessionInfo(建立时间、活跃时间、步数) | 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)。清理顺序是有讲究的:
- 先
await stack.aclose()—— 优雅退出 MCP 会话; - 再
run_in_executor(executor, env.close)—— 在创建它的那个线程里关环境; - 最后
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 线程模型:三种执行器
服务端同时维护三种执行器,各有分工。
| 执行器 | 规模 | 谁用 | 定义处 |
|---|---|---|---|
_executor | max_workers=32 | HTTP /reset、/step 这类一次性请求 | http_server.py:255 |
| 每会话 executor | max_workers=1 | WebSocket 会话的同步 reset/step | http_server.py:385 |
_shared_session_executor | max_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)在构造函数里就跑。逻辑很短:
max_concurrent_envs <= 1→ 直接放行;- 拆开
functools.partial,拿到真正的类; - 如果工厂不是类(是个函数),先造一个临时环境读属性再
close()掉(http_server.py:296-299); - 没有
SUPPORTS_CONCURRENT_SESSIONS=True→ 抛ConcurrencyConfigurationError。
错误信息本身就是使用说明(src/openenv/core/env_server/exceptions.py:32-37):要么把 max_concurrent_envs 调回 1,要么确认环境真的隔离了状态再打开开关。
容量满了怎么办
SessionCapacityError 带着 active_sessions 和 max_sessions 两个字段(exceptions.py:42),在两条通道上被翻译成不同格式:
| 通道 | 表现 |
|---|---|
/ws | WSErrorResponse,code = CAPACITY_REACHED(http_server.py:1666-1675) |
/mcp | JSON-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_env、browsergym_env、tbench2_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_TYPE | type 不认识 |
VALIDATION_ERROR | Pydantic 校验失败(带 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_app | http_server.py:1699 | 看 ENABLE_WEB_INTERFACE 环境变量,决定要不要带 Gradio Web UI |
create_fastapi_app | http_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:73 的 ENV 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.py | HTTPEnvServer |
| 工厂校验 | 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.py | register_routes |
| WebSocket 主循环 | src/openenv/core/env_server/http_server.py | websocket_endpoint |
| kwargs 过滤 | src/openenv/core/env_server/http_server.py | _get_valid_kwargs |
| 应用工厂 | src/openenv/core/env_server/http_server.py | create_app、create_fastapi_app |
| GET 端点声明式注册 | src/openenv/core/env_server/route_config.py | GetEndpointConfig、register_get_endpoints |
| 异常族 | src/openenv/core/env_server/exceptions.py | SessionCapacityError、ConcurrencyConfigurationError、EnvironmentFactoryError |
| 并发配置类型 | src/openenv/core/env_server/types.py | ConcurrencyConfig、ServerCapacityStatus、SessionInfo |
| Web UI | src/openenv/core/env_server/web_interface.py | create_web_interface_app、WebInterfaceManager |