AgentEvolver 源码解读 · 第 6 章:环境服务与基建
> 前五章皆在训练进程之「内」。然 agent 每步 env.step 打出去的,是一台独立的环境服务。此章看「墙外」:env_service/ 如何用 FastAPI + Ray 把 AppWorld 这类环境封装成可远程调用的沙箱,以及贯穿全程的数据模型(schema / trajectory)。读完便知:何以最小化训练要先 python env_service/env_service.py 起服务,再用 curl 探 /healthz。
---
一、Why:为何要把环境拆成独立服务
训练主进程要绷紧 GPU 跑 rollout;环境(尤其 AppWorld,起一个 world 实例很重)若塞进同一进程,既抢资源又难并发。AgentEvolver 的解法:环境服务独立成 HTTP 服务,Ray actor 作远程沙箱,多实例并行。主进程里的 EnvClient 只发 REST 请求,不关心环境怎么跑。隔离带来三利:
1. 并发:一个服务可同时持数十个 instance_id,对应一个 batch 的多条轨迹。
2. 容错:EnvClient 有重试与 fallback,环境偶崩不拖垮 RL 主循环。
3. 可换:换环境只需实现 BaseEnv 并 @Registry.register,主进程零改动。
---
二、服务端:EnvService + FastAPI + Ray
env_service.py 是服务的全部。启动方式:python env_service/env_service.py --env appworld --port 8080。
2.1 进程内单例与 Ray
class EnvService:
def __init__(self):
if not ray.is_initialized():
ray.init(address='local') # 本地 Ray,远程环境跑在 actor 上
self.env_actors = {} # instance_id -> Ray actor
self.remote_env = {} # env_type -> 动态生成的 RemoteEnv 类
self.last_access_time = {}
self.max_idle_time = 3600 # 闲置 1 小时自动回收
async def cleanup_inactive_instances(self):
# 后台 loop 每隔 cleanup_interval=300s 扫一遍,超时实例 release
2.2 动态注册:import_and_register_env
def import_and_register_env(env_name, env_file=None):
env_file = env_file or f"{env_name}_env"
env_module_path = os.path.join(SERVER_DIR, "env_service/environments", env_name, f"{env_file}.py")
spec = importlib.util.spec_from_file_location(f"{env_name}_env", env_module_path)
module = importlib.util.module_from_spec(spec) ; sys.modules[spec.name] = module
spec.loader.exec_module(module)
envir_class = getattr(module, f"{env_name.capitalize()}Env") # 如 AppworldEnv
Registry.register(env_name)(envir_class)
return envir_class
启动脚本先 import_and_register_env(args.env),把 AppworldEnv 挂进 Registry;之后所有请求按 env_type 取类。新增环境 = 新建 environments//_env.py 并继承 BaseEnv,无需动服务主干。2.3 远程沙箱:@ray.remote RemoteEnv
get_remote_env_cls 按 env_type 动态造一个 Ray actor 类(每个 env_type 只造一次,缓存于 self.remote_env):
@ray.remote
class RemoteEnv:
def __init__(self, task_id, instance_id, params):
module = importlib.import_module(f"env_service.environments.{env_type}.{env_type}_env")
self.env = getattr(module, f"{env_type.capitalize()}Env")(task_id, instance_id, params)
def get_init_state(self, params): return self.env.get_init_state(params)
def step(self, action, params): return self.env.step(action, params)
def evaluate(self, messages, params): return self.env.evaluate(messages, params)
def get_info(self, messages, params): return self.env.get_info(messages, params)
def close(self): return self.env.close()
create_instance 据此 .remote(task_id, instance_id, params) 拉起一个 actor,并立即 .get_init_state.remote(params) 取初始状态(含 system 提示 + 任务 instruction)。step / evaluate / release 皆经 .remote() 异步执行——每个环境实例即一个独立 Ray actor,天然隔离。2.4 六条 REST 路由
服务用 FastAPI 暴露六接口(ServiceRequest 统一收 env_type/task_id/instance_id/messages/params):| 路由 | 作用 | 关键逻辑 |
|---|---|---|
GET /healthz | 探活 | 直接 Response("OK", 200) |
POST /get_env_profile | 取任务列表 | Registry.get(env_type).get_query_list(split) |
POST /create | 建实例 | create_instance → 返回 init_state(含 query) |
POST /step | 走一步 | env_service.step(instance_id, action, params) |
POST /evaluate | 算分 | env_service.evaluate(instance_id, messages, params) |
POST /get_info | 取工具说明 | 返回 tools_info(AppWorld 的 API 速览) |
POST /release | 释放实例 | ray.kill(actor) + 删登记 |
HTTPException,并附完整 traceback("".join(traceback.format_exception(...)))——排错时极有用。lifespan 在启动/关闭时管一个 cleanup_loop 后台协程,自动回收闲置实例。---
三、客户端:EnvClient 与重试兜底
env_client.py 是训练侧入口,纯 requests + 自定义 retry_call:
class EnvClient:
def __init__(self, base_url="http://localhost:8000"):
self.timeout = 150.0 + random.uniform(50, 200) # 抖动超时,避免集体卡死
def _make_request(self, endpoint, env_type, task_id, instance_id, messages, params):
url = f"{self.base_url}/{endpoint.lstrip('/')}"
data = {"env_type":env_type,"task_id":task_id,"instance_id":instance_id,
"messages":messages or {},"params":params or {}}
response = requests.post(url, json=data, timeout=self.timeout)
response.raise_for_status() ; return response.json()
def step(self, instance_id, action, params, max_retry=3):
fallback = {"state":[{"role":"assistant","content":"Step failed (timeout or exception),please retry"}],
"reward":0,"is_terminated":False,"info":{...}}
def call(): return self._make_request("step", instance_id=instance_id, messages=action, params=params)["data"]
return retry_call(call, max_retry=max_retry, fail_return=fallback, ...)
设计要点:
- 每个方法都带
fallback字典——环境崩时返回一个「假装这步失败、不终止」的占位状态,让主循环继续而非整体挂掉。 retry_call指数退避(random.uniform(3,10)秒),失败达上限返回fail_return,绝不抛到 RL 主进程。safe_log写错误日志用fsync强落盘,且包try/except——日志本身失败也不能影响训练。这正是生产级 RL 基建的谨慎。
AgentFlow.execute 里的 env 正是 EnvClient 实例(由 main_ppo 注入),故 env.step(instance_id, {...}) 实则一次 HTTP POST。---
四、环境基类与 AppWorld 实现
4.1 BaseEnv:抽象的契约
base.py 定义六方法抽象接口:get_init_state / step / evaluate / close / get_info / get_query_list(静态)。任何新环境只需实现这六个,Registry + 服务即自动接管。4.2 AppworldEnv:状态化的 Python REPL
environments/appworld/appworld_env.py 是默认环境,继承 BaseEnv 且 @Registry.register("appworld"):@Registry.register("appworld")
class AppworldEnv(BaseEnv):
def get_init_state(self, params):
self.world = AppWorld(task_id=self.task_id, experiment_name=self.instance_id)
self.query = self._get_query(params)
return {"state":[{"role":"system","content":self.query},
{"role":"user","content":self.world.task.instruction}],
"info":{"instance_id":self.instance_id,"task_id":self.task_id}}
PROMPT_TEMPLATE/simple_prompt:两段 jinja2 模板,把「你是自主 AI 助理、用 Python REPL 调 AppWorld API」的长示例与关键指令(查看 API 文档、分页遍历、调用complete_task(answer=...)等)渲染进 system 提示。simple模式用精简版。step:收ActionMessage→transition解析代码并执行 → 查world.task_completed()→ 若终止则evaluate算分:
def step(self, action, params=None):
action_msg = ActionMessage(**action)
state_message = self.transition(action_msg)
is_terminated = self.world.task_completed()
reward = self.evaluate(params) if is_terminated else 0.0
return {"state":[state_message.simple_dict], "reward":reward,
"is_terminated":is_terminated, "info":{}}
transition:核心解析器。若 model 输出含 ````python\n(.*?)```代码块,抽出来;或从tool_calls取appworld工具的arguments["code"],调appworld_excute→self.world.execute(code)(即真实 Python REPL,变量跨步保留)。执行异常则把错误回信塞进result,逼模型自纠。非appworld工具名一律报「tool not found」——环境只认这一道门。- evaluate
:调 AppWorld 官方 evaluate_task,sparse=True时「零失败即 1、否则 0」;否则返回passes/(passes+failures)的稠密分。sparse正是第 1 章actor_rollout_ref.rollout.sparse的落点。 - get_query_list
:静态方法,load_task_ids(split)直接取 AppWorld 官方任务清单(train/val/test)——/get_env_profile与 TaskManager 都靠它拿任务池。
用 full_code_regex = r"%%ZCCODE8%%" 抽代码;ignore_multiple_calls=True 时只取首个代码块(每步一动作,契合 simple_prompt 第 4 条「每步只写一小块代码」)。---
五、数据模型:两份 schema
项目里有两套 schema,别混淆:
5.1 env_service/trajectory.py(环境侧消息协议)
定义 Role 枚举(system/user/tool/assistant/context_assistant/summary_assistant)、ToolCall、Message、ActionMessage(role=assistant)、StateMessage(role=user,带回信)、Reward、Sample、Trajectory。ActionMessage/StateMessage 是 env.step 的输入输出载体,simple_dict 负责序列化成 REST 友好字典。5.2 agentevolver/schema(训练侧数据模型)
task.py:Task(task_id / env_type / open_query / query / ground_truth / evaluator)——ground_truth 与 evaluator 是 TaskManager 第 2 阶段合成的;TaskObjective(task+confidence+reward)用于探索期。
trajectory.py:Reward(outcome / success_rate / madness / description)——恰是 AgentFlow 算分打包的对象;Trajectory(steps/query/is_terminated/reward/metadata);Sample(RL 数据集正样本:messages + 全套 input_ids/prompt_ids/response_ids + 各 loss_mask + position_ids + reward_scores)。truncate_output_ids 在 group_tokenize 后裁超长 response(prompt 超长则直接 raise,逼人调 max_prompt_length)。
二者桥接:agent_flow 用 agentevolver.schema 的 Reward/Sample/Trajectory;env_service 用 env_service.trajectory 的 ActionMessage/StateMessage。训练进程通过 EnvClient 的 REST JSON 在两者之间搬运,互不打扰。---
六、端到端时序(一回合)
___CODE_BLOCK_9___python ...___CODE_BLOCK_10___
convert_tool_to_user_message(utils/utils.py)在 agent_flow 里把 role="tool" 的回信转成 qwen 格式 user 消息,确保多轮上下文里环境反馈以「用户」角色呈现——与 AppWorld 模板里 USER: Output: ... 的对话结构严丝合缝。---
七、小结
第 6 章的骨架:环境服务 =
EnvService(Ray 本地 + env_actors 字典 + 闲置回收)跑 FastAPI,六路由经 import_and_register_env 动态挂 BaseEnv 子类;每个实例是一个 @ray.remote RemoteEnv actor,天然隔离、可并发。EnvClient 用 retry_call + fallback 把环境故障挡在 RL 主循环之外。AppworldEnv 是默认环境——jinja2 模板注入「自主助理」人设,world.execute 当状态化 Python REPL,evaluate_task 出稀疏/稠密分。两份 schema 各管一侧:环境侧 ActionMessage/StateMessage、训练侧 Reward/Sample/Trajectory`。一句话:环境服务是 agent 的「外部世界」,用 HTTP 隔、用 Ray 隔、用 Registry 解耦;训练进程只管发请求、收回信,世界的重活在墙外 actor 里悄悄跑完。至此六章源码解读毕。自顶向下:工程骨架(1)→ 任务自生成(2)→ 经验导航(3)→ 步骤归因(4)→ 执行流与上下文(5)→ 环境服务与基建(6)。源码之下,AgentEvolver 之「自进化」全貌已现。