Skip to content

ep(bench): add mscclpp high-throughput (rank-major) backend to the unified Python EP benchmark - #860

Merged
Binyang Li (Binyang2014) merged 205 commits into
feature/epfrom
qinghuazhou/unified_ep_bench_ht_python
Aug 22, 2026
Merged

ep(bench): add mscclpp high-throughput (rank-major) backend to the unified Python EP benchmark#860
Binyang Li (Binyang2014) merged 205 commits into
feature/epfrom
qinghuazhou/unified_ep_bench_ht_python

Conversation

@seagater

Copy link
Copy Markdown
Contributor

Summary

Adds the mscclpp high-throughput (HT) EP backend to the unified in-process
Python benchmark, on top of the review-comment follow-ups from #858. Benchmarks
all five backends through one shared harness with identical dispatch→combine
timing: mscclpp LL, mscclpp HT (TOKEN_MAJOR / RANK_MAJOR), NCCL-EP,
DeepEP V2, FlashInfer.

Changes

  • New mscclpp-ht backend (--backend mscclpp-ht): drives
    MoECommunicator in HIGH_THROUGHPUT mode. Supports --ep-layout token_major (default) and rank_major; --validate checks rank-major
    reduces bit-exactly to the token-major reference.
  • Lives in ep_bench_mscclpp.py alongside the LL setup_mscclpp (both drive
    the same MoECommunicator API), sharing parse_kineto_kernels.
  • CUDA-graph capture enabled for HT: the harness captures the cached
    dispatch (previous_handle= → no host-side notify wait) + combine as one
    graph, reusing the harness-owned _capture_paired_graph and
    --graph-group-size.
  • Runtime resolved to the official rank-major HT support (Support rank major #857) + multi-node
    NVL/NVLS (Add multi-node NVL/NVLS algorithm support #855) from feature/ep; supersedes the branch's local rank-major prototype.

Validation (GB200 NVL72)

  • HT cuda-graph captured at 1 and 2 nodes, both TOKEN_MAJOR and RANK_MAJOR,
    d8704 e512 k8 t4096 — all single-graph, no eager fallback. Kernel Total (D+C):
    1n ~934 µs, 2n ~1406 µs; rank-major ≈ token-major (no measurable overhead).
  • The four LL/other backends still run eager and cuda-graph correctly through
    the merged harness (numbers unchanged from ep(bench): unified in-process Python EP benchmark — review-comment follow-ups #858).

Changho Hwang (chhwang) and others added 30 commits April 20, 2026 18:32
Port DeepEP's high-throughput MoE dispatch/combine kernels onto MSCCL++
as an optional build target `mscclpp_ep_cpp`, gated by -DMSCCLPP_BUILD_EXT_EP
(OFF by default). Sources are lifted from DeepEP branch
`chhwang/dev-atomic-add-cleanup` and rebased onto upstream MSCCL++ APIs;
the NVSHMEM / IBGDA dependencies are replaced with `PortChannel` +
`MemoryChannel` + the new `Connection::atomicAdd` primitive.

Scope
-----
Intranode (NVLink-only):
  * `Buffer` ctor/dtor: cudaMalloc nvl workspace, export IPC handle,
    allocate FIFO + peer-pointer tables, start `ProxyService`.
  * `sync()`: import peer IPC handles, upload peer pointer table,
    build `MemoryDevice2DeviceSemaphore` + `MemoryChannel` per peer.
  * `get_dispatch_layout`, `intranode_dispatch`, `intranode_combine`
    ported verbatim (torch::Tensor ABI preserved).

Internode HT (NVLink + RDMA):
  * `sync()` RDMA branch: cudaMalloc RDMA buffer + `bootstrap->barrier()`
    (replacing NVSHMEM symmetric-heap allocation); register with
    `all_transport`, exchange via `sendMemory`/`recvMemory`, build 12 IB
    QPs/peer + 16 semaphores/peer + 16 port channels/peer.
  * Full `internode.cu` port (notify_dispatch / dispatch / cached_notify
    / combine / get_dispatch_layout). The 4 raw `ChannelTrigger` atomic
    sites are rewritten to call the new
    `PortChannelDeviceHandle::atomicAdd(offset, value)` API; the single
    `nvshmem_fence()` is replaced with `__threadfence_system()` (remote
    visibility guaranteed by the subsequent port-channel barrier).
  * `internode_dispatch` / `internode_combine` host code ported, with
    the torch tensor marshalling and CPU spin-wait on mapped counters.

Low-latency (pure RDMA):
  * Not ported. `low_latency_dispatch`, `low_latency_combine`,
    `clean_low_latency_buffer`, `get_next_low_latency_combine_buffer`
    throw `std::runtime_error`; the Python frontend refuses to
    construct a Buffer with `low_latency_mode=True`.

Python layer
------------
* New pybind11 + libtorch Python extension `mscclpp_ep_cpp` (separate
  from the nanobind `_mscclpp` because the EP ABI carries
  `torch::Tensor` / `at::cuda::CUDAStream`).
* `mscclpp.ext.ep.Buffer` mirrors `deep_ep.Buffer`; exchanges device
  IDs, IPC handles and the bootstrap UniqueId over the user's
  `torch.distributed` process group before calling `sync()`.
* `mscclpp.ext` auto-imports `ep` if the extension is built.

Build
-----
* `src/ext/ep/CMakeLists.txt`: finds Python + Torch; warns and skips if
  `CMAKE_PREFIX_PATH` doesn't point at `torch.utils.cmake_prefix_path`.
  Falls back to Torch's bundled pybind11 if a standalone pybind11 is not
  installed. Links `libtorch_python` explicitly (without it, `import
  mscclpp_ep_cpp` fails with `undefined symbol: THPDtypeType`).
* Top-level `CMakeLists.txt` exposes the `MSCCLPP_BUILD_EXT_EP` option
  (default OFF).

Tests
-----
* `test/python/ext/ep/test_ep_smoke.py`: skipped if the extension isn't
  built. Covers Config round-trip, low-latency size hint, and the LL
  construction guard. Multi-rank functional tests still to do on H100.

Notes
-----
* Builds against the preceding "atomic add" commit which adds
  `Connection::atomicAdd` and `PortChannelDeviceHandle::atomicAdd` to
  upstream MSCCL++.
* Intranode path verified end-to-end (build + import + smoke tests).
* Internode HT is code-complete but requires real IB hardware to
  validate; see `src/ext/ep/README.md` for the detailed port plan and
  remaining LL migration.
Port DeepEP's pure-RDMA low-latency (LL) MoE kernels from
csrc/kernels/internode_ll.cu (branch chhwang/dev-atomic-add-cleanup)
into the MSCCL++ EP extension. NVSHMEM / IBGDA device primitives are
replaced with MSCCL++ PortChannelDeviceHandle operations:

  nvshmemx_barrier_all_block()            -> port-channel signal+wait ring
  nvshmemi_ibgda_put_nbi_warp(...)        -> lane-0 PortChannel.put(...)
  nvshmemi_ibgda_amo_nonfetch_add(...)    -> lane-0 PortChannel.atomicAdd(...)

The atomicAdd path relies on the MSCCL++ Connection::atomicAdd /
PortChannelDeviceHandle::atomicAdd API cherry-picked from branch
chhwang/new-atomic-add; the LL dispatch path uses a signed delta
(-num_tokens_sent - 1) which the new int64_t signature supports.

Changes:
* New file src/ext/ep/kernels/internode_ll.cu (~530 lines) with the
  three kernels clean_low_latency_buffer, dispatch<kUseFP8,...>,
  combine<...> plus their launchers. rdma_buffer_ptr is threaded
  through the launchers so the kernel can translate virtual addresses
  into registered-memory offsets expected by MSCCL++.
* kernels/api.cuh: replace the single stub signature with full LL
  launcher prototypes.
* buffer.cc: replace the four LL throw-stubs
  (clean_low_latency_buffer, low_latency_dispatch,
  low_latency_combine, get_next_low_latency_combine_buffer) with
  torch-Tensor implementations ported from DeepEP/csrc/deep_ep.cpp.
* Drop src/ext/ep/internode_stub.cc and its CMake entry.
* python/mscclpp/ext/ep/buffer.py: remove the low_latency_mode=True
  NotImplementedError guard; update docstring.
* test/python/ext/ep/test_ep_smoke.py: rename
  test_low_latency_rejected -> test_low_latency_buffer_construct
  to reflect that LL construction is now accepted.
* src/ext/ep/README.md: update status matrix, document the
  NVSHMEM -> MSCCL++ translation table, and list the known
  limitations.

This is a structural port: the kernels compile, link, and pass the
single-rank smoke tests, but end-to-end behaviour on multi-node H100
is not yet validated. Two known caveats:

  1. Performance will NOT match IBGDA because MSCCL++ port channels
     use a CPU proxy; this port is for functional parity, not latency.
  2. Buffer::sync() in LL mode only connects peers that share the
     same local GPU id (DeepEP convention), so the LL kernels assume
     a one-GPU-per-node topology (num_ranks == num_rdma_ranks).
     Multi-GPU-per-node LL layouts will need a follow-up in sync().

Tested:
  cmake --build build -j --target mscclpp_ep_cpp   # builds clean
  pytest test/python/ext/ep/test_ep_smoke.py        # 3 passed
Three issues blocked end-to-end intranode validation across multiple
ranks. This commit fixes them and adds a 2/4/8-rank functional test.

1. Combine receiver: OOB __shared__ read

   In the combine receiver warp, the wait loop evaluated
   `channel_tail_idx[recv_lane_id] <= expected_head` before the
   `expected_head >= 0` guard. `channel_tail_idx` is a shared array
   of size `kNumRanks`, but the loop runs on all 32 lanes of a warp,
   so lanes with `recv_lane_id >= kNumRanks` indexed out of bounds.
   compute-sanitizer reported "Invalid __shared__ read of size 4
   bytes" at combine<bf16,2,768>+0xdd0, surfaced asynchronously as
   cudaErrorIllegalAddress at the kernel launch site. Swap the
   operands so the rank-bounds check short-circuits the shared read.

2. Python bindings: UniqueId ABI

   `mscclpp::UniqueId` is a `std::array<uint8_t, N>` which pybind11
   auto-converts to a Python `list`, silently overriding any
   `py::class_<UniqueId>` wrapper. Expose `create_unique_id` /
   `connect` as lambdas that produce/consume `py::bytes` and memcpy
   into a local `UniqueId`. Also coerce `bytes`->`bytearray` at the
   Python call site for `sync()` whose signature expects
   `pybind11::bytearray`.

3. Python frontend: communicator required for NVL-only sync

   `Buffer::sync()` uses `communicator->connect(ipc_config, ...)` on
   the pure-NVLink path, so the communicator must be initialized
   even when `num_rdma_ranks == 1` and `low_latency_mode == False`.
   Always broadcast the unique id and call `runtime.connect()`
   before `sync()`.

Validation on a single H100x8 node via torchrun:
- 2 ranks: dispatch 195 tokens, combine diff=0
- 4 ranks: dispatch 371 tokens, combine diff=0
- 8 ranks: dispatch 456 tokens, combine diff=0

Test harness added at test/python/ext/ep/test_intranode_multirank.py.
The `internode` kernels index device-side port channel handles as
`port_channel_handles[channel_id * num_ranks + peer_rank]`, where
`peer_rank` is a global rank in [0, num_ranks). `Buffer::sync` was
building that table by iterating `std::unordered_map<int, MemoryId>`
(and similarly for connections/semaphores), which yields hash order
rather than ascending rank order. Once the cross-node fan-out grew
beyond a single peer, a local rank's trigger for peer `r` landed on
the semaphore/memory pair of a different peer, so RDMA puts and
atomic tail updates went to the wrong destination and the forwarder
spun on a tail counter that never advanced.

Changes:
  - Build `sema_ids` and `port_channel_handles` by iterating
    `for (int r = 0; r < num_ranks; ++r)` and looking up the
    connection / memory id for rank `r`, skipping ranks excluded by
    low-latency mode (inserting a placeholder handle so the stride
    stays `num_ranks`).
  - Tag the RDMA-phase `sendMemory`/`recvMemory`/`connect` calls with
    `kRdmaTag = 1` so they do not collide with NVL-phase tag-0
    traffic between the same pair of ranks.
  - Drop an unused `r` local in the NVL setup loop.

With this fix and a matched `libmscclpp.so` on both nodes, the
2-node x 8-GPU internode HT dispatch path completes successfully
(`[dispatch] OK`). Combine is still under investigation.

Also adds `test/python/ext/ep/test_internode_multirank.py`, a
torchrun-based 2-node functional test that exercises
`get_dispatch_layout` -> `internode_dispatch` -> `internode_combine`
and validates per-source-rank token values end-to-end.
Two issues prevented internode HT combine from completing on 2x8 H100:

1. Wrong prefix matrices passed to internode_combine. Combine runs in the
   reverse direction of dispatch, so it must consume the receiver-side
   matrices returned by dispatch (recv_rdma_channel_prefix_matrix,
   recv_rdma_rank_prefix_sum, recv_gbl_channel_prefix_matrix), not the
   sender-side rdma_channel_prefix_matrix / gbl_channel_prefix_matrix.
   This matches DeepEP's deep_ep/buffer.py::internode_combine handle
   unpacking. Without the fix the NVL forwarder's 'NVL check' timed out
   because token_start_idx/token_end_idx were computed against the wrong
   per-channel layout.

2. Cross-rank race between dispatch and combine. Even with the correct
   matrices, launching combine immediately after dispatch deadlocked the
   forwarder NVL check (tail stuck one short of expected_head) because
   peers still had in-flight dispatch proxy traffic while fast ranks had
   already started combine. A torch.cuda.synchronize() + dist.barrier()
   between the two calls makes the test pass deterministically on 16
   ranks (combine diff == 0, max|expected| up to 60.0).

The barrier in the test is a workaround; the real fix belongs in
Buffer::internode_dispatch / Buffer::internode_combine so the
dispatch->combine handoff fully fences outstanding proxy work across
ranks. Marked with an XXX comment in the test.
Refresh status docs and comments now that internode HT dispatch and
combine have been validated end-to-end on 2 nodes x 8 H100 GPUs via
test/python/ext/ep/test_internode_multirank.py (all 16 ranks recover
their per-rank token payloads with zero diff).

- src/ext/ep/README.md: consolidate the previously duplicated README
  into a single document; mark intranode and internode HT dispatch and
  combine as validated in the status table; add a 'Running the tests'
  section with torchrun examples for both the intranode and the 2x8
  internode setups; record the dispatch->combine
  torch.cuda.synchronize() + dist.barrier() requirement under Known
  limitations; mark Phase 2 DONE and keep Phase 3 (LL) as structural
  port, untested.

- python/mscclpp/ext/ep/buffer.py: update the module docstring and the
  Buffer constructor docstring to say internode HT is validated and
  clarify that LL mode is untested on multi-node hardware.

- src/ext/ep/buffer.cc: drop the stale 'NVSHMEM support not yet ported'
  and 'low-latency paths still stubbed' comments. mscclpp_ep does not
  use NVSHMEM at all (PortChannel/MemoryChannel replace it), and the LL
  paths are a structural port that is present but untested, not stubbed.
  Note validation on 2x H100x8 in the internode section header.
- Buffer::sync no longer drops non-same-GPU-id peers in low_latency_mode.
  DeepEP's original filter was safe because its LL path used NVSHMEM; this
  port drives LL via PortChannel so the kernel indexes
  port_channel_handles[local_expert*num_ranks + dst_rank] for every
  dst_rank. All peers now get a real memory/connection/semaphore/port
  channel entry.
- Add test/python/ext/ep/test_low_latency_multirank.py (LL dispatch+combine
  functional round-trip, BF16 only). Works cross-node in DeepEP's
  1-GPU-per-node topology.
- Known limitation documented in src/ext/ep/README.md and the test docstring:
  intra-node 8-GPU LL currently hangs because every peer transfer routes
  through the CPU proxy over IB loopback between distinct HCAs on the same
  host, and (separately) CudaIpcConnection::atomicAdd is a 64-bit op which
  mis-aligns the 32-bit rdma_recv_count slots when used for same-node
  peers. Proper fix needs a mixed-transport LL variant (MemoryChannel for
  same-node, PortChannel for cross-node) or 64-bit counters.
Gated behind MSCCLPP_EP_BENCH=1 to keep correctness runs fast. Reports
per-iter latency (max across ranks, CUDA-event timed) and aggregate
effective bandwidth (sum across ranks, dispatch+combine payload bytes).
Tunable via MSCCLPP_EP_BENCH_WARMUP / _ITERS / _TOKENS / _HIDDEN.

Bench reuses the Buffer allocated for the correctness phase and
self-skips if the requested hidden exceeds the per-peer NVL/RDMA budget.
Previously the optional benchmark measured full round-trip latency. Split
it to time dispatch alone (N iters) and combine alone (N iters reusing
one dispatch output), reporting per-phase latency (max across ranks) and
aggregate effective bandwidth (sum across ranks).

Applies to intranode HT, internode HT, and the (currently unreachable on
intra-node 8-GPU) LL test. Internode HT keeps the sync+barrier guard
between dispatch and combine but excludes it from either phase's timing.
…o int64

The low-latency dispatch/combine kernels signal recv counts via MSCCL++
PortChannel.atomicAdd, which lowers to IB IBV_WR_ATOMIC_FETCH_AND_ADD.
That opcode requires the remote address to be 8-byte aligned, but
LowLatencyLayout packed the per-expert signaling slots as int32. Odd
slots landed at offset %8 == 4; the NIC silently dropped those atomics
and the target rank spun forever in recv_hook (observed: even->odd
direction works, odd->even does not, across all tested topologies
including 2-rank intra-node, 8-rank intra-node, and 2-node 1-GPU-each).

Widen dispatch_rdma_recv_count_buffer / combine_rdma_recv_flag_buffer to
int64_t, update clean kernel + kernel signatures + next_clean pointers
accordingly, and add int64_t overloads for st_na_release /
ld_acquire_sys_global in utils.cuh.

Also drop the bogus self CUDA-IPC connection in Buffer::sync() that was
previously skewing the cross-rank buildAndAddSemaphore handshake order;
the kernel's same-rank branch uses a direct warp copy and never touches
the self port-channel slot (filled with a zero-initialized placeholder
so the [local_expert*num_ranks + dst_rank] indexing still holds).
Dropping the self ipc_cfg connection caused cudaErrorInvalidResourceHandle
on multi-node launches. Keep the self connection (needed by other code
paths that assume every rank is in the connections map) but continue to
skip the self slot in the semaphore + port-channel construction loops so
the kernel's [local_expert*num_ranks + dst_rank] indexing hits only peer
handles; the self slot is a zero-initialized placeholder since the
kernel's same-rank branch uses a direct warp copy.
The prior commit skipped r==rank in the semaphore and port-channel
build loops on the theory that the self-slot handshake skew was the
cause of LL direction asymmetry. That was wrong (the real bug was
int32 atomic alignment), and skipping self breaks other code paths
that assume every rank slot is represented -- cross-node HT and LL
failed with cudaErrorInvalidResourceHandle at the first barrier after
Buffer init. Restore the self-inclusive loop.
When all ranks live on the same host (num_rdma_ranks == 1), the LL
kernels now bypass PortChannel/IB-loopback entirely. In Buffer::sync()
we additionally:
  - allGather IPC handles for each rank's rdma_buffer_ptr and
    cudaIpcOpenMemHandle them into peer_rdma_bases[]
  - build per-peer MemoryChannels over CUDA IPC connections (tag=2)
    used only for the LL barrier ring

The three LL kernels (clean / dispatch / combine) gain a kIpcPath
template parameter and two extra args (peer_rdma_bases,
memory_channel_handles). At each peer op:
  - put -> peer-mapped warp copy over NVLink
  - atomicAdd-like flag store -> single-writer st_na_release on peer ptr
  - signal/wait barrier -> MemoryChannel signal/wait

Cross-node LL (num_rdma_ranks > 1) is untouched; the IPC setup block is
a no-op. The host launch wrappers select the variant via use_ipc_path.
Each local expert sends one copy per dispatched token back to its owner,
so the bytes actually on the wire during combine match dispatch. The
previous num_tokens×hidden under-counted by ~num_topk×, making combine
BW look artificially low next to dispatch.
- Report both per-rank and aggregate BW to align with NCCL-EP's ep_bench
  (which reports per-rank GB/s).
- Accept MSCCLPP_EP_LL_TOKENS/HIDDEN/TOPK/EXPERTS_PER_RANK env overrides
  so we can match external benchmark problem sizes (NCCL-EP LL defaults
  are num_tokens=128, hidden=7168, top_k=8).
Same alignment with NCCL-EP ep_bench as the LL test: report both
per-rank (agg/num_ranks) and aggregate throughput.
LL dispatch/combine are latency-bound at typical problem sizes: for
num_experts=32 the previous grid was cell_div(32,3)=11 blocks, i.e. 8%
of a 132-SM H100. The recv-side bodies already stride tokens by sm_id,
so extra blocks parallelize token work linearly. Extra blocks past
num_experts are gated out of the send/count phases by the existing
'responsible_expert_idx < num_experts' check.

Cap at the device's SM count (cooperative launch + launch_bounds(960,1)
allow one block per SM).
On the PortChannel (cross-node) path the extra blocks don't help: the
dispatch recv loop strides tokens per-warp-group (not per-SM), and the
additional blocks instead add cooperative-grid sync overhead and
increase concurrent host-proxy FIFO traffic. Measured cross-node
dispatch regressed from 1013us to 3063us when the unconditional grid
bump was active.

Keep the scaled grid for the IPC path (intra-node), where combine-recv
and dispatch token striding scale with sm_id and the 1.2-1.3x speedup
reproduces.
The LL combine benchmark was cloning the ~58 MB dispatch recv buffer
('recv_x.clone()') on every timed iteration, adding ~20 us of D2D
memcpy per sample and masking kernel-level changes. It also called
torch.empty() for the output inside the loop. Both now live outside
the timed region; the kernel is invoked against a persistent bench_out
and the recv_x produced by the most recent dispatch.
NCCL-EP's LL dispatch/combine kernel uses (numWarpGroups=1,
numWarpsPerGroup=32) when num_experts <= device_num_sms, giving each
SM ownership of a single expert and 32 warps to cooperate on its
recv-side per-(expert, src_rank) work. We were using (3, 10) — 3
experts per SM, 10 warps per (expert, rank) pair — which left a
significant amount of recv-side parallelism on the table because each
warp had to walk ~3x more tokens sequentially.

Switching to (1, 32) for both dispatch and combine matches NCCL-EP's
structure for typical EP sizes (num_experts in {32, 64, 256}) where
num_experts <= 132 SMs.

The static_assert kNumMaxTopK + 1 <= kNumWarpGroups * kNumWarpsPerGroup
still holds (9 <= 32) and the wider block also lets the staging loop
process the hidden-dim with one int4 per thread (hidden_bf16_int4=896
fits easily in 992 working threads).
Cross-node LL regressed when (1, 32) was applied uniformly: dispatch
1031us -> 1570us, combine 2553us -> 3484us. Larger grid means more
concurrent putWithSignal calls onto the host-proxy FIFO and a costlier
cg::this_grid().sync() between phases, both of which dominate the IB
path even though more SMs help the recv-side compute.

Make (kNumWarpGroups, kNumWarpsPerGroup) path-dependent: (1, 32) when
use_ipc_path, (3, 10) otherwise. Restores cross-node performance and
keeps the intra-node win.
- Add MSCCLPP_EP_BENCH_EXPERTS / _TOPK env knobs so the bench phase can
  match NCCL-EP's `ep_bench -a ht` defaults (256 experts, top-8). The
  functional check above continues to use the smaller (num_ranks*4
  experts, topk=4) configuration.

- Switch BW accounting from recv_tokens*hidden to bench_tokens*hidden,
  matching NCCL-EP's `RDMA_send` per-rank byte count. The previous
  formula counted DeepEP's expanded recv layout (one row per
  (token,src_rank) pair), inflating reported GB/s ~5x and making
  cross-stack comparisons misleading.
Same change as the intra-node bench (commit 4ed6f22), applied to the
cross-node test:

- Add MSCCLPP_EP_BENCH_EXPERTS / _TOPK env knobs so the bench phase can
  match NCCL-EP's `ep_bench -a ht` defaults (256 experts, top-8).
- Switch BW accounting from recv_tokens*hidden to bench_tokens*hidden,
  matching NCCL-EP's `RDMA_send` per-rank byte count.
Each mscclpp::ProxyService spawns one host-side proxy thread that
drains its FIFO and posts IB work requests. With LL combine pushing
~1k put + 60 atomicAdd FIFO entries per iter, that single thread is
the wall-clock bottleneck on cross-node runs.

Split the channel set across kNumProxyServices=4 separate services
so the host-side dispatch parallelism scales linearly. SemaphoreIds
and MemoryIds are scoped to a ProxyService, so:

- addMemory() is broadcast to every service in the same global order
  so a single MemoryId still identifies the memory everywhere.
- Each (peer_rank, channel_idx) is assigned to one proxy_idx via
  round-robin; the resulting PortChannel is built on that proxy and
  inherits its FIFO. The kernel is unchanged: the flat handle array
  routes the right way automatically.

No kernel-level changes, no tuning of QP count, no new env knobs.
… 1 on Blackwell)

Override at runtime with MSCCLPP_EP_NUM_PROXIES.
N=8 is the knee on H100+IB; N>=12 collapses from CPU oversubscription.
Intra-node LL is unchanged.
Add dist.barrier() + dist.destroy_process_group() in a finally block so
non-zero ranks don't poll the TCPStore after rank 0 (the store server)
exits, which produced noisy 'recvValue failed / Connection was likely
closed' stack traces from ProcessGroupNCCL's HeartbeatMonitor.

Also pass device_id to init_process_group in the internode test to
silence 'Guessing device ID based on global rank' warnings.
Aligns with NCCL-EP's ep_bench convention (BW computed from average time
across ranks). Previously we reported only the max time and computed BW
per-rank, which made our numbers more pessimistic than NCCL-EP's.
…ess import

Complete the previous consolidation commit (434121a), which only recorded the
deletion of ep_bench_mscclpp_ht.py because a stale pathspec made git add abort
before staging the real changes. This adds setup_mscclpp_ht into
ep_bench_mscclpp.py and repoints the harness import/parser map at it, so the
mscclpp-ht backend resolves again.
@azure-pipelines

Copy link
Copy Markdown
Azure Pipelines:
There may be pipelines that require an authorized user to comment /azp run to run.

Copilot AI 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.

Pull request overview

This PR extends the unified in-process Python EP benchmark harness to include an MSCCL++ high-throughput (HT) backend, while also simplifying the benchmark stack by removing the older standalone C++/CUPTI benchmark path. The result is a single Python harness intended to time dispatch→combine consistently across multiple EP implementations.

Changes:

  • Add a new --backend mscclpp-ht implementation (HT mode) in ep_bench_mscclpp.py, and wire it into the unified runner.
  • Refactor CUDA-graph capture to be harness-owned via a single “paired graph” helper, with backends providing capture-safe ops via a uniform {dispatch, combine, teardown, barrier, graph} contract.
  • Add per-backend Kineto kernel-name parsing helpers and centralize substring-based kernel-time summation in ep_bench_common.py.

Reviewed changes

Copilot reviewed 10 out of 10 changed files in this pull request and generated 2 comments.

Show a summary per file
File Description
test/python/ep/run_ep_bench.py Deleted legacy multi-process driver that shelled out to external binaries.
test/python/ep/run_ep_bench_python.py Adds mscclpp-ht, harness-owned single-graph capture, backend registry refactor, and pluggable Kineto parsing.
test/python/ep/mscclpp_ep_bench.cu Deleted standalone C++ LL benchmark binary source.
test/python/ep/ep_bench_nccl.py Moves graph capture responsibility to harness via graph_spec; adds Kineto parse helper.
test/python/ep/ep_bench_mscclpp.py Adds HT backend (setup_mscclpp_ht), adds Kineto parsing, and adapts LL backend to the new harness contract.
test/python/ep/ep_bench_flashinfer.py Adapts to harness-owned capture; adds Kineto parse helper; provides capture spec with pre_replay barrier.
test/python/ep/ep_bench_deepep.py Adapts to harness-owned capture; updates graph-capture gating and adds Kineto parse helper.
test/python/ep/ep_bench_common.py Adds shared sum_matching_kernel_us() helper for per-backend Kineto parsing.
test/python/ep/cupti_kernel_timer.cpp Deleted in-process CUPTI timer implementation.
test/python/ep/CMakeLists.txt Deleted standalone build for the removed C++ benchmark + CUPTI helper.

Comment thread test/python/ep/ep_bench_mscclpp.py Outdated
Comment thread test/python/ep/run_ep_bench_python.py Outdated
Address review comment: _kineto_kernel_us reads as a confusing name. Rename
to torch_profiler_kernel_us, which describes what it does (times the
dispatch/combine kernels with torch.profiler). Pure rename, no behavior
change.
…p8 note)

- Default --graph-group-size 10 -> 50 (reviewer: "Maybe increase to 50 by
  default?"): capture 50 dispatch->combine iterations per graph by default,
  further amortizing launch overhead / launch skew.
- Unify the internal name to graph_group_size everywhere (reviewer: "Different
  with iteration_per_group?" and "Why hard code to 1 here?"): the harness used
  iters_per_graph internally while the CLI arg is --graph-group-size, which read
  as two different concepts. Rename _capture_paired_graph / run_backend params
  and the effective-value local to graph_group_size; the "=1" default now reads
  as "no grouping" (1 iteration captured), which is why it is 1 when a backend
  is not graph-captured.
- Correct the FP8 wording (reviewer: "check if nccl support fp8 right?"):
  NCCL-EP DOES support FP8 (nccl_ep device code has token_data_type 0=FP8/uint8,
  calculate_fp8_scales, use_fp8). The previous "NCCL-EP path is bf16 only" text
  implied the library cannot; reword the --dispatch-dtype help and the guard to
  say the NCCL-EP path IN THIS BENCHMARK is BF16-only (the harness does not plumb
  NCCL-EP dispatch scales yet), not the library.

Verified: default --cuda-graph captures with graph_group_size=50, per-iteration
numbers unchanged; the fp8 guard prints the corrected message.
…nel sync)

Port the high-throughput RANK_MAJOR dispatch layout onto the current
barrier-channel HT runtime. RANK_MAJOR places each dispatched token at a
fixed [num_ranks, max_tokens_per_rank, hidden] slot grouped by source rank
instead of the compacted DeepEP prefix offset; the recv-pool combine path is
layout-agnostic and unchanged. Layout is threaded through dispatch.cu
(templated dispatchKernel), api.cuh, ht_runtime.{hpp,cc}, bindings.cpp, and
high_throughput.py. Defaults to TOKEN_MAJOR for compatibility.
…rror

In single-graph CUDA-graph mode both dispatch and combine replay inside
dispatch_fn() and combine_fn is a no-op, so the host-observed combine span is
~0us. Clamp comb_us to 1e-3 us so the downstream throughput division
(comb_bytes / c_avg) cannot raise ZeroDivisionError. Kernel-only kineto still
reports the true per-phase combine time.
…t a bug

Reword the DeepEP CUDA-graph gate comment: the RDMA/IB scale-out (GIN/IBGDA)
path is not graph-capturable because DeepEP internode transport drives
NVSHMEM/IBGDA put-signal operations that are illegal inside a CUDA graph (CUDA
719 in symmetric.hpp). This is a documented DeepEP internode limitation, not a
harness bug; we disable capture when GIN is active and run that path eagerly.
…omment

DeepEP V2 (ElasticBuffer) scale-out uses NCCL GIN (GPU-Initiated Networking,
backed by GDAKI/DOCA GPUNetIO on this stack), not the legacy NVSHMEM/IBGDA
Buffer path. Fix the earlier comment that misattributed the graph-capture gate
to NVSHMEM/IBGDA put-signal ops and symmetric.hpp. The real on-stream blocker
is the dispatch CPU sync for exact recv-token counts (do_cpu_sync); cached
dispatch forces it False, which is why the NVLink/MNNVL path is capture-safe.
Per review, keep a single CLI flag for the number of dispatch->combine
iterations captured inside one CUDA graph. Rename the arg dest to
iters_per_graph, drop the --graph-group-size alias, and update the validation
message, help text, comments, and the captured-graph log line accordingly. The
internal _capture_paired_graph/run_backend graph_group_size parameter (which
receives the value) is unchanged.
Binyang Li (Binyang2014) and others added 11 commits August 21, 2026 23:14
Commit automatically merged files first; conflict resolutions remain in the working tree for review.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Port throughput rank-major dispatch and the unified HT benchmark onto the refactored EP runtime while removing obsolete pre-refactor runtime files.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Use throughput naming for shared prepare, dispatch, and reduce-combine implementations that support both token-major and rank-major layouts.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Use generic prepare and notify entry points, derive throughput channel count from configuration, expose the receive pool through dispatchOutputBuffer, and remove legacy Python backend modules.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Keep test_latency_multirank.py as the maintained latency-mode multi-rank test.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Use one backend name per library, select MSCCL++ and NCCL-EP latency or throughput with --mode, and tune DeepEP's unified ElasticBuffer path with --num-sms.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Construct CommGroup directly from the benchmark's MPI or torch communicator.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Replace environment block-count overrides with --num-blocks and defer to the runtime's latency and throughput defaults when omitted.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Map --num-sms to MSCCL++'s total communication block count, including reserved control blocks, while retaining DeepEP's native SM-budget semantics.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
Clarify that latency num_blocks includes two reserved scheduler/control blocks while throughput uses all blocks as workers.

Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
@Binyang2014
Binyang Li (Binyang2014) merged commit 3d02152 into feature/ep Aug 22, 2026
5 checks passed
@Binyang2014
Binyang Li (Binyang2014) deleted the qinghuazhou/unified_ep_bench_ht_python branch August 22, 2026 03:55
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants