Skip to content
Merged
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
35 changes: 35 additions & 0 deletions mooncake-integration/store/mooncake_perf_points.def
Original file line number Diff line number Diff line change
Expand Up @@ -179,3 +179,38 @@ PERF_KEY_DEF(MASTER_BG_DISCARD_EXPIRED, "master_service.cpp::EvictionThreadFun
PERF_KEY_DEF(MASTER_BG_SNAPSHOT_PERSIST, "master_service.cpp::SnapshotThreadFunc", "SnapshotPersist")
PERF_KEY_DEF(MASTER_BG_CLIENT_MONITOR, "master_service.cpp::ClientMonitorFunc", "ClientMonitorScan")
PERF_KEY_DEF(MASTER_BG_CLIENT_UNMOUNT, "master_service.cpp::ClientMonitorFunc", "ExpiredClientUnmount")

// ============================================================
// === vLLM MooncakeStoreConnector 路径打点(vllm 0.26.1rc0)===
// === 覆盖 S2-S9 入口 + T1-T8 下沉 + Client 服务层 ===
// ============================================================

// === store_py.cpp Python 绑定层(vllm 直接调用入口,S2-S9)===
PERF_KEY_DEF(STORE_PY_SETUP, "store_py.cpp::setup", "Setup")
PERF_KEY_DEF(STORE_PY_REGISTER_BUFFER, "store_py.cpp::register_buffer", "RegisterBuffer")
PERF_KEY_DEF(STORE_PY_BATCH_PUT_MULTI, "store_py.cpp::batch_put_from_multi_buffers","BatchPutMultiBuf")
PERF_KEY_DEF(STORE_PY_BATCH_GET_INTO_MULTI, "store_py.cpp::batch_get_into_multi_buffers","BatchGetIntoMultiBuf")
PERF_KEY_DEF(STORE_PY_BATCH_IS_EXIST, "store_py.cpp::batch_is_exist", "BatchIsExist")
PERF_KEY_DEF(STORE_PY_BATCH_GET_REPLICA_DESC, "store_py.cpp::batch_get_replica_desc", "BatchGetReplicaDesc")
PERF_KEY_DEF(STORE_PY_REMOVE_ALL, "store_py.cpp::remove_all", "RemoveAll")
PERF_KEY_DEF(STORE_PY_CLOSE, "store_py.cpp::close", "Close")

// === RealClient 核心逻辑层(下沉路径,T1-T8)===
PERF_KEY_DEF(RC_SETUP_REAL, "real_client.cpp::setup_real", "SetupReal")
PERF_KEY_DEF(RC_SETUP_INTERNAL, "real_client.cpp::setup_internal", "SetupInternal")
PERF_KEY_DEF(RC_TEARDOWN_ALL, "real_client.cpp::tearDownAll", "TeardownAll")
PERF_KEY_DEF(RC_TEARDOWN_ALL_INTERNAL, "real_client.cpp::tearDownAll_internal", "TeardownAllInternal")
PERF_KEY_DEF(RC_REMOVE_ALL, "real_client.cpp::removeAll", "RemoveAll")
PERF_KEY_DEF(RC_REMOVE_ALL_INTERNAL, "real_client.cpp::removeAll_internal", "RemoveAllInternal")
PERF_KEY_DEF(RC_BATCH_IS_EXIST, "real_client.cpp::batchIsExist", "BatchIsExist")
PERF_KEY_DEF(RC_BATCH_IS_EXIST_INTERNAL, "real_client.cpp::batchIsExist_internal", "BatchIsExistInternal")
PERF_KEY_DEF(RC_REGISTER_BUFFER, "real_client.cpp::register_buffer", "RegisterBuffer")
PERF_KEY_DEF(RC_REGISTER_BUFFER_INTERNAL, "real_client.cpp::register_buffer_internal","RegisterBufferInternal")
PERF_KEY_DEF(RC_BATCH_PUT_MULTI, "real_client.cpp::batch_put_from_multi_buffers", "BatchPutMultiBuf")
PERF_KEY_DEF(RC_BATCH_PUT_MULTI_INTERNAL, "real_client.cpp::batch_put_from_multi_buffers_internal","BatchPutMultiBufInternal")
PERF_KEY_DEF(RC_BATCH_GET_INTO_MULTI, "real_client.cpp::batch_get_into_multi_buffers", "BatchGetIntoMultiBuf")
PERF_KEY_DEF(RC_BATCH_GET_INTO_MULTI_INTERNAL, "real_client.cpp::batch_get_into_multi_buffers_internal","BatchGetIntoMultiBufInternal")
PERF_KEY_DEF(RC_BATCH_GET_REPLICA_DESC, "real_client.cpp::batch_get_replica_desc", "BatchGetReplicaDesc")

// === Client 服务层(vllm 路径下沉,部分已有打点)===
PERF_KEY_DEF(CLIENT_BATCH_QUERY, "client_service.cpp::BatchQuery", "BatchQuery")
96 changes: 86 additions & 10 deletions mooncake-integration/store/store_py.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#include "memory_alloc.h"
#include "ssd_register_client.h"
#include "device/accelerator_registry.h"
#include "mooncake_logging.h" // MC_LOG

#include <cstdlib> // for atexit
#include <memory>
Expand Down Expand Up @@ -2084,19 +2085,25 @@ PYBIND11_MODULE(store, m) {
const std::string &tenant_id = "default",
bool enable_client_http_server = false,
int client_http_port = DEFAULT_CLIENT_HTTP_PORT) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_SETUP,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
auto real_client = self.init_real_client();
std::shared_ptr<mooncake::TransferEngine> transfer_engine =
nullptr;
if (!engine.is_none()) {
transfer_engine =
engine.cast<std::shared_ptr<TransferEngine>>();
}
return real_client->setup_real(
auto ret = real_client->setup_real(
local_hostname, metadata_server, global_segment_size,
local_buffer_size, protocol, rdma_devices,
master_server_addr, transfer_engine, "", enable_ssd_offload,
ssd_offload_path, tenant_id, enable_client_http_server,
client_http_port);
pt.End(ret == 0 ? 0 : -1);
// MC_LOG 在下沉层 setup_real 输出(Q1b)
return ret;
},
py::arg("local_hostname"), py::arg("metadata_server"),
py::arg("global_segment_size"), py::arg("local_buffer_size"),
Expand All @@ -2109,6 +2116,9 @@ PYBIND11_MODULE(store, m) {
.def(
"setup",
[](MooncakeStorePyWrapper &self, const py::dict &config_dict) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_SETUP,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
auto real_client = self.init_real_client();

// Convert py::dict to ConfigDict (all values as strings)
Expand All @@ -2120,8 +2130,11 @@ PYBIND11_MODULE(store, m) {
}

auto result = real_client->setup_internal(config);
return result.has_value() ? 0
: static_cast<int>(result.error());
int ret = result.has_value() ? 0
: static_cast<int>(result.error());
pt.End(ret == 0 ? 0 : -1);
// MC_LOG 在下沉层 setup_real 输出(Q1b)
return ret;
},
py::arg("config"),
"Setup the store with a configuration dictionary.\n"
Expand Down Expand Up @@ -2230,8 +2243,20 @@ PYBIND11_MODULE(store, m) {
.def(
"remove_all",
[](MooncakeStorePyWrapper &self, bool force) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_REMOVE_ALL,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
auto t0 = std::chrono::steady_clock::now();
py::gil_scoped_release release;
return self.store_->removeAll(force);
auto ret = self.store_->removeAll(force);
auto t1 = std::chrono::steady_clock::now();
auto elapsed_us = std::chrono::duration_cast<std::chrono::microseconds>(
t1 - t0).count();
pt.End(ret == 0 ? 0 : -1);
// 下沉 removeAll 不加 MC_LOG,入口层输出汇总(Q1b)
MC_LOG(INFO) << "[remove_all] elapsed_us=" << elapsed_us
<< " success=" << (ret == 0 ? 1 : 0);
return ret;
},
py::arg("force") = false,
"Remove all objects from the store. If force=True, skip lease "
Expand All @@ -2255,17 +2280,43 @@ PYBIND11_MODULE(store, m) {
"batch_is_exist",
[](MooncakeStorePyWrapper &self,
const std::vector<std::string> &keys) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_IS_EXIST,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
py::gil_scoped_release release;
return self.store_->batchIsExist(keys);
auto ret = self.store_->batchIsExist(keys);
pt.End(0);
// MC_LOG 在下沉层 batchIsExist 输出 per-key(Q1b)
return ret;
},
py::arg("keys"),
"Check if multiple objects exist. Returns list of results: 1 if "
"exists, 0 if not exists, -1 if error")
.def("close",
[](MooncakeStorePyWrapper &self) {
if (!self.store_) return 0;
SpDiag::PerfPoint pt(PerfKey::STORE_PY_CLOSE,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
auto t0 = std::chrono::steady_clock::now();
if (!self.store_) {
auto t1 = std::chrono::steady_clock::now();
auto elapsed_us = std::chrono::duration_cast<std::chrono::microseconds>(
t1 - t0).count();
pt.End(0);
// 无下沉 MC_LOG,入口层输出汇总(Q1b)
MC_LOG(INFO) << "[close] elapsed_us=" << elapsed_us
<< " success=1";
return 0;
}
int rc = self.store_->tearDownAll();
self.store_.reset();
auto t1 = std::chrono::steady_clock::now();
auto elapsed_us = std::chrono::duration_cast<std::chrono::microseconds>(
t1 - t0).count();
pt.End(rc == 0 ? 0 : -1);
// 下沉 tearDownAll 不加 MC_LOG,入口层输出汇总(Q1b)
MC_LOG(INFO) << "[close] elapsed_us=" << elapsed_us
<< " success=" << (rc == 0 ? 1 : 0);
return rc;
})
.def("health_check", &MooncakeStorePyWrapper::health_check,
Expand Down Expand Up @@ -2649,10 +2700,16 @@ PYBIND11_MODULE(store, m) {
"register_buffer",
[](MooncakeStorePyWrapper &self, uintptr_t buffer_ptr,
size_t size) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_REGISTER_BUFFER,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
// Register memory buffer for RDMA operations
void *buffer = reinterpret_cast<void *>(buffer_ptr);
py::gil_scoped_release release;
return self.store_->register_buffer(buffer, size);
auto ret = self.store_->register_buffer(buffer, size);
pt.End(ret == 0 ? 0 : -1);
// MC_LOG 在下沉层 register_buffer 输出汇总(Q1b)
return ret;
},
py::arg("buffer_ptr"), py::arg("size"),
"Register a memory buffer for direct access operations")
Expand Down Expand Up @@ -2897,13 +2954,20 @@ PYBIND11_MODULE(store, m) {
const std::vector<std::vector<uintptr_t>> &all_buffer_ptrs,
const std::vector<std::vector<size_t>> &all_sizes,
const ReplicateConfig &config = ReplicateConfig{}) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_PUT_MULTI,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
if (!self.is_client_initialized()) {
LOG(ERROR) << "Client is not initialized";
pt.End(-1);
return std::vector<int>{};
}
py::gil_scoped_release release;
return self.store_->batch_put_from_multi_buffers(
auto ret = self.store_->batch_put_from_multi_buffers(
keys, CastAddrs2Ptrs(all_buffer_ptrs), all_sizes, config);
pt.End(ret.empty() ? -1 : 0);
// MC_LOG 在下沉层 *_internal 输出 汇总+per-key(Q1b)
return ret;
},
py::arg("keys"), py::arg("all_buffer_ptrs"), py::arg("all_sizes"),
py::arg("config") = ReplicateConfig{},
Expand All @@ -2917,10 +2981,16 @@ PYBIND11_MODULE(store, m) {
const std::vector<std::vector<uintptr_t>> &all_buffer_ptrs,
const std::vector<std::vector<size_t>> &all_sizes,
bool prefer_alloc_in_same_node = false) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_GET_INTO_MULTI,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
py::gil_scoped_release release;
return self.store_->batch_get_into_multi_buffers(
auto ret = self.store_->batch_get_into_multi_buffers(
keys, CastAddrs2Ptrs(all_buffer_ptrs), all_sizes,
prefer_alloc_in_same_node);
pt.End(ret.empty() ? -1 : 0);
// MC_LOG 在下沉层 *_internal 输出 汇总+per-key(Q1b)
return ret;
},
py::arg("keys"), py::arg("all_buffer_ptrs"), py::arg("all_sizes"),
py::arg("prefer_alloc_in_same_node") = false,
Expand All @@ -2938,8 +3008,14 @@ PYBIND11_MODULE(store, m) {
"batch_get_replica_desc",
[](MooncakeStorePyWrapper &self,
const std::vector<std::string> &keys) {
SpDiag::PerfPoint pt(PerfKey::STORE_PY_BATCH_GET_REPLICA_DESC,
SpDiag::PerfLevel::KEY_MODULE);
pt.Start();
py::gil_scoped_release release;
return self.store_->batch_get_replica_desc(keys);
auto ret = self.store_->batch_get_replica_desc(keys);
pt.End(0);
// MC_LOG 在下沉层 batch_get_replica_desc 输出 per-key(Q1b)
return ret;
},
py::arg("keys"))
.def(
Expand Down
10 changes: 10 additions & 0 deletions mooncake-store/benchmarks/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,16 @@ target_link_libraries(
stress_cluster_bench PRIVATE mooncake_store transfer_engine asio_shared
gflags::gflags glog::glog pthread)

# Benchmark for vLLM Store Connector path
# Triggers: batch_put_from_multi_buffers / batchIsExist /
# batch_get_into_multi_buffers (and setup/register_buffer/tearDownAll).
# Used to verify SpDiag perf points and MC_LOG logging.
add_executable(store_connector_bench store_connector_bench.cpp)
target_link_libraries(
store_connector_bench PRIVATE mooncake_store transfer_engine asio_shared
gflags::gflags glog::glog pthread)


# Benchmark for RealClient::get_into_ranges with configurable value size,
# fragments per key and keys per query.
add_executable(stress_cluster_ranges_bench stress_cluster_ranges_bench.cpp)
Expand Down
Loading
Loading