From 9861c110bde6bd2ebcebeeb96b078c00ba4e323b Mon Sep 17 00:00:00 2001 From: dev Date: Tue, 29 Sep 2026 20:33:28 +0000 Subject: [PATCH 1/2] spec(LOAD-MODELOPT-NVFP4-BORROW): the two-handle NVFP4 device resident leak The fp4 resident publishes one device block on two owning handles, and every Marlin repack builder released only one of them, so each repacked weight kept a full packed+scale device copy alive for the process lifetime. Record the defect, the single-owner remedy, the red-first test that proves the release frees, and the checker clause the new contract needs. FOLLOWING_AGENTS_PROTOCOL Following-Agents-Protocol: true AI-Assisted: true Assisted-by: AGENT:opencode-go/deepseek-v4.1-flash [pi] --- .../ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR.md | 19 +++ .agents/specs/nvfp4-one-owner-per-resident.md | 115 ++++++++++++++++++ 2 files changed, 134 insertions(+) create mode 100644 .agents/issues/LOAD-MODELOPT-NVFP4-BORROW/ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR.md create mode 100644 .agents/specs/nvfp4-one-owner-per-resident.md diff --git a/.agents/issues/LOAD-MODELOPT-NVFP4-BORROW/ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR.md b/.agents/issues/LOAD-MODELOPT-NVFP4-BORROW/ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR.md new file mode 100644 index 000000000..45dab13a2 --- /dev/null +++ b/.agents/issues/LOAD-MODELOPT-NVFP4-BORROW/ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR.md @@ -0,0 +1,19 @@ +ID: ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR +Title: Every Marlin NVFP4 repack leaks its packed and scale device buffers, because Nvfp4Weight publishes one allocation on two owning handles +Row: LOAD-MODELOPT-NVFP4-BORROW +State: OPEN +Kind: bug +GitHub: - +Mirror: PENDING +Availability: FULL +Created: 2026-09-29 +Updated: 2026-09-29 +Closed: - + +## Problem + +Nvfp4Weight carries TWO owning handles per buffer: the type-specific d_packed/d_scale members (include/vllm/model_executor/models/qwen3_5_weights.h:707-708) and the generic raw-twin slots packed.d_dev/scale.d_dev (qwen3_5_weights.h:195) that AdoptDeviceBytesAsHost keys on (src/vllm/model_executor/models/qwen3_5_weights.cpp:420-429). ResidentNvfp4 publishes the SAME allocation on both: w.d_packed = shared_ptr(p, Free) followed by w.packed.d_dev = w.d_packed (dense_nvfp4_gemm.h:319-347, twin at src/vllm/model_executor/models/qwen3_5.cpp:1449). Releasing a repacked weight means dropping both handles, but every Marlin repack builder drops only the type-specific pair: dense_nvfp4_gemm.h:448-449, dense_nvfp4_gemm.h:667-670 (both operands of the pair), qwen3_5.cpp:2952-2953, qwen3_5.cpp:3132-3135, qwen3_5.cpp:6890-6891 and laguna.cpp:650-655. The surviving packed.d_dev alias holds the control block, so each repacked weight keeps a full packed+scale device copy alive for the process lifetime; on an MoE model that is one leak per expert per projection. The repack is precisely the step that is supposed to keep peak weight memory flat (dense_nvfp4_gemm.h:395-397: "we do this repack lazily on first forward and then FREE the fp4 originals"), so the leak silently undoes the memory design on exactly the cards that need it: the NVFP4 arms are served on 24 GiB consumer Blackwell (sm_120a), where a second full copy of the fp4 originals does not fit beside the repacked weights. Cause is grounded in the source above, not inferred from a failed allocation. + +## Resolution + +- diff --git a/.agents/specs/nvfp4-one-owner-per-resident.md b/.agents/specs/nvfp4-one-owner-per-resident.md new file mode 100644 index 000000000..48084f074 --- /dev/null +++ b/.agents/specs/nvfp4-one-owner-per-resident.md @@ -0,0 +1,115 @@ +# `Nvfp4Weight` publishes one device allocation on two owning handles — ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR + +`ResidentNvfp4` hands the same device block to `Nvfp4Weight::d_packed` and to +`OwnedTensor::d_dev`. Every Marlin repack builder released only the first, so the +second kept a full packed+scale device copy alive for the process lifetime — once +per repacked weight, and once per expert per projection on an MoE model. The +repack is the step that is supposed to keep peak weight memory flat, so the +surviving alias silently undoes the memory design on exactly the cards that need +it: the NVFP4 arms serve on 24 GiB consumer Blackwell (`sm_120a`), where a second +full copy of the fp4 originals does not fit beside the repacked weights. + +Issue: [ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR](../issues/LOAD-MODELOPT-NVFP4-BORROW/ISSUE-LOCAL-01M3QDT6A8JEM8WNJTPR16MXHR.md). +Owning row: `LOAD-MODELOPT-NVFP4-BORROW` ([engine-matrix.md](../engine-matrix.md)), +the row that owns `ResidentNvfp4`'s upload and adoption path and the spec +[`load-modelopt-nvfp4-borrow.md`](load-modelopt-nvfp4-borrow.md). + +## The defect, grounded + +| Where (line anchors at this branch's base, `b45a94273`) | What | +|---|---| +| `include/vllm/model_executor/models/qwen3_5_weights.h:707-708` | `Nvfp4Weight::d_packed`/`d_scale`, documented as "the shared_ptr deleter frees through the vt Backend". | +| `.../qwen3_5_weights.h:195` | `OwnedTensor::d_dev`, the generic raw-twin slot `AdoptDeviceBytesAsHost` keys on (`src/vllm/model_executor/models/qwen3_5_weights.cpp:420-429`). | +| `.../dense_nvfp4_gemm.h:319-347` | `ResidentNvfp4` sets `w.d_packed = shared_ptr(p, Free)` and then `w.packed.d_dev = w.d_packed`: **two owning handles, one control block.** | +| `src/vllm/model_executor/models/qwen3_5.cpp:1449-1476` | The private twin of the same function, the same two-handle shape. | +| `.../dense_nvfp4_gemm.h:448-449`, `:667-670` | `BuildMarlinDenseResident` and `BuildMarlinDensePairResident` release the type-specific pair only. | +| `src/vllm/model_executor/models/qwen3_5.cpp:2952-2953`, `:3132-3135`, `:6890-6891` | The 27B dense and MoE repack builders, the same release shape. | +| `src/vllm/model_executor/models/laguna.cpp:650-655` | Six resets per expert, the same shape. | + +The alias is not an accident of one builder: it is the design (`packed.d_dev` +is what `AdoptDeviceBytesAsHost` can act on), and the fix is to make it the ONLY +owner rather than to teach six call sites to drop two handles. + +## The design call, and what it costs + +Alternatives rejected: + +1. **Keep both handles and reset both at every caller.** Correct at each site and + wrong for the codebase: six release sites, one of which (a new builder) will + forget. The bug is the duplicated ownership, not the six callers. +2. **Make `d_dev` observe `d_packed` without owning it** (a raw pointer or a + `weak_ptr`). `AdoptDeviceBytesAsHost` needs the block alive while the adopted + host view exists, and the residency machinery keys on `d_dev` everywhere else; + splitting the lifetime rules between two slots reintroduces the same class of + mistake. +3. **Have `ReleaseResident` free through the deleter directly.** It cannot: the + adopted `bytes` view holds a copy of the shared_ptr on a host-addressable + device, so the block outlives the weight's handle by construction. The view + must be dropped first, which is what `ReleaseResident` does. + +## Design + +One owner per allocation: `packed.d_dev` / `scale.d_dev`. The type-specific +`d_packed`/`d_scale` members are deleted, so a reader that still expects them +fails to **compile** — the migration cannot be half-done. + +Two additions carry the release: + +- `OwnedTensor::HostViewIsDeviceTwin()` — true when `d_dev` is set, `bytes` is an + adopted (borrowed, `mmap_src == nullptr`) view, and it points at that same + block. This is the predicate for "the host view is an alias of this tensor's + own device allocation", i.e. it holds the block alive. +- `Nvfp4Weight::ReleaseResident()` — drops an adopted twin view first, then + resets `d_dev` on both tensors. Called at all six release sites. + +`ResidentNvfp4` (both copies) uploads into `packed.d_dev`/`scale.d_dev` +directly. The old upload guard `!w.d_packed` becomes `!w.packed.d_dev`; that is +equivalent because the upload is the only writer of either slot for an +`Nvfp4Weight`'s packed/scale tensors (`grep -rn 'd_dev = ' src include`). + +**Behavior is unchanged: same bytes, same adoption point, same release points.** +The only new effect is that the adopted alias view is dropped before the block is +freed, which is required for the free to happen at all. + +## What this does NOT claim + +No end-to-end GiB number is measured here. The unit gate proves that a release +through the weight's own resident state frees both buffers; it does not measure +the process RSS of a real repacking serve on a real checkpoint. That measurement +needs a model and a GPU and is the natural follow-up, not a precondition for the +correctness fix. + +## Tests + +`tests/vllm/test_load_direct_upload.cpp`, which already owns the `ResidentNvfp4` +accounting/adoption cases over the `FakeBackend` + `ObservableMapping` harness: + +- New: `fp4 resident: releasing the resident frees both device buffers`. Uploads + a borrowed fp4 weight on a non-host-addressable fake device (no adoption, so + the alias is the only extra reference), then performs the release and asserts + `b.frees == 2`. + - **RED pre-fix**, captured: + `CHECK( b.frees == 2 )` → `values: CHECK( 0 == 2 )`, because the release + reset only the type-specific pair. + - **GREEN post-fix**: the release is `w.ReleaseResident()`. +- The three existing cases are updated to the single-owner shape + (`w.packed.d_dev` / `w.scale.d_dev` in place of `w.d_packed` / `w.d_scale`), + and their final `CHECK(b.frees == 2)` (the weight out of scope) is unchanged. + +What the gate does not prove: that a real CUDA driver frees at that instant. The +`FakeBackend` counts `Free` calls; it cannot observe the driver. + +## Gates + +- `cmake --build build --target test_load_direct_upload && ./build/tests/test_load_direct_upload` + — RED before, GREEN after, all 17 cases. +- The full `ctest --test-dir build` on the affected tree. +- `scripts/agent-preflight.sh` (the record gates for the changed files). +- The build itself is the completeness check for the deleted members. + +## Stop conditions + +- A caller of `d_packed`/`d_scale` outside the sites listed above: stop and + re-scope, because the deletion is no longer complete. +- A regression in the adoption cases (host-addressable device): stop, because + the release ordering is wrong. From db8f1f2d0bc1ccbb6e91573a3a3920a536bd25ca Mon Sep 17 00:00:00 2001 From: dev Date: Sun, 27 Sep 2026 20:27:32 +0000 Subject: [PATCH 2/2] fix(LOAD-MODELOPT-NVFP4-BORROW): one owner per device resident MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Nvfp4Weight published one allocation on two owning handles: its d_packed/d_scale and the packed/scale d_dev alias that AdoptDeviceBytesAsHost keys on. Releasing a resident meant dropping both, and every Marlin repack builder dropped only d_packed — so each repacked weight (and each expert on MoE models) leaked a full packed+scale device copy for the process lifetime. d_dev is already the generic raw-twin slot the residency machinery keys on, so the type-specific handles are redundant. ResidentNvfp4 uploads into packed.d_dev/scale.d_dev directly and ReleaseResident() is the single release; it drops an adopted host view first, since that view holds the same block alive on host-addressable devices. Behavior is unchanged: same bytes, same adoption, same release points. The new case in test_load_direct_upload pins the release the repack builders perform: RED before this change (CHECK(b.frees == 2) read 0, the alias kept the blocks), GREEN after; 17/17 cases, 206 assertions. scripts/check-fp4-resident-consistency.py pinned the old two-handle contract and would have failed on this tree. Its publication clause now requires a backend-owned allocation, and a new clause requires ReleaseResident() to reset BOTH d_dev slots, dropping an adopted twin view before each reset. Its mutation suite gains the release cases. FOLLOWING_AGENTS_PROTOCOL Following-Agents-Protocol: true AI-Assisted: true Assisted-by: AGENT:opencode-go/deepseek-v4.1-flash [pi] --- .../model_executor/models/dense_nvfp4_gemm.h | 30 ++- .../model_executor/models/qwen3_5_weights.h | 30 ++- scripts/check-fp4-resident-consistency.py | 193 ++++++++++++++---- src/vllm/model_executor/models/laguna.cpp | 9 +- .../models/nemotron_h_device.cpp | 2 +- src/vllm/model_executor/models/qwen3_5.cpp | 42 ++-- .../test_check_fp4_resident_consistency.py | 178 ++++++++++++---- tests/vllm/test_load_direct_upload.cpp | 79 +++++-- 8 files changed, 410 insertions(+), 153 deletions(-) diff --git a/include/vllm/model_executor/models/dense_nvfp4_gemm.h b/include/vllm/model_executor/models/dense_nvfp4_gemm.h index fb2c2fdf4..3c624dd9d 100644 --- a/include/vllm/model_executor/models/dense_nvfp4_gemm.h +++ b/include/vllm/model_executor/models/dense_nvfp4_gemm.h @@ -317,38 +317,35 @@ struct Nvfp4Dev { }; inline Nvfp4Dev ResidentNvfp4(Dev d, const Nvfp4Weight& w) { - if (!w.d_packed) { + if (!w.packed.d_dev) { const size_t pb = w.packed.bytes.size(); void* p = d.b.Alloc(pb); // ENG-LOAD-DIRECT-UPLOAD (issue #150). `LoadCtNvfp4W4A16`/`LoadCtMxfp4W4A16` // /`LoadCtNvfp4Raw` BORROW `packed` and `scale` from the safetensors mmap, // so this is the one host->device move of those bytes and it must be // accounted and followed by the same post-upload residency step every other - // qualifying weight gets. Publishing the allocation on the OwnedTensor is - // what lets `AdoptDeviceBytesAsHost` run at all (it keys on `d_dev`); the - // two handles share one control block, so the buffer is still freed exactly - // once, through the vt Backend. + // qualifying weight gets. Publishing on `d_dev` is what lets + // `AdoptDeviceBytesAsHost` run (it keys on that slot), and it is the only + // owner (Nvfp4Weight::ReleaseResident). vllm::load_stats::AddDeviceUpload(pb); d.b.Copy(d.q, p, w.packed.bytes.data(), pb); Backend* bk = &d.b; - w.d_packed = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); - w.packed.d_dev = w.d_packed; + w.packed.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); AdoptDeviceBytesAsHost(d.b, w.packed); } - if (!w.d_scale) { + if (!w.scale.d_dev) { const size_t sb = w.scale.bytes.size(); void* p = d.b.Alloc(sb); vllm::load_stats::AddDeviceUpload(sb); d.b.Copy(d.q, p, w.scale.bytes.data(), sb); Backend* bk = &d.b; - w.d_scale = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); - w.scale.d_dev = w.d_scale; + w.scale.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); AdoptDeviceBytesAsHost(d.b, w.scale); } Nvfp4Dev r; - r.packed = MakeTensor(w.d_packed.get(), DType::kI8, d.q.device, {w.n, w.k / 2}); + r.packed = MakeTensor(w.packed.d_dev.get(), DType::kI8, d.q.device, {w.n, w.k / 2}); // Scale grid is [N, K/group_size]: K/16 for NVFP4, K/32 for MXFP4. - r.scale = MakeTensor(w.d_scale.get(), DType::kI8, d.q.device, {w.n, w.k / w.group_size}); + r.scale = MakeTensor(w.scale.d_dev.get(), DType::kI8, d.q.device, {w.n, w.k / w.group_size}); return r; } @@ -445,8 +442,7 @@ inline void BuildMarlinDenseResident(Dev d, const Nvfp4Weight& w, d.b.Copy(d.q, mr.g, &g, sizeof(float)); } d.b.Synchronize(d.q); // repack done -> safe to free the fp4 originals - w.d_packed.reset(); - w.d_scale.reset(); + w.ReleaseResident(); mr.ready = true; } @@ -664,10 +660,8 @@ inline void BuildMarlinDensePairResident(Dev d, const Nvfp4Weight& gw, d.b.Synchronize(d.q); // repack done -> safe to free staging + fp4 originals d.b.Free(tmp_w); d.b.Free(tmp_s); - gw.d_packed.reset(); - gw.d_scale.reset(); - uw.d_packed.reset(); - uw.d_scale.reset(); + gw.ReleaseResident(); + uw.ReleaseResident(); mr.ready = true; } diff --git a/include/vllm/model_executor/models/qwen3_5_weights.h b/include/vllm/model_executor/models/qwen3_5_weights.h index 34e2b13c2..e67156aa0 100644 --- a/include/vllm/model_executor/models/qwen3_5_weights.h +++ b/include/vllm/model_executor/models/qwen3_5_weights.h @@ -179,13 +179,21 @@ struct OwnedTensor { // `residency_policy().release_host_weights_after_upload` // (platforms/interface.h; BACKEND-PLATFORM item 2). Logically const: the // tensor's VALUE is unchanged, only the now-dead host mirror is freed. This - // mirrors the existing mutable lazy-device-upload residency design (d_dev/ - // d_packed above are populated on a const weight). swap-with-empty (not + // mirrors the existing mutable lazy-device-upload residency design (d_dev + // above is populated on a const weight). swap-with-empty (not // clear()) guarantees the std::vector capacity is actually deallocated. After // release View()/bytes must not be read; shape/dtype metadata is retained and // Empty() remains false so device-resident dispatch continues to see it. void ReleaseHost() const; + // True when the host view is an adopted alias of this tensor's own device + // allocation (host-addressable devices: AdoptDeviceBytesAsHost re-points + // `bytes` at it and the block becomes the view's keep-alive). + bool HostViewIsDeviceTwin() const { + return d_dev != nullptr && bytes.borrowed() && mmap_src == nullptr && + bytes.data() == static_cast(d_dev.get()); + } + // Lazily-populated device-resident copies (CUDA forward only; null on host or // before first use). Uploaded ONCE and reused across every forward step so the // model's bf16/f32 weights (embed table, norms, attention/GDN projections, @@ -677,6 +685,16 @@ struct Nvfp4Weight { int64_t k = 0; // in_features (K % 16 == 0) bool Empty() const { return packed.Empty(); } + // Drop this weight's device resident. `packed.d_dev`/`scale.d_dev` are the + // only owners, so the resets free; an adopted host view (see + // HostViewIsDeviceTwin) must go first because it holds the same block alive. + void ReleaseResident() const { + if (packed.HostViewIsDeviceTwin()) packed.ReleaseHost(); + if (scale.HostViewIsDeviceTwin()) scale.ReleaseHost(); + packed.d_dev.reset(); + scale.d_dev.reset(); + } + // Block-scale FORMAT. Default = NVFP4: group_size 16, fp8-e4m3 `scale`, a // per-tensor `scale2` global. is_mxfp4 selects compressed-tensors MXFP4 // (`mxfp4-pack-quantized`): group_size 32, E8M0 (UE8M0) `scale` [N, K/32], NO @@ -702,13 +720,11 @@ struct Nvfp4Weight { // True when the activation-quant globals were loaded (27B true-W4A4 path). bool IsTrueW4A4() const { return alpha > 0.0F; } - // Lazily-populated device-resident copies (CUDA forward only; null on host or - // before first use). The shared_ptr deleter frees through the vt Backend. - mutable std::shared_ptr d_packed; - mutable std::shared_ptr d_scale; + // Device residency lives on `packed.d_dev`/`scale.d_dev` (the generic raw-twin + // slot): one owner per allocation, no second Nvfp4Weight handle. // Lazily-populated SWIZZLED weight block scale for the cutlass sm120a fp4 GEMM // path (VT_NVFP4_CUTLASS): [round_up(n,128), round_up(k/16,4)] in the cutlass - // atom layout, computed once from d_scale via vt::SwizzleBlockscale. + // `scale.d_dev` (the generic raw-twin slot) via vt::SwizzleBlockscale. mutable std::shared_ptr d_scale_sw; // vLLM/FlashInfer-compatible model-owned f32 alpha for the true-W4A4 CUTLASS // path. Uploaded once from the persistent `alpha` member; the diagnostic host diff --git a/scripts/check-fp4-resident-consistency.py b/scripts/check-fp4-resident-consistency.py index 1e9731ba6..ee8641a04 100755 --- a/scripts/check-fp4-resident-consistency.py +++ b/scripts/check-fp4-resident-consistency.py @@ -10,9 +10,14 @@ 1. `load_stats::AddDeviceUpload(nb)` — account the move. These bytes are counted NOWHERE else: they were borrowed on the way in, so the host-copy counter never saw them. - 2. `w..d_dev = w.d_` — PUBLISH the allocation on the - OwnedTensor. `AdoptDeviceBytesAsHost` keys on `d_dev` and returns immediately - when it is null, so without this the next statement is a silent no-op. + 2. `w..d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); })` + — PUBLISH the allocation on the OwnedTensor. `AdoptDeviceBytesAsHost` keys on + `d_dev` and returns immediately when it is null, so without this the next + statement is a silent no-op. `d_dev` is ALSO the one owner of those bytes: + the type-specific `Nvfp4Weight::d_` fields this used to publish + alongside are gone, because a second owning handle is what leaked a full + packed+scale device copy per repacked weight (the repack builders reset only + one of the two). The release clause below pins that. 3. `AdoptDeviceBytesAsHost(d.b, w.)` — the post-upload residency step: release the consumed source pages, and on a host-addressable device (GB10 unified memory under Vulkan) adopt the device allocation as the host view so @@ -44,8 +49,10 @@ block to `w..bytes.size()` (so `AddDeviceUpload(0)`, or the other buffer's byte count, is not a pass); (b) COPIED — the upload reads `w..bytes.data()` (not the other buffer's); - (c) PUBLISHED — `w..d_dev = >` (so `= nullptr` and - a foreign handle are not a pass); + (c) PUBLISHED — `w..d_dev = `: a `std::shared_ptr` + construction whose deleter returns the block through the vt + Backend (`...Free(...)`). `= nullptr`, `= {}`, and a bare + `= ..d_dev` (which owns nothing) are not a pass; (d) ADOPTED — `AdoptDeviceBytesAsHost(..., w.)` on THIS function's weight parameter (so adopting another object's buffer is not a pass); (e) ORDERED — (c) precedes (d), which is what makes (d) more than a no-op; @@ -57,6 +64,17 @@ `Copy` — lifts it above (c) as well and was already caught by (e), so only the sink shape shows what (f) adds. + (g) RELEASED — `Nvfp4Weight::ReleaseResident()` resets `packed.d_dev` AND + `scale.d_dev`, and where it drops an adopted host twin for a + buffer that drop precedes that buffer's reset. THIS IS THE + CLAUSE THE LEAK VIOLATED: the repack builders reset only the + type-specific handle, so the surviving `d_dev` alias kept the + block for the process lifetime. It is checked in + include/vllm/model_executor/models/qwen3_5_weights.h, the one + definition both upload copies release through, because the + release is not inside the upload function the clauses above + scope to. + TEXT THE COMPILER NEVER SEES IS NOT A PASS. Every clause runs against `checker_text.normalize_source`, which blanks `//` and `/* */` comments, `#if 0` / `#if false` regions and `if (false)` / `if (0)` branches before matching, keeping @@ -68,9 +86,10 @@ a real build configuration, not a disguised deletion. WHAT THIS GATE DOES *NOT* DO, stated plainly so the record does not imply more. It is -a STRUCTURAL check over text: it proves the six statements are present in code the -compiler keeps, bound to the right buffer of the right object, and that the adoption -follows both the copy and the publication. It cannot prove they are CORRECT at run +a STRUCTURAL check over text: it proves the per-buffer statements and the single +release are present in code the compiler keeps, bound to the right buffer of the right +object, and that the adoption follows both the copy and the publication. It cannot +prove they are CORRECT at run time — that the pointer published is the one that was uploaded, that the byte count matches the allocation, or that the copy transferred the right bytes. Nor does it model the preprocessor: a statement moved under a build-configuration `#ifdef` is @@ -109,6 +128,13 @@ Path("src/vllm/model_executor/models/qwen3_5.cpp"), ) +# The ONE release of an Nvfp4Weight's device resident. Both upload copies above are +# released through it, so a release that forgets a buffer leaks from BOTH. +RELEASE_SOURCE = Path("include/vllm/model_executor/models/qwen3_5_weights.h") + +# The one release function both upload copies are released through. +RELEASE_FUNCTION = "ReleaseResident" + # The buffers an Nvfp4Weight uploads. Each needs the full step, in its OWN block. BUFFERS = ("packed", "scale") @@ -134,15 +160,30 @@ _NULLISH = ("nullptr", "NULL", "{}", "0") -def _device_handle(buffer: str) -> str: - """`packed` -> `d_packed`: the Nvfp4Weight's own handle for that buffer.""" - return "d_" + buffer - - def _line_no(text: str, pos: int) -> int: return text.count("\n", 0, pos) + 1 +def _statement_end(text: str, start: int) -> int: + """Index of the `;` that ends the statement beginning at `start`, tracking + bracket depth so a semicolon INSIDE a lambda body does not truncate it. + + The published right-hand side is now `std::shared_ptr(p, [bk](void* q) + { bk->Free(q); })`, whose inner brace holds a `;`. A flat `[^;]+;` capture — + what this checker used while the rhs was a bare member name — stops there and + sees no deleter, so the check would have failed on the real, correct shape.""" + depth = 0 + for i in range(start, len(text)): + c = text[i] + if c in "([{": + depth += 1 + elif c in ")]}": + depth -= 1 + elif c == ";" and depth == 0: + return i + return -1 + + def function_bodies(text: str) -> list[tuple[int, str, str]]: """Every `ResidentNvfp4(Dev, const Nvfp4Weight&)` as (line_no, recv, body), the body delimited by brace matching from the definition's opening brace.""" @@ -163,16 +204,19 @@ def checked_bodies(text: str) -> list[tuple[int, str, str]]: def buffer_block(body: str, recv: str, buffer: str) -> str | None: - """The `if (!.d_) { ... }` upload block for ONE buffer, or None. + """The `if (!..d_dev) { ... }` upload block for ONE buffer, or + None. This is the scope every clause is checked in. Checking the whole body instead - was the original bug: a single surviving statement served both buffers.""" + was the original bug: a single surviving statement served both buffers. The + guard is the single-owner slot rather than the deleted `d_` handle: + the upload is skipped exactly when the device resident already exists.""" guard = re.compile( r"if\s*\(\s*!\s*" + re.escape(recv) + r"\s*\.\s*" - + re.escape(_device_handle(buffer)) - + r"\s*\)\s*\{" + + re.escape(buffer) + + r"\s*\.\s*d_dev\s*\)\s*\{" ) m = guard.search(body) if m is None: @@ -211,19 +255,28 @@ def _copied(block: str, recv: str, buffer: str) -> int | None: def _published(block: str, recv: str, buffer: str) -> int | None: - """Offset of `..d_dev = .d_>`, or None.""" + """Offset of `..d_dev = `, or None. + + The published value must OWN the block: a `std::shared_ptr` construction whose + deleter returns it through the vt Backend. A bare `= nullptr`, `= {}` or an + alias of some other object's `d_dev` publishes no owner, so + `AdoptDeviceBytesAsHost` either returns immediately or the block outlives the + weight with nobody to free it.""" for m in re.finditer( re.escape(recv) + r"\s*\.\s*" + re.escape(buffer) - + r"\s*\.\s*d_dev\s*=\s*(?P[^;]+);", + + r"\s*\.\s*d_dev\s*=\s*", block, ): - rhs = m.group("rhs").strip() + end = _statement_end(block, m.end()) + if end < 0: + continue + rhs = block[m.end() : end].strip() if rhs in _NULLISH: continue - if re.search( - re.escape(recv) + r"\s*\.\s*" + re.escape(_device_handle(buffer)) + r"\b", rhs + if re.search(r"\bshared_ptr\s*<\s*void\s*>\s*\(", rhs) and re.search( + r"\bFree\s*\(", rhs ): return m.start() return None @@ -250,7 +303,7 @@ def body_violations(body: str, recv: str = "w") -> list[str]: block = buffer_block(body, recv, buffer) if block is None: problems.append( - f"no `if (!{recv}.{_device_handle(buffer)})` upload block for `{buffer}`. " + f"no `if (!{recv}.{buffer}.d_dev)` upload block for `{buffer}`. " f"Each buffer's post-upload step is checked inside its OWN block, " f"because body-wide matching lets one buffer's statements satisfy the " f"other's. If this copy was deliberately restructured, update " @@ -276,9 +329,10 @@ def body_violations(body: str, recv: str = "w") -> list[str]: pub = _published(block, recv, buffer) if pub is None: problems.append( - f"`{buffer}` does not PUBLISH its device allocation on the OwnedTensor " - f"(`{recv}.{buffer}.d_dev = {recv}.{_device_handle(buffer)}`); " - f"AdoptDeviceBytesAsHost keys on `d_dev` and would return immediately" + f"`{buffer}` does not PUBLISH a backend-owned device allocation on " + f"`{recv}.{buffer}.d_dev` (a std::shared_ptr whose deleter returns " + f"the block through the vt Backend); AdoptDeviceBytesAsHost keys " + f"on `d_dev` and would return immediately" ) adopt = _adopted(block, recv, buffer) if adopt is None: @@ -322,6 +376,64 @@ def file_violations(text: str, label: str) -> list[str]: ] +_RELEASE_DEF = re.compile( + r"\bvoid\s+" + re.escape(RELEASE_FUNCTION) + r"\s*\(\s*\)\s*const\s*\{" +) + + +def release_body(text: str) -> tuple[int, str] | None: + """The `ReleaseResident() const` body as (line_no, normalized body), or None.""" + norm = normalize_source(text) + m = _RELEASE_DEF.search(norm) + if m is None: + return None + end = _match_braces(norm, m.end()) + return _line_no(text, m.start()), norm[m.end() : end - 1] + + +def release_violations(text: str, label: str) -> list[str]: + """Every way `ReleaseResident` leaves a device resident un-owned, and the one + ordering that makes the release a no-op on a host-addressable device.""" + found = release_body(text) + if found is None: + return [ + f"{label}: no `void {RELEASE_FUNCTION}() const` definition. It is the ONE " + f"release of an Nvfp4Weight's device resident (both upload copies above " + f"use it); if it was deliberately renamed or moved, update RELEASE_SOURCE " + f"in scripts/check-fp4-resident-consistency.py in the same change." + ] + line_no, body = found + problems: list[str] = [] + for buffer in BUFFERS: + reset = re.search( + r"\b" + re.escape(buffer) + r"\s*\.\s*d_dev\s*\.\s*reset\s*\(\s*\)", body + ) + if reset is None: + problems.append( + f"`{buffer}.d_dev` is never reset: after a Marlin repack dropped the " + f"fp4 originals, this block stays owned for the process lifetime — " + f"one leaked packed+scale device copy per repacked weight and per " + f"expert. ReleaseResident must reset BOTH buffers" + ) + continue + twin = re.search(r"\b" + re.escape(buffer) + r"\s*\.\s*ReleaseHost\s*\(\s*\)", body) + if twin is None: + problems.append( + f"`{buffer}`'s adopted host twin is never dropped: on a " + f"host-addressable device that view holds the same block alive, so " + f"the reset below frees nothing and the release leaks exactly as it " + f"did when the type-specific handle was the only one reset" + ) + elif twin.start() > reset.start(): + problems.append( + f"`{buffer}`'s adopted host twin is dropped AFTER its `d_dev` " + f"reset: the view holds the block alive on a host-addressable " + f"device, so the reset frees nothing and the later drop touches a " + f"view whose owner is gone" + ) + return [f"{label}:{line_no}: {p}" for p in problems] + + def main() -> int: violations: list[str] = [] checked = 0 @@ -334,21 +446,32 @@ def main() -> int: checked += len(checked_bodies(text)) violations.extend(file_violations(text, str(rel))) + rel = str(RELEASE_SOURCE) + path = ROOT / RELEASE_SOURCE + if not path.exists(): + print(f"ERROR: {rel} not found", file=sys.stderr) + return 1 + violations.extend( + release_violations(path.read_text(encoding="utf-8", errors="ignore"), rel) + ) + if violations: print( - "ERROR: an fp4 resident upload drops part of the ENG-LOAD-DIRECT-UPLOAD " - "post-upload step (issue #150):", + "ERROR: an fp4 resident upload or its single release drops part of the " + "ENG-LOAD-DIRECT-UPLOAD post-upload step (issue #150):", file=sys.stderr, ) for v in violations: print(f" - {v}", file=sys.stderr) print( - "Every ResidentNvfp4 must COUNT its upload, COPY from that buffer, PUBLISH " - "the allocation on the OwnedTensor, and then run AdoptDeviceBytesAsHost — " - "for BOTH `packed` and `scale`, each inside its own `if (!w.d_)` " - "block. The shared copy is pinned at run time by " - "tests/vllm/test_load_direct_upload.cpp; this gate exists because the " - "qwen3_5.cpp duplicate sits in an anonymous namespace no test can reach.", + "Every ResidentNvfp4 must COUNT its upload, COPY from that buffer, " + "PUBLISH the allocation on the OwnedTensor, and then run " + "AdoptDeviceBytesAsHost — for BOTH `packed` and `scale`, each inside its " + "own `if (!w..d_dev)` block — and `ReleaseResident()` must reset BOTH " + "`d_dev` slots, dropping an adopted twin view first. The shared copy is " + "pinned at run time by tests/vllm/test_load_direct_upload.cpp; this gate " + "exists because the qwen3_5.cpp duplicate sits in an anonymous namespace " + "no test can reach.", file=sys.stderr, ) return 1 @@ -356,7 +479,7 @@ def main() -> int: print( f"OK: {checked} ResidentNvfp4 definition(s) across {len(SOURCES)} file(s) count " f"the upload, copy, publish d_dev, and adopt — per buffer, for both packed and " - f"scale." + f"scale — and ReleaseResident() resets both d_dev slots." ) return 0 diff --git a/src/vllm/model_executor/models/laguna.cpp b/src/vllm/model_executor/models/laguna.cpp index c196ba42d..db54252d1 100644 --- a/src/vllm/model_executor/models/laguna.cpp +++ b/src/vllm/model_executor/models/laguna.cpp @@ -647,12 +647,9 @@ void BuildLagunaMoeMarlinResident(vllm::dense_nvfp4::Dev d, const LagunaMoeWeigh // GEMV/CPU paths that read these bytes can never run in this process. for (int e = 0; e < E; ++e) { const size_t se = static_cast(e); - moe.experts_gate_fp4[se].d_packed.reset(); - moe.experts_gate_fp4[se].d_scale.reset(); - moe.experts_up_fp4[se].d_packed.reset(); - moe.experts_up_fp4[se].d_scale.reset(); - moe.experts_down_fp4[se].d_packed.reset(); - moe.experts_down_fp4[se].d_scale.reset(); + moe.experts_gate_fp4[se].ReleaseResident(); + moe.experts_up_fp4[se].ReleaseResident(); + moe.experts_down_fp4[se].ReleaseResident(); moe.experts_gate_fp4[se].packed.ReleaseHost(); moe.experts_gate_fp4[se].scale.ReleaseHost(); moe.experts_up_fp4[se].packed.ReleaseHost(); diff --git a/src/vllm/model_executor/models/nemotron_h_device.cpp b/src/vllm/model_executor/models/nemotron_h_device.cpp index 54f26ea0e..04027c0f4 100644 --- a/src/vllm/model_executor/models/nemotron_h_device.cpp +++ b/src/vllm/model_executor/models/nemotron_h_device.cpp @@ -1019,7 +1019,7 @@ DBuf DeviceLmHeadD(Dev d, const NemotronHHostWeights& host, // value-based gate: that arm computes the SAME logits, while re-uploading // 198.18 MB — the whole [131072, 2688] packed operand plus its group scales — // on EVERY call, because `LmHeadNvfp4View` hands out a stack temporary and - // `ResidentNvfp4` caches on `w.d_packed`, a member of the weight it was given. + // `ResidentNvfp4` caches on `w.packed.d_dev`, a member of the weight it was given. // // The predicate and not the seam's `fallback_gemms` counter, deliberately. // `MutableW4A16Stats()` is a plain non-atomic static shared by every consumer diff --git a/src/vllm/model_executor/models/qwen3_5.cpp b/src/vllm/model_executor/models/qwen3_5.cpp index 0900a64a6..38f6c3d3f 100644 --- a/src/vllm/model_executor/models/qwen3_5.cpp +++ b/src/vllm/model_executor/models/qwen3_5.cpp @@ -1447,36 +1447,34 @@ struct Nvfp4Dev { // every forward step, so subsequent calls reuse the resident copy — no per-op // weight staging. CUDA path only; the deleter frees through the vt Backend. Nvfp4Dev ResidentNvfp4(Dev d, const Nvfp4Weight& w) { - if (!w.d_packed) { + if (!w.packed.d_dev) { const size_t pb = w.packed.bytes.size(); void* p = d.b.Alloc(pb); // ENG-LOAD-DIRECT-UPLOAD (issue #150): the 27B `LoadCtNvfp4Raw` weights // BORROW packed/scale from the safetensors mmap, so this is their one // host->device move. Account it and run the same post-upload residency step // every other qualifying weight gets, exactly as dense_nvfp4_gemm.h's - // shared ResidentNvfp4 does. Publishing the allocation on the OwnedTensor - // is what lets AdoptDeviceBytesAsHost run (it keys on `d_dev`); the two - // handles share one control block, so the buffer is freed exactly once. + // shared ResidentNvfp4 does. Publishing on `d_dev` is what lets + // `AdoptDeviceBytesAsHost` run (it keys on that slot), and it is the only + // owner (Nvfp4Weight::ReleaseResident). vllm::load_stats::AddDeviceUpload(pb); d.b.Copy(d.q, p, w.packed.bytes.data(), pb); Backend* bk = &d.b; - w.d_packed = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); - w.packed.d_dev = w.d_packed; + w.packed.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); AdoptDeviceBytesAsHost(d.b, w.packed); } - if (!w.d_scale) { + if (!w.scale.d_dev) { const size_t sb = w.scale.bytes.size(); void* p = d.b.Alloc(sb); vllm::load_stats::AddDeviceUpload(sb); d.b.Copy(d.q, p, w.scale.bytes.data(), sb); Backend* bk = &d.b; - w.d_scale = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); - w.scale.d_dev = w.d_scale; + w.scale.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); AdoptDeviceBytesAsHost(d.b, w.scale); } Nvfp4Dev r; - r.packed = MakeTensor(w.d_packed.get(), DType::kI8, d.q.device, {w.n, w.k / 2}); - r.scale = MakeTensor(w.d_scale.get(), DType::kI8, d.q.device, {w.n, w.k / 16}); + r.packed = MakeTensor(w.packed.d_dev.get(), DType::kI8, d.q.device, {w.n, w.k / 2}); + r.scale = MakeTensor(w.scale.d_dev.get(), DType::kI8, d.q.device, {w.n, w.k / 16}); return r; } @@ -1495,7 +1493,7 @@ Tensor ResidentNvfp4ScaleSwizzled(Dev d, const Nvfp4Weight& w) { auto round_up = [](int64_t x, int64_t y) { return (x + y - 1) / y * y; }; const int64_t Np = round_up(w.n, 128), Kp = round_up(w.k / 16, 4); if (!w.d_scale_sw) { - Nvfp4Dev dw = ResidentNvfp4(d, w); // ensures d_scale (linear device copy) + Nvfp4Dev dw = ResidentNvfp4(d, w); // ensures scale.d_dev (linear device copy) void* p = d.b.Alloc(static_cast(Np * Kp)); Backend* bk = &d.b; w.d_scale_sw = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); @@ -2316,7 +2314,7 @@ Tensor ResidentDeviceAlpha(Dev d, const float* host_alpha, Tensor ResidentNvfp4Alpha(Dev d, const Nvfp4Weight& w) { VT_CHECK(w.IsTrueW4A4(), "qwen3_5 NVFP4 device alpha: true-W4A4 required"); - VT_CHECK(w.d_packed && w.d_scale && w.d_scale_sw, + VT_CHECK(w.packed.d_dev && w.scale.d_dev && w.d_scale_sw, "qwen3_5 NVFP4 device alpha: incomplete weight resident state"); return ResidentDeviceAlpha(d, &w.alpha, w.d_alpha, "qwen3_5 NVFP4 device alpha: invalid scalar"); @@ -2949,8 +2947,7 @@ void BuildMarlinDenseResident(Dev d, const Nvfp4Weight& w, MarlinDenseResident& const float g = vt::cuda::MarlinNvfp4ProcessGlobalScale(w.scale2, sf); d.b.Copy(d.q, mr.g, &g, sizeof(float)); d.b.Synchronize(d.q); // repack done -> safe to free the fp4 originals - w.d_packed.reset(); - w.d_scale.reset(); + w.ReleaseResident(); mr.ready = true; } @@ -3129,10 +3126,8 @@ void BuildMarlinDensePairResident(Dev d, const Nvfp4Weight& gw, const Nvfp4Weigh d.b.Synchronize(d.q); // repack done -> safe to free staging + fp4 originals d.b.Free(tmp_w); d.b.Free(tmp_s); - gw.d_packed.reset(); - gw.d_scale.reset(); - uw.d_packed.reset(); - uw.d_scale.reset(); + gw.ReleaseResident(); + uw.ReleaseResident(); mr.ready = true; } @@ -6887,12 +6882,9 @@ void BuildMoeMarlinResident(Dev d, const MoeBlockWeights& w, const HfConfig& cfg /*marlin_committed=*/MarlinMoeEnabled(), /*host_free_env=*/host_free_on); for (int e = 0; e < E; ++e) { const size_t se = static_cast(e); - w.expert_gate_fp4[se].d_packed.reset(); - w.expert_gate_fp4[se].d_scale.reset(); - w.expert_up_fp4[se].d_packed.reset(); - w.expert_up_fp4[se].d_scale.reset(); - w.expert_down_fp4[se].d_packed.reset(); - w.expert_down_fp4[se].d_scale.reset(); + w.expert_gate_fp4[se].ReleaseResident(); + w.expert_up_fp4[se].ReleaseResident(); + w.expert_down_fp4[se].ReleaseResident(); if (release_host) { w.expert_gate_fp4[se].packed.ReleaseHost(); w.expert_gate_fp4[se].scale.ReleaseHost(); diff --git a/tests/scripts/test_check_fp4_resident_consistency.py b/tests/scripts/test_check_fp4_resident_consistency.py index bf960b9cd..a1e536592 100644 --- a/tests/scripts/test_check_fp4_resident_consistency.py +++ b/tests/scripts/test_check_fp4_resident_consistency.py @@ -53,24 +53,22 @@ # exactly one thing in it. GOOD = """ Nvfp4Dev ResidentNvfp4(Dev d, const Nvfp4Weight& w) { - if (!w.d_packed) { + if (!w.packed.d_dev) { const size_t pb = w.packed.bytes.size(); void* p = d.b.Alloc(pb); vllm::load_stats::AddDeviceUpload(pb); d.b.Copy(d.q, p, w.packed.bytes.data(), pb); Backend* bk = &d.b; - w.d_packed = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); - w.packed.d_dev = w.d_packed; + w.packed.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); AdoptDeviceBytesAsHost(d.b, w.packed); } - if (!w.d_scale) { + if (!w.scale.d_dev) { const size_t sb = w.scale.bytes.size(); void* p = d.b.Alloc(sb); vllm::load_stats::AddDeviceUpload(sb); d.b.Copy(d.q, p, w.scale.bytes.data(), sb); Backend* bk = &d.b; - w.d_scale = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); - w.scale.d_dev = w.d_scale; + w.scale.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); }); AdoptDeviceBytesAsHost(d.b, w.scale); } Nvfp4Dev r; @@ -78,6 +76,22 @@ } """ +# The ONE release both upload copies above are released through. A miniature of +# `Nvfp4Weight::ReleaseResident`: the two `d_dev` resets are the ownership, the two +# guarded `ReleaseHost()` drops are the adopted host twin that holds the same block +# alive on a host-addressable device and therefore has to go FIRST. +GOOD_RELEASE = """ +void ReleaseResident() const { + if (packed.HostViewIsDeviceTwin()) packed.ReleaseHost(); + if (scale.HostViewIsDeviceTwin()) scale.ReleaseHost(); + packed.d_dev.reset(); + scale.d_dev.reset(); +} +""" + +PACKED_PUBLISH = " w.packed.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); });\n" +SCALE_PUBLISH = " w.scale.d_dev = std::shared_ptr(p, [bk](void* q) { bk->Free(q); });\n" + def body_of(text: str) -> str: """The single ResidentNvfp4 body, through the SAME normalization `file_violations` @@ -156,10 +170,10 @@ def test_each_buffer_gets_its_own_block(self) -> None: def test_a_missing_upload_block_is_a_violation(self) -> None: # A restructure the checker cannot scope must go RED and say so, never pass # by silently falling back to body-wide matching. - text = mutate("if (!w.d_scale) {", "if (true) {") + text = mutate("if (!w.scale.d_dev) {", "if (true) {") problems = body_violations(body_of(text)) self.assertEqual(len(problems), 1, problems) - self.assertIn("no `if (!w.d_scale)` upload block", problems[0]) + self.assertIn("no `if (!w.scale.d_dev)` upload block", problems[0]) class InvariantTests(unittest.TestCase): @@ -222,25 +236,28 @@ def test_uploading_the_other_buffers_bytes_fails(self) -> None: def test_dropping_the_packed_publication_fails(self) -> None: # Mutation (2): `d_dev` is never published, so AdoptDeviceBytesAsHost returns # immediately and the adoption below it is a silent no-op. - text = mutate(" w.packed.d_dev = w.d_packed;\n") + text = mutate(PACKED_PUBLISH) problems = body_violations(body_of(text)) self.assertTrue(any("`packed` does not PUBLISH" in p for p in problems), problems) self.assertFalse(any("`scale` does not PUBLISH" in p for p in problems), problems) def test_dropping_the_scale_publication_fails(self) -> None: - text = mutate(" w.scale.d_dev = w.d_scale;\n") + text = mutate(SCALE_PUBLISH) problems = body_violations(body_of(text)) self.assertTrue(any("`scale` does not PUBLISH" in p for p in problems), problems) def test_publishing_a_null_handle_fails(self) -> None: # The statement is still there; it publishes nothing. AdoptDeviceBytesAsHost # keys on `d_dev` and returns on null, so this is the deletion in disguise. - text = mutate(" w.packed.d_dev = w.d_packed;\n", " w.packed.d_dev = nullptr;\n") + text = mutate(PACKED_PUBLISH, " w.packed.d_dev = nullptr;\n") problems = body_violations(body_of(text)) self.assertTrue(any("`packed` does not PUBLISH" in p for p in problems), problems) - def test_publishing_the_other_buffers_handle_fails(self) -> None: - text = mutate(" w.scale.d_dev = w.d_scale;\n", " w.scale.d_dev = w.d_packed;\n") + def test_publishing_a_bare_alias_owns_nothing_fails(self) -> None: + # The statement is present and non-null; it just claims some other object's + # device resident without owning it and without a deleter, which is the + # two-owner shape this row deleted, one indirection further along. + text = mutate(SCALE_PUBLISH, " w.scale.d_dev = w.packed.d_dev;\n") problems = body_violations(body_of(text)) self.assertTrue(any("`scale` does not PUBLISH" in p for p in problems), problems) @@ -272,8 +289,8 @@ def test_publishing_after_the_adoption_fails(self) -> None: # Mutation (4): the ORDERING. Publishing d_dev after the adopt call leaves the # adoption looking at a null handle — present, and useless. text = mutate( - " w.packed.d_dev = w.d_packed;\n AdoptDeviceBytesAsHost(d.b, w.packed);\n", - " AdoptDeviceBytesAsHost(d.b, w.packed);\n w.packed.d_dev = w.d_packed;\n", + PACKED_PUBLISH + " AdoptDeviceBytesAsHost(d.b, w.packed);\n", + " AdoptDeviceBytesAsHost(d.b, w.packed);\n" + PACKED_PUBLISH, ) problems = body_violations(body_of(text)) self.assertTrue(any("publishes `d_dev` AFTER" in p for p in problems), problems) @@ -288,13 +305,11 @@ def test_adopting_before_the_upload_copy_fails(self) -> None: text = mutate( " d.b.Copy(d.q, p, w.packed.bytes.data(), pb);\n" " Backend* bk = &d.b;\n" - " w.d_packed = std::shared_ptr(p, [bk](void* q) { bk->Free(q); });\n" - " w.packed.d_dev = w.d_packed;\n" - " AdoptDeviceBytesAsHost(d.b, w.packed);\n", + + PACKED_PUBLISH + + " AdoptDeviceBytesAsHost(d.b, w.packed);\n", " Backend* bk = &d.b;\n" - " w.d_packed = std::shared_ptr(p, [bk](void* q) { bk->Free(q); });\n" - " w.packed.d_dev = w.d_packed;\n" - " AdoptDeviceBytesAsHost(d.b, w.packed);\n" + + PACKED_PUBLISH + + " AdoptDeviceBytesAsHost(d.b, w.packed);\n" " d.b.Copy(d.q, p, w.packed.bytes.data(), pb);\n", ) problems = body_violations(body_of(text)) @@ -310,8 +325,8 @@ def test_dropping_all_six_statements_fails(self) -> None: for line in ( " vllm::load_stats::AddDeviceUpload(pb);\n", " vllm::load_stats::AddDeviceUpload(sb);\n", - " w.packed.d_dev = w.d_packed;\n", - " w.scale.d_dev = w.d_scale;\n", + PACKED_PUBLISH, + SCALE_PUBLISH, " AdoptDeviceBytesAsHost(d.b, w.packed);\n", " AdoptDeviceBytesAsHost(d.b, w.scale);\n", ): @@ -374,17 +389,11 @@ def test_a_commented_upload_counter_is_not_a_pass(self) -> None: self.assertFalse(any("`packed`'s host->device upload is NOT counted" in p for p in problems), problems) def test_a_commented_publication_is_not_a_pass(self) -> None: - text = mutate( - " w.packed.d_dev = w.d_packed;\n", - " // w.packed.d_dev = w.d_packed;\n", - ) + text = mutate(PACKED_PUBLISH, " // " + PACKED_PUBLISH.strip() + "\n") self.assertTrue(any("`packed` does not PUBLISH" in p for p in body_violations(body_of(text)))) def test_an_if_0_publication_is_not_a_pass(self) -> None: - text = mutate( - " w.scale.d_dev = w.d_scale;\n", - "#if 0\n w.scale.d_dev = w.d_scale;\n#endif\n", - ) + text = mutate(SCALE_PUBLISH, "#if 0\n" + SCALE_PUBLISH + "#endif\n") self.assertTrue(any("`scale` does not PUBLISH" in p for p in body_violations(body_of(text)))) def test_a_commented_upload_copy_is_not_a_pass(self) -> None: @@ -402,7 +411,7 @@ def test_ordinary_comments_beside_live_statements_still_pass(self) -> None: text = mutate( " vllm::load_stats::AddDeviceUpload(pb);\n", " // ENG-LOAD-DIRECT-UPLOAD (issue #150): account the move, publish\n" - " // the allocation on the OwnedTensor (w.packed.d_dev = w.d_packed),\n" + " // the allocation on the OwnedTensor (w.packed.d_dev = shared_ptr),\n" " /* then AdoptDeviceBytesAsHost(d.b, w.packed) releases the pages. */\n" " vllm::load_stats::AddDeviceUpload(pb);\n", ) @@ -439,6 +448,95 @@ def test_a_whole_definition_inside_if_0_is_not_COUNTED_as_checked(self) -> None: self.assertIn("no `ResidentNvfp4", problems[0]) +class ReleaseTests(unittest.TestCase): + """Clause (g): the one release, and the ordering that makes it a release. + + Pre-fix, the repack builders reset ONLY the type-specific handle, so the + surviving `d_dev` alias kept a full packed+scale device copy for the process + lifetime. These cases pin the clause that sees that shape and each half of its + repair, on a miniature and on the REAL header both upload copies release + through.""" + + def test_the_real_release_shape_passes(self) -> None: + self.assertEqual(mod.release_violations(GOOD_RELEASE, "mini"), []) + + def test_dropping_the_packed_reset_fails(self) -> None: + text = mutate(" packed.d_dev.reset();\n", "", GOOD_RELEASE) + problems = mod.release_violations(text, "mini") + self.assertTrue(any("`packed.d_dev` is never reset" in p for p in problems), problems) + self.assertFalse(any("`scale.d_dev` is never reset" in p for p in problems), problems) + + def test_dropping_the_scale_reset_fails(self) -> None: + text = mutate(" scale.d_dev.reset();\n", "", GOOD_RELEASE) + problems = mod.release_violations(text, "mini") + self.assertTrue(any("`scale.d_dev` is never reset" in p for p in problems), problems) + + def test_dropping_the_packed_twin_drop_fails(self) -> None: + text = mutate( + " if (packed.HostViewIsDeviceTwin()) packed.ReleaseHost();\n", "", GOOD_RELEASE + ) + problems = mod.release_violations(text, "mini") + self.assertTrue(any("`packed`'s adopted host twin is never dropped" in p for p in problems), problems) + + def test_dropping_the_scale_twin_drop_fails(self) -> None: + text = mutate( + " if (scale.HostViewIsDeviceTwin()) scale.ReleaseHost();\n", "", GOOD_RELEASE + ) + problems = mod.release_violations(text, "mini") + self.assertTrue(any("`scale`'s adopted host twin is never dropped" in p for p in problems), problems) + + def test_dropping_the_twin_after_the_reset_fails(self) -> None: + # The ordering the release depends on: the adopted view holds the block, so + # a `d_dev` reset first frees nothing and the drop below touches a view + # whose owner is gone. + text = mutate( + " if (scale.HostViewIsDeviceTwin()) scale.ReleaseHost();\n", "", GOOD_RELEASE + ) + text = mutate( + " scale.d_dev.reset();\n", + " scale.d_dev.reset();\n if (scale.HostViewIsDeviceTwin()) scale.ReleaseHost();\n", + text, + ) + problems = mod.release_violations(text, "mini") + self.assertTrue(any("dropped AFTER its `d_dev` reset" in p for p in problems), problems) + + def test_a_missing_release_definition_is_a_violation_not_a_pass(self) -> None: + problems = mod.release_violations("int main() { return 0; }", "some/file.h") + self.assertEqual(len(problems), 1) + self.assertIn("no `void ReleaseResident", problems[0]) + + def test_the_LIVE_release_resets_both_dev_slots(self) -> None: + text = (ROOT / mod.RELEASE_SOURCE).read_text(encoding="utf-8", errors="ignore") + self.assertIsNotNone(mod.release_body(text), "ReleaseResident() is unreachable") + self.assertEqual(mod.release_violations(text, str(mod.RELEASE_SOURCE)), []) + + def test_LIVE_scale_reset_dropped_goes_red(self) -> None: + # THE LEAK, against the real header: the pre-fix builders dropped only the + # type-specific handle, so this is the shape a regression takes. + rel = mod.RELEASE_SOURCE + text = (ROOT / rel).read_text(encoding="utf-8", errors="ignore") + anchor = " scale.d_dev.reset();\n" + self.assertEqual(text.count(anchor), 1, anchor) + problems = mod.release_violations(text.replace(anchor, "", 1), str(rel)) + self.assertTrue(any("`scale.d_dev` is never reset" in p for p in problems), problems) + + def test_LIVE_packed_reset_dropped_goes_red(self) -> None: + rel = mod.RELEASE_SOURCE + text = (ROOT / rel).read_text(encoding="utf-8", errors="ignore") + anchor = " packed.d_dev.reset();\n" + self.assertEqual(text.count(anchor), 1, anchor) + problems = mod.release_violations(text.replace(anchor, "", 1), str(rel)) + self.assertTrue(any("`packed.d_dev` is never reset" in p for p in problems), problems) + + def test_LIVE_scale_twin_dropped_goes_red(self) -> None: + rel = mod.RELEASE_SOURCE + text = (ROOT / rel).read_text(encoding="utf-8", errors="ignore") + anchor = " if (scale.HostViewIsDeviceTwin()) scale.ReleaseHost();\n" + self.assertEqual(text.count(anchor), 1, anchor) + problems = mod.release_violations(text.replace(anchor, "", 1), str(rel)) + self.assertTrue(any("`scale`'s adopted host twin is never dropped" in p for p in problems), problems) + + class LiveTreeTests(unittest.TestCase): def test_the_checker_passes_on_the_current_tree(self) -> None: r = subprocess.run( @@ -446,10 +544,11 @@ def test_the_checker_passes_on_the_current_tree(self) -> None: ) self.assertEqual(r.returncode, 0, r.stdout + r.stderr) - def test_both_copies_are_actually_reached(self) -> None: + def test_both_copies_and_the_release_are_actually_reached(self) -> None: # The gate is worthless if SOURCES drifts off the real files: it would print - # OK over nothing. Assert both named files really contain a checked body, - # with both per-buffer upload blocks inside it. + # OK over nothing. Assert every named file really contains a checked body, + # with both per-buffer upload blocks inside it, and that the release source + # really defines the one release. for rel in mod.SOURCES: text = (ROOT / rel).read_text(encoding="utf-8", errors="ignore") bodies = checked_bodies(text) @@ -457,6 +556,8 @@ def test_both_copies_are_actually_reached(self) -> None: _, recv, body = bodies[0] for buffer in mod.BUFFERS: self.assertIsNotNone(buffer_block(body, recv, buffer), f"{rel}:{buffer}") + release_text = (ROOT / mod.RELEASE_SOURCE).read_text(encoding="utf-8", errors="ignore") + self.assertIsNotNone(mod.release_body(release_text), str(mod.RELEASE_SOURCE)) def test_disguised_deletions_of_the_LIVE_duplicate_go_red(self) -> None: # The findings against the REAL qwen3_5.cpp text, not the miniature. The @@ -465,7 +566,7 @@ def test_disguised_deletions_of_the_LIVE_duplicate_go_red(self) -> None: rel = Path("src/vllm/model_executor/models/qwen3_5.cpp") text = (ROOT / rel).read_text(encoding="utf-8", errors="ignore") adopt = " AdoptDeviceBytesAsHost(d.b, w.packed);\n" - publish = " w.packed.d_dev = w.d_packed;\n" + publish = PACKED_PUBLISH for anchor, replacement in ( (adopt, " // AdoptDeviceBytesAsHost(d.b, w.packed);\n"), (adopt, " /* AdoptDeviceBytesAsHost(d.b, w.packed); */\n"), @@ -494,9 +595,8 @@ def test_sinking_the_LIVE_upload_copy_below_its_adoption_goes_red(self) -> None: anchor = ( " d.b.Copy(d.q, p, w.packed.bytes.data(), pb);\n" " Backend* bk = &d.b;\n" - " w.d_packed = std::shared_ptr(p, [bk](void* q) { bk->Free(q); });\n" - " w.packed.d_dev = w.d_packed;\n" - " AdoptDeviceBytesAsHost(d.b, w.packed);\n" + + PACKED_PUBLISH + + " AdoptDeviceBytesAsHost(d.b, w.packed);\n" ) self.assertEqual(text.count(anchor), 1) lines = anchor.splitlines(keepends=True) diff --git a/tests/vllm/test_load_direct_upload.cpp b/tests/vllm/test_load_direct_upload.cpp index 71046d1c6..f93a6f94e 100644 --- a/tests/vllm/test_load_direct_upload.cpp +++ b/tests/vllm/test_load_direct_upload.cpp @@ -645,10 +645,10 @@ TEST_CASE("adopt: the windowed-release flag still governs the direct-upload rele // (`dense_nvfp4_gemm.h` and the private one in `qwen3_5.cpp`): // // 1. `load_stats::AddDeviceUpload(nb)` — account the move; -// 2. `w.packed.d_dev = w.d_packed` — PUBLISH the allocation on the -// OwnedTensor, which is the only reason `AdoptDeviceBytesAsHost` can act -// on it at all (that function keys on `d_dev` and returns immediately when -// it is null); +// 2. `w.packed.d_dev = ` — publish the allocation on the +// OwnedTensor, which is BOTH the only reason `AdoptDeviceBytesAsHost` can +// act on it (that function keys on `d_dev`) and the ONE owner of those +// bytes: an Nvfp4Weight has no second device handle to forget to release; // 3. `AdoptDeviceBytesAsHost(d.b, w.packed)` — the post-upload residency step // every other qualifying weight already got. // @@ -738,8 +738,8 @@ TEST_CASE("fp4 resident: ResidentNvfp4 COUNTS its upload, publishes d_dev, and a vllm::Nvfp4Weight w = BorrowedFp4Weight(dims, mp, ms); REQUIRE(w.packed.bytes.size() == dims.packed_bytes); REQUIRE(w.scale.bytes.size() == dims.scale_bytes); - REQUIRE(w.d_packed == nullptr); - REQUIRE(w.d_scale == nullptr); + REQUIRE(w.packed.d_dev == nullptr); + REQUIRE(w.scale.d_dev == nullptr); const vllm::load_stats::Counters before = vllm::load_stats::Snapshot(); vllm::dense_attn::Dev d{b, q}; @@ -758,17 +758,15 @@ TEST_CASE("fp4 resident: ResidentNvfp4 COUNTS its upload, publishes d_dev, and a // (2) PUBLICATION. The allocation is on the OwnedTensor, not only on the // Nvfp4Weight's own handle. Without this the adoption below cannot happen // at all: AdoptDeviceBytesAsHost returns immediately on a null `d_dev`. - REQUIRE(w.d_packed != nullptr); - REQUIRE(w.d_scale != nullptr); - CHECK(w.packed.d_dev.get() == w.d_packed.get()); - CHECK(w.scale.d_dev.get() == w.d_scale.get()); + REQUIRE(w.packed.d_dev != nullptr); + REQUIRE(w.scale.d_dev != nullptr); // (3) ADOPTION. The host view IS the device allocation, and the weight is no // longer a direct-upload borrow. CHECK(w.packed.bytes.borrowed()); CHECK(w.scale.bytes.borrowed()); - CHECK(static_cast(w.packed.bytes.data()) == w.d_packed.get()); - CHECK(static_cast(w.scale.bytes.data()) == w.d_scale.get()); + CHECK(static_cast(w.packed.bytes.data()) == w.packed.d_dev.get()); + CHECK(static_cast(w.scale.bytes.data()) == w.scale.d_dev.get()); CHECK(w.packed.mmap_src == nullptr); CHECK(w.scale.mmap_src == nullptr); CHECK(w.packed.mmap_src_bytes == 0u); @@ -792,13 +790,51 @@ TEST_CASE("fp4 resident: ResidentNvfp4 COUNTS its upload, publishes d_dev, and a CHECK(ms.byte_at_drop == 0); // (5) The returned device views are the uploaded buffers. - CHECK(dev.packed.data == w.d_packed.get()); - CHECK(dev.scale.data == w.d_scale.get()); + CHECK(dev.packed.data == w.packed.d_dev.get()); + CHECK(dev.scale.data == w.scale.d_dev.get()); } - // ONE control block per buffer, despite the two handles (`d_packed` and - // `packed.d_dev`, plus the adopted `bytes` keep-alive aliasing it): the device - // memory is freed exactly once, through the vt Backend. + // ONE control block per buffer: `packed.d_dev` is the ONE owner and the + // adopted `bytes` view aliases the same block, so the device memory is freed + // exactly once, through the vt Backend. + CHECK(b.frees == 2); +} + +// THE LEAK THE SINGLE OWNER REPAIRS. After a Marlin repack the fp4 originals are +// dead weight, and the repack builders are their only owner. Pre-fix, the +// release they performed dropped only the type-specific `d_packed`/`d_scale` +// pair, and `packed.d_dev` was a SECOND handle on the same control block, so the +// block survived every repack for the process lifetime -- once per repacked +// weight and once per expert per projection on an MoE model. That silently +// undid the repack's whole point (`dense_nvfp4_gemm.h`: "then FREE the fp4 +// originals"), on exactly the 24 GiB cards the NVFP4 arms are served from. +// The invariant below is the one the leak violated: a release through the +// weight's own resident state must actually free BOTH buffers. +TEST_CASE("fp4 resident: releasing the resident frees both device buffers") { + ForcedResidencyArm arm; + ScopedEnvVar adopt_default("VT_ADOPT_DEVICE_BYTES", "1"); + // A discrete (non-host-addressable) device: no adoption, so before the fix + // `packed.d_dev` and the type-specific handle were two references to one + // block. This is the arm every consumer Blackwell card takes. + FakeBackend b(/*host_addressable=*/false); + vt::Queue q = b.CreateQueue(); + const Fp4Dims dims = MakeFp4Dims(); + + ObservableMapping mp; + ObservableMapping ms; + vllm::Nvfp4Weight w = BorrowedFp4Weight(dims, mp, ms); + vllm::dense_attn::Dev d{b, q}; + const vllm::dense_nvfp4::Nvfp4Dev dev = vllm::dense_nvfp4::ResidentNvfp4(d, w); + REQUIRE(dev.packed.data != nullptr); + REQUIRE(dev.scale.data != nullptr); + CHECK(b.frees == 0); + + // The one release a repack builder performs once the repack has landed. RED + // before the single-owner change: the builders reset the type-specific pair + // and the surviving `packed.d_dev`/`scale.d_dev` alias kept both blocks, so + // this read 0 rather than 2. + w.ReleaseResident(); + CHECK(b.frees == 2); } @@ -866,15 +902,14 @@ TEST_CASE("fp4 resident: a host-addressable device adopts an OWNED fp4 mirror to CHECK(after.device_upload_bytes - before.device_upload_bytes == dims.packed_bytes + dims.scale_bytes); - REQUIRE(w.d_packed != nullptr); - CHECK(w.packed.d_dev.get() == w.d_packed.get()); - CHECK(w.scale.d_dev.get() == w.d_scale.get()); + REQUIRE(w.packed.d_dev != nullptr); + CHECK(w.scale.d_dev != nullptr); // Adopted: ONE copy, and it is the device one. CHECK(w.packed.bytes.borrowed()); CHECK(w.scale.bytes.borrowed()); - CHECK(static_cast(w.packed.bytes.data()) == w.d_packed.get()); - CHECK(static_cast(w.scale.bytes.data()) == w.d_scale.get()); + CHECK(static_cast(w.packed.bytes.data()) == w.packed.d_dev.get()); + CHECK(static_cast(w.scale.bytes.data()) == w.scale.d_dev.get()); CHECK(w.packed.bytes.size() == dims.packed_bytes); CHECK(w.scale.bytes.size() == dims.scale_bytes); CHECK_FALSE(w.packed.host_released);