Skip to content

[PD][mori] Align prefill transfer control plane for unified control plane - #36160

Merged
ShangmingCai merged 5 commits into
sgl-project:mainfrom
inkcherry:pd/mori-worker-control-plane-align
Aug 27, 2026
Merged

ShangmingCai merged 5 commits into
sgl-project:mainfrom
inkcherry:pd/mori-worker-control-plane-align

Conversation

@inkcherry

@inkcherry inkcherry commented Aug 24, 2026 •

Copy link
Copy Markdown
Contributor

changes for requests of #34510

PR changes :

  • Move Mori prefill transfer control from Sender to transfer_worker / Manager
  • Queue work by room only — no Sender reference in queued tasks
  • Notify / failure handled centrally on the Manager side
  • Sender is enqueue + poll only
  • Transport (RDMA) unchanged

Motivation

Modifications

Accuracy Tests

Speed Tests and Profiling

Checklist

Review and Merge Process

  1. Ping Merge Oncalls to start the process. See the PR Merge Process.
  2. Get approvals from CODEOWNERS and other reviewers.
  3. Trigger CI tests with comments or contact authorized users to do so.
    • Common commands include /tag-and-rerun-ci, /tag-run-ci-label, /rerun-failed-ci
  4. After green CI and required approvals, ask Merge Oncalls or people with Write permission to merge the PR.

CI States

Latest PR Test (Base): ✅ Run #32890274134
Latest PR Test (Extra): ✅ Run #32970353352
Latest PR Test (AMD ROCm 7.2): ❌ Run #32890274261

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>
@jambow0320

Copy link
Copy Markdown
Contributor

A few comments:

1. MoriFailureExceptionMixin

Could this be dropped? Right now only MoriKVSender inherits it, while MoriKVReceiver still carries its own identical copy of failure_exception. Mooncake's mixin exists because it is genuinely shared between its sender and receiver, but that is not the case here.

Also, in NIXL the sender and receiver implementations of failure_exception differ (the sender additionally drains kv_mgr.exceptions, sets _send_failed, and prefers _send_error, while the receiver does neither clear() nor a conclude_state latch). So when we eventually extract the common logic, I would keep failure_exception on the sender and the receiver separately rather than introducing a shared mixin.

2. _room_status_notified is never cleaned up

The dict is only ever initialized, read, and set to True — there is no removal anywhere, and the inherited CommonKVSender.clear() does not know about it. Two consequences:

  • it grows unbounded over the lifetime of the process;
  • if a later request reuses the same bootstrap_room, the leftover True suppresses that request's terminal notification entirely, and decode only fails after SGLANG_DISAGGREGATION_WAITING_TIMEOUT (300s by default).

Previously this was self.status_notified on the sender, so it went away with the object and needed no cleanup.

Could we clear it in clear()? I would lean towards overriding clear() in MoriKVSender for now, since the other backends do not have _room_status_notified. Both terminal paths reach clear() (directly on success, and via failure_exception() on failure), so that should cover it.

3. Duplicate decode notification on the SLA path

_wait_transfer_completion concludes the failure and then returns IN_PROGRESS:

sla_tripped = True
self._conclude_room_failure(room, f"KV transfer exceeded SLA {sla_ms}ms")
return StatusCode.IN_PROGRESS

and the caller treats any non-SUCCESS as a failure and concludes again:

rc = self._wait_transfer_completion(statuses, room)
if rc != StatusCode.SUCCESS:
    self._conclude_room_failure(room, self._collect_transfer_failure_reason(statuses))
    return

So _conclude_room_failure is invoked twice for the same chunk, and the second one would emit a second Failed to decode (with the generic "unknown reason" text, since no status has actually failed at SLA time). As far as I can tell, _room_status_notified is what currently absorbs this. It may be cleaner to not conclude inside _wait_transfer_completion and let the caller conclude once, or to return a distinct sentinel for the SLA case.

@jambow0320 jambow0320 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Left a couple of minor comments inline.

Comment on lines +1637 to +1638
def abort(self):
self._finalize_failure("Aborted by AbortReq.")
super().abort()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we drop this override? def abort(self): super().abort() is equivalent to not defining it at all.

One thing worth flagging: it also drops the decode notification. Previously abort() went through _finalize_failure("Aborted by AbortReq.") and pushed Failed to decode; CommonKVSender.abort() only records locally, so a prefill-side abort now leaves decode waiting out SGLANG_DISAGGREGATION_WAITING_TIMEOUT (300s). mori was the only backend that had this.

Fine to land as-is — I'm adding a unified peer notification to CommonKVSender.abort() which will restore it for all three backends.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sure, removed.

Comment on lines +1617 to +1620
def clear(self) -> None:
with self.kv_mgr._room_notify_lock:
super().clear()
self.kv_mgr._room_status_notified.pop(self.bootstrap_room, None)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we narrow the lock to just the pop?

def clear(self) -> None:
    super().clear()
    with self.kv_mgr._room_notify_lock:
        self.kv_mgr._room_status_notified.pop(self.bootstrap_room, None)

CommonKVSender.clear() pops kv_mgr.transfer_infos, which is guarded by transfer_lock everywhere else in this file, and _notify_decode_for_room already takes _room_notify_lock → transfer_lock. No deadlock today, but the ordering can be break later.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

modified

@ShangmingCai

ShangmingCai commented Aug 25, 2026 •

Copy link
Copy Markdown
Collaborator

/tag-and-rerun-ci

@github-actions github-actions Bot added the run-ci CI: run the baseline test suite on this PR label Aug 25, 2026
@ShangmingCai

Copy link
Copy Markdown
Collaborator

/rerun-failed-ci

@ShangmingCai
ShangmingCai merged commit 8ad7641 into sgl-project:main Aug 27, 2026
406 of 486 checks passed
jambow0320 added a commit to jambow0320/sglang that referenced this pull request Aug 30, 2026
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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bypass-fastfail run-ci CI: run the baseline test suite on this PR

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants