diff --git a/docs/guides/single-controller.md b/docs/guides/single-controller.md index d50bac5d28..391ff704df 100644 --- a/docs/guides/single-controller.md +++ b/docs/guides/single-controller.md @@ -234,6 +234,9 @@ Do not carry `max_num_epochs: -1` across either. [ppo.md](./ppo.md#asynchronous- The SC path is still under active development. Feature gaps are tracked in [issue #2625](https://github.com/NVIDIA-NeMo/RL/issues/2625). Notable items: +- Multimodal/VLM GRPO is supported with Megatron generation. Set + `policy.is_vlm: true`; see the + [CLEVR Single-Controller recipe](../../examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.yaml). - Multi-Teacher On-Policy Distillation (MOPD) is supported for text-only NeMo Gym rollouts; multimodal/VLM MOPD is not yet supported. See [Multi-Teacher On-Policy Distillation](../about/algorithms/mopd.md#running-mopd). diff --git a/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-16n8g-megatron-tp4ep4-async-gym-video.v1.yaml b/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-16n8g-megatron-tp4ep4-async-gym-video.v1.yaml index 20799328e5..4dafdca018 100644 --- a/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-16n8g-megatron-tp4ep4-async-gym-video.v1.yaml +++ b/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-16n8g-megatron-tp4ep4-async-gym-video.v1.yaml @@ -49,6 +49,14 @@ policy: generation: bad_words: [] mcore_generation_config: + buffer_size_gb: 8 + use_cuda_graphs_for_non_decode_steps: false + moe_pad_experts_for_cuda_graph_inference: false + moe_router_dtype: fp32 + vision_embedding_cache_max_bytes: 536870912 + video_num_frames: ${data.default.num_frames} + video_temporal_patch_size: ${data.default.video_temporal_patch_size} + video_target_num_patches: ${data.default.video_target_num_patches} image_dynamic_resolution: true logprobs_mode: raw_logprobs megatron_inference_wrapper: megatron.core.inference.model_inference_wrappers.multimodal.nemotron_omni_inference_wrapper.NemotronOmniInferenceWrapper @@ -91,6 +99,9 @@ data: max_input_seq_length: ${policy.max_total_sequence_length} num_workers: 0 default: + num_frames: 32 + video_sampling_style: nemotron_vl + video_temporal_patch_size: 2 video_target_num_patches: 1024 video_maintain_aspect_ratio: true env: diff --git a/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.yaml b/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.yaml new file mode 100644 index 0000000000..02392c7f29 --- /dev/null +++ b/examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.yaml @@ -0,0 +1,18 @@ +defaults: ./vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron_generation.v1.yaml +grpo: + async_grpo: null + val_period: 0 + val_at_start: false + val_at_end: false +data_plane: + enabled: true + simple: + num_storage_units: 16 +async_rl: + sampler: + name: in_order + max_lookahead_versions: 1 + recompute_kv_cache_after_weight_updates: false + min_groups_for_streaming_train: ${grpo.num_prompts_per_step} + max_inflight_prompts: ${mul:${grpo.num_prompts_per_step}, 2} + max_buffered_rollouts: ${mul:${grpo.num_prompts_per_step}, 2} diff --git a/examples/run_grpo_single_controller.py b/examples/run_grpo_single_controller.py index 3a151016a2..b06b7ccf9c 100644 --- a/examples/run_grpo_single_controller.py +++ b/examples/run_grpo_single_controller.py @@ -126,7 +126,12 @@ def main() -> None: maybe_configure_data_plane_env(config.data_plane) init_ray() - tokenizer = get_tokenizer(config.policy["tokenizer"]) + processor = None + if config.policy.get("is_vlm"): + processor = get_tokenizer(config.policy["tokenizer"], get_processor=True) + tokenizer = processor.tokenizer + else: + tokenizer = get_tokenizer(config.policy["tokenizer"]) assert config.policy["generation"] is not None, ( "A generation config is required for SC-driven async GRPO" ) @@ -144,7 +149,9 @@ def main() -> None: if bool(config.env.get("should_use_nemo_gym")): setup_nemo_gym_config(config, tokenizer) - actor_args, setup_timing_metrics = setup_single_controller(config, tokenizer) + actor_args, setup_timing_metrics = setup_single_controller( + config, tokenizer, processor=processor + ) print("🚀 Launching SingleControllerActor") sc = SingleControllerActor.remote( diff --git a/nemo_rl/algorithms/single_controller_utils/setup.py b/nemo_rl/algorithms/single_controller_utils/setup.py index d753bc602c..a3411f11bf 100644 --- a/nemo_rl/algorithms/single_controller_utils/setup.py +++ b/nemo_rl/algorithms/single_controller_utils/setup.py @@ -70,6 +70,7 @@ ) from nemo_rl.algorithms.utils import set_seed from nemo_rl.data.collate_fn import rl_collate_fn +from nemo_rl.data.multimodal_utils import WIRE_MULTIMODAL_FIELDS from nemo_rl.data.utils import load_dataloader_state, setup_response_data from nemo_rl.data_plane import ( DATA_PLANE_CHECKPOINT_SCHEMA_VERSION, @@ -997,9 +998,10 @@ def setup_single_controller( policy_config["pretrained_checkpoint"] = checkpointing_pretrained # Token capture: validate the supported combination loudly at setup - # (NeMo-Gym rollout path, vLLM backend, async_engine=true) and give - # capture-enabled vLLM workers a venv that carries nemo_gym (the - # worker hosts Gym's capture core + adapter in-process). + # (NeMo-Gym rollout path, vLLM backend, async_engine=true). The vLLM + # worker venv always carries nemo_gym (see VLLM_EXECUTABLE in + # ray_actor_environment_registry.py), so nothing here needs to change the + # worker's environment. token_capture_cfg = master_config.token_capture if token_capture_cfg.enabled: if not should_use_nemo_gym(master_config): @@ -1020,14 +1022,6 @@ def setup_single_controller( "policy.generation.vllm_cfg.async_engine=true (the capture " "host is the worker's in-process HTTP server)" ) - from nemo_rl.distributed.ray_actor_environment_registry import ( - ACTOR_ENVIRONMENT_REGISTRY, - ) - from nemo_rl.distributed.virtual_cluster import PY_EXECUTABLES - - ACTOR_ENVIRONMENT_REGISTRY[ - "nemo_rl.models.generation.vllm.vllm_worker_async.VllmAsyncGenerationWorker" - ] = PY_EXECUTABLES.VLLM_GYM # Fill the derived ledger-hosting fields (see TokenCaptureConfig): a # per-run control-plane bearer token and the process-shared capture @@ -1074,6 +1068,8 @@ def setup_single_controller( # ========================== # TODO: add validate dataset wiring. use_nemo_gym = should_use_nemo_gym(master_config) + data_tokenizer = processor if processor is not None else tokenizer + is_vlm = processor is not None if use_nemo_gym and generation_config["backend"] not in ("vllm", "megatron"): raise NotImplementedError( "SC NeMo-Gym integration currently supports the vllm and megatron backends only; got " @@ -1084,13 +1080,18 @@ def setup_single_controller( if use_nemo_gym: # NeMo-Gym creates the env actor outside setup_response_data; we wire # it in after generation is up (it needs the OpenAI server URLs). - response_data = setup_response_data(tokenizer, data_config, env_configs=None) + response_data = setup_response_data( + data_tokenizer, data_config, env_configs=None, is_vlm=is_vlm + ) assert len(response_data) == 2 dataset, _val_dataset = response_data env_handles: dict[str, EnvironmentInterface] = {} else: response_data = setup_response_data( - tokenizer, data_config, env_configs=master_config.env + data_tokenizer, + data_config, + env_configs=master_config.env, + is_vlm=is_vlm, ) assert len(response_data) == 4 dataset, _val_dataset, env_handles, _val_env_handles = response_data @@ -1154,6 +1155,7 @@ def setup_single_controller( megatron_reserved_url = None megatron_port_holder = None reserved_http_server_port = None + weight_synchronizer: Optional[WeightSynchronizer] = None if megatron_backend: generation_config["model_name"] = master_config.policy["model_name"] @@ -1313,7 +1315,6 @@ def _build_generation_then_trainer( build_tasks["trainer"] = _build_trainer_and_value # Submit build tasks and get results - weight_synchronizer: Optional[WeightSynchronizer] = None try: with ThreadPoolExecutor(max_workers=len(build_tasks)) as executor: submitted = {k: executor.submit(fn) for k, fn in build_tasks.items()} @@ -1340,7 +1341,9 @@ def _build_generation_then_trainer( train_cluster=train_cluster, inference_cluster=inference_cluster, refit_buffer_size_gb=policy_config.get("refit_buffer_size_gb"), + refit_timeout_s=master_config.async_rl.generation_fleet_health.refit_timeout_s, ) + generation.weight_synchronizer = weight_synchronizer weight_synchronizer.init_communicator() setup_timing_metrics.collective_init_time_s = time.perf_counter() - t0 t0 = time.perf_counter() @@ -1427,12 +1430,19 @@ def _build_generation_then_trainer( # SingleController reuses one partition for the run. Warm every known # tensor field before rollout, policy, and teacher writers become # concurrent; TransferQueue otherwise registers field names lazily. + partition_fields = fields_with_optional_routed_experts( + SC_ROLLOUT_SCHEMA_FIELDS, + enabled=router_replay_enabled(policy_config), + ) + if processor is not None: + partition_fields.extend( + field + for field in sorted(WIRE_MULTIMODAL_FIELDS) + if field not in partition_fields + ) dp_client.register_partition( partition_id=partition_id, - fields=fields_with_optional_routed_experts( - SC_ROLLOUT_SCHEMA_FIELDS, - enabled=router_replay_enabled(policy_config), - ), + fields=partition_fields, num_samples=( master_config.async_rl.max_buffered_rollouts * algo_cfg.num_generations_per_prompt @@ -1457,13 +1467,19 @@ def _build_generation_then_trainer( ) group_size = algo_cfg.num_generations_per_prompt num_rollout_samples = master_config.async_rl.max_buffered_rollouts * group_size + partition_fields = fields_with_optional_routed_experts( + DP_TRAIN_FIELDS, + enabled=r3_enabled and not token_capture_cfg.defer_routed_experts_to_policy, + ) + if processor is not None: + partition_fields.extend( + field + for field in sorted(WIRE_MULTIMODAL_FIELDS) + if field not in partition_fields + ) dp_client.register_partition( partition_id=partition_id, - fields=fields_with_optional_routed_experts( - DP_TRAIN_FIELDS, - enabled=r3_enabled - and not token_capture_cfg.defer_routed_experts_to_policy, - ), + fields=partition_fields, num_samples=num_rollout_samples, consumer_tasks=["prev_lp", "ref_lp", "train"], grpo_group_size=group_size, @@ -1478,23 +1494,7 @@ def _build_generation_then_trainer( # Host Gym's capture core in every vLLM DP leader (in-worker DP # client + TQTokenSink + the single install_capture call), and give # workers the initial weight version to stamp on captured calls. - try: - generation.setup_token_capture( - dp_config, token_capture_cfg.staging_partition - ) - except Exception as error: - if "No module named 'nemo_gym'" in str(error): - # Worker venvs are cached by actor class name - # (nemo_rl/utils/venvs.py), so a venv prebuilt before token - # capture predates the nemo_gym extra and is reused as-is. - raise RuntimeError( - "token_capture.enabled requires nemo_gym inside the vLLM " - "worker venv, but the cached worker venv predates it. " - "Rebuild worker venvs (NRL_FORCE_REBUILD_VENVS=true) or " - "delete $NEMO_RL_VENV_DIR/nemo_rl.models.generation.vllm." - "vllm_worker_async.VllmAsyncGenerationWorker and rerun." - ) from error - raise + generation.setup_token_capture(dp_config, token_capture_cfg.staging_partition) generation.set_rollout_weight_version(0) if weight_synchronizer is None: @@ -1509,6 +1509,7 @@ def _build_generation_then_trainer( refit_buffer_size_gb=policy_config.get("refit_buffer_size_gb"), refit_timeout_s=master_config.async_rl.generation_fleet_health.refit_timeout_s, ) + generation.weight_synchronizer = weight_synchronizer weight_synchronizer.init_communicator() setup_timing_metrics.collective_init_time_s = time.perf_counter() - t0 diff --git a/nemo_rl/data/multimodal_utils.py b/nemo_rl/data/multimodal_utils.py index e9bebd4e44..4f047bb21e 100644 --- a/nemo_rl/data/multimodal_utils.py +++ b/nemo_rl/data/multimodal_utils.py @@ -1093,6 +1093,18 @@ def encode_multimodal_for_wire( ) +# Model inputs some remote-code processors omit from ``model_input_names`` even +# though their forward requires them. Consumed by +# ``extract_multimodal_model_inputs``; membership here does NOT imply the field +# is wire-registered (see ``PACKED_/PER_TOKEN_MULTIMODAL_FIELDS``). +UNDECLARED_MULTIMODAL_MODEL_INPUTS = ( + "imgs_sizes", + "num_frames", + "pixel_values_flat", + "image_num_patches", +) + + def get_multimodal_keys_from_processor(processor) -> list[str]: """Get keys of the multimodal data that can be used as model inputs. @@ -1216,12 +1228,7 @@ def extract_multimodal_model_inputs( # TODO(rohitrango): Let ProcessorInterface declare model-specific media inputs. # Some remote-code processors omit these inputs from model_input_names even # though their model forward requires them. - for key in ( - "imgs_sizes", - "num_frames", - "pixel_values_flat", - "image_num_patches", - ): + for key in UNDECLARED_MULTIMODAL_MODEL_INPUTS: if key in processed and key not in multimodal_keys: multimodal_keys.append(key) for key in multimodal_keys: @@ -1233,7 +1240,7 @@ def extract_multimodal_model_inputs( f"Processor model input {key!r} must be a torch.Tensor, got " f"{type(value).__name__}." ) - if key == "imgs_sizes": + if key in ("imgs_sizes", "num_frames"): value = value.to(dtype=torch.int32) extracted[key] = PackedTensor( value, diff --git a/nemo_rl/data/processors.py b/nemo_rl/data/processors.py index ef8744afb7..77cb14dae2 100644 --- a/nemo_rl/data/processors.py +++ b/nemo_rl/data/processors.py @@ -460,9 +460,8 @@ def vlm_hf_data_processor( from nemo_rl.data.datasets.response_datasets.refcoco import format_refcoco_dataset from nemo_rl.data.multimodal_utils import ( PackedTensor, - get_dim_to_pack_along, + extract_multimodal_model_inputs, get_multimodal_default_settings_from_processor, - get_multimodal_keys_from_processor, resolve_to_image, uses_image_placeholder, ) @@ -608,50 +607,13 @@ def vlm_hf_data_processor( # add this for backward compatibility user_message["token_ids"] = message["input_ids"][0] - # add all keys and values to the user message, and the list of keys - multimodal_keys = list(get_multimodal_keys_from_processor(processor)) - # Current Nemotron Omni processors emit imgs_sizes. Historical MMPR - # checkpoints instead emit a batch of fixed-size image tiles and only - # declare pixel_values. Treat each tile as one dynamic-resolution image so - # the Nemotron Omni path can patchify it and preserve the processor's exact - # placeholder count. - if ( - uses_placeholder - and "pixel_values" in message - and "imgs_sizes" not in message - and message["pixel_values"].ndim == 4 - ): - pixel_values = message["pixel_values"] - num_tiles, _, height, width = pixel_values.shape - message["imgs_sizes"] = torch.tensor( - [[height, width]] * num_tiles, dtype=torch.long - ) - - # imgs_sizes is not always declared in model_input_names by bundled image - # processors, so append it explicitly when present. RADIO uses temporal - # patching even for still images and requires one num_frames=1 entry per - # image/tile. - if "imgs_sizes" in message and "imgs_sizes" not in multimodal_keys: - multimodal_keys.append("imgs_sizes") - if "imgs_sizes" in message and "num_frames" not in message: - message["num_frames"] = torch.ones(len(message["imgs_sizes"]), dtype=torch.long) - if "num_frames" in message and "num_frames" not in multimodal_keys: - multimodal_keys.append("num_frames") - for key in multimodal_keys: - if key in message: - user_message[key] = PackedTensor( - message[key], - dim_to_pack=get_dim_to_pack_along(processor, key), - pad_to_max_shape=uses_placeholder and key == "pixel_values", - ) - - # specifically for gemma, we need to add token_type_ids to the user message as a sequence-type value - if "token_type_ids" in message: - user_message["token_type_ids"] = message["token_type_ids"][0] - - # for qwen2.5-vl (transformers>=5.3), mm_token_type_ids tells the model which tokens are text/image/video for 3D RoPE - if "mm_token_type_ids" in message: - user_message["mm_token_type_ids"] = message["mm_token_type_ids"][0] + # Single source of truth for media extraction: the MMPR imgs_sizes + # fallback, RADIO num_frames synthesis, PackedTensor wrapping (incl. the + # imgs_sizes int32 cast) and the gemma / qwen2.5-vl sequence-type maps all + # live in ``extract_multimodal_model_inputs``, which the NeMo-Gym path also + # uses. One implementation is what stops the two paths from handing the + # same model differently-typed inputs. + user_message.update(extract_multimodal_model_inputs(processor, message)) ### append to user message message_log.append(user_message) diff --git a/nemo_rl/data_plane/interfaces.py b/nemo_rl/data_plane/interfaces.py index ce8084c4d3..9f527c4b2c 100644 --- a/nemo_rl/data_plane/interfaces.py +++ b/nemo_rl/data_plane/interfaces.py @@ -312,7 +312,22 @@ def slice(self, start: int, stop: int) -> "KVBatchMeta": ) def concat(self, *others: "KVBatchMeta") -> "KVBatchMeta": - """Append ``others`` to ``self``. All metas must share ``partition_id``.""" + """Append metadata from the same partition. + + Sample IDs are concatenated in argument order, while fields are + unioned in first-seen order. Sequence lengths and tags are retained + only when every input provides them. + + Args: + *others: Metadata batches whose ``partition_id`` matches this + batch. + + Returns: + A new metadata batch containing all input rows. + + Raises: + ValueError: If any input has a different ``partition_id``. + """ if any(o.partition_id != self.partition_id for o in others): raise ValueError("KVBatchMeta.concat: partition_ids must match") all_m = (self, *others) @@ -325,9 +340,14 @@ def concat(self, *others: "KVBatchMeta") -> "KVBatchMeta": ) all_have_tags = all(m.tags is not None for m in all_m) tags = [t for m in all_m for t in (m.tags or [])] if all_have_tags else None - return self._replace( + merged_fields = list( + dict.fromkeys(field for meta in all_m for field in (meta.fields or [])) + ) + result = self._replace( sample_ids=sample_ids, sequence_lengths=seq_lens, tags=tags ) + result.fields = merged_fields or None + return result def drop(self, indices: "Sequence[int]") -> "KVBatchMeta | None": """Complement of :meth:`subset`. Returns ``None`` when all rows are dropped.""" diff --git a/nemo_rl/data_plane/worker_mixin.py b/nemo_rl/data_plane/worker_mixin.py index 92a01628ec..a5309d453f 100644 --- a/nemo_rl/data_plane/worker_mixin.py +++ b/nemo_rl/data_plane/worker_mixin.py @@ -319,11 +319,9 @@ def setup_data_plane(self, cfg: DataPlaneRuntimeConfig) -> None: # entrypoints then delegate to the same ``train`` / ``get_logprobs`` # that carry the flag into ``models/megatron/data.py``. # - # ``train_microbatch_presharded`` is the exception: it lands in - # ``_train_microbatch_body``, which passes none of the capability flags - # and never attaches the media-token validity mask. That path is - # SingleController-only, so ``train_microbatch`` raises for a - # multimodal model rather than training on rows it mis-describes. + # ``train_microbatch_presharded`` lands in ``_train_microbatch_body``, + # which applies the same media-token validity mask and model packing/CP + # capability flags as the regular training path. if self._dp_client is not None: return self._route_fallback_counts = Counter() diff --git a/nemo_rl/distributed/ray_actor_environment_registry.py b/nemo_rl/distributed/ray_actor_environment_registry.py index 1d27c8ff41..a2788019ab 100644 --- a/nemo_rl/distributed/ray_actor_environment_registry.py +++ b/nemo_rl/distributed/ray_actor_environment_registry.py @@ -17,8 +17,13 @@ from nemo_rl.distributed.virtual_cluster import PY_EXECUTABLES USE_SYSTEM_EXECUTABLE = os.environ.get("NEMO_RL_PY_EXECUTABLES_SYSTEM", "0") == "1" +# vLLM workers always get the vllm + nemo_gym extras. Token capture +# (token_capture.enabled) needs nemo_gym inside the worker, and worker venvs +# are cached by actor class name, so the extras must be fixed here rather than +# swapped in at runtime (a venv prebuilt with plain `--extra vllm` would be +# reused as-is and the nemo_gym import would fail). VLLM_EXECUTABLE = ( - PY_EXECUTABLES.SYSTEM if USE_SYSTEM_EXECUTABLE else PY_EXECUTABLES.VLLM + PY_EXECUTABLES.SYSTEM if USE_SYSTEM_EXECUTABLE else PY_EXECUTABLES.VLLM_GYM ) SGLANG_EXECUTABLE = ( PY_EXECUTABLES.SYSTEM if USE_SYSTEM_EXECUTABLE else PY_EXECUTABLES.SGLANG diff --git a/nemo_rl/distributed/virtual_cluster.py b/nemo_rl/distributed/virtual_cluster.py index e65ef299dc..3f8ae284a6 100644 --- a/nemo_rl/distributed/virtual_cluster.py +++ b/nemo_rl/distributed/virtual_cluster.py @@ -27,6 +27,8 @@ ) from ray.util.scheduling_strategies import PlacementGroupSchedulingStrategy +from nemo_rl.utils.venvs import add_hf_modules_cache_to_pythonpath + logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) @@ -73,8 +75,12 @@ class PY_EXECUTABLES: # Use NeMo-Gym dependencies NEMO_GYM = f"uv run --locked --extra nemo_gym --directory {git_root}" - # vLLM worker hosting Gym's token capture (token_capture.enabled): the - # worker imports nemo_gym's dependency-free capture core + vLLM adapter. + # Default env for the vLLM generation workers (see + # ray_actor_environment_registry.py). It carries nemo_gym so the worker can + # host Gym's token capture (token_capture.enabled) without swapping the + # worker's env at runtime: worker venvs are cached by actor class name, so + # a venv prebuilt with plain `--extra vllm` would be reused as-is and the + # nemo_gym import would fail. VLLM_GYM = f"uv run --locked --extra vllm --extra nemo_gym --directory {git_root}" # Use NeMo-RL direct dependencies and SGLang. @@ -318,7 +324,11 @@ def init_ray(log_dir: Optional[str] = None) -> None: if _k.startswith(("PMIX_", "PMI_", "MPI_", "OMPI_", "SLURM_")): os.environ.pop(_k, None) - env_vars = dict(os.environ) + # Ray actors deserialize constructor arguments before importing NeMo-RL. + # Put Hugging Face's generated ``transformers_modules`` package on the + # cluster-wide PYTHONPATH so trust_remote_code objects can be unpickled at + # that boundary. This covers both V1 worker groups and direct V2/SC actors. + env_vars = add_hf_modules_cache_to_pythonpath(dict(os.environ)) env_vars.pop("RAY_EXPERIMENTAL_NOSET_CUDA_VISIBLE_DEVICES", None) runtime_env = { diff --git a/nemo_rl/environments/nemo_gym_multimodal.py b/nemo_rl/environments/nemo_gym_multimodal.py index 3e65b88ba6..4987288f57 100644 --- a/nemo_rl/environments/nemo_gym_multimodal.py +++ b/nemo_rl/environments/nemo_gym_multimodal.py @@ -33,7 +33,6 @@ PackedTensor, extract_input_media_sources_from_responses_messages, extract_multimodal_model_inputs, - get_dim_to_pack_along, get_responses_content_part_url, image_to_data_url, media_sources_equal, @@ -824,11 +823,6 @@ def nemo_gym_example_to_video_datum_spec( if "imgs_sizes" in processed and "num_frames" not in processed: processed["num_frames"] = torch.tensor([len(frame_items)], dtype=torch.int32) user_message.update(extract_multimodal_model_inputs(processor, processed)) - if "num_frames" in processed: - user_message["num_frames"] = PackedTensor( - processed["num_frames"].to(dtype=torch.int32), - dim_to_pack=get_dim_to_pack_along(processor, "num_frames"), - ) length = len(user_message["token_ids"]) loss_multiplier = 1.0 diff --git a/nemo_rl/experience/interfaces.py b/nemo_rl/experience/interfaces.py index 9457d98658..29320e93f2 100644 --- a/nemo_rl/experience/interfaces.py +++ b/nemo_rl/experience/interfaces.py @@ -61,3 +61,4 @@ class PromptGroupRecord: metadata: dict[str, Any] completions: list["Completion"] rollout_metrics: dict[str, Any] + loss_multiplier: float = 1.0 diff --git a/nemo_rl/experience/payload.py b/nemo_rl/experience/payload.py index c981bc08c5..c6de8933b8 100644 --- a/nemo_rl/experience/payload.py +++ b/nemo_rl/experience/payload.py @@ -22,6 +22,10 @@ from tensordict import TensorDict from nemo_rl.data.interfaces import LLMMessageLogType, VLMMessageLogType +from nemo_rl.data.multimodal_utils import ( + encode_multimodal_for_wire, + multimodal_row_tags, +) from nemo_rl.data_plane.codec import pack_jagged_fields from nemo_rl.data_plane.column_io import TOKEN_ALIGNED_FIELDS from nemo_rl.data_plane.schema import ( @@ -105,10 +109,12 @@ def record_to_train_batch( flags for configured advantage penalties. Returns: - BatchedDataDict with input_ids, input_lengths, generation_logprobs, - token_mask, an all-ones sample_mask, the raw mask_sample and truncated - flags, prompt_ids_for_adv, total_reward, violation counts, and optional - routed experts and message-violation masks. + BatchedDataDict with input IDs and lengths, generation log probabilities, + token and prompt-level sample masks, raw ``mask_sample`` and ``truncated`` + flags, prompt IDs for advantage computation, rewards, and violation + counts. Optional fields include routed experts, message-violation masks, + and any packed or per-token multimodal model inputs carried by the + completions. """ # Lazy imports: grpo and llm_message_utils transitively pull # experience.rollouts, so importing at module top risks a cycle. @@ -157,7 +163,7 @@ def record_to_train_batch( ) mask_sample = _mask_sample_flags(c.env_extras for c in completions) truncated = torch.tensor([c.truncated for c in completions], dtype=torch.bool) - sample_mask = torch.ones(n, dtype=torch.float32) + sample_mask = torch.full((n,), float(record.loss_multiplier), dtype=torch.float32) train_data: dict[str, Any] = { "input_ids": flat["token_ids"], @@ -176,6 +182,7 @@ def record_to_train_batch( if include_message_violation_fields: train_data[INVALID_TOOL_CALL_MASK] = flat[INVALID_TOOL_CALL_MASK] train_data[MALFORMED_THINKING_MASK] = flat[MALFORMED_THINKING_MASK] + train_data.update(flat.get_multimodal_dict(as_tensors=False)) return BatchedDataDict[Any](train_data) @@ -195,8 +202,11 @@ def pack_payload( prompt_idx: Stable dataset prompt index stamped on every row's tag. Returns: - sample_ids of the form {group_id}_g{i}, a jagged-packed TensorDict, and per-row - tags carrying weight_version plus any per-row violation counts. + Sample IDs of the form ``{group_id}_g{i}``, a jagged-packed TensorDict + containing tensor fields and encoded multimodal wire fields, and + per-row tags. Tags carry the weight version, prompt index, violation + counts, and ``__row_shapes`` metadata required to reconstruct + packed multimodal rows. """ lengths = train_batch["input_lengths"] n = int(lengths.shape[0]) @@ -206,16 +216,23 @@ def pack_payload( if isinstance(v, torch.Tensor) or (isinstance(v, np.ndarray) and v.dtype == object) } + multimodal = BatchedDataDict[Any](train_batch).get_multimodal_dict(as_tensors=False) + for key, value in multimodal.items(): + wire_value = encode_multimodal_for_wire(key, value) + if wire_value is not None: + tensor_fields[key] = wire_value fields_td = pack_jagged_fields( tensor_fields, lengths=lengths, token_aligned_fields=TOKEN_ALIGNED_FIELDS ) sample_ids = [f"{group_id}_g{i}" for i in range(n)] violations = train_batch.get(_VIOLATION_COUNTS_KEY, [{}] * n) + multimodal_tags = multimodal_row_tags(multimodal, n) or [{} for _ in range(n)] tags = [ { "weight_version": weight_version, "prompt_idx": prompt_idx, **violations[i], + **multimodal_tags[i], } for i in range(n) ] diff --git a/nemo_rl/experience/rollout_manager.py b/nemo_rl/experience/rollout_manager.py index b2e16b7cd0..e5b5286fba 100644 --- a/nemo_rl/experience/rollout_manager.py +++ b/nemo_rl/experience/rollout_manager.py @@ -31,6 +31,7 @@ TQReplayBuffer, ) from nemo_rl.data.interfaces import DatumSpec, LLMMessageLogType +from nemo_rl.data.llm_message_utils import batched_message_log_to_flat_message from nemo_rl.data_plane.schema import MASK_SAMPLE from nemo_rl.distributed.batched_data_dict import BatchedDataDict from nemo_rl.environments.interfaces import EnvironmentInterface @@ -454,6 +455,7 @@ async def run_rollout( metadata={"task_name": input_sample["task_name"]}, completions=completions, rollout_metrics=rollout_metrics, + loss_multiplier=float(input_sample.get("loss_multiplier", 1.0)), ) async def _run_single_rollout( @@ -633,9 +635,14 @@ async def _generate_response( Returns: Tuple of (assistant_message, input_lengths, gen_metrics) """ - # Prepare generation input - input_ids = torch.cat([m["token_ids"] for m in message_log]).unsqueeze(0) - input_lengths = torch.tensor([input_ids.shape[1]], dtype=torch.int32) + # Flatten both tokens and model-ready multimodal inputs. Building this + # from token_ids alone leaves expanded media placeholders in the prompt + # without the pixel tensors Megatron needs to project. + flat_messages, input_lengths = batched_message_log_to_flat_message( + [message_log], + pad_value_dict={"token_ids": self._tokenizer.pad_token_id}, + ) + input_ids = flat_messages["token_ids"] generation_input_data = BatchedDataDict[GenerationDatumSpec]( { "input_ids": input_ids, @@ -643,6 +650,9 @@ async def _generate_response( "stop_strings": [stop_strings], } ) + generation_input_data.update( + flat_messages.get_multimodal_dict(as_tensors=False) + ) # Generate response # TODO: update generate_async to return a single item directly @@ -877,6 +887,7 @@ async def run_rollout( metadata={"task_name": "nemo_gym"}, completions=completions, rollout_metrics=rollout_metrics, + loss_multiplier=float(input_sample.get("loss_multiplier", 1.0)), ) def _validate_init_params(self) -> None: @@ -1926,6 +1937,7 @@ async def _generate_for_finalization_attempt( fallback_weight_version=start_version, prompt_idx=record.prompt_idx, mask_sample=mask_sample, + loss_multiplier=record.loss_multiplier, ) from nemo_rl.experience.rollout_reassembler_actor import ( assert_metadata_only, diff --git a/nemo_rl/experience/rollout_reassembler.py b/nemo_rl/experience/rollout_reassembler.py index 551586d8c7..2b5818c8a2 100644 --- a/nemo_rl/experience/rollout_reassembler.py +++ b/nemo_rl/experience/rollout_reassembler.py @@ -364,6 +364,7 @@ def finalize_group( mask_sample: list[bool], fallback_weight_version: int, prompt_idx: int, + loss_multiplier: float = 1.0, ) -> FinalizedGroup: """Publish exactly N canonical rows for one prompt group. @@ -374,7 +375,9 @@ def finalize_group( advantage-stage flag the native ``pack_payload`` path emits from each ``Completion``; it rides along unchanged so the train pump's environment masking reads the same field on both paths (placeholder - rows already train nothing through ``sample_mask`` 0). ``truncated`` + rows already train nothing through ``sample_mask`` 0). + ``loss_multiplier`` supplies the dataset-level weight for every valid + row, matching the ordinary ``record_to_train_batch`` path. ``truncated`` is not carried from the dispatcher -- the receipt path has no real tokens to measure it from at dispatch time -- so it is computed here instead, from each row's rebuilt length against ``max_seq_len``. @@ -511,7 +514,7 @@ def finalize_group( input_ids[i, :length] = torch.tensor(row.token_ids, dtype=torch.int64) token_mask[i, :length] = torch.tensor(row.token_mask, dtype=torch.float32) logprobs[i, :length] = torch.tensor(row.logprobs, dtype=torch.float32) - sample_mask[i] = 1.0 + sample_mask[i] = float(loss_multiplier) train_batch = { "input_ids": input_ids, diff --git a/nemo_rl/experience/rollout_reassembler_actor.py b/nemo_rl/experience/rollout_reassembler_actor.py index 61505f49d0..57d91953d2 100644 --- a/nemo_rl/experience/rollout_reassembler_actor.py +++ b/nemo_rl/experience/rollout_reassembler_actor.py @@ -63,6 +63,8 @@ class ReassemblyRequest: # train pump reads the same ``mask_sample`` field as the native path # (SingleController reads it unconditionally). mask_sample: tuple[bool, ...] + # Dataset-level loss weight shared by every completion in this prompt group. + loss_multiplier: float = 1.0 @dataclass(frozen=True) @@ -152,6 +154,7 @@ def finalize(self, request: ReassemblyRequest) -> FinalizedGroup: mask_sample=list(request.mask_sample), fallback_weight_version=request.fallback_weight_version, prompt_idx=request.prompt_idx, + loss_multiplier=request.loss_multiplier, ) assert_metadata_only(result) return result diff --git a/nemo_rl/models/generation/megatron/config.py b/nemo_rl/models/generation/megatron/config.py index 98d0733a80..eb3a41f025 100644 --- a/nemo_rl/models/generation/megatron/config.py +++ b/nemo_rl/models/generation/megatron/config.py @@ -44,6 +44,9 @@ class MCoreGenerationSpecificArgs(TypedDict): # - 'block': graphs are owned at the enclosing block (TransformerBlock / HybridBlock). # Only meaningful when cuda_graph_impl='local'. inference_cuda_graph_scope: NotRequired[str] + # Required for EP>1 + inference CUDA graphs, except when using the + # `inference_optimized` transformer implementation. + moe_pad_experts_for_cuda_graph_inference: NotRequired[bool] materialize_only_last_token_logits: bool enable_chunked_prefill: bool @@ -91,8 +94,6 @@ class MCoreGenerationSpecificArgs(TypedDict): # FP8/MXFP8 for the dedicated (non-colocated) inference model; # merged into its `megatron_cfg` by `merged_inference_megatron_cfg`. fp8_cfg: NotRequired[Fp8Config] - # Merged into megatron_cfg for gen workers; required for EP>1 + local CUDA graphs. - moe_pad_experts_for_cuda_graph_inference: NotRequired[bool] class MCoreGenerationConfig(GenerationConfig): diff --git a/nemo_rl/models/policy/workers/megatron_policy_worker.py b/nemo_rl/models/policy/workers/megatron_policy_worker.py index 3b1dae699f..6a83b4036d 100644 --- a/nemo_rl/models/policy/workers/megatron_policy_worker.py +++ b/nemo_rl/models/policy/workers/megatron_policy_worker.py @@ -1455,6 +1455,9 @@ def train_microbatch( explicitly in ``finish_train_step``. Returns nothing: gradients land in ``param.main_grad`` and per-microbatch metrics accumulate in the open-step state until ``finish_train_step`` surfaces them. + + Multimodal validity-mask and model-owned packing/CP behavior match the + regular ``train`` path. """ state = self._assert_step_open() try: diff --git a/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni.sh b/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni.sh index c8f916ef05..023f025775 100755 --- a/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni.sh +++ b/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni.sh @@ -20,6 +20,20 @@ PROJECT_ROOT=$(realpath "${SCRIPT_DIR}/../..") cd "${PROJECT_ROOT}" +# run_test [fast] +# - "run_test fast " = always runs (both fast and full modes) +# - "run_test " = only runs in full mode; skipped when FAST=1 +run_test() { + if [[ "$1" == "fast" ]]; then + shift + time "$@" + elif [[ "${FAST:-0}" == "1" ]]; then + echo "FAST: Skipping: $*" + else + time "$@" + fi +} + GPU_COUNT=$(nvidia-smi --query-gpu=index --format=csv,noheader | wc -l) if (( GPU_COUNT < 2 )); then echo "SKIP: Nemotron Omni functional tests require at least two GB200 GPUs" @@ -29,8 +43,8 @@ fi # Both tests colocate TP2/EP2 training and generation on two GB200 GPUs. export CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES:-0,1}" -time uv run --no-sync bash ./tests/functional/nemotron_omni_clevr_megatron_1n2g.sh -time uv run --no-sync bash ./tests/functional/nemotron_omni_gym_video_megatron_1n2g.sh +run_test fast uv run --no-sync bash ./tests/functional/nemotron_omni_clevr_megatron_1n2g.sh +run_test fast uv run --no-sync bash ./tests/functional/nemotron_omni_gym_video_megatron_1n2g.sh cd "${PROJECT_ROOT}/tests" if compgen -G ".coverage*" > /dev/null; then diff --git a/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni_Single_Controller.sh b/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni_Single_Controller.sh new file mode 100755 index 0000000000..36f15552dd --- /dev/null +++ b/tests/functional/L1_Functional_Tests_GB200_Megatron_Omni_Single_Controller.sh @@ -0,0 +1,53 @@ +#!/bin/bash +# Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +set -xeuo pipefail + +SCRIPT_DIR=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" &> /dev/null && pwd) +PROJECT_ROOT=$(realpath "${SCRIPT_DIR}/../..") + +cd "${PROJECT_ROOT}" + +# run_test [fast] +# - "run_test fast " = always runs (both fast and full modes) +# - "run_test " = only runs in full mode; skipped when FAST=1 +run_test() { + if [[ "$1" == "fast" ]]; then + shift + time "$@" + elif [[ "${FAST:-0}" == "1" ]]; then + echo "FAST: Skipping: $*" + else + time "$@" + fi +} + +GPU_COUNT=$(nvidia-smi --query-gpu=index --format=csv,noheader | wc -l) +if (( GPU_COUNT < 2 )); then + echo "SKIP: Nemotron Omni SingleController functional tests require at least two GB200 GPUs" + exit 0 +fi + +# SingleController is non-colocated: one GPU trains the frozen-decoder policy +# and one GPU hosts Megatron generation. +export CUDA_VISIBLE_DEVICES="${CUDA_VISIBLE_DEVICES:-0,1}" + +run_test fast uv run --no-sync bash ./tests/functional/nemotron_omni_clevr_megatron_single_controller_1n2g.sh +run_test fast uv run --no-sync bash ./tests/functional/nemotron_omni_gym_video_megatron_single_controller_1n2g.sh + +cd "${PROJECT_ROOT}/tests" +if compgen -G ".coverage*" > /dev/null; then + coverage combine .coverage* +fi diff --git a/tests/functional/nemotron_omni_clevr_megatron_single_controller_1n2g.sh b/tests/functional/nemotron_omni_clevr_megatron_single_controller_1n2g.sh new file mode 100755 index 0000000000..9de2f6ba70 --- /dev/null +++ b/tests/functional/nemotron_omni_clevr_megatron_single_controller_1n2g.sh @@ -0,0 +1,174 @@ +#!/usr/bin/env bash +# Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved. + +set -euo pipefail + +SCRIPT_DIR=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd) +PROJECT_ROOT=$(realpath "${SCRIPT_DIR}/../..") + +if [[ -z "${HF_TOKEN:-}" ]]; then + echo "SKIP: HF_TOKEN is required for the Omni checkpoint" + exit 0 +fi + +GPU_COUNT=$(nvidia-smi --query-gpu=index --format=csv,noheader | wc -l) +if (( GPU_COUNT < 2 )); then + echo "SKIP: Omni CLEVR SingleController smoke requires at least two visible GPUs" + exit 0 +fi +DETECTED_CUDA_ARCH=$(nvidia-smi --query-gpu=compute_cap --format=csv,noheader -i 0) +export TORCH_CUDA_ARCH_LIST="${TORCH_CUDA_ARCH_LIST:-${DETECTED_CUDA_ARCH}}" +MEGATRON_TRANSFORMER_IMPL="${MEGATRON_TRANSFORMER_IMPL:-inference_optimized}" +MEGATRON_CUDA_GRAPH_IMPL="${MEGATRON_CUDA_GRAPH_IMPL:-local}" +if [[ "${MEGATRON_CUDA_GRAPH_IMPL}" == "local" ]]; then + INFERENCE_CUDA_GRAPH_SCOPE=block + NUM_CUDA_GRAPHS=-1 +else + INFERENCE_CUDA_GRAPH_SCOPE=none + NUM_CUDA_GRAPHS=0 +fi +if [[ "${MEGATRON_TRANSFORMER_IMPL}" != "inference_optimized" && + "${MEGATRON_CUDA_GRAPH_IMPL}" == "local" ]]; then + MOE_PAD_EXPERTS_FOR_CG=true +else + MOE_PAD_EXPERTS_FOR_CG=false +fi + +EXP_NAME=$(basename "$0" .sh) +EXP_DIR="${SCRIPT_DIR}/${EXP_NAME}" +LOG_DIR="${EXP_DIR}/logs" +DATA_ROOT="${EXP_DIR}/data" +TRAIN_PATH="${DATA_ROOT}/train.jsonl" +VAL_PATH="${DATA_ROOT}/val.jsonl" +JSON_METRICS="${EXP_DIR}/metrics.json" +RUN_LOG="${EXP_DIR}/run.log" +rm -rf "${EXP_DIR}" +mkdir -p "${LOG_DIR}" "${DATA_ROOT}" + +cd "${PROJECT_ROOT}" +export PYTHONPATH="${PROJECT_ROOT}:${PYTHONPATH:-}" + +# Match the non-SingleController L1 fixture. +TRAIN_PATH="${TRAIN_PATH}" VAL_PATH="${VAL_PATH}" uv run --no-sync python - <<'PY' +import base64 +import io +import json +import os + +from PIL import Image + +buffer = io.BytesIO() +Image.new("RGB", (224, 224), color="red").save(buffer, format="PNG") +image_url = "data:image/png;base64," + base64.b64encode(buffer.getvalue()).decode() + + +def sample(index: int) -> dict: + return { + "messages": [ + { + "role": "user", + "content": [ + {"type": "image", "image": image_url}, + { + "type": "text", + "text": f"Sample {index}: What color is the image?", + }, + ], + }, + {"role": "assistant", "content": "red"}, + ] + } + + +for path, count in ((os.environ["TRAIN_PATH"], 64), (os.environ["VAL_PATH"], 2)): + with open(path, "w") as output: + for index in range(count): + output.write(json.dumps(sample(index)) + "\n") +PY + +# SingleController requires disaggregated generation. One frozen-decoder model +# fits on each GB200, so split the two visible GPUs 1 trainer + 1 generator. +uv run --no-sync python examples/run_grpo_single_controller.py \ + --config examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron_generation.v1.yaml \ + cluster.num_nodes=1 \ + cluster.gpus_per_node=2 \ + ++cluster.segment_size=1 \ + policy.megatron_cfg.env_vars.TORCH_CUDA_ARCH_LIST="${TORCH_CUDA_ARCH_LIST}" \ + policy.megatron_cfg.tensor_model_parallel_size=1 \ + policy.megatron_cfg.expert_model_parallel_size=1 \ + policy.megatron_cfg.expert_tensor_parallel_size=1 \ + policy.megatron_cfg.context_parallel_size=1 \ + policy.megatron_cfg.sequence_parallel=true \ + policy.megatron_cfg.activation_checkpointing=true \ + ++policy.megatron_cfg.freeze_config.freeze_language_model=true \ + +policy.megatron_cfg.bias_dropout_fusion=false \ + policy.megatron_cfg.optimizer.optimizer_cpu_offload=false \ + policy.megatron_cfg.optimizer.optimizer_offload_fraction=0.0 \ + ++policy.megatron_cfg.optimizer.params_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.main_grads_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.main_params_dtype=float16 \ + ++policy.megatron_cfg.optimizer.exp_avg_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.exp_avg_sq_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.store_param_remainders=false \ + policy.generation.backend=megatron \ + policy.generation.colocated.enabled=false \ + policy.generation.colocated.resources.num_nodes=1 \ + policy.generation.colocated.resources.gpus_per_node=1 \ + policy.generation.max_new_tokens=128 \ + policy.generation.mcore_generation_config.tensor_model_parallel_size=1 \ + policy.generation.mcore_generation_config.expert_model_parallel_size=1 \ + policy.generation.mcore_generation_config.expert_tensor_parallel_size=1 \ + ++policy.generation.mcore_generation_config.context_parallel_size=1 \ + ++policy.generation.mcore_generation_config.moe_router_dtype=fp32 \ + policy.generation.mcore_generation_config.transformer_impl="${MEGATRON_TRANSFORMER_IMPL}" \ + policy.generation.mcore_generation_config.sequence_parallel=true \ + policy.generation.mcore_generation_config.refit_backend=nccl \ + policy.generation.mcore_generation_config.buffer_size_gb=2 \ + policy.generation.mcore_generation_config.cuda_graph_impl="${MEGATRON_CUDA_GRAPH_IMPL}" \ + policy.generation.mcore_generation_config.inference_cuda_graph_scope="${INFERENCE_CUDA_GRAPH_SCOPE}" \ + policy.generation.mcore_generation_config.num_cuda_graphs="${NUM_CUDA_GRAPHS}" \ + policy.generation.mcore_generation_config.use_cuda_graphs_for_non_decode_steps=false \ + policy.generation.mcore_generation_config.moe_pad_experts_for_cuda_graph_inference="${MOE_PAD_EXPERTS_FOR_CG}" \ + policy.generation.mcore_generation_config.enable_chunked_prefill=true \ + ++policy.generation.mcore_generation_config.async_sched_mode=async \ + policy.generation.mcore_generation_config.max_model_len=1024 \ + policy.generation.mcore_generation_config.max_tokens=1024 \ + policy.max_total_sequence_length=1024 \ + data.train.dataset_name=ResponseDataset \ + ++data.train.data_path="${TRAIN_PATH}" \ + data.train.split=train \ + data.validation.dataset_name=ResponseDataset \ + ++data.validation.data_path="${VAL_PATH}" \ + data.validation.split=train \ + data.num_workers=0 \ + grpo.async_grpo=null \ + grpo.num_prompts_per_step=1 \ + grpo.num_generations_per_prompt=2 \ + grpo.max_num_steps=1 \ + grpo.val_period=0 \ + grpo.val_at_start=false \ + grpo.val_at_end=false \ + policy.train_global_batch_size=2 \ + policy.train_micro_batch_size=1 \ + ++data_plane.enabled=true \ + ++data_plane.impl=transfer_queue \ + ++data_plane.backend=simple \ + ++data_plane.claim_meta_poll_interval_s=0.5 \ + ++data_plane.simple.num_storage_units=2 \ + ++async_rl.sampler.name=in_order \ + ++async_rl.sampler.max_lookahead_versions=1 \ + ++async_rl.recompute_kv_cache_after_weight_updates=false \ + ++async_rl.min_groups_for_streaming_train=1 \ + ++async_rl.max_inflight_prompts=2 \ + ++async_rl.max_buffered_rollouts=2 \ + logger.tensorboard_enabled=true \ + logger.log_dir="${LOG_DIR}" \ + logger.wandb_enabled=false \ + logger.monitor_gpus=false \ + checkpointing.enabled=false \ + "$@" 2>&1 | tee "${RUN_LOG}" + +uv run --no-sync tests/json_dump_tb_logs.py "${LOG_DIR}" --output_path "${JSON_METRICS}" +uv run --no-sync tests/check_metrics.py "${JSON_METRICS}" \ + 'max(data["train/gen_kl_error"]) < 0.05' \ + 'all_finite(data["train/reward"])' diff --git a/tests/functional/nemotron_omni_gym_video_megatron_1n2g.sh b/tests/functional/nemotron_omni_gym_video_megatron_1n2g.sh index 3fd6aad267..54d84dd1b0 100755 --- a/tests/functional/nemotron_omni_gym_video_megatron_1n2g.sh +++ b/tests/functional/nemotron_omni_gym_video_megatron_1n2g.sh @@ -133,9 +133,9 @@ uv run --no-sync python examples/nemo_gym/run_grpo_nemo_gym.py \ ++policy.generation.mcore_generation_config.video_temporal_patch_size=2 \ ++policy.generation.mcore_generation_config.video_target_num_patches=256 \ policy.max_total_sequence_length=1024 \ - +data.default.num_frames=8 \ - +data.default.video_sampling_style=nemotron_vl \ - +data.default.video_temporal_patch_size=2 \ + data.default.num_frames=8 \ + data.default.video_sampling_style=nemotron_vl \ + data.default.video_temporal_patch_size=2 \ +data.default.min_generation_tokens=128 \ data.default.video_target_num_patches=256 \ data.train.data_path="${TRAIN_PATH}" \ diff --git a/tests/functional/nemotron_omni_gym_video_megatron_single_controller_1n2g.sh b/tests/functional/nemotron_omni_gym_video_megatron_single_controller_1n2g.sh new file mode 100755 index 0000000000..86d24273d8 --- /dev/null +++ b/tests/functional/nemotron_omni_gym_video_megatron_single_controller_1n2g.sh @@ -0,0 +1,174 @@ +#!/usr/bin/env bash +# Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved. + +set -euo pipefail + +SCRIPT_DIR=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd) +PROJECT_ROOT=$(realpath "${SCRIPT_DIR}/../..") + +if [[ -z "${HF_TOKEN:-}" ]]; then + echo "SKIP: HF_TOKEN is required for the Omni checkpoint" + exit 0 +fi + +GPU_COUNT=$(nvidia-smi --query-gpu=index --format=csv,noheader | wc -l) +if (( GPU_COUNT < 2 )); then + echo "SKIP: Omni Gym-video SingleController smoke requires at least two GPUs" + exit 0 +fi +DETECTED_CUDA_ARCH=$(nvidia-smi --query-gpu=compute_cap --format=csv,noheader -i 0) +export TORCH_CUDA_ARCH_LIST="${TORCH_CUDA_ARCH_LIST:-${DETECTED_CUDA_ARCH}}" +MEGATRON_TRANSFORMER_IMPL="${MEGATRON_TRANSFORMER_IMPL:-inference_optimized}" +MEGATRON_CUDA_GRAPH_IMPL="${MEGATRON_CUDA_GRAPH_IMPL:-local}" +if [[ "${MEGATRON_CUDA_GRAPH_IMPL}" == "local" ]]; then + INFERENCE_CUDA_GRAPH_SCOPE=block + NUM_CUDA_GRAPHS=-1 +else + INFERENCE_CUDA_GRAPH_SCOPE=none + NUM_CUDA_GRAPHS=0 +fi +if [[ "${MEGATRON_TRANSFORMER_IMPL}" != "inference_optimized" && + "${MEGATRON_CUDA_GRAPH_IMPL}" == "local" ]]; then + MOE_PAD_EXPERTS_FOR_CG=true +else + MOE_PAD_EXPERTS_FOR_CG=false +fi + +EXP_NAME=$(basename "$0" .sh) +EXP_DIR="${SCRIPT_DIR}/${EXP_NAME}" +LOG_DIR="${EXP_DIR}/logs" +DATA_ROOT="${EXP_DIR}/data" +VIDEO_PATH="${DATA_ROOT}/red.mp4" +RAW_TRAIN_PATH="${DATA_ROOT}/train-raw.jsonl" +RAW_VAL_PATH="${DATA_ROOT}/val-raw.jsonl" +TRAIN_PATH="${DATA_ROOT}/train-gym.jsonl" +VAL_PATH="${DATA_ROOT}/val-gym.jsonl" +JSON_METRICS="${EXP_DIR}/metrics.json" +RUN_LOG="${EXP_DIR}/run.log" +rm -rf "${EXP_DIR}" +mkdir -p "${LOG_DIR}" "${DATA_ROOT}" + +cd "${PROJECT_ROOT}" +export PYTHONPATH="${PROJECT_ROOT}:${PYTHONPATH:-}" +export NRL_VIDEO_BACKEND=torchcodec +export NRL_VIDEO_SAMPLING_STYLE=nemotron_vl +export NRL_VIDEO_TEMPORAL_PATCH_SIZE=2 + +bash tools/install_audio_deps.sh +ffmpeg -hide_banner -loglevel error -y \ + -f lavfi -i color=c=red:s=224x224:r=8:d=2 \ + -c:v libx264 -pix_fmt yuv420p "${VIDEO_PATH}" + +for sample_id in $(seq 1 64); do + jq -nc \ + --arg prompt "Sample ${sample_id}: What color fills the video? A. Red B. Blue" \ + --arg video "${VIDEO_PATH}" \ + '{prompt: $prompt, video: $video, answer: "A", verifier: "mcqa"}' +done > "${RAW_TRAIN_PATH}" +for sample_id in $(seq 1 2); do + jq -nc \ + --arg prompt "Validation ${sample_id}: What color fills the video? A. Red B. Blue" \ + --arg video "${VIDEO_PATH}" \ + '{prompt: $prompt, video: $video, answer: "A", verifier: "mcqa"}' +done > "${RAW_VAL_PATH}" + +uv run --no-sync examples/nemo_gym/prepare_video_dataset.py convert \ + --input "${RAW_TRAIN_PATH}" \ + --output "${TRAIN_PATH}" +uv run --no-sync examples/nemo_gym/prepare_video_dataset.py convert \ + --input "${RAW_VAL_PATH}" \ + --output "${VAL_PATH}" + +# SingleController requires disaggregated generation. One frozen-decoder model +# fits on each GB200, so split the two visible GPUs 1 trainer + 1 generator. +uv run --no-sync python examples/run_grpo_single_controller.py \ + --config examples/configs/recipes/vlm/vlm_grpo-nemotron-omni-30ba3b-16n8g-megatron-tp4ep4-async-gym-video.v1.yaml \ + cluster.num_nodes=1 \ + cluster.gpus_per_node=2 \ + ++cluster.segment_size=1 \ + policy.megatron_cfg.env_vars.TORCH_CUDA_ARCH_LIST="${TORCH_CUDA_ARCH_LIST}" \ + policy.megatron_cfg.tensor_model_parallel_size=1 \ + policy.megatron_cfg.pipeline_model_parallel_size=1 \ + policy.megatron_cfg.expert_model_parallel_size=1 \ + policy.megatron_cfg.expert_tensor_parallel_size=1 \ + policy.megatron_cfg.context_parallel_size=1 \ + policy.megatron_cfg.sequence_parallel=true \ + policy.megatron_cfg.activation_checkpointing=true \ + ++policy.megatron_cfg.freeze_config.freeze_language_model=true \ + policy.megatron_cfg.optimizer.optimizer_cpu_offload=false \ + policy.megatron_cfg.optimizer.optimizer_offload_fraction=0.0 \ + ++policy.megatron_cfg.optimizer.params_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.main_grads_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.main_params_dtype=float16 \ + ++policy.megatron_cfg.optimizer.exp_avg_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.exp_avg_sq_dtype=bfloat16 \ + ++policy.megatron_cfg.optimizer.store_param_remainders=false \ + policy.generation.backend=megatron \ + ++policy.generation.bad_words=null \ + policy.generation.colocated.enabled=false \ + policy.generation.colocated.resources.num_nodes=1 \ + policy.generation.colocated.resources.gpus_per_node=1 \ + policy.generation.max_new_tokens=128 \ + policy.generation.mcore_generation_config.expose_http_server=true \ + policy.generation.mcore_generation_config.tensor_model_parallel_size=1 \ + policy.generation.mcore_generation_config.expert_model_parallel_size=1 \ + policy.generation.mcore_generation_config.expert_tensor_parallel_size=1 \ + ++policy.generation.mcore_generation_config.context_parallel_size=1 \ + ++policy.generation.mcore_generation_config.moe_router_dtype=fp32 \ + policy.generation.mcore_generation_config.transformer_impl="${MEGATRON_TRANSFORMER_IMPL}" \ + policy.generation.mcore_generation_config.sequence_parallel=true \ + policy.generation.mcore_generation_config.refit_backend=nccl \ + policy.generation.mcore_generation_config.buffer_size_gb=2 \ + policy.generation.mcore_generation_config.cuda_graph_impl="${MEGATRON_CUDA_GRAPH_IMPL}" \ + policy.generation.mcore_generation_config.inference_cuda_graph_scope="${INFERENCE_CUDA_GRAPH_SCOPE}" \ + policy.generation.mcore_generation_config.num_cuda_graphs="${NUM_CUDA_GRAPHS}" \ + policy.generation.mcore_generation_config.use_cuda_graphs_for_non_decode_steps=false \ + ++policy.generation.mcore_generation_config.moe_pad_experts_for_cuda_graph_inference="${MOE_PAD_EXPERTS_FOR_CG}" \ + policy.generation.mcore_generation_config.enable_chunked_prefill=true \ + ++policy.generation.mcore_generation_config.async_sched_mode=async \ + policy.generation.mcore_generation_config.enable_prefix_caching=true \ + policy.generation.mcore_generation_config.max_model_len=1024 \ + policy.generation.mcore_generation_config.max_tokens=1024 \ + ++policy.generation.mcore_generation_config.video_num_frames=8 \ + ++policy.generation.mcore_generation_config.video_temporal_patch_size=2 \ + ++policy.generation.mcore_generation_config.video_target_num_patches=256 \ + policy.max_total_sequence_length=1024 \ + data.default.num_frames=8 \ + data.default.video_sampling_style=nemotron_vl \ + data.default.video_temporal_patch_size=2 \ + +data.default.min_generation_tokens=128 \ + data.default.video_target_num_patches=256 \ + data.train.data_path="${TRAIN_PATH}" \ + data.validation.data_path="${VAL_PATH}" \ + grpo.deduplicate_multimodal_data=false \ + grpo.async_grpo=null \ + grpo.num_prompts_per_step=1 \ + grpo.num_generations_per_prompt=2 \ + grpo.max_num_steps=1 \ + grpo.val_period=0 \ + grpo.val_at_start=false \ + grpo.val_at_end=false \ + policy.train_global_batch_size=2 \ + policy.train_micro_batch_size=1 \ + ++data_plane.enabled=true \ + ++data_plane.impl=transfer_queue \ + ++data_plane.backend=simple \ + ++data_plane.claim_meta_poll_interval_s=0.5 \ + ++data_plane.simple.num_storage_units=2 \ + ++async_rl.sampler.name=in_order \ + ++async_rl.sampler.max_lookahead_versions=1 \ + ++async_rl.recompute_kv_cache_after_weight_updates=false \ + ++async_rl.min_groups_for_streaming_train=1 \ + ++async_rl.max_inflight_prompts=2 \ + ++async_rl.max_buffered_rollouts=2 \ + logger.tensorboard_enabled=true \ + logger.log_dir="${LOG_DIR}" \ + logger.wandb_enabled=false \ + logger.monitor_gpus=false \ + checkpointing.enabled=false \ + "$@" 2>&1 | tee "${RUN_LOG}" + +uv run --no-sync tests/json_dump_tb_logs.py "${LOG_DIR}" --output_path "${JSON_METRICS}" +uv run --no-sync tests/check_metrics.py "${JSON_METRICS}" \ + 'max(data["train/gen_kl_error"]) < 0.05' \ + 'all_finite(data["train/reward"])' diff --git a/tests/test_suites/nightly_gb200.txt b/tests/test_suites/nightly_gb200.txt index 9ab02ba24c..d2c8c9acf6 100644 --- a/tests/test_suites/nightly_gb200.txt +++ b/tests/test_suites/nightly_gb200.txt @@ -32,6 +32,7 @@ tests/test_suites/llm/grpo-qwen3-30ba3b-4n4g-megatron-qa-nvfp4-w4a4-real.sh tests/test_suites/vlm/vlm_grpo-qwen2.5-vl-3b-instruct-clevr-1n4g-dtensor2tp1.v1.sh tests/test_suites/vlm/vlm_grpo-qwen2.5-vl-3b-instruct-clevr-1n4g-megatrontp1.v1.sh tests/test_suites/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron_generation.v1.sh +tests/test_suites/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.sh # Deepscaler (short tests) tests/test_suites/llm/grpo-deepscaler-1.5b-1n4g-8K.sh diff --git a/tests/test_suites/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.sh b/tests/test_suites/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.sh new file mode 100755 index 0000000000..d0b8676b1a --- /dev/null +++ b/tests/test_suites/vlm/vlm_grpo-nemotron-omni-30ba3b-clevr-8n4g-megatron-single-controller-async.v1.sh @@ -0,0 +1,58 @@ +#!/bin/bash +# Copyright (c) 2026, NVIDIA CORPORATION. All rights reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +SCRIPT_DIR=$( cd -- "$( dirname -- "${BASH_SOURCE[0]}" )" &> /dev/null && pwd) +source "$SCRIPT_DIR/common.env" + +# Match the non-SingleController convergence run while exercising the +# SingleController/TransferQueue orchestration path. + +# ===== BEGIN CONFIG ===== +NUM_NODES=8 +GPUS_PER_NODE=4 +SEGMENT_SIZE=2 +STEPS_PER_RUN=50 +MAX_STEPS=50 +NUM_RUNS=$(( (MAX_STEPS + STEPS_PER_RUN - 1) / STEPS_PER_RUN )) +NUM_MINUTES=120 +# ===== END CONFIG ===== + +exit_if_max_steps_reached + +cd "$PROJECT_ROOT" + +uv run examples/run_grpo_single_controller.py \ + --config "$CONFIG_PATH" \ + grpo.max_num_steps="$MAX_STEPS" \ + policy.megatron_cfg.scheduler.lr_warmup_iters=10 \ + logger.log_dir="$LOG_DIR" \ + logger.wandb_enabled=True \ + logger.wandb.project=nemo-rl \ + logger.wandb.name="$EXP_NAME" \ + logger.monitor_gpus=True \ + logger.tensorboard_enabled=True \ + checkpointing.enabled=False \ + checkpointing.checkpoint_dir="$CKPT_DIR" \ + "$@" \ + 2>&1 | tee "$RUN_LOG" + +uv run tests/json_dump_tb_logs.py "$LOG_DIR" --output_path "$JSON_METRICS" + +uv run tests/check_metrics.py "$JSON_METRICS" \ + 'all_finite(data["train/loss"])' \ + 'all_finite(data["train/grad_norm"])' \ + 'min(data["train/grad_norm"]) > 0' \ + 'all_finite(data["train/token_mult_prob_error"])' \ + 'mean(data["train/reward"], range_start=-10) > 0.6' diff --git a/tests/unit/data/datasets/test_mmpr_tiny.py b/tests/unit/data/datasets/test_mmpr_tiny.py index 7ec4c34a53..d4e0da4b96 100644 --- a/tests/unit/data/datasets/test_mmpr_tiny.py +++ b/tests/unit/data/datasets/test_mmpr_tiny.py @@ -307,7 +307,7 @@ def test_historical_tiled_processor_gets_media_metadata(self, tiny_image_path): torch.tensor([[224, 224], [224, 224], [224, 224]]), ) assert torch.equal( - user_message["num_frames"].as_tensor(), torch.ones(3, dtype=torch.long) + user_message["num_frames"].as_tensor(), torch.ones(3, dtype=torch.int32) ) def test_prompted_text_contains_boxed_literal_and_no_raw_dataset_string( diff --git a/tests/unit/data_plane/test_kvbatchmeta.py b/tests/unit/data_plane/test_kvbatchmeta.py index a8dc3bc822..49272eaea5 100644 --- a/tests/unit/data_plane/test_kvbatchmeta.py +++ b/tests/unit/data_plane/test_kvbatchmeta.py @@ -247,6 +247,34 @@ def test_tags_none_when_either_side_missing_in_concat(): assert with_tags.concat(without).tags is None +def test_concat_unions_payload_fields_in_first_seen_order(): + text = KVBatchMeta( + partition_id="p", + task_name="train", + sample_ids=["a"], + fields=["input_ids", "input_lengths"], + ) + multimodal = KVBatchMeta( + partition_id="p", + task_name="train", + sample_ids=["b"], + fields=[ + "input_ids", + "pixel_values", + "image_grid_thw", + ], + ) + + joined = text.concat(multimodal) + + assert joined.fields == [ + "input_ids", + "input_lengths", + "pixel_values", + "image_grid_thw", + ] + + # ── Realistic tags from the rollout-shapes helper ── diff --git a/tests/unit/data_plane/test_rollout_reassembler.py b/tests/unit/data_plane/test_rollout_reassembler.py index bd6b965c2f..27816aad80 100644 --- a/tests/unit/data_plane/test_rollout_reassembler.py +++ b/tests/unit/data_plane/test_rollout_reassembler.py @@ -215,6 +215,7 @@ def test_finalize_group_publishes_n_rows_with_placeholder(tq_client, partitions) mask_sample=[True, False], fallback_weight_version=9, prompt_idx=0, + loss_multiplier=0.25, ) assert not finalized.dropped assert finalized.meta is not None @@ -229,7 +230,7 @@ def test_finalize_group_publishes_n_rows_with_placeholder(tq_client, partitions) rows = _fetch_rows(tq_client, rollout_ids) sample_mask = torch.as_tensor(rows["sample_mask"]).flatten() - assert sample_mask.tolist() == [1.0, 0.0] + assert sample_mask.tolist() == [0.25, 0.0] input_ids = torch.as_tensor(rows["input_ids"][0]).flatten() assert input_ids[:valid_len].tolist() == expected.token_ids # Placeholder borrows the valid sibling's prompt for baseline grouping. diff --git a/tests/unit/distributed/test_virtual_cluster.py b/tests/unit/distributed/test_virtual_cluster.py index 4a17a69fcf..a832b69e5a 100644 --- a/tests/unit/distributed/test_virtual_cluster.py +++ b/tests/unit/distributed/test_virtual_cluster.py @@ -283,6 +283,31 @@ def test_maybe_configure_data_plane_env_then_init_ray_threads_env_vars(): assert env_vars["MC_ENABLE_DEST_DEVICE_AFFINITY"] == "1" +def test_init_ray_adds_hf_modules_cache_to_cluster_pythonpath(): + """Direct actors must import trust_remote_code classes while unpickling.""" + from nemo_rl.distributed.virtual_cluster import init_ray + + with ( + patch("ray.init") as mock_ray_init, + patch("ray.cluster_resources") as mock_cluster_resources, + ): + mock_cluster_resources.return_value = {"nrl_tag_0": 1} + env = { + "CUDA_VISIBLE_DEVICES": "0", + "HF_MODULES_CACHE": "/hf/modules", + "PYTHONPATH": "/project", + } + with patch.dict(os.environ, env, clear=True): + init_ray() + + env_vars = mock_ray_init.call_args_list[0][1]["runtime_env"]["env_vars"] + assert env_vars["HF_MODULES_CACHE"] == "/hf/modules" + assert env_vars["PYTHONPATH"].split(os.pathsep) == [ + "/hf/modules", + "/project", + ] + + def test_init_ray_alone_has_no_data_plane_awareness(): """Every non-data-plane launcher's call (bare init_ray(), no preceding maybe_configure_data_plane_env) must not touch mooncake env vars -- diff --git a/tests/unit/environments/test_nemo_gym.py b/tests/unit/environments/test_nemo_gym.py index 51cb57c6f0..0c57afd293 100644 --- a/tests/unit/environments/test_nemo_gym.py +++ b/tests/unit/environments/test_nemo_gym.py @@ -507,6 +507,7 @@ def apply_chat_template(self, messages, *, tokenize, **kwargs): assert datum is not None user_message = datum["message_log"][0] assert user_message["num_frames"].as_tensor().tolist() == [4] + assert user_message["num_frames"].as_tensor().dtype == torch.int32 assert user_message["imgs_sizes"].as_tensor().dtype == torch.int32 extra_env_info = datum["extra_env_info"] outbound_content = extra_env_info["responses_create_params"]["input"][0]["content"] diff --git a/tests/unit/experience/test_payload.py b/tests/unit/experience/test_payload.py index 49de9e044c..0ba33dd432 100644 --- a/tests/unit/experience/test_payload.py +++ b/tests/unit/experience/test_payload.py @@ -16,6 +16,8 @@ import torch +from nemo_rl.data.multimodal_utils import PackedTensor +from nemo_rl.data_plane.codec import materialize from nemo_rl.data_plane.schema import ( INVALID_TOOL_CALL_MASK, MALFORMED_THINKING_MASK, @@ -82,7 +84,9 @@ def _completion( ) -def _record(completions: list[Completion]) -> PromptGroupRecord: +def _record( + completions: list[Completion], *, loss_multiplier: float = 1.0 +) -> PromptGroupRecord: return PromptGroupRecord( prompt_idx=0, prompt=[ @@ -96,6 +100,7 @@ def _record(completions: list[Completion]) -> PromptGroupRecord: metadata={"task_name": "test"}, completions=completions, rollout_metrics={}, + loss_multiplier=loss_multiplier, ) @@ -243,6 +248,63 @@ def test_record_to_train_batch_omits_routed_experts_when_absent() -> None: assert "routed_experts" not in fields +def test_multimodal_packed_tensor_round_trips_through_tq_payload() -> None: + completions = [ + _completion(route_start=10, reward=1.0, with_routes=False), + _completion(route_start=30, reward=2.0, with_routes=False), + ] + media = torch.arange(8, dtype=torch.float32).reshape(2, 4) + completions[0].message_log[0]["pixel_values"] = PackedTensor(media, dim_to_pack=0) + + train_batch = record_to_train_batch( + _record(completions), + pad_value_dict={"token_ids": 0, "input_ids": 0}, + include_message_violation_fields=False, + ) + assert isinstance(train_batch["pixel_values"], PackedTensor) + + _, fields, tags = pack_payload( + train_batch, + weight_version=3, + group_id="group", + prompt_idx=17, + ) + assert "pixel_values" in fields + assert tags[0]["pixel_values__row_shapes"]["shapes"] == [[2, 4]] + assert tags[1]["pixel_values__row_shapes"]["shapes"] == [] + + restored = materialize(fields, tags=tags) + restored_media = restored["pixel_values"] + assert isinstance(restored_media, PackedTensor) + assert len(restored_media) == 2 + assert restored_media.logical_segment_counts_by_row() == [1, 0] + assert torch.equal(restored_media.as_tensor(), media) + + +def test_per_token_multimodal_field_is_packed_with_sequence_lengths() -> None: + train_batch = { + "input_lengths": torch.tensor([3, 2], dtype=torch.int32), + "input_ids": torch.tensor([[10, 11, 12], [20, 21, 0]]), + "token_type_ids": torch.tensor([[0, 1, 1], [0, 1, 0]]), + } + + _, fields, tags = pack_payload( + train_batch, + weight_version=3, + group_id="group", + prompt_idx=17, + ) + + assert [row.tolist() for row in fields["token_type_ids"].unbind()] == [ + [0, 1, 1], + [0, 1], + ] + assert tags == [ + {"weight_version": 3, "prompt_idx": 17}, + {"weight_version": 3, "prompt_idx": 17}, + ] + + def test_record_to_train_batch_carries_raw_masks_without_applying_them() -> None: record = _record( [ @@ -287,6 +349,33 @@ def test_record_to_train_batch_carries_raw_masks_without_applying_them() -> None assert torch.equal(fields["truncated"], train_batch["truncated"]) +def test_record_to_train_batch_broadcasts_prompt_loss_multiplier() -> None: + record = _record( + [ + _completion(route_start=10, reward=1.0), + _completion(route_start=30, reward=2.0), + ], + loss_multiplier=0.25, + ) + + train_batch = record_to_train_batch( + record, + pad_value_dict={"token_ids": 0, "input_ids": 0}, + include_message_violation_fields=False, + ) + + expected = torch.full((2,), 0.25) + assert torch.equal(train_batch["sample_mask"], expected) + + _, fields, _ = pack_payload( + train_batch, + weight_version=3, + group_id="group", + prompt_idx=17, + ) + assert torch.equal(fields["sample_mask"], expected) + + def _failed_completion() -> Completion: """A trajectory whose first generation raised: prompt only, no routes.""" return Completion( diff --git a/tests/unit/experience/test_rollout_generation_failures.py b/tests/unit/experience/test_rollout_generation_failures.py index f6317d3a8e..3b19149495 100644 --- a/tests/unit/experience/test_rollout_generation_failures.py +++ b/tests/unit/experience/test_rollout_generation_failures.py @@ -88,6 +88,8 @@ def __init__(self, input_ids: torch.Tensor) -> None: class _FakeTokenizer: + pad_token_id = 0 + def decode(self, ids, skip_special_tokens=True): del skip_special_tokens return f"<{len(ids)} tokens>" diff --git a/tests/unit/experience/test_rollout_manager.py b/tests/unit/experience/test_rollout_manager.py index fe8726250e..ea7f78c9e2 100644 --- a/tests/unit/experience/test_rollout_manager.py +++ b/tests/unit/experience/test_rollout_manager.py @@ -29,6 +29,7 @@ import tempfile import uuid from copy import deepcopy +from types import SimpleNamespace import pytest import torch @@ -40,6 +41,7 @@ from nemo_rl.data.collate_fn import rl_collate_fn from nemo_rl.data.datasets.response_datasets import NemoGymDataset from nemo_rl.data.interfaces import DatumSpec +from nemo_rl.data.multimodal_utils import PackedTensor from nemo_rl.data.processors import nemo_gym_data_processor from nemo_rl.distributed.batched_data_dict import BatchedDataDict from nemo_rl.experience.interfaces import ( @@ -51,6 +53,7 @@ ) from nemo_rl.experience.rollout_manager import ( AsyncNemoGymRolloutImpl, + AsyncRolloutImpl, RolloutManager, RolloutOutcome, RolloutRetryPolicy, @@ -92,6 +95,71 @@ async def apply(): return _run(apply()) +def test_generate_response_forwards_message_log_media_to_generation() -> None: + captured: dict[str, BatchedDataDict] = {} + + class _Generation: + async def generate_async(self, data): + captured["data"] = data + input_len = int(data["input_lengths"][0]) + yield ( + 0, + BatchedDataDict( + { + "output_ids": torch.cat( + (data["input_ids"], torch.tensor([[42]])), dim=1 + ), + "unpadded_sequence_lengths": torch.tensor([input_len + 1]), + "logprobs": torch.zeros(1, input_len + 1), + } + ), + ) + + manager = object.__new__(AsyncRolloutImpl) + manager._policy_generation = _Generation() + manager._tokenizer = SimpleNamespace( + pad_token_id=0, + decode=lambda *_args, **_kwargs: "answer", + ) + manager._timeouts = SimpleNamespace(generation_s=10.0) + pixel_values = PackedTensor(torch.ones(2, 3, 4, 4), dim_to_pack=0) + imgs_sizes = PackedTensor(torch.tensor([[4, 4], [4, 4]]), dim_to_pack=0) + message_log = [ + { + "role": "user", + "content": "image", + "token_ids": torch.tensor([1, 2, 3]), + "pixel_values": pixel_values, + "imgs_sizes": imgs_sizes, + }, + { + "role": "assistant", + "content": "follow-up", + "token_ids": torch.tensor([4, 5]), + }, + ] + + assistant_message, input_lengths, _ = _run( + manager._generate_response(message_log, [""]) + ) + + generation_data = captured["data"] + assert generation_data["input_ids"].tolist() == [[1, 2, 3, 4, 5]] + assert generation_data["input_lengths"].tolist() == [5] + assert generation_data["stop_strings"] == [[""]] + assert isinstance(generation_data["pixel_values"], PackedTensor) + assert isinstance(generation_data["imgs_sizes"], PackedTensor) + assert torch.equal( + generation_data["pixel_values"].as_tensor(), pixel_values.as_tensor() + ) + assert torch.equal( + generation_data["imgs_sizes"].as_tensor(), imgs_sizes.as_tensor() + ) + assert input_lengths.tolist() == [5] + assert assistant_message["content"] == "answer" + assert assistant_message["token_ids"].tolist() == [42] + + class _FakeBuffer: """Minimal TQReplayBuffer stand-in that records reserve/commit calls.""" @@ -965,7 +1033,10 @@ def test_async_rollout_manager( - completions hold independent (not aliased) message_log objects """ vllm_generation, tokenizer, task_to_env, _, _ = multi_step_setup_vllm_async - input_sample = single_multi_step_calculator_input_sample + input_sample = { + **single_multi_step_calculator_input_sample, + "loss_multiplier": 0.25, + } num_generations = 2 max_seq_len = 1024 max_rollout_turns = input_sample["extra_env_info"]["max_steps"] + 1 @@ -989,6 +1060,7 @@ def test_async_rollout_manager( f"Expected {num_generations} completions, got {len(record.completions)}" ) assert record.prompt_idx == input_sample["idx"] + assert record.loss_multiplier == input_sample["loss_multiplier"] for i, completion in enumerate(record.completions): assert isinstance(completion, Completion) @@ -1243,6 +1315,7 @@ def test_async_nemo_gym_rollout_manager( f"Expected {num_generations} completions, got {len(record.completions)}" ) assert record.prompt_idx == 0 + assert record.loss_multiplier == single_prompt["loss_multiplier"] for i, completion in enumerate(record.completions): assert isinstance(completion, Completion) @@ -1444,7 +1517,9 @@ def reserve( ) -def _receipt_record(rollout_ids, receipts, instance_configs=None): +def _receipt_record( + rollout_ids, receipts, instance_configs=None, *, loss_multiplier=1.0 +): instance_configs = instance_configs or [None] * len(rollout_ids) completions = [ Completion( @@ -1467,6 +1542,7 @@ def _receipt_record(rollout_ids, receipts, instance_configs=None): metadata={"task_name": "nemo_gym"}, completions=completions, rollout_metrics={}, + loss_multiplier=loss_multiplier, ) @@ -1505,6 +1581,7 @@ async def run_rollout(self, _sample, *, rollout_ids=None): rollout_ids, [{"rollout_id": rid} for rid in rollout_ids], instance_configs=instance_configs, + loss_multiplier=float(_sample.get("loss_multiplier", 1.0)), ) mgr._impl = _CaptureImpl() @@ -1530,7 +1607,11 @@ def test_mints_ids_and_returns_metadata_request(self): buf = _FakeCaptureBuffer() mgr = _make_capture_manager(buf) - request = _run(mgr.generate_for_finalization({"prompt": "p"}, target_step=5)) + request = _run( + mgr.generate_for_finalization( + {"prompt": "p", "loss_multiplier": 0.25}, target_step=5 + ) + ) # Rollout ids were minted from the reserved group id and threaded # end to end: reserve -> impl -> metadata-only actor request. @@ -1543,6 +1624,7 @@ def test_mints_ids_and_returns_metadata_request(self): assert [r["rollout_id"] for r in request.receipts] == expected_ids assert request.rewards == (0.5, 0.5) assert request.mask_sample == (False, False) + assert request.loss_multiplier == 0.25 assert request.fallback_weight_version == 7 # Finalization and commit are exclusively owned by the controller's # actor-pool path; the manager leaves the reservation unready. diff --git a/tests/unit/experience/test_rollout_manager_router_replay.py b/tests/unit/experience/test_rollout_manager_router_replay.py index 6f5390364b..ec41573cc3 100644 --- a/tests/unit/experience/test_rollout_manager_router_replay.py +++ b/tests/unit/experience/test_rollout_manager_router_replay.py @@ -38,6 +38,8 @@ def _fallback_routes(count: int) -> torch.Tensor: class _FakeTokenizer: + pad_token_id = 0 + def decode(self, token_ids: torch.Tensor, skip_special_tokens: bool) -> str: del token_ids, skip_special_tokens return "generated" diff --git a/tests/unit/experience/test_rollout_reassembler_actor.py b/tests/unit/experience/test_rollout_reassembler_actor.py index e73e6d14eb..4a21747e22 100644 --- a/tests/unit/experience/test_rollout_reassembler_actor.py +++ b/tests/unit/experience/test_rollout_reassembler_actor.py @@ -15,7 +15,8 @@ from __future__ import annotations -from dataclasses import fields +from dataclasses import fields, replace +from unittest.mock import MagicMock import pytest import torch @@ -25,6 +26,7 @@ from nemo_rl.experience.rollout_reassembler_actor import ( _FORBIDDEN_RPC_KEYS, ReassemblyRequest, + RolloutReassemblerActor, assert_metadata_only, ) @@ -71,6 +73,34 @@ def test_finalizer_request_and_result_are_metadata_only() -> None: assert_metadata_only(result) +def test_finalize_forwards_loss_multiplier_to_reassembler() -> None: + actor_cls = RolloutReassemblerActor.__ray_metadata__.modified_class + actor = object.__new__(actor_cls) + actor._finalizer = MagicMock() + result = FinalizedGroup( + meta=None, + group_min_wv=4, + group_max_wv=4, + staging_keys=[], + dropped=True, + drop_reason="test", + ) + actor._finalizer.finalize_group.return_value = result + request = replace(_request(), loss_multiplier=0.25) + + assert actor.finalize(request) is result + actor._finalizer.finalize_group.assert_called_once_with( + "group", + ["group_g0"], + [request.receipts[0]], + [1.0], + mask_sample=[False], + fallback_weight_version=4, + prompt_idx=0, + loss_multiplier=0.25, + ) + + @pytest.mark.parametrize( "payload", [ @@ -100,6 +130,7 @@ def test_rpc_dataclass_fields_are_classified() -> None: "fallback_weight_version", "prompt_idx", "mask_sample", + "loss_multiplier", } assert {f.name for f in fields(FinalizedGroup)} == { "meta", diff --git a/tests/unit/models/generation/test_megatron_generation_parse.py b/tests/unit/models/generation/test_megatron_generation_parse.py index e177cf7491..40114f7509 100644 --- a/tests/unit/models/generation/test_megatron_generation_parse.py +++ b/tests/unit/models/generation/test_megatron_generation_parse.py @@ -236,6 +236,11 @@ def test_http_server_port_reservation(monkeypatch): assert holder._sock.fileno() == -1 assert reserved_socket.getsockname()[1] == port + # Still accepting after the holder closed its copy: the port was + # never released across the handoff. + with socket.create_connection(("127.0.0.1", port), timeout=5): + pass + # MCore closes the handed-off fd and gives every frontend replica # its own SO_REUSEPORT listener. Such a listener can join the reuse # group while this test stub still holds the adopted duplicate. diff --git a/tests/unit/models/policy/test_megatron_split_state.py b/tests/unit/models/policy/test_megatron_split_state.py index c3e023bc4e..990397b4bc 100644 --- a/tests/unit/models/policy/test_megatron_split_state.py +++ b/tests/unit/models/policy/test_megatron_split_state.py @@ -410,6 +410,27 @@ def test_finish_without_begin_raises(self, mock_module_symbols): class TestTrainMicrobatch: + def test_forwards_multimodal_iterator_capabilities(self, mock_module_symbols): + from nemo_rl.algorithms.loss.interfaces import LossType + + w = _make_worker(LossType.TOKEN_LEVEL) + w.media_placeholder_token_id = 42 + w.delegate_pack_to_model = True + w.delegate_mtp_loss_mask_to_model = True + batch = _fake_batch() + + with patch( + f"{WORKER_MOD}.attach_media_token_validity_mask" + ) as attach_validity_mask: + w.begin_train_step(loss_fn=w._test_loss_fn) + w.train_microbatch(batch) + + attach_validity_mask.assert_called_once_with(batch, 42) + iterator_kwargs = mock_module_symbols["gmi"].call_args.kwargs + assert iterator_kwargs["delegate_pack_to_model"] is True + assert iterator_kwargs["delegate_mtp_loss_mask_to_model"] is True + assert iterator_kwargs["model_slices_context_parallel_inputs"] is False + def test_wraps_forward_backward_in_no_sync(self, mock_module_symbols): """The single most important assertion in this file. Without the no_sync wrap, mcore DDP dispatches a per-call cross-DP reduce on diff --git a/tests/unit/models/value/test_dtensor_value_worker.py b/tests/unit/models/value/test_dtensor_value_worker.py index 5ca4129652..5ab1c429d5 100644 --- a/tests/unit/models/value/test_dtensor_value_worker.py +++ b/tests/unit/models/value/test_dtensor_value_worker.py @@ -546,7 +546,9 @@ def test_value_worker_train_decreases_loss(value_setup): losses.append(float(loss_tensor.mean().item())) value.finish_training() - assert losses[-1] <= losses[0] + 1e-3, ( + # This is a small fixed-batch smoke, not a convergence test. Allow minor + # optimizer/model-version jitter while still catching a meaningful loss jump. + assert losses[-1] <= losses[0] + 2e-3, ( f"Value loss should not increase after 3 steps; got {losses}" ) diff --git a/tests/unit/single_controller/test_entrypoint.py b/tests/unit/single_controller/test_entrypoint.py index c4436936db..c3b931e037 100644 --- a/tests/unit/single_controller/test_entrypoint.py +++ b/tests/unit/single_controller/test_entrypoint.py @@ -89,7 +89,7 @@ def main_context(monkeypatch: pytest.MonkeyPatch) -> SimpleNamespace: monkeypatch.setattr( run_grpo_single_controller, "setup_single_controller", - lambda *_args: (actor_args, SetupTimingMetrics()), + lambda *_args, **_kwargs: (actor_args, SetupTimingMetrics()), ) monkeypatch.setattr( run_grpo_single_controller.SingleControllerActor, @@ -170,3 +170,30 @@ def test_main_configures_generation_for_trained_mtp( assert ( main_context.config.policy["generation"] is main_context.configured_generation ) + + +def test_main_passes_processor_for_vlm( + main_context: SimpleNamespace, + monkeypatch: pytest.MonkeyPatch, +) -> None: + processor = SimpleNamespace(tokenizer="vlm-tokenizer") + get_tokenizer = MagicMock(return_value=processor) + setup_single_controller = MagicMock( + return_value=(main_context.actor_args, SetupTimingMetrics()) + ) + main_context.config.policy["is_vlm"] = True + monkeypatch.setattr(run_grpo_single_controller, "get_tokenizer", get_tokenizer) + monkeypatch.setattr( + run_grpo_single_controller, + "setup_single_controller", + setup_single_controller, + ) + + run_grpo_single_controller.main() + + get_tokenizer.assert_called_once_with( + main_context.config.policy["tokenizer"], get_processor=True + ) + setup_single_controller.assert_called_once_with( + main_context.config, "vlm-tokenizer", processor=processor + ) diff --git a/tests/unit/single_controller/test_setup.py b/tests/unit/single_controller/test_setup.py index 9b1ad8e5cc..25d695f5ad 100644 --- a/tests/unit/single_controller/test_setup.py +++ b/tests/unit/single_controller/test_setup.py @@ -59,6 +59,7 @@ from nemo_rl.algorithms.single_controller_utils.config import ( validate_single_controller_config, ) +from nemo_rl.data.multimodal_utils import WIRE_MULTIMODAL_FIELDS from nemo_rl.data_plane import DATA_PLANE_CHECKPOINT_SCHEMA_VERSION from nemo_rl.data_plane.schema import SC_ROLLOUT_SCHEMA_FIELDS from nemo_rl.experience.rollouts import EffortLevelsConfig @@ -400,8 +401,11 @@ def __init__(self, **kwargs): assert teacher_topology is None -def test_single_controller_mopd_recipe_resolves_to_runtime_contract(): +def test_single_controller_mopd_recipe_resolves_to_runtime_contract(monkeypatch): """The inherited recipe resolves exactly as the SC entrypoint consumes it.""" + # The parent recipe locates its fixture data below HF_HOME. This test only + # validates config resolution, so it needs a stable path, not real data. + monkeypatch.setenv("HF_HOME", "/tmp/nemo-rl-test-hf") register_omegaconf_resolvers() repo_root = Path(__file__).resolve().parents[3] recipe = repo_root / ( @@ -1045,16 +1049,37 @@ def test_env_handles_sourced_from_setup_response_data(self, patched_factories): """setup_response_data receives master_config.env and supplies env handles.""" math_env_cfg = {"some": "value"} mc = _make_master_config(env={"math": math_env_cfg}) + tokenizer = MagicMock(pad_token_id=0) - actor_args, _ = setup_single_controller(mc, MagicMock(pad_token_id=0)) + actor_args, _ = setup_single_controller(mc, tokenizer) - _, call_kwargs = patched_factories["setup_response_data"].call_args + call_args, call_kwargs = patched_factories["setup_response_data"].call_args + assert call_args[0] is tokenizer assert call_kwargs["env_configs"] == {"math": math_env_cfg} + assert call_kwargs["is_vlm"] is False assert actor_args.env_handles is patched_factories["env_handles"] + def test_vlm_processor_used_for_data_and_environment_setup(self, patched_factories): + mc = _make_master_config(env={"clevr-cogent": {"some": "value"}}) + tokenizer = MagicMock(pad_token_id=0) + processor = MagicMock(tokenizer=tokenizer) + processor.model_input_names = ["input_ids", "pixel_values", "image_grid_thw"] + + actor_args, _ = setup_single_controller(mc, tokenizer, processor=processor) + + call_args, call_kwargs = patched_factories["setup_response_data"].call_args + assert call_args[0] is processor + assert call_kwargs["env_configs"] == {"clevr-cogent": {"some": "value"}} + assert call_kwargs["is_vlm"] is True + warmup_fields = actor_args.dp_client.register_partition.call_args.kwargs[ + "fields" + ] + assert WIRE_MULTIMODAL_FIELDS <= set(warmup_fields) + def test_weight_sync_factory_args(self, patched_factories): """create_weight_synchronizer receives policy / generation / topology.""" mc = _make_master_config(colocated=False, backend="vllm") + mc.async_rl.generation_fleet_health.refit_timeout_s = 42.0 tokenizer = MagicMock(pad_token_id=0) setup_single_controller(mc, tokenizer) @@ -1064,6 +1089,11 @@ def test_weight_sync_factory_args(self, patched_factories): assert factory_kwargs["generation"] is patched_factories["fake_gen"] assert factory_kwargs["generation_backend"] == "vllm" assert factory_kwargs["colocated"] is False + assert factory_kwargs["refit_timeout_s"] == 42.0 + assert ( + patched_factories["fake_gen"].weight_synchronizer + is patched_factories["create_weight_synchronizer"].return_value + ) def test_custom_partition_id(self, patched_factories): mc = _make_master_config() @@ -1167,8 +1197,13 @@ def test_nemo_gym_wires_env_handle(self, patched_factories): patch.object(sc_setup_mod, "router_replay_enabled", return_value=False), ): tokenizer = MagicMock(pad_token_id=0) - actor_args, _ = setup_single_controller(mc, tokenizer) + processor = MagicMock(tokenizer=tokenizer) + actor_args, _ = setup_single_controller(mc, tokenizer, processor=processor) + data_args, data_kwargs = patched_factories["setup_response_data"].call_args + assert data_args[0] is processor + assert data_kwargs["env_configs"] is None + assert data_kwargs["is_vlm"] is True mock_spinup.assert_called_once_with( env_configs=mc.env, base_urls=patched_factories["fake_gen"].dp_openai_server_base_urls, @@ -1181,6 +1216,10 @@ def test_nemo_gym_wires_env_handle(self, patched_factories): token_capture=None, ) assert actor_args.env_handles["nemo_gym"] is fake_gym_actor + warmup_fields = actor_args.dp_client.register_partition.call_args.kwargs[ + "fields" + ] + assert WIRE_MULTIMODAL_FIELDS <= set(warmup_fields) def test_token_capture_always_creates_finalizer_actor_pool(self, patched_factories): mc = _make_master_config(backend="vllm") @@ -1203,6 +1242,8 @@ def test_token_capture_always_creates_finalizer_actor_pool(self, patched_factori None, ) fake_actors = [MagicMock(name=f"finalizer_{index}") for index in range(3)] + tokenizer = MagicMock(pad_token_id=9) + processor = MagicMock(tokenizer=tokenizer) with ( patch.object(sc_setup_mod, "should_use_nemo_gym", return_value=True), @@ -1215,7 +1256,7 @@ def test_token_capture_always_creates_finalizer_actor_pool(self, patched_factori return_value=fake_actors, ) as mock_create_finalizer_actors, ): - actor_args, _ = setup_single_controller(mc, MagicMock(pad_token_id=9)) + actor_args, _ = setup_single_controller(mc, tokenizer, processor=processor) (actor_dp_config, actor_config), actor_kwargs = ( mock_create_finalizer_actors.call_args @@ -1227,6 +1268,9 @@ def test_token_capture_always_creates_finalizer_actor_pool(self, patched_factori assert actor_kwargs == {"num_workers": 3} assert actor_args.finalizer_actors == fake_actors assert not hasattr(actor_args.rollout_manager, "_finalizer") + partition_calls = actor_args.dp_client.register_partition.call_args_list + assert WIRE_MULTIMODAL_FIELDS <= set(partition_calls[0].kwargs["fields"]) + assert WIRE_MULTIMODAL_FIELDS.isdisjoint(partition_calls[1].kwargs["fields"]) def test_setup_timing_populated_for_noncolocated_vllm(self, patched_factories): """Non-colocated vLLM records every per-phase field.""" @@ -1533,13 +1577,15 @@ def _spinup_gym(**_): mock_megatron.return_value.finish_generation.assert_called_once_with() if gym: # Gym spins up on the reserved URL, before the served-address - # cross-check — so the mismatch leg sees it too. + # cross-check — so the mismatch leg sees it too. The initial refit + # must happen during that wait because it starts Megatron's server. _, spinup_kwargs = mock_spinup.call_args assert spinup_kwargs["base_urls"] == [reserved_url] # The initial refit ran in setup, against the collective brought up # there; the served-address check reads the URLs it populated. weight_sync.init_communicator.assert_called_once_with() weight_sync.sync_weights.assert_called_once_with() + assert mock_megatron.return_value.weight_synchronizer is weight_sync else: mock_spinup.assert_not_called() # Native: the actor's startup sync performs the initial refit. diff --git a/tests/unit/single_controller/test_watchdog_pump.py b/tests/unit/single_controller/test_watchdog_pump.py index 0ae1f0ce77..b15659508d 100644 --- a/tests/unit/single_controller/test_watchdog_pump.py +++ b/tests/unit/single_controller/test_watchdog_pump.py @@ -339,7 +339,7 @@ def test_a_confirmed_death_stands_the_trainers_deadline_down(self): shard_count=2, policy=FleetHealthPolicy(unhealthy_threshold=99) ) ctrl = self._with_fleet(monitor, worker_alive=[True, False]) - asyncio.run(_run_probe_ticks(ctrl, 1)) + asyncio.run(ctrl._probe_generation_fleet()) assert ctrl._stood_down == [0, 1], ( "every policy worker must be told to stand its refit deadline down once a " f"generation shard is confirmed gone; saw {ctrl._stood_down}"