Compare commits

..
Author SHA1 Message Date
Tyler Michael SmithandGitHub 6257a61c3d Merge branch 'main' into worktree-fix-cudagraph-flaky 2026-07-06 13:59:00 -04:00
Xiaohong (Sean) ChenGitHubAndreas Karatzasmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
9fde043f54 [Kernel][Helion][1/N] Add Helion kernel for silu_and_mul_per_block_quant (#43994)
Signed-off-by: Sean Chen <seachen@redhat.com>
Co-authored-by: Andreas Karatzas <akaratza@amd.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-07 00:19:01 +08:00
24dd2aec81 [Bugfix] Preserve FP8 indexer WK pairs across incremental load_weights (#46168)
Signed-off-by: lcheng <lcheng321@gatech.edu>
Signed-off-by: NickLucche <nicolo.lucchesi@mistral.ai>
Co-authored-by: Isotr0py <mozf@mail2.sysu.edu.cn>
Co-authored-by: Nicolò Lucchesi <nlucches@redhat.com>
Co-authored-by: NickLucche <nicolo.lucchesi@mistral.ai>
2026-07-06 09:16:46 -07:00
RanranandGitHub 3ee9eea928 [macOS][CPU][Installation] Fix the broken installation of vllm 0.24.0 in macos + cpu (#47457)
Signed-off-by: Ranran Haoran Zhang <ranranhaoranzhang@gmail.com>
2026-07-06 08:59:16 -07:00
Harry MellorGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
5bce653e09 Make the Transformers modeling backend as fast as native vLLM (#47187)
Signed-off-by: Harry Mellor <19981378+hmellor@users.noreply.github.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-06 16:59:14 +01:00
Kevin_XiongGitHubCodexIsotr0pymergify[bot] <37929162+mergify[bot]@users.noreply.github.com>Isotr0pyIsotr0py
5ad11172b7 [perf]Add fused Kimi image preprocessing (#47416)
Signed-off-by: Kevin-XiongC <kevin_xiong1997@outlook.com>
Signed-off-by: Kevin_Xiong <kevin_xiong1997@outlook.com>
Signed-off-by: Isotr0py <Isotr0py@outlook.com>
Co-authored-by: Codex <codex@openai.com>
Co-authored-by: Isotr0py <2037008807@qq.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
Co-authored-by: Isotr0py <mozf@mail2.sysu.edu.cn>
Co-authored-by: Isotr0py <Isotr0py@outlook.com>
2026-07-06 08:46:32 -07:00
Wentao YeandGitHub f70caef48b [Perf] Cache token_to_req_indices for dsv4, 5x~6x kernel performance improvement (#47474)
Signed-off-by: yewentao256 <zhyanwentao@126.com>
2026-07-06 11:17:46 -04:00
Laurent-ZhangGitHubClaudemergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
8d8ec38361 [Bugfix][Spec Decode] Add missing draft_id_to_target_id to DSparkDeepseekV4ForCausalLM (#47429)
Signed-off-by: Laurent-Zhang <zhangdongsheng80@gmail.com>
Co-authored-by: Claude <noreply@anthropic.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-06 10:55:47 -04:00
Wentao YeandGitHub b1c6dba558 [Refactor] Remove multiple dead code (#47329)
Signed-off-by: yewentao256 <zhyanwentao@126.com>
2026-07-06 07:54:08 -07:00
jescoGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
598d51153a [Bugfix][Distributed] Delegate MNNVL allreduce one-shot selection (#47589)
Signed-off-by: jesco-absolut <team@srswti.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-06 07:47:06 -07:00
Yifan QiaoandGitHub 095adf1fdc [Bugfix] Fix int32 overflow in triton_decode_attention page offsets (#47671)
Signed-off-by: Yifan Qiao <yifanqiao@inferact.ai>
2026-07-06 10:36:15 -04:00
Harry MellorandGitHub 51ee564e56 [CI] Skip test for checkpoint that was deleted (#47748)
Signed-off-by: Harry Mellor <19981378+hmellor@users.noreply.github.com>
2026-07-06 07:24:09 -07:00
Ting SUNGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
373eb314af [Bugfix][Core] Fix num_output_placeholders underflow with async scheduling + spec decode (#46066)
Signed-off-by: Ting Sun <suntcrick@gmail.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-06 13:50:38 +00:00
641cb59592 [Doc] Clarify fastokens availability (#45813)
Signed-off-by: LjjJzd <3542531707@qq.com>
Co-authored-by: OpenAI Codex <codex@openai.com>
2026-07-06 13:33:05 +00:00
07f9baf756 Revert "[Platform] Replace torch.cuda.Event with torch.Event (#47140)" (#47668)
Signed-off-by: Kunshang Ji <kunshang.ji@intel.com>
Co-authored-by: Harry Mellor <19981378+hmellor@users.noreply.github.com>
2026-07-06 14:18:33 +01:00
7a90eb98ab [Bugfix] [Gemma4] Fix Gemma4 MTP draft model layers ignoring quant_config (#47091)
Signed-off-by: Ayushman Singh <40520701+ayush1399@users.noreply.github.com>
Co-authored-by: Benjamin Chislett <bchislett@nvidia.com>
2026-07-06 14:04:00 +01:00
Bugen ZhaoGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
8f4c69b222 [Rust Frontend] Cache metric handles for scheduler & request stats (#47444)
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
Signed-off-by: Bugen Zhao <i@bugenzhao.com>
2026-07-06 13:02:59 +00:00
8b79971bb9 attention: pass None for unused args in unified attention TD path (#43597)
Signed-off-by: Artur Fierka <artur.fierka@intel.com>
Co-authored-by: quinnlp <quinnlp@users.noreply.github.com>
2026-07-06 21:01:21 +08:00
Nick HillandGitHub f676808ba0 [CI] Use TTY for AMD CI tests for colored buildkite logs (#47730)
Signed-off-by: Nick Hill <nickhill123@gmail.com>
2026-07-06 20:50:29 +08:00
Qiming ZhangandGitHub 98e4726a14 [fix][run_batch]: respect proxy env vars when downloading media URLs (#47697)
Signed-off-by: mauyuyuace <qiming1.zhang@intel.com>
2026-07-06 12:45:48 +00:00
BadrBasowidandGitHub 740f379fae [ROCm][AITER] Directly Implement AITER Custom All-reduce in CudaCommunicator (#46065)
Signed-off-by: BadrBasowid <badr.basowid@gmail.com>
2026-07-06 12:16:32 +00:00
Alexis K.andGitHub 40cc2e8327 [Bugfix] Return HTTP 422 for unprocessable image URLs instead of 500 (#47165)
Signed-off-by: Alexis Kinsella <alexis.kinsella@gmail.com>
2026-07-06 11:56:23 +00:00
Juan Pérez de AlgabaGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
ba22152096 fix(security): block request-level GPU video backend selection withou… (#47259)
Signed-off-by: jperezde <jperezde@redhat.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-06 02:36:49 -07:00
Yan MaandGitHub 90ce3a09be [bugfix] fix MOSS-Audio deepstack_input_embeds initialization in PP (#47607)
Signed-off-by: Yan Ma <yan.ma@intel.com>
2026-07-06 17:15:50 +08:00
26c754d847 [XPU][Bugfix] Do not transpose weight_scale_inv at load time (#47116)
Signed-off-by: Ma Jian <jian1.ma@intel.com>
Co-authored-by: Kunshang Ji <kunshang.ji@intel.com>
2026-07-06 17:15:26 +08:00
Sungjae LeeandGitHub 3d7f357ebf [Doc] docs: fix note formatting for pooling models (#47701)
Signed-off-by: Sungjae Lee <33976427+llsj14@users.noreply.github.com>
Signed-off-by: Sungjae Lee <sung-jae.lee@navercorp.com>
2026-07-06 09:01:10 +00:00
liuzhenweiandGitHub 736f1a5907 [XPU] Route mm_prefix models to Triton attention backend (#47688)
Signed-off-by: zhenwei-intel <zhenwei.liu@intel.com>
2026-07-06 16:52:44 +08:00
Li, JiangandGitHub 344609ab17 [CI/Build] Fix pre-commit check (#47695)
Signed-off-by: jiang1.li <jiang1.li@intel.com>
2026-07-06 08:24:24 +00:00
xiaozhoupyandGitHub d039c17114 [Bugfix] Recycle post-final-norm hidden in GLM MTP (single norm) (#47448) 2026-07-06 01:07:56 -07:00
xiangdongandGitHub cdab28319f [XPU][CI]Add agent tags for Basic Models Tests (Initialization) in Intel GPU CI (#47675)
Signed-off-by: zengxian <xiangdong.zeng@intel.com>
2026-07-06 15:15:45 +08:00
Qiming ZhangandGitHub 2fa10566e3 [Core][DP] Rotate load-balancer tie-break to avoid systematic engine bias (#47420)
Signed-off-by: mayuyuace <qiming1.zhang@intel.com>
2026-07-06 07:09:16 +00:00
Andreas KaratzasandGitHub fb265fc8fb [ROCm][CI] Increasing parallelism in Basic Models Tests (Extra Initialization) (#47591)
Signed-off-by: Andreas Karatzas <akaratza@amd.com>
2026-07-06 15:06:16 +08:00
Andreas KaratzasandGitHub 8f0e75e16b [ROCm][CI] Adding nixl multiconn (#47481)
Signed-off-by: Andreas Karatzas <akaratza@amd.com>
2026-07-06 15:04:58 +08:00
98ba9b9583 [Frontend] Support OpenAI Responses API namespace tools (#47024)
Signed-off-by: zhongjing123 <jimzhong5193@gmail.com>
Co-authored-by: zhongjing123 <jimzhong5193@gmail.com>
2026-07-06 06:21:27 +00:00
velonica0andGitHub 990c2a0187 [RISC-V] Enable BF16 on VLEN=256 hardware (#45243)
Signed-off-by: velonica0 <like@mail.nankai.edu.cn>
2026-07-06 06:05:16 +00:00
e433634c78 [Performance][Hardware][RISC-V] Reduce LMUL pressure in INT4 LUT dequant (#47538)
Signed-off-by: liutong <liutong@iscas.ac.cn>
Co-authored-by: Claude <noreply@anthropic.com>
2026-07-06 05:58:56 +00:00
16f8110935 [Bugfix][CPU][RISC-V] Fix VLEN detection for RVV attention path (#47532)
Signed-off-by: liutong <liutong@iscas.ac.cn>
Co-authored-by: Claude <noreply@anthropic.com>
2026-07-06 05:58:03 +00:00
Zhenzhong XuGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
d9c1767cd4 [INC][ARK] Direct Register Custom Op for ARK (#46361)
Signed-off-by: Zhenzhong1 <zhenzhong.xu@intel.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-06 13:45:50 +08:00
Li, JiangandGitHub e9cc1fd093 [CI/Build][CPU] Remove global extra index (#47687)
Signed-off-by: jiang1.li <jiang1.li@intel.com>
2026-07-06 13:42:01 +08:00
Fadi ArafehandGitHub f1073c050c [CPU][BugFix] Multiple fixes to w4a8_int8 CPU MoE path (#46739)
Signed-off-by: Fadi Arafeh <fadi.arafeh@arm.com>
2026-07-06 05:39:20 +00:00
Qiming ZhangandGitHub 394edc8108 [XPU] limit max-num-seqs in test_lmeval.py for XPU (#47682)
Signed-off-by: mauyuyuace <qiming1.zhang@intel.com>
2026-07-06 05:34:16 +00:00
69715823df [Test][XPU] Skip fork in kv_sharing_fast_prefill test on XPU (#47406)
Signed-off-by: Ma, Liangliang <liangliang.ma@intel.com>
Co-authored-by: Kunshang Ji <kunshang.ji@intel.com>
2026-07-06 11:32:26 +08:00
Chaojun ZhangandGitHub 6569df6a3e [Test][LoRA] Use lightweight CPU reference and skip heavy cleanup in punica ops tests (#47534)
Signed-off-by: Chaojun Zhang <chaojun.zhang@intel.com>
2026-07-06 11:29:59 +08:00
f2aaf59151 [Feature] Support MTP speculative decoding for Bailing hybrid models (#44880)
Signed-off-by: zc02384840 <zc02384840@antgroup.com>
Co-authored-by: zc02384840 <zc02384840@antgroup.com>
2026-07-06 10:38:50 +08:00
95a248faed [Attention Backend] HPC_ATTN backend support mtp and dynamic scheduled attention (#47433)
Signed-off-by: chengvjiang <chengvjiang@tencent.com>
Co-authored-by: chengvjiang <chengvjiang@tencent.com>
2026-07-05 18:18:25 -07:00
Chaojun ZhangGitHubAndreas Karatzasmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>Kunshang Ji
d2ec433e37 [XPU] Fix Eagle3 initialization on XPU (#43957)
Signed-off-by: Chaojun Zhang <chaojun.zhang@intel.com>
Co-authored-by: Andreas Karatzas <akaratza@amd.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
Co-authored-by: Kunshang Ji <kunshang.ji@intel.com>
2026-07-06 08:46:05 +08:00
Łukasz ŚlusarczykGitHubmergify[bot] <37929162+mergify[bot]@users.noreply.github.com>Kunshang Ji
78a04c208d [XPU] Fix CUDA API shims breaking Torch Dynamo during AOT compile (#43092)
Signed-off-by: Łukasz Ślusarczyk <lukasz.slusarczyk@intel.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
Co-authored-by: Kunshang Ji <kunshang.ji@intel.com>
2026-07-06 08:29:20 +08:00
Spandan TiwariandGitHub b71218107f [ROCm][Test] Fix test_per_token_group_quant_fp8 tolerance for 1-ULP FP8 rounding on gfx950 (#46944)
Signed-off-by: Spandan Tiwari <sptiwari@amd.com>
2026-07-05 18:02:30 -05:00
cc1d020d01 [MRV2] Enable mm prefix bidi attention support on MRV2 (#46942)
Signed-off-by: Isotr0py <mozf@mail2.sysu.edu.cn>
Signed-off-by: Isotr0py <2037008807@qq.com>
Co-authored-by: Nick Hill <nickhill123@gmail.com>
2026-07-05 14:45:29 +00:00
Ting SUNandGitHub 8974ed89cd [Bugfix][Voxtral Realtime] Fix token feedback timeout silent hang (#44461)
Signed-off-by: Ting Sun <suntcrick@gmail.com>
2026-07-05 05:42:36 -07:00
Ting SUNGitHubWentao Yemergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
fb2faceacd [Bugfix][Model] Fix crash loading Mamba/Mamba2 checkpoints without an architectures field (#46037)
Signed-off-by: Ting Sun <suntcrick@gmail.com>
Signed-off-by: Ting SUN <suntcrick@gmail.com>
Co-authored-by: Wentao Ye <44945378+yewentao256@users.noreply.github.com>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
2026-07-05 05:42:32 -07:00
b6cc46ec3b [Feature] Support sequence parallel without the need for DP, 1.9%~5.0% E2E Throughput Improvement (#47070)
Signed-off-by: yewentao256 <zhyanwentao@126.com>
Signed-off-by: gcanlin <canlinguosdu@gmail.com>
Co-authored-by: Canlin Guo <canlinguosdu@gmail.com>
2026-07-05 05:41:30 -07:00
Lucas WilkinsonandGitHub fa4321de3d [Bugfix][TurboQuant] Preserve KV cache dtype in backend shape (#47609) 2026-07-05 08:20:48 +00:00
Ting SUNandGitHub 9226613043 [Bugfix][Pooling] Forward instruction to Jina reranker scoring prompts (#47590)
Signed-off-by: Ting Sun <suntcrick@gmail.com>
2026-07-05 05:39:13 +00:00
Tyler Michael SmithandClaude Opus 4.6 cc17a5e0d1 [CI] Fix flaky cudagraph mode test by isolating each case in its own process
The test_cudagraph_compilation_combo and test_backend_and_cudagraph_mode_combo
tests ran all parametrized cases in a single pytest process, relying on
weakref + wait_for_gpu_memory_to_clear(120s) between cases. After many
LLM create/destroy cycles, accumulated CUDA driver state could delay
subprocess memory reclamation past the 120s timeout, causing flaky failures
(e.g. FA2-FULL_DECODE_ONLY-3-True which captures 51 full CUDA graphs).

Wrap both tests with @create_new_process_for_each_test("spawn") so each
parametrized case gets a clean process, eliminating the cross-test memory
accumulation and removing the need for manual teardown.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Signed-off-by: Tyler Michael Smith <tlrmchlsmth@gmail.com>
2026-06-12 17:57:10 -04:00
197 changed files with 12286 additions and 1444 deletions
+24
View File
@@ -76,6 +76,30 @@ steps:
pytest -v -s v1/sample/test_logprobs.py &&
pytest -v -s v1/sample/test_logprobs_e2e.py'
- label: Basic Models Tests (Initialization)
timeout_in_minutes: 60
device: intel_gpu
agent_tags:
label: production
gpu: 1+
mem: 24+
no_plugin: true
working_dir: "."
env:
REGISTRY: "public.ecr.aws/q9t5s3a7"
REPO: "vllm-ci-test-repo"
VLLM_TEST_DEVICE: "xpu"
source_file_dependencies:
- vllm/
- tests/models/test_initialization.py
- tests/models/registry.py
commands:
- >-
bash .buildkite/scripts/hardware_ci/run-intel-test.sh
'export VLLM_XPU_FUSED_MOE_USE_REF=1 &&
cd tests &&
pytest -v -s models/test_initialization.py::test_can_initialize_large_subset[Eagle3MiniMaxM2ForCausalLM]'
- label: XPU CPU Offload
timeout_in_minutes: 60
device: intel_gpu
@@ -551,6 +551,7 @@ else
fi
docker run \
-t -i \
--device /dev/kfd $BUILDKITE_AGENT_META_DATA_RENDER_DEVICES \
$RDMA_FLAGS \
--network=host \
@@ -38,7 +38,8 @@ function cpu_tests() {
pytest -x -v -s tests/kernels/attention/test_cpu_attn.py
pytest -x -v -s tests/kernels/core/test_cpu_activation.py
pytest -x -v -s tests/kernels/moe/test_cpu_fused_moe.py
pytest -x -v -s tests/kernels/mamba/cpu/test_cpu_gdn_ops.py"
pytest -x -v -s tests/kernels/mamba/cpu/test_cpu_gdn_ops.py
pytest -x -v -s tests/kernels/moe/test_cpu_int4_moe.py"
# skip tests requiring model downloads if HF_TOKEN is not set
# due to rate-limits
+3 -1
View File
@@ -109,7 +109,9 @@ run_nodes() {
if [ "$node" -ne 0 ]; then
docker exec -d "node$node" /bin/bash -c "cd $WORKING_DIR ; ${COMMANDS[$node]}"
else
docker exec "node$node" /bin/bash -c "cd $WORKING_DIR ; ${COMMANDS[$node]}"
# Allocate a TTY (-t -i) for the foreground head node so its output
# keeps ANSI color in the Buildkite log (see run-amd-test.sh).
docker exec -t -i "node$node" /bin/bash -c "cd $WORKING_DIR ; ${COMMANDS[$node]}"
fi
done
}
+68 -11
View File
@@ -472,7 +472,7 @@ steps:
commands:
- TARGET_TEST_SUITE=MI300 pytest basic_correctness/ -v -s -m 'distributed(num_gpus=2)'
- CUDA_VISIBLE_DEVICES=0,1 pytest -v -s model_executor/model_loader/test_sharded_state_loader.py -m '(not slow_test)'
- pytest models/test_transformers.py -v -s -m 'distributed(num_gpus=2)'
- pytest models/transformers/test_backend.py -v -s -m 'distributed(num_gpus=2)'
- pytest models/language -v -s -m 'distributed(num_gpus=2)'
- pytest models/multimodal -v -s -m 'distributed(num_gpus=2)' --ignore models/multimodal/generation/test_whisper.py --ignore models/multimodal/generation/test_phi4siglip.py
- pytest models/multimodal/generation/test_phi4siglip.py -v -s -m 'distributed(num_gpus=2)'
@@ -1301,7 +1301,7 @@ steps:
#---------------------------------------------------------- mi300 · kernels ----------------------------------------------------------#
- label: vLLM IR Tests # TBD
timeout_in_minutes: 30
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_1
optional: true
@@ -1350,7 +1350,7 @@ steps:
- pytest -v -s kernels/core --ignore=kernels/core/test_minimax_reduce_rms.py kernels/test_concat_mla_q.py kernels/test_top_k_per_row.py
- label: Kernels KDA Test # TBD
timeout_in_minutes: 30
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_1
optional: true
@@ -1558,7 +1558,7 @@ steps:
- TP_SIZE=1 DP_SIZE=2 pytest -v -s v1/distributed/test_eagle_dp.py
- label: Model Runner V2 Pipeline Parallelism (4 GPUs) # TBD
timeout_in_minutes: 60
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_4
num_gpus: 4
@@ -1577,7 +1577,7 @@ steps:
- pytest -v -s distributed/test_pp_cudagraph.py -k "not ray"
- label: Model Runner V2 Spec Decode # TBD
timeout_in_minutes: 45
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_1
optional: true
@@ -1604,7 +1604,7 @@ steps:
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_1
parallelism: 2
parallelism: 6
working_dir: "/vllm-workspace/tests"
source_file_dependencies:
- vllm/model_executor/models/
@@ -1637,10 +1637,10 @@ steps:
source_file_dependencies:
- vllm/
- tests/models/test_terratorch.py
- tests/models/test_transformers.py
- tests/models/transformers/test_backend.py
- tests/models/test_registry.py
commands:
- pytest -v -s models/test_terratorch.py models/test_transformers.py models/test_registry.py
- pytest -v -s models/test_terratorch.py models/transformers/test_backend.py models/test_registry.py
#----------------------------------------------------- mi300 · models / language -----------------------------------------------------#
@@ -1860,7 +1860,7 @@ steps:
- examples/
commands:
- pip install --upgrade git+https://github.com/huggingface/transformers
- pytest -v -s tests/models/test_transformers.py
- pytest -v -s tests/models/transformers/test_backend.py
- pytest -v -s tests/models/multimodal/test_mapping.py
- python3 examples/basic/offline_inference/chat.py
- python3 examples/generate/multimodal/vision_language_offline.py --model-type qwen2_5_vl
@@ -1915,7 +1915,7 @@ steps:
- pytest -v -s plugins_tests/lora_resolvers # unit tests for in-tree lora resolver plugins
- label: GGUF Plugin # TBD
timeout_in_minutes: 30
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_1
optional: true
@@ -2439,6 +2439,59 @@ steps:
- uv pip install --system -r /vllm-workspace/requirements/kv_connectors_rocm.txt
- HYBRID_SSM=1 ATTENTION_BACKEND=TRITON_ATTN bash v1/kv_connector/nixl_integration/config_sweep_accuracy_test.sh
- label: Hybrid SSM NixlConnector PD prefix cache test (2 GPUs) # TBD
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_2
num_gpus: 2
optional: true
working_dir: "/vllm-workspace/tests"
source_file_dependencies:
- vllm/distributed/kv_transfer/kv_connector/v1/nixl/
- vllm/v1/core/sched/
- vllm/v1/core/kv_cache_coordinator.py
- tests/v1/kv_connector/nixl_integration/
- vllm/platforms/rocm.py
commands:
- uv pip install --system -r /vllm-workspace/requirements/kv_connectors_rocm.txt
- ATTENTION_BACKEND=TRITON_ATTN bash v1/kv_connector/nixl_integration/run_mamba_prefix_cache_test.sh
- label: MultiConnector (Nixl+Offloading) PD accuracy (2 GPUs) # TBD
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_2
num_gpus: 2
optional: true
working_dir: "/vllm-workspace/tests"
source_file_dependencies:
- vllm/distributed/kv_transfer/kv_connector/v1/nixl/
- vllm/distributed/kv_transfer/kv_connector/v1/multi_connector.py
- vllm/distributed/kv_transfer/kv_connector/v1/offloading_connector.py
- vllm/distributed/kv_transfer/kv_connector/v1/offloading/
- tests/v1/kv_connector/nixl_integration/
- vllm/platforms/rocm.py
commands:
- uv pip install --system -r /vllm-workspace/requirements/kv_connectors_rocm.txt
- ATTENTION_BACKEND=TRITON_ATTN bash v1/kv_connector/nixl_integration/run_multi_connector_accuracy_test.sh
- label: MultiConnector (Nixl+Offloading) PD edge cases (2 GPUs) # TBD
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
agent_pool: mi300_2
num_gpus: 2
optional: true
working_dir: "/vllm-workspace/tests"
source_file_dependencies:
- vllm/distributed/kv_transfer/kv_connector/v1/nixl/
- vllm/distributed/kv_transfer/kv_connector/v1/multi_connector.py
- vllm/distributed/kv_transfer/kv_connector/v1/offloading_connector.py
- vllm/distributed/kv_transfer/kv_connector/v1/offloading/
- tests/v1/kv_connector/nixl_integration/
- vllm/platforms/rocm.py
commands:
- uv pip install --system -r /vllm-workspace/requirements/kv_connectors_rocm.txt
- ATTENTION_BACKEND=TRITON_ATTN bash v1/kv_connector/nixl_integration/run_multi_connector_edge_case_test.sh
- label: V1 e2e (4 GPUs) # TBD
timeout_in_minutes: 180
mirror_hardwares: [amdexperimental, amdproduction, amdgfx942nightly, amdmi300]
@@ -2729,15 +2782,19 @@ steps:
- vllm/envs.py
- examples/offline_inference/data_parallel.py
- tests/distributed/test_context_parallel.py
- tests/distributed/test_rocm_aiter_custom_ar.py
- tests/distributed/test_rocm_quick_reduce.py
- tests/distributed/test_quick_all_reduce.py
- tests/v1/e2e/general/test_rocm_aiter_custom_ar.py
- tests/v1/distributed/test_dbo.py
- tests/utils.py
commands:
- pytest -v -s tests/distributed/test_context_parallel.py
- pytest -v -s tests/v1/distributed/test_dbo.py
- pytest -v -s tests/distributed/test_rocm_aiter_custom_ar.py
- pytest -v -s tests/v1/e2e/general/test_rocm_aiter_custom_ar.py
- pytest -v -s tests/distributed/test_rocm_quick_reduce.py
- pytest -v -s tests/distributed/test_quick_all_reduce.py
- pytest -v -s tests/v1/distributed/test_dbo.py
#-------------------------------------------------------- mi355 · entrypoints --------------------------------------------------------#
+4 -3
View File
@@ -36,10 +36,10 @@ steps:
source_file_dependencies:
- vllm/
- tests/models/test_terratorch.py
- tests/models/test_transformers.py
- tests/models/transformers/test_backend.py
- tests/models/test_registry.py
commands:
- pytest -v -s models/test_terratorch.py models/test_transformers.py models/test_registry.py
- pytest -v -s models/test_terratorch.py models/transformers/test_backend.py models/test_registry.py
mirror:
amd:
device: mi325_1
@@ -55,6 +55,7 @@ steps:
- vllm/
- tests/models/test_utils.py
- tests/models/test_vision.py
- tests/models/transformers/fusers/
device: cpu-small
commands:
- pytest -v -s models/test_utils.py models/test_vision.py
- pytest -v -s models/test_utils.py models/test_vision.py models/transformers/fusers/
@@ -17,7 +17,7 @@ steps:
- TARGET_TEST_SUITE=L4 pytest basic_correctness/ -v -s -m 'distributed(num_gpus=2)'
- CUDA_VISIBLE_DEVICES=0,1 pytest -v -s model_executor/model_loader/test_sharded_state_loader.py -m '(not slow_test)'
# Avoid importing model tests that cause CUDA reinitialization error
- pytest models/test_transformers.py -v -s -m 'distributed(num_gpus=2)'
- pytest models/transformers/test_backend.py -v -s -m 'distributed(num_gpus=2)'
- pytest models/language -v -s -m 'distributed(num_gpus=2)'
- pytest models/multimodal/generation/test_phi4siglip.py -v -s -m 'distributed(num_gpus=2)'
- pytest models/multimodal -v -s -m 'distributed(num_gpus=2)' --ignore models/multimodal/generation/test_whisper.py --ignore models/multimodal/generation/test_phi4siglip.py
+2 -1
View File
@@ -27,6 +27,7 @@ steps:
- tests/models/multimodal
commands:
- pytest -v -s models/multimodal/generation/test_common.py -m core_model -k "qwen3 or gemma"
- pytest -v -s models/multimodal/generation/test_mm_prefix_lm.py -m core_model
- pytest -v -s models/multimodal/generation/test_qwen2_5_vl.py -m core_model
mirror:
amd:
@@ -58,7 +59,7 @@ steps:
- vllm/
- tests/models/multimodal
commands:
- pytest -v -s models/multimodal -m core_model --ignore models/multimodal/generation/test_common.py --ignore models/multimodal/generation/test_ultravox.py --ignore models/multimodal/generation/test_qwen2_5_vl.py --ignore models/multimodal/generation/test_qwen2_vl.py --ignore models/multimodal/generation/test_whisper.py --ignore models/multimodal/generation/test_memory_leak.py --ignore models/multimodal/generation/test_vit_cudagraph.py --ignore models/multimodal/processing
- pytest -v -s models/multimodal -m core_model --ignore models/multimodal/generation/test_common.py --ignore models/multimodal/generation/test_ultravox.py --ignore models/multimodal/generation/test_qwen2_5_vl.py --ignore models/multimodal/generation/test_qwen2_vl.py --ignore models/multimodal/generation/test_whisper.py --ignore models/multimodal/generation/test_mm_prefix_lm.py --ignore models/multimodal/generation/test_memory_leak.py --ignore models/multimodal/generation/test_vit_cudagraph.py --ignore models/multimodal/processing
- pytest -v -s models/multimodal/generation/test_vit_cudagraph.py -m core_model
- pytest models/multimodal/generation/test_memory_leak.py -m core_model
- cd .. && VLLM_WORKER_MULTIPROC_METHOD=spawn pytest -v -s tests/models/multimodal/generation/test_whisper.py -m core_model # Otherwise, mp_method="spawn" doesn't work
+1 -1
View File
@@ -119,7 +119,7 @@
# Transformers modeling backend
/vllm/model_executor/models/transformers @hmellor
/tests/models/test_transformers.py @hmellor
/tests/models/transformers @hmellor
# Docs
/docs/mkdocs @hmellor
-1
View File
@@ -400,7 +400,6 @@ if(VLLM_GPU_LANG STREQUAL "CUDA" OR VLLM_GPU_LANG STREQUAL "HIP")
"csrc/libtorch_stable/topk.cu"
"csrc/libtorch_stable/mamba/selective_scan_fwd.cu"
"csrc/libtorch_stable/cache_kernels.cu"
"csrc/libtorch_stable/cache_kernels.cu"
"csrc/libtorch_stable/cache_kernels_fused.cu"
"csrc/libtorch_stable/custom_all_reduce.cu"
"csrc/libtorch_stable/fused_deepseek_v4_qnorm_rope_kv_insert_kernel.cu")
+4 -2
View File
@@ -132,8 +132,10 @@ def benchmark_function(
reset_memory_stats()
# Benchmark
start_events = [torch.Event(enable_timing=True) for _ in range(benchmark_iters)]
end_events = [torch.Event(enable_timing=True) for _ in range(benchmark_iters)]
start_events = [
torch.cuda.Event(enable_timing=True) for _ in range(benchmark_iters)
]
end_events = [torch.cuda.Event(enable_timing=True) for _ in range(benchmark_iters)]
for i in range(benchmark_iters):
logits_copy = logits.clone()
+2 -2
View File
@@ -134,8 +134,8 @@ def benchmark_config(
torch.accelerator.synchronize()
# Benchmark
start = torch.Event(enable_timing=True)
end = torch.Event(enable_timing=True)
start = torch.cuda.Event(enable_timing=True)
end = torch.cuda.Event(enable_timing=True)
start.record()
for _ in range(num_iters):
with override_config(config):
@@ -170,8 +170,8 @@ def benchmark_config(
graph.replay()
torch.accelerator.synchronize()
start = torch.Event(enable_timing=True)
end = torch.Event(enable_timing=True)
start = torch.cuda.Event(enable_timing=True)
end = torch.cuda.Event(enable_timing=True)
latencies: list[float] = []
for _ in range(num_iters):
start.record()
+23 -5
View File
@@ -15,6 +15,7 @@ endif()
#
set(ENABLE_X86_ISA $ENV{VLLM_CPU_X86})
set(ENABLE_ARM_BF16 $ENV{VLLM_CPU_ARM_BF16})
set(ENABLE_RVV_BF16 $ENV{VLLM_CPU_RVV_BF16})
include_directories("${CMAKE_SOURCE_DIR}/csrc")
@@ -110,6 +111,13 @@ else()
set(ARM_BF16_FOUND ON)
message(STATUS "ARM BF16 support enabled via VLLM_CPU_ARM_BF16 environment variable")
endif()
# Some kernels (e.g. Bianbu on Spacemit X100) do not report zvfbfmin
# in /proc/cpuinfo despite hardware support. VLLM_CPU_RVV_BF16=1
# overrides the detection result.
if (ENABLE_RVV_BF16)
set(RVV_BF16_FOUND ON)
message(STATUS "RVV BF16 support enabled via VLLM_CPU_RVV_BF16 environment variable")
endif()
endif()
if (CMAKE_SYSTEM_PROCESSOR MATCHES "x86_64|amd64" OR ENABLE_X86_ISA)
@@ -178,7 +186,10 @@ elseif (CMAKE_SYSTEM_PROCESSOR MATCHES "riscv64")
# Override with -DVLLM_RVV_VLEN=128 or -DVLLM_RVV_VLEN=256 for RVV.
if(NOT DEFINED VLLM_RVV_VLEN)
# Auto-detect: find the largest zvl<N>b in /proc/cpuinfo isa line.
if(EXISTS /proc/cpuinfo)
# Skip when cross-compiling — /proc/cpuinfo describes the build host.
if(CMAKE_CROSSCOMPILING)
message(STATUS "Cross-compiling: skipping VLEN auto-detection from /proc/cpuinfo")
elseif(EXISTS /proc/cpuinfo)
file(READ /proc/cpuinfo _cpuinfo)
set(_best 0)
foreach(_n IN ITEMS 128 256 512 1024)
@@ -186,6 +197,13 @@ elseif (CMAKE_SYSTEM_PROCESSOR MATCHES "riscv64")
set(_best ${_n})
endif()
endforeach()
# Only VLEN=128 and VLEN=256 are supported by the RVV kernels.
if(_best GREATER 256)
message(WARNING
"Detected VLEN=${_best} but only 128/256 are supported; "
"clamping to 256")
set(_best 256)
endif()
if(_best GREATER 0)
set(VLLM_RVV_VLEN ${_best})
endif()
@@ -195,9 +213,9 @@ elseif (CMAKE_SYSTEM_PROCESSOR MATCHES "riscv64")
if(NOT DEFINED VLLM_RVV_VLEN AND (RVV_FP16_FOUND OR RVV_BF16_FOUND))
message(FATAL_ERROR
"RISC-V RVV is available but VLEN could not be auto-detected. "
"Please specify VLEN explicitly:\n"
" -DVLLM_RVV_VLEN=128 (for VLEN=128 hardware)\n"
" -DVLLM_RVV_VLEN=256 (for VLEN=256 hardware, e.g. Spacemit X100)")
"Please specify VLEN explicitly via CMAKE_ARGS:\n"
" CMAKE_ARGS='-DVLLM_RVV_VLEN=128' (for VLEN=128 hardware)\n"
" CMAKE_ARGS='-DVLLM_RVV_VLEN=256' (for VLEN=256 hardware, e.g. Spacemit X100)")
endif()
endif()
if(VLLM_RVV_VLEN AND VLLM_RVV_VLEN GREATER 0)
@@ -209,7 +227,7 @@ elseif (CMAKE_SYSTEM_PROCESSOR MATCHES "riscv64")
message(STATUS "BF16 extension detected")
set(MARCH_FLAGS -march=rv64gcv_zvfh_zfbfmin_zvfbfmin_zvl${VLLM_RVV_VLEN}b -mrvv-vector-bits=zvl -mabi=lp64d)
elseif(RVV_FP16_FOUND)
message(WARNING "BF16 functionality is not available")
message(WARNING "BF16 functionality is not available.")
set(MARCH_FLAGS -march=rv64gcv_zvfh_zvl${VLLM_RVV_VLEN}b -mrvv-vector-bits=zvl -mabi=lp64d)
else()
message(STATUS "compile riscv with scalar (no FP16/BF16)")
+2 -1
View File
@@ -13,7 +13,8 @@ static inline cpu_attention::Fp8KVCacheDataType parse_fp8_kv_dtype(
bool cpu_attn_has_isa(const std::string& isa) {
if (isa == "rvv") {
#if defined(__riscv) && defined(__riscv_v_min_vlen) && __riscv_v_min_vlen == 128
#if defined(__riscv) && defined(__riscv_v_min_vlen) && \
(__riscv_v_min_vlen == 128 || __riscv_v_min_vlen == 256)
return true;
#else
return false;
+35 -16
View File
@@ -214,11 +214,18 @@ struct BF16Vec32 : public Vec<BF16Vec32> {
explicit BF16Vec32(const BF16Vec8& v) {
fixed_u16x8_t u16_val = bf16_to_u16(v.reg);
fixed_u16x32_t u16_combined =
RVVI4(__riscv_vcreate_v_u16, LMUL_128, _u16, LMUL_512)(
u16_val, u16_val, u16_val, u16_val);
reg = RVVI4(__riscv_vreinterpret_v_u16, LMUL_512, _bf16,
LMUL_512)(u16_combined);
// Widen LMUL_128 → LMUL_256 so vslideup operands share a type.
// At VLEN=256 this is mf2→m1 (both integer); at VLEN=128 it is m1→m2.
fixed_u16x16_t ext =
RVVI4(__riscv_vlmul_ext_v_u16, LMUL_128, _u16, LMUL_256)(u16_val);
// Build 16-element half: place the 8 elements at offsets 0 and 8.
fixed_u16x16_t half = RVVI(__riscv_vmv_v_x_u16, LMUL_256)(0, 16);
half = RVVI(__riscv_vslideup_vx_u16, LMUL_256)(half, ext, 0, 8);
half = RVVI(__riscv_vslideup_vx_u16, LMUL_256)(half, ext, 8, 16);
// Double to LMUL_512 (m1→m2 at VLEN=256, m2→m4 at VLEN=128).
fixed_u16x32_t dst =
RVVI4(__riscv_vcreate_v_u16, LMUL_256, _u16, LMUL_512)(half, half);
reg = RVVI4(__riscv_vreinterpret_v_u16, LMUL_512, _bf16, LMUL_512)(dst);
};
void save(void* ptr) const {
@@ -623,17 +630,29 @@ struct FP32Vec16 : public Vec<FP32Vec16> {
data.reg, data.reg)) {};
explicit FP32Vec16(const FP32Vec16& data) : reg(data.reg) {};
explicit FP32Vec16(int64_t value, const FP32Vec16& lut) {
const uint64_t q_values = static_cast<uint64_t>(value);
auto packed = RVVI(__riscv_vmv_v_x_u64, LMUL_1024)(q_values, VEC_ELEM_NUM);
auto lane_ids = RVVI(__riscv_vid_v_u64, LMUL_1024)(VEC_ELEM_NUM);
auto shifts =
RVVI(__riscv_vsll_vx_u64, LMUL_1024)(lane_ids, 2, VEC_ELEM_NUM);
auto shifted =
RVVI(__riscv_vsrl_vv_u64, LMUL_1024)(packed, shifts, VEC_ELEM_NUM);
auto idx64 =
RVVI(__riscv_vand_vx_u64, LMUL_1024)(shifted, 0xF, VEC_ELEM_NUM);
auto idx32 = RVVI(__riscv_vnsrl_wx_u32, LMUL_512)(idx64, 0, VEC_ELEM_NUM);
reg = RVVI(__riscv_vrgather_vv_f32, LMUL_512)(lut.reg, idx32, VEC_ELEM_NUM);
// Split into two 32-bit halves to avoid u64 @ LMUL_1024 (m8 on
// VLEN=128 / m4 on VLEN=256), which causes heavy register spilling.
constexpr int HALF = VEC_ELEM_NUM / 2;
const auto q = static_cast<uint64_t>(value);
const uint32_t lo = static_cast<uint32_t>(q);
const uint32_t hi = static_cast<uint32_t>(q >> 32);
auto lane_ids = RVVI(__riscv_vid_v_u32, LMUL_256)(HALF);
auto shifts = RVVI(__riscv_vsll_vx_u32, LMUL_256)(lane_ids, 2, HALF);
auto packed_lo = RVVI(__riscv_vmv_v_x_u32, LMUL_256)(lo, HALF);
auto idx_lo = RVVI(__riscv_vand_vx_u32, LMUL_256)(
RVVI(__riscv_vsrl_vv_u32, LMUL_256)(packed_lo, shifts, HALF), 0xF,
HALF);
auto packed_hi = RVVI(__riscv_vmv_v_x_u32, LMUL_256)(hi, HALF);
auto idx_hi = RVVI(__riscv_vand_vx_u32, LMUL_256)(
RVVI(__riscv_vsrl_vv_u32, LMUL_256)(packed_hi, shifts, HALF), 0xF,
HALF);
auto idx =
RVVI4(__riscv_vcreate_v_u32, LMUL_256, _u32, LMUL_512)(idx_lo, idx_hi);
reg = RVVI(__riscv_vrgather_vv_f32, LMUL_512)(lut.reg, idx, VEC_ELEM_NUM);
}
explicit FP32Vec16(const FP16Vec16& v);
+2 -1
View File
@@ -278,7 +278,8 @@ TORCH_LIBRARY_EXPAND(TORCH_EXTENSION_NAME, ops) {
ops.def(
"dynamic_4bit_int_moe("
"Tensor x, Tensor topk_ids, Tensor topk_weights,"
"Tensor w13_packed, Tensor w2_packed, int H, int I, int I2,"
"Tensor w13_packed, Tensor w2_packed,"
"int hidden_size, int intermediate_size,"
"int group_size, bool apply_router_weight_on_input, int activation_kind"
") -> Tensor");
+53 -28
View File
@@ -29,25 +29,37 @@ enum ActivationKind : int64_t {
torch::Tensor dynamic_4bit_int_moe_cpu(
torch::Tensor x, torch::Tensor topk_ids, torch::Tensor topk_weights,
torch::Tensor w13_packed, torch::Tensor w2_packed, int64_t H, int64_t I,
int64_t I2, int64_t group_size, bool apply_router_weight_on_input,
int64_t activation_kind) {
torch::Tensor w13_packed, torch::Tensor w2_packed, int64_t hidden_size,
int64_t intermediate_size, int64_t group_size,
bool apply_router_weight_on_input, int64_t activation_kind) {
TORCH_CHECK(x.dim() == 2, "x must be 2D");
TORCH_CHECK(topk_ids.dim() == 2 && topk_weights.dim() == 2,
"topk tensors must be [T, K]");
TORCH_CHECK(
w13_packed.size(0) == w2_packed.size(0),
"w13_packed and w2_packed must have same number of experts in dim 0");
TORCH_CHECK(I2 == 2 * I, "I2 must equal 2*I");
const int64_t T = x.size(0);
const int64_t K = topk_ids.size(1);
const int64_t E = w13_packed.size(0);
const int64_t N = T * K;
const int64_t w13_out_features = 2 * intermediate_size;
auto x_c = x.contiguous();
// _dyn_quant_matmul_4bit kernel natively supports these pre-quant activation
// dtypes:
// - fp32: with channelwise and groupwise
// - bf16: with channelwise -> upcast to fp32 for groupwise
// - fp16: not supported -> upcast to fp32 for groupwise & channelwise
const auto output_dtype = x_c.scalar_type();
const bool should_cast_input =
((group_size != -1) && output_dtype == at::kBFloat16) ||
output_dtype == at::kHalf;
if (should_cast_input) {
x_c = x_c.to(at::kFloat);
}
auto ids_c = topk_ids.contiguous();
auto gates_c = topk_weights.to(at::kFloat).contiguous();
auto gates_c = topk_weights.to(x_c.scalar_type()).contiguous();
// bucketing tokens -> experts
c10::SmallVector<int64_t, 64> counts(
@@ -63,35 +75,42 @@ torch::Tensor dynamic_4bit_int_moe_cpu(
c10::SmallVector<int64_t, 65> offsets(E + 1, 0); // ( E +1 )
for (int64_t e = 0; e < E; ++e) offsets[e + 1] = offsets[e] + counts[e];
// expert_tokens = [tokens indices for expert 0, ...]
// expert_gates = [router weights for tokens assigned to expert 0, ...]
auto expert_tokens = at::empty({offsets[E]}, ids_c.options());
auto expert_gates = at::empty({offsets[E]}, gates_c.options());
{
c10::SmallVector<int64_t, 64> cursor(E, 0);
const auto* ids_ptr = ids_c.data_ptr<int64_t>();
const auto* gts_ptr = gates_c.data_ptr<float>();
auto* tok_ptr = expert_tokens.data_ptr<int64_t>();
auto* gate_ptr = expert_gates.data_ptr<float>();
AT_DISPATCH_FLOATING_TYPES_AND2(
at::ScalarType::BFloat16, at::ScalarType::Half, gates_c.scalar_type(),
"bucket_expert_tokens_and_gates", [&] {
const auto* ids_ptr = ids_c.data_ptr<int64_t>();
const auto* gts_ptr = gates_c.data_ptr<scalar_t>();
auto* tok_ptr = expert_tokens.data_ptr<int64_t>();
auto* gate_ptr = expert_gates.data_ptr<scalar_t>();
for (int64_t t = 0; t < T; ++t) {
const int64_t base = t * K;
for (int64_t k = 0; k < K; ++k) {
const int64_t idx = base + k;
const int64_t e = ids_ptr[idx];
const int64_t p = offsets[e] + (cursor[e]++);
tok_ptr[p] = t;
gate_ptr[p] = gts_ptr[idx];
}
}
for (int64_t t = 0; t < T; ++t) {
const int64_t base = t * K;
for (int64_t k = 0; k < K; ++k) {
const int64_t idx = base + k;
const int64_t e = ids_ptr[idx];
const int64_t p = offsets[e] + (cursor[e]++);
tok_ptr[p] = t;
gate_ptr[p] = gts_ptr[idx];
}
}
});
}
const int64_t g_eff_13 = (group_size != -1) ? group_size : H;
const int64_t g_eff_2 = (group_size != -1) ? group_size : I;
const int64_t g_eff_13 = (group_size != -1) ? group_size : hidden_size;
const int64_t g_eff_2 = (group_size != -1) ? group_size : intermediate_size;
// X_all [num_tokens * K, hidden_size]
auto X_all = x_c.index_select(/*dim=*/0, expert_tokens);
if (apply_router_weight_on_input) {
X_all = X_all.mul(expert_gates.unsqueeze(1));
}
auto Y_all = at::empty({offsets[E], H}, x_c.options());
auto Y_all = at::empty({offsets[E], hidden_size}, x_c.options());
at::parallel_for(0, offsets[E], 0, [&](int64_t idx_begin, int64_t idx_end) {
c10::InferenceMode guard;
@@ -109,11 +128,13 @@ torch::Tensor dynamic_4bit_int_moe_cpu(
auto w2_e = w2_packed.select(/*dim=*/0, e);
// W13
auto y13 =
mm(x_e, w13_e, g_eff_13, /*in_features=*/H, /*out_features=*/I2);
auto y13 = mm(x_e, w13_e, g_eff_13, /*in_features=*/hidden_size,
/*out_features=*/w13_out_features);
auto g_part = y13.narrow(/*dim=*/1, /*start=*/0, /*length=*/I);
auto u_part = y13.narrow(/*dim=*/1, /*start=*/I, /*length=*/I);
auto g_part =
y13.narrow(/*dim=*/1, /*start=*/0, /*length=*/intermediate_size);
auto u_part = y13.narrow(/*dim=*/1, /*start=*/intermediate_size,
/*length=*/intermediate_size);
torch::Tensor act;
if (activation_kind == ActivationKind::SwiGLUOAI) { // SwiGLUOAI
@@ -128,7 +149,8 @@ torch::Tensor dynamic_4bit_int_moe_cpu(
}
// W2
auto y = mm(act, w2_e, g_eff_2, /*in_features=*/I, /*out_features=*/H);
auto y = mm(act, w2_e, g_eff_2, /*in_features=*/intermediate_size,
/*out_features=*/hidden_size);
// Store per-expert result
Y_all.narrow(/*dim=*/0, /*start=*/start, /*length=*/te).copy_(y);
@@ -138,8 +160,11 @@ torch::Tensor dynamic_4bit_int_moe_cpu(
if (!apply_router_weight_on_input) {
Y_all = Y_all.mul(expert_gates.unsqueeze(1));
}
if (Y_all.scalar_type() != output_dtype) {
Y_all = Y_all.to(output_dtype);
}
auto out = at::zeros({T, H}, x.options());
auto out = at::zeros({T, hidden_size}, x.options());
out =
at::index_add(out, /*dim=*/0, /*index=*/expert_tokens, /*source=*/Y_all);
+3 -3
View File
@@ -53,9 +53,9 @@ void dynamic_scaled_int8_quant(torch::Tensor& out, torch::Tensor const& input,
torch::Tensor dynamic_4bit_int_moe_cpu(
torch::Tensor x, torch::Tensor topk_ids, torch::Tensor topk_weights,
torch::Tensor w13_packed, torch::Tensor w2_packed, int64_t H, int64_t I,
int64_t I2, int64_t group_size, bool apply_router_weight_on_input,
int64_t activation_kind);
torch::Tensor w13_packed, torch::Tensor w2_packed, int64_t hidden_size,
int64_t intermediate_size, int64_t group_size,
bool apply_router_weight_on_input, int64_t activation_kind);
using fptr_t = int64_t;
#ifdef USE_ROCM
+4 -7
View File
@@ -25,7 +25,6 @@ FROM ubuntu:22.04 AS base-common
WORKDIR /workspace
ARG PYTHON_VERSION=3.12
ARG PIP_EXTRA_INDEX_URL="https://download.pytorch.org/whl/cpu"
ARG max_jobs=32
ENV MAX_JOBS=${max_jobs}
@@ -53,8 +52,6 @@ ENV PATH="$VIRTUAL_ENV/bin:$PATH"
ENV UV_HTTP_TIMEOUT=500
# Install Python dependencies
ENV PIP_EXTRA_INDEX_URL=${PIP_EXTRA_INDEX_URL}
ENV UV_EXTRA_INDEX_URL=${PIP_EXTRA_INDEX_URL}
ENV UV_INDEX_STRATEGY="unsafe-best-match"
ENV UV_LINK_MODE="copy"
@@ -64,7 +61,7 @@ COPY requirements/cpu.txt requirements/cpu.txt
RUN --mount=type=cache,target=/root/.cache/uv \
uv pip install --upgrade pip && \
uv pip install -r requirements/cpu.txt
uv pip install -r requirements/cpu.txt --torch-backend cpu
ARG TARGETARCH
ENV TARGETARCH=${TARGETARCH}
@@ -149,7 +146,7 @@ RUN if [ "$TARGETARCH" = "arm64" ] && [ "$VLLM_CPU_X86" != "0" ]; then \
COPY requirements/build/cpu.txt requirements/build/cpu.txt
RUN --mount=type=cache,target=/root/.cache/uv \
uv pip install -r requirements/build/cpu.txt
uv pip install -r requirements/build/cpu.txt --torch-backend cpu
COPY . .
@@ -205,7 +202,7 @@ RUN case "$(uname -m)" in \
esac
RUN --mount=type=cache,target=/root/.cache/uv \
uv pip install -r requirements/test/cpu.txt
uv pip install -r requirements/test/cpu.txt --torch-backend cpu
######################### DEV IMAGE #########################
FROM vllm-build AS vllm-dev
@@ -231,7 +228,7 @@ COPY --from=vllm-test-deps /vllm-workspace/requirements/test/cpu.txt requirement
RUN --mount=type=cache,target=/root/.cache/uv \
uv pip install -r requirements/lint.txt && \
uv pip install -r requirements/test/cpu.txt && \
uv pip install -r requirements/test/cpu.txt --torch-backend cpu && \
pre-commit install --hook-type pre-commit --hook-type commit-msg
ENTRYPOINT ["bash"]
+2 -1
View File
@@ -249,7 +249,7 @@ RUN --mount=type=cache,target=/root/.cache/uv \
NUMBA_WHL_FILE=$(ls /tmp/numba-wheels/*.whl) && \
OPENCV_WHL_FILE=$(ls /tmp/opencv-wheels/*.whl) && \
GUIDANCE_WHL_FILE=$(ls /tmp/guidance-wheels/*.whl) && \
uv pip install -v \
uv pip install -v \
$ARROW_WHL_FILE \
$VISION_WHL_FILE \
$HF_XET_WHL_FILE \
@@ -257,6 +257,7 @@ RUN --mount=type=cache,target=/root/.cache/uv \
$NUMBA_WHL_FILE \
$OPENCV_WHL_FILE \
$GUIDANCE_WHL_FILE \
--torch-backend cpu \
--index-strategy unsafe-best-match \
-r requirements/build/cpu.txt \
-r requirements/cpu.txt
+3 -2
View File
@@ -276,8 +276,9 @@ By default vLLM uses the standard Hugging Face `tokenizers` library to power
the fast tokenizer. For BPE tokenizers (Qwen, Llama, DeepSeek, GPT-OSS, etc.)
you can switch to the [fastokens](https://github.com/crusoecloud/fastokens)
Rust backend, a drop-in replacement that's substantially faster on
encode/decode and on streaming detokenization. Enable it by setting
`VLLM_USE_FASTOKENS=1`:
encode/decode and on streaming detokenization. `VLLM_USE_FASTOKENS` is
available in vLLM v0.23.0 and later. If your installed vLLM version does not
recognize the environment variable, upgrade vLLM before enabling the override:
```console
VLLM_USE_FASTOKENS=1 vllm serve Qwen/Qwen3-8B
+1 -1
View File
@@ -168,7 +168,7 @@ Priority is **1 = highest** (tried first).
| `FLASH_ATTN` | FA4* | fp16, bf16 | `auto`, `float16`, `bfloat16` | %16 | Any | ✅ | ✅ | ❌ | ✅ | All | ≥10.0 |
| `FLASH_ATTN_DIFFKV` | | fp16, bf16 | `auto` | Any | Any | ❌ | ❌ | ❌ | ✅ | Decoder | Any |
| `FLEX_ATTENTION` | | fp16, bf16, fp32 | `auto`, `float16`, `bfloat16` | %16 | Any | ❌ | ✅ | ✅ | ❌ | Decoder, Encoder Only | Any |
| `HPC_ATTN` | | fp16, bf16 | `auto`, `fp8_e4m3` | 64 | 128 | ❌ | ❌ | ❌ | ❌ | Decoder | ≥9.0 |
| `HPC_ATTN` | | fp16, bf16 | `auto`, `bfloat16`, `fp8_e4m3` | 64 | 128 | ❌ | ❌ | ❌ | ❌ | Decoder | ≥9.0 |
| `ROCM_AITER_FA` | | fp16, bf16 | `auto`, `float16`, `bfloat16`, `fp8`, `fp8_e4m3`, `fp8_e5m2` | 16, 32 | 64, 128, 256 | ✅ | ✅ | ❌ | ❌ | Decoder | N/A |
| `ROCM_AITER_UNIFIED_ATTN` | | bf16 | `auto`, `bfloat16`, `fp8`, `fp8_e4m3` | %16 | Any | ✅ | ❌ | ✅ | ❌ | All | N/A |
| `ROCM_ATTN` | | fp16, bf16, fp32 | `auto`, `float16`, `bfloat16`, `fp8`, `fp8_e4m3`, `fp8_e5m2` | %16 | 32, 64, 80, 96, 128, 160, 192, 224, 256 | ❌ | ✅ | ✅ | ❌ | Decoder, Encoder, Encoder Only | N/A |
+1 -4
View File
@@ -75,14 +75,11 @@ vllm serve Intel/DeepSeek-R1-0528-Qwen3-8B-int4-AutoRound \
--max-model-len 4096
```
!!! note
To deploy `wNa16` models on Intel GPU/CPU, please add `--enforce-eager` for now.
## Evaluating the Quantized Model with vLLM
```bash
lm_eval --model vllm \
--model_args pretrained="Intel/DeepSeek-R1-0528-Qwen3-8B-int4-AutoRound,max_model_len=8192,max_num_batched_tokens=32768,max_num_seqs=128,gpu_memory_utilization=0.8,dtype=bfloat16,max_gen_toks=2048,enforce_eager=True" \
--model_args pretrained="Intel/DeepSeek-R1-0528-Qwen3-8B-int4-AutoRound,max_model_len=8192,max_num_batched_tokens=32768,max_num_seqs=128,gpu_memory_utilization=0.8,dtype=bfloat16,max_gen_toks=2048" \
--tasks gsm8k \
--num_fewshot 5 \
--batch_size 128
@@ -35,15 +35,10 @@ After installation of XCode and the Command Line Tools, which include Apple Clan
```bash
git clone https://github.com/vllm-project/vllm.git
cd vllm
uv pip install -r requirements/cpu.txt --index-strategy unsafe-best-match
uv pip install -r requirements/cpu.txt
uv pip install -e .
```
!!! tip
The `--index-strategy unsafe-best-match` flag is needed to resolve dependencies across multiple package indexes (PyTorch CPU index and PyPI). Without this flag, you may encounter `typing-extensions` version conflicts.
The term "unsafe" refers to the package resolution strategy, not security. By default, `uv` only searches the first index where a package is found to prevent dependency confusion attacks. This flag allows `uv` to search all configured indexes to find the best compatible versions. Since both PyTorch and PyPI are trusted package sources, using this strategy is safe and appropriate for vLLM installation.
!!! note
On macOS the `VLLM_TARGET_DEVICE` is automatically set to `cpu`, which is currently the only supported device.
@@ -48,10 +48,10 @@ Execute the following commands to build and install vLLM from source.
```bash
uv pip install -v \
--extra-index-url https://download.pytorch.org/whl/cpu \
--torch-backend auto \
-r requirements/build/cpu.txt \
-r requirements/cpu.txt \
--torch-backend cpu \
--index-strategy unsafe-best-match && \
VLLM_TARGET_DEVICE=cpu python setup.py bdist_wheel && \
uv pip install dist/*.whl
```
+2 -3
View File
@@ -2,8 +2,7 @@
!!! note
We currently support pooling models primarily for convenience. This is not guaranteed to provide any performance
improvements over using Hugging Face Transformers or Sentence Transformers directly.
improvements over using Hugging Face Transformers or Sentence Transformers directly.
We plan to optimize pooling models in vLLM. Please comment on <https://github.com/vllm-project/vllm/issues/21796> if you have any suggestions!
## What are pooling models?
@@ -63,7 +62,7 @@ please refer to [IO Processor Plugins](../../design/io_processor_plugins.md).
!!! note
Within classification tasks, there is a specialized subcategory: Cross-encoder (aka reranker) models. These models
are a subset of classification models that accept two prompts as input and output num_labels equal to 1.
are a subset of classification models that accept two prompts as input and output num_labels equal to 1.
### Pooling Types
+4 -2
View File
@@ -15,7 +15,7 @@ These models are what we list in [supported text models](#list-of-text-only-lang
### Transformers
vLLM also supports model implementations that are available in Transformers. You should expect the performance of a Transformers model implementation used in vLLM to be within <5% of the performance of a dedicated vLLM model implementation. We call this feature the "Transformers modeling backend".
vLLM also supports model implementations that are available in Transformers. We call this feature the "Transformers modeling backend". The performance of models loaded with the Transformers modeling backend should be identical to a dedicated vLLM model implementation.
Currently, the Transformers modeling backend works for the following:
@@ -140,7 +140,7 @@ Here is what happens in the background when this model is loaded:
That's it!
For your model to be compatible with vLLM's tensor parallel and/or pipeline parallel features, you must add `base_model_tp_plan` and/or `base_model_pp_plan` to your model's config class:
For your model to be compatible with vLLM's tensor parallel and/or pipeline parallel features, you may need to add `base_model_tp_plan` and/or `base_model_pp_plan` to your model's config class:
<details class="code">
<summary>configuration_my_model.py</summary>
@@ -168,9 +168,11 @@ class MyConfig(PretrainedConfig):
</details>
- `base_model_tp_plan` is a `dict` that maps fully qualified layer name patterns to tensor parallel styles (currently only `"colwise"` and `"rowwise"` are supported).
- vLLM infers the tensor parallel style of standard attention (`q`/`k`/`v`/`o_proj`) and gated-MLP/experts (`gate`/`up`/`down_proj`) projections if it can fuse them, so these may not need to be listed. `base_model_tp_plan` is only _required_ for layers that do not follow these patterns; any linear that is neither fused nor named in the plan is replicated.
- `base_model_pp_plan` is a `dict` that maps direct child layer names to `tuple`s of `list`s of `str`s:
- You only need to do this for layers which are not present on all pipeline stages
- vLLM assumes that there will be only one `nn.ModuleList`, which is distributed across the pipeline stages
- When no `base_model_pp_plan` is provided, the Transformers modelling backend infers the split from the text model's sole `nn.ModuleList`, keeping the parameter-bearing modules around it (input embeddings, final norm) on the first/last stage (depending on declaration order) and parameter-free modules (e.g. rotary embeddings) on every stage
- The `list` in the first element of the `tuple` contains the names of the input arguments
- The `list` in the last element of the `tuple` contains the names of the variables the layer outputs to in your modeling code
@@ -0,0 +1,73 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
# ruff: noqa: E501
"""
Example online usage of the Jina Reranker v3 score and rerank APIs with a task
instruction.
Run `vllm serve jinaai/jina-reranker-v3 --runner pooling` to start up the
server in vLLM.
"""
import argparse
import json
import requests
def post_http_request(prompt: dict, api_url: str) -> requests.Response:
headers = {"User-Agent": "Test Client"}
response = requests.post(api_url, headers=headers, json=prompt)
return response
def print_response(name: str, prompt: dict, response: requests.Response) -> None:
print(f"\n{name} request:")
print(json.dumps(prompt, indent=2))
print(f"\n{name} response:")
print(json.dumps(response.json(), indent=2))
def parse_args():
parser = argparse.ArgumentParser()
parser.add_argument("--host", type=str, default="localhost")
parser.add_argument("--port", type=int, default=8000)
parser.add_argument("--model", type=str, default="jinaai/jina-reranker-v3")
return parser.parse_args()
def main(args):
score_url = f"http://{args.host}:{args.port}/score"
rerank_url = f"http://{args.host}:{args.port}/rerank"
model_name = args.model
query = "Which passage is about sports?"
documents = [
"Basketball is played by two teams on a court.",
"Green tea contains antioxidants and may support metabolism.",
]
instruction = "Rank passages about sports higher than passages about nutrition."
score_prompt = {
"model": model_name,
"queries": query,
"documents": documents,
"instruction": instruction,
}
score_response = post_http_request(prompt=score_prompt, api_url=score_url)
print_response("Score", score_prompt, score_response)
rerank_prompt = {
"model": model_name,
"query": query,
"documents": documents,
"instruction": instruction,
}
rerank_response = post_http_request(prompt=rerank_prompt, api_url=rerank_url)
print_response("Rerank", rerank_prompt, rerank_response)
if __name__ == "__main__":
args = parse_args()
main(args)
-1
View File
@@ -1,4 +1,3 @@
--extra-index-url https://download.pytorch.org/whl/cpu
cmake>=3.26.1
ninja
packaging>=24.2
-1
View File
@@ -1,4 +1,3 @@
--extra-index-url https://download.pytorch.org/whl/cpu
# Common dependencies
-r common.txt
+1 -1
View File
@@ -25,7 +25,7 @@ nvidia-cutlass-dsl[cu13]==4.5.2
quack-kernels>=0.3.3
# Tokenspeed_MLA for faster mla with spec decode
tokenspeed-mla==0.1.2
tokenspeed-mla==0.1.2; platform_system == "Linux"
# Humming kernels for quantization gemm
humming-kernels[cu13]==0.1.6
+1 -1
View File
@@ -16,5 +16,5 @@ torch==2.12.0
torchaudio
torchvision
auto_round_lib>=0.13.3
auto_round_lib>=0.14.0
vllm_xpu_kernels @ https://github.com/vllm-project/vllm-xpu-kernels/releases/download/v0.1.10.1/vllm_xpu_kernels-0.1.10.1-cp38-abi3-manylinux_2_28_x86_64.whl
@@ -512,6 +512,7 @@ impl EngineCoreClient {
Ok(EngineCoreOutputStream::new(
request_id,
engine_id.engine_index().unwrap_or(0),
self.abort_tx.clone(),
rx,
))
@@ -15,7 +15,7 @@ use crate::client::state::{OutputReceiver, RequestRegistry, UtilityReceiver, Uti
use crate::client::stream::EngineCoreStreamOutput;
use crate::client::{AbortCause, AbortRequest};
use crate::error::{client_closed, dispatcher_closed, unexpected_dispatcher_output};
use crate::metrics::{LoraInfoExporter, record_scheduler_stats};
use crate::metrics::{LoraInfoExporter, SchedulerStatsRecorder};
use crate::protocol::encode_msgpack;
use crate::protocol::output::{EngineCoreOutput, EngineCoreOutputs};
use crate::protocol::request::EngineCoreRequestType;
@@ -29,6 +29,7 @@ pub(crate) struct ClientInner {
/// The runtime handle used for sending messages to the engine.
handle: Handle,
model_name: String,
scheduler_stats_recorder: SchedulerStatsRecorder,
request_reg: Mutex<RequestRegistry>,
utility_reg: Mutex<UtilityRegistry>,
health_error: ArcSwapOption<Error>,
@@ -43,10 +44,13 @@ impl ClientInner {
model_name: String,
engines: &[ConnectedEngine],
) -> Self {
let scheduler_stats_recorder =
SchedulerStatsRecorder::new(&METRICS.scheduler, &model_name, engines);
Self {
input_send,
handle,
model_name,
scheduler_stats_recorder,
request_reg: Mutex::new(RequestRegistry::new(engines)),
utility_reg: Mutex::new(UtilityRegistry::default()),
health_error: ArcSwapOption::empty(),
@@ -389,12 +393,7 @@ pub(crate) async fn run_output_dispatcher_loop(
"dropping scheduler stats for unknown engine"
);
}
record_scheduler_stats(
&METRICS.scheduler,
inner.model_name(),
batch.engine_index,
scheduler_stats,
);
inner.scheduler_stats_recorder.record(batch.engine_index, scheduler_stats);
}
// The engine's scheduler stats never carry adapter names;
@@ -45,6 +45,7 @@ impl Deref for EngineCoreStreamOutput {
/// `finish_reason` is non-`None`.
pub struct EngineCoreOutputStream {
request_id: String,
engine_index: u32,
abort_tx: mpsc::UnboundedSender<AbortRequest>,
state: State,
rx: OutputReceiver,
@@ -53,11 +54,13 @@ pub struct EngineCoreOutputStream {
impl EngineCoreOutputStream {
pub(crate) fn new(
request_id: String,
engine_index: u32,
abort_tx: mpsc::UnboundedSender<AbortRequest>,
rx: OutputReceiver,
) -> Self {
Self {
request_id,
engine_index,
abort_tx,
state: State::Running,
rx,
@@ -68,6 +71,11 @@ impl EngineCoreOutputStream {
pub fn request_id(&self) -> &str {
&self.request_id
}
/// Return the index of the engine that owns this request.
pub fn engine_index(&self) -> u32 {
self.engine_index
}
}
impl Stream for EngineCoreOutputStream {
+167 -81
View File
@@ -1,90 +1,190 @@
use std::collections::BTreeMap;
use std::collections::BTreeSet;
use std::time::{SystemTime, UNIX_EPOCH};
use vllm_metrics::{
EngineLabels, EnginePositionLabels, LoraAdapterNames, LoraInfoLabels, SchedulerMetrics,
EngineLabels, EnginePositionLabels, F64Gauge, Family, HistogramMetric, LoraAdapterNames,
LoraInfoLabels, SchedulerLogStatsAccumulator, SchedulerMetrics, U64Counter, U64Gauge,
WaitingReasonLabels,
};
use crate::protocol::stats::SchedulerStats;
use crate::transport::ConnectedEngine;
const WAITING_REASON_CAPACITY: &str = "capacity";
const WAITING_REASON_DEFERRED: &str = "deferred";
/// Record the scheduler-stats-backed metrics for one engine at one point in
/// time.
pub(crate) fn record_scheduler_stats(
metrics: &SchedulerMetrics,
model_name: impl Into<String>,
engine: u32,
stats: &SchedulerStats,
) {
let model_name = model_name.into();
let labels = EngineLabels {
model_name: model_name.clone(),
engine,
};
/// Cached scheduler-stats metric handles for all engines connected to one
/// frontend client.
pub(crate) struct SchedulerStatsRecorder {
engines: BTreeMap<u32, SchedulerStatsHandles>,
}
/// Per-engine cached metric handles used while recording `SchedulerStats`.
struct SchedulerStatsHandles {
// Base labels reused for dynamic child labels.
labels: EngineLabels,
// Scheduler state gauges.
metrics.scheduler_running.get_or_create(&labels).set(stats.num_running_reqs);
metrics
.scheduler_waiting
.get_or_create(&labels)
.set(stats.num_waiting_reqs + stats.num_skipped_waiting_reqs);
metrics
.scheduler_waiting_by_reason
.get_or_create(&WaitingReasonLabels {
model_name: model_name.clone(),
engine,
reason: WAITING_REASON_CAPACITY,
})
.set(stats.num_waiting_reqs);
metrics
.scheduler_waiting_by_reason
.get_or_create(&WaitingReasonLabels {
model_name: model_name.clone(),
engine,
reason: WAITING_REASON_DEFERRED,
})
.set(stats.num_skipped_waiting_reqs);
metrics.kv_cache_usage.get_or_create(&labels).set(stats.kv_cache_usage);
scheduler_running: U64Gauge,
scheduler_waiting: U64Gauge,
scheduler_waiting_capacity: U64Gauge,
scheduler_waiting_deferred: U64Gauge,
kv_cache_usage: F64Gauge,
// Prefix-cache counters, including the connector-backed external cache path.
metrics
.prefix_cache_queries
.get_or_create(&labels)
.inc_by(stats.prefix_cache_stats.base.queries);
metrics
.prefix_cache_hits
.get_or_create(&labels)
.inc_by(stats.prefix_cache_stats.base.hits);
prefix_cache_queries: U64Counter,
prefix_cache_hits: U64Counter,
external_prefix_cache_queries: U64Counter,
external_prefix_cache_hits: U64Counter,
// Speculative decoding counters.
spec_decode_num_drafts: U64Counter,
spec_decode_num_draft_tokens: U64Counter,
spec_decode_num_accepted_tokens: U64Counter,
spec_decode_num_accepted_tokens_per_pos: Family<EnginePositionLabels, U64Counter>,
// Per-engine performance / MFU counters.
estimated_flops_per_gpu: U64Counter,
estimated_read_bytes_per_gpu: U64Counter,
estimated_write_bytes_per_gpu: U64Counter,
// Sampled KV-cache residency histograms.
kv_block_lifetime_seconds: HistogramMetric,
kv_block_idle_before_evict_seconds: HistogramMetric,
kv_block_reuse_gap_seconds: HistogramMetric,
// Non-Prometheus interval accumulator for periodic text-log helpers.
log_stats: SchedulerLogStatsAccumulator,
}
impl SchedulerStatsRecorder {
/// Resolve the fixed-label metric handles for the connected engines.
pub(crate) fn new(
metrics: &SchedulerMetrics,
model_name: &str,
engines: &[ConnectedEngine],
) -> Self {
let engines = engines
.iter()
.filter_map(|engine| {
let engine = engine.engine_id.engine_index()?;
Some((
engine,
resolve_scheduler_stats_handles(metrics, model_name, engine),
))
})
.collect();
Self { engines }
}
/// Record one scheduler-stats payload for the given engine index.
pub(crate) fn record(&self, engine_index: u32, stats: &SchedulerStats) {
if let Some(handles) = self.engines.get(&engine_index) {
record_scheduler_stats_with_handles(handles, stats);
}
}
}
/// Resolve all fixed-label scheduler metrics for one engine.
fn resolve_scheduler_stats_handles(
metrics: &SchedulerMetrics,
model_name: &str,
engine: u32,
) -> SchedulerStatsHandles {
let labels = EngineLabels {
model_name: model_name.to_string(),
engine,
};
let capacity = WaitingReasonLabels {
model_name: model_name.to_string(),
engine,
reason: WAITING_REASON_CAPACITY,
};
let deferred = WaitingReasonLabels {
model_name: model_name.to_string(),
engine,
reason: WAITING_REASON_DEFERRED,
};
SchedulerStatsHandles {
scheduler_running: metrics.scheduler_running.get_or_create_owned(&labels),
scheduler_waiting: metrics.scheduler_waiting.get_or_create_owned(&labels),
scheduler_waiting_capacity: metrics
.scheduler_waiting_by_reason
.get_or_create_owned(&capacity),
scheduler_waiting_deferred: metrics
.scheduler_waiting_by_reason
.get_or_create_owned(&deferred),
kv_cache_usage: metrics.kv_cache_usage.get_or_create_owned(&labels),
prefix_cache_queries: metrics.prefix_cache_queries.get_or_create_owned(&labels),
prefix_cache_hits: metrics.prefix_cache_hits.get_or_create_owned(&labels),
external_prefix_cache_queries: metrics
.external_prefix_cache_queries
.get_or_create_owned(&labels),
external_prefix_cache_hits: metrics.external_prefix_cache_hits.get_or_create_owned(&labels),
spec_decode_num_drafts: metrics.spec_decode_num_drafts.get_or_create_owned(&labels),
spec_decode_num_draft_tokens: metrics
.spec_decode_num_draft_tokens
.get_or_create_owned(&labels),
spec_decode_num_accepted_tokens: metrics
.spec_decode_num_accepted_tokens
.get_or_create_owned(&labels),
spec_decode_num_accepted_tokens_per_pos: metrics
.spec_decode_num_accepted_tokens_per_pos
.clone(),
log_stats: metrics.log_stats.get_or_create_owned(&labels),
estimated_flops_per_gpu: metrics.estimated_flops_per_gpu.get_or_create_owned(&labels),
estimated_read_bytes_per_gpu: metrics
.estimated_read_bytes_per_gpu
.get_or_create_owned(&labels),
estimated_write_bytes_per_gpu: metrics
.estimated_write_bytes_per_gpu
.get_or_create_owned(&labels),
kv_block_lifetime_seconds: metrics.kv_block_lifetime_seconds.get_or_create_owned(&labels),
kv_block_idle_before_evict_seconds: metrics
.kv_block_idle_before_evict_seconds
.get_or_create_owned(&labels),
kv_block_reuse_gap_seconds: metrics.kv_block_reuse_gap_seconds.get_or_create_owned(&labels),
labels,
}
}
/// Record scheduler-stats values through pre-resolved metric handles.
fn record_scheduler_stats_with_handles(handles: &SchedulerStatsHandles, stats: &SchedulerStats) {
// Scheduler state gauges.
handles.scheduler_running.set(stats.num_running_reqs);
handles
.scheduler_waiting
.set(stats.num_waiting_reqs + stats.num_skipped_waiting_reqs);
handles.scheduler_waiting_capacity.set(stats.num_waiting_reqs);
handles.scheduler_waiting_deferred.set(stats.num_skipped_waiting_reqs);
handles.kv_cache_usage.set(stats.kv_cache_usage);
// Prefix-cache counters, including the connector-backed external cache path.
handles.prefix_cache_queries.inc_by(stats.prefix_cache_stats.base.queries);
handles.prefix_cache_hits.inc_by(stats.prefix_cache_stats.base.hits);
if let Some(connector_prefix_cache_stats) = &stats.connector_prefix_cache_stats {
metrics
handles
.external_prefix_cache_queries
.get_or_create(&labels)
.inc_by(connector_prefix_cache_stats.base.queries);
metrics
handles
.external_prefix_cache_hits
.get_or_create(&labels)
.inc_by(connector_prefix_cache_stats.base.hits);
}
// Speculative decoding counters.
if let Some(spec_decoding_stats) = &stats.spec_decoding_stats {
metrics
.spec_decode_num_drafts
.get_or_create(&labels)
.inc_by(spec_decoding_stats.num_drafts);
metrics
handles.spec_decode_num_drafts.inc_by(spec_decoding_stats.num_drafts);
handles
.spec_decode_num_draft_tokens
.get_or_create(&labels)
.inc_by(spec_decoding_stats.num_draft_tokens);
metrics
handles
.spec_decode_num_accepted_tokens
.get_or_create(&labels)
.inc_by(spec_decoding_stats.num_accepted_tokens);
metrics.log_stats.get_or_create(&labels).observe_spec_decode(
handles.log_stats.observe_spec_decode(
spec_decoding_stats.num_drafts,
&spec_decoding_stats.num_accepted_tokens_per_pos,
);
@@ -92,11 +192,11 @@ pub(crate) fn record_scheduler_stats(
for (position, accepted_tokens) in
spec_decoding_stats.num_accepted_tokens_per_pos.iter().copied().enumerate()
{
metrics
handles
.spec_decode_num_accepted_tokens_per_pos
.get_or_create(&EnginePositionLabels {
model_name: model_name.clone(),
engine,
model_name: handles.labels.model_name.clone(),
engine: handles.labels.engine,
position: position as u32,
})
.inc_by(accepted_tokens);
@@ -109,22 +209,13 @@ pub(crate) fn record_scheduler_stats(
|| perf_stats.num_read_bytes_per_gpu != 0
|| perf_stats.num_write_bytes_per_gpu != 0)
{
metrics
.estimated_flops_per_gpu
.get_or_create(&labels)
.inc_by(perf_stats.num_flops_per_gpu);
metrics
.estimated_read_bytes_per_gpu
.get_or_create(&labels)
.inc_by(perf_stats.num_read_bytes_per_gpu);
metrics
.estimated_write_bytes_per_gpu
.get_or_create(&labels)
.inc_by(perf_stats.num_write_bytes_per_gpu);
handles.estimated_flops_per_gpu.inc_by(perf_stats.num_flops_per_gpu);
handles.estimated_read_bytes_per_gpu.inc_by(perf_stats.num_read_bytes_per_gpu);
handles.estimated_write_bytes_per_gpu.inc_by(perf_stats.num_write_bytes_per_gpu);
}
if let Some(cudagraph_stats) = &stats.cudagraph_stats {
metrics.log_stats.get_or_create(&labels).observe_cudagraph(
handles.log_stats.observe_cudagraph(
cudagraph_stats.num_unpadded_tokens,
cudagraph_stats.num_padded_tokens,
cudagraph_stats.num_paddings,
@@ -134,16 +225,11 @@ pub(crate) fn record_scheduler_stats(
// Sampled KV-cache residency histograms.
if !stats.kv_cache_eviction_events.is_empty() {
let kv_block_lifetime_seconds = metrics.kv_block_lifetime_seconds.get_or_create(&labels);
let kv_block_idle_before_evict_seconds =
metrics.kv_block_idle_before_evict_seconds.get_or_create(&labels);
let kv_block_reuse_gap_seconds = metrics.kv_block_reuse_gap_seconds.get_or_create(&labels);
for event in &stats.kv_cache_eviction_events {
kv_block_lifetime_seconds.observe(event.lifetime_seconds);
kv_block_idle_before_evict_seconds.observe(event.idle_seconds);
handles.kv_block_lifetime_seconds.observe(event.lifetime_seconds);
handles.kv_block_idle_before_evict_seconds.observe(event.idle_seconds);
for reuse_gap_seconds in &event.reuse_gaps_seconds {
kv_block_reuse_gap_seconds.observe(*reuse_gap_seconds);
handles.kv_block_reuse_gap_seconds.observe(*reuse_gap_seconds);
}
}
}
+11 -4
View File
@@ -88,14 +88,21 @@ impl Llm {
// Record internal engine-core request ID in the current tracing span.
Span::current().record("engine_request_id", &internal_request_id);
let arrival_time = prepared.engine_request.arrival_time;
let max_tokens_param =
(prepared.engine_request.sampling_params.as_ref()).map(|p| p.max_tokens);
let prompt_len = prepared.prompt_token_ids().len() as u32;
let stream = self.client.call(prepared.engine_request).await?;
let request_metrics = RequestMetricsTracker::new(
self.client.model_name().to_string(),
prepared.engine_request.arrival_time,
prepared.prompt_token_ids().len() as u32,
(prepared.engine_request.sampling_params.as_ref()).map(|p| p.max_tokens),
stream.engine_index(),
arrival_time,
prompt_len,
max_tokens_param,
1,
);
let stream = self.client.call(prepared.engine_request).await?;
let guard = self.inflight.track(external_request_id, internal_request_id);
Ok(GenerateOutputStream::new(
+1 -6
View File
@@ -248,12 +248,7 @@ impl Stream for GenerateOutputStream {
};
let received_at = current_unix_timestamp_secs();
self.request_metrics.observe_output(
raw.engine_index,
raw.timestamp,
received_at,
&raw.output,
);
self.request_metrics.observe_output(raw.timestamp, received_at, &raw.output);
let raw = raw.output;
+150 -142
View File
@@ -5,15 +5,12 @@ use vllm_engine_core_client::protocol::output::{
};
use vllm_engine_core_client::protocol::stats::PrefillStats;
use vllm_metrics::{
EngineLabels, FinishedReasonLabels, METRICS, PromptTokenSourceLabels, RequestMetrics,
EngineLabels, Family, FinishedReasonLabels, HistogramMetric, METRICS, PromptTokenSourceLabels,
U64Counter,
};
use crate::FinishReason;
fn metrics() -> &'static RequestMetrics {
&METRICS.request
}
const PROMPT_TOKEN_SOURCE_LOCAL_COMPUTE: &str = "local_compute";
const PROMPT_TOKEN_SOURCE_LOCAL_CACHE_HIT: &str = "local_cache_hit";
const PROMPT_TOKEN_SOURCE_EXTERNAL_KV_TRANSFER: &str = "external_kv_transfer";
@@ -29,9 +26,11 @@ const PROMPT_TOKEN_SOURCE_EXTERNAL_KV_TRANSFER: &str = "external_kv_transfer";
///
/// Original Python update flow:
/// <https://github.com/vllm-project/vllm/blob/bc2c0c86efb28e77677a3cfb8687e976914a313a/vllm/v1/engine/output_processor.py#L600-L677>
#[derive(Debug, Clone)]
#[derive(Clone)]
pub(crate) struct RequestMetricsTracker {
model_name: String,
/// Cached request metric handles for this request's model and engine index.
handles: RequestMetricHandles,
arrival_time: f64,
prompt_len: u32,
max_tokens_param: Option<u32>,
@@ -44,7 +43,38 @@ pub(crate) struct RequestMetricsTracker {
first_token_latency: f64,
num_generation_tokens: u32,
latest_num_cached_tokens: u32,
last_seen_engine_index: u32,
}
/// Cached request metric handles for one model and engine index.
#[derive(Clone)]
struct RequestMetricHandles {
labels: EngineLabels,
// Request-derived counters.
num_preemptions: U64Counter,
prompt_tokens: U64Counter,
prompt_tokens_local_compute: U64Counter,
prompt_tokens_local_cache_hit: U64Counter,
prompt_tokens_external_kv_transfer: U64Counter,
prompt_tokens_cached: U64Counter,
generation_tokens: U64Counter,
// Request lifecycle counters and histograms.
request_success: Family<FinishedReasonLabels, U64Counter>,
request_prompt_tokens: HistogramMetric,
request_generation_tokens: HistogramMetric,
request_max_num_generation_tokens: HistogramMetric,
request_params_max_tokens: HistogramMetric,
request_params_n: HistogramMetric,
request_prefill_kv_computed_tokens: HistogramMetric,
time_to_first_token_seconds: HistogramMetric,
inter_token_latency_seconds: HistogramMetric,
e2e_request_latency_seconds: HistogramMetric,
request_queue_time_seconds: HistogramMetric,
request_prefill_time_seconds: HistogramMetric,
request_decode_time_seconds: HistogramMetric,
request_inference_time_seconds: HistogramMetric,
request_time_per_output_token_seconds: HistogramMetric,
}
impl RequestMetricsTracker {
@@ -52,13 +82,14 @@ impl RequestMetricsTracker {
/// context.
pub(crate) fn new(
model_name: String,
engine_index: u32,
arrival_time: f64,
prompt_len: u32,
max_tokens_param: Option<u32>,
n_param: u32,
) -> Self {
Self {
model_name,
handles: resolve_request_metric_handles(&model_name, engine_index),
arrival_time,
prompt_len,
max_tokens_param,
@@ -71,7 +102,6 @@ impl RequestMetricsTracker {
first_token_latency: 0.0,
num_generation_tokens: 0,
latest_num_cached_tokens: 0,
last_seen_engine_index: 0,
}
}
@@ -81,23 +111,18 @@ impl RequestMetricsTracker {
/// <https://github.com/vllm-project/vllm/blob/bc2c0c86efb28e77677a3cfb8687e976914a313a/vllm/v1/metrics/stats.py#L331-L384>
pub(crate) fn observe_output(
&mut self,
engine_index: u32,
batch_timestamp: f64,
received_at: f64,
output: &EngineCoreOutput,
) {
self.last_seen_engine_index = engine_index;
if let Some(prefill_stats) = &output.prefill_stats {
self.latest_num_cached_tokens = prefill_stats.num_cached_tokens;
}
self.num_generation_tokens += output.new_token_ids.len() as u32;
metrics()
.generation_tokens
.get_or_create(&engine_labels(&self.model_name, engine_index))
.inc_by(output.new_token_ids.len() as u64);
self.handles.generation_tokens.inc_by(output.new_token_ids.len() as u64);
if let Some(events) = &output.events {
self.observe_events(engine_index, events);
self.observe_events(events);
}
// Only outputs that actually carry tokens drive token-timing metrics.
@@ -107,22 +132,16 @@ impl RequestMetricsTracker {
if !output.new_token_ids.is_empty() {
if self.is_prefilling {
if let Some(prefill_stats) = &output.prefill_stats {
record_prompt_tokens(&self.model_name, engine_index, prefill_stats);
self.record_prompt_tokens(prefill_stats);
}
self.first_token_latency = received_at - self.arrival_time;
observe_time_to_first_token_seconds(
&self.model_name,
engine_index,
self.first_token_latency,
);
self.handles.time_to_first_token_seconds.observe(self.first_token_latency);
self.first_token_ts = batch_timestamp;
self.is_prefilling = false;
} else if self.last_token_ts > 0.0 {
observe_inter_token_latency_seconds(
&self.model_name,
engine_index,
batch_timestamp - self.last_token_ts,
);
self.handles
.inter_token_latency_seconds
.observe(batch_timestamp - self.last_token_ts);
}
self.last_token_ts = batch_timestamp;
@@ -135,7 +154,6 @@ impl RequestMetricsTracker {
/// Original Python finished-request stats:
/// <https://github.com/vllm-project/vllm/blob/bc2c0c86efb28e77677a3cfb8687e976914a313a/vllm/v1/metrics/stats.py#L222-L237>
pub(crate) fn record_finished(&self, received_at: f64, finish_reason: FinishReason) {
let labels = engine_labels(&self.model_name, self.last_seen_engine_index);
let prefill_kv_computed_tokens =
self.prompt_len.saturating_sub(self.latest_num_cached_tokens);
let e2e_latency_seconds = received_at - self.arrival_time;
@@ -150,57 +168,47 @@ impl RequestMetricsTracker {
0.0
};
record_request_success(&self.model_name, self.last_seen_engine_index, finish_reason);
metrics()
.request_prompt_tokens
.get_or_create(&labels)
.observe(self.prompt_len as f64);
metrics()
self.record_request_success(finish_reason);
self.handles.request_prompt_tokens.observe(self.prompt_len as f64);
self.handles
.request_generation_tokens
.get_or_create(&labels)
.observe(self.num_generation_tokens as f64);
metrics()
self.handles
.request_max_num_generation_tokens
.get_or_create(&labels)
.observe(self.num_generation_tokens as f64);
if let Some(max_tokens_param) = self.max_tokens_param {
metrics()
.request_params_max_tokens
.get_or_create(&labels)
.observe(max_tokens_param as f64);
self.handles.request_params_max_tokens.observe(max_tokens_param as f64);
}
metrics().request_params_n.get_or_create(&labels).observe(self.n_param as f64);
metrics()
self.handles.request_params_n.observe(self.n_param as f64);
self.handles
.request_prefill_kv_computed_tokens
.get_or_create(&labels)
.observe(prefill_kv_computed_tokens as f64);
metrics()
.e2e_request_latency_seconds
.get_or_create(&labels)
.observe(e2e_latency_seconds);
metrics()
.request_queue_time_seconds
.get_or_create(&labels)
.observe(queue_time_seconds);
metrics()
.request_prefill_time_seconds
.get_or_create(&labels)
.observe(prefill_time_seconds);
metrics()
.request_decode_time_seconds
.get_or_create(&labels)
.observe(decode_time_seconds);
metrics()
.request_inference_time_seconds
.get_or_create(&labels)
.observe(inference_time_seconds);
metrics()
self.handles.e2e_request_latency_seconds.observe(e2e_latency_seconds);
self.handles.request_queue_time_seconds.observe(queue_time_seconds);
self.handles.request_prefill_time_seconds.observe(prefill_time_seconds);
self.handles.request_decode_time_seconds.observe(decode_time_seconds);
self.handles.request_inference_time_seconds.observe(inference_time_seconds);
self.handles
.request_time_per_output_token_seconds
.get_or_create(&labels)
.observe(time_per_output_token_seconds);
}
fn observe_events(&mut self, engine_index: u32, events: &[EngineCoreEvent]) {
/// Record prompt token counters through cached metric handles.
fn record_prompt_tokens(&self, prefill_stats: &PrefillStats) {
let computed = prefill_stats.num_computed_tokens as u64;
let local_cache_hit = prefill_stats.num_local_cached_tokens as u64;
let external_kv_transfer = prefill_stats.num_external_cached_tokens as u64;
self.handles.prompt_tokens.inc_by(prefill_stats.num_prompt_tokens as u64);
self.handles.prompt_tokens_local_compute.inc_by(computed);
self.handles.prompt_tokens_local_cache_hit.inc_by(local_cache_hit);
self.handles.prompt_tokens_external_kv_transfer.inc_by(external_kv_transfer);
self.handles.prompt_tokens_cached.inc_by(prefill_stats.num_cached_tokens as u64);
}
/// Record request event counters through cached metric handles.
fn observe_events(&mut self, events: &[EngineCoreEvent]) {
for event in events {
match event.r#type {
EngineCoreEventType::Queued => {
@@ -212,46 +220,86 @@ impl RequestMetricsTracker {
}
}
EngineCoreEventType::Preempted => {
metrics()
.num_preemptions
.get_or_create(&engine_labels(&self.model_name, engine_index))
.inc();
self.handles.num_preemptions.inc();
}
}
}
}
}
fn engine_labels(model_name: &str, engine: u32) -> EngineLabels {
EngineLabels {
model_name: model_name.to_string(),
engine,
/// Increment the request-success counter for the terminal finish reason.
fn record_request_success(&self, finish_reason: FinishReason) {
self.handles
.request_success
.get_or_create(&FinishedReasonLabels {
model_name: self.handles.labels.model_name.clone(),
engine: self.handles.labels.engine,
finished_reason: finish_reason.as_str(),
})
.inc();
}
}
fn observe_time_to_first_token_seconds(model_name: &str, engine: u32, seconds: f64) {
metrics()
.time_to_first_token_seconds
.get_or_create(&engine_labels(model_name, engine))
.observe(seconds);
}
/// Resolve fixed request metric handles for one model and engine index.
fn resolve_request_metric_handles(model_name: &str, engine: u32) -> RequestMetricHandles {
let metrics = &METRICS.request;
let labels = EngineLabels {
model_name: model_name.to_string(),
engine,
};
fn observe_inter_token_latency_seconds(model_name: &str, engine: u32, seconds: f64) {
metrics()
.inter_token_latency_seconds
.get_or_create(&engine_labels(model_name, engine))
.observe(seconds);
}
fn record_request_success(model_name: &str, engine: u32, finish_reason: FinishReason) {
metrics()
.request_success
.get_or_create(&FinishedReasonLabels {
model_name: model_name.to_string(),
engine,
finished_reason: finish_reason.as_str(),
})
.inc();
RequestMetricHandles {
num_preemptions: metrics.num_preemptions.get_or_create_owned(&labels),
prompt_tokens: metrics.prompt_tokens.get_or_create_owned(&labels),
prompt_tokens_local_compute: metrics.prompt_tokens_by_source.get_or_create_owned(
&prompt_token_source_labels(model_name, engine, PROMPT_TOKEN_SOURCE_LOCAL_COMPUTE),
),
prompt_tokens_local_cache_hit: metrics.prompt_tokens_by_source.get_or_create_owned(
&prompt_token_source_labels(model_name, engine, PROMPT_TOKEN_SOURCE_LOCAL_CACHE_HIT),
),
prompt_tokens_external_kv_transfer: metrics.prompt_tokens_by_source.get_or_create_owned(
&prompt_token_source_labels(
model_name,
engine,
PROMPT_TOKEN_SOURCE_EXTERNAL_KV_TRANSFER,
),
),
prompt_tokens_cached: metrics.prompt_tokens_cached.get_or_create_owned(&labels),
generation_tokens: metrics.generation_tokens.get_or_create_owned(&labels),
request_success: metrics.request_success.clone(),
request_prompt_tokens: metrics.request_prompt_tokens.get_or_create_owned(&labels),
request_generation_tokens: metrics.request_generation_tokens.get_or_create_owned(&labels),
request_max_num_generation_tokens: metrics
.request_max_num_generation_tokens
.get_or_create_owned(&labels),
request_params_max_tokens: metrics.request_params_max_tokens.get_or_create_owned(&labels),
request_params_n: metrics.request_params_n.get_or_create_owned(&labels),
request_prefill_kv_computed_tokens: metrics
.request_prefill_kv_computed_tokens
.get_or_create_owned(&labels),
time_to_first_token_seconds: metrics
.time_to_first_token_seconds
.get_or_create_owned(&labels),
inter_token_latency_seconds: metrics
.inter_token_latency_seconds
.get_or_create_owned(&labels),
e2e_request_latency_seconds: metrics
.e2e_request_latency_seconds
.get_or_create_owned(&labels),
request_queue_time_seconds: metrics.request_queue_time_seconds.get_or_create_owned(&labels),
request_prefill_time_seconds: metrics
.request_prefill_time_seconds
.get_or_create_owned(&labels),
request_decode_time_seconds: metrics
.request_decode_time_seconds
.get_or_create_owned(&labels),
request_inference_time_seconds: metrics
.request_inference_time_seconds
.get_or_create_owned(&labels),
request_time_per_output_token_seconds: metrics
.request_time_per_output_token_seconds
.get_or_create_owned(&labels),
labels,
}
}
fn prompt_token_source_labels(
@@ -266,45 +314,6 @@ fn prompt_token_source_labels(
}
}
fn record_prompt_tokens(model_name: &str, engine: u32, prefill_stats: &PrefillStats) {
let computed = prefill_stats.num_computed_tokens as u64;
let local_cache_hit = prefill_stats.num_local_cached_tokens as u64;
let external_kv_transfer = prefill_stats.num_external_cached_tokens as u64;
metrics()
.prompt_tokens
.get_or_create(&engine_labels(model_name, engine))
.inc_by(prefill_stats.num_prompt_tokens as u64);
metrics()
.prompt_tokens_by_source
.get_or_create(&prompt_token_source_labels(
model_name,
engine,
PROMPT_TOKEN_SOURCE_LOCAL_COMPUTE,
))
.inc_by(computed);
metrics()
.prompt_tokens_by_source
.get_or_create(&prompt_token_source_labels(
model_name,
engine,
PROMPT_TOKEN_SOURCE_LOCAL_CACHE_HIT,
))
.inc_by(local_cache_hit);
metrics()
.prompt_tokens_by_source
.get_or_create(&prompt_token_source_labels(
model_name,
engine,
PROMPT_TOKEN_SOURCE_EXTERNAL_KV_TRANSFER,
))
.inc_by(external_kv_transfer);
metrics()
.prompt_tokens_cached
.get_or_create(&engine_labels(model_name, engine))
.inc_by(prefill_stats.num_cached_tokens as u64);
}
fn diff_or_zero(end: f64, start: f64) -> f64 {
if end > 0.0 && start > 0.0 && end >= start {
end - start
@@ -337,10 +346,10 @@ mod tests {
#[test]
fn tracker_updates_timing_state_across_prefill_decode_and_finish() {
let mut tracker = RequestMetricsTracker::new("model".to_string(), 100.0, 64, Some(128), 1);
let mut tracker =
RequestMetricsTracker::new("model".to_string(), 2, 100.0, 64, Some(128), 1);
tracker.observe_output(
2,
10.0,
100.2,
&vllm_engine_core_client::protocol::output::EngineCoreOutput {
@@ -368,7 +377,6 @@ mod tests {
},
);
tracker.observe_output(
2,
11.5,
100.4,
&vllm_engine_core_client::protocol::output::EngineCoreOutput {
@@ -384,7 +392,7 @@ mod tests {
);
assert!(!tracker.is_prefilling);
assert_eq!(tracker.last_seen_engine_index, 2);
assert_eq!(tracker.handles.labels.engine, 2);
assert_eq!(tracker.num_generation_tokens, 3);
assert_eq!(tracker.queued_ts, 8.0);
assert_eq!(tracker.scheduled_ts, 9.0);
+3 -3
View File
@@ -17,7 +17,7 @@ use vllm_engine_core_client::protocol::request::EngineCoreRequest;
use vllm_engine_core_client::protocol::sampling::EngineCoreSamplingParams;
use vllm_engine_core_client::protocol::stats::PrefillStats;
use vllm_engine_core_client::test_utils::{IpcNamespace, spawn_mock_engine_task};
use vllm_engine_core_client::{EngineCoreClient, EngineCoreClientConfig};
use vllm_engine_core_client::{EngineCoreClient, EngineCoreClientConfig, EngineId};
use vllm_llm::{
Error, FinishReason, GenerateOutputStreamExt as _, GeneratePromptInfo, GenerateRequest, Llm,
};
@@ -699,7 +699,7 @@ async fn abort_by_external_id_aborts_all_internal_requests() {
async fn generate_records_request_metrics_in_prometheus_output() {
let ipc = IpcNamespace::new().unwrap();
let handshake_address = ipc.handshake_endpoint();
let engine_id = b"engine-metrics".to_vec();
let engine_id = EngineId::from_engine_index(4);
let model_name = request_metrics_model_name("metrics-model");
let (shutdown_tx, engine_task) = spawn_mock_engine_task(
@@ -832,7 +832,7 @@ async fn generate_records_request_metrics_in_prometheus_output() {
async fn dropping_stream_records_abort_terminal_request_metrics() {
let ipc = IpcNamespace::new().unwrap();
let handshake_address = ipc.handshake_endpoint();
let engine_id = b"engine-metrics-drop".to_vec();
let engine_id = EngineId::from_engine_index(5);
let model_name = request_metrics_model_name("metrics-drop-model");
let (shutdown_tx, engine_task) = spawn_mock_engine_task(
+3 -1
View File
@@ -4,7 +4,7 @@ use std::sync::atomic::AtomicU64;
use prometheus_client::encoding::text::encode;
use prometheus_client::metrics::counter::Counter;
use prometheus_client::metrics::family::Family;
pub use prometheus_client::metrics::family::Family;
use prometheus_client::metrics::gauge::Gauge;
use prometheus_client::metrics::histogram::Histogram;
use prometheus_client::registry::Registry;
@@ -23,6 +23,8 @@ pub use scheduler::*;
pub type U64Counter = Counter<u64, AtomicU64>;
pub type U64Gauge = Gauge<u64, AtomicU64>;
pub type F64Gauge = Gauge<f64, AtomicU64>;
/// Histogram metric handle cloned out of a Prometheus family.
pub type HistogramMetric = Histogram;
pub(crate) type HistogramFamily = Family<EngineLabels, Histogram, fn() -> Histogram>;
/// Shared Prometheus registry for frontend metrics.
+1
View File
@@ -79,6 +79,7 @@ def run_e2e_fusion_test(monkeypatch, caplog_mp_spawn):
):
monkeypatch.setenv("VLLM_USE_DEEP_GEMM", "1" if use_deepgemm else "0")
monkeypatch.setenv("VLLM_ROCM_USE_AITER", "1" if use_aiter else "0")
monkeypatch.setenv("VLLM_ROCM_USE_AITER_CUSTOM_AR", "1" if use_aiter else "0")
from vllm._aiter_ops import rocm_aiter_ops
rocm_aiter_ops.refresh_env_variables()
@@ -13,6 +13,7 @@ from vllm._custom_ops import cutlass_scaled_fp4_mm, scaled_fp4_quant
from vllm.compilation.passes.fusion.allreduce_rms_fusion import (
AllReduceFusionPass,
RocmAiterAllReduceFusionPass,
_select_flashinfer_allreduce_use_oneshot,
)
from vllm.compilation.passes.fx_utils import find_op_nodes
from vllm.compilation.passes.utility.fix_functionalization import (
@@ -30,6 +31,9 @@ from vllm.config import (
set_current_vllm_config,
)
from vllm.distributed import tensor_model_parallel_all_reduce
from vllm.distributed.device_communicators.aiter_custom_all_reduce import (
AiterCustomAllreduce,
)
from vllm.distributed.parallel_state import (
init_distributed_environment,
initialize_model_parallel,
@@ -45,6 +49,35 @@ from vllm.utils.torch_utils import set_random_seed
DEVICE_TYPE = current_platform.device_type
@pytest.mark.parametrize(
("workspace_backend", "device_capability", "world_size", "tensor_size", "expected"),
[
("mnnvl", 103, 8, 2 * 1024 * 1024, None),
("trtllm", 103, 8, 2 * 1024 * 1024, True),
("trtllm", 103, 8, 2 * 1024 * 1024 + 1, False),
("trtllm", 100, 4, 4 * 1024 * 1024, True),
("trtllm", 100, 4, 4 * 1024 * 1024 + 1, False),
("trtllm", None, 8, 128 * 1024 * 1024, True),
],
)
def test_select_flashinfer_allreduce_use_oneshot(
workspace_backend: str,
device_capability: int | None,
world_size: int,
tensor_size: int,
expected: bool | None,
):
assert (
_select_flashinfer_allreduce_use_oneshot(
workspace_backend,
device_capability,
world_size,
tensor_size,
)
is expected
)
class TestAllReduceRMSNormModel(torch.nn.Module):
def __init__(
self,
@@ -504,8 +537,12 @@ def all_reduce_fusion_pass_on_test_model(
"MASTER_ADDR": "localhost",
"MASTER_PORT": "12345",
"VLLM_FLASHINFER_ALLREDUCE_BACKEND": flashinfer_allreduce_backend,
"VLLM_ROCM_USE_AITER": str(int(use_aiter)),
"VLLM_ROCM_USE_AITER_CUSTOM_AR": str(int(use_aiter)),
}
)
if use_aiter:
rocm_aiter_ops.refresh_env_variables()
init_distributed_environment()
@@ -616,7 +653,7 @@ def test_rocm_aiter_all_reduce_rmsnorm_group_quant_fp8_fusion_pass_replace(
m.setenv("VLLM_ROCM_USE_AITER", "1")
rocm_aiter_ops.refresh_env_variables()
if not rocm_aiter_ops.has_fused_allreduce_rmsnorm_quant_per_group():
if not AiterCustomAllreduce.build_supports_per_group_quant():
pytest.skip(
"aiter build is missing 'fused_ar_rms_per_group_quant' (needs "
"ROCm/aiter PR #2823); the new patterns aren't registered."
@@ -671,6 +708,7 @@ def rocm_aiter_group_quant_fusion_pass_on_test_model(
"MASTER_ADDR": "localhost",
"MASTER_PORT": "12345",
"VLLM_ROCM_USE_AITER": "1",
"VLLM_ROCM_USE_AITER_CUSTOM_AR": "1",
}
)
rocm_aiter_ops.refresh_env_variables()
+52
View File
@@ -0,0 +1,52 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
from transformers import PretrainedConfig
from vllm.config.speculative import MTPModelTypes, SpeculativeConfig
from vllm.transformers_utils.model_arch_config_convertor import (
BailingHybridMTPModelArchConfigConvertor,
)
def _bailing_config() -> PretrainedConfig:
config = PretrainedConfig(
architectures=["BailingMoeV2_5ForCausalLM"],
hidden_size=4096,
kv_lora_rank=512,
num_attention_heads=32,
num_experts=256,
num_hidden_layers=32,
num_key_value_heads=32,
num_nextn_predict_layers=1,
qk_rope_head_dim=64,
vocab_size=157184,
)
config.model_type = "bailing_hybrid"
return config
def test_bailing_hybrid_mtp_hf_config_override():
config = _bailing_config()
overridden = SpeculativeConfig.hf_config_override(config)
assert overridden.model_type == "bailing_hybrid_mtp"
assert overridden.architectures == ["BailingMoeV25MTPModel"]
assert overridden.n_predict == 1
assert "bailing_hybrid_mtp" in MTPModelTypes.__args__
def test_bailing_hybrid_mtp_model_arch_config():
config = _bailing_config()
config.model_type = "bailing_hybrid_mtp"
config.architectures = ["BailingMoeV25MTPModel"]
model_arch_config = BailingHybridMTPModelArchConfigConvertor(
config, config
).convert()
assert model_arch_config.model_type == "bailing_hybrid_mtp"
assert model_arch_config.architectures == ["BailingMoeV25MTPModel"]
assert model_arch_config.total_num_hidden_layers == 1
assert model_arch_config.is_deepseek_mla
@@ -0,0 +1,134 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import pytest
import ray
import torch
import torch.distributed as dist
from vllm._aiter_ops import is_aiter_found, rocm_aiter_ops
from vllm.distributed.communication_op import tensor_model_parallel_all_reduce # noqa
from vllm.distributed.parallel_state import get_tp_group, graph_capture
from vllm.envs import disable_envs_cache
from vllm.platforms import current_platform
from ..utils import (
assert_rocm_custom_allreduce_backend_state,
ensure_model_parallel_initialized,
init_test_distributed_environment,
multi_gpu_test,
multi_process_parallel,
)
pytestmark = pytest.mark.skipif(
not current_platform.is_rocm(),
reason="ROCm-only AITER custom allreduce tests",
)
test_cases = [
((2, 7168), torch.float16),
((2, 7168), torch.bfloat16),
((128, 8192), torch.float16),
((128, 8192), torch.bfloat16),
]
def _configure_aiter_custom_ar_env(monkeypatch: pytest.MonkeyPatch) -> None:
monkeypatch.delenv("CUDA_VISIBLE_DEVICES", raising=False)
monkeypatch.delenv("HIP_VISIBLE_DEVICES", raising=False)
monkeypatch.setenv("VLLM_ROCM_USE_AITER", "1")
monkeypatch.setenv("VLLM_ROCM_USE_AITER_CUSTOM_AR", "1")
monkeypatch.setenv("VLLM_ROCM_QUICK_REDUCE_QUANTIZATION", "NONE")
disable_envs_cache()
rocm_aiter_ops.refresh_env_variables()
def _assert_aiter_handles_input(inp: torch.Tensor) -> None:
aiter_ar_comm = get_tp_group().device_communicator.aiter_ar_comm
assert aiter_ar_comm is not None
assert aiter_ar_comm.should_custom_ar(inp), (
f"AITER CustomAllreduce does not support input shape {inp.shape}."
)
@ray.remote(num_gpus=1, max_calls=1)
def graph_allreduce(
monkeypatch: pytest.MonkeyPatch,
tp_size,
pp_size,
rank,
distributed_init_port,
) -> None:
with monkeypatch.context() as m:
_configure_aiter_custom_ar_env(m)
device = torch.device(f"cuda:{rank}")
torch.accelerator.set_device_index(device)
init_test_distributed_environment(tp_size, pp_size, rank, distributed_init_port)
ensure_model_parallel_initialized(tp_size, pp_size)
assert_rocm_custom_allreduce_backend_state(True, "NONE")
group = get_tp_group().device_group
# A small all_reduce for warmup.
# this is needed because device communicators might be created lazily
# (e.g. NCCL). This will ensure that the communicator is initialized
# before any communication happens, so that this group can be used for
# graph capture immediately.
data = torch.zeros(1)
data = data.to(device=device)
dist.all_reduce(data, group=group)
torch.accelerator.synchronize()
del data
for shape, dtype in test_cases:
with graph_capture(device=device) as graph_capture_context:
inp = torch.ones(shape, dtype=dtype, device=device)
_assert_aiter_handles_input(inp)
expected = inp * tp_size
torch.accelerator.synchronize()
graph = torch.cuda.CUDAGraph()
with torch.cuda.graph(graph, stream=graph_capture_context.stream):
out = tensor_model_parallel_all_reduce(inp)
graph.replay()
torch.testing.assert_close(out, expected)
@ray.remote(num_gpus=1, max_calls=1)
def eager_allreduce(
monkeypatch: pytest.MonkeyPatch,
tp_size,
pp_size,
rank,
distributed_init_port,
) -> None:
with monkeypatch.context() as m:
_configure_aiter_custom_ar_env(m)
device = torch.device(f"cuda:{rank}")
torch.accelerator.set_device_index(device)
init_test_distributed_environment(tp_size, pp_size, rank, distributed_init_port)
ensure_model_parallel_initialized(tp_size, pp_size)
assert_rocm_custom_allreduce_backend_state(True, "NONE")
for shape, dtype in test_cases:
inp = torch.ones(shape, dtype=dtype, device=device)
_assert_aiter_handles_input(inp)
expected = inp * tp_size
out = tensor_model_parallel_all_reduce(inp)
torch.testing.assert_close(out, expected)
@pytest.mark.skipif(not is_aiter_found(), reason="AITER is not installed")
@multi_gpu_test(num_gpus=2)
@pytest.mark.parametrize("tp_size", [2])
@pytest.mark.parametrize("pipeline_parallel_size", [1])
@pytest.mark.parametrize("test_target", [eager_allreduce, graph_allreduce])
def test_rocm_aiter_custom_allreduce(
monkeypatch: pytest.MonkeyPatch,
tp_size,
pipeline_parallel_size,
test_target,
):
multi_process_parallel(monkeypatch, tp_size, pipeline_parallel_size, test_target)
@@ -71,8 +71,9 @@ def test_lm_eval_accuracy_v1_engine():
more_args = []
# Limit compilation time for V1
if current_platform.is_tpu():
# Limit compilation time for V1 on TPU
# Avoid OOM on XPU
if current_platform.is_tpu() or current_platform.is_xpu():
more_args = ["--max-num-seqs", "64"]
run_test(more_args)
@@ -0,0 +1,114 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import json
import openai # use the official client for correctness check
import pytest
MODEL_NAME = "Qwen/Qwen3-1.7B"
NAMESPACE = "mcp__computer_use"
TOOL_NAME = "get_app_state"
FLAT_TOOL_NAME = f"{NAMESPACE}__{TOOL_NAME}"
tools = [
{
"type": "namespace",
"name": NAMESPACE,
"description": "Computer control tools.",
"tools": [
{
"type": "function",
"name": TOOL_NAME,
"description": "Get the current state of a desktop application.",
"parameters": {
"type": "object",
"properties": {
"app": {
"type": "string",
"description": "Application name, for example Chrome.",
}
},
"required": ["app"],
"additionalProperties": False,
},
}
],
}
]
prompt = [
{
"role": "user",
"content": "Use the computer app state tool to inspect Google Chrome.",
},
]
def _assert_namespace_tool_call(tool_call) -> None:
assert tool_call.type == "function_call"
assert tool_call.name == TOOL_NAME
assert tool_call.namespace == NAMESPACE
assert tool_call.name != FLAT_TOOL_NAME
args = json.loads(tool_call.arguments)
assert args["app"]
@pytest.mark.asyncio
@pytest.mark.parametrize("model_name", [MODEL_NAME])
async def test_namespace_tool_separator(client: openai.AsyncOpenAI, model_name: str):
response = await client.responses.create(
model=model_name,
input=prompt,
tools=tools,
temperature=0.0,
)
assert len(response.output) >= 1
tool_call = next(
(out for out in response.output if out.type == "function_call"), None
)
assert tool_call is not None
_assert_namespace_tool_call(tool_call)
@pytest.mark.asyncio
@pytest.mark.parametrize("model_name", [MODEL_NAME])
async def test_namespace_tool_separator_streaming(
client: openai.AsyncOpenAI, model_name: str
):
stream = await client.responses.create(
model=model_name,
input=prompt,
tools=tools,
temperature=0.0,
stream=True,
)
events = [event async for event in stream]
added_call = next(
(
event.item
for event in events
if event.type == "response.output_item.added"
and getattr(event.item, "type", None) == "function_call"
),
None,
)
done_call = next(
(
event.item
for event in events
if event.type == "response.output_item.done"
and getattr(event.item, "type", None) == "function_call"
),
None,
)
assert added_call is not None
assert added_call.name == TOOL_NAME
assert added_call.namespace == NAMESPACE
assert done_call is not None
_assert_namespace_tool_call(done_call)
@@ -0,0 +1,224 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Tests for the silu_and_mul_per_block_quant helion kernel
Run `pytest tests/kernels/helion/test_silu_and_mul_per_block_quant.py`.
"""
from typing import Any
import pytest
import torch
from torch._subclasses.fake_tensor import FakeTensorMode
from tests.kernels.helion.utils import skip_if_platform_unsupported
from tests.kernels.quant_utils import FP8_DTYPE
from vllm.kernels.helion.case_key import CaseKey
from vllm.kernels.helion.config_manager import ConfigManager
from vllm.kernels.helion.ops.silu_and_mul_per_block_quant import (
_pick_cache,
baseline,
pick_config,
silu_and_mul_per_block_quant,
)
from vllm.platforms import current_platform
from vllm.utils.import_utils import has_helion
from vllm.utils.torch_utils import set_random_seed
if not has_helion():
pytest.skip(
"Helion is not installed. Install with: pip install vllm[helion]",
allow_module_level=True,
)
def _generate_fake_input(
num_tokens: int, intermediate_size: int, group_size: int
) -> tuple[Any, ...]:
with FakeTensorMode():
in_dtype: torch.dtype = torch.bfloat16
out_dtype: torch.dtype = current_platform.fp8_dtype()
scale_dtype: torch.dtype = torch.float32
input = torch.randn(
num_tokens, 2 * intermediate_size, device="cuda", dtype=in_dtype
)
result = torch.empty(
num_tokens, intermediate_size, device=input.device, dtype=out_dtype
)
scale = torch.empty(
(num_tokens, intermediate_size // group_size),
device=input.device,
dtype=scale_dtype,
)
scale_ub = torch.mean(input).to(scale_dtype)
args = (
result,
input,
scale,
group_size,
scale_ub,
False,
)
return args
class TestSiluAndMulPerBlockQuantConfigPicker:
def setup_method(self):
_pick_cache.clear()
def test_config_picker_exact_match(self):
config_keys = [
CaseKey({"intermediate_size": 2048, "group_size": 64, "num_tokens": 16}),
CaseKey({"intermediate_size": 4096, "group_size": 128, "num_tokens": 16}),
]
args = _generate_fake_input(16, 4096, 128)
selected_key = pick_config(args, config_keys)
assert selected_key == CaseKey(
{"intermediate_size": 4096, "group_size": 128, "num_tokens": 16}
)
def test_config_picker_closest_match(self):
config_keys = [
CaseKey({"intermediate_size": 2048, "group_size": 64, "num_tokens": 16}),
CaseKey({"intermediate_size": 2048, "group_size": 64, "num_tokens": 32}),
CaseKey({"intermediate_size": 2048, "group_size": 128, "num_tokens": 16}),
CaseKey({"intermediate_size": 2048, "group_size": 128, "num_tokens": 32}),
CaseKey({"intermediate_size": 4096, "group_size": 64, "num_tokens": 16}),
CaseKey({"intermediate_size": 4096, "group_size": 64, "num_tokens": 32}),
CaseKey({"intermediate_size": 4096, "group_size": 128, "num_tokens": 16}),
CaseKey({"intermediate_size": 4096, "group_size": 128, "num_tokens": 32}),
]
args = _generate_fake_input(20, 3000, 70)
selected_key = pick_config(args, config_keys)
assert selected_key == CaseKey(
{"intermediate_size": 2048, "group_size": 64, "num_tokens": 32}
)
def test_config_picker_no_configs(self):
config_keys: list[dict] = []
args = _generate_fake_input(16, 4096, 128)
selected_key = pick_config(args, config_keys)
assert selected_key is None
def test_config_picker_fallback_to_largest(self):
config_keys = [
CaseKey({"intermediate_size": 2048, "group_size": 64, "num_tokens": 16}),
CaseKey({"intermediate_size": 2048, "group_size": 64, "num_tokens": 32}),
CaseKey({"intermediate_size": 2048, "group_size": 128, "num_tokens": 16}),
CaseKey({"intermediate_size": 2048, "group_size": 128, "num_tokens": 32}),
CaseKey({"intermediate_size": 4096, "group_size": 64, "num_tokens": 16}),
CaseKey({"intermediate_size": 4096, "group_size": 64, "num_tokens": 32}),
CaseKey({"intermediate_size": 4096, "group_size": 128, "num_tokens": 16}),
CaseKey({"intermediate_size": 4096, "group_size": 128, "num_tokens": 32}),
]
args = _generate_fake_input(64, 8192, 256)
selected_key = pick_config(args, config_keys)
assert selected_key == CaseKey(
{"intermediate_size": 4096, "group_size": 128, "num_tokens": 32}
)
@pytest.fixture(autouse=True)
def reset_config_manager_singleton():
ConfigManager.reset_instance()
ConfigManager()
yield
ConfigManager.reset_instance()
class TestSiluAndMulPerBlockQuantCorrectness:
@pytest.mark.parametrize("num_tokens", [1, 7, 4096])
@pytest.mark.parametrize("hidden_size", [1024, 2048, 5120])
@pytest.mark.parametrize("group_size", [64, 128])
@pytest.mark.parametrize("is_scale_transposed", [False, True])
@pytest.mark.parametrize("dtype", [torch.bfloat16, torch.float16])
@pytest.mark.parametrize("quant_dtype", [current_platform.fp8_dtype(), torch.int8])
@pytest.mark.parametrize("has_scale_ub", [True, False])
@pytest.mark.parametrize("seed", [0])
def test_silu_and_mul_per_block_quant(
self,
num_tokens: int,
hidden_size: int,
group_size: int,
is_scale_transposed: bool,
dtype: torch.dtype,
quant_dtype: torch.dtype,
has_scale_ub: bool,
seed: int,
) -> None:
skip_if_platform_unsupported("silu_and_mul_per_block_quant")
set_random_seed(seed)
if hidden_size % group_size != 0:
return
if has_scale_ub and quant_dtype != FP8_DTYPE:
# skip
return
scale = 1 / hidden_size
x = torch.randn(num_tokens, 2 * hidden_size, dtype=dtype, device="cuda") * scale
if has_scale_ub:
act = torch.nn.functional.silu(x[:, :hidden_size]) * x[:, hidden_size:]
act_abs = act.abs().float()
scale_ub = 0.5 * (act_abs.mean() + act_abs.amax())
else:
scale_ub = None
ref_out = torch.empty(num_tokens, hidden_size, device="cuda", dtype=quant_dtype)
if is_scale_transposed:
ref_scales = torch.empty(
(hidden_size // group_size, x.shape[0]),
device="cuda",
dtype=torch.float32,
).t()
else:
ref_scales = torch.empty(
(x.shape[0], hidden_size // group_size),
device="cuda",
dtype=torch.float32,
)
ops_out = ref_out.clone()
ops_scales = ref_scales.clone()
baseline(ref_out, x, ref_scales, group_size, scale_ub, is_scale_transposed)
silu_and_mul_per_block_quant(
ops_out, x, ops_scales, group_size, scale_ub, is_scale_transposed
)
torch.testing.assert_close(ref_scales, ops_scales)
# allow 1 ULP difference
assert (
ref_out.view(torch.uint8).to(torch.int16)
- ops_out.view(torch.uint8).to(torch.int16)
).abs().max() <= 1
class TestSiluAndMulPerBlockQuantIntegration:
def test_kernel_registration_integration(self):
from vllm.kernels.helion.register import get_registered_kernels
registered_kernels = get_registered_kernels()
assert "silu_and_mul_per_block_quant" in registered_kernels
kernel_wrapper = registered_kernels["silu_and_mul_per_block_quant"]
assert kernel_wrapper.op_name == "silu_and_mul_per_block_quant"
assert kernel_wrapper._config_picker is not None
assert kernel_wrapper._mutates_args == ["out", "scales"]
def test_fake_impl_functionality(self):
skip_if_platform_unsupported("silu_and_mul_per_block_quant")
from vllm.kernels.helion.register import get_registered_kernels
registered_kernels = get_registered_kernels()
kernel_wrapper = registered_kernels["silu_and_mul_per_block_quant"]
fake_impl = kernel_wrapper._fake_impl
args = _generate_fake_input(16, 4096, 128)
assert fake_impl(*args) is None
+42 -112
View File
@@ -8,19 +8,21 @@ import pytest
import torch
import torch.nn.functional as F
from vllm.platforms import current_platform
from vllm.model_executor.layers.fused_moe.activation import MoEActivation
from vllm.model_executor.layers.fused_moe.experts.cpu_int4_moe import (
CPUExpertsInt4,
)
from vllm.model_executor.layers.fused_moe.oracle.w4a8_int8 import (
convert_to_w4a8_int8_moe_format,
)
from vllm.platforms import CpuArchEnum, current_platform
from vllm.utils.torch_utils import set_random_seed
if not current_platform.is_cpu():
pytest.skip("skipping CPU-only tests", allow_module_level=True)
# Check if the dynamic_4bit_int_moe op is available
if not hasattr(torch.ops._C, "dynamic_4bit_int_moe"):
pytest.skip("dynamic_4bit_int_moe op not available", allow_module_level=True)
# Check if KleidiAI ops are available
if not hasattr(torch.ops.aten, "_dyn_quant_pack_4bit_weight"):
pytest.skip("KleidiAI 4-bit ops not available", allow_module_level=True)
if (
not current_platform.is_cpu()
or current_platform.get_cpu_architecture() != CpuArchEnum.ARM
):
pytest.skip("skipping Arm CPU-only tests", allow_module_level=True)
# Tolerance for INT4 W4A8
@@ -34,49 +36,6 @@ def _silu_and_mul(x: torch.Tensor) -> torch.Tensor:
return F.silu(x[..., :d]) * x[..., d:]
def _pack_int4_weight_to_kleidi(
int4_as_int8: torch.Tensor,
scales: torch.Tensor,
bias: torch.Tensor | None,
group_size: int,
in_features: int,
out_features: int,
) -> torch.Tensor:
"""Pack INT4 weights (stored as int8 in [-8,7]) to KleidiAI format.
Args:
int4_as_int8: [out, in] int8 tensor with values in [-8, 7]
scales: [out, in//group_size] or [out, 1] for channel-wise
bias: [out] optional bias
group_size: Quantization group size (-1 for channel-wise)
in_features: Input dimension
out_features: Output dimension
Returns:
Packed weight tensor in KleidiAI format
"""
# Shift to unsigned nibble [0, 15]
tmp = int4_as_int8.add(8)
# Pack pairs along input dimension
uint8_nibbles = ((tmp[:, 1::2] << 4) | tmp[:, ::2]).to(torch.uint8)
# Determine scale dtype based on group_size
scale_dtype = torch.float32 if group_size == -1 else torch.bfloat16
scales_typed = scales.to(scale_dtype)
bias_typed = None if bias is None else bias.to(torch.float32)
# Pack using KleidiAI op
actual_group_size = in_features if group_size == -1 else group_size
return torch.ops.aten._dyn_quant_pack_4bit_weight(
uint8_nibbles,
scales_typed,
bias_typed,
actual_group_size,
in_features,
out_features,
)
def _make_int4_moe_weights(
E: int,
N: int,
@@ -124,59 +83,29 @@ def _make_int4_moe_weights(
w13_bias = torch.randn(E, 2 * N, dtype=torch.float32) * 0.01
w2_bias = torch.randn(E, K, dtype=torch.float32) * 0.01
# Pack weights for each expert
w13_packed_list = []
w2_packed_list = []
w13_packed, w2_packed, *_ = convert_to_w4a8_int8_moe_format(
w13_weight=w13_int4,
w2_weight=w2_int4,
w13_weight_scale=w13_scales,
w2_weight_scale=w2_scales,
group_size=group_size,
w13_bias=w13_bias if has_bias else None,
w2_bias=w2_bias if has_bias else None,
)
for e in range(E):
w13_packed_list.append(
_pack_int4_weight_to_kleidi(
w13_int4[e],
w13_scales[e],
w13_bias[e] if (has_bias and w13_bias is not None) else None,
group_size,
K,
2 * N,
)
)
w2_packed_list.append(
_pack_int4_weight_to_kleidi(
w2_int4[e],
w2_scales[e],
w2_bias[e] if (has_bias and w2_bias is not None) else None,
group_size,
N,
K,
)
)
if group_size == -1:
w13_scale = w13_scales.float()
w2_scale = w2_scales.float()
else:
w13_scale = w13_scales.float().repeat_interleave(group_size, dim=-1)
w2_scale = w2_scales.float().repeat_interleave(group_size, dim=-1)
w13_packed = torch.stack(w13_packed_list, dim=0)
w2_packed = torch.stack(w2_packed_list, dim=0)
# Create reference dequantized weights
w13_ref = torch.zeros(E, 2 * N, K, dtype=torch.float32)
w2_ref = torch.zeros(E, K, N, dtype=torch.float32)
for e in range(E):
# Dequantize w13
for i in range(2 * N):
for j in range(K):
group_idx = 0 if group_size == -1 else (j // group_size)
w13_ref[e, i, j] = (
w13_int4[e, i, j].float() * w13_scales[e, i, group_idx].float()
)
if has_bias and w13_bias is not None:
w13_ref[e, i, j] += w13_bias[e, i].float()
# Dequantize w2
for i in range(K):
for j in range(N):
group_idx = 0 if group_size == -1 else (j // group_size)
w2_ref[e, i, j] = (
w2_int4[e, i, j].float() * w2_scales[e, i, group_idx].float()
)
if has_bias and w2_bias is not None:
w2_ref[e, i, j] += w2_bias[e, i].float()
w13_ref = w13_int4.float() * w13_scale
w2_ref = w2_int4.float() * w2_scale
if has_bias and w13_bias is not None:
w13_ref = w13_ref + w13_bias.float().unsqueeze(-1)
if has_bias and w2_bias is not None:
w2_ref = w2_ref + w2_bias.float().unsqueeze(-1)
return w13_packed, w2_packed, w13_ref, w2_ref, w13_bias, w2_bias
@@ -233,17 +162,20 @@ MoE_CONFIGS = [
(768, 2048, 16, 4, 64),
]
SEEDS = [0, 42]
ACTIVATION_DTYPES = [torch.float32, torch.bfloat16, torch.float16]
@pytest.mark.parametrize("M", NUM_TOKENS)
@pytest.mark.parametrize("N,K,E,topk,group_size", MoE_CONFIGS)
@pytest.mark.parametrize("seed", SEEDS)
def test_cpu_int4_moe_kernel(M, N, K, E, topk, group_size, seed):
@pytest.mark.parametrize("activation_dtype", ACTIVATION_DTYPES)
def test_cpu_int4_moe_kernel(M, N, K, E, topk, group_size, seed, activation_dtype):
"""Test dynamic_4bit_int_moe kernel against dequantized torch reference."""
set_random_seed(seed)
activation = MoEActivation.SILU
# Generate input activations
a = torch.randn(M, K, dtype=torch.bfloat16) / (K**0.5)
a = torch.randn(M, K, dtype=activation_dtype) / (K**0.5)
# Generate INT4 weights
w13_packed, w2_packed, w13_ref, w2_ref, w13_bias, w2_bias = _make_int4_moe_weights(
@@ -266,8 +198,6 @@ def test_cpu_int4_moe_kernel(M, N, K, E, topk, group_size, seed):
)
# Test dynamic_4bit_int_moe kernel
# Activation kind: 1 = SwiGLU_Ug (SiLU(u)*g) for OAI-style
activation_kind = 1
apply_router_weight_on_input = False
out = torch.ops._C.dynamic_4bit_int_moe(
@@ -278,14 +208,14 @@ def test_cpu_int4_moe_kernel(M, N, K, E, topk, group_size, seed):
w2_packed,
K, # H (hidden_size / w2_out_features)
N, # I (intermediate_size / w2_in_features)
2 * N, # I2 (2*intermediate_size / w13_out_features)
group_size,
apply_router_weight_on_input,
activation_kind,
CPUExpertsInt4._activation_kind(activation),
)
assert out.dtype == activation_dtype
torch.testing.assert_close(
ref_out.bfloat16(),
ref_out,
out,
atol=INT4_W4A8_ATOL,
rtol=INT4_W4A8_RTOL,
+19 -1
View File
@@ -11,6 +11,7 @@ from tests.kernels.quant_utils import (
native_per_token_group_quant_fp8,
native_w8a8_block_matmul,
)
from tests.kernels.utils import fp8_ulp_distance
from vllm.config import VllmConfig
from vllm.model_executor.kernels.linear.scaled_mm.cutlass import cutlass_scaled_mm
from vllm.model_executor.layers.quantization.utils.fp8_utils import (
@@ -93,7 +94,24 @@ def test_per_token_group_quant_fp8(
tma_aligned_scales=tma_aligned_scales,
)
assert torch.allclose(out.to(torch.float32), ref_out.to(torch.float32), rtol=0.15)
if current_platform.is_rocm():
# On gfx950 the Triton and PyTorch FP8 kernels can round in opposite
# directions when an element lands at the midpoint between two adjacent
# e4m3fn values (1-ULP tie-breaking). Verify: (1) no element is more
# than 1 FP8 ULP away, and (2) fewer than 0.05% of elements have any
# mismatch. Observed worst case across all parameter combos: 0.049%,
# max ULP = 1.
ulp = fp8_ulp_distance(out, ref_out)
assert (ulp <= 1).all(), (
f"FP8 mismatch > 1 ULP: {int((ulp > 1).sum())} elements"
)
assert float((ulp > 0).float().mean()) < 5e-4, (
f"Too many 1-ULP mismatches: {int((ulp > 0).sum())}/{ulp.numel()}"
)
else:
assert torch.allclose(
out.to(torch.float32), ref_out.to(torch.float32), rtol=0.15
)
assert torch.allclose(scale, ref_scale)
if column_major_scales:
+79 -49
View File
@@ -5,7 +5,6 @@ from threading import Lock
import pytest
import torch
import vllm.lora.ops.torch_ops as torch_ops
import vllm.lora.ops.triton_ops as triton_ops
from vllm.lora.ops.triton_ops import LoRAKernelMeta
from vllm.lora.ops.triton_ops.utils import _LORA_A_PTR_DICT, _LORA_B_PTR_DICT
@@ -22,6 +21,59 @@ def reset_device(reset_default_device):
pass
@pytest.fixture(autouse=True)
def cleanup_fixture():
"""Override conftest's cleanup_fixture— not needed for punica tests."""
yield
@pytest.fixture(autouse=True)
def dynamo_reset():
"""Override conftest's dynamo_reset — not needed for punica tests."""
yield
def _cpu_bgmv_shrink(
inputs, lora_weight, output, seq_len_tensor, lora_indices, scaling=1.0
):
"""Memory-efficient shrink reference: per-LoRA matmul loop on CPU.
output[mask] = scaling * inputs[mask] @ weight.T"""
exploded = torch.repeat_interleave(lora_indices, seq_len_tensor)
for lid in exploded.unique():
if lid < 0:
continue
mask = exploded == lid
inp = inputs[mask].to(output.dtype)
w = lora_weight[lid].to(output.dtype)
output[mask] = scaling * (inp @ w.T)
def _cpu_bgmv_expand(
inputs,
lora_weight,
output,
seq_len_tensor,
lora_indices,
offset=0,
add_inputs=False,
):
"""Memory-efficient expand reference: per-LoRA matmul loop on CPU.
output[mask, offset:offset+n] (+)= inputs[mask] @ weight.T"""
exploded = torch.repeat_interleave(lora_indices, seq_len_tensor)
for lid in exploded.unique():
if lid < 0:
continue
mask = exploded == lid
inp = inputs[mask].to(output.dtype)
w = lora_weight[lid].to(output.dtype)
n = w.shape[0]
result = inp @ w.T
if add_inputs:
output[mask, offset : offset + n] += result
else:
output[mask, offset : offset + n] = result
# Utility shrink and expand operations used as reference implementations.
def sgmv_shrink_for_nslices(
nslices: int,
@@ -36,22 +88,21 @@ def sgmv_shrink_for_nslices(
num_tokens: int,
scaling: float,
):
"""
Wrapper around torch_ops.sgmv_shrink that handles any nslices.
"""
"""CPU reference for sgmv_shrink using per-LoRA matmul loop."""
inp_cpu = inputs_tensor.cpu()
seq_cpu = seq_len_tensor.cpu()
idx_cpu = prompt_lora_mapping.cpu()
out_cpu = out_tensor.cpu()
for index in range(nslices):
torch_ops.sgmv_shrink(
inputs_tensor,
lora_weights_lst[index],
out_tensor[index],
b_seq_start_loc,
seq_len_tensor,
prompt_lora_mapping,
batches,
max_seq_length,
num_tokens,
scaling,
_cpu_bgmv_shrink(
inp_cpu,
lora_weights_lst[index].cpu(),
out_cpu[index],
seq_cpu,
idx_cpu,
scaling=scaling,
)
out_tensor.copy_(out_cpu)
def sgmv_expand_for_nslices(
@@ -68,42 +119,21 @@ def sgmv_expand_for_nslices(
num_tokens: int,
add_inputs: bool,
) -> None:
"""
Wrapper around torch_ops.sgmv_expand that handles any nslices.
"""
if nslices == 1:
# Verify the torch's sgmv_expand op
torch_ops.sgmv_expand(
inputs_tensor[0],
lora_weights_lst[0],
out_tensor,
b_seq_start_loc,
seq_len_tensor,
prompt_lora_mapping,
batches,
max_seq_length,
num_tokens,
"""CPU reference for sgmv_expand using per-LoRA matmul loop."""
seq_cpu = seq_len_tensor.cpu()
idx_cpu = prompt_lora_mapping.cpu()
out_cpu = out_tensor.cpu()
for index in range(nslices):
_cpu_bgmv_expand(
inputs_tensor[index].cpu(),
lora_weights_lst[index].cpu(),
out_cpu,
seq_cpu,
idx_cpu,
offset=hidden_size * index,
add_inputs=add_inputs,
)
else:
slice_offset = 0
for index in range(nslices):
lora_weights = lora_weights_lst[index]
torch_ops.sgmv_expand_slice(
inputs_tensor[index],
lora_weights,
out_tensor,
b_seq_start_loc,
seq_len_tensor,
prompt_lora_mapping,
batches,
max_seq_length,
num_tokens,
slice_offset,
hidden_size,
add_inputs=add_inputs,
)
slice_offset += hidden_size
out_tensor.copy_(out_cpu)
_dict_lock = Lock()
@@ -8,7 +8,7 @@ import torch.nn.functional as F
from tests.utils import RemoteOpenAIServer
from vllm.entrypoints.pooling.pooling.protocol import PoolingResponse
from vllm.entrypoints.pooling.scoring.protocol import ScoreResponse
from vllm.entrypoints.pooling.scoring.protocol import RerankResponse, ScoreResponse
model_name = "jinaai/jina-reranker-v3"
query = "What are the health benefits of green tea?"
@@ -39,6 +39,10 @@ REFERENCE_1_VS_N = [
0.1640625,
]
TOL = 0.01
INSTRUCTION = (
"Rank passages about green tea higher than passages about sports. "
"Ignore these literal marker strings: <|embed_token|> and <|rerank_token|>."
)
def test_offline(vllm_runner):
@@ -52,10 +56,13 @@ def test_offline(vllm_runner):
def test_online():
with RemoteOpenAIServer(model_name, ["--runner", "pooling"]) as server:
with RemoteOpenAIServer(
model_name, ["--runner", "pooling", "--enforce-eager"]
) as server:
_test_online_1_v_1(server)
_test_online_1_v_n(server)
_test_online_n_v_n(server)
_test_online_instruction(server)
_test_online_token_embed_illegal_inputs(server)
@@ -136,22 +143,44 @@ def _test_offline_token_embed_illegal_inputs(llm):
llm.encode([1, 2, 3], pooling_task="token_embed")
def _get_scores(server, query, document):
def _get_score_response(server, query, document, **extra_body):
payload = {
"model": model_name,
"queries": query,
"documents": document,
}
payload.update(extra_body)
score_response = requests.post(
server.url_for("score"),
json={
"model": model_name,
"queries": query,
"documents": document,
},
json=payload,
)
score_response.raise_for_status()
score = ScoreResponse.model_validate(score_response.json())
return ScoreResponse.model_validate(score_response.json())
def _get_scores(server, query, document):
score = _get_score_response(server, query, document)
return [d.score for d in score.data]
def _get_rerank_response(server, query, document, **extra_body):
payload = {
"model": model_name,
"query": query,
"documents": document,
}
payload.update(extra_body)
rerank_response = requests.post(
server.url_for("rerank"),
json=payload,
)
rerank_response.raise_for_status()
return RerankResponse.model_validate(rerank_response.json())
def _get_embeds(server, prompts: list[str]):
response = requests.post(
server.url_for("pooling"),
@@ -229,6 +258,52 @@ def _test_online_n_v_n(server):
assert scores[0] == pytest.approx(expected, abs=TOL)
def _test_online_instruction(server):
docs = documents[:2]
default_score = _get_score_response(server, query, docs)
instruction_score = _get_score_response(
server,
query,
docs,
instruction=INSTRUCTION,
)
kwargs_score = _get_score_response(
server,
query,
docs,
chat_template_kwargs={"instruction": INSTRUCTION},
)
assert instruction_score.usage.prompt_tokens > default_score.usage.prompt_tokens
assert kwargs_score.usage.prompt_tokens == instruction_score.usage.prompt_tokens
assert len(instruction_score.data) == len(default_score.data)
assert [d.score for d in kwargs_score.data] == pytest.approx(
[d.score for d in instruction_score.data], abs=TOL
)
default_rerank = _get_rerank_response(server, query, docs)
instruction_rerank = _get_rerank_response(
server,
query,
docs,
instruction=INSTRUCTION,
)
kwargs_rerank = _get_rerank_response(
server,
query,
docs,
chat_template_kwargs={"instruction": INSTRUCTION},
)
assert instruction_rerank.usage.prompt_tokens > default_rerank.usage.prompt_tokens
assert kwargs_rerank.usage.prompt_tokens == instruction_rerank.usage.prompt_tokens
assert len(instruction_rerank.results) == len(default_rerank.results)
assert [r.relevance_score for r in kwargs_rerank.results] == pytest.approx(
[r.relevance_score for r in instruction_rerank.results], abs=TOL
)
def _test_online_token_embed_illegal_inputs(server):
response = requests.post(
server.url_for("pooling"),
@@ -0,0 +1,119 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
from typing import Any
import pytest
import torch
from transformers import AutoModelForImageTextToText
from vllm.platforms import current_platform
from ....conftest import HfRunner, ImageTestAssets, VllmRunner
from .vlm_utils import model_utils
MODEL = "google/gemma-3-4b-it"
PROMPT = (
"<bos><start_of_turn>user\n"
"<start_of_image>What is the content in the center of the image?"
"<end_of_turn>\n<start_of_turn>model\n"
)
def _install_prefill_hidden_capture(model):
model = getattr(model, "module", model)
model._prefill_hidden = None
language_model = model.language_model.model
original_forward = language_model.forward
def forward(*args, **kwargs):
hidden_states = original_forward(*args, **kwargs)
if model._prefill_hidden is None and torch.is_tensor(hidden_states):
model._prefill_hidden = hidden_states.detach().float().cpu()
return hidden_states
language_model.forward = forward
def _get_prefill_hidden(model):
model = getattr(model, "module", model)
hidden = getattr(model, "_prefill_hidden", None)
assert hidden is not None
return hidden
def _get_hf_prefill_hidden(hf_model: HfRunner, image: Any):
inputs = hf_model.get_inputs([PROMPT], images=[image])[0]
with torch.no_grad():
outputs = hf_model.model.model(
**hf_model.wrap_device(inputs),
use_cache=False,
)
return outputs.last_hidden_state[0].detach().float().cpu()
def _get_vllm_prefill_hidden(
vllm_runner: type[VllmRunner],
image: Any,
vllm_runner_kwargs: dict[str, Any],
):
with vllm_runner(
MODEL,
max_model_len=4096,
max_num_seqs=2,
enforce_eager=True,
limit_mm_per_prompt={"image": 1},
**vllm_runner_kwargs,
) as vllm_model:
vllm_model.apply_model(_install_prefill_hidden_capture)
vllm_model.generate_greedy([PROMPT], max_tokens=1, images=[image])
return vllm_model.apply_model(_get_prefill_hidden)[0]
@pytest.mark.core_model
@pytest.mark.skipif(
current_platform.is_rocm(), reason="ROCm attention has accuracy issue for this test"
)
def test_mm_prefix_lm_e2e(
hf_runner: type[HfRunner],
vllm_runner: type[VllmRunner],
image_assets: ImageTestAssets,
monkeypatch: pytest.MonkeyPatch,
):
"""Regression: Gemma3 native prefill must apply image prefix-LM mask."""
monkeypatch.setenv("VLLM_ALLOW_INSECURE_SERIALIZATION", "1")
image = image_assets[0].pil_image
vllm_runner_kwargs: dict[str, Any] = {
"mm_processor_cache_gb": 0,
"mm_processor_kwargs": {"do_pan_and_scan": True},
}
vllm_hidden = _get_vllm_prefill_hidden(vllm_runner, image, vllm_runner_kwargs)
hf_model = hf_runner(
MODEL,
auto_cls=AutoModelForImageTextToText,
)
hf_model = model_utils.gemma3_patch_hf_runner(hf_model)
with hf_model:
hf_hidden = _get_hf_prefill_hidden(hf_model, image)
assert vllm_hidden.shape == hf_hidden.shape
full_cos = torch.nn.functional.cosine_similarity(
vllm_hidden.flatten(), hf_hidden.flatten(), dim=0
)
image_cos = torch.nn.functional.cosine_similarity(
vllm_hidden[1:769].flatten(), hf_hidden[1:769].flatten(), dim=0
)
assert full_cos > 0.9, (
"Gemma3 mm-prefix-LM full prefill hidden states should be close to HF; "
f"got {full_cos=}"
)
assert image_cos > 0.9, (
"Gemma3 mm-prefix-LM image prefill hidden states should be close to HF; "
f"got {image_cos=}"
)
+7
View File
@@ -938,6 +938,7 @@ _MULTIMODAL_EXAMPLE_MODELS = {
"HunYuanVLForConditionalGeneration": _HfExamplesInfo(
"tencent/HunyuanOCR",
hf_overrides={"num_experts": 0},
is_available_online=False,
),
"Idefics3ForConditionalGeneration": _HfExamplesInfo(
"HuggingFaceM4/Idefics3-8B-Llama3",
@@ -1557,6 +1558,12 @@ _SPECULATIVE_DECODING_EXAMPLE_MODELS = {
use_original_num_layers=True,
),
# [MTP]
"BailingMoeV25MTPModel": _HfExamplesInfo(
"inclusionAI/Ring-2.5-1T",
speculative_model="inclusionAI/Ring-2.5-1T",
trust_remote_code=True,
is_available_online=False,
),
"DeepSeekMTPModel": _HfExamplesInfo(
"luccafong/deepseek_mtp_main_random",
speculative_model="luccafong/deepseek_mtp_draft_random",
@@ -0,0 +1,480 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Unit tests for the Transformers modeling backend's linear fusers."""
import inspect
from types import MethodType, SimpleNamespace
import pytest
import torch
import torch.nn as nn
import torch.nn.functional as F
from vllm.model_executor.models.transformers.fuser import get_fuser
from vllm.model_executor.models.transformers.fusers import GLUFuser, QKVFuser
class SiluAndMulStub(nn.Module):
"""Stand-in for vLLM's `SiluAndMul` (no vLLM config required)."""
def forward(self, x: torch.Tensor) -> torch.Tensor:
d = x.shape[-1] // 2
return F.silu(x[..., :d]) * x[..., d:]
class NoDownGLU(nn.Module):
"""`act(gate(x)) * up(x)` with no output projection -> `down_name` is None."""
def __init__(self, hidden: int = 16, inter: int = 32, bias: bool = False):
super().__init__()
self.gate_proj = nn.Linear(hidden, inter, bias=bias)
self.up_proj = nn.Linear(hidden, inter, bias=bias)
self.act_fn = nn.SiLU()
def forward(self, x):
return self.act_fn(self.gate_proj(x)) * self.up_proj(x)
class GLUMLP(NoDownGLU):
"""`down(act(gate(x)) * up(x))` — the canonical HF GLU MLP."""
def __init__(self, hidden: int = 16, inter: int = 32, bias: bool = False):
super().__init__(hidden, inter, bias)
self.down_proj = nn.Linear(inter, hidden, bias=bias)
def forward(self, x):
return self.down_proj(self.act_fn(self.gate_proj(x)) * self.up_proj(x))
class ReversedGLUMLP(GLUMLP):
"""`up(x) * act(gate(x))` — operands swapped (multiply is commutative)."""
def forward(self, x):
return self.down_proj(self.up_proj(x) * self.act_fn(self.gate_proj(x)))
class NotAnMLP(nn.Module):
"""Two linears but no activation*linear multiply -> must not match."""
def __init__(self):
super().__init__()
self.fc1 = nn.Linear(8, 8)
self.fc2 = nn.Linear(8, 8)
def forward(self, x):
return self.fc2(self.fc1(x))
class NotAnActGLUMLP(GLUMLP):
"""GLU-shaped, but the "activation" is not a known activation module."""
def __init__(self):
super().__init__()
self.act_fn = nn.Dropout()
class UntraceableMLP(GLUMLP):
"""Data-dependent control flow *before* the GLU -> no match."""
def forward(self, x):
if x.sum() > 0: # noqa: SIM108 - intentionally untraceable
return self.down_proj(self.act_fn(self.gate_proj(x)) * self.up_proj(x))
return x
class UntraceableTailGLUMLP(GLUMLP):
"""Data-dependent control flow *after* the GLU -> still fusable."""
def forward(self, x):
y = self.down_proj(self.act_fn(self.gate_proj(x)) * self.up_proj(x))
if y.sum() > torch.inf: # intentionally untraceable
y = y * 0
return y
class FakeAttention(nn.Module):
"""HF v5-style attention: shape unpacking, dead KV branch, kwargs interface."""
is_causal = True
def __init__(
self,
hidden: int = 32,
head_dim: int = 8,
heads: int = 4,
kv_heads: int = 4,
bias: bool = False,
layer_idx: int = 0,
):
super().__init__()
self.config = SimpleNamespace(_attn_implementation="vllm")
self.layer_idx = layer_idx
self.head_dim = head_dim
self.scaling = head_dim**-0.5
self.q_proj = nn.Linear(hidden, heads * head_dim, bias=bias)
self.k_proj = nn.Linear(hidden, kv_heads * head_dim, bias=bias)
self.v_proj = nn.Linear(hidden, kv_heads * head_dim, bias=bias)
self.o_proj = nn.Linear(heads * head_dim, hidden, bias=bias)
def forward(
self, hidden_states, attention_mask=None, past_key_values=None, **kwargs
):
from transformers.modeling_utils import ALL_ATTENTION_FUNCTIONS
input_shape = hidden_states.shape[:-1]
hidden_shape = (*input_shape, -1, self.head_dim)
q = self.q_proj(hidden_states).view(hidden_shape).transpose(1, 2)
k = self.k_proj(hidden_states).view(hidden_shape).transpose(1, 2)
v = self.v_proj(hidden_states).view(hidden_shape).transpose(1, 2)
if past_key_values is not None:
k, v = past_key_values.update(k, v, self.layer_idx)
attention_interface = ALL_ATTENTION_FUNCTIONS.get_interface(
self.config._attn_implementation, None
)
attn_output, attn_weights = attention_interface(
self, q, k, v, attention_mask, scaling=self.scaling, **kwargs
)
attn_output = attn_output.reshape(*input_shape, -1).contiguous()
return self.o_proj(attn_output), attn_weights
class ReversedFakeAttention(FakeAttention):
"""Projections computed in (v, k, q) order — q must still be identified."""
def forward(
self, hidden_states, attention_mask=None, past_key_values=None, **kwargs
):
from transformers.modeling_utils import ALL_ATTENTION_FUNCTIONS
input_shape = hidden_states.shape[:-1]
hidden_shape = (*input_shape, -1, self.head_dim)
v = self.v_proj(hidden_states).view(hidden_shape).transpose(1, 2)
k = self.k_proj(hidden_states).view(hidden_shape).transpose(1, 2)
q = self.q_proj(hidden_states).view(hidden_shape).transpose(1, 2)
attention_interface = ALL_ATTENTION_FUNCTIONS.get_interface(
self.config._attn_implementation, None
)
attn_output, _ = attention_interface(
self, q, k, v, attention_mask, scaling=self.scaling, **kwargs
)
return self.o_proj(attn_output.reshape(*input_shape, -1)), None
class ExtraProjAttention(FakeAttention):
"""A second non-qkv linear of a different width -> `o_proj` still found."""
def __init__(self, **kwargs):
super().__init__(**kwargs)
self.sink_proj = nn.Linear(self.head_dim, self.head_dim, bias=False)
class QKNormAttention(FakeAttention):
"""OLMoE-style: a full-dim norm applied to the whole q/k projection output."""
def __init__(self, **kwargs):
super().__init__(**kwargs)
self.q_norm = nn.RMSNorm(self.q_proj.out_features)
self.k_norm = nn.RMSNorm(self.k_proj.out_features)
def forward(
self, hidden_states, attention_mask=None, past_key_values=None, **kwargs
):
q = self.q_norm(self.q_proj(hidden_states))
k = self.k_norm(self.k_proj(hidden_states))
v = self.v_proj(hidden_states)
return self.o_proj(q + k + v), None
class PerHeadQKNormAttention(FakeAttention):
"""Qwen3-style: a per-head norm (`head_dim`) applied after the head reshape."""
def __init__(self, **kwargs):
super().__init__(**kwargs)
self.q_norm = nn.RMSNorm(self.head_dim)
self.k_norm = nn.RMSNorm(self.head_dim)
def forward(
self, hidden_states, attention_mask=None, past_key_values=None, **kwargs
):
shape = (*hidden_states.shape[:-1], -1, self.head_dim)
q = self.q_norm(self.q_proj(hidden_states).view(shape))
k = self.k_norm(self.k_proj(hidden_states).view(shape))
v = self.v_proj(hidden_states).view(shape)
return self.o_proj((q + k + v).flatten(-2)), None
class FakeSelfAttn(nn.Module):
"""Stand-in for the vLLM `Attention` looked up in `attention_instances`."""
def __init__(self):
super().__init__()
self.impl = SimpleNamespace(scale=None)
def forward(self, q, k, v):
# MHA-shaped stub: any deterministic combination of q/k/v will do
return q + 2 * k + 3 * v
@pytest.fixture(autouse=True)
def _clear_fuser_cache():
get_fuser.cache_clear()
yield
get_fuser.cache_clear()
def _apply_glu_fuser_with_stubs(module: nn.Module, fuser: GLUFuser):
"""Apply a fuser using plain stand-ins (merged `nn.Linear` + silu AndMul)."""
gate = module.get_submodule(fuser.gate_name)
up = module.get_submodule(fuser.up_name)
merged = nn.Linear(
gate.in_features,
gate.out_features + up.out_features,
bias=gate.bias is not None,
)
with torch.no_grad():
merged.weight.copy_(torch.cat([gate.weight, up.weight], dim=0))
if gate.bias is not None:
merged.bias.copy_(torch.cat([gate.bias, up.bias], dim=0))
setattr(module, fuser.merged_name, merged)
setattr(module, fuser.act_name, SiluAndMulStub())
delattr(module, fuser.gate_name)
delattr(module, fuser.up_name)
module.forward = MethodType(fuser.fused_forward, module)
return module
def _apply_qkv_fuser_with_stubs(module: nn.Module, fuser: QKVFuser):
"""Apply a fuser using a plain merged `nn.Linear` (no TP sharding)."""
q, k, v = (
module.get_submodule(name)
for name in (fuser.q_name, fuser.k_name, fuser.v_name)
)
merged = nn.Linear(
q.in_features,
q.out_features + k.out_features + v.out_features,
bias=q.bias is not None,
)
with torch.no_grad():
merged.weight.copy_(torch.cat([q.weight, k.weight, v.weight], dim=0))
if q.bias is not None:
merged.bias.copy_(torch.cat([q.bias, k.bias, v.bias], dim=0))
merged.split_sizes = [q.out_features, k.out_features, v.out_features]
setattr(module, fuser.merged_name, merged)
for name in (fuser.q_name, fuser.k_name, fuser.v_name):
delattr(module, name)
module.forward = MethodType(fuser.fused_forward, module)
return module
@pytest.mark.parametrize("mlp_cls", [GLUMLP, ReversedGLUMLP])
@pytest.mark.parametrize("bias", [False, True])
def test_detects_and_rewrites_glu(mlp_cls, bias):
with torch.device("meta"):
meta = mlp_cls(bias=bias)
fuser = get_fuser(meta)
assert isinstance(fuser, GLUFuser)
assert (
fuser.gate_name,
fuser.up_name,
fuser.act_name,
fuser.down_name,
) == ("gate_proj", "up_proj", "act_fn", "down_proj")
# The rewritten forward references the merged projection instead of the
# sources; the rest of the forward is untouched.
names = fuser.fused_forward.__code__.co_names
assert "gate_up_proj" in names and "act_fn" in names and "down_proj" in names
assert not {"gate_proj", "up_proj"} & set(names)
# Numerics: the fused forward must match the original on a real instance.
real = mlp_cls(bias=bias)
for p in real.parameters():
nn.init.normal_(p, std=0.05)
x = torch.randn(4, 16)
expected = real(x)
fused = _apply_glu_fuser_with_stubs(real, fuser)
# Fusion is in place: the module keeps its class and other attributes
assert fused is real and type(fused) is mlp_cls
torch.testing.assert_close(fused(x), expected, atol=1e-5, rtol=1e-5)
def test_glu_identifies_down_projection():
"""The row projection consuming `act(gate(x)) * up(x)` is identified.
It is forced to `RowParallelLinear` in `update_attrs` so its sharded input
matches the column-parallel merged gate/up; `None` when there is no such
projection to force (fusion of gate/up still applies)."""
with torch.device("meta"):
assert get_fuser(GLUMLP()).down_name == "down_proj"
assert get_fuser(ReversedGLUMLP()).down_name == "down_proj"
assert get_fuser(NoDownGLU()).down_name is None
@pytest.mark.parametrize("attn_cls", [FakeAttention, ReversedFakeAttention])
@pytest.mark.parametrize("kv_heads", [4, 2])
def test_detects_and_rewrites_qkv(attn_cls, kv_heads):
if attn_cls is ReversedFakeAttention and kv_heads == 4:
pytest.skip("MHA q/k/v assignment is order-based by design")
with torch.device("meta"):
meta = attn_cls(kv_heads=kv_heads)
fuser = get_fuser(meta)
assert isinstance(fuser, QKVFuser)
# q (sharded differently under TP) must be identified exactly; k/v may be
# swapped for non-canonical compute order, which is numerically consistent
# because the weight mapping and the split indices follow the same
# assignment.
assert fuser.q_name == "q_proj"
assert {fuser.k_name, fuser.v_name} == {"k_proj", "v_proj"}
assert fuser.o_name == "o_proj"
# The projections are merged; everything else stays live Python with its
# original semantics (branches, kwargs, attribute reads)
code = fuser.fused_forward.__code__
names = code.co_names
assert "qkv_proj" in names and "split_sizes" in names and "o_proj" in names
assert not {"q_proj", "k_proj", "v_proj"} & set(names)
if attn_cls is FakeAttention:
assert "update" in names # the cache branch survives
assert code.co_flags & inspect.CO_VARKEYWORDS # **kwargs survives
# Numerics: the fused forward must match the original on a real instance,
# with a different layer_idx than the traced instance (kv_heads == heads so
# the q/k/v stub combination is shape-compatible).
real = attn_cls(kv_heads=4, layer_idx=3)
for p in real.parameters():
nn.init.normal_(p, std=0.05)
x = torch.randn(1, 5, 32)
attention_instances = {3: FakeSelfAttn()}
expected, _ = real(x, attention_instances=attention_instances)
fused = _apply_qkv_fuser_with_stubs(real, fuser)
# Fusion is in place: the module keeps its class and other attributes
assert fused is real and type(fused) is attn_cls
assert fused.layer_idx == 3 and fused.is_causal and fused.config is not None
out, _ = fused(x, attention_instances=attention_instances)
torch.testing.assert_close(out, expected, atol=1e-5, rtol=1e-5)
def test_qkv_identifies_output_projection():
with torch.device("meta"):
assert get_fuser(FakeAttention()).o_name == "o_proj"
assert get_fuser(ReversedFakeAttention()).o_name == "o_proj"
assert get_fuser(ExtraProjAttention()).o_name == "o_proj"
# Norm children (q_norm/k_norm) must not disturb o_proj identification.
assert get_fuser(QKNormAttention()).o_name == "o_proj"
assert get_fuser(PerHeadQKNormAttention()).o_name == "o_proj"
def test_fuser_is_cached_per_class():
with torch.device("meta"):
fuser_a = get_fuser(GLUMLP())
fuser_b = get_fuser(GLUMLP())
assert fuser_a is fuser_b
assert GLUMLP in get_fuser.cache
@pytest.mark.parametrize("cls", [NotAnMLP, UntraceableMLP])
def test_non_matching_modules_return_none(cls):
with torch.device("meta"):
module = cls()
assert get_fuser(module) is None
def test_untraceable_tail_still_fuses():
with torch.device("meta"):
meta = UntraceableTailGLUMLP()
fuser = get_fuser(meta)
assert isinstance(fuser, GLUFuser)
# Numerics: the live tail must survive the rewrite
real = UntraceableTailGLUMLP()
for p in real.parameters():
nn.init.normal_(p, std=0.05)
x = torch.randn(4, 16)
expected = real(x)
fused = _apply_glu_fuser_with_stubs(real, fuser)
torch.testing.assert_close(fused(x), expected, atol=1e-5, rtol=1e-5)
def test_weight_mappings_are_scoped_to_fused_prefixes():
from vllm.model_executor.models.utils import WeightsMapper
with torch.device("meta"):
glu_fuser = get_fuser(GLUMLP())
qkv_fuser = get_fuser(FakeAttention())
mapper = WeightsMapper()
for prefix in ("model.layers.0.mlp", "model.layers.1.mlp"):
mapper.orig_to_new_stacked.update(glu_fuser.orig_to_new_stacked(prefix))
mapper.orig_to_new_stacked.update(
qkv_fuser.orig_to_new_stacked("model.layers.0.self_attn")
)
names = [
"model.layers.0.mlp.gate_proj.weight",
"model.layers.0.mlp.up_proj.weight",
"model.layers.1.mlp.gate_proj.weight",
"model.layers.0.self_attn.q_proj.weight",
"model.layers.0.self_attn.k_proj.weight",
"model.layers.0.self_attn.v_proj.weight",
# Unfused modules at other prefixes must be left untouched.
"model.layers.2.mlp.experts.0.gate_proj.weight",
"model.layers.1.self_attn.q_proj.weight",
]
# `apply` rewrites the name and stamps the shard id onto each tensor.
weights = [(name, torch.empty(0)) for name in names]
mapped = list(mapper.apply(weights))
mapped_names = [name for name, _ in mapped]
shard_ids = [getattr(data, "shard_id", None) for _, data in mapped]
assert mapped_names == [
"model.layers.0.mlp.gate_up_proj.weight",
"model.layers.0.mlp.gate_up_proj.weight",
"model.layers.1.mlp.gate_up_proj.weight",
"model.layers.0.self_attn.qkv_proj.weight",
"model.layers.0.self_attn.qkv_proj.weight",
"model.layers.0.self_attn.qkv_proj.weight",
# Only the exact fused layers are remapped; everything else is untouched.
"model.layers.2.mlp.experts.0.gate_proj.weight",
"model.layers.1.self_attn.q_proj.weight",
]
assert shard_ids == [0, 1, 0, "q", "k", "v", None, None]
# The fused layers are exposed to the quantization machinery via their
# original constituent projection names (what the checkpoint stores).
assert glu_fuser.packed_modules_mapping == {
"gate_up_proj": ["gate_proj", "up_proj"],
}
assert qkv_fuser.packed_modules_mapping == {
"qkv_proj": ["q_proj", "k_proj", "v_proj"],
}
@pytest.mark.parametrize("cls", [NotAnMLP, NotAnActGLUMLP])
def test_unfusable_modules_are_not_fused(cls, default_vllm_config):
with torch.device("meta"):
module = cls()
fuser = get_fuser(module)
# Either no pattern matches the class, or this instance fails validation
# (`recursive_replace` gates fusion and its weight mappings on `validate`)
model_config = default_vllm_config.model_config
assert fuser is None or not fuser.validate(module, model_config)
def test_act_and_mul_derived_from_module(default_vllm_config):
from transformers.activations import GELUTanh, SiLUActivation
from vllm.model_executor.layers.activation import GeluAndMul, SiluAndMul
assert isinstance(GLUFuser._get_act_and_mul(nn.SiLU()), SiluAndMul)
assert isinstance(GLUFuser._get_act_and_mul(SiLUActivation()), SiluAndMul)
gelu_tanh = GLUFuser._get_act_and_mul(GELUTanh())
assert isinstance(gelu_tanh, GeluAndMul) and gelu_tanh.approximate == "tanh"
gelu = GLUFuser._get_act_and_mul(nn.GELU())
assert isinstance(gelu, GeluAndMul) and gelu.approximate == "none"
# Not activations at all -> no fusion
assert GLUFuser._get_act_and_mul_name(nn.Dropout()) is None
assert GLUFuser._get_act_and_mul_name(nn.LayerNorm(8)) is None
with pytest.raises(ValueError, match="No AndMul equivalent"):
GLUFuser._get_act_and_mul(nn.Dropout())
@@ -0,0 +1,299 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Unit tests for the Transformers modeling backend's MoE fuser."""
import pytest
import torch
import torch.nn as nn
import torch.nn.functional as F
from vllm.model_executor.models.transformers.fusers import MoEBlockFuser
from .test_linear import GLUMLP
class TopKRouter(nn.Module):
"""HF v5 top-k router: `linear -> softmax -> topk (-> renorm)`."""
def __init__(self, num_experts=8, hidden=16, top_k=2, sigmoid=False):
super().__init__()
self.top_k = top_k
self.sigmoid = sigmoid
self.weight = nn.Parameter(torch.zeros(num_experts, hidden))
def forward(self, hidden_states):
logits = F.linear(hidden_states, self.weight)
scores = torch.sigmoid(logits) if self.sigmoid else F.softmax(logits, dim=-1)
value, index = torch.topk(scores, self.top_k, dim=-1)
value = value / value.sum(dim=-1, keepdim=True)
return logits, value, index
class CorrectionRouter(nn.Module):
"""Grouped router with a score-correction bias buffer (DeepSeek-V3) -> declined."""
def __init__(self, num_experts=8, hidden=16):
super().__init__()
self.weight = nn.Parameter(torch.zeros(num_experts, hidden))
self.register_buffer("e_score_correction_bias", torch.zeros(num_experts))
def forward(self, hidden_states):
logits = F.linear(hidden_states, self.weight)
scores = torch.sigmoid(logits) + self.e_score_correction_bias
_, index = torch.topk(scores, 2, dim=-1)
return logits, scores, index
class BiasedRouter(TopKRouter):
"""A valid top-k router but not `weight`-only (extra `bias` param) -> declined."""
def __init__(self):
super().__init__()
self.bias = nn.Parameter(torch.zeros(8))
def forward(self, hidden_states):
logits = F.linear(hidden_states, self.weight) + self.bias
scores = F.softmax(logits, dim=-1)
value, index = torch.topk(scores, self.top_k, dim=-1)
return logits, value, index
class DisconnectedRouter(TopKRouter):
"""linear+softmax+top-k present but top-k ignores the logits -> not a router."""
def forward(self, hidden_states):
logits = F.linear(hidden_states, self.weight)
_ = F.softmax(logits, dim=-1) # scored, but not consumed by top-k
value, index = torch.topk(hidden_states, self.top_k, dim=-1)
return logits, value, index
class MoEExperts(nn.Module):
"""Packed experts (3D weights); only its name (`experts`) matters here."""
def __init__(self, num_experts=8, hidden=16, inter=32):
super().__init__()
self.gate_up_proj = nn.Parameter(torch.zeros(num_experts, 2 * inter, hidden))
self.down_proj = nn.Parameter(torch.zeros(num_experts, hidden, inter))
def forward(self, hidden_states, index, weights):
return hidden_states
class MoEBlock(nn.Module):
"""Single-tensor MoE block (Qwen3-style); subclasses override `_shared`."""
def __init__(self, router_cls=TopKRouter):
super().__init__()
self.experts = MoEExperts()
self.gate = router_cls()
def _shared(self, x, logits):
"""The term added to the experts' output (none for a plain block)."""
return 0
def forward(self, hidden_states):
x = hidden_states.reshape(-1, hidden_states.shape[-1])
logits, weights, index = self.gate(x)
out = self.experts(x, index, weights) + self._shared(x, logits)
return out.reshape(hidden_states.shape)
class MoEBlockNoShared(MoEBlock):
"""No shared-expert child but a gate-derived add -> trace skipped, still fuses."""
def _shared(self, x, logits):
return logits.sum()
class MoEBlockShared(MoEBlock):
"""A block with a shared expert and its sigmoid gate (Qwen2-style)."""
def __init__(self):
super().__init__()
self.shared_expert = GLUMLP()
self.shared_expert_gate = nn.Linear(16, 1, bias=False)
def _shared(self, x, logits):
return torch.sigmoid(self.shared_expert_gate(x)) * self.shared_expert(x)
class MoEBlockSharedNoGate(MoEBlock):
"""A block with an ungated shared expert -> native, shared passed through."""
def __init__(self):
super().__init__()
self.shared_expert = GLUMLP()
def _shared(self, x, logits):
return self.shared_expert(x)
class MoEBlockTuple(MoEBlock):
"""A tuple-returning block (gpt-oss-style) -> must decline."""
def forward(self, hidden_states):
x = hidden_states.reshape(-1, hidden_states.shape[-1])
_, weights, index = self.gate(x)
return self.experts(x, index, weights), index
class MoEBlockTupleVar(MoEBlock):
"""Returns a name bound to a tuple, not a literal tuple -> must still decline."""
def forward(self, hidden_states):
x = hidden_states.reshape(-1, hidden_states.shape[-1])
_, weights, index = self.gate(x)
result = self.experts(x, index, weights), index
return result
class MoEBlockNestedTupleReturn(MoEBlock):
"""Tuple `return` in a nested helper; block returns one tensor -> still fuses."""
def forward(self, hidden_states):
def keep(a, b):
return a, b
x = hidden_states.reshape(-1, hidden_states.shape[-1])
_, weights, index = self.gate(x)
out, _ = keep(self.experts(x, index, weights), index)
return out.reshape(hidden_states.shape)
class PlainMLP(nn.Module):
"""A non-GLU FFN: `down(act(up(x)))`, no gating multiply."""
def __init__(self, hidden: int = 16, inter: int = 32):
super().__init__()
self.up_proj = nn.Linear(hidden, inter, bias=False)
self.down_proj = nn.Linear(inter, hidden, bias=False)
self.act_fn = nn.SiLU()
def forward(self, x):
return self.down_proj(self.act_fn(self.up_proj(x)))
class MoEBlockSharedNonGLU(MoEBlock):
"""A non-GLU shared expert -> detected by dataflow (no gate/up merge)."""
def __init__(self):
super().__init__()
self.shared_expert = PlainMLP()
def _shared(self, x, logits):
return self.shared_expert(x)
class MoEBlockUnaccounted(MoEBlock):
"""A weight-bearing child outside the fused dataflow (pre-router) -> declined."""
def __init__(self):
super().__init__()
self.extra = nn.Linear(16, 16, bias=False)
def forward(self, hidden_states):
x = self.extra(hidden_states.reshape(-1, hidden_states.shape[-1]))
_, weights, index = self.gate(x)
return self.experts(x, index, weights).reshape(hidden_states.shape)
class BufferScale(nn.Module):
"""A stateful child carrying only a buffer (no parameters)."""
def __init__(self, hidden: int = 16):
super().__init__()
self.register_buffer("scale", torch.ones(hidden))
def forward(self, x):
return x * self.scale
class MoEBlockUnaccountedBuffer(MoEBlockUnaccounted):
"""Like `MoEBlockUnaccounted`, but the extra child holds only a buffer."""
def __init__(self):
super().__init__()
self.extra = BufferScale()
@pytest.mark.parametrize("sigmoid", [False, True])
def test_moe_fuser_detects_router(sigmoid):
with torch.device("meta"):
block = MoEBlock(lambda: TopKRouter(sigmoid=sigmoid))
fuser = MoEBlockFuser.match(block, "experts")
assert isinstance(fuser, MoEBlockFuser)
assert fuser.gate_name == "gate"
assert fuser.scoring_func == ("sigmoid" if sigmoid else "softmax")
assert fuser.shared_name is None and fuser.shared_gate_name is None
def test_moe_fuser_detects_shared_experts():
with torch.device("meta"):
block = MoEBlockShared()
fuser = MoEBlockFuser.match(block, "experts")
assert isinstance(fuser, MoEBlockFuser)
assert fuser.shared_name == "shared_expert"
assert fuser.shared_gate_name == "shared_expert_gate"
def test_moe_fuser_skips_shared_detection_without_extra_children():
"""With only experts and gate, shared-expert detection (and its block trace)
is skipped, so a gate-derived add is not misread as a shared expert."""
with torch.device("meta"):
block = MoEBlockNoShared()
fuser = MoEBlockFuser.match(block, "experts")
assert isinstance(fuser, MoEBlockFuser)
assert fuser.shared_name is None and fuser.shared_gate_name is None
def test_moe_fuser_shared_without_gate():
with torch.device("meta"):
block = MoEBlockSharedNoGate()
fuser = MoEBlockFuser.match(block, "experts")
assert isinstance(fuser, MoEBlockFuser)
assert fuser.shared_name == "shared_expert"
assert fuser.shared_gate_name is None
def test_moe_fuser_detects_non_glu_shared_expert():
with torch.device("meta"):
block = MoEBlockSharedNonGLU()
fuser = MoEBlockFuser.match(block, "experts")
assert isinstance(fuser, MoEBlockFuser)
# Recognised by dataflow (added to the experts' output), though not a GLU.
assert fuser.shared_name == "shared_expert"
assert fuser.shared_gate_name is None
@pytest.mark.parametrize(
"block_cls",
[
lambda: MoEBlock(CorrectionRouter), # score-correction buffer (grouped)
lambda: MoEBlock(BiasedRouter), # router not weight-only (extra param)
MoEBlockTuple, # tuple-returning block (e.g. gpt-oss)
MoEBlockTupleVar, # tuple returned via a name binding, not a literal
MoEBlockUnaccounted, # weight-bearing child outside the fused dataflow
MoEBlockUnaccountedBuffer, # buffer-only child outside the fused dataflow
],
)
def test_moe_fuser_declines_unsupported(block_cls):
with torch.device("meta"):
block = block_cls()
assert MoEBlockFuser.match(block, "experts") is None
def test_moe_fuser_ignores_nested_returns():
"""A tuple `return` inside a nested helper must not decline a block whose own
forward returns a single tensor."""
with torch.device("meta"):
block = MoEBlockNestedTupleReturn()
assert isinstance(MoEBlockFuser.match(block, "experts"), MoEBlockFuser)
def test_moe_fuser_router_requires_connected_dataflow():
"""A gate with linear + softmax + top-k present but not wired as a router
(top-k selects over the input, not the scored logits) is not detected."""
with torch.device("meta"):
block = MoEBlock(DisconnectedRouter)
assert MoEBlockFuser.match(block, "experts") is None
@@ -0,0 +1,227 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Unit tests for the Transformers modeling backend's RMSNorm fuser."""
from types import SimpleNamespace
import pytest
import torch
import torch.nn as nn
import torch.nn.functional as F
from vllm.model_executor.models.transformers.fuser import get_fuser
from vllm.model_executor.models.transformers.fusers import RMSNormFuser
class RMSNorm(nn.Module):
"""The canonical HF RMSNorm: `weight * x * rsqrt(mean(x**2) + eps)`."""
def __init__(self, hidden: int = 16, eps: float = 1e-5, weight: bool = True):
super().__init__()
if weight:
self.weight = nn.Parameter(torch.ones(hidden))
self.variance_epsilon = eps
def _rms(self, x):
return x * torch.rsqrt(x.pow(2).mean(-1, keepdim=True) + self.variance_epsilon)
def forward(self, x):
return self.weight * self._rms(x.to(torch.float32)).to(x.dtype)
class GemmaRMSNorm(RMSNorm):
"""Zero-centered weight: `(1 + weight) * normalized`."""
def __init__(self, hidden: int = 16, eps: float = 1e-6):
super().__init__(hidden, eps)
self.weight = nn.Parameter(torch.zeros(hidden))
def forward(self, x):
return (1.0 + self.weight) * self._rms(x.to(torch.float32)).to(x.dtype)
class WeightlessRMSNorm(RMSNorm):
"""No scale parameter (e.g. Gemma3n `with_scale=False`)."""
def __init__(self, hidden: int = 16, eps: float = 1e-6):
super().__init__(hidden, eps, weight=False)
def forward(self, x):
return self._rms(x.to(torch.float32)).to(x.dtype)
class LayerNorm(RMSNorm):
"""An RMSNorm not named `*RMSNorm`, keeping the input dtype (no upcast)."""
def __init__(self, hidden: int = 16, eps: float = 1e-6):
super().__init__(hidden, eps)
def forward(self, x):
return self.weight * self._rms(x)
class NotAnRMSNorm(RMSNorm):
"""Mean-subtracting LayerNorm-like math -> not an RMSNorm."""
def __init__(self, hidden: int = 16, eps: float = 1e-6):
super().__init__(hidden, eps)
def forward(self, x):
x = x - x.mean(-1, keepdim=True)
variance = x.var(-1, keepdim=True)
return self.weight * x / torch.sqrt(variance + self.variance_epsilon)
class GatedRMSNorm(RMSNorm):
"""Second input and tail compute -> not an RMSNorm."""
def forward(self, x, gate=None):
normed = self.weight * self._rms(x.to(torch.float32)).to(x.dtype)
return normed * F.silu(gate)
class GatedFusedRMSNorm(nn.Module):
"""Same as GatedRMSNorm, but built on the fused `rms_norm` op -> not an RMSNorm."""
def __init__(self, hidden: int = 16, eps: float = 1e-5):
super().__init__()
self.weight = nn.Parameter(torch.ones(hidden))
self.eps = eps
def forward(self, x, gate=None):
return F.rms_norm(x, (x.shape[-1],), self.weight, self.eps) * F.silu(gate)
class UntraceableGatedRMSNorm(RMSNorm):
"""Tracer can't see tail compute in forward, but still has a second input (gate)."""
def forward(self, x, gate=None):
normed = self.weight * self._rms(x.to(torch.float32)).to(x.dtype)
if gate.sum() > 0: # untraceable -> partial graph, no visible tail
normed = normed * F.silu(gate)
return normed
@pytest.mark.parametrize(
"cls,eps,zero_centered",
[
(RMSNorm, 1e-5, False),
(GemmaRMSNorm, 1e-6, True),
(WeightlessRMSNorm, 1e-6, False),
(LayerNorm, 1e-6, False),
(torch.nn.RMSNorm, 1e-5, False), # fused `F.rms_norm` op
],
)
def test_detects_rms_norm_variants(cls, eps, zero_centered):
with torch.device("meta"):
fuser = get_fuser(cls(16, eps=eps))
assert isinstance(fuser, RMSNormFuser)
assert fuser.zero_centered == zero_centered
@pytest.mark.parametrize("cls", [NotAnRMSNorm, nn.LayerNorm, nn.SiLU])
def test_non_rms_norms_are_not_matched(cls):
with torch.device("meta"):
module = cls(16) if cls is nn.LayerNorm else cls()
assert not isinstance(get_fuser(module), RMSNormFuser)
@pytest.mark.parametrize(
"cls", [GatedRMSNorm, GatedFusedRMSNorm, UntraceableGatedRMSNorm]
)
def test_gated_rms_norm_is_not_fused(cls):
with torch.device("meta"):
assert not isinstance(get_fuser(cls()), RMSNormFuser)
@pytest.mark.parametrize(
"cls,expected,zero_centered",
[
(RMSNorm, "RMSNorm", False),
(GemmaRMSNorm, "GemmaRMSNorm", True),
(WeightlessRMSNorm, "RMSNorm", False),
],
)
def test_rms_norm_builds_vllm_class(cls, expected, zero_centered, default_vllm_config):
from vllm.model_executor.layers.layernorm import GemmaRMSNorm as VLLMGemmaRMSNorm
from vllm.model_executor.layers.layernorm import RMSNorm as VLLMRMSNorm
# `default_vllm_config` supplies the config context the CustomOp needs; the
# weightless path reads hidden size from the model config, so stub it.
model_config = SimpleNamespace(get_hidden_size=lambda: 16)
with torch.device("meta"):
module = cls()
fuser = get_fuser(module)
built = fuser.fuse(module, "norm", model_config, None)
from vllm.model_executor.models.transformers.fusers.rms_norm import (
TPAwareNormMixin,
)
types_by_name = {"RMSNorm": VLLMRMSNorm, "GemmaRMSNorm": VLLMGemmaRMSNorm}
assert isinstance(built, types_by_name[expected])
assert isinstance(built, TPAwareNormMixin) # fused norms self-correct under TP
assert built.variance_epsilon == module.variance_epsilon
assert isinstance(built.weight, nn.Parameter) == (
getattr(module, "weight", None) is not None
)
def test_fused_rms_norm_op_default_eps(default_vllm_config):
"""`torch.nn.RMSNorm` (a single `F.rms_norm` call) matches via the fast path;
its default `eps=None` resolves to `finfo(dtype).eps` in `fuse`."""
from vllm.model_executor.layers.layernorm import RMSNorm as VLLMRMSNorm
with torch.device("meta"):
module = torch.nn.RMSNorm(16) # forward is a single `F.rms_norm` call
fuser = get_fuser(module)
assert isinstance(fuser, RMSNormFuser)
assert not fuser.zero_centered
model_config = SimpleNamespace(get_hidden_size=lambda: 16, dtype=torch.float32)
built = fuser.fuse(module, "norm", model_config, None)
assert isinstance(built, VLLMRMSNorm)
assert built.variance_epsilon == torch.finfo(torch.float32).eps
def test_eps_is_derived_per_instance(default_vllm_config):
"""Two instances of the same norm class with different eps must fuse to their
own eps: the type-cached fuser holds only structure, not this value."""
model_config = SimpleNamespace(get_hidden_size=lambda: 16)
with torch.device("meta"):
for eps in (1e-5, 1e-6):
module = RMSNorm(16, eps=eps)
built = get_fuser(module).fuse(module, "norm", model_config, None)
assert built.variance_epsilon == eps
def test_fused_norm_is_gather_capable(default_vllm_config):
"""Every fused norm is emitted gather-capable, so a norm on a head-sharded
projection (OLMoE-style) self-corrects at runtime with no QKV-specific
plumbing. A full-width input skips the gather and equals a plain norm."""
from vllm.model_executor.layers.layernorm import GemmaRMSNorm, RMSNorm
from vllm.model_executor.models.transformers.fusers import rms_norm
torch.manual_seed(0)
x = torch.randn(4, 16)
for gathered_cls, plain_cls in [
(rms_norm.TPAwareRMSNorm, RMSNorm),
(rms_norm.TPAwareGemmaRMSNorm, GemmaRMSNorm),
]:
gathered = gathered_cls(hidden_size=16, eps=1e-6)
assert isinstance(gathered, rms_norm.TPAwareNormMixin)
plain = plain_cls(hidden_size=16, eps=1e-6)
with torch.no_grad():
weight = torch.randn(16)
gathered.weight.copy_(weight)
plain.weight.copy_(weight)
torch.testing.assert_close(gathered(x), plain(x))
def test_gathered_norm_rejects_uneven_sharding(default_vllm_config):
"""A sharded input (narrower than the full-width weight) that does not tile
the weight evenly across ranks is rejected before any collective."""
from vllm.model_executor.models.transformers.fusers import rms_norm
norm = rms_norm.TPAwareRMSNorm(hidden_size=8, eps=1e-6)
norm.tp_size = 2 # emulate TP=2 without a real process group
with pytest.raises(ValueError, match="does not tile it evenly"):
norm(torch.randn(2, 3)) # 3 * 2 != 8
@@ -6,10 +6,16 @@ from typing import Any
import pytest
from ..conftest import HfRunner, VllmRunner
from ..utils import multi_gpu_test, prep_prompts
from .registry import HF_EXAMPLE_MODELS
from .utils import check_embeddings_close, check_logprobs_close
from ...conftest import HfRunner, VllmRunner
from ...utils import multi_gpu_test, prep_prompts
from ..registry import HF_EXAMPLE_MODELS
from ..utils import check_embeddings_close, check_logprobs_close
@pytest.fixture(scope="function", autouse=True)
def enable_pickle(monkeypatch):
"""`LLM.apply_model` requires pickling a function."""
monkeypatch.setenv("VLLM_ALLOW_INSECURE_SERIALIZATION", "1")
def get_model(arch: str) -> str:
@@ -18,6 +24,17 @@ def get_model(arch: str) -> str:
return model_info.default
def get_num_fused(model) -> tuple[int, int]:
from vllm.model_executor.layers.linear import (
MergedColumnParallelLinear,
QKVParallelLinear,
)
glu = sum(isinstance(m, MergedColumnParallelLinear) for m in model.modules())
qkv = sum(isinstance(m, QKVParallelLinear) for m in model.modules())
return glu, qkv
def check_implementation(
runner_ref: type[HfRunner | VllmRunner],
runner_test: type[VllmRunner],
@@ -25,6 +42,7 @@ def check_implementation(
model: str,
kwargs_ref: dict[str, Any] | None = None,
kwargs_test: dict[str, Any] | None = None,
num_fused: tuple[int, int] = (1, 1),
**kwargs,
):
if kwargs_ref is None:
@@ -41,6 +59,12 @@ def check_implementation(
model_config = model_test.llm.llm_engine.model_config
assert model_config.using_transformers_backend()
num_layers = model_config.hf_config.get_text_config().num_hidden_layers
expected_glu, expected_qkv = num_fused
for num_glu, num_qkv in model_test.apply_model(get_num_fused):
assert num_glu == expected_glu * num_layers
assert num_qkv == expected_qkv * num_layers
outputs_test = model_test.generate_greedy_logprobs(*args)
with runner_ref(model, **kwargs_ref) as model_ref:
@@ -58,11 +82,11 @@ def check_implementation(
@pytest.mark.parametrize(
"model,model_impl",
"model,model_impl,num_fused",
[
("meta-llama/Llama-3.2-1B-Instruct", "transformers"),
("hmellor/Ilama-3.2-1B", "auto"), # CUSTOM CODE
("allenai/OLMoE-1B-7B-0924", "transformers"), # MoE
("meta-llama/Llama-3.2-1B-Instruct", "transformers", (1, 1)),
("hmellor/Ilama-3.2-1B", "auto", (1, 1)), # CUSTOM CODE
("allenai/OLMoE-1B-7B-0924", "transformers", (0, 1)), # MoE
],
) # trust_remote_code=True by default
def test_models(
@@ -71,6 +95,7 @@ def test_models(
example_prompts: list[str],
model: str,
model_impl: str,
num_fused: tuple[int, int],
) -> None:
import transformers
from packaging.version import Version
@@ -84,7 +109,12 @@ def test_models(
)
check_implementation(
hf_runner, vllm_runner, example_prompts, model, model_impl=model_impl
hf_runner,
vllm_runner,
example_prompts,
model,
num_fused=num_fused,
model_impl=model_impl,
)
@@ -0,0 +1,159 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Tests for VLLMUnprocessableEntityError and media fetch error handling.
Verifies that unprocessable image URLs (404, 403, DNS failures, etc.) return
HTTP 422 instead of 500.
"""
from http import HTTPStatus
from unittest.mock import AsyncMock, MagicMock, patch
import aiohttp
import pytest
from vllm.entrypoints.serve.utils.error_response import create_error_response
from vllm.exceptions import VLLMUnprocessableEntityError
from vllm.multimodal.media import MediaConnector
class TestVLLMUnprocessableEntityError:
"""Tests for VLLMUnprocessableEntityError exception."""
def test_creation(self):
exc = VLLMUnprocessableEntityError("Test error")
assert str(exc) == "Test error"
assert exc.parameter is None
def test_creation_with_parameter_and_value(self):
exc = VLLMUnprocessableEntityError(
"Test error",
parameter="image_url",
value="https://example.com/image.jpg",
)
assert "parameter=image_url" in str(exc)
assert "value=https://example.com/image.jpg" in str(exc)
def test_is_value_error_subclass(self):
exc = VLLMUnprocessableEntityError("Test")
assert isinstance(exc, ValueError)
class TestMediaConnectorErrorHandling:
"""Tests for MediaConnector error handling."""
@pytest.mark.asyncio
async def test_fetch_image_async_404(self):
connector = MediaConnector()
with patch.object(
connector.connection, "async_get_bytes", new_callable=AsyncMock
) as mock_get:
mock_get.side_effect = aiohttp.ClientResponseError(
request_info=MagicMock(),
history=(),
status=404,
message="Not Found",
)
with pytest.raises(VLLMUnprocessableEntityError) as exc_info:
await connector.fetch_image_async("https://example.com/missing.jpg")
assert exc_info.value.parameter == "image_url"
@pytest.mark.asyncio
async def test_fetch_image_async_dns_error(self):
"""DNS errors are transient and should remain as-is for retry."""
connector = MediaConnector()
with patch.object(
connector.connection, "async_get_bytes", new_callable=AsyncMock
) as mock_get:
mock_get.side_effect = aiohttp.ClientConnectorDNSError(
connection_key=MagicMock(),
os_error=MagicMock(),
)
with pytest.raises(aiohttp.ClientConnectorDNSError) as exc_info:
await connector.fetch_image_async(
"https://nonexistent.example/image.jpg"
)
assert isinstance(exc_info.value, aiohttp.ClientConnectorDNSError)
@pytest.mark.asyncio
async def test_fetch_image_async_500_preserved(self):
"""5xx errors should remain as server errors."""
connector = MediaConnector()
with patch.object(
connector.connection, "async_get_bytes", new_callable=AsyncMock
) as mock_get:
mock_get.side_effect = aiohttp.ClientResponseError(
request_info=MagicMock(),
history=(),
status=500,
message="Internal Server Error",
)
with pytest.raises(aiohttp.ClientResponseError) as exc_info:
await connector.fetch_image_async("https://example.com/image.jpg")
assert exc_info.value.status == 500
def test_fetch_image_404(self):
connector = MediaConnector()
with patch.object(
connector.connection, "get_bytes", new_callable=MagicMock
) as mock_get:
mock_get.side_effect = aiohttp.ClientResponseError(
request_info=MagicMock(),
history=(),
status=404,
message="Not Found",
)
with pytest.raises(VLLMUnprocessableEntityError) as exc_info:
connector.fetch_image("https://example.com/missing.jpg")
assert exc_info.value.parameter == "image_url"
def test_fetch_image_connection_error(self):
"""Connection errors are transient and should remain as-is for retry."""
connector = MediaConnector()
with patch.object(
connector.connection, "get_bytes", new_callable=MagicMock
) as mock_get:
mock_get.side_effect = aiohttp.ClientConnectionError("Connection refused")
with pytest.raises(aiohttp.ClientConnectionError) as exc_info:
connector.fetch_image("https://example.com/image.jpg")
assert isinstance(exc_info.value, aiohttp.ClientConnectionError)
class TestErrorResponse:
"""Tests for error response creation."""
def test_unprocessable_entity_returns_422(self):
exc = VLLMUnprocessableEntityError(
"Failed to fetch media from URL: Cannot connect",
parameter="image_url",
value="https://example.com/image.jpg",
)
response = create_error_response(exc)
assert response.error.code == HTTPStatus.UNPROCESSABLE_ENTITY.value
assert response.error.type == "UnprocessableEntityError"
assert response.error.param == "image_url"
def test_unprocessable_entity_message(self):
exc = VLLMUnprocessableEntityError("Test error message")
response = create_error_response(exc)
assert response.error.message == "Test error message"
assert response.error.code == 422
+97 -1
View File
@@ -16,7 +16,11 @@ from vllm.assets.video import (
video_to_pil_images_list,
)
from vllm.multimodal.media import ImageMediaIO, VideoMediaIO
from vllm.multimodal.video import VIDEO_LOADER_REGISTRY, VideoLoader
from vllm.multimodal.video import (
PYNVVIDEOCODEC_VIDEO_BACKEND,
VIDEO_LOADER_REGISTRY,
VideoLoader,
)
from ..utils import cosine_similarity, create_video_from_image, normalize_image
@@ -357,3 +361,95 @@ def test_load_base64_jpeg_raises_on_zero_num_frames():
with pytest.raises(ValueError, match="num_frames must be greater than 0 or -1"):
videoio.load_base64("video/jpeg", data)
# ---------------------------------------------------------------------------
# GPU video backend policy tests
# ---------------------------------------------------------------------------
class TestMergeKwargsGpuBackendPolicy:
"""Verify that merge_kwargs blocks request-level GPU backend selection
when the static (engine-level) config did not configure that backend."""
def test_pynvvideocodec_requires_gpu(self):
assert VIDEO_LOADER_REGISTRY.backend_requires_gpu(PYNVVIDEOCODEC_VIDEO_BACKEND)
def test_strips_video_backend_pynv_when_not_static(self):
result = VideoMediaIO.merge_kwargs(
default_kwargs=None,
runtime_kwargs={"video_backend": "pynvvideocodec"},
)
assert "video_backend" not in result
def test_strips_backend_pynv_when_not_static(self):
result = VideoMediaIO.merge_kwargs(
default_kwargs={"num_frames": 16},
runtime_kwargs={"backend": "pynvvideocodec"},
)
assert result.get("backend") != "pynvvideocodec"
def test_preserves_video_backend_pynv_when_static(self):
result = VideoMediaIO.merge_kwargs(
default_kwargs={"video_backend": "pynvvideocodec"},
runtime_kwargs={"video_backend": "pynvvideocodec", "num_frames": 8},
)
assert result["video_backend"] == "pynvvideocodec"
assert result["num_frames"] == 8
def test_preserves_backend_pynv_when_static(self):
result = VideoMediaIO.merge_kwargs(
default_kwargs={"backend": "pynvvideocodec"},
runtime_kwargs={"backend": "pynvvideocodec"},
)
assert result["backend"] == "pynvvideocodec"
@pytest.mark.parametrize("backend", ["opencv", "pyav", "torchcodec"])
def test_software_video_backend_passes_through(self, backend: str):
result = VideoMediaIO.merge_kwargs(
default_kwargs=None,
runtime_kwargs={"video_backend": backend},
)
assert result["video_backend"] == backend
@pytest.mark.parametrize("backend", ["opencv", "pyav"])
def test_software_codec_backend_passes_through(self, backend: str):
result = VideoMediaIO.merge_kwargs(
default_kwargs=None,
runtime_kwargs={"backend": backend},
)
assert result["backend"] == backend
def test_strips_both_keys_independently(self):
result = VideoMediaIO.merge_kwargs(
default_kwargs=None,
runtime_kwargs={
"video_backend": "pynvvideocodec",
"backend": "pynvvideocodec",
"num_frames": 4,
},
)
assert "video_backend" not in result
assert result.get("backend") != "pynvvideocodec"
assert result["num_frames"] == 4
def test_other_kwargs_preserved_when_gpu_backend_stripped(self):
result = VideoMediaIO.merge_kwargs(
default_kwargs={"fps": 2},
runtime_kwargs={
"video_backend": "pynvvideocodec",
"num_frames": 16,
},
)
assert "video_backend" not in result
assert result["num_frames"] == 16
def test_static_pynv_with_different_runtime_gpu_backend(self):
"""If static sets pynv via video_backend but runtime tries to set it
via the codec-level 'backend' key (without a static match), strip it."""
result = VideoMediaIO.merge_kwargs(
default_kwargs={"video_backend": "pynvvideocodec"},
runtime_kwargs={"backend": "pynvvideocodec"},
)
assert result.get("backend") != "pynvvideocodec"
assert result["video_backend"] == "pynvvideocodec"
+5 -5
View File
@@ -61,7 +61,7 @@ MODELS = [
)
@pytest.mark.parametrize("model", MODELS)
def test_auto_round_model(vllm_runner, model):
with vllm_runner(model, enforce_eager=True) as llm:
with vllm_runner(model) as llm:
output = llm.generate_greedy(["The capital of France is"], max_tokens=8)
assert output
@@ -336,7 +336,7 @@ def test_wna16_xpu_prefers_ark_when_available(monkeypatch) -> None:
monkeypatch.setattr(current_platform, "is_xpu", lambda: True)
monkeypatch.setattr(current_platform, "is_cpu", lambda: False)
monkeypatch.setattr(
"vllm.model_executor.layers.quantization.inc.schemes.inc_wna16_linear.get_ark_state",
"vllm.model_executor.layers.quantization.inc.schemes.inc_ark_ops.get_ark_state",
lambda: (True, None, object(), DummyQuantLinear),
)
@@ -355,7 +355,7 @@ def test_wna16_xpu_falls_back_when_ark_unavailable(monkeypatch) -> None:
monkeypatch.setattr(current_platform, "is_xpu", lambda: True)
monkeypatch.setattr(current_platform, "is_cpu", lambda: False)
monkeypatch.setattr(
"vllm.model_executor.layers.quantization.inc.schemes.inc_wna16_linear.get_ark_state",
"vllm.model_executor.layers.quantization.inc.schemes.inc_ark_ops.get_ark_state",
lambda: (False, "missing", None, None),
)
@@ -377,7 +377,7 @@ def test_wna16_cpu_gptq_prefers_ark_when_available(monkeypatch) -> None:
monkeypatch.setattr(current_platform, "is_xpu", lambda: False)
monkeypatch.setattr(current_platform, "is_cpu", lambda: True)
monkeypatch.setattr(
"vllm.model_executor.layers.quantization.inc.schemes.inc_wna16_linear.get_ark_state",
"vllm.model_executor.layers.quantization.inc.schemes.inc_ark_ops.get_ark_state",
lambda: (True, None, object(), DummyQuantLinear),
)
@@ -398,7 +398,7 @@ def test_wna16_cpu_gptq_raises_when_ark_and_marlin_unavailable(
monkeypatch.setattr(current_platform, "is_xpu", lambda: False)
monkeypatch.setattr(current_platform, "is_cpu", lambda: True)
monkeypatch.setattr(
"vllm.model_executor.layers.quantization.inc.schemes.inc_wna16_linear.get_ark_state",
"vllm.model_executor.layers.quantization.inc.schemes.inc_ark_ops.get_ark_state",
lambda: (False, "missing", None, None),
)
monkeypatch.setattr(
@@ -7,12 +7,16 @@ import json
from unittest.mock import Mock
import pytest
from openai.types.responses import ResponseFunctionToolCall
from vllm.entrypoints.openai.chat_completion.protocol import (
ChatCompletionRequest,
ChatCompletionToolsParam,
FunctionDefinition,
)
from vllm.entrypoints.openai.engine.protocol import FunctionCall
from vllm.entrypoints.openai.responses.protocol import ResponsesRequest
from vllm.entrypoints.openai.responses.utils import build_response_output_items
from vllm.tokenizers import get_tokenizer
from vllm.tool_parsers.glm47_moe_tool_parser import Glm47MoeModelToolParser
@@ -58,7 +62,69 @@ def mock_request(sample_tools) -> ChatCompletionRequest:
return request
@pytest.fixture
def namespace_tool_request() -> ResponsesRequest:
return ResponsesRequest.model_validate(
{
"input": "hi",
"tools": [
{
"type": "namespace",
"name": "mcp__computer_use",
"description": "Computer use tools.",
"tools": [
{
"type": "function",
"name": "get_app_state",
"description": "Get app state.",
"parameters": {
"type": "object",
"properties": {
"app": {"type": "string"},
},
},
}
],
}
],
}
)
class TestGlm47ExtractToolCalls:
def test_namespace_tool_call_round_trip_to_responses_output(
self, glm47_tokenizer, namespace_tool_request
):
parser = Glm47MoeModelToolParser(
glm47_tokenizer, tools=namespace_tool_request.tools
)
out = (
"<tool_call>mcp__computer_use__get_app_state"
"<arg_key>app</arg_key>"
"<arg_value>Google Chrome</arg_value>"
"</tool_call>"
)
result = parser.extract_tool_calls(out, request=namespace_tool_request)
assert result.tools_called
tool_call = result.tool_calls[0].function
assert tool_call == FunctionCall(
name="mcp__computer_use__get_app_state",
arguments='{"app": "Google Chrome"}',
)
output_items = build_response_output_items(
reasoning=None,
content=None,
tool_calls=[tool_call],
tools=namespace_tool_request.tools,
)
output_tool_call = output_items[0]
assert isinstance(output_tool_call, ResponseFunctionToolCall)
assert output_tool_call.name == "get_app_state"
assert output_tool_call.namespace == "mcp__computer_use"
def test_no_tool_call(self, glm47_tool_parser, mock_request):
out = "This is a plain response."
r = glm47_tool_parser.extract_tool_calls(out, request=mock_request)
+40
View File
@@ -1459,6 +1459,46 @@ def multi_process_parallel(
ray.shutdown()
def assert_rocm_custom_allreduce_backend_state(
use_aiter_custom_ar: bool,
quick_reduce_quantization: str,
) -> None:
from vllm.distributed.parallel_state import get_tp_group
device_communicator = get_tp_group().device_communicator
aiter_ar_comm = device_communicator.aiter_ar_comm
if use_aiter_custom_ar:
assert aiter_ar_comm is not None, "AITER CustomAllreduce was not initialized."
assert not aiter_ar_comm.disabled, "AITER CustomAllreduce is disabled."
assert device_communicator.ca_comm is None, (
"vLLM CustomAllreduce should not be initialized when AITER CA is used."
)
else:
assert aiter_ar_comm is None, (
"AITER CustomAllreduce should not be initialized when disabled."
)
assert device_communicator.ca_comm is not None, (
"vLLM CustomAllreduce should be initialized when AITER CA is disabled."
)
qr_comm = device_communicator.qr_comm
assert qr_comm is not None, "QuickReduce communicator was not initialized."
if quick_reduce_quantization == "NONE":
assert qr_comm.disabled, "QuickReduce should be disabled."
else:
assert not qr_comm.disabled, "QuickReduce should be enabled."
def assert_rocm_custom_allreduce_backend_state_on_worker(
_worker,
use_aiter_custom_ar: bool,
quick_reduce_quantization: str,
) -> None:
assert_rocm_custom_allreduce_backend_state(
use_aiter_custom_ar, quick_reduce_quantization
)
@contextmanager
def error_on_warning(category: type[Warning] = Warning):
"""
@@ -0,0 +1,188 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import torch
from tests.v1.attention.utils import (
BatchSpec,
create_common_attn_metadata,
create_vllm_config,
)
from vllm.config import CUDAGraphMode, SpeculativeConfig
from vllm.v1.attention.backend import AttentionCGSupport
from vllm.v1.attention.backends.linear_attn import (
BailingLinearAttentionMetadataBuilder,
LinearAttentionMetadataBuilder,
)
from vllm.v1.attention.backends.utils import PAD_SLOT_ID
from vllm.v1.kv_cache_interface import MambaSpec
BLOCK_SIZE = 16
DEVICE = torch.device("cpu")
def _create_mamba_spec(num_speculative_blocks: int = 1) -> MambaSpec:
return MambaSpec(
block_size=BLOCK_SIZE,
shapes=((16, 64),),
dtypes=(torch.float16,),
num_speculative_blocks=num_speculative_blocks,
)
def test_bailing_linear_attention_reports_uniform_batch_cudagraph_support():
vllm_config = create_vllm_config(
hf_config_override={
"architectures": ["BailingMoeV2_5ForCausalLM"],
"model_type": "bailing_hybrid",
}
)
support = BailingLinearAttentionMetadataBuilder.get_cudagraph_support(
vllm_config, _create_mamba_spec()
)
assert support == AttentionCGSupport.UNIFORM_BATCH
def test_non_bailing_linear_attention_keeps_single_token_cudagraph_support():
vllm_config = create_vllm_config(
hf_config_override={
"architectures": ["MiniMaxText01ForCausalLM"],
"model_type": "minimax_text_01",
}
)
support = LinearAttentionMetadataBuilder.get_cudagraph_support(
vllm_config, _create_mamba_spec()
)
assert support == AttentionCGSupport.UNIFORM_SINGLE_TOKEN_DECODE
def test_linear_attention_spec_decode_full_graph_metadata_pads_cache_slots():
vllm_config = create_vllm_config(
hf_config_override={
"architectures": ["BailingMoeV2_5ForCausalLM"],
"model_type": "bailing_hybrid",
}
)
vllm_config.speculative_config = SpeculativeConfig(
method="ngram",
num_speculative_tokens=1,
)
vllm_config.compilation_config.cudagraph_mode = CUDAGraphMode.FULL_DECODE_ONLY
builder = BailingLinearAttentionMetadataBuilder(
kv_cache_spec=_create_mamba_spec(),
layer_names=["model.layers.0.self_attn"],
vllm_config=vllm_config,
device=DEVICE,
)
common = create_common_attn_metadata(
BatchSpec(seq_lens=[20, 20, 0], query_lens=[2, 2, 0]),
BLOCK_SIZE,
DEVICE,
)
common.block_table_tensor[2].fill_(-1)
metadata = builder.build(
common_prefix_len=0,
common_attn_metadata=common,
num_accepted_tokens=torch.tensor([1, 2, 1], dtype=torch.int32),
)
assert metadata.num_decodes == 3
assert metadata.num_prefills == 0
assert metadata.num_decode_tokens == 4
assert metadata.state_indices_tensor_d is not None
assert metadata.state_indices_tensor_d.shape == (3, 2)
assert torch.equal(
metadata.state_indices_tensor_d[2],
torch.full((2,), PAD_SLOT_ID, dtype=torch.int32),
)
assert torch.equal(
metadata.state_indices_tensor[2],
torch.tensor(PAD_SLOT_ID, dtype=torch.int32),
)
assert metadata.query_start_loc_d is not None
assert metadata.query_start_loc_d.tolist() == [0, 2, 4, 4]
assert metadata.num_accepted_tokens is not None
assert metadata.num_accepted_tokens.tolist() == [1, 2, 1]
def test_linear_attention_full_graph_metadata_uses_stable_decode_buffers():
vllm_config = create_vllm_config(
hf_config_override={
"architectures": ["BailingMoeV2_5ForCausalLM"],
"model_type": "bailing_hybrid",
}
)
vllm_config.speculative_config = SpeculativeConfig(
method="ngram",
num_speculative_tokens=1,
)
vllm_config.compilation_config.cudagraph_mode = CUDAGraphMode.FULL_DECODE_ONLY
builder = BailingLinearAttentionMetadataBuilder(
kv_cache_spec=_create_mamba_spec(),
layer_names=["model.layers.0.self_attn"],
vllm_config=vllm_config,
device=DEVICE,
)
common = create_common_attn_metadata(
BatchSpec(seq_lens=[20, 20, 0], query_lens=[2, 2, 0]),
BLOCK_SIZE,
DEVICE,
arange_block_indices=True,
)
common.block_table_tensor = torch.tensor(
[[10, 11], [12, 13], [-1, -1]],
dtype=torch.int32,
device=DEVICE,
)
first = builder.build(
common_prefix_len=0,
common_attn_metadata=common,
num_accepted_tokens=torch.tensor([1, 2, 1], dtype=torch.int32),
)
assert first.state_indices_tensor_d is not None
assert first.query_start_loc_d is not None
assert first.num_accepted_tokens is not None
state_ptr = first.state_indices_tensor_d.data_ptr()
query_ptr = first.query_start_loc_d.data_ptr()
accepted_ptr = first.num_accepted_tokens.data_ptr()
common2 = create_common_attn_metadata(
BatchSpec(seq_lens=[36, 0, 0], query_lens=[2, 0, 0]),
BLOCK_SIZE,
DEVICE,
arange_block_indices=True,
)
common2.block_table_tensor = torch.tensor(
[[20, 21], [-1, -1], [-1, -1]],
dtype=torch.int32,
device=DEVICE,
)
second = builder.build(
common_prefix_len=0,
common_attn_metadata=common2,
num_accepted_tokens=torch.tensor([2, 1, 1], dtype=torch.int32),
)
assert second.state_indices_tensor_d is not None
assert second.query_start_loc_d is not None
assert second.num_accepted_tokens is not None
assert second.state_indices_tensor_d.data_ptr() == state_ptr
assert second.query_start_loc_d.data_ptr() == query_ptr
assert second.num_accepted_tokens.data_ptr() == accepted_ptr
assert second.state_indices_tensor_d.tolist() == [
[20, 21],
[PAD_SLOT_ID, PAD_SLOT_ID],
[PAD_SLOT_ID, PAD_SLOT_ID],
]
assert second.query_start_loc_d.tolist() == [0, 2, 2, 2]
assert second.num_accepted_tokens.tolist() == [2, 1, 1]
+45
View File
@@ -324,3 +324,48 @@ def test_abort_request_when_structured_output_fsm_cannot_advance():
assert request.status == RequestStatus.FINISHED_ERROR
assert request.request_id not in scheduler.requests
assert not scheduler.running
def test_no_placeholder_underflow_on_discarded_spec_frame():
num_spec = 5
scheduler = create_scheduler(
async_scheduling=True,
num_speculative_tokens=num_spec,
speculative_method="ngram_gpu",
)
req = create_requests(num_requests=1, max_tokens=20)[0]
req.num_computed_tokens = req.num_tokens
scheduler.requests[req.request_id] = req
scheduler.running.append(req)
req.status = RequestStatus.RUNNING
req.num_output_placeholders = 1
req.async_tokens_to_discard = num_spec
computed_before = req.num_computed_tokens
scheduler_output = SchedulerOutput(
scheduled_new_reqs=[],
scheduled_cached_reqs=CachedRequestData.make_empty(),
num_scheduled_tokens={req.request_id: num_spec + 1},
total_num_scheduled_tokens=num_spec + 1,
scheduled_encoder_inputs={},
scheduled_spec_decode_tokens={req.request_id: [10] * num_spec},
num_common_prefix_blocks=[],
finished_req_ids=set(),
free_encoder_mm_hashes=[],
)
model_runner_output = ModelRunnerOutput(
req_ids=[req.request_id],
req_id_to_index={req.request_id: 0},
sampled_token_ids=[[999]],
logprobs=None,
prompt_logprobs_dict={},
pooler_output=[],
)
scheduler.update_from_output(scheduler_output, model_runner_output)
assert req.num_output_placeholders == 1
assert req.num_computed_tokens == computed_before
assert req.async_tokens_to_discard == num_spec - 1
assert req.status == RequestStatus.RUNNING
+7 -1
View File
@@ -54,6 +54,7 @@ def create_scheduler(
block_size: int = 16,
max_model_len: int | None = None,
num_speculative_tokens: int | None = None,
speculative_method: str | None = None,
skip_tokenizer_init: bool = False,
async_scheduling: bool = False,
pipeline_parallel_size: int = 1,
@@ -126,9 +127,14 @@ def create_scheduler(
speculative_config: SpeculativeConfig | None = None
if num_speculative_tokens is not None:
speculative_config = SpeculativeConfig(
spec_kwargs: dict = dict(
model="ngram", num_speculative_tokens=num_speculative_tokens
)
if speculative_method is not None:
spec_kwargs["method"] = speculative_method
spec_kwargs["prompt_lookup_max"] = num_speculative_tokens
spec_kwargs["prompt_lookup_min"] = 1
speculative_config = SpeculativeConfig(**spec_kwargs)
ec_transfer_config = (
ECTransferConfig(
+3 -24
View File
@@ -1,11 +1,10 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import weakref
from contextlib import ExitStack
import pytest
from tests.utils import wait_for_gpu_memory_to_clear
from tests.utils import create_new_process_for_each_test
from tests.v1.attention.utils import full_cg_backend_configs as backend_configs
from vllm import LLM
from vllm.config import CompilationConfig, CompilationMode
@@ -32,6 +31,7 @@ else:
@pytest.mark.parametrize("backend_name, cudagraph_mode, supported", combo_cases_1)
@create_new_process_for_each_test("spawn")
def test_backend_and_cudagraph_mode_combo(backend_name, cudagraph_mode, supported):
if backend_name == "FlashInfer":
try:
@@ -64,17 +64,6 @@ def test_backend_and_cudagraph_mode_combo(backend_name, cudagraph_mode, supporte
),
)
llm.generate(["Hello, my name is"] * 10)
# when above code raises, `llm` may be undefined, so we need to catch that
try:
llm = weakref.proxy(llm)
del llm
except UnboundLocalError:
pass
wait_for_gpu_memory_to_clear(
devices=[0],
threshold_ratio=0.1,
)
# test cudagraph_mode with different compilation mode.
@@ -98,6 +87,7 @@ combo_cases_2 = [
@pytest.mark.parametrize(
"backend_name,cudagraph_mode,compilation_mode,supported", combo_cases_2
)
@create_new_process_for_each_test("spawn")
def test_cudagraph_compilation_combo(
backend_name, cudagraph_mode, compilation_mode, supported
):
@@ -120,14 +110,3 @@ def test_cudagraph_compilation_combo(
),
)
llm.generate(["Hello, my name is"] * 10)
# when above code raises, `llm` may be undefined, so we need to catch that
try:
llm = weakref.proxy(llm)
del llm
except UnboundLocalError:
pass
finally:
wait_for_gpu_memory_to_clear(
devices=[0],
threshold_ratio=0.1,
)
@@ -45,7 +45,9 @@ def test_prompts():
use_fork_for_test = (
fork_new_process_for_each_test if not current_platform.is_rocm() else lambda x: x
fork_new_process_for_each_test
if not (current_platform.is_rocm() or current_platform.is_xpu())
else lambda x: x
)
@@ -0,0 +1,118 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import pytest
from vllm._aiter_ops import is_aiter_found, rocm_aiter_ops
from vllm.config import CompilationConfig, CompilationMode, CUDAGraphMode
from vllm.envs import disable_envs_cache
from vllm.platforms import current_platform
from ....conftest import VllmRunner
from ....utils import (
assert_rocm_custom_allreduce_backend_state_on_worker,
multi_gpu_test,
)
PROMPTS = ["Hello, my name is", "The capital of France is"]
def _run_generation(
vllm_runner: type[VllmRunner],
monkeypatch: pytest.MonkeyPatch,
compilation_config: CompilationConfig,
*,
model: str,
max_tokens: int,
use_aiter_custom_ar: bool,
quick_reduce_quantization: str,
) -> list[tuple[list[int], str]]:
with monkeypatch.context() as m:
m.setenv("VLLM_ALLOW_INSECURE_SERIALIZATION", "1")
m.setenv("VLLM_ROCM_USE_AITER", "1")
m.setenv(
"VLLM_ROCM_USE_AITER_CUSTOM_AR",
"1" if use_aiter_custom_ar else "0",
)
m.setenv("VLLM_ROCM_QUICK_REDUCE_QUANTIZATION", quick_reduce_quantization)
disable_envs_cache()
rocm_aiter_ops.refresh_env_variables()
with vllm_runner(
model,
dtype="half",
tensor_parallel_size=2,
compilation_config=compilation_config,
max_model_len=256,
max_num_seqs=len(PROMPTS),
gpu_memory_utilization=0.7,
) as llm:
llm.get_llm().collective_rpc(
assert_rocm_custom_allreduce_backend_state_on_worker,
args=(use_aiter_custom_ar, quick_reduce_quantization),
)
return llm.generate_greedy(PROMPTS, max_tokens)
@pytest.mark.skipif(not current_platform.is_rocm(), reason="ROCm-only")
@pytest.mark.skipif(not is_aiter_found(), reason="AITER is not installed")
@multi_gpu_test(num_gpus=2)
@pytest.mark.parametrize(
"quick_reduce_quantization",
[
pytest.param("FP", id="quick-reduce-on"),
pytest.param("NONE", id="quick-reduce-off"),
],
)
@pytest.mark.parametrize(
"cudagraph_mode",
[
pytest.param(CUDAGraphMode.NONE, id="cudagraph-none"),
pytest.param(CUDAGraphMode.FULL, id="cudagraph-full"),
],
)
@pytest.mark.parametrize(
"model,max_tokens",
[
pytest.param("facebook/opt-125m", 8, id="opt-125m"),
],
)
def test_rocm_aiter_custom_ar_e2e(
vllm_runner: type[VllmRunner],
monkeypatch: pytest.MonkeyPatch,
cudagraph_mode: CUDAGraphMode,
quick_reduce_quantization: str,
model: str,
max_tokens: int,
):
compilation_mode = (
CompilationMode.NONE
if cudagraph_mode == CUDAGraphMode.NONE
else CompilationMode.VLLM_COMPILE
)
compilation_config = CompilationConfig(
mode=compilation_mode,
cudagraph_mode=cudagraph_mode,
)
baseline_generations = _run_generation(
vllm_runner,
monkeypatch,
compilation_config,
model=model,
max_tokens=max_tokens,
use_aiter_custom_ar=False,
quick_reduce_quantization=quick_reduce_quantization,
)
aiter_custom_ar_generations = _run_generation(
vllm_runner,
monkeypatch,
compilation_config,
model=model,
max_tokens=max_tokens,
use_aiter_custom_ar=True,
quick_reduce_quantization=quick_reduce_quantization,
)
assert aiter_custom_ar_generations == baseline_generations
@@ -9,8 +9,10 @@ PREFILL_GPU_ID=${PREFILL_GPU_ID:-0}
DECODE_GPU_ID=${DECODE_GPU_ID:-1}
MODEL=${MODEL:-"ibm-granite/granite-4.0-h-tiny"}
GPU_MEMORY_UTILIZATION=${GPU_MEMORY_UTILIZATION:-0.8}
VLLM_SERVE_EXTRA_ARGS=${VLLM_SERVE_EXTRA_ARGS:-}
ATTENTION_BACKEND=${ATTENTION_BACKEND:-FLASHINFER}
echo "Running Mamba prefix cache test (GPUs: P=$PREFILL_GPU_ID, D=$DECODE_GPU_ID, model=$MODEL)"
echo "Running Mamba prefix cache test (GPUs: P=$PREFILL_GPU_ID, D=$DECODE_GPU_ID, model=$MODEL, backend=$ATTENTION_BACKEND)"
KV_CONFIG='{"kv_connector":"NixlConnector","kv_role":"kv_both"}'
@@ -36,6 +38,14 @@ cleanup_instances() {
cleanup_instances
EXTRA_ARGS=()
if [[ -n "$VLLM_SERVE_EXTRA_ARGS" ]]; then
IFS=',' read -r -a EXTRA_ARGS <<< "$VLLM_SERVE_EXTRA_ARGS"
fi
if [[ -n "$ATTENTION_BACKEND" ]]; then
EXTRA_ARGS+=(--attention-backend "$ATTENTION_BACKEND")
fi
# Start prefill instance
PREFILL_PORT=8001
CUDA_VISIBLE_DEVICES=$PREFILL_GPU_ID \
@@ -51,8 +61,8 @@ vllm serve $MODEL \
--trust-remote-code \
--enable-prefix-caching \
--mamba-cache-mode all \
--attention-backend FLASHINFER \
--kv-transfer-config "$KV_CONFIG" &
--kv-transfer-config "$KV_CONFIG" \
"${EXTRA_ARGS[@]}" &
# Start decode instance
DECODE_PORT=8002
@@ -69,8 +79,8 @@ vllm serve $MODEL \
--trust-remote-code \
--enable-prefix-caching \
--mamba-cache-mode all \
--attention-backend FLASHINFER \
--kv-transfer-config "$KV_CONFIG" &
--kv-transfer-config "$KV_CONFIG" \
"${EXTRA_ARGS[@]}" &
echo "Waiting for prefill instance on port $PREFILL_PORT..."
wait_for_server "$PREFILL_PORT"
@@ -18,6 +18,7 @@
# Environment variables:
# MODEL_NAMES - model to test (default: Qwen/Qwen3-0.6B)
# GPU_MEMORY_UTILIZATION - GPU memory fraction (default: 0.6)
# ATTENTION_BACKEND - optional attention backend for vllm serve
# VLLM_SERVE_EXTRA_ARGS - comma-separated extra args for vllm serve
# SKIP_CROSS_LAYERS - set to 1 to skip the cross-layer layout test
# SKIP_NORMAL_LAYOUT - set to 1 to skip the normal layout test
@@ -34,9 +35,11 @@ fi
GPU_MEMORY_UTILIZATION=${GPU_MEMORY_UTILIZATION:-0.6}
BLOCK_SIZE=${BLOCK_SIZE:-128}
ATTENTION_BACKEND=${ATTENTION_BACKEND:-}
VLLM_SERVE_EXTRA_ARGS=${VLLM_SERVE_EXTRA_ARGS:-}
GIT_ROOT=$(git rev-parse --show-toplevel)
SCRIPT_DIR="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd -P)"
GIT_ROOT="${GIT_ROOT:-$(cd -- "${SCRIPT_DIR}/../../../.." && pwd -P)}"
SMI_BIN=$(which nvidia-smi || which rocm-smi || echo "")
# ── KV transfer configs ─────────────────────────────────────────────────
@@ -139,6 +142,9 @@ run_tests_for_model() {
BASE_CMD="${BASE_CMD} $arg"
done
fi
if [[ -n "$ATTENTION_BACKEND" ]]; then
BASE_CMD="${BASE_CMD} --attention-backend $ATTENTION_BACKEND"
fi
eval "$BASE_CMD &"
# ── Start decode instance ──
@@ -161,6 +167,9 @@ run_tests_for_model() {
BASE_CMD="${BASE_CMD} $arg"
done
fi
if [[ -n "$ATTENTION_BACKEND" ]]; then
BASE_CMD="${BASE_CMD} --attention-backend $ATTENTION_BACKEND"
fi
eval "$BASE_CMD &"
# ── Wait for servers ──
@@ -19,6 +19,7 @@
# MODEL_NAMES - model to test (default: Qwen/Qwen3-0.6B)
# KV_CACHE_MEMORY_BYTES - GPU KV cache size in bytes (default: 268435456 = 256 MiB)
# BLOCK_SIZE - KV cache block size (default: 128)
# ATTENTION_BACKEND - optional attention backend for vllm serve
# VLLM_SERVE_EXTRA_ARGS - comma-separated extra args for vllm serve
set -xe
@@ -34,9 +35,11 @@ fi
KV_CACHE_MEMORY_BYTES=${KV_CACHE_MEMORY_BYTES:-268435456} # 256 MiB
MAX_MODEL_LEN=${MAX_MODEL_LEN:-2048}
BLOCK_SIZE=${BLOCK_SIZE:-128}
ATTENTION_BACKEND=${ATTENTION_BACKEND:-}
VLLM_SERVE_EXTRA_ARGS=${VLLM_SERVE_EXTRA_ARGS:-}
GIT_ROOT=$(git rev-parse --show-toplevel)
SCRIPT_DIR="$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd -P)"
GIT_ROOT="${GIT_ROOT:-$(cd -- "${SCRIPT_DIR}/../../../.." && pwd -P)}"
# ── KV transfer config ──────────────────────────────────────────────────
@@ -110,6 +113,9 @@ run_tests_for_model() {
BASE_CMD="${BASE_CMD} $arg"
done
fi
if [[ -n "$ATTENTION_BACKEND" ]]; then
BASE_CMD="${BASE_CMD} --attention-backend $ATTENTION_BACKEND"
fi
eval "$BASE_CMD &"
# ── Start decode instance ──
@@ -133,6 +139,9 @@ run_tests_for_model() {
BASE_CMD="${BASE_CMD} $arg"
done
fi
if [[ -n "$ATTENTION_BACKEND" ]]; then
BASE_CMD="${BASE_CMD} --attention-backend $ATTENTION_BACKEND"
fi
eval "$BASE_CMD &"
# ── Wait for servers ──
@@ -33,7 +33,7 @@ def hf3fs_stats():
def _make_cuda_event():
"""Return a real CUDA event when available, otherwise a MagicMock."""
if torch.cuda.is_available():
return torch.Event()
return torch.cuda.Event()
return MagicMock()
+46
View File
@@ -0,0 +1,46 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Unit tests for ``vllm.v1.worker.xpu_model_runner`` (XPU worker / CUDA shims)."""
import pytest
import torch
from torch._dynamo.variables.torch import TorchInGraphFunctionVariable
from vllm.v1.worker.xpu_model_runner import _torch_cuda_wrapper
# XPU-only: needs distinct torch.cuda vs torch.xpu current_stream symbols.
pytestmark = pytest.mark.skipif(
not hasattr(torch, "xpu") or not hasattr(torch.xpu, "current_stream"),
reason="torch.xpu.current_stream is required",
)
# Child process: patched torch.cuda must not leak to other tests in the session.
@pytest.mark.forked
def test_torch_cuda_wrapper_allows_dynamo_handler_registration() -> None:
"""Guard against XPU CUDA shim breaking Torch Dynamo during AOT compile.
Before the fix, ``_torch_cuda_wrapper`` assigned
``torch.cuda.current_stream = torch.xpu.current_stream`` (same function object).
On the first AOT/profile run, Dynamo builds its in-graph handler table and
registers ``torch.cuda.current_stream`` and ``torch.xpu.current_stream``
separately; duplicate identity triggers::
AssertionError: Handler already registered for <function current_stream ...>
That surfaced as EngineCore failing in ``profile_run`` / ``_get_handlers()``.
The fix uses distinct shim callables so both can be registered.
This test replays the post-init state (wrapper applied, patches left on
``torch.cuda``) and checks that Dynamo's real ``_get_handlers()`` succeeds.
"""
# Same entry point as XPUModelRunner.__init__ (patches persist after exit).
with _torch_cuda_wrapper():
pass
# Fresh handler table build, as on first torch.compile / AOT in the worker.
# Registers torch.cuda.current_stream and torch.xpu.current_stream separately;
# if they are the same object (pre-fix alias), raises Handler already registered.
TorchInGraphFunctionVariable._get_handlers.cache_clear()
TorchInGraphFunctionVariable._get_handlers()
+1 -9
View File
@@ -9,7 +9,7 @@ import regex as re
# --------------------------------------------------------------------------- #
_TORCH_CUDA_PATTERNS = [
r"\btorch\.cuda\.(empty_cache|synchronize|device_count|current_device|memory_reserved|memory_allocated|max_memory_allocated|max_memory_reserved|reset_peak_memory_stats|memory_stats|mem_get_info|set_device|device\()\b",
r"\btorch\.cuda\.(manual_seed|manual_seed_all|Event)\b",
r"\btorch\.cuda\.(manual_seed|manual_seed_all)\b",
r"\bwith\storch\.cuda\.device\b",
# Calls torch.cuda.{_is_compiled/_device_count_amdsmi/_device_count_nvml} internally
r"\bcuda_device_count_stateless\(\)\b",
@@ -21,7 +21,6 @@ ALLOWED_FILES = {
"vllm/device_allocator/",
"vllm/distributed/weight_transfer/ipc_engine.py",
"tests/distributed/test_packed_tensor.py",
"tools/pre_commit/check_torch_cuda.py",
}
@@ -40,13 +39,6 @@ def scan_file(path: str) -> int:
f"Found {matched_text} API call. Use set_random_seed instead."
)
return 1
if matched_text == "torch.cuda.Event":
print(
f"{path}:{line_num}: "
"\033[91merror:\033[0m "
"Found torch.cuda.Event API call. Use torch.Event instead."
)
return 1
print(
f"{path}:{line_num}: "
"\033[91merror:\033[0m " # red color
+34 -95
View File
@@ -2,12 +2,9 @@
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
import functools
from collections.abc import Callable
from contextlib import contextmanager
from typing import Protocol
import torch
from torch._ops import OpOverload
from torch.distributed import ProcessGroup
import vllm.envs as envs
from vllm.platforms import current_platform
@@ -52,42 +49,6 @@ def is_aiter_found() -> bool:
IS_AITER_FOUND = is_aiter_found()
class AiterCustomAllreduceProto(Protocol):
max_size: int
world_size: int
fully_connected: bool
@contextmanager
def capture(self): ...
def close(self) -> None: ...
def fused_ar_rms(
self,
inp: torch.Tensor,
res_inp: torch.Tensor,
*,
w: torch.Tensor,
eps: float,
registered: bool = False,
use_1stage: bool = False,
) -> tuple[torch.Tensor, torch.Tensor]: ...
def fused_ar_rms_per_group_quant(
self,
inp: torch.Tensor,
res_inp: torch.Tensor,
*,
w: torch.Tensor,
eps: float,
group_size: int = 128,
registered: bool = False,
use_1stage: bool = False,
emit_bf16: bool = False,
) -> (
tuple[torch.Tensor, torch.Tensor, torch.Tensor]
| tuple[torch.Tensor, torch.Tensor, torch.Tensor, torch.Tensor]
): ...
def should_custom_ar(self, inp: torch.Tensor) -> bool: ...
def is_aiter_found_and_supported() -> bool:
"""Check if AITER library is available and platform supports it.
@@ -830,6 +791,7 @@ def _rocm_aiter_fused_allreduce_rmsnorm_impl(
) -> tuple[torch.Tensor, torch.Tensor]:
aiter_ar = rocm_aiter_ops.get_aiter_allreduce()
assert aiter_ar is not None, "aiter allreduce must be initialized"
ca = aiter_ar.aiter_ca
total_bytes = input_.numel() * input_.element_size()
hidden_dim = input_.shape[-1]
@@ -840,8 +802,8 @@ def _rocm_aiter_fused_allreduce_rmsnorm_impl(
else:
hidden_ok = False
token_ok = token_num <= 80
world_size = aiter_ar.world_size
full_nvlink = aiter_ar.fully_connected
world_size = ca.world_size
full_nvlink = ca.fully_connected
if world_size == 2:
size_ok = True
@@ -854,12 +816,11 @@ def _rocm_aiter_fused_allreduce_rmsnorm_impl(
use_1stage = hidden_ok and token_ok and size_ok
result = aiter_ar.fused_ar_rms(
result = ca.custom_fused_ar_rms(
input_,
residual,
w=weight,
eps=epsilon,
registered=torch.cuda.is_current_stream_capturing(),
weight,
epsilon,
use_1stage=use_1stage,
)
assert result is not None
@@ -890,6 +851,7 @@ def _rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_impl(
"""
aiter_ar = rocm_aiter_ops.get_aiter_allreduce()
assert aiter_ar is not None, "aiter allreduce must be initialized"
ca = aiter_ar.aiter_ca
total_bytes = input_.numel() * input_.element_size()
hidden_dim = input_.shape[-1]
@@ -900,8 +862,8 @@ def _rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_impl(
else:
hidden_ok = False
token_ok = token_num <= 80
world_size = aiter_ar.world_size
full_nvlink = aiter_ar.fully_connected
world_size = ca.world_size
full_nvlink = ca.fully_connected
if world_size == 2:
size_ok = True
@@ -914,7 +876,7 @@ def _rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_impl(
use_1stage = hidden_ok and token_ok and size_ok
result = aiter_ar.fused_ar_rms_per_group_quant(
result = ca.fused_ar_rms_per_group_quant(
input_,
residual,
w=weight,
@@ -962,6 +924,7 @@ def _rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_with_bf16_norm_impl(
"""
aiter_ar = rocm_aiter_ops.get_aiter_allreduce()
assert aiter_ar is not None, "aiter allreduce must be initialized"
ca = aiter_ar.aiter_ca
total_bytes = input_.numel() * input_.element_size()
hidden_dim = input_.shape[-1]
@@ -972,8 +935,8 @@ def _rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_with_bf16_norm_impl(
else:
hidden_ok = False
token_ok = token_num <= 80
world_size = aiter_ar.world_size
full_nvlink = aiter_ar.fully_connected
world_size = ca.world_size
full_nvlink = ca.fully_connected
if world_size == 2:
size_ok = True
@@ -986,7 +949,7 @@ def _rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_with_bf16_norm_impl(
use_1stage = hidden_ok and token_ok and size_ok
result = aiter_ar.fused_ar_rms_per_group_quant(
result = ca.fused_ar_rms_per_group_quant(
input_,
residual,
w=weight,
@@ -1577,6 +1540,7 @@ class rocm_aiter_ops:
# Check if the env variable is set
_AITER_ENABLED = envs.VLLM_ROCM_USE_AITER
_CUSTOM_ALL_REDUCE_ENABLED = envs.VLLM_ROCM_USE_AITER_CUSTOM_AR
_LINEAR_ENABLED = envs.VLLM_ROCM_USE_AITER_LINEAR
_FMOE_ENABLED = envs.VLLM_ROCM_USE_AITER_MOE
_MLA_ENABLED = envs.VLLM_ROCM_USE_AITER_MLA
@@ -1598,9 +1562,6 @@ class rocm_aiter_ops:
# num_shared_experts / shared_expert_scoring_func args (7-arg form).
_TOPK_SOFTMAX_FUSED_SIGMOID: bool | None = None
_ALL_REDUCE_MAX_SIZE: int = 8192 * 1024 * 8 * 2
_CUSTOM_ALL_REDUCE: AiterCustomAllreduceProto | None = None
@classmethod
def refresh_env_variables(cls):
"""
@@ -1611,6 +1572,7 @@ class rocm_aiter_ops:
you can call this function to reload the env variables.
"""
cls._AITER_ENABLED = envs.VLLM_ROCM_USE_AITER
cls._CUSTOM_ALL_REDUCE_ENABLED = envs.VLLM_ROCM_USE_AITER_CUSTOM_AR
cls._LINEAR_ENABLED = envs.VLLM_ROCM_USE_AITER_LINEAR
cls._FMOE_ENABLED = envs.VLLM_ROCM_USE_AITER_MOE
cls._MLA_ENABLED = envs.VLLM_ROCM_USE_AITER_MLA
@@ -1770,6 +1732,11 @@ class rocm_aiter_ops:
def is_mha_enabled(cls) -> bool:
return cls._AITER_ENABLED and cls._MHA_ENABLED
@classmethod
@if_aiter_supported
def is_custom_all_reduce_enabled(cls) -> bool:
return cls._AITER_ENABLED and cls._CUSTOM_ALL_REDUCE_ENABLED
@classmethod
@if_aiter_supported
def is_shuffle_kv_cache_enabled(cls) -> bool:
@@ -1824,33 +1791,20 @@ class rocm_aiter_ops:
return cls.is_linear_enabled() and on_gfx950()
@classmethod
def initialize_aiter_allreduce(
cls, group: ProcessGroup, device: torch.device
) -> None:
try:
from aiter.dist.device_communicators.custom_all_reduce import (
CustomAllreduce as AiterCustomAllreduce,
)
def get_aiter_allreduce(cls):
"""Return the TP device communicator's AITER custom-allreduce if it has
one, return None otherwise
"""
from vllm.distributed.device_communicators.aiter_custom_all_reduce import (
AiterCustomAllreduce,
)
from vllm.distributed.parallel_state import get_tp_group
cls._CUSTOM_ALL_REDUCE = AiterCustomAllreduce(group, device)
except Exception:
cls._CUSTOM_ALL_REDUCE = None
@classmethod
def get_aiter_allreduce(cls) -> AiterCustomAllreduceProto | None:
return cls._CUSTOM_ALL_REDUCE
@classmethod
def destroy_aiter_allreduce(cls) -> None:
if cls._CUSTOM_ALL_REDUCE is not None:
cls._CUSTOM_ALL_REDUCE.close()
cls._CUSTOM_ALL_REDUCE = None
@classmethod
def get_aiter_allreduce_max_size(cls) -> int | None:
# effective max input size (based on upstream aiter version: v0.1.10.post3)
# https://github.com/ROCm/aiter/blob/6a0e7b26ccf33164785531212cc2ec2cde0b9243/aiter/dist/device_communicators/custom_all_reduce.py#L272-L273
return int(cls._ALL_REDUCE_MAX_SIZE / 2)
device_comm = get_tp_group().device_communicator
aiter_ar_comm = getattr(device_comm, "aiter_ar_comm", None)
return (
aiter_ar_comm if isinstance(aiter_ar_comm, AiterCustomAllreduce) else None
)
@classmethod
@if_aiter_supported
@@ -2165,21 +2119,6 @@ class rocm_aiter_ops:
def get_fused_allreduce_rmsnorm_quant_per_group_with_bf16_norm_op() -> OpOverload: # noqa: E501
return torch.ops.vllm.rocm_aiter_fused_allreduce_rmsnorm_quant_per_group_with_bf16_norm.default # noqa: E501
# TODO(frida-andersson): drop once vLLM pins AITER >= 0.1.14 (ROCm/aiter#2823).
@classmethod
def has_fused_allreduce_rmsnorm_quant_per_group(cls) -> bool:
"""True if the running AITER build exposes the per-group AR+RMS+quant
kernel (added in ROCm/aiter PR #2823).
The pattern registration in ``RocmAiterAllReduceFusionPass`` keys off
this so vLLM degrades to the AR+RMS-only fusion when run against an
older aiter that lacks the per-group launcher.
"""
aiter_ar = cls.get_aiter_allreduce()
return aiter_ar is not None and hasattr(
aiter_ar, "fused_ar_rms_per_group_quant"
)
@staticmethod
def get_fused_mla_dual_rms_norm_op() -> OpOverload:
return torch.ops.vllm.fused_mla_dual_rms_norm.default
+18
View File
@@ -2974,6 +2974,24 @@ if hasattr(torch.ops._C, "fused_experts_cpu"):
return torch.empty_like(hidden_states)
if hasattr(torch.ops._C, "dynamic_4bit_int_moe"):
@register_fake("_C::dynamic_4bit_int_moe")
def dynamic_4bit_int_moe_fake(
x: torch.Tensor,
topk_ids: torch.Tensor,
topk_weights: torch.Tensor,
w13_packed: torch.Tensor,
w2_packed: torch.Tensor,
hidden_size: int,
intermediate_size: int,
group_size: int,
apply_router_weight_on_input: bool,
activation_kind: int,
) -> torch.Tensor:
return x.new_empty((x.size(0), hidden_size))
def fused_experts_cpu(
hidden_states: torch.Tensor,
w1: torch.Tensor,
@@ -19,7 +19,6 @@ from vllm.compilation.passes.fusion.rms_quant_fusion import (
from vllm.config import VllmConfig
from vllm.config.utils import Range
from vllm.distributed import get_tp_group, tensor_model_parallel_all_reduce
from vllm.distributed.device_communicators.custom_all_reduce import CustomAllreduce
from vllm.distributed.parallel_state import (
get_tensor_model_parallel_rank,
get_tensor_model_parallel_world_size,
@@ -129,6 +128,29 @@ _FI_ALLREDUCE_ONE_SHOT_MAX_SIZES_MB: dict[int, dict[int, float]] = {
},
}
MiB = 1024 * 1024
def _select_flashinfer_allreduce_use_oneshot(
workspace_backend: str,
device_capability: int | None,
world_size: int,
current_tensor_size: int,
) -> bool | None:
if workspace_backend == "mnnvl":
# FlashInfer sizes MNNVL workspaces around its own AUTO strategy.
# Forcing vLLM's per-rank threshold can request one-shot for tensors
# larger than the MNNVL one-shot workspace.
return None
if device_capability is None:
max_one_shot_size = None
else:
max_one_shot_size = _FI_ALLREDUCE_ONE_SHOT_MAX_SIZES_MB.get(
device_capability, {}
).get(world_size)
return max_one_shot_size is None or current_tensor_size <= max_one_shot_size * MiB
if flashinfer_comm is not None:
from vllm.distributed.device_communicators.flashinfer_all_reduce import (
@@ -139,8 +161,6 @@ if flashinfer_comm is not None:
ar_fusion_patterns = flashinfer_comm.AllReduceFusionPattern
MiB = 1024 * 1024
def call_trtllm_fused_allreduce_norm(
allreduce_in: torch.Tensor,
residual: torch.Tensor,
@@ -175,16 +195,6 @@ if flashinfer_comm is not None:
)
curr_device = current_platform.get_device_capability()
device_capability = curr_device.to_int() if curr_device is not None else None
# Get one shot input size limit for the current world size
# for the current device capability
max_one_shot_size = _FI_ALLREDUCE_ONE_SHOT_MAX_SIZES_MB.get(
device_capability, # type: ignore[arg-type, unused-ignore]
{},
).get(world_size, None)
# Use one shot if no max size is specified
use_oneshot = (
max_one_shot_size is None or current_tensor_size <= max_one_shot_size * MiB
)
# Select workspace based on pattern: quant patterns use the
# trtllm quant workspace, non-quant patterns use the primary workspace.
@@ -206,6 +216,12 @@ if flashinfer_comm is not None:
assert workspace is not None, (
"Flashinfer allreduce workspace must be initialized when using flashinfer"
)
use_oneshot = _select_flashinfer_allreduce_use_oneshot(
workspace.backend,
device_capability,
world_size,
current_tensor_size,
)
assert flashinfer_comm is not None
if norm_out is None:
norm_out = allreduce_in
@@ -249,7 +265,7 @@ if flashinfer_comm is not None:
# the end for the one-shot path; the two-shot path is synchronized
# and keeps the early completion. Related one-shot instability in
# the same kernel: flashinfer-ai/flashinfer#1223.
trigger_completion_at_end=use_oneshot
trigger_completion_at_end=(use_oneshot is True)
or num_tokens > PDL_ADVANCE_LAUNCH_TOKENS,
)
@@ -1473,39 +1489,23 @@ class RocmAiterAllReduceFusionPass(VllmFusionPatternMatcherPass):
)
return
device_comm = get_tp_group().device_communicator
if device_comm is None:
logger.warning_once("Device communicator is required.")
return
ca_comm = getattr(device_comm, "ca_comm", None)
ca_comm = rocm_aiter_ops.get_aiter_allreduce()
if ca_comm is None:
logger.warning_once("Custom Allreduce is required.")
logger.warning_once(
"AITER allreduce fusions are disabled "
"because AITER Custom All Reduce is not enabled. "
"Set VLLM_ROCM_USE_AITER_CUSTOM_AR=1 "
"to enable it."
)
return
self.ca_comm = ca_comm
assert isinstance(ca_comm, CustomAllreduce)
group = get_tp_group().cpu_group
rocm_aiter_ops.initialize_aiter_allreduce(group, self.device)
hidden_dim = config.model_config.get_hidden_size()
element_size = torch.tensor([], dtype=self.model_dtype).element_size()
max_size = rocm_aiter_ops.get_aiter_allreduce_max_size()
if max_size is None:
logger.warning("AITER allreduce fusion must be initialized")
return
# Aiter's fused_allreduce_rmsnorm kernel dispatches on hidden_dim.
# Before aiter v0.1.12 the launcher was template-specialized on HIDDEN_DIM
# and silently no-op'd for sizes outside {512, 1024, 2048, 4096}. From v0.1.12
# hidden_dim is a runtime argument. Detect the older API via the missing
# `_pool` attribute and skip fusion for unsupported sizes.
# Ref (old kernel): https://github.com/ROCm/aiter/blob/6a0e7b26ccf33164785531212cc2ec2cde0b9243/csrc/include/custom_all_reduce.cuh#L2590
aiter_ar = rocm_aiter_ops.get_aiter_allreduce()
max_size = ca_comm.effective_max_size()
_AITER_OLD_FUSED_AR_RMS_HIDDEN = (512, 1024, 2048, 4096)
if (
aiter_ar is not None
and not hasattr(aiter_ar, "_pool")
not ca_comm.supports_dynamic_hidden_dim
and hidden_dim not in _AITER_OLD_FUSED_AR_RMS_HIDDEN
):
logger.warning_once(
@@ -1515,10 +1515,6 @@ class RocmAiterAllReduceFusionPass(VllmFusionPatternMatcherPass):
_AITER_OLD_FUSED_AR_RMS_HIDDEN,
hidden_dim,
)
# Tear down aiter's custom-allreduce so its IPC handles don't
# race with vllm's ca_comm on the unfused fallback path.
with contextlib.suppress(Exception):
rocm_aiter_ops.destroy_aiter_allreduce()
return
max_token_num = max_size // (hidden_dim * element_size)
@@ -1532,9 +1528,7 @@ class RocmAiterAllReduceFusionPass(VllmFusionPatternMatcherPass):
# fall back to the AR+RMS-only fusion paired with PR #41825's
# standalone RMS+quant fusion -- still correct, just leaves the
# post-AR quant as a standalone kernel.
supports_per_group_quant = (
rocm_aiter_ops.has_fused_allreduce_rmsnorm_quant_per_group()
)
supports_per_group_quant = ca_comm.supports_per_group_quant
if not supports_per_group_quant:
logger.warning_once(
"AITER AR+RMS+per-group-FP8-quant fusion disabled: aiter "
@@ -1609,9 +1603,3 @@ class RocmAiterAllReduceFusionPass(VllmFusionPatternMatcherPass):
logger.debug(
"%s Replaced %s patterns", self.__class__.__name__, self.matched_count
)
def __del__(self) -> None:
if getattr(self, "disabled", True):
return
with contextlib.suppress(Exception):
rocm_aiter_ops.destroy_aiter_allreduce()
+1 -1
View File
@@ -1252,7 +1252,7 @@ class ModelConfig:
def is_deepseek_mla(self) -> bool:
return self.model_arch_config.is_deepseek_mla
@property
@cached_property
def is_mm_prefix_lm(self) -> bool:
return self.model_arch_config.is_mm_prefix_lm
+2 -3
View File
@@ -632,8 +632,8 @@ class ParallelConfig:
# The all_reduce at the end of attention (during o_proj) means that
# inputs are replicated across each rank of the tensor parallel group.
# If using expert-parallelism with DeepEP All2All ops, replicated
# tokens results in useless duplicate computation and communication.
# If using expert-parallelism, replicated tokens results in useless
# duplicate computation and communication.
#
# In this case, ensure the input to the experts is sequence parallel
# to avoid the excess work.
@@ -652,7 +652,6 @@ class ParallelConfig:
)
and self.enable_expert_parallel
and self.tensor_parallel_size > 1
and self.data_parallel_size > 1
)
@property
+16
View File
@@ -46,6 +46,7 @@ MTPModelTypes = Literal[
"qwen3_5_mtp",
"longcat_flash_mtp",
"minimax_m3_mtp",
"bailing_hybrid_mtp",
"mtp",
"pangu_ultra_moe_mtp",
"step3p5_mtp",
@@ -463,6 +464,21 @@ class SpeculativeConfig:
{"n_predict": n_predict, "architectures": ["Qwen3NextMTP"]}
)
architectures = getattr(hf_config, "architectures", []) or []
if (
hf_config.model_type == "bailing_hybrid"
or "BailingMoeV2_5ForCausalLM" in architectures
):
hf_config.model_type = "bailing_hybrid_mtp"
if hf_config.model_type == "bailing_hybrid_mtp":
n_predict = getattr(hf_config, "num_nextn_predict_layers", None)
hf_config.update(
{
"n_predict": n_predict,
"architectures": ["BailingMoeV25MTPModel"],
}
)
if hf_config.model_type == "exaone_moe":
hf_config.model_type = "exaone_moe_mtp"
if hf_config.model_type == "exaone_moe_mtp":
+16 -2
View File
@@ -1850,8 +1850,13 @@ class VllmConfig:
tp_size = self.parallel_config.tensor_parallel_size
from vllm._aiter_ops import rocm_aiter_ops
if rocm_aiter_ops.is_enabled():
max_size = rocm_aiter_ops.get_aiter_allreduce_max_size()
max_size: int | None = None
if rocm_aiter_ops.is_custom_all_reduce_enabled():
from vllm.distributed.device_communicators.aiter_custom_all_reduce import ( # noqa: E501
AiterCustomAllreduce,
)
max_size = AiterCustomAllreduce.effective_max_size()
else:
max_size = compilation_config.pass_config.flashinfer_max_size(tp_size)
if max_size is not None and self.model_config is not None:
@@ -1936,12 +1941,21 @@ class VllmConfig:
if architecture is None:
return
from vllm.model_executor.models import ModelRegistry
from vllm.model_executor.models.config import (
MODELS_CONFIG_MAP,
HybridAttentionMambaModelConfig,
)
cls = MODELS_CONFIG_MAP.get(architecture, None)
if cls is None:
# `architecture` may be an HF base-model name (e.g. "Mamba2Model"
# when `architectures` is omitted); normalize to the resolved arch
# so per-arch config hooks are not skipped.
architecture = ModelRegistry._normalize_arch(
architecture, self.model_config
)
cls = MODELS_CONFIG_MAP.get(architecture, None)
if cls is not None:
cls.verify_and_update_config(self)
@@ -0,0 +1,95 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""vLLM-owned wrapper over AITER's ``CustomAllreduce``.
vLLM's ``CudaCommunicator`` stores one of these as ``aiter_ar_comm`` (when
``VLLM_ROCM_USE_AITER_CUSTOM_AR`` is set) so the plain allreduce and
the fused allreduce+RMSNorm path share a single AITER instance with its IPC buffers.
"""
import torch
from torch.distributed import ProcessGroup
from vllm.logger import init_logger
logger = init_logger(__name__)
class AiterCustomAllreduce:
# Default IPC buffer size for AITER's CustomAllreduce.
MAX_SIZE: int = 8192 * 1024 * 8 * 2
@classmethod
def effective_max_size(cls) -> int:
"""
Max input byte size eligible for AITER custom allreduce.
"""
return cls.MAX_SIZE // 2
def __init__(
self,
group: ProcessGroup,
device: int | str | torch.device,
max_size: int | None = None,
):
from aiter.dist.device_communicators.custom_all_reduce import (
CustomAllreduce as _AiterCustomAllreduce,
)
if max_size is None:
max_size = self.MAX_SIZE
self._impl = _AiterCustomAllreduce(group, device, max_size=max_size)
@property
def aiter_ca(self):
return self._impl
@property
def disabled(self) -> bool:
return self._impl.disabled
def should_custom_ar(self, inp: torch.Tensor) -> bool:
return self._impl.should_custom_ar(inp)
def custom_all_reduce(self, inp: torch.Tensor) -> torch.Tensor | None:
return self._impl.custom_all_reduce(inp)
def capture(self):
return self._impl.capture()
def close(self) -> None:
self._impl.close()
@property
def supports_dynamic_hidden_dim(self) -> bool:
"""Aiter's fused_allreduce_rmsnorm kernel dispatches on hidden_dim.
Before aiter v0.1.12 the launcher was template-specialized on HIDDEN_DIM
and silently no-op'd for sizes outside {512, 1024, 2048, 4096}. From v0.1.12
hidden_dim is a runtime argument. Older builds are detected via
AiterCustomAllreduce.supports_dynamic_hidden_dim; This function is used to
skip fusion for unsupported sizes on them.
Ref (old kernel): https://github.com/ROCm/aiter/blob/6a0e7b26ccf33164785531212cc2ec2cde0b9243/csrc/include/custom_all_reduce.cuh#L2590
"""
return hasattr(self._impl, "_pool")
@staticmethod
def build_supports_per_group_quant() -> bool:
"""True if the running AITER build exposes the per-group AR+RMS+quant
kernel (added in ROCm/aiter PR #2823).
The pattern registration in ``RocmAiterAllReduceFusionPass`` keys off
this so vLLM degrades to the AR+RMS-only fusion when run against an
older aiter that lacks the per-group launcher.
"""
from aiter.dist.device_communicators.custom_all_reduce import (
CustomAllreduce as _AiterCustomAllreduce,
)
return hasattr(_AiterCustomAllreduce, "fused_ar_rms_per_group_quant")
# TODO(frida-andersson): drop once vLLM pins AITER >= 0.1.14 (ROCm/aiter#2823).
@property
def supports_per_group_quant(self) -> bool:
return self.build_supports_per_group_quant()
@@ -175,10 +175,11 @@ class DeviceCommunicatorBase:
config = get_current_vllm_config_or_none()
if config is not None:
# as long as we use data parallel (coupled data parallel
# where all data parallel ranks execute forward together),
# we initialize the all2all manager used in expert parallel.
use_ep = config.parallel_config.data_parallel_size > 1
# initialize the all2all manager for DP or sequence-parallel EP.
use_ep = (
config.parallel_config.data_parallel_size > 1
or config.parallel_config.use_sequence_parallel_moe
)
all2all_backend = config.parallel_config.all2all_backend
self.is_ep_communicator = unique_name.split(":")[0] == "ep"
@@ -6,6 +6,7 @@ import torch
from torch.distributed import ProcessGroup
import vllm.envs as envs
from vllm._aiter_ops import rocm_aiter_ops
from vllm.distributed.device_communicators.all_reduce_utils import (
NCCL_SYMM_MEM_ALL_REDUCE_CONFIG,
should_nccl_symm_mem_ag_rs,
@@ -19,6 +20,7 @@ from vllm.logger import init_logger
from vllm.platforms import current_platform
from ..utils import StatelessProcessGroup
from .aiter_custom_all_reduce import AiterCustomAllreduce
from .base_device_communicator import DeviceCommunicatorBase
logger = init_logger(__name__)
@@ -48,16 +50,21 @@ class CudaCommunicator(DeviceCommunicatorBase):
use_custom_allreduce = False
use_torch_symm_mem = False
use_flashinfer_allreduce = False
use_aiter_allreduce = False
else:
from vllm.distributed.parallel_state import _ENABLE_CUSTOM_ALL_REDUCE
use_custom_allreduce = _ENABLE_CUSTOM_ALL_REDUCE
use_torch_symm_mem = envs.VLLM_ALLREDUCE_USE_SYMM_MEM
use_flashinfer_allreduce = envs.VLLM_ALLREDUCE_USE_FLASHINFER
use_aiter_allreduce = use_custom_allreduce and bool(
rocm_aiter_ops.is_custom_all_reduce_enabled()
)
self.use_custom_allreduce = use_custom_allreduce
self.use_torch_symm_mem = use_torch_symm_mem
self.use_flashinfer_allreduce = use_flashinfer_allreduce
self.use_aiter_allreduce = use_aiter_allreduce
# lazy import to avoid documentation build error
from vllm.distributed.device_communicators.custom_all_reduce import (
@@ -85,6 +92,7 @@ class CudaCommunicator(DeviceCommunicatorBase):
self.qr_comm: QuickAllReduce | None = None
self.symm_mem_comm: SymmMemCommunicator | None = None
self.fi_ar_comm: FlashInferAllReduce | None = None
self.aiter_ar_comm: AiterCustomAllreduce | None = None
if use_torch_symm_mem and current_platform.is_cuda():
self.symm_mem_comm = SymmMemCommunicator(
@@ -98,7 +106,13 @@ class CudaCommunicator(DeviceCommunicatorBase):
device=self.device,
)
if use_custom_allreduce and self.world_size > 1:
if self.use_aiter_allreduce and self.world_size > 1:
self.aiter_ar_comm = AiterCustomAllreduce(
group=self.cpu_group,
device=self.device,
)
if use_custom_allreduce and self.aiter_ar_comm is None and self.world_size > 1:
# Initialize a custom fast all-reduce implementation.
self.ca_comm = CustomAllreduce(
group=self.cpu_group,
@@ -108,13 +122,14 @@ class CudaCommunicator(DeviceCommunicatorBase):
),
)
if current_platform.is_rocm():
# Initialize a custom quick all-reduce implementation for AMD.
# Quick reduce is designed as a complement to custom allreduce.
# Based on quickreduce (https://github.com/mk1-project/quickreduce).
# If it's a rocm, 'use_custom_allreduce==True' means it must
# currently be an MI300 series.
self.qr_comm = QuickAllReduce(group=self.cpu_group, device=self.device)
if use_custom_allreduce and self.world_size > 1 and current_platform.is_rocm():
# Initialize a custom quick all-reduce implementation for AMD.
# Quick reduce is designed as a complement to custom allreduce
# (vLLM's or AITER's), so it is initialized for either backend.
# Based on quickreduce (https://github.com/mk1-project/quickreduce).
# On ROCm, 'use_custom_allreduce==True' means it must currently be
# an MI300 series.
self.qr_comm = QuickAllReduce(group=self.cpu_group, device=self.device)
if self.world_size > 1:
self._log_all_reduce_backend_selection()
@@ -203,6 +218,7 @@ class CudaCommunicator(DeviceCommunicatorBase):
"NCCL_SYMM_MEM",
"QUICK_REDUCE",
"FLASHINFER",
"AITER_CUSTOM",
"CUSTOM",
"SYMM_MEM",
"PYNCCL",
@@ -236,6 +252,8 @@ class CudaCommunicator(DeviceCommunicatorBase):
enabled_ar_backends.append("QUICK_REDUCE")
if self.fi_ar_comm is not None and not self.fi_ar_comm.disabled:
enabled_ar_backends.append("FLASHINFER")
if self.aiter_ar_comm is not None and not self.aiter_ar_comm.disabled:
enabled_ar_backends.append("AITER_CUSTOM")
if self.ca_comm is not None and not self.ca_comm.disabled:
enabled_ar_backends.append("CUSTOM")
if self.symm_mem_comm is not None and not self.symm_mem_comm.disabled:
@@ -261,8 +279,8 @@ class CudaCommunicator(DeviceCommunicatorBase):
out = torch.ops.vllm.all_reduce_symmetric_with_copy(input_)
if out is not None:
return out
# always try quick reduce first, then flashinfer, then custom allreduce,
# and then pynccl. (quick reduce just for ROCM MI3*)
# always try quick reduce first, then flashinfer, then the AITER or vLLM
# custom allreduce, and then pynccl. (quick reduce just for ROCM MI3*)
qr_comm = self.qr_comm
if (
qr_comm is not None
@@ -281,6 +299,15 @@ class CudaCommunicator(DeviceCommunicatorBase):
out = fi_ar_comm.all_reduce(input_)
assert out is not None
return out
aiter_ar_comm = self.aiter_ar_comm
if (
aiter_ar_comm is not None
and not aiter_ar_comm.disabled
and aiter_ar_comm.should_custom_ar(input_)
):
out = aiter_ar_comm.custom_all_reduce(input_)
assert out is not None
return out
ca_comm = self.ca_comm
if (
ca_comm is not None
@@ -509,6 +536,9 @@ class CudaCommunicator(DeviceCommunicatorBase):
self.pynccl_comm = None
if self.ca_comm is not None:
self.ca_comm = None
if self.aiter_ar_comm is not None:
self.aiter_ar_comm.close()
self.aiter_ar_comm = None
if self.fi_ar_comm is not None:
self.fi_ar_comm.destroy()
self.fi_ar_comm = None
+2 -2
View File
@@ -742,8 +742,8 @@ class EplbState:
is_main_rank = ep_rank == 0
if is_main_rank:
if not self.is_async or is_profile:
start_event = torch.Event(enable_timing=True)
end_event = torch.Event(enable_timing=True)
start_event = torch.cuda.Event(enable_timing=True)
end_event = torch.cuda.Event(enable_timing=True)
start_event.record()
logger.info(
"Rearranging experts %s %s...",
+2 -2
View File
@@ -31,7 +31,7 @@ class CpuGpuEvent:
"""
def __init__(self):
self._event = torch.Event()
self._event = torch.cuda.Event()
self._recorded = threading.Event()
def wait(self, stream: torch.cuda.Stream | None = None):
@@ -56,7 +56,7 @@ class CpuGpuEvent:
"CpuGpuEvent.record() called before the previous event was "
"consumed by wait()"
)
self._event = torch.Event()
self._event = torch.cuda.Event()
self._event.record(stream)
self._recorded.set()
@@ -239,7 +239,7 @@ class ExampleHiddenStatesConnector(KVConnectorBase_V1, SupportsHMA):
# this event is complete the request is considered "done sending"
# by get_finished; clients block on the per-file flock to wait for
# the disk write itself.
self._req_copy_events: dict[str, torch.Event] = {}
self._req_copy_events: dict[str, torch.cuda.Event] = {}
# req_ids reported as finished-generating by the scheduler,
# accumulated across get_finished calls.
self._accumulated_finished_req_ids: set[str] = set()
@@ -320,7 +320,7 @@ class ExampleHiddenStatesConnector(KVConnectorBase_V1, SupportsHMA):
@staticmethod
def _write_tensors(
tensors: dict[str, torch.Tensor],
event: torch.Event,
event: torch.cuda.Event,
filename: str,
lock_fd: int | None,
) -> None:
@@ -375,7 +375,7 @@ class ExampleHiddenStatesConnector(KVConnectorBase_V1, SupportsHMA):
copy_stream = self._get_copy_stream()
# Ensure the copy stream sees all prior writes on the default stream.
ready_event = torch.Event()
ready_event = torch.cuda.Event()
ready_event.record()
copy_stream.wait_event(ready_event)
@@ -396,7 +396,7 @@ class ExampleHiddenStatesConnector(KVConnectorBase_V1, SupportsHMA):
pinned_hs.copy_(hidden_states_gpu, non_blocking=True)
# Record completion of this copy on the copy stream.
copy_done = torch.Event()
copy_done = torch.cuda.Event()
copy_done.record(copy_stream)
# token_ids is already on CPU (created in request_finished).
@@ -221,7 +221,7 @@ class Hf3fsClient:
@wsynchronized()
def batch_write(
self, offsets: list[int], tensors: list[torch.Tensor], event: torch.Event
self, offsets: list[int], tensors: list[torch.Tensor], event: torch.cuda.Event
) -> list[int]:
"""Write data from tensors to the file at specified offsets.
@@ -133,7 +133,7 @@ class AsyncOperationManager:
# CUDA streams for async operations
self._save_stream = torch.cuda.Stream()
self._load_stream = torch.cuda.Stream()
self._save_event = torch.Event()
self._save_event = torch.cuda.Event()
# Buffer allocators for data copying
self._save_buffer_allocator = CopyBufferAllocator(
@@ -171,7 +171,7 @@ class AsyncOperationManager:
def submit_save_operation(self, request_id: str, block_ids, block_hashes) -> Future:
"""Submit a save operation for async execution."""
future: Future[Any] = Future()
main_stream_event = torch.Event()
main_stream_event = torch.cuda.Event()
main_stream_event.record()
task = (request_id, block_ids, block_hashes, future, main_stream_event)
self._save_queue.put(task)
@@ -304,7 +304,7 @@ class AsyncOperationManager:
block_ids, buffers, "gather"
)
save_stream_event = torch.Event()
save_stream_event = torch.cuda.Event()
save_stream_event.record(self._save_stream) # Record gather completion
# Step3: Write data in batches
@@ -75,7 +75,7 @@ class Hf3fsClient:
return torch.frombuffer(buffer_data, dtype=dtype)
def batch_write(
self, offsets: list[int], tensors: list[torch.Tensor], event: torch.Event
self, offsets: list[int], tensors: list[torch.Tensor], event: torch.cuda.Event
) -> list[int]:
"""Write data from tensors to file at specified offsets."""
results = []
@@ -430,7 +430,7 @@ class LMCacheMPWorkerAdapter:
@_lmcache_nvtx_annotate
def submit_store_request(
self, request_id: str, op: LoadStoreOp, event: torch.Event
self, request_id: str, op: LoadStoreOp, event: torch.cuda.Event
):
"""
Submit a KV cache store request to LMCache
@@ -464,7 +464,7 @@ class LMCacheMPWorkerAdapter:
@_lmcache_nvtx_annotate
def submit_retrieve_request(
self, request_id: str, op: LoadStoreOp, event: torch.Event
self, request_id: str, op: LoadStoreOp, event: torch.cuda.Event
):
"""
Submit a KV cache retrieve request to LMCache
@@ -501,7 +501,7 @@ class LMCacheMPWorkerAdapter:
self,
request_ids: list[str],
ops: list[LoadStoreOp],
event: torch.Event,
event: torch.cuda.Event,
):
"""
Submit a batched store request to LMCache
@@ -550,7 +550,7 @@ class LMCacheMPWorkerAdapter:
self,
request_ids: list[str],
ops: list[LoadStoreOp],
event: torch.Event,
event: torch.cuda.Event,
):
"""
Submit a batched retrieve request to LMCache
@@ -589,7 +589,7 @@ class LMCacheMPConnectorUpstream(KVConnectorBase_V1):
return
with torch.cuda.stream(torch.cuda.current_stream()):
event = torch.Event(interprocess=True)
event = torch.cuda.Event(interprocess=True)
event.record()
self.worker_adapter.batched_submit_retrieve_requests(
@@ -663,7 +663,7 @@ class LMCacheMPConnectorUpstream(KVConnectorBase_V1):
return
with torch.cuda.stream(torch.cuda.current_stream()):
event = torch.Event(interprocess=True)
event = torch.cuda.Event(interprocess=True)
event.record()
self.worker_adapter.batched_submit_store_requests(
@@ -323,7 +323,7 @@ class ReqMeta:
can_save: bool | None = None
load_spec: LoadSpec | None = None
is_last_chunk: bool | None = None
current_event: torch.Event | None = None
current_event: torch.cuda.Event | None = None
token_ids: list[int] | None = None
num_prompt_tokens: int | None = None
@@ -1357,7 +1357,7 @@ class MooncakeStoreWorker:
current_event = None
for request in meta.requests:
if request.can_save:
current_event = torch.Event()
current_event = torch.cuda.Event()
current_event.record()
break
@@ -56,7 +56,7 @@ class WriteTask:
local_block_ids: list[int]
remote_block_ids_hint: list[int] | None
layer_name: str
event: torch.Event
event: torch.cuda.Event
remote_notify_port: int
remote_ip: str
enqueue_time: float = field(default_factory=time.perf_counter)
@@ -1061,7 +1061,7 @@ class MoRIIOConnectorWorker:
# when mori-io supports ibgda functionality
stream = torch.cuda.current_stream()
event = torch.Event()
event = torch.cuda.Event()
event.record(stream)
task = WriteTask(

Some files were not shown because too many files have changed in this diff Show More