From ffa77e925132eb6d8fbfae1dc0397141d7a50494 Mon Sep 17 00:00:00 2001 From: Geekkevin23333 <9502067+Geekkevin23333@user.noreply.gitee.com> Date: Mon, 28 Sep 2026 10:59:24 +0800 Subject: [PATCH 1/2] feat: add R2E-Gym patch execution reward --- docs/r2e_patch_execution.md | 24 ++++++ env/r2e_patch_harness.py | 140 ++++++++++++++++++++++++++++++++ examples/r2e_gym_async_rl.py | 37 ++++++++- scripts/test_cpu_ci.sh | 4 +- tests/test_r2e_patch_harness.py | 34 ++++++++ workflow/r2e_gym.py | 84 ++++++++++++------- 6 files changed, 291 insertions(+), 32 deletions(-) create mode 100644 docs/r2e_patch_execution.md create mode 100644 env/r2e_patch_harness.py create mode 100644 tests/test_r2e_patch_harness.py diff --git a/docs/r2e_patch_execution.md b/docs/r2e_patch_execution.md new file mode 100644 index 0000000..a923f06 --- /dev/null +++ b/docs/r2e_patch_execution.md @@ -0,0 +1,24 @@ +# R2E-Gym patch execution + +The default R2E-Gym workflow generates an issue and scores its text. It does not execute a patch. Set `R2E_PATCH_EXECUTION=1` to switch the same workflow to unified-diff generation and container test reward. The model must return a `git diff` style patch. Each turn runs tests before and after applying the patch in a fresh task checkout. Reward is 1 only when the baseline fails and the patched tests pass; failed execution is returned as tool feedback for a possible revision. + +The public dataset index contains `repo_name`, `commit_hash`, and `docker_image`, but no repository checkout or test command. Supply these explicitly before enabling patch mode: + +```sh +export R2E_PATCH_EXECUTION=1 +export R2E_REPO_ROOT=/absolute/path/to/task-repositories +export R2E_TEST_COMMAND_JSON='["pytest", "-q"]' +export R2E_PATCH_TIMEOUT=300 +``` + +`R2E_REPO_ROOT/` must be a local Git repository containing the requested commit. `R2E_TEST_COMMAND_JSON` is an argument array run inside the dataset's Docker image with the isolated checkout mounted at `/workspace`. Use an image and test command appropriate to your R2E-Gym task; the dataset's `expected_output_json` is not itself an executable test specification. Docker must be available on each rollout worker. The harness disables container networking and drops Linux capabilities. Never point patch mode at a privileged Docker daemon when processing untrusted model output. + +For R2E-Gym `.sif` task images containing `/testbed` and `/r2e_tests`, use Singularity: + +```sh +export R2E_PATCH_BACKEND=singularity +export R2E_SIF_ROOT=/absolute/path/to/sif-images +export R2E_TEST_COMMAND_JSON='["bash", "/testbed/run_tests.sh"]' +``` + +The entrypoint resolves each image as `_.sif` under `R2E_SIF_ROOT`. Every rollout worker must have Singularity and access to these images. diff --git a/env/r2e_patch_harness.py b/env/r2e_patch_harness.py new file mode 100644 index 0000000..7318578 --- /dev/null +++ b/env/r2e_patch_harness.py @@ -0,0 +1,140 @@ +"""Execute a unified diff against an R2E-Gym repository in Docker.""" +from __future__ import annotations + +import subprocess +import tempfile +import shlex +import shutil +from pathlib import Path +from typing import Any, Sequence + + +class R2EPatchExecutionHarness: + requires_repo_path = True + def __init__(self, timeout: int = 300): + self.timeout = max(1, int(timeout)) + + def execute( + self, + patch: str, + workspace: str | Path, + test_command: Sequence[str], + docker_image: str, + commit_hash: str, + ) -> dict[str, Any]: + """Apply a patch to an isolated commit checkout, then run tests in Docker. + + The caller must supply a trusted image and command. Model output is only + used as patch data; it never becomes a shell command. + """ + source = Path(workspace).resolve() + if not source.is_dir() or not (source / ".git").exists(): + raise ValueError(f"workspace is not a git checkout: {source}") + if not patch.strip() or not docker_image or not commit_hash or not test_command: + raise ValueError("patch, docker_image, commit_hash and test_command are required") + with tempfile.TemporaryDirectory(prefix="r2e-patch-") as td: + checkout = Path(td) / "repo" + clone = subprocess.run(["git", "clone", "--quiet", "--no-hardlinks", "--", str(source), str(checkout)], + capture_output=True, text=True, timeout=self.timeout) + if clone.returncode: + raise RuntimeError(f"cannot clone task repository: {clone.stderr}") + reset = subprocess.run(["git", "checkout", "--detach", commit_hash], cwd=checkout, + capture_output=True, text=True, timeout=self.timeout) + if reset.returncode: + raise ValueError(f"commit is unavailable in workspace: {commit_hash}") + patch_file = Path(td) / "candidate.patch" + patch_file.write_text(patch, encoding="utf-8") + command = ["docker", "run", "--rm", "--network", "none", "--cap-drop", "ALL", + "-v", f"{checkout}:/workspace", "-w", "/workspace", docker_image, + *test_command] + try: + baseline = subprocess.run(command, capture_output=True, text=True, timeout=self.timeout) + except subprocess.TimeoutExpired: + raise RuntimeError("R2E baseline test timed out") + check = subprocess.run(["git", "apply", "--check", "--", str(patch_file)], cwd=checkout, + capture_output=True, text=True, timeout=self.timeout) + if check.returncode: + return {"patch_applied": False, "tests_passed": False, + "returncode": check.returncode, "stdout": check.stdout, "stderr": check.stderr} + apply = subprocess.run(["git", "apply", "--", str(patch_file)], cwd=checkout, + capture_output=True, text=True, timeout=self.timeout) + if apply.returncode: + return {"patch_applied": False, "tests_passed": False, + "returncode": apply.returncode, "stdout": apply.stdout, "stderr": apply.stderr} + try: + result = subprocess.run(command, capture_output=True, text=True, timeout=self.timeout) + except subprocess.TimeoutExpired: + return {"patch_applied": True, "tests_passed": False, "returncode": None, + "stdout": "", "stderr": "test timeout"} + return {"patch_applied": True, "baseline_failed": baseline.returncode != 0, + "tests_passed": baseline.returncode != 0 and result.returncode == 0, + "returncode": result.returncode, "stdout": result.stdout[-10000:], + "stderr": result.stderr[-10000:]} + + +class R2ESingularityPatchHarness: + """Execute patches in R2E-Gym SIF images containing /testbed.""" + + requires_repo_path = False + + def __init__(self, timeout: int = 300): + self.timeout = max(1, int(timeout)) + + def execute( + self, + patch: str, + workspace: str | Path | None, + test_command: Sequence[str], + docker_image: str, + commit_hash: str, + ) -> dict[str, Any]: + del workspace, commit_hash + image = Path(docker_image).resolve() + if not image.is_file() or image.suffix != ".sif": + raise ValueError(f"R2E SIF image is unavailable: {image}") + if not patch.strip() or not test_command: + raise ValueError("patch and test_command are required") + if shutil.which("singularity") is None: + raise RuntimeError("Singularity is unavailable on this worker") + with tempfile.TemporaryDirectory(prefix="r2e-sif-") as td: + sandbox = Path(td) / "sandbox" + build = subprocess.run( + ["singularity", "build", "--sandbox", str(sandbox), str(image)], + capture_output=True, text=True, timeout=self.timeout, + ) + if build.returncode: + raise RuntimeError(f"cannot build R2E sandbox: {build.stderr[-3000:]}") + testbed = sandbox / "testbed" + if not (testbed / ".git").exists(): + raise ValueError("R2E image does not contain a /testbed git checkout") + patch_path = testbed / ".r2e_candidate.patch" + patch_path.write_text(patch, encoding="utf-8") + tests_link = testbed / "r2e_tests" + if not tests_link.exists() and (sandbox / "r2e_tests").exists(): + tests_link.symlink_to("/r2e_tests") + prefix = ["singularity", "exec", "--writable", "--no-home", str(sandbox), + "bash", "-lc"] + setup = "export GIT_CONFIG_GLOBAL=/tmp/r2e_gitconfig; git config --global --add safe.directory /testbed; cd /testbed && " + command = setup + shlex.join(test_command) + try: + baseline = subprocess.run(prefix + [command], capture_output=True, + text=True, timeout=self.timeout) + except subprocess.TimeoutExpired: + raise RuntimeError("R2E baseline test timed out") + for action in ("git apply --check .r2e_candidate.patch", "git apply .r2e_candidate.patch"): + result = subprocess.run(prefix + [setup + action], capture_output=True, + text=True, timeout=self.timeout) + if result.returncode: + return {"patch_applied": False, "tests_passed": False, + "returncode": result.returncode, "stdout": result.stdout[-10000:], + "stderr": result.stderr[-10000:]} + try: + result = subprocess.run(prefix + [command], capture_output=True, + text=True, timeout=self.timeout) + except subprocess.TimeoutExpired: + return {"patch_applied": True, "tests_passed": False, + "returncode": None, "stdout": "", "stderr": "test timeout"} + return {"patch_applied": True, "baseline_failed": baseline.returncode != 0, + "tests_passed": baseline.returncode != 0 and result.returncode == 0, + "returncode": result.returncode, "stdout": result.stdout[-10000:], + "stderr": result.stderr[-10000:]} diff --git a/examples/r2e_gym_async_rl.py b/examples/r2e_gym_async_rl.py index 7995a0c..1a4b78b 100644 --- a/examples/r2e_gym_async_rl.py +++ b/examples/r2e_gym_async_rl.py @@ -4,6 +4,7 @@ import json import os +import shutil import sys from collections import Counter from datetime import timedelta @@ -20,6 +21,10 @@ from RL_Framework import AsyncRLTrainer, parse_args_and_load_config from RL_Framework.engine.device_utils import distributed_backend, set_device from RL_Framework.env.r2e_gym_reward import r2e_gym_reward_fn +from RL_Framework.env.r2e_patch_harness import ( + R2EPatchExecutionHarness, + R2ESingularityPatchHarness, +) from RL_Framework.workflow.r2e_gym import R2EGymWorkflow @@ -83,6 +88,21 @@ def main(): print(f"Scheduler: {config.heterogeneous_rollout.scheduling.scheduler_type}") dataset = load_dataset("json", data_files=data_path, split="train") + patch_mode = os.environ.get("R2E_PATCH_EXECUTION", "0") == "1" + patch_backend = os.environ.get("R2E_PATCH_BACKEND", "docker").lower() + if patch_mode: + if patch_backend not in {"docker", "singularity"}: + raise ValueError("R2E_PATCH_BACKEND must be docker or singularity") + if shutil.which(patch_backend) is None: + raise RuntimeError(f"R2E patch mode requires {patch_backend} on this worker") + repo_root = os.environ.get("R2E_REPO_ROOT") + sif_root = os.environ.get("R2E_SIF_ROOT") + test_command_json = os.environ.get("R2E_TEST_COMMAND_JSON") + if (patch_backend == "docker" and not repo_root) or (patch_backend == "singularity" and not sif_root) or not test_command_json: + raise ValueError("Patch mode requires R2E_TEST_COMMAND_JSON and the selected backend's repo/image root") + test_command = json.loads(test_command_json) + if not isinstance(test_command, list) or not test_command or not all(isinstance(part, str) for part in test_command): + raise ValueError("R2E_TEST_COMMAND_JSON must be a nonempty JSON string array") def preprocess(example): prompt = (example.get("prompt") or "").strip() @@ -92,7 +112,7 @@ def preprocess(example): or "" ).strip() prompt_id = f"{example.get('repo_name', 'repo')}:{example.get('commit_hash', '')}" - return { + result = { "prompt_id": prompt_id, "prompt": prompt, "task_text": target_issue, @@ -103,6 +123,18 @@ def preprocess(example): "expected_output_json": example.get("expected_output_json", "{}"), "modified_files": _normalize_modified_files(example.get("modified_files")), } + if patch_mode: + repo_name = example.get("repo_name", "") + if not repo_name or os.path.basename(repo_name) != repo_name: + raise ValueError(f"Invalid repo_name: {repo_name!r}") + if patch_backend == "docker": + result["repo_path"] = os.path.join(repo_root, repo_name) + else: + result["docker_image"] = os.path.join( + sif_root, f"{repo_name}_{str(example.get('commit_hash', ''))[:8]}.sif" + ) + result["test_command"] = test_command + return result dataset = dataset.map(preprocess) if is_main_process: @@ -134,6 +166,9 @@ def preprocess(example): temperature=config.temperature, top_p=config.top_p, n_samples=config.n_samples, + patch_harness=(R2ESingularityPatchHarness if patch_backend == "singularity" else R2EPatchExecutionHarness)( + timeout=int(os.environ.get("R2E_PATCH_TIMEOUT", "300")) + ) if patch_mode else None, ) trainer = AsyncRLTrainer(config) diff --git a/scripts/test_cpu_ci.sh b/scripts/test_cpu_ci.sh index d13d0a4..c3ceb95 100644 --- a/scripts/test_cpu_ci.sh +++ b/scripts/test_cpu_ci.sh @@ -11,4 +11,6 @@ python -m pytest -q --strict-markers --junitxml=reports/cpu-tests.xml \ tests/test_cmlfq_cost_scheduler.py \ tests/test_hetero_cmlfq_integration.py \ tests/test_rollout_engine.py \ - tests/test_cpu_offload_backend.py + tests/test_cpu_offload_backend.py \ + tests/test_r2e_patch_harness.py \ + tests/test_r2e_gym_workflow.py diff --git a/tests/test_r2e_patch_harness.py b/tests/test_r2e_patch_harness.py new file mode 100644 index 0000000..09b26a4 --- /dev/null +++ b/tests/test_r2e_patch_harness.py @@ -0,0 +1,34 @@ +import subprocess +from pathlib import Path + +from env.r2e_patch_harness import R2EPatchExecutionHarness + + +def test_patch_harness_uses_isolated_checkout_and_docker(tmp_path: Path, monkeypatch): + repo = tmp_path / "source" + repo.mkdir() + subprocess.run(["git", "init", "-q"], cwd=repo, check=True) + (repo / "value.py").write_text("VALUE = 1\n") + subprocess.run(["git", "add", "."], cwd=repo, check=True) + subprocess.run(["git", "-c", "user.name=t", "-c", "user.email=t@t", "commit", "-qm", "base"], cwd=repo, check=True) + commit = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=repo, text=True).strip() + (repo / "value.py").write_text("VALUE = 2\n") + patch = subprocess.check_output(["git", "diff", "--", "value.py"], cwd=repo, text=True) + subprocess.run(["git", "checkout", "--", "value.py"], cwd=repo, check=True) + real_run = subprocess.run + calls = 0 + + def run(command, **kwargs): + nonlocal calls + if command[0] == "docker": + calls += 1 + mount = command[command.index("-v") + 1].split(":/workspace")[0] + expected = "VALUE = 1\n" if calls == 1 else "VALUE = 2\n" + assert (Path(mount) / "value.py").read_text() == expected + return subprocess.CompletedProcess(command, 1 if calls == 1 else 0, "", "") + return real_run(command, **kwargs) + + monkeypatch.setattr(subprocess, "run", run) + result = R2EPatchExecutionHarness().execute(patch, repo, ["pytest", "-q"], "example:test", commit) + assert result["patch_applied"] and result["baseline_failed"] and result["tests_passed"] + assert (repo / "value.py").read_text() == "VALUE = 1\n" diff --git a/workflow/r2e_gym.py b/workflow/r2e_gym.py index a9d67fc..f7cda55 100644 --- a/workflow/r2e_gym.py +++ b/workflow/r2e_gym.py @@ -10,6 +10,10 @@ import torch from transformers import PreTrainedTokenizerBase +from RL_Framework.env.r2e_patch_harness import ( + R2EPatchExecutionHarness, + R2ESingularityPatchHarness, +) from RL_Framework.env.r2e_gym_reward import ( evaluate_issue, @@ -44,6 +48,7 @@ def __init__( temperature: float = 1.0, top_p: float = 1.0, n_samples: int = 1, + patch_harness: R2EPatchExecutionHarness | R2ESingularityPatchHarness | None = None, ): self.reward_fn = reward_fn self.tokenizer = tokenizer @@ -59,6 +64,7 @@ def __init__( self.temperature = temperature self.top_p = top_p self.n_samples = n_samples + self.patch_harness = patch_harness def _encode(self, text: str) -> list[int]: return self.tokenizer.encode(text or "", add_special_tokens=False) @@ -120,14 +126,43 @@ def _clip_tokens_to_sequence_budget( return tokens[:remaining], logprobs[:remaining] def _build_initial_prompt(self, row: dict[str, Any]) -> str: - task_prompt = row.get("prompt") or row.get("problem_statement") or row.get("task_text") or "" + task_prompt = ( + row.get("problem_statement") or row.get("task_text") or "" + if self.patch_harness + else row.get("prompt") or row.get("problem_statement") or row.get("task_text") or "" + ) task_prompt = self._truncate_text_to_tokens(task_prompt, self.max_prompt_tokens) + system_prompt = ( + "Return only a unified git diff that fixes the reported issue. " + "Start with diff --git. Do not use markdown fences." + if self.patch_harness else SYSTEM_PROMPT + ) messages = [ - {"role": "system", "content": SYSTEM_PROMPT}, + {"role": "system", "content": system_prompt}, {"role": "user", "content": task_prompt}, ] return self._apply_chat_template(messages) + async def _score(self, completion: str, row: dict[str, Any]) -> dict[str, Any]: + if self.patch_harness: + required = ["test_command", "docker_image"] + if self.patch_harness.requires_repo_path: + required.extend(("repo_path", "commit_hash")) + missing = [key for key in required if not row.get(key)] + if missing: + raise ValueError(f"R2E patch mode missing task fields: {missing}") + execution = await asyncio.to_thread( + self.patch_harness.execute, completion, row.get("repo_path"), + row["test_command"], row["docker_image"], row.get("commit_hash", ""), + ) + return {"reward": float(execution["tests_passed"]), "execution": execution} + return evaluate_issue( + completion=completion, + target_issue=row.get("target_issue", row.get("task_text", "")), + expected_output_json=row.get("expected_output_json"), + modified_files=row.get("modified_files"), + ) + async def run_episode( self, engine: Any, @@ -176,10 +211,8 @@ async def run_episode( "temperature": self.temperature, "top_p": self.top_p, "n": 1, - # The output contract has an explicit closing delimiter. - # Stopping there prevents repetitive tails from consuming - # the remaining 30K generation budget. - "stop": ["[/ISSUE]"], + # Issue generation has an explicit closing delimiter. + "stop": [] if self.patch_harness else ["[/ISSUE]"], "include_stop_str_in_output": True, "seed": int.from_bytes( hashlib.sha256( @@ -210,16 +243,14 @@ async def run_episode( if output_tokens: segments.append((output_tokens, output_logprobs, 1)) - metrics = evaluate_issue( - completion=output_text, - target_issue=data.get("target_issue", data.get("task_text", "")), - expected_output_json=data.get("expected_output_json"), - modified_files=data.get("modified_files"), - ) + metrics = await self._score(output_text, data) if turn >= self.max_turns - 1 or metrics["reward"] >= self.stop_reward: break - feedback = format_validator_feedback(metrics) + feedback = ( + f"### Patch execution\n{metrics['execution']}\nRevise the patch." + if self.patch_harness else format_validator_feedback(metrics) + ) feedback_tokens = self._encode(feedback) used_after_output = sum(len(tokens) for tokens, _, _ in segments) feedback_tokens, feedback_logprobs = self._clip_tokens_to_sequence_budget( @@ -233,7 +264,7 @@ async def run_episode( segments.append((feedback_tokens, feedback_logprobs, 0)) generated_tokens = sum(len(tokens) for tokens, _, _ in segments[1:]) tool_event = { - "tool_type": "r2e_issue_validator", + "tool_type": "r2e_patch_executor" if self.patch_harness else "r2e_issue_validator", "output": feedback, "status": "success" if metrics["reward"] >= self.stop_reward else "failure", "payload_tokens": len(feedback_tokens), @@ -278,7 +309,7 @@ async def run_episode( finish_cmlfq(cmlfq_request_id, total_output_tokens) full_completion = self.tokenizer.decode(all_input_ids, skip_special_tokens=True) - reward = self.reward_fn( + reward = metrics["reward"] if self.patch_harness else self.reward_fn( prompt=prompt_text, completion=final_completion or full_completion, target_issue=data.get("target_issue", data.get("task_text", "")), @@ -380,7 +411,7 @@ async def eval_one(index: int) -> dict[str, Any]: "top_p": 1.0, "n": 1, "prompt_id": prompt_id, - "stop": ["[/ISSUE]"], + "stop": [] if self.patch_harness else ["[/ISSUE]"], "include_stop_str_in_output": True, } if request_id: @@ -388,12 +419,7 @@ async def eval_one(index: int) -> dict[str, Any]: response = await engine.generate(**generate_kwargs) completion = response.get("text", "") generated_tokens += len(self._encode(completion)) - metrics = evaluate_issue( - completion=completion, - target_issue=row.get("target_issue", row.get("task_text", "")), - expected_output_json=row.get("expected_output_json"), - modified_files=row.get("modified_files"), - ) + metrics = await self._score(completion, row) turn_metrics.append(metrics) if ( not use_feedback @@ -401,7 +427,10 @@ async def eval_one(index: int) -> dict[str, Any]: or metrics["reward"] >= self.stop_reward ): break - feedback = format_validator_feedback(metrics) + feedback = ( + f"### Patch execution\n{metrics['execution']}\nRevise the patch." + if self.patch_harness else format_validator_feedback(metrics) + ) feedback_tokens = len(self._encode(feedback)) generated_tokens += feedback_tokens route_tool_return = getattr(engine, "route_cmlfq_tool_return", None) @@ -409,7 +438,7 @@ async def eval_one(index: int) -> dict[str, Any]: route_tool_return( request_id, { - "tool_type": "r2e_issue_validator", + "tool_type": "r2e_patch_executor" if self.patch_harness else "r2e_issue_validator", "output": feedback, "status": ( "success" @@ -427,12 +456,7 @@ async def eval_one(index: int) -> dict[str, Any]: finish_cmlfq = getattr(engine, "finish_cmlfq_request", None) if request_id and callable(finish_cmlfq): finish_cmlfq(request_id, generated_tokens) - metrics = evaluate_issue( - completion=completion, - target_issue=row.get("target_issue", row.get("task_text", "")), - expected_output_json=row.get("expected_output_json"), - modified_files=row.get("modified_files"), - ) + metrics = turn_metrics[-1] if turn_metrics else await self._score(completion, row) return { "ok": True, "index": int(index), From 2a6e20f367868793e5b1ecf6967b6697874ee8ce Mon Sep 17 00:00:00 2001 From: Geekkevin23333 <9502067+Geekkevin23333@user.noreply.gitee.com> Date: Mon, 28 Sep 2026 11:01:38 +0800 Subject: [PATCH 2/2] fix: avoid tokenizer package import during CPU workflow tests --- workflow/r2e_gym.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/workflow/r2e_gym.py b/workflow/r2e_gym.py index f7cda55..d2f0a5c 100644 --- a/workflow/r2e_gym.py +++ b/workflow/r2e_gym.py @@ -6,10 +6,12 @@ import hashlib import os import random -from typing import Any, Callable +from typing import TYPE_CHECKING, Any, Callable import torch -from transformers import PreTrainedTokenizerBase + +if TYPE_CHECKING: + from transformers import PreTrainedTokenizerBase from RL_Framework.env.r2e_patch_harness import ( R2EPatchExecutionHarness, R2ESingularityPatchHarness,