nemo_rl.data_plane.adapters.tq_mooncake_checkpoint#

Owner-distributed Mooncake checkpoints for TransferQueue.

Normal Mooncake PUTs remain memory-only. At an explicit TQ checkpoint, the controller commands the existing workers through their Ray actor handles. Workers discover disjoint key slices and route ownership metadata through Ray object references. Each payload is written directly from its owning process’s CPU segment: no payload crosses Ray, no native GET or staging copy is needed, and the coordinator never collects detailed object addresses/sizes. No new actors, registry, or normal PUT/CLEAR bookkeeping are introduced.

Restore reverses the process: the current Mooncake clients read size-balanced sets of durable objects and upsert them into their own preferred segments before TQ restores controller metadata. Saved client identities are not reused across restarts.

During save, objects selected by the controller snapshot must remain unchanged until every owner’s ACK. Generation may continue writing unrelated fresh keys, but overwrites and clears of selected objects must wait. Their borrowed CPU allocations must remain valid: do not move their replicas, unmount their segments, or close their store clients. Hard pinning prevents eviction, not these explicit mutations; same-host peer addresses are never dereferenced. During restore, keep writers and clears stopped through every owner’s verification and ACK. All intended restore clients must connect before tq.load_checkpoint is called. Checkpoint files must remain immutable throughout restore. Each owner validates its indexes and key absence before writing. Storage load completes before TQ installs its controller; it does not independently reread the controller snapshot. Failures may leave partial or overwritten objects and do not roll back peers’ restores. A failed restore is unusable; start with a fresh restore environment.

Checkpoints use checksum-free format v3. Earlier development formats are not supported.

Module Contents#

Classes#

_StoredObject

_SaveGroup

_ParticipantInfo

_ParticipantRequest

_ManifestObject

_ShardIndex

_CheckpointBuffer

_CheckpointParticipant

Executes checkpoint file I/O inside one existing Mooncake client.

_CheckpointManagerMixin

Add explicit checkpoint operations to TQ’s lazily imported manager.

Functions#

_fsync_directory

Make prior entry creation/rename operations durable in path.

_checkpoint_settings

_checkpoint_enabled

_storage_layout

_validate_metadata_mode

_validate_checkpoint_runtime

_has_no_named_dims

True when the tensor carries no named dimensions.

_physical_keys

Return every produced Mooncake key referenced by a TQ controller cut.

_controller_path

_controller_keys

_unregister_or_quarantine

_registered_buffer

_write_buffer

_read_buffer

_checkpoint_batches

Bound native calls by keys and bytes; allow an oversized singleton.

_checkpoint_buffer

_check_batch_results

Native GET returns byte counts; UPSERT returns zero statuses, per key.

_batch_replicas

_status_is_complete

_complete_memory_replicas

_complete_memory_buffers

_controller_session

Stable identity for one live TQ controller, including its endpoints.

_request_body

_safe_shard_name

_manifest_object_from_mapping

_shard_index_from_mapping

_parse_shard_indexes

_read_shard_indexes

_write_shard_index

_local_replica_config

_fanout_requests

Send metadata-only commands to existing actors; never RPC back to self.

_require_exact_responses

_live_participants

Describe the explicitly supplied workers, plus the calling process.

_write_manifest

_save_storage_checkpoint

Route metadata references; payload and detailed plans stay distributed.

_load_manifest

_require_clean_store

_load_sharded_checkpoint

_load_storage_checkpoint

Restore payload on current clients before TQ loads its controller.

configure_checkpoint_workers

Bind existing actor handles for this process’s checkpoint coordinator.

run_checkpoint_command

Execute an actor command using only its already-attached local store.

install_tq_mooncake_checkpoint_plugin

Install the explicit-checkpoint storage manager once.

Data#

API#

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._STORAGE_DIR#

‘mooncake_storage’

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._MANIFEST_FILE#

‘manifest.json’

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._DEFAULT_TIMEOUT_S#

200.0

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._BATCH_KEYS#

400

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._BATCH_BYTES#

None

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._BUFFER_ALIGNMENT#

256

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._QUARANTINED_BUFFERS: list[mmap.mmap]#

[]

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._fsync_directory(path: pathlib.Path) → None#

Make prior entry creation/rename operations durable in path.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._checkpoint_settings(
config: Any,
) → collections.abc.Mapping[str, Any]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._checkpoint_enabled(config: Any) → bool#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._storage_layout(config: Any) → dict[str, Any]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._validate_metadata_mode(manager: Any) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._validate_checkpoint_runtime(manager: Any) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._has_no_named_dims(tensor: torch.Tensor) → bool#

True when the tensor carries no named dimensions.

torch < 2.13 exposes the named-tensor API and reports (None, None) for an unnamed 2-D tensor; torch 2.13 removed named tensors along with the names attribute, so every tensor is unnamed there.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._physical_keys(
controller_state: collections.abc.Mapping[str, Any],
) → list[str]#

Return every produced Mooncake key referenced by a TQ controller cut.

class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._StoredObject#
key: str#

None

size: int#

None

address: int#

None

class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._SaveGroup#
controller_session: str#

None

owner: str#

None

objects: list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._StoredObject]#

None

class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantInfo#
participant_id: str#

None

controller_session: str#

None

segment_name: str#

None

transport_endpoint: str#

None

classmethod from_mapping(
value: collections.abc.Mapping[str, Any],
) → nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantInfo#
class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantRequest#
participant: nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantInfo#

None

body: dict[str, Any]#

None

class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ManifestObject#
key: str#

None

shard: str#

None

offset: int#

None

size: int#

None

saved_owner: str#

None

class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex#
shard: str#

None

object_count: int#

None

payload_bytes: int#

None

saved_owner: str#

None

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._controller_path(checkpoint_dir: pathlib.Path) → pathlib.Path#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._controller_keys(checkpoint_dir: pathlib.Path) → list[str]#
class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointBuffer#
payload: mmap.mmap#

None

pointer: int#

None

quarantined: bool#

False

classmethod allocate(
size: int,
) → nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointBuffer#
close() → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._unregister_or_quarantine(
store: Any,
buffer: nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointBuffer,
*,
label: str,
) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._registered_buffer(
store: Any,
buffer: nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointBuffer,
*,
size: int,
label: str,
) → Iterator[int]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._write_buffer(
output: Any,
buffer: mmap.mmap | memoryview,
size: int,
*,
label: str,
offset: int = 0,
) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._read_buffer(
source: Any,
buffer: mmap.mmap,
size: int,
*,
label: str,
offset: int = 0,
) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._checkpoint_batches(
sizes: list[int],
) → Iterator[tuple[int, int, list[int]]]#

Bound native calls by keys and bytes; allow an oversized singleton.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._checkpoint_buffer(
store: Any,
sizes: list[int],
*,
label: str,
) → Iterator[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointBuffer]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._check_batch_results(
results: Any,
expected: list[int],
*,
operation: str,
) → None#

Native GET returns byte counts; UPSERT returns zero statuses, per key.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._batch_replicas(
store: Any,
keys: list[str],
) → collections.abc.Mapping[str, Any]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._status_is_complete(status: Any) → bool#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._complete_memory_replicas(
descriptors: Any,
) → list[tuple[str, int]]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._complete_memory_buffers(descriptors: Any) → list[Any]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._controller_session(manager: Any) → str#

Stable identity for one live TQ controller, including its endpoints.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._request_body(
manager: Any,
operation: str,
**payload: Any,
) → dict[str, Any]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._safe_shard_name(value: Any) → str#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._manifest_object_from_mapping(
value: collections.abc.Mapping[str, Any],
) → nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ManifestObject#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._shard_index_from_mapping(
value: Any,
) → nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._parse_shard_indexes(
values: Any,
) → list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._read_shard_indexes(
root: pathlib.Path,
values: Any,
) → list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ManifestObject]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._write_shard_index(
storage_dir: pathlib.Path,
*,
shard: str,
entries: list[dict[str, Any]],
saved_owner: str,
) → nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._local_replica_config(
manager: Any,
segment_name: str,
) → Any#
class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointParticipant(manager: Any)#

Executes checkpoint file I/O inside one existing Mooncake client.

Initialization

_dispatch(
body: collections.abc.Mapping[str, Any],
) → dict[str, Any]#
_discover_save(
body: collections.abc.Mapping[str, Any],
) → dict[str, Any]#

Query each key once and leave detailed plans with their Ray owners.

_resolve_save_objects(
body: collections.abc.Mapping[str, Any],
) → list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._StoredObject]#
_checkpoint_root(
body: collections.abc.Mapping[str, Any],
) → pathlib.Path#
_verify_restored_objects(
entries: list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ManifestObject],
current_endpoints: set[str],
) → None#
_save_shard(
body: collections.abc.Mapping[str, Any],
) → dict[str, Any]#
_load_shards(
body: collections.abc.Mapping[str, Any],
) → dict[str, Any]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._fanout_requests(
requests: list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantRequest],
*,
workers: collections.abc.Mapping[str, Any],
local: nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointParticipant | None,
timeout_s: float,
) → dict[str, dict[str, Any]]#

Send metadata-only commands to existing actors; never RPC back to self.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._require_exact_responses(
requests: list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantRequest],
responses: collections.abc.Mapping[str, Any],
*,
operation: str,
) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._live_participants(
manager: Any,
) → tuple[list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ParticipantInfo], dict[str, Any]]#

Describe the explicitly supplied workers, plus the calling process.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._write_manifest(
storage_dir: pathlib.Path,
*,
config: Any,
entries: list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex],
) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._save_storage_checkpoint(
manager: Any,
checkpoint_dir: str,
) → None#

Route metadata references; payload and detailed plans stay distributed.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._load_manifest(
checkpoint_root: pathlib.Path,
) → tuple[list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex], dict[str, Any]]#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._require_clean_store(store: Any, keys: list[str]) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._load_sharded_checkpoint(
manager: Any,
root: pathlib.Path,
shards: list[nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._ShardIndex],
) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._load_storage_checkpoint(
manager: Any,
checkpoint_dir: str,
) → None#

Restore payload on current clients before TQ loads its controller.

class nemo_rl.data_plane.adapters.tq_mooncake_checkpoint._CheckpointManagerMixin(
controller_info: Any,
config: dict[str, Any],
)#

Add explicit checkpoint operations to TQ’s lazily imported manager.

Initialization

config: dict[str, Any]#

None

async save_checkpoint(checkpoint_dir: str) → None#
async load_checkpoint(checkpoint_dir: str) → None#
nemo_rl.data_plane.adapters.tq_mooncake_checkpoint.configure_checkpoint_workers(workers: list[Any]) → None#

Bind existing actor handles for this process’s checkpoint coordinator.

Call after all intended owners have attached, before save or restore. Do not include the calling actor: its shard is executed directly, including when restoring inside SingleController’s constructor.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint.run_checkpoint_command(
body: collections.abc.Mapping[str, Any],
) → dict[str, Any] | None#

Execute an actor command using only its already-attached local store.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint.install_tq_mooncake_checkpoint_plugin() → None#

Install the explicit-checkpoint storage manager once.

nemo_rl.data_plane.adapters.tq_mooncake_checkpoint.__all__#

[‘configure_checkpoint_workers’, ‘install_tq_mooncake_checkpoint_plugin’, ‘run_checkpoint_command’]