Skip to content

operatorx: multi-device gemm/moe (tp/dp/ep) under InferenceX recipes - #3526

Merged
hbarclay merged 6 commits into
hbarclay/operatorxfrom
hbarclay/operatorx-multidev
Sep 28, 2026
Merged

hbarclay merged 6 commits into
hbarclay/operatorxfrom
hbarclay/operatorx-multidev

Conversation

@hbarclay

@hbarclay hbarclay commented Sep 28, 2026 •

Copy link
Copy Markdown
Collaborator

Multi-device (single-node) GEMM and MoE for OperatorX, stacked on #3452 (base hbarclay/operatorx).

Schema

  • parallel: {"tp", "dp", "ep"} on gemm and moe args (core/parallel.py); args keep the full layer shape; no parallel = one device (same code path).
  • gemm: tp only — the weight is split over K, the all-reduce is part of the op.
  • moe: says how tokens (dp) and experts (tp sharded / ep partitioned) are divided, not how it's computed. vLLM mapping: tensor_parallel_size=tp, data_parallel_size=dp, enable_expert_parallel = ep > 1 (ep ∈ {1, tp×dp}). tokens is per DP group.

How a case runs (as InferenceX serves it)

  • recipes.py: planner matches each case to the InferenceX srt-slurm recipe for its source checkpoint on the runner's hardware (expanded with infx.srt_slurm / srtctl) → image, launch env, vllm serve args (+ the checkpoint's config.json, no weights). Per-framework topology mapping is one table entry (SPLITS), so SGLang can slot in beside vLLM.
  • runners/common/vllm/engine.py: each rank runs the recipe args through vLLM's EngineArgs / create_engine_config with external_launcher, then the worker's init_worker_distributed_environment. vLLM decides collectives, MoE dispatch (prepare/finalize), kernels, CUDA-graph sizes. Replaces the hand-built single-rank context.
  • Layers: RowParallelLinear; the MoE block under vLLM's parallel config (shared experts TP-sharded like DeepseekV2MLP, sequence-parallel chunking like DeepseekV2MoE, reduce_results=True); DP token counts from coordinate_batch_across_dp; graph capture inside graph_capture().
  • runners/common/ranks.py: ranks agree at every branch that launches collectives (build/trial failures, graph-capture fallback, throttle retries, warmup rounds) so one rank can't strand the others; timing reduces each iteration across ranks with CollectiveX's ep_harness._reduce_vec (MAX).
  • ci.py: one shard per (split, recipe); OPERATORX_PARALLEL / OPERATORX_ENGINE_ARGS / OPERATORX_MODEL_CONFIG; CollectiveX's network-env scrub + NCCL_CUMEM_ENABLE=1; per-step MASTER_PORT; timeout -k 30 around the rank srun; AMD no longer limited to one GPU (torch backend stays single-device). Old agentic/*.sh export-scraping recipe loader removed.
  • Workflow: planner inits utils/srt-slurm and runs with srtctl's light deps.

Testlists

gemm_parallel (50) / moe_parallel (46): every split InferenceX serves DSV4-Pro, DSV4.1-Flash, Kimi-K3, MiniMax-M3 with today — TP2/4/8, DP4+EP4, DP8+EP8, DP8, and SGLang's TP+EP (2/4/8) — at 8 and 256 tokens.

Status (tailscale-h200 dev box, no CI yet)

  • Both lists planned for cluster:h200-dgxc with ci.plan + main's recipes (14 cells) and run on 8×H200: 94 ok, 2 unsupported, 0 error. The 2 are AMD's MXFP4 MiniMax MoE (no NVIDIA kernel, same as single-device). Covers TP2/4/8, TP+EP 2/4/8, DP4+EP4, DP8, DP8+EP8, incl. cells under the DSV4.1-Flash and MiniMax-M3 H200 recipes (their images; recipe caps such as CUDA graphs ≤ 64 tokens honored).
  • vLLM picks NoDPEP for TP(+EP) and NaiveDPEP (AG/RS) for DP.
  • Every multi-device row records ranks: per-rank latency, profiled timeline and telemetry, cross-rank floor and skew (TP8 GEMM: 17.0 µs slowest-rank, 16.2 µs floor, ~1 µs skew).
  • Single-device MXFP8 MiniMax MoE on the new engine path: 1648 µs vs 1580 µs in the last H200 sweep.
  • Fixed on the way: wall-clock warmup ran different iteration counts per rank → collective hang; ranks now agree each round.

Open

  • This branch predates main's srt-slurm recipes, so CI runs from it find no recipes (cases fall back to default images) until it gets main's benchmarks/single_node/srt-slurm-recipes + infx/srt_slurm.
  • DSV4 DEP's MegaMoE (moe-backend deep_gemm_amxf4_mega_moe) is chosen in DSV4's model code, not FusedMoE; the recipe arg is passed but this harness builds FusedMoE, so DEP runs vLLM's generic EP path.
  • Quantization still comes from the op's operand descriptors (_scheme), not yet from the checkpoint's quantization_config.
  • Test runs will dispatch with -f ingest=false.

🤖 Generated with Claude Code

@functionstackx

Copy link
Copy Markdown
Collaborator

Sorry, over the weekend, there was 2 major refactors to clean up the technical debt accumalated over the past 11 months of moving at the speed of light. We don't see any major refactors in the forthseeable future besides cleaning up AMD multinode AgentX pile of bash. As much, due to the refactors, u would need to ask your agent to rebase from remote main@latest. Thank you in advance for ur understanding

@hbarclay
hbarclay force-pushed the hbarclay/operatorx-multidev branch 3 times, most recently from 82f35b9 to 894251b Compare September 28, 2026 02:51
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@hbarclay
hbarclay force-pushed the hbarclay/operatorx-multidev branch from f3e1e77 to b39576f Compare September 28, 2026 04:06
hbarclay and others added 5 commits September 28, 2026 00:15
…ion config

vLLM rejects max_cudagraph_capture_size given both as its own argument and in the
compilation config, which a non-eager engine now receives.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Attention sets enable_jit_warmup only where the image's KernelConfig has it; a MoE case
on an image whose vLLM lacks the API the generic block uses is unsupported, not an error
(the DeepSeek-V4 DEP recipe image runs vLLM 0.17).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…don't split along it

An attention case on the DeepSeek-V4 DEP recipe (DP attention, MegaMoE experts) kept
the recipe's MoE backend but lost its expert parallelism, which MegaMoE requires.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…point's width

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@hbarclay
hbarclay marked this pull request as ready for review September 28, 2026 17:32
@hbarclay
hbarclay merged commit c4b765e into hbarclay/operatorx Sep 28, 2026
13 of 25 checks passed
@hbarclay
hbarclay deleted the hbarclay/operatorx-multidev branch September 28, 2026 17:32
@claude

claude Bot commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Claude finished @hbarclay's task in 2m 30s —— View job


Review of PR #3526: multi-device GEMM/MoE (TP/DP/EP)

  • Gather context (PR head 425b44b vs. base hbarclay/operatorx)
  • Review core changes (parallel, ranks, timing, engine, vLLM linear/moe, ci.py sharding and launch)
  • Post findings

Most branches that launch collectives now check that every rank agrees before continuing: build and run go through together, graph-capture fallback through agree, throttle retries and the warmup loop through any_. The sharding per (split, recipe, layer/engine) and the per-cell MASTER_PORT also look right. I found three problems that can still hang a shard or report wrong data:

1. 🟡 runners/common/telemetry.py:314: the ranks block can come from a different attempt than latency_us (inline comment posted)
measure() keeps the best attempt, but ranks.LAST is overwritten on every attempt. So after a throttle retry, ranks.latency_us_min, skew_us and per_rank[].latency_us describe the last attempt, not the one reported. Fix: store dict(ranks.LAST) in best and restore it after the loop.

2. 🟡 runners/common/profiling.py:151-159 (and _markers_pass): ranks don't agree on the profiler pass
Suppose torch.profiler.profile(...) (or the NVTX pass) fails to start on one rank. That rank returns {"error": ...} without calling kernel_fn, while the others replay TP all-reduce or EP dispatch iters times and block until timeout -k 30 kills the shard. This breaks the invariant in ranks.py. Fix: enter the profiler, then if not ranks.agree(started): return {...}, then replay.

3. 🟡 ci.py:736-737: counters mode will probably deadlock when world_size > 1
Each rank's ncu uses the default kernel replay. The opx… NVTX ranges now contain cross-rank kernels: the all-reduce inside RowParallelLinear(reduce_results=True), MoE dispatch/combine, and the SP all-gather. Replaying one of these on a single rank waits forever on its peers. plan() doesn't reject mode="counters" together with multi-device world sizes. Fix: add if mode == "counters" and world_sizes != [1]: raise ValueError(...) in plan(), or switch to --replay-mode application (or range replay with --communicator tcp --lockstep-kernel-launch) for multi-rank runs.

This PR is already merged into hbarclay/operatorx, so these would need a follow-up on that branch. I didn't make any code changes and didn't run anything; this review only reads the diff.

@@ -328,7 +312,7 @@ def measure(op, time_once) -> tuple[float, dict]:
_dump(op, attempt, summary, samples)
if best is None or median_us < best[0]:
best = (median_us, summary, sleep_s)

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.

🟡 WARNING: The ranks block can describe a different attempt than latency_us.

Why it matters: Each time_once call goes through ranks.iterations(), and that overwrites ranks.LAST. This loop keeps the best attempt's median, but ranks.summary() (called later in runners/*/runner.py) reads LAST, which holds the last attempt. So after a throttle retry where the last attempt isn't the best one, ranks.latency_us_min, skew_us and per_rank[].latency_us come from another run than the reported latency_us. For example, latency_us_min can end up larger than latency_us.

Fix: Save LAST together with the best attempt and restore it after the loop:

        if best is None or median_us < best[0]:
            best = (median_us, summary, sleep_s, dict(ranks.LAST))
    ...
    median_us, summary, sleep_s, last = best
    ranks.LAST.clear()
    ranks.LAST.update(last)

@claude claude Bot 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.

Beyond the inline findings, I also checked the recipe-selection fallback in recipes.py (find(), lines 76-84): when no variant matches a case's actual tp/dp/ep split, it does fall back to the first variant at any topology for that checkpoint, but this is not silent — ci.py records the outcome via the recipe_match field (ci.py:233), so a mismatched recipe is surfaced rather than hidden.

Extended reasoning...

This operatorx PR reworks multi-rank GEMM/MoE execution (CUDA-graph capture, rank-agreement barriers in ranks.py, vLLM engine integration) with a confirmed inline finding that a fault during CUDA-graph capture on one rank (linear.py _launcher, called outside any ranks.together block) can hang sibling ranks instead of failing the run cleanly. I additionally traced the recipes.py checkpoint-fallback path that another reviewer flagged as a candidate silent-mismatch bug, and confirmed ci.py surfaces the match quality via recipe_match rather than swallowing it, so that path is not an independent bug. No security-sensitive surface (auth/crypto/permissions) is touched; the main risk is the confirmed distributed-hang finding, which already warrants a human look.

Comment on lines 311 to 313
except Exception as e:
if _is_fault(e):
raise

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.

🔴 In multi-rank runs, a CUDA/device fault during CUDA-graph capture on one rank hangs every other rank instead of failing the run. _launcher (called directly from run() in nvidia/runner.py:50, outside any ranks.together block) only calls ranks.agree(err is None) on the non-fault branch; when _is_fault(e) is true at line 312-313 it re-raises immediately without ever calling ranks.agree. Other ranks that succeeded capture then block forever inside their own ranks.agree() collective, which needs every rank to participate. Fix: make the fault path also participate in a rank-agreement collective (or wrap the whole launcher in ranks.together) before re-raising, for both call sites: linear.py:312-313 (shared by gemm and moe via moe.py's import) and attention.py:562-563.

Why this was flagged

Trigger: on an OOM or other device fault (torch.OutOfMemoryError, "CUDA error", "illegal memory", "device-side assert") during torch.cuda.graph() capture on one rank of a multi-rank TP/DP run, while sibling ranks capture successfully - plausible at the large batch/token sizes this PR's testlists exercise (gemm_parallel/moe_parallel at TP8/DP8). Entry: nvidia/runner.py:50 impl.launcher(ctx), called directly, not inside ranks.together. Bug: linear.py:312-313 if _is_fault(e): raise propagates the exception without calling ranks.agree, while other ranks proceed to linear.py:316 ranks.agree(err is None), a collective all_reduce that now never completes. Base branch had no multi-rank path at all, so this hang is new. ranks.py's own docstring states "Every rank must take the same branch wherever a collective launches, or the others hang" - this path violates that invariant. Same pattern in attention.py's _launcher (line ~562-566).

Verification: Severity: normal. The candidate is real. In _launcher (operatorx/runners/common/vllm/linear.py:311-316), the fault branch if _is_fault(e): raise (line 312-313) re-raises WITHOUT any cross-rank collective, whereas the normal branch reaches if not ranks.agree(err is None) (line 316), a blocking all_reduce (ranks.py:24-31). _launcher is called at nvidia/runner.py:50 from run(), which is…

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

Status: Done

Development

Successfully merging this pull request may close these issues.

2 participants