diff --git a/scripts/benchmark/rl/README.md b/scripts/benchmark/rl/README.md new file mode 100644 index 000000000..5001ca886 --- /dev/null +++ b/scripts/benchmark/rl/README.md @@ -0,0 +1,41 @@ +# G1 PPO CUDA benchmark + +Measure EmbodiChain PPO on **Dexsim Default CUDA** and **Dexsim Newton/MJWarp CUDA**. +The official G1 task and PPO configuration supply the environment, policy, and optimizer. + +Run each backend in a separate process with a new output directory: + +```bash +for backend in default newton; do + python -m scripts.benchmark.rl.locomotion_ppo \ + --backend "$backend" --output "outputs/g1-ppo/$backend" +done +``` + +The default workload uses 4096 environments, 24 steps per rollout, 3 warmup updates, +and 10 measured updates. Use `--envs 32 --warmup 1 --updates 3` for a smoke run. +`--task-dir` selects another G1 configuration directory; `--seed` and `--cpu-threads` +control the training seed and PyTorch CPU thread count. +Newton uses 20 solver iterations (`--solver-iterations`) and the task's 50 line-search +iterations. The selected environment configuration is saved as `env.yaml`. + +`result.json` contains synchronized rollout/optimizer timings, losses, backend and +runtime details, environment/policy creation time, and memory samples. PPO SPS is +total measured transitions divided by total rollout plus optimizer time. Creation +is timed separately; warmup, validation, memory sampling, and checkpoint writes +are excluded from PPO timings. `report.md` summarizes the run; +`model.pt` contains the policy, observation normalizers, optimizer, and trainer counters. +The Markdown report includes creation/rollout/update times, throughput and its +per-update coefficient of variation (CV), process RAM/VRAM peaks, and GPU allocator +memory. CV is the population standard deviation divided by the mean of per-update +SPS. Process VRAM is sampled at phase boundaries; allocator peaks are measured +inside each timed phase. +Process VRAM matches one worker PID in NVIDIA's host PID namespace. A worker in +a nested PID namespace requires readable host procfs to resolve that PID; +private or inaccessible procfs keeps the sample unknown and reports `n/a`. + +Run G1/Go2 reset and inference checks, then G1 training/checkpoint tests on both backends: + +```bash +pytest tests/gym/envs/tasks/test_locomotion_cuda.py --run-gpu -v --junitxml=outputs/locomotion-tests.xml +``` diff --git a/scripts/benchmark/rl/locomotion_ppo.py b/scripts/benchmark/rl/locomotion_ppo.py new file mode 100644 index 000000000..e5f0e7157 --- /dev/null +++ b/scripts/benchmark/rl/locomotion_ppo.py @@ -0,0 +1,402 @@ +# ---------------------------------------------------------------------------- +# Copyright (c) 2021-2026 DexForce Technology Co., Ltd. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# ---------------------------------------------------------------------------- +"""Measure G1 EmbodiChain PPO rollout and optimizer throughput on CUDA.""" + +from __future__ import annotations + +import argparse +from collections.abc import Callable +import importlib.metadata +import json +import math +from pathlib import Path +import statistics +import time +from typing import Any + +__all__ = ["main"] + +BACKENDS = { + "default": "Dexsim Default CUDA", + "newton": "Dexsim Newton/MJWarp CUDA", +} +TASK_DIR = ( + Path(__file__).resolve().parents[3] + / "embodichain_tasks/configs/tasks/locomotion/velocity/g1_flat" +) + + +def _load_config(task_dir: Path, backend: str, envs: int, steps: int) -> dict: + import yaml + + suffix = ".newton" if backend == "newton" else "" + config = yaml.safe_load((task_dir / f"agents/ppo{suffix}.yaml").read_text()) + config["trainer"].update( + gym_config=str((task_dir / f"env{suffix}.yaml").resolve()), + num_envs=envs, + buffer_size=steps, + ) + mini_batches = config["algorithm"]["cfg"]["num_mini_batches"] + if envs * steps < mini_batches or envs * steps % mini_batches: + raise ValueError("Rollout size must be divisible by num_mini_batches") + config["algorithm"]["cfg"]["batch_size"] = envs * steps // mini_batches + return config + + +def _build_trainer(runtime: Any, config: dict, output: Path) -> Any: + from embodichain.learning.rl.algo import build_algo + from embodichain.learning.rl.utils.trainer import Trainer + + algorithm = build_algo( + "ppo", config["algorithm"]["cfg"], runtime.policy, runtime.device + ) + return Trainer( + policy=runtime.policy, + env=runtime.env, + algorithm=algorithm, + buffer_size=config["trainer"]["buffer_size"], + batch_size=algorithm.cfg.batch_size, + writer=None, + eval_freq=0, + save_freq=0, + checkpoint_dir=str(output), + exp_name="g1-ppo-benchmark", + use_wandb=False, + ) + + +def _check_backend(env: Any, name: str) -> Any: + import torch + + sim = env.unwrapped.sim + if sim.physics.name != name or torch.device(sim.device).type != "cuda": + raise RuntimeError("Requested CUDA physics backend is not active") + if name == "newton": + from dexsim.engine.newton_physics.backend_registry import get_newton_backend + + backend = get_newton_backend(sim._world) + if ( + backend.solver.use_mujoco_cpu + or not backend.solver.mjw_data.qpos.device.is_cuda + ): + raise RuntimeError("Newton/MJWarp is not using CUDA") + return backend + if not (sim._world_config.enable_gpu_sim and sim._world_config.direct_gpu_api): + raise RuntimeError("Default backend is not using direct GPU simulation") + return None + + +def _measure(operation: Callable[[], Any]) -> tuple[Any, dict[str, float]]: + import psutil + import torch + + torch.cuda.synchronize() + process = psutil.Process() + cpu_before = process.memory_info().rss + gpu_before = torch.cuda.memory_allocated() + torch.cuda.reset_peak_memory_stats() + start = time.perf_counter() + value = operation() + torch.cuda.synchronize() + elapsed = time.perf_counter() - start + return value, { + "wall_s": elapsed, + "cpu_delta_mb": (process.memory_info().rss - cpu_before) / 2**20, + "gpu_delta_mb": (torch.cuda.memory_allocated() - gpu_before) / 2**20, + "peak_gpu_mb": torch.cuda.max_memory_allocated() / 2**20, + } + + +def _summarize(result: dict) -> dict: + measured = [row for row in result["updates"] if not row["warmup"]] + if not measured or result["status"] != "passed": + return {} + rollout = sum(row["rollout"]["wall_s"] for row in measured) + update = sum(row["ppo_update"]["wall_s"] for row in measured) + transitions = result["num_envs"] * result["steps"] * len(measured) + rates = [ + result["num_envs"] + * result["steps"] + / (row["rollout"]["wall_s"] + row["ppo_update"]["wall_s"]) + for row in measured + ] + samples = result.get("resources", {}).get("samples", []) + gpu_samples = [row["gpu_process_mib"] for row in samples] + return { + "rollout_sps": transitions / rollout, + "ppo_sps": transitions / (rollout + update), + "rollout_ms": rollout * 1000 / len(measured), + "update_ms": update * 1000 / len(measured), + "ppo_sps_cv_pct": statistics.pstdev(rates) / statistics.mean(rates) * 100, + "cpu_rss_peak_mib": max( + ( + max(row["cpu_rss_mib"], row["cpu_rss_lifetime_peak_mib"]) + for row in samples + ), + default=None, + ), + "gpu_process_sampled_peak_mib": ( + max(gpu_samples) if gpu_samples and None not in gpu_samples else None + ), + } + + +def _write_report(output: Path, result: dict) -> None: + summary = _summarize(result) + label = BACKENDS[result["backend"]] + runtime = result.get("runtime", {}) + + def fmt(value: float | None) -> str: + return "n/a" if value is None else f"{value:.3f}" + + lines = [ + "# G1 EmbodiChain PPO", + "", + f"{label}; {result['num_envs']} environments × {result['steps']} steps; " + f"{result['warmup']} warmup + {result['measured_updates']} measured updates.", + "", + f"GPU: {runtime.get('gpu', 'n/a')}; CUDA: {runtime.get('cuda', 'n/a')}; " + f"DexSim: {runtime.get('dexsim_commit', 'n/a')}; seed: {result.get('seed', 'n/a')}.", + "", + "Memory is in MiB. GPU delta/peak use the PyTorch allocator; " + "resources in result.json also record process VRAM at phase boundaries.", + "", + f"Process RAM (RSS peak): {fmt(summary.get('cpu_rss_peak_mib'))} MiB; " + f"process VRAM (sampled peak): {fmt(summary.get('gpu_process_sampled_peak_mib'))} MiB.", + "", + "## Time & Memory", + "", + "| Phase | cost_time_ms | cpu_delta_mb | gpu_delta_mb | peak_gpu_mb |", + "| --- | --- | --- | --- | --- |", + ] + for phase in ("create_runtime", "rollout", "ppo_update"): + if phase == "create_runtime": + rows = [result[phase]] if phase in result else [] + else: + rows = [row[phase] for row in result["updates"] if not row["warmup"]] + if summary and rows: + values = [ + statistics.mean(row["wall_s"] for row in rows) * 1000, + statistics.mean(row["cpu_delta_mb"] for row in rows), + statistics.mean(row["gpu_delta_mb"] for row in rows), + max(row["peak_gpu_mb"] for row in rows), + ] + lines.append( + f"| {phase} | " + " | ".join(f"{v:.3f}" for v in values) + " |" + ) + else: + lines.append(f"| {phase} | n/a | n/a | n/a | n/a |") + lines += [ + "", + "## Success & Other Metrics", + "", + "| Backend | Status | Success rate | Rollout SPS | PPO SPS | PPO SPS CV (%) |", + "| --- | --- | --- | --- | --- | --- |", + f"| {label} | {result['status']} | n/a | " + + ( + f"{summary['rollout_sps']:.1f} | {summary['ppo_sps']:.1f} | {summary['ppo_sps_cv_pct']:.2f}" + if summary + else "n/a | n/a | n/a" + ) + + " |", + "", + "## Leaderboard", + "", + "| Algorithm | Backend | Success rank |", + "| --- | --- | --- |", + f"| EmbodiChain PPO | {label} | n/a |", + "", + "This workload measures short training throughput; episode success is not evaluated.", + ] + (output / "report.md").write_text("\n".join(lines) + "\n") + + +def _run(args: argparse.Namespace, result: dict) -> None: + import torch + import yaml + import dexsim + from embodichain.lab.gym.utils.registration import ( + discover_task_packages, + execute_init_hooks, + ) + from embodichain.learning.rl.runtime import build_gym_policy_runtime + from scripts.benchmark.rl.resource_usage import ProcessResources + from scripts.benchmark.rl.runtime import set_random_seed + + device = torch.device("cuda:0") + if not torch.cuda.is_available(): + raise RuntimeError("CUDA is required") + torch.cuda.set_device(device) + torch.set_num_threads(args.cpu_threads) + set_random_seed(args.seed, device) + discover_task_packages() + execute_init_hooks() + config = _load_config(args.task_dir, args.backend, args.envs, args.steps) + env_config = yaml.safe_load(Path(config["trainer"]["gym_config"]).read_text()) + env_config["num_envs"] = args.envs + if args.backend == "newton": + env_config["physics_config"]["solver_cfg"][ + "iterations" + ] = args.solver_iterations + env_path = args.output / "env.yaml" + env_path.write_text(yaml.safe_dump(env_config, sort_keys=False)) + config["trainer"]["gym_config"] = str(env_path.resolve()) + result["configuration"] = { + "policy": config["policy"], + "algorithm": config["algorithm"], + "physics": env_config["physics_config"], + } + result["runtime"] = { + "dexsim_commit": getattr(dexsim, "__commit_id__", None), + "torch": torch.__version__, + "cuda": torch.version.cuda, + "gpu": torch.cuda.get_device_name(device), + "cpu_threads": torch.get_num_threads(), + "packages": { + name: importlib.metadata.version(name) + for name in ("newton", "warp-lang", "mujoco-warp") + }, + } + resources = ProcessResources(str(torch.cuda.get_device_properties(device).uuid)) + runtime = None + try: + resources.sample("before_create") + runtime, result["create_runtime"] = _measure( + lambda: build_gym_policy_runtime( + config, + device=device, + num_envs=args.envs, + headless=True, + renderer=config["trainer"]["renderer"], + gpu_id=0, + seed=args.seed, + config_dir=args.task_dir, + ) + ) + backend = _check_backend(runtime.env, args.backend) + trainer = _build_trainer(runtime, config, args.output) + trainer.algorithm.bind_schedule(total_updates=args.warmup + args.updates) + resources.sample("after_create") + for index in range(args.warmup + args.updates): + warming = index < args.warmup + if index == args.warmup: + before = [p.detach().clone() for p in runtime.policy.parameters()] + resources.sample("after_warmup") + if backend is not None and backend.cuda_graph_status != "captured": + raise RuntimeError( + "Newton CUDA graph was not captured during warmup" + ) + _, rollout = _measure(trainer._collect_rollout) + batch = trainer.buffer.get(flatten=False) + if not all( + bool(torch.isfinite(batch[key]).all()) + for key in ("obs", "critic_obs", "action", "reward") + ): + raise RuntimeError("Nonfinite rollout") + losses, update = _measure(lambda: trainer.algorithm.update(batch)) + trainer.num_updates += 1 + if not all(math.isfinite(float(value)) for value in losses.values()): + raise RuntimeError("Nonfinite PPO loss") + if backend is not None and bool( + backend.solver.mjw_data.overflow.numpy().any() + ): + raise RuntimeError("MJWarp capacity or solver iteration limit reached") + result["updates"].append( + { + "warmup": warming, + "rollout": rollout, + "ppo_update": update, + "losses": {k: float(v) for k, v in losses.items()}, + } + ) + resources.sample("warmup" if warming else f"measured_{index}") + delta = max( + float((old - new.detach()).abs().max()) + for old, new in zip(before, runtime.policy.parameters()) + ) + if ( + not math.isfinite(delta) + or delta <= 0 + or not all( + bool(torch.isfinite(p).all()) for p in runtime.policy.parameters() + ) + ): + raise RuntimeError("PPO did not produce finite parameter updates") + result["parameter_max_abs_delta"] = delta + trainer.save_checkpoint(str(args.output / "model.pt")) + finally: + result["resources"] = resources.result() + if runtime is not None: + runtime.close() + + +def main() -> None: + """Run one backend and save timings, a checkpoint, and a Markdown report.""" + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--backend", choices=BACKENDS, required=True) + parser.add_argument("--output", type=Path, required=True) + parser.add_argument("--task-dir", type=Path, default=TASK_DIR) + parser.add_argument("--envs", type=int, default=4096) + parser.add_argument("--steps", type=int, default=24) + parser.add_argument("--warmup", type=int, default=3) + parser.add_argument("--updates", type=int, default=10) + parser.add_argument("--seed", type=int, default=42) + parser.add_argument("--cpu-threads", type=int, default=4) + parser.add_argument( + "--solver-iterations", + type=int, + default=20, + help="Newton solver iteration limit", + ) + args = parser.parse_args() + if ( + min( + args.envs, + args.steps, + args.warmup, + args.updates, + args.cpu_threads, + args.solver_iterations, + ) + < 1 + ): + parser.error("Workload sizes and CPU thread count must be positive") + args.output.mkdir(parents=True, exist_ok=False) + result = { + "protocol": "g1-embodichain-ppo", + "backend": args.backend, + "num_envs": args.envs, + "steps": args.steps, + "warmup": args.warmup, + "measured_updates": args.updates, + "seed": args.seed, + "status": "failed", + "updates": [], + } + try: + _run(args, result) + result["status"] = "passed" + finally: + result["summary"] = _summarize(result) + (args.output / "result.json").write_text( + json.dumps(result, indent=2, allow_nan=False) + "\n" + ) + _write_report(args.output, result) + print(json.dumps(result["summary"], indent=2)) + + +if __name__ == "__main__": + main() diff --git a/scripts/benchmark/rl/resource_usage.py b/scripts/benchmark/rl/resource_usage.py new file mode 100644 index 000000000..e7801b04a --- /dev/null +++ b/scripts/benchmark/rl/resource_usage.py @@ -0,0 +1,133 @@ +# ---------------------------------------------------------------------------- +# Copyright (c) 2021-2026 DexForce Technology Co., Ltd. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# ---------------------------------------------------------------------------- +"""Sample process RAM and device-scoped process VRAM outside benchmark timing.""" + +from __future__ import annotations + +import os +from pathlib import Path +import subprocess +from typing import Any + +__all__ = ["ProcessResources"] + +# Linux PID_NS_INIT_INO identifies the initial PID namespace. +_INITIAL_PID_NS_INODE = 0xEFFFFFFC + + +def _nvidia_smi_pid() -> int | None: + """Resolve the worker PID only when procfs exposes NVIDIA's host namespace.""" + try: + if Path("/proc/self/ns/pid").stat().st_ino == _INITIAL_PID_NS_INODE: + return os.getpid() + if Path("/proc/1/ns/pid").stat().st_ino != _INITIAL_PID_NS_INODE: + return None + return int(Path("/proc/self").readlink().name) + except (OSError, ValueError): + return None + + +def _process_gpu_mib(output: str, pid: int | None, gpu_uuid: str) -> float | None: + """Use an exact GPU/PID match; missing or unavailable values stay unknown.""" + if pid is None: + return None + values = [] + for line in output.splitlines(): + parts = [part.strip() for part in line.split(",")] + if len(parts) != 3: + continue + process_pid, memory, uuid = parts + if uuid.lower().removeprefix("gpu-") != gpu_uuid.lower().removeprefix("gpu-"): + continue + try: + if int(process_pid) == pid: + value = float(memory) + if value >= 0 and value < float("inf"): + values.append(value) + except ValueError: + continue + return max(values) if values else None + + +class ProcessResources: + """Record phase-boundary samples for this Linux CUDA worker, in MiB.""" + + def __init__(self, gpu_uuid: str) -> None: + self.gpu_uuid = gpu_uuid + self.gpu_pid = _nvidia_smi_pid() + self.samples: list[dict[str, Any]] = [] + + def sample(self, label: str) -> dict[str, Any]: + """Synchronize, then sample process and allocator memory outside timing.""" + import psutil + import resource + import torch + + torch.cuda.synchronize() + item: dict[str, Any] = { + "label": label, + "cpu_rss_mib": psutil.Process().memory_info().rss / 2**20, + # Linux ru_maxrss is KiB and covers the worker lifetime so far. + "cpu_rss_lifetime_peak_mib": resource.getrusage( + resource.RUSAGE_SELF + ).ru_maxrss + / 1024, + "torch_allocated_mib": torch.cuda.memory_allocated() / 2**20, + "torch_reserved_mib": torch.cuda.memory_reserved() / 2**20, + "gpu_process_mib": None, + "gpu_error": None, + } + try: + pss = getattr(psutil.Process().memory_full_info(), "pss", None) + item["cpu_pss_mib"] = None if pss is None else pss / 2**20 + except (psutil.AccessDenied, OSError): + item["cpu_pss_mib"] = None + try: + query = subprocess.run( + [ + "nvidia-smi", + "--query-compute-apps=pid,used_gpu_memory,gpu_uuid", + "--format=csv,noheader,nounits", + ], + capture_output=True, + text=True, + check=True, + timeout=5, + ) + item["gpu_process_mib"] = _process_gpu_mib( + query.stdout, self.gpu_pid, self.gpu_uuid + ) + if item["gpu_process_mib"] is None: + item["gpu_error"] = ( + "Worker host PID is unresolved; not interpreted as zero" + if self.gpu_pid is None + else "No exact GPU/PID memory sample; not interpreted as zero" + ) + except (OSError, subprocess.SubprocessError) as error: + item["gpu_error"] = str(error) + self.samples.append(item) + return item + + def result(self) -> dict[str, Any]: + """Return provenance and all observed values, including missing samples.""" + return { + "schema": "process-memory", + "unit": "MiB", + "gpu_method": "nvidia-smi compute-process memory, exact UUID/PID", + "sampling": "phase boundaries outside timed regions", + "scope": "scene creation and training, sampled between phases", + "samples": self.samples, + } diff --git a/tests/benchmark/test_locomotion_ppo.py b/tests/benchmark/test_locomotion_ppo.py new file mode 100644 index 000000000..4f9043106 --- /dev/null +++ b/tests/benchmark/test_locomotion_ppo.py @@ -0,0 +1,284 @@ +# ---------------------------------------------------------------------------- +# Copyright (c) 2021-2026 DexForce Technology Co., Ltd. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# ---------------------------------------------------------------------------- +"""Locomotion workload, throughput aggregation, and memory/report tests.""" + +from __future__ import annotations + +import json +import os +from pathlib import Path +import sys +from types import SimpleNamespace + +import pytest +import torch + +from scripts.benchmark.rl import locomotion_ppo as benchmark +from scripts.benchmark.rl import resource_usage +from scripts.benchmark.rl.locomotion_ppo import ( + TASK_DIR, + _load_config, + _summarize, + _write_report, +) +from scripts.benchmark.rl.resource_usage import ProcessResources, _process_gpu_mib + + +def _result(status: str = "passed") -> dict: + def row(warming: bool, rollout: float, update: float) -> dict: + return { + "warmup": warming, + **{ + phase: { + "wall_s": duration, + "cpu_delta_mb": 1, + "gpu_delta_mb": 2, + "peak_gpu_mb": 3, + } + for phase, duration in (("rollout", rollout), ("ppo_update", update)) + }, + } + + return { + "backend": "default", + "status": status, + "num_envs": 32, + "steps": 24, + "warmup": 1, + "measured_updates": 2, + "create_runtime": row(False, 1000, 0)["rollout"], + "updates": [row(True, 100, 100), row(False, 1, 1), row(False, 3, 1)], + } + + +def test_throughput_uses_total_time_and_excludes_warmup() -> None: + summary = _summarize(_result()) + assert summary["ppo_sps"] == pytest.approx(32 * 24 * 2 / 6) + assert summary["rollout_sps"] == pytest.approx(32 * 24 * 2 / 4) + assert summary["rollout_ms"] == 2000 + assert summary["update_ms"] == 1000 + assert summary["ppo_sps_cv_pct"] == pytest.approx(100 / 3) + + +@pytest.mark.parametrize("second_vram", [1200, None]) +def test_report_separates_process_memory_from_allocator(tmp_path, second_vram) -> None: + result = _result() + result["resources"] = { + "samples": [ + { + "cpu_rss_mib": 50, + "cpu_rss_lifetime_peak_mib": 80, + "gpu_process_mib": 1000, + }, + { + "cpu_rss_mib": 70, + "cpu_rss_lifetime_peak_mib": 65, + "gpu_process_mib": second_vram, + }, + ] + } + summary = _summarize(result) + assert summary["cpu_rss_peak_mib"] == 80 + assert summary["gpu_process_sampled_peak_mib"] == second_vram + _write_report(tmp_path, result) + report = (tmp_path / "report.md").read_text() + assert "Process RAM (RSS peak): 80.000 MiB" in report + assert ( + f"process VRAM (sampled peak): {'n/a' if second_vram is None else '1200.000'} MiB" + in report + ) + assert "| create_runtime | 1000000.000 |" in report + + +@pytest.mark.parametrize("backend", ["default", "newton"]) +def test_config_keeps_official_policy_and_scales_minibatches(backend: str) -> None: + config = _load_config(TASK_DIR, backend, 32, 24) + assert config["trainer"]["num_envs"] == 32 + assert config["trainer"]["buffer_size"] == 24 + assert config["algorithm"]["cfg"]["batch_size"] == 192 + assert config["policy"]["obs_groups"] == {"actor": ["policy"], "critic": ["critic"]} + assert config["policy"]["actor_obs_normalization"] + assert config["policy"]["critic_obs_normalization"] + + +@pytest.mark.parametrize("envs,steps", [(1, 1), (3, 5)]) +def test_config_rejects_incomplete_minibatches(envs: int, steps: int) -> None: + with pytest.raises(ValueError, match="divisible"): + _load_config(TASK_DIR, "default", envs, steps) + + +@pytest.mark.parametrize("status", ["passed", "failed"]) +def test_report_does_not_rank_training_as_task_success(tmp_path, status: str) -> None: + result = _result(status) + _write_report(tmp_path, result) + report = (tmp_path / "report.md").read_text() + assert sum(line.startswith("| --- |") for line in report.splitlines()) == 3 + assert f"Dexsim Default CUDA | {status} | n/a |" in report + assert "EmbodiChain PPO | Dexsim Default CUDA | n/a |" in report + if status == "failed": + assert _summarize(result) == {} + assert "| rollout | n/a | n/a | n/a | n/a |" in report + + +def test_memory_sample_matches_device_and_process() -> None: + output = "11, 100, GPU-a\n11, 900, GPU-b\n12, 700, GPU-a\n" + assert _process_gpu_mib(output, 11, "a") == 100 + assert _process_gpu_mib(output, 13, "a") is None + assert _process_gpu_mib(output, None, "a") is None + assert _process_gpu_mib("11, N/A, GPU-a", 11, "a") is None + + +def test_resource_output_omits_process_and_device_identifiers() -> None: + result = ProcessResources("local-device-id").result() + assert "gpu_uuid" not in result + assert "pid_candidates" not in result + assert "local-device-id" not in str(result) + + +@pytest.mark.no_sim +@pytest.mark.parametrize( + "proc_namespace,expected", + [ + ("native", 1000), + ("host", 1000), + ("private", None), + ("unreadable", None), + ("malformed", None), + ], +) +def test_resource_sample_resolves_one_provider_pid( + monkeypatch: pytest.MonkeyPatch, proc_namespace: str, expected: float | None +) -> None: + """Exclude a colliding container PID and retain unknown namespace samples.""" + import psutil + + original_stat = Path.stat + original_read_text = Path.read_text + + def stat(path: Path, *args, **kwargs): + if str(path) == "/proc/self/ns/pid": + return SimpleNamespace( + st_ino=0xEFFFFFFC if proc_namespace == "native" else 0xF0000002 + ) + if str(path) == "/proc/1/ns/pid": + if proc_namespace == "unreadable": + raise PermissionError("namespace unavailable") + return SimpleNamespace( + st_ino=0xEFFFFFFC if proc_namespace != "private" else 0xF0000001 + ) + return original_stat(path, *args, **kwargs) + + def read_text(path: Path, *args, **kwargs): + if str(path) == "/proc/self/status": + return "NSpid:\t4300\t73\n" + return original_read_text(path, *args, **kwargs) + + monkeypatch.setattr(Path, "stat", stat) + monkeypatch.setattr(Path, "read_text", read_text) + monkeypatch.setattr( + Path, + "readlink", + lambda path: Path("unknown" if proc_namespace == "malformed" else "4300"), + ) + monkeypatch.setattr( + os, "getpid", lambda: 4300 if proc_namespace == "native" else 73 + ) + monkeypatch.setattr( + psutil, + "Process", + lambda: SimpleNamespace( + memory_info=lambda: SimpleNamespace(rss=0), + memory_full_info=lambda: SimpleNamespace(pss=0), + ), + ) + monkeypatch.setattr(torch.cuda, "synchronize", lambda: None) + monkeypatch.setattr(torch.cuda, "memory_allocated", lambda: 0) + monkeypatch.setattr(torch.cuda, "memory_reserved", lambda: 0) + monkeypatch.setattr( + resource_usage.subprocess, + "run", + lambda *args, **kwargs: SimpleNamespace( + stdout="4300, 1000, GPU-a\n73, 7000, GPU-a\n4300, 9000, GPU-b\n" + ), + ) + sample = ProcessResources("a").sample("test") + assert sample["gpu_process_mib"] == expected + assert (sample["gpu_error"] is None) == (expected is not None) + + +@pytest.mark.no_sim +def test_creation_failure_retains_resource_samples( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Write failure reports with samples taken before runtime construction.""" + from embodichain.lab.gym.utils import registration + from embodichain.learning.rl import runtime as policy_runtime + from scripts.benchmark.rl import runtime as benchmark_runtime + + env_config = tmp_path / "input.yaml" + env_config.write_text("physics_config:\n solver_cfg: {}\n") + output = tmp_path / "result" + monkeypatch.setattr( + sys, + "argv", + ["locomotion_ppo", "--backend", "newton", "--output", str(output)], + ) + monkeypatch.setattr( + benchmark, + "_load_config", + lambda *args: { + "trainer": {"gym_config": str(env_config), "renderer": "hybrid"}, + "policy": {}, + "algorithm": {}, + }, + ) + monkeypatch.setattr(torch.cuda, "is_available", lambda: True) + monkeypatch.setattr(torch.cuda, "set_device", lambda device: None) + monkeypatch.setattr(torch.cuda, "get_device_name", lambda device: "test GPU") + monkeypatch.setattr( + torch.cuda, + "get_device_properties", + lambda device: SimpleNamespace(uuid="test-device"), + ) + monkeypatch.setattr(torch, "set_num_threads", lambda count: None) + monkeypatch.setattr(benchmark_runtime, "set_random_seed", lambda *args: None) + monkeypatch.setattr(registration, "discover_task_packages", lambda: None) + monkeypatch.setattr(registration, "execute_init_hooks", lambda: None) + monkeypatch.setattr(benchmark, "_measure", lambda operation: (operation(), {})) + + def sample(resources: ProcessResources, label: str) -> dict: + item = {"label": label, "cpu_rss_mib": 50.0} + resources.samples.append(item) + return item + + monkeypatch.setattr(ProcessResources, "sample", sample) + failure = RuntimeError("runtime construction failed") + + def fail_creation(*args, **kwargs): + raise failure + + monkeypatch.setattr(policy_runtime, "build_gym_policy_runtime", fail_creation) + with pytest.raises(RuntimeError) as caught: + benchmark.main() + assert caught.value is failure + result = json.loads((output / "result.json").read_text()) + assert result["status"] == "failed" + assert result["resources"]["samples"] == [ + {"label": "before_create", "cpu_rss_mib": 50.0} + ] + assert ( + "| Dexsim Newton/MJWarp CUDA | failed |" in (output / "report.md").read_text() + ) diff --git a/tests/gym/envs/tasks/test_locomotion_cuda.py b/tests/gym/envs/tasks/test_locomotion_cuda.py new file mode 100644 index 000000000..03f7ff426 --- /dev/null +++ b/tests/gym/envs/tasks/test_locomotion_cuda.py @@ -0,0 +1,310 @@ +# ---------------------------------------------------------------------------- +# Copyright (c) 2021-2026 DexForce Technology Co., Ltd. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# ---------------------------------------------------------------------------- + +"""G1/Go2 CUDA inference, partial reset, and PPO checkpoint integration tests.""" + +from __future__ import annotations + +from collections.abc import Mapping +from pathlib import Path + +import json +import os +import subprocess +import sys + +import numpy as np +import pytest +import torch + +from embodichain.lab.gym.utils.registration import ( + discover_task_packages, + execute_init_hooks, +) +from embodichain.learning.rl.evaluation import ( + convert_policy_action_for_env, + infer_policy_action, +) +from embodichain.learning.rl.runtime import build_gym_policy_runtime +from embodichain.learning.rl.utils import ( + dict_to_tensordict, + flatten_observation_groups, +) +from scripts.benchmark.rl.locomotion_ppo import ( + _load_config, + _build_trainer, + _check_backend, +) + +ROOT = Path(__file__).resolve().parents[4] +TASK_DIR = ROOT / "embodichain_tasks/configs/tasks/locomotion/velocity" + + +def _subprocess_env() -> dict[str, str]: + """Make repository benchmark modules available to fresh workers.""" + env = os.environ.copy() + paths = [str(ROOT)] + if env.get("PYTHONPATH"): + paths.append(env["PYTHONPATH"]) + env["PYTHONPATH"] = os.pathsep.join(paths) + return env + + +def _finite(value: object) -> bool: + """Check tensor and array leaves of a nested observation.""" + if isinstance(value, torch.Tensor): + return bool(torch.isfinite(value).all()) + if isinstance(value, np.ndarray) and np.issubdtype(value.dtype, np.number): + return bool(np.isfinite(value).all()) + if isinstance(value, Mapping): + return all(_finite(child) for child in value.values()) + if isinstance(value, (tuple, list)): + return all(_finite(child) for child in value) + return True + + +@pytest.mark.subprocess_sim +@pytest.mark.gpu +@pytest.mark.requires_sim +@pytest.mark.requires_tasks +@pytest.mark.parametrize("robot", ["g1", "go2"]) +@pytest.mark.parametrize("backend", ["default", "newton"]) +def test_locomotion_cuda_reset_step_and_partial_reset(robot: str, backend: str) -> None: + """Exercise each robot/backend combination in its own process.""" + result = subprocess.run( + [sys.executable, str(Path(__file__).resolve()), robot, backend], + capture_output=True, + text=True, + timeout=120, + cwd=ROOT, + env=_subprocess_env(), + ) + assert result.returncode == 0, result.stdout + result.stderr + + +def _run_locomotion_smoke(robot: str, backend: str) -> None: + """Check inference and selected-row reset with the official task config.""" + discover_task_packages() + execute_init_hooks() + + config_dir = TASK_DIR / f"{robot}_flat" + config = _load_config(config_dir, backend, 8, 24) + runtime = build_gym_policy_runtime( + config, + device=torch.device("cuda:0"), + num_envs=8, + headless=True, + renderer=config["trainer"]["renderer"], + gpu_id=0, + seed=42, + config_dir=config_dir, + ) + try: + env = runtime.env + raw = env.unwrapped + newton = _check_backend(env, backend) + observation, _ = env.reset(seed=42) + assert _finite(observation) + runtime.policy.eval() + with torch.inference_mode(): + for _ in range(32): + action = infer_policy_action( + runtime.policy, + observation, + device=torch.device("cuda:0"), + num_envs=8, + ) + assert action.shape[0] == 8 + assert _finite(action) + observation, reward, terminated, truncated, _ = env.step( + convert_policy_action_for_env(env, action) + ) + assert _finite((observation, reward)) + assert reward.shape[0] == 8 + assert terminated.shape[0] == truncated.shape[0] == 8 + before = raw._elapsed_steps.clone() + state_fields = ("qpos", "qvel", "root_pose", "root_vel") + before_state = { + name: getattr(raw.robot.body_data, name).clone() for name in state_fields + } + if newton is not None: + assert newton.cuda_graph_status == "captured" + observation, _ = env.reset( + options={"reset_ids": torch.tensor([0, 3], device=raw.device)} + ) + assert _finite(observation) + assert bool((raw._elapsed_steps[[0, 3]] == 0).all()) + remaining = [1, 2, 4, 5, 6, 7] + assert torch.equal(raw._elapsed_steps[remaining], before[remaining]) + for name in state_fields: + torch.testing.assert_close( + getattr(raw.robot.body_data, name)[remaining], + before_state[name][remaining], + rtol=0, + atol=0, + ) + finally: + runtime.close() + + +@pytest.mark.subprocess_sim +@pytest.mark.gpu +@pytest.mark.requires_sim +@pytest.mark.requires_tasks +@pytest.mark.parametrize("backend", ["default", "newton"]) +def test_g1_ppo_checkpoint_resume(tmp_path: Path, backend: str) -> None: + """Train through the benchmark CLI, then resume in a separate process.""" + output = tmp_path / backend + commands = [ + [ + sys.executable, + "-m", + "scripts.benchmark.rl.locomotion_ppo", + "--backend", + backend, + "--output", + str(output), + "--envs", + "32", + "--steps", + "24", + "--warmup", + "1", + "--updates", + "3", + ], + [sys.executable, str(Path(__file__).resolve()), "resume", backend, str(output)], + ] + for command in commands: + result = subprocess.run( + command, + capture_output=True, + text=True, + timeout=180, + cwd=ROOT, + env=_subprocess_env(), + ) + assert result.returncode == 0, result.stdout + result.stderr + result = json.loads((output / "result.json").read_text()) + assert result["status"] == "passed" + assert result["summary"]["ppo_sps"] > 0 + assert result["parameter_max_abs_delta"] > 0 + + +def _resume_checkpoint(backend: str, output: Path) -> None: + discover_task_packages() + execute_init_hooks() + + device = torch.device("cuda:0") + config_dir = TASK_DIR / "g1_flat" + config = _load_config(config_dir, backend, 8, 24) + config["trainer"]["gym_config"] = str((output / "env.yaml").resolve()) + runtime = build_gym_policy_runtime( + config, + device=device, + num_envs=8, + headless=True, + renderer=config["trainer"]["renderer"], + gpu_id=0, + seed=42, + config_dir=config_dir, + ) + try: + _check_backend(runtime.env, backend) + saved = torch.load(output / "model.pt", map_location="cpu", weights_only=False) + assert saved["num_updates"] == 4 + assert saved["global_step"] == 32 * 24 * 4 + trainer = _build_trainer(runtime, config, output) + runtime.policy.load_state_dict(saved["policy"]) + trainer.algorithm.optimizer.load_state_dict(saved["optimizer"]) + for actual, expected in ( + (runtime.policy.state_dict(), saved["policy"]), + (trainer.algorithm.optimizer.state_dict(), saved["optimizer"]), + ): + torch.testing.assert_close( + actual, expected, rtol=0, atol=0, check_device=False + ) + optimizer_steps = [ + float(state["step"]) for state in saved["optimizer"]["state"].values() + ] + assert optimizer_steps + runtime.policy.eval() + observation, _ = runtime.env.reset(seed=42) + with torch.no_grad(): + for _ in range(32): + action = infer_policy_action( + runtime.policy, observation, device=device, num_envs=8 + ) + assert _finite(action) + observation, reward, *_ = runtime.env.step( + convert_policy_action_for_env(runtime.env, action) + ) + assert _finite((observation, reward)) + runtime.policy.train() + current_observation = dict_to_tensordict(observation, device) + trainer.collector.obs_td = current_observation + expected_actor_obs = flatten_observation_groups( + current_observation, runtime.policy.actor_obs_groups + ).clone() + expected_critic_obs = flatten_observation_groups( + current_observation, runtime.policy.critic_obs_groups + ).clone() + trainer.global_step = saved["global_step"] + trainer.num_updates = saved["num_updates"] + trainer.best_eval_value = saved["best_eval_value"] + summary = trainer.train(trainer.global_step + 8 * 24) + torch.testing.assert_close( + trainer.buffer.buffer["obs"][:, 0], expected_actor_obs, rtol=0, atol=0 + ) + torch.testing.assert_close( + trainer.buffer.buffer["critic_obs"][:, 0], + expected_critic_obs, + rtol=0, + atol=0, + ) + losses = [ + value + for key, value in summary["last_train_metrics"].items() + if key.startswith("train/") + ] + assert losses and all(np.isfinite(value) for value in losses) + trainer.save_checkpoint(str(output / "resumed.pt")) + resumed = torch.load( + output / "resumed.pt", map_location="cpu", weights_only=False + ) + assert resumed["num_updates"] == 5 + assert resumed["global_step"] == saved["global_step"] + 8 * 24 + assert any( + not torch.equal(saved["policy"][key], value) + for key, value in resumed["policy"].items() + ) + assert all(_finite(value) for value in resumed["policy"].values()) + after_steps = [ + float(state["step"]) for state in resumed["optimizer"]["state"].values() + ] + assert all( + after > before + for before, after in zip(optimizer_steps, after_steps, strict=True) + ) + finally: + runtime.close() + + +if __name__ == "__main__": + if sys.argv[1] == "resume": + _resume_checkpoint(sys.argv[2], Path(sys.argv[3])) + else: + _run_locomotion_smoke(sys.argv[1], sys.argv[2])