Repository navigation
[PD] Align defensive protocol behavior across Mooncake, NIXL, and Mori - #35281
Conversation
Co-authored-by: Cursor <cursoragent@cursor.com>
There was a problem hiding this comment.
Reviewed each area against the current sources and the Common planner. Confirmed as claimed:
- Area 1: the Mori ABORT handler mirrors NIXL (record failure, mark
Failed, no ACK), and recognizing ABORT beforeMORI_GUARDmatches the Sender, sinceCommonKVReceiver._send_abort_notificationemits the unguarded 4-frame[ABORT, room, ip, port]; - Area 1: the heartbeat prerequisites all exist.
prefill_info_table,session_pool,addr_to_rooms_tracker,heartbeat_interval, andmax_failuresare initialized in the CommonKVManager Decode branch,CommonKVReceiver.__init__populatesaddr_to_rooms_tracker, and Mori already maintainsroom_to_bootstrap_addr/_cleanup_room_tracking; - Area 4: the MHA/MLA slicing is branch-for-branch identical to
get_mha_kv_ptrs_with_pp/get_mla_kv_ptrs_with_pp(same-length fast path, modulus draft V-offset, global offset otherwise), andis_hybrid_mla_backendis set on the Common base (common/conn.py:160), so thesend_kvcachebranch is safe; - Area 2, NIXL side: the worker guard matches Mooncake's exact condition plus the
_staging_outstandingremoval; NIXL has no chunk tracing, so there is no trace-finish counterpart to miss.
Three findings and one note, grouped by the PR's areas.
2. The NIXL guard stops one line short
The success conclude block pops transfer_infos (nixl/conn.py:1366 before this PR) while request_status keeps Success until the Sender's clear(). A stale chunk processed after conclusion (a deferred staging re-enqueue landing behind the last chunk does this) therefore passes the new guard, then trips the retained assert room in self.transfer_infos directly below it; the worker's except records the exception and flips a concluded-Success room to Failed while poll() can still observe it. Mooncake's protection for exactly this window is the graceful reqs_to_be_processed = ... if kv_chunk.room in self.transfer_infos else []. Replacing the assert with that fallback finishes the claimed alignment and removes the flip-to-Failed window.
Minor, same area: Mori _run_chunk's new missing-room return is silent where the NIXL guard logs a debug skip; a matching debug line helps correlate with the Decode side. Behavior is otherwise fine, since poll() already finalizes on a missing room and Mori's update_status refuses to resurrect cleared rooms.
3. The Mooncake/NIXL commits keep the loop alive but do not validate
The PR body describes "minimum frame/tag/guard validation", but the Mooncake/NIXL commits wrap the loop body in except Exception: logger.exception(...) without validating anything. That keeps the thread serving, but it is not the Mori model, which validates guard and frame count before acting and logs a concise warning. Two practical costs:
- Exceptions raised mid-processing, after state was already mutated, are also swallowed. In Mooncake
decode_threadthe Success path mutatesprefill_response_trackerbefore the staging-handler block, so a raise there leaves the room parked inTransferringuntilwaiting_timeout, and the log cannot distinguish a malformed frame from a real bug; zmq.ZMQErroris also caught, so a closed socket or terminated context turns the loop into a hot spin, and every foreign packet costs a full traceback (NIXL's rejection mechanism is still the GUARD assert, now inside the except).
except zmq.ZMQError: break plus a rate-limited warning for parse failures would match Mori more faithfully. Mori's own loops share the ZMQError spin, so fixing all three here or explicitly deferring it to the common extraction both work; worth stating which this PR intends.
5. One initialization window remains
register_to_bootstrap() already runs inside CommonKVManager.__init__, so a Decode can reach the socket before any Mooncake-specific field exists, and self._staging_ctx is still created after start_prefill_thread() (the PrefillStagingContext() assignment further down the Prefill branch). A WATERMARK arriving in that window hits handle_watermark_msg(self._staging_ctx, ...) on a missing attribute:
- Before this PR, that killed the thread;
- After this PR, the new blanket except swallows it, and the WATERMARK is silently dropped, which stalls staging flow control.
Moving the _staging_ctx assignment above start_prefill_thread(), next to the three session fields, closes the window.
4. Layout note for later stages
Inherited rather than introduced: the modulus heuristic misclassifies when dst_total_layers is a multiple of the Prefill-local layer count with draft KV present (for example pp=2 with draft depth equal to half the target depth), in which case dst_v_offset = dst_total_layers indexes into the draft region. get_mha_kv_ptrs_with_pp has the same hole, so porting it faithfully is the right call for this PR; it belongs in the Stage-2 planner-extraction notes.
Interaction with #34977
Checked every hunk against the wire-test surface: this PR does not touch TransferInfo/KVArgsRegisterInfo.from_zmq, is_dummy semantics, or any frame layout. The Mooncake bootstrap still reads msg[3]/msg[7], and NIXL still strips GUARD before from_zmq; the parsing code is only re-indented into the try blocks. The truth-table tests in #34977 call from_zmq directly, so the two PRs land independently in either order; no rebase needed on my side.
rishabhsinha17
left a comment
There was a problem hiding this comment.
Re-read at 8388c35 after the merge with main. Status of the earlier findings:
- Area 2: resolved.
transfer_infos.get(room)with the missing-room return replaces the assert, which closes the concluded-Success-to-Failed window. - Area 3: the loop-wrapping
except Exceptionblocks are gone in this revision, so the mid-processing swallow and the ZMQError hot spin no longer apply here. The PR body's "keep Mooncake/NIXL control loops serving after malformed messages" still describes the earlier revision; worth updating to match. - Area 5: still open.
start_prefill_thread()now runs at mooncake/conn.py:214 with the three session fields ahead of it, butself._staging_ctxis still assigned at line 242, after the thread is live. A WATERMARK that arrives in that window hitshandle_watermark_msg(self._staging_ctx, ...)on a missing attribute. Moving the_staging_ctxassignment above the thread start closes it.
#34977 interaction unchanged: the new hunks touch only the Mori ABORT handler's msg[0]/msg[1] reads, not from_zmq or the frame layouts, so the truth-table tests still land independently.
| except Exception: | ||
| logger.exception("Failed to parse transfer info message") | ||
|
|
||
| def _handle_abort_notification(self, msg: List[bytes]) -> bool: |
There was a problem hiding this comment.
Do we need an abort ack for Mori as well?
There was a problem hiding this comment.
Eventually yes.
What blocks it is that MoriKVManager has no _staging_outstanding, so the shared drain-ack helper cannot be reused as is:
def _maybe_ack_drained_abort(self, room: int) -> None:
if self._staging_outstanding.get(room, 0) > 0:
return
...Closing the gap needs changes on both sides: the prefill side has to register the ack target and emit the ack once the room drains, and mori's _start_decode_thread has to recognize ABORT_ACK
This PR only does the first step — making the prefill side recognize the decode-side ABORT at all, which previously was dropped as malformed because MORI_GUARD was validated first. I will handle the full ack path in a follow-up PR.
| elif self.disaggregation_mode == DisaggregationMode.DECODE: | ||
| self.room_to_bootstrap_addr: Dict[int, str] = {} | ||
| self._start_decode_thread() | ||
| self._start_heartbeat_checker_thread() |
There was a problem hiding this comment.
No. #36160 only touches the PREFILL branch of Mori. This change is in DECODE branch and just wires in the existing common heartbeat checker so decode can detect a dead prefill node and fail the affected rooms instead of waiting out timeout.
|
/tag-and-rerun-ci |
# Conflicts: # python/sglang/srt/disaggregation/mori/conn.py
Resolved two conflicts introduced by upstream: * mori/conn.py: sgl-project#29133 landed an equivalent Mori ABORT handler (_TAG_ABORT / _handle_abort_message) that holds transfer_lock across the status read-modify-write and returns early on an already-Failed room. Kept upstream's handler and dropped this branch's _handle_abort_notification, which duplicated it without the lock. sgl-project#29133 also covers the Mori _submit_kv_transfer cleared-room guard this branch carried. * nixl/conn.py: kept the cleared-room guard's room_transfer_infos local alongside the packed_source_by_dcp_rank map added by sgl-project#35762; the two changes are adjacent but independent.
|
/tag-and-rerun-ci |
sgl-project#35281) Co-authored-by: Shangming Cai <csmthu@gmail.com>
Background
#34692 added the missing NIXL Prefill bootstrap timeout. This PR continues Step 1 of #34510 by aligning clear, low-risk defensive protocol behavior across Mooncake, NIXL, and Mori without changing the Common state machine, wire format, completion channel, or Transport architecture.
These changes preserve the current backend structure. They only add missing behavior or fix local race/layout issues. Extracting the common protocol and Transport remains work for later stages.
Outline
This PR groups the changes into five areas:
1. Mori peer and request failure handling
Related commits:
b285ea79ce—align Mori Decode heartbeat62c13f66ba—align Mori abort handling with common protocolMori Decode did not start the existing Common heartbeat checker, so a dead Prefill would not trigger route/cache eviction or active-room failure. This PR only wires in the existing
_start_heartbeat_checker_thread(), matching Mooncake and NIXL.Common Receiver sends an
ABORTframe without a Mori guard, while Mori Prefill previously validatedMORI_GUARDfirst and discarded ABORT as malformed. This PR recognizes the existing Common ABORT frame before guard validation:Failed;This only fixes ABORT delivery. It does not define ACK, cancel, or quiescence semantics. Terminal-state and safe-release behavior remains deferred to the later RFC stages.
2. Skip stale work after room teardown
Related commit:
46056556a5—skip cleared rooms in NIXL and Mori workersChunks already queued do not disappear when the Sender clears a room.
transfer_infosassertion;add_transfer_request()and subsequent terminalization.Mooncake already checks at
transfer_worker()entry whether the room still exists or isFailed. On a match, it finishes the current chunk trace, removes residual_staging_outstandingstate, and continues without reading that room's transfer metadata.This PR makes the other two backends discard queued tasks for missing rooms at worker entry:
_staging_outstandingcount, and continues;_run_chunk().The stale task therefore cannot read deleted
transfer_infos, submit to the Transport, advance state, or trigger terminalization. Existingclear()paths remain responsible for normal room cleanup; these guards only consume queue items that arrive after teardown.When the common protocol is extracted, this should be replaced by unified idempotent room teardown/tombstone semantics. For now, the backend-local guards avoid expanding the scope into Common state-machine changes.
3. Isolate malformed control messages
Scope update: this change has been reverted from the final diff.
Since this branch was created, #35049 and #35360 have added more complex deferred-release and ABORT/ABORT_ACK state transitions to the same Mooncake/NIXL control loops. Keeping the original broad per-message exception boundary while reconciling those changes would mix malformed-frame handling with partially applied abort, socket, and engine failures.
Rather than expanding this PR to redesign those failure boundaries, we reverted the change and plan to handle malformed control messages together with the later protocol-unification work, based on the latest control paths.
Related commit:a716ceaac7—harden Mooncake and NIXL control message parsingMooncake/NIXL control loops serve multiple rooms, but previously indexed, decoded, or unpacked frames directly. An empty message, truncated ABORT, invalid guard, or unexpected frame count could terminate the long-lived thread.Mori already treats one message as the exception boundary: its Prefill bootstrap loop and Decode status loop validate guard/frame input insidewhile True, log an error for the current message, and continue receiving the next one. A bad frame therefore does not remove control-plane service from later rooms.This PR aligns Mooncake and NIXL with that behavior by adding minimum frame/tag/guard validation at the Mooncake Prefill/Decode and NIXL Prefill/Decode-staging message entries. Empty messages, truncated fields, and invalid encodings are logged and skipped without terminating the long-lived control thread.4. Align Mori KV layout planning
Related commits:
48c5e1569d—align Mori speculative MHA layout with common planner081186c4c8—align Mori same-PP and hybrid MLA layout planningTogether, these commits fix Mori descriptor slicing that differed from the existing Common planner.
The first part handles Decode-only speculative KV layout:
Mori previously split descriptors in half and could select a draft region as main V. This PR uses the same V-offset rule as
get_mha_kv_ptrs_with_pp().The second part adds:
These fixes remain inside Mori for now because they operate on Mori
MemoryDesclists, and planner/materialization are still coupled in Mori Manager. A later stage can extract target/draft, same-PP, and MHA/MLA logical layer selection into a Common helper, leaving Mori responsible only for mapping the common plan toMemoryDescandbatch_write().5. Initialize Mooncake state before starting the control thread
Related commit:
12a88dbcfe—initialize Mooncake session state before control threadMooncake Prefill previously started its receive thread before creating state that the thread accesses:
If Decode registration arrived immediately after thread startup, it could hit an initialization race. This PR only moves these three fields before
start_prefill_thread()while preserving all other queue, worker, probe, and receive ordering.This is intentionally a minimal ordering fix rather than a Step 1 reorganization of backend thread lifecycle. When the common Manager/Transport is extracted, control-thread initialization, startup, shared state, and shutdown should be managed together.
Testing
CPU red/green regression tests are kept on a dedicated test branch:
Coverage includes:
The tests pass with the fixes and reproduce the original failures against the corresponding unfixed source.
CI States
Latest PR Test (Base): ✅ Run #33297379569
Latest PR Test (Extra): ❌ Run #33297379432
Latest PR Test (AMD ROCm 7.2): ❌ Run #33297379594