[Lumen-RL] Receive trained weights over RCCL, transactionally - #2298
i-chaochen wants to merge 16 commits into
Conversation
🏷️ CI GuideRuns automatically on every eligible PR before approval:
Heavy model tests:
|
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Critical weight-transaction and collective-coordination issues remain unresolved.
Get a fresh assessment by requesting another Copilot review.
Review effort: Lite
Findings: 5
Open (8)
Bucket decoder accepts duplicate names and overlapping byte ranges · New Uncoordinated receive failures can deadlock collective broadcasts · New Commit coverage uses source names instead of canonical parameter names · New FP8 coverage is recorded when requantization writes nothing · New Abort fails to clear pending expert relayout state · New Failed DP engines omit required per-TP rank results · New Update failures omit error responses and cause synchronous timeouts · New Contract test silently skips when the external checkout is unavailable · New
What changed in this PR
Adds transactional RCCL weight streaming, collective RPC routing, and capability discovery for RLHF rollouts.
Changes:
- Adds independent process groups and RDMA weight reception.
- Adds transactional weight application and serving fences.
- Adds collective RPC routing and capability negotiation.
- Expands CPU-only transport and wire-format tests.
| File | Description |
|---|---|
tests/test_rdma_weight_receiver.py |
Tests RDMA wire decoding and contracts. |
tests/test_collective_rpc_transport.py |
Tests TP RPC transport. |
tests/test_collective_rpc_dp.py |
Tests DP routing and correlation. |
tests/test_collective_rpc_dispatch.py |
Tests RPC dispatch behavior. |
tests/test_capabilities.py |
Tests capability negotiation. |
atom/utils/independent_process_group.py |
Creates isolated distributed groups. |
atom/rollout/weight_updater.py |
Applies transactional weight updates. |
atom/rollout/rdma_weight_receiver.py |
Receives and decodes RCCL weight streams. |
atom/rollout/model_runner_ext.py |
Integrates RLHF runner extensions. |
atom/rollout/capabilities.py |
Reports worker capabilities. |
atom/rollout/async_engine.py |
Exposes RPC and capability APIs. |
atom/model_engine/engine_utility.py |
Handles RPC dispatch and utility responses. |
atom/model_engine/engine_core_mgr.py |
Routes DP-level RPC responses. |
atom/model_engine/collective_rpc.py |
Defines RPC wire types and routing. |
atom/model_engine/capabilities.py |
Aggregates capabilities. |
atom/model_engine/async_proc.py |
Collects per-rank RPC replies. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
ZhangDanyang-AMD
left a comment
There was a problem hiding this comment.
The end goal is to run a 200-step ATOM + Megatron Qwen3 disaggregated case and match vLLM Disaggregated 2-Node Megatron + vLLM RDMA Deployment then run Running the eight examples from the release image cases 1 and 3 for 1 step to confirm other functionality is unaffected.
| """ | ||
| if not hasattr(self, "_expert_relayout_pending"): | ||
| self._expert_relayout_pending = {} | ||
| return self._expert_relayout_pending |
There was a problem hiding this comment.
abort_weight_update clears only the packed accumulator. A later reload can inherit expert-relayout state from the failed transaction. See 1177-1187
| # Release DP load bookkeeping now: an aborted seq may never emit a | ||
| # finished STREAM output, so relying on the finish path alone would leak | ||
| # its in-flight count. _release_seq_load is idempotent. | ||
| self._release_seq_load(req_id) |
There was a problem hiding this comment.
engine_utility.py:312-329 writes an unconditional response while it is consumed synchronously so non-RPC responses enter the legacy queue.
7be3e9e to
a67bc0b
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved critical RPC routing, RDMA stream validation, transaction coordination, and process-group isolation issues remain.
Review effort: Lite
Findings: 5
Open (8)
Concurrent RPC replies are misrouted between request waiters · New Early transaction rejection can desynchronize collective broadcasts · New End marker accepts invalid metadata and payload sizes · New DP stride uses logical instead of physical worker count · New Supplied stores are not namespaced by group name · New RDMA lifecycle methods are missing from capability discovery · New FP8 support is advertised based on method presence · New Contract test silently skips when the external checkout is unavailable
Resolved since last review (7)
Abort fails to clear pending expert relayout state FP8 coverage is recorded when requantization writes nothing Commit coverage uses source names instead of canonical parameter names Uncoordinated receive failures can deadlock collective broadcasts Bucket decoder accepts duplicate names and overlapping byte ranges Update failures omit error responses and cause synchronous timeouts Failed DP engines omit required per-TP rank results
| for rank, output_queue in enumerate(self.rpc_outputs_queues): | ||
| results.append( | ||
| self._await_rank_reply(rank, output_queue, func_name, payload, deadline) |
There was a problem hiding this comment.
Agreed; the docstring promised something the transport can't do. Every caller reads the same per-rank reply queues, so two outstanding calls take each other's replies, drop them as stale, and time out. Fixed in a23ad9de8 by serializing: the request is sent and every rank's reply collected under one lock, so a second call's request reaches the workers only after the first has its answers. I serialized rather than adding a router because the only caller is the EngineCore busy loop, which issues one call at a time, and nothing else reads these queues (call_func and KV aggregation have their own). The docstring now says what request_id is for: recognising a late reply from a call that already timed out, so it isn't taken as the next call's answer. Test: test_concurrent_callers_each_get_their_own_replies.
| _ADVERTISED_METHODS = ( | ||
| # weight sync | ||
| "update_weights", | ||
| "update_weights_from_shm", | ||
| "update_weights_from_ipc", | ||
| # memory lifecycle | ||
| "release_memory", | ||
| "resume_memory", | ||
| "clear_kv_cache", | ||
| # hidden-state extraction | ||
| "configure_hidden_states", | ||
| # capability discovery itself | ||
| "get_worker_capabilities", | ||
| ) |
There was a problem hiding this comment.
Agreed, fixed in a23ad9d. init_rdma_weight_group, receive_weights_rdma and destroy_rdma_weight_group are now advertised, plus get_weight_update_status, which an orchestrator needs to see the fence after a failed stream. Each is still reported only where the runner defines it. Covered by test_the_rdma_lifecycle_is_advertised_alongside_the_feature. On #2298, where the receiver lives, test_the_lifecycle_the_real_mixins_define_is_advertised also runs discovery over the real mixins, so a misspelt name can't silently drop out.
| if callable(getattr(self, "_is_fp8_param", None)): | ||
| features.add("fp8_weight_update") | ||
| if callable(getattr(self, "receive_weights_rdma", None)): |
There was no way to invoke an arbitrary ModelRunner method on every rank. `call_func` broadcasts to all TP ranks but returns only rank 0's value, and `_UTILITY_HANDLERS` is a closed registry, so a consumer outside the engine could not reach a worker method at all. The generic path: - `RpcPayload`/`RpcResult` in a new `collective_rpc` module, kept free of heavy imports so the dispatch and negotiation layers stay importable without an AITER build -- which is also how the non-GPU runner sees them. - `AsyncIOProc` recognises a lone `RpcPayload` and answers on every path. A missing method, a raising target and an unpicklable return each used to strand the caller in an untimed queue get; all three now come back as an `RpcResult` carrying `error`. - A reply channel per rank, so "every rank finished" is observable instead of inferred from rank 0. - `AsyncIOProcManager.collective_rpc` collects one reply per rank against a deadline, probing liveness so a dead rank is named promptly rather than at the timeout. - `CoreManager.collective_rpc` fans out across DP engines and correlates replies by request id. `broadcast_utility_command_sync` matches by queue position, so two overlapping callers take each other's replies; every other utility command keeps that legacy path untouched. - `AsyncLLMEngine.collective_rpc`, mirroring vLLM's signature. Also fixes three pre-existing hangs of one shape -- a command that produced no `UTILITY_RESPONSE` left `broadcast_utility_command_sync` blocked for its full 300s and then reported a timeout naming the command but not the cause. An unrecognised command was logged and dropped; `update_weights` never answered; `abort_request` never answered, including on its early returns. Capability discovery rides the same path: a worker reports its protocol version, position, methods and features, and the engine intersects features across ranks. Intersection rather than union is load-bearing, because a feature present on only some ranks cannot be driven by a collective -- the ranks that had it would block in one collective while the rest went elsewhere, deadlocking the group with no error. Methods are advertised from a fixed list rather than a dir() sweep, so a refactor is not a capability change and private helpers are not offered. 75 new CPU-only tests. Full suite 5637 -> 5712 passed, with failures, skips and errors unchanged. Co-authored-by: Cursor <cursoragent@cursor.com>
Review follow-ups on the generic collective_rpc. A DP engine that is missing or fails now contributes one failed RpcResult per TP rank rather than a single tp_rank=-1 placeholder. The result list is always DP x TP long, so a position still identifies a rank. abort_request is fire-and-forget again. Making it always answer left replies nobody reads on CoreManager's shared response queue, and broadcast_utility_command_sync, which matches replies by position, takes them as its own. The sync path now refuses fire-and-forget commands up front instead, so waiting on one fails at once rather than after 300s. A utility handler that raises now sends an error reply before the exception propagates, and broadcast_utility_command_sync raises with that cause. update_weights answered only on success, so a rejected tensor left the caller waiting out its timeout for an engine that was already gone. Co-authored-by: Cursor <cursoragent@cursor.com>
main now ships a forward's block tables as appends to rows each worker caches (#2383). Every message a worker dequeues passes that decoder, and a generic call comes through it untouched. The one way to break it is collective_rpc("forward", ...): that goes around the encoder, which has to see every forward, while each worker's decoder still resets its rows and fails the next scheduled forward on a missing prefix. collective_rpc now refuses that name. The per-engine rank count is read with defaults, since main's new distributed-store tests build CoreManager from a bare stand-in config. Co-authored-by: Cursor <cursoragent@cursor.com>
a67bc0b to
79d2ea4
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved critical and moderate correctness issues remain in reload fencing, RPC handling, and capability reporting.
Review effort: Lite
Findings: 7
Open (10)
Successful non-commit reloads leave the weight update fence active · New Commit can succeed before synchronization faults are handled · New Supplied stores are not namespaced by group name DP stride uses logical instead of physical worker count End marker accepts invalid metadata and payload sizes Early transaction rejection can desynchronize collective broadcasts Concurrent RPC replies are misrouted between request waiters FP8 support is advertised based on method presence RDMA lifecycle methods are missing from capability discovery Contract test silently skips when the external checkout is unavailable
| # in place, so without this a failed stream would keep serving from a | ||
| # mix of two versions -- visible much later as unexplained divergence | ||
| # rather than as an error here. | ||
| self.assert_weight_update_ready() |
There was a problem hiding this comment.
Valid, fixed in d0785b7. The direct, SHM and IPC paths now lift the fence when a reload finishes (update_weights on return, SHM/IPC on the is_last bucket), provided no transaction is open. They verify no coverage, so the lift is skipped when a fused parameter is still waiting on shards or the reload matched nothing, and it is logged as a warning saying so. Verified coverage on those paths is the planned integrity follow-up.
| runner.abort_weight_update(expected_version, exc) | ||
| raise | ||
|
|
||
| torch.cuda.synchronize(device) |
Replies are reported per DP rank, DP-major then TP. Under pipeline parallelism CoreManager runs one engine per stage instead, each holding a slice of the layers, so every stage came back labelled as another DP rank with nothing to say which stage a result was from. DisaggCoreManager's two engines, the prefill and decode halves of one replica, were mislabelled the same way. Neither has a caller today -- RDMA weight sync and capability discovery both assume a DP x TP grid -- so the call is refused up front with the reason, rather than extended with a pipeline rank that nothing would read. Co-authored-by: Cursor <cursoragent@cursor.com>
A truthy but non-string method passed the handler's check and reached getattr on every TP worker. getattr's default covers AttributeError only, so the TypeError escaped _run_generic_rpc and ended the worker's loop: one malformed call cost the engine its workers. The handler now requires a string before it broadcasts, and the worker's lookup sits inside error handling of its own, so a name that cannot be resolved comes back as an RpcResult error either way. Co-authored-by: Cursor <cursoragent@cursor.com>
With barrier=True every rank waits for all the others before answering. One that died never arrives, so the survivors sat in Barrier.wait() with no timeout and never replied. The manager waited on them one rank at a time until the deadline and then reported timeouts -- for the dead rank too, whose turn came after the budget was spent and so never reached the liveness check. While waiting on a barrier call, the manager now breaks the barrier as soon as any rank is found dead, and each survivor answers with that instead of waiting. A dead rank is named even once the deadline has passed. A rank that is only slow is left to the deadline: it can still arrive, and a broken barrier stays broken for every later call. Co-authored-by: Cursor <cursoragent@cursor.com>
The feature was advertised whenever _is_fp8_param existed, and every RLHFModelRunner inherits that from WeightUpdaterMixin, so BF16 runners claimed FP8 weight updates too. It is now derived from the loaded parameters through the same helper, so the report agrees with what a weight update would actually do. Co-authored-by: Cursor <cursoragent@cursor.com>
79d2ea4 to
10dd72a
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved moderate issues can cause distributed hangs, rank misconfiguration, stream desynchronization, or incorrect update reporting.
Review effort: Lite
Findings: 7
Open (11)
Commit can succeed before synchronization faults are handled Successful non-commit reloads leave the weight update fence active Supplied stores are not namespaced by group name DP stride uses logical instead of physical worker count End marker accepts invalid metadata and payload sizes Early transaction rejection can desynchronize collective broadcasts Concurrent RPC replies are misrouted between request waiters SHM and IPC callers ignore FP8 update results · New FP8 support is advertised based on method presence RDMA lifecycle methods are missing from capability discovery Contract test silently skips when the external checkout is unavailable
Pre-existing code, but the collective_rpc change sits right below it, so reviewdog's diff-context filter reports it on this PR and fails the ruff check. The catch is intended: abort_request runs from client-disconnect handlers, which must not raise. Co-authored-by: Cursor <cursoragent@cursor.com>
The dead-rank check ran only on a poll that was about to keep waiting, so the deadline and died branches returned first. A rank found dead on the last poll -- or one whose turn came after the budget was spent -- left the others waiting at the barrier after the call had already given up on them, and every later call to them hung. The check now comes first in the empty-poll branch, ahead of both returns, and runs on every call rather than only barrier ones: a barrier call that timed out on a slow rank leaves the rest waiting, and if that rank dies afterwards, only a later call is there to notice. It is still a no-op until some rank has actually died. Co-authored-by: Cursor <cursoragent@cursor.com>
The handler checked the request id only for truthiness, so an unhashable one -- a list, say -- was broadcast, echoed back by every worker, and then raised inside RpcResponseRouter.route() on the manager's output thread. Nothing catches it there, so the thread ended and that engine's outputs went unread from then on. The handler's own error reply echoed the same id, so refusing it there alone was not enough; and an id-less collective_rpc reply fell through to the shared queue, where the next synchronous caller would have taken it as its own. The handler now requires a non-empty string id before it broadcasts, and the manager drops any collective_rpc reply without one instead of routing or queueing it. Co-authored-by: Cursor <cursoragent@cursor.com>
10dd72a to
283e6dc
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved reload-safety, stream-validation, rank-calculation, capability-discovery, and process-group issues remain.
Review effort: Lite
Findings: 8
Open (12)
Prefill path bypasses reload readiness fence · New Commit can succeed before synchronization faults are handled Successful non-commit reloads leave the weight update fence active Supplied stores are not namespaced by group name DP stride uses logical instead of physical worker count End marker accepts invalid metadata and payload sizes Early transaction rejection can desynchronize collective broadcasts Concurrent RPC replies are misrouted between request waiters SHM and IPC callers ignore FP8 update results FP8 support is advertised based on method presence RDMA lifecycle methods are missing from capability discovery Contract test silently skips when the external checkout is unavailable
collective_rpc refused forward but not exit, so a caller could run ModelRunner.exit() on every TP worker -- tearing down its NCCL groups, KV connector and CUDA graphs -- after which busy_loop's own "exit" check ended each worker's loop, behind a reply that read as success. The same path also reached async_proc_aggregation, which drains the KV connector's finished transfers into a reply the KV aggregator never sees, so the requests waiting on them would stall. These are now one reserved set, built from the worker loop's own KV names so a new one cannot be missed, and refused before anything is enqueued. Co-authored-by: Cursor <cursoragent@cursor.com>
283e6dc to
3fe9a00
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved process-group, transactional coordination, coverage, and capability issues block approval.
Review effort: Lite
Findings: 9
Open (13)
Pass correct arguments when creating RDMA process groups · New Prefill path bypasses reload readiness fence Commit can succeed before synchronization faults are handled Successful non-commit reloads leave the weight update fence active Supplied stores are not namespaced by group name DP stride uses logical instead of physical worker count End marker accepts invalid metadata and payload sizes Early transaction rejection can desynchronize collective broadcasts Concurrent RPC replies are misrouted between request waiters SHM and IPC callers ignore FP8 update results FP8 support is advertised based on method presence RDMA lifecycle methods are missing from capability discovery Contract test silently skips when the external checkout is unavailable
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved critical and moderate issues affect RDMA safety, transactional reload correctness, RPC routing, and capability discovery.
Review effort: Lite
Findings: 8
Open (11)
Unknown commands can deadlock peers by skipping sized broadcasts · New Rejected weight updates can abort an existing transaction · New Partial or empty reloads incorrectly lift the serving fence · New Pass correct arguments when creating RDMA process groups Commit can succeed before synchronization faults are handled Successful non-commit reloads leave the weight update fence active Supplied stores are not namespaced by group name Concurrent RPC replies are misrouted between request waiters SHM and IPC callers ignore FP8 update results FP8 support is advertised based on method presence RDMA lifecycle methods are missing from capability discovery
Resolved since last review (5)
| if command != _CMD_BUCKET or metadata_bytes <= 0 or payload_bytes <= 0: | ||
| # Unlike a bad bucket, this cannot be drained: without sizes | ||
| # there is no telling what the sender broadcasts next. | ||
| raise RuntimeError( | ||
| f"invalid RDMA weight header: command={command} " | ||
| f"metadata_bytes={metadata_bytes} payload_bytes={payload_bytes}" |
| except Exception as exc: | ||
| # Fence before re-raising: the parameters are now a mix of versions, so | ||
| # serving must stop until a later full reload succeeds. | ||
| runner.abort_weight_update(expected_version, exc) |
| matched_nothing = updated == 0 and counts["skipped"] > 0 | ||
| self._lift_partial_reload_fence( | ||
| "direct", complete=not incomplete and not matched_nothing | ||
| ) |
…ifecycle The TP-level collective_rpc said request ids let several calls be outstanding, but every caller reads the same per-rank reply queues, so two concurrent calls took each other's replies, dropped them as stale and timed out. Its one caller is the engine busy loop, one call at a time, so calls are now serialized under a lock, and the docstring says what the id is actually for: recognising a late reply from a call that already gave up. Capability discovery reported rdma_weight_receive with none of the methods that drive it. The advertised list now names the RDMA lifecycle and the weight-update status an orchestrator reads, each still listed only where the runner has it. Co-authored-by: Cursor <cursoragent@cursor.com>
Lets each TP worker join an independent process group alongside the trainer, so weights land straight in resident GPU memory -- no safetensors round-trip, no shared folder, no host bounce. The wire format is frozen by the sending side and matched byte for byte: a 4x int64 header [command, metadata_bytes, payload_bytes, version], then a uint8 JSON metadata list, then a uint8 payload each entry views into. Decoded tensors are views over the payload rather than copies, because a transfer already sized in tens of GB cannot afford to double its peak footprint. A test reads the sender's own constants from source so a drift on either side fails here instead of decoding 61 GB into garbage. Two things make this more than a memcpy: Fused parameters. `qkv_proj` and `gate_up_proj` combine several checkpoint tensors, and one fused parameter's shards can arrive in different buckets, so state has to survive across them. `update_weights`' application loop is extracted into `_apply_named_tensors` and shared, leaving finalisation -- expert relayout, KV clear, accumulator teardown -- with the caller. A bucketed stream does that once at commit; doing it per bucket would relayout a half-built parameter. In-place application. A stream that fails halfway leaves the model a mix of two versions, and inference would keep serving, quietly wrong. So a stream is a transaction: begin/apply/commit/abort, with monotonic versions so a replayed or out-of-order stream is refused, and commit-time coverage verification against named_parameters(). `forward` now calls `assert_weight_update_ready()`, which refuses to serve mid-reload or after a partial one -- the cost of a false stop is a raised error, against rollouts from a half-updated model surfacing much later as unexplained divergence. Coverage accounting credits `weight_scale` values that requantisation writes as a side effect. They are derived, never sent, so without that they look permanently missing and verification would reject every stream. `init_independent_process_group` is vendored into atom/utils rather than imported from the RL framework: a rollout container carrying only ATOM must still be able to build the group. 11 new tests. Full suite 5712 -> 5723 passed, failures/skips/errors unchanged; #2028's 34 weight-sync tests still pass, which is the check that matters for the shared-loop extraction. Co-authored-by: Cursor <cursoragent@cursor.com>
Review follow-ups on the RDMA receiver and its transaction. Coverage was recorded under the incoming checkpoint name and checked against named_parameters(). On every fused route the two differ -- q_proj lands in qkv_proj, one expert's gate_proj in w13_weight -- so verify_full_load rejected every reload of a fused model, Qwen3 included. Coverage is now kept by parameter object. A packed parameter counts once all of its shards have landed; an expert buffer once every expert's every shard has, read off the relayout bookkeeping before the relayout consumes it. A tied parameter counts under either name. FP8 requantisation now reports whether it wrote. A shape it could not shard, or a quant type it did not know, counted as written and credited the scale as well. abort clears the pending expert relayout, and begin starts clean. Left behind, an aborted stream's slices were relaid out again by the next reload, or failed it on shards the aborted stream never sent. A rank whose bucket fails now keeps receiving until the end marker rather than leaving the broadcasts, which hung the trainer and every other rank in the next one. The decoder rejects repeated names and overlapping byte ranges, and a reload refuses a weight sent twice. The sender contract test no longer reads a hardcoded checkout path: it runs against LUMENRL_ROOT when that is set, and also pins the header field order. Co-authored-by: Cursor <cursoragent@cursor.com>
…ce on recovery Review follow-ups on the receiver and its transaction. A begin that refused the stream -- a replayed or older version, or a reload already open -- ran before the drain logic, so the rank left without receiving it while the trainer and every other rank waited in the next broadcast. It now runs inside it: the stream is received, nothing is applied, and the refusal is raised at the end marker like any failure. The end marker was accepted with any sizes. The contract is [END, 0, 0, version]; one carrying sizes now fails the stream instead of committing it. The DP stride was the logical TP width, but an engine runs tp_world_size x prefill_context_parallel_size workers. Under PCP one DP engine's workers took the next one's ranks, and under simulated TP ranks went unclaimed. The stride is now the engine's own worker count, the one ModelRunner places devices by. Only a commit lifted the fence an aborted stream leaves, so recovering through the direct, SHM or IPC path left every later forward refused. A legacy reload that finishes -- its last bucket applied, no fused parameter left waiting on shards, and not one that matched nothing -- now lifts it, with a warning that those paths verify no coverage. Co-authored-by: Cursor <cursoragent@cursor.com>
…l requantisations receive_weight_stream synchronized the device after its try block, so an asynchronous fault in the bucket writes, or in commit's own finalisation, surfaced after commit had declared the version good, and nothing fenced serving. The synchronize is now inside the try, and a fault there aborts. The SHM and IPC loops counted every FP8 requantisation as updated, though the requantiser returns without writing on a shape it cannot shard or a quant type it does not know. A reload of nothing else then reported success and lifted the fence over the old weights. They now branch on its result, as the transactional path already did. The helper's docstring said group_name namespaces the store. The isolation between groups is _new_process_group_helper's, which keys everything under the group name whatever store it is given; the extra prefix is only for a store created here, as init_process_group does, and the trainer's copy does the same. Said so, rather than changing the key layout on one end of the group only. A contract test composes the real mixins, so an advertised name the receiver spells differently cannot drop out of discovery unnoticed. Co-authored-by: Cursor <cursoragent@cursor.com>
d0785b7 to
dc72ba1
Compare
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Unresolved RPC timeout, stream-handling, transaction-abort, fencing, and CI contract-validation issues remain.
Review effort: Lite
Findings: 7
Open (9)
Sequential reply collection causes false timeouts · New Partial or empty reloads incorrectly lift the serving fence Rejected weight updates can abort an existing transaction Unknown commands can deadlock peers by skipping sized broadcasts Commit can succeed before synchronization faults are handled Successful non-commit reloads leave the weight update fence active Concurrent RPC replies are misrouted between request waiters FP8 support is advertised based on method presence RDMA lifecycle methods are missing from capability discovery
| for rank, output_queue in enumerate(self.rpc_outputs_queues): | ||
| results.append( | ||
| self._await_rank_reply( | ||
| rank, output_queue, func_name, payload, deadline | ||
| ) |



Motivation
Stacked on #2297
ATOM rollout has no direct weight path from a trainer: weights arrive via
safetensors on shared storage or a ZMQ CUDA-IPC hop. This lets each TP worker
join an independent process group alongside the trainer, so weights land
straight in resident GPU memory.
Technical Details
The wire format is frozen by the sending side and matched byte for byte: a
4x int64 header
[command, metadata_bytes, payload_bytes, version], a uint8JSON metadata list, then a uint8 payload each entry views into. Decoded tensors
are views over the payload rather than copies, since a transfer already sized in
tens of GB cannot afford to double its peak footprint. A test reads the sender's
own constants from source, so drift on either side fails here rather than
decoding tens of GB into garbage.
Two things make this more than a memcpy.
Fused parameters:
qkv_projandgate_up_projcombine several checkpointtensors, and one fused parameter's shards can arrive in different buckets, so
state must survive across them.
update_weights' application loop is extractedinto
_apply_named_tensorsand shared, leaving finalisation -- expert relayout,KV clear, accumulator teardown -- with the caller. A bucketed stream does that
once at commit; per bucket it would relayout a half-built parameter.
In-place application: a stream that fails halfway leaves the model a mix of two
versions, and inference would keep serving, quietly wrong. So a stream is a
transaction -- begin/apply/commit/abort -- with monotonic versions so a replayed
or out-of-order stream is refused, and commit-time coverage verification against
named_parameters().forwardnow callsassert_weight_update_ready(), whichrefuses to serve mid-reload or after a partial one. The cost of a false stop is
a raised error; the cost of not stopping is rollouts from a half-updated model,
surfacing much later as unexplained divergence.
Coverage accounting credits
weight_scalevalues that requantisation writes asa side effect. They are derived, never sent, so without that they look
permanently missing and verification would reject every stream.
init_independent_process_groupis vendored intoatom/utilsrather thanimported from the RL framework, so a rollout container carrying only ATOM can
still build the group.
Note the coverage check currently gates the RDMA path only; extending it to the
shm and IPC transports is follow-up work.
Test Plan
11 new CPU-only tests covering the wire-format round trip, view-not-copy, and
rejection of every corrupt-metadata shape, plus a test pinning the sender's
constants. Re-ran #2028's weight-sync tests, which are the real check on the
shared-loop extraction.
Test Result
black --check .clean;ruffno findings on touched filesNot yet validated on hardware. The 9-rank run needs two GPU nodes
simultaneously and has not been schedulable yet.
Submission Checklist