# Copyright (c) 2022-2026, The Isaac Lab Project Developers (https://github.com/isaac-sim/IsaacLab/blob/main/CONTRIBUTORS.md).
# All rights reserved.
#
# SPDX-License-Identifier: BSD-3-Clause
"""Replication queue, :func:`replicate` drain, and :class:`ReplicateSession` sugar."""
from __future__ import annotations
import importlib
from collections.abc import Callable, Iterable
from typing import TYPE_CHECKING, Any
from isaaclab.utils.backend_utils import FactoryBase
from isaaclab.utils.string import string_to_callable
from isaaclab.utils.version import has_kit
from .clone_plan import make_clone_plan
from .cloner_cfg import DEFAULT_ENV_TEMPLATE
from .cloner_strategies import sequential
from .usd import UsdReplicateContext
if TYPE_CHECKING:
import torch
from pxr import Usd
from .clone_plan import ClonePlan
REPLICATION_QUEUE: list[Any] = []
"""Asset cfgs registered by :func:`queue_replication` and drained by :func:`replicate`.
The queue only records *which* cfgs participate in cloning; how each cfg is cloned is
resolved at dispatch from :attr:`~isaaclab.assets.AssetBaseCfg.cloning_contexts` or the
active backend's default stack.
"""
def queue_replication(cfg: Any) -> None:
"""Register ``cfg`` for cloning when :func:`replicate` next runs.
Args:
cfg: Asset cfg with resolved ``prim_path``.
"""
REPLICATION_QUEUE.append(cfg)
def replicate(plan: ClonePlan, *, stage: Usd.Stage, replicate_physics: bool = True) -> None:
"""Drain :data:`REPLICATION_QUEUE` against ``plan``, dispatch each backend, publish the plan.
Physics contexts come from :attr:`~isaaclab.assets.AssetBaseCfg.cloning_contexts` when
set, otherwise from the backend's ``PHYSICS_CONTEXT`` class.
:class:`~isaaclab.cloner.UsdReplicateContext` is added automatically when the cfg has a
spawner and Kit is available. Explicit contexts are honored regardless of Kit availability.
With ``replicate_physics=False`` physics contexts are dropped; USD replication still fires
when the spawner+Kit condition is met or the cfg explicitly requests it.
Cfgs absent from ``plan.cfg_rows`` are silently skipped. Backend contexts run in
ascending ``replicate_priority`` order. The queue is cleared up front, so a backend
failure cannot leak stale entries into the next call. Every context receives the plan's
explicitly declared shared assets when it is constructed.
Args:
plan: Replication layout to dispatch.
stage: USD stage to author replicated prim specs into.
replicate_physics: Whether physics replication clones each environment. If False,
cloning is USD-only; an asset whose contexts are all physics-based is not cloned.
"""
from isaaclab.sim import SimulationContext # noqa: PLC0415
queued = REPLICATION_QUEUE.copy()
REPLICATION_QUEUE.clear()
backend_package = FactoryBase._get_package_name(FactoryBase._get_backend())
backend_physics_ctx = importlib.import_module(f"{backend_package}.cloner").PHYSICS_CONTEXT
# Group queued cfgs by backend, taking the union of row indices each backend owns.
# In the homogeneous plan every cfg maps to row 0, so multiple queue_replication
# calls (e.g. one per body type in RigidObjectCollection) all contribute {0} and the set
# union keeps it as a single row — no redundant copy specs are authored.
kit_available = has_kit()
backend_rows: dict[type, set[int]] = {}
for cfg in queued:
rows = plan.cfg_rows.get(id(cfg))
if rows is None:
continue
if cfg.cloning_contexts is None:
contexts = [backend_physics_ctx]
else:
contexts = [string_to_callable(c) if isinstance(c, str) else c for c in cfg.cloning_contexts]
if not replicate_physics:
contexts = [c for c in contexts if c is UsdReplicateContext]
ctx_set = dict.fromkeys(contexts)
if cfg.spawn is not None and kit_available:
ctx_set.setdefault(UsdReplicateContext, None)
for BackendCtxCls in ctx_set:
backend_rows.setdefault(BackendCtxCls, set()).update(rows)
backend_ctxs: dict[type, Any] = {}
for BackendCtxCls, row_set in backend_rows.items():
ctx = BackendCtxCls(stage, global_paths=plan.global_paths)
backend_ctxs[BackendCtxCls] = ctx
row_list = sorted(row_set)
ctx.queue_mapping(
[plan.sources[i] for i in row_list],
[plan.destinations[i] for i in row_list],
plan.env_ids,
plan.clone_mask[row_list],
positions=plan.positions,
)
for ctx in sorted(backend_ctxs.values(), key=lambda ctx: ctx.replicate_priority):
ctx.replicate()
SimulationContext.instance().set_clone_plan(plan)
[docs]
class ReplicateSession:
"""Folds :func:`make_clone_plan` and :func:`replicate` into a ``with`` block.
``__enter__`` builds the plan (and mutates each cfg's ``spawn_path``); asset
constructors inside the block register their cfgs into
:data:`REPLICATION_QUEUE`; ``__exit__`` drains and dispatches.
Example:
.. code-block:: python
with cloner.ReplicateSession(cfgs, num_clones=128, env_spacing=2.0, device="cuda:0", stage=sim.stage):
for cfg in cfgs:
cfg.class_type(cfg)
"""
[docs]
def __init__(
self,
cfgs: Iterable[Any],
num_clones: int,
env_spacing: float,
device: str,
*,
stage: Usd.Stage,
global_paths: tuple[str, ...] = (),
clone_strategy: Callable = sequential,
valid_set: torch.Tensor | None = None,
replicate_physics: bool = True,
env_template: str = DEFAULT_ENV_TEMPLATE,
):
"""Capture arguments for :func:`make_clone_plan` and :func:`replicate`.
Args:
cfgs: Asset cfgs with resolved ``prim_path``.
num_clones: Number of target envs.
env_spacing: Grid spacing between env origins [m].
device: Torch device for plan tensors.
stage: USD stage to author replicated prim specs into.
global_paths: Complete shared-asset roots declared by the composition root. Defaults to none.
clone_strategy: Prototype-to-env assignment function.
valid_set: Optional ``[num_combos, num_groups]`` long tensor of valid
prototype combinations; ``None`` uses the full cartesian product.
replicate_physics: Whether physics replication clones each environment;
forwarded to :func:`replicate`.
env_template: Path template for a replicated env prim, ``{}`` marking the env index.
"""
self._cfgs = cfgs
self._stage = stage
self._replicate_physics = replicate_physics
self._kwargs = dict(
num_clones=num_clones,
env_spacing=env_spacing,
device=device,
global_paths=global_paths,
clone_strategy=clone_strategy,
valid_set=valid_set,
env_template=env_template,
)
self._plan: ClonePlan | None = None
def __enter__(self) -> ReplicateSession:
self._plan = make_clone_plan(self._cfgs, **self._kwargs)
return self
def __exit__(self, exc_type, exc_value, traceback) -> None:
if exc_type is None:
assert self._plan is not None
replicate(self._plan, stage=self._stage, replicate_physics=self._replicate_physics)
else:
# Drop cfgs registered before the failure so the next session is clean.
REPLICATION_QUEUE.clear()
@property
def plan(self) -> ClonePlan:
"""The :class:`~isaaclab.cloner.ClonePlan` produced in :meth:`__enter__`."""
if self._plan is None:
raise RuntimeError("ReplicateSession.plan is only available inside the with block.")
return self._plan