Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
96 commits
Select commit Hold shift + click to select a range
7c79757
feat(data-plane): track data-plane time, latency percentiles and byte…
ZhiyuLi-Nvidia Aug 12, 2026
163dcd3
feat(data-plane): emit per-step data-plane metrics from the sync trainer
ZhiyuLi-Nvidia Aug 13, 2026
3bb969d
refactor(data-plane): address /simplify review of the metrics diff
ZhiyuLi-Nvidia Aug 13, 2026
9355d19
perf(data-plane): cut metrics overhead, add optional wire-hash check
ZhiyuLi-Nvidia Aug 21, 2026
738e863
fix(data-plane): fingerprint jagged leaves; simplify the hash to one …
ZhiyuLi-Nvidia Aug 22, 2026
6208e81
perf(data-plane): scope jagged fingerprints to the buffer, not a rect…
ZhiyuLi-Nvidia Aug 22, 2026
3c23f1d
fix(data-plane): count fields dropped as incomparable
ZhiyuLi-Nvidia Aug 22, 2026
8a7fe10
refactor(data-plane): trim the observability diff after review
ZhiyuLi-Nvidia Aug 22, 2026
2d5ffcf
perf(data-plane): size nested leaves from their packed values buffer
ZhiyuLi-Nvidia Aug 22, 2026
0b47cd2
refactor(data-plane): name the real hazard in _tensor_bytes
ZhiyuLi-Nvidia Aug 22, 2026
e1b9745
feat(data-plane): enable data-plane metrics by default
ZhiyuLi-Nvidia Aug 22, 2026
9a5bdcd
perf(data-plane): hold one byte total per partition, not one per key
ZhiyuLi-Nvidia Aug 22, 2026
0297a33
fix(data-plane): one unit per dimension, and a tail metric that moves
ZhiyuLi-Nvidia Aug 22, 2026
3bc1ff0
fix(data-plane): report the latency split in ms, not as a ratio
ZhiyuLi-Nvidia Aug 22, 2026
fb54cc3
fix(data-plane): scope the metric prefix to the driver that produced it
ZhiyuLi-Nvidia Aug 22, 2026
dc30eed
feat(data-plane): aggregate metrics across processes, and report thei…
ZhiyuLi-Nvidia Aug 22, 2026
a51618e
fix(data-plane): bill the fan-out to observability, not just the wrapper
ZhiyuLi-Nvidia Aug 22, 2026
224320c
fix(data-plane): a max below its own median, and a fraction above 1
ZhiyuLi-Nvidia Aug 22, 2026
fba0047
fix(data-plane): make every series say whether it is a delta or a level
ZhiyuLi-Nvidia Aug 22, 2026
3b2eb78
fix(data-plane): free only the keys a clear actually held
ZhiyuLi-Nvidia Aug 22, 2026
a2fee81
fix(data-plane): clamp percentiles where they are derived, not at one…
ZhiyuLi-Nvidia Aug 22, 2026
7d4d73b
refactor(data-plane): cleanup pass, and two bugs it turned up
ZhiyuLi-Nvidia Aug 22, 2026
6baa3d0
feat(data-plane): log a per-op breakdown table, not just series
ZhiyuLi-Nvidia Aug 22, 2026
06eddec
fix(data-plane): put the latency split in the cluster view too
ZhiyuLi-Nvidia Aug 22, 2026
75445bc
fix(data-plane): per-op cluster time also reads as elapsed when it is…
ZhiyuLi-Nvidia Aug 22, 2026
86f19ee
fix(data-plane): report per-op time per call, which is the invariant one
ZhiyuLi-Nvidia Aug 22, 2026
b061865
fix(data-plane): split the latency per call, like the mean it explains
ZhiyuLi-Nvidia Aug 22, 2026
1dd8b8a
fix(data-plane): gate percentiles per quantile; per-row hashes when r…
ZhiyuLi-Nvidia Aug 22, 2026
740ad3b
fix(data-plane): report p90, not p99, as the per-op tail
ZhiyuLi-Nvidia Aug 22, 2026
0f188e5
feat(data-plane): chart time shares, not 32 per-op series
ZhiyuLi-Nvidia Aug 22, 2026
cd12f8b
refactor(data-plane): name the metric for its denominator, share -> t…
ZhiyuLi-Nvidia Aug 22, 2026
70847a5
refactor(data-plane): time_pct -> pct_of_dataplane
ZhiyuLi-Nvidia Aug 22, 2026
303ef2d
refactor(data-plane): spell out percent in the metric names
ZhiyuLi-Nvidia Aug 22, 2026
82b5bd5
docs(data-plane): the transfer term IS the bandwidth term
ZhiyuLi-Nvidia Aug 22, 2026
814396d
fix(data-plane): report hash counters and measuring cost in both scopes
ZhiyuLi-Nvidia Aug 22, 2026
5d10ed6
fix(data-plane): log before the step commits; track per-op transfer v…
ZhiyuLi-Nvidia Aug 23, 2026
9e68df4
fix(data-plane): stop the hash guard failing every per-row field
ZhiyuLi-Nvidia Aug 23, 2026
275a44b
fix(data-plane): make a hash mismatch adjudicable, and measure what i…
ZhiyuLi-Nvidia Aug 23, 2026
05c60b6
refactor(data-plane): apply the /simplify review
ZhiyuLi-Nvidia Aug 23, 2026
33d8077
refactor(data-plane): narrow this PR back to the data plane
ZhiyuLi-Nvidia Aug 23, 2026
1ebb49f
fix(data-plane): satisfy pyrefly on the three types the diff loosened
ZhiyuLi-Nvidia Aug 24, 2026
bfa2662
style(data-plane): sort the test imports
ZhiyuLi-Nvidia Aug 24, 2026
c0a96bf
fix(data-plane): salt hash fingerprints on the host
ZhiyuLi-Nvidia Aug 24, 2026
e496314
feat(data-plane): log data-plane metrics from the single-controller loop
ZhiyuLi-Nvidia Aug 24, 2026
795da38
test(data-plane): simplify observability tests; guard the metrics path
ZhiyuLi-Nvidia Aug 27, 2026
b154db7
test(data-plane): gate every data-plane nightly on the wire guard
ZhiyuLi-Nvidia Aug 27, 2026
50af02c
refactor(data-plane): drop the latency/bandwidth fit
ZhiyuLi-Nvidia Aug 27, 2026
e142ddf
feat(data-plane): report jagged pack/unpack time
ZhiyuLi-Nvidia Aug 29, 2026
c374c95
refactor(data-plane): report charted durations in seconds
ZhiyuLi-Nvidia Aug 30, 2026
dea2ac3
fix(data-plane): make the wire guard see a row's shape
ZhiyuLi-Nvidia Aug 30, 2026
a82313e
fix(data-plane): log single-controller metrics before the step is com…
ZhiyuLi-Nvidia Sep 7, 2026
7d5123c
fix(data-plane): drop the PROMOTE_1D_FIELDS branch from _from_wire
ZhiyuLi-Nvidia Sep 7, 2026
31ea144
refactor(data-plane): share the duplicated observability wiring
ZhiyuLi-Nvidia Sep 9, 2026
5d1274c
chore(data-plane): satisfy ruff, ruff-format and pyrefly on the obser…
ZhiyuLi-Nvidia Sep 9, 2026
2f35378
fix(data-plane): satisfy pyrefly in the observability module
ZhiyuLi-Nvidia Sep 9, 2026
3ac8b4a
fix(data-plane): release fingerprints when another process clears the…
Sep 10, 2026
c3ff70c
fix(data-plane): reconcile the accounting for samples another process…
Sep 10, 2026
2a3a758
refactor(data-plane): drop the reconcile warning, keep the release
Sep 10, 2026
e804973
feat(data-plane): make the wire guard check across processes
Sep 10, 2026
eda30d5
Revert "feat(data-plane): make the wire guard check across processes"
Sep 10, 2026
c2db57d
feat(data-plane): carry the wire-in digest beside the field it describes
Sep 10, 2026
f621a7e
test(data-plane): rewrite the hash tests against the mirror columns
Sep 10, 2026
2849445
fix(data-plane): one definition of the mirror column list
Sep 10, 2026
bad05cb
test(data-plane): one sanity check, not five
Sep 10, 2026
f22f37f
test(data-plane): the later writer stamps too
Sep 10, 2026
42a257e
perf(data-plane): fold the field digests on tensors, and say so in th…
Sep 10, 2026
4e9b62f
test(data-plane): assert the wire guard looked, not just that it foun…
Sep 10, 2026
f5d48d1
test(data-plane): gate on what the guard did, not on what it did not
Sep 10, 2026
12d36cc
docs(data-plane): frac_of_step is carried work, not exclusive time
Sep 10, 2026
5b73368
test(data-plane): the guard is a test flag, not a recipe default
Sep 11, 2026
dac5f98
fix(data-plane): the cluster's wall time is the max process, not the sum
Sep 11, 2026
42a1ecf
docs(data-plane): say what frac_of_step reduces, and what the guard c…
Sep 11, 2026
7ba6461
docs(data-plane): the guard's measured cost, from a 15-step A/B
Sep 11, 2026
ed21f68
docs(data-plane): do not attribute the step delta to the guard
Sep 11, 2026
07c6ca2
docs(data-plane): the guard's cost, measured on one node
Sep 11, 2026
15d7155
docs: track doc for the metrics-accuracy investigation
Sep 11, 2026
c3deb05
docs(data-plane): say why the slowest process is the number that matters
Sep 11, 2026
1c3b180
Revert the iterate track doc
Sep 11, 2026
a68253b
fix(data-plane): name the unit on self/frac, and pin max over mean
Sep 11, 2026
e2fbf2f
refactor(data-plane): name the reduction, drop the orphaned parameter
Sep 11, 2026
384bb8c
docs(data-plane): self/frac is a share of process-time, not of wall_s
Sep 11, 2026
60f95bc
test(data-plane): stop restating the inherited default in 24 recipes
Sep 11, 2026
7285c25
test(data-plane): drop the log-ordering guard, superseded upstream
Sep 13, 2026
65c7202
style(data-plane): satisfy ruff-format
Sep 14, 2026
ccdd299
fix(data-plane): remove a definition the rebase duplicated
Sep 14, 2026
b548e51
fix(data-plane): restore the observability block on the PPO SC recipe
Sep 15, 2026
426fefe
fix(data-plane): scope the slowest-process wall time to the step
Sep 15, 2026
52149ea
fix(data-plane): give the step metrics one owner and one predicate
Sep 15, 2026
299c6ff
fix(logger): buffer the breakdown table with the step's metrics
Sep 15, 2026
85e8eb1
docs(data-plane): correct the metric descriptions and widen the hash …
Sep 15, 2026
bf67e07
refactor(data-plane): keep one owner for the previous step reading
Sep 15, 2026
9a80cf0
fix(data-plane): do not look up the single-process key on a merged sn…
Sep 16, 2026
cd980c9
refactor(data-plane): merge the step window through the existing max …
Sep 16, 2026
6f3b86b
fix(data-plane): make a failing metrics panel loud on its first failure
Sep 16, 2026
b713773
fix(data-plane): degrade one metric rather than the whole panel
Sep 16, 2026
d60550e
fix(data-plane): satisfy pyrefly on the step window and the failure c…
Sep 16, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions examples/configs/grpo_math_1B.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -597,8 +597,9 @@ data_plane:
# serialized groups and cost throughput;
# higher buys little at linear HBM cost.
# GDR needs that headroom to pay off.
# observability: # NotRequired
# enabled: false
observability: # per-op data-plane timing/volume
enabled: true # per-op timing/volume; cost is below measurement noise
verify_tensor_hash: false # debug: per-row hash of each row's values+dtype+shape, wire-in vs wire-out

# Multi-Teacher On-Policy Distillation (MOPD): distills from one or more teacher
# models into the policy via token-level teacher-minus-student logprob advantages,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,9 @@ data_plane:
local_buffer_size: 4294967296
reuse_registered_buffers: true
staging_buffer_size: 268435456
observability:
enabled: true
verify_tensor_hash: false

# SC async-RL runtime knobs, replacing the nulled ppo.async_ppo block.
async_rl:
Expand Down
44 changes: 44 additions & 0 deletions nemo_rl/algorithms/grpo_sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,10 @@
from nemo_rl.data.llm_message_utils import batched_message_log_to_flat_message
from nemo_rl.data.multimodal_utils import present_multimodal_fields
from nemo_rl.data_plane.interfaces import KVBatchMeta
from nemo_rl.data_plane.observability import (
log_step_metrics,
metrics_never_fail_the_step,
)
from nemo_rl.data_plane.schema import DP_CALIB_INPUT_FIELDS, DP_TRAIN_FIELDS
from nemo_rl.distributed.batched_data_dict import BatchedDataDict
from nemo_rl.environments.interfaces import EnvironmentInterface
Expand Down Expand Up @@ -381,6 +385,42 @@ def _compute_seq_logprob_error_metrics(
return masking_data["sample_mask"], seq_logprob_error_metrics


def _log_data_plane_metrics(
policy: Any, logger: Logger, step: int, total_step_time: float
) -> None:
"""Log this step's data-plane cost. Never raises.

On by default, so this runs every step of every recipe.
"""
with metrics_never_fail_the_step(step):
_log_data_plane_metrics_impl(policy, logger, step, total_step_time)


def _log_data_plane_metrics_impl(
policy: Any, logger: Logger, step: int, total_step_time: float
) -> None:
"""Log this step's data-plane cost. No-op unless observability is enabled.

The policy computes both the metrics and the scope they cover, because
the baselines they are differenced against belong with the client whose
counters they baseline. This end owns only where they are logged.

The prefix names the scope because the two differ by a lot: the driver
issues about one op of each kind per step while the bulk traffic is the
workers' per-DP-rank ``get_samples``. Note that even the cluster view
omits the rollout actor, which builds its own client and is not on the
worker group -- so ``kv_first_write`` is not in these totals.
"""
get_metrics = getattr(policy, "get_data_plane_step_metrics", None)
if not callable(get_metrics):
return # not a data-plane policy
result = get_metrics(total_step_time)
if result is None:
return # observability disabled -> plain adapter
metrics, scope = result
log_step_metrics(logger, metrics, step, scope)


def grpo_train_sync(
policy: ColocatablePolicyInterface,
policy_generation: GenerationInterface,
Expand Down Expand Up @@ -1375,6 +1415,10 @@ def grpo_train_sync(
logger.log_metrics(
performance_metrics, total_steps + 1, prefix="performance"
)
# Before the step_finished=True log below, which commits the step:
# anything logged against a committed step is dropped by wandb, so
# these series were computed, printed, and silently discarded.
_log_data_plane_metrics(policy, logger, total_steps + 1, total_time)
logger.log_metrics(
timing_metrics,
total_steps + 1,
Expand Down
43 changes: 43 additions & 0 deletions nemo_rl/algorithms/single_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,11 @@
from nemo_rl.data.multimodal_utils import present_multimodal_fields
from nemo_rl.data_plane import DATA_PLANE_CHECKPOINT_SCHEMA_VERSION, KVBatchMeta
from nemo_rl.data_plane.async_utils import call_data_plane
from nemo_rl.data_plane.observability import (
is_metrics_client,
log_step_metrics,
metrics_never_fail_the_step,
)
from nemo_rl.data_plane.schema import (
DP_CALIB_INPUT_FIELDS,
DP_TRAIN_FIELDS,
Expand Down Expand Up @@ -1752,6 +1757,39 @@ async def _cleanup_consumed_metas_unlocked(
if errors:
raise BaseExceptionGroup("post-train DataPlane cleanup failed", errors)

def _log_data_plane_metrics(self, total_step_time: float) -> None:
"""Log this step's data-plane cost. Never raises.

On by default, so this runs every step of every recipe. Mirrors
``grpo_sync._log_data_plane_metrics``.
"""
with metrics_never_fail_the_step(self._train_steps):
self._log_data_plane_metrics_impl(total_step_time)

def _log_data_plane_metrics_impl(self, total_step_time: float) -> None:
"""Log this step's data-plane cost. No-op unless observability is enabled.

The synchronous loop logs these series from ``_log_data_plane_metrics``
in ``grpo_sync``. Without the same call here the single-controller path
builds the metrics client, pays for its counters on every op, and emits
nothing -- the failure is silent, because an empty dashboard looks the
same as a data plane that cost nothing.

Driver scope only, and the prefix says so. This client issues the
advantage stage's get, the put that writes the advantages back, and
the post-train clear; the bulk traffic is
the trainer and generation workers' own clients, in their own
processes with their own counters, so ``comm_volume_mb`` here is well
under what the job actually moved. ``grpo_sync`` gets a cluster view by
fanning out over its policy worker group; this loop has no such group to
fan out over, so driver scope is all there is here.
"""
if not is_metrics_client(self._dp_client):
return # observability disabled -> plain adapter

metrics = self._dp_client.get_step_metrics(total_step_time)
log_step_metrics(self._logger, metrics, self._train_steps, "driver")

@staticmethod
def _group_ids_from_meta(meta: KVBatchMeta) -> list[str]:
"""Return stable prompt-group IDs in canonical sample order."""
Expand Down Expand Up @@ -3077,6 +3115,11 @@ async def _train_pump(self) -> None:
self._logger.log_metrics(
step_metrics, step=self._train_steps, prefix="train"
)
# Must precede the step_finished=True log below. That log commits
# the wandb step, and wandb silently discards anything logged
# against a step it has already committed -- no exception, no
# failed return, just an empty chart. grpo_sync had the same bug.
self._log_data_plane_metrics(total_time)
# step_finished=True here since this is the final log of our current step.
self._logger.log_metrics(
timing_metrics,
Expand Down
Loading
Loading