nemo_rl.experience.rollout_reassembler_actor#

CPU Ray actors for metadata-only token-capture finalization.

Module Contents#

Classes#

ReassemblyRequest

Metadata-only input for one prompt group’s finalization.

RolloutReassemblerActorConfig

Internal constructor values shared by every finalizer actor.

RolloutReassemblerActor

Own a connect-only TQ client and lightweight finalizer in one process.

Functions#

assert_metadata_only

Reject tensors and known heavy row fields reachable from an RPC graph.

create_rollout_reassembler_actors

Construct and validate the pool after TQ partitions are registered.

Data#

API#

nemo_rl.experience.rollout_reassembler_actor._FORBIDDEN_RPC_KEYS#

‘frozenset(…)’

class nemo_rl.experience.rollout_reassembler_actor.ReassemblyRequest#

Metadata-only input for one prompt group’s finalization.

group_id: str#

None

rollout_ids: tuple[str, ...]#

None

canonical_sample_ids: tuple[str, ...]#

None

receipts: tuple[Optional[dict[str, Any]], ...]#

None

rewards: tuple[float, ...]#

None

fallback_weight_version: int#

None

prompt_idx: int#

None

mask_sample: tuple[bool, ...]#

None

loss_multiplier: float#

1.0

class nemo_rl.experience.rollout_reassembler_actor.RolloutReassemblerActorConfig#

Internal constructor values shared by every finalizer actor.

partition_id: str#

None

staging_partition: str#

None

pad_token_id: int#

None

router_replay_enabled: bool#

None

defer_routed_experts_to_policy: bool#

None

max_seq_len: int#

None

nemo_rl.experience.rollout_reassembler_actor.assert_metadata_only(value: Any, *, path: str = 'rpc') None#

Reject tensors and known heavy row fields reachable from an RPC graph.

class nemo_rl.experience.rollout_reassembler_actor.RolloutReassemblerActor(
dp_config: nemo_rl.data_plane.DataPlaneConfig,
config: nemo_rl.experience.rollout_reassembler_actor.RolloutReassemblerActorConfig,
)#

Own a connect-only TQ client and lightweight finalizer in one process.

Initialization

mooncake_checkpoint(
body: dict[str, Any],
) dict[str, Any] | None#

Run an owner-local checkpoint command; return metadata, never payloads.

check_dependencies() None#

Import the finalization API before the controller starts rollouts.

finalize(
request: nemo_rl.experience.rollout_reassembler_actor.ReassemblyRequest,
) nemo_rl.experience.rollout_reassembler.FinalizedGroup#

Finalize one request without allowing tensor payloads across Ray RPC.

nemo_rl.experience.rollout_reassembler_actor.create_rollout_reassembler_actors(
dp_config: nemo_rl.data_plane.DataPlaneConfig,
config: nemo_rl.experience.rollout_reassembler_actor.RolloutReassemblerActorConfig,
*,
num_workers: int,
) list[Any]#

Construct and validate the pool after TQ partitions are registered.