Repository navigation
[PD] Centralize drain-aware abort acknowledgements - #37077
ShangmingCai merged 16 commits into
Conversation
fed2b63 to
1b5910b
Compare
There was a problem hiding this comment.
Thanks @livingshade! LGTM, verified on MI355*2 node and left minor comments.
| self.bootstrap_room, failure_reason, is_from_another_rank=is_propagated | ||
| ) | ||
|
|
||
| def _send_abort_notification(self): |
There was a problem hiding this comment.
I validated it on MI355X E2E: 45 focused tests passed, the registered Mori E2E passed for TP1→1 and TP2→4, and a cross-node 1P1D abort test remained healthy with a successful follow-up request.
There was a problem hiding this comment.
we can userun-mori-pd label if the AMD CI still has some issue. see
#37077
jambow0320
left a comment
There was a problem hiding this comment.
Two thoughts, both about scope rather than correctness.
1. Can we move the pre-arm into CommonKVReceiver._send_abort_notification
_send_abort_notification has a second caller: _check_waiting_timeout (common/conn.py:1617, reached from poll() in all three backends). It sends the ABORT and sets abort_notified = True without going through Scheduler._abort_request, so register_deferred_abort_room is never called for it.
This PR already fixes this for Mori. The same logic in the base method would fix mooncake and nixl too:
# CommonKVReceiver
def _send_abort_notification(self):
if self.kv_mgr.enable_deferred_decode_kv_release and self.init_time is not None:
self.kv_mgr.register_deferred_abort_room(self.bootstrap_room)
for bootstrap_info in self.bootstrap_infos:
...self.init_time already encodes what metadata_published encodes — it stays None until the end of send_metadata in all three receivers (mooncake:2600, mori:2055, nixl:2968) and _check_waiting_timeout already uses it as that marker — so the new field metadata_published may not be needed. (However, adding a field for readability is also reasonable.)
With putting 'register_deferred_abort_room' into '_send_abort_notification', scheduler.py:4918-4930 (the re-arm right after receiver.abort()) can be deleted. That in turn makes reset_deferred_abort_room and the setdefault override unnecessary.
2. Consider nixl's "no ack" fallback instead of the drain queue
nixl handles the same situation — writes possibly still in flight when completion can no longer be confirmed — by not acking at all:
# nixl/conn.py:1390-1393
self.record_failure(room, str(e))
self.update_status(room, KVPoll.Failed)
# No ack here on purpose: the DONE barrier bails on the first
# ERR, so siblings may still be writing; fall back to the timeout.It needs no extra threads or queue.
The current drain path can stop prefill-side transfers: _enqueue_transfer_drain calls put without block=False, so a transfer worker blocks once the queue is full, while status.Wait() has no timeout parameter and _drain_worker retries indefinitely. _num_shards (default 8) sizes the transfer workers, the drain workers and the queue, so roughly 16 transfers that never reach a terminal status are enough to block every transfer worker — and shards are keyed by room % _num_shards, so healthy rooms stop too.
1353231 to
b0e466f
Compare
52bded8 to
11dc85d
Compare
…lease Take main up to the drain-ACK fan-out change and fan the generation-tagged ACK out to every decode rank that aborted the room: the prefill registry now keeps one AckTarget per decode endpoint instead of one per room, and the drain ACKs all of them. Mori uses the same registry. The transfer_infos snapshot fan-out from sgl-project#41402 is dropped: those peers carry no abort generation, and decode discards generation-less ACKs. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…lease Take main up to the defer-on-every-failure change. _send_abort_notification keeps both force_arm and the armed generation. Fix the auto-merged NIXL worker failure path: sgl-project#41404 uncounts the chunk and ACKs once every handle settled, but the branch still poisoned the room unconditionally afterwards. Decode's ABORT for that failure arrives after the poison, so every settled NIXL failure waited out the release timeout. Poison only when a handle may still be writing. The shared worker-failure poison scenario now provokes an unsettled handle for NIXL, and a settled-failure test covers the ACK with its generation. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A decode that did not arm the ACK tracker, e.g. a prealloc abort before metadata was published, sends ABORT without a generation and does not wait for an ACK. Warning about mixed SGLang versions there is a false alarm in a same-version deployment; an older decode that does wait still falls back to its release timeout. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- Pass the armed generation to note_abort_ack in the host-release test that came in from main. - Initialize _abort_generation on the bare receiver fixture: without it the unarmed ABORT send raised inside the best-effort try and nothing was sent. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
/tag-and-rerun-ci |
sgl-project#41023 kept Mori out of the default-on deferred release because its bootstrap thread marked an aborted room Failed without ever sending an ABORT_ACK. Mori now holds the ack until the transfer worker or drainer quiesces the room, so it can opt in like Mooncake and NIXL. With SGLANG_DISAGGREGATION_DEFERRED_DECODE_KV_RELEASE at its default, this turns deferred release on for Mori; set it to 0 to keep the immediate release. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With deferred release on, a chunk whose room was aborted between submit and the post-submit check waited inline in _wait_transfer_completion, which has no cutoff at the default SGLANG_MORI_TRANSFER_TIMEOUT_MS=0. Hand it to the drain queue instead and return not-quiescent: the drainer still waits for the statuses before decrementing and ACKing, so safety is unchanged, and the shard worker moves on. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A failure whose statuses were still in progress was concluded only after the drainer's status.Wait() returned, so a stuck transfer left the room un-Failed: decode learned of it only from its own waiting timeout, and later chunks for the room kept being submitted. Conclude at the drain handoff instead and keep only the ACK deferred. The drain queue now carries (chunk, statuses). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
|
/rerun-failed-ci |
|
/rerun-failed-ci |
…ncake manager sgl-project#37077 added `_deferred_ack_poisoned_rooms` to CommonKVManager and clears it in update_status() on a new room lifecycle, so a manager built via __new__ must carry it. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
TestMooncakeEarlySend drives the unbound MooncakeKVManager.transfer_worker against a hand-rolled SimpleNamespace. That fake was complete when written, but two changes have since landed on main that made transfer_worker read attributes it does not carry: - sgl-project#42051 added self.state_layout_rejections, read under session_lock for every non-dummy req. - sgl-project#37077 added self.poison_deferred_ack_room(), called from the worker's except path. Merging main into the branch picked up the new conn.py without touching the fake, and the merge was clean because the two sides touch different files. test_worker_waits_before_reading_kv and test_worker_preserves_no_event_path then died with AttributeError on state_layout_rejections, which the except handler turned into a second AttributeError on poison_deferred_ack_room. Only the fixture was stale: the real manager initializes state_layout_rejections on the prefill path in __init__ and inherits poison_deferred_ack_room from CommonKVManager, so no shipping path was affected. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Upstream main changed two base-class contracts that the UniFlow backend relied on: - sgl-project#42348 stopped passing dest_tp_ranks and pp_rank when the prefill scheduler builds a KV sender, so UniflowKVSender now takes the base signature. - sgl-project#37077 centralized drain-aware abort acknowledgements with a per-room generation. UniFlow decode now sends its armed generation in the GUARD-framed ABORT, prefill echoes it in the ABORT_ACK through _send_abort_ack(room, AckTarget), and decode consumes ACKs with handle_abort_ack_message. An ABORT without a generation gets no ACK, matching the common path. The unit tests follow the new state and wire format, and add a case for a generation-less ABORT.
Upstream main changed two base-class contracts that the UniFlow backend relied on: - sgl-project#42348 stopped passing dest_tp_ranks and pp_rank when the prefill scheduler builds a KV sender, so UniflowKVSender now takes the base signature. - sgl-project#37077 centralized drain-aware abort acknowledgements with a per-room generation. UniFlow decode now sends its armed generation in the GUARD-framed ABORT, prefill echoes it in the ABORT_ACK through _send_abort_ack(room, AckTarget), and decode consumes ACKs with handle_abort_ack_message. An ABORT without a generation gets no ACK, matching the common path. The unit tests follow the new state and wire format, and add a case for a generation-less ABORT.

Motivation
When decode aborts or fails a request, prefill may still be RDMA-writing into the decode's KV pages; releasing them early lets a late write corrupt the request that reuses them. Deferred decode KV release (default-on since #41023) holds the pages until every prefill rank ACKs that its writes drained.
This PR adds the drain-aware ACK to Mori and moves the ACK protocol into the common layer shared by Mooncake, NIXL and Mori.
Built on #40645, #41402 and #41404 and resolved potental conflicts.
Related to #33861 and #34510.
Modifications
Mori
supports_deferred_decode_kv_release = True).Common
The common layer is extracted from what Mori established:
AbortNotification/AbortAckwire format. ABORT carries a per-room generation armed before sending, and decode only counts ACKs of the current generation.transfer_infossnapshot, which has no generation.Tests
Commonand backend-specific test.Speed Tests and Profiling
No change on the normal transfer path beyond per-chunk counter bookkeeping.
Accuracy Tests
No change to model forward or sampling. 1P1D H100 (Mooncake),
TestDisaggregationAccuracy: gsm8k 0.83.Unit tests (CPU): test_deferred_decode_kv_release.py (common protocol, tracker, fan-out, decode release), plus backend-specific test_mooncake_deferred_kv_release.py, test_nixl_deferred_kv_release.py and test_mori_deferred_kv_release.py.
1P1D e2e on H100 (TP1 / TP1 over IB), Mooncake and NIXL.
test_disaggregation_chunked_prefill_abort,TestDisaggregationMooncakeFailureMori: AMD CI stage-b-test-large-8-gpu-mi35x-disaggregation-amd (includes test_mori_transfer_engine_e2e.py).
Checklist
CI States
Latest PR Test (Base): ✅ Run #37089529256⚠️ Not enabled -- add
Latest PR Test (Extra):
run-ci-extralabel to opt in.Latest PR Test (AMD ROCm 10): ❌ Run #37089529181