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#
All inputs SingleControllerActor needs, built driver-side by setup_single_controller(). |
Functions#
Load and validate an authoritative native TQ checkpoint when present. |
|
Warm all SingleController partitions before concurrent data-plane use. |
|
Validate teacher GPU geometry and return its deduplicated node count. |
|
Allocate student clusters while leaving validated nodes for teachers. |
|
Spin up the generation backend (vLLM, SGLang, or Megatron). |
|
Finish loading and starting the deferred generation. |
|
Build the TQ-mediated trainer (driver-side TQPolicy). |
|
Build the TQ-mediated PPO critic (driver-side TQValue). |
|
Spin up the NeMo-Gym shard set against the reserved vLLM URLs. |
|
Return the per-backend max sequence length. |
|
Clamp max_num_steps to max_num_epochs * len(dataloader). |
|
Set train_iters from max_num_steps after its dataloader clamp. |
|
Route generation through fleet health, when it is enabled and supported. |
|
Per-shard OpenAI base URLs, or None when the backend exposes no servers. |
|
Start the NeMo-Gym-facing router, if enabled. |
|
Build the advantage estimator from whichever algorithm’s factory applies. |
|
Translate |
|
Load every unique teacher’s LM head onto each student worker. |
|
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,
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,
Warm all SingleController partitions before concurrent data-plane use.
VLM token capture (
include_multimodal_fieldswith 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,
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,
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,
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,
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,
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],
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,
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,
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,
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,
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,
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,
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,
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,
Translate
async_rl.rollout_failureinto 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],
Load every unique teacher’s LM head onto each student worker.
One RPC per unique checkpoint (
TeacherWorkerGroup.teacher_index, assigned bycreate_teacher_worker_groups), so the student ends up with one LM-head shard per teacher, keyed by that same index – the key every row’sOPD_FULL_TEACHER_INDEX_FIELDtag resolves against at training time (seereconstruct_opd_full_teacher_logits).Callers gate this on the
hidden_statespayload; thelogitspayload 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',
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).