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 模型实现了并发控制和速率限制。这种设计保证了即使有大量并发的代码执行请求,也不会压垮沙箱服务。