Skip to content
10 changes: 10 additions & 0 deletions examples/configs/grpo_math_1B.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -565,6 +565,16 @@ data_plane:
local_buffer_size: 4294967296 # 4 GiB/process
reuse_registered_buffers: true # reuse RDMA-registered buffers
staging_buffer_size: 268435456 # 256 MiB/pool slot; bigger transfers bypass the pool
use_gdr: false # GPU-memory RDMA staging in CUDA clients
gdr_staging_buffer_mb: 1024 # HBM reserved per active GDR client;
# 8 clients/node is ~8 GiB off the training
# budget. Must exceed the largest per-sample,
# per-field payload: below that, GDR reads of
# CPU-written values fail inside TransferQueue.
# Lower values split fetches into more
# serialized groups and cost throughput;
# higher buys little at linear HBM cost.
# GDR needs that headroom to pay off.
# observability: # NotRequired
# enabled: false

Expand Down
20 changes: 18 additions & 2 deletions nemo_rl/data_plane/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -437,6 +437,8 @@ data_plane:
local_buffer_size: 4294967296 # 4 GiB/process
reuse_registered_buffers: true # reuse RDMA-registered buffers
staging_buffer_size: 268435456 # 256 MiB/pool slot; bigger transfers bypass the pool
use_gdr: false # GPU-memory RDMA staging in CUDA clients
gdr_staging_buffer_mb: 1024 # persistent MiB per active GDR client
Comment thread
zyzhou5 marked this conversation as resolved.
Comment thread
zyzhou5 marked this conversation as resolved.
# observability: # NotRequired
# enabled: false
```
Expand All @@ -450,8 +452,22 @@ with no warning either way.
Backend choice:
- **`simple`** — ZMQ-backed; lowest setup overhead. Default for tests
and small runs.
- **`mooncake_cpu`** — Mooncake transfer engine; higher throughput at
scale. Required for multi-node clusters with large bulk volume.
- **`mooncake_cpu`** — Mooncake's RDMA-only transfer engine. By default,
tensors transfer through registered CPU staging. Set
`mooncake_cpu.use_gdr: true` to let CUDA-initialized clients use
Comment thread
terrykong marked this conversation as resolved.
TransferQueue's GDR staging path. CPU-only clients, such as a
SingleController producer, continue to use CPU RDMA. GDR changes the
client-side tensor transfer and staging path; queued objects still reside in
Mooncake-managed host-memory segments.

The CPU host staging pool's `staging_buffer_size` is independent of the GDR
buffer. `gdr_staging_buffer_mb` is the persistent GPU staging capacity per
active CUDA client and defaults to 1024 MiB. Transfers through that per-client
buffer are serialized. Aggregate fetches may exceed the buffer and are split
into groups. In the mixed CPU-producer/GDR-receiver flow used by
SingleController, however, each individual tensor must currently fit because
the CPU PUT path does not create the chunk metadata required by an oversized
GDR GET.

Capacity rule of thumb (any backend):

Expand Down
48 changes: 48 additions & 0 deletions nemo_rl/data_plane/adapters/transfer_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import importlib
import ipaddress
import json
import logging
import os
import resource
import socket
Expand Down Expand Up @@ -55,6 +56,8 @@
data_plane_supports_checkpointing,
)

LOGGER = logging.getLogger(__name__)

# ──────────────────────────────────────────────────────────────────────────
# Backend init — lifted from rl-arena/arena/backends.py.
# ──────────────────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -684,6 +687,8 @@ def _init_tq(cfg: DataPlaneConfig) -> None:
"metadata_server": f"{local_ip}:50050",
"master_server_address": f"{local_ip}:50051",
**_mooncake_transport_config(),
"use_gdr": bool(mooncake_cfg.use_gdr),
"gdr_staging_buffer_mb": int(mooncake_cfg.gdr_staging_buffer_mb),
},
},
}
Expand Down Expand Up @@ -771,6 +776,12 @@ def _from_wire(td: TensorDict) -> TensorDict:
class TQDataPlaneClient(DataPlaneClient):
"""Adapter façade — maps NeMo-RL calls onto TransferQueue's public API."""

# Class-level so ``put_samples`` stays readable on an instance built
# without ``__init__`` — ``object.__new__`` in tests, or a process that
# unpickles a client without running the constructor.
_gdr_requested: bool = False
_gdr_put_confirmed: bool = False

def __init__(self, cfg: DataPlaneConfig, *, bootstrap: bool = True) -> None:
"""Construct a TQ-backed client.

Expand Down Expand Up @@ -817,6 +828,13 @@ def __init__(self, cfg: DataPlaneConfig, *, bootstrap: bool = True) -> None:

self._backend = cfg["backend"]
self._supports_checkpointing = data_plane_supports_checkpointing(cfg)
# GDR is a mooncake_cpu-only transport knob, so key it off the backend
# directly rather than off any incidental per-backend flag.
self._gdr_requested = self._backend == "mooncake_cpu" and bool(
backend_config(cfg).use_gdr
)
self._gdr_put_confirmed = False

if bootstrap:
_init_tq(cfg)
else:
Expand Down Expand Up @@ -1031,6 +1049,31 @@ def put_samples(
wire_fields = detached_fields
field_names = [str(key) for key in detached_fields.keys()]

confirm_gdr_put = bool(
self._gdr_requested
and not self._gdr_put_confirmed
and torch.cuda.is_initialized()
and wire_fields is not None
and any(
isinstance(wire_fields.get(key), torch.Tensor)
for key in wire_fields.keys()
)
)
if confirm_gdr_put:
# Checked before the put, not after: TQ fixes GDR eligibility when
# the client attaches, so this is decidable up front — and once
# `kv_batch_put` returns, the rows are already durable and the
# controller has been notified, so raising then would strand them.
tq_client = tq.get_client()
storage_manager = getattr(tq_client, "storage_manager", None)
storage_client = getattr(storage_manager, "storage_client", None)
gdr_staging = getattr(storage_client, "_gdr_staging", None)
if not getattr(storage_client, "use_gdr", False) or gdr_staging is None:
raise RuntimeError(
"GDR was requested for a CUDA-initialized TransferQueue "
"client, but TransferQueue selected CPU RDMA for tensor PUTs"
)

self._mark_data_operation_started()
# TQ's wire vocabulary is `keys=` — translation point.
tq.kv_batch_put(
Expand All @@ -1039,6 +1082,11 @@ def put_samples(
fields=wire_fields,
tags=user_tags,
)
if confirm_gdr_put:
LOGGER.info(
"TransferQueue GDR tensor PUT active (partition=%s)", partition_id
)
self._gdr_put_confirmed = True
Comment thread
zyzhou5 marked this conversation as resolved.

return KVBatchMeta(
partition_id=partition_id,
Expand Down
9 changes: 8 additions & 1 deletion nemo_rl/data_plane/interfaces.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
from pathlib import Path
from typing import Annotated, Any, Callable, Literal, NotRequired, Sequence, TypedDict

from pydantic import BaseModel, Field
from pydantic import BaseModel, Field, PositiveInt
from tensordict import TensorDict

DATA_PLANE_CHECKPOINT_SCHEMA_VERSION = 2
Expand Down Expand Up @@ -80,6 +80,11 @@ class MooncakeCpuConfig(BaseModel, extra="allow"):
admitted and never shrink — so raise it only when a per-key payload (one
sample of one field) genuinely exceeds it, not for headroom.

``use_gdr`` lets CUDA-initialized clients transfer through TransferQueue's
persistent GPU staging buffer. ``gdr_staging_buffer_mb`` is the positive
HBM capacity of that buffer per active GDR client. CPU-only clients keep
using the registered host-buffer path.

Every RDMA rail on the host is offered to mooncake (see ``rdma_devices``).
That is only safe with ``MC_ENABLE_DEST_DEVICE_AFFINITY=1``, which pins each
transfer's peer rail to the local one by name; on a rail-isolated RoCE
Expand All @@ -91,6 +96,8 @@ class MooncakeCpuConfig(BaseModel, extra="allow"):
local_buffer_size: int = 4294967296 # 4 GiB per client process
reuse_registered_buffers: bool = True
staging_buffer_size: int = 268435456 # 256 MiB per pool slot
use_gdr: bool = False
gdr_staging_buffer_mb: PositiveInt = 1024
Comment thread
zyzhou5 marked this conversation as resolved.


class DataPlaneConfig(TypedDict):
Expand Down
17 changes: 15 additions & 2 deletions nemo_rl/data_plane/worker_mixin.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,10 +36,9 @@
import numpy as np
import torch

FetchPolicy = Literal["auto", "independent", "leader_broadcast"]

from nemo_rl.data.llm_message_utils import attach_message_log_view
from nemo_rl.data.multimodal_utils import PackedTensor
from nemo_rl.data_plane.interfaces import LocalDataPlaneConfig, backend_config
from nemo_rl.data_plane.schema import (
ELEM_COUNTS_PER_GB,
GLOBAL_FORWARD_PAD_SEQLEN,
Expand All @@ -64,6 +63,8 @@
DataPlaneRuntimeConfig,
)

FetchPolicy = Literal["auto", "independent", "leader_broadcast"]


def _broadcast_batched_data_dict(
data: Optional[BatchedDataDict[Any]],
Expand Down Expand Up @@ -327,6 +328,18 @@ def setup_data_plane(self, cfg: DataPlaneRuntimeConfig) -> None:
self._route_fallback_counts = Counter()
from nemo_rl.data_plane import build_data_plane_client

# ``LocalDataPlaneConfig`` is the process-local plane: no TQ, no
# mooncake, so no GDR to order against a CUDA context.
if (
not isinstance(cfg, LocalDataPlaneConfig)
and cfg["backend"] == "mooncake_cpu"
and backend_config(cfg).use_gdr
and not torch.cuda.is_initialized()
):
raise RuntimeError(
"CUDA must be initialized before attaching TransferQueue with GDR"
)

# bootstrap=False — the driver already created the named
# controller actor; this process attaches as a client.
self._dp_client = build_data_plane_client(cfg, bootstrap=False)
Expand Down
7 changes: 7 additions & 0 deletions tests/unit/data_plane/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,13 @@ Generated audit of every test function under `tests/unit/data_plane/` with a one
- `test_local_node_ip_returns_empty_on_exception` — DNS exception → returns empty string (no crash).
- `test_mc_tcp_bind_address_overwrites_existing` — TQDataPlaneClient `__init__` uses direct assignment (not `setdefault`).

## `test_mooncake_gdr.py` (4 tests)

- `test_init_tq_forwards_nested_gdr_config_and_keeps_rdma` — Nested GDR settings reach TransferQueue without changing #2935's RDMA/all-rail transport.
- `test_cpu_only_client_may_attach_with_gdr_config` — A CPU-only producer may attach because GDR eligibility is client-local.
- `test_gdr_tensor_put_is_confirmed_once_and_never_falls_back` — A CUDA client's tensor PUT raises rather than downgrading to CPU RDMA, and logs the GDR confirmation exactly once.
- `test_gdr_receiver_requires_cuda_initialized` — A policy receiver must initialize CUDA before attaching its GDR client.

## `test_message_log_decompose.py` (11 tests)

- `test_decompose_message_log_basic_shapes` — Basic shapes of decompose output.
Expand Down
19 changes: 18 additions & 1 deletion tests/unit/data_plane/test_backend_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,12 +61,19 @@ def test_checkpointing_capability_defaults_to_unsupported(
def test_nested_block_is_used() -> None:
cfg = _cfg(
"mooncake_cpu",
mooncake_cpu={"global_segment_size": 111, "reuse_registered_buffers": False},
mooncake_cpu={
"global_segment_size": 111,
"reuse_registered_buffers": False,
"use_gdr": True,
"gdr_staging_buffer_mb": 512,
},
)
resolved = backend_config(cfg)
assert isinstance(resolved, MooncakeCpuConfig)
assert resolved.global_segment_size == 111
assert resolved.reuse_registered_buffers is False
assert resolved.use_gdr is True
assert resolved.gdr_staging_buffer_mb == 512


def test_absent_block_falls_back_to_model_defaults() -> None:
Expand All @@ -81,6 +88,16 @@ def test_absent_block_falls_back_to_model_defaults() -> None:
assert resolved.local_buffer_size == 4294967296 # 4 GiB per client process
# The opt-out flag defaults on, so omitting it must not disable the pool.
assert resolved.reuse_registered_buffers is True
assert resolved.use_gdr is False
assert resolved.gdr_staging_buffer_mb == 1024


@pytest.mark.parametrize("value", [0, -1])
def test_gdr_staging_size_rejects_non_positive_values(value: int) -> None:
with pytest.raises(pydantic.ValidationError):
backend_config(
_cfg("mooncake_cpu", mooncake_cpu={"gdr_staging_buffer_mb": value})
)


def test_accepts_an_already_coerced_model() -> None:
Expand Down
Loading
Loading