nemo_rl.data_plane.factory#

Single entrypoint that maps a :class:DataPlaneConfig to a client.

Module Contents#

Functions#

data_plane_enabled

Whether the data plane is on. None (key absent) means off.

select_sync_trainer

Pick the synchronous trainer based on data_plane.enabled.

make_policy_factory

The policy_factory for setup(), or None for a plain Policy.

maybe_configure_data_plane_env

Set backend env vars that must be identical in every process.

build_data_plane_client

Construct the configured data-plane client.

API#

nemo_rl.data_plane.factory.data_plane_enabled(
cfg: nemo_rl.data_plane.interfaces.DataPlaneConfig | None,
) → bool#

Whether the data plane is on. None (key absent) means off.

nemo_rl.data_plane.factory.select_sync_trainer(
master_config: nemo_rl.algorithms.grpo.MasterConfig,
*,
label: str = 'GRPO',
) → Callable[..., Any]#

Pick the synchronous trainer based on data_plane.enabled.

Shared by every launcher so trainer choice cannot drift between them. Pairs with :func:make_policy_factory: turning the data plane on means picking both the TQ-mediated trainer and a TQPolicy, and a launcher that picked only one would fail after a full model load.

Parameters:
  • master_config – The resolved config; only data_plane is read.

  • label – Algorithm name for the progress line (e.g. "VLM GRPO").

nemo_rl.data_plane.factory.make_policy_factory(
cfg: nemo_rl.data_plane.interfaces.DataPlaneConfig | None,
) → Optional[Callable[..., Any]]#

The policy_factory for setup(), or None for a plain Policy.

Lives at the launcher level so the legacy trainer stays data-plane-agnostic (architectural invariant — see tests/unit/data_plane/test_architecture_invariants.py).

nemo_rl.data_plane.factory.maybe_configure_data_plane_env(
cfg: nemo_rl.data_plane.interfaces.DataPlaneConfig | None,
) → None#

Set backend env vars that must be identical in every process.

Call this on the driver before init_ray(): init_ray snapshots the driver’s environment into runtime_env["env_vars"] and hands it to every Ray worker, which are fresh processes, so the value is in place before they run any engine code.

The binding constraint is not init_ray but that this must run before anything imports the backend’s engine, which snapshots its configuration as it loads. :func:~nemo_rl.data_plane.adapters.transfer_queue_env.configure_engine_env raises rather than silently no-op’ing if that is violated.

Subprocesses the driver did not spawn through Ray inherit whatever the driver’s environment held at fork, so they are covered as long as this ran first.

No-op when the data plane is disabled or the backend has no such knobs.

Parameters:

cfg – Data-plane config, or None when the data plane is off.

nemo_rl.data_plane.factory.build_data_plane_client(
cfg: nemo_rl.data_plane.interfaces.DataPlaneRuntimeConfig | None,
*,
bootstrap: bool = True,
checkpointing: bool = False,
) → nemo_rl.data_plane.interfaces.DataPlaneClient#

Construct the configured data-plane client.

Dispatches on the configured implementation. TransferQueue supports cross-process transfer; the local adapter keeps colocated SFT batches in one process. Raises if data_plane is disabled — the legacy trainer (nemo_rl.algorithms.grpo.grpo_train) should be used in that case rather than a NoOp fallback here.

Parameters:
  • cfg – Data-plane config; must have enabled=True.

  • bootstrap – True on the driver — bootstraps the TQ controller. False on worker processes — connects to the existing controller (avoids creating a second named actor).

  • checkpointing – Prepare storage for saving or restoring data-plane state. Derived by the caller from its existing checkpoint settings and resume path. Only used at bootstrap; workers inherit the mode from TQ.

Returns:

A configured DataPlaneClient; wrapped in

class:

MetricsDataPlaneClient when observability is enabled.