Skip to content
Open
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
### 设计文档
- [模块架构与关联关系](design/module_architecture.md) - 各模块职责、依赖方向、控制流与数据流,附 Mermaid 图
- [基本概念](design/basic_concepts.md) - Storage、Instance Group、Instance、Block、CacheLocation 等核心概念
- [Client SDK I/O 契约](design/client_sdk_io_contract.md) - deadline 语义、buffer 生命周期、各后端取消能力矩阵
- [ReportEvent 增量上报与权威快照设计](design/report_event_snapshot_uri_version.md) - 增量/快照协同、提交屏障、故障恢复、性能取舍与 Subscriber 集成
- [ReportEvent / GetHostCacheState 小 block 性能记录](design/report_event_performance.md) - local/Redis 指标解释、锁与可见性语义、有界并发、容量基准及后续优化边界
- [高可用与选主机制](design/ha_leader_elector.md) - HA 架构、LeaderElector 状态机、CoordinationBackend、Leader 发现
Expand Down
46 changes: 46 additions & 0 deletions docs/design/client_sdk_io_contract.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
# Client SDK I/O 契约(deadline / 取消 / Buffer 生命周期)
Comment thread
lpdink marked this conversation as resolved.

## 1. 核心契约(一句话)

**deadline(绝对 steady_clock 毫秒,0=无限制)到达后,实现方不得再触碰调用方的 local buffers。**

各后端按自身能力履约(见 §3 矩阵)。做不到"绝不碰"的后端必须如实声明为 soft 级。

## 2. 调用方用法

- 读/写路径均传真实 deadline:`deadline_ms_from_now(sdk_get/put_timeout_ms)`。
- 写路径的租约(write_timeout_seconds)不由 connector 计算或传递;SdkWrapper 内部对
外部 deadline 与内部 timeout_config 取 min,**min 结果原样下发进 SDK 作为准入
deadline**(caller 传 0 时即内部预算)——wrapper 的等待上限与 SDK 的准入依据
始终是同一个时间点,内层自律不越过外层租约。
- deadline 在各进程内各自计算,不跨进程/跨 rank 传递(见 Known Limitations)。

## 3. 后端履约矩阵

| 后端 | 有界性 | 准入检查 | Buffer 级别 | 备注 |
|---|---|---|---|---|
| LocalFile / NFS | 逐 block | 逐 block | hard | abort 路径 sync GPU stream |
| HF3FS | hf3fs_wait_for_ios(abs_timeout) | 逐 block | hard | 数据先落我们自己的 shm;超时不执行 CopyIovs |
| Mooncake | 逐 key 前置检查 | 逐 key | soft | 上游无取消语义;超时路径输出可归因日志 |
| TairMempool (PACE) | PACE 内部超时 + deadline 透传 | PACE 无队列准入 | hard(默认 staging) | PACE 已有 cancel_and_drain + BufferUseGuard |

## 4. Known Limitations

1. **各进程自算 DDL 导致租约可能越界**:worker 自算起点晚于 scheduler 拿到租约的时刻,
极端情况下写入可能越过 KVCM 租约。取舍理由:多机时间不一致在 happy path 上更危险,
租约越界不在 happy path。
2. **写租约时间原点偏差**:KVCM 服务端在处理 StartWriteCache 时即开始计时
(write_location_manager.cc:187),且会 cap(cache_manager.cc:1090
kMaxWriteTimeoutSeconds=1800)。客户端算出的 DDL 可能晚于服务端真实 expire。
这是 KVCM 服务端既有问题,本次不改。
3. **ManagerClient / RTPLLMClient 层无 DDL**:这两层公开接口保持传 0,走 SdkWrapper 兜底。
不确定外部用户是谁,暂不加。
4. **超时路径不等在飞任务;普通错误路径有界等待**:RunWithTimeoutParallel 在
deadline 到期时立即返回,在飞任务的 I/O 不被取消也不被等待(hard 后端的逐块
准入会尽快停下;soft 后端见 §3)。普通错误路径则不同:错误往往发生在 deadline
之前,SDK 仍在契约允许的窗口内写 caller buffer,因此错误路径会先置 stop 拦截
仍在排队的 group(不再发起新 I/O),再有界等待在飞 peer 至多到 deadline 才返回
——否则 caller 拿到错误后立即复用/释放 buffer,与在飞 DMA 构成数据竞争。
5. **HF3FS 超时泄漏 iov/IOR**:当 DeadlineExpired 导致 WaitIos 超时时,不释放已提交
的 I/O 的 iov 缓冲区和 IOR(释放会导致 UAF)。泄漏规模 = 一次超时的读写调用涉及的
iov 大小。线上应在 WARN 日志中观测泄漏频率,若高频则需引入 3FS 取消机制。
14 changes: 9 additions & 5 deletions kv_cache_manager/client/include/transfer_client.h
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,17 @@ class TransferClient {
const InitParams &init_params,
const SharedMemoryRegistration &shared_memory_registration);

// deadline_ms: 绝对时间点(steady_clock 毫秒),到点后本次调用不再触碰 block_buffers;
// 0 = 调用方不施加 deadline,退回 client 配置的静态超时预算。
// 置于 trace_info 之后以保持既有位置参数调用兼容。
virtual ClientErrorCode LoadKvCaches(const UriStrVec &uri_str_vec,
const BlockBuffers &block_buffers,
std::shared_ptr<TransferTraceInfo> trace_info = nullptr) = 0;
virtual std::pair<ClientErrorCode, UriStrVec>
SaveKvCaches(const UriStrVec &uri_str_vec,
const BlockBuffers &block_buffers,
std::shared_ptr<TransferTraceInfo> trace_info = nullptr) = 0;
std::shared_ptr<TransferTraceInfo> trace_info = nullptr,
int64_t deadline_ms = 0) = 0;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve the existing TransferClient virtual ABI

When an existing C++ client loads the updated public kv_cache_manager_client.so without being recompiled, its virtual call still passes only the old arguments; default arguments are applied only during compilation, so the new implementation reads an unspecified register as deadline_ms and may reject I/O immediately or use an arbitrary deadline. The same ABI mismatch affects SaveKvCaches; retain the old virtual entry points and add deadline-aware overloads, or otherwise require and document a coordinated rebuild.

Useful? React with 👍 / 👎.

virtual std::pair<ClientErrorCode, UriStrVec> SaveKvCaches(const UriStrVec &uri_str_vec,
const BlockBuffers &block_buffers,
std::shared_ptr<TransferTraceInfo> trace_info = nullptr,
int64_t deadline_ms = 0) = 0;

protected:
TransferClient() = default;
Expand Down
2 changes: 2 additions & 0 deletions kv_cache_manager/client/pybind/py_client_binding.cc
Original file line number Diff line number Diff line change
Expand Up @@ -186,12 +186,14 @@ PYBIND11_MODULE(kvcm_py_client, module) {
py::arg("uri_str_vec"),
py::arg("block_buffers"),
py::arg("trace_info") = nullptr,
py::arg("deadline_ms") = 0,
py::call_guard<py::gil_scoped_release>())
.def("SaveKvCaches",
&kvcm::TransferClient::SaveKvCaches,
py::arg("uri_str_vec"),
py::arg("block_buffers"),
py::arg("trace_info") = nullptr,
py::arg("deadline_ms") = 0,
py::call_guard<py::gil_scoped_release>());

} // namespace kv_cache_manager
1 change: 1 addition & 0 deletions kv_cache_manager/client/src/internal/sdk/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ cc_library(
"//conditions:default": [],
}),
hdrs = [
"deadline_util.h",
"local_file_sdk.h",
"sdk_factory.h",
"sdk_interface.h",
Expand Down
29 changes: 29 additions & 0 deletions kv_cache_manager/client/src/internal/sdk/deadline_util.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
#pragma once

#include <chrono>
#include <optional>

namespace kv_cache_manager {

// deadline_ms:绝对时间点(steady_clock 毫秒),0 = 无 deadline。
inline int64_t SteadyClockMs() {
return std::chrono::duration_cast<std::chrono::milliseconds>(std::chrono::steady_clock::now().time_since_epoch())
.count();
}

// 有 deadline 且已过期。
inline bool DeadlineExpired(int64_t deadline_ms) { return deadline_ms > 0 && SteadyClockMs() >= deadline_ms; }

// 剩余毫秒;无 deadline 或已过期时返回 nullopt(需区分两者时配合 DeadlineExpired)。
inline std::optional<int64_t> RemainingMs(int64_t deadline_ms) {
if (deadline_ms <= 0) {
return std::nullopt;
}
const int64_t now_ms = SteadyClockMs();
if (now_ms >= deadline_ms) {
return std::nullopt;
}
return deadline_ms - now_ms;
}

} // namespace kv_cache_manager
38 changes: 30 additions & 8 deletions kv_cache_manager/client/src/internal/sdk/hf3fs_sdk.cc
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include <sys/types.h>
#include <unistd.h>

#include "kv_cache_manager/client/src/internal/sdk/deadline_util.h"
#include "kv_cache_manager/client/src/internal/sdk/hf3fs_gpu_util_alias.h"
#include "kv_cache_manager/client/src/internal/sdk/hf3fs_mempool.h"
#include "kv_cache_manager/common/logger.h"
Expand Down Expand Up @@ -64,23 +65,34 @@ ClientErrorCode Hf3fsSdk::Init(const std::shared_ptr<SdkBackendConfig> &sdk_back
return ER_OK;
}

ClientErrorCode Hf3fsSdk::Get(const std::vector<DataStorageUri> &remote_uris, const BlockBuffers &local_buffers) {
ClientErrorCode
Hf3fsSdk::Get(const std::vector<DataStorageUri> &remote_uris, const BlockBuffers &local_buffers, int64_t deadline_ms) {
if (remote_uris.size() != local_buffers.size()) {
KVCM_LOG_WARN(
"get failed, size mismatch, uris size: %lu, buffers size: %lu", remote_uris.size(), local_buffers.size());
return ER_INVALID_PARAMS;
}

// 逐 block 准入检查:deadline 已过期则不再为后续 block
// 创建 Hf3fsUsrbioClient / 发起 I/O。无 deadline(0)时 DeadlineExpired() 恒为 false。
for (size_t i = 0; i < remote_uris.size(); ++i) {
if (Get(remote_uris[i], local_buffers[i]) != ER_OK) {
if (DeadlineExpired(deadline_ms)) {
KVCM_LOG_WARN("get skipped, deadline expired, path: %s, done blocks: %zu/%zu, remaining blocks skipped "
"to protect caller buffer",
remote_uris[i].GetPath().c_str(),
i,
remote_uris.size());
return ER_SDK_TIMEOUT;
}
if (Get(remote_uris[i], local_buffers[i], deadline_ms) != ER_OK) {
return ER_SDKREAD_ERROR;
}
}

return ER_OK;
}

ClientErrorCode Hf3fsSdk::Get(const DataStorageUri &uri, const BlockBuffer &block_buffer) const {
ClientErrorCode Hf3fsSdk::Get(const DataStorageUri &uri, const BlockBuffer &block_buffer, int64_t deadline_ms) const {
if (block_buffer.iovs.empty()) {
return ER_OK;
}
Expand All @@ -101,7 +113,7 @@ ClientErrorCode Hf3fsSdk::Get(const DataStorageUri &uri, const BlockBuffer &bloc

auto usrbio_client = std::make_shared<Hf3fsUsrbioClient>(
BuildHf3fsFileConfig(path), read_iov_handle_, write_iov_handle_, usrbio_api_);
if (!usrbio_client->Read(iovs)) {
if (!usrbio_client->Read(iovs, deadline_ms)) {
KVCM_LOG_WARN("get failed, read failed, path: %s", path.c_str());
return ER_SDKREAD_ERROR;
}
Expand All @@ -111,7 +123,8 @@ ClientErrorCode Hf3fsSdk::Get(const DataStorageUri &uri, const BlockBuffer &bloc

ClientErrorCode Hf3fsSdk::Put(const std::vector<DataStorageUri> &remote_uris,
const BlockBuffers &local_buffers,
std::shared_ptr<std::vector<DataStorageUri>> actual_remote_uris) {
std::shared_ptr<std::vector<DataStorageUri>> actual_remote_uris,
int64_t deadline_ms) {
if (remote_uris.size() != local_buffers.size()) {
KVCM_LOG_WARN(
"put failed, size mismatch, uris size: %lu, buffers size: %lu", remote_uris.size(), local_buffers.size());
Expand All @@ -122,16 +135,25 @@ ClientErrorCode Hf3fsSdk::Put(const std::vector<DataStorageUri> &remote_uris,
return ER_SDKALLOC_ERROR;
}

// 逐 block 准入检查:deadline 已过期则不再为后续 block
// 创建 Hf3fsUsrbioClient / 发起 I/O。无 deadline(0)时 DeadlineExpired() 恒为 false。
for (size_t i = 0; i < remote_uris.size(); ++i) {
if (Put(remote_uris[i], local_buffers[i]) != ER_OK) {
if (DeadlineExpired(deadline_ms)) {
KVCM_LOG_WARN("put skipped, deadline expired, path: %s, done blocks: %zu/%zu, remaining blocks skipped",
remote_uris[i].GetPath().c_str(),
i,
remote_uris.size());
return ER_SDK_TIMEOUT;
}
if (Put(remote_uris[i], local_buffers[i], deadline_ms) != ER_OK) {
return ER_SDKWRITE_ERROR;
}
}

return ER_OK;
}

ClientErrorCode Hf3fsSdk::Put(const DataStorageUri &uri, const BlockBuffer &block_buffer) const {
ClientErrorCode Hf3fsSdk::Put(const DataStorageUri &uri, const BlockBuffer &block_buffer, int64_t deadline_ms) const {
if (block_buffer.iovs.empty()) {
return ER_OK;
}
Expand All @@ -152,7 +174,7 @@ ClientErrorCode Hf3fsSdk::Put(const DataStorageUri &uri, const BlockBuffer &bloc

auto usrbio_client = std::make_shared<Hf3fsUsrbioClient>(
BuildHf3fsFileConfig(path), read_iov_handle_, write_iov_handle_, usrbio_api_);
if (!usrbio_client->Write(iovs)) {
if (!usrbio_client->Write(iovs, deadline_ms)) {
KVCM_LOG_WARN("put failed, write failed, path: %s", path.c_str());
return ER_SDKWRITE_ERROR;
}
Expand Down
11 changes: 7 additions & 4 deletions kv_cache_manager/client/src/internal/sdk/hf3fs_sdk.h
Original file line number Diff line number Diff line change
Expand Up @@ -17,18 +17,21 @@ class Hf3fsSdk : public SdkInterface {
SdkType Type() override { return SdkType::HF3FS; }
ClientErrorCode Init(const std::shared_ptr<SdkBackendConfig> &sdk_backend_config,
const std::shared_ptr<StorageConfig> &storage_config) override;
ClientErrorCode Get(const std::vector<DataStorageUri> &remote_uris, const BlockBuffers &local_buffers) override;
ClientErrorCode Get(const std::vector<DataStorageUri> &remote_uris,
const BlockBuffers &local_buffers,
int64_t deadline_ms) override;
ClientErrorCode Put(const std::vector<DataStorageUri> &remote_uris,
const BlockBuffers &local_buffers,
std::shared_ptr<std::vector<DataStorageUri>> actual_remote_uris) override;
std::shared_ptr<std::vector<DataStorageUri>> actual_remote_uris,
int64_t deadline_ms) override;

protected:
ClientErrorCode Alloc(const std::vector<DataStorageUri> &remote_uris,
std::vector<DataStorageUri> &alloc_uris) override;

private:
ClientErrorCode Get(const DataStorageUri &uri, const BlockBuffer &block_buffer) const;
ClientErrorCode Put(const DataStorageUri &uri, const BlockBuffer &block_buffer) const;
ClientErrorCode Get(const DataStorageUri &uri, const BlockBuffer &block_buffer, int64_t deadline_ms) const;
ClientErrorCode Put(const DataStorageUri &uri, const BlockBuffer &block_buffer, int64_t deadline_ms) const;

bool CheckConfig(const Hf3fsSdkConfig &hf3fs_config) const;
void DeleteRemainingIovShm() const;
Expand Down
Loading
Loading