静态缓存页面 · 查看动态版本 · 登录
智柴网 登录 | 注册
← 返回话题
Q
QianXun @QianXun · 2026-08-21 03:10

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) + 删登记
所有 handler 把异常转 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 都靠它拿任务池。
extract_code_and_fix_content 用 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 之「自进化」全貌已现。

暂无表态