From 8f631da7634fc7b3d4a1cd95935d85ed01829f6d Mon Sep 17 00:00:00 2001
From: zTz01 <1773266173@qq.com>
Date: Wed, 8 Jul 2026 10:59:46 +0800
Subject: [PATCH 01/10] feat(kv): add Fluxon KV integration for SGLang
---
...06\346\210\220\350\256\276\350\256\241.md" | 381 ++
fluxon_py/kvclient/fluxon.py | 898 ++++-
fluxon_py/kvclient/kvclient_interface.py | 148 +-
.../src/facade/transfer_engine.rs | 18 +-
.../src/lib.rs | 53 +-
fluxon_rs/fluxon_fs/src/agent.rs | 19 +-
.../src/agent_service/transfer_agent.rs | 1662 +++++----
fluxon_rs/fluxon_fs/src/cache_controller.rs | 2 +-
fluxon_rs/fluxon_fs_s3_gateway/src/lib.rs | 5 +-
.../fluxon_kv/src/client_kv_api/delete.rs | 100 +-
.../src/client_kv_api/external_api.rs | 2110 ++++++++++-
fluxon_rs/fluxon_kv/src/client_kv_api/get.rs | 791 +++-
.../client_kv_api/local_reserve_rebalance.rs | 308 ++
fluxon_rs/fluxon_kv/src/client_kv_api/mod.rs | 1407 ++++++-
.../fluxon_kv/src/client_kv_api/msg_pack.rs | 413 +-
fluxon_rs/fluxon_kv/src/client_kv_api/put.rs | 1666 ++++++++-
.../fluxon_kv/src/client_seg_pool/mod.rs | 19 +-
fluxon_rs/fluxon_kv/src/config.rs | 73 +-
.../fluxon_kv/src/external_client_api/mod.rs | 1776 +++++++--
fluxon_rs/fluxon_kv/src/kv_test.rs | 36 +-
fluxon_rs/fluxon_kv/src/lib.rs | 56 +-
.../fluxon_kv/src/master_kv_router/delete.rs | 5 +-
.../fluxon_kv/src/master_kv_router/get.rs | 254 +-
.../fluxon_kv/src/master_kv_router/mod.rs | 699 +++-
.../src/master_kv_router/msg_pack.rs | 477 ++-
.../src/master_kv_router/placement.rs | 6 +-
.../fluxon_kv/src/master_kv_router/put.rs | 1397 ++++++-
.../lease_manager_test.rs | 82 +-
fluxon_rs/fluxon_kv/src/memholder/lifetime.rs | 22 +-
fluxon_rs/fluxon_kv/src/memholder/mod.rs | 112 +-
fluxon_rs/fluxon_kv/src/metrics.rs | 195 +-
fluxon_rs/fluxon_kv/src/observe_kvope.rs | 1 +
.../rpcresp_kvresult_convert/msg_and_error.rs | 33 +
.../rpcresp_kvresult_convert.rs | 251 +-
fluxon_rs/fluxon_ops/build.rs | 20 +-
fluxon_rs/fluxon_ops/src/lib.rs | 433 ++-
fluxon_rs/fluxon_pyo3/src/error.rs | 25 +-
fluxon_rs/fluxon_pyo3/src/lib.rs | 3316 ++++++++++++++++-
fluxon_rs/fluxon_util/src/dev_config.rs | 2 +-
fluxon_rs/fluxon_util/src/lib.rs | 14 +-
fluxon_rs/fluxon_util/src/log.rs | 46 +-
fluxon_rs/fluxon_util/tests/log_mgmt.rs | 12 +-
42 files changed, 17235 insertions(+), 2108 deletions(-)
create mode 100644 "fluxon_doc_cn/design/sglang_fluxon_kv\351\233\206\346\210\220\350\256\276\350\256\241.md"
create mode 100644 fluxon_rs/fluxon_kv/src/client_kv_api/local_reserve_rebalance.rs
diff --git "a/fluxon_doc_cn/design/sglang_fluxon_kv\351\233\206\346\210\220\350\256\276\350\256\241.md" "b/fluxon_doc_cn/design/sglang_fluxon_kv\351\233\206\346\210\220\350\256\276\350\256\241.md"
new file mode 100644
index 0000000..8a0dc6d
--- /dev/null
+++ "b/fluxon_doc_cn/design/sglang_fluxon_kv\351\233\206\346\210\220\350\256\276\350\256\241.md"
@@ -0,0 +1,381 @@
+# SGLang Fluxon KV 集成设计
+
+## 设计目标
+
+本文说明开源仓库中 SGLang HiCache 接入 Fluxon KV 的 hostless 实现设计,聚焦接口契约、状态归属和生命周期边界。本文不是通用 `FlatDict` KV API 的完整说明。
+
+当前调用链主要分为三条主线:
+
+- 写入侧:`local_fast_put_start -> SGLang native write -> local_fast_put_commit`。
+- 读取侧:`get_start -> get_transfer -> SGLang native restore -> release_views`。
+- 放弃读取侧 restore 时:`cancel_get_transfer` 释放 `get_start` 持有的资源。
+
+Fluxon 仍然按 key 保存 opaque value bytes。SGLang 负责 page key、page layout、MHA/MLA/Mamba value 内部格式、GPU KV cache index 和 native kernel 调度;Fluxon 负责 key-version 路由、put 准入与 reservation、owner-local reserve、host value 地址、holder 生命周期和跨 owner 数据面。
+
+这个集成把 SGLang 逻辑上的 L2/L3 缓存落到 Fluxon 统一管理的本机/远端 KV 层中。这样可以用同一套 owner、holder、commit/release 语义管理 KV page,减少传统 L2 host cache 与 L3 backend 分属不同系统时产生的同机重复缓存和生命周期割裂。
+
+常见部署下,一台机器启动一个 Fluxon owner;同机多个 GPU 对应的多个 SGLang worker 进程通过 external/client attach 到同一个 owner shared segment。相比一个 SGLang 进程一个后端 segment、进程间不共享后端 segment 的形态,Fluxon 把同机 host/shared memory、owner local reserve slots 和 holder 生命周期放到同一个 owner 生命周期模型里。这样多个 SGLang worker 可以通过同一个 owner segment 获得本机快速可见性和受控的本地可写内存供给,减少各进程为了各自安全边界重复持有 KV、固定预留 segment 或独立维护 pin/release 状态带来的浪费。
+
+## 范围边界
+
+| 范围 | 当前结论 |
+| --- | --- |
+| SGLang hostless 写入 | 已接入。SGLang 通过 `local_fast_put_start` 取得一批可写 host value 地址,native kernel 写入后再调用 `local_fast_put_commit` 提交。 |
+| SGLang hostless 读取 | 使用 `get_start/get_transfer/release_views`。`get_start` 先计算连续可恢复前缀,`get_transfer` 再把可恢复前缀转换成 readable `plan_ptr`。 |
+| Fluxon value layout | Fluxon 不理解 KV page 内部布局;只按 `key + value_len` 管理连续字节。 |
+| SGLang node 状态 | SGLang 的 `storage_*` 字段是调度层状态,不等同于 Fluxon master route;跨节点复用以 Fluxon commit future 成功为准。 |
+
+## 总体架构
+
+```mermaid
+flowchart LR
+ A["SGLang HiCache radix node"] --> B["HiCacheFluxon backend"]
+ B --> C["FluxonKVCacheStore"]
+ C --> D["fluxon_pyo3::KvClient"]
+ D --> E["Fluxon external / owner"]
+ E --> F["Fluxon master"]
+
+ B --> G["sgl-kernel kvcacheio"]
+ D --> H["plan_ptr blob
magic/count/value_ptrs"]
+ H --> G
+ G --> I["GPU KV cache"]
+ E --> J["owner shared segment
local reserve slots / MemHolder"]
+ J --> H
+```
+
+这里有两个层次:
+
+- KV 语义层:SGLang 传入 page key,Fluxon 对外保存 `key -> value`。
+- hostless 数据层:Fluxon 返回 `plan_ptr`,SGLang native kernel 根据 blob 里的 `value_ptrs[]` 直接执行 GPU/host 数据传输。
+
+从缓存物理层级看,Fluxon 把传统 L2/L3 逻辑抽象落到 local side / remote side:local side 覆盖本机 GPU KV、owner shared segment、owner local reserve slots 和 `MemHolder`;remote side 覆盖跨 owner 或跨机器的数据面。多个 SGLang worker attach 到同一个 owner segment 时,同机 KV bytes 不需要按 worker 进程重复保存在多个后端 segment 中。
+
+`plan_ptr` 只是一轮 backup 或 restore 的短生命周期 carrier。它不能作为 key、缓存地址、跨进程句柄或长期状态保存。
+
+## 公共契约
+
+本节只列 SGLang HiCache hostless 接入 Fluxon KV 时直接依赖的接口。
+
+| 接口 | 层级 | 契约 |
+| --- | --- | --- |
+| `wait_local_segments_ready()` | Fluxon + SGLang 集成 | 返回当前进程可见的 local segment mapping,供 SGLang 做 CUDA host registration。 |
+| `local_fast_put_start(keys, value_len, opts)` | SGLang hostless 写入 | 为一批等长 values 准备可写地址,返回 put `plan_ptr`。 |
+| `local_fast_put_commit(plan_ptr)` | SGLang hostless 写入 | 在 SGLang native kernel 写完 `value_ptrs[]` 后消费 put plan,把对应 slots 提交为 Fluxon KV route,返回 `KvFuture`。 |
+| `put_abort(plan_ptr)` | SGLang hostless 写入 | 在 commit 前释放 put plan、key reservation 和 local reserve slot lease。 |
+| `GetStartResult` | SGLang hostless 读取 | 描述连续命中前缀、可传输长度、atomic group 命中数和第一个 miss 位置。 |
+| `GetStartHandle` | SGLang hostless 读取 | 持有一次 get-start 结果和 backend handle;必须被 `get_transfer` 消费或被 `cancel_get_transfer` 取消。 |
+| `get_start(keys, prefix_best_effort, atomic_group_lens)` | SGLang hostless 读取 | 按 key 顺序计算连续命中的 prefix,并返回 `GetStartHandle`。 |
+| `get_transfer(handle)` | SGLang hostless 读取 | 消费 handle 的可传输前缀,执行必要 transfer,并返回 readable `plan_ptr`。 |
+| `cancel_get_transfer(handle)` | SGLang hostless 读取 | 放弃未 transfer 的 `GetStartHandle`,释放 get-start 期间持有的 owner/external 资源。 |
+| `release_views(plan_ptr)` | SGLang hostless 读取 | 释放 get-transfer 产生的 readable plan,丢弃其持有的 holder 引用。 |
+
+`PutOptionalArgs` 在 SGLang hostless 写入路径中的语义如下:
+
+| 字段 | SGLang 使用方式 |
+| --- | --- |
+| `reject_if_inflight_same_key` | 固定开启,避免同一 page key 并发写回造成重复 inflight put。 |
+| `reject_if_exist_same_key` | 固定开启,SGLang 把重复 key 当作已写回或冲突重试处理。 |
+| `write_through` | 当前配置决定提交策略;调用方显式传入时 Fluxon 必须按字段语义执行。 |
+| `lease_id` | 当前 SGLang hostless 主线不依赖 lease。 |
+
+## Key 与组件命名
+
+SGLang 传给 Fluxon 的 key 必须先经过 backend namespace 处理:
+
+```text
+storage_key = key_prefix + ":" + logical_key
+logical_key = page_hash + optional_component_suffix + config_suffix + optional_extra_backend_tag
+```
+
+规则:
+
+- page hash 是 SGLang prefix 复用和 Fluxon KV 存取的共同语义 ID。
+- `PoolName.KV` 使用默认 component;Mamba 等额外 component 通过 suffix 区分。
+- `config_suffix` 编入模型名、TP/PP 等会影响 page layout 的维度,避免不同运行配置复用同一批 physical values。
+- `extra_backend_tag` 用于同一集群内隔离实验或实例,不改变 Fluxon KV 的值格式。
+
+Fluxon 只看最终 `storage_key`。page 内部如何拆成 K/V layer、MLA tensor 或 Mamba state,由 SGLang kernel 参数解释。
+
+## Segment Registration
+
+hostless 读写依赖 SGLang 进程可访问的 Fluxon owner segment 已经完成 CUDA host registration。
+
+常见部署中,同一台机器上的多个 SGLang worker 连接同一个 Fluxon owner,并映射同一个 owner shared segment。每个 SGLang 进程仍需要在自己的 CUDA context 中完成 host registration;底层内存归属、holder 引用和回收由 owner 统一管理。
+
+```mermaid
+sequenceDiagram
+ participant S as SGLang HiCache
+ participant B as HiCacheFluxon
+ participant P as Fluxon Python store
+ participant R as fluxon_pyo3
+ participant O as owner segment mapping
+ participant C as CUDA runtime
+
+ S->>B: register_mem_pool_host / register_mem_host_pool_v2
+ B->>P: wait_local_segments_ready()
+ P->>R: wait_local_segments_ready()
+ R->>O: wait mapped range
+ O-->>R: segment_label, write_ptr, read_ptr, len, generation
+ R-->>P: segment list
+ P-->>B: dict list
+ B->>C: cudaHostRegister(write_ptr/read_ptr, len)
+```
+
+`wait_local_segments_ready()` 返回的 item 至少包含:
+
+| 字段 | 含义 |
+| --- | --- |
+| `segment_label` | owner 本地一般为 `cpu:0`;external attach owner 时为 `external_owner:0`。 |
+| `write_ptr` | 当前进程可写映射地址。 |
+| `read_ptr` | 当前进程可读映射地址,存在时也可注册。 |
+| `len` | 映射长度。 |
+| `generation` | owner 启动代际,用于拒绝过期 holder 或 mapping。 |
+| `node_id` | segment 所属 Fluxon node。 |
+
+SGLang external-client 模式要求看到 `external_owner:*` segment。注册失败时必须同步报错,不能退回到未注册 host memory 的 direct H2D path。
+
+## Plan Blob ABI
+
+`plan_ptr` 是 Fluxon 返回给 SGLang 的临时句柄,本质上是一段 plan blob 的首地址。SGLang native kernel 通过 `plan_ptr` 找到 blob,再从 blob 里读取本次 batch 对应的 value 地址表。
+
+plan blob 是 Fluxon 在 `local_fast_put_start(...)` 或 `get_transfer(...)` 时创建的一段连续 host memory,格式固定:
+
+```c
+uint64_t magic; // 固定校验值,确认 plan_ptr 指向 Fluxon plan blob
+uint64_t count; // value_ptrs 的数量,也就是本次 batch 的 page 数
+uint64_t value_ptrs[count]; // 每个 page 对应的 Fluxon value 起始地址
+```
+
+如果 SGLang 一次写入或恢复 10 个 page,Fluxon 会创建一个 blob,并返回一个 `plan_ptr`:
+
+```text
+plan_ptr -> blob 起始地址
+
+blob[0] = magic
+blob[1] = 10
+blob[2] = value_ptr_0
+blob[3] = value_ptr_1
+...
+blob[11] = value_ptr_9
+```
+
+`magic` 不是 KV value 地址,只是固定校验值;`value_ptr_0 ... value_ptr_9` 才是 Fluxon 为这些 page 准备的 value 起始地址。它们是当前进程可访问的绝对地址,不是偏移量。
+
+`value_len` 不写入 blob。Fluxon 只负责按 `value_len` 分配每个 value 的连续字节区间,并把起始地址放进 `value_ptrs[]`;每个 value 内部如何切成 K/V、layer、MLA 或 Mamba state,由 SGLang 调用 write/restore kernel 时显式传入 layout 参数。
+
+`plan_ptr` 只在当前进程、当前 batch 生命周期内有效。`local_fast_put_commit(plan_ptr)`、`put_abort(plan_ptr)` 或 `release_views(plan_ptr)` 后,Fluxon 会清理对应 plan,SGLang 不能继续使用这个 `plan_ptr`。
+
+## Backup 时序
+
+hostless backup 的核心约束是:Fluxon 负责 GPU KV cache 之下的本机/远端 KV 层的地址分配、route 提交和生命周期管理,但真正的 KV bytes 由 SGLang native kernel 从 GPU KV cache 写入。因此 Fluxon 不能在收到 key 后立刻发布 KV route;它必须先完成 key reservation、put id 分配和 owner-local reserve slot claim,把稳定可写的 `value_ptrs[]` 通过 `plan_ptr` 返回给 SGLang。SGLang native kernel 写完这些地址后,`local_fast_put_commit` 才能把这些 slots 提交为 resident values,并发布 Fluxon KV route。
+
+这条写路径拆成两个 Fluxon API 阶段:
+
+1. `local_fast_put_start(keys, value_len)`:只做写入准入、key reservation、put id 分配和可写 value 地址准备,返回 `plan_ptr(value_ptrs)`;此时 value bytes 还没有写完,不能作为可读 KV route 暴露。
+2. `local_fast_put_commit(plan_ptr)`:在 SGLang native kernel 完成写入后消费 put plan,提交 slot / transfer / route,并返回 `KvFuture`;只有 future 成功后,SGLang 才能把 node 标记为 `storage_backed`。
+
+写入阶段的拆分主要是为了同时满足两个约束:
+
+- 数据面效率:SGLang 不需要先把 GPU KV page 包装成通用 KV payload 再交给后端,而是直接用 kernel 写入 Fluxon 返回的 value 地址。
+- 可见性安全:`put_start` 阶段只预留地址,不发布 route;避免其它 get 读到尚未写完或尚未 commit 的 value。
+
+```mermaid
+sequenceDiagram
+ participant U as UnifiedRadixCache
+ participant B as HiCacheFluxon
+ participant F as Fluxon store
+ participant K as sgl-kernel
+ participant O as Fluxon owner
+ participant M as Fluxon master
+
+ U->>B: local_fast_put_start(missing_keys, value_len)
+ B->>F: local_fast_put_start(storage_keys, value_len, opts)
+ F->>O: claim local reserve slots / external owner offsets
+ O-->>F: plan_ptr(value_ptrs)
+ B-->>U: plan_ptr
+ U->>K: write_*_to_fluxon_values(plan_ptr, page_indices, layout ptrs)
+ K-->>U: writes queued on CUDA stream
+ U->>U: record local_ready_event
+ U->>B: local_fast_put_commit(plan_ptr) after event ready
+ B->>F: local_fast_put_commit(plan_ptr)
+ F->>O: record precommit visible / transfer_end or put_done
+ O->>M: commit route
+ F-->>B: KvFuture
+ U->>U: scheduler poll future
+```
+
+hostless backup 默认不依赖 put 前 exists 扫描。重复 key 或在途 key 由 `local_fast_put_start` 的 `reject_if_exist_same_key` 和 `reject_if_inflight_same_key` 准入语义处理;SGLang 上层按冲突错误做重试或跳过。
+
+`local_fast_put_start(keys, value_len)` 的要求:
+
+- `keys` 不能为空。
+- `value_len` 必须大于 0,且同一批 keys 共享同一个 value size。
+- SGLang 必须在 `local_fast_put_commit` 前完成 native write;写入失败时必须调用 `put_abort`。
+- `local_fast_put_commit` 只能调用一次;调用后 plan 从 registry 清理,后续只能等待返回的 `KvFuture`。
+
+commit 请求会带上 `len`、`src_offset` 和 target 信息。Fluxon 用这些字段判断本次 value 是否落在当前进程可访问的 owner segment 或 owner-local reserve slot 中;如果需要本地可见索引,`put_done` 会返回 owner 分配的 `local_cache_holder_id`。`MemoryInfo` 和 holder 生命周期在下文说明。
+
+## Restore 时序
+
+hostless restore 的核心约束是:SGLang 只能恢复有序 page keys 的连续前缀,并且不能切开一个 radix node 对应的 atomic group。Fluxon 需要先在本机/远端 KV 层里判断这批 keys 的可恢复边界,再把真正可恢复的部分转换成 SGLang kernel 可读取的 `plan_ptr(value_ptrs)`。因此 restore 被拆成 `get_start` 和 `get_transfer` 两个阶段。
+
+`get_start(keys, prefix_best_effort, atomic_group_lens)` 只做恢复规划:按 key 顺序做 local visible check / owner get start,计算 page 级连续命中前缀 `raw_prefix_hit_len`,再按 `atomic_group_lens` 向下收敛成 `transferable_len`。这个阶段回答“本次最多能恢复多少”,但不要求 SGLang 立刻分配 GPU KV pages,也不暴露 readable plan。
+
+`get_transfer(handle)` 在 SGLang 决定恢复后执行:它消费 `get_start` 返回的 handle,等待或完成必要 transfer,只 materialize `keys[..transferable_len]`,并返回持有 holder 引用的 readable `plan_ptr`。随后 SGLang native kernel 使用 `restore_*_from_fluxon_values(...)`,把 Fluxon value memory 拷回 GPU KV cache。
+
+读取阶段拆成 planning 和 materialization 两步,主要是为了保证:
+
+- prefix 安全:中间 page miss 时,只恢复连续命中的完整前缀,不构造带洞的 GPU KV 状态。
+- atomic group 安全:`transferable_len` 不会切开 radix node group,避免恢复半个 node。
+- 资源效率:SGLang 在知道可恢复边界后再分配 GPU KV pages,避免先分配再发现 miss。
+- 生命周期安全:`get_transfer` 返回的 plan 持有 holder 引用,直到 `release_views(plan_ptr)` 后才释放,保证 kernel restore 期间 value 地址稳定。
+
+```mermaid
+sequenceDiagram
+ participant U as UnifiedRadixCache
+ participant B as HiCacheFluxon
+ participant F as Fluxon store
+ participant E as external client
+ participant O as Fluxon owner / external
+ participant K as sgl-kernel
+
+ U->>B: get_start(page_keys, atomic_group_lens)
+ B->>F: get_start(storage_keys, prefix_best_effort, atomic_group_lens)
+ F->>E: batch_get_start(keys)
+ E->>O: ExternalBatchGetStartReq
+ O-->>E: handle + raw_prefix_hit_len
+ E-->>F: backend handle + raw_prefix_hit_len
+ F->>F: build GetStartResult(raw_prefix_hit_len, transferable_len, ...)
+ F-->>B: GetStartHandle + GetStartResult
+ B-->>U: transferable_len / prefix_hit_groups / first_miss_index
+ alt transferable_len > 0 and caller chooses restore
+ U->>B: get_transfer(handle)
+ B->>F: get_transfer(handle)
+ F->>E: batch_get_transfer(handle, transferable keys)
+ E->>O: ExternalBatchGetTransferReq
+ O-->>E: holders / transfer results
+ E-->>F: plan_ptr(value_ptrs, holders kept alive)
+ F-->>B: plan_ptr
+ B-->>U: plan_ptr
+ U->>K: restore_*_from_fluxon_values(plan_ptr, prefix page indices, layout ptrs)
+ K-->>U: H2D queued on CUDA stream
+ U->>B: release_views(plan_ptr) after restore finalizer
+ else caller gives up restore
+ U->>B: cancel_get_transfer(handle)
+ B->>F: cancel_get_transfer(handle)
+ F->>E: cancel_batch_get_start(handle)
+ end
+```
+
+`GetStartResult` 的关键字段如下:
+
+| 字段 | 含义 |
+| --- | --- |
+| `raw_prefix_hit_len` | 按 key 顺序连续命中的 page 数,未按 atomic group 收敛。 |
+| `transferable_len` | 可以交给 `get_transfer` 的 page 数;它不会切开 atomic group。 |
+| `prefix_hit_groups` | 完整命中的 atomic group 数。 |
+| `first_miss_index` | 第一个 miss page 的 index;全部命中时为 `None`。 |
+| `first_miss_group_index` | 第一个 miss 所在 atomic group;全部命中时为 `None`。 |
+| `all_hit` | `transferable_len == len(keys)`。 |
+
+生命周期规则:
+
+- `get_start` 成功后,调用方必须二选一:`get_transfer(handle)` 或 `cancel_get_transfer(handle)`。
+- `get_transfer(handle)` 成功后,handle 已被消费;后续由 returned `plan_ptr` 和 `release_views(plan_ptr)` 管理。
+- `release_views(plan_ptr)` 必须在 native restore 完成后执行,即使 native restore 失败也要释放。
+- `get_start` 只命中部分前缀时,SGLang 只能恢复 `transferable_len` 覆盖的完整 atomic groups,不能构造半个 atomic group 的 GPU KV 状态。
+- `get_transfer` 返回 miss / KeyNotFound 时,SGLang 必须放弃本次 restore 并执行 rollback。
+
+## Fluxon 本地可见索引与生命周期
+
+Fluxon client/external 侧会维护当前进程可直接访问的 value 索引,以及 get/put plan 持有的 holder 引用。SGLang hostless 路径主要涉及下面几类状态:
+
+| 状态 | 创建入口 | 生命周期 |
+| --- | --- | --- |
+| precommit local visible | `local_fast_put_commit` 开始后,由 owner-local reserve slot 对应的 `MemoryInfo` 记录到 `precommit_local_visible_info` | 表示本次 put 的 value 已在当前进程可读,但后台 put commit 尚未最终完成;成功后转为 `get_cached_info` 中的 committed entry,失败或取消时移除。 |
+| committed local visible info | put commit 成功后,由 SGLang 进程内的 Fluxon external/client 将 `MemoryInfo` 记录到 `get_cached_info` | 保存 key、put version、`holder_id`、`offset`、`len` 和 owner node;后续 `get_start/get_transfer` 可以复用这份 `MemoryInfo` 构造 readable plan。 |
+| get-transfer holder | `get_transfer` 成功后绑定到 readable plan,并由 plan 持有引用 | `release_views(plan_ptr)` 后释放引用;plan 生命周期内 holder 保证对应 value 不被释放。 |
+
+收到 `local_cache_holder_id` 后,SGLang 进程内的 Fluxon external/client 使用 `holder_id`、`offset` 和 `len` 构造 `MemoryInfo`,并记录到自身 `get_cached_info`。`MemoryInfo` 记录当前进程访问该 value 所需的地址信息和释放动作;后续 `get_start` 命中 `get_cached_info` 时,可以直接把这份 `MemoryInfo` 纳入本次 get 结果,`get_transfer` 再把这些 value 地址写入 readable plan。底层内存的回收由 owner route、holder 引用和 owner-local reserve grant 生命周期共同约束。
+
+`precommit_local_visible_info` 只覆盖 commit 进行中的短窗口。put commit 成功并确认 holder 后,会移除 precommit entry 并记录 committed entry;如果 commit 失败,precommit entry 必须被清理,不能作为后续 storage-backed 恢复来源。
+
+## Owner 本地写入预留池
+
+Owner 本地写入预留池是 owner segment 中为 SGLang hostless put 预先划分的 writable slots。`local_fast_put_start` 从这些 slots 中为本次 put 分配地址,并把地址写入 `plan_ptr(value_ptrs)` 返回给 SGLang native kernel。此时 slot 只处于 reserved/prepared 状态,还没有绑定为 Fluxon KV 的正式 `key -> value` route。
+
+这个 pool 是共享 owner segment 上的弹性本地可写内存供给层。它让多个 SGLang worker 都能快速取得受 owner 生命周期管理的 `value_ptrs[]`,同时避免为每个 worker 固定切出长期独占的后端 segment。reserve slot 不足时可以按需求补充 grant,空闲后再按 cooldown 回收。
+
+SGLang native kernel 写完 `value_ptrs[]` 后,`local_fast_put_commit` 才把这些 slots 提交为 resident values,并完成 Fluxon KV route commit。commit 前如果 native write 失败,`put_abort(plan_ptr)` 会释放这些 reserved slots。
+
+这条路径把写入拆成两个阶段:
+
+1. `local_fast_put_start` 完成 key reservation、put id 分配和本地 slot claim,返回 `plan_ptr(value_ptrs)`。
+2. SGLang native kernel 写完 `value_ptrs[]` 后,`local_fast_put_commit` 再把这些 slot 转为 resident values,并完成 Fluxon KV 的 put 提交。
+
+这里的 `local_fast_*` 是 Python/PyO3 public hostless plan API。真正的 value bytes 由 SGLang native kernel 写入 `value_ptrs[]` 指向的地址;Fluxon 在 commit 阶段只消费 put plan、处理必要 transfer 或 direct done,并完成 route commit。
+
+对象含义:
+
+| 对象 | 含义 |
+| --- | --- |
+| grant | owner 侧一次申请的大块本地内存,当前固定为 `512 MiB`。 |
+| slot | grant 内按 `slot_size` 切分的小块;一个 slot 承载一个 Fluxon value。 |
+| slot lease | `local_fast_put_start` 为本次 batch 临时 claim 到的一组 slots;失败或 abort 时必须释放。 |
+| value pointer | slot 的起始地址,会写入 plan blob 的 `value_ptrs[]`,供 SGLang kernel 直接写入。 |
+| resident value | `local_fast_put_commit` 后由 slot 构造出的本地可读 value。 |
+| route | master/owner 确认后的 key 到 value 位置映射;route 成功后该 value 才是全局可见的 KV replica。 |
+
+slot 生命周期:
+
+```text
+Free
+ -> Prepared // local_fast_put_start claim slot
+ -> PendingLocalVisible // local_fast_put_commit 开始,本地 resident value 已记录为 pending visible
+ -> Committed // put_done 成功,route 引用该 slot
+ -> Free // route 和 holder 引用都释放后回收
+```
+
+如果 native write 失败,调用方必须执行 `put_abort(plan_ptr)`,Prepared slots 会回到 Free。`local_fast_put_commit` 成功返回后,slot 是否能释放由 route 引用和 holder 引用共同决定;只要 master/owner route 或 `MemHolder` 仍引用该 slot,底层 grant 就不能释放。
+
+当前容量策略:
+
+| 项 | 当前实现 |
+| --- | --- |
+| grant 物理粒度 | `OWNER_LOCAL_RESERVE_GRANT_QUANTUM_BYTES = 512 * 1024 * 1024` |
+| 最小 slot size | `4 KiB` |
+| slot size 计算 | `max(value_len, 4 KiB).next_power_of_two()` |
+| slot 上限 | `slot_size <= 512 MiB` |
+| refill 触发 | 当前 slot class free slots 不足时登记 pending demand 并唤醒 rebalance actor。 |
+| 默认等待 | soft wait `10 ms`,hard timeout `1 s`。 |
+| shrink 单位 | 整个 grant;不做 live grant compaction。 |
+
+底层物理释放收束在 grant 级别。单个 committed slot 只是 grant 内逻辑索引,不直接拥有释放整块 mmap/registered memory 的权力。
+
+## SGLang Node Storage 状态
+
+SGLang 侧的 node metadata 不等价于 Fluxon master route 状态。当前四个字段建议按下面语义解释:
+
+| 字段 | true 的含义 | 清理时机 |
+| --- | --- | --- |
+| `storage_staged` | 该 node 有一批 Fluxon hostless backup 正在 staged 路径中。 | `KvFuture` 完成或失败后清空。 |
+| `storage_local_ready` | CUDA write 已完成,SGLang 已调用 `local_fast_put_commit`,但返回的 `KvFuture` 还未 ack;该状态只表示本次 hostless backup 已进入 Fluxon commit 流程,不表示 KV route 已经全局确认。 | async ack 结束后清空。 |
+| `storage_pending` | Fluxon `KvFuture` 还没结束。 | future 成功或失败后清空。 |
+| `storage_backed` | Fluxon 后台提交成功,KV route 已确认可作为 shared backing。 | 该 node 被删除或失效时清空。 |
+
+因此,SGLang 可以用 `storage_staged/storage_local_ready` 判断本次 hostless backup 已经推进到本地写入或 commit 阶段;但跨节点复用和长期共享必须等 `KvFuture` 成功,并以 `storage_backed` 为准。
+
+TP 场景下,每个 rank 仍有各自的 radix tree 和恢复决策。`get_start` 的结果只描述当前 rank 这批 keys 的可恢复前缀;如果一个 rank miss、另一个 rank hit,上层必须按 SGLang 的 TP restore 约束处理一致性,不能把单 rank 的部分成功当作完整 request 已恢复。
+
+## 失败处理
+
+| 场景 | 必须动作 |
+| --- | --- |
+| `local_fast_put_start` 后 native write 失败 | 调用 `put_abort(plan_ptr)`,释放 key reservation 和 local reserve slot lease。 |
+| `local_fast_put_commit` 返回 future 后后台失败 | SGLang 清理 `storage_staged/storage_pending/storage_local_ready`,必要时删除已 evicted 的 dead leaf。 |
+| `get_start` 后放弃 restore | 调用 `cancel_get_transfer(handle)`,释放 get-start 持有的 owner/external 资源。 |
+| `get_start` 只命中部分前缀 | 只允许恢复 `transferable_len` 覆盖的完整 atomic groups;后续 page 按 miss 处理。 |
+| `get_transfer` 报 key miss | SGLang rollback 当前 restore,不继续构造半个 node 的 GPU 恢复。 |
+| `get_transfer` 成功后 native restore 失败 | 先 `release_views(plan_ptr)`,再执行 SGLang rollback;此时 handle 已被消费。 |
+| CUDA host registration 失败 | direct path 同步失败,不能降级为未注册 host memory。 |
+| `plan_ptr` 类型用错 | `local_fast_put_commit`、`put_abort`、`release_views` 都按 registry entry 类型校验并 fail fast。 |
diff --git a/fluxon_py/kvclient/fluxon.py b/fluxon_py/kvclient/fluxon.py
index 1325e3d..ae9df9a 100644
--- a/fluxon_py/kvclient/fluxon.py
+++ b/fluxon_py/kvclient/fluxon.py
@@ -3,6 +3,7 @@
This module provides a concrete implementation using the PyO3 Rust bindings.
"""
+from dataclasses import dataclass
from typing import Union, Optional, Callable, Any, Dict, List, Tuple
import ctypes
import os
@@ -25,7 +26,7 @@
from .kvclient_interface import KvClient
from .kvclient_interface import KvLeaseApi, KvRpcApi, PutOptionalArgs, FlatDict
from .backend_fallback_close import unregister_store_from_cleanup
-from .kvclient_interface import KvFuture, MemHolder
+from .kvclient_interface import GetStartHandle, GetStartResult, KvFuture, MemHolder
from .nonzerocopy_encode import (
DLPacked,
INTERNAL_DLPACK_META_KEY,
@@ -138,11 +139,114 @@ def _resolve_side_transfer_worker_python() -> str:
return sys.executable
+@dataclass(frozen=True)
+class _RegisteredBufferDescriptor:
+ ptr: int
+ size: int
+ device_kind: str = "host"
+ device_id: str = ""
+ layout: str = "raw"
+ metadata: Optional[Dict[str, Any]] = None
+
+ @property
+ def end(self) -> int:
+ return self.ptr + self.size
+
+ def contains(self, ptr: int, size: int) -> bool:
+ req_end = ptr + size
+ return ptr >= self.ptr and req_end <= self.end
+
+ def as_dict(self) -> Dict[str, Any]:
+ return {
+ "ptr": self.ptr,
+ "size": self.size,
+ "device_kind": self.device_kind,
+ "device_id": self.device_id,
+ "layout": self.layout,
+ "metadata": dict(self.metadata or {}),
+ }
+
+
def _map_nospace_to_storagefull(err: ApiError) -> ApiError:
"""Normalize storage-capacity errors without depending on backend internals."""
return err
+def _get_start_prefix_hit_groups(
+ raw_prefix_hit_len: int,
+ group_lens: Tuple[int, ...],
+) -> int:
+ prefix_hit_groups = 0
+ transferable_len = 0
+ for group_len in group_lens:
+ next_transferable_len = transferable_len + group_len
+ if next_transferable_len > raw_prefix_hit_len:
+ break
+ transferable_len = next_transferable_len
+ prefix_hit_groups += 1
+ return prefix_hit_groups
+
+
+def _get_start_group_index_for_key_index(
+ group_lens: Tuple[int, ...],
+ key_index: int,
+) -> Optional[int]:
+ cursor = 0
+ for group_index, group_len in enumerate(group_lens):
+ cursor += group_len
+ if key_index < cursor:
+ return group_index
+ return None
+
+
+def _build_get_start_result_from_backend_payload(
+ payload: Dict[str, Any],
+ keys: List[str],
+ prefix_best_effort: bool,
+ normalized_group_lens: Optional[List[int]],
+) -> GetStartResult:
+ result_keys = tuple(keys)
+ group_lens = (
+ (len(result_keys),)
+ if normalized_group_lens is None
+ else tuple(normalized_group_lens)
+ )
+ raw_prefix_hit_len = int(payload["raw_prefix_hit_len"])
+ if raw_prefix_hit_len < 0 or raw_prefix_hit_len > len(result_keys):
+ raise RuntimeError(
+ "get_start returned invalid raw_prefix_hit_len: "
+ f"raw_prefix_hit_len={raw_prefix_hit_len} keys={len(result_keys)}"
+ )
+ prefix_hit_groups = _get_start_prefix_hit_groups(
+ raw_prefix_hit_len,
+ group_lens,
+ )
+ if not prefix_best_effort and prefix_hit_groups != len(group_lens):
+ transferable_len = 0
+ prefix_hit_groups = 0
+ else:
+ transferable_len = sum(group_lens[:prefix_hit_groups])
+ first_miss_index = (
+ None if raw_prefix_hit_len == len(result_keys) else raw_prefix_hit_len
+ )
+ first_miss_group_index = (
+ None
+ if first_miss_index is None
+ else _get_start_group_index_for_key_index(group_lens, first_miss_index)
+ )
+ return GetStartResult(
+ keys=result_keys,
+ raw_prefix_hit_len=raw_prefix_hit_len,
+ transferable_len=transferable_len,
+ prefix_hit_groups=prefix_hit_groups,
+ atomic_group_lens=tuple(normalized_group_lens) if normalized_group_lens is not None else None,
+ prefix_best_effort=prefix_best_effort,
+ first_miss_index=first_miss_index,
+ first_miss_group_index=first_miss_group_index,
+ all_hit=transferable_len == len(result_keys),
+ )
+
+
def _error_to_ret_code(err: ApiError) -> int:
if hasattr(err, "code") and callable(err.code):
try:
@@ -299,7 +403,11 @@ def __init__(self, config: FluxonKvClientConfig):
self._client: Optional[fluxon_pyo3.KvClient] = None
self._config = config
self._init_error: Optional[ApiError] = None
+ self._registered_buffer_descriptors: List[_RegisteredBufferDescriptor] = []
cluster_name = config.fluxonkv_spec_cluster_name
+ self._batch_concurrency = 128
+ if self._batch_concurrency <= 0:
+ raise ValueError("batch_concurrency must be > 0")
self._blocking_put_outer_total_log_window = _BlockingPutOuterTotalLogWindow(
f"FluxonKVCacheStore[{cluster_name}]"
)
@@ -368,11 +476,17 @@ def put(
reject_if_inflight_same_key = (
bool(opts.reject_if_inflight_same_key) if opts is not None else False
)
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
inner_res = self._client.put(
key,
ptrs,
lease_id=lease_id,
reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
)
if not inner_res.is_ok():
err = inner_res.unwrap_error()
@@ -416,11 +530,17 @@ def put_blocking(
reject_if_inflight_same_key = (
bool(opts.reject_if_inflight_same_key) if opts is not None else False
)
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
inner_res = self._client.put_blocking(
key,
ptrs,
lease_id=lease_id,
reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
)
if not inner_res.is_ok():
return Result.new_error(inner_res.unwrap_error())
@@ -448,6 +568,199 @@ def get_blocking(self, key: str) -> Result[MemHolder, ApiError]:
except ApiError as e:
return Result.new_error(e)
+ @staticmethod
+ def _normalize_batch_result_list(batch_result: Any, expected_len: int, op_name: str) -> List[Any]:
+ if isinstance(batch_result, Result):
+ if not batch_result.is_ok():
+ raise RuntimeError(f"{op_name} backend error: {batch_result.unwrap_error()}")
+ batch_result = batch_result.unwrap()
+ if not isinstance(batch_result, list):
+ raise RuntimeError(f"{op_name} returned non-list: {type(batch_result)}")
+ if len(batch_result) != expected_len:
+ raise RuntimeError(
+ f"{op_name} returned unexpected length: expected={expected_len} got={len(batch_result)}"
+ )
+ return list(batch_result)
+
+ def batch_put_blocking(
+ self,
+ keys: List[str],
+ values: List[FlatDict],
+ opts: Optional[PutOptionalArgs] = None,
+ concurrency: Optional[int] = None,
+ ) -> List[Result[OkNone, ApiError]]:
+ if len(keys) != len(values):
+ raise ValueError("batch_put_blocking requires keys and values to have the same length")
+ if len(keys) == 0:
+ return []
+ if self._client is None:
+ err = GeneralError(message="Store not initialized when batch_put_blocking(). Call setup() first.")
+ return [Result.new_error(err) for _ in keys]
+
+ lease_id: Optional[int] = opts.lease_id if opts is not None else None
+ reject_if_inflight_same_key = (
+ bool(opts.reject_if_inflight_same_key) if opts is not None else False
+ )
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
+
+ keepalive_groups: List[List[bytes]] = []
+ dlpack_groups: List[List[object]] = []
+ ptr_groups: List[List[tuple[int, int, int, int, int, Optional[int]]]] = []
+ try:
+ for value in values:
+ keepalive: List[bytes] = []
+ dlpack_capsules: List[object] = []
+ ptr_groups.append(build_flat_dict_ptrs(value, keepalive, dlpack_capsules))
+ keepalive_groups.append(keepalive)
+ dlpack_groups.append(dlpack_capsules)
+
+ inner_res = self._client.batch_put_blocking(
+ keys,
+ ptr_groups,
+ lease_id=lease_id,
+ reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
+ concurrency=concurrency if concurrency is not None else self._batch_concurrency,
+ )
+ if not inner_res.is_ok():
+ err = inner_res.unwrap_error()
+ return [Result.new_error(err) for _ in keys]
+
+ batch_results = self._normalize_batch_result_list(
+ inner_res.unwrap(), len(keys), "batch_put_blocking"
+ )
+ submit_out: List[Result[OkNone, ApiError]] = []
+ for idx, item in enumerate(batch_results):
+ if isinstance(item, Result):
+ submit_out.append(item)
+ continue
+ if item is None:
+ submit_out.append(Result.new_ok(OkNone()))
+ continue
+ if isinstance(item, int) and item == 0:
+ submit_out.append(Result.new_ok(OkNone()))
+ continue
+ if isinstance(item, int):
+ submit_out.append(
+ Result.new_error(
+ GeneralError(
+ message=(
+ "batch_put_blocking returned backend code "
+ f"{item} for key {keys[idx]!r}"
+ )
+ )
+ )
+ )
+ continue
+ submit_out.append(
+ Result.new_error(
+ GeneralError(
+ message=f"unexpected batch_put result type: {type(item)}"
+ )
+ )
+ )
+ return submit_out
+ except ApiError as e:
+ return [Result.new_error(e) for _ in keys]
+ except Exception as e:
+ return [Result.new_error(GeneralError(f"batch_put_blocking failed: {e}")) for _ in keys]
+ finally:
+ for keepalive in keepalive_groups:
+ keepalive.clear()
+ for dlpack_capsules in dlpack_groups:
+ dlpack_capsules.clear()
+
+ def batch_get_blocking(
+ self,
+ keys: List[str],
+ concurrency: Optional[int] = None,
+ ) -> List[Result[Union[Any, MemHolder], ApiError]]:
+ if len(keys) == 0:
+ return []
+ if self._client is None:
+ err = GeneralError(message="Store not initialized when batch_get_blocking(). Call setup() first.")
+ return [Result.new_error(err) for _ in keys]
+
+ try:
+ inner_res = self._client.batch_get_blocking(
+ keys,
+ concurrency=concurrency if concurrency is not None else self._batch_concurrency,
+ )
+ if not inner_res.is_ok():
+ err = inner_res.unwrap_error()
+ return [Result.new_error(err) for _ in keys]
+
+ batch_results = self._normalize_batch_result_list(
+ inner_res.unwrap(), len(keys), "batch_get_blocking"
+ )
+ out: List[Result[Union[Any, MemHolder], ApiError]] = []
+ for idx, item in enumerate(batch_results):
+ if isinstance(item, Result):
+ out.append(item)
+ continue
+ if item is None:
+ out.append(
+ Result.new_error(
+ GeneralError(message=f"batch_get_blocking returned None for key {keys[idx]!r}")
+ )
+ )
+ continue
+ out.append(Result.new_ok(item))
+ return out
+ except ApiError as e:
+ return [Result.new_error(e) for _ in keys]
+ except Exception as e:
+ return [Result.new_error(GeneralError(f"batch_get_blocking failed: {e}")) for _ in keys]
+
+ def put_payload_from_ptr(
+ self,
+ key: str,
+ payload_ptr: int,
+ payload_size: int,
+ opts: Optional[PutOptionalArgs] = None,
+ ) -> Result[KvFuture, ApiError]:
+ if self._client is None:
+ return Result.new_error(
+ GeneralError(
+ message="Store not initialized when put_payload_from_ptr(). Call setup() first."
+ )
+ )
+
+ keepalive: List[bytes] = []
+ try:
+ ptrs = _build_payload_field_ptrs(payload_ptr, payload_size, keepalive)
+ lease_id: Optional[int] = opts.lease_id if opts is not None else None
+ reject_if_inflight_same_key = (
+ bool(opts.reject_if_inflight_same_key) if opts is not None else False
+ )
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
+ inner_res = self._client.put(
+ key,
+ ptrs,
+ lease_id=lease_id,
+ reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
+ )
+ if not inner_res.is_ok():
+ return Result.new_error(_map_nospace_to_storagefull(inner_res.unwrap_error()))
+ inner_future = inner_res.unwrap()
+ assert inner_future is not None
+ outer_future = _FluxonPutFuture(inner_future, keepalive, [])
+ keepalive = []
+ return Result.new_ok(outer_future)
+ except ApiError as e:
+ return Result.new_error(e)
+ finally:
+ keepalive.clear()
+
def put_payload_from_ptr_blocking(
self,
key: str,
@@ -469,11 +782,17 @@ def put_payload_from_ptr_blocking(
reject_if_inflight_same_key = (
bool(opts.reject_if_inflight_same_key) if opts is not None else False
)
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
inner_res = self._client.put_blocking(
key,
ptrs,
lease_id=lease_id,
reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
)
if not inner_res.is_ok():
return Result.new_error(inner_res.unwrap_error())
@@ -566,17 +885,297 @@ def get_payload_into_ptr_blocking(
del holder
return Result.new_ok(payload_size)
+ def register_buffer(
+ self,
+ ptr: int,
+ size: int,
+ device_kind: str = "host",
+ device_id: str = "",
+ layout: str = "raw",
+ metadata: Optional[Dict[str, Any]] = None,
+ ) -> Result[OkNone, ApiError]:
+ if self._client is None:
+ return Result.new_error(
+ GeneralError(
+ message="Store not initialized when register_buffer(). Call setup() first."
+ )
+ )
+ if not isinstance(ptr, int):
+ return Result.new_error(
+ InvalidArgumentError(message=f"ptr must be int; got {type(ptr)}")
+ )
+ if not isinstance(size, int):
+ return Result.new_error(
+ InvalidArgumentError(message=f"size must be int; got {type(size)}")
+ )
+ if ptr < 0:
+ return Result.new_error(
+ InvalidArgumentError(message=f"ptr must be >= 0; got {ptr}")
+ )
+ if size < 0:
+ return Result.new_error(
+ InvalidArgumentError(message=f"size must be >= 0; got {size}")
+ )
+ if not isinstance(device_kind, str):
+ return Result.new_error(
+ InvalidArgumentError(
+ message=f"device_kind must be str; got {type(device_kind)}"
+ )
+ )
+ if not isinstance(device_id, str):
+ return Result.new_error(
+ InvalidArgumentError(
+ message=f"device_id must be str; got {type(device_id)}"
+ )
+ )
+ if not isinstance(layout, str):
+ return Result.new_error(
+ InvalidArgumentError(
+ message=f"layout must be str; got {type(layout)}"
+ )
+ )
+ if metadata is not None and not isinstance(metadata, dict):
+ return Result.new_error(
+ InvalidArgumentError(
+ message=f"metadata must be dict or None; got {type(metadata)}"
+ )
+ )
+ try:
+ inner_res = self._client.register_buffer(ptr, size)
+ if not inner_res.is_ok():
+ return Result.new_error(inner_res.unwrap_error())
+ _ = inner_res.unwrap()
+ self._registered_buffer_descriptors.append(
+ _RegisteredBufferDescriptor(
+ ptr=int(ptr),
+ size=int(size),
+ device_kind=device_kind,
+ device_id=device_id,
+ layout=layout,
+ metadata=dict(metadata or {}),
+ )
+ )
+ return Result.new_ok(OkNone())
+ except ApiError as e:
+ return Result.new_error(e)
+
+ def batch_put_from(
+ self,
+ keys: List[str],
+ payload_ptrs: List[int],
+ payload_sizes: List[int],
+ opts: Optional[PutOptionalArgs] = None,
+ ) -> List[int]:
+ if len(keys) != len(payload_ptrs) or len(keys) != len(payload_sizes):
+ raise ValueError(
+ "batch_put_from requires keys, payload_ptrs, and payload_sizes to have the same length"
+ )
+ if len(keys) == 0:
+ return []
+
+ if self._client is None:
+ code = _error_to_ret_code(
+ GeneralError(
+ message="Store not initialized when batch_put_from(). Call setup() first."
+ )
+ )
+ return [code] * len(keys)
+
+ lease_id: Optional[int] = opts.lease_id if opts is not None else None
+ reject_if_inflight_same_key = (
+ bool(opts.reject_if_inflight_same_key) if opts is not None else False
+ )
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
+
+ try:
+ inner_res = self._client.batch_put_from(
+ keys,
+ payload_ptrs,
+ payload_sizes,
+ lease_id=lease_id,
+ reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
+ )
+ if inner_res.is_ok():
+ submit_results = list(inner_res.unwrap())
+ if len(submit_results) != len(keys):
+ raise RuntimeError(
+ "batch_put_from returned unexpected result length: "
+ f"expected={len(keys)} got={len(submit_results)}"
+ )
+ return [int(item) for item in submit_results]
+ err = inner_res.unwrap_error()
+ code = _error_to_ret_code(err)
+ return [code] * len(keys)
+ except ApiError as e:
+ code = _error_to_ret_code(e)
+ return [code] * len(keys)
+ except Exception as e:
+ code = _error_to_ret_code(GeneralError(f"batch_put_from failed: {e}"))
+ return [code] * len(keys)
+
def get_size(self, key: str) -> Result[int, ApiError]:
"""Get the size of a stored value (non-blocking)."""
return self._client.get_size(key)
- def is_exist(self, key: str) -> Result[bool, ApiError]:
+ def is_exist(self, key: str, allow_local_snapshot: bool = False) -> Result[bool, ApiError]:
"""Check if a key exists in the store (non-blocking)."""
try:
- return self._client.is_exist(key)
+ if self._client is None:
+ return Result.new_error(
+ GeneralError(
+ message="Store not initialized when is_exist(). Call setup() first."
+ )
+ )
+ return self._client.is_exist(key, allow_local_snapshot=allow_local_snapshot)
except Exception as e:
return Result.new_error(GeneralError(f"Existence check failed: {str(e)}"))
+ def batch_get_into(
+ self,
+ keys: List[str],
+ payload_ptrs: List[int],
+ payload_capacities: List[int],
+ ) -> List[int]:
+ if len(keys) != len(payload_ptrs) or len(keys) != len(payload_capacities):
+ raise ValueError(
+ "batch_get_into requires keys, payload_ptrs, and payload_capacities to have the same length"
+ )
+ if len(keys) == 0:
+ return []
+
+ if self._client is not None:
+ inner_res = self._client.batch_get_into(keys, payload_ptrs, payload_capacities)
+ if inner_res.is_ok():
+ return list(inner_res.unwrap())
+ err = inner_res.unwrap_error()
+ return [_error_to_ret_code(err)] * len(keys)
+
+ results: List[int] = []
+ for key, ptr, size in zip(keys, payload_ptrs, payload_capacities):
+ get_result = self.get_payload_into_ptr_blocking(key, ptr, size)
+ if get_result.is_ok():
+ results.append(int(get_result.unwrap()))
+ else:
+ results.append(_error_to_ret_code(get_result.unwrap_error()))
+ return results
+
+ def batch_is_exist(
+ self,
+ keys: List[str],
+ allow_local_snapshot: bool = False,
+ ) -> List[int]:
+ if len(keys) == 0:
+ return []
+
+ if self._client is None:
+ code = _error_to_ret_code(
+ GeneralError(
+ message="Store not initialized when batch_is_exist(). Call setup() first."
+ )
+ )
+ return [code] * len(keys)
+
+ try:
+ inner_res = self._client.batch_is_exist(
+ keys,
+ allow_local_snapshot=allow_local_snapshot,
+ )
+ if not inner_res.is_ok():
+ code = _error_to_ret_code(inner_res.unwrap_error())
+ return [code] * len(keys)
+ batch_results = self._normalize_batch_result_list(
+ inner_res.unwrap(), len(keys), "batch_is_exist"
+ )
+ out: List[int] = []
+ for idx, item in enumerate(batch_results):
+ if not isinstance(item, int):
+ raise GeneralError(
+ message=(
+ f"batch_is_exist returned non-int item for key {keys[idx]!r}: "
+ f"{type(item)}"
+ )
+ )
+ out.append(int(item))
+ return out
+ except ApiError as e:
+ code = _error_to_ret_code(e)
+ return [code] * len(keys)
+ except Exception as e:
+ code = _error_to_ret_code(GeneralError(f"batch_is_exist failed: {e}"))
+ return [code] * len(keys)
+
+ def get_meta(self, key: str) -> Result[Dict[str, Any], ApiError]:
+ """Query key metadata and one live placement without fetching payload bytes."""
+ try:
+ inner = self._client.get_meta(key)
+ if not inner.is_ok():
+ return Result.new_error(inner.unwrap_error())
+ meta = inner.unwrap()
+ assert isinstance(meta, dict), f"get_meta returned non-dict: {type(meta)}"
+ return Result.new_ok(meta)
+ except Exception as e:
+ return Result.new_error(GeneralError(f"GetMeta failed for key '{key}': {str(e)}"))
+
+ def batch_get_meta(self, keys: List[str]) -> List[Dict[str, Any]]:
+ if len(keys) == 0:
+ return []
+
+ if self._client is not None and hasattr(self._client, "batch_get_meta"):
+ inner_res = self._client.batch_get_meta(keys)
+ if inner_res.is_ok():
+ rows = inner_res.unwrap()
+ assert isinstance(rows, list), (
+ f"batch_get_meta returned non-list: {type(rows)}"
+ )
+ return list(rows)
+ err = inner_res.unwrap_error()
+ code = _error_to_ret_code(err)
+ return [
+ {
+ "exists": False,
+ "len": 0,
+ "node_id": "",
+ "src_addr": 0,
+ "src_base_addr": 0,
+ "segment_device_id": "",
+ "segment_device_desc": "",
+ "replica_count": 0,
+ "transport_error": True,
+ "error_code": -code,
+ "error_json": str(err),
+ }
+ for _ in keys
+ ]
+
+ rows: List[Dict[str, Any]] = []
+ for key in keys:
+ meta_res = self.get_meta(key)
+ if meta_res.is_ok():
+ rows.append(meta_res.unwrap())
+ else:
+ err = meta_res.unwrap_error()
+ rows.append(
+ {
+ "exists": False,
+ "len": 0,
+ "node_id": "",
+ "src_addr": 0,
+ "src_base_addr": 0,
+ "segment_device_id": "",
+ "segment_device_desc": "",
+ "replica_count": 0,
+ "transport_error": True,
+ "error_code": -_error_to_ret_code(err),
+ "error_json": str(err),
+ }
+ )
+ return rows
+
def count_prefix(self, prefix: str) -> Result[int, ApiError]:
"""Count number of keys with the given prefix.
@@ -806,6 +1405,247 @@ def get_cluster_name(self) -> str:
raise RuntimeError("Store not initialized")
return str(self._client.cluster_name())
+ def wait_local_segments_ready(self) -> List[dict[str, Any]]:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when wait_local_segments_ready(). Call setup() first."
+ )
+ inner_res = self._client.wait_local_segments_ready()
+ if not inner_res.is_ok():
+ raise RuntimeError(
+ f"wait_local_segments_ready backend error: {inner_res.unwrap_error()}"
+ )
+ inner_res = inner_res.unwrap()
+ if not isinstance(inner_res, list):
+ raise RuntimeError(
+ "wait_local_segments_ready must return a list of segment mappings"
+ )
+ out: List[dict[str, Any]] = []
+ for item in inner_res:
+ if not isinstance(item, dict):
+ raise RuntimeError(
+ "wait_local_segments_ready segment item must be a dict"
+ )
+ out.append(dict(item))
+ return out
+
+ def local_fast_put_start(
+ self,
+ keys: List[str],
+ value_len: int,
+ opts: Optional[PutOptionalArgs] = None,
+ ) -> int:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when local_fast_put_start(). Call setup() first."
+ )
+ if len(keys) == 0:
+ raise ValueError("local_fast_put_start requires at least one key")
+ if not isinstance(value_len, int):
+ raise ValueError(f"value_len must be int; got {type(value_len)}")
+ if value_len <= 0:
+ raise ValueError(f"value_len must be > 0; got {value_len}")
+ reject_if_inflight_same_key = (
+ bool(opts.reject_if_inflight_same_key) if opts is not None else False
+ )
+ reject_if_exist_same_key = (
+ bool(opts.reject_if_exist_same_key) if opts is not None else False
+ )
+ write_through = bool(opts.write_through) if opts is not None else True
+ inner_res = self._client.local_fast_put_start(
+ keys,
+ value_len,
+ reject_if_inflight_same_key=reject_if_inflight_same_key,
+ reject_if_exist_same_key=reject_if_exist_same_key,
+ write_through=write_through,
+ )
+ if not inner_res.is_ok():
+ err = inner_res.unwrap_error()
+ if isinstance(err, Exception):
+ raise err
+ raise RuntimeError(f"local_fast_put_start backend error: {err}")
+ plan_ptr = inner_res.unwrap()
+ if not isinstance(plan_ptr, int) or plan_ptr <= 0:
+ raise RuntimeError(f"local_fast_put_start returned invalid plan_ptr: {plan_ptr!r}")
+ return int(plan_ptr)
+
+ def local_fast_put_commit(self, plan_ptr: int) -> KvFuture:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when local_fast_put_commit(). Call setup() first."
+ )
+ inner_res = self._client.local_fast_put_commit(plan_ptr)
+ if not inner_res.is_ok():
+ raise RuntimeError(
+ f"local_fast_put_commit backend error: {inner_res.unwrap_error()}"
+ )
+ inner_future = inner_res.unwrap()
+ if inner_future is None:
+ raise RuntimeError("local_fast_put_commit returned empty future")
+ return _FluxonBatchRetCodeFuture(inner_future, [plan_ptr])
+
+ def put_abort(self, plan_ptr: int) -> None:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when put_abort(). Call setup() first."
+ )
+ inner_res = self._client.put_abort(plan_ptr)
+ if not inner_res.is_ok():
+ raise RuntimeError(f"put_abort backend error: {inner_res.unwrap_error()}")
+ _ = inner_res.unwrap()
+
+ def get_views(
+ self,
+ keys: List[str],
+ concurrency: Optional[int] = None,
+ ) -> int:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when get_views(). Call setup() first."
+ )
+ if len(keys) == 0:
+ raise ValueError("get_views requires at least one key")
+ inner_res = self._client.get_views(
+ keys,
+ concurrency=concurrency if concurrency is not None else self._batch_concurrency,
+ )
+ if not inner_res.is_ok():
+ raise RuntimeError(f"get_views backend error: {inner_res.unwrap_error()}")
+ plan_ptr = inner_res.unwrap()
+ if not isinstance(plan_ptr, int) or plan_ptr <= 0:
+ raise RuntimeError(f"get_views returned invalid plan_ptr: {plan_ptr!r}")
+ return int(plan_ptr)
+
+ def release_views(self, plan_ptr: int) -> None:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when release_views(). Call setup() first."
+ )
+ inner_res = self._client.release_views(plan_ptr)
+ if not inner_res.is_ok():
+ raise RuntimeError(f"release_views backend error: {inner_res.unwrap_error()}")
+ _ = inner_res.unwrap()
+
+ def get_start(
+ self,
+ keys: List[str],
+ prefix_best_effort: bool = True,
+ atomic_group_lens: Optional[List[int]] = None,
+ ) -> GetStartHandle:
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when get_start(). Call setup() first."
+ )
+ if len(keys) == 0:
+ raise ValueError("get_start requires at least one key")
+ normalized_group_lens: Optional[List[int]] = None
+ if atomic_group_lens is not None:
+ normalized_group_lens = [int(length) for length in atomic_group_lens]
+ if any(length <= 0 for length in normalized_group_lens):
+ raise ValueError("get_start atomic_group_lens entries must be > 0")
+ if sum(normalized_group_lens) != len(keys):
+ raise ValueError(
+ "get_start atomic_group_lens must sum to keys length: "
+ f"sum={sum(normalized_group_lens)} keys={len(keys)}"
+ )
+
+ started_at_ns = time.monotonic_ns()
+ inner_res = self._client.get_start(
+ list(keys),
+ bool(prefix_best_effort),
+ normalized_group_lens,
+ self._batch_concurrency,
+ )
+ if not inner_res.is_ok():
+ raise RuntimeError(f"get_start backend error: {inner_res.unwrap_error()}")
+ payload = inner_res.unwrap()
+ if not isinstance(payload, dict):
+ raise RuntimeError(f"get_start returned non-dict payload: {type(payload)}")
+ backend_handle = int(payload["handle"])
+ result = _build_get_start_result_from_backend_payload(
+ payload,
+ list(keys),
+ bool(prefix_best_effort),
+ normalized_group_lens,
+ )
+ logging.info(
+ "FluxonKVCacheStore get_start result: keys=%d raw_prefix_hit_len=%d "
+ "transferable_len=%d prefix_hit_groups=%d all_hit=%s "
+ "first_miss_index=%s first_miss_group_index=%s "
+ "prefix_best_effort=%s duration_ms=%.3f",
+ len(keys),
+ result.raw_prefix_hit_len,
+ result.transferable_len,
+ result.prefix_hit_groups,
+ result.all_hit,
+ result.first_miss_index,
+ result.first_miss_group_index,
+ result.prefix_best_effort,
+ (time.monotonic_ns() - started_at_ns) / 1_000_000.0,
+ )
+ return GetStartHandle(
+ keys=tuple(keys),
+ result=result,
+ created_at_ns=started_at_ns,
+ backend_token=id(self),
+ backend_handle=backend_handle,
+ )
+
+ def cancel_get_transfer(self, handle: GetStartHandle) -> None:
+ if not isinstance(handle, GetStartHandle):
+ raise TypeError(f"cancel_get_transfer requires GetStartHandle, got {type(handle)}")
+ if handle.backend_token is not None and handle.backend_token != id(self):
+ raise RuntimeError(
+ "cancel_get_transfer handle belongs to a different FluxonKVCacheStore"
+ )
+ if handle.closed:
+ return
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when cancel_get_transfer(). Call setup() first."
+ )
+ inner_res = self._client.cancel_get_transfer(int(handle.backend_handle))
+ if not inner_res.is_ok():
+ raise RuntimeError(
+ f"cancel_get_transfer backend error: {inner_res.unwrap_error()}"
+ )
+ handle.closed = True
+
+ def get_transfer(
+ self,
+ handle: GetStartHandle,
+ concurrency: Optional[int] = None,
+ ) -> int:
+ if not isinstance(handle, GetStartHandle):
+ raise TypeError(f"get_transfer requires GetStartHandle, got {type(handle)}")
+ if handle.backend_token is not None and handle.backend_token != id(self):
+ raise RuntimeError("get_transfer handle belongs to a different FluxonKVCacheStore")
+ if handle.closed:
+ raise RuntimeError("get_transfer handle has been closed")
+ result = handle.result
+ if result.transferable_len == 0:
+ raise RuntimeError(
+ "get_transfer requires a non-empty transferable prefix: "
+ f"transferable_len={result.transferable_len} total={len(result.keys)} "
+ f"raw_prefix_hit_len={result.raw_prefix_hit_len} "
+ f"first_miss_index={result.first_miss_index} "
+ f"first_miss_group_index={result.first_miss_group_index}"
+ )
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when get_transfer(). Call setup() first."
+ )
+ _ = concurrency
+ inner_res = self._client.get_transfer(handle.backend_handle)
+ if not inner_res.is_ok():
+ handle.closed = True
+ raise RuntimeError(f"get_transfer backend error: {inner_res.unwrap_error()}")
+ plan_ptr = inner_res.unwrap()
+ handle.closed = True
+ if not isinstance(plan_ptr, int) or plan_ptr <= 0:
+ raise RuntimeError(f"get_transfer returned invalid plan_ptr: {plan_ptr!r}")
+ return int(plan_ptr)
+
def get_etcd_config(self) -> List[str]:
if self._client is None:
raise RuntimeError("Store not initialized")
@@ -889,6 +1729,19 @@ def metrics_snapshot(self) -> MetricSnapshot:
return MetricSnapshot(per_segment=normalized)
+ def observability_snapshot_async(self) -> KvFuture:
+ """Return a future for Fluxon locality and IO counters."""
+ if self._client is None:
+ raise RuntimeError(
+ "Store not initialized when observability_snapshot_async(). Call setup() first."
+ )
+ res = self._client.observability_snapshot_async()
+ if not res.is_ok():
+ raise RuntimeError(
+ f"observability_snapshot_async backend error: {res.unwrap_error()}"
+ )
+ return res.unwrap()
+
# --- Fluxon-kv lease helpers (synchronous) ---
def allocate_lease(self, ttl_seconds: int) -> Result[int, ApiError]:
try:
@@ -948,7 +1801,7 @@ def __init__(self, inner_future: Any) -> None:
self._inner = inner_future
def is_waiting(self) -> bool:
- return bool(getattr(self._inner, "is_waiting")())
+ return bool(self._inner.is_waiting())
def _decode_wait_success(
self,
@@ -985,7 +1838,7 @@ def __init__(self, inner_future: Any) -> None:
self._inner = inner_future
def is_waiting(self) -> bool:
- return bool(getattr(self._inner, "is_waiting")())
+ return bool(self._inner.is_waiting())
def wait(self) -> Result[bytes, ApiError]:
res = self._inner.wait()
@@ -1022,7 +1875,7 @@ def __del__(self) -> None:
self._dlpack_capsules = []
def is_waiting(self) -> bool:
- return bool(getattr(self._inner, "is_waiting")())
+ return bool(self._inner.is_waiting())
def wait(self) -> Result[Union[Any, MemHolder], ApiError]:
from ..api_error import OkNone, Result as PyResult # type: ignore
@@ -1035,3 +1888,36 @@ def wait(self) -> Result[Union[Any, MemHolder], ApiError]:
_ = res.unwrap()
return PyResult.new_ok(OkNone()) # type: ignore
+
+
+class _FluxonBatchRetCodeFuture(KvFuture):
+ """Future wrapper for batch APIs that resolve to one integer ret-code per key."""
+
+ def __init__(self, inner_future: Any, keepalive: List[object]) -> None:
+ self._inner = inner_future
+ self._keepalive = keepalive
+
+ def __del__(self) -> None:
+ self._keepalive = []
+
+ def is_waiting(self) -> bool:
+ return bool(self._inner.is_waiting())
+
+ def wait(self) -> Result[List[int], ApiError]:
+ res = self._inner.wait()
+ self._keepalive = []
+ if not res.is_ok():
+ return Result.new_error(res.unwrap_error())
+ raw = res.unwrap()
+ if not isinstance(raw, list):
+ return Result.new_error(
+ GeneralError(message=f"batch future returned non-list payload: {type(raw)}")
+ )
+ out: List[int] = []
+ for item in raw:
+ if not isinstance(item, int):
+ return Result.new_error(
+ GeneralError(message=f"batch future returned non-int item: {type(item)}")
+ )
+ out.append(int(item))
+ return Result.new_ok(out)
diff --git a/fluxon_py/kvclient/kvclient_interface.py b/fluxon_py/kvclient/kvclient_interface.py
index f50db0f..e126eec 100644
--- a/fluxon_py/kvclient/kvclient_interface.py
+++ b/fluxon_py/kvclient/kvclient_interface.py
@@ -21,6 +21,47 @@
FlatDict = Dict[str, Union[int, float, bool, str, bytes, DLPacked]]
+@dataclass(frozen=True)
+class GetStartResult:
+ """
+ Result of a group-prefix best-effort get_start().
+
+ Semantics:
+ - ``keys`` is the caller-provided ordered page-key sequence.
+ - ``raw_prefix_hit_len`` is the page-level continuous hit prefix.
+ - ``transferable_len`` is rounded down to complete atomic groups and is the
+ only prefix that can be consumed by get_transfer().
+ """
+
+ keys: Tuple[str, ...]
+ raw_prefix_hit_len: int
+ transferable_len: int
+ prefix_hit_groups: int
+ atomic_group_lens: Optional[Tuple[int, ...]]
+ prefix_best_effort: bool
+ first_miss_index: Optional[int]
+ first_miss_group_index: Optional[int]
+ all_hit: bool
+
+
+@dataclass
+class GetStartHandle:
+ """
+ Opaque-ish handle returned by get_start() and consumed by get_transfer().
+
+ Callers must pass this handle to cancel_get_transfer() when abandoning it
+ without calling get_transfer(). Fluxon backends may keep strong holder
+ references alive while this handle is live.
+ """
+
+ keys: Tuple[str, ...]
+ result: GetStartResult
+ created_at_ns: int
+ backend_token: Optional[int] = None
+ backend_handle: int = 0
+ closed: bool = False
+
+
@dataclass
class PutOptionalArgs:
"""
@@ -29,9 +70,16 @@ class PutOptionalArgs:
- lease_id: attach the written key to a lease on commit.
- reject_if_inflight_same_key: ask Fluxon to fail-fast when the same key is already
being written by another inflight put.
+ - reject_if_exist_same_key: ask Fluxon to fail-fast when the key already has a
+ committed live replica.
+ - write_through: keep synchronous remote-placement semantics when the backend
+ supports an async write-back path. Defaults to True to match SGLang
+ HiCache's default write policy.
"""
lease_id: Optional[int] = None
reject_if_inflight_same_key: bool = False
+ reject_if_exist_same_key: bool = False
+ write_through: bool = True
def support_mooncake(self) -> Tuple[bool, List[str]]:
"""
@@ -48,6 +96,10 @@ def support_mooncake(self) -> Tuple[bool, List[str]]:
unsupported.append("lease_id")
if self.reject_if_inflight_same_key:
unsupported.append("reject_if_inflight_same_key")
+ if self.reject_if_exist_same_key:
+ unsupported.append("reject_if_exist_same_key")
+ if self.write_through:
+ unsupported.append("write_through")
return (len(unsupported) == 0, unsupported)
@@ -139,7 +191,7 @@ def put_blocking(
_ = wait_result.unwrap()
return Result.new_ok(OkNone())
- def get_blocking(self, key: str) -> Result["MemHolder", ApiError]:
+ def get_blocking(self, key: str) -> Result[Union[Any, "MemHolder"], ApiError]:
"""Synchronously retrieve a value by key.
Default implementation delegates to ``get()`` followed by
@@ -151,6 +203,100 @@ def get_blocking(self, key: str) -> Result["MemHolder", ApiError]:
return Result.new_error(get_result.unwrap_error())
return get_result.unwrap().wait()
+ def batch_put_blocking(
+ self,
+ keys: List[str],
+ values: List[FlatDict],
+ opts: Optional[PutOptionalArgs] = None,
+ concurrency: Optional[int] = None,
+ ) -> List[Result[OkNone, ApiError]]:
+ """Synchronously store a batch of key-value pairs."""
+ if len(keys) != len(values):
+ raise ValueError("batch_put_blocking requires keys and values to have the same length")
+ _ = concurrency
+ return [self.put_blocking(key, value, opts=opts) for key, value in zip(keys, values)]
+
+ def batch_get_blocking(
+ self,
+ keys: List[str],
+ concurrency: Optional[int] = None,
+ ) -> List[Result[Union[Any, "MemHolder"], ApiError]]:
+ """Synchronously retrieve a batch of keys."""
+ _ = concurrency
+ return [self.get_blocking(key) for key in keys]
+
+ def local_fast_put_start(
+ self,
+ keys: List[str],
+ value_len: int,
+ opts: Optional[PutOptionalArgs] = None,
+ ) -> int:
+ _ = keys
+ _ = value_len
+ _ = opts
+ raise NotImplementedError(
+ "local_fast_put_start is only implemented by backends with native plan_ptr support"
+ )
+
+ def local_fast_put_commit(self, plan_ptr: int) -> "KvFuture":
+ _ = plan_ptr
+ raise NotImplementedError(
+ "local_fast_put_commit is only implemented by backends with native plan_ptr support"
+ )
+
+ def put_abort(self, plan_ptr: int) -> None:
+ _ = plan_ptr
+ raise NotImplementedError(
+ "put_abort is only implemented by backends with native plan_ptr support"
+ )
+
+ def get_views(
+ self,
+ keys: List[str],
+ concurrency: Optional[int] = None,
+ ) -> int:
+ _ = keys
+ _ = concurrency
+ raise NotImplementedError(
+ "get_views is only implemented by backends with native plan_ptr support"
+ )
+
+ def release_views(self, plan_ptr: int) -> None:
+ _ = plan_ptr
+ raise NotImplementedError(
+ "release_views is only implemented by backends with native plan_ptr support"
+ )
+
+ def get_start(
+ self,
+ keys: List[str],
+ prefix_best_effort: bool = True,
+ atomic_group_lens: Optional[List[int]] = None,
+ ) -> GetStartHandle:
+ _ = keys
+ _ = prefix_best_effort
+ _ = atomic_group_lens
+ raise NotImplementedError(
+ "get_start is only implemented by backends with native prefix get support"
+ )
+
+ def get_transfer(
+ self,
+ handle: GetStartHandle,
+ concurrency: Optional[int] = None,
+ ) -> int:
+ _ = handle
+ _ = concurrency
+ raise NotImplementedError(
+ "get_transfer is only implemented by backends with native prefix get support"
+ )
+
+ def cancel_get_transfer(self, handle: GetStartHandle) -> None:
+ _ = handle
+ raise NotImplementedError(
+ "cancel_get_transfer is only implemented by backends with native prefix get support"
+ )
+
@abstractmethod
def get_size(self, key: str) -> Result[int, ApiError]:
"""Get the size of a stored value (non-blocking)."""
diff --git a/fluxon_rs/fluxon_commu/src/facade/transfer_engine.rs b/fluxon_rs/fluxon_commu/src/facade/transfer_engine.rs
index 878e5c6..e5353a5 100644
--- a/fluxon_rs/fluxon_commu/src/facade/transfer_engine.rs
+++ b/fluxon_rs/fluxon_commu/src/facade/transfer_engine.rs
@@ -74,12 +74,10 @@ impl ClosedLocalSegmentLeaseRegistry {
where
G: Send + Sync + 'static,
{
- let boxed = self
- .guards
- .lock()
- .await
- .remove(&handle)
- .ok_or_else(|| format!("closed sdk local segment lease handle {handle} not found"))?;
+ let boxed =
+ self.guards.lock().await.remove(&handle).ok_or_else(|| {
+ format!("closed sdk local segment lease handle {handle} not found")
+ })?;
boxed.downcast::().map(|guard| *guard).map_err(|_| {
format!(
"closed sdk local segment lease handle {handle} has unexpected runtime guard type"
@@ -461,7 +459,7 @@ impl ClientTransferEngineCore {
len,
seg_guard,
)
- .await
+ .await
}
}
@@ -482,7 +480,11 @@ impl ClientTransferEngineCore {
let initial_local_segment_guard = match seg_guard {
Some(guard) => Some(guard),
None if runtime.supports_local_segment_transfer() => {
- let local_addr = if peer_src_or_target { target_addr } else { src_addr };
+ let local_addr = if peer_src_or_target {
+ target_addr
+ } else {
+ src_addr
+ };
match runtime.ensure_local_segment_guard(local_addr, None).await {
Ok(guard) => Some(guard),
Err(_) => None,
diff --git a/fluxon_rs/fluxon_commu_closed_sdk_consumer/src/lib.rs b/fluxon_rs/fluxon_commu_closed_sdk_consumer/src/lib.rs
index 6fab54e..caad34b 100644
--- a/fluxon_rs/fluxon_commu_closed_sdk_consumer/src/lib.rs
+++ b/fluxon_rs/fluxon_commu_closed_sdk_consumer/src/lib.rs
@@ -11,9 +11,9 @@ use fluxon_commu_contract::{
ClosedRuntimeCallRawObservedOutputView, ClosedRuntimeClusterEventStreamItem,
ClosedRuntimeClusterManagerCall, ClosedRuntimeClusterManagerResponse,
ClosedRuntimeClusterRdmaResolvedConfigStreamItem, ClosedRuntimeDesiredTransferPeer,
- ClosedRuntimeDispatchRequestView,
- ClosedRuntimeDispatchResponse, ClosedRuntimeDispatchTransportPolicy, ClosedRuntimeError,
- ClosedRuntimeHandle, ClosedRuntimeHostCallbackHandle, ClosedRuntimeP2pCall,
+ ClosedRuntimeDispatchRequestView, ClosedRuntimeDispatchResponse,
+ ClosedRuntimeDispatchTransportPolicy, ClosedRuntimeError, ClosedRuntimeHandle,
+ ClosedRuntimeHostCallbackHandle, ClosedRuntimeP2pCall,
ClosedRuntimeP2pCallRawObservedRequestView, ClosedRuntimeP2pResponse,
ClosedRuntimeP2pSendResponseRawRequestView, ClosedRuntimePeerGen, ClosedRuntimeRawSlice,
ClosedRuntimeRequest, ClosedRuntimeResponse, ClosedRuntimeTransferEngineCall,
@@ -491,11 +491,15 @@ impl WireBodyPartsOwner {
let (raw_lengths, raw_payload) = match raw_bytes.len() {
0 => (WireBodyRawLengths::Empty, WireBodyRawPayload::Empty),
1 => {
- let part = raw_bytes.into_iter().next().expect("single raw part missing");
- let len =
- u32::try_from(part.len()).map_err(|_| ClosedSdkConsumerError::RuntimeDecode {
+ let part = raw_bytes
+ .into_iter()
+ .next()
+ .expect("single raw part missing");
+ let len = u32::try_from(part.len()).map_err(|_| {
+ ClosedSdkConsumerError::RuntimeDecode {
detail: format!("wire raw part too large for u32 length: {}", part.len()),
- })?;
+ }
+ })?;
(
WireBodyRawLengths::Single([len]),
WireBodyRawPayload::Single(part),
@@ -849,8 +853,7 @@ fn decode_call_raw_observed_output_view(
return Err(ClosedSdkConsumerError::RuntimeDecode {
detail: format!(
"closed SDK call_raw_observed serialize_part overflow: serialize_len={} full_len={}",
- message_view.body.serialize_part.len,
- message_view.body.full_body.len,
+ message_view.body.serialize_part.len, message_view.body.full_body.len,
),
});
}
@@ -860,21 +863,19 @@ fn decode_call_raw_observed_output_view(
.ok_or_else(|| ClosedSdkConsumerError::RuntimeDecode {
detail: "closed SDK call_raw_observed raw_bytes length overflow".to_string(),
})?;
- let expected_full_len =
- message_view
- .body
- .serialize_part
- .len
- .checked_add(raw_total)
- .ok_or_else(|| ClosedSdkConsumerError::RuntimeDecode {
- detail: "closed SDK call_raw_observed body length overflow".to_string(),
- })?;
+ let expected_full_len = message_view
+ .body
+ .serialize_part
+ .len
+ .checked_add(raw_total)
+ .ok_or_else(|| ClosedSdkConsumerError::RuntimeDecode {
+ detail: "closed SDK call_raw_observed body length overflow".to_string(),
+ })?;
if expected_full_len != message_view.body.full_body.len {
return Err(ClosedSdkConsumerError::RuntimeDecode {
detail: format!(
"closed SDK call_raw_observed body length mismatch: expected={} full_len={}",
- expected_full_len,
- message_view.body.full_body.len,
+ expected_full_len, message_view.body.full_body.len,
),
});
}
@@ -923,9 +924,7 @@ fn decode_call_raw_observed_output_view(
frame_recv_done_ts_us: message_view.local_observe.frame_recv_done_ts_us,
dispatch_enqueued_ts_us: message_view.local_observe.dispatch_enqueued_ts_us,
dispatch_started_ts_us: message_view.local_observe.dispatch_started_ts_us,
- complete_pending_call_ts_us: message_view
- .local_observe
- .complete_pending_call_ts_us,
+ complete_pending_call_ts_us: message_view.local_observe.complete_pending_call_ts_us,
},
},
observe: fluxon_commu_contract::ClosedRuntimeRpcCallTransportObserveTrace {
@@ -1550,8 +1549,8 @@ async fn invoke_completion_async_with_keepalive(
) -> i32,
) -> Result<(i32, Bytes), ClosedSdkConsumerError> {
let (sender, receiver) = tokio::sync::oneshot::channel::<(i32, Bytes)>();
- let user_data = Box::into_raw(Box::new(RuntimeCompletionState { sender, keepalive }))
- .cast::();
+ let user_data =
+ Box::into_raw(Box::new(RuntimeCompletionState { sender, keepalive })).cast::();
let submit_status = submit(user_data, Some(runtime_completion_callback));
if submit_status != 0 {
unsafe {
@@ -2082,7 +2081,9 @@ pub async fn p2p_call_raw_observed(
)
.await?;
match status_code {
- FLUXON_COMMU_CLOSED_RUNTIME_RESULT_OK => decode_call_raw_observed_output_view(payload.as_ref()),
+ FLUXON_COMMU_CLOSED_RUNTIME_RESULT_OK => {
+ decode_call_raw_observed_output_view(payload.as_ref())
+ }
FLUXON_COMMU_CLOSED_RUNTIME_RESULT_ERR => {
let error = bitcode::decode::(payload.as_ref()).map_err(
|decode_error| ClosedSdkConsumerError::RuntimeDecode {
diff --git a/fluxon_rs/fluxon_fs/src/agent.rs b/fluxon_rs/fluxon_fs/src/agent.rs
index a482616..45e15e7 100644
--- a/fluxon_rs/fluxon_fs/src/agent.rs
+++ b/fluxon_rs/fluxon_fs/src/agent.rs
@@ -1408,13 +1408,24 @@ impl FluxonFsAgent {
.id
.to_string();
let cache_root_base = if self.kv_framework.is_external_mode() {
- self.kv_framework
+ let shared_file_path = self
+ .kv_framework
.external_client_api_view()
.external_client_api()
.inner()
- .large_file_paths()
- .fs_disk_cache_base_dir()
- .map_err(|err| format!("invalid external large_file_paths: {}", err))?
+ .shared_file_path();
+ if shared_file_path.is_empty() {
+ return Err("external shared_file_path is empty".to_string());
+ }
+ let cache_root = Path::new(&shared_file_path).join("fluxon_fs_disk_cache");
+ fs::create_dir_all(&cache_root).map_err(|err| {
+ format!(
+ "invalid external shared_file_path for fluxon fs disk cache: {} ({})",
+ cache_root.display(),
+ err
+ )
+ })?;
+ cache_root
} else {
self.kv_framework
.client_seg_pool_view()
diff --git a/fluxon_rs/fluxon_fs/src/agent_service/transfer_agent.rs b/fluxon_rs/fluxon_fs/src/agent_service/transfer_agent.rs
index 1738ade..ca54a71 100644
--- a/fluxon_rs/fluxon_fs/src/agent_service/transfer_agent.rs
+++ b/fluxon_rs/fluxon_fs/src/agent_service/transfer_agent.rs
@@ -9,28 +9,23 @@ use std::time::{Duration, Instant};
use fluxon_fs_core::config::{
FS_AGENT_TRANSFER_STREAM_CLOSE_RPC_PATH, FS_AGENT_TRANSFER_STREAM_NEXT_RPC_PATH,
- FS_AGENT_TRANSFER_STREAM_OPEN_RPC_PATH,
- FS_MASTER_TRANSFER_SCHEDULER_HEARTBEAT_RPC_PATH, FS_MASTER_TRANSFER_SCHEDULER_RESULT_RPC_PATH,
- FluxonFsTransferBatchCollectInfoWire, FluxonFsTransferBatchKind,
- FluxonFsTransferCollectInfoKind, FluxonFsTransferDispositionWire,
- FluxonFsTransferFailedFileReasonKindWire,
- FluxonFsTransferReadStreamCloseWire, FluxonFsTransferReadStreamNextResultWire,
- FluxonFsTransferReadStreamNextWire, FluxonFsTransferReadStreamOpenResultWire,
- FluxonFsTransferReadStreamOpenWire,
- FluxonFsTransferSkipEntryKind, FluxonFsTransferSkipEntryWire,
- FluxonFsTransferManifestEntryWire, FluxonFsTransferManifestWire,
- FluxonFsTransferScanMode,
- FluxonFsTransferScanEventAckWire, FluxonFsTransferScanEventKindWire,
- FluxonFsTransferScanEventWire, FluxonFsTransferScanLaunchResultWire,
+ FS_AGENT_TRANSFER_STREAM_OPEN_RPC_PATH, FS_MASTER_TRANSFER_SCHEDULER_HEARTBEAT_RPC_PATH,
+ FS_MASTER_TRANSFER_SCHEDULER_RESULT_RPC_PATH, FluxonFsTransferBatchCollectInfoWire,
+ FluxonFsTransferBatchKind, FluxonFsTransferCollectInfoKind, FluxonFsTransferDispositionWire,
+ FluxonFsTransferFailedFileReasonKindWire, FluxonFsTransferManifestEntryWire,
+ FluxonFsTransferManifestWire, FluxonFsTransferReadStreamCloseWire,
+ FluxonFsTransferReadStreamNextResultWire, FluxonFsTransferReadStreamNextWire,
+ FluxonFsTransferReadStreamOpenResultWire, FluxonFsTransferReadStreamOpenWire,
FluxonFsTransferScanAssignmentWire, FluxonFsTransferScanBatchWire,
- FluxonFsTransferScanChildUnitWire, FluxonFsTransferScanFrontier,
+ FluxonFsTransferScanChildUnitWire, FluxonFsTransferScanEventAckWire,
+ FluxonFsTransferScanEventKindWire, FluxonFsTransferScanEventWire, FluxonFsTransferScanFrontier,
FluxonFsTransferScanFrontierDirEntry, FluxonFsTransferScanFrontierEntry,
- FluxonFsTransferScanResultWire,
- FluxonFsTransferSymlinkNoticeEntryWire, FluxonFsTransferWorkerCollectInfoResultWire,
- FluxonFsTransferWorkerAssignmentWire, FluxonFsTransferWorkerFileResultWire,
- FluxonFsTransferWorkerFailedFileResultWire,
- FluxonFsTransferWorkerHeartbeatResultWire, FluxonFsTransferWorkerHeartbeatTelemetryWire,
- FluxonFsTransferWorkerHeartbeatWire,
+ FluxonFsTransferScanLaunchResultWire, FluxonFsTransferScanMode, FluxonFsTransferScanResultWire,
+ FluxonFsTransferSkipEntryKind, FluxonFsTransferSkipEntryWire,
+ FluxonFsTransferSymlinkNoticeEntryWire, FluxonFsTransferWorkerAssignmentWire,
+ FluxonFsTransferWorkerCollectInfoResultWire, FluxonFsTransferWorkerFailedFileResultWire,
+ FluxonFsTransferWorkerFileResultWire, FluxonFsTransferWorkerHeartbeatResultWire,
+ FluxonFsTransferWorkerHeartbeatTelemetryWire, FluxonFsTransferWorkerHeartbeatWire,
FluxonFsTransferWorkerLaunchResultWire, FluxonFsTransferWorkerResultAckWire,
FluxonFsTransferWorkerResultWire, FluxonFsTransferWorkerStopReasonWire,
transfer_collect_info_output_relpath,
@@ -39,8 +34,8 @@ use fluxon_fs_core::retry::{
BackoffConfig, DEFAULT_WARN_INTERVAL_SECS, WarnConfig, next_backoff, should_warn,
};
use fluxon_kv::rpcresp_kvresult_convert::msg_and_error::{ApiError, KvError};
-use fluxon_kv::user_api::flat_dict::{FlatDict, FlatValue};
use fluxon_kv::user_api::FluxonUserApi;
+use fluxon_kv::user_api::flat_dict::{FlatDict, FlatValue};
use parking_lot::{Condvar, Mutex};
use super::{
@@ -202,16 +197,13 @@ fn transfer_scan_session_state() -> &'static Mutex {
TRANSFER_SCAN_SESSION_STATE.get_or_init(|| Mutex::new(TransferScanSessionState::default()))
}
-fn cleanup_expired_transfer_scan_sessions(
- state: &mut TransferScanSessionState,
- now_unix_ms: i64,
-) {
- state
- .root_dir_listing_sessions
- .retain(|_, session| session.lease_expire_unix_ms <= 0 || session.lease_expire_unix_ms > now_unix_ms);
- state
- .subtree_streaming_sessions
- .retain(|_, session| session.lease_expire_unix_ms <= 0 || session.lease_expire_unix_ms > now_unix_ms);
+fn cleanup_expired_transfer_scan_sessions(state: &mut TransferScanSessionState, now_unix_ms: i64) {
+ state.root_dir_listing_sessions.retain(|_, session| {
+ session.lease_expire_unix_ms <= 0 || session.lease_expire_unix_ms > now_unix_ms
+ });
+ state.subtree_streaming_sessions.retain(|_, session| {
+ session.lease_expire_unix_ms <= 0 || session.lease_expire_unix_ms > now_unix_ms
+ });
}
fn same_root_continuation_scan_unit(
@@ -301,10 +293,7 @@ fn flush_pending_root_direct_files_batch(
return Ok(None);
}
let batch = build_direct_files_only_batch_from_entries_with_batch_id(
- direct_files_only_batch_id_for_partition(
- assignment,
- session.next_direct_files_batch_index,
- ),
+ direct_files_only_batch_id_for_partition(assignment, session.next_direct_files_batch_index),
assignment,
assignment.root_relpath.clone(),
std::mem::take(&mut session.pending_direct_files),
@@ -313,7 +302,8 @@ fn flush_pending_root_direct_files_batch(
)?;
session.pending_direct_bytes = 0;
session.next_direct_files_batch_index = session.next_direct_files_batch_index.saturating_add(1);
- session.emitted_direct_files_batch_count = session.emitted_direct_files_batch_count.saturating_add(1);
+ session.emitted_direct_files_batch_count =
+ session.emitted_direct_files_batch_count.saturating_add(1);
Ok(Some(batch))
}
@@ -414,7 +404,8 @@ fn open_transfer_root_dir_listing_session(
root_dir_abs: &str,
assignment: &FluxonFsTransferScanAssignmentWire,
) -> Result