nemo_rl.algorithms.single_controller_utils.setup#

Driver-side factory for the SingleController (async-RL) training path.

setup builds the full SingleControllerActorArgs on the driver and the caller passes it to SingleControllerActor.remote. Everything lives on the driver because driver-side TQPolicy owns the worker group directly — running this inside another Ray actor nests runtime_envs and breaks Ray’s resource resolution (see the PR #2692 follow-up).

Module Contents#

Classes#

SingleControllerActorArgs

All inputs SingleControllerActor needs, built driver-side by setup_single_controller().

Functions#

_maybe_restore_native_data_plane_checkpoint

Load and validate an authoritative native TQ checkpoint when present.

_register_single_controller_partitions

Warm all SingleController partitions before concurrent data-plane use.

_non_colocated_teacher_node_count

Validate teacher GPU geometry and return its deduplicated node count.

_build_clusters

Allocate student clusters while leaving validated nodes for teachers.

_build_generation

Spin up the generation backend (vLLM, SGLang, or Megatron).

_finish_deferred_generation

Finish loading and starting the deferred generation.

_build_trainer

Build the TQ-mediated trainer (driver-side TQPolicy).

_build_value

Build the TQ-mediated PPO critic (driver-side TQValue).

_spinup_gym

Spin up the NeMo-Gym shard set against the reserved vLLM URLs.

_generation_max_seq_len

Return the per-backend max sequence length.

_clamp_max_num_steps

Clamp max_num_steps to max_num_epochs * len(dataloader).

_maybe_inject_megatron_train_iters

Set train_iters from max_num_steps after its dataloader clamp.

_maybe_attach_fleet_health

Route generation through fleet health, when it is enabled and supported.

_shard_base_urls

Per-shard OpenAI base URLs, or None when the backend exposes no servers.

_maybe_start_generation_router

Start the NeMo-Gym-facing router, if enabled.

_build_advantage_estimator

Build the advantage estimator from whichever algorithm’s factory applies.

_build_retry_policy

Translate async_rl.rollout_failure into the rollout layer’s policy object.

_load_opd_full_teacher_lm_heads

Load every unique teacher’s LM head onto each student worker.

setup_single_controller

Build the full SC actor args driver-side.

API#

class nemo_rl.algorithms.single_controller_utils.setup.SingleControllerActorArgs#

All inputs SingleControllerActor needs, built driver-side by setup_single_controller().

Passed as a single arg to SingleControllerActor.remote. Heavy objects are cloudpickled in; Mooncake restore stays actor-side so its process-local memory segment is attached before checkpoint data is loaded.

gen_handle: Any#

None

trainer_handle: Any#

None

env_handles: dict[str, nemo_rl.environments.interfaces.EnvironmentInterface]#

None

train_cluster: nemo_rl.distributed.virtual_cluster.RayVirtualCluster#

None

inference_cluster: nemo_rl.distributed.virtual_cluster.RayVirtualCluster#

None

dp_client: nemo_rl.data_plane.DataPlaneClient#

None

dataloader: torchdata.stateful_dataloader.StatefulDataLoader#

None

weight_synchronizer: nemo_rl.weight_sync.WeightSynchronizer#

None

advantage_estimator: Any#

None

loss_fn: nemo_rl.algorithms.loss.interfaces.LossFunction#

None

rollout_manager: nemo_rl.experience.rollout_manager.RolloutManager#

None

tq_buffer: nemo_rl.algorithms.async_utils.replay_buffer.TQReplayBuffer#

None

partition_id: str#

None

save_state: nemo_rl.algorithms.grpo.GRPOSaveState#

None

last_checkpoint_path: Optional[str]#

None

finalizer_actors: list[Any]#

None

advantage_actors: list[Any]#

None

data_plane_checkpoint_metadata: Optional[nemo_rl.algorithms.async_utils.replay_buffer.DataPlaneCheckpointMetadata]#

None

partition_includes_multimodal_fields: bool#

False

bootstrap_identity: Optional[nemo_rl.algorithms.single_controller_utils.rollout_checkpoint.BootstrapCompatibilityIdentity]#

None

rollout_checkpoint_load_metrics: Optional[dict[str, float]]#

None

fleet_monitor: Optional[nemo_rl.models.generation.fleet_health.GenerationFleetHealth]#

None

generation_router: Optional[ray.actor.ActorHandle[nemo_rl.models.generation.generation_router.GenerationRouterImpl]]#

None

teacher_worker_groups: Optional[dict[str, Any]]#

None

alias_to_group_alias: Optional[dict[str, str]]#

None

value_handle: Optional[nemo_rl.models.value.tq_value.TQValue]#

None

value_loss_fn: Optional[nemo_rl.algorithms.loss.interfaces.LossFunction]#

None

nemo_rl.algorithms.single_controller_utils.setup._maybe_restore_native_data_plane_checkpoint(
*,
load_checkpoint: Callable[[str | pathlib.Path], dict[str, Any]],
last_checkpoint_path: Optional[str],
save_state: nemo_rl.algorithms.grpo.GRPOSaveState,
partition_id: str,
sampler_name: str,
opd_full_teacher_checkpoints: Optional[list[str]] = None,
) → Optional[nemo_rl.algorithms.async_utils.replay_buffer.DataPlaneCheckpointMetadata]#

Load and validate an authoritative native TQ checkpoint when present.

The replay metadata file is the format marker. Checkpoints without any replay artifact resume trainer state with an empty replay buffer; legacy tensor-bearing replay files are rejected rather than silently ignored. Rollout tensors are never serialized into a controller-side replay checkpoint.

nemo_rl.algorithms.single_controller_utils.setup._register_single_controller_partitions(
dp_client: nemo_rl.data_plane.DataPlaneClient,
*,
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
partition_id: str,
include_multimodal_fields: bool,
) → None#

Warm all SingleController partitions before concurrent data-plane use.

VLM token capture (include_multimodal_fields with capture enabled) adds the media columns the vLLM worker stages beside each captured call to the staging partition.

nemo_rl.algorithms.single_controller_utils.setup._non_colocated_teacher_node_count(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → int#

Validate teacher GPU geometry and return its deduplicated node count.

nemo_rl.algorithms.single_controller_utils.setup._build_clusters(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → tuple[nemo_rl.distributed.virtual_cluster.RayVirtualCluster, nemo_rl.distributed.virtual_cluster.RayVirtualCluster, Optional[dict[str, tuple[str, int]]]]#

Allocate student clusters while leaving validated nodes for teachers.

Colocated (Megatron generation only) shares one cluster for policy and generation and returns it for both arms; other backends split nodes into train + inference clusters.

nemo_rl.algorithms.single_controller_utils.setup._build_generation(
inference_cluster: nemo_rl.distributed.virtual_cluster.RayVirtualCluster,
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
*,
defer_model_load: bool = False,
reserved_http_server_ports: Optional[dict[int, int]] = None,
tokenizer: Optional[transformers.tokenization_utils_base.PreTrainedTokenizerBase] = None,
processor: Optional[transformers.AutoProcessor] = None,
) → tuple[Any, float]#

Spin up the generation backend (vLLM, SGLang, or Megatron).

Parameters:
  • inference_cluster – Ray virtual cluster the generation workers run on.

  • master_config – SC MasterConfig.

  • defer_model_load – If True (for the NeMo-Gym flow), reserve OpenAI server URLs without loading weights; caller runs gen.load_and_start() later (vLLM only).

  • reserved_http_server_ports – OpenAI server ports pre-published to NeMo-Gym, keyed by the distributed rank that adopts each one (Megatron only).

  • tokenizer – Tokenizer for the dedicated Megatron inference policy (Megatron only).

  • processor – Optional AutoProcessor for VLM paths (Megatron only).

Returns:

A tuple of (generation object, wall time spent in this call). The generation object is a VllmGeneration, SGLangGeneration, or MegatronGeneration.

nemo_rl.algorithms.single_controller_utils.setup._finish_deferred_generation(
generation: Any,
) → tuple[Any, float]#

Finish loading and starting the deferred generation.

Parameters:

generation – The deferred generation object.

Returns:

A tuple of (finished generation object, wall time spent in this call).

nemo_rl.algorithms.single_controller_utils.setup._build_trainer(
train_cluster: nemo_rl.distributed.virtual_cluster.RayVirtualCluster,
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
tokenizer,
processor,
*,
weights_path: Optional[pathlib.Path],
optimizer_path: Optional[pathlib.Path],
checkpointing: bool,
reserved_http_server_ports: Optional[dict[int, int]] = None,
) → tuple[Any, float]#

Build the TQ-mediated trainer (driver-side TQPolicy).

Parameters:
  • train_cluster – Ray virtual cluster the trainer workers run on.

  • master_config – SC MasterConfig.

  • tokenizer – Tokenizer used by the policy.

  • processor – Optional AutoProcessor for VLM paths.

  • weights_path – Checkpointed policy weights to resume from, or None.

  • optimizer_path – Checkpointed optimizer state to resume from, or None.

  • checkpointing – Whether data-plane checkpoint save or restore is needed.

  • reserved_http_server_ports – Pre-published OpenAI server ports for NeMo Gym, keyed by the colocated Megatron trainer rank that adopts each one.

Returns:

A tuple of (TQPolicy trainer, wall time spent in this call).

nemo_rl.algorithms.single_controller_utils.setup._build_value(
train_cluster: nemo_rl.distributed.virtual_cluster.RayVirtualCluster,
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
tokenizer: transformers.tokenization_utils_base.PreTrainedTokenizerBase,
*,
weights_path: Optional[pathlib.Path],
optimizer_path: Optional[pathlib.Path],
) → tuple[nemo_rl.models.value.tq_value.TQValue, float]#

Build the TQ-mediated PPO critic (driver-side TQValue).

Parameters:
  • train_cluster – Ray virtual cluster the critic shares with the trainer.

  • master_config – SC MasterConfig.

  • tokenizer – Tokenizer used by the value model.

  • weights_path – Checkpointed value weights to resume from, or None.

  • optimizer_path – Checkpointed value optimizer state to resume from, or None.

Returns:

A tuple of (TQValue critic, wall time spent in this call).

nemo_rl.algorithms.single_controller_utils.setup._spinup_gym(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
base_urls: list[str],
tokenizer: transformers.tokenization_utils_base.PreTrainedTokenizerBase,
) → tuple[nemo_rl.environments.nemo_gym.NemoGymShardSet, float]#

Spin up the NeMo-Gym shard set against the reserved vLLM URLs.

Parameters:
  • master_config – SC MasterConfig.

  • base_urls – Reserved vLLM OpenAI server URLs.

  • tokenizer – Installed on the actor at spinup rather than passed per rollout call. See NemoGym.set_tokenizer.

Returns:

A tuple of (NeMo-Gym shard set, wall time spent in this call).

nemo_rl.algorithms.single_controller_utils.setup._generation_max_seq_len(generation_config) → int#

Return the per-backend max sequence length.

vllm uses vllm_cfg.max_model_len; sglang uses sglang_cfg.context_length; megatron uses mcore_generation_config.max_model_len.

nemo_rl.algorithms.single_controller_utils.setup._clamp_max_num_steps(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
dataloader: torchdata.stateful_dataloader.StatefulDataLoader,
) → None#

Clamp max_num_steps to max_num_epochs * len(dataloader).

nemo_rl.algorithms.single_controller_utils.setup._maybe_inject_megatron_train_iters(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → None#

Set train_iters from max_num_steps after its dataloader clamp.

nemo_rl.algorithms.single_controller_utils.setup._maybe_attach_fleet_health(
generation: Any,
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → Optional[nemo_rl.models.generation.fleet_health.GenerationFleetHealth]#

Route generation through fleet health, when it is enabled and supported.

Returns:

The monitor the SingleController should drive, or None when fleet health is disabled or the backend does not support it.

nemo_rl.algorithms.single_controller_utils.setup._shard_base_urls(
generation: Any,
) → Optional[list[Optional[str]]]#

Per-shard OpenAI base URLs, or None when the backend exposes no servers.

nemo_rl.algorithms.single_controller_utils.setup._maybe_start_generation_router(
base_urls: list[Optional[str]],
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → Any#

Start the NeMo-Gym-facing router, if enabled.

Parameters:
  • base_urls – OpenAI server URLs for the router.

  • master_config – SingleController MasterConfig.

Returns:

The router actor handle, or None when the router is disabled.

nemo_rl.algorithms.single_controller_utils.setup._build_advantage_estimator(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → Any#

Build the advantage estimator from whichever algorithm’s factory applies.

nemo_rl.algorithms.single_controller_utils.setup._build_retry_policy(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
) → nemo_rl.experience.rollout_manager.RolloutRetryPolicy#

Translate async_rl.rollout_failure into the rollout layer’s policy object.

nemo_rl.algorithms.single_controller_utils.setup._load_opd_full_teacher_lm_heads(
trainer: Any,
teacher_worker_groups: dict[str, Any],
) → None#

Load every unique teacher’s LM head onto each student worker.

One RPC per unique checkpoint (TeacherWorkerGroup.teacher_index, assigned by create_teacher_worker_groups), so the student ends up with one LM-head shard per teacher, keyed by that same index – the key every row’s OPD_FULL_TEACHER_INDEX_FIELD tag resolves against at training time (see reconstruct_opd_full_teacher_logits).

Callers gate this on the hidden_states payload; the logits payload ships the projected distribution and needs no teacher LM head.

Parameters:
  • trainer – The driver-side policy whose workers hold the student model.

  • teacher_worker_groups – Deduplicated teacher groups, keyed by primary alias.

nemo_rl.algorithms.single_controller_utils.setup.setup_single_controller(
master_config: nemo_rl.algorithms.single_controller_utils.config.MasterConfig,
tokenizer: transformers.tokenization_utils_base.PreTrainedTokenizerBase,
*,
processor: Optional[transformers.AutoProcessor] = None,
partition_id: str = 'rollout_data',
) → tuple[nemo_rl.algorithms.single_controller_utils.setup.SingleControllerActorArgs, nemo_rl.algorithms.metric_utils.SetupTimingMetrics]#

Build the full SC actor args driver-side.

Parameters:
  • master_config – SC MasterConfig.

  • tokenizer – Tokenizer used by the policy.

  • processor – Optional AutoProcessor for VLM paths.

  • partition_id – TQ partition the rollout writer + sampler share.

Returns:

A tuple of (pre-built SC actor args, driver-side per-phase timings logged by the SC actor).