feat(rollout): integrate MInf ledger capture with checkpointing - #4129
Conversation
Signed-off-by: Laura Dang <laurad@nvidia.com>
…very.sh Signed-off-by: Laura Dang <lauradang.2000@gmail.com>
Signed-off-by: Laura Dang <lauradang.2000@gmail.com>
Signed-off-by: Laura Dang <lauradang.2000@gmail.com>
Signed-off-by: Laura Dang <laurad@nvidia.com>
Signed-off-by: Laura Dang <laurad@nvidia.com>
Signed-off-by: Laura Dang <laurad@nvidia.com>
Signed-off-by: Laura Dang <laurad@nvidia.com>
… cache Extract the vLLM worker's _fetch_chain_prefix and _resolve_admission_prefix into tq_token_sink.py as ChainPrefixCache and resolve_admission_prefix, and make the worker methods one-line delegates. TQMegatronPromptPreparer now resolves staging chains through the same pair, so both backends share one cached TQ read (256 entries keyed by the chain's last staging key, deepest cached key bounds the fetch to the uncached suffix). The preparer reads the splice boundary from the request-metadata keys the Megatron chat endpoint writes (prefix_splice_suffix_token_ids, prefix_splice_boundary_token_id). The key strings are spelled out here rather than imported so this module stays importable in the vLLM worker and finalizer environments; a test asserts parity with Megatron's constants when Megatron is importable. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
Signed-off-by: Laura Dang <laurad@nvidia.com>
033f7a7 to
5bf90aa
Compare
|
Auto-sync is disabled for ready for review pull requests in this repository. Workflows must be run manually. Contributors can view more details about this message here. |
❌ Submodule Fast-Forward Check FailedCheck based on commit: 5bf90aa (PR #4129 from ❌ Submodules that need attention:Gym: ❌ Commits have DIVERGED from a common ancestor Please ensure all submodule commits are fast-forwards of the main branch before merging. |
lauradang
left a comment
There was a problem hiding this comment.
Team review of PR #4129 (7 agents: rl-expert, expert-gym, expert-megatron-lm, bug-finder, test-agent, design-reviewer, devil-advocate). Static review only: this host is macOS without GPUs and the lockfile is Linux-only, so no tests or linters were run; every claim about upstream behavior is linked to source at the pinned SHAs.
Not landable at HEAD — five deterministic blockers, all inline:
- The Megatron capture path programs against Megatron-LM hooks (
payload_stager,prompt_preparer,PREFIX_SPLICE_*) that exist only on unmerged Megatron-LM #7015; the Megatron-Bridge pointer is unchanged, sosetup_token_captureraises on every Megatron capture run (megatron_worker.py). - The Gym pointer
c3d7cf4fis not reachable from any NVIDIA-NeMo/Gym ref and Gymmainis wire-incompatible with this RL (3rdparty/Gym-workspace/Gym). PY_EXECUTABLES.MCORE_GYMfails the existingtest_every_extras_py_executable_is_wired_to_an_actor(virtual_cluster.py).- The runtime
ACTOR_ENVIRONMENT_REGISTRYswap is inert in the official container, whose prebuiltMegatronPolicyWorkervenv lacksnemo_gym(setup.py). test_prefix_splice_keys_match_megatron_constantshard-fails the Nemo_Gym shard, and the barecheckpointing.save_data_plane=trueoverride fails the streaming-recovery L1 test.
Lint: not run here; two findings (import order in megatron_generation.py, a stray blank line in test_checkpointing.py) would fail pre-commit run --all-files in CI.
The design itself is sound where it is checkable: the durability boundary is identical across backends, TQMegatronTokenStager.stage() fails closed, ChainPrefixCache / resolve_admission_prefix are a clean dedup, the set_generation_epoch rank-0 fan-out is correct, and all Gym API usage matches c3d7cf4f.
Informational (not staged as findings, below the confidence bar): at the #7015 tip both engine hooks run synchronously on the MP-coordinator engine loop, so each captured request costs one or two ray.get TQ round trips during which that DP replica does not step; vLLM offloads the same fetch via asyncio.to_thread. Worth a capture-on/off throughput number once the pin lands.
Devil's advocate: 24 confirmed, 1 disputed (a claimed isort collapse of the reassembler import), 7 downgraded below threshold.
Generated by Claude Code
- Declare the nemo_gym extra on MegatronPolicyWorker in actor_environments.py (the venv source of truth) instead of swapping ACTOR_ENVIRONMENT_REGISTRY at runtime, which the prebuilt container venv ignored; drop PY_EXECUTABLES.MCORE_GYM and the ModuleNotFoundError string-match remediation. - Gate backend=megatron token capture at setup on the MInf capture hook protocols (RequestPayloadStager / RequestPromptPreparer) from Megatron-LM PR #7015, failing with a NotImplementedError that names the dependency while the Megatron-Bridge pin predates it. - Point the Gym submodule at Gym PR NVIDIA-NeMo#2823 (fc08bf19), which is reachable from NVIDIA-NeMo/Gym, fast-forwards from Gym main, and nests ng_capture under request_metadata as the Megatron endpoint requires. - Skip test_prefix_splice_keys_match_megatron_constants when the pinned megatron-core lacks the constants, and importorskip megatron.core in the Megatron hosting test, so the Nemo_Gym shard skips instead of aborting. - Restore the ++ Hydra override for checkpointing.save_data_plane in the streaming recovery script (the key is absent from the Gym config chain). - Document the Megatron-LM #7015 dependency in the design doc. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test a9b17f2 |
- Add design-docs/token-capture-ledger.md to the docs toctree and drop links to guides that do not exist yet (Sphinx treats both as errors). - Read vllm_cfg / mcore_generation_config through a dict cast in the token-capture validation so pyrefly does not reject the TypedDict keys. - Apply ruff formatting to megatron_worker.py and test_checkpointing.py. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
✅ Submodule Fast-Forward Check ResultsCheck based on commit: faf4102 (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
✅ Submodule Fast-Forward Check ResultsCheck based on commit: 7813604 (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
…e hosting - vLLM worker: pass the resolved prefix straight to Gym's begin_call and drop the duplicated staging_chain / prev_len checks; tests match Gym's error text and read the prefix off the ActiveCall. - Delete the driver-side #7015 hook gate (_require_minf_capture_hooks) and its tests; pin the requirement as an mcore-lane unit test that fails, rather than skips, if a Megatron-Bridge bump drops the engine hooks. - Import the prefix-splice field names from megatron-core instead of keeping copies in tq_token_sink. - Remove the unused LOGGER / logging import in megatron_generation. - Collect test_interfaces.py in the three vLLM L0 lanes. - Add a backend parity test that drives both the vLLM and the Megatron capture glue through Gym's worked_example rollout (steady and mid-request refit) and requires byte-identical staged rows and commit coordinates. - Sibling-recovery functional test gates phase 2 on finalize/invalid_row_rate and finalize/capture_poisoned_rollouts staying at zero. - New SingleController nightly sibling of the Megatron-inference async Gym recipe with token_capture.enabled=true and the same finalizer gates. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test bce284e |
✅ Submodule Fast-Forward Check ResultsCheck based on commit: bce284e (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
The InferenceClient's ZMQ socket is read by its listener task on the inference loop thread, and ZMQ sockets are not thread safe. Marshal set_generation_epoch onto that loop via run_coroutine_threadsafe, the same way _sleep()/_wake() already do. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
✅ Submodule Fast-Forward Check ResultsCheck based on commit: d9a559f (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
_assemble_receipt validated the whole manifest in one list comprehension and shipped every raw row regardless. With one malformed row the finalizer's RolloutReceipt.model_validate failed on that same row and rejected the receipt as invalid_receipt with no staging keys, so the good rows' staged TQ entries leaked until the restart sweep and the terminal_selection=None path was unreachable. Validate row by row, mask the rollout as invalid_manifest_row, and ship only the rows that parsed. The finalizer now rejects as rollout_failed:invalid_manifest_row with the good rows' staging keys. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
TQMegatronPromptPreparer and TQMegatronTokenStager lived in nemo_rl/data_plane/tq_token_sink.py and pulled in nemo_rl.models.generation.openai_server_utils at module scope. That made every importer of tq_token_sink, including the finalizer actor, execute nemo_rl/models/generation/__init__.py and load the vLLM and TRT-LLM config modules. Move both classes to nemo_rl/models/generation/megatron/token_capture.py, mirroring how the vLLM backend keeps its capture glue in vllm_worker_async.py and imports only TQTokenSink/TQTokenSource from data_plane. Drop the local MegatronPayloadStageResult copy in favor of Megatron-LM's RequestPayloadStageResult, imported inside the method like RequestPromptPreparationResult already is. ChainPrefixCache and resolve_admission_prefix stay in tq_token_sink since both backends use them. Importing rollout_reassembler no longer loads nemo_rl.models.generation. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test 4652e32 |
Resolve the uv.lock conflict by taking main's lockfile and relocking with the Dockerfile-pinned uv (0.11.28). The remaining lock deltas against main (hydra-core floor, flashinfer 0.6.18.post1 on the flashinfer index, and its cutlass-dsl transitive) come from this branch's Megatron-Bridge bump. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test 1dfa5bf |
✅ Submodule Fast-Forward Check ResultsCheck based on commit: 1dfa5bf (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
Replace the remaining bare "ng_capture" / "ng_commit_coords" literals in TQMegatronPromptPreparer.prepare_prompt and TQMegatronTokenStager with NG_CAPTURE_FIELD / NG_COMMIT_COORDS_FIELD exported by nemo_gym.token_id_capture, so a Gym rename fails at import time instead of silently declining every request. Addresses NVIDIA-NeMo#4129 (comment) Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
✅ Submodule Fast-Forward Check ResultsCheck based on commit: 002023c (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
CI's lint job fails when a file type-checks clean but is absent from pyrefly.toml's project-includes; the new Megatron capture module was missing, which failed Lint check and skipped every downstream test job on the previous PR head. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test 5f049a3 |
✅ Submodule Fast-Forward Check ResultsCheck based on commit: 5f049a3 (PR #4129 from ✅ Submodules that are properly updated:Gym: ✅ PR branch is ahead of main branch (fast-forward) All submodule changes look good! ✨ |
Conflict resolutions: - Megatron-Bridge submodule: take main's 1f8873bb (NVIDIA-NeMo#4139); it already contains the PR's 4d472695 gather_output fix. - community_import.py / test_community_import.py: take main, which removed the _prefer_nvrx_for_dist_ckpt_save shim the PR had extended with a signature check. - megatron_policy_worker.py: take main's unconditional FileSystemWriterAsync import and direct cleanup_tensor_caches() call. - L1_Functional_Tests_SingleController: main split the script into _1/_2/_3; the PR's Megatron sibling-recovery variant now lives in _3 next to its vLLM twin. - pyproject.toml: keep the PR's comment on the flashinfer index URL spelling. - uv.lock: relocked with uv 0.11.28 (Dockerfile pin); result matches main. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test 40920dd |
…dget - test_backend_capture_glue_reproduces_the_gym_worked_example compared coords against every CallRecord field except mode/response_id, but CallRecord also carries attribution-only fields (admitted_at, output_fingerprint, continuation_fingerprint, fingerprint_version) that CommitCoords never has, so the lookup raised KeyError on whichever one the set yielded first. Compare on the fields the two models share. - The new token-capture nightly cost 48 GPU-hours and pushed the suite from 4698 to 4746, over the 4720 cap asserted by test_nightly_compute_stays_below_4720_hours. Cut it to 3 steps / 82 min (21 GPU-hours) so the suite lands at 4719. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
|
/ok to test 628b7bd |
terrykong
left a comment
There was a problem hiding this comment.
approved. @lauradang merge pending validating this with a real run
…-replay Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
PR NVIDIA-NeMo#4129 landed on main as squash e073dab, whose tree matches the 628b7bd head merged in the parent commit. Conflicts were resolved by three-way merging against e073dab so the result is main plus only the routing-indices-replay changes. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Laura Dang <laurad@nvidia.com>
Summary
CaptureAdmission, token-free lineage ledger, and the shared TQ staging contractOffloadedRequestPayloadthroughRolloutTokenCaptureandTQTokenSinkCallRecordonly after that acknowledgementRolloutReceiptat rollout end, then fetch, verify, and linearize staged records into the canonical GRPO sample during finalizationToken capture flow
The two serving backends use different prompt-preparation and generation hooks, but converge at the same durability boundary. A model call becomes eligible for lineage resolution only after its canonical token delta is staged in TQ. Capture failures write a failure row instead of a committed call record, causing finalization to reject or mask the rollout rather than train on incomplete lineage.
The source diagram is tracked in
docs/assets/token-capture-ledger-queue-data-flow.dot; the full invariants, terminal-selection rules, cleanup behavior, and fail-closed semantics are documented indocs/design-docs/token-capture-ledger.md.Dependencies
amahishi/partial-rollout-telemetry-v3OptionalRolloutReceipt.terminal_selection: NVIDIA-NeMo/Gym#2823 (submodule currently pinned to its head37dc751f; must land on Gymmainbefore this merges, then re-pin and regenerateuv.lock)RequestPayloadStager/RequestPromptPreparer, typedRequestPromptPreparationResult, andoffload_paramscarryingng_capture,template_prefix_token_ids,eos_token_id): NVIDIA/Megatron-LM#7015, pinned transitively through NVIDIA-NeMo/Megatron-Bridge#6101. Both dependencies must land before this PR merges, then the Bridge submodule must be re-pinned to its merged commit anduv.lockregenerated.Testing
Local test limitation
The RL lockfile supports Linux x86_64/aarch64 only, so the focused RL pytest files cannot run natively on this macOS host. An ephemeral dependency run progressed through PyTorch, Transformers, and pyzmq before reaching the NVIDIA-only
pynvmlruntime dependency; Linux CI provides the authoritative RL test coverage.