Repository navigation
[PD] Share the prefill->decode failure notification across backends - #36612
Conversation
Align Mori prefill-side KV transfer with Mooncake/NIXL: enqueue room-keyed TransferKVChunk work units, run submit/wait/notify in MoriKVManager._process_transfer_chunk, and thin MoriKVSender to enqueue-only + status poll. Keeps early-send wait_event on the chunk and centralizes decode notification in the manager for a shared control-plane extraction path. Co-authored-by: Cursor <cursoragent@cursor.com>
| other endpoints never receive the remaining chunks either, so telling | ||
| only the endpoint that raised leaves the rest waiting for the timeout. | ||
| """ | ||
| infos = getattr(self, "transfer_infos", {}).get(bootstrap_room) |
There was a problem hiding this comment.
Avoid using getattr, it will silently hide some bugs.
There was a problem hiding this comment.
Fixed, thanks! I also went through the PR again and cleaned up the code to make sure it follows the repo's rules.
livingshade
left a comment
There was a problem hiding this comment.
LGTM, only some corner cases picked up by claude.
| # Raise only once every handle of this batch settled, not on the | ||
| # first ERR: a sibling still in PROC keeps writing into the | ||
| # decode's KV pages, and the failure path below tells the decode | ||
| # those pages are free. |
There was a problem hiding this comment.
Probably should add a timeout here? If the remote node down it might take very long to exit the loop. Also might need to change time.sleep(0) to something like time.sleep(0.00001)to avoid busy loop.
Difference from old path: Previous implementation is acceptable because only one room would experience decode time out of 300s, but here it will block all handles. So other room is also being blocked. I would say current implementation has a much larger blast radius using claudish.
There was a problem hiding this comment.
Fixed: bounded to 5s once a handle reports ERR, and on timeout the room fails locally without notifying decode, so it falls back to the waiting timeout exactly as before. Worth noting the no-ERR hang predates this PR: the old loop also spun until every handle was DONE, so only the mixed ERR+PROC case is a regression.
| bootstrap_room | ||
| ) | ||
| if ( | ||
| expected_response_num is None |
There was a problem hiding this comment.
We might want to add a warning log here?
There was a problem hiding this comment.
Fixed: split from the partial-response return and now warns.
| """Returns the status that was emitted, or None for a cleared room. | ||
|
|
||
| A room is sharded onto a single transfer worker, so this runs at most | ||
| once per room without any at-most-once bookkeeping. ``targets`` defaults |
There was a problem hiding this comment.
At most once is wrong for mooncake in its transfer_worker, probably should change to idempotent
There was a problem hiding this comment.
Fixed the claim.
Merging main after sgl-project#36160 landed re-added `_room_status_notified` and `_room_notify_lock` to `MoriKVManager.__init__` while every user of them stayed deleted, leaving dead state and an F821 on the `Dict` import that went with them.
A transfer worker serves several rooms, so waiting for every handle of a failed batch to settle stalled all of them whenever a sibling stayed in PROC until NIXL noticed its peer was gone. Back off and give up after 5s: the room fails locally, but decode is not told its pages are free, so a late write still cannot land on a reused page. The all-DONE path keeps its tight spin. `apply_prefill_status` folded a missing expected response count into the same silent return as a partial one; it now warns, since the alternative symptom is a 300s decode timeout with no cause in the log. `conclude_transfer` claimed at-most-once delivery. One worker owns a room, but a staging chunk deferred past the last one can fail after the room concluded Success, and reporting that failure is deliberate.
…tifi # Conflicts: # python/sglang/srt/disaggregation/mori/conn.py
The multiplicity is what a reader needs at this line; the deferred-chunk mechanism and the choice to report the later failure belong in the PR discussion, not next to the signature.
| logger.warning( | ||
| "No expected prefill response count for room %s", | ||
| bootstrap_room, | ||
| prefill_rank, | ||
| ) |
There was a problem hiding this comment.
should fix here, only one %s
The message lost its second placeholder while the argument stayed, so logging raised on format and the warning never reached the log.
| bootstrap_room, failure_reason or "KV transfer failed" | ||
| ) | ||
|
|
||
| self.update_status(bootstrap_room, status) |
There was a problem hiding this comment.
Maybe move this after _room_notify_targets?We should postpone update_status incase there is a race for self.transfer_infos.get(bootstrap_room) and .clear()
There was a problem hiding this comment.
Fixed. I also traced the incorrect ordering back to Mori's _notify_decode_for_room. When that logic was extracted into the common path, the same ordering was inadvertently propagated to all backends.
I've updated my local AGENTS.md and reviewed the PR again. For this kind of cross-backend refactoring, Mooncake should be treated as the reference implementation, since the other backends may already contain backend-specific inconsistencies or minor bugs.
|
|
/rerun-group disaggregation |
|
Results for 🚀 🚀 🚀 🚀 🚀 🚀 🚀 🚀 |
conclude_transfer looked up the room's decode endpoints after publishing the terminal status. Publishing it is what lets the scheduler thread poll the sender and call clear(), which drops transfer_infos, so the lookup could come back empty and the room would notify nobody -- a completed transfer then dies on the decode's 300s waiting timeout. Resolve the endpoints first; the send only needs that snapshot. The NIXL worker settled the batch inside the barrier alone, so an exception raised while handles were still being posted told decode its KV pages were free while a sibling was still writing into them. Both waits are now one _await_handles, parameterised by whether the batch is already known broken, and the exception path runs it too. A handle whose state cannot be read counts as running, so it no longer notifies. NixlKVReceiver.failure_exception did not clear. The decode bootstrap failure path retires a request through that method alone, so the room's manager state -- including the addr_to_rooms_tracker entry now dropped in clear() -- leaked on every failed handshake. Mooncake and mori already clear there; the conclude_state latch has to precede it, since check_status defaults a cleared room back to WaitingForInput. prefetch_staging_reqs normalised is_dummy through callable() because NIXL exposed it as a method. It is a bool field on all three backends now.
|
/rerun-group disaggregation |
|
/tag-and-rerun-ci |
|
Results for 🚀 🚀 🚀 🚀 🚀 🚀 🚀 🚀 |
Pick up sgl-project#40645, sgl-project#40711, the shared prefill->decode status plumbing (sgl-project#36612) and runtime role switching (sgl-project#28403). Conflict resolution: - common: keep upstream's update_status structure and still reset the deferred-ACK state when a room lifecycle starts. CommonKVSender.clear() follows sgl-project#40645: ACK now when nothing is in flight, else keep the target. - mooncake: the PR's unconditional register-then-try-ACK already covers sgl-project#40645 and sgl-project#40711. - nixl: keep upstream's settle-then-conclude exception path and still poison the ACK target, since that chunk stays counted. - mori: re-apply the drain-aware ACK path on upstream's conclude_transfer / conclude_failure API and list-returning _submit_kv_transfer; register the drain threads with teardown() and stop them with a sentinel. - tests: update stubs for the new NIXL exception path and for the deferred-ACK fields CommonKVManager now always owns. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

Motivation
NIXL has no prefill→decode failure notification at all: a broken transfer only surfaces on the decode side after
SGLANG_DISAGGREGATION_WAITING_TIMEOUT(300s by default). Mooncake and Mori both notify, but each through its own send path, wire format and terminal-state bookkeeping, so the semantics have drifted apart.Stacked on #36160.
Modifications
1. Shared control plane in
CommonKVManagerPrefill side:
conclude_transfer/conclude_failurerecord the reason (first root cause wins), mark the roomFailed, resolve the destination endpoints and send. ASuccessrequested after a failure was already recorded is downgraded, so decode never seesSuccessfor a transfer that already broke. Both returnNonefor an already-cleared room, so a late conclude cannot re-create it inrequest_statusor leave a failure record behind._room_notify_targetsreturns every non-dummy decode endpoint of the room.Decode side:
apply_prefill_statusaggregatesSuccessacross prefill ranks and records failures. Mooncake's staging last-scatter hook moved in with it, guarded byenable_staging.Wire:
_encode_kv_status_message/parse_kv_status_messageplus two class attributes (kv_status_msg_tag,kv_status_msg_carries_reason) describe each backend's frame layout. All three wires stay byte-for-byte unchanged — Mooncake keeps its 3 untagged frames[room, status, prefill_rank], Mori keeps[MORI_GUARD, room, status, prefill_rank, reason], and NIXL's new message is taggedKV_STATUSand carries the reason.2. NIXL gains the notification
The worker's completion barrier now raises only once every handle of the batch has settled, instead of on the first
ERR. A sibling still inPROCkeeps writing into the decode's KV pages, and the failure notification is exactly what tells the decode those pages are free.The decode listener starts unconditionally (it was only opened on staging or deferred KV release before) and recognises
KV_STATUS.TransferInfo.is_dummybecomes a field computed infrom_zmqinstead of a method, matching Mooncake and Mori so the shared target resolution reads one attribute across all three.3. Mooncake and Mori drop their own implementations
sync_status_to_decode_endpoint,notify_decode_status,_notify_decode_for_room,_conclude_room_failure,_compute_prefill_unique_rank, the per-room notify latch and its lock, bothupdate_statusoverrides, Mori's_cleanup_room_trackingandMoriKVSender.clear()are all replaced by the shared path.Fixes found along the way
update_statusdid not treatFailedas terminal.KVPoll.Failed = 0sorts lowest, so the oldmax(current, status)promoted it back toTransferringorSuccess, and any non-Failedstatus could re-create an already-cleared room. The policy now lives inCommonKVManager— onlyBootstrappingandWaitingForInputmay create a room, andFailedis terminal — which lets the NIXL and Mori overrides, each of which worked around a different half of this, be deleted.bootstrap_roomwould then report as its own root cause.addr_to_rooms_trackeris populated inCommonKVReceiver.__init__but was torn down by three different backend mechanisms. NIXL only did it on theSuccesspath ofpoll(), leaking the entry — andtransfer_statuses— for every failed room. Both now happen inclear(), which the scheduler calls on either terminal outcome.Behavior changes
required_prefill_response_numcounts only non-dummy prefill ranks, so those messages were never awaited.Testing
CPU regression tests are kept on a dedicated test branch:
branch: https://github.com/jambow0320/sglang/tree/rfc-pd-test
path:
test/registered/unit/disaggregation/rfc-test/test_pd_failure_notification.pyThey cover the wire round-trip for all three frame layouts, the tolerant parser, the
Success→Faileddowngrade, the cleared-room guard, fan-out to non-dummy endpoints only, the decode-side aggregation, the sharedupdate_statuspolicy, and the NIXL worker's settle-then-raise barrier in both the quiesced and still-running cases. Mori is covered by stubbing the AMD-only package at import time.Still needs GPU validation: an injected NIXL
ERRproducing end-to-end fast failure, multi-rank fan-out, and Mori on AMD hardware.CI States
Latest PR Test (Base): ✅ Run #34388571724
Latest PR Test (Extra): ❌ Run #34388571285
Latest PR Test (AMD ROCm 10): ❌ Run #34388571486