From 09c05ce18067cf1d4827bcf0344f90f55b4f74f2 Mon Sep 17 00:00:00 2001 From: Ketor Date: Sun, 13 Sep 2026 11:31:22 +0800 Subject: [PATCH] fix: restore runtime RDMA and native hybrid cache support Prepare v2.27.0 with runtime provider packaging, version-scoped SGLang/vLLM hybrid-state patches, and a real-GPU restoration regression. Native client/server ABI and wire formats are unchanged. --- .github/workflows/ci.yml | 25 ++ CHANGELOG.md | 17 + Dockerfile | 4 +- VERSION | 2 +- docs/CONNECTORS.md | 51 ++- docs/DEPLOY.md | 4 +- integration/common/pyproject.toml | 2 +- .../sglang-hybrid-mamba-dsa-indexer.patch | 27 +- .../tests/gpu_hybrid_state_roundtrip.py | 365 ++++++++++++++++++ integration/lmcache/pyproject.toml | 4 +- ...m-glm53flash-kpool-tail-slot-mapping.patch | 21 + integration/vllm/pyproject.toml | 4 +- 12 files changed, 496 insertions(+), 30 deletions(-) create mode 100755 integration/hicache/tests/gpu_hybrid_state_roundtrip.py create mode 100644 integration/vllm/patches/vllm-glm53flash-kpool-tail-slot-mapping.patch diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 93c94e5..79ba1e4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -238,6 +238,31 @@ jobs: DFKV_SERVER_URING=1 ./build/rdma_lease_loopback_test --gtest_output=xml:lease-uring.xml ./build/rdma_lease_client_test --gtest_output=xml:lease-client.xml python3 -c 'import xml.etree.ElementTree as E; roots=[E.parse(p).getroot() for p in ("lease-sync.xml","lease-uring.xml","lease-client.xml")]; assert all(int(r.attrib["tests"]) > 0 and not list(r.iter("skipped")) and not list(r.iter("failure")) for r in roots)' + - name: Exercise RDMA through the shipped runtime image + if: steps.rxe_probe.outputs.ok == 'true' + run: | + set -eu + docker build --target runtime -t dfkv-runtime-rdma . + id=$(docker run -d --privileged --network host \ + --ulimit memlock=-1:-1 -e DFKV_RDMA_DEV=rxe0 \ + dfkv-runtime-rdma --dir /tmp/dfkv-data \ + --port 12000 --rdma-port 12001 \ + --metrics-port 12010 --metrics-bind 127.0.0.1 \ + --cap 67108864 --store-engine file) + trap 'docker logs "$id"; docker rm -f "$id" >/dev/null' EXIT + ready=0 + for attempt in $(seq 1 60); do + if docker exec "$id" dfkvctl stat 127.0.0.1:12000 \ + | grep -q 'dfkv_server_ready 1'; then + ready=1 + break + fi + test "$(docker inspect -f '{{.State.Running}}' "$id")" = true + sleep 1 + done + test "$ready" = 1 + docker exec -e DFKV_RDMA=1 -e DFKV_RDMA_DEV=rxe0 "$id" \ + dfkv_smoke --members n=127.0.0.1:12001 --size 1048576 - name: Build + run RDMA datapath test under ThreadSanitizer (hard gate when rxe works) if: steps.rxe_probe.outputs.ok == 'true' run: | diff --git a/CHANGELOG.md b/CHANGELOG.md index d43218c..1e0e1e0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,23 @@ ## Unreleased +### v2.27.0 — Runtime RDMA providers and native hybrid-state integration + +- Include `ibverbs-providers` in the runtime image; exposing RDMA devices does + not make them discoverable when the userspace hardware plugins are absent. +- Exercise real RDMA PUT/GET through the built runtime image when the CI + hardware probe succeeds, rather than relying on version or TCP-only checks. +- Refresh the shipped SGLang hybrid Mamba/DSA patch for the pinned native + GLM-5.3-Flash image, preserving indexer sidecars and actual scaled-MLA host + geometry with the current allocator API. +- Add a bounded real-GPU offload/poison/restore regression for KV, indexer, + temporal and convolution state, without model weights. +- Ship a targeted native vLLM GLM-5.3-Flash V2-runner patch: exclude the + single-block kpool tail from generic position-indexed slot mapping, retaining + the dedicated circular metadata builder and existing allocation policy. +- Document strict patch applicability and fresh cache generations. Native + client/server wire formats and connector runtime code are unchanged. + ### v2.26.4 — SGLang DeepSeek-V4.1 HiCache deployment guidance - Add the verified preview-image HiCache/L3 recipe and physical side-pool diff --git a/Dockerfile b/Dockerfile index de275c2..f975d22 100644 --- a/Dockerfile +++ b/Dockerfile @@ -9,7 +9,7 @@ # DFKV_STATIC_LIBSTDCXX folds libstdc++/libgcc into the artifacts (those are NOT # glibc; static-linking them removes a separate runtime dep). libibverbs CANNOT # be static (it dlopen()s provider drivers at runtime), so the run node still -# needs rdma-core / libibverbs installed. +# needs rdma-core / libibverbs and their hardware provider plugins installed. # ubuntu:22.04, pinned so a release tag always rebuilds from identical bytes. ARG DFKV_BASE_IMAGE=docker.io/library/ubuntu@sha256:3b06811b2afd352be909dd088a004166d665dc76d38b13eada33522a9d915c6f FROM ${DFKV_BASE_IMAGE} AS build @@ -29,7 +29,7 @@ RUN cmake -S . -B build -G Ninja -DCMAKE_BUILD_TYPE=Release -DDFKV_BUILD_TESTS=O FROM ${DFKV_BASE_IMAGE} AS runtime RUN apt-get update && apt-get install -y --no-install-recommends \ - rdma-core libibverbs1 && rm -rf /var/lib/apt/lists/* + rdma-core libibverbs1 ibverbs-providers && rm -rf /var/lib/apt/lists/* COPY --from=build /out/ /usr/local/ RUN ldconfig EXPOSE 12000 diff --git a/VERSION b/VERSION index dabe377..295b40e 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -2.26.4 \ No newline at end of file +2.27.0 \ No newline at end of file diff --git a/docs/CONNECTORS.md b/docs/CONNECTORS.md index 7fac504..c1aefd1 100644 --- a/docs/CONNECTORS.md +++ b/docs/CONNECTORS.md @@ -532,23 +532,41 @@ sglang serve /models/glm-5.2-nvfp4 --served-model-name glm-5.2 \ MR,无 payload memcpy。 - MLA 下插件自动单对象、无 rank 后缀、`backup_skip`(仅 tp_rank0 写)。decode 共享前缀配同 members。 - **多池模型**(Mamba/SWA/DeepSeek-V4)用 v2 PoolTransfer 接口(插件已实现)。 - DSA/DeepSeekV4 主 `kv` 池是无数据的 LogicalHostPool(`get_page_buffer_meta→None`), - 插件对其 `batch_set_v1` 写空 marker 锚定命中前缀、`batch_get_v1` no-op,真实 KV 走 v2 侧池。 + 部分 DeepSeek-V4 等引擎的主 `kv` 是无数据 LogicalHostPool(`get_page_buffer_meta→None`): + 插件写 marker 锚定前缀,主池读取不搬运字节,实际数据走 v2 侧池。普通 MLA/DSA + 的主 `kv` 可以是真实物理池;必须按运行时布局核对,不能将所有 DSA 都当作 marker。 - **Mamba+DSA 混合模型必须同时注册 `KV`、`MAMBA`、`INDEXER`。** 只恢复 latent 与 recurrent state、遗漏历史 indexer 数据,会出现“所有已请求 对象均命中,但长提示答案错误”;不能通过增加生成长度或把命中计数当作正确性 来跳过检查。SGLang 的普通 DSA 路径正确注册 sidecar,不代表混合路径也已注册。 对尚缺该接线的引擎,随包提供 [混合 DSA indexer 修复补丁](../integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch)。 - 它修复引擎的 host-pool entry 与 `SidecarPoolSpec` 两处声明,不改变 dfkv - value/wire 格式。补丁的实测源文件是 - `python/sglang/srt/mem_cache/hybrid_cache/hybrid_pool_assembler.py`, - 原始 SHA256 为 `48456e7ce2cf20f839d100333404c90ba1d65b370a7f8265e9bc044ec18c2216`; - 应在对应 SGLang 源码树先执行 `git apply --check`,再应用、重建或部署受控 - overlay。不要盲目覆盖其它引擎版本。启动池描述必须包含 `INDEXER`,并对原始 - 长提示执行完整进程重启后的 L3 回载与答案校验;旧的缺 indexer 缓存不能视为 - 上游 SGLang main(2026-09-08 核对)仍未在 `build_hybrid_mamba_stack` 注册 - INDEXER,该补丁对当前上游同样适用,并非仅限旧镜像。 + 它补齐 host-pool entry 与 `SidecarPoolSpec`,并让 MLA 宿主页使用设备池实际 + `kv_cache_dim`,避免带缩放布局被默认维度截断;不改变 dfkv value/wire 格式。 + 当前补丁针对官方 `lmsysorg/sglang:glm-5.3-flash` 的 + `0.0.0.dev1+gf609d677b` 构建,镜像 digest 为 + `sha256:a2c0f7d4d9ebce97a2707c5415081d284d741db1033a1008a955453b9b5255bf`。 + 目标文件是 `python/sglang/srt/mem_cache/hybrid_cache/hybrid_pool_assembler.py`, + 原始 SHA256 为 `6689d642d7868e66a34cd95ef09e5d710ef6f40fd053ca2f1f7d04afe99223ff`。 + 先核对源文件,再严格检查和应用补丁;禁止用 fuzzy apply 跨版本套用。 + 例如在该镜像内将 dfkv 仓库挂载到 `/dfkv` 后: + + ```bash + git -C /sgl-workspace/sglang apply --check \ + /dfkv/integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch + git -C /sgl-workspace/sglang apply \ + /dfkv/integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch + python3 /dfkv/integration/hicache/tests/gpu_hybrid_state_roundtrip.py \ + --include-scaled-mla + ``` + + [GPU 状态回归](../integration/hicache/tests/gpu_hybrid_state_roundtrip.py)使用真实池、 + 分配器和 CUDA 传输:卸载后改写全部设备字节,再恢复到不同槽位,逐字节验证 + KV、INDEXER、temporal 与 conv 状态。它不加载模型权重,也不代替整模型、 + 多 rank、L3 或完整进程重启验收。部署时必须记录补丁和派生镜像身份, + 启动池描述须包含 `INDEXER`;使用新的 `model_revision`,不能复用曾遗漏 + indexer 或截断状态的旧缓存。其它引擎版本应重新核对接口与状态恢复,不能 + 从此目标构建外推适用性。 - **identity/layout 必须协同发布。** namespace 使用 SGLang runtime 给出的精确 `model_name` + `sglang-hicache/raw-v1`;同一模型的 pool/hash/并行坐标/component 进入 canonical object key。dfkv value 只有 raw bytes,不会检查 page size、 @@ -850,6 +868,17 @@ lsmod | grep nvidia_peermem 逻辑 block 数折叠物理 tile 轴,完整收集每个逻辑 block 的所有 bytes,不能 直接把逻辑 block ID 用作 kernel-tile 索引。 +#### GLM Flash 原生 V2 runner 的循环尾状态映射 + +`vllm/vllm-openai:glm53-flash` 的 `g385dce36b` 将 `KpoolTailSpec` 错当作 +普通位置索引缓存:长 prefill 在通用 slot-mapping 内核中越界,异步错误可能 +随后表现为 KDA 投影的 cuBLAS 失败。无 dfkv 连接器也会触发。 +随包提供 [循环尾映射补丁](../integration/vllm/patches/vllm-glm53flash-kpool-tail-slot-mapping.patch), +仅让该类型跳过通用映射;已有 `KpoolTailMetadataBuilder` 仍按每请求的 +首块与 `position % kpool` 生成真实尾状态位置。保留原表宽、APC、模型和请求长度。 +补丁针对上述源码版本,应用前运行 `git apply --check`,不盲目应用到其他版本。 +原短—短—长请求故障序列在修复后通过;这项引擎回归不替代 L3 字节完整性和性能验收。 + #### 混合模型的 KV 加载故障恢复 diff --git a/docs/DEPLOY.md b/docs/DEPLOY.md index 26aa380..37cc6f5 100644 --- a/docs/DEPLOY.md +++ b/docs/DEPLOY.md @@ -40,7 +40,7 @@ deploy/package_release.sh build release tar tzf "release/dfkv-$(cat VERSION)-linux-x86_64.tar.gz" ldd build/libdfkv.so | grep ibverbs ``` -依赖:`libibverbs-dev`(构建期)+ 运行节点装 `rdma-core`。无 RDMA 也可去掉 `-DDFKV_WITH_RDMA` 构建纯 TCP 版。 +依赖:`libibverbs-dev`(构建期);运行节点须安装 libibverbs 及对应硬件 provider。Ubuntu/Debian 运行包为 `rdma-core libibverbs1 ibverbs-providers`,只装前两项不足以发现设备。发布的 runtime 镜像包含这三项;设备透传不能替代用户态 provider。无 RDMA 也可去掉 `-DDFKV_WITH_RDMA` 构建纯 TCP 版。 > QP 信息走 TCP bootstrap 交换(非 librdmacm),所以只依赖 libibverbs,不需要 librdmacm。 > CacheLib/Navy 不是 v2.0.0 的构建依赖,也没有占位 backend/CMake 开关; > 评估结论、兼容性差距和重开条件见 @@ -51,7 +51,7 @@ ldd build/libdfkv.so | grep ibverbs > 直接构建的二进制不能部署到 glibc 2.35。CI 的 portable job 和仓库根 > `Dockerfile` 都固定 Ubuntu 22.04 + RDMA + 静态 libstdc++,并构建同一 > canonical tarball。`DFKV_STATIC_LIBSTDCXX` 只静态链接 libstdc++/libgcc; -> libibverbs 仍动态加载 provider,运行节点必须安装 rdma-core。 +> libibverbs 仍动态加载 provider,运行节点必须同时具备对应插件。镜像验收要在可用 RDMA 设备上完成真实 PUT/GET;仅 `--version` 或 TCP 冒烟不能证明 RDMA 可用。 ## 2. 每节点:分发 + 缓存目录 diff --git a/integration/common/pyproject.toml b/integration/common/pyproject.toml index b770fe2..0e60f5d 100644 --- a/integration/common/pyproject.toml +++ b/integration/common/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "dfkv-common" -version = "2.26.4" +version = "2.27.0" description = "Canonical namespace and pool-key schema shared by dfkv connectors" requires-python = ">=3.9" diff --git a/integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch b/integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch index 1a8db55..31137de 100644 --- a/integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch +++ b/integration/hicache/patches/sglang-hybrid-mamba-dsa-indexer.patch @@ -1,15 +1,23 @@ --- a/python/sglang/srt/mem_cache/hybrid_cache/hybrid_pool_assembler.py +++ b/python/sglang/srt/mem_cache/hybrid_cache/hybrid_pool_assembler.py -@@ -695,6 +695,8 @@ - storage_backend_extra_config: Optional[dict] = None, - enable_storage_metrics: bool = False, +@@ -724,7 +724,7 @@ ) -> tuple[HostPoolGroup, HybridCacheController]: -+ from sglang.srt.mem_cache.memory_pool import DSATokenToKVPool -+ transfer_layer_num = len(full_layer_mapping | mamba_layer_mapping) mamba_allocator = params.req_to_token_pool.mamba_allocator +- from sglang.srt.mem_cache.memory_pool import HybridLinearKVPool ++ from sglang.srt.mem_cache.memory_pool import DSATokenToKVPool, HybridLinearKVPool + mtp_draft_device_pools = tuple( -@@ -748,6 +750,24 @@ + pool.full_kv_pool if isinstance(pool, HybridLinearKVPool) else pool +@@ -739,6 +739,7 @@ + kv_pool=kv_pool, + page_size=params.page_size, + use_mla=use_mla, ++ override_kv_cache_dim=kv_pool.kv_cache_dim if use_mla else None, + host_size=kv_host_size, + mtp_draft_device_pools=mtp_draft_device_pools, + ) +@@ -778,6 +779,25 @@ device_free_fn=mamba_allocator.free, ), ] @@ -20,7 +28,7 @@ + kv_pool, + kv_host_pool, + get_memory().hicache_mem_layout, -+ allocator_type=_get_allocator_type(server_args), ++ allocator_type=_get_allocator_type(), + ) + entries.append( + build_pool_entry( @@ -29,12 +37,13 @@ + device_pool=kv_pool, + layer_mapping=full_layer_mapping, + transfer_layer_num=transfer_layer_num + len(mtp_draft_device_pools), ++ packed_draft_device_pools=mtp_draft_device_pools, + ) + ) host_pool_group = HostPoolGroup(entries) cache_controller = HybridCacheController( params.token_to_kv_pool_allocator, -@@ -1311,6 +1331,11 @@ +@@ -1345,6 +1365,11 @@ storage_backend_extra_config=storage_backend_extra_config, enable_storage_metrics=enable_storage_metrics, ) @@ -46,7 +55,7 @@ return StackBuildResult( host_pool_group=host_pool_group, cache_controller=cache_controller, -@@ -1318,9 +1343,10 @@ +@@ -1352,9 +1377,10 @@ ComponentType.FULL: host_pool_group.get_pool(PoolName.KV), ComponentType.MAMBA: host_pool_group.get_pool(PoolName.MAMBA), }, diff --git a/integration/hicache/tests/gpu_hybrid_state_roundtrip.py b/integration/hicache/tests/gpu_hybrid_state_roundtrip.py new file mode 100755 index 0000000..2115673 --- /dev/null +++ b/integration/hicache/tests/gpu_hybrid_state_roundtrip.py @@ -0,0 +1,365 @@ +#!/usr/bin/env python3 +"""Real-CUDA regression for hybrid HiCache state restoration, no model weights. + +Run inside the supported SGLang image with the dfkv checkout at /dfkv: + python3 /dfkv/integration/hicache/tests/gpu_hybrid_state_roundtrip.py \ + --include-scaled-mla +An explicit --assembler-override /path/to/hybrid_pool_assembler.py permits a +candidate comparison without modifying the installed SGLang package. + +Exit 0: every requested device-byte comparison passed. Exit 1: data mismatch. +Exit 2: unavailable capability/setup/transfer error. Exit 124: watchdog timeout. +Stdout contains one JSON document; library diagnostics go to stderr. +This exercises actual L1/L2 pools, allocators, controller and CUDA transfers; +it does not replace model, all-rank, L3, MTP or speculative-state acceptance. +""" + +import argparse +import contextlib +import hashlib +import importlib.util +import json +import os +from pathlib import Path +import sys +import threading +import time +import traceback +from types import SimpleNamespace + + +def digest(tensor): + return hashlib.sha256(tensor.numpy().tobytes()).hexdigest() + + +def require(value, message): + if not value: + raise RuntimeError(message) + + +def forbid_eviction(*args, **kwargs): + raise RuntimeError("Tiny fixture unexpectedly requested cache eviction") + + +def run_case(torch, assembler, case, server_args): + from sglang.srt.configs.mamba_utils import ( + KimiLinearCacheParams, + KimiLinearStateShape, + Mamba2StateDType, + ) + from sglang.srt.mem_cache.allocator.mamba import MambaSlotAllocator + from sglang.srt.mem_cache.allocator.paged import PagedTokenToKVPoolAllocator + from sglang.srt.mem_cache.cache_init_params import CacheInitParams + from sglang.srt.mem_cache.hicache_storage import PoolName, PoolTransfer + from sglang.srt.mem_cache.memory_pool import ( + DSATokenToKVPool, + HybridLinearKVPool, + MambaPool, + ) + from sglang.srt.mem_cache.unified_cache.component_type import ComponentType + + page_size, token_capacity, slots = 64, 256, 4 + full_layers, linear_layers = [1, 3], [0, 2] + kv_dim = 656 if case["layout"] == "scaled-mla" else 576 + case["phase"] = "construct_device_pools" + mamba = MambaPool( + size=slots, + spec_state_size=0, + cache_params=KimiLinearCacheParams( + shape=KimiLinearStateShape.create( + tp_world_size=1, num_heads=2, head_dim=16, conv_kernel_size=4 + ), + dtype=Mamba2StateDType(conv=torch.bfloat16, temporal=torch.float32), + layers=linear_layers, + ), + mamba_layer_ids=linear_layers, + device="cuda", + ) + hybrid = HybridLinearKVPool( + size=token_capacity, + dtype=torch.float8_e4m3fn, + page_size=page_size, + head_num=1, + head_dim=576, + full_attention_layer_ids=full_layers, + device="cuda", + mamba_pool=mamba, + use_mla=True, + use_dsa=True, + kv_lora_rank=512, + qk_rope_head_dim=64, + index_head_dim=128, + kv_cache_dim=kv_dim, + ) + kv = hybrid.full_kv_pool + require(isinstance(kv, DSATokenToKVPool), "Hybrid did not construct a real DSA pool") + require(kv.kv_cache_dim == kv_dim, "Requested device MLA geometry was not honored") + require(kv.dsa_kv_cache_store_fp8 == (kv_dim == 656), "Unexpected DSA FP8 storage mode") + kv_allocator = PagedTokenToKVPoolAllocator( + token_capacity, page_size, torch.float8_e4m3fn, "cuda", hybrid, False + ) + mamba_allocator = MambaSlotAllocator(slots, "cuda") + # Request/tree metadata only. Pools, slot allocation and all transfers are native. + req_pool = SimpleNamespace( + mamba_pool=mamba, + mamba_allocator=mamba_allocator, + mamba_map={global_id: local_id for local_id, global_id in enumerate(linear_layers)}, + ) + params = CacheInitParams( + disable=False, + req_to_token_pool=req_pool, + token_to_kv_pool_allocator=kv_allocator, + page_size=page_size, + ) + cache = SimpleNamespace( + req_to_token_pool=req_pool, + token_to_kv_pool_allocator=kv_allocator, + evict_host=forbid_eviction, + evict_for_alloc=forbid_eviction, + ) + case["phase"] = "assemble_native_strategy" + strategy = assembler._MambaStrategy() + require( + strategy.matches(hybrid, {ComponentType.FULL, ComponentType.MAMBA}), + "Native hybrid strategy does not match actual pools", + ) + stack = strategy.build( + cache=cache, + kvcache=hybrid, + params=params, + server_args=server_args, + load_cache_event=threading.Event(), + storage_backend=None, + ) + group, controller = stack.host_pool_group, stack.cache_controller + case["assembled_pools"] = [entry.name.value for entry in group.entries] + case["sidecars"] = [ + {"pool": spec.pool_name.value, "indices_from_pool": spec.indices_from_pool.value} + for spec in stack.sidecars + ] + case["geometry"] = { + "page_size": page_size, + "token_capacity": token_capacity, + "mamba_slot_capacity": slots, + "device_kv_dim": kv.kv_cache_dim, + "host_kv_dim": group.get_pool(PoolName.KV).kv_cache_dim, + "full_layer_mapping": hybrid.full_attention_layer_id_mapping, + "mamba_layer_mapping": req_pool.mamba_map, + "transfer_layer_num": stack.transfer_layer_num, + } + source_kv = kv_allocator.alloc(2 * page_size) + source_mamba = mamba_allocator.alloc(2) + require(source_kv is not None and source_mamba is not None, "Source allocation failed") + source_pages = source_kv.reshape(-1, page_size)[:, 0] // page_size + source_indices = {"kv": source_kv, "indexer": source_pages, "mamba": source_mamba} + buffers = [] + for local_id, global_id in enumerate(full_layers): + buffers.append((f"kv.layer_{global_id}", "kv", kv.kv_buffer[local_id])) + buffers.append((f"indexer.layer_{global_id}", "indexer", kv.index_k_with_scale_buffer[local_id])) + for local_id, global_id in enumerate(linear_layers): + buffers.append((f"temporal.layer_{global_id}", "mamba", mamba.mamba_cache.temporal[local_id])) + for conv_id, conv in enumerate(mamba.mamba_cache.conv): + buffers.append((f"conv_{conv_id}.layer_{global_id}", "mamba", conv[local_id])) + case["device_state_allocation_bytes"] = sum(t.numel() * t.element_size() for _, _, t in buffers) + require(case["device_state_allocation_bytes"] < 4 * 1024**2, "Fixture exceeded 4 MiB state bound") + expected = {} + case["phase"] = "fill_known_bytes" + for ordinal, (name, kind, tensor) in enumerate(buffers): + require(tensor.is_cuda and tensor.is_contiguous(), f"Unsupported physical buffer: {name}") + raw = tensor.view(torch.uint8).reshape(tensor.shape[0], -1) + # 255 is reserved for poison. Patterns differ across states/layers/rows. + pattern = ((torch.arange(raw.numel(), device="cuda", dtype=torch.int64) * 17 + ordinal * 29) % 251).to(torch.uint8) + raw.copy_(pattern.reshape_as(raw)) + expected[name] = raw.index_select(0, source_indices[kind]).cpu() + case["states"][name] = { + "status": "not_restored", + "bytes": expected[name].numel(), + "expected_sha256": digest(expected[name]), + "dtype": str(tensor.dtype), + "buffer_shape": list(tensor.shape), + } + torch.cuda.synchronize() + + # Follow the strategy's declared sidecars exactly. In the original module + # INDEXER is absent, so it is not copied; its *data comparison* must fail. + backup_extra = [PoolTransfer(name=PoolName.MAMBA, device_indices=source_mamba)] + backup_extra.extend( + PoolTransfer(name=spec.pool_name, indices_from_pool=spec.indices_from_pool, hit_policy=spec.hit_policy) + for spec in stack.sidecars + ) + case["phase"] = "offload_device_to_host" + host_kv = controller.write(source_kv, extra_pools=backup_extra) + require(host_kv is not None, "Host allocation failed") + controller.ack_write_queue[-1].finish_event.synchronize() + case["offload_completed"] = True + case["phase"] = "poison_all_device_state" + for name, kind, tensor in buffers: + raw = tensor.view(torch.uint8).reshape(tensor.shape[0], -1) + raw.fill_(255) + poisoned = raw.cpu() + require(bool(torch.all(poisoned == 255)), f"Poison did not reach actual device buffer: {name}") + case["states"][name]["poison_verified"] = True + torch.cuda.synchronize() + + restore_extra = [ + PoolTransfer( + name=transfer.name, + host_indices=transfer.host_indices, + indices_from_pool=transfer.indices_from_pool, + hit_policy=transfer.hit_policy, + ) + for transfer in backup_extra + ] + case["phase"] = "restore_host_to_new_device_slots" + destination_kv = controller.load(host_kv, extra_pools=restore_extra) + require(destination_kv is not None, "Restore device allocation failed") + destination_mamba = restore_extra[0].device_indices + require(destination_mamba is not None, "Restore Mamba allocation failed") + require(not bool(torch.isin(destination_kv, source_kv).any()), "Restore reused occupied source KV") + require(not bool(torch.isin(destination_mamba, source_mamba).any()), "Restore reused occupied source Mamba") + controller.start_loading() + controller.ack_load_queue[-1].finish_event.synchronize() + case["restore_completed"] = True + destination_indices = { + "kv": destination_kv, + "indexer": destination_kv.reshape(-1, page_size)[:, 0] // page_size, + "mamba": destination_mamba, + } + case["source_pages"] = source_pages.cpu().tolist() + case["destination_pages"] = destination_indices["indexer"].cpu().tolist() + case["source_mamba_slots"] = source_mamba.cpu().tolist() + case["destination_mamba_slots"] = destination_mamba.cpu().tolist() + case["phase"] = "compare_actual_device_bytes" + for name, kind, tensor in buffers: + raw = tensor.view(torch.uint8).reshape(tensor.shape[0], -1) + actual = raw.index_select(0, destination_indices[kind]).cpu() + wanted = expected[name] + mismatch = (actual != wanted).flatten() + bad = torch.nonzero(mismatch, as_tuple=False).flatten() + first = int(bad[0]) if bad.numel() else None + case["states"][name].update( + status="pass" if first is None else "fail", + actual_sha256=digest(actual), + equal_bytes=int(actual.numel() - bad.numel()), + mismatch_bytes=int(bad.numel()), + poison_bytes_remaining=int((actual == 255).sum()), + first_mismatch=None if first is None else { + "byte_offset": first, + "expected": int(wanted.flatten()[first]), + "actual": int(actual.flatten()[first]), + }, + ) + case["status"] = "pass" if all(row["status"] == "pass" for row in case["states"].values()) else "fail" + case["phase"] = "release_host_pools" + group.destroy() + case["phase"] = "complete" + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--assembler-override", type=Path, help="Explicit candidate assembler .py; default imports installed native module") + parser.add_argument("--include-scaled-mla", action="store_true", help="Also exercise 656-byte scaled-FP8 MLA rows after the official 576-byte raw-FP8 layout") + parser.add_argument("--timeout-seconds", type=int, default=120) + args = parser.parse_args() + if not 1 <= args.timeout_seconds <= 600: + parser.error("--timeout-seconds must be between 1 and 600") + output_stream = sys.stdout + started = time.monotonic() + result = { + "schema": 1, + "status": "error", + "backend": "direct", + "host_layout": "page_first_direct", + "assembler_override": str(args.assembler_override) if args.assembler_override else None, + "coverage_limits": [ + "Single CUDA device, two full and two KDA layers, two pages and two state slots; no distributed rank coordination.", + "Checks all stored bytes including index-K quantization scales, MLA payload/scales and KDA temporal/conv state; does not execute model projections or inference.", + "Native L1/L2 only; no model weights, storage backend, dfkv/network traffic, MTP, speculative state or compressed index-kpool tails.", + ], + "cases": [], + } + done = threading.Event() + + def watchdog(): + if not done.wait(args.timeout_seconds): + result.update(status="error", error="watchdog_timeout", elapsed_seconds=time.monotonic() - started) + payload = (json.dumps(result, sort_keys=True) + "\n").encode() + os.write(output_stream.fileno(), payload) + os._exit(124) + + threading.Thread(target=watchdog, daemon=True).start() + code = 2 + try: + with contextlib.redirect_stdout(sys.stderr): + import torch + import sglang + from sglang.srt.runtime_context import get_context, get_parallel + from sglang.srt.server_args import ServerArgs + + require(torch.cuda.is_available() and torch.version.cuda is not None, "A real NVIDIA CUDA runtime/device is required") + require(os.environ.get("SGLANG_MOONCAKE_CUSTOM_MEM_POOL") is None, "Custom Mooncake allocator is outside this bounded no-network smoke") + torch.cuda.set_device(0) + result["runtime"] = { + "sglang_version": getattr(sglang, "__version__", None), + "sglang_path": sglang.__file__, + "torch_version": torch.__version__, + "cuda_version": torch.version.cuda, + "device_name": torch.cuda.get_device_name(0), + } + server_args = ServerArgs( + model_path="unused-no-model-weights", + device="cuda", + hicache_ratio=2.0, + hicache_size=0, + hicache_mem_layout="page_first_direct", + hicache_io_backend="direct", + hicache_storage_backend=None, + hicache_write_policy="write_through", + hicache_host_memory_mode="cache", + ) + # Config projection only: unlike publish(), this does not resolve a + # model config or load/download weights. Parallel overrides below + # supply the only topology metadata needed by this single-rank path. + get_context().set_server_args(server_args) + if args.assembler_override: + path = args.assembler_override.resolve(strict=True) + spec = importlib.util.spec_from_file_location("sglang_hybrid_assembler_override", path) + require(spec is not None and spec.loader is not None, "Cannot load assembler override") + assembler = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = assembler + spec.loader.exec_module(assembler) + else: + from sglang.srt.mem_cache.hybrid_cache import hybrid_pool_assembler as assembler + result["assembler_path"] = str(Path(assembler.__file__).resolve()) + result["assembler_sha256"] = hashlib.sha256(Path(assembler.__file__).read_bytes()).hexdigest() + with get_parallel().override( + tp_rank=0, tp_size=1, pp_rank=0, pp_size=1, + attn_tp_rank=0, attn_tp_size=1, attn_cp_rank=0, attn_cp_size=1, + attn_dp_rank=0, attn_dp_size=1, dcp_enabled=False, + attn_dcp_rank=0, attn_dcp_size=1, + ): + layouts = ["raw-fp8"] + (["scaled-mla"] if args.include_scaled_mla else []) + for layout in layouts: + case = {"layout": layout, "status": "error", "states": {}} + result["cases"].append(case) + try: + run_case(torch, assembler, case, server_args) + except Exception as exc: + case.update(status="error", error=f"{type(exc).__name__}: {exc}", traceback=traceback.format_exc()) + # A failed native CUDA operation can leave the context + # poisoned. Do not continue into another layout or pass. + break + statuses = [case["status"] for case in result["cases"]] + result["status"] = "error" if "error" in statuses else "fail" if "fail" in statuses else "pass" + code = {"pass": 0, "fail": 1, "error": 2}[result["status"]] + except Exception as exc: + result.update(status="error", error=f"{type(exc).__name__}: {exc}", traceback=traceback.format_exc()) + finally: + result["elapsed_seconds"] = time.monotonic() - started + done.set() + print(json.dumps(result, sort_keys=True), file=output_stream, flush=True) + return code + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/integration/lmcache/pyproject.toml b/integration/lmcache/pyproject.toml index f100da1..e232a47 100644 --- a/integration/lmcache/pyproject.toml +++ b/integration/lmcache/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "dfkv-connector" -version = "2.26.4" +version = "2.27.0" description = "LMCache RemoteConnector for the dfkv KV cache (ctypes over libdfkv.so)" readme = "README.md" requires-python = ">=3.9" @@ -14,7 +14,7 @@ authors = [{ name = "Wine93", email = "wine93.info@gmail.com" }] # runtime by path (DFKV_LIB / remote_storage_plugin.dfkv.lib). No bundled .so, # no CPython extension, so the wheel is platform-independent. dependencies = [ - "dfkv-common==2.26.4", + "dfkv-common==2.27.0", "lmcache", "torch", ] diff --git a/integration/vllm/patches/vllm-glm53flash-kpool-tail-slot-mapping.patch b/integration/vllm/patches/vllm-glm53flash-kpool-tail-slot-mapping.patch new file mode 100644 index 0000000..e51c3f6 --- /dev/null +++ b/integration/vllm/patches/vllm-glm53flash-kpool-tail-slot-mapping.patch @@ -0,0 +1,21 @@ +--- a/vllm/v1/worker/gpu/model_runner.py ++++ b/vllm/v1/worker/gpu/model_runner.py +@@ -68,6 +68,7 @@ + from vllm.v1.core.sched.output import GrammarOutput, SchedulerOutput + from vllm.v1.kv_cache_interface import ( + CircularBufferSpec, ++ KpoolTailSpec, + KVCacheConfig, + MambaSpec, + UniformTypeKVCacheSpecs, +@@ -578,7 +579,9 @@ + layer_spec = ( + spec.first_spec if isinstance(spec, UniformTypeKVCacheSpecs) else spec + ) +- slot_mapping_enabled.append(not isinstance(layer_spec, CircularBufferSpec)) ++ slot_mapping_enabled.append( ++ not isinstance(layer_spec, (CircularBufferSpec, KpoolTailSpec)) ++ ) + # Let each cache type account for CP. Attention KV is DCP-sharded, + # while Mamba/GDN recurrent state is replicated across DCP ranks. + max_num_blocks = spec.max_num_blocks_per_req( diff --git a/integration/vllm/pyproject.toml b/integration/vllm/pyproject.toml index dfcc959..14ee108 100644 --- a/integration/vllm/pyproject.toml +++ b/integration/vllm/pyproject.toml @@ -4,12 +4,12 @@ build-backend = "setuptools.build_meta" [project] name = "dfkv-vllm" -version = "2.26.4" +version = "2.27.0" description = "Direct vLLM KVConnectorBase_V1 connector for dfkv (GPUDirect RDMA, no LMCache)" requires-python = ">=3.12" # vllm + torch are provided by the runtime image; not pinned here. The connector # itself is pure Python (ctypes over libdfkv.so), so there is no native build. -dependencies = ["dfkv-common==2.26.4"] +dependencies = ["dfkv-common==2.27.0"] # Telemetry is opt-in: the OTel SDK is only needed when DFKV_METRICS_ENABLED / # DFKV_TRACE_ENABLED is set. Without this extra the connector stays dependency-