-
Notifications
You must be signed in to change notification settings - Fork 280
CollectiveX: kv-transfer suite — NIXL + MoRI-IO KV-cache handoff benchmark (stacked on #2489) #2510
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
Oseltamivir
wants to merge
51
commits into
main
Choose a base branch
from
collectivex-kv-transfer
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
51 commits
Select commit
Hold shift + click to select a range
6959aa3
CollectiveX: add the kv-transfer suite (NIXL + MoRI-IO KV-cache hando…
Oseltamivir 2156e6d
CollectiveX: self-raise the memlock rlimit before KV pool registration
Oseltamivir db04dc2
CollectiveX: enable mi355x kv-transfer via MoRI-IO
Oseltamivir 91d25fc
CollectiveX: pin KV legs to the operator's RDMA selectors
Oseltamivir 986c679
CollectiveX: enable h200 kv-transfer via NIXL
Oseltamivir 3b84f54
CollectiveX: add the mooncake backend and the mnnvl fabric lane
Oseltamivir cb3bf94
CollectiveX: open the launcher identity gates to kv backends and vali…
Oseltamivir c3f1ef9
CollectiveX: retarget the kv-mla scenario to DeepSeek-V4-Pro's shape
Oseltamivir 0db55e7
CollectiveX: transfer DeepSeek-V4-Pro's cache as vLLM allocates it, a…
Oseltamivir d7d3eb6
CollectiveX: keep sweep_matrix stdlib-only and give the test job numpy
Oseltamivir fa156b4
CollectiveX: give the sweep dispatch a suites selector
Oseltamivir c4b5a37
CollectiveX: retry kv wheel installs with --break-system-packages
Oseltamivir d2b3f78
CollectiveX: pin mooncake off its NVLink-IPC transport on same-fabric…
Oseltamivir e84ec10
CollectiveX: stop forcing mooncake onto MNNVL on gb-nv kv legs
Oseltamivir a915197
CollectiveX: sweep only the DeepSeek-V4-Pro shape
Oseltamivir 0e0cac1
CollectiveX: emit the kv artifact version as a number
Oseltamivir fec40e3
CollectiveX: densify the kv batch axis and add a sampling trial
Oseltamivir 40d9203
CollectiveX: extend the kv ISL ladder to 8k-512k under a burst descri…
Oseltamivir 7b93d7d
CollectiveX: extend the kv batch ladder to 32
Oseltamivir 4e9390c
CollectiveX: batch ladder to 64 with a kv case guard on gb-nv
Oseltamivir f2f36c3
CollectiveX: backfill-friendly allocation asks for kv legs
Oseltamivir 11fe18d
CollectiveX: mooncake on mi355x as a push-only row from AMD's atom-de…
Oseltamivir 7e788ce
CollectiveX: give mooncake's per-call guard headroom under contention
Oseltamivir 679ca0a
CollectiveX kv: keep the two smallest batches so the batch axis keeps…
Oseltamivir dc95b1f
CollectiveX: verify every request in a kv burst and harden the budget…
Oseltamivir 1b229bd
CollectiveX: densify the kv-transfer grid and scale the time budgets …
Oseltamivir 59598bd
CollectiveX sweep: plumb skip_queue_pr through to the shard runs-on l…
Oseltamivir c6e2d26
CollectiveX: request the nodes:N slot label so the priority controlle…
Oseltamivir eccaa52
CollectiveX: target runner pools by cluster label and name jobs by ru…
Oseltamivir 79f9497
CollectiveX: name Slurm allocations after the GHA runner
Oseltamivir b827588
CollectiveX: guarantee a five-rung batch ladder at every kv grid point
Oseltamivir e8b40e1
CollectiveX: move the kv sweep to a power-of-two batch ladder
Oseltamivir b3955a5
CollectiveX: size the kv gloo control-plane timeout to the hang guard
Oseltamivir ca92e9e
CollectiveX: cap the mi355x mooncake kv pool at its registration wall
Oseltamivir 304b401
CollectiveX: raise the shard job ceiling above the kv long-pole alloc…
Oseltamivir f60552e
CollectiveX: enable gb300 kv-transfer over the XDR rails
Oseltamivir 030a8c1
CollectiveX: size the gb300 kv time budget to its measured pacing
Oseltamivir b577852
CollectiveX: enable b300 kv-transfer over the RoCE GPU rails
Oseltamivir d0b3f5e
CollectiveX: exclude the sick b300 nodes from allocation
Oseltamivir ca90653
CollectiveX: retry b300 allocations that fail the network profile
Oseltamivir 5e99c94
CollectiveX: rebuild kv-transfer on vLLM's packed block-major geometry
Oseltamivir b0cc0df
CollectiveX: regain the cuda UCX transports under a host-inherited po…
Oseltamivir f02af33
CollectiveX: drop a host-inherited cuda-less UCX_TLS list instead of …
Oseltamivir 7f4f691
CollectiveX: register the KV pool in pieces below b300's cuda MR wall
Oseltamivir 93feafc
CollectiveX: cap the b300 KV pool under the pods' total-registration …
Oseltamivir 50bc89c
CollectiveX: pin the b300 mooncake NIC to one rail
Oseltamivir d1a1227
CollectiveX: record kv row provenance and correct the b300 stall mech…
Oseltamivir 2dea42a
CollectiveX: let a kv_device pin narrow UCX below the operator inventory
Oseltamivir a6a036f
CollectiveX: let the kv_device pin override host-inherited UCX_NET_DE…
Oseltamivir 0c4d307
CollectiveX: keep kv-transfer documents out of the EP bandwidth summary
Oseltamivir dc68d14
CollectiveX: document the b300 nixl one-rail pin
Oseltamivir File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,97 @@ | ||
| #!/usr/bin/env python3 | ||
| """Backend contract for the KV-cache transfer suite. | ||
|
|
||
| The harness owns the data (kv_pool pools, pattern fill, verification) and the | ||
| protocol (rank 0 = target, rank 1 = initiator, lockstep barriers); an adapter | ||
| owns registration, connection, and posting. Transfers are one-sided from the | ||
| initiator, so completion is host-visible and timing is wall clock around | ||
| post-to-complete; no CUDA events, because no local kernel participates. | ||
| `pull` (READ) is the vLLM NixlConnector shape, `push` (WRITE) the SGLang disagg | ||
| shape; the measured quantity is the same completion either way. Connection | ||
| payloads ride the harness object exchange, never adapter side channels. | ||
| """ | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| import time | ||
|
|
||
|
|
||
| class KVBackend: | ||
| """One transfer library on one rank. Subclasses implement the five hooks.""" | ||
|
|
||
| name = "abstract" | ||
| #: maturity mirrors EPBackend.maturity ("production" | "candidate"). | ||
| maturity = "candidate" | ||
| library_version: str | None = None | ||
|
|
||
| def __init__(self, args, role: str, device) -> None: | ||
| self.args = args | ||
| self.role = role | ||
| self.device = device | ||
|
|
||
| # -- lifecycle ------------------------------------------------------------ | ||
| def register(self, pool, bulk, reg_layout=None) -> None: | ||
| """Register the pool + bulk tensors with the library. | ||
|
|
||
| ``reg_layout`` is the pool's shared region layout — (base, | ||
| packed_bytes, nbytes) triples, contiguous from zero and valid for | ||
| every planned config (run_kv._harmonize). Adapters may use it to | ||
| split one oversized registration into pieces cut on the descriptor | ||
| grid, so no descriptor straddles two pieces; ignoring it is valid. | ||
| """ | ||
| raise NotImplementedError | ||
|
|
||
| def publish(self) -> dict: | ||
| """Payload the peer needs to reach this rank (addresses, packed descs).""" | ||
| raise NotImplementedError | ||
|
|
||
| def connect(self, peer: dict) -> None: | ||
| """Consume the peer's payload; after this, transfers may be prepared.""" | ||
| raise NotImplementedError | ||
|
|
||
| def teardown(self) -> None: # pragma: no cover - adapter-specific | ||
| pass | ||
|
|
||
| # -- transfers (initiator only) -------------------------------------------- | ||
| def make_paged(self, cfg: dict, op: str, local_tables, remote_tables): | ||
| """Return (post, wait, prep_seconds) for one request's paged KV. | ||
|
|
||
| ``post()`` submits the whole descriptor list asynchronously; ``wait()`` | ||
| blocks until it completes — split so a batch of requests overlaps like | ||
| a decode step admitting several requests at once. Preparation cost | ||
| (descriptor build + handle creation) is amortized by engines through | ||
| prepped-handle reuse, so it is reported separately, never inside the | ||
| timed transfer. | ||
| """ | ||
| raise NotImplementedError | ||
|
|
||
| def make_bulk(self, nbytes: int, op: str): | ||
| """Return (post, wait, prep_seconds) for one contiguous transfer of | ||
| ``nbytes`` — the single-descriptor contiguous baseline row (logical | ||
| payload over host-observed completion; not a proven physical wire | ||
| rate — backends may split large operations internally).""" | ||
| raise NotImplementedError | ||
|
|
||
|
|
||
| def time_bursts(transfers, warmup: int, reps: int) -> tuple[list[float], list[float]]: | ||
| """(burst_ms, request_ms), warmups dropped. ``transfers`` is a list of | ||
| (post, wait) pairs — one per request. A burst posts every request, then | ||
| drains the waits in posting order; burst_ms is post-of-first to | ||
| completion-of-last, and request_ms records each individual request's | ||
| host-observed completion offset from the burst start. Because the waits | ||
| drain in posting order, a request's mark upper-bounds its true completion | ||
| (a later request that finished early is observed at its wait's turn).""" | ||
| burst_ms: list[float] = [] | ||
| request_ms: list[float] = [] | ||
| for rep in range(warmup + reps): | ||
| start = time.perf_counter() | ||
| for post, _ in transfers: | ||
| post() | ||
| marks = [] | ||
| for _, wait in transfers: | ||
| wait() | ||
| marks.append((time.perf_counter() - start) * 1e3) | ||
| if rep >= warmup: | ||
| burst_ms.append(marks[-1]) | ||
| request_ms.extend(marks) | ||
| return burst_ms, request_ms |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Skip-queue ignored without node slots
Medium Severity
skip_queue_pris nested underNODE_SLOT_SCHEDULER_ENABLED, so theci-skip-queue-pr-*label is only requested when the node-slot flag is on. With that flag unset or false, a filledskip_queue_prfalls through to the three-labelruns-onpath and the job queues normally.skip_queue_prand node-slot matching are independent; the other sweep templates attach the skip-queue label whenever the priority scheduler is on.Reviewed by Cursor Bugbot for commit c6e2d26. Configure here.