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,514 changes: 1,389 additions & 125 deletions mooncake-store/benchmarks/stress_cluster_bench.cpp

Large diffs are not rendered by default.

809 changes: 809 additions & 0 deletions mooncake-store/go/examples/dummy_clients_test/main.go

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions mooncake-store/go/mooncakestore/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,4 +31,5 @@ var (
ErrBatchOp = errors.New("mooncakestore: batch operation failed")
ErrHostname = errors.New("mooncakestore: get hostname failed")
ErrInvalidArgument = errors.New("mooncakestore: invalid argument")
ErrBufferNotFound = errors.New("mooncakestore: buffer not found")
)
99 changes: 95 additions & 4 deletions mooncake-store/go/mooncakestore/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,11 @@ import (
"unsafe"
)

const (
MOONCAKE_CLIENT_REAL = C.MOONCAKE_CLIENT_REAL
MOONCAKE_CLIENT_DUMMY = C.MOONCAKE_CLIENT_DUMMY
)

// Store wraps a Mooncake Store client handle.
type Store struct {
handle C.mooncake_store_t
Expand All @@ -34,7 +39,13 @@ type Store struct {
// New creates a new Store instance. Call Setup before performing operations,
// and Close when done.
func New() (*Store, error) {
h := C.mooncake_store_create()
return NewWithType(C.MOONCAKE_CLIENT_REAL)
}

// NewWithType creates a new Store instance with the specified client type.
// Use MOONCAKE_CLIENT_REAL or MOONCAKE_CLIENT_DUMMY.
func NewWithType(clientType C.mooncake_client_type_t) (*Store, error) {
h := C.mooncake_store_create(clientType)
if h == nil {
return nil, ErrStoreNil
}
Expand All @@ -43,17 +54,24 @@ func New() (*Store, error) {

// Setup initialises the store client and connects to the cluster.
//
// Parameters:
// For real client (MOONCAKE_CLIENT_REAL), the relevant parameters are:
// - localHostname: hostname/IP of this node
// - metadataServer: metadata server URL (e.g. "http://host:8080/metadata")
// - globalSegmentSize: size of the global memory segment in bytes
// - localBufferSize: size of the local transfer buffer in bytes
// - protocol: transport protocol ("tcp" or "rdma")
// - deviceName: RDMA device name (empty for TCP or auto-discovery)
// - masterServerAddr: master service address (e.g. "host:50051")
//
// For dummy client (MOONCAKE_CLIENT_DUMMY), the relevant parameters are:
// - localBufferSize: size of the local transfer buffer in bytes
// - memPoolSize: size of the memory pool in bytes
// - serverAddress: server address for dummy client
// - ipcSocketPath: IPC socket path for dummy client
func (s *Store) Setup(localHostname, metadataServer string,
globalSegmentSize, localBufferSize uint64,
protocol, deviceName, masterServerAddr string) error {
protocol, deviceName, masterServerAddr string,
memPoolSize uint64, serverAddress, ipcSocketPath string) error {
if s.handle == nil {
return ErrStoreNil
}
Expand All @@ -63,16 +81,21 @@ func (s *Store) Setup(localHostname, metadataServer string,
cProtocol := C.CString(protocol)
cDeviceName := C.CString(deviceName)
cMasterAddr := C.CString(masterServerAddr)
cServerAddress := C.CString(serverAddress)
cIpcSocketPath := C.CString(ipcSocketPath)
defer C.free(unsafe.Pointer(cLocalHostname))
defer C.free(unsafe.Pointer(cMetadataServer))
defer C.free(unsafe.Pointer(cProtocol))
defer C.free(unsafe.Pointer(cDeviceName))
defer C.free(unsafe.Pointer(cMasterAddr))
defer C.free(unsafe.Pointer(cServerAddress))
defer C.free(unsafe.Pointer(cIpcSocketPath))

ret := C.mooncake_store_setup(s.handle,
cLocalHostname, cMetadataServer,
C.uint64_t(globalSegmentSize), C.uint64_t(localBufferSize),
cProtocol, cDeviceName, cMasterAddr)
cProtocol, cDeviceName, cMasterAddr,
C.uint64_t(memPoolSize), cServerAddress, cIpcSocketPath)
if ret != 0 {
return ErrSetupFailed
}
Expand Down Expand Up @@ -538,3 +561,71 @@ func (s *Store) UnregisterBuffer(ptr uintptr) error {
}
return nil
}

// ---------------------------------------------------------------------------
// DummyClient-specific: Query registered buffers
// ---------------------------------------------------------------------------

type RegisteredBufferInfo struct {
Ptr uintptr
Size uint64
IsLocal bool
}

func (s *Store) RegisteredBufferCount() (int, error) {
if s.handle == nil {
return 0, ErrStoreNil
}
count := C.mooncake_store_get_registered_buffer_count(s.handle)
return int(count), nil
}

func (s *Store) RegisteredBufferAt(index int) (*RegisteredBufferInfo, error) {
if s.handle == nil {
return nil, ErrStoreNil
}
if index < 0 {
return nil, ErrInvalidArgument
}
var size C.size_t
ptr := C.mooncake_store_get_registered_buffer_at(s.handle,
C.size_t(index), &size)
if ptr == nil {
return nil, ErrBufferNotFound
}
return &RegisteredBufferInfo{
Ptr: uintptr(ptr),
Size: uint64(size),
}, nil
}

// UnregisterAllBuffers 注销所有已注册的 buffer
// 用于程序结束前清理资源,避免内存泄漏
// 注意:此函数会遍历所有已注册 buffer 并逐个调用 UnregisterBuffer
func (s *Store) UnregisterAllBuffers() (int, error) {
if s.handle == nil {
return 0, ErrStoreNil
}
count, err := s.RegisteredBufferCount()
if err != nil {
return 0, err
}
unregistered := 0
for i := 0; i < count; i++ {
bufInfo, err := s.RegisteredBufferAt(i)
if err != nil {
continue
}
if err := s.UnregisterBuffer(bufInfo.Ptr); err == nil {
unregistered++
}
}
return unregistered, nil
}

func (s *Store) IsHotCachePtr(ptr uintptr) bool {
if s.handle == nil || ptr == 0 {
return false
}
return C.mooncake_store_is_hot_cache_ptr(s.handle, unsafe.Pointer(ptr)) == 1
}
8 changes: 8 additions & 0 deletions mooncake-store/include/dummy_client.h
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,14 @@ class DummyClient : public PyClient {

[[nodiscard]] std::string get_hostname() const;

struct RegisteredBufferInfo {
void *ptr;
size_t size;
bool is_local;
};

std::vector<RegisteredBufferInfo> get_registered_buffers() const;

// Check if a pointer falls within the hot cache shm region
bool is_hot_cache_ptr(const void *ptr) const {
if (!hot_cache_base_) return false;
Expand Down
4 changes: 4 additions & 0 deletions mooncake-store/include/real_client.h
Original file line number Diff line number Diff line change
Expand Up @@ -429,6 +429,10 @@ class RealClient : public PyClient {
const std::vector<size_t> &sizes, const ReplicateConfig &config,
int32_t device_id, const UUID &client_id);

tl::expected<void, ErrorCode> put_from_dummy_helper(
const std::string &key, uint64_t dummy_buffer, size_t size,
const ReplicateConfig &config, int32_t device_id, const UUID &client_id);

std::vector<tl::expected<void, ErrorCode>>
batch_put_from_multi_buffers_dummy_helper(
const std::vector<std::string> &keys,
Expand Down
24 changes: 22 additions & 2 deletions mooncake-store/include/store_c.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,12 @@ extern "C" {

typedef void *mooncake_store_t;

enum mooncake_client_type {
MOONCAKE_CLIENT_REAL = 0,
MOONCAKE_CLIENT_DUMMY = 1,
};
typedef enum mooncake_client_type mooncake_client_type_t;

struct mooncake_replicate_config {
size_t replica_num;
int with_soft_pin;
Expand All @@ -45,7 +51,7 @@ typedef struct mooncake_replicate_config mooncake_replicate_config_t;
// Lifecycle
// ---------------------------------------------------------------------------

mooncake_store_t mooncake_store_create();
mooncake_store_t mooncake_store_create(mooncake_client_type_t client_type);

void mooncake_store_destroy(mooncake_store_t store);

Expand All @@ -54,7 +60,10 @@ int mooncake_store_setup(mooncake_store_t store, const char *local_hostname,
uint64_t global_segment_size,
uint64_t local_buffer_size, const char *protocol,
const char *device_name,
const char *master_server_addr);
const char *master_server_addr,
uint64_t mem_pool_size,
const char *server_address,
const char *ipc_socket_path);

int mooncake_store_init_all(mooncake_store_t store, const char *protocol,
const char *device_name,
Expand Down Expand Up @@ -125,6 +134,17 @@ int mooncake_store_register_buffer(mooncake_store_t store, void *buffer,

int mooncake_store_unregister_buffer(mooncake_store_t store, void *buffer);

// ---------------------------------------------------------------------------
// DummyClient-specific: Query registered buffers
// ---------------------------------------------------------------------------

size_t mooncake_store_get_registered_buffer_count(mooncake_store_t store);

void *mooncake_store_get_registered_buffer_at(mooncake_store_t store,
size_t index, size_t *size_out);

int mooncake_store_is_hot_cache_ptr(mooncake_store_t store, const void *ptr);

#ifdef __cplusplus
}
#endif
Expand Down
20 changes: 18 additions & 2 deletions mooncake-store/src/dummy_client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1106,6 +1106,20 @@ std::string DummyClient::get_hostname() const {
return "";
}

std::vector<DummyClient::RegisteredBufferInfo>
DummyClient::get_registered_buffers() const {
std::vector<RegisteredBufferInfo> result;
if (!shm_helper_) return result;

const auto& shms = shm_helper_->get_shms();
for (const auto& shm : shms) {
if (shm->registered && shm->base_addr && shm->size > 0) {
result.push_back({shm->base_addr, shm->size, shm->is_local});
}
}
return result;
}

std::vector<int> DummyClient::batch_put_from(
const std::vector<std::string>& keys, const std::vector<void*>& buffer_ptrs,
const std::vector<size_t>& sizes, const ReplicateConfig& config) {
Expand Down Expand Up @@ -1133,8 +1147,10 @@ std::vector<int> DummyClient::batch_put_from(

int DummyClient::put_from(const std::string& key, void* buffer, size_t size,
const ReplicateConfig& config) {
// TODO: implement this function
return -1;
uint64_t buf_addr = reinterpret_cast<uint64_t>(buffer);
auto result = invoke_rpc<&RealClient::put_from_dummy_helper, void>(
key, buf_addr, size, config, device_id_, client_id_);
return to_py_ret(result);
}

std::vector<int64_t> DummyClient::batch_get_into(
Expand Down
49 changes: 47 additions & 2 deletions mooncake-store/src/real_client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -4073,6 +4073,33 @@ RealClient::batch_put_from_dummy_helper(
return batch_put_from_internal(keys, buffers_result.value(), sizes, config);
}

tl::expected<void, ErrorCode> RealClient::put_from_dummy_helper(
const std::string &key, uint64_t dummy_buffer, size_t size,
const ReplicateConfig &config, int32_t device_id, const UUID &client_id) {
#ifdef USE_ASCEND_DIRECT
if (!ContextManager::getInstance().setCurrentContextByPhysicalId(
device_id)) {
LOG(ERROR) << "Failed to set context for physical device " << device_id;
return tl::unexpected(ErrorCode::INVALID_PARAMS);
}
#endif

std::shared_lock<std::shared_mutex> lock(dummy_client_mutex_);
auto it = shm_contexts_.find(client_id);
if (it == shm_contexts_.end()) {
LOG(ERROR) << "client_id=" << client_id << ", error=shm_not_mapped";
return tl::unexpected(ErrorCode::INVALID_PARAMS);
}
auto &context = it->second;

auto buffers_result = map_dummy_addrs_to_real_ptrs(
context, {dummy_buffer}, {size}, client_id);
if (!buffers_result) {
return tl::unexpected(buffers_result.error());
}
return put_from_internal(key, buffers_result.value()[0], size, config);
}

std::vector<tl::expected<void, ErrorCode>> RealClient::batch_put_from_internal(
const std::vector<std::string> &keys, const std::vector<void *> &buffers,
const std::vector<size_t> &sizes, const ReplicateConfig &config) {
Expand Down Expand Up @@ -5903,11 +5930,29 @@ void RealClient::dummy_client_monitor_func() {
// Update the client status to NEED_REMOUNT
if (!expired_clients.empty()) {
for (auto &client_id : expired_clients) {
{
std::shared_lock<std::shared_mutex> lock(
dummy_client_mutex_);
if (shm_contexts_.find(client_id) ==
shm_contexts_.end()) {
LOG(INFO)
<< "client_id=" << client_id
<< ", action=shm_already_unmapped_by_other_path";
continue;
}
}

// Unmap mapped_shms associated with this client
tl::expected<void, ErrorCode> result;
if (globalConfig().ascend_agent_mode) {
ascend_unmap_shm_internal(client_id);
result = ascend_unmap_shm_internal(client_id);
} else {
unmap_shm_internal(client_id);
result = unmap_shm_internal(client_id);
}
if (!result) {
// Client already unmapped (e.g., by other thread or earlier cleanup)
LOG(INFO) << "client_id=" << client_id
<< ", action=client_already_unmapped";
}
}
}
Expand Down
1 change: 1 addition & 0 deletions mooncake-store/src/real_client_main.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ void RegisterClientRpcService(coro_rpc::coro_rpc_server &server,
server.register_handler<&RealClient::getSize_internal>(&real_client);
server.register_handler<&RealClient::batch_put_from_dummy_helper>(
&real_client);
server.register_handler<&RealClient::put_from_dummy_helper>(&real_client);
server.register_handler<
&RealClient::batch_put_from_multi_buffers_dummy_helper>(&real_client);
server.register_handler<&RealClient::upsert_dummy_helper>(&real_client);
Expand Down
Loading
Loading