Skip to content
Open
Show file tree
Hide file tree
Changes from all 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