Skip to content
Open
Show file tree
Hide file tree
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 Aug 6, 2026
2156e6d
CollectiveX: self-raise the memlock rlimit before KV pool registration
Oseltamivir Aug 6, 2026
db04dc2
CollectiveX: enable mi355x kv-transfer via MoRI-IO
Oseltamivir Aug 6, 2026
91d25fc
CollectiveX: pin KV legs to the operator's RDMA selectors
Oseltamivir Aug 6, 2026
986c679
CollectiveX: enable h200 kv-transfer via NIXL
Oseltamivir Aug 6, 2026
3b84f54
CollectiveX: add the mooncake backend and the mnnvl fabric lane
Oseltamivir Aug 7, 2026
cb3bf94
CollectiveX: open the launcher identity gates to kv backends and vali…
Oseltamivir Aug 7, 2026
c3f1ef9
CollectiveX: retarget the kv-mla scenario to DeepSeek-V4-Pro's shape
Oseltamivir Aug 7, 2026
0db55e7
CollectiveX: transfer DeepSeek-V4-Pro's cache as vLLM allocates it, a…
Oseltamivir Aug 7, 2026
d7d3eb6
CollectiveX: keep sweep_matrix stdlib-only and give the test job numpy
Oseltamivir Aug 7, 2026
fa156b4
CollectiveX: give the sweep dispatch a suites selector
Oseltamivir Aug 7, 2026
c4b5a37
CollectiveX: retry kv wheel installs with --break-system-packages
Oseltamivir Aug 7, 2026
d2b3f78
CollectiveX: pin mooncake off its NVLink-IPC transport on same-fabric…
Oseltamivir Aug 7, 2026
e84ec10
CollectiveX: stop forcing mooncake onto MNNVL on gb-nv kv legs
Oseltamivir Aug 7, 2026
a915197
CollectiveX: sweep only the DeepSeek-V4-Pro shape
Oseltamivir Aug 7, 2026
0e0cac1
CollectiveX: emit the kv artifact version as a number
Oseltamivir Aug 7, 2026
fec40e3
CollectiveX: densify the kv batch axis and add a sampling trial
Oseltamivir Aug 7, 2026
40d9203
CollectiveX: extend the kv ISL ladder to 8k-512k under a burst descri…
Oseltamivir Aug 8, 2026
7b93d7d
CollectiveX: extend the kv batch ladder to 32
Oseltamivir Aug 8, 2026
4e9390c
CollectiveX: batch ladder to 64 with a kv case guard on gb-nv
Oseltamivir Aug 8, 2026
f2f36c3
CollectiveX: backfill-friendly allocation asks for kv legs
Oseltamivir Aug 8, 2026
11fe18d
CollectiveX: mooncake on mi355x as a push-only row from AMD's atom-de…
Oseltamivir Aug 10, 2026
7e788ce
CollectiveX: give mooncake's per-call guard headroom under contention
Oseltamivir Aug 10, 2026
679ca0a
CollectiveX kv: keep the two smallest batches so the batch axis keeps…
Oseltamivir Aug 12, 2026
dc95b1f
CollectiveX: verify every request in a kv burst and harden the budget…
Oseltamivir Aug 12, 2026
1b229bd
CollectiveX: densify the kv-transfer grid and scale the time budgets …
Oseltamivir Aug 27, 2026
59598bd
CollectiveX sweep: plumb skip_queue_pr through to the shard runs-on l…
Oseltamivir Aug 27, 2026
c6e2d26
CollectiveX: request the nodes:N slot label so the priority controlle…
Oseltamivir Aug 27, 2026
eccaa52
CollectiveX: target runner pools by cluster label and name jobs by ru…
Oseltamivir Aug 27, 2026
79f9497
CollectiveX: name Slurm allocations after the GHA runner
Oseltamivir Aug 27, 2026
b827588
CollectiveX: guarantee a five-rung batch ladder at every kv grid point
Oseltamivir Aug 27, 2026
e8b40e1
CollectiveX: move the kv sweep to a power-of-two batch ladder
Oseltamivir Aug 28, 2026
b3955a5
CollectiveX: size the kv gloo control-plane timeout to the hang guard
Oseltamivir Aug 28, 2026
ca92e9e
CollectiveX: cap the mi355x mooncake kv pool at its registration wall
Oseltamivir Aug 28, 2026
304b401
CollectiveX: raise the shard job ceiling above the kv long-pole alloc…
Oseltamivir Aug 28, 2026
f60552e
CollectiveX: enable gb300 kv-transfer over the XDR rails
Oseltamivir Aug 28, 2026
030a8c1
CollectiveX: size the gb300 kv time budget to its measured pacing
Oseltamivir Aug 29, 2026
b577852
CollectiveX: enable b300 kv-transfer over the RoCE GPU rails
Oseltamivir Aug 30, 2026
d0b3f5e
CollectiveX: exclude the sick b300 nodes from allocation
Oseltamivir Aug 30, 2026
ca90653
CollectiveX: retry b300 allocations that fail the network profile
Oseltamivir Aug 30, 2026
5e99c94
CollectiveX: rebuild kv-transfer on vLLM's packed block-major geometry
Oseltamivir Aug 30, 2026
b0cc0df
CollectiveX: regain the cuda UCX transports under a host-inherited po…
Oseltamivir Aug 30, 2026
f02af33
CollectiveX: drop a host-inherited cuda-less UCX_TLS list instead of …
Oseltamivir Aug 30, 2026
7f4f691
CollectiveX: register the KV pool in pieces below b300's cuda MR wall
Oseltamivir Aug 30, 2026
93feafc
CollectiveX: cap the b300 KV pool under the pods' total-registration …
Oseltamivir Aug 30, 2026
50bc89c
CollectiveX: pin the b300 mooncake NIC to one rail
Oseltamivir Aug 31, 2026
d1a1227
CollectiveX: record kv row provenance and correct the b300 stall mech…
Oseltamivir Aug 31, 2026
2dea42a
CollectiveX: let a kv_device pin narrow UCX below the operator inventory
Oseltamivir Aug 31, 2026
a6a036f
CollectiveX: let the kv_device pin override host-inherited UCX_NET_DE…
Oseltamivir Aug 31, 2026
0c4d307
CollectiveX: keep kv-transfer documents out of the EP bandwidth summary
Oseltamivir Aug 31, 2026
dc68d14
CollectiveX: document the b300 nixl one-rail pin
Oseltamivir Aug 31, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
75 changes: 62 additions & 13 deletions .github/workflows/collectivex-sweep.yml
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,16 @@ on:
type: choice
default: ''
options: ['', dequant]
suites:
description: "Comma-list of suites to resolve (ep-core, kv-transfer); blank = all suites"
type: string
default: ''
skip_queue_pr:
description: "PR number carrying an authorized skip_queue label; shards then request the ci-skip-queue-pr-<N> label so the priority controller admits them ahead of the queue. Blank = normal queueing"
type: string
default: ''
concurrency:
group: cx-${{ github.ref }}-${{ inputs.backend }}-${{ inputs.only_sku }}
group: cx-${{ github.ref }}-${{ inputs.backend }}-${{ inputs.only_sku }}-${{ inputs.suites }}
cancel-in-progress: false

jobs:
Expand Down Expand Up @@ -61,6 +69,7 @@ jobs:
INPUT_EXCLUDE_SKUS: ${{ inputs.exclude_skus }}
INPUT_EP_SIZES: ${{ inputs.ep_sizes }}
INPUT_MODES: ${{ inputs.modes }}
INPUT_SUITES: ${{ inputs.suites }}
RUN_ID: ${{ github.run_id }}
RUN_ATTEMPT: ${{ github.run_attempt }}
run: |
Expand All @@ -70,6 +79,7 @@ jobs:
[ -n "$INPUT_EXCLUDE_SKUS" ] && args+=(--exclude-skus "$INPUT_EXCLUDE_SKUS")
[ -n "$INPUT_EP_SIZES" ] && args+=(--ep-sizes "$INPUT_EP_SIZES")
[ -n "$INPUT_MODES" ] && args+=(--modes "$INPUT_MODES")
[ -n "$INPUT_SUITES" ] && args+=(--suites "$INPUT_SUITES")
python3 sweep_matrix.py "${args[@]}" --out matrix_full.json >/dev/null
python3 - "$GITHUB_OUTPUT" <<'PY'
import hashlib
Expand Down Expand Up @@ -114,22 +124,61 @@ jobs:
runs-on: >-
${{ fromJSON(
vars.PRIORITY_SCHEDULER_ENABLED == 'true' &&
format(
'["self-hosted",{0},{1},{2}]',
toJSON(matrix.sku),
toJSON(format(
'ci-job-{0}-{1}',
needs.setup.outputs.priority,
matrix.queue-token
)),
toJSON(format('ci-attempt-{0}', github.run_attempt))
(
vars.NODE_SLOT_SCHEDULER_ENABLED == 'true' &&
(
inputs.skip_queue_pr != '' &&
format(
'["self-hosted",{0},{1},{2},{3},{4}]',
toJSON(matrix.runner),
toJSON(format('nodes:{0}', matrix.nodes)),
toJSON(format(
'ci-job-{0}-{1}',
needs.setup.outputs.priority,
matrix.queue-token
)),
toJSON(format('ci-attempt-{0}', github.run_attempt)),
toJSON(format('ci-skip-queue-pr-{0}', inputs.skip_queue_pr))
) ||
format(
'["self-hosted",{0},{1},{2},{3}]',
toJSON(matrix.runner),
toJSON(format('nodes:{0}', matrix.nodes)),
toJSON(format(
'ci-job-{0}-{1}',
needs.setup.outputs.priority,
matrix.queue-token
)),
toJSON(format('ci-attempt-{0}', github.run_attempt))
)

Copy link
Copy Markdown

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_pr is nested under NODE_SLOT_SCHEDULER_ENABLED, so the ci-skip-queue-pr-* label is only requested when the node-slot flag is on. With that flag unset or false, a filled skip_queue_pr falls through to the three-label runs-on path and the job queues normally. skip_queue_pr and node-slot matching are independent; the other sweep templates attach the skip-queue label whenever the priority scheduler is on.

Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit c6e2d26. Configure here.

) ||
format(
'["self-hosted",{0},{1},{2}]',
toJSON(matrix.runner),
toJSON(format(
'ci-job-{0}-{1}',
needs.setup.outputs.priority,
matrix.queue-token
)),
toJSON(format('ci-attempt-{0}', github.run_attempt))
)
) ||
format('[{0}]', toJSON(matrix.sku))
format('[{0}]', toJSON(matrix.runner))
) }}
name: p${{ needs.setup.outputs.priority }} | ${{ matrix.sku }} ${{ matrix.backend }} shard ${{ matrix.id }}
timeout-minutes: 350
name: p${{ needs.setup.outputs.priority }} | ${{ matrix.runner }} ${{ matrix.backend }} shard ${{ matrix.id }}
# Must sit above every launcher's largest Slurm allocation plus setup,
# or GitHub cancels a healthy shard before the launcher's own guards can
# act: the kv gb200 mnnvl leg holds a 460 minute allocation with a 420
# minute per-case guard inside it (run 33150394862 was cancelled here at
# 350 minutes while doing honest work), and the kv gb300 legs hold 690
# minutes with a 660 minute guard because gb300 paces ~1.8x gb200 at the
# top isls. EP shards finish far earlier and are unaffected by the
# ceiling.
timeout-minutes: 720
env:
COLLX_BENCH: ${{ matrix.backend }}
COLLX_MODE: ${{ matrix.mode }}
COLLX_IMAGE_OVERRIDE: ${{ matrix.image_ref || '' }}
COLLX_NODES: ${{ matrix.nodes }}
COLLX_GPUS_PER_NODE: ${{ matrix.gpus_per_node }}
COLLX_SCALE_UP_DOMAIN: ${{ matrix.scale_up_domain }}
Expand Down
3 changes: 3 additions & 0 deletions .github/workflows/test-collectivex.yml
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@ jobs:
# ep_flashinfer's combine-model switch compares wheel versions with
# packaging.version; the benchmark image ships it, this runner does not.
pip install packaging
# The kv-transfer workload/geometry tests are numpy math (the cpu
# torch wheel does not pull numpy in).
pip install numpy

# Every skip in this suite is torch-gated, so a missing torch turns the oracle
# checks into silent passes. Fail here instead, where the cause is obvious.
Expand Down
36 changes: 33 additions & 3 deletions experimental/CollectiveX/README.md
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
# CollectiveX

CollectiveX is an experimental MoE expert-parallel communication benchmark. It measures dispatch,
combine, and paired roundtrip latency across EP libraries and accelerator systems, then uploads
neutral result artifacts.
CollectiveX is an experimental inference-communication benchmark. Its EP suite measures MoE
dispatch, combine, and paired roundtrip latency across EP libraries and accelerator systems; its
KV-transfer suite measures disaggregated-serving KV-cache handoffs across transfer libraries and
fabrics. Both upload neutral result artifacts.

CollectiveX schedules benchmarks, executes them on real allocations, and uploads the neutral
artifacts each run emits. It does not validate those artifacts, promote, rank, recommend, select, or
Expand Down Expand Up @@ -149,6 +150,35 @@ scale-up ranks per domain; GB EP16 remains MNNVL scale-up and therefore uses LSA
SKU/backend/EP cell is attempted is a capability fact; whether it succeeded is decided by the
benchmark's return code.

## KV-Cache Transfer Suite

`kv-transfer` legs run 2 nodes x 1 GPU — the per-worker prefill/decode pair a
disaggregated deployment actually forms — and move bursts of 1 to 32 concurrent
requests' paged KV as per-request layer-major descriptor lists over seed-keyed
random block tables (the post-fragmentation layout vLLM and SGLang post; a
burst posts every request's prepped transfer, then awaits them all), plus one
contiguous bulk row as the wire-speed ceiling. The workload is transcribed from
what vLLM allocates for the model it serves: `kv-dsv4` = DeepSeek-V4-Pro's mixed
cache (30 Compressed Sparse Attention layers at 4 tokens per 576 B entry plus
their 132 B indexer entries, 31 Heavily Compressed Attention layers at 128
tokens per entry, and the 128-token sliding-window cache on all 61 layers; fp8
by architecture); ISL
8k to 512k at page sizes 16 and 64 tokens; `pull` (READ, vLLM NixlConnector) and `push` (WRITE,
SGLang disagg) both timed from the initiator with offset-pattern verification on
the destination pool in both directions. Backends: `nixl` (what Dynamo, vLLM,
and SGLang ship), `mooncake` (NVIDIA-only; the wheel links libcuda at import),
and `mori-io` (AMD's native engine) where the registry's `kv_backends` map
enables them; no entry, no legs, mirroring `ll_backends`. A backend entry may
restrict ops, pin an image, or set a NIC filter (mooncake on mi355x is
push-only from AMD's atom-dev image over the GPU-paired Pollara NIC; upstream
ionic RDMA READ is broken; both b300 backends are pinned to one rail —
mooncake because the image's engine draws peer NICs blindly and cross-rail
draws stall ~1 s, nixl because UCX's own two-rail READ selection is unstable
across runs while one rail holds 49 GB/s at p95/p50 ≤ 1.01 — so b300 kv rows
are one-rail measurements). Fabrics: `rdma`
(torch pools) and, on GB racks, `mnnvl` (cuMem FABRIC pools; see the
methodology for the bulk-vs-paged lane inversion that row exists to publish).

## Workflow And Artifacts

`.github/workflows/collectivex-sweep.yml` has two jobs. `setup` generates a public-SKU matrix
Expand Down
4 changes: 4 additions & 0 deletions experimental/CollectiveX/bandwidth.py
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,10 @@ def render(documents: list[dict]) -> str:
"marks an extrapolated alpha, and rungs failing the correctness gate are excluded.",
"",
]
# kv-transfer documents have their own row model (per-transfer, no
# tokens_per_rank or routing); this renderer reads only EP-suite rows.
documents = [d for d in documents
if d["identity"]["case_factors"]["case"].get("suite") != "kv-transfer"]
for document in sorted(documents, key=_sort_key):
case = document["identity"]["case_factors"]["case"]
ep = _ep(document)
Expand Down
97 changes: 97 additions & 0 deletions experimental/CollectiveX/bench/kv_backend.py
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
Loading