跳转至

agent_loop.py — 这是 agent_loop 模块的核心文件

文件路径: verl/experimental/agent_loop/agent_loop.py

文件概述

这是 agent_loop 模块的核心文件,定义了整个 Agent 循环系统的基础架构。它包含以下关键组件:

  1. AsyncLLMServerManager - 异步 LLM 推理服务器的负载均衡管理器
  2. AgentLoopOutput / _InternalAgentLoopOutput - Agent 循环输出的数据结构
  3. AgentLoopBase - Agent 循环的抽象基类
  4. AgentLoopWorker - 执行批量推理的工作节点
  5. AgentLoopManager - 管理多个 Worker 的调度器

整个文件大约 900+ 行,是理解 verl 如何进行异步推理(rollout)的关键入口。

关键代码讲解

1. AsyncLLMServerManager - LLM 服务器管理器

这个类管理多个 OpenAI 兼容的 LLM 推理服务器,提供负载均衡和会话粘滞功能。

class AsyncLLMServerManager:
    def __init__(self, config, server_handles, max_cache_size=10000):
        self.server_handles = server_handles
        random.shuffle(self.server_handles)

        # 最少请求数负载均衡(最小堆)
        self.weighted_serveres = [[0, idx, server] for idx, server in enumerate(self.server_handles)]
        heapq.heapify(self.weighted_serveres)

        # LRU 缓存:将 request_id 映射到同一服务器(会话粘滞)
        self.request_id_to_server = LRUCache(maxsize=max_cache_size)

负载均衡策略:使用最小堆实现"最少请求"负载均衡——每次选择当前请求数最少的服务器。

会话粘滞:多轮对话中,同一个 request_id 会被路由到同一台服务器,以利用前缀缓存(Prefix Caching)加速。

def _choose_server(self, request_id):
    if request_id in self.request_id_to_server:
        return self.request_id_to_server[request_id]
    _, _, server = self.weighted_serveres[0]
    self.weighted_serveres[0][0] += 1
    heapq.heapreplace(self.weighted_serveres, self.weighted_serveres[0])
    self.request_id_to_server[request_id] = server
    return server

2. AgentLoopOutput - 输出数据结构

Agent 循环的标准输出格式:

class AgentLoopOutput(BaseModel):
    prompt_ids: list[int]        # 提示词 token IDs
    response_ids: list[int]      # 响应 token IDs(包括 LLM 生成和工具响应)
    response_mask: list[int]     # 响应掩码:1=LLM生成,0=工具响应/填充
    response_logprobs: Optional[list[float]] = None  # 对数概率
    routed_experts: Optional[Any] = None             # MoE 路由专家
    multi_modal_data: Optional[dict] = None          # 多模态数据
    reward_score: Optional[float] = None             # 奖励分数
    num_turns: int = 0           # 对话轮数
    metrics: AgentLoopMetrics    # 性能指标
    extra_fields: dict = {}      # 额外字段

response_mask 的含义(多轮对话场景):

responses:     |<- LLM生成 ->|<- 工具调用 ->|<- LLM生成 ->|<- 填充 ->|
response_mask: | 1, 1, ..., 1 | 0, 0, ..., 0 | 1, 1, ..., 1 | 0, ..., 0|
  • 1 表示 LLM 自己生成的 token(会参与策略梯度计算)
  • 0 表示工具返回的 token 或填充(不参与梯度计算)

3. AgentLoopBase - Agent 循环基类

所有 Agent 循环实现的抽象基类:

class AgentLoopBase(ABC):
    def __init__(self, trainer_config, server_manager, tokenizer, processor, dataset_cls, data_config):
        self.config = trainer_config.config
        self.rollout_config, _ = _get_rollout_and_model_config(self.config)
        self.server_manager = server_manager
        self.tokenizer = tokenizer
        self.processor = processor  # 用于多模态(VLM)
        # ...

    @abstractmethod
    async def run(self, sampling_params, **kwargs) -> AgentLoopOutput:
        """子类必须实现:执行一次完整的 Agent 循环"""
        raise NotImplementedError

基类提供了两个关键工具方法:

  • process_vision_info() - 从消息中提取图片/视频等多模态数据
  • apply_chat_template() - 应用聊天模板,将消息转换为 token IDs

4. 注册机制

通过装饰器实现 Agent 循环的注册:

_agent_loop_registry: dict[str, dict] = {}

def register(agent_name: str):
    def decorator(subclass):
        fqdn = f"{subclass.__module__}.{subclass.__qualname__}"
        _agent_loop_registry[agent_name] = {"_target_": fqdn}
        return subclass
    return decorator

注册后的 Agent 可以通过 Hydra 的 instantiate 动态创建:

agent_loop = hydra.utils.instantiate(
    config=_agent_loop_registry[agent_name],
    trainer_config=..., server_manager=..., ...
)

5. AgentLoopWorker - 工作节点

Worker 负责批量并行执行 Agent 循环:

class AgentLoopWorker:
    async def generate_sequences(self, batch: DataProto) -> DataProto:
        # 为 batch 中的每个样本创建异步任务
        tasks = []
        for i in range(len(batch)):
            kwargs = {k: v[i] for k, v in batch.non_tensor_batch.items()}
            tasks.append(asyncio.create_task(
                self._run_agent_loop(sampling_params, trajectory_info[i], **kwargs)
            ))
        # 并行执行所有 Agent 循环
        outputs = await asyncio.gather(*tasks)
        output = self._postprocess(outputs)
        return output

后处理流程(_agent_loop_postprocess): - prompt_ids 左填充(对齐到 prompt_length) - response_ids 右填充(对齐到 response_length) - 拼接得到 input_ids = prompt + response - 计算 attention_mask 和 position_ids

6. AgentLoopManager - 管理器

管理多个 Worker,负责分发批次数据:

class AgentLoopManager:
    @classmethod
    async def create(cls, config, worker_group=None, ...):
        instance = cls(config, worker_group, ...)
        await instance._initialize_llm_servers()  # 启动推理服务器
        await instance._init_agent_loop_workers()  # 创建 Worker
        return instance

支持两种模式: - 混合模式(Hybrid):推理服务器与训练引擎共享 GPU - 独立模式(Standalone):推理服务器使用独立 GPU(如 one-step-off、fully async)

核心类/函数列表

名称 类型 说明
AsyncLLMServerManager 类 LLM 服务器负载均衡管理器
AgentLoopMetrics 数据类 性能指标
AgentLoopOutput 数据类 Agent 循环标准输出
_InternalAgentLoopOutput 数据类 内部填充后的输出
AgentLoopBase 抽象基类 Agent 循环基类
AgentLoopWorker 类 批量推理工作节点
AgentLoopManager 类 Worker 调度管理器
register() 函数 Agent 循环注册装饰器
get_trajectory_info() 异步函数 获取轨迹追踪信息

与其他模块的关系

AgentLoopManager
  ├── 管理多个 AgentLoopWorker
  │     ├── 使用 AsyncLLMServerManager(管理 vLLM/SGLang 推理服务器)
  │     └── 实例化 AgentLoopBase 子类(SingleTurnAgentLoop / ToolAgentLoop)
  │           └── 产出 AgentLoopOutput
  ├── 被 RayPPOTrainer 调用(同步模式)
  ├── 被 OneStepOffRayTrainer 扩展(OneStepOffAgentLoopManager)
  └── 被 FullyAsyncRollouter 扩展(FullyAsyncAgentLoopManager)

小结

agent_loop.py 是 verl 异步推理架构的核心。它定义了一套完整的 Agent 循环框架:从服务器管理、负载均衡、Agent 注册机制,到批量并行执行和后处理。所有的 Agent 实现(单轮、工具调用、部分回滚等)都继承自这里定义的 AgentLoopBase,所有的训练器(同步、一步离策略、全异步)都通过 AgentLoopManager 或其子类来调度推理。