Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
2a22365
transfer: shm submit/result rings and a node-local shm TE process; pe…
zhuofan1123 Sep 21, 2026
87c02b8
hash: feed numpy buffers to the hasher directly instead of torch.from…
zhuofan1123 Sep 21, 2026
e7c3f59
radixshmem: one global YAML config and the radix-server bootstrap
zhuofan1123 Sep 21, 2026
1a8c098
radixshmem: CPU tier on a radix-server (index engine, planner, SlotSt…
zhuofan1123 Sep 21, 2026
227681d
kvmanager/kvtask: the radixshmem path and node-local DP
zhuofan1123 Sep 21, 2026
4e0717d
integration: sglang node-local DP and DSv4 SWA sizing; vllm adapter f…
zhuofan1123 Sep 21, 2026
fdbc40f
tests/docs: radixshmem engine, data-plane and e2e tests; changelog
zhuofan1123 Sep 21, 2026
f596705
shm_channel: carry the worker-measured durations over the result ring
zhuofan1123 Sep 21, 2026
fa50952
transfer_manager: the shm TE inherits CUDA_VISIBLE_DEVICES unless the…
zhuofan1123 Sep 21, 2026
0a2cbaf
radixshmem: attach to the operator's radix-server; FlexKV brings the …
zhuofan1123 Sep 22, 2026
69ca1af
radixshmem: the RHT registration chunk is the radix-server's --regist…
zhuofan1123 Sep 22, 2026
aa26bb5
radixshmem: harden the shm TE ring, PUT planning and the numpy hash p…
zhuofan1123 Sep 22, 2026
c5ba5ed
radixshmem: drop the shm TE ring; run on FlexKV's own process model (…
zhuofan1123 Sep 22, 2026
d1d7836
radixshmem: configure the server by name (FLEXKV_RADIXSHMEM_SERVER_NA…
zhuofan1123 Sep 22, 2026
fb3367e
radixshmem: retry the attach only while nothing answers; the cluster …
zhuofan1123 Sep 22, 2026
94f6381
radixshmem: say at start-up that cpu_cache_gb does not size the CPU tier
zhuofan1123 Sep 22, 2026
44d38b1
radixshmem: one --transfer-dev with every HCA in the cluster example
zhuofan1123 Sep 22, 2026
4e736de
radixshmem: attach the index at start-up and retry a peer RHT shard t…
zhuofan1123 Sep 22, 2026
d952491
mooncake store: use the public batch_is_exist and read put()'s per-ke…
zhuofan1123 Sep 22, 2026
1f0d88f
radixshmem: a failed cluster-wide flush no longer recycles slots the …
zhuofan1123 Sep 28, 2026
67302ab
radixshmem: follow the shmradix -> radixshmem package rename
zhuofan1123 Oct 9, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -91,3 +91,6 @@ ssd_cache*/
benchmarks/nvcomp_benchmarks/inputs
benchmarks/nvcomp_benchmarks/runs
uv.lock

# mooncake source tree cloned by install.sh --enable-p2p
.mooncake-build/
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Feature

Universal:
- Add radixshmem mode (`FLEXKV_ENABLE_RADIXSHMEM=1`; requires [radixshmem](https://github.com/ai-dynamo/radixshmem) 3608a23 or later): the CPU tier is a per-node `radix-server` run by the operator (`radix-server --name /flexkv --data-bytes 64G [--swa-ratio R] [cluster flags]`) that owns the radix index, the SlotStore and the cross-node RDMA pull. FlexKV attaches with `radixshmem.RadixClient(name, Geometry)`: it hands the server its slot geometry (tokens per block, bytes per CPU block / SWA page, SWA window, slot alignment), adopts the slot counts the server plans from its budget into `CacheConfig` (`shm_radix_bootstrap.adopt_radix_server`; `cpu_cache_gb` has no effect in this mode) and uses the server's SlotStore as the CPU pool in the TE and every transfer worker; cross-node reuse is `RadixClient.pull_async` from the prefetch path whenever the server runs as a cluster. `FLEXKV_RADIXSHMEM_SERVER_NAME` (default `/flexkv`) names the server to attach to; the endpoint is the one radixshmem derives from the name, and the ready wait (600 s) and the prefetch limits are fixed defaults. Process model as in the other modes: engine mode for one DP, the node's KVServer for dp_size > 1 or several `FLEXKV_INSTANCE_NUM` engines. Planning is `RadixShmemCacheEngine` (`flexkv/cache/radix_shmem_planner.py`, a `GlobalCacheEngine` subclass over the `_prepare_request` / `_build_cpu_cache_engine` hooks) on `CacheEngineRadixShmem` (`flexkv/cache/radix_shmem_engine.py`); a plan cancelled before launch or failing during planning returns its slots and match pin. CPU tier only (`enable_ssd`, `enable_remote`, `enable_p2p_*` must be off). Reference: `docs/radixshmem/config_zh.md`, launch scripts in `examples/radixshmem/`.
- `KVServer.create_server(inherit_env=False)` passes `PYTHONPATH`, `LD_LIBRARY_PATH` and `PATH` to the server child along with the `FLEXKV_*` variables.
- `gen_hashes` / `Hasher.update` hash numpy buffers directly (`c_ext.gen_hashes_numpy` / `update_numpy`, safe to call concurrently); `gen_hashes_numpy` validates dtype (int64 tokens, uint64 hashes), contiguity and sizes.

Targeting SGLang:
- The native FlexKV backend is available in upstream SGLang `v0.5.16` and later; no patch is required ([sglang#29701](https://github.com/sgl-project/sglang/pull/29701))
- Add DeepSeek-V4 support for heterogeneous C4/C128/indexer KV groups, FullKV + SWA dual caches, attention/indexer compress-state sidecars, and layerwise restore ([#225](https://github.com/taco-project/FlexKV/pull/225))
Expand Down
54 changes: 54 additions & 0 deletions csrc/bindings.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
#include <cuda_runtime.h>
#include <fcntl.h>
#include <nvtx3/nvToolsExt.h>
#include <pybind11/numpy.h>
#include <pybind11/pybind11.h>
#include <pybind11/stl.h>
#include <sys/mman.h>
Expand Down Expand Up @@ -481,6 +482,48 @@ PYBIND11_MODULE(c_ext, m) {
m.def("gen_hashes", &flexkv::gen_hashes, "Generate hashes for a tensor",
py::arg("hasher"), py::arg("token_ids"), py::arg("tokens_per_block"),
py::arg("block_hashes"));
m.def(
"gen_hashes_numpy",
[](flexkv::Hasher &hasher, py::array token_ids, int tokens_per_block,
py::array block_hashes) {
// numpy-buffer variant of gen_hashes; bypasses torch.from_numpy.
// Same contract as gen_hashes(Tensor): int64 tokens, uint64 hashes,
// both C-contiguous, one hash per whole block. Checked here because a
// raw buffer reinterpreted as int64 would otherwise be read past its
// end (an int32 array is half as long as the loop assumes).
if (tokens_per_block <= 0)
throw py::value_error("gen_hashes_numpy: tokens_per_block must be > 0");
if (!py::isinstance<py::array_t<std::int64_t>>(token_ids))
throw py::type_error(
"gen_hashes_numpy: token_ids must be an int64 array, got dtype " +
py::str(token_ids.dtype()).cast<std::string>());
if (!py::isinstance<py::array_t<std::uint64_t>>(block_hashes))
throw py::type_error(
"gen_hashes_numpy: block_hashes must be a uint64 array, got dtype " +
py::str(block_hashes.dtype()).cast<std::string>());
if (!(token_ids.flags() & py::array::c_style) ||
!(block_hashes.flags() & py::array::c_style))
throw py::value_error(
"gen_hashes_numpy: token_ids and block_hashes must be C-contiguous");
py::buffer_info tok = token_ids.request();
py::buffer_info bh = block_hashes.request(true);
if (bh.size * static_cast<py::ssize_t>(tokens_per_block) > tok.size)
throw py::value_error(
"gen_hashes_numpy: block_hashes has " + std::to_string(bh.size) +
" blocks of " + std::to_string(tokens_per_block) +
" tokens but token_ids holds only " + std::to_string(tok.size) +
" tokens");
const int64_t *tok_ptr = static_cast<const int64_t *>(tok.ptr);
flexkv::HashType *bh_ptr = static_cast<flexkv::HashType *>(bh.ptr);
for (py::ssize_t i = 0; i < bh.size; i++) {
hasher.update(tok_ptr + i * tokens_per_block,
tokens_per_block * sizeof(int64_t));
bh_ptr[i] = hasher.digest();
}
},
"Generate block hashes directly from numpy buffers", py::arg("hasher"),
py::arg("token_ids"), py::arg("tokens_per_block"),
py::arg("block_hashes"));

py::class_<flexkv::SSDIOCTX>(m, "SSDIOCTX")
.def(
Expand Down Expand Up @@ -766,6 +809,17 @@ PYBIND11_MODULE(c_ext, m) {
py::overload_cast<const void *, size_t>(&flexkv::Hasher::update),
"Update the hasher with pointer and size", py::arg("input"),
py::arg("size"))
.def(
"update_numpy",
[](flexkv::Hasher &self, py::array arr) {
// Hash the numpy buffer directly, bypassing torch.from_numpy
// (whose tensor conversion is not concurrency-safe).
py::buffer_info info = arr.request();
self.update(info.ptr,
static_cast<size_t>(info.size * info.itemsize));
},
"Update the hasher directly from a numpy array buffer",
py::arg("input"))
.def("digest", &flexkv::Hasher::digest, "Return the hash value");
#ifdef FLEXKV_ENABLE_CFS
py::class_<flexkv::Pcfs>(m, "Pcfs")
Expand Down
137 changes: 137 additions & 0 deletions docs/radixshmem/config_zh.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
# radixshmem 模式配置

FlexKV 以 radixshmem 作为 CPU 层(索引 + SlotStore + 跨节点拉取)。配置分两处:

| 谁 | 载体 | 内容 |
|---|---|---|
| 运维 | `radix-server` 命令行,每节点一个进程 | 名字、SlotStore 字节预算与 SWA 占比、hugepage、传输引擎、集群成员(etcd、网卡、rank)、索引调优 |
| FlexKV | 两个环境变量 | 是否启用、attach 哪个 server |

FlexKV 只实例化 `radixshmem.RadixClient`:第一个 client 把几何(每 block 的 token 数、一个 CPU block 和一个 SWA page
的字节数、SWA 窗口、slot 对齐)交给 server,server 按 `--data-bytes` 和 `--swa-ratio` 规划各池的 slot 数并发布,
每个 FlexKV 进程 attach 时把 slot 数采纳到 `CacheConfig`(`num_cpu_blocks`、`swa.num_slots`)。CPU 层容量由
`radix-server --data-bytes` 决定,`cpu_cache_gb` 在该模式下不起作用(KVManager 启动时打印一条 WARNING 提示)。

实现:`flexkv/server/shm_radix_bootstrap.py`(几何、attach、采纳、固定参数)。
radixshmem 侧接口见 radixshmem 仓库 `python/README.md`。

## 1. 环境变量

| 变量 | 默认 | 说明 |
|---|---|---|
| `FLEXKV_ENABLE_RADIXSHMEM` | `0` | `1` 启用。在 `flexkv` 首次 import 前设置。 |
| `FLEXKV_RADIXSHMEM_SERVER_NAME` | `/flexkv` | attach 的 `radix-server --name`,即索引 shm 名,以 `/` 开头。gRPC 端点是 radixshmem 由名字派生的 `unix:///dev/shm/<name>.sock`。 |
| `FLEXKV_CPU_LAYOUT` | | 必须是 `BLOCKFIRST`:一个 SlotStore slot 就是一个连续的 CPU block。 |
| `FLEXKV_INSTANCE_NUM` / `FLEXKV_INSTANCE_ID` | `1` / `0` | 同一节点上多个推理引擎共享一个 radix-server 时区分实例(第 4.4 节)。 |

该模式只承担 CPU 层:`enable_ssd`、`enable_remote`、`enable_p2p_cpu`、`enable_p2p_ssd` 必须关闭,启动时校验。
跨节点复用由 radix-server 完成(etcd + RDMA),server 以集群参数启动时自动开启。

进程模型与 FlexKV 其它模式一致:`dp_size=1` 且单实例时 KVTaskEngine 在引擎进程内,TE 是它的子进程;
`dp_size>1` 或多实例时每节点一个 KVServer,DP 进程是它的 client,KVServer 里的 KVTaskEngine 和 TE attach radix-server。

## 2. 固定参数

attach 的其余参数是 `flexkv/server/shm_radix_bootstrap.py` 里的常量:

| 常量 | 值 | 含义 |
|---|---|---|
| `READY_TIMEOUT_S` | 600 | 等 server 可达且 ready 的总时长,覆盖 server 晚起、SlotStore prefault、集群 rendezvous;server 的 `--bootstrap-timeout` 不要超过它。 |
| `PREFETCH_TIMEOUT_MS` | 5000 | 一次 prefetch 拉取的服务端超时,到期后 job 以本地命中的部分完成。 |
| `PREFETCH_MAX_INFLIGHT` | 128 | 每个 KVTaskEngine 在飞的 peer 拉取上限,达到后新的 prefetch 跳过 peer 查询。 |
| `MAX_OUTSTANDING` | 256 | `RadixClient` 未领取 job 的上限。 |

## 3. 几何与 slot 数

FlexKV 交给 server 的几何(`shm_radix_bootstrap.expected_geometry` → `radixshmem.Geometry`):

| 字段 | 来源 |
|---|---|
| `block_size` | `CacheConfig.tokens_per_block` |
| `full_slot_bytes` | 按 `StorageEngine` 的 BLOCKFIRST 布局算出的一个 CPU block 字节数 |
| `swa_slot_bytes` / `swa_window_blocks` | `CacheConfig.swa` 开启时一个 SWA page 的字节数与窗口块数;未开启则没有 SWA 池 |
| `slot_align` | 不超过 4096 且整除每个池 slot 字节数的最大二次幂,使 SlotStore stride 等于 block 字节数 |
| `register_chunk_tokens` | `0`,即 server 的 `--register-chunk-tokens`(默认 4096 token);采纳后按 `tokens // tokens_per_block`(至少 1)换算成 block 数 |

server 的规划:`swa_slots = floor(swa_ratio × data_bytes / swa_stride)`,`full_slots = (data_bytes − SWA 占用) / full_stride`。
任一池为 0、模型有 SWA 而 `--swa-ratio` 为 0,configure 时拒绝,FlexKV 报 `cannot serve FlexKV's geometry`。

**采纳**(`adopt_radix_server`,KVManager 与 KVTaskEngine 各调一次,幂等):`pools.full.num_slots` 写进
`CacheConfig.num_cpu_blocks`,`pools.swa.num_slots` 写进 `CacheConfig.swa.num_slots`,同时记录 `register_chunk_tokens`。
日志形如 `adopted radix-server /flexkv's geometry: FULL 8605 slots ..., SWA 1024 slots ...; RHT registration chunk 4096 tokens = 64 blocks`。

**校验**(`check_geometry`,TE 等 attach 方):server 发布的 `block_size`、各池 `slot_bytes`、SlotStore stride、SWA 窗口
须与自身布局一致,否则报错退出。slot 数是 server 的,不在校验范围。

同一 server 上的多个 client 须带相同的几何(相同模型、page size、SWA 配置);不同的几何被 server 以 `GeometryMismatch`
拒绝,FlexKV 报 `already serves another geometry`。

## 4. 启动

### 4.1 单机

```bash
# 运维,每节点一次;nohup / systemd 皆可。有 SWA 池的模型(DSv4 等)给 --swa-ratio。
radix-server --name /flexkv --data-bytes 64G --swa-ratio 0.5

# 推理引擎侧
export FLEXKV_ENABLE_RADIXSHMEM=1
export FLEXKV_CPU_LAYOUT=BLOCKFIRST
# 不设 FLEXKV_RADIXSHMEM_SERVER_NAME 即 attach /flexkv
```

server 起来后打印 `Waiting for a client's geometry`;FlexKV 的第一个进程 attach 时交出几何,server 建区域后 ready。
server 可晚于引擎启动,FlexKV 在 `ready_timeout_s` 内重试连接。

### 4.2 多机(一个集群)

每个节点各起一个 server,相同的 `--cluster-id` 和 `--registry`;`--rpc-interface`(或 `--rpc-address`)给对端拨入的 IP,
`--node-name` 空时为 `node<ip>`:

```bash
radix-server --name /flexkv --data-bytes 64G --swa-ratio 0.5 \
--expected-min-nodes 4 --num-rht-shards 4 --rht-slots 4 \
--registry etcd://10.0.0.1:2379 --cluster-id prod_a \
--rpc-interface bond0 --index-dev mlx5_bond_0 --gid-idx 3 \
--transfer-dev mlx5_0,mlx5_1,mlx5_2,mlx5_3,mlx5_4,mlx5_5,mlx5_6,mlx5_7 --bootstrap-timeout 600
```

集群一致的几何字段(`block_size`、池集合、每池 `slot_bytes`、SWA 窗口、`slot_align`、`register_chunk_tokens`)由第一个
拿到几何的节点发布到 etcd `radix/<cluster_id>/geometry/<node>`,其余节点采纳;各节点的 slot 数可以不同。FlexKV 侧每个节点
同一个 `FLEXKV_RADIXSHMEM_SERVER_NAME`。

### 4.3 同机多节点(测试)

两个 server 在一台机器上:不同的 `--name`、不同的 `--node-name`、`--rpc-address 127.0.0.1`。两个 FlexKV 进程各自的
`FLEXKV_RADIXSHMEM_SERVER_NAME` 指向自己的 server,并各给一个 `FLEXKV_SERVER_RECV_PORT`。

### 4.4 一节点多引擎共享一个 server

两个独立的推理引擎(各自的 FlexKV、各自的 GPU)attach 同一个 radix-server,互相命中对方存的 KV:

```bash
# 引擎 A # 引擎 B
FLEXKV_INSTANCE_NUM=2 FLEXKV_INSTANCE_ID=0 FLEXKV_INSTANCE_NUM=2 FLEXKV_INSTANCE_ID=1
```

同一个 `FLEXKV_RADIXSHMEM_SERVER_NAME`。`instance_num > 1` 走 server-client 模式:`instance 0` 的 dp0 内嵌本节点的 KVServer(或用
`FLEXKV_SERVER_LAUNCH_MODE=external` 单独启动),它的 TE 等 `instance_num × gpus_per_node` 张 GPU 都注册后 ready,
所以两个引擎都要启动。两边模型 / page size / SWA 配置必须相同(第 3 节)。node-local DP(多机 DP attention)下不支持多实例。

## 5. 命名

| 对象 | 名字 |
|---|---|
| index shm | `--name`;集群模式下 radixshmem 追加 `_<node_name>`,attach 方只需 `--name` |
| SlotStore shm | `<name>_data`(`--data-name` 可改) |
| gRPC socket | `/dev/shm/<name>.sock`,FlexKV 按名字派生;server 端保持默认 `--endpoint` |
| etcd 键空间 | `radix/<cluster_id>/...` |

## 6. 报错含义

- `FLEXKV_RADIXSHMEM_SERVER_NAME` 须以 `/` 开头且无空白。
- `CacheConfig`:`FLEXKV_CPU_LAYOUT != BLOCKFIRST`;打开了 `enable_ssd`、`enable_remote`、`enable_p2p_cpu` 或 `enable_p2p_ssd`;
SWA 开启但 `window_blocks < 1`。
- attach:600 s 内连不上 server 报 `no radix-server named ... reachable ...`(附启动命令);server 有响应但 attach
失败(如 `server is closed`)立即报错;server 一直在等几何或配置失败报 `not ready within ...`(附 server 当时的
`mode` 和 `last_error`);几何冲突见第 3 节。
17 changes: 17 additions & 0 deletions examples/radixshmem/radix_server_multi_node.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
#!/bin/bash
# The operator's radix-server on ONE node of a cluster; run it on every node with
# the same --cluster-id / --registry. Peers dial the IP resolved from
# --rpc-interface (node identity defaults to node<ip>); the index control plane
# uses --index-dev, the KV bytes move over --transfer-dev. The first node that
# receives a geometry publishes it to etcd, the others adopt it; slot counts
# may differ per node (different budgets). FlexKV's ready_timeout_s must cover
# --bootstrap-timeout.
set -eu
radix-server --name "${NAME:-/flexkv}" \
--data-bytes "${DATA_BYTES:-64G}" \
${SWA_RATIO:+--swa-ratio "$SWA_RATIO"} \
--expected-min-nodes "${NODES:-2}" --num-rht-shards "${NODES:-2}" --rht-slots 4 \
--registry "${REGISTRY:-etcd://10.0.0.1:2379}" --cluster-id "${CLUSTER_ID:-flexkv_prod}" \
--rpc-interface "${RPC_INTERFACE:-bond0}" --index-dev "${INDEX_DEV:-mlx5_bond_0}" --gid-idx "${GID_IDX:-3}" \
--transfer-dev "${TRANSFER_DEV:-mlx5_0}" \
--bootstrap-timeout "${BOOTSTRAP_TIMEOUT:-600}"
11 changes: 11 additions & 0 deletions examples/radixshmem/radix_server_single_node.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
#!/bin/bash
# The operator's radix-server for one node, no cluster: nothing about the model
# on the command line. FlexKV's first client brings the geometry (page size,
# bytes per block, SWA page + window); the server plans the slot counts from the
# budget and FlexKV adopts them (docs/radixshmem/config_zh.md section 3).
# --swa-ratio is needed only for models with an SWA pool (DeepSeek-V4 etc.).
set -eu
radix-server --name "${NAME:-/flexkv}" \
--data-bytes "${DATA_BYTES:-64G}" \
${SWA_RATIO:+--swa-ratio "$SWA_RATIO"} \
${HUGEPAGE_PATH:+--hugepage-path "$HUGEPAGE_PATH"}
Loading