跳到主要内容

数据截至 (上游 commit 5359534c6f00)

第 6 章 · 奖励、harness 与数据采集

本章讲训练侧的三层:环境内怎么算奖励(Rubric / Transform)、怎么把环境驱动成一条轨迹(harness)、怎么把轨迹存成数据集(collect)。设计依据是 rfcs/004-rubrics.mdrfcs/005-agentic-harnesses.md


6.1 一条铁律:奖励只能诞生在环境里

这是 OpenEnv 反复强调的不变量。理由不难懂:奖励是领域知识。什么叫「这段代码写得对」「这步棋走得好」,只有环境自己知道。如果让训练框架在外面拍脑袋加分,同一个环境换个框架就换个分数,实验根本没法比。

所以这一章的三层结构可以这么理解:

┌────────────────────────────────────────────────┐
│ 环境内部 │
│ Rubric / Transform ← 奖励在这里产生 │
└───────────────────┬────────────────────────────┘
│ Observation.reward

┌────────────────────────────────────────────────┐
│ harness 层 │
│ 驱动 rollout + _resolve_env_reward 交叉校验 │
└───────────────────┬────────────────────────────┘
│ HarnessRolloutResult

┌────────────────────────────────────────────────┐
│ collect 层 │
│ 落盘成 JSONL 数据集,可推到 HF Hub │
└────────────────────────────────────────────────┘

6.2 环境内算分:两种工具,不同定位

TransformRubric
输入只有 observationaction + observation
定位观察后处理管道评分器
借鉴自TorchRL transformPyTorch nn.Module
组合方式CompositeTransform 串联属性赋值自动注册子评分器
位置interfaces.py:115rubrics/base.py:17

简单场景用 Transform 就够(比如 envs/coding_env/server/transforms.pyCodeSafetyTransform 扫危险模式打惩罚)。复杂场景才上 Rubric。


6.3 Rubric:把评分器写成 nn.Module

直觉

如果你写过 PyTorch,Rubric 会非常眼熟:

PyTorchRubric
实现 forward()实现 forward(action, observation) -> float
子模块属性赋值即注册子 rubric 属性赋值即注册
named_modules()named_rubrics()
register_forward_hookregister_forward_hook
state_dict()state_dict()

自动注册怎么做到的

重写 __setattr__(src/openenv/core/rubrics/base.py:50-54):

def __setattr__(self, name: str, value: Any) -> None:
if isinstance(value, Rubric):
self._rubric_children[name] = value
object.__setattr__(self, name, value)

__init__ 里用 object.__setattr__ 初始化那几个内部字段,避免自己触发自己(base.py:44-48)。

组合示意

# 示意,非源码
class CodeRubric(Rubric):
def __init__(self):
super().__init__()
self.syntax = SyntaxRubric() # 赋值即注册为子评分器
self.tests = TestPassRubric()

def forward(self, action, observation):
return 0.3 * self.syntax(action, observation) + 0.7 * self.tests(action, observation)

rubric = CodeRubric()
reward = rubric(action, obs)
for name, r in rubric.named_rubrics():
print(name, r.last_score) # "syntax 1.0" / "tests 0.5" —— 可逐项审查

last_score_finish_forward() 自动记录(base.py:95-103),这是可观测性的关键:训练时能看到每个子项各贡献了多少。Environment 的文档字符串里就演示了这个用法(interfaces.py:164-167)。

同步/异步双路

__call__(base.py:56-76)先用 inspect.iscoroutinefunction(self.forward) 判断,再走不同分支。注意注释里的一处讲究:同步路径的前置钩子必须在 forward() 之前调,而异步路径因为 forward() 返回的只是协程(还没执行),钩子可以在 _call_async 里补调。

现成的 rubric

类别文件
容器SequentialGateWeightedSumRubricListRubricDictrubrics/containers.py
轨迹TrajectoryRubricExponentialDiscountingTrajectoryRubricrubrics/trajectory.py
模型裁判LLMJudgerubrics/llm_judge.py

延迟奖励:TrajectoryRubric

有些奖励只有在局末才知道:棋赢没赢、计划成没成。TrajectoryRubric(rubrics/trajectory.py:23)的做法是:

  • 每步都累积 (action, observation)_trajectory;
  • done 之前返回 intermediate_reward(默认 0.0);
  • 局末调 score_trajectory() 算总分,再由 compute_step_rewards() 决定怎么把分摊回各步(信用分配)。

文档里明确写了两条注意事项(trajectory.py:33-37):轨迹存CPU 内存,环境如果观察里有 GPU 张量必须先搬到 CPU;超长 episode(几千步)会吃不少内存。

环境怎么用 Rubric

基类给了四个受保护钩子(interfaces.py:258-346):_apply_rubric / _apply_rubric_async / _reset_rubric / _reset_rubric_async。文档字符串里直接给了用法样板:

def step(self, action, ...):
# ... 执行动作、构造 observation ...
observation.reward = self._apply_rubric(action, observation)
return observation

同样是环境作者主动调,框架不替你调——奖励逻辑的所有权始终在环境手里。


6.4 harness 层:把环境驱动成一条轨迹

定义在 src/openenv/core/harness/__init__.py,状态是「实验中」——模块文档开头就说它在 RFC 005 定稿前不属于稳定 API(harness/__init__.py:5-6)。

它要解决的小问题

有了环境和模型,还缺中间那段:谁负责「把工具清单塞进 prompt → 采样模型 → 解析 tool call → 调环境 → 把结果拼回对话 → 循环」?这段循环各家 RL 框架都在重写。harness 把它抽象出来。

四个角色

ResourceSessionFactory ──create()──▶ ResourceSession

│ initial_messages / list_tools
│ call_tool / verify / close

ModelStep ◀──────── HarnessAdapter.run_white_box()
(采样模型) │

HarnessRolloutResult
角色是什么位置
ResourceSession一次 rollout 独占的环境会话harness/__init__.py:112
ResourceSessionFactory造会话的工厂:144
HarnessAdapter驱动循环的策略:157
ModelStep采样模型下一轮的可调用协议:101

白盒 vs 黑盒

HarnessAdapter 要求实现两个方法:

方法谁掌握采样实现
run_white_box训练器掌握(能拿到 token id 和 logprob)MCPHarnessAdapter(:482)
run_black_boxharness 自己掌握(比如一个 CLI 工具)CLIHarnessAdapter(:595)

两个实现各自把对方的方法实现成 NotImplementedError 并附上指路说明(:589-592:613-615)。这是个诚实的做法:接口统一,但不假装自己两种都能干。

白盒路径为什么重要?看 ModelStepResult(:67)带的字段:prompt_idscompletion_idslogprobs。GRPO 这类算法需要这些才能算重要性比。黑盒 harness(想象一个不透明的编码 agent CLI)给不出来,只能做评测。

MCPHarnessAdapter 的循环

run_white_box(:485)的骨架:

messages = session.initial_messages()
tools = session.list_tools()

for turn in range(max_turns):
采样 → assistant 消息入队,累计 prompt_ids/completion_ids/logprobs
没有 tool_calls? → done,退出
对每个 tool_call:
超出 max_total_tool_calls? → 标记 truncated,直接返回
session.call_tool(...) → 记入 tool_trace
结果作为 role="tool" 消息入队
tool_result.done? → 退出

限额由 HarnessRunLimits(:91)给:max_turns(默认 10)、max_tool_calls_per_turnmax_total_tool_callssampling

黑盒路径的桥:SessionMCPBridge

SessionMCPBridge(:394)把一个 ResourceSession 包装成进程内的 MCP JSON-RPC 端点。于是一个只会说 MCP 的外部 harness,不需要网络就能驱动 OpenEnv 会话。CLIHarnessAdapter 拿到 bridge 后交给用户传入的 runner(:617-624)。

把现成客户端接进来:StepEnvSessionAdapter

StepEnvSessionAdapter(:243)把任何 reset/step/state 客户端变成 ResourceSession。几个实现细节:

  • 构造时如果客户端有 .sync() 就先转同步(:268-271)——harness 整层是同步的;
  • 工具名撞保留字直接 ValueError(:281-286),边界在这一层再守一次;
  • call_tool 把工具名+参数交给 action_builder 变成动作,step() 之后顺手读一次 state(:361-372)。

仓库里有四个环境提供了自己的 session factory:envs/openspiel_env/harness.pyenvs/browsergym_env/harness.pyenvs/opencode_env/harness.pyenvs/reasoning_gym_env/harness.py


6.5 _resolve_env_reward:铁律的执行者

这是本章最值得单独讲的一个函数(harness/__init__.py:218)。

它在防什么

奖励可以从两个地方冒出来:

  1. 工具调用结果的 metadata(环境在 step 时给的);
  2. session.verify() 返回的 VerifyResult.env_reward(局末汇总时给的)。

如果编排层能随意选一个、或者干脆自己造一个,那「奖励只在环境内产生」就成了一句空话。

它怎么防

① 从 tool_trace 倒着找,取最近一个非 None 的 reward → trace_reward
② 取 verify.env_reward → verify_reward
③ 两个都有,且 math.isclose 判定不相等 → 直接抛 ValueError
④ 优先返回 trace_reward,其次 verify_reward
⑤ 两个都没有 → 抛 "rollout did not produce an environment reward"

第 ③ 步的报错文案是:

raise ValueError(
"verify.env_reward must forward the environment reward from the rollout"
)

真实源码 harness/__init__.py:241-243,容差是 rel_tol=1e-9, abs_tol=1e-6

妙在哪: 这是一个交叉校验verify() 是环境作者写的钩子,理论上可以在里面改写奖励;这个函数用「两个来源必须一致」把这条路堵死了。VerifyResult 的文档也把话说在前面(:41-42):「必须转发环境内已经产生的奖励,不得在编排层合成新奖励」。

第 ⑤ 步同样重要:没有奖励就是错误,不是 0。静默补 0 会让整批数据看起来正常,实际全是垃圾。


6.6 build_harness_rollout_func:接进 TRL

build_harness_rollout_func()(harness/__init__.py:636)返回一个 (prompts, trainer) -> dict 的函数,这正是 TRL 期望的 rollout function 形状。

返回的字典有五个键:prompt_idscompletion_idslogprobs<reward_key>(默认 env_reward)、verify_metrics

每条 prompt 一个独立会话,try/finally 保证 session.close()(:644-687)。奖励那一列走的就是 _resolve_env_reward


6.7 collect:把轨迹存成数据集

src/openenv/core/harness/collect.py 建在 harness 之上,目标是生成 SFT/蒸馏用的数据集。

输出格式

results.jsonl,一行一局。字段表在 src/openenv/core/harness/README.md:71-82,对应 EpisodeRecord(collect.py:63):

字段说明
episode_id<prefix>-<6位序号>,断点续跑靠它
messages对话记录,TRL SFTTrainer 直接吃
reward环境奖励,_resolve_env_reward 得来(collect.py:92)
done / tool_trace / metrics / verify_metrics / artifacts / task / extra其余元数据

三个实用设计

断点续跑。 collected_episode_ids()(collect.py:141)读回已落盘的 id 集合,run(resume=True) 跳过它们(collect.py:270-274)。逐行解析,坏行直接跳过而不是整体崩(collect.py:152-155)。

任务对齐。 续跑跳过某个 id 时,仍然会先调一次 _next_task()continue(collect.py:269-274),注释写明是为了「即使 resume 跳过 id,也保持任务与 episode 的对齐」。

质量筛选。 should_keep: Callable[[EpisodeRecord], bool],默认全留;CLI 默认筛 reward >= 0,--keep-losses 可关。

教师模型是可换的

build_model_step()(collect.py:366)把任意 LLMClient 适配成 ModelStepcreate_llm_client()(src/openenv/core/llm_client.py:359)支持 OpenAI、Anthropic、以及 OpenAI 兼容的自建端点。

CLI 还提供 --provider scripted——一个不需要 API key 的脚本化教师(_build_scripted_model_step,cli/commands/collect.py:111),专门用来冒烟测试整条管线。

一条命令

来自 src/openenv/core/harness/README.md:19-24:

openenv collect openspiel:tic_tac_toe \
--base-url https://<user>-<space>.hf.space \
--output-dir /tmp/ttt-scripted \
-n 10 --provider scripted

push_to_hf_hub()(collect.py:499)会自动生成带 YAML frontmatter 的 README,让 HF 数据集查看器把 results.jsonl 识别成 train split。


6.8 evals:另一条评测线

src/openenv/core/evals/ 是个更薄的抽象:EvalHarness(evals/base.py:11)+ EvalConfig/EvalResult(evals/types.py),外加一个 Inspect 框架的适配器 inspect_harness.py(可选依赖 openenv[inspect])。

它和 Rubric 的区别 RFC 001 说得很清楚:rubric 是数据无关、逐样本的;eval 是数据相关、要聚合的


6.9 关键细节与坑

  • harness 层是同步的。 ResourceSession 的方法全是普通 def,StepEnvSessionAdapter 会把 async 客户端转成 .sync()。要并发就在更外层开进程/线程。
  • Rubric 的异步支持是「按需」的。 __call__ 检测到 forward 是协程函数就返回协程;调用方(比如 Environment._apply_rubric_async)负责 await(interfaces.py:305-309)。同步调用点拿到协程会出问题。
  • MCPHarnessAdapter 的终止条件有两个。 模型不再发 tool call,或者某个 tool_result.done 为真。超限时只标 metrics["truncated"] = True,不抛异常(harness/__init__.py:545)。
  • EpisodeRecord.to_dictjson.dumps(..., default=str) 兜底。 不可序列化的对象会变成字符串(collect.py:107),不会炸,但也会静默丢失结构。
  • collect 的失败是逐局记账的。 单局异常计入 num_failed 并打印,不中断整轮(collect.py:293-298)。

6.10 代码地图

主题文件符号
评分器基类src/openenv/core/rubrics/base.pyRubric__setattr__named_rubricslast_score
评分器容器src/openenv/core/rubrics/containers.pySequentialGateWeightedSum
延迟奖励src/openenv/core/rubrics/trajectory.pyTrajectoryRubricExponentialDiscountingTrajectoryRubric
模型裁判src/openenv/core/rubrics/llm_judge.pyLLMJudge
环境侧钩子src/openenv/core/env_server/interfaces.py_apply_rubric_reset_rubric_apply_transform
观察后处理src/openenv/core/env_server/base_transforms.pyCompositeTransform
会话协议src/openenv/core/harness/__init__.pyResourceSessionResourceSessionFactory
harness 驱动src/openenv/core/harness/__init__.pyMCPHarnessAdapterCLIHarnessAdapter
奖励交叉校验src/openenv/core/harness/__init__.py_resolve_env_reward_tool_result_reward
客户端接入src/openenv/core/harness/__init__.pyStepEnvSessionAdapterSessionMCPBridge
TRL 接口src/openenv/core/harness/__init__.pybuild_harness_rollout_func
数据集采集src/openenv/core/harness/collect.pyCollectRunnerEpisodeRecordRolloutSerializer
教师适配src/openenv/core/harness/collect.pybuild_model_steppush_to_hf_hub
LLM 客户端src/openenv/core/llm_client.pycreate_llm_clientOpenAIClientAnthropicClient
评测抽象src/openenv/core/evals/base.pyEvalHarness
设计提案rfcs/004-rubrics.mdrfcs/005-agentic-harnesses.md