跳转至

sandbox_fusion_tools.py — verl/tools/sandbox_fusion_tools.py

文件路径

verl/tools/sandbox_fusion_tools.py

文件概述

SandboxFusionTool 是一个代码执行工具,允许 LLM 在安全的沙箱环境中执行代码。它使用 Sandbox Fusion 服务来运行代码,支持 Python 等语言。这个文件还包含了基于 Ray 的并发执行框架(速率限制、执行池),这个框架也被其他工具文件复用。

应用场景:在 RL 训练中,LLM 可能需要编写并执行代码来解决编程题、数学题等。这个工具提供了安全的代码执行能力。

关键代码讲解

1. 令牌桶速率限制器(Ray Actor)

@ray.remote(concurrency_groups={"acquire": 1, "release": 10})
class TokenBucketWorker:
    def __init__(self, rate_limit: int):
        self.rate_limit = rate_limit
        self.current_count = 0
        self._semaphore = threading.Semaphore(rate_limit)

    @ray.method(concurrency_group="acquire")
    def acquire(self):
        self._semaphore.acquire()
        self.current_count += 1

    @ray.method(concurrency_group="release")
    def release(self):
        self._semaphore.release()
        self.current_count -= 1

这是一个 Ray Actor(分布式对象),用于全局速率限制: - 使用信号量(Semaphore)控制并发数。 - acquire 方法的并发组限制为 1(串行获取令牌),release 限制为 10(可并行释放)。 - current_count 用于监控当前有多少令牌被占用。

2. 执行工作器

class ExecutionWorker:
    def __init__(self, enable_global_rate_limit=True, rate_limit=10):
        self.rate_limit_worker = self._init_rate_limit(rate_limit) if enable_global_rate_limit else None

    def _init_rate_limit(self, rate_limit):
        return TokenBucketWorker.options(name="rate-limiter", get_if_exists=True).remote(rate_limit)

    def execute(self, fn: Callable[..., T], *fn_args, **fn_kwargs) -> T:
        with ExitStack() as stack:
            stack.callback(self.rate_limit_worker.release.remote)
            ray.get(self.rate_limit_worker.acquire.remote())
            try:
                return fn(*fn_args, **fn_kwargs)
            except Exception as e:
                logger.warning(f"Error when executing code: {e}")

执行工作器封装了"先获取令牌 -> 执行任务 -> 释放令牌"的流程。使用 ExitStack 确保即使执行出错,令牌也会被释放。get_if_exists=True 确保全局只有一个速率限制器(单例模式)。

3. 初始化执行池

def init_execution_pool(num_workers, enable_global_rate_limit=True, rate_limit=10, mode=PoolMode.ThreadMode):
    if mode == PoolMode.ThreadMode:
        return (
            ray.remote(ExecutionWorker)
            .options(max_concurrency=num_workers)
            .remote(enable_global_rate_limit=enable_global_rate_limit, rate_limit=rate_limit)
        )

将 ExecutionWorker 变成 Ray Actor,设置最大并发数为 num_workers。这样可以并行处理多个代码执行请求。

4. SandboxFusionTool 构造函数

class SandboxFusionTool(BaseTool):
    def __init__(self, config: dict, tool_schema: OpenAIFunctionToolSchema):
        super().__init__(config, tool_schema)
        self._instance_dict = {}
        self.num_workers = config.get("num_workers", 10)
        self.rate_limit = config.get("rate_limit", 10)
        self.default_timeout = config.get("default_timeout", 30)
        self.default_language = config.get("default_language", "python")
        self.sandbox_fusion_url = config.get("sandbox_fusion_url", "")
        self.memory_limit_mb = config.get("memory_limit_mb", 1024)
        # ... 初始化执行池 ...

可配置参数包括:工作线程数、速率限制、超时时间、默认编程语言、沙箱服务地址、内存限制。

5. 执行代码

@rollout_trace_op
async def execute(self, instance_id, parameters, **kwargs):
    code = parameters.get("code", "")
    timeout = parameters.get("timeout", self.default_timeout)
    language = parameters.get("language", self.default_language)
    result = await self.execution_pool.execute.remote(self.execute_code, instance_id, code, timeout, language)
    return ToolResponse(text=None if result is None else str(result)), None, None

def execute_code(self, instance_id, code, timeout=30, language="python"):
    result_status, metadata = _process_single_case(
        0, None, None, self.sandbox_fusion_url, code, timeout, self.memory_limit_mb, language
    )
    if metadata["run_status"] == "Finished":
        actual_output = metadata["stdout"] + metadata["stderr"]
        return ToolResponse(text=actual_output)
    else:
        return ToolResponse(text="no stdout here")

执行流程: 1. 从参数中提取代码、超时时间、编程语言。 2. 通过执行池(带速率限制)提交到 Sandbox Fusion 服务。 3. _process_single_case 是实际调用沙箱 API 的函数。 4. 如果执行成功,返回标准输出 + 错误输出;否则返回提示信息。

核心类/函数列表

类/函数 作用
PoolMode 执行池模式枚举(线程/进程)
TokenBucketWorker Ray Actor,令牌桶速率限制器
ExecutionWorker 带速率限制的任务执行器
init_execution_pool 初始化 Ray 执行池
SandboxFusionTool 沙箱代码执行工具
execute_code 实际调用沙箱服务执行代码

与其他模块的关系

  • 继承自 BaseTool。
  • 依赖 verl.utils.reward_score.sandbox_fusion.utils._process_single_case 调用沙箱 API。
  • 并发框架被复用:search_tool.py 和 image_zoom_in_tool.py 都有类似的 TokenBucketWorker 和执行池设计(代码结构几乎相同)。
  • 通过 tool_registry.py 被注册和初始化。

小结

SandboxFusionTool 提供了安全的远程代码执行能力,并通过 Ray 的 Actor 模型实现了并发控制和速率限制。这种设计保证了即使有大量并发的代码执行请求,也不会压垮沙箱服务。