Repository navigation
Feat/radixshmem - #300
Draft
zhuofan1123 wants to merge 21 commits into
Draft
Feat/radixshmem#300zhuofan1123 wants to merge 21 commits into
zhuofan1123 wants to merge 21 commits into
Conversation
…r-process op/graph id ranges Every DP rank's cache engine talks to ONE TransferEngine subprocess per node over shared memory instead of each rank owning a TE: * flexkv/transfer/shm_channel.py: fixed-slot submit ring (fragmented pickles) and fixed-width CompletedOp result ring per client, plus a futex control block the TE parks on when idle. * flexkv/transfer/shm_channel_handle.py: TransferManagerShmChannelHandle (client side) and te_shm_main, which brings up the TransferManager, polls all N submit rings, forwards graphs and routes completions back by client. * flexkv/transfer_manager.py: TransferManagerShmTEProcess owns that subprocess; CUDA_VISIBLE_DEVICES is cleared for it only when the deployment spans more than one GPU. * flexkv/common/transfer.py: TransferOp.set_op_id_range and TransferOpGraph.set_graph_id_range give each CE process a disjoint id range so submissions to the shared TE never collide. * flexkv/transfer/transfer_engine.py: no FlexKV peer transfer worker in radixshmem mode; the radix-server pulls peer blocks itself. Reconciled with main during the rebase: * main replaced the id locks with itertools.count(); the ranges are now ``start + counter % size`` on top of that counter, still lock-free. * CompletedOp gained wait_ms/xfer_ms/e2e_ms on main (#297); the result ring does not carry them, so shm-channel completions report 0.0. * KVServer and TransferManagerOnRemote follow main #287 and inherit CUDA_VISIBLE_DEVICES: server-client mode and MPS are still in use, and neither launcher is on the radixshmem path (that path spawns TransferManagerShmTEProcess, which keeps its own total_gpus > 1 rule). The branch's 17ccce7 revert of #287 (radixshmem_rebase_v1.bak-20260921) is not carried over. Rebased from radixshmem_rebase_v1 (68 commits, kept as radixshmem_rebase_v1.bak-20260921) onto main 738ddc1 as one squash-merge, then split by subsystem; this is part 1/7. Co-authored-by: linhu-nv <linhu@nvidia.com> Co-authored-by: Iris Ge <ige@nvidia.com> Co-authored-by: Hao Xu <xuhao269634@163.com>
…_numpy torch.from_numpy is not safe to call concurrently from several planner threads. c_ext gains Hasher.update_numpy and gen_hashes_numpy, which read the pybind11 buffer directly, and flexkv.common.hash_utils uses them. Part 2/7 of the radixshmem rebase (see part 1 for provenance).
* flexkv/common/radixshmem_config.py: RadixShmemConfig, loaded once from FLEXKV_RADIXSHMEM_CONFIG_PATH; server/client/index/distributed sections, index.register_chunk_size defaulting to 4096 / tokens_per_block. * flexkv/common/config.py, flexkv/integration/config.py: FLEXKV_ENABLE_RADIXSHMEM / enable_radixshmem, node-local DP (ModelConfig.local_dp_size, RankInfo.local_dp_client_id) for every tier, and the sglang dp_rank check under radixshmem with dp_size > 1. * flexkv/server/shm_radix_bootstrap.py: create the radix regions and the embedded radix-server subprocess; size the CPU FULL / SWA slot pools from the cache config so the SlotStore stride equals FlexKV's block. * examples/radixshmem_configs/*.yaml, docs/radixshmem/config_zh.md. Part 3/7 of the radixshmem rebase (see part 1 for provenance). Co-authored-by: Hao Xu <xuhao269634@163.com> Co-authored-by: Iris Ge <ige@nvidia.com> Co-authored-by: linhu-nv <linhu@nvidia.com> Co-authored-by: teeebin <your-email@example.com>
…ore data plane) * flexkv/cache/radix_shmem_engine.py: CacheEngineRadixShmem, a RadixClient on the node's shared index; pinned matches (ShmRadixMatch), publish-after- transfer inserts (StagedRadixInsert), FULL and SWA components, peer prefetch via RadixClient.pull_async. * flexkv/cache/radix_shmem_planner.py: RadixShmemCacheEngine, a GlobalCacheEngine subclass that plans GET / PUT / PREFETCH on that tier and hands the prefetch job to KVTaskEngine on a RadixPlanHandle. * flexkv/cache/cache_engine.py: the hooks the subclass needs -- _build_cpu_cache_engine, and the shared GET/PUT prologue _prepare_request returning a RequestWindow. * flexkv/storage/storage_engine.py, flexkv/storage/allocator.py: in radixshmem mode the CPU FULL / SWA pools are views of the radix-server's SlotStore (_attach_radix_pool); workers re-attach them by name through SlotStoreTensorHandle in worker_data. Registered in main's single (PoolEndpoint, device_id) handle registry. Part 4/7 of the radixshmem rebase (see part 1 for provenance). Co-authored-by: Hao Xu <xuhao269634@163.com> Co-authored-by: Iris Ge <ige@nvidia.com> Co-authored-by: linhu-nv <linhu@nvidia.com> Co-authored-by: teeebin <your-email@example.com>
* KVManager: with enable_radixshmem, partition the op/graph id space per DP client, let local DP client 0 bootstrap the radix regions and the single shm TE subprocess, and have every client attach through the shm channel. * KVTaskEngine: pick RadixShmemCacheEngine, complete job-backed PREFETCH tasks from the RadixClient job instead of from graph completion, and skip the legacy cross-node TransferManagerOnRemote when a node-local shm TE replaces it (ModelConfig.local_dp_size set). Part 5/7 of the radixshmem rebase (see part 1 for provenance). Co-authored-by: Hao Xu <xuhao269634@163.com> Co-authored-by: Iris Ge <ige@nvidia.com> Co-authored-by: linhu-nv <linhu@nvidia.com>
…ast paths * sglang connector: node-local DP for every tier, size the host SWA slot from the DSv4 sidecar groups before the TE learns them from the GPU registration, reject a missing dp_rank under radixshmem with dp_size > 1. * vllm adapter: slotted task dataclasses, getattr-based namespace extraction and count_nonzero on the match/put masks. * tests/test_sglang_store_protocol.py: the legacy prefetch-result test now sets _swa_kv_pool on its bare connector like its siblings; prefetch_async reads it since the joint Full+SWA prefetch rule (it failed on the branch too). Part 6/7 of the radixshmem rebase (see part 1 for provenance). Co-authored-by: Hao Xu <xuhao269634@163.com> Co-authored-by: linhu-nv <linhu@nvidia.com>
tests/radixshmem/ holds the engine + planner + peer-pull suite and the two e2e scripts (single node, prefetch across nodes); tests/test_shm_channel.py covers the rings. CHANGELOG entry for the radixshmem CPU tier; ignore the mooncake tree cloned by install.sh. Part 7/7 of the radixshmem rebase (see part 1 for provenance). Co-authored-by: Hao Xu <xuhao269634@163.com> Co-authored-by: linhu-nv <linhu@nvidia.com>
main #297 has every worker report wait_ms / xfer_ms / e2e_ms on the CompletedOp, and KVTaskEngine exports them as the flexkv_py_transfer_{wait,xfer,e2e}_duration_seconds histograms. The shm result ring packed a CompletedOp into a fixed 30 B record without them, so under radixshmem -- where every completion crosses the ring -- the three fields decoded as 0.0 and record_transfer_duration() skipped them as "never timed". The three histograms stayed empty on that path. Append the three values as f64 to the record (30 -> 54 B; still inside the 64 B result slot, and COMPLETED_OP_WIRE_SIZE / the slot-size assert follow the struct). block_results still does not cross the ring, as before. Test: the local round trip now asserts the durations survive, and a new case checks the record fits the default slot and keeps the failed flag alongside them.
… parent is pinned narrower than it serves TransferManagerShmTEProcess cleared CUDA_VISIBLE_DEVICES for the TE subprocess whenever total_gpus > 1, so that a TE spawned by a DP rank pinned to one device (vLLM DP) could still cudaIpcOpenMemHandle from every rank's GPU. With sglang that clearing is wrong: sglang gives all TP workers one namespace (e.g. CUDA_VISIBLE_DEVICES=2,3 with --base-gpu-id 0), the workers register logical ids 0/1 inside it, and the TE, renumbered to the full set, opened those handles from physical GPUs 0 and 1 while the model ran on 2/3. Transfers only succeeded through NVLink peer access -- verified on H20-GPU-24 with DeepSeek-V4-Flash TP2 (flexkv-test/tp2-cvd run C). Decide from the namespace instead of the GPU count: keep the parent's CUDA_VISIBLE_DEVICES when it lists at least as many devices as this node's TE serves (instance_num x gpus_per_node), and clear it only when the parent sees fewer (the per-rank pinned layout, where the registered ids are physical). The single-GPU case keeps its old behaviour (inherit). Same rule main #287 applies to KVServer / TransferManagerOnRemote. tests/test_shm_te_cvd.py covers the policy table.
…geometry and adopts the slot counts
radixshmem e5ce067 turned the radix-server into a process that starts with
nothing model-specific (`radix-server --name /flexkv --data-bytes 64G
[--swa-ratio R] [cluster flags]`) and takes its slot shape from a client's
Geometry; the slot counts are planned from the server's byte budget. FlexKV
follows: it no longer creates, sizes or supervises a radix-server.
* shm_radix_bootstrap: RadixServerProcess, build_radix_server_config and
the name derivation are gone. RadixGeometry is FlexKV's side of the
geometry (tokens per block, bytes per CPU block and SWA page, SWA window,
slot alignment; no counts) with to_shmradix(); attach_radix_client(name,
geometry=...) builds shmradix.RadixClient(name, Geometry), retries while
the server is not reachable yet, then wait_ready(); a server serving
another geometry or one whose budget holds no slot is refused with the
cause. check_geometry compares block size, slot bytes, strides and the
SWA window (not counts); adopt_geometry writes the server's counts into
CacheConfig.num_cpu_blocks / swa.num_slots.
* KVManager: every DP process attaches with the geometry and adopts the
counts before anything sizes a pool; local dp 0 only spawns the shm TE.
distributed_node_id is the server's rank. CacheEngineRadixShmem and the
TE attach the same way (idempotent Configure), the TE also re-checks the
regions. Peer reuse follows the server (world_size > 1) instead of a YAML
flag.
* radixshmem_config: the YAML is `server` (name, endpoint, ready_timeout_s)
and `client`; the former cluster / data / index sections are refused with
a pointer to the radix-server flags. FLEXKV_RADIX_SERVER_LAUNCH_MODE,
FLEXKV_RADIX_NODE_NAME and FLEXKV_RADIX_RPC_ADDRESS are removed. TE
channel names derive from the server name.
* tests: the unit suite starts ServerConfig servers with a budget sized so
the planned counts equal the test's, and brings the geometry through the
engine; new cases cover adoption, the geometry refusal and a server that
starts after the client. The e2e tests start `python -m shmradix.cli`
themselves (per node in the two-node test) and attach through the new
YAML.
* docs/radixshmem/config_zh.md rewritten (server flags vs FlexKV YAML,
geometry hand-off and adoption, multi-engine sharing, migration table);
examples replaced by radixshmem.yaml + radix_server_{single,multi}_node.sh.
The TE adopts the counts again after its own recompute_cache_block_counts
(which sizes from cpu_cache_gb and undid the KVManager's adoption on DSv4,
where layer_groups make the recompute effective).
Verified on H20-GPU-24: 56 unit tests (in-process ServerConfig servers),
the single-node e2e (dp 1 and 2 on real GPUs, radix-server started by the
test), and under sglang with DeepSeek-V4-Flash TP2 the attach, geometry
hand-off, count adoption (FULL 8605 -> 8610, SWA 1024 -> 1052 from
`radix-server --data-bytes 32G --swa-ratio 0.5`), connector cluster query
and TE attach; the request round trip after the TE fix is still to be run
(the node's GPUs were taken by other jobs).
…er-chunk-tokens FlexKV used to bring its own registration chunk: REGISTER_CHUNK_TOKENS = 4096 became index.register_chunk_size (in blocks) on the server config it built. radixshmem f939910 takes the chunk in tokens (Geometry.register_chunk_tokens, radix-server --register-chunk-tokens, default 4096) and publishes it with the geometry, so FlexKV now follows the server instead of dictating a value: - RadixGeometry.register_chunk_tokens (default 0 = the server's) is forwarded verbatim by to_shmradix(); expected_geometry leaves it at 0. - adopt_geometry also takes over the published register_chunk_tokens and its size in FlexKV blocks (radixshmem's rule, tokens // tokens_per_block and at least 1: register_chunk_blocks()); CacheEngineRadixShmem exposes both as attributes; the attach and adopt logs name them. - check_geometry compares a pinned value with the server's and warns when the server's chunk is not a whole number of FlexKV blocks. - tests: the adoption and engine tests run against a server started with a non-default chunk (2048 / 64 tokens) and cover the pinned and mismatch paths; docs/radixshmem/config_zh.md and CHANGELOG updated.
…ath; consistent changelog Review follow-ups on feat/radixshmem (base main 738ddc1): - shm_channel: LAYERWISE has a wire index. _TT_NAMES lacked it, so the first layerwise completion raised KeyError inside the TE's result thread and every rank on the node stopped receiving completions. An unknown name now goes out as None with one warning; a submit record that does not unpickle is dropped with an error instead of being re-read on every poll; a graph the TE cannot submit is reported back to its channel as failed (op_id -1) and a delivery error on one channel no longer stops the others. - radixshmem PUT: CacheEngineRadixShmem.take clamps a request to the pool's size (radixshmem's allocate_slots raises ValueError above it) and answers empty on a refusal; _plan_put returns the taken FULL/SWA slots and releases the match pin when planning fails after the take, then re-raises. Before, a server planning fewer SWA slots than the window leaked FULL slots on every PUT and left the matched prefix pinned for good. - gen_hashes_numpy checks its buffers (int64 tokens, uint64 hashes, both C-contiguous, enough tokens for the requested blocks) and gen_hashes raises TypeError on non-int64 input as the torch path did, instead of reading past an int32 buffer. Hashes of int64 input are unchanged. - KVManager.shutdown works on an instance whose __init__ did not reach the radixshmem setup (_shutdown_radix_shmem_children uses getattr); the CI unit file tests/test_kvmanager_client_api.py failed with AttributeError. - CHANGELOG: the radixshmem entries describe the final state only (operator-run server, server/client YAML, peer reuse from the server's world_size, the shm TE ring, numpy hashing). The stale bullets about the embedded server, the cluster/data/index YAML pass-through, `distributed` and two never-written docs are gone; the requirement is radixshmem f939910 (e5ce067 was amended away and is unreachable). Tests: tests/test_shm_channel.py (every TransferType round-trips, unknown type, poisoned record), tests/radixshmem/test_radix_shmem_engine.py (take clamps, PUT planning failure rolls back), new unit-scope tests/test_hash_utils.py (chain equivalence with Hasher, dtype and buffer checks). 472 passed across the CI unit files and the radixshmem suites.
…engine mode or KVServer) The radixshmem path had a transport of its own: one node-local TE subprocess (TransferManagerShmTEProcess / te_shm_main) that every DP rank fed over shared-memory submit / result rings (flexkv/transfer/shm_channel.py, shm_channel_handle.py), with per-process op / graph id ranges so the submissions could not collide. That is gone. radixshmem mode now runs on the process model FlexKV already has: - dp_size 1, one engine: the KVTaskEngine lives in the engine process with its TE subprocess (TransferManagerInterProcessHandle), as in the default mode. - dp_size > 1 or FLEXKV_INSTANCE_NUM > 1: server-client mode, one KVServer per node (embedded by the first local DP rank, or started on its own with FLEXKV_SERVER_LAUNCH_MODE=external) whose KVTaskEngine and TE attach the radix-server. server_client_mode is decided exactly as on main again. - shm_radix_bootstrap.adopt_radix_server(model_config, cache_config) attaches with FlexKV's geometry, takes the planned slot counts and the cluster rank over into CacheConfig and detaches. The KVManager runs it before it starts a KVServer or a KVTaskEngine, and KVTaskEngine.__init__ runs it again (idempotent), so an externally started KVServer sizes its pools from the server too. - KVServer.create_server(inherit_env=False) also hands PYTHONPATH, LD_LIBRARY_PATH and PATH to the child: it runs the parent's interpreter and has to import what the parent imports (shmradix, when radixshmem is on PYTHONPATH rather than in site-packages); the embedded server used to die on `import shmradix` there. Removed: flexkv/transfer/shm_channel.py and shm_channel_handle.py, TransferManagerShmTEProcess, shm_te_clears_cuda_visible_devices, the TransferManagerHandle mode="shm", KVManager._init_radix_shmem_path / _spawn_shm_te / _shutdown_radix_shmem_children, the shm_te_* KVTaskEngine parameters, the TransferOp / TransferOpGraph id ranges (flexkv/common/transfer.py is main's again), RadixShmemConfig.te_server_id (default_endpoint replaces the socket derivation), tests/test_shm_channel.py and tests/test_shm_te_cvd.py. docs/radixshmem/config_zh.md and the CHANGELOG describe the process model; the e2e test docstrings follow. The ring design stays reachable at aa26bb5 (branch feat/radixshmem.bak-shm-te-ring). Verified: 455 CPU tests (the CI unit files, the radixshmem suite, the sglang store protocol); the GPU e2e tests/radixshmem/test_e2e_radix_shmem.py passes for dp_size 1 (engine mode) and 2 (KVServer) with byte-exact round trips.
…ME); fixed attach defaults FlexKV's side of radixshmem mode is one setting: FLEXKV_RADIXSHMEM_SERVER_NAME (default /flexkv), the --name of the node's radix-server (shm_radix_bootstrap.radix_server_name(); a shm name: starts with '/', no whitespace). The gRPC endpoint is the socket radixshmem derives from the name (unix:///dev/shm/<name>.sock). The ready wait (READY_TIMEOUT_S = 600 s: a late server start, the SlotStore prefault, a cluster rendezvous), the prefetch pull deadline (PREFETCH_TIMEOUT_MS = 5000), the peer pulls in flight per KVTaskEngine (PREFETCH_MAX_INFLIGHT = 128) and the RadixClient job cap (MAX_OUTSTANDING = 256) are constants in flexkv/server/shm_radix_bootstrap.py. attach_radix_client, adopt_radix_server and radix_server_is_distributed take the server name; the planner reads the constants. flexkv/common/radixshmem_config.py (the YAML behind FLEXKV_RADIXSHMEM_CONFIG_PATH) is gone, and so is examples/radixshmem_configs/radixshmem.yaml; the launch scripts live in examples/radixshmem/. Tests: the YAML tests are replaced by tests of the variable, the e2e tests pass the server name through the environment, and the RDMA cluster test gives its two servers distinct names (hence distinct sockets) instead of an endpoint override. docs/radixshmem/config_zh.md and the CHANGELOG describe the variable and the fixed defaults. Verified: 442 CPU tests (the CI unit files, the radixshmem suite, the sglang store protocol); the GPU e2e tests/radixshmem/test_e2e_radix_shmem.py passes for dp_size 1 and 2.
…question brings the geometry attach_radix_client retried every error from the RadixClient constructor as "server not up yet" at DEBUG level for READY_TIMEOUT_S and then reported "no radix-server reachable", which pointed operators at the wrong cause when a live server had refused the attach (INTERNAL: server is closed) or FlexKV had called it wrongly (TypeError). Now only radixshmem's UNAVAILABLE RuntimeError (nothing at the socket) and an RPC TimeoutError are retried; anything else is logged with the label and the socket and raised as is. The timeout message after wait_ready reads the server's state at that moment (client.status()) instead of the constructor-time info. radix_server_is_distributed(model_config, cache_config, ...) attaches with FlexKV's geometry like every other attach, so a server nobody has configured yet is configured by the question and answers from world_size instead of timing out with "nobody handed it a geometry". The sglang connector passes its model and cache config; the call sits after the SWA layer groups are set, so the geometry matches the KVManager's. Tests: attach errors other than UNAVAILABLE come back at once, UNAVAILABLE is retried until the deadline; an unconfigured server answers the cluster question and a geometry-less attach on it times out naming the cause. Verified: 444 CPU tests; GPU e2e dp_size 1 and 2 pass.
In radixshmem mode the CPU tier is the radix-server's SlotStore, sized by its --data-bytes (and --swa-ratio); the slot counts it plans replace num_cpu_blocks / swa.num_slots. KVManager now logs that when it starts: shm_radix_bootstrap.cpu_sizing_notice() gives the level and the text, a WARNING with the configured value when the deployment set cpu_cache_gb (CacheConfig._user_cpu_cache_gb), INFO for a bare default. Unit test for the notice; the doc's overview mentions the warning.
The multi-node radix-server example lists mlx5_0..mlx5_7 in one comma-separated --transfer-dev (radixshmem MR !18 adds the comma form; repeating the flag still works). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
…hat is not up yet attach_radix_client(attach_index=True) brings client.index up before returning instead of on first use. On a cluster that opens the RDMA queue pairs to every peer's RHT shard holder; right after the rendezvous a holder may not accept yet and radixshmem gives up on the connect with "RhtConsumer: failed to connect RHT holder". That error is now retried every 5 s until the ready deadline; any other error is raised at once. CacheEngineRadixShmem and the TE's TransferManager turn it on; callers that only read info (adopting the slot counts, asking whether the server is distributed) leave it off and open no RDMA state. Seen on the two-node cp4x4 bed: node 24's engine attached about 60 s after "Cluster ready (2 nodes)" and failed the XRC connect to node 25's holder while four FlexKV processes attached at once. With the retry the attach went through after 7-8 rounds; a later run needed none.
…y result MooncakeStoreClient.exists() called the SDK's private _batch_exist, which mooncake-transfer-engine 0.3.12 no longer has (batch_is_exist is the public name, and batch_exists_impl already used it). The single-key put() compared the list batch_put_from returns with 0 and so reported every successful put as a failure; check the one element like batch_put does. Neither is on FlexKV's transfer hot path (that uses batch_put / batch_get / batch_exists), found while bringing up the shared Mooncake Store arm of the cp4x4 comparison (4 sglang instances, DSv4-Flash, two nodes).
…tree owns CacheEngineRadixShmem.insert() calls _tree.flush() after _tree.insert() has taken the slots (auto_recycle=True). When the flush raised, e.g. with "DcInitiator::flush WC error status=10 vendor=136" from an RDMA write to a peer's RHT shard, StagedRadixInsert.publish() caught the exception and recycled the same slots into the mempool. The slots then sat in the tree and on the free list at once, so a later PUT could overwrite KV that the tree still served for the old prefix. The flush only makes the blocks routable from other nodes; the local insert has already succeeded. insert() now logs the flush failure as a warning and returns normally, so publish() recycles only when the insert itself fails. test_failed_cluster_publish_keeps_the_slots_in_the_tree makes flush raise and checks that the inserted slots stay matchable and are not handed out again by take(). It fails without this change. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
radixshmem on GitHub (ai-dynamo/radixshmem, 3608a23) renamed its Python package from shmradix to radixshmem; the API is otherwise unchanged. Import radixshmem everywhere, rename RadixGeometry.to_shmradix / _ensure_shmradix to match, start the test server with python -m radixshmem.cli, and require 3608a23 or later in the changelog.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
No description provided.