From be71c7899872b01e5803645dcbfbeaf51ff4f44b Mon Sep 17 00:00:00 2001 From: dev Date: Tue, 29 Sep 2026 20:49:22 +0000 Subject: [PATCH 1/3] spec(ENG-HOST-EMBEDDING): the host-resident token table A dense forward that owns a device embedding table pays [vocab, H] of device memory for it even when the gather could run in host RAM. Record the row, the structured spec the record checker requires of it, the seam that takes both arms, the llama.cpp ggml_get_rows shape this is a secondary oracle port of, and the golden vectors the host arm is gated on. The row addition moves the engine-matrix counts and the pinned ENGINE_ROWS, and carries the CLAIM-ENG-HOST-EMBEDDING record the ACTIVE state requires. FOLLOWING_AGENTS_PROTOCOL Following-Agents-Protocol: true AI-Assisted: true Assisted-by: AGENT:opencode-go/deepseek-v4.1-flash [pi] --- .agents/claims/CLAIM-ENG-HOST-EMBEDDING.md | 5 + .agents/engine-matrix.md | 5 +- .../ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md | 19 +++ .agents/specs/host-embedding.md | 141 ++++++++++++++++++ scripts/check-agent-record.py | 4 +- tests/scripts/test_agent_record.py | 42 ++++++ 6 files changed, 213 insertions(+), 3 deletions(-) create mode 100644 .agents/claims/CLAIM-ENG-HOST-EMBEDDING.md create mode 100644 .agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md create mode 100644 .agents/specs/host-embedding.md diff --git a/.agents/claims/CLAIM-ENG-HOST-EMBEDDING.md b/.agents/claims/CLAIM-ENG-HOST-EMBEDDING.md new file mode 100644 index 000000000..be8c4abd3 --- /dev/null +++ b/.agents/claims/CLAIM-ENG-HOST-EMBEDDING.md @@ -0,0 +1,5 @@ +# CLAIM-ENG-HOST-EMBEDDING + +| Claim | Row IDs | Agent | Worktree / remote dir | Branch | Owned scope | State | Last update | +|---|---|---|---|---|---|---|---| +| `CLAIM-ENG-HOST-EMBEDDING` | `ENG-HOST-EMBEDDING` (`ACTIVE`) | pi (deepseek-v4.1-flash), helper role — holds the row's spec, the host arm and its test; a fresh reviewer over the immutable head still owes its own pass | isolated worktree `/tmp/wt-embed`, Release CPU build plus the local `sm_120a` card (RTX PRO 4000 Blackwell, 24 GiB, a NON-fleet device). NO fleet lease, NO oracle run, NO benchmark, NO model weights loaded | `row/ENG-HOST-EMBEDDING`, PR [#3356](https://github.com/mudler/vllm.cpp/pull/3356) | Owns ONLY: `.agents/specs/host-embedding.md`, the `ENG-HOST-EMBEDDING` row and its two counts in `.agents/engine-matrix.md`, this claim, issue `ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM`, and the implementation it carries (`host_embedding.{h,cpp}`, the `EmbedGather` call sites, `docs/ENVIRONMENT.md`, `tests/vllm/models/test_host_embedding.cpp`). EXCLUDES every other row's records and every model/capability the seam is merely wired into; EXCLUDES the end-to-end VRAM/throughput measurement, which the row owes | `ACTIVE` | 2026-09-29 — row claimed; spec, implementation, test and docs land in ONE pull request per the recorded Git-integration preference | diff --git a/.agents/engine-matrix.md b/.agents/engine-matrix.md index 152b527c1..24c35b88f 100644 --- a/.agents/engine-matrix.md +++ b/.agents/engine-matrix.md @@ -47,8 +47,8 @@ forensics: roadmap_v1.md and the parity ledger. | Serving, API, CLI, library | 39 | 10 | 3 | 0 | 4 | 12 | 3 | 3 | 4 | | LoRA and adapters | 2 | 0 | 0 | 0 | 0 | 1 | 0 | 0 | 1 | | Long context and attention | 11 | 5 | 0 | 0 | 1 | 0 | 1 | 0 | 4 | -| Loading, tokenizer, config | 15 | 3 | 4 | 0 | 1 | 1 | 2 | 2 | 2 | -| **Total** | **179** | **35** | **19** | **4** | **13** | **44** | **10** | **12** | **41** | +| Loading, tokenizer, config | 16 | 3 | 4 | 0 | 1 | 2 | 2 | 2 | 2 | +| **Total** | **180** | **35** | **19** | **4** | **13** | **45** | **10** | **12** | **41** | ## Engine core and scheduling @@ -279,6 +279,7 @@ claims it. |---|---|---|---|---|---|---|---|---| | `LOAD-SAFETENSORS` | General safetensors loading, stacked parameters, weights mapping. Binding memory scan finds all source mappings live through load plus a persistent CPU mirror of selected tensors; windowed progressive `madvise(DONTNEED)` release now drops each copied-then-dead source range during the copy loop so the mirror build no longer double-resides with the full source mmap | T0 | streamed target-device construction `vllm/model_executor/model_loader/base_loader.py:43-82`; incremental `safe_open` yield `vllm/model_executor/model_loader/weight_utils.py:905-954`; immediate parameter copy `vllm/model_executor/models/utils.py:170-180,252-279` | all mappings opened/retained `src/vllm/entrypoints/model_loader.cpp:47-63,303-312`; `MAP_PRIVATE` reader `src/vllm/model_executor/model_loader/safetensors_reader.cpp:43-70`; windowed release primitive/gate `src/vllm/model_executor/model_loader/safetensors_reader.cpp:285-340`; copy-helper instrumentation `src/vllm/model_executor/models/qwen3_5_weights.cpp:88,104,119,171,190,226,231` + `qwen3_5_dense_weights.cpp:64,80,109,159,164,263,323` | [reader contracts + windowed-release cases](../tests/vllm/test_safetensors.cpp#L68) green (smaps-Rss drop, byte-identity, gate semantics, neighbor safety); exact 27B accounting finds **1,155 tensors / 24,610,136,064 B (22.920 GiB)** persistent host bytes. **VmHWM A/B MEASURED** at `cb2d310` (root `~/work/vllm.cpp-windowed-load/cb2d310c…518/evidence`, one flock, single 27B load per arm): OFF **48,285,916 kB** vs ON **24,750,704 kB = VmRSS** (**−23.54 GB**, load transient eliminated), ON-arm smoke 6/6; ledger row 2026-07-15. Exact-grid memory axes remain `FAILED` (projected PASS); direct-device streaming gate still open | [specs/safetensors-windowed-load.md](specs/safetensors-windowed-load.md) | `PARTIAL` | - | | `ENG-LOAD-DIRECT-UPLOAD` | Cut the SECOND copy out of a safetensors load: a weight the device consumes VERBATIM VIEWS the read-only shard mmap (`OwnedBytes::Borrow`, keep-alive carried on `StTensor::mapping`) instead of being copied into an owned host buffer, so `ResidentWeight` uploads straight from the file mapping and the load moves those bytes ONCE. Backend- and arch-agnostic: it lives in the shared dense loader helpers, so every arch that routes through them inherits it. Additive leaf below `LOAD-SAFETENSORS`; the windowed source-page release is preserved and now also runs after the device upload. Non-verbatim paths (transpose, dtype conversion, dequant, concatenation, GGUF load-time repack) are unchanged and keep their copy. Also lands the measurement half of issue #150: per-phase load timing plus host-copy / borrowed / device-upload byte counters (`VT_LOAD_STATS`) | T1 | one move per byte: streamed shard iterator `vllm/model_executor/model_loader/weight_utils.py:905-954`; immediate copy into the destination parameter `vllm/model_executor/models/utils.py:252-279`; driver `vllm/model_executor/model_loader/base_loader.py:43-82` @ `555967922` | refcounted mapping `src/vllm/model_executor/model_loader/safetensors_reader.cpp:46-80,231-262`; byte counters `:337-400`; `BorrowStTensorBytes` + post-upload adoption/release `src/vllm/model_executor/models/qwen3_5_weights.cpp`; qualifying call sites `include/vllm/model_executor/models/dense_weight_loaders.h` (`LoadBf16Direct`, `LoadCtNvfp4W4A16`, `LoadCtMxfp4W4A16`), `src/vllm/model_executor/models/qwen3_5_dense_weights.cpp` (`LoadModelBf16Direct`, `LoadCtNvfp4Raw`), `src/vllm/model_executor/models/qwen3_5_weights.cpp` (35B `LoadBf16Direct`); upload counter `include/vllm/model_executor/models/dense_attn_block.h`; fp4 resident upload counter + post-upload residency `include/vllm/model_executor/models/dense_nvfp4_gemm.h` (`ResidentNvfp4`) and `src/vllm/model_executor/models/qwen3_5.cpp` (private `ResidentNvfp4`); shared checker text normalization `scripts/checker_text.py` | [mechanism gate](../tests/vllm/test_load_direct_upload.cpp#L1) 14/14, 183 assertions (178 through round 4) — the 6 borrow-mechanism cases, 4 post-upload residency cases that pin the adopt branch, its RELEASE-BEFORE-REASSIGN ordering and the `VT_ADOPT_DEVICE_BYTES` decoupling, 3 fp4-resident cases that pin `ResidentNvfp4`'s upload accounting, its `d_dev` publication, its adoption and the previously unexercised NON-borrowed/host-addressable regime, and 1 case that pins the GENERAL adopt branch's `MADV_DONTNEED` by OBSERVING RESIDENCY (`mincore()` over the host mirror's interior pages, glibc pinned to the sbrk arena so `free()` cannot return them by itself) — the RSS half of this row, which every value assertion was blind to. RED under eleven mutations, each applied alone with the binary rebuilt (a failed build aborts rather than re-running a stale one), tree md5-verified restored and re-GREEN after each. Mutations are named by the TEST CASE they land on, never by a line ref, because carried-forward line refs are how the counts below were twice recorded wrong: drop the size-identity check 1 case / 5 assertions; skip the borrow in `LoadBf16Direct` 2 / 7; delete the adopt branch 5 / 19; move the release after the `bytes` reassignment 5 / 7 (incl. the ordering assertions read from inside the munmap); move the release after the host-addressable early return but BEFORE the env one 2 / 3; move it after the `VT_ADOPT_DEVICE_BYTES` early return, i.e. after BOTH, 3 / 4; delete the general branch's `MADV_DONTNEED` 1 / 1; `ResidentNvfp4` drop both `AddDeviceUpload` 3 / 3; drop both `d_dev` publications 3 / 20; drop both `AdoptDeviceBytesAsHost` 3 / 16; all six at once 3 / 23. FOUR of those rows had been recorded wrong (1/1, 2/2, 3/7 and 3/3) — every one a count measured against an older suite and copied forward rather than re-measured; all eleven are measured on this tree. The `qwen3_5.cpp` duplicate of `ResidentNvfp4` is in an anonymous namespace no test can call, so it is held to the same invariant by [`scripts/check-fp4-resident-consistency.py`](../scripts/check-fp4-resident-consistency.py) (42-case mutation suite), which checks PER BUFFER inside each buffer's own `if (!w.d_)` upload block — body-wide matching let one surviving `AddDeviceUpload` satisfy both buffers, so dropping exactly one counter passed (reproduced against the real `qwen3_5.cpp` text: old checker exit 0, new exit 1). Nine drop-one/substitute-one mutations of the LIVE duplicate now go RED. Round 5 closed two holes in it. It matched RAW source, so a statement left behind as a COMMENT, inside `#if 0`, or inside `if (false)` read as present while the compiler saw a deletion; the clause matchers now run on [`scripts/checker_text.py`](../scripts/checker_text.py)'s `normalize_source`, which blanks all three IN PLACE so byte offsets and line numbers survive — `strip_comments` had existed as two byte-identical private copies (check-runner-routing-consistency.py, check-surface-coverage.py) and both now import the shared helper instead of a third copy being written. And the ordering clause required PUBLISH-before-ADOPT but not COPY-before-ADOPT, so a body whose adoption runs BEFORE its copy passed although the upload then reads pages the adoption already released; new clause (f) READ-FIRST covers it. WHICH MOTION produces that order is part of the claim, because the two are different mutations: SINKING `d.b.Copy(...)` to below `AdoptDeviceBytesAsHost` leaves the `d_dev` publication in place so ONLY clause (f) bites (MEASURED on the live duplicate, old exit 0 / new exit 1 — the 0/1 row below), whereas HOISTING the adoption above the copy also lifts it above the publication and trips the PRE-EXISTING clause (e), which the old checker already caught (MEASURED 1/1) and which is therefore no evidence for (f). The equivalent mutations of the SHARED `dense_nvfp4_gemm.h` copy are red at run time, MEASURED one at a time with the binary DELETED before each rebuild and the header md5-verified (`4665255f7af6f52254367f9118ad92ee`) before and after each: sink `packed`'s `Copy` 2 cases / 2 assertions, sink `scale`'s 2 / 2, sink BOTH 2 / 4 — the failures are `w..bytes.data()[0] == kSrcPattern` and `AllBytesMatchPattern(...)` in `fp4 resident: ResidentNvfp4 COUNTS its upload…` and `fp4 resident: a host-addressable device adopts an OWNED fp4 mirror too`, TWO per mutated buffer, so an ODD count is not obtainable; hoist `packed`'s adoption 3 / 8 and hoist both 3 / 16, wider because a null `d_dev` makes the adoption return immediately and also takes the NON-host-addressable case. This clause had been recorded as **2 cases / 3 assertions**, which is the count of the *release after the host-addressable early return, BEFORE the env one* row carried onto a different mutation — the same failure mode as the four rows above, found by a fifth reviewer and re-measured here rather than copied. MEASURED on the LIVE `qwen3_5.cpp` on disk, one mutation at a time, md5-verified restored after each (35b5ea250f490105d579f4ffb573aa36 before and after all eleven) — old checker exit / new checker exit: delete the packed adoption 1/1, `//` it out 0/1, `/* */` it out 0/1, `#if 0` around it 0/1, `if (false)` around it 0/1, `//` the packed upload counter 0/1, `//` the packed `d_dev` publication 0/1, `#if 0` around that publication 0/1, `if (false)` around it 0/1, SINK the packed `Copy` to below its adoption 0/1. That gate is STRUCTURAL: it guards the duplicate against DELETION — including deletion disguised as a comment, an `#if 0` or a never-taken branch — against gross substitution and against mis-ordering, NOT against corruption, and it does not model the preprocessor (extending it to arbitrary `#ifdef` conditions is DECLINED: `#ifdef VT_CUTLASS_NVFP4` around the live adoption stays exit 0 in both checkers, because a build configuration is not a disguised deletion). Run-time proof exists only for the shared copy. The `mincore()` residency case now also carries a RUNTIME allocator guard — `#if __GLIBC__` proves the headers, not that glibc's allocator is running — which goes red (1 case / 1 assertion) when a purging allocator is simulated; the `mallopt` calls leave glibc's `no_dyn_threshold` set for the process either way, MEASURED as not specific to `M_TRIM_THRESHOLD` (glibc 2.39, `mallinfo2().hblks` probe: no mallopt LIVE, `M_TRIM_THRESHOLD` DISABLED, `M_MMAP_MAX` DISABLED), so it is recorded rather than removed. Round-4 re-run on the dev box, CLEAN Vulkan/llvmpipe Release build: `test_load_direct_upload` 14/14 (178), `test_safetensors` 34/34 (79), `test_qwen36_weights` 7/7 (45), `test_vulkan_backend` 35/35 (2107), `test_backend_cross_device` 11/11 (132), `test_opt_paged_engine` **6/6 token-exact (96/96), 0 declines, device type 3**. Round-5 re-run on the same box, branch rebased onto `a0fa12c7`, all identical except `test_load_direct_upload` **14/14 (183)** — the five added assertions are the RSS-reclaim case's runtime allocator guard, no case-count change — plus `test_checker_text` 33/33 and `scripts/gen-vulkan-spirv.py --check` `committed SPIR-V is up to date` run with the pinned `~/tools/glslang-16.5.0/bin/glslang` (without it the script exits 1 on this tree, so the green is a real compile-and-compare and not a skip). Round 6 changed NO compiled file — this row, the spec, the `(f)` clause docstring and one checker-test name — and re-ran the whole battery on the branch rebased onto `e3cc4f64`, CLEAN Release Vulkan build, 0 warnings: every count above reproduces byte-for-byte, `check-fp4-resident-consistency` rc 0, its suite 42, `test_checker_text` 33, `test_check_surface_coverage` 46, `test_check_runner_routing_consistency` 31, SPIR-V check up to date, preflight all green. GB10 Vulkan gates on the changed tree: `test_vulkan_backend` 35/35 (2650), `test_backend_cross_device` 11/11 (132), `test_opt_paged_engine` **6/6 prompts token-exact (96/96), 0 declines, device type 3**. GB10 CUDA (cutlass+triton) full `ctest` **383/393**, BOTH SACRED gates PASS (`test_qwen36_paged_engine`, `test_qwen27_paged_engine`); all 10 failures reproduce with the same signature on a clean `origin/main` CUDA build. MEASURED 27B bf16 load **1.54x warm / 1.61x cold**, bytes moved **100.196 -> 81.260 GiB** | [specs/load-direct-upload.md](specs/load-direct-upload.md) | `ACTIVE` | `CLAIM-ENG-LOAD-DIRECT-UPLOAD` | +| `ENG-HOST-EMBEDDING` | Keep the token table in HOST RAM and gather the requested rows on the CPU: `VT_HOST_EMBEDDING=1` runs the gather through the CPU `vt::Embedding` kernel (ONE ROW per id, so bf16/f16/f32 and every GGUF block format work) and copies the `[T,H]` result to the device, so the table never occupies device memory. A bf16 table takes a byte-copy fast path. `EmbedGather` is the one call every dense forward that owns a device table uses and it owns both arms, so the device upload happens only when the host arm declines (flag off, host bytes gone), with a one-shot `[host-embed] DISABLED` line naming the reason. An untied embedding frees the table's whole device residency; a tied head keeps it for the GEMM and only the gather moves. vllm.cpp-original; the async runner's device-resident id override is consumed before the gather so a stale host id vector cannot be gathered | T2 | absent in pin: vLLM gathers from a DEVICE-resident `VocabParallelEmbedding` parameter (`vllm/model_executor/layers/vocab_parallel_embedding.py`); the host-side per-row gather is llama.cpp's `ggml_get_rows` / `ggml_compute_forward_get_rows_q` (`ggml/src/ggml-cpu/ops.cpp:4850` @ b10451) | `include/vllm/model_executor/models/host_embedding.h`; `src/vllm/model_executor/models/host_embedding.cpp`; call sites across the Qwen3.5 family, MuseGlimmer, the shared Qwen3 dense driver and the classic dense families (Gemma 1-4, GLM4, Granite, MiniCPM 1/3, OLMo2, OPT, Phi, Phi3, StableLM, Command-R, DeepSeek-V2, GLM-MoE-DSA, Dots3-Note, Nemotron-H, Voxtral) | `tests/vllm/models/test_host_embedding.cpp`; `docs/ENVIRONMENT.md` | [specs/host-embedding.md](specs/host-embedding.md) | `ACTIVE` | `CLAIM-ENG-HOST-EMBEDDING` | | `LOAD-SAFETENSORS-DIRECT-DENSE` | Layer-bounded target-device loading for ordinary plain-BF16 Qwen3.5 dense safetensors on discrete CUDA; additive leaf below `LOAD-SAFETENSORS`, preserving windowed source release and owned-shard lifetime. Plain weights, stacked raw-NK owners, tied logits, logical resident state and same-queue layer staging are implemented and locally CUDA-gated. H32 Triton AOT, plain-BF16 decode graphs and ratio-4 FA2 repair the transplanted hot path; the broader row remains speed-gating | T1 | target-device model construction `vllm/model_executor/model_loader/base_loader.py:43-82`; incremental safetensors yield `vllm/model_executor/model_loader/weight_utils.py:820-954`; stacked/tied parameters `vllm/model_executor/models/qwen3_5.py:276-303,483-492` | load/dispatch `src/vllm/model_executor/models/qwen3_5_dense_weights.cpp:52-133,187-246,334-472`; discrete/unified classifier `src/vt/cuda/cuda_backend.cu:251-265`; plain execution/residency and graph selection `src/vllm/model_executor/models/qwen3_5.cpp`, `src/vllm/model_executor/models/qwen3_5_dense.cpp`; H32 recurrence `src/vt/cuda/cuda_gdn.cu`; ratio-4 FA2 `src/vt/cuda/cuda_flash_attn_fa2.cu`, `src/vt/cuda/cuda_paged_attn.cu`; exact benchmark corpus/output capture `examples/bench/bench_core.h:61-68,357-386,409-589`; reference metrics/token collector `tools/bench/vllm_closed_loop_metrics.py`; guarded driver and summarizer `tools/bench/run_qwen35_4b_compare.sh`, `tools/bench/summarize_qwen35_4b_compare.py` | [Real 4B graph/direct ON/OFF/eager gate](../tests/vllm/models/test_qwen35_plain_weights.cpp#L80) passes **3/3, 1672/1672**; [H32/GDN tests](../tests/vt/test_ops_gdn.cpp#L1) **10/10** flag and **66/66, 4242/4242** full; [paged-attention tests](../tests/vt/test_ops_paged_attn.cpp#L1354) **25/25, 454474/454474**. Final 18-leg root `/tmp/qwen35-main-final-fa2-20260725` plus stable vLLM confirmation: direct ON/OFF/vLLM-0.25 total **5769.99/5660.70/5849.80 tok/s**, output **638.03/625.94/646.85**, TPOT/ITL **43.72/43.84/38.55 ms**, peak PSS **2.406/8.592/7.662 GiB**, stable PSS **0.759/8.589/4.029 GiB**, VRAM **12850.7/12843.3/12942.7 MiB**. ON=OFF output IDs 128/128 every pair; ON is **+1.93%** total and cuts peak/stable PSS **72.0%/91.2%**. H32 AOT/graph/FA2 A/B gains are **+4.59%/+0.39%/+1.60%**. Graph-node trace has 453 launches and 200972 child kernels; local FA2 is 180.28 us/call vs vLLM 178.40, so the remaining **0.9864x** throughput and TPOT gap is host/engine-side. Sanitizer availability and external 27B/35B remain open | [plain-BF16 direct-load spike](specs/qwen35-plain-bf16-direct-load.md); [2026-07-25 evidence](../docs/bench-evidence/qwen35-4b-main-repair-20260725.md) | `GATING` | - | | `ENG-MOE-HOSTFREE` | MoE Marlin resident host-weight release: after `BuildMoeMarlinResident` uploads+repacks the routed experts to the device Marlin resident, free the per-expert fp4 HOST mirror (`OwnedTensor` packed+scale bytes) and `madvise(MADV_DONTNEED)` the pages back to the OS — returns the ~16.9 GiB steady 35B host double-store (`LoadNvfp4Raw` `MakeOwned` copies kept resident forever). Guarded to the committed Marlin path (`MarlinMoeEnabled()`; retained for the `VT_NVFP4_MARLIN=0` wmma fallback that re-reads them); `VT_MOE_HOST_FREE=0` A/B rollback. Realizes `release_host_weights_after_upload` for the dominant host consumer. **Item-2 (2026-07-19, `CLAIM-BACKEND-PLATFORM-2`): the host-free decision is now CONSUMED from `GetPlatform(d.q.device.type).residency_policy()` via `vllm::platforms::ShouldReleaseHostWeights(policy, MarlinMoeEnabled()/*kernel-path*/, VT_MOE_HOST_FREE/*env*/)` — `CudaPlatform` flag flipped false→true (reproduces today EXACTLY); `MarlinMoeEnabled()` stays the orthogonal KERNEL-PATH safety gate.** | T0 | streamed target-device weight construction `vllm/model_executor/model_loader/base_loader.py:43-82`; residency capability `vllm/platforms/interface.py:134-229`; Marlin MoE in-place repack (no host mirror) `vllm/model_executor/layers/quantization/utils/marlin_utils_fp4.py:375-434` | free region `src/vllm/model_executor/models/qwen3_5.cpp:3743-3781`; `OwnedTensor::ReleaseHost` decl `include/vllm/model_executor/models/qwen3_5_weights.h:65` + impl (madvise+swap) `src/vllm/model_executor/models/qwen3_5_weights.cpp:24`; public hook `Qwen3_5Model::PrepareMarlinResident` `src/vllm/model_executor/models/qwen3_5.cpp:4472` | `tests/vllm/test_qwen36_weights.cpp:273` (`ReleaseHost` frees buffer+capacity) + `:314` (`PrepareMarlinResident` release under Marlin / retention under `VT_NVFP4_MARLIN=0`, DGX release 27/27 + retention 15/15); DGX A/B (`VT_MOE_HOST_FREE`) 35B STEADY serving PSS **20.17→3.53 GiB** (root `dgx:~/work/mem35-hostfree`); token-neutral 315/315 + 235/235; c2 smoke clean; memcheck clean. Whole-window load-phase PEAK bounded by the `ENG-MOE-LOADSTREAM` follow-up. Ledger [parity-ledger.md#L521](parity-ledger.md#L521) | [moe-marlin-host-free.md](specs/moe-marlin-host-free.md) | `DONE` | `ac77bec` | | `ENG-MOE-LOADSTREAM` | 35B load-phase PEAK-PSS interleave (follow-up to `ENG-MOE-HOSTFREE`): the steady free returns the routed-expert host mirror only AFTER the whole model loads, so whole-window `peak_pss`/`peak_rss` (~19.8 GiB) is still set by all N layers' ~256 experts host-coexisting at load. DEFER the routed-expert host copies (`LoadQwen3_5Moe` loads each layer WITHOUT experts + installs a per-layer `load_layer_experts` streaming closure that owns the mmap'd shards) and materialize ONE layer's experts inside `PrepareMarlinResident` immediately before that layer's device Marlin build + host free — so at most one layer's experts coexist on the host (peak ~ one layer, not all N). Device residents byte-identical (same source bytes, same per-layer build order). Non-CUDA/`VT_NVFP4_MARLIN=0`/no-Marlin build falls back to bulk host materialization (their forward reads the host bytes). 27B is a different loader (`LoadQwen3_5Dense`, true-W4A4) → unaffected; GGUF/synthetic/borrowed pass no shards owner → eager. **Item-2 (2026-07-19, `CLAIM-BACKEND-PLATFORM-2`): the per-layer interleave gate is now CONSUMED from `GetPlatform(queue.device.type).residency_policy()` via `vllm::platforms::ShouldInterleaveLoadStream(policy, MarlinMoeEnabled())` — reproduces the old `queue.device.type != kCUDA` OR `!MarlinMoeEnabled()` gate EXACTLY (unified/CPU retain-host ⇒ policy false ⇒ materialize-all fallback); the ~4 GiB load-peak win is preserved.** | T0 | streamed target-device construction `vllm/model_executor/model_loader/base_loader.py:43-82`; in-place Marlin repack, no host mirror `vllm/model_executor/layers/quantization/utils/marlin_utils_fp4.py:375-434`; residency capability `vllm/platforms/interface.py:134-229` | deferred field `include/vllm/model_executor/models/qwen3_5_weights.h:307` (`Qwen3_5MoeWeights::load_layer_experts`); loader defer + closure `src/vllm/model_executor/models/qwen3_5_weights.cpp:331,399` (`LoadMoeExpertsInto`/`LoadQwen3_5Moe`); shared shards owner `include/vllm/model_executor/models/model_registry.h:60` + `src/vllm/model_executor/models/model_registry.cpp:213` (`ModelSource::FromSafetensorsOwned`) + `src/vllm/entrypoints/model_loader.cpp:365` (`LoadFromDir`); per-layer interleave + `MaterializeAllDeferredExperts` `src/vllm/model_executor/models/qwen3_5.cpp:4498,4508` (`PrepareMarlinResident`) | CPU coexistence-bound contract `tests/vllm/test_qwen36_weights.cpp:324` ("deferred routed-expert load: move-safe closure + bounded coexistence", peak==1 across the per-layer materialize→free loop, move-safe closure); clean `-Werror` CPU build 0 warn, full ctest (3 HTTP/engine-proc parallel-port flakes pass isolated), tools 164/164. **DGX PROVEN** (`~/work/vllm.cpp-mem35-loadstream` new vs `-parent` 7a1a6d6 eager, production flags CUTLASS sm120a+Marlin+FA2 sm_121a, one flock): 35B load-to-ready **peak RSS (VmHWM) 21.43 GiB → 4.19 GiB (−17.24 GiB / −80%, below vLLM 13.3 GiB)**; token BYTE-IDENTICAL both binaries — 35B `test_qwen36_paged_engine` 315/315 + 27B `test_qwen27_paged_engine` 235/235; **27B UNAFFECTED** (peak RSS 24.8 GiB, matches its baseline — dense loader, no deferral); `compute-sanitizer memcheck` on the deferred load path **0 errors / 315 assertions** (no use-after-free of the freed host bytes); weights unit 127 assertions (CPU coexistence peak==1 + DGX residency). `benchmark_binding=false` — the orchestrator re-grids the binding `peak_pss`/`peak_rss` axes to confirm the FAIL→PASS flip | [moe-expert-load-stream.md](specs/moe-expert-load-stream.md) | `ANCHOR-BACKFILL` | CLAIM-MEM35-LOADSTREAM | diff --git a/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md b/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md new file mode 100644 index 000000000..3d8ea0ac1 --- /dev/null +++ b/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md @@ -0,0 +1,19 @@ +ID: ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM +Title: The token table has no host-resident arm: every dense forward uploads [vocab, H] to the device even though the gather could run in host RAM +Row: ENG-HOST-EMBEDDING +State: OPEN +Kind: feature +GitHub: - +Mirror: PENDING +Availability: FULL +Created: 2026-09-29 +Updated: 2026-09-29 +Closed: - + +## Problem + +A dense forward that owns a device embedding table pays [vocab, H] of device memory for it and gathers there. On the cards this project targets that is real: a 151k x 5120 bf16 table is 1.5 GiB of a 24 GiB pool, and on a tied-head model the table is kept for the lm_head GEMM regardless. The gather itself is tiny on the host: one row per id through the CPU vt::Embedding kernel, then one [T, H] copy to the device, which is llama.cpp`s ggml_get_rows shape (ggml/src/ggml-cpu/ops.cpp:4850 @ b10451) rather than vLLM`s device-side gather. Nothing in the tree offered this: `ResidentWeight`/`EmbedGather` always uploaded, so a text-only or memory-bound serve had no way to keep the table host-side. The change adds `VT_HOST_EMBEDDING=1` (host gather, one row per id, bf16/f16/f32 and every GGUF block format, bf16 byte-copy fast path) behind the existing `EmbedGather` seam, so every dense forward that owns a device table gets it for free and the device behaviour is unchanged when the flag is off or the host bytes are gone. It is vllm.cpp-original, so the secondary oracle is llama.cpp`s ggml_get_rows, and the correctness bar is the same pinned IQ4_NL/Q5_0 golden vectors the device-side vt::Embedding gate already uses. + +## Resolution + +- diff --git a/.agents/specs/host-embedding.md b/.agents/specs/host-embedding.md new file mode 100644 index 000000000..114cfad4f --- /dev/null +++ b/.agents/specs/host-embedding.md @@ -0,0 +1,141 @@ +# Host-resident token table: gather the embedding rows on the CPU — ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM + +A dense forward that owns a device embedding table uploads `[vocab, H]` once and +gathers there. On a 24 GiB card that is 1.5 GiB of pool for a 151k x 5120 bf16 +table, and on a tied head the table is kept for its GEMM whatever the gather does. +The gather itself is small: one row per id. + +Issue: [ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM](../issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md). +Owning row: `ENG-HOST-EMBEDDING` ([engine-matrix.md](../engine-matrix.md)), added +with this spec. + +## Scope + +IN: the host-resident arm of the token-table gather — the `VT_HOST_EMBEDDING` +switch, the CPU gather through `vt::Embedding`, the `[T,H]` copy to the device, +the `EmbedGather` seam that owns both arms, the wake-up of every dense forward +that owns a device table, `docs/ENVIRONMENT.md`, and the tests that pin the arm. + +OUT: the tied-head GEMM and every other consumer of the table (the table's host +bytes stay available for them); the decode-graph arms whose ids live on a device +tensor (they keep the device gather by design); any change to the device arm's +bytes or dtype; and the end-to-end VRAM/throughput measurement, which is a row +deliverable recorded under `## Owed`, not a precondition for the arm. + +## Upstream chain + +**vLLM has no equivalent.** Its embedding is a device-resident +`VocabParallelEmbedding` parameter and the gather runs on the device. This is a +vllm.cpp-original capability, so the primary oracle has nothing to mirror and the +shape comes from a **secondary oracle**: llama.cpp's `ggml_get_rows` / +`ggml_compute_forward_get_rows_q` (`ggml/src/ggml-cpu/ops.cpp:4850` @ `b10451`), +the dequantizing host gather that reads ONE ROW per id and never materializes the +table. That is the same function `tests/vt/test_ops_embedding_quant.cpp` cites for +the device-side `vt::Embedding` op, so both arms are compared against one +denominator. + +## Our baseline + +Before this change every dense forward that owns a device table staged it with +`ResidentWeight(d, table, {vocab, H})` and gathered with `vt::Embedding(d.q, ...)`. +For an fp8/non-CUDA-device path `ResidentWeight` aliases the host bytes, but on a +staging backend (CUDA, non-unified) it uploads `[vocab, H]` and keeps it resident +for the process. There was no way to keep the table host-side, and no seam that +owned both arms. + +## Port map + +| Piece | Where | +|---|---| +| The seam and the flag | `include/vllm/model_executor/models/host_embedding.h`, `src/vllm/model_executor/models/host_embedding.cpp` | +| The wake-up call sites | the Qwen3.5 family, MuseGlimmer, the shared Qwen3 dense driver, and the classic dense families (Gemma 1-4, GLM4, Granite, MiniCPM 1/3, OLMo2, OPT, Phi, Phi3, StableLM, Command-R, DeepSeek-V2, GLM-MoE-DSA, Dots3-Note, Nemotron-H, Voxtral) | +| The async id override consumed before the gather | `src/vllm/model_executor/models/qwen3_5_internal.h` (`detail::TakeDeviceTokenIds`) | +| The documented knobs | `docs/ENVIRONMENT.md` (`VT_HOST_EMBEDDING`, `VT_HOST_EMBED_TRACE`) | +| Build registration | the new TU in `CMakeLists.txt` | + +## Tests to port + +The oracle fixtures already exist and are reused rather than re-derived: +`tests/vt/iq4nl_q5_0_golden_vectors.h` carries the pinned IQ4_NL/Q5_0 golden +vectors (real bytes of the shipped `per_layer_token_embd.weight` and the pinned +oracle's own `dequantize_row_iq4_nl` output) that pin the device-side op +bit-exactly. The host arm is held to the same vectors, so the two arms are +compared against one denominator, and the CPU per-row decode is the llama.cpp +`ggml_get_rows` behaviour this arm ports. + +## Design + +`EmbedGather(d, out, token_ids, table, vocab, H, what)` is the ONE call a dense +forward makes when it owns a device table. It owns both arms: + +1. **Host arm** (`VT_HOST_EMBEDDING=1`): the table stays in `table.bytes` and is + never uploaded. A process-lifetime CPU queue runs `vt::Embedding` with a + `ViewOn` of the host bytes, which decodes one row per gathered id for every + table residency the loaders produce (bf16/f16/f32 and the GGUF block formats). + A bf16 table into a bf16 output takes a byte-copy fast path. The `[T,H]` + staging buffer is copied to `out`. + - The async runner may have spliced this step's sampled token into a + device-resident ids buffer, so the host arm consumes + `detail::TakeDeviceTokenIds()` and reads the ids back before the gather. + - `VT_HOST_EMBED_TRACE=1` prints the arm and `T`. + - The flag is read ONCE per process (a function-static), so a serving process + cannot switch arms mid-run; a same-binary A/B is two processes. +2. **Device arm** (the shipped behaviour): `ResidentWeight` upload + + `ApplyDeviceTokenIds` + `vt::Embedding`, unchanged. Taken when the flag is off, + when the host bytes are gone (`bytes.empty()` or `host_released`), or when the + host arm declines for any reason; a one-shot `[host-embed] DISABLED` line names + the reason. + +## Dependencies + +- `vt::Embedding`'s CPU kernel and the CPU backend (in-tree, no new library). +- No GPU is needed for the host arm's correctness; a staging backend is needed to + exercise the `d_dev == nullptr` half. +- The async id override (`detail::TakeDeviceTokenIds`) must exist for the serving + loop; it is already in-tree from ENG-ASYNC-SCHED W4. + +## Work breakdown + +- **W1** (this PR): the seam, the flag, the CPU gather, the CPU-queue singleton, + the bf16 fast path, the async-override consumption, the wake-up of the dense + forward family, the docs, and the test. +- **W2** (owed): the end-to-end device-memory number and a decode-throughput A/B + on the target 24 GiB `sm_120a` card, both arms in one binary. Not a precondition + for the arm; the reason it is not in W1 is that a helper PR makes no speed claim. +- **W3** (later, if wanted): a config surface (`--offload-config`'s `vllm_cpp` + key) instead of an environment variable, if the operator prefers the config + document over env. + +## Risks and decisions + +| Risk / decision | Handling | +|---|---| +| A host gather adds a synchronize on the serving loop | The async runner's device-resident ids are consumed BEFORE the gather, so the host arm reads the spliced ids back rather than racing them; the decode-graph arms that hold ids on device keep the device gather and are listed in the header | +| A tied head keeps the table resident, so the memory saving is smaller | Recorded in the knob's documentation and in the header: the saving is the table's device residency for an untied embedding, and only the gather moves for a tied one | +| The flag is process-static, so a test binary cannot exercise both arms | The test binary enables it before `main` and is flag-ON by construction; the off arm is a plain early return and every other model suite runs it | +| `d_dev == nullptr` does not prove the arm ran on the CPU backend | `ResidentWeight` aliases host bytes when `is_cpu()`, so the test captures the one-shot `[host-embed]` banner and checks `HostEmbedInto`'s return value; the mutation that disables the arm makes 4 assertions fail | +| The host table's bytes are released by another path | The host arm declines on `bytes.empty()` / `host_released`, and `EmbedGather` then takes the device arm; the decline is a test case | + +## Evidence + +- The host arm's rows equal the pinned oracle through the same golden vectors the + device op is gated on; the arm is proven to have run (captured banner + return + value), not merely to have produced right numbers. +- The table is never uploaded on the host arm (`d_dev == nullptr`). +- `tests/vt/test_ops_embedding_quant` (6/6, 1637 assertions) is unchanged. + +## Gates + +- `ctest --test-dir build -R test_host_embedding`. +- `tests/vt/test_ops_embedding_quant` stays green (the device op is unchanged). +- `python3 scripts/check-env-doc.py` — the two new env vars are documented in + `docs/ENVIRONMENT.md` in this change (the checker's only remaining complaint is + the pre-existing `VT_VK_*` gap). +- The full `ctest --test-dir build`. +- `scripts/agent-preflight.sh --staged`. + +## Owed + +- The W2 measurement above (a 27B NVFP4/Q8mix GGUF on the 24 GiB `sm_120a` card, + the flag on and off, device memory and decode throughput), owed to the row. +- A dedicated config key (W3) only if the operator asks for one. diff --git a/scripts/check-agent-record.py b/scripts/check-agent-record.py index 78332cabd..99b4f9306 100644 --- a/scripts/check-agent-record.py +++ b/scripts/check-agent-record.py @@ -389,7 +389,9 @@ # family, SERVE-RECIPE-ARGS / -REQUEST-LENGTH-GUARD, LOAD-GGUF-MMPROJ and the # attention-window row. Bumped because a new row EXISTS, never to make a # transition pass. -ENGINE_ROWS = 179 +# 180 since 2026-09-29: +1 row (ENG-HOST-EMBEDDING, the host-resident token +# table). Bumped because a new row EXISTS, never to make a transition pass. +ENGINE_ROWS = 180 ENGINE_SUMMARY_SECTIONS = ( ("Engine and scheduling", "Engine core and scheduling"), diff --git a/tests/scripts/test_agent_record.py b/tests/scripts/test_agent_record.py index 93ad0a401..f99c2dac4 100644 --- a/tests/scripts/test_agent_record.py +++ b/tests/scripts/test_agent_record.py @@ -1135,6 +1135,48 @@ def test_the_row_names_its_issue_and_its_spec(self) -> None: index = tracked_issues(self) self.assertIn("issues/81)", index) +class HostEmbeddingRowIsCounted(unittest.TestCase): + """The host-resident token table row is counted and carries its records. + + Adding the row bumped the engine-matrix counts AND the pinned `ENGINE_ROWS`, + which the checker enforces as an anti-drift pin. This class is the executable + half of that bump: the row must exist exactly once, name its spec and its + claim, and its spec must carry the structured sections the checker requires + of an ACTIVE row's spec. Deleting the row, its spec link, its claim owner or + a required section each turns one of these red. + """ + + ROW = "ENG-HOST-EMBEDDING" + + def test_the_row_exists_once_in_the_engine_matrix(self) -> None: + text = (ROOT / ".agents/engine-matrix.md").read_text(encoding="utf-8") + matching = [ + line for line in text.splitlines() if line.startswith(f"| `{self.ROW}` |") + ] + self.assertEqual(len(matching), 1, f"{self.ROW} must appear exactly once") + + def test_the_row_names_its_spec_and_its_claim(self) -> None: + text = (ROOT / ".agents/engine-matrix.md").read_text(encoding="utf-8") + row = next(l for l in text.splitlines() if l.startswith(f"| `{self.ROW}` |")) + self.assertIn("specs/host-embedding.md", row) + self.assertIn("`CLAIM-ENG-HOST-EMBEDDING`", row) + + def test_the_spec_carries_the_structured_sections(self) -> None: + spec = (ROOT / ".agents/specs/host-embedding.md").read_text(encoding="utf-8") + for heading in ( + "## Scope", + "## Upstream chain", + "## Our baseline", + "## Port map", + "## Tests to port", + "## Dependencies", + "## Work breakdown", + "## Risks and decisions", + "## Gates", + ): + self.assertIn(heading, spec, heading) + + class CanonicalIssueRecordTests(unittest.TestCase): def record(self, number: int) -> object: return agent_record.issue_records.IssueRecord( From 3cd5700ca8369484a759241dd868962e35fb96bb Mon Sep 17 00:00:00 2001 From: dev Date: Sun, 27 Sep 2026 19:47:05 +0000 Subject: [PATCH 2/3] feat(ENG-HOST-EMBEDDING): gather the token table on the host VT_HOST_EMBEDDING=1 keeps the embedding table in host RAM and runs the gather on a cached CPU queue, then copies the [T,H] rows to the device. The CPU vt::Embedding kernel decodes one row per gathered id, so every table residency the loaders produce works: bf16/f16/f32 and the GGUF block formats. A bf16 table takes a byte-copy fast path. The async runner's device-resident id override is consumed before the gather. `EmbedGather` is the one call every dense forward that owns a device table uses; it owns both arms, so the upload happens only when the host arm declines. Wired across the Qwen3.5 family, the shared Qwen3 dense driver, MuseGlimmer, and the classic dense families (Gemma 1-4, GLM4, Granite, MiniCPM, OLMo2, OPT, Phi, Phi3, StableLM, Command-R, DeepSeek-V2, GLM-MoE-DSA, Dots3-Note, Nemotron-H, Voxtral). Decode-graph arms whose ids already live on device keep the device gather. The device gather is unchanged and runs when the flag is off or the host bytes are gone. VT_HOST_EMBED_TRACE logs each call. The new test is the arm's first coverage and has two halves: the rows are compared against the same pinned IQ4_NL/Q5_0 oracle vectors the device-side op test uses, AND the table must stay off the device (`d_dev == nullptr`). A gather that uploaded first passes the first half and fails the second. FOLLOWING_AGENTS_PROTOCOL Following-Agents-Protocol: true AI-Assisted: true Assisted-by: AGENT:opencode-go/deepseek-v4.1-flash [pi] --- CMakeLists.txt | 1 + docs/ENVIRONMENT.md | 2 + .../model_executor/models/host_embedding.h | 43 +++ src/vllm/model_executor/models/commandr.cpp | 6 +- .../model_executor/models/deepseek_v2.cpp | 12 +- .../models/dots3_note_device.cpp | 8 +- src/vllm/model_executor/models/gemma.cpp | 6 +- src/vllm/model_executor/models/gemma2.cpp | 6 +- src/vllm/model_executor/models/gemma3.cpp | 6 +- src/vllm/model_executor/models/gemma4.cpp | 6 +- src/vllm/model_executor/models/gemma4_mm.cpp | 8 +- src/vllm/model_executor/models/glm4.cpp | 6 +- .../models/glm_moe_dsa_forward.cpp | 20 +- src/vllm/model_executor/models/granite.cpp | 6 +- .../model_executor/models/host_embedding.cpp | 121 ++++++++ src/vllm/model_executor/models/minicpm.cpp | 6 +- src/vllm/model_executor/models/minicpm3.cpp | 6 +- .../model_executor/models/muse_glimmer.cpp | 6 +- .../model_executor/models/muse_glimmer_mm.cpp | 8 +- .../models/nemotron_h_device.cpp | 7 +- src/vllm/model_executor/models/olmo2.cpp | 6 +- src/vllm/model_executor/models/opt.cpp | 6 +- src/vllm/model_executor/models/phi.cpp | 6 +- src/vllm/model_executor/models/phi3.cpp | 6 +- src/vllm/model_executor/models/qwen3.cpp | 54 +--- src/vllm/model_executor/models/qwen3_5.cpp | 87 +++--- src/vllm/model_executor/models/stablelm.cpp | 6 +- src/vllm/model_executor/models/voxtral.cpp | 9 +- tests/CMakeLists.txt | 6 + tests/vllm/models/test_host_embedding.cpp | 271 ++++++++++++++++++ 30 files changed, 577 insertions(+), 170 deletions(-) create mode 100644 include/vllm/model_executor/models/host_embedding.h create mode 100644 src/vllm/model_executor/models/host_embedding.cpp create mode 100644 tests/vllm/models/test_host_embedding.cpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 88716c8ad..d37977806 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -780,6 +780,7 @@ add_library(vllm STATIC src/vllm/model_executor/models/qwen3_5_gguf_weights.cpp src/vllm/model_executor/models/qwen3_gguf_weights.cpp src/vllm/model_executor/models/qwen3_5.cpp + src/vllm/model_executor/models/host_embedding.cpp src/vllm/model_executor/models/qwen3_5_common.cpp src/vllm/model_executor/models/qwen3_5_dense.cpp src/vllm/model_executor/models/qwen3_5_moe.cpp diff --git a/docs/ENVIRONMENT.md b/docs/ENVIRONMENT.md index d514d76cd..68dbfa92d 100644 --- a/docs/ENVIRONMENT.md +++ b/docs/ENVIRONMENT.md @@ -233,6 +233,8 @@ portable/reference path. In normal operation leave them unset. | `VT_ATTN_PREAMBLE_COOP` | off | `=1` selects the warp-per-item cooperative attention preamble arm (`AttnQkNormRopeGateCoopK`) on ROCm, mapping one warp per token item instead of the donor walk (`AttnQkNormRopeGateK`); read once per process like the sibling arms | | `VT_GDN_COLPERM_KEEP_QUANT` | off | `=1` keeps the column-permuted `ssm_out`/`out_proj` tensor as Q5_K in tiled order (no `ReorderVCols`) and permutes the 4096-element GEMV input at runtime instead; the column reorder cuts across Q5_K block boundaries, so the weight cannot be permuted in place. Saves ~4x weight bandwidth (Q5_K ~5 MB vs bf16 20 MB per call) | | `VT_GDN_ROWPERM_KEEP_QUANT` | off | `=1` keeps the row-permuted V-head GDN projections (in the tiled order the row permutation produces) as K-quant instead of expanding to bf16 at load; the runtime gather supplies the permutation. Opt-in; the default reorders then expands | +| `VT_HOST_EMBEDDING` | off | `=1` keeps the token table in host RAM and gathers the requested rows on a cached CPU queue, then copies the `[T,H]` result to the device. The CPU `vt::Embedding` kernel decodes ONE ROW per gathered id, so every table residency the loaders produce works (bf16/f16/f32 and the GGUF block formats); a bf16 table takes a byte-copy fast path. The async runner's device-resident id override is consumed before the gather. Falls back to the device gather when the flag is off or the host bytes are gone, with a one-shot `[host-embed] DISABLED` line naming the reason. Used by every dense forward that owns a device table (the Qwen3.5 family, MuseGlimmer, the shared Qwen3 dense driver, and the classic dense families: Gemma 1-4, GLM4, Granite, MiniCPM 1/3, OLMo2, OPT, Phi, Phi3, StableLM, Command-R, DeepSeek-V2, GLM-MoE-DSA, Dots3-Note, Nemotron-H, Voxtral); an untied embedding frees the table's device residency, a tied head keeps it for the GEMM. Decode-graph arms that hold their ids on device keep the device gather | +| `VT_HOST_EMBED_TRACE` | off | `=1` prints one `[host-embed] T=… override=…` line per host-side embedding gather, for telling which path a step took. Diagnostic only | | `VT_GDN_SCAN_COOP` | off | `=1` selects the warp-per-row cooperative GDN scan arm (`GdnScanCoopK`) on ROCm, mapping one warp per output row instead of the donor walk (`GdnScanK`); read once per process like the sibling arms | | `VT_GDN_PACKED_DECODE` | on (CUDA GDN) | Unpacked GDN decode path | | `VT_GDN_DECODE_BV` | `32` (CUDA GDN decode experiment) | Exact `16` selects the byte-identical 16-value fused-recurrence tile; unset and every other spelling keep the 32-value schedule. Experimental opt-in; no release or cross-hardware default change | diff --git a/include/vllm/model_executor/models/host_embedding.h b/include/vllm/model_executor/models/host_embedding.h new file mode 100644 index 000000000..7be7dadb8 --- /dev/null +++ b/include/vllm/model_executor/models/host_embedding.h @@ -0,0 +1,43 @@ +// VT_HOST_EMBEDDING: gather the embedding rows on the CPU and copy the [T,H] +// result to the device, so the token table never occupies device memory. The +// CPU `vt::Embedding` kernel decodes ONE ROW per gathered id — the same per-row +// discipline as llama.cpp's ggml_get_rows — for every table residency the +// loaders produce: bf16/f16/f32 and the GGUF block-quant formats. +// +// Returns false (the caller uses the device path) when the flag is off or the +// table's host bytes are gone. A forward whose embedding is ALREADY a host +// gather does not need this; it exists for the forwards that own a device +// table. An untied embedding frees the table's whole device residency; a tied +// head keeps the table resident for its GEMM, so the flag then only moves the +// gather off the device. +// +// WHO CALLS IT. Every dense forward that owns a device table and takes its ids +// as a host vector calls `EmbedGather`. The sites that deliberately do NOT are +// the ones where a host gather cannot help or would add a synchronize the path +// exists to remove: decode-graph arms whose ids already live in a device tensor +// (qwen3_moe, gemma3, deepseek_v2, nemotron_h paged, qwen4_exp), the mm embed +// hooks whose ids arrive on device (qwen3_vl, dots3_note), the draft heads +// (qwen3_dflash/dspark), and the custom-residency loaders (kimi_linear). +#pragma once + +#include +#include + +#include "vllm/model_executor/models/dense_device_glue.h" // Dev, DBuf, OwnedTensor + +namespace vllm { +namespace dense_attn { + +bool HostEmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids, + const OwnedTensor& table, int64_t vocab, int64_t H); + +// The gather a forward should call when it owns a device table: the host arm +// above, else `ResidentWeight` + the async id override + `vt::Embedding`. The +// table and its shape are handed over UNRESOLVED so the upload happens only when +// the host arm declines. `what` names the caller in the override's shape check. +void EmbedGather(Dev d, DBuf& out, const std::vector& token_ids, + const OwnedTensor& table, int64_t vocab, int64_t H, + const char* what); + +} // namespace dense_attn +} // namespace vllm diff --git a/src/vllm/model_executor/models/commandr.cpp b/src/vllm/model_executor/models/commandr.cpp index 72827bd68..91ab193b1 100644 --- a/src/vllm/model_executor/models/commandr.cpp +++ b/src/vllm/model_executor/models/commandr.cpp @@ -29,6 +29,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -194,9 +195,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "commandr embed"); } StepInputs si = BuildStepInputs(d, positions, attn_meta, config); diff --git a/src/vllm/model_executor/models/deepseek_v2.cpp b/src/vllm/model_executor/models/deepseek_v2.cpp index 740898c0c..4e6789dc7 100644 --- a/src/vllm/model_executor/models/deepseek_v2.cpp +++ b/src/vllm/model_executor/models/deepseek_v2.cpp @@ -77,6 +77,7 @@ #include "vllm/model_executor/moe_placement_seam.h" #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/ResidentWeight glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" #include "vllm/model_executor/models/mla_attention.h" #include "vllm/model_executor/models/qwen3_5_internal.h" // detail::DeviceTokenIds @@ -577,14 +578,9 @@ void EmbedInto(Dev d, DBuf& hidden, const Tensor& ids, void EmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids, const DeepseekV2Weights& weights) { - const int64_t T = static_cast(token_ids.size()); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - // #1305: the eager arms of these two registrations had the SAME defect as the - // graph arm — they embedded the host vector and never looked at the device - // mirror. The override's copy is enqueued on the main queue, so it is ordered - // AFTER the combine that produced it rather than racing it. - detail::ApplyDeviceTokenIds(d.b, d.q, dids.ptr(), T, "deepseek v2 embed"); - EmbedInto(d, hidden, dids.t(), weights); + const DeepseekV2Params& p = weights.params; + EmbedGather(d, hidden, token_ids, weights.embed_tokens, p.vocab_size, + p.hidden_size, "deepseek v2 embed"); } // The CAPTURABLE region: everything after the embedding — the MLA step metadata diff --git a/src/vllm/model_executor/models/dots3_note_device.cpp b/src/vllm/model_executor/models/dots3_note_device.cpp index 21ff56e7a..a2de243c1 100644 --- a/src/vllm/model_executor/models/dots3_note_device.cpp +++ b/src/vllm/model_executor/models/dots3_note_device.cpp @@ -166,6 +166,7 @@ #include "vllm/model_executor/model_loader/safetensors_reader.h" #include "vllm/model_executor/models/deepseek_v2.h" // MlaStep / BuildMlaStep #include "vllm/model_executor/models/dense_attn_block.h" +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/dense_weight_loaders.h" #include "vllm/model_executor/models/mla_attention.h" #include "vt/backend.h" @@ -1222,11 +1223,8 @@ ForwardLogits Dots3NoteModel::ForwardDevice( d.b.Copy(d.q, hidden_buf.ptr(), emb.data, static_cast(T * H) * vt::SizeOf(DType::kBF16)); } else { - DBuf ids(d, DType::kI32, {T}, token_ids.data()); - Tensor tab = ResidentWeight(d, dw.embed_tokens, {vocab, H}); - Tensor h = hidden_buf.t(); - Tensor idt = ids.t(); - vt::Embedding(d.q, h, tab, idt); + EmbedGather(d, hidden_buf, token_ids, dw.embed_tokens, vocab, H, + "dots3-note embed"); } // ── the residual stream (deepseek_v2.py:1262-1345, unchanged by dots3) ──── diff --git a/src/vllm/model_executor/models/gemma.cpp b/src/vllm/model_executor/models/gemma.cpp index 8f1b3c256..ecbcf71f9 100644 --- a/src/vllm/model_executor/models/gemma.cpp +++ b/src/vllm/model_executor/models/gemma.cpp @@ -21,6 +21,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpGeluMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -201,9 +202,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed then scale by sqrt(hidden) cast to bf16 (gemma.py:288-295). DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "gemma embed"); } const float nsqrt = std::sqrt(static_cast(H)); const double normalizer = static_cast(vt::BF16ToF32(vt::F32ToBF16(nsqrt))); diff --git a/src/vllm/model_executor/models/gemma2.cpp b/src/vllm/model_executor/models/gemma2.cpp index f8bc732d1..7c42ba102 100644 --- a/src/vllm/model_executor/models/gemma2.cpp +++ b/src/vllm/model_executor/models/gemma2.cpp @@ -34,6 +34,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpGeluMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -308,9 +309,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed then scale by sqrt(hidden) cast to bf16 (gemma2.py:276-283). DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "gemma2 embed"); } const float nsqrt = std::sqrt(static_cast(H)); const double normalizer = static_cast(vt::BF16ToF32(vt::F32ToBF16(nsqrt))); diff --git a/src/vllm/model_executor/models/gemma3.cpp b/src/vllm/model_executor/models/gemma3.cpp index f29e91d4d..843b097c8 100644 --- a/src/vllm/model_executor/models/gemma3.cpp +++ b/src/vllm/model_executor/models/gemma3.cpp @@ -40,6 +40,7 @@ #include "vllm/model_executor/layers/attention/attention.h" #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpGeluMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // Pool #include "vllm/model_executor/models/gemma3_decode_graph.h" #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits @@ -469,9 +470,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // matching torch's bf16-scalar multiply. DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "gemma3 embed"); } StepInputs si = BuildStepInputs(d, positions, attn_meta, config); return ForwardLayers(d, std::move(hidden), si, attn_meta, attn_kv, weights, config, diff --git a/src/vllm/model_executor/models/gemma4.cpp b/src/vllm/model_executor/models/gemma4.cpp index 8b14d0d3f..265fba583 100644 --- a/src/vllm/model_executor/models/gemma4.cpp +++ b/src/vllm/model_executor/models/gemma4.cpp @@ -47,6 +47,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpGeluMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // Pool #include "vllm/model_executor/models/gemma4_moe.h" #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits @@ -464,9 +465,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // one did H2D-into-`tok` then D2D-into-`hidden`). Recorded under `## Owed`. d.b.Copy(d.q, tok.ptr(), inputs_embeds_override->data, tok.bytes()); } else { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, tok.t(), dtab, dids.t()); + EmbedGather(d, tok, token_ids, weights.embed_tokens, vocab, H, + "gemma4 embed"); const float nsqrt = std::sqrt(static_cast(H)); const double normalizer = static_cast(vt::BF16ToF32(vt::F32ToBF16(nsqrt))); diff --git a/src/vllm/model_executor/models/gemma4_mm.cpp b/src/vllm/model_executor/models/gemma4_mm.cpp index a39ae13b3..f5da5899d 100644 --- a/src/vllm/model_executor/models/gemma4_mm.cpp +++ b/src/vllm/model_executor/models/gemma4_mm.cpp @@ -29,6 +29,7 @@ #include #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/ResidentWeight +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/gemma4.h" #include "vllm/model_executor/models/model_registry.h" #include "vllm/model_executor/models/qwen3_5.h" // GdnStateCache, PagedKvCache @@ -147,11 +148,8 @@ std::vector EmbedScaledBf16(Dev d, const Gemma4Weights& weights, const int64_t H = config.hidden_size; const int64_t T = static_cast(ids.size()); DBuf emb(d, DType::kBF16, {T, H}); - { - Tensor tab = ResidentWeight(d, weights.embed_tokens, {config.vocab_size, H}); - DBuf dids(d, DType::kI32, {T}, ids.data()); - vt::Embedding(d.q, emb.t(), tab, dids.t()); - } + EmbedGather(d, emb, ids, weights.embed_tokens, config.vocab_size, H, + "gemma4 mm embed"); const float nsqrt = std::sqrt(static_cast(H)); const double normalizer = static_cast(vt::BF16ToF32(vt::F32ToBF16(nsqrt))); vt::MulScalar(d.q, emb.t(), emb.t(), normalizer); diff --git a/src/vllm/model_executor/models/glm4.cpp b/src/vllm/model_executor/models/glm4.cpp index 0a9489064..19b5623e8 100644 --- a/src/vllm/model_executor/models/glm4.cpp +++ b/src/vllm/model_executor/models/glm4.cpp @@ -29,6 +29,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue (Dev/DBuf/...) +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -223,9 +224,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed: hidden[T,H] bf16 = embed_tokens[token_ids]. DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "glm4 embed"); } DBuf res(d, DType::kBF16, {T, H}); diff --git a/src/vllm/model_executor/models/glm_moe_dsa_forward.cpp b/src/vllm/model_executor/models/glm_moe_dsa_forward.cpp index bdb9b9245..c3b2c1ad4 100644 --- a/src/vllm/model_executor/models/glm_moe_dsa_forward.cpp +++ b/src/vllm/model_executor/models/glm_moe_dsa_forward.cpp @@ -57,6 +57,7 @@ #include "vllm/model_executor/moe_placement_seam.h" #include "vllm/model_executor/layers/linear.h" #include "vllm/model_executor/models/dense_attn_block.h" +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/deepseek_v2.h" // MlaStep / BuildMlaStep #include "vllm/model_executor/models/mla_attention.h" #include "vllm/platforms/interface.h" @@ -656,27 +657,12 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, expert_stream::ExpertStreamStepGuard step_guard; DBuf hidden(d, DType::kBF16, {T, p.hidden_size}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - // #1305/#2544/#2596 — TAKE the asynchronous runner's DEVICE identifiers and - // splice them over the ids just uploaded from the host. On the async serving - // path the runner's combine writes each decode row's sampled token into the - // DEVICE buffer on the main queue and leaves `token_ids` deliberately stale - // for decode rows (`v1/worker/gpu/runner.cpp`, the mirror arm, which is the - // DEFAULT on CUDA — integrated as well as discrete). Without this line this - // model embedded that stale host vector, so every step after the first - // generated from token id 0 while the run still returned rc=0. The copy is - // enqueued on the main queue, so it is ordered AFTER the combine that - // produced it rather than racing it, and it is a no-op returning false on - // every path that is not the asynchronous CUDA runner. - detail::ApplyDeviceTokenIds(d.b, d.q, dids.ptr(), T, "glm-dsa embed"); - Tensor htab = ResidentWeight(d, weights.embed_tokens, - {p.vocab_size, p.hidden_size}); - Tensor h = hidden.t(); // The embedding table is BLOCK-QUANTIZED (Q4_K on this artifact) and stays // that way: `vt::Embedding` dequantizes ONE ROW per gathered id, mirroring // `ggml_compute_forward_get_rows_q`. Expanding a `[154880, 6144]` table to // bf16 would be 1.77 GiB for the sake of a gather. - vt::Embedding(d.q, h, htab, dids.t()); + dense_attn::EmbedGather(d, hidden, token_ids, weights.embed_tokens, + p.vocab_size, p.hidden_size, "glm-dsa embed"); return ForwardLayers(d, hidden.t(), positions, am, attn_kv, weights, logits_indices); } diff --git a/src/vllm/model_executor/models/granite.cpp b/src/vllm/model_executor/models/granite.cpp index d2ad868d2..dbcbd0790 100644 --- a/src/vllm/model_executor/models/granite.cpp +++ b/src/vllm/model_executor/models/granite.cpp @@ -27,6 +27,7 @@ #include "vllm/model_executor/layers/quantization/compressed_tensors/schemes/nvfp4.h" // MakeMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -221,9 +222,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed then scale by embedding_multiplier (granite.py:313). DBuf res(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, res.t(), dtab, dids.t()); + EmbedGather(d, res, token_ids, weights.embed_tokens, vocab, H, + "granite embed"); } vt::MulScalar(d.q, res.t(), res.t(), embedding_multiplier); diff --git a/src/vllm/model_executor/models/host_embedding.cpp b/src/vllm/model_executor/models/host_embedding.cpp new file mode 100644 index 000000000..0bd1306d5 --- /dev/null +++ b/src/vllm/model_executor/models/host_embedding.cpp @@ -0,0 +1,121 @@ +// The VT_HOST_EMBEDDING gather. See host_embedding.h for the contract. +#include "vllm/model_executor/models/host_embedding.h" + +#include +#include +#include +#include + +#include "vllm/model_executor/models/qwen3_5_internal.h" // detail::TakeDeviceTokenIds +#include "vllm/model_executor/models/dense_attn_block.h" // ResidentWeight (device arm) + +namespace vllm { +namespace dense_attn { +namespace { + +vt::Queue& HostEmbedQueue() { + // Process-lifetime singleton: backend registries may outlive this TU's static + // destructors, so this queue is deliberately never destroyed. + static vt::Queue* q = + new vt::Queue(vt::GetBackend(vt::DeviceType::kCPU).CreateQueue()); + return *q; +} + +} // namespace + +bool HostEmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids, + const OwnedTensor& table, int64_t vocab, int64_t H) { + static const bool on = [] { + const char* e = std::getenv("VT_HOST_EMBEDDING"); + return e != nullptr && e[0] == '1'; + }(); + static bool logged = false; + if (!on) return false; + const int64_t T = static_cast(token_ids.size()); + if (T == 0) return true; + if (table.bytes.empty() || table.host_released) { + if (!logged) { + logged = true; + std::fprintf(stderr, + "[host-embed] DISABLED: token table host bytes unavailable " + "(dtype=%s)\n", + vt::Name(table.dtype)); + } + return false; + } + // ENG-ASYNC-SCHED W4: the async runner may have spliced this step's sampled + // token into a device-resident ids buffer, making `token_ids` stale for + // decode rows. Take that override (if any) and read the ids back; otherwise + // the host vector is authoritative and no device round-trip runs. + std::vector resolved = token_ids; + const detail::DeviceTokenIds override_ids = detail::TakeDeviceTokenIds(); + if (override_ids.ids != nullptr) { + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + d.b.Copy(d.q, dids.ptr(), override_ids.ids, + static_cast(override_ids.count) * sizeof(int32_t)); + dids.Download(d, resolved.data()); + } + if (std::getenv("VT_HOST_EMBED_TRACE") != nullptr) { + std::fprintf(stderr, "[host-embed] T=%lld override=%d\n", + static_cast(T), override_ids.ids != nullptr ? 1 : 0); + } + const vt::DType out_dtype = hidden.t().dtype; + if (table.dtype == DType::kBF16 && out_dtype == DType::kBF16) { + // Byte-copy fast path: the gathered rows are already in the output dtype. + std::vector rows(static_cast(T) * static_cast(H)); + const uint16_t* src = reinterpret_cast(table.bytes.data()); + for (int64_t t = 0; t < T; ++t) { + const int32_t id = resolved[static_cast(t)]; + VT_CHECK(id >= 0 && id < vocab, "host embedding: token id out of range"); + std::memcpy(rows.data() + static_cast(t) * static_cast(H), + src + static_cast(id) * static_cast(H), + static_cast(H) * sizeof(uint16_t)); + } + d.b.Copy(d.q, hidden.ptr(), rows.data(), rows.size() * sizeof(uint16_t)); + if (!logged) { + logged = true; + std::fprintf(stderr, "[host-embed] bf16 row copy (T=%lld)\n", + static_cast(T)); + } + return true; + } + const vt::Device cpu{}; + vt::Tensor table_view = + table.ViewOn(const_cast(table.bytes.data()), cpu, {vocab, H}); + std::vector staging(static_cast(T) * + static_cast(H) * vt::SizeOf(out_dtype)); + vt::Tensor out = MakeTensor(staging.data(), out_dtype, cpu, {T, H}); + vt::Tensor ids = MakeTensor(resolved.data(), DType::kI32, cpu, {T}); + vt::Embedding(HostEmbedQueue(), out, table_view, ids); + if (!logged) { + logged = true; + std::fprintf(stderr, "[host-embed] CPU row gather (T=%lld dtype=%s out=%s)\n", + static_cast(T), vt::Name(table.dtype), + vt::Name(out_dtype)); + } + d.b.Copy(d.q, hidden.ptr(), staging.data(), staging.size()); + return true; +} + +void EmbedGather(Dev d, DBuf& out, const std::vector& token_ids, + const OwnedTensor& table, int64_t vocab, int64_t H, + const char* what) { + if (HostEmbedInto(d, out, token_ids, table, vocab, H)) return; + const int64_t T = static_cast(token_ids.size()); + Tensor dtab = ResidentWeight(d, table, {vocab, H}); + // ROW-SERVE-ASYNC-DENSE-MIRROR (ENG-ASYNC-SCHED W4, #1305): on the async + // serving loop the runner's device combine splices each decode row's sampled + // token into a device-resident id buffer on the MAIN QUEUE while the host + // `token_ids` vector stays stale BY DESIGN — materializing it on the host is + // the synchronize that path removes. Splicing the override over the upload here + // is ordered AFTER the combine rather than racing it, and it is CLEARED on + // first use so a second, unrelated embed cannot be handed this step's rows. + // Null on every other path, where the copy is a no-op and the caller is + // byte-identical to its pre-#1305 self. + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + detail::ApplyDeviceTokenIds(d.b, d.q, dids.ptr(), T, what); + vt::Embedding(d.q, out.t(), dtab, dids.t()); +} + +} // namespace dense_attn +} // namespace vllm diff --git a/src/vllm/model_executor/models/minicpm.cpp b/src/vllm/model_executor/models/minicpm.cpp index 5949a6541..89245e4c8 100644 --- a/src/vllm/model_executor/models/minicpm.cpp +++ b/src/vllm/model_executor/models/minicpm.cpp @@ -26,6 +26,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -225,9 +226,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed then scale by scale_emb (minicpm.py:441-443). DBuf res(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, res.t(), dtab, dids.t()); + EmbedGather(d, res, token_ids, weights.embed_tokens, vocab, H, + "minicpm embed"); } vt::MulScalar(d.q, res.t(), res.t(), scale_emb); diff --git a/src/vllm/model_executor/models/minicpm3.cpp b/src/vllm/model_executor/models/minicpm3.cpp index 98117965e..f53f04bcb 100644 --- a/src/vllm/model_executor/models/minicpm3.cpp +++ b/src/vllm/model_executor/models/minicpm3.cpp @@ -33,6 +33,7 @@ #include "vllm/model_executor/layers/attention/mla_chunked_context.h" #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/ResidentWeight glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/deepseek_v2.h" // BuildMlaBatchSplit/MlaBatchSplit #include "vllm/model_executor/models/device_pool.h" // Pool() #include "vllm/model_executor/models/mla_attention.h" @@ -245,9 +246,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed then scale by scale_emb (minicpm.py:441-443). DBuf res(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, res.t(), dtab, dids.t()); + EmbedGather(d, res, token_ids, weights.embed_tokens, vocab, H, + "minicpm3 embed"); } vt::MulScalar(d.q, res.t(), res.t(), p.scale_emb); diff --git a/src/vllm/model_executor/models/muse_glimmer.cpp b/src/vllm/model_executor/models/muse_glimmer.cpp index c78f6b2c7..698de716e 100644 --- a/src/vllm/model_executor/models/muse_glimmer.cpp +++ b/src/vllm/model_executor/models/muse_glimmer.cpp @@ -54,6 +54,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -392,10 +393,9 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // a buffer it reuses across steps (#2300). d.b.Copy(d.q, hidden.ptr(), inputs_embeds->data, hidden.bytes()); } else { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); DBuf emb(d, DType::kBF16, {T, H}); - vt::Embedding(d.q, emb.t(), dtab, dids.t()); + EmbedGather(d, emb, token_ids, weights.embed_tokens, vocab, H, + "muse_glimmer embed"); vt::RmsNorm(d.q, hidden.t(), emb.t(), ones_hidden.t(), vt::RmsNormArgs{g.norm_eps, /*gemma=*/false}); } diff --git a/src/vllm/model_executor/models/muse_glimmer_mm.cpp b/src/vllm/model_executor/models/muse_glimmer_mm.cpp index 0583fd24e..9ecc02a63 100644 --- a/src/vllm/model_executor/models/muse_glimmer_mm.cpp +++ b/src/vllm/model_executor/models/muse_glimmer_mm.cpp @@ -46,6 +46,7 @@ #include #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/ResidentWeight +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/model_registry.h" #include "vllm/model_executor/models/muse_glimmer.h" #include "vllm/model_executor/models/qwen3_5.h" // PagedKvCache, GdnStateCache @@ -261,11 +262,12 @@ std::vector MuseGlimmerMergeMultimodalEmbeds( { std::vector ones_host(static_cast(H), vt::F32ToBF16(1.0f)); DBuf ones(d, DType::kBF16, {H}, ones_host.data()); - Tensor tab = ResidentWeight(d, weights.embed_tokens, {t.vocab_size, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); DBuf emb(d, DType::kBF16, {T, H}); DBuf normed(d, DType::kBF16, {T, H}); - vt::Embedding(d.q, emb.t(), tab, dids.t()); + // This function downloads the normed rows anyway, so the host gather also + // skips building the device table on this arm. + EmbedGather(d, emb, token_ids, weights.embed_tokens, t.vocab_size, H, + "muse_glimmer mm embed"); vt::RmsNorm(d.q, normed.t(), emb.t(), ones.t(), vt::RmsNormArgs{t.rms_norm_eps, /*gemma=*/false}); normed.Download(d, bits.data()); diff --git a/src/vllm/model_executor/models/nemotron_h_device.cpp b/src/vllm/model_executor/models/nemotron_h_device.cpp index 54f26ea0e..23add2149 100644 --- a/src/vllm/model_executor/models/nemotron_h_device.cpp +++ b/src/vllm/model_executor/models/nemotron_h_device.cpp @@ -104,6 +104,7 @@ // header rather than dense_device_glue.h — `ResidentWeight`, the lazy // upload-once seam this row converts NemotronH's dense weights onto. #include "vllm/model_executor/models/dense_attn_block.h" +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/moe_placement_seam.h" #include "vt/backend.h" #include "vt/ops.h" @@ -1085,10 +1086,8 @@ std::vector NemotronHDeviceForward(const NemotronHHostWeights& host, for (int32_t id : ids) { VT_CHECK(id >= 0 && id < V, "NemotronH device forward: token id out of range"); } - DBuf it(d, DType::kI32, {T}, ids.data()); - d.b.Synchronize(d.q); // `ids` is a local; see UploadAs for why this waits. - Tensor tab = ResidentWeight(d, host.embeddings); - vt::Embedding(d.q, residual.t(), tab, it.t()); + EmbedGather(d, residual, token_ids, host.embeddings, V, H, + "nemotron_h embed"); } vt::RmsNormArgs nargs; diff --git a/src/vllm/model_executor/models/olmo2.cpp b/src/vllm/model_executor/models/olmo2.cpp index af80ebc6a..a41f7c16c 100644 --- a/src/vllm/model_executor/models/olmo2.cpp +++ b/src/vllm/model_executor/models/olmo2.cpp @@ -31,6 +31,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue (Dev/DBuf/...) +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -279,9 +280,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, // Embed: hidden[T,H] bf16 = embed_tokens[token_ids]. DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "olmo2 embed"); } StepInputs si = BuildStepInputs(d, positions, attn_meta, config); diff --git a/src/vllm/model_executor/models/opt.cpp b/src/vllm/model_executor/models/opt.cpp index f0b93a759..6d459ef7e 100644 --- a/src/vllm/model_executor/models/opt.cpp +++ b/src/vllm/model_executor/models/opt.cpp @@ -41,6 +41,7 @@ #include #include "vllm/model_executor/models/dense_attn_block.h" // Dev/DBuf/ResidentWeight/KvSlice glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool/ActivePool (shared) #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vllm/platforms/interface.h" @@ -260,9 +261,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, OPTStepInputs si = BuildOPTStepInputs(d, positions, attn_meta); DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "opt embed"); const int64_t p_rows = weights.embed_positions.shape[0]; Tensor ptab = ResidentWeight(d, weights.embed_positions, {p_rows, H}); diff --git a/src/vllm/model_executor/models/phi.cpp b/src/vllm/model_executor/models/phi.cpp index 84c16c683..7deff4377 100644 --- a/src/vllm/model_executor/models/phi.cpp +++ b/src/vllm/model_executor/models/phi.cpp @@ -32,6 +32,7 @@ #include #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -193,9 +194,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "phi embed"); } StepInputs si = BuildStepInputs(d, positions, attn_meta, config); diff --git a/src/vllm/model_executor/models/phi3.cpp b/src/vllm/model_executor/models/phi3.cpp index 1009ded6e..7532c3289 100644 --- a/src/vllm/model_executor/models/phi3.cpp +++ b/src/vllm/model_executor/models/phi3.cpp @@ -22,6 +22,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -178,9 +179,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, dw.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, dw.embed_tokens, vocab, H, + "phi3 embed"); } DBuf res(d, DType::kBF16, {T, H}); diff --git a/src/vllm/model_executor/models/qwen3.cpp b/src/vllm/model_executor/models/qwen3.cpp index 15ba698ee..0b8c75a2b 100644 --- a/src/vllm/model_executor/models/qwen3.cpp +++ b/src/vllm/model_executor/models/qwen3.cpp @@ -52,6 +52,7 @@ #include "vllm/model_executor/layers/quantization/compressed_tensors/schemes/nvfp4.h" // LinearMethod seam #include "vllm/model_executor/models/decode_graph_sizes.h" // DecodeGraphSizes/PadToCaptureSize #include "vllm/model_executor/models/dense_attn_block.h" // shared AttnBlock + device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/dense_nvfp4_gemm.h" // NVFP4 W4A16 dispatch #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool/ActivePool (shared) #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits @@ -216,54 +217,17 @@ void GatherRows(Dev d, void* dst, const Tensor& src, const std::vector& d.b.Copy(d.q, dp + s * rb, sp + static_cast(idx[s]) * rb, rb); } -// ROW-SERVE-ASYNC-DENSE-MIRROR (ENG-ASYNC-SCHED W4 / the #31 P0, ported to the -// classic dense family): overwrite the REAL prefix of a freshly uploaded input-id -// buffer with the device-resident ids the async runner's combine produced. The -// exact analogue of qwen3_5.cpp's ApplyDeviceTokenIdsOverride — that TU wired it -// for the gate models (MoE + 27B dense); this is the identical consumer for the -// SHARED pure-dense driver (Qwen3ForCausalLM and every registry that routes -// through Qwen3DenseModel / EmbedInto: InternLM2, Mistral, Llama). -// -// WHY: on the async serving loop (AsyncLLM depth-2) the sampled token is NOT -// written to token_ids_cpu synchronously; the runner's device combine splices each -// decode row's real token into the device input-ids on the MAIN QUEUE while the -// host `token_ids` vector stays stale. The default host upload below then RACES -// that device write (unsynchronized device-write/host-read), nondeterministically -// embedding the stale/zero placeholder -> token-0 degeneration. Copying the -// device ids over the DBuf prefix here is main-queue-ordered AFTER the combine, so -// the embed never does the racing host read — exactly upstream (states.py:64 -// device-resident prev_sampled_token_ids + gpu_model_runner.py GPU gather). -// -// The override is published by the registry forward's detail::DeviceTokenIdsScope -// and CONSUMED here on first use; null on every path except the CUDA async runner, -// so with no override this is byte-identical to the pre-fix host upload. -// #1305: the take-and-clear and the bounds-checked copy this used to spell out -// are `detail::ApplyDeviceTokenIds` (`qwen3_5_internal.h`), one body for the four -// models that consume the scope. Behaviour, ordering and the refusal message are -// unchanged; only the copy count is. -static void ApplyDeviceTokenIdsOverride(Dev d, DBuf& dids, int64_t T) { - detail::ApplyDeviceTokenIds(d.b, d.q, dids.ptr(), T, "qwen3 dense embed"); -} - -// Embed: hidden[T,H] bf16 = embed_tokens[token_ids] (device-resident table). KEPT -// OUTSIDE THE CUDA-GRAPH (mirrors qwen3_moe.cpp / qwen3_5.cpp EmbedInto): the CUDA -// Embedding op allocates a device bounds-check flag (cudaMalloc/cudaFree) and syncs -// the stream, both illegal inside a capture region — and it consumes the HOST +// Embed: hidden[T,H] bf16 = embed_tokens[token_ids]. KEPT OUTSIDE THE +// CUDA-GRAPH (mirrors qwen3_moe.cpp / qwen3_5.cpp EmbedInto): the CUDA Embedding +// op allocates a device bounds-check flag (cudaMalloc/cudaFree) and syncs the +// stream, both illegal inside a capture region — and it consumes the HOST // token_ids. The graph driver runs this per step into its PERSISTENT hidden buffer, -// then captures/replays ForwardLayers over that fixed hidden address. +// then captures/replays ForwardLayers over that fixed hidden address. The table +// residency and the async device-id override are `EmbedGather`'s. void EmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids, const Qwen3DenseWeights& weights, const HfConfig& config) { - const int64_t T = static_cast(token_ids.size()); - Tensor dtab = ResidentWeight(d, weights.embed_tokens, - {config.vocab_size, config.hidden_size}); - // ROW-SERVE-ASYNC-DENSE-MIRROR: when the async runner has already placed this - // step's input ids on the device (and spliced each decode row's sampled token - // into them there), embed straight from that buffer. `token_ids` is stale for - // decode rows in that case BY DESIGN — materializing it on the host is the - // synchronize the async path removes — so its real prefix is overwritten here. - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - ApplyDeviceTokenIdsOverride(d, dids, T); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, config.vocab_size, + config.hidden_size, "qwen3 dense embed"); } // The CAPTURABLE region: everything AFTER the embedding — the residual stream diff --git a/src/vllm/model_executor/models/qwen3_5.cpp b/src/vllm/model_executor/models/qwen3_5.cpp index 0900a64a6..d6e62c4dc 100644 --- a/src/vllm/model_executor/models/qwen3_5.cpp +++ b/src/vllm/model_executor/models/qwen3_5.cpp @@ -22,6 +22,7 @@ #include "vllm/model_executor/models/dense_exl3_linear.h" // MODEL-QWEN35-EXL3 (#2495): the EXL3 linear seam #include "vllm/model_executor/models/dense_fp8_block_gemm.h" // MODEL-FP8-BLOCK-LINEAR (#1189 M4) #include "vllm/model_executor/models/dense_device_glue.h" +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool/AuxPool/ActivePool (shared) #include "vt/tenstorrent/tenstorrent_device.h" // DebugDeviceReadbackF32 (TT-only debug seam) @@ -807,6 +808,7 @@ using v1::GDNAttentionMetadata; // migrating them is a separate change with its own gate. using dense_attn::DBuf; using dense_attn::Dev; +using dense_attn::HostEmbedInto; using dense_attn::MakeTensor; using dense_attn::Reshape; using dense_attn::ResolveDevicePoolPolicy; @@ -8287,6 +8289,7 @@ void RunDenseLayerPaged(Dev d, const Qwen3_5DenseLayerWeights& layer, // runs the (dense or MoE) decoder layer + final norm over it. `embed_tokens` is // the shared target embedding; `target_hidden_states` is the target model's // post-final-norm bf16 [T,H] output (the drafter's hidden-state tap). + DBuf MtpHeadHidden(Dev device, const Qwen3_5MTPWeights& weights, const HfConfig& config, const OwnedTensor& embed_tokens, const std::vector& input_ids, @@ -8295,11 +8298,14 @@ DBuf MtpHeadHidden(Dev device, const Qwen3_5MTPWeights& weights, const int64_t vocab_size = config.vocab_size; const float eps = static_cast(config.rms_norm_eps); - Tensor embedding_table = Qwen3_5EmbeddingTable(device.b, device.q, embed_tokens, - vocab_size, hidden_size); - DBuf device_ids(device, DType::kI32, {tokens}, input_ids.data()); DBuf embedding(device, DType::kBF16, {tokens, hidden_size}); - vt::Embedding(device.q, embedding.t(), embedding_table, device_ids.t()); + if (!HostEmbedInto(device, embedding, input_ids, embed_tokens, vocab_size, + hidden_size)) { + Tensor embedding_table = Qwen3_5EmbeddingTable( + device.b, device.q, embed_tokens, vocab_size, hidden_size); + DBuf device_ids(device, DType::kI32, {tokens}, input_ids.data()); + vt::Embedding(device.q, embedding.t(), embedding_table, device_ids.t()); + } Tensor embedding_norm_weight = ResidentWeight(device, weights.pre_fc_norm_embedding, {hidden_size}); @@ -8760,17 +8766,19 @@ static void EmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids std::fprintf(stderr, "[TT-FWD] EmbedInto T=%lld H=%lld\n", static_cast(T), static_cast(H)); const int64_t vocab = config.vocab_size; - Tensor dtab = - Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); - // ENG-ASYNC-SCHED W4: when the async runner has already placed this step's - // input ids on the device (and spliced each decode row's sampled token into - // them there), embed straight from that buffer. `token_ids` is stale for - // decode rows in that case BY DESIGN — materializing it on the host is the - // synchronize W4 removes — so it must not be uploaded here. Its SIZE is still - // authoritative: the runner sized the device buffer from the same step. - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - ApplyDeviceTokenIdsOverride(d, dids, T); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + if (!HostEmbedInto(d, hidden, token_ids, weights.embed_tokens, vocab, H)) { + Tensor dtab = + Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); + // ENG-ASYNC-SCHED W4: when the async runner has already placed this step's + // input ids on the device (and spliced each decode row's sampled token into + // them there), embed straight from that buffer. `token_ids` is stale for + // decode rows in that case BY DESIGN — materializing it on the host is the + // synchronize W4 removes — so it must not be uploaded here. Its SIZE is still + // authoritative: the runner sized the device buffer from the same step. + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + ApplyDeviceTokenIdsOverride(d, dids, T); + vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + } } // DFlash DF-AUX-TAPS (SPEC-DFLASH D1) — capture the residual-stream value at a @@ -9430,12 +9438,14 @@ std::vector Qwen3_5Model::ForwardDense(const std::vector& token_ Dev d{vt::GetBackend(queue.device.type), queue}; const float eps = static_cast(config.rms_norm_eps); - // Embed: hidden = embed_tokens[token_ids] (bf16, device-resident). res = 0. - Tensor dtab = - Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); + // Embed: hidden = embed_tokens[token_ids]. res = 0. DBuf hidden(d, ActDType(d), {T, H}); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + if (!HostEmbedInto(d, hidden, token_ids, weights.embed_tokens, vocab, H)) { + Tensor dtab = + Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + } DBuf res(d, ResidualDType(d), {T, H}); res.Zero(d); @@ -9479,11 +9489,14 @@ std::vector Qwen3_5Model::ForwardMoeHidden( Dev d{vt::GetBackend(queue.device.type), queue}; const float eps = static_cast(config.rms_norm_eps); - Tensor dtab = - Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, config.vocab_size, H); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); DBuf hidden(d, ActDType(d), {T, H}); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + if (!HostEmbedInto(d, hidden, token_ids, weights.embed_tokens, + config.vocab_size, H)) { + Tensor dtab = Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, + config.vocab_size, H); + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + } DBuf res(d, ResidualDType(d), {T, H}); res.Zero(d); @@ -9519,15 +9532,17 @@ std::vector Qwen3_5DenseModel::ForwardDense( Dev d{vt::GetBackend(queue.device.type), queue}; const float eps = static_cast(config.rms_norm_eps); - // Embed: hidden = embed_tokens[token_ids] (bf16, device-resident). res = 0. + // Embed: hidden = embed_tokens[token_ids]. res = 0. // For a TEXT-only step the three mRoPE position streams are identical, so the // partial NeoX RoPE in FullAttnBlock degenerates to 1-D RoPE over `positions` // (notes §2). The vision tower / image-video merger are DEFERRED. - Tensor dtab = - Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); DBuf hidden(d, ActDType(d), {T, H}); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + if (!HostEmbedInto(d, hidden, token_ids, weights.embed_tokens, vocab, H)) { + Tensor dtab = + Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + } DBuf res(d, ResidualDType(d), {T, H}); res.Zero(d); @@ -9566,11 +9581,14 @@ std::vector Qwen3_5DenseModel::ForwardDenseHidden( Dev d{vt::GetBackend(queue.device.type), queue}; const float eps = static_cast(config.rms_norm_eps); - Tensor dtab = - Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, config.vocab_size, H); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); DBuf hidden(d, ActDType(d), {T, H}); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + if (!HostEmbedInto(d, hidden, token_ids, weights.embed_tokens, + config.vocab_size, H)) { + Tensor dtab = Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, + config.vocab_size, H); + DBuf dids(d, DType::kI32, {T}, token_ids.data()); + vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + } DBuf res(d, ResidualDType(d), {T, H}); res.Zero(d); @@ -9927,6 +9945,9 @@ static void DenseEmbedInto(Dev d, DBuf& hidden, const int64_t T = static_cast(token_ids.size()); const int64_t H = config.hidden_size; const int64_t vocab = config.vocab_size; + if (HostEmbedInto(d, hidden, token_ids, weights.embed_tokens, vocab, H)) { + return; + } Tensor dtab = Qwen3_5EmbeddingTable(d.b, d.q, weights.embed_tokens, vocab, H); DBuf dids(d, DType::kI32, {T}, token_ids.data()); diff --git a/src/vllm/model_executor/models/stablelm.cpp b/src/vllm/model_executor/models/stablelm.cpp index 028028abe..024e073e6 100644 --- a/src/vllm/model_executor/models/stablelm.cpp +++ b/src/vllm/model_executor/models/stablelm.cpp @@ -25,6 +25,7 @@ #include "vllm/model_executor/layers/linear.h" // UnquantizedMlpGateUpMethod seam #include "vllm/model_executor/models/dense_attn_block.h" // shared device glue +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/device_pool.h" // DevicePool/Pool #include "vllm/model_executor/models/qwen3_5_common.h" // HostLogits #include "vt/backend.h" @@ -193,9 +194,8 @@ DBuf ForwardBody(Dev d, const std::vector& token_ids, DBuf hidden(d, DType::kBF16, {T, H}); { - Tensor dtab = ResidentWeight(d, weights.embed_tokens, {vocab, H}); - DBuf dids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), dtab, dids.t()); + EmbedGather(d, hidden, token_ids, weights.embed_tokens, vocab, H, + "stablelm embed"); } StepInputs si = BuildStepInputs(d, positions, attn_meta, config); diff --git a/src/vllm/model_executor/models/voxtral.cpp b/src/vllm/model_executor/models/voxtral.cpp index b8b67867a..f6d33ca50 100644 --- a/src/vllm/model_executor/models/voxtral.cpp +++ b/src/vllm/model_executor/models/voxtral.cpp @@ -25,6 +25,7 @@ #include "vllm/model_executor/model_loader/safetensors_reader.h" #include "vllm/model_executor/models/decode_graph_sizes.h" // DecodeGraphSizes/PadToCaptureSize #include "vllm/model_executor/models/dense_attn_block.h" // AttnBlock, BuildStepInputs, ResidentWeight +#include "vllm/model_executor/models/host_embedding.h" // VT_HOST_EMBEDDING gather #include "vllm/model_executor/models/dense_weight_loaders.h" #include "vllm/model_executor/models/qwen3_vl_text.h" // Qwen3VLMergeMultimodal (modality-agnostic merge) #include "vllm/model_executor/models/voxtral_loader_internal.h" // the two mmap-reading loader steps (#772) @@ -191,11 +192,9 @@ using dense_attn::MakeTensor; // then captures/replays ForwardLastLogits over that fixed hidden address. void VoxtralEmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids, const Qwen3DenseWeights& weights, const HfConfig& config) { - const int64_t T = static_cast(token_ids.size()); - Tensor tab = - ResidentWeight(d, weights.embed_tokens, {config.vocab_size, config.hidden_size}); - DBuf ids(d, DType::kI32, {T}, token_ids.data()); - vt::Embedding(d.q, hidden.t(), tab, ids.t()); + dense_attn::EmbedGather(d, hidden, token_ids, weights.embed_tokens, + config.vocab_size, config.hidden_size, + "voxtral embed"); } // Overwrite dst's CONTENTS from src WITHOUT changing dst.data() when the sizes diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index e2c9d553a..7e9b53b4a 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -4097,6 +4097,12 @@ vllm_cpp_add_test(test_qwen3_5_lm_head_dtypes vllm_cpp_add_test(test_qwen3_5_dense_load_residency vllm/models/test_qwen3_5_dense_load_residency.cpp) +# ENG-HOST-EMBEDDING: the host-resident token table. The host arm must produce +# the pinned-oracle rows AND leave the table off the device (d_dev null), which +# is the half a gather that uploaded first would pass. The flag is process-static, +# so the binary enables it before the first gather and is flag-ON by construction. +vllm_cpp_add_test(test_host_embedding vllm/models/test_host_embedding.cpp) + # QUANT-QWEN38-27B-NVFP4-ARM W4 (#821, .agents/specs/qwen38-27b-quant-arms.md): # the compressed-tensors `mixed-precision` accounting + scheme-resolution gate # for `unsloth/Qwen3.8-27B-NVFP4`. Hermetic — it reads the two COMMITTED diff --git a/tests/vllm/models/test_host_embedding.cpp b/tests/vllm/models/test_host_embedding.cpp new file mode 100644 index 000000000..7f5c3c88d --- /dev/null +++ b/tests/vllm/models/test_host_embedding.cpp @@ -0,0 +1,271 @@ +// VT_HOST_EMBEDDING: the HOST arm of the token-table gather. +// +// The arm exists to keep `[vocab, H]` out of device memory: the table stays in +// host RAM, the requested rows are gathered through the CPU `vt::Embedding` +// kernel (ONE ROW per id, the `ggml_get_rows` discipline), and only the `[T,H]` +// result is copied to the device. The gate therefore has two halves, and both +// matter: +// +// 1. the rows are RIGHT, byte-compared against the same pinned IQ4_NL/Q5_0 +// oracle vectors the device-side `vt::Embedding` gate uses +// (`tests/vt/iq4nl_q5_0_golden_vectors.h`), and +// 2. the table never reached the device — `OwnedTensor::d_dev` stays null. A +// gather that uploads the table first passes (1) and fails the entire point. +// +// The flag is read ONCE per process (`HostEmbedInto`'s function-static), so this +// binary turns it on in a global initializer, before main and thus before any +// TEST_CASE. The flag-OFF arm is a plain early return, and every other model +// suite in the tree exercises it (none sets the variable); what is pinned here for +// the boundary between the arms is the DECLINE, which is the released-host-bytes +// case below. +#include + +#include +#include +#include +#include +#include +#include +#include + +#if defined(__unix__) || defined(__APPLE__) +#include +#define VLLM_HOST_EMBED_CAPTURE 1 +#endif + +#include "vllm/model_executor/models/dense_device_glue.h" // Dev, DBuf +#include "vllm/model_executor/models/host_embedding.h" // HostEmbedInto, EmbedGather +#include "vllm/model_executor/models/owned_bytes.h" +#include "vllm/model_executor/models/qwen3_5_weights.h" // OwnedTensor +#include "vt/dtype.h" +#include "vt/ops.h" + +#include "vt/iq4nl_q5_0_golden_vectors.h" + +namespace { + +struct EnableHostEmbedding { + EnableHostEmbedding() { + ::setenv("VT_HOST_EMBEDDING", "1", 1); + // The arm prints ONE banner per process (`logged`), so the first case is the + // one that can read it. It is what proves the arm RAN: the CPU device arm + // aliases the host bytes and sets no `d_dev`, so `d_dev == nullptr` alone + // cannot tell the two arms apart on this backend. + ::setenv("VT_HOST_EMBED_TRACE", "1", 1); + } +}; +const EnableHostEmbedding g_enable_host_embedding; + +#ifdef VLLM_HOST_EMBED_CAPTURE +// Everything written to stderr while `body` runs, as a string. Same shape as the +// qwen4_exp MoE tap's capture: the arm's banner is the observable, so a test that +// cannot read it must not report that the arm ran. +template +std::string CaptureStderr(F&& body) { + std::fflush(stderr); + int saved = ::dup(2); + char path[] = "/tmp/vllm_hostembed_XXXXXX"; + int fd = ::mkstemp(path); + REQUIRE(saved >= 0); + REQUIRE(fd >= 0); + ::dup2(fd, 2); + body(); + std::fflush(stderr); + ::dup2(saved, 2); + ::close(saved); + ::lseek(fd, 0, SEEK_SET); + std::string out; + char buf[4096]; + ssize_t n = 0; + while ((n = ::read(fd, buf, sizeof(buf))) > 0) out.append(buf, static_cast(n)); + ::close(fd); + ::unlink(path); + return out; +} +#endif + +uint32_t BitsToF32(float f) { + uint32_t bits = 0; + std::memcpy(&bits, &f, sizeof(bits)); + return bits; +} +float BitsToFloat(uint32_t bits) { + float f = 0.0F; + std::memcpy(&f, &bits, sizeof(f)); + return f; +} + +// The loader's residency for a gather table: an OWNED host buffer with the +// descriptor a forward reads. No device handles — creating them is the one thing +// the host arm must not do. +vllm::OwnedTensor HostTable(vt::DType dt, int64_t rows, int64_t k, + const void* bytes, size_t nbytes) { + std::vector copy(nbytes); + std::memcpy(copy.data(), bytes, nbytes); + vllm::OwnedTensor t; + t.dtype = dt; + t.rank = 2; + t.shape[0] = rows; + t.shape[1] = k; + t.bytes = vllm::OwnedBytes(std::move(copy)); + return t; +} + +} // namespace + +TEST_CASE("host embedding: IQ4_NL rows match the pinned oracle, table stays off-device") { + constexpr int64_t k = 160; // 5 whole IQ4_NL blocks per row + constexpr int64_t rows = 2; + vllm::OwnedTensor table = + HostTable(vt::DType::kIQ4_NL, rows, k, vllm_test::kIq4nlGoldenBlocks, + sizeof(vllm_test::kIq4nlGoldenBlocks)); + // Repeats and a backwards step, so a gather that ignored the id or walked the + // table in order cannot pass (the same ids the device-side op test uses). + const std::vector ids{1, 0, 1, 1, 0}; + const int64_t T = static_cast(ids.size()); + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kBF16, {T, k}); +#ifdef VLLM_HOST_EMBED_CAPTURE + // (0) THE ARM RAN. This is the assertion that separates the host arm from the + // device arm on a CPU backend, where `ResidentWeight` aliases the host bytes + // and the `d_dev` check below is trivially true either way. With the flag off + // this capture is empty and the case fails. + const std::string trace = CaptureStderr([&] { + vllm::dense_attn::EmbedGather(d, out, ids, table, rows, k, + "test_host_embedding"); + }); + CHECK(trace.find("[host-embed]") != std::string::npos); + CHECK(trace.find("CPU row gather") != std::string::npos); +#else + vllm::dense_attn::EmbedGather(d, out, ids, table, rows, k, + "test_host_embedding"); +#endif + + // (2) THE POINT OF THE ARM on a backend that stages: no device resident was + // created. On the CPU backend the device arm aliases instead of staging, so + // this half is not the discriminating one here — (0) is. + CHECK(table.d_dev == nullptr); + CHECK(table.HasHostBytes()); + + // (1) The rows, against the pinned oracle, rounded once into bf16 — the same + // comparison `test_ops_embedding_quant.cpp` makes on the device arm. + const auto* got = static_cast(out.ptr()); + for (int64_t t = 0; t < T; ++t) { + for (int64_t j = 0; j < k; ++j) { + CAPTURE(t); + CAPTURE(j); + const float want = BitsToFloat( + vllm_test::kIq4nlGoldenBits[static_cast(ids[static_cast(t)]) * + static_cast(k) + + static_cast(j)]); + CHECK(got[static_cast(t) * static_cast(k) + + static_cast(j)] == vt::F32ToBF16(want)); + } + } +} + +TEST_CASE("host embedding: a bf16 table takes the byte-copy fast path") { + constexpr int64_t k = 8; + constexpr int64_t rows = 3; + std::vector values(static_cast(rows * k)); + for (size_t i = 0; i < values.size(); ++i) { + values[i] = vt::F32ToBF16(static_cast(i) * 0.5F - 2.0F); + } + vllm::OwnedTensor table = + HostTable(vt::DType::kBF16, rows, k, values.data(), values.size() * 2); + const std::vector ids{2, 0, 1, 2}; + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kBF16, + {static_cast(ids.size()), k}); + // The arm ACCEPTS (the process-wide banner already printed in the first case, + // so the return value is the assertion here rather than the trace). + CHECK(vllm::dense_attn::HostEmbedInto(d, out, ids, table, rows, k)); + vllm::dense_attn::EmbedGather(d, out, ids, table, rows, k, + "test_host_embedding"); + + CHECK(table.d_dev == nullptr); + const auto* got = static_cast(out.ptr()); + for (size_t t = 0; t < ids.size(); ++t) { + for (int64_t j = 0; j < k; ++j) { + CAPTURE(t); + CAPTURE(j); + CHECK(got[t * static_cast(k) + static_cast(j)] == + values[static_cast(ids[t]) * static_cast(k) + + static_cast(j)]); + } + } +} + +TEST_CASE("host embedding: an f16 table converts through the CPU kernel") { + constexpr int64_t k = 4; + constexpr int64_t rows = 2; + std::vector halved(static_cast(rows * k)); + for (size_t i = 0; i < halved.size(); ++i) { + halved[i] = vt::F32ToF16(static_cast(i) + 0.25F); + } + vllm::OwnedTensor table = + HostTable(vt::DType::kF16, rows, k, halved.data(), halved.size() * 2); + const std::vector ids{1, 0}; + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kF32, + {static_cast(ids.size()), k}); + CHECK(vllm::dense_attn::HostEmbedInto(d, out, ids, table, rows, k)); + vllm::dense_attn::EmbedGather(d, out, ids, table, rows, k, + "test_host_embedding"); + + CHECK(table.d_dev == nullptr); + const auto* got = static_cast(out.ptr()); + for (size_t t = 0; t < ids.size(); ++t) { + for (int64_t j = 0; j < k; ++j) { + CAPTURE(t); + CAPTURE(j); + const float want = + vt::F16ToF32(halved[static_cast(ids[t]) * static_cast(k) + + static_cast(j)]); + CHECK(BitsToF32(got[t * static_cast(k) + static_cast(j)]) == + BitsToF32(want)); + } + } +} + +TEST_CASE("host embedding: an out-of-range id is refused") { + constexpr int64_t k = 4; + constexpr int64_t rows = 2; + const std::vector values(static_cast(rows * k), + vt::F32ToBF16(1.0F)); + vllm::OwnedTensor table = + HostTable(vt::DType::kBF16, rows, k, values.data(), values.size() * 2); + const std::vector ids{2}; // rows == 2, so 2 is out of range + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kBF16, {1, k}); + CHECK_THROWS_AS(vllm::dense_attn::EmbedGather(d, out, ids, table, rows, k, + "test_host_embedding"), + std::runtime_error); +} + +TEST_CASE("host embedding: a released host table declines the host arm") { + constexpr int64_t k = 4; + constexpr int64_t rows = 2; + const std::vector values(static_cast(rows * k), + vt::F32ToBF16(1.0F)); + vllm::OwnedTensor table = + HostTable(vt::DType::kBF16, rows, k, values.data(), values.size() * 2); + table.host_released = true; + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kBF16, {1, k}); + const std::vector ids{0}; + + // The boundary between the arms: the host arm must DECLINE, so `EmbedGather` + // falls through to the device arm rather than gathering bytes that are gone. + CHECK_FALSE(vllm::dense_attn::HostEmbedInto(d, out, ids, table, rows, k)); +} From d686c8582cddc8d2cf15c05c3be3509966ccb541 Mon Sep 17 00:00:00 2001 From: dev Date: Wed, 30 Sep 2026 07:41:34 +0100 Subject: [PATCH 3/3] fix(ENG-HOST-EMBEDDING): bound the host arm's device-id override The upstream review found that `HostEmbedInto` copied `override_ids.count` identifiers into a `[T]` buffer with no bound, so a two-element override with T=1 wrote 8 bytes into 4. The device arm already refuses that mismatch in `detail::ApplyDeviceTokenIds`, which also tolerates a shorter override as the padded case. The host arm now splices through that same body before downloading the identifiers: a longer override throws with the caller's name, and a shorter one replaces its prefix while the host upload's tail stays. Two cases in `test_host_embedding` pin both boundaries, which the suite previously never exercised at all. The same review found the test's global initializer called POSIX `::setenv` unconditionally. The target is registered unconditionally, so a CPU-only MSVC build could not compile it. It now uses `vllm_test::SetEnv` from `tests/support/test_env.h`. The override cases reach the internal `detail::DeviceTokenIdsScope` declaration, so the target takes `${CMAKE_SOURCE_DIR}/src` on its include path like the other detail-seam gates. Verified on a Visual Studio 2022 Build Tools 17.14 x64 Release CPU-only build: `test_host_embedding` compiles, runs 7/7 cases and 862/862 assertions, and `test_ops_embedding_quant` stays at 1637/1637. Running any vllm-linked binary on this host needs a one-line uncommitted unblock of a pre-existing MSVC static-init crash in `tev1_registry.cpp` (added by upstream commit 25f99e28d); that bug is not part of this change. Refs ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM FOLLOWING_AGENTS_PROTOCOL Following-Agents-Protocol: true AI-Assisted: true Assisted-by: opencode-go:deepseek-v4.1-flash [pi] --- .../ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md | 9 ++- .agents/specs/host-embedding.md | 12 ++- .../model_executor/models/host_embedding.cpp | 12 ++- tests/CMakeLists.txt | 4 + tests/vllm/models/test_host_embedding.cpp | 81 ++++++++++++++++++- 5 files changed, 111 insertions(+), 7 deletions(-) diff --git a/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md b/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md index 3d8ea0ac1..2729888ed 100644 --- a/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md +++ b/.agents/issues/ENG-HOST-EMBEDDING/ISSUE-LOCAL-01M3QETSTM8X8AJKBM9QCTGBKM.md @@ -16,4 +16,11 @@ A dense forward that owns a device embedding table pays [vocab, H] of device mem ## Resolution -- +- 2026-09-30: review repairs on the W1 PR (mudler/vllm.cpp#3356). The host arm's + async-id override now routes through the shared `detail::ApplyDeviceTokenIds` + body instead of its own unchecked Copy, so an override longer than the embed + input is refused rather than written past the `[T]` buffer; two cases in + `tests/vllm/models/test_host_embedding.cpp` pin the oversized refusal and the + shorter-prefix tail preservation. The test's global initializer uses the + portable `vllm_test::SetEnv` (`tests/support/test_env.h`) instead of POSIX + `::setenv`, which an unconditional target cannot compile under MSVC. diff --git a/.agents/specs/host-embedding.md b/.agents/specs/host-embedding.md index 114cfad4f..e51b65b32 100644 --- a/.agents/specs/host-embedding.md +++ b/.agents/specs/host-embedding.md @@ -49,7 +49,7 @@ owned both arms. |---|---| | The seam and the flag | `include/vllm/model_executor/models/host_embedding.h`, `src/vllm/model_executor/models/host_embedding.cpp` | | The wake-up call sites | the Qwen3.5 family, MuseGlimmer, the shared Qwen3 dense driver, and the classic dense families (Gemma 1-4, GLM4, Granite, MiniCPM 1/3, OLMo2, OPT, Phi, Phi3, StableLM, Command-R, DeepSeek-V2, GLM-MoE-DSA, Dots3-Note, Nemotron-H, Voxtral) | -| The async id override consumed before the gather | `src/vllm/model_executor/models/qwen3_5_internal.h` (`detail::TakeDeviceTokenIds`) | +| The async id override consumed before the gather | `src/vllm/model_executor/models/qwen3_5_internal.h` (`detail::TakeDeviceTokenIds`, then `detail::ApplyDeviceTokenIds` for the splice) | | The documented knobs | `docs/ENVIRONMENT.md` (`VT_HOST_EMBEDDING`, `VT_HOST_EMBED_TRACE`) | | Build registration | the new TU in `CMakeLists.txt` | @@ -76,7 +76,11 @@ forward makes when it owns a device table. It owns both arms: staging buffer is copied to `out`. - The async runner may have spliced this step's sampled token into a device-resident ids buffer, so the host arm consumes - `detail::TakeDeviceTokenIds()` and reads the ids back before the gather. + `detail::TakeDeviceTokenIds()` and reads the ids back before the gather. The + splice runs through the same `detail::ApplyDeviceTokenIds` body the device + arm uses, so an override longer than the embed input is refused with the + caller's name instead of writing past the `[T]` buffer, and a SHORTER + override replaces exactly its prefix while the padded host tail stays. - `VT_HOST_EMBED_TRACE=1` prints the arm and `T`. - The flag is read ONCE per process (a function-static), so a serving process cannot switch arms mid-run; a same-binary A/B is two processes. @@ -115,6 +119,7 @@ forward makes when it owns a device table. It owns both arms: | The flag is process-static, so a test binary cannot exercise both arms | The test binary enables it before `main` and is flag-ON by construction; the off arm is a plain early return and every other model suite runs it | | `d_dev == nullptr` does not prove the arm ran on the CPU backend | `ResidentWeight` aliases host bytes when `is_cpu()`, so the test captures the one-shot `[host-embed]` banner and checks `HostEmbedInto`'s return value; the mutation that disables the arm makes 4 assertions fail | | The host table's bytes are released by another path | The host arm declines on `bytes.empty()` / `host_released`, and `EmbedGather` then takes the device arm; the decline is a test case | +| The runner and the model disagree about this step's row count | The override's count is bounded against the embed input by `detail::ApplyDeviceTokenIds` on BOTH arms; a longer override throws with `what` naming the caller, a shorter one is the padded case and keeps the host upload's tail. Two test cases pin the boundary | ## Evidence @@ -122,6 +127,9 @@ forward makes when it owns a device table. It owns both arms: device op is gated on; the arm is proven to have run (captured banner + return value), not merely to have produced right numbers. - The table is never uploaded on the host arm (`d_dev == nullptr`). +- The async override's shape is bounded: a longer-than-`T` override throws, and a + shorter one splices its prefix over the host upload while preserving the tail + (`tests/vllm/models/test_host_embedding.cpp`, the two override cases). - `tests/vt/test_ops_embedding_quant` (6/6, 1637 assertions) is unchanged. ## Gates diff --git a/src/vllm/model_executor/models/host_embedding.cpp b/src/vllm/model_executor/models/host_embedding.cpp index 0bd1306d5..d9aa8be77 100644 --- a/src/vllm/model_executor/models/host_embedding.cpp +++ b/src/vllm/model_executor/models/host_embedding.cpp @@ -50,9 +50,17 @@ bool HostEmbedInto(Dev d, DBuf& hidden, const std::vector& token_ids, std::vector resolved = token_ids; const detail::DeviceTokenIds override_ids = detail::TakeDeviceTokenIds(); if (override_ids.ids != nullptr) { + // The device arm below splices through `detail::ApplyDeviceTokenIds`, which + // bounds the override against the embed input and enqueues the Copy on the + // queue. The host arm must not keep a second, unchecked copy of that rule: + // a count larger than `T` means the runner and the model disagree about + // this step, and the raw Copy that used to sit here wrote past the T-row + // buffer instead of refusing. The helper also tolerates a SHORTER override, + // which is the padded case: the first `count` rows are replaced and the + // host upload's tail is preserved. DBuf dids(d, DType::kI32, {T}, token_ids.data()); - d.b.Copy(d.q, dids.ptr(), override_ids.ids, - static_cast(override_ids.count) * sizeof(int32_t)); + detail::ApplyDeviceTokenIds(d.b, d.q, dids.ptr(), T, override_ids, + "host embedding"); dids.Download(d, resolved.data()); } if (std::getenv("VT_HOST_EMBED_TRACE") != nullptr) { diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 7e9b53b4a..22f6c59eb 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -4102,6 +4102,10 @@ vllm_cpp_add_test(test_qwen3_5_dense_load_residency # is the half a gather that uploaded first would pass. The flag is process-static, # so the binary enables it before the first gather and is flag-ON by construction. vllm_cpp_add_test(test_host_embedding vllm/models/test_host_embedding.cpp) +# The override cases publish through `detail::DeviceTokenIdsScope`, which lives +# in the not-installed `src/vllm/model_executor/models/qwen3_5_internal.h` beside +# the other `detail::` seams; same reach as test_expert_stream_steps. +target_include_directories(test_host_embedding PRIVATE ${CMAKE_SOURCE_DIR}/src) # QUANT-QWEN38-27B-NVFP4-ARM W4 (#821, .agents/specs/qwen38-27b-quant-arms.md): # the compressed-tensors `mixed-precision` accounting + scheme-resolution gate diff --git a/tests/vllm/models/test_host_embedding.cpp b/tests/vllm/models/test_host_embedding.cpp index 7f5c3c88d..bd9d701d4 100644 --- a/tests/vllm/models/test_host_embedding.cpp +++ b/tests/vllm/models/test_host_embedding.cpp @@ -36,22 +36,27 @@ #include "vllm/model_executor/models/dense_device_glue.h" // Dev, DBuf #include "vllm/model_executor/models/host_embedding.h" // HostEmbedInto, EmbedGather #include "vllm/model_executor/models/owned_bytes.h" +#include "vllm/model_executor/models/qwen3_5_internal.h" // detail::DeviceTokenIdsScope #include "vllm/model_executor/models/qwen3_5_weights.h" // OwnedTensor #include "vt/dtype.h" #include "vt/ops.h" +#include "support/test_env.h" #include "vt/iq4nl_q5_0_golden_vectors.h" namespace { struct EnableHostEmbedding { EnableHostEmbedding() { - ::setenv("VT_HOST_EMBEDDING", "1", 1); + // `support/test_env.h` is the portable setter (`_putenv_s` on MSVC): the + // unconditional POSIX `::setenv` this used to call does not compile in a + // CPU-only MSVC build, and this target is unconditional. + vllm_test::SetEnv("VT_HOST_EMBEDDING", "1"); // The arm prints ONE banner per process (`logged`), so the first case is the // one that can read it. It is what proves the arm RAN: the CPU device arm // aliases the host bytes and sets no `d_dev`, so `d_dev == nullptr` alone // cannot tell the two arms apart on this backend. - ::setenv("VT_HOST_EMBED_TRACE", "1", 1); + vllm_test::SetEnv("VT_HOST_EMBED_TRACE", "1"); } }; const EnableHostEmbedding g_enable_host_embedding; @@ -269,3 +274,75 @@ TEST_CASE("host embedding: a released host table declines the host arm") { // falls through to the device arm rather than gathering bytes that are gone. CHECK_FALSE(vllm::dense_attn::HostEmbedInto(d, out, ids, table, rows, k)); } + +TEST_CASE("host embedding: a device override splices its rows over the host upload") { + constexpr int64_t k = 4; + constexpr int64_t rows = 3; + // Row r is the constant bf16(r + 1), so a resolved id is readable from the + // output bytes alone, with no second gather to trust. + std::vector values(static_cast(rows * k)); + for (int64_t r = 0; r < rows; ++r) { + for (int64_t j = 0; j < k; ++j) { + values[static_cast(r * k + j)] = + vt::F32ToBF16(static_cast(r + 1)); + } + } + vllm::OwnedTensor table = + HostTable(vt::DType::kBF16, rows, k, values.data(), values.size() * 2); + + const std::vector host_ids{0, 1, 2}; + // ONE id for the first row: the shared splice replaces exactly `count` rows, + // so the output must read [1, 1, 2]. The tail keeps the host upload's rows, + // which is the padded case a "replace everything" copy would break. + const std::vector override_ids{1}; + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kBF16, {3, k}); + + { + vllm::detail::DeviceTokenIdsScope scope( + override_ids.data(), static_cast(override_ids.size())); + CHECK(vllm::dense_attn::HostEmbedInto(d, out, host_ids, table, rows, k)); + } + + const auto* got = static_cast(out.ptr()); + const int32_t want[3] = {1, 1, 2}; + for (int64_t t = 0; t < 3; ++t) { + for (int64_t j = 0; j < k; ++j) { + CAPTURE(t); + CAPTURE(j); + CHECK(got[static_cast(t) * static_cast(k) + + static_cast(j)] == + values[static_cast(want[static_cast(t)]) * + static_cast(k) + + static_cast(j)]); + } + } +} + +TEST_CASE("host embedding: an override longer than the embed input is refused") { + constexpr int64_t k = 4; + constexpr int64_t rows = 2; + const std::vector values(static_cast(rows * k), + vt::F32ToBF16(1.0F)); + vllm::OwnedTensor table = + HostTable(vt::DType::kBF16, rows, k, values.data(), values.size() * 2); + + const std::vector host_ids{0}; // T = 1 row + const std::vector override_ids{1, 0}; // 2 > T + + vt::Queue q = vt::GetBackend(vt::DeviceType::kCPU).CreateQueue(); + vllm::dense_attn::Dev d{vt::GetBackend(q.device.type), q}; + vllm::dense_attn::DBuf out(d, vt::DType::kBF16, {1, k}); + + vllm::detail::DeviceTokenIdsScope scope( + override_ids.data(), static_cast(override_ids.size())); + // The host upload is one row; splicing two would write past it. The device + // arm's shared `ApplyDeviceTokenIds` bounds the override and throws with the + // caller's name, and the host arm now routes through that same check instead + // of issuing its own unchecked Copy. + CHECK_THROWS_AS( + vllm::dense_attn::HostEmbedInto(d, out, host_ids, table, rows, k), + std::runtime_error); +}