diff --git a/scripts/reinforcement_learning/rsl_rl/train.py b/scripts/reinforcement_learning/rsl_rl/train.py index bee6d5dc..e2aa49f0 100644 --- a/scripts/reinforcement_learning/rsl_rl/train.py +++ b/scripts/reinforcement_learning/rsl_rl/train.py @@ -38,6 +38,35 @@ parser.add_argument( "--ray-proc-id", "-rid", type=int, default=None, help="Automatically configured by Ray integration, otherwise None." ) +# --- null-space preference critic ------------------------------------------------------------- +parser.add_argument("--beta", type=float, default=None, help="Preference step budget. 0 == baseline PPO.") +parser.add_argument( + "--pref_source", type=str, default=None, choices=["zero", "noise", "action_rate", "ee_height", "terms"], + help="Preference stream: zero (sanity A), noise (sanity B), action_rate (noise-bait probe), " + "terms (scripted predicates).", +) +parser.add_argument("--pref_noise_std", type=float, default=None, help="Std for --pref_source=noise.") +parser.add_argument( + "--pref_terms", type=str, default=None, + help="Comma-separated RewardManager term names for --pref_source=terms.", +) +parser.add_argument( + "--critic_arch", type=str, default=None, choices=["shared", "separate"], + help="Second value head as a widened shared trunk, or an independent critic MLP.", +) +parser.add_argument("--gamma_pref", type=float, default=None, help="Discount for the preference return.") +parser.add_argument( + "--pref_mask_noise", type=lambda v: v.lower() not in ("0", "false", "no"), default=None, + help="Layer 1: keep the preference gradient off the exploration-noise params (default true).", +) +parser.add_argument( + "--pref_detach_noise_features", action="store_true", default=False, + help="Layer 2: also detach the trunk features feeding the gSDE noise head in the pref surrogate.", +) +parser.add_argument( + "--projection_mode", type=str, default=None, choices=["gradient", "advantage", "sum"], + help="gradient = faithful null-space projection; advantage/sum = ablations.", +) # append RSL-RL cli arguments cli_args.add_rsl_rl_args(parser) # append AppLauncher cli args @@ -87,6 +116,16 @@ from rsl_rl.runners import DistillationRunner, OnPolicyRunner +from uwlab_rl.rsl_rl.nullspace import ( + ActionRatePreference, + DualCriticOnPolicyRunner, + EndEffectorHeightPreference, + DualRewardVecEnvWrapper, + GaussianNoisePreference, + RewardManagerTermsPreference, + ZeroPreference, +) + from isaaclab.envs import ( DirectMARLEnv, DirectMARLEnvCfg, @@ -126,6 +165,28 @@ def main(env_cfg: ManagerBasedRLEnvCfg | DirectRLEnvCfg | DirectMARLEnvCfg, agen args_cli.max_iterations if args_cli.max_iterations is not None else agent_cfg.max_iterations ) + # --- null-space preference critic CLI overrides --- + # Applied before sanitize_rsl_rl_cfg, which only strips keys for algorithm classes it can + # resolve inside rsl_rl.algorithms; NullspacePPO lives in uwlab_rl, so these survive. + if args_cli.beta is not None: + agent_cfg.algorithm.beta = args_cli.beta + if args_cli.gamma_pref is not None: + agent_cfg.algorithm.gamma_pref = args_cli.gamma_pref + if args_cli.projection_mode is not None: + agent_cfg.algorithm.projection_mode = args_cli.projection_mode + if args_cli.pref_mask_noise is not None: + agent_cfg.algorithm.pref_mask_noise = args_cli.pref_mask_noise + if args_cli.pref_detach_noise_features: + agent_cfg.algorithm.pref_detach_noise_features = True + if args_cli.critic_arch is not None: + agent_cfg.policy.critic_arch = args_cli.critic_arch + if args_cli.pref_source is not None: + agent_cfg.pref_source = args_cli.pref_source + if args_cli.pref_noise_std is not None: + agent_cfg.pref_noise_std = args_cli.pref_noise_std + if args_cli.pref_terms is not None: + agent_cfg.pref_term_names = tuple(n.strip() for n in args_cli.pref_terms.split(",") if n.strip()) + # make config compatible with installed rsl-rl version agent_cfg = cli_args.sanitize_rsl_rl_cfg(agent_cfg) @@ -200,13 +261,61 @@ def main(env_cfg: ManagerBasedRLEnvCfg | DirectRLEnvCfg | DirectMARLEnvCfg, agen env = gym.wrappers.RecordVideo(env, **video_kwargs) # wrap around environment for rsl-rl - env = RslRlVecEnvWrapper(env, clip_actions=agent_cfg.clip_actions) + if agent_cfg.class_name == "DualCriticOnPolicyRunner": + # The dual-critic path needs a second reward stream. DualRewardVecEnvWrapper reads the + # RewardManager's already-materialised per-term buffer rather than restructuring the + # manager -- OmniReset's `progress_context` term returns zeros but caches state that the + # reward, terminations, reset curriculum and data-collection configs all read back, so + # splitting or reweighting the manager silently corrupts the reward. + source_name = getattr(agent_cfg, "pref_source", "zero") + if source_name == "zero": + pref_source = ZeroPreference() + elif source_name == "noise": + pref_source = GaussianNoisePreference( + std=getattr(agent_cfg, "pref_noise_std", 1.0), seed=agent_cfg.seed + ) + elif source_name == "action_rate": + # Noise-bait probe: a preference maximally satisfiable by shrinking exploration. + pref_source = ActionRatePreference() + elif source_name == "ee_height": + # High-conflict preference: fights the lift the task requires. + pref_source = EndEffectorHeightPreference() + elif source_name == "terms": + term_names = list(getattr(agent_cfg, "pref_term_names", ())) + if not term_names: + raise ValueError("--pref_source=terms requires --pref_terms=") + pref_source = RewardManagerTermsPreference(term_names) + else: + raise ValueError(f"Unknown pref_source: {source_name}") + print(f"[INFO] Preference reward source: {source_name}") + env = DualRewardVecEnvWrapper(env, pref_source=pref_source, clip_actions=agent_cfg.clip_actions) + else: + env = RslRlVecEnvWrapper(env, clip_actions=agent_cfg.clip_actions) + + # In-job wandb with a stable run id (gen_convergence_yaml.py sets WANDB_RUN_ID + WANDB_RESUME=allow): + # every preemption restart resumes the SAME wandb run instead of opening a new one. On resume wandb + # loads the run's stored config, and rsl_rl's WandbSummaryWriter then re-sends its own config with + # values that change on every start (log_dir, env_cfg.log_dir, resume settings). wandb raises + # ConfigError on a changed value unless allow_val_change=True, which would kill the start before + # training -- so let config updates overwrite, in this process only. + if agent_cfg.logger == "wandb" and os.environ.get("WANDB_RUN_ID"): + import wandb.sdk.wandb_config as _wandb_config + + _config_update = _wandb_config.Config.update + + def _update_allow_change(self, d, allow_val_change=None): + return _config_update(self, d, allow_val_change=True) + + _wandb_config.Config.update = _update_allow_change + print(f"[INFO] wandb: resuming run id {os.environ['WANDB_RUN_ID']} (config updates may overwrite)") # create runner from rsl-rl if agent_cfg.class_name == "OnPolicyRunner": runner = OnPolicyRunner(env, agent_cfg.to_dict(), log_dir=log_dir, device=agent_cfg.device) elif agent_cfg.class_name == "DistillationRunner": runner = DistillationRunner(env, agent_cfg.to_dict(), log_dir=log_dir, device=agent_cfg.device) + elif agent_cfg.class_name == "DualCriticOnPolicyRunner": + runner = DualCriticOnPolicyRunner(env, agent_cfg.to_dict(), log_dir=log_dir, device=agent_cfg.device) else: raise ValueError(f"Unsupported runner class: {agent_cfg.class_name}") # write git state to logs diff --git a/source/uwlab_assets/uwlab_assets/__init__.py b/source/uwlab_assets/uwlab_assets/__init__.py index 0bf6a5f1..8ded0ebd 100644 --- a/source/uwlab_assets/uwlab_assets/__init__.py +++ b/source/uwlab_assets/uwlab_assets/__init__.py @@ -65,7 +65,12 @@ def resolve_cloud_path(path: str) -> str: return path rel = _extract_relative_path(path) - cache_dir = os.path.join(os.path.expanduser("~"), ".cache", "uwlab", "assets") + # Overridable via UWLAB_ASSET_CACHE_DIR. On a cluster ``$HOME`` is ephemeral, so the default + # would re-download ~7 GB of USD assets on every job; pointing this at a persistent (writable) + # mount makes the first job populate the cache and every later job hit it. + cache_dir = os.getenv("UWLAB_ASSET_CACHE_DIR") or os.path.join( + os.path.expanduser("~"), ".cache", "uwlab", "assets" + ) local = os.path.join(cache_dir, rel) if os.path.isfile(local): diff --git a/source/uwlab_rl/test/nullspace/test_grad_clip_and_pref_norm.py b/source/uwlab_rl/test/nullspace/test_grad_clip_and_pref_norm.py new file mode 100644 index 00000000..f5db3ed6 --- /dev/null +++ b/source/uwlab_rl/test/nullspace/test_grad_clip_and_pref_norm.py @@ -0,0 +1,217 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Tests for the gradient-clip coupling fix (NOTES 29). + +The bug: rsl_rl clips the whole policy as ONE gradient vector. With a dual critic, a large +preference value loss enters the same norm as the actor's gradient, so clipping scales the actor's +update down -- at beta=0 too, where the preference is meant to have no effect at all. Under +``pref_source=action_rate`` the preference value loss reached 1e5-1e6 against +``max_grad_norm=1.0``, which throttled the actor and pinned the adaptive learning rate at its cap. + +``projection.py`` imports nothing but torch, so it is loaded by path (as in ``test_projection.py``) +and none of these tests need Isaac Sim. +""" + +from __future__ import annotations + +import importlib.util +import io +import pathlib + +import pytest +import torch + +_PATH = pathlib.Path(__file__).resolve().parents[2] / "uwlab_rl" / "rsl_rl" / "nullspace" / "projection.py" +_spec = importlib.util.spec_from_file_location("nsc_projection_gradclip", _PATH) +_mod = importlib.util.module_from_spec(_spec) +_spec.loader.exec_module(_mod) +partition_policy_params = _mod.partition_policy_params +clip_grad_norm_by_group = _mod.clip_grad_norm_by_group + +MAX_NORM = 1.0 # the OmniReset agent config's max_grad_norm + + +class _DualPolicy(torch.nn.Module): + """Stand-in carrying DualCriticActorCritic's real parameter names.""" + + def __init__(self, separate_pref_critic: bool = True) -> None: + super().__init__() + self.actor = torch.nn.Linear(4, 2) + self.log_std = torch.nn.Parameter(torch.zeros(2)) # gSDE noise scale -- belongs to the actor + self.critic = torch.nn.Linear(4, 1) + if separate_pref_critic: + self.critic_pref = torch.nn.Linear(4, 1) + + +def _fill(params, value: float) -> None: + for p in params: + p.grad = torch.full_like(p, value) + + +def _flat(params) -> torch.Tensor: + return torch.cat([p.grad.flatten() for p in params]).clone() + + +def _groups(policy): + return partition_policy_params(policy.named_parameters()) + + +# -- the bug, reproduced, and the fix ------------------------------------------------------------ + + +def test_global_clip_lets_the_pref_critic_throttle_the_actor(): + """Documents the bug: under one global clip the actor's update shrinks because the + PREFERENCE CRITIC's gradient is large -- even though the actor's own gradient is small.""" + p = _DualPolicy() + g = _groups(p) + _fill(g["actor"] + g["critic"], 0.1) # actor norm ~0.35: under the limit on its own + _fill(g["critic_pref"], 1e3) # the action_rate preference critic's scale + before = _flat(g["actor"]) + assert before.norm() < MAX_NORM + + torch.nn.utils.clip_grad_norm_(p.parameters(), MAX_NORM) # upstream behaviour + + assert _flat(g["actor"]).norm() / before.norm() < 1e-2 # crushed by a different module + + +def test_per_group_clip_leaves_the_actor_untouched_by_the_pref_critic(): + p = _DualPolicy() + g = _groups(p) + _fill(g["actor"] + g["critic"], 0.1) + _fill(g["critic_pref"], 1e3) + before = _flat(g["actor"]) + + norms = clip_grad_norm_by_group(g, MAX_NORM) + + assert torch.allclose(_flat(g["actor"]), before) # below its own limit -> not scaled at all + assert _flat(g["critic_pref"]).norm() == pytest.approx(MAX_NORM, rel=1e-4) # clipped on its own + assert norms["actor"] == pytest.approx(before.norm().item(), rel=1e-5) # pre-clip norms reported + assert norms["critic_pref"] > 1e3 + + +def test_actor_update_is_invariant_to_pref_critic_scale(): + """The property the acceptance run checks at the outcome level: at beta=0 the actor's clipped + gradient must not depend on how large the preference critic's gradient is. Under the global + clip it does; under per-group clipping it does not.""" + per_group, global_ = [], [] + for pref_scale in (0.0, 1.0, 1e3, 1e6): + for mode, out in (("per_group", per_group), ("global", global_)): + p = _DualPolicy() + torch.manual_seed(0) + g = _groups(p) + _fill(g["actor"] + g["critic"], 0.1) + _fill(g["critic_pref"], pref_scale) + if mode == "per_group": + clip_grad_norm_by_group(g, MAX_NORM) + else: + torch.nn.utils.clip_grad_norm_(p.parameters(), MAX_NORM) + out.append(_flat(g["actor"])) + + assert all(torch.allclose(a, per_group[0]) for a in per_group) + assert not all(torch.allclose(a, global_[0]) for a in global_) + + +def test_per_group_still_clips_an_oversized_actor(): + """Per-group clipping must not silently disable clipping for the actor.""" + p = _DualPolicy() + g = _groups(p) + _fill(g["actor"], 10.0) + _fill(g["critic"] + g["critic_pref"], 0.1) + + clip_grad_norm_by_group(g, MAX_NORM) + + assert _flat(g["actor"]).norm() == pytest.approx(MAX_NORM, rel=1e-4) + + +# -- partition ------------------------------------------------------------------------------------ + + +def test_partition_follows_actor_parameters_and_avoids_the_prefix_trap(): + """``critic`` is a prefix of ``critic_pref``; the preference critic must not land in the task + critic's group, and ``log_std`` must stay with the actor (as in ``actor_parameters``).""" + p = _DualPolicy() + ids = {k: {id(x) for x in v} for k, v in _groups(p).items()} + + assert ids["actor"] == {id(p.actor.weight), id(p.actor.bias), id(p.log_std)} + assert ids["critic"] == {id(p.critic.weight), id(p.critic.bias)} + assert ids["critic_pref"] == {id(p.critic_pref.weight), id(p.critic_pref.bias)} + + +def test_shared_critic_arch_leaves_pref_group_empty_and_is_safe(): + """``critic_arch='shared'`` has no separate preference critic; params without a gradient are + skipped rather than erroring.""" + p = _DualPolicy(separate_pref_critic=False) + g = _groups(p) + assert g["critic_pref"] == [] + _fill(g["critic"], 0.1) # actor deliberately left with grad=None + + norms = clip_grad_norm_by_group(g, MAX_NORM) + + assert norms["critic_pref"] == 0.0 and norms["actor"] == 0.0 and norms["critic"] > 0.0 + + +# -- preference reward normalisation ------------------------------------------------------------- + +_networks = pytest.importorskip("rsl_rl.networks") +EDVN = _networks.EmpiricalDiscountedVariationNormalization + + +def test_pref_reward_normalizer_brings_action_rate_scale_rewards_to_order_one(): + torch.manual_seed(0) + n = EDVN(shape=1, gamma=0.99) + n.train() + for _ in range(300): + raw = -(torch.rand(4096, 1) * 3e3) # -||a - a_prev||^2 at the magnitudes that were observed + out = n(raw) + + assert raw.abs().mean() > 1e3 + assert 1e-3 < out.abs().mean().item() < 10.0 + + +def test_zero_preference_stays_exactly_zero(): + """``pref_source='zero'`` must be untouched: no NaN from a vanishing std, no drift off zero.""" + n = EDVN(shape=1, gamma=0.99) + n.train() + for _ in range(200): + out = n(torch.zeros(4096, 1)) + + assert torch.isfinite(out).all() + assert torch.equal(out, torch.zeros_like(out)) + + +def test_normalizer_works_under_inference_mode_and_survives_a_checkpoint_round_trip(): + """rsl_rl collects rollouts inside ``torch.inference_mode()`` (on_policy_runner.py:101), so the + normaliser's buffers are rewritten there. They must still be readable for logging, serialisable + by ``torch.save``, and loadable into a freshly built policy on resume.""" + n = EDVN(shape=1, gamma=0.99) + n.train() + with torch.inference_mode(): + for _ in range(20): + n(-(torch.rand(512, 1) * 1e3)) + scale = float(n.emp_norm._std) # read outside inference mode, as update() logging does + assert scale > 1.0 + + buf = io.BytesIO() + torch.save(n.state_dict(), buf) # what the runner's save() does + buf.seek(0) + fresh = EDVN(shape=1, gamma=0.99) + fresh.load_state_dict(torch.load(buf)) # what a resumed run does, outside inference mode + assert float(fresh.emp_norm._std) == pytest.approx(scale) + + fresh.train() + with torch.inference_mode(): + out = fresh(-(torch.rand(512, 1) * 1e3)) + assert torch.isfinite(out).all() + + +def test_normalizer_on_the_policy_is_checkpointed_and_adds_no_parameters(): + """Attached to the policy it rides along in ``policy.state_dict()`` -- and, holding only + buffers, it cannot enter a clip group or the optimizer.""" + p = _DualPolicy() + p.pref_reward_normalizer = EDVN(shape=1, gamma=0.99) + + assert any(k.startswith("pref_reward_normalizer.emp_norm.") for k in p.state_dict()) + assert not any(n.startswith("pref_reward_normalizer") for n, _ in p.named_parameters()) + grouped = {id(x) for v in _groups(p).values() for x in v} + assert grouped == {id(x) for x in p.parameters()} # every parameter grouped exactly once diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/__init__.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/__init__.py new file mode 100644 index 00000000..5786aaab --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/__init__.py @@ -0,0 +1,35 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Dual-critic PPO with a null-space-projected preference gradient.""" + +from .dual_actor_critic import DualCriticActorCritic +from .dual_storage import DualRolloutStorage +from .nullspace_ppo import NullspacePPO +from .projection import project_nullspace, project_nullspace_masked +from .reward_split import ( + ActionRatePreference, + DualRewardVecEnvWrapper, + EndEffectorHeightPreference, + GaussianNoisePreference, + PreferenceRewardSource, + RewardManagerTermsPreference, + ZeroPreference, +) +from .runner import DualCriticOnPolicyRunner + +__all__ = [ + "ActionRatePreference", + "DualCriticActorCritic", + "DualCriticOnPolicyRunner", + "DualRewardVecEnvWrapper", + "DualRolloutStorage", + "EndEffectorHeightPreference", + "GaussianNoisePreference", + "NullspacePPO", + "PreferenceRewardSource", + "RewardManagerTermsPreference", + "ZeroPreference", + "project_nullspace", + "project_nullspace_masked", +] diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/dual_actor_critic.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/dual_actor_critic.py new file mode 100644 index 00000000..2a53cc96 --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/dual_actor_critic.py @@ -0,0 +1,180 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Actor-critic with two value heads: one for task return, one for preference return. + +Two architectures, both worth measuring (see design notes §3.2): + +``shared`` one critic trunk, final layer widened 1 -> 2 outputs. Cheapest, and matches the + "shared trunk, second output head" of the design. Caveat: the preference value loss + backpropagates through the shared trunk, so the two heads are coupled *even at β=0*. + That coupling is exactly what sanity run A measures. +``separate`` a second independent critic MLP. Strictly more parameters, but zero coupling at + β=0, so the β=0 overlay is guaranteed by construction rather than by experiment. + +The base class' ``evaluate`` keeps returning the *task* value with shape (N, 1) so that any +untouched upstream code path behaves exactly as before. +""" + +from __future__ import annotations + +import torch +from tensordict import TensorDict +from typing import Any + +from rsl_rl.modules import ActorCritic +from rsl_rl.networks import MLP + +CRITIC_ARCHES = ("shared", "separate") + + +class DualCriticActorCritic(ActorCritic): + """Asymmetric actor-critic (critic sees privileged obs) with task and preference value heads.""" + + def __init__( + self, + obs: TensorDict, + obs_groups: dict[str, list[str]], + num_actions: int, + critic_hidden_dims: tuple[int] | list[int] = [256, 256, 256], + activation: str = "elu", + critic_arch: str = "shared", + **kwargs: dict[str, Any], + ) -> None: + if critic_arch not in CRITIC_ARCHES: + raise ValueError(f"critic_arch must be one of {CRITIC_ARCHES}, got {critic_arch!r}") + + super().__init__( + obs, + obs_groups, + num_actions, + critic_hidden_dims=critic_hidden_dims, + activation=activation, + **kwargs, + ) + + self.critic_arch = critic_arch + + num_critic_obs = sum(obs[g].shape[-1] for g in obs_groups["critic"]) + + if critic_arch == "shared": + # Replace the 1-output critic built by the base class with a 2-output one. + self.critic = MLP(num_critic_obs, 2, critic_hidden_dims, activation) + print(f"Critic MLP (shared trunk, 2 heads): {self.critic}") + else: + self.critic_pref = MLP(num_critic_obs, 1, critic_hidden_dims, activation) + print(f"Critic MLP (separate preference critic): {self.critic_pref}") + + # -- parameter partition ------------------------------------------------------------------- + # The projection is applied to the *actor* gradient only; the value losses train the critics + # normally. EmpiricalNormalization holds buffers, not parameters, so a name-prefix split is + # exact here. + + #: Parameter names that scale exploration noise rather than the deployed mean action. + NOISE_PARAM_NAMES = ("std", "log_std") + + def actor_parameters(self) -> list[torch.nn.Parameter]: + """Actor MLP plus the exploration-noise parameters (``std`` / ``log_std`` for gSDE).""" + return [p for n, p in self.named_parameters() if not n.startswith("critic")] + + def critic_parameters(self) -> list[torch.nn.Parameter]: + return [p for n, p in self.named_parameters() if n.startswith("critic")] + + def is_noise_param_name(self, name: str) -> bool: + return name.split(".")[0] in self.NOISE_PARAM_NAMES + + def actor_param_noise_mask(self) -> list[bool]: + """Per-entry flags over ``actor_parameters()``: True where the parameter is noise scale. + + The preference objective is a statement about *deployed behaviour*, and the policy is + deployed on the mean action -- gSDE noise does not exist at deployment time. So the + preference gradient has no business touching these. + """ + return [self.is_noise_param_name(n) for n, _ in self.named_parameters() if not n.startswith("critic")] + + # -- preference-scoped log probability ------------------------------------------------------ + + def log_prob_mean_path(self, obs: TensorDict, actions: torch.Tensor) -> torch.Tensor: + """Log-prob whose gradient reaches the network **only through the mean action**. + + Numerically identical to the usual log-prob -- ``detach`` changes no values -- so the PPO + importance ratio built from it is the same number, and "one ratio, one clip" still holds. + Only the backward graph differs. + + Why this is needed (Layer 2): gSDE's action variance is + ``mm(features**2, exp(log_std)**2)`` where ``features = actor[:-1](obs)`` + (rsl_rl/modules/actor_critic.py:72, :283). Excluding ``log_std`` from the preference + gradient (Layer 1) therefore does *not* close the path -- the preference can still shrink + exploration by reshaping the trunk features that feed the noise head. Detaching both + inputs to the variance closes it. + """ + obs = self.get_actor_obs(obs) + obs = self.actor_obs_normalizer(obs) + mean = self.actor(obs) + + if self.noise_std_type == "gsde": + features = self.actor[:-1](obs).detach() + std = torch.sqrt( + torch.mm(features**2, torch.exp(self.log_std.detach()) ** 2) + self.distribution.epsilon + ) + elif self.noise_std_type == "scalar": + std = self.std.detach().expand_as(mean) + else: # "log" + std = torch.exp(self.log_std.detach()).expand_as(mean) + + return torch.distributions.Normal(mean, std).log_prob(actions).sum(dim=-1) + + # -- diagnostics --------------------------------------------------------------------------- + + def noise_magnitude(self) -> torch.Tensor: + """Mean realised action std of the current distribution (the exploration-leak canary).""" + return self.action_std.mean().detach() + + @torch.no_grad() + def noise_decomposition(self, obs: TensorDict) -> tuple[float, float]: + """Split realised noise into its two multiplicative factors. + + gSDE variance is ``mm(f(obs)**2, exp(log_std)**2)`` with ``f = actor[:-1]``, so + + log(realised std) ≈ log‖f(obs)‖ + log σ + trunk path direct path + (Layer 2) (Layer 1) + + Layer 1 masks **only the second factor**, so an aggregate fall in realised noise cannot + distinguish "the trunk reshaped its features" from "sigma shrank". Reporting them + separately makes the Layer 2 decision mechanical: + + sigma flat, ‖f‖ flat -> no leak + sigma flat, ‖f‖ collapsing -> trunk path, Layer 2 indicated + sigma moving at all -> Layer 1 implementation bug (it is excluded from g_pref + by construction, so nothing else can move it) + + Returns ``(mean ‖f(obs)‖, mean sigma)``. + """ + if self.noise_std_type != "gsde": + std = self.std if self.noise_std_type == "scalar" else torch.exp(self.log_std) + return 1.0, float(std.mean()) + x = self.actor_obs_normalizer(self.get_actor_obs(obs)) + feat = self.actor[:-1](x) + return float(feat.norm(dim=-1).mean()), float(torch.exp(self.log_std).mean()) + + # -- value heads --------------------------------------------------------------------------- + + def _critic_features(self, obs: TensorDict) -> torch.Tensor: + obs = self.get_critic_obs(obs) + return self.critic_obs_normalizer(obs) + + def evaluate_dual(self, obs: TensorDict, **kwargs: dict[str, Any]) -> tuple[torch.Tensor, torch.Tensor]: + """Return ``(value_task, value_pref)``, each shaped (N, 1).""" + x = self._critic_features(obs) + if self.critic_arch == "shared": + out = self.critic(x) + return out[..., 0:1], out[..., 1:2] + return self.critic(x), self.critic_pref(x) + + def evaluate(self, obs: TensorDict, **kwargs: dict[str, Any]) -> torch.Tensor: + """Task value only, shape (N, 1) -- preserves the upstream contract.""" + x = self._critic_features(obs) + if self.critic_arch == "shared": + return self.critic(x)[..., 0:1] + return self.critic(x) diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/dual_storage.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/dual_storage.py new file mode 100644 index 00000000..e26d0e98 --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/dual_storage.py @@ -0,0 +1,169 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Rollout storage carrying two reward/value/return/advantage streams. + +The two streams get **separately normalised advantages**. That is the detail that makes the +reward model's arbitrary output scale cancel, so β controls direction budget only and does not +have to be re-tuned per task or per reward-model retrain. + +Upstream normalises once over the whole (num_steps, num_envs) batch rather than per minibatch +(``normalize_advantage_per_mini_batch=False`` in the OmniReset agent config), so each stream is +normalised over its own full batch here. +""" + +from __future__ import annotations + +import torch +from collections.abc import Generator +from tensordict import TensorDict + +from rsl_rl.storage import RolloutStorage + + +class DualRolloutStorage(RolloutStorage): + """RolloutStorage plus a parallel preference stream.""" + + class Transition(RolloutStorage.Transition): + def __init__(self) -> None: + super().__init__() + self.rewards_pref: torch.Tensor | None = None + self.values_pref: torch.Tensor | None = None + + def __init__( + self, + training_type: str, + num_envs: int, + num_transitions_per_env: int, + obs: TensorDict, + actions_shape: tuple[int] | list[int], + device: str = "cpu", + ) -> None: + super().__init__(training_type, num_envs, num_transitions_per_env, obs, actions_shape, device) + + if training_type != "rl": + raise ValueError("DualRolloutStorage only supports training_type='rl'.") + + z = lambda: torch.zeros(num_transitions_per_env, num_envs, 1, device=self.device) # noqa: E731 + self.rewards_pref = z() + self.values_pref = z() + self.returns_pref = z() + self.advantages_pref = z() + + def add_transitions(self, transition: Transition) -> None: + # NOTE: the base implementation increments self.step, so capture it first. + step = self.step + super().add_transitions(transition) + self.rewards_pref[step].copy_(transition.rewards_pref.view(-1, 1)) + self.values_pref[step].copy_(transition.values_pref) + + # -- GAE --------------------------------------------------------------------------------- + + @staticmethod + def _gae( + rewards: torch.Tensor, + values: torch.Tensor, + dones: torch.Tensor, + last_values: torch.Tensor, + gamma: float, + lam: float, + num_steps: int, + ) -> tuple[torch.Tensor, torch.Tensor]: + """Standard GAE. Returns (returns, advantages), unnormalised.""" + returns = torch.zeros_like(values) + advantage = 0 + for step in reversed(range(num_steps)): + next_values = last_values if step == num_steps - 1 else values[step + 1] + next_is_not_terminal = 1.0 - dones[step].float() + delta = rewards[step] + next_is_not_terminal * gamma * next_values - values[step] + advantage = delta + next_is_not_terminal * gamma * lam * advantage + returns[step] = advantage + values[step] + return returns, returns - values + + def compute_returns_dual( + self, + last_values: torch.Tensor, + last_values_pref: torch.Tensor, + gamma: float, + lam: float, + gamma_pref: float, + lam_pref: float, + normalize_advantage: bool = True, + ) -> None: + """Two independent GAE passes with their own discounts, normalised separately. + + ``gamma_pref << gamma`` is expected to help: manner is mostly local while success is + long-horizon. It is a configurable ablation, not a hidden default. + """ + n = self.num_transitions_per_env + + self.returns, self.advantages = self._gae( + self.rewards, self.values, self.dones, last_values, gamma, lam, n + ) + self.returns_pref, self.advantages_pref = self._gae( + self.rewards_pref, self.values_pref, self.dones, last_values_pref, gamma_pref, lam_pref, n + ) + + if normalize_advantage: + self.advantages = self._normalize(self.advantages) + # Guarded: with a zero preference reward (sanity run A) the advantages are identically + # zero, and dividing by their std would turn 0/1e-8 into noise. Leave them at zero. + self.advantages_pref = self._normalize(self.advantages_pref) + + @staticmethod + def _normalize(adv: torch.Tensor) -> torch.Tensor: + std = adv.std() + if not torch.isfinite(std) or std < 1e-8: + return torch.zeros_like(adv) + return (adv - adv.mean()) / (std + 1e-8) + + # -- minibatches ------------------------------------------------------------------------- + + def mini_batch_generator(self, num_mini_batches: int, num_epochs: int = 8) -> Generator: + """Same as upstream, plus the preference target values / advantages / returns. + + Reimplemented rather than wrapped because the upstream generator draws its own random + permutation; both streams must be indexed by the *same* permutation. + """ + batch_size = self.num_envs * self.num_transitions_per_env + mini_batch_size = batch_size // num_mini_batches + indices = torch.randperm(num_mini_batches * mini_batch_size, requires_grad=False, device=self.device) + + observations = self.observations.flatten(0, 1) + actions = self.actions.flatten(0, 1) + values = self.values.flatten(0, 1) + returns = self.returns.flatten(0, 1) + old_actions_log_prob = self.actions_log_prob.flatten(0, 1) + advantages = self.advantages.flatten(0, 1) + old_mu = self.mu.flatten(0, 1) + old_sigma = self.sigma.flatten(0, 1) + + values_pref = self.values_pref.flatten(0, 1) + returns_pref = self.returns_pref.flatten(0, 1) + advantages_pref = self.advantages_pref.flatten(0, 1) + + for _ in range(num_epochs): + for i in range(num_mini_batches): + batch_idx = indices[i * mini_batch_size : (i + 1) * mini_batch_size] + yield ( + observations[batch_idx], + actions[batch_idx], + values[batch_idx], + advantages[batch_idx], + returns[batch_idx], + old_actions_log_prob[batch_idx], + old_mu[batch_idx], + old_sigma[batch_idx], + (None, None), + None, + # preference stream + values_pref[batch_idx], + advantages_pref[batch_idx], + returns_pref[batch_idx], + ) + + def recurrent_mini_batch_generator(self, num_mini_batches: int, num_epochs: int = 8) -> Generator: + raise NotImplementedError( + "DualRolloutStorage does not support recurrent policies. The OmniReset agent config uses " + "a feedforward ActorCritic (gSDE noise), so this path is unused; implement it if that changes." + ) diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/nullspace_ppo.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/nullspace_ppo.py new file mode 100644 index 00000000..f2d3f663 --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/nullspace_ppo.py @@ -0,0 +1,543 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""PPO with a dual critic and a null-space-projected preference gradient. + +Structure of one update, and why each piece is where it is: + +* **One importance ratio, one clip.** Upstream already computes ``ratio`` once; both the task and + preference surrogates are built from that same clipped ratio. This is not two PPO updates. +* **The projection is applied to gradients, not advantages.** The first-order guarantee in the + design is a statement about ``g_task`` and ``g_pref``; combining advantages before the surrogate + is a per-sample reweighting that does *not* satisfy it. Phase 0 measured the PPO update at + ~1% of iteration wall-clock (collection dominates ~99%), so the extra backward costs ~1%. + ``projection_mode="advantage"`` implements the cheap approximation as an ablation. +* **Separately normalised advantages** (in :class:`DualRolloutStorage`) are what make β + dimensionless and independent of the preference reward's scale. +""" + +from __future__ import annotations + +import importlib.metadata as metadata +import inspect +import torch +import torch.nn as nn +from tensordict import TensorDict + +from rsl_rl.algorithms import PPO +from rsl_rl.networks import EmpiricalDiscountedVariationNormalization + +from .dual_storage import DualRolloutStorage +from .projection import ( + clip_grad_norm_by_group, + flat_mask, + flatten_grads, + partition_policy_params, + project_nullspace, + project_nullspace_masked, + unflatten_to, +) + +PROJECTION_MODES = ("gradient", "advantage", "sum") + + +class NullspacePPO(PPO): + """PPO whose preference objective is projected into the null space of the task gradient.""" + + def __init__( + self, + policy, # noqa: ANN001 -- DualCriticActorCritic + beta: float = 0.0, + gamma_pref: float | None = None, + lam_pref: float | None = None, + projection_mode: str = "gradient", + pref_value_loss_coef: float | None = None, + pref_mask_noise: bool = True, + pref_detach_noise_features: bool = False, + log_pref_alignment: bool = True, + grad_clip_mode: str = "per_group", + normalize_pref_reward: bool = True, + **kwargs, + ) -> None: + # IsaacLab's RslRlPpoAlgorithmCfg carries fields that only exist in rsl-rl >= 4.x + # (e.g. `share_cnn_encoders`, `optimizer`). Upstream's `sanitize_rsl_rl_cfg` strips those + # by resolving the class out of `rsl_rl.algorithms` -- which cannot find NullspacePPO, + # so they arrive here intact and PPO 3.1.2 rejects them. Drop them explicitly, loudly. + accepted = set(inspect.signature(PPO.__init__).parameters) - {"self", "policy"} + dropped = sorted(set(kwargs) - accepted) + if dropped: + print( + f"[NullspacePPO] Dropping algorithm config keys unsupported by the installed " + f"rsl-rl ({metadata.version('rsl-rl-lib')}): {dropped}" + ) + super().__init__(policy, **{k: v for k, v in kwargs.items() if k in accepted}) + + if projection_mode not in PROJECTION_MODES: + raise ValueError(f"projection_mode must be one of {PROJECTION_MODES}, got {projection_mode!r}") + if self.symmetry is not None: + raise NotImplementedError("NullspacePPO does not support symmetry augmentation.") + if self.rnd is not None: + raise NotImplementedError( + "NullspacePPO does not support RND. RND folds a second reward stream into a single " + "critic by weighted sum -- that is the baseline this method is compared against." + ) + if not hasattr(policy, "evaluate_dual"): + raise TypeError("NullspacePPO requires a DualCriticActorCritic (missing evaluate_dual).") + + self.beta = beta + # Manner is mostly local, success is long-horizon; default to matching the task discount so + # that any difference is an explicit experimental choice rather than a hidden default. + self.gamma_pref = self.gamma if gamma_pref is None else gamma_pref + self.lam_pref = self.lam if lam_pref is None else lam_pref + self.projection_mode = projection_mode + self.pref_value_loss_coef = ( + self.value_loss_coef if pref_value_loss_coef is None else pref_value_loss_coef + ) + + # --- exploration-leak scoping ----------------------------------------------------- + # Every scripted manner preference (mechanical power, EE speed, action smoothness) is + # monotonically improved by shrinking action noise, while perturbing noise around a local + # optimum of the mean policy is ~second order in task return. Large first-order preference + # gradient against ~zero first-order task gradient means the projector does not merely + # permit that direction -- it *selects* it. That is entropy collapse arriving through the + # exact channel built to find task-neutral directions. + # + # gSDE noise is an optimisation parameter, not a deployed behavioural property (the policy + # is deployed on the mean action), so no legitimate preference is given up by masking it. + self.pref_mask_noise = pref_mask_noise # Layer 1: exclude noise params from g_pref + self.pref_detach_noise_features = pref_detach_noise_features # Layer 2: cut the trunk path + self._pref_allowed: torch.Tensor | None = None # lazily built flat subspace mask + # pref_removed_frac is the screening instrument for preference/task conflict, so it is + # logged in EVERY arm as a time series -- including beta=0, where nothing is applied. + self.log_pref_alignment = log_pref_alignment + + # --- gradient-norm clipping scope (NOTES 29) ---------------------------------------- + # Upstream clips the whole policy as one vector. For a dual critic that couples the + # preference critic to the actor: its value loss enters the same norm, and a large one + # scales the actor's update down -- at beta=0 too. Clip each group on its own budget. + if grad_clip_mode not in ("per_group", "global"): + raise ValueError(f"grad_clip_mode must be 'per_group' or 'global', got {grad_clip_mode!r}") + self.grad_clip_mode = grad_clip_mode + + # --- preference reward scale (NOTES 29) --------------------------------------------- + # The preference critic regresses onto returns of an arbitrary-scale reward; with + # action_rate its loss reached 1e5-1e6. A positive running scale leaves the normalised + # preference advantages unchanged, so this fixes the critic without changing what the + # actor sees. Registered on the policy so its running stats are checkpointed with it. + # The task stream is deliberately NOT normalised. + self.normalize_pref_reward = normalize_pref_reward + if normalize_pref_reward: + self.policy.pref_reward_normalizer = EmpiricalDiscountedVariationNormalization( + shape=1, gamma=self.gamma_pref + ).to(self.device) + + self.transition = DualRolloutStorage.Transition() + self._diag: dict[str, float] = {} + + # -- storage --------------------------------------------------------------------------- + + def init_storage(self, training_type, num_envs, num_transitions_per_env, obs, actions_shape) -> None: # noqa: ANN001 + self.storage = DualRolloutStorage( + training_type, num_envs, num_transitions_per_env, obs, actions_shape, self.device + ) + + # -- rollout --------------------------------------------------------------------------- + + def act(self, obs: TensorDict) -> torch.Tensor: + if self.policy.is_recurrent: + self.transition.hidden_states = self.policy.get_hidden_states() + self.transition.actions = self.policy.act(obs).detach() + # One critic forward for both heads. + v_task, v_pref = self.policy.evaluate_dual(obs) + self.transition.values = v_task.detach() + self.transition.values_pref = v_pref.detach() + self.transition.actions_log_prob = self.policy.get_actions_log_prob(self.transition.actions).detach() + self.transition.action_mean = self.policy.action_mean.detach() + self.transition.action_sigma = self.policy.action_std.detach() + self.transition.observations = obs + return self.transition.actions + + def process_env_step( + self, obs: TensorDict, rewards: torch.Tensor, dones: torch.Tensor, extras: dict + ) -> None: + """Record a transition. + + Reimplemented rather than delegated because the time-out bootstrap must be applied **per + stream**, each with its own discount and value head. Episodes here are 16 s / 160 control + steps and ``time_out`` is the dominant termination, so this path fires constantly -- using + the task discount for the preference return would be wrong at every episode boundary. + """ + self.policy.update_normalization(obs) + + rewards_pref = extras.get("reward_pref") + if rewards_pref is None: + raise KeyError( + "extras['reward_pref'] missing. Wrap the env in DualRewardVecEnvWrapper so the " + "preference stream is emitted alongside the task reward." + ) + + self.transition.rewards = rewards.clone() + self.transition.rewards_pref = rewards_pref.clone().to(self.device) + if self.normalize_pref_reward: + # Before the time-out bootstrap below: that adds gamma_pref * values_pref, and the + # preference critic's values live in this normalised scale once it trains on it. + self.transition.rewards_pref = self.policy.pref_reward_normalizer( + self.transition.rewards_pref.view(-1, 1) + ).view(-1) + self.transition.dones = dones + + if "time_outs" in extras: + time_outs = extras["time_outs"].unsqueeze(1).to(self.device) + self.transition.rewards += self.gamma * torch.squeeze(self.transition.values * time_outs, 1) + self.transition.rewards_pref += self.gamma_pref * torch.squeeze( + self.transition.values_pref * time_outs, 1 + ) + + self.storage.add_transitions(self.transition) + self.transition.clear() + self.policy.reset(dones) + + def compute_returns(self, obs: TensorDict) -> None: + with torch.no_grad(): + last_v_task, last_v_pref = self.policy.evaluate_dual(obs) + self.storage.compute_returns_dual( + last_v_task.detach(), + last_v_pref.detach(), + gamma=self.gamma, + lam=self.lam, + gamma_pref=self.gamma_pref, + lam_pref=self.lam_pref, + normalize_advantage=not self.normalize_advantage_per_mini_batch, + ) + + # -- update ---------------------------------------------------------------------------- + + def _surrogate(self, advantages: torch.Tensor, ratio: torch.Tensor) -> torch.Tensor: + """PPO clipped surrogate loss for one advantage stream, sharing the given ratio.""" + adv = torch.squeeze(advantages) + surrogate = -adv * ratio + surrogate_clipped = -adv * torch.clamp(ratio, 1.0 - self.clip_param, 1.0 + self.clip_param) + return torch.max(surrogate, surrogate_clipped).mean() + + def _value_loss(self, value: torch.Tensor, target: torch.Tensor, returns: torch.Tensor) -> torch.Tensor: + if not self.use_clipped_value_loss: + return (returns - value).pow(2).mean() + clipped = target + (value - target).clamp(-self.clip_param, self.clip_param) + return torch.max((value - returns).pow(2), (clipped - returns).pow(2)).mean() + + def update(self) -> dict[str, float]: + stats = { + "value_function": 0.0, + "value_function_pref": 0.0, + "surrogate": 0.0, + "surrogate_pref": 0.0, + "entropy": 0.0, + # Guard metrics -- logged in EVERY arm including beta=0, because the exploration leak + # is only visible as drift *relative to the beta=0 reference*. + "guard_entropy": 0.0, + "guard_noise_std": 0.0, + # Decomposition of realised noise into its two multiplicative factors. Layer 1 masks + # only `sigma`; a fall in the aggregate cannot say which path moved. + "guard_feat_norm": 0.0, + "guard_sigma": 0.0, + } + diag_sums: dict[str, float] = {} + + generator = self.storage.mini_batch_generator(self.num_mini_batches, self.num_learning_epochs) + actor_params = self.policy.actor_parameters() + param_groups = partition_policy_params(self.policy.named_parameters()) + grad_norm_sums: dict[str, float] = {} + + for ( + obs_batch, + actions_batch, + target_values_batch, + advantages_batch, + returns_batch, + old_actions_log_prob_batch, + old_mu_batch, + old_sigma_batch, + hidden_states_batch, + masks_batch, + target_values_pref_batch, + advantages_pref_batch, + returns_pref_batch, + ) in generator: + + if self.normalize_advantage_per_mini_batch: + with torch.no_grad(): + advantages_batch = self._norm(advantages_batch) + advantages_pref_batch = self._norm(advantages_pref_batch) + + self.policy.act(obs_batch, masks=masks_batch, hidden_state=hidden_states_batch[0]) + actions_log_prob_batch = self.policy.get_actions_log_prob(actions_batch) + value_batch, value_pref_batch = self.policy.evaluate_dual( + obs_batch, masks=masks_batch, hidden_state=hidden_states_batch[1] + ) + mu_batch = self.policy.action_mean + sigma_batch = self.policy.action_std + entropy_batch = self.policy.entropy + + self._adapt_learning_rate(mu_batch, sigma_batch, old_mu_batch, old_sigma_batch) + + # One ratio, one clip, shared by both objectives. + ratio = torch.exp(actions_log_prob_batch - torch.squeeze(old_actions_log_prob_batch)) + surrogate_task = self._surrogate(advantages_batch, ratio) + + if self.pref_detach_noise_features: + # Layer 2: same ratio *value*, different backward graph -- the preference + # objective reaches the network only through the mean action, so it cannot be + # satisfied by reshaping the trunk features that set the gSDE noise scale. + lp_pref = self.policy.log_prob_mean_path(obs_batch, actions_batch) + ratio_pref = torch.exp(lp_pref - torch.squeeze(old_actions_log_prob_batch)) + self._assert_ratio_equivalence(ratio, ratio_pref) + else: + ratio_pref = ratio + surrogate_pref = self._surrogate(advantages_pref_batch, ratio_pref) + + value_loss = self._value_loss(value_batch, target_values_batch, returns_batch) + value_loss_pref = self._value_loss(value_pref_batch, target_values_pref_batch, returns_pref_batch) + + # Critic losses + entropy bonus. The surrogates are handled separately below because + # their gradients must be projected before they reach the actor. + # + # NOTE: the entropy bonus deliberately sits HERE and not inside either surrogate. It + # is a regulariser on the optimisation, not a preference about behaviour, so it stays + # on the task side of the split -- it must keep its unrestricted gradient path to the + # noise parameters even when the preference term is masked off them. + loss_rest = ( + self.value_loss_coef * value_loss + + self.pref_value_loss_coef * value_loss_pref + - self.entropy_coef * entropy_batch.mean() + ) + + self.optimizer.zero_grad() + + if self.projection_mode == "gradient": + actor_grad, diag = self._projected_actor_grad(surrogate_task, surrogate_pref, actor_params) + loss_rest.backward() + for p, g in zip(actor_params, unflatten_to(actor_grad, actor_params)): + p.grad = g.clone() if p.grad is None else p.grad + g + elif self.projection_mode == "sum": + # ---- Arm B′: dual critic, separately normalised advantages, NO projection ---- + # This isolates the two contributions the design claims separately: scale + # invariance (from per-stream advantage normalisation) and the projection itself. + # Without B′ the Pareto plot cannot tell them apart -- if C ties B′, the + # contribution is scale-free preference control and the projection is a safety + # belt, which is a different claim from the one currently drafted. + # Diagnostics are still computed so pref_removed_frac remains comparable. + diag = self._alignment_diagnostics(surrogate_task, surrogate_pref, actor_params) + (surrogate_task + self.beta * surrogate_pref + loss_rest).backward() + else: # "advantage" + # Cheap approximation: combine the per-sample advantages *before* the surrogate. + # This is a per-sample reweighting and does NOT satisfy the first-order guarantee. + adv = torch.squeeze(advantages_batch) + self.beta * torch.squeeze(advantages_pref_batch) + surrogate_combined = self._surrogate(adv.unsqueeze(-1), ratio) + diag = self._alignment_diagnostics(surrogate_task, surrogate_pref, actor_params) + (surrogate_combined + loss_rest).backward() + + if self.is_multi_gpu: + self.reduce_parameters() + + if self.grad_clip_mode == "per_group": + gnorms = clip_grad_norm_by_group(param_groups, self.max_grad_norm) + else: # "global": upstream behaviour, kept only to reproduce pre-fix runs (NOTES 29) + total = nn.utils.clip_grad_norm_(self.policy.parameters(), self.max_grad_norm) + gnorms = {"global": float(total)} + self.optimizer.step() + for k, v in gnorms.items(): + grad_norm_sums[k] = grad_norm_sums.get(k, 0.0) + v + + stats["value_function"] += value_loss.item() + stats["value_function_pref"] += value_loss_pref.item() + stats["surrogate"] += surrogate_task.item() + stats["surrogate_pref"] += surrogate_pref.item() + stats["entropy"] += entropy_batch.mean().item() + stats["guard_entropy"] += entropy_batch.mean().item() + stats["guard_noise_std"] += self.policy.noise_magnitude().item() + fnorm, sigma = self.policy.noise_decomposition(obs_batch) + stats["guard_feat_norm"] += fnorm + stats["guard_sigma"] += sigma + for k, v in diag.items(): + diag_sums[k] = diag_sums.get(k, 0.0) + v + + num_updates = self.num_learning_epochs * self.num_mini_batches + for k in stats: + stats[k] /= num_updates + for k, v in diag_sums.items(): + stats[f"proj/{k}"] = v / num_updates + # Pre-clip gradient norm per group: makes a cross-group throttle directly visible + # instead of inferring it from the adaptive learning rate (NOTES 29). + for k, v in grad_norm_sums.items(): + stats[f"grad_norm/{k}"] = v / num_updates + stats["beta"] = self.beta + stats["pref_mask_noise"] = float(self.pref_mask_noise) + stats["pref_detach_noise_features"] = float(self.pref_detach_noise_features) + stats["grad_clip_per_group"] = float(self.grad_clip_mode == "per_group") + stats["normalize_pref_reward"] = float(self.normalize_pref_reward) + if self.normalize_pref_reward: + stats["pref_reward_scale"] = float(self.policy.pref_reward_normalizer.emp_norm._std) + + # Per-stream reward means. The runner's reward bookkeeping only tracks the scalar it gets + # from env.step (the task stream), so log the preference stream here -- attribution + # between the two curves is the point of keeping them separate (design notes §3.4). + stats["reward_task_mean"] = self.storage.rewards.mean().item() + stats["reward_pref_mean"] = self.storage.rewards_pref.mean().item() + + self.storage.clear() + self._diag = stats + return stats + + # -- helpers --------------------------------------------------------------------------- + + def _preference_subspace(self, actor_params: list) -> torch.Tensor | None: + """Flat mask of the coordinates the preference term may move (None = all of them).""" + if not self.pref_mask_noise: + return None + if self._pref_allowed is None: + noise_flags = self.policy.actor_param_noise_mask() + self._pref_allowed = ~flat_mask(actor_params, noise_flags) + return self._pref_allowed + + def _alignment_diagnostics( + self, surrogate_task: torch.Tensor, surrogate_pref: torch.Tensor, actor_params: list + ) -> dict[str, float]: + """Measure task/preference gradient alignment **without applying** the preference. + + Logged in every arm -- including β=0 and the unprojected ablations -- because + ``pref_removed_frac`` is the screening instrument for whether a preference genuinely + conflicts with the task, and it must be a *time series*: alignment moves over training, + and the late-training regime (g_task → 0, projector → identity) is where the + self-scheduling argument lives. A single converged number is not evidence for that. + + Costs one extra backward (~1% of iteration wall-clock, since collection dominates ~99%). + The gradient actually applied is unaffected. + """ + if not self.log_pref_alignment: + return {} + g_task = flatten_grads( + torch.autograd.grad(surrogate_task, actor_params, retain_graph=True, allow_unused=True), + actor_params, + ) + g_pref = flatten_grads( + torch.autograd.grad(surrogate_pref, actor_params, retain_graph=True, allow_unused=True), + actor_params, + ) + allowed = self._preference_subspace(actor_params) + if allowed is not None: + g_task, g_pref = g_task[allowed], g_pref[allowed] + # beta=1 purely to read out the geometry; nothing is applied. + _, diag = project_nullspace(g_task, g_pref, beta=1.0) + return diag + + def _projected_actor_grad( + self, surrogate_task: torch.Tensor, surrogate_pref: torch.Tensor, actor_params: list + ) -> tuple[torch.Tensor, dict[str, float]]: + """Flat actor gradient with the preference component projected orthogonally to the task.""" + g_task = flatten_grads( + torch.autograd.grad(surrogate_task, actor_params, retain_graph=True, allow_unused=True), + actor_params, + ) + if self.beta == 0.0: + # At β=0 the projection is the identity, so the APPLIED gradient is exactly baseline + # PPO plus an untouched second critic head. The alignment read-out is still taken + # (extra backward, nothing applied) so the β=0 reference has the same time series as + # every other arm -- otherwise there is nothing to compare the others against. + diag = self._alignment_diagnostics(surrogate_task, surrogate_pref, actor_params) + diag.setdefault("g_task_norm", torch.linalg.vector_norm(g_task).item()) + diag["pref_contrib_ratio"] = 0.0 + return g_task, diag + + g_pref = flatten_grads( + torch.autograd.grad(surrogate_pref, actor_params, retain_graph=True, allow_unused=True), + actor_params, + ) + + allowed = self._preference_subspace(actor_params) + if allowed is None: + # Unmasked: the known-degenerate configuration. Kept runnable on purpose so the + # exploration collapse can be demonstrated once rather than argued about. + return project_nullspace(g_task, g_pref, self.beta) + + combined, diag = project_nullspace_masked(g_task, g_pref, self.beta, allowed) + self._assert_noise_params_untouched(combined, g_task, allowed) + return combined, diag + + def _assert_noise_params_untouched( + self, combined: torch.Tensor, g_task: torch.Tensor, allowed: torch.Tensor + ) -> None: + """Layer 1's claim, checked directly: preference contributes EXACTLY zero to noise params. + + Under Layer 1 masking, ``sigma`` is excluded from ``g_pref`` by construction, so nothing + in the preference path can move it. If it moves, that is an implementation bug -- not a + leak -- and the two must never be confused when reading the guard metrics. Asserting it + here is stronger than inferring flatness from a cross-run plot. + """ + if getattr(self, "_mask_checked", False): + return + self._mask_checked = True + excluded = ~allowed + if excluded.any() and not torch.equal(combined[excluded], g_task[excluded]): + d = (combined[excluded] - g_task[excluded]).abs().max().item() + raise RuntimeError( + f"Layer 1 masking violated: preference moved the noise parameters " + f"(max |diff| {d:.3e}). The projection must be *restricted* to the mean-action " + "subspace, not applied to a zero-padded g_pref." + ) + n_excluded = int(excluded.sum()) + print( + f"[NullspacePPO] Layer 1 active: {n_excluded} noise params " + f"({100 * n_excluded / excluded.numel():.2f}% of actor) excluded from g_pref; verified untouched." + ) + + def _assert_ratio_equivalence(self, ratio: torch.Tensor, ratio_pref: torch.Tensor) -> None: + """Check once that Layer 2 changed only the gradient path, not the objective's value. + + If these ever differ numerically we are no longer running "one ratio, one clip" -- we are + running two different PPO objectives, and the clipping would disagree between streams. + """ + if getattr(self, "_ratio_checked", False): + return + self._ratio_checked = True + if not torch.allclose(ratio.detach(), ratio_pref.detach(), rtol=1e-4, atol=1e-6): + d = (ratio - ratio_pref).abs().max().item() + raise RuntimeError( + f"Layer-2 preference ratio diverged from the task ratio (max |diff| {d:.3e}). " + "log_prob_mean_path must reproduce the policy's log-prob exactly and differ only " + "in its backward graph." + ) + print("[NullspacePPO] Layer 2 active: preference ratio value-identical, mean-path gradient only.") + + def _adapt_learning_rate(self, mu, sigma, old_mu, old_sigma) -> None: # noqa: ANN001 + """Unchanged from upstream: KL-adaptive LR on the *policy* distribution only.""" + if self.desired_kl is None or self.schedule != "adaptive": + return + with torch.inference_mode(): + kl = torch.sum( + torch.log(sigma / old_sigma + 1.0e-5) + + (torch.square(old_sigma) + torch.square(old_mu - mu)) / (2.0 * torch.square(sigma)) + - 0.5, + axis=-1, + ) + kl_mean = torch.mean(kl) + if self.is_multi_gpu: + torch.distributed.all_reduce(kl_mean, op=torch.distributed.ReduceOp.SUM) + kl_mean /= self.gpu_world_size + if self.gpu_global_rank == 0: + if kl_mean > self.desired_kl * 2.0: + self.learning_rate = max(1e-5, self.learning_rate / 1.5) + elif self.desired_kl / 2.0 > kl_mean > 0.0: + self.learning_rate = min(1e-2, self.learning_rate * 1.5) + if self.is_multi_gpu: + lr_tensor = torch.tensor(self.learning_rate, device=self.device) + torch.distributed.broadcast(lr_tensor, src=0) + self.learning_rate = lr_tensor.item() + for param_group in self.optimizer.param_groups: + param_group["lr"] = self.learning_rate + + @staticmethod + def _norm(adv: torch.Tensor) -> torch.Tensor: + std = adv.std() + if not torch.isfinite(std) or std < 1e-8: + return torch.zeros_like(adv) + return (adv - adv.mean()) / (std + 1e-8) diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/projection.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/projection.py new file mode 100644 index 00000000..1ff73a7a --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/projection.py @@ -0,0 +1,193 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Null-space projection of a preference gradient against a task gradient. + +The method in one line: move along the task gradient as usual, and add only the component of the +preference gradient that is orthogonal to it, so the preference term produces no *first-order* +change in task value. + + ĝ_task = g_task / ‖g_task‖ + g = g_task + β · (I − ĝ_task ĝ_taskᵀ) · g_pref + +Everything here operates on flat vectors so it is independent of the parameter layout, and is +pure/side-effect-free so it can be unit-tested without a simulator. +""" + +from __future__ import annotations + +import torch + +# Floor on ‖g_task‖ used when normalizing. The projector degenerates as g_task → 0 (exactly where +# the method matters most, once task success has converged), so this is a real numerical guard, +# not decoration -- see the "known costs" section of the design notes. +DEFAULT_EPS = 1e-12 + + +def flatten_grads(grads: list[torch.Tensor | None], params: list[torch.Tensor]) -> torch.Tensor: + """Flatten a list of per-parameter grads into one vector, treating None as zeros. + + ``torch.autograd.grad(..., allow_unused=True)`` returns None for parameters that did not + participate in the graph (e.g. gSDE's ``log_std`` when the surrogate does not touch it). + """ + return torch.cat([ + (g if g is not None else torch.zeros_like(p)).reshape(-1) for g, p in zip(grads, params) + ]) + + +def unflatten_to(vec: torch.Tensor, params: list[torch.Tensor]) -> list[torch.Tensor]: + """Split a flat vector back into tensors shaped like ``params``.""" + out, i = [], 0 + for p in params: + n = p.numel() + out.append(vec[i : i + n].view_as(p)) + i += n + return out + + +def flat_mask(params: list[torch.Tensor], flags: list[bool]) -> torch.Tensor: + """Boolean mask over the flattened parameter vector, True where ``flags`` is True.""" + return torch.cat([ + torch.full((p.numel(),), bool(f), dtype=torch.bool, device=p.device) + for p, f in zip(params, flags) + ]) + + +def project_nullspace_masked( + g_task: torch.Tensor, + g_pref: torch.Tensor, + beta: float, + allowed: torch.Tensor, + eps: float = DEFAULT_EPS, +) -> tuple[torch.Tensor, dict[str, float]]: + """Project within a *subspace* of the parameters, leaving the rest untouched by preference. + + ``allowed`` is a boolean mask selecting the coordinates the preference term is permitted to + move (the mean-action parameters). The projection is computed against ``g_task`` **restricted + to that same subspace**, and the preference contribution is zero everywhere else. + + Why restrict rather than zero ``g_pref`` and project in the full space: the projector subtracts + a multiple of ``ĝ_task``, and ``g_task`` has non-zero components on the excluded coordinates. + Projecting in the full space would therefore feed a correction term back onto exactly the + parameters we are trying to protect. Restricting the whole operation to the subspace is what + actually guarantees the preference never moves them. + + The scoping is a deliberate weakening of the guarantee, and belongs in the method section: + no first-order change in task value *as mediated by the deployed (mean-action) parameters*. + """ + g_sub_task, g_sub_pref = g_task[allowed], g_pref[allowed] + combined_sub, diag = project_nullspace(g_sub_task, g_sub_pref, beta, eps) + + g = g_task.clone() + g[allowed] = combined_sub + diag["masked_frac"] = 1.0 - (allowed.sum().item() / max(allowed.numel(), 1)) + return g, diag + + +def project_nullspace( + g_task: torch.Tensor, + g_pref: torch.Tensor, + beta: float, + eps: float = DEFAULT_EPS, +) -> tuple[torch.Tensor, dict[str, float]]: + """Combine task and preference gradients with the preference term projected orthogonally. + + Args: + g_task: flat task gradient. + g_pref: flat preference gradient. + beta: preference step budget. ``beta == 0`` returns ``g_task`` exactly (bitwise), which is + what makes the β=0 sanity run a true no-op. + eps: floor on ‖g_task‖. + + Returns: + (combined gradient, diagnostics). Diagnostics are cheap scalars worth logging every + update: they are how you tell "preference is being squeezed out" from "preference is + driving the update", and they are the early-warning signal for reward hacking. + """ + task_norm = torch.linalg.vector_norm(g_task) + pref_norm = torch.linalg.vector_norm(g_pref) + + # cosine similarity before projection: +1 means preference already agrees with the task + # (projection removes almost everything), 0 means it is already orthogonal (projection is a + # no-op and preference is "free"), -1 means it directly opposes success. + denom = task_norm * pref_norm + cos_before = (torch.dot(g_task, g_pref) / denom).item() if denom > eps else 0.0 + + if beta == 0.0: + return g_task, { + "g_task_norm": task_norm.item(), + "g_pref_norm": pref_norm.item(), + "cos_before": cos_before, + "pref_removed_frac": 0.0, + "pref_contrib_ratio": 0.0, + } + + # Normalize the task direction, flooring the denominator. + g_hat = g_task / task_norm.clamp_min(eps) + # Remove the component of g_pref along g_hat. + g_pref_perp = g_pref - torch.dot(g_hat, g_pref) * g_hat + + perp_norm = torch.linalg.vector_norm(g_pref_perp) + # Fraction of the preference gradient the projection discarded. → 1 means preference is + # (anti)parallel to the task and the null space affords it nothing. + removed = (1.0 - (perp_norm / pref_norm).item()) if pref_norm > eps else 0.0 + + g = g_task + beta * g_pref_perp + + return g, { + "g_task_norm": task_norm.item(), + "g_pref_norm": pref_norm.item(), + "cos_before": cos_before, + "pref_removed_frac": removed, + # How much of the final step the preference term accounts for. This is the quantity that + # self-schedules: small early (task gradient dominates), growing as g_task → 0. + "pref_contrib_ratio": (beta * perp_norm / task_norm.clamp_min(eps)).item(), + } + + +# -- gradient-norm clipping scope (NOTES 29) ------------------------------------------------- + + +def partition_policy_params(named_params) -> dict[str, list[torch.nn.Parameter]]: # noqa: ANN001 + """Split a dual actor-critic's parameters into the groups that must not share one + gradient-norm budget. + + ``actor`` the actor MLP plus exploration-noise parameters (``std`` / ``log_std``); + ``critic`` the task value head -- or, with ``critic_arch='shared'``, the 2-head trunk; + ``critic_pref`` the separate preference critic (empty with ``critic_arch='shared'``). + + Name-prefix based, consistent with ``DualCriticActorCritic.actor_parameters`` (everything not + starting with ``critic`` is actor). ``critic_pref`` is tested first because ``critic`` is a + prefix of it. + """ + groups: dict[str, list[torch.nn.Parameter]] = {"actor": [], "critic": [], "critic_pref": []} + for name, param in named_params: + if name.startswith("critic_pref"): + groups["critic_pref"].append(param) + elif name.startswith("critic"): + groups["critic"].append(param) + else: + groups["actor"].append(param) + return groups + + +def clip_grad_norm_by_group( + groups: dict[str, list[torch.nn.Parameter]], max_norm: float +) -> dict[str, float]: + """``clip_grad_norm_`` applied to each group on its own budget. + + Upstream rsl_rl clips the whole policy as one vector. With two critics that couples them to + the actor: a large preference value loss enters the same norm and scales the actor's update + down, at beta=0 too. Clipping per group removes that channel. + + Returns each group's pre-clip norm (0.0 for a group with no gradients) -- the diagnostic that + makes a cross-group throttle directly observable. + """ + norms: dict[str, float] = {} + for name, params in groups.items(): + with_grad = [p for p in params if p.grad is not None] + if not with_grad: + norms[name] = 0.0 + continue + norms[name] = float(torch.nn.utils.clip_grad_norm_(with_grad, max_norm)) + return norms diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/reward_split.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/reward_split.py new file mode 100644 index 00000000..07e5e937 --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/reward_split.py @@ -0,0 +1,195 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Split the environment's scalar reward into a task stream and a preference stream. + +Design constraint discovered in Phase 0 recon: OmniReset's ``RewardsCfg`` contains +``progress_context``, a term that returns zeros but caches the goal distances that ``r_dist``, +``r_success``, three termination terms, the reset-distribution success monitor and the parent +project's data-collection configs all read back off it. ``RewardManager`` also *skips any term +whose weight is 0.0* without calling it. So splitting the manager into two managers, reordering +terms, or zeroing weights are all ways to silently corrupt the reward. + +We therefore never touch the RewardManager. It already materialises per-term rewards in +``_step_reward`` of shape (num_envs, num_terms); we read that buffer and mask it. + +Note on scaling: ``reward_manager.compute`` sets ``_step_reward[:, i] = func_i * weight_i`` and +adds ``func_i * weight_i * dt`` to the reward buffer. So ``_step_reward`` has the weight *already +applied* and ``dt`` divided out -- the contribution of term i to the returned reward is +``_step_reward[:, i] * dt``, with no second weight multiply. +""" + +from __future__ import annotations + +import torch +from abc import ABC, abstractmethod + +from isaaclab_rl.rsl_rl import RslRlVecEnvWrapper + + +class PreferenceRewardSource(ABC): + """Where the preference reward stream comes from. + + Phase 1 uses Zero (sanity run A) and GaussianNoise (sanity run B); Phase 2 uses + RewardManagerTerms with scripted manner predicates; Phase 3 swaps in a distilled critic. + """ + + #: True when this source draws from terms that are *already inside* the RewardManager, and + #: whose contribution must therefore be subtracted out of the task stream to avoid + #: double-counting. False for exogenous sources (noise, a learned model). + subtracts_from_task: bool = False + + def initialize(self, env) -> None: # noqa: ANN001 + """Optional hook once the env exists (resolve term names to indices, etc.).""" + + @abstractmethod + def compute(self, env, total_reward: torch.Tensor) -> torch.Tensor: # noqa: ANN001 + """Return the per-env preference reward for this step, shape (num_envs,).""" + + +class ZeroPreference(PreferenceRewardSource): + """No preference signal. Sanity run A: β=0 must reproduce the baseline.""" + + def compute(self, env, total_reward: torch.Tensor) -> torch.Tensor: # noqa: ANN001 + return torch.zeros_like(total_reward) + + +class GaussianNoisePreference(PreferenceRewardSource): + """Zero-mean Gaussian preference reward. Sanity run B. + + Tests that the projection *neutralises an uninformative preference* rather than merely + injecting variance: with β=1 the curves must still overlay the baseline. + """ + + def __init__(self, std: float = 1.0, seed: int | None = None) -> None: + self.std = std + self.seed = seed + self._gen: torch.Generator | None = None + + def initialize(self, env) -> None: # noqa: ANN001 + if self.seed is not None: + self._gen = torch.Generator(device=env.unwrapped.device).manual_seed(self.seed) + + def compute(self, env, total_reward: torch.Tensor) -> torch.Tensor: # noqa: ANN001 + return torch.normal( + mean=0.0, + std=self.std, + size=total_reward.shape, + device=total_reward.device, + generator=self._gen, + ) + + +class ActionRatePreference(PreferenceRewardSource): + """The **noise-bait probe**: preference = negative action-rate norm. + + This is literally ``r_smooth``'s second term (``action_rate_l2_clamped``), chosen because it is + *maximally* satisfiable by shrinking exploration noise and barely satisfiable any other way. + + It exists to test the scoping of the null-space constraint, and it is strictly sharper than the + Gaussian sanity run: zero-mean noise produces an *unsystematic* preference gradient that will + not preferentially shrink exploration, so that run passes even with the leak wide open. + + Expected result when correctly scoped: **almost no compliance gain**, because both routes are + closed -- the noise parameters are masked out of the preference gradient, and the mean policy + is held by the task constraint. Large apparent compliance means the leak is still open; + escalate to Layer 2 (``pref_detach_noise_features``). + + Exogenous (not subtracted from the task stream): the task's own ``action_rate`` term stays + where it is, so the arms remain comparable to the baseline on task reward. + """ + + def compute(self, env, total_reward: torch.Tensor) -> torch.Tensor: # noqa: ANN001 + am = env.unwrapped.action_manager + return -torch.clamp(torch.sum(torch.square(am.action - am.prev_action), dim=1), 0, 1e4) + + +class EndEffectorHeightPreference(PreferenceRewardSource): + """**High-conflict** preference: keep the end-effector low. + + Deliberately fights the task. Cupcake-on-plate requires lifting the object and placing it on + a raised receptacle, so "stay low" cannot be satisfied without giving up task success. + + Why the Phase 2 set needs at least one of these: the action-rate bait measured + ``pref_removed_frac ~ 1e-3`` -- i.e. ~99.9% of the preference gradient survives projection. + In a parameter space of a few hundred thousand dimensions near-orthogonality is the default, + so for such preferences the projected arm is numerically almost identical to the unprojected + one, and the Pareto plot cannot separate them **by construction**. The projection can only + demonstrate its value on preferences that genuinely conflict. + + Use ``pref_removed_frac`` to screen candidates before committing the Phase 2 matrix, and pick + a set that spans a range rather than clustering near zero. + """ + + def __init__(self, body_name: str = "wrist_3_link") -> None: + self.body_name = body_name + self._body_id: int | None = None + + def initialize(self, env) -> None: # noqa: ANN001 + robot = env.unwrapped.scene["robot"] + ids, names = robot.find_bodies(self.body_name) + if not ids: + raise ValueError(f"body {self.body_name!r} not found on the robot") + self._body_id = ids[0] + print(f"[pref] end-effector height preference on body {names[0]!r}") + + def compute(self, env, total_reward: torch.Tensor) -> torch.Tensor: # noqa: ANN001 + robot = env.unwrapped.scene["robot"] + z_w = robot.data.body_pos_w[:, self._body_id, 2] + # Height above this env's own origin, so the signal is identical across the env grid. + return -(z_w - env.unwrapped.scene.env_origins[:, 2]) + + +class RewardManagerTermsPreference(PreferenceRewardSource): + """Preference reward = the sum of named RewardManager terms (Phase 2 scripted predicates). + + The named terms stay in the manager (so evaluation order and the ``progress_context`` + dependency chain are untouched) but their contribution is moved out of the task stream. + """ + + subtracts_from_task = True + + def __init__(self, term_names: list[str]) -> None: + self.term_names = list(term_names) + self._idx: torch.Tensor | None = None + + def initialize(self, env) -> None: # noqa: ANN001 + rm = env.unwrapped.reward_manager + known = list(rm.active_terms) + missing = [n for n in self.term_names if n not in known] + if missing: + raise ValueError( + f"Preference terms {missing} are not active RewardManager terms. Active: {known}" + ) + self._idx = torch.tensor( + [known.index(n) for n in self.term_names], dtype=torch.long, device=env.unwrapped.device + ) + + def compute(self, env, total_reward: torch.Tensor) -> torch.Tensor: # noqa: ANN001 + rm = env.unwrapped.reward_manager + # _step_reward already has per-term weights applied and dt divided out. + return rm._step_reward[:, self._idx].sum(dim=-1) * env.unwrapped.step_dt + + +class DualRewardVecEnvWrapper(RslRlVecEnvWrapper): + """Emits the task reward as the usual scalar and the preference reward via ``extras``. + + The preference stream rides in ``extras["reward_pref"]`` so that the rest of the rsl_rl + plumbing (runner loop signature, env interface) is unchanged; only our PPO subclass reads it. + """ + + PREF_KEY = "reward_pref" + + def __init__(self, env, pref_source: PreferenceRewardSource | None = None, **kwargs) -> None: # noqa: ANN001 + super().__init__(env, **kwargs) + self.pref_source = pref_source or ZeroPreference() + self.pref_source.initialize(self) + + def step(self, actions: torch.Tensor): # noqa: ANN201 + obs, rew, dones, extras = super().step(actions) + + r_pref = self.pref_source.compute(self, rew) + r_task = rew - r_pref if self.pref_source.subtracts_from_task else rew + + extras[self.PREF_KEY] = r_pref + return obs, r_task, dones, extras diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/runner.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/runner.py new file mode 100644 index 00000000..bfbe85ed --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace/runner.py @@ -0,0 +1,61 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""On-policy runner that can construct the dual-critic policy and null-space PPO. + +Upstream's ``_construct_algorithm`` resolves class names with a bare ``eval()`` evaluated in +*its own module namespace*, so classes defined outside ``rsl_rl.runners.on_policy_runner`` are +unreachable by name. This subclass overrides only that resolution step; everything else -- +the rollout loop, logging, checkpointing -- is inherited unchanged. +""" + +from __future__ import annotations + +from tensordict import TensorDict + +from rsl_rl.runners import OnPolicyRunner +from rsl_rl.utils import resolve_obs_groups # noqa: F401 (kept for parity with upstream imports) + +from .dual_actor_critic import DualCriticActorCritic +from .nullspace_ppo import NullspacePPO + +#: Class names this runner can resolve, in addition to whatever upstream supports. +REGISTRY = { + "DualCriticActorCritic": DualCriticActorCritic, + "NullspacePPO": NullspacePPO, +} + + +class DualCriticOnPolicyRunner(OnPolicyRunner): + """OnPolicyRunner that resolves the null-space classes by name.""" + + def _construct_algorithm(self, obs: TensorDict) -> NullspacePPO: + policy_cfg = dict(self.policy_cfg) + alg_cfg = dict(self.alg_cfg) + + policy_class_name = policy_cfg.pop("class_name", "DualCriticActorCritic") + alg_class_name = alg_cfg.pop("class_name", "NullspacePPO") + + unknown = [n for n in (policy_class_name, alg_class_name) if n not in REGISTRY] + if unknown: + raise ValueError( + f"{unknown} not resolvable by DualCriticOnPolicyRunner. Known: {sorted(REGISTRY)}. " + "Use the stock OnPolicyRunner for upstream classes." + ) + + # RND is rejected by NullspacePPO; surface that here rather than deep in the algorithm. + if alg_cfg.get("rnd_cfg") is not None: + raise NotImplementedError("RND is not supported alongside the dual critic.") + alg_cfg.pop("rnd_cfg", None) + alg_cfg.pop("symmetry_cfg", None) + + policy = REGISTRY[policy_class_name]( + obs, self.cfg["obs_groups"], self.env.num_actions, **policy_cfg + ).to(self.device) + + alg = REGISTRY[alg_class_name]( + policy, device=self.device, **alg_cfg, multi_gpu_cfg=self.multi_gpu_cfg + ) + + alg.init_storage("rl", self.env.num_envs, self.num_steps_per_env, obs, [self.env.num_actions]) + return alg diff --git a/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace_cfg.py b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace_cfg.py new file mode 100644 index 00000000..8754e31f --- /dev/null +++ b/source/uwlab_rl/uwlab_rl/rsl_rl/nullspace_cfg.py @@ -0,0 +1,125 @@ +# Copyright (c) 2026, Null-space preference critic project. +# SPDX-License-Identifier: BSD-3-Clause + +"""Configuration for the dual-critic / null-space-projection PPO agent.""" + +from __future__ import annotations + +from isaaclab.utils import configclass +from isaaclab_rl.rsl_rl import RslRlOnPolicyRunnerCfg, RslRlPpoAlgorithmCfg + +from .rl_cfg import RslRlFancyActorCriticCfg + + +@configclass +class RslRlNullspaceActorCriticCfg(RslRlFancyActorCriticCfg): + """Actor-critic with a second value head for the preference return.""" + + class_name: str = "DualCriticActorCritic" + + critic_arch: str = "shared" + """``shared``: one critic trunk widened to 2 outputs (cheap; the preference value loss + backpropagates through the shared trunk, so the heads are coupled even at beta=0). + ``separate``: an independent preference critic MLP -- no coupling through shared weights. + That alone does NOT guarantee the beta=0 overlay: until ``grad_clip_mode='per_group'`` the + actor and both critics shared one gradient-norm clip, which coupled them regardless of + architecture (NOTES 29). Run sanity check A against both.""" + + +@configclass +class RslRlNullspacePpoAlgorithmCfg(RslRlPpoAlgorithmCfg): + """PPO whose preference gradient is projected into the null space of the task gradient.""" + + class_name: str = "NullspacePPO" + + beta: float = 0.0 + """Preference step budget. 0.0 applies no preference gradient to the actor. + + It reproduced baseline PPO only with ``pref_source='zero'``. Under a non-zero preference the + preference critic's gradient reached the actor through the shared gradient-norm clip, so + beta=0 was not a no-op (NOTES 29). With ``grad_clip_mode='per_group'`` it is, for any + ``pref_source``.""" + + gamma_pref: float | None = None + """Discount for the preference return. None mirrors ``gamma``. Manner is mostly local while + success is long-horizon, so gamma_pref << gamma is expected to help -- kept as an explicit + ablation rather than a hidden default.""" + + lam_pref: float | None = None + """GAE lambda for the preference stream. None mirrors ``lam``.""" + + projection_mode: str = "gradient" + """``gradient``: project g_pref orthogonally to g_task (faithful; costs one extra backward, + measured at ~1% of iteration wall-clock since collection dominates). + ``advantage``/``sum``: combine the scalar objectives instead -- the weighted-sum baseline and + the cheap approximation that drops the first-order guarantee. Ablations only.""" + + pref_value_loss_coef: float | None = None + """Value-loss coefficient for the preference critic. None mirrors ``value_loss_coef``.""" + + pref_mask_noise: bool = True + """Layer 1. Exclude the exploration-noise parameters (gSDE ``log_std``) from the preference + gradient, and run the projection inside the remaining (mean-action) subspace. + + Default ON. Every scripted manner preference is monotonically improved by shrinking action + noise, while noise perturbations are ~second order in task return near a local optimum -- so + an unmasked projector actively *selects* exploration collapse as a "task-neutral" direction. + gSDE noise does not exist at deployment (the policy is deployed on the mean action), so + masking gives up no legitimate preference. Set False only for the deliberate demonstration + run of the degenerate mode.""" + + pref_detach_noise_features: bool = False + """Layer 2 escalation. Also detach the trunk features that feed the gSDE noise head inside + the preference surrogate. + + Layer 1 alone does NOT close the leak: the noise scale is + ``mm(actor[:-1](obs)**2, exp(log_std)**2)``, so the preference can still shrink exploration by + reshaping the trunk. Enable if guard entropy / noise magnitude drift against the beta=0 + reference after Layer 1.""" + + grad_clip_mode: str = "per_group" + """How ``max_grad_norm`` is applied. ``per_group`` clips the actor, the task critic and the + preference critic each on its own budget; ``global`` clips them as one vector (upstream). + + ``global`` is a bug for a dual-critic agent and is kept only to reproduce pre-fix runs + (NOTES 29). The preference critic's gradient enters the same norm as the actor's, so a large + preference value loss scales the actor's update down -- at beta=0 as well. With + ``pref_source='action_rate'`` that loss reached 1e5-1e6 against ``max_grad_norm=1.0``: the + actor was throttled, the adaptive schedule pinned the learning rate at its 1e-2 cap, and + because the preference scales with action noise the throttle tracked exploration.""" + + normalize_pref_reward: bool = True + """Scale the preference reward by a running std of its discounted sum (Pathak et al.; rsl_rl's + ``EmpiricalDiscountedVariationNormalization``) before it reaches the preference critic. + + Per-stream advantage normalisation already makes the *actor's* preference signal scale-free, + but the preference critic regressed onto raw returns, so its loss carried the reward's + arbitrary scale. A positive running factor leaves the normalised preference advantages + unchanged, so this fixes the critic without altering what the actor sees. The task stream is + never normalised: it must stay identical to upstream for the beta=0 null to mean anything.""" + + +@configclass +class RslRlNullspaceRunnerCfg(RslRlOnPolicyRunnerCfg): + """Runner that can resolve the null-space classes by name.""" + + class_name: str = "DualCriticOnPolicyRunner" + + pref_source: str = "zero" + """Preference reward stream: + ``zero`` sanity run A; + ``noise`` sanity run B -- an *uninformative* preference the projection must not + amplify into variance; + ``action_rate`` the noise-bait probe -- a preference that is maximally satisfiable by + shrinking exploration and barely satisfiable any other way. Correctly + scoped it should yield almost NO compliance gain; large apparent gain + means the exploration leak is still open; + ``terms`` Phase 2 scripted manner predicates read out of the RewardManager.""" + + pref_noise_std: float = 1.0 + """Std of the Gaussian preference reward when ``pref_source='noise'``.""" + + pref_term_names: tuple[str, ...] = () + """RewardManager term names forming the preference stream when ``pref_source='terms'``. + These stay in the manager (evaluation order untouched) but are subtracted from the task + stream so they are not double-counted.""" diff --git a/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/__init__.py b/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/__init__.py index 1e240e13..25819fe7 100644 --- a/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/__init__.py +++ b/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/__init__.py @@ -88,6 +88,18 @@ }, ) +gym.register( + id="OmniReset-Ur5eRobotiq2f85-RelCartesianOSC-State-Nullspace-v0", + entry_point="isaaclab.envs:ManagerBasedRLEnv", + disable_env_checker=True, + kwargs={ + # Same env as the -State-v0 baseline; only the agent differs, so a beta=0 run is + # comparable to the baseline line-for-line. + "env_cfg_entry_point": f"{__name__}.rl_state_cfg:Ur5eRobotiq2f85RelCartesianOSCTrainCfg", + "rsl_rl_cfg_entry_point": f"{agents.__name__}.rsl_rl_cfg:Nullspace_PPORunnerCfg", + }, +) + gym.register( id="OmniReset-Ur5eRobotiq2f85-RelCartesianOSC-State-Finetune-v0", entry_point="isaaclab.envs:ManagerBasedRLEnv", diff --git a/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/agents/rsl_rl_cfg.py b/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/agents/rsl_rl_cfg.py index 8f6db98e..613ad0b3 100644 --- a/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/agents/rsl_rl_cfg.py +++ b/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/config/ur5e_robotiq_2f85/agents/rsl_rl_cfg.py @@ -6,6 +6,10 @@ from isaaclab.utils import configclass from isaaclab_rl.rsl_rl import RslRlOnPolicyRunnerCfg, RslRlPpoAlgorithmCfg +from uwlab_rl.rsl_rl.nullspace_cfg import ( + RslRlNullspaceActorCriticCfg, + RslRlNullspacePpoAlgorithmCfg, +) from uwlab_rl.rsl_rl.rl_cfg import ( BehaviorCloningCfg, OffPolicyAlgorithmCfg, @@ -81,3 +85,46 @@ class Base_DAggerRunnerCfg(Base_PPORunnerCfg): ) ), ) + + +@configclass +class Nullspace_PPORunnerCfg(Base_PPORunnerCfg): + """Dual-critic PPO with the preference gradient projected into the task gradient's null space. + + Inherits every PPO hyperparameter from Base_PPORunnerCfg so that a beta=0 run differs from the + baseline only by the plumbing under test. The preference stream itself is selected at runtime + (--pref_source / --beta on train.py), which is what lets sanity runs A and B share one config. + """ + + experiment_name = "ur5e_robotiq_2f85_omnireset_nullspace" + class_name = "DualCriticOnPolicyRunner" + + policy = RslRlNullspaceActorCriticCfg( + init_noise_std=1.0, + actor_obs_normalization=True, + critic_obs_normalization=True, + actor_hidden_dims=[512, 256, 128, 64], + critic_hidden_dims=[512, 256, 128, 64], + activation="elu", + noise_std_type="gsde", + state_dependent_std=False, + critic_arch="shared", + ) + + algorithm = RslRlNullspacePpoAlgorithmCfg( + value_loss_coef=1.0, + use_clipped_value_loss=True, + normalize_advantage_per_mini_batch=False, + clip_param=0.2, + entropy_coef=0.006, + num_learning_epochs=5, + num_mini_batches=4, + learning_rate=1.0e-4, + schedule="adaptive", + gamma=0.99, + lam=0.95, + desired_kl=0.01, + max_grad_norm=1.0, + beta=0.0, + projection_mode="gradient", + ) diff --git a/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/mdp/utils.py b/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/mdp/utils.py index 86fe3f40..807dd675 100644 --- a/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/mdp/utils.py +++ b/source/uwlab_tasks/uwlab_tasks/manager_based/manipulation/omnireset/mdp/utils.py @@ -326,7 +326,10 @@ def get_temp_dir(rank: int | None = None) -> str: uid = os.getuid() job_id = os.getenv("SLURM_JOB_ID") or os.getenv("PBS_JOBID") or "local" - download_dir = os.path.join("/tmp", "uwlab", str(uid), str(job_id), f"rank_{rank}") + # Base is overridable via UWLAB_TMP_DIR: on shared boxes /tmp/uwlab may be owned by another + # user (unwritable), so point this at a writable workspace dir instead. + base = os.getenv("UWLAB_TMP_DIR") or os.path.join("/tmp", "uwlab") + download_dir = os.path.join(base, str(uid), str(job_id), f"rank_{rank}") os.makedirs(download_dir, mode=0o700, exist_ok=True) return download_dir