diff --git a/docs/sphinx/source/adr/ADR-0013-cuda-mps-single-rank-execution-sharing.md b/docs/sphinx/source/adr/ADR-0013-cuda-mps-single-rank-execution-sharing.md new file mode 100644 index 000000000..7dc32a5c4 --- /dev/null +++ b/docs/sphinx/source/adr/ADR-0013-cuda-mps-single-rank-execution-sharing.md @@ -0,0 +1,83 @@ +--- +orphan: true +--- + +# ADR-0013 CUDA MPS Single-Rank Execution Sharing + +语言: 简体中文 + +- Status: Proposed +- Date: 2026-10-04 +- Owners: UniLab training runtime maintainers +- Supersedes: None +- Superseded by: None + +## Context + +off-policy 训练已经要求每个 rank 的 learner 与 collector 使用同一张 rank-local GPU, +但二者仍是独立 CUDA context,由驱动按多 context 粗粒度分时仲裁。Discussion #1800 +的参考主机证据显示,显式启用 NVIDIA CUDA MPS 可以让这两个 context 共享 GPU 并重叠 +kernel,从而改善端到端吞吐;其证据范围是单主机、单 rank、MJWarp。 + +现有配置没有表达该 host execution mode,用户只能在部署层隐式设置环境变量。这既不能 +fail closed,也没有可审计的 runtime evidence。 + +## Decision + +- 新增共享 off-policy owner 设置 + `training.cuda_process_sharing: null | mps`,默认 `null` 保持现状。 +- 初始有效范围限定为 Linux、NVIDIA CUDA、`training.sim_backend=mjwarp`、单主机、 + `world_size=1` 的 SAC/FastSAC/FlashSAC/WarpSAC 共享路径。 +- `mps` 是显式请求:所有前置条件在 env probe、learner construction 与 collector + spawn 之前验证,失败时给出第一个未满足条件和启动既有 daemon 的命令,不 fallback。 +- rank-local 物理一致性按 GPU UUID 判断,而不是 CUDA ordinal。 +- UniLab 只验证既有 MPS control daemon 与 control socket/FIFO,不 start/stop/repair + daemon,不引入 `auto`、SM percentage 或 OS nice/affinity 配置。 +- `run_config.json` 记录配置值;runtime manifest v1 增加 producer 诊断字段 + `cuda_process_sharing`,记录 configured/effective、devices、UUID、control pipe、 + server PID 与 validated 状态。该字段先保持 v1 诊断扩展,不提升为稳定公共字段。 +- 多 GPU DP 继续由 #2063 单独决策;本 ADR 的 per-rank 参数与 evidence shape 保留 + 加法扩展空间。 + +## Stable Contracts + +- Owner 配置:三个 off-policy 主配置中的 `training.cuda_process_sharing`。 +- Fail-closed probe:`src/unilab/training/cuda_process_sharing.py`。 +- Builder wiring:`src/unilab/scripts/train_offpolicy.py` 在构造 env/runner 前 probe, + 并把 evidence 合并进 runner 的 runtime manifest。 +- 测试与文档:probe fake、owner default、builder fail-closed、manifest evidence、 + 英文/中文生产指南。 + +## Alternatives Considered + +- 继续要求用户只在 shell 设置 MPS 环境变量:无法验证拓扑,也没有结构化证据。 +- 由 UniLab 启动/停止 daemon:把共享 host 服务生命周期混入训练进程,且无法安全服务 + 并发/多用户运行。 +- CPU priority/affinity:Discussion #1800 的测量显示其不解决跨进程 CUDA context + arbitration,且引入未证实的公共配置复杂度。 + +## Consequences + +- 显式 `mps` 请求失败时不会有默认多 context fallback;这避免静默性能和拓扑漂移。 +- MPS daemon、pipe/log 目录与多用户隔离仍是部署者责任。 +- 非 MJWarp backend 和 DP 请求会失败,不形成支持声明。 +- `env_steps_per_sync` 仍是训练语义 owner 设置,不由 execution-sharing mode 修改。 + +## Evidence In Repo + +- `src/unilab/training/cuda_process_sharing.py` +- `src/unilab/scripts/train_offpolicy.py` +- `src/unilab/conf/sac/config.yaml` +- `src/unilab/conf/flashsac/config.yaml` +- `src/unilab/conf/warpsac/config.yaml` +- `tests/training/test_cuda_process_sharing.py` +- `tests/algos/test_offpolicy_double_buffer_runner.py` +- `docs/sphinx/source/en/2-user_guide/1-training/7-tensor_runtime_production.md` +- `docs/sphinx/source/zh_CN/2-user_guide/1-training/7-tensor_runtime_production.md` + +## Related Documents + +- {doc}`ADR Index ` +- {doc}`RL Infrastructure 开发标准 ` +- Discussion #1800 +- Issue #2063 diff --git a/docs/sphinx/source/adr/README.md b/docs/sphinx/source/adr/README.md index d75d38a52..0180d74b4 100644 --- a/docs/sphinx/source/adr/README.md +++ b/docs/sphinx/source/adr/README.md @@ -26,6 +26,7 @@ orphan: true | [ADR-0010 Fixed Model Variant Ownership Boundary](ADR-0010-fixed-model-variant-ownership-boundary.md) | Fixed variants / cross-repository boundary | Proposed | | [ADR-0011 Torch-Only Manager-Based Runtime](ADR-0011-torch-only-manager-based-runtime.md) | Manager runtime / tensor lifecycle | Superseded | | [ADR-0012 Sole Tensor Manager And Scoped Backends](ADR-0012-sole-tensor-manager-and-scoped-backends.md) | Manager runtime / backend scope | Proposed | +| [ADR-0013 CUDA MPS Single-Rank Execution Sharing](ADR-0013-cuda-mps-single-rank-execution-sharing.md) | Training runtime / execution sharing | Proposed | ## ADR Governance diff --git a/docs/sphinx/source/en/2-user_guide/1-training/4-tensor_runtime.md b/docs/sphinx/source/en/2-user_guide/1-training/4-tensor_runtime.md index 46f52c0f0..71cecc0fa 100644 --- a/docs/sphinx/source/en/2-user_guide/1-training/4-tensor_runtime.md +++ b/docs/sphinx/source/en/2-user_guide/1-training/4-tensor_runtime.md @@ -17,12 +17,15 @@ autotuner and does not add multi-GPU scaling. | `training.collector_metrics_interval` | `1` | Positive integer, at most `10000` | | `training.replay_ingress_depth` | `2` | Positive integer, at most `16` | | `training.replay_ingress_slot_rows` | `null` | `1` through `algo.num_envs`; `null` means `algo.num_envs` | +| `training.cuda_process_sharing` | `null` | `null` or explicit `mps` for single-rank MJWarp | | Learner rows per synchronization | `algo.batch_size * algo.updates_per_step` | Constrained by the CUDA memory budget | -Values must be exact positive integers. Booleans, strings, floating-point -numbers, zero, and negative values fail closed. The G1 Motion Tracking / MJWarp -owner intentionally overrides only `collector_metrics_interval` to `100`; all -other tensor-runtime defaults remain as listed above. +Values other than `training.cuda_process_sharing` must be exact positive +integers. Booleans, strings, floating-point numbers, zero, and negative values +fail closed. `training.cuda_process_sharing` accepts only `null` or the explicit +string `mps`. The G1 Motion Tracking / MJWarp owner intentionally overrides only +`collector_metrics_interval` to `100`; all other tensor-runtime defaults remain +as listed above. ## Replay ingress tradeoffs @@ -68,6 +71,8 @@ The runtime manifest records the effective evidence needed to audit a run: - `inference_memory_budget` records the bounded CUDA inference-ring budget. - `tensor_memory_budget` records the combined CUDA inference, replay storage, replay ingress, learner batch, and workspace budget calculation. +- `cuda_process_sharing` records validated execution-sharing evidence when the + owner explicitly requests `mps`; it is absent for the default `null` mode. The manifest is written before spawn when a budget decision is made. If an unsafe combination is rejected, the error identifies the offending setting and diff --git a/docs/sphinx/source/en/2-user_guide/1-training/7-tensor_runtime_production.md b/docs/sphinx/source/en/2-user_guide/1-training/7-tensor_runtime_production.md index 412178644..007ddead0 100644 --- a/docs/sphinx/source/en/2-user_guide/1-training/7-tensor_runtime_production.md +++ b/docs/sphinx/source/en/2-user_guide/1-training/7-tensor_runtime_production.md @@ -42,6 +42,7 @@ the MJWarp owner inherits the MuJoCo FlashSAC owner and resolves to: | Replay-ingress depth | `training.replay_ingress_depth=2` | | Replay-ingress slot rows | `null`, resolving to `algo.num_envs` | | Collector metric interval | `training.collector_metrics_interval=100` | +| CUDA process sharing | `training.cuda_process_sharing=null` | | Scene | `src/unilab/assets/robots/g1/scene_flat.xml` | | Motion | `motions/g1/dance1_subject2_part.npz` | @@ -87,6 +88,56 @@ export CUDA_VISIBLE_DEVICES= All CUDA ordinals inside trainer processes are relative to that mask. Do not add more devices: M11 is single-GPU only. +## CUDA MPS Execution Sharing + +`training.cuda_process_sharing` is an explicit execution-sharing mode for the +already-required rank-local topology. It is not a backend switch and does not +replace `--sim mjwarp` or owner YAML selection. + +The default remains `null`: learner and collector keep independent CUDA contexts. +For the supported single-host, single-rank MJWarp off-policy topology, request +CUDA MPS with: + +```bash +training.cuda_process_sharing=mps +``` + +Start an existing control daemon under user-owned directories before training: + +```bash +export CUDA_MPS_PIPE_DIRECTORY=/absolute/path/mps/pipe +export CUDA_MPS_LOG_DIRECTORY=/absolute/path/mps/log +mkdir -p "$CUDA_MPS_PIPE_DIRECTORY" "$CUDA_MPS_LOG_DIRECTORY" +nvidia-cuda-mps-control -d +``` + +UniLab validates, but never starts or stops, that daemon. An `mps` request must +be Linux/NVIDIA CUDA, use MJWarp, resolve one physical GPU by UUID for learner +and collector, have `world_size=1`, and reach a live control socket/FIFO and +control daemon. Any unmet prerequisite fails before environment probing, +learner construction, or collector spawn; there is no silent multi-context +fallback. The error names the first prerequisite and the daemon start command. + +`run_config.json` records the configured owner value. A valid run's +`run_summary.json` embeds `runtime_manifest.cuda_process_sharing` with the +configured/effective mode, learner and collector devices, physical UUID +evidence, control pipe, server PID, and validation state. The section is a +runtime-manifest v1 producer diagnostic, not a stable scalar contract. + +MPS changes only GPU execution sharing. It does not change `env_steps_per_sync`, +learner/collector placement, or training semantics. Multi-GPU DP remains +unsupported until the separate DP gate and daemon topology decision are +completed. CUDA MPS is host- and deployment-dependent: in shared containers, +multi-user hosts, or restricted runners, control may be unavailable, and an +explicit request fails closed rather than silently degrading. + +Stop the daemon after use: + +```bash +export CUDA_MPS_PIPE_DIRECTORY=/absolute/path/mps/pipe +echo quit | nvidia-cuda-mps-control +``` + Check the runtime that will execute the benchmark: ```bash diff --git a/docs/sphinx/source/zh_CN/2-user_guide/1-training/4-tensor_runtime.md b/docs/sphinx/source/zh_CN/2-user_guide/1-training/4-tensor_runtime.md index e327d6df0..180db5902 100644 --- a/docs/sphinx/source/zh_CN/2-user_guide/1-training/4-tensor_runtime.md +++ b/docs/sphinx/source/zh_CN/2-user_guide/1-training/4-tensor_runtime.md @@ -17,10 +17,12 @@ SAC、FlashSAC 与 WarpSAC 会在探测环境或构造 learner 之前解析所 | `training.collector_metrics_interval` | `1` | 正整数,最大 `10000` | | `training.replay_ingress_depth` | `2` | 正整数,最大 `16` | | `training.replay_ingress_slot_rows` | `null` | `1` 到 `algo.num_envs`;`null` 表示 `algo.num_envs` | +| `training.cuda_process_sharing` | `null` | `null` 或单 rank MJWarp 显式 `mps` | | Learner 每次同步的行数 | `algo.batch_size * algo.updates_per_step` | 受 CUDA 显存预算约束 | -所有值必须是精确的正整数。布尔值、字符串、浮点数、零和负数都会 fail -closed。G1 Motion Tracking / MJWarp owner 仅有意将 +除 `training.cuda_process_sharing` 外,所有值必须是精确的正整数。布尔值、字符串、 +浮点数、零和负数都会 fail closed。`training.cuda_process_sharing` 只接受 `null` +或显式字符串 `mps`。G1 Motion Tracking / MJWarp owner 仅有意将 `collector_metrics_interval` 覆盖为 `100`;其他 tensor-runtime 默认值仍保持 上表所示。 @@ -64,6 +66,8 @@ Runtime manifest 记录审计 run 所需的有效证据: - `inference_memory_budget` 记录有边界的 CUDA inference-ring 预算。 - `tensor_memory_budget` 记录 CUDA inference、replay storage、replay ingress、 learner batch 与 workspace 的组合预算计算。 +- owner 显式请求 `mps` 时,`cuda_process_sharing` 记录已验证的 + execution-sharing 证据;默认 `null` 模式下该字段不存在。 预算决策在 spawn 前写入 manifest。如果不安全的组合被拒绝,错误会指出具体 设置,且相关进程不会被启动。 diff --git a/docs/sphinx/source/zh_CN/2-user_guide/1-training/7-tensor_runtime_production.md b/docs/sphinx/source/zh_CN/2-user_guide/1-training/7-tensor_runtime_production.md index c28223be2..ef3d8f8a7 100644 --- a/docs/sphinx/source/zh_CN/2-user_guide/1-training/7-tensor_runtime_production.md +++ b/docs/sphinx/source/zh_CN/2-user_guide/1-training/7-tensor_runtime_production.md @@ -41,6 +41,7 @@ MJWarp owner 继承 FlashSAC MuJoCo owner,并解析为: | Replay-ingress depth | `training.replay_ingress_depth=2` | | Replay-ingress slot rows | `null`,解析为 `algo.num_envs` | | Collector metric interval | `training.collector_metrics_interval=100` | +| CUDA 进程共享 | `training.cuda_process_sharing=null` | | Scene | `src/unilab/assets/robots/g1/scene_flat.xml` | | Motion | `motions/g1/dance1_subject2_part.npz` | @@ -82,6 +83,51 @@ export CUDA_VISIBLE_DEVICES= trainer 进程内所有 CUDA ordinal 都相对于该 mask。不要添加更多设备:M11 仅 覆盖单 GPU。 +## CUDA MPS 执行共享 + +`training.cuda_process_sharing` 是针对既有 rank-local 拓扑的显式 +execution-sharing 模式。它不是 backend 开关,也不替代 `--sim mjwarp` 或 owner +YAML 选择。 + +默认 `null` 表示 learner 与 collector 保持独立 CUDA context。对受支持的单主机、 +单 rank MJWarp off-policy 拓扑,使用以下设置请求 CUDA MPS: + +```bash +training.cuda_process_sharing=mps +``` + +训练前在用户拥有的目录中启动既有 control daemon: + +```bash +export CUDA_MPS_PIPE_DIRECTORY=/absolute/path/mps/pipe +export CUDA_MPS_LOG_DIRECTORY=/absolute/path/mps/log +mkdir -p "$CUDA_MPS_PIPE_DIRECTORY" "$CUDA_MPS_LOG_DIRECTORY" +nvidia-cuda-mps-control -d +``` + +UniLab 只验证、绝不 start/stop 该 daemon。`mps` 请求必须是 Linux/NVIDIA CUDA、 +使用 MJWarp、learner 与 collector 按 UUID 解析到同一张物理 GPU、`world_size=1`, +且能访问 live control socket/FIFO 与 control daemon。任一条件不满足都会在 +environment probe、learner construction 或 collector spawn 之前失败;没有静默 +multi-context fallback。错误会指出第一个未满足条件以及 daemon 启动命令。 + +`run_config.json` 记录配置值。有效 run 的 `run_summary.json` 内嵌 +`runtime_manifest.cuda_process_sharing`,包含 configured/effective 模式、learner +与 collector device、物理 UUID 证据、control pipe、server PID 和 validation 状态。 +该 section 是 runtime-manifest v1 的 producer diagnostic,不是稳定标量契约。 + +MPS 只改变 GPU execution sharing,不改变 `env_steps_per_sync`、 +learner/collector placement 或训练语义。多 GPU DP 在独立的 DP gate 与 daemon +拓扑决策完成前保持不支持。CUDA MPS 依赖 host 与部署方式:在共享容器、多用户 +主机或受限 runner 中 control 可能不可用,显式请求会 fail closed,不会静默降级。 + +使用后停止 daemon: + +```bash +export CUDA_MPS_PIPE_DIRECTORY=/absolute/path/mps/pipe +echo quit | nvidia-cuda-mps-control +``` + 检查实际执行 benchmark 的 runtime: ```bash diff --git a/src/unilab/conf/flashsac/config.yaml b/src/unilab/conf/flashsac/config.yaml index 796f85142..ed0067e7e 100644 --- a/src/unilab/conf/flashsac/config.yaml +++ b/src/unilab/conf/flashsac/config.yaml @@ -65,6 +65,9 @@ training: wandb_notes: null wandb_mode: null sim_backend: mujoco + # null preserves independent CUDA contexts; 'mps' explicitly requests CUDA + # process sharing and fails closed if the host daemon is not already valid. + cuda_process_sharing: null nan_guard: enabled: true buffer_size: 100 diff --git a/src/unilab/conf/sac/config.yaml b/src/unilab/conf/sac/config.yaml index 1cb3aa1e7..f83ac1664 100644 --- a/src/unilab/conf/sac/config.yaml +++ b/src/unilab/conf/sac/config.yaml @@ -54,6 +54,9 @@ training: wandb_notes: null wandb_mode: null sim_backend: mujoco + # null preserves independent CUDA contexts; 'mps' explicitly requests CUDA + # process sharing and fails closed if the host daemon is not already valid. + cuda_process_sharing: null nan_guard: enabled: true buffer_size: 100 diff --git a/src/unilab/conf/warpsac/config.yaml b/src/unilab/conf/warpsac/config.yaml index 495f8ee86..4d3efb61b 100644 --- a/src/unilab/conf/warpsac/config.yaml +++ b/src/unilab/conf/warpsac/config.yaml @@ -69,6 +69,9 @@ training: wandb_notes: null wandb_mode: null sim_backend: mujoco + # null preserves independent CUDA contexts; 'mps' explicitly requests CUDA + # process sharing and fails closed if the host daemon is not already valid. + cuda_process_sharing: null nan_guard: enabled: true buffer_size: 100 diff --git a/src/unilab/scripts/train_offpolicy.py b/src/unilab/scripts/train_offpolicy.py index 042b6b18b..af06b5eb2 100644 --- a/src/unilab/scripts/train_offpolicy.py +++ b/src/unilab/scripts/train_offpolicy.py @@ -38,6 +38,7 @@ configure_backend_process_device, pin_genesis_device_before_cuda_init, resolve_backend_env_device_id, + resolve_backend_process_device, ) from unilab.training import ( assert_offpolicy_task_choice_matches_algo, @@ -49,6 +50,10 @@ resolve_nan_guard_cfg, should_run_playback, ) +from unilab.training.cuda_process_sharing import ( + CudaProcessSharingEvidence, + probe_cuda_process_sharing, +) from unilab.training.experiment import ExperimentTracker from unilab.training.onnx_export import export_policy_onnx, verify_policy_onnx from unilab.utils.checkpoint import ( @@ -151,12 +156,6 @@ def build_offpolicy_play_env_cfg_override(algo_name: str, cfg: DictConfig) -> di def build_runner(algo_name: str, cfg: DictConfig, log_dir: str | None = None): """Build algorithm runner from unified Hydra config.""" - env_factory = registry_env_factory(str(cfg.training.task_name), str(cfg.training.sim_backend)) - from uni_rl.offpolicy.thread_budget import ( - apply_torch_thread_runtime, - resolve_torch_thread_runtime, - ) - # Cold-path DP CPU partition: each rank's collector owns one contiguous # CPU block (single rank keeps the legacy unset behavior). The ids only # reach the collector env override — never the num_envs=1 probe envs, @@ -180,6 +179,24 @@ def build_runner(algo_name: str, cfg: DictConfig, log_dir: str | None = None): ) if bound_device is not None: rank_device = bound_device + collector_device = resolve_backend_process_device( + str(cfg.training.sim_backend), + rank_device, + ) + cuda_process_sharing = probe_cuda_process_sharing( + getattr(cfg.training, "cuda_process_sharing", None), + rank_device, + collector_device, + backend=str(cfg.training.sim_backend), + world_size=dp_world_size, + ) + + env_factory = registry_env_factory(str(cfg.training.task_name), str(cfg.training.sim_backend)) + from uni_rl.offpolicy.thread_budget import ( + apply_torch_thread_runtime, + resolve_torch_thread_runtime, + ) + # Every rank is single-GPU; rank-local visibility routes the backend payload. env_cfg_override = apply_backend_env_device_override( build_offpolicy_env_cfg_override(algo_name, cfg), @@ -283,9 +300,23 @@ def build_runner(algo_name: str, cfg: DictConfig, log_dir: str | None = None): else: raise ValueError(f"Unsupported algo: {algo_name}") + _attach_cuda_process_sharing_manifest(runner, cuda_process_sharing) return runner +def _attach_cuda_process_sharing_manifest( + runner: Any, + evidence: CudaProcessSharingEvidence, +) -> None: + """Merge validated mode evidence into the producer-owned manifest.""" + + runtime_manifest = getattr(runner, "runtime_manifest", None) + if isinstance(runtime_manifest, dict): + runtime_manifest["cuda_process_sharing"] = evidence.manifest() + return + runner.runtime_manifest = {"cuda_process_sharing": evidence.manifest()} + + def play_offpolicy( algo_name: str, cfg: DictConfig, @@ -407,6 +438,7 @@ def play_offpolicy( def main(cfg: DictConfig) -> None: enable_faulthandler() + import torch reject_removed_device_config(OmegaConf.select(cfg, "training.devices", default=None)) rank = current_dp_rank() @@ -460,8 +492,6 @@ def main(cfg: DictConfig) -> None: if rank == 0 and world_size > 1: supervisor = DpRankSupervisor(world_size=world_size, log_dir=log_dir) - import torch - tracker = None if not cfg.training.play_only and rank == 0: tracker = ExperimentTracker( diff --git a/src/unilab/training/cuda_process_sharing.py b/src/unilab/training/cuda_process_sharing.py new file mode 100644 index 000000000..14a3342d2 --- /dev/null +++ b/src/unilab/training/cuda_process_sharing.py @@ -0,0 +1,383 @@ +"""Fail-closed validation for CUDA process-sharing execution modes. + +The off-policy learner and collector are separate processes on one rank-local +GPU. CUDA MPS is an execution-sharing mode for that required topology: it lets +their independent CUDA contexts overlap kernels instead of being arbitrated by +coarse driver time-slicing. + +This module never starts or stops an MPS daemon. An explicit request is either +validated against the host state or rejected before UniLab constructs a learner, +collector, or environment. +""" + +from __future__ import annotations + +import os +import platform +import stat +import subprocess +from collections.abc import Callable, Sequence +from dataclasses import dataclass, replace +from pathlib import Path +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from collections.abc import MutableMapping + +SUPPORTED_CUDA_PROCESS_SHARING_BACKENDS = frozenset({"mjwarp"}) +_DEFAULT_PIPE_DIRECTORY = "/tmp/nvidia-mps" +_MPS_CONTROL_QUERY_TIMEOUT_SEC = 5.0 + + +@dataclass(frozen=True) +class CudaProcessSharingEvidence: + """Structured per-rank evidence for an explicitly requested mode.""" + + configured: str | None + effective: str | None + validated: bool + learner_device: str | None = None + collector_device: str | None = None + learner_gpu_uuid: str | None = None + collector_gpu_uuid: str | None = None + control_pipe: str | None = None + server_pid: int | None = None + + def manifest(self) -> dict[str, Any]: + """Return the JSON-shaped diagnostics stored in a runtime manifest.""" + return { + "configured": self.configured, + "effective": self.effective, + "validated": self.validated, + "learner_device": self.learner_device, + "collector_device": self.collector_device, + "learner_gpu_uuid": self.learner_gpu_uuid, + "collector_gpu_uuid": self.collector_gpu_uuid, + "control_pipe": self.control_pipe, + "server_pid": self.server_pid, + } + + +def _configured_mode(requested: Any) -> str | None: + if requested is None: + return None + if isinstance(requested, str) and requested.strip().lower() == "mps": + return "mps" + raise ValueError( + "Unsupported training.cuda_process_sharing=" + f"{requested!r}; expected null or 'mps'. An explicit 'mps' request " + "never falls back to multi-context CUDA execution." + ) + + +def _cuda_index(device: str, requested: str) -> int: + value = device.strip() + if value.lower() == "cuda": + return 0 + base, separator, index_text = value.partition(":") + if base.lower() != "cuda" or not separator or not index_text: + raise ValueError( + "training.cuda_process_sharing='mps' requires CUDA learner and " + f"collector devices; got learner_device={device!r}." + ) + try: + index = int(index_text) + except ValueError as exc: + raise ValueError( + f"training.cuda_process_sharing='mps' received invalid CUDA device {device!r}." + ) from exc + if index < 0: + raise ValueError( + f"training.cuda_process_sharing='mps' received invalid CUDA device {device!r}." + ) + return index + + +def _torch_device_uuid( + torch_module: Any, + index: int, + *, + device_kind: str, + requested: str, +) -> str: + if not torch_module.cuda.is_available(): + raise ValueError( + "training.cuda_process_sharing='mps' requires CUDA, but CUDA is " + f"unavailable for the {device_kind} device cuda:{index}." + ) + if torch_module.version.hip is not None: + raise ValueError( + "training.cuda_process_sharing='mps' supports NVIDIA CUDA only; " + "this process reports a ROCm/HIP Torch build." + ) + device_count = int(torch_module.cuda.device_count()) + if index >= device_count: + raise ValueError( + f"training.cuda_process_sharing='mps' {device_kind} device " + f"cuda:{index} is out of range; torch.cuda.device_count()={device_count}." + ) + properties = torch_module.cuda.get_device_properties(index) + uuid = str(getattr(properties, "uuid", "") or "").strip() + if not uuid: + raise ValueError( + "training.cuda_process_sharing='mps' could not resolve the " + f"{device_kind} CUDA device UUID; required mode={requested!r}." + ) + return _canonical_gpu_uuid(uuid) + + +def _parse_csv_row(row: str) -> list[str]: + return [field.strip() for field in row.split(",")] + + +def _canonical_gpu_uuid(value: str) -> str: + return value.strip().upper().removeprefix("GPU-").replace("-", "") + + +def _nvidia_uuid( + run_command: Callable[..., Any], + visible_entries: Sequence[str], + index: int, +) -> str | None: + """Map a visibility-relative index to its host UUID, if evidence exists.""" + try: + result = run_command( + ["nvidia-smi", "--query-gpu=index,uuid", "--format=csv,noheader"], + text=True, + capture_output=True, + timeout=_MPS_CONTROL_QUERY_TIMEOUT_SEC, + check=True, + ) + except Exception: + return None + rows = [_parse_csv_row(row) for row in str(result.stdout).splitlines() if row.strip()] + if not rows or any(len(row) != 2 for row in rows): + return None + by_host_index = {row[0]: _canonical_gpu_uuid(row[1]) for row in rows if row[0].isdigit()} + if index < len(visible_entries): + token = _canonical_gpu_uuid(visible_entries[index]) + if token.isdigit(): + return by_host_index.get(token) + if token in by_host_index.values(): + return token + return by_host_index.get(str(index)) + + +def _server_evidence( + run_command: Callable[..., Any], + control_pipe: Path, +) -> int | None: + """Query the existing control daemon; never launch or terminate one.""" + environment = os.environ.copy() + environment["CUDA_MPS_PIPE_DIRECTORY"] = str(control_pipe.parent) + try: + result = run_command( + ["nvidia-cuda-mps-control", "get-server-list"], + text=True, + capture_output=True, + timeout=_MPS_CONTROL_QUERY_TIMEOUT_SEC, + check=True, + env=environment, + ) + except Exception: + return None + queried_pid: int | None = None + for line in str(result.stdout).splitlines(): + line = line.strip() + if line.isdigit(): + queried_pid = int(line) + break + pid_file = control_pipe.parent / "nvidia-cuda-mps-control.pid" + try: + file_pid = int(pid_file.read_text(encoding="utf-8").strip()) + except (OSError, ValueError): + file_pid = None + return queried_pid if queried_pid is not None else file_pid + + +def _start_command(pipe_directory: Path) -> str: + log_directory = os.environ.get("CUDA_MPS_LOG_DIRECTORY") + log_suffix = ( + log_directory if log_directory is not None else f"{pipe_directory.parent.as_posix()}/log" + ) + return ( + "mkdir -p " + f"{pipe_directory.as_posix()} {log_suffix} && " + f"export CUDA_MPS_PIPE_DIRECTORY={pipe_directory.as_posix()} && " + f"export CUDA_MPS_LOG_DIRECTORY={log_suffix} && " + "nvidia-cuda-mps-control -d" + ) + + +def _control_pipe_kind(path: Path) -> str | None: + try: + mode = path.stat().st_mode + except OSError: + return None + if stat.S_ISSOCK(mode): + return "unix_socket" + if stat.S_ISFIFO(mode): + return "fifo" + return None + + +def _visible_entries(raw: str | None) -> tuple[str, ...]: + if raw is None or not raw.strip() or raw.strip() == "all": + return () + return tuple(entry.strip() for entry in raw.split(",") if entry.strip()) + + +def probe_cuda_process_sharing( + requested: Any, + learner_device: str, + collector_device: str | None, + *, + backend: str, + world_size: int = 1, + torch_module: Any | None = None, + run_command: Callable[..., Any] = subprocess.run, +) -> CudaProcessSharingEvidence: + """Validate one rank's requested CUDA process-sharing mode. + + Args: + requested: Owner value ``null`` or ``mps``. + learner_device: Rank-local learner CUDA device. + collector_device: Backend-bound collector CUDA device, when applicable. + backend: Configured owner backend identity. + world_size: Current off-policy DP world size. MPS is single-rank only. + torch_module: Injectable Torch module for deterministic tests. + run_command: Injectable subprocess runner for deterministic tests. + + Returns: + Evidence suitable for a per-rank runtime-manifest diagnostic section. + + Raises: + ValueError: The first unmet prerequisite for an explicit request. + """ + + configured = _configured_mode(requested) + if configured is None: + return CudaProcessSharingEvidence( + configured=None, + effective=None, + validated=False, + learner_device=str(learner_device), + collector_device=str(collector_device) if collector_device is not None else None, + ) + + normalized_backend = str(backend).strip().lower() + if normalized_backend not in SUPPORTED_CUDA_PROCESS_SHARING_BACKENDS: + raise ValueError( + "training.cuda_process_sharing='mps' supports only the mjwarp " + f"backend in this release; got training.sim_backend={backend!r}." + ) + if platform.system() != "Linux": + raise ValueError( + "training.cuda_process_sharing='mps' requires Linux; " + f"got platform.system()={platform.system()!r}." + ) + if isinstance(world_size, bool) or not isinstance(world_size, int) or world_size != 1: + raise ValueError( + "training.cuda_process_sharing='mps' supports single-host, " + f"single-rank training only; got world_size={world_size!r}." + ) + if collector_device is None: + raise ValueError( + "training.cuda_process_sharing='mps' requires a CUDA collector " + "process device; the mjwarp learner/collector rank-local binding " + "did not produce one." + ) + + learner_index = _cuda_index(str(learner_device), configured) + collector_index = _cuda_index(str(collector_device), configured) + pipe_directory = Path(os.environ.get("CUDA_MPS_PIPE_DIRECTORY") or _DEFAULT_PIPE_DIRECTORY) + control_pipe = pipe_directory / "control" + visible_entries = _visible_entries(os.environ.get("CUDA_VISIBLE_DEVICES")) + evidence = CudaProcessSharingEvidence( + configured=configured, + effective=configured, + validated=True, + learner_device=f"cuda:{learner_index}", + collector_device=f"cuda:{collector_index}", + control_pipe=str(control_pipe), + ) + + if torch_module is None: + import torch + + torch_module = torch + learner_uuid = _torch_device_uuid( + torch_module, + learner_index, + device_kind="learner", + requested=configured, + ) + collector_uuid = _nvidia_uuid(run_command, visible_entries, collector_index) + evidence = replace( + evidence, + learner_gpu_uuid=learner_uuid, + collector_gpu_uuid=collector_uuid, + ) + + if learner_uuid is not None and collector_uuid is not None and learner_uuid != collector_uuid: + raise ValueError( + "training.cuda_process_sharing='mps' requires the learner and " + "collector to share one physical GPU. Resolved UUIDs: " + f"learner={learner_uuid!r}, collector={collector_uuid!r}." + ) + if collector_uuid is None: + raise ValueError( + "training.cuda_process_sharing='mps' could not resolve the " + f"collector CUDA UUID for {collector_device!r}; run nvidia-smi " + "--query-gpu=index,uuid --format=csv and verify the rank-local mask." + ) + + if not control_pipe.exists(): + raise ValueError( + "training.cuda_process_sharing='mps' found no control pipe at " + f"{control_pipe}. Start an existing host daemon with: " + f"{_start_command(pipe_directory)}" + ) + pipe_kind = _control_pipe_kind(control_pipe) + if pipe_kind is None: + raise ValueError( + "training.cuda_process_sharing='mps' requires a live control socket " + f"or FIFO at {control_pipe}; that path is not a daemon control pipe." + ) + if not os.access(control_pipe, os.R_OK | os.W_OK): + raise ValueError( + "training.cuda_process_sharing='mps' cannot access control pipe " + f"{control_pipe}; the launching user must be able to read and write " + "the daemon's explicitly owned pipe directory." + ) + server_pid = _server_evidence(run_command, control_pipe) + if server_pid is None: + raise ValueError( + "training.cuda_process_sharing='mps' could not reach the control " + f"daemon through {control_pipe}. Verify the daemon and retry: " + f"{_start_command(pipe_directory)}" + ) + + return replace( + evidence, + learner_gpu_uuid=learner_uuid, + collector_gpu_uuid=collector_uuid, + server_pid=server_pid, + ) + + +def apply_cuda_process_sharing_manifest( + runtime_manifest: MutableMapping[str, Any], + evidence: CudaProcessSharingEvidence, +) -> None: + """Record process-sharing evidence in a producer-owned manifest mapping.""" + + runtime_manifest["cuda_process_sharing"] = evidence.manifest() + + +__all__ = [ + "CudaProcessSharingEvidence", + "SUPPORTED_CUDA_PROCESS_SHARING_BACKENDS", + "apply_cuda_process_sharing_manifest", + "probe_cuda_process_sharing", +] diff --git a/tests/algos/test_offpolicy_double_buffer_runner.py b/tests/algos/test_offpolicy_double_buffer_runner.py index b0e7ed108..f337b8f8c 100644 --- a/tests/algos/test_offpolicy_double_buffer_runner.py +++ b/tests/algos/test_offpolicy_double_buffer_runner.py @@ -20,6 +20,8 @@ TensorRuntimeSettings, ) +from unilab.training import cuda_process_sharing + _ROOT = Path(__file__).parent.parent.parent _CONF_DIR = _ROOT / "src" / "unilab" / "conf" @@ -88,6 +90,41 @@ def _fake_env_factory(num_envs, env_cfg_override): return _FakeEnv() +def _cuda_torch_module(monkeypatch: pytest.MonkeyPatch, uuid: str = "GPU-a"): + monkeypatch.setattr(torch.cuda, "is_available", lambda: True) + monkeypatch.setattr(torch.cuda, "device_count", lambda: 1) + monkeypatch.setattr( + torch.cuda, + "get_device_properties", + lambda _index: SimpleNamespace(uuid=uuid), + ) + + +class _FakeCudaProcessSharingEvidence: + configured = "mps" + effective = "mps" + validated = True + learner_device = "cuda:0" + collector_device = "cuda:0" + learner_gpu_uuid = None + collector_gpu_uuid = None + control_pipe = "/tmp/unilab-mps/control" + server_pid = 2768293 + + def manifest(self): + return { + "configured": self.configured, + "effective": self.effective, + "validated": self.validated, + "learner_device": self.learner_device, + "collector_device": self.collector_device, + "learner_gpu_uuid": self.learner_gpu_uuid, + "collector_gpu_uuid": self.collector_gpu_uuid, + "control_pipe": self.control_pipe, + "server_pid": self.server_pid, + } + + def test_offpolicy_config_has_one_replay_path(): cfg = _offpolicy_cfg() assert cfg.training.replay_prefetch_mode == "one_tick" @@ -98,6 +135,13 @@ def test_offpolicy_config_has_one_replay_path(): assert cfg.training.replay_ingress_slot_rows is None +@pytest.mark.parametrize("algo", ["sac", "flashsac", "warpsac"]) +def test_offpolicy_owners_default_cuda_process_sharing_off(algo: str): + cfg = _offpolicy_cfg(algo=algo) + + assert cfg.training.cuda_process_sharing is None + + def test_flashsac_scoped_tensor_benchmark_reduces_metric_flush_frequency(): cfg = _offpolicy_cfg( ["task=g1_motion_tracking/mjwarp"], @@ -119,6 +163,65 @@ def test_warpsac_declares_public_tensor_runtime_knobs(): assert cfg.training.replay_ingress_slot_rows is None +@pytest.mark.parametrize("algo", ["sac", "flashsac", "warpsac"]) +def test_cuda_process_sharing_request_fails_before_env_materialization( + monkeypatch: pytest.MonkeyPatch, + algo: str, +): + module = _offpolicy() + _cuda_torch_module(monkeypatch) + cfg = _offpolicy_cfg( + ["task=g1_walk_flat/mjwarp", "training.cuda_process_sharing=mps"], + algo=algo, + ) + + def reject_factory(*args, **kwargs): + del args, kwargs + raise AssertionError("invalid MPS topology must fail before env creation") + + monkeypatch.setattr(module, "registry_env_factory", reject_factory) + monkeypatch.setattr( + module, + "configure_backend_process_device", + lambda _backend, device: device, + ) + monkeypatch.setattr(cuda_process_sharing, "_nvidia_uuid", lambda *_args, **_kwargs: "A") + with pytest.raises( + ValueError, + match="could not reach the control daemon|found no control pipe", + ): + module.build_runner(algo, cfg) + + +def test_valid_cuda_process_sharing_evidence_enters_runner_manifest( + monkeypatch: pytest.MonkeyPatch, +): + module = _offpolicy() + _cuda_torch_module(monkeypatch) + cfg = _offpolicy_cfg(["task=g1_walk_flat/mjwarp", "training.cuda_process_sharing=mps"]) + evidence = _FakeCudaProcessSharingEvidence() + monkeypatch.setattr(module, "registry_env_factory", lambda *args, **kwargs: _fake_env_factory) + monkeypatch.setattr( + module, + "configure_backend_process_device", + lambda _backend, device: device, + ) + monkeypatch.setattr( + module, + "probe_cuda_process_sharing", + lambda *args, **kwargs: evidence, + ) + + import uni_rl.algos.fast_sac.double_buffer as owner_module + + monkeypatch.setattr(owner_module, "FastSACLearner", _FakeLearner) + monkeypatch.setattr(owner_module, "DoubleBufferOffPolicyRunner", _FakeRunner) + + runner = module.build_runner("sac", cfg) + + assert runner.runtime_manifest["cuda_process_sharing"] == evidence.manifest() + + @pytest.mark.parametrize("mode", ["invalid_mode", "same_tick"]) def test_non_one_tick_prefetch_is_rejected_before_dispatch(mode: str): cfg = _offpolicy_cfg([f"training.replay_prefetch_mode={mode}"]) diff --git a/tests/training/test_cuda_process_sharing.py b/tests/training/test_cuda_process_sharing.py new file mode 100644 index 000000000..ff396d7cc --- /dev/null +++ b/tests/training/test_cuda_process_sharing.py @@ -0,0 +1,257 @@ +"""Fail-closed CUDA process-sharing probe tests.""" + +from __future__ import annotations + +import platform +import socket +import stat +from types import SimpleNamespace +from typing import Any + +import pytest + +from unilab.training.cuda_process_sharing import probe_cuda_process_sharing + + +class _Completed: + def __init__(self, stdout: str = "", returncode: int = 0) -> None: + self.stdout = stdout + self.returncode = returncode + + +def _torch(uuid: str) -> Any: + properties = SimpleNamespace(uuid=uuid) + return SimpleNamespace( + version=SimpleNamespace(hip=None), + cuda=SimpleNamespace( + is_available=lambda: True, + device_count=lambda: 1, + get_device_properties=lambda _index: properties, + ), + ) + + +def _fake_commands(gpu_uuid: str = "GPU-a", *, server_pid: int = 2768293) -> Any: + def run_command(command: list[str], **kwargs: Any) -> _Completed: + del kwargs + if command[0] == "nvidia-smi": + return _Completed(f"0, {gpu_uuid}\n") + if command == ["nvidia-cuda-mps-control", "get-server-list"]: + return _Completed(f"{server_pid}\n") + raise AssertionError(f"unexpected command: {command}") + + return run_command + + +@pytest.fixture +def linux_mps(monkeypatch: pytest.MonkeyPatch, tmp_path: Any): + monkeypatch.setattr(platform, "system", lambda: "Linux") + monkeypatch.setenv("CUDA_MPS_PIPE_DIRECTORY", str(tmp_path / "pipe")) + monkeypatch.setenv("CUDA_MPS_LOG_DIRECTORY", str(tmp_path / "log")) + monkeypatch.setenv("CUDA_VISIBLE_DEVICES", "GPU-a") + control = tmp_path / "pipe" / "control" + control.parent.mkdir(parents=True) + control.touch() + control.chmod(0o666) + return control + + +def _socket_control(path: Any) -> None: + path.unlink() + listener = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) + listener.bind(str(path)) + path.chmod(0o666) + + +def test_disabled_default_does_not_probe_host(monkeypatch: pytest.MonkeyPatch) -> None: + def reject(*args: Any, **kwargs: Any) -> None: + del args, kwargs + raise AssertionError("disabled mode must not probe CUDA or MPS") + + monkeypatch.setattr(platform, "system", lambda: "not-linux") + evidence = probe_cuda_process_sharing( + None, + "cpu", + None, + backend="mujoco", + torch_module=None, + run_command=reject, + ) + + assert evidence.manifest() == { + "configured": None, + "effective": None, + "validated": False, + "learner_device": "cpu", + "collector_device": None, + "learner_gpu_uuid": None, + "collector_gpu_uuid": None, + "control_pipe": None, + "server_pid": None, + } + + +def test_valid_single_rank_request_records_server_evidence(linux_mps: Any) -> None: + _socket_control(linux_mps) + + evidence = probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=_fake_commands(), + ) + + assert evidence.manifest() == { + "configured": "mps", + "effective": "mps", + "validated": True, + "learner_device": "cuda:0", + "collector_device": "cuda:0", + "learner_gpu_uuid": "A", + "collector_gpu_uuid": "A", + "control_pipe": str(linux_mps), + "server_pid": 2768293, + } + + +def test_valid_identity_compares_physical_uuid_not_ordinal(linux_mps: Any) -> None: + _socket_control(linux_mps) + + evidence = probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=_fake_commands(), + ) + + assert evidence.validated + assert evidence.learner_gpu_uuid == evidence.collector_gpu_uuid + + +@pytest.mark.parametrize( + ("overrides", "message"), + [ + ({"requested": "invalid"}, "expected null or 'mps'"), + ({"backend": "mujoco"}, "only the mjwarp backend"), + ({"world_size": 2}, "single-host, single-rank"), + ({"collector_device": None}, "requires a CUDA collector"), + ({"learner_device": "cpu"}, "requires CUDA learner and collector"), + ], +) +def test_first_prerequisites_fail_closed( + linux_mps: Any, + monkeypatch: pytest.MonkeyPatch, + overrides: dict[str, Any], + message: str, +) -> None: + del monkeypatch + kwargs: dict[str, Any] = { + "requested": "mps", + "learner_device": "cuda:0", + "collector_device": "cuda:0", + "backend": "mjwarp", + "world_size": 1, + "torch_module": _torch("GPU-a"), + "run_command": _fake_commands(), + } + kwargs.update(overrides) + + with pytest.raises(ValueError, match=message): + probe_cuda_process_sharing(**kwargs) + + +def test_non_linux_rejects(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setattr(platform, "system", lambda: "Darwin") + + with pytest.raises(ValueError, match="requires Linux"): + probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=_fake_commands(), + ) + + +def test_unavailable_cuda_rejects(linux_mps: Any) -> None: + unavailable = _torch("GPU-a") + unavailable.cuda.is_available = lambda: False + + with pytest.raises(ValueError, match="CUDA is unavailable"): + probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=unavailable, + run_command=_fake_commands(), + ) + + +def test_gpu_identity_mismatch_rejects(linux_mps: Any) -> None: + _socket_control(linux_mps) + + with pytest.raises(ValueError, match=r"learner='A'.*collector='B'"): + probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=_fake_commands(gpu_uuid="GPU-b"), + ) + + +def test_missing_control_pipe_rejects_with_start_command( + linux_mps: Any, +) -> None: + linux_mps.unlink() + + with pytest.raises(ValueError, match="nvidia-cuda-mps-control -d"): + probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=_fake_commands(), + ) + + +def test_regular_file_is_not_a_control_pipe(linux_mps: Any) -> None: + assert not stat.S_ISSOCK(linux_mps.stat().st_mode) + + with pytest.raises(ValueError, match="not a daemon control pipe"): + probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=_fake_commands(), + ) + + +def test_unreachable_daemon_rejects_without_fallback(linux_mps: Any) -> None: + _socket_control(linux_mps) + + def failing_daemon(command: list[str], **kwargs: Any) -> _Completed: + del kwargs + if command[0] == "nvidia-smi": + return _Completed("0, GPU-a\n") + return _Completed("Cannot find MPS control daemon process", returncode=1) + + with pytest.raises(ValueError, match="could not reach the control daemon"): + probe_cuda_process_sharing( + "mps", + "cuda:0", + "cuda:0", + backend="mjwarp", + torch_module=_torch("GPU-a"), + run_command=failing_daemon, + ) diff --git a/tests/utils/test_experiment_tracking.py b/tests/utils/test_experiment_tracking.py index 24c616d99..c0b65706f 100644 --- a/tests/utils/test_experiment_tracking.py +++ b/tests/utils/test_experiment_tracking.py @@ -352,7 +352,7 @@ def test_experiment_tracker_writes_local_run_files(tmp_path, monkeypatch): task_name="G1MotionTracking", sim_backend="mujoco", training_cfg={"logger": "tensorboard"}, - full_cfg={"training": {"logger": "tensorboard"}}, + full_cfg={"training": {"logger": "tensorboard", "cuda_process_sharing": None}}, device="cuda", collector_device="cpu", seed_info={ @@ -375,6 +375,7 @@ def test_experiment_tracker_writes_local_run_files(tmp_path, monkeypatch): assert run_config["run"]["configured_seed_source"] == "algo.seed" assert run_config["run"]["effective_seed"] == 5 assert run_config["run"]["hardware"] == hardware + assert run_config["config"]["training"]["cuda_process_sharing"] is None assert run_summary["final_mean_reward"] == 12.3 assert run_summary["completed_iterations"] == 10 assert run_summary["configured_seed"] == 5