From d6dbd0094cfafd0529ab755389f00f54c70d2454 Mon Sep 17 00:00:00 2001 From: Zhiyu Li Date: Mon, 7 Sep 2026 14:13:30 -0700 Subject: [PATCH] fix(data-plane): log SC metrics before the step is committed, not after `log_metrics(..., step_finished=True)` calls `run.log(..., commit=True)`, which closes that wandb step. wandb accepts a later log against a closed step, returns normally, and discards it -- no exception, no falsy return. `single_controller` called `_log_data_plane_metrics` immediately *after* that commit, so every `data_plane/*` series it produced was computed and dropped. The console still printed "data plane: ...ms" because that print lives inside the discarded call, so the run looked healthy. `grpo_sync` had the identical bug and was fixed when a real run showed 85 logged keys and zero `data_plane/*`. This path was never exercised by a run, so it stayed broken -- its own docstring says it mirrors `grpo_sync`, and it mirrored this too. The test asserts the invariant across every algorithm module rather than one call site, since a third wiring would otherwise repeat it. It compares the *call* (rindex) rather than the first occurrence, which is the method definition and would make the assertion vacuously true -- the first draft of this test passed against the unfixed file. Signed-off-by: Zhiyu Li --- nemo_rl/algorithms/single_controller.py | 6 ++- tests/unit/data_plane/test_observability.py | 45 +++++++++++++++++++++ 2 files changed, 50 insertions(+), 1 deletion(-) diff --git a/nemo_rl/algorithms/single_controller.py b/nemo_rl/algorithms/single_controller.py index 04daf5ab07..f32ee2e470 100644 --- a/nemo_rl/algorithms/single_controller.py +++ b/nemo_rl/algorithms/single_controller.py @@ -2183,6 +2183,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, @@ -2190,7 +2195,6 @@ async def _train_pump(self) -> None: prefix="timing/train", step_finished=True, ) - self._log_data_plane_metrics(total_time) self._timer.reset() # min sample version refers to the version each consumed sample was diff --git a/tests/unit/data_plane/test_observability.py b/tests/unit/data_plane/test_observability.py index f1151f92a1..7358cc2919 100644 --- a/tests/unit/data_plane/test_observability.py +++ b/tests/unit/data_plane/test_observability.py @@ -1263,3 +1263,48 @@ def test_inspection_snapshot_does_not_steal_codec_time(): 0.010, rel=1e-3 ) client.close() + + +def test_no_algorithm_logs_data_plane_after_committing_the_step(): + """``log_metrics(..., step_finished=True)`` commits the wandb step, and + wandb discards anything logged against a step it has already committed -- + without raising, and without a falsy return to check. + + This has bitten twice. ``grpo_sync`` was caught only because a real run + showed 85 logged keys and zero ``data_plane/*``. ``single_controller`` + carried the same code -- its docstring says it mirrors ``grpo_sync`` -- + and so carried the same bug, unnoticed, because no run exercised it. + + Asserting the invariant for every algorithm rather than for one call site + is the point: a third wiring would otherwise repeat it. Source order is + the only observable, because the drop happens inside wandb where a fake + logger sees a perfectly ordinary call. + """ + import pathlib + + import nemo_rl + + algorithms = pathlib.Path(nemo_rl.__file__).parent / "algorithms" + checked = [] + for path in sorted(algorithms.glob("*.py")): + # Comments mention the flag too, so they cannot be part of the search. + source = "\n".join( + line + for line in path.read_text().splitlines() + if not line.lstrip().startswith("#") + ) + if "_log_data_plane_metrics(" not in source: + continue + if "step_finished=True" not in source: + continue + checked.append(path.name) + # rindex, not index: the first occurrence is the *definition*, which + # naturally precedes everything. The call site is what has to come + # before the commit, and it is the last occurrence. + assert source.rindex("_log_data_plane_metrics(") < source.index( + "step_finished=True" + ), ( + f"{path.name}: data-plane metrics are logged after the " + "step_finished=True commit, so wandb will discard them" + ) + assert len(checked) >= 2, f"expected sync and single-controller, got {checked}"