From 619f900b124694c238717f68cdc0078ed99296cf Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Thu, 9 Jul 2026 16:05:11 +0000 Subject: [PATCH 1/5] feat(store): retain object groups with server-owned leases --- mooncake-integration/store/store_py.cpp | 16 +++++ mooncake-store/include/client_service.h | 3 + mooncake-store/include/master_client.h | 3 + mooncake-store/include/master_service.h | 4 ++ mooncake-store/include/real_client.h | 6 ++ mooncake-store/include/rpc_service.h | 4 ++ mooncake-store/src/client_service.cpp | 10 +++ mooncake-store/src/master_client.cpp | 11 +++ mooncake-store/src/master_service.cpp | 71 ++++++++++++++++++++ mooncake-store/src/real_client.cpp | 22 ++++++ mooncake-store/src/rpc_service.cpp | 8 +++ mooncake-store/tests/master_service_test.cpp | 45 +++++++++++++ 12 files changed, 203 insertions(+) diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index 15ca9bac56..a70cbb16fd 100644 --- a/mooncake-integration/store/store_py.cpp +++ b/mooncake-integration/store/store_py.cpp @@ -2289,6 +2289,22 @@ PYBIND11_MODULE(store, m) { py::arg("keys"), "Check if multiple objects exist. Returns list of results: 1 if " "exists, 0 if not exists, -1 if error") + .def( + "retain_groups", + [](MooncakeStorePyWrapper &self, + const std::vector &group_ids, uint64_t ttl_ms) { + if (!self.real_client_) { + return std::vector( + group_ids.size(), + static_cast(ErrorCode::INVALID_PARAMS)); + } + py::gil_scoped_release release; + return self.real_client_->retainGroups(group_ids, ttl_ms); + }, + py::arg("group_ids"), py::arg("ttl_ms"), + "Retain all complete objects in each group for the requested TTL. " + "Returns 1 if retained, 0 if absent or incomplete, and a negative " + "error code on failure.") .def("close", [](MooncakeStorePyWrapper &self) { if (!self.store_) return 0; diff --git a/mooncake-store/include/client_service.h b/mooncake-store/include/client_service.h index 621b4edd70..48e5265427 100644 --- a/mooncake-store/include/client_service.h +++ b/mooncake-store/include/client_service.h @@ -363,6 +363,9 @@ class Client { std::vector> BatchIsExist( const std::vector& keys); + std::vector> RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms); + /** * @brief Create a copy task to copy an object's replicas to target segments * @param key Object key diff --git a/mooncake-store/include/master_client.h b/mooncake-store/include/master_client.h index 7b54123044..ef85a26854 100644 --- a/mooncake-store/include/master_client.h +++ b/mooncake-store/include/master_client.h @@ -120,6 +120,9 @@ class MasterClient { [[nodiscard]] std::vector> BatchExistKey( const std::vector& object_keys); + [[nodiscard]] std::vector> RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms); + /** * @brief Calculate Store-observed cache reuse metrics * @param object_keys None diff --git a/mooncake-store/include/master_service.h b/mooncake-store/include/master_service.h index 42fd27efa3..ab015da249 100644 --- a/mooncake-store/include/master_service.h +++ b/mooncake-store/include/master_service.h @@ -205,6 +205,10 @@ class MasterService { std::vector> BatchExistKey( const std::vector& keys, const std::string& tenant_id); + std::vector> RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms, + const std::string& tenant_id); + /** * @brief Fetch all keys for a single tenant. * @return ErrorCode::OK if exists diff --git a/mooncake-store/include/real_client.h b/mooncake-store/include/real_client.h index 870fa9ea7e..05cd4a1a11 100644 --- a/mooncake-store/include/real_client.h +++ b/mooncake-store/include/real_client.h @@ -314,6 +314,9 @@ class RealClient : public PyClient { */ std::vector batchIsExist(const std::vector &keys); + std::vector retainGroups(const std::vector &group_ids, + uint64_t ttl_ms); + /** * @brief Get the size of an object * @param key Key of the object @@ -655,6 +658,9 @@ class RealClient : public PyClient { std::vector> batchIsExist_internal( const std::vector &keys); + std::vector> retainGroups_internal( + const std::vector &group_ids, uint64_t ttl_ms); + tl::expected getSize_internal(const std::string &key); std::shared_ptr get_buffer_internal( diff --git a/mooncake-store/include/rpc_service.h b/mooncake-store/include/rpc_service.h index 8e2cfa4895..06be74d799 100644 --- a/mooncake-store/include/rpc_service.h +++ b/mooncake-store/include/rpc_service.h @@ -44,6 +44,10 @@ class WrappedMasterService { const std::vector& keys, const std::string& tenant_id = "default"); + std::vector> RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms, + const std::string& tenant_id = "default"); + tl::expected< std::unordered_map, boost::hash>, ErrorCode> diff --git a/mooncake-store/src/client_service.cpp b/mooncake-store/src/client_service.cpp index 446503c91a..474022d2dd 100644 --- a/mooncake-store/src/client_service.cpp +++ b/mooncake-store/src/client_service.cpp @@ -3013,6 +3013,16 @@ std::vector> Client::BatchIsExist( return response; } +std::vector> Client::RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms) { + auto response = master_client_.RetainGroups(group_ids, ttl_ms); + if (response.size() != group_ids.size()) { + return std::vector>( + group_ids.size(), tl::unexpected(ErrorCode::RPC_FAIL)); + } + return response; +} + void* Client::GetBaseAddr() { return transfer_engine_->getBaseAddr(); } tl::expected Client::MountLocalDiskSegment( diff --git a/mooncake-store/src/master_client.cpp b/mooncake-store/src/master_client.cpp index e2d9db347f..fedd478f77 100644 --- a/mooncake-store/src/master_client.cpp +++ b/mooncake-store/src/master_client.cpp @@ -32,6 +32,11 @@ struct RpcNameTraits<&WrappedMasterService::BatchExistKey> { static constexpr const char* value = "BatchExistKey"; }; +template <> +struct RpcNameTraits<&WrappedMasterService::RetainGroups> { + static constexpr const char* value = "RetainGroups"; +}; + template <> struct RpcNameTraits<&WrappedMasterService::GetReplicaList> { static constexpr const char* value = "GetReplicaList"; @@ -469,6 +474,12 @@ std::vector> MasterClient::BatchExistKey( return result; } +std::vector> MasterClient::RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms) { + return invoke_batch_rpc<&WrappedMasterService::RetainGroups, bool>( + group_ids.size(), group_ids, ttl_ms, tenant_id_); +} + tl::expected MasterClient::CalcCacheStats() { return invoke_rpc<&WrappedMasterService::CalcCacheStats, diff --git a/mooncake-store/src/master_service.cpp b/mooncake-store/src/master_service.cpp index bf55656980..3f0fe75e57 100644 --- a/mooncake-store/src/master_service.cpp +++ b/mooncake-store/src/master_service.cpp @@ -2176,6 +2176,77 @@ std::vector> MasterService::BatchExistKey( return results; } +std::vector> MasterService::RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms, + const std::string& tenant_id) { + std::vector> results(group_ids.size()); + if (group_ids.empty()) { + return results; + } + if (ttl_ms == 0) { + std::fill(results.begin(), results.end(), + tl::unexpected(ErrorCode::INVALID_PARAMS)); + return results; + } + + const std::string normalized_tenant = NormalizeRequestTenantId(tenant_id); + std::vector> indices_by_shard(kNumShards); + for (size_t i = 0; i < group_ids.size(); ++i) { + indices_by_shard[getShardIndex(group_ids[i])].push_back(i); + } + + const size_t start_shard = RandomIndex(kNumShards); + for (size_t scanned = 0; scanned < kNumShards; ++scanned) { + const size_t shard_idx = + (start_shard + kNumShards - scanned) % kNumShards; + const auto& group_indices = indices_by_shard[shard_idx]; + if (group_indices.empty()) { + continue; + } + + std::shared_lock shared_lock(snapshot_mutex_); + MetadataShardAccessorRO shard(this, shard_idx); + auto tenant_it = shard->tenants.find(normalized_tenant); + if (tenant_it == shard->tenants.end()) { + for (const size_t i : group_indices) { + results[i] = false; + } + continue; + } + + const auto& tenant_state = tenant_it->second; + for (const size_t i : group_indices) { + auto group_it = tenant_state.group_members.find(group_ids[i]); + if (group_it == tenant_state.group_members.end() || + group_it->second.empty()) { + results[i] = false; + continue; + } + + bool complete = true; + for (const auto& member_key : group_it->second) { + auto member_it = tenant_state.metadata.find(member_key); + if (member_it == tenant_state.metadata.end() || + !member_it->second.IsValid() || + !member_it->second.HasReplica(&Replica::fn_is_completed)) { + complete = false; + break; + } + } + if (!complete) { + results[i] = false; + continue; + } + + for (const auto& member_key : group_it->second) { + tenant_state.metadata.at(member_key).GrantLease(ttl_ms, 0); + } + results[i] = true; + } + } + return results; +} + auto MasterService::GetAllKeys(const std::string& tenant_id) -> tl::expected, ErrorCode> { std::vector all_keys; diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index fe00d17224..fadf1ac836 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -2053,6 +2053,19 @@ std::vector RealClient::batchIsExist( return results; } +std::vector RealClient::retainGroups( + const std::vector &group_ids, uint64_t ttl_ms) { + auto internal_results = retainGroups_internal(group_ids, ttl_ms); + std::vector results; + results.reserve(internal_results.size()); + for (const auto &result : internal_results) { + results.push_back(result.has_value() + ? (result.value() ? 1 : 0) + : toInt(result.error())); + } + return results; +} + tl::expected RealClient::getSize_internal( const std::string &key) { if (!client_) { @@ -4698,6 +4711,15 @@ std::vector> RealClient::batchIsExist_internal( return client_->BatchIsExist(keys); } +std::vector> RealClient::retainGroups_internal( + const std::vector &group_ids, uint64_t ttl_ms) { + if (!client_ || ttl_ms == 0) { + return std::vector>( + group_ids.size(), tl::unexpected(ErrorCode::INVALID_PARAMS)); + } + return client_->RetainGroups(group_ids, ttl_ms); +} + int RealClient::put_from_with_metadata(const std::string &key, void *buffer, void *metadata_buffer, size_t size, size_t metadata_size, diff --git a/mooncake-store/src/rpc_service.cpp b/mooncake-store/src/rpc_service.cpp index d5de362f54..439f20786f 100644 --- a/mooncake-store/src/rpc_service.cpp +++ b/mooncake-store/src/rpc_service.cpp @@ -76,6 +76,12 @@ std::vector> WrappedMasterService::BatchExistKey( return result; } +std::vector> WrappedMasterService::RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms, + const std::string& tenant_id) { + return master_service_.RetainGroups(group_ids, ttl_ms, tenant_id); +} + tl::expected< std::unordered_map, boost::hash>, ErrorCode> @@ -1349,6 +1355,8 @@ void RegisterRpcService( &wrapped_master_service); server.register_handler<&mooncake::WrappedMasterService::BatchExistKey>( &wrapped_master_service); + server.register_handler<&mooncake::WrappedMasterService::RetainGroups>( + &wrapped_master_service); server.register_handler<&mooncake::WrappedMasterService::ServiceReady>( &wrapped_master_service); server.register_handler< diff --git a/mooncake-store/tests/master_service_test.cpp b/mooncake-store/tests/master_service_test.cpp index 46512a5749..94de677fe7 100644 --- a/mooncake-store/tests/master_service_test.cpp +++ b/mooncake-store/tests/master_service_test.cpp @@ -1441,6 +1441,51 @@ TEST_F(MasterServiceTest, } } +TEST_F(MasterServiceTest, RetainGroupsExtendsTheWholeGroupLease) { + auto service_config = + MasterServiceConfig::builder().set_default_kv_lease_ttl(100).build(); + constexpr size_t kSegmentSize = 4 * 1024 * 1024; + constexpr size_t kObjectSize = 2 * 1024 * 1024; + std::unique_ptr service_(new MasterService(service_config)); + [[maybe_unused]] const auto context = + PrepareSimpleSegment(*service_, "retained_group_segment", + kDefaultSegmentBase, kSegmentSize); + const UUID client_id = generate_uuid(); + + const std::string key_a = "retained_group_key_a"; + const std::string key_b = "retained_group_key_b"; + const std::string group_id = FindGroupIdOnDifferentShard(key_a); + ReplicateConfig config; + config.replica_num = 1; + config.group_ids = std::vector{group_id}; + PutCompletedObject(*service_, client_id, key_a, config, kObjectSize); + PutCompletedObject(*service_, client_id, key_b, config, kObjectSize); + + auto retained = service_->RetainGroups({group_id, "missing_group"}, 1000, + "default"); + ASSERT_EQ(retained.size(), 2); + ASSERT_TRUE(retained[0].has_value()); + EXPECT_TRUE(retained[0].value()); + ASSERT_TRUE(retained[1].has_value()); + EXPECT_FALSE(retained[1].value()); + + std::this_thread::sleep_for(std::chrono::milliseconds(200)); + ReplicateConfig trigger_config; + trigger_config.replica_num = 1; + auto trigger_result = service_->PutStart( + client_id, "trigger_retained_group_eviction", "default", kObjectSize, + trigger_config); + ASSERT_FALSE(trigger_result.has_value()); + EXPECT_EQ(ErrorCode::NO_AVAILABLE_HANDLE, trigger_result.error()); + EXPECT_TRUE(service_->GetReplicaList(key_a, "default").has_value()); + EXPECT_TRUE(service_->GetReplicaList(key_b, "default").has_value()); + + auto invalid = service_->RetainGroups({group_id}, 0, "default"); + ASSERT_EQ(invalid.size(), 1); + ASSERT_FALSE(invalid[0].has_value()); + EXPECT_EQ(ErrorCode::INVALID_PARAMS, invalid[0].error()); +} + TEST_F(MasterServiceTest, GroupedEvictionSkipsUnsafeMembersAndEvictsSafePeers) { constexpr size_t kSegmentSize = 4 * 1024 * 1024; constexpr size_t kObjectSize = 2 * 1024 * 1024; From c72ad373be35c2901eca4fcb6b902de2169d2484 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Thu, 9 Jul 2026 16:15:16 +0000 Subject: [PATCH 2/5] fix(store): retain groups before members arrive --- mooncake-integration/store/store_py.cpp | 5 +- mooncake-store/include/master_service.h | 7 +- mooncake-store/src/master_service.cpp | 107 +++++++++++++------ mooncake-store/tests/master_service_test.cpp | 9 +- 4 files changed, 89 insertions(+), 39 deletions(-) diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index a70cbb16fd..b0b1b92753 100644 --- a/mooncake-integration/store/store_py.cpp +++ b/mooncake-integration/store/store_py.cpp @@ -2302,9 +2302,8 @@ PYBIND11_MODULE(store, m) { return self.real_client_->retainGroups(group_ids, ttl_ms); }, py::arg("group_ids"), py::arg("ttl_ms"), - "Retain all complete objects in each group for the requested TTL. " - "Returns 1 if retained, 0 if absent or incomplete, and a negative " - "error code on failure.") + "Retain current and future objects in each group for the requested " + "TTL. Returns 1 if accepted and a negative error code on failure.") .def("close", [](MooncakeStorePyWrapper &self) { if (!self.store_) return 0; diff --git a/mooncake-store/include/master_service.h b/mooncake-store/include/master_service.h index ab015da249..14468339dc 100644 --- a/mooncake-store/include/master_service.h +++ b/mooncake-store/include/master_service.h @@ -1222,11 +1222,15 @@ class MasterService { std::unordered_map> group_members; // group_id → set of keys + std::unordered_map + group_retention_deadlines; bool Empty() const { return metadata.empty() && processing_keys.empty() && replication_tasks.empty() && offloading_tasks.empty() && - promotion_tasks.empty() && group_members.empty(); + promotion_tasks.empty() && group_members.empty() && + group_retention_deadlines.empty(); } }; @@ -1401,6 +1405,7 @@ class MasterService { const std::string& tenant_id, const std::string& key, const std::string& group_id); + void PruneExpiredGroupRetentions(); std::unordered_map::iterator EraseMetadata( TenantState& tenant_state, std::unordered_map::iterator it, diff --git a/mooncake-store/src/master_service.cpp b/mooncake-store/src/master_service.cpp index 3f0fe75e57..7e3f77f030 100644 --- a/mooncake-store/src/master_service.cpp +++ b/mooncake-store/src/master_service.cpp @@ -1606,11 +1606,32 @@ void MasterService::RegisterGroupMember(TenantState& tenant_state, return; } const auto& normalized_tenant = NormalizeTenantIdRef(tenant_id); - std::unique_lock lock(group_routing_mutex_); - object_group_ids_[MakeTenantScopedKey(normalized_tenant, key)] = group_id; - groups_needing_lease_refresh_.insert( - MakeTenantScopedKey(normalized_tenant, group_id)); + { + std::unique_lock lock(group_routing_mutex_); + object_group_ids_[MakeTenantScopedKey(normalized_tenant, key)] = + group_id; + groups_needing_lease_refresh_.insert( + MakeTenantScopedKey(normalized_tenant, group_id)); + } tenant_state.group_members[group_id].insert(key); + + auto retention_it = tenant_state.group_retention_deadlines.find(group_id); + if (retention_it == tenant_state.group_retention_deadlines.end()) { + return; + } + const auto now = std::chrono::system_clock::now(); + if (retention_it->second <= now) { + tenant_state.group_retention_deadlines.erase(retention_it); + return; + } + auto metadata_it = tenant_state.metadata.find(key); + if (metadata_it != tenant_state.metadata.end()) { + const auto remaining_ms = + std::chrono::duration_cast( + retention_it->second - now) + .count(); + metadata_it->second.GrantLease(static_cast(remaining_ms), 0); + } } void MasterService::UnregisterGroupMember(TenantState& tenant_state, @@ -1951,9 +1972,12 @@ void MasterService::TaskCleanupThreadFunc() { } std::shared_lock shared_lock(snapshot_mutex_); - auto write_access = task_manager_.get_write_access(); - write_access.prune_expired_tasks(); - write_access.prune_finished_tasks(); + { + auto write_access = task_manager_.get_write_access(); + write_access.prune_expired_tasks(); + write_access.prune_finished_tasks(); + } + PruneExpiredGroupRetentions(); } LOG(INFO) << "Task cleanup thread stopped"; } @@ -2205,48 +2229,69 @@ std::vector> MasterService::RetainGroups( } std::shared_lock shared_lock(snapshot_mutex_); - MetadataShardAccessorRO shard(this, shard_idx); - auto tenant_it = shard->tenants.find(normalized_tenant); - if (tenant_it == shard->tenants.end()) { - for (const size_t i : group_indices) { - results[i] = false; + MetadataShardAccessorRW shard(this, shard_idx); + auto& tenant_state = shard->tenants[normalized_tenant]; + const auto now = std::chrono::system_clock::now(); + for (auto it = tenant_state.group_retention_deadlines.begin(); + it != tenant_state.group_retention_deadlines.end();) { + if (it->second <= now) { + it = tenant_state.group_retention_deadlines.erase(it); + } else { + ++it; } - continue; } - - const auto& tenant_state = tenant_it->second; for (const size_t i : group_indices) { + if (group_ids[i].empty()) { + results[i] = tl::unexpected(ErrorCode::INVALID_PARAMS); + continue; + } + const auto deadline = now + std::chrono::milliseconds(ttl_ms); + auto& retained_until = + tenant_state.group_retention_deadlines[group_ids[i]]; + retained_until = std::max(retained_until, deadline); + auto group_it = tenant_state.group_members.find(group_ids[i]); if (group_it == tenant_state.group_members.end() || group_it->second.empty()) { - results[i] = false; + results[i] = true; continue; } - - bool complete = true; for (const auto& member_key : group_it->second) { auto member_it = tenant_state.metadata.find(member_key); - if (member_it == tenant_state.metadata.end() || - !member_it->second.IsValid() || - !member_it->second.HasReplica(&Replica::fn_is_completed)) { - complete = false; - break; + if (member_it != tenant_state.metadata.end()) { + member_it->second.GrantLease(ttl_ms, 0); } } - if (!complete) { - results[i] = false; - continue; - } - - for (const auto& member_key : group_it->second) { - tenant_state.metadata.at(member_key).GrantLease(ttl_ms, 0); - } results[i] = true; } } return results; } +void MasterService::PruneExpiredGroupRetentions() { + const auto now = std::chrono::system_clock::now(); + for (size_t shard_idx = 0; shard_idx < kNumShards; ++shard_idx) { + MetadataShardAccessorRW shard(this, shard_idx); + for (auto tenant_it = shard->tenants.begin(); + tenant_it != shard->tenants.end();) { + auto& tenant_state = tenant_it->second; + for (auto it = tenant_state.group_retention_deadlines.begin(); + it != tenant_state.group_retention_deadlines.end();) { + if (it->second <= now) { + it = tenant_state.group_retention_deadlines.erase(it); + } else { + ++it; + } + } + if (tenant_state.Empty()) { + tenant_it = shard->tenants.erase(tenant_it); + } else { + ++tenant_it; + } + } + } +} + auto MasterService::GetAllKeys(const std::string& tenant_id) -> tl::expected, ErrorCode> { std::vector all_keys; diff --git a/mooncake-store/tests/master_service_test.cpp b/mooncake-store/tests/master_service_test.cpp index 94de677fe7..e0c9e32f79 100644 --- a/mooncake-store/tests/master_service_test.cpp +++ b/mooncake-store/tests/master_service_test.cpp @@ -1441,7 +1441,7 @@ TEST_F(MasterServiceTest, } } -TEST_F(MasterServiceTest, RetainGroupsExtendsTheWholeGroupLease) { +TEST_F(MasterServiceTest, RetainGroupsProtectsFutureGroupMembers) { auto service_config = MasterServiceConfig::builder().set_default_kv_lease_ttl(100).build(); constexpr size_t kSegmentSize = 4 * 1024 * 1024; @@ -1458,8 +1458,6 @@ TEST_F(MasterServiceTest, RetainGroupsExtendsTheWholeGroupLease) { ReplicateConfig config; config.replica_num = 1; config.group_ids = std::vector{group_id}; - PutCompletedObject(*service_, client_id, key_a, config, kObjectSize); - PutCompletedObject(*service_, client_id, key_b, config, kObjectSize); auto retained = service_->RetainGroups({group_id, "missing_group"}, 1000, "default"); @@ -1467,7 +1465,10 @@ TEST_F(MasterServiceTest, RetainGroupsExtendsTheWholeGroupLease) { ASSERT_TRUE(retained[0].has_value()); EXPECT_TRUE(retained[0].value()); ASSERT_TRUE(retained[1].has_value()); - EXPECT_FALSE(retained[1].value()); + EXPECT_TRUE(retained[1].value()); + + PutCompletedObject(*service_, client_id, key_a, config, kObjectSize); + PutCompletedObject(*service_, client_id, key_b, config, kObjectSize); std::this_thread::sleep_for(std::chrono::milliseconds(200)); ReplicateConfig trigger_config; From 1b9ffd5e62157c1b75e623dc421120b07e059c4f Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Thu, 9 Jul 2026 16:32:41 +0000 Subject: [PATCH 3/5] feat(store): log group retention outcomes --- mooncake-store/src/rpc_service.cpp | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/mooncake-store/src/rpc_service.cpp b/mooncake-store/src/rpc_service.cpp index 439f20786f..f322e1293f 100644 --- a/mooncake-store/src/rpc_service.cpp +++ b/mooncake-store/src/rpc_service.cpp @@ -1,5 +1,7 @@ #include "rpc_service.h" +#include + #include #include @@ -79,7 +81,14 @@ std::vector> WrappedMasterService::BatchExistKey( std::vector> WrappedMasterService::RetainGroups( const std::vector& group_ids, uint64_t ttl_ms, const std::string& tenant_id) { - return master_service_.RetainGroups(group_ids, ttl_ms, tenant_id); + auto result = master_service_.RetainGroups(group_ids, ttl_ms, tenant_id); + const auto accepted = std::count_if( + result.begin(), result.end(), + [](const auto& item) { return item.has_value() && item.value(); }); + LOG(INFO) << "RetainGroups accepted " << accepted << "/" + << group_ids.size() << " groups for ttl_ms=" << ttl_ms + << ", tenant_id=" << tenant_id; + return result; } tl::expected< From 56974e371ce39b463dcd461b30b895e7e919cde3 Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Thu, 9 Jul 2026 20:30:06 +0000 Subject: [PATCH 4/5] feat(store): bound group retention admission --- mooncake-integration/store/store_py.cpp | 3 +- mooncake-store/include/master_config.h | 16 ++++ mooncake-store/include/master_service.h | 3 + mooncake-store/include/types.h | 3 + mooncake-store/src/master_service.cpp | 33 ++++++- mooncake-store/tests/master_service_test.cpp | 97 ++++++++++++++++++++ 6 files changed, 150 insertions(+), 5 deletions(-) diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index b0b1b92753..7ed6e9d3bb 100644 --- a/mooncake-integration/store/store_py.cpp +++ b/mooncake-integration/store/store_py.cpp @@ -2303,7 +2303,8 @@ PYBIND11_MODULE(store, m) { }, py::arg("group_ids"), py::arg("ttl_ms"), "Retain current and future objects in each group for the requested " - "TTL. Returns 1 if accepted and a negative error code on failure.") + "TTL. Returns 1 if accepted, 0 if admission is full, and a " + "negative error code on failure.") .def("close", [](MooncakeStorePyWrapper &self) { if (!self.store_) return 0; diff --git a/mooncake-store/include/master_config.h b/mooncake-store/include/master_config.h index fd893953f8..62fa997c06 100644 --- a/mooncake-store/include/master_config.h +++ b/mooncake-store/include/master_config.h @@ -694,6 +694,8 @@ class MasterServiceConfigBuilder { private: uint64_t default_kv_lease_ttl_ = DEFAULT_DEFAULT_KV_LEASE_TTL; uint64_t default_kv_soft_pin_ttl_ = DEFAULT_KV_SOFT_PIN_TTL_MS; + size_t max_retained_groups_ = DEFAULT_MAX_RETAINED_GROUPS; + uint64_t max_group_retention_ttl_ms_ = DEFAULT_MAX_GROUP_RETENTION_TTL_MS; bool allow_evict_soft_pinned_objects_ = DEFAULT_ALLOW_EVICT_SOFT_PINNED_OBJECTS; double eviction_ratio_ = DEFAULT_EVICTION_RATIO; @@ -760,6 +762,16 @@ class MasterServiceConfigBuilder { return *this; } + MasterServiceConfigBuilder& set_max_retained_groups(size_t max_groups) { + max_retained_groups_ = max_groups; + return *this; + } + + MasterServiceConfigBuilder& set_max_group_retention_ttl_ms(uint64_t ttl) { + max_group_retention_ttl_ms_ = ttl; + return *this; + } + MasterServiceConfigBuilder& set_allow_evict_soft_pinned_objects( bool allow) { allow_evict_soft_pinned_objects_ = allow; @@ -1032,6 +1044,8 @@ class MasterServiceConfig { public: uint64_t default_kv_lease_ttl = DEFAULT_DEFAULT_KV_LEASE_TTL; uint64_t default_kv_soft_pin_ttl = DEFAULT_KV_SOFT_PIN_TTL_MS; + size_t max_retained_groups = DEFAULT_MAX_RETAINED_GROUPS; + uint64_t max_group_retention_ttl_ms = DEFAULT_MAX_GROUP_RETENTION_TTL_MS; bool allow_evict_soft_pinned_objects = DEFAULT_ALLOW_EVICT_SOFT_PINNED_OBJECTS; double eviction_ratio = DEFAULT_EVICTION_RATIO; @@ -1202,6 +1216,8 @@ inline MasterServiceConfig MasterServiceConfigBuilder::build() const { MasterServiceConfig config; config.default_kv_lease_ttl = default_kv_lease_ttl_; config.default_kv_soft_pin_ttl = default_kv_soft_pin_ttl_; + config.max_retained_groups = max_retained_groups_; + config.max_group_retention_ttl_ms = max_group_retention_ttl_ms_; config.allow_evict_soft_pinned_objects = allow_evict_soft_pinned_objects_; config.eviction_ratio = eviction_ratio_; config.eviction_high_watermark_ratio = eviction_high_watermark_ratio_; diff --git a/mooncake-store/include/master_service.h b/mooncake-store/include/master_service.h index 14468339dc..db3d78b623 100644 --- a/mooncake-store/include/master_service.h +++ b/mooncake-store/include/master_service.h @@ -1547,6 +1547,9 @@ class MasterService { // Lease related members const uint64_t default_kv_lease_ttl_; // in milliseconds const uint64_t default_kv_soft_pin_ttl_; // in milliseconds + const size_t max_retained_groups_; + const uint64_t max_group_retention_ttl_ms_; + std::atomic retained_group_count_{0}; const bool allow_evict_soft_pinned_objects_; // Eviction related members diff --git a/mooncake-store/include/types.h b/mooncake-store/include/types.h index bee73e2b1e..6013da1343 100644 --- a/mooncake-store/include/types.h +++ b/mooncake-store/include/types.h @@ -86,6 +86,9 @@ static constexpr uint64_t DEFAULT_DEFAULT_KV_LEASE_TTL = 5000; // in milliseconds static constexpr uint64_t DEFAULT_KV_SOFT_PIN_TTL_MS = 30 * 60 * 1000; // 30 minutes +static constexpr size_t DEFAULT_MAX_RETAINED_GROUPS = 65536; +static constexpr uint64_t DEFAULT_MAX_GROUP_RETENTION_TTL_MS = + 24 * 60 * 60 * 1000; // 24 hours static constexpr bool DEFAULT_ALLOW_EVICT_SOFT_PINNED_OBJECTS = true; static constexpr double DEFAULT_EVICTION_RATIO = 0.05; static constexpr double DEFAULT_EVICTION_HIGH_WATERMARK_RATIO = 0.95; diff --git a/mooncake-store/src/master_service.cpp b/mooncake-store/src/master_service.cpp index 7e3f77f030..a059c5fae4 100644 --- a/mooncake-store/src/master_service.cpp +++ b/mooncake-store/src/master_service.cpp @@ -192,6 +192,8 @@ MasterService::MasterService(const MasterServiceConfig& config) }), default_kv_lease_ttl_(config.default_kv_lease_ttl), default_kv_soft_pin_ttl_(config.default_kv_soft_pin_ttl), + max_retained_groups_(config.max_retained_groups), + max_group_retention_ttl_ms_(config.max_group_retention_ttl_ms), allow_evict_soft_pinned_objects_(config.allow_evict_soft_pinned_objects), eviction_ratio_(config.eviction_ratio), eviction_high_watermark_ratio_(config.eviction_high_watermark_ratio), @@ -1622,6 +1624,7 @@ void MasterService::RegisterGroupMember(TenantState& tenant_state, const auto now = std::chrono::system_clock::now(); if (retention_it->second <= now) { tenant_state.group_retention_deadlines.erase(retention_it); + retained_group_count_.fetch_sub(1, std::memory_order_relaxed); return; } auto metadata_it = tenant_state.metadata.find(key); @@ -2207,11 +2210,15 @@ std::vector> MasterService::RetainGroups( if (group_ids.empty()) { return results; } - if (ttl_ms == 0) { + if (ttl_ms == 0 || ttl_ms > max_group_retention_ttl_ms_) { std::fill(results.begin(), results.end(), tl::unexpected(ErrorCode::INVALID_PARAMS)); return results; } + if (max_retained_groups_ == 0) { + std::fill(results.begin(), results.end(), false); + return results; + } const std::string normalized_tenant = NormalizeRequestTenantId(tenant_id); std::vector> indices_by_shard(kNumShards); @@ -2236,6 +2243,7 @@ std::vector> MasterService::RetainGroups( it != tenant_state.group_retention_deadlines.end();) { if (it->second <= now) { it = tenant_state.group_retention_deadlines.erase(it); + retained_group_count_.fetch_sub(1, std::memory_order_relaxed); } else { ++it; } @@ -2246,9 +2254,23 @@ std::vector> MasterService::RetainGroups( continue; } const auto deadline = now + std::chrono::milliseconds(ttl_ms); - auto& retained_until = - tenant_state.group_retention_deadlines[group_ids[i]]; - retained_until = std::max(retained_until, deadline); + auto retention_it = + tenant_state.group_retention_deadlines.find(group_ids[i]); + if (retention_it == tenant_state.group_retention_deadlines.end()) { + const auto previous = retained_group_count_.fetch_add( + 1, std::memory_order_relaxed); + if (previous >= max_retained_groups_) { + retained_group_count_.fetch_sub(1, + std::memory_order_relaxed); + results[i] = false; + continue; + } + retention_it = tenant_state.group_retention_deadlines + .emplace(group_ids[i], deadline) + .first; + } else { + retention_it->second = std::max(retention_it->second, deadline); + } auto group_it = tenant_state.group_members.find(group_ids[i]); if (group_it == tenant_state.group_members.end() || @@ -2279,6 +2301,8 @@ void MasterService::PruneExpiredGroupRetentions() { it != tenant_state.group_retention_deadlines.end();) { if (it->second <= now) { it = tenant_state.group_retention_deadlines.erase(it); + retained_group_count_.fetch_sub(1, + std::memory_order_relaxed); } else { ++it; } @@ -7877,6 +7901,7 @@ void MasterService::MetadataSerializer::Reset() { for (auto& shard : service_->metadata_shards_) { shard.tenants.clear(); } + service_->retained_group_count_.store(0, std::memory_order_relaxed); { std::unique_lock lock( service_->group_routing_mutex_); diff --git a/mooncake-store/tests/master_service_test.cpp b/mooncake-store/tests/master_service_test.cpp index e0c9e32f79..d44d9c9f1b 100644 --- a/mooncake-store/tests/master_service_test.cpp +++ b/mooncake-store/tests/master_service_test.cpp @@ -1487,6 +1487,103 @@ TEST_F(MasterServiceTest, RetainGroupsProtectsFutureGroupMembers) { EXPECT_EQ(ErrorCode::INVALID_PARAMS, invalid[0].error()); } +TEST_F(MasterServiceTest, RetainGroupsEnforcesAdmissionBounds) { + auto service_config = MasterServiceConfig::builder() + .set_max_retained_groups(1) + .set_max_group_retention_ttl_ms(100) + .build(); + std::unique_ptr service_(new MasterService(service_config)); + + auto accepted = service_->RetainGroups({"retained_group_a"}, 80, "default"); + ASSERT_EQ(accepted.size(), 1); + ASSERT_TRUE(accepted[0].has_value()); + EXPECT_TRUE(accepted[0].value()); + + auto extended = + service_->RetainGroups({"retained_group_a"}, 100, "default"); + ASSERT_EQ(extended.size(), 1); + ASSERT_TRUE(extended[0].has_value()); + EXPECT_TRUE(extended[0].value()); + + auto rejected = service_->RetainGroups({"retained_group_b"}, 80, "default"); + ASSERT_EQ(rejected.size(), 1); + ASSERT_TRUE(rejected[0].has_value()); + EXPECT_FALSE(rejected[0].value()); + + auto too_long = + service_->RetainGroups({"retained_group_a"}, 101, "default"); + ASSERT_EQ(too_long.size(), 1); + ASSERT_FALSE(too_long[0].has_value()); + EXPECT_EQ(ErrorCode::INVALID_PARAMS, too_long[0].error()); +} + +TEST_F(MasterServiceTest, RetainGroupsAdmissionCapIsConcurrent) { + constexpr size_t kMaxRetainedGroups = 8; + constexpr size_t kRequests = 32; + auto service_config = MasterServiceConfig::builder() + .set_max_retained_groups(kMaxRetainedGroups) + .build(); + std::unique_ptr service_(new MasterService(service_config)); + std::atomic accepted{0}; + std::vector threads; + threads.reserve(kRequests); + + for (size_t i = 0; i < kRequests; ++i) { + threads.emplace_back([&, i] { + auto result = service_->RetainGroups( + {"concurrent_retained_group_" + std::to_string(i)}, 1000, + "default"); + if (result[0].has_value() && result[0].value()) { + accepted.fetch_add(1, std::memory_order_relaxed); + } + }); + } + for (auto& thread : threads) { + thread.join(); + } + + EXPECT_EQ(accepted.load(std::memory_order_relaxed), kMaxRetainedGroups); +} + +TEST_F(MasterServiceTest, RetainedGroupBecomesEvictableAfterTtl) { + auto service_config = + MasterServiceConfig::builder().set_default_kv_lease_ttl(20).build(); + constexpr size_t kSegmentSize = 4 * 1024 * 1024; + constexpr size_t kObjectSize = 2 * 1024 * 1024; + std::unique_ptr service_(new MasterService(service_config)); + [[maybe_unused]] const auto context = + PrepareSimpleSegment(*service_, "expired_retained_group_segment", + kDefaultSegmentBase, kSegmentSize); + const UUID client_id = generate_uuid(); + + const std::string key_a = "expired_retained_group_key_a"; + const std::string key_b = "expired_retained_group_key_b"; + const std::string group_id = FindGroupIdOnDifferentShard(key_a); + ReplicateConfig config; + config.replica_num = 1; + config.group_ids = std::vector{group_id}; + + auto retained = service_->RetainGroups({group_id}, 60, "default"); + ASSERT_EQ(retained.size(), 1); + ASSERT_TRUE(retained[0].has_value()); + ASSERT_TRUE(retained[0].value()); + PutCompletedObject(*service_, client_id, key_a, config, kObjectSize); + PutCompletedObject(*service_, client_id, key_b, config, kObjectSize); + + std::this_thread::sleep_for(std::chrono::milliseconds(150)); + ReplicateConfig trigger_config; + trigger_config.replica_num = 1; + auto trigger_result = + service_->PutStart(client_id, "trigger_expired_retained_group_eviction", + "default", kObjectSize, trigger_config); + ASSERT_FALSE(trigger_result.has_value()); + EXPECT_EQ(ErrorCode::NO_AVAILABLE_HANDLE, trigger_result.error()); + + std::this_thread::sleep_for(std::chrono::milliseconds(200)); + EXPECT_FALSE(service_->ExistKey(key_a, "default").value_or(true)); + EXPECT_FALSE(service_->ExistKey(key_b, "default").value_or(true)); +} + TEST_F(MasterServiceTest, GroupedEvictionSkipsUnsafeMembersAndEvictsSafePeers) { constexpr size_t kSegmentSize = 4 * 1024 * 1024; constexpr size_t kObjectSize = 2 * 1024 * 1024; From f9b03d4c7170205ba4b26d0c5abd35d0b947681b Mon Sep 17 00:00:00 2001 From: Ishan Dhanani Date: Fri, 10 Jul 2026 09:24:17 +0000 Subject: [PATCH 5/5] style(store): format retention POC --- mooncake-store/include/master_service.h | 3 +-- mooncake-store/src/real_client.cpp | 5 ++--- mooncake-store/src/rpc_service.cpp | 5 ++--- mooncake-store/tests/master_service_test.cpp | 15 +++++++-------- 4 files changed, 12 insertions(+), 16 deletions(-) diff --git a/mooncake-store/include/master_service.h b/mooncake-store/include/master_service.h index db3d78b623..04bb49c1f0 100644 --- a/mooncake-store/include/master_service.h +++ b/mooncake-store/include/master_service.h @@ -1222,8 +1222,7 @@ class MasterService { std::unordered_map> group_members; // group_id → set of keys - std::unordered_map + std::unordered_map group_retention_deadlines; bool Empty() const { diff --git a/mooncake-store/src/real_client.cpp b/mooncake-store/src/real_client.cpp index fadf1ac836..c803107ee2 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -2059,9 +2059,8 @@ std::vector RealClient::retainGroups( std::vector results; results.reserve(internal_results.size()); for (const auto &result : internal_results) { - results.push_back(result.has_value() - ? (result.value() ? 1 : 0) - : toInt(result.error())); + results.push_back(result.has_value() ? (result.value() ? 1 : 0) + : toInt(result.error())); } return results; } diff --git a/mooncake-store/src/rpc_service.cpp b/mooncake-store/src/rpc_service.cpp index f322e1293f..3272860faa 100644 --- a/mooncake-store/src/rpc_service.cpp +++ b/mooncake-store/src/rpc_service.cpp @@ -85,9 +85,8 @@ std::vector> WrappedMasterService::RetainGroups( const auto accepted = std::count_if( result.begin(), result.end(), [](const auto& item) { return item.has_value() && item.value(); }); - LOG(INFO) << "RetainGroups accepted " << accepted << "/" - << group_ids.size() << " groups for ttl_ms=" << ttl_ms - << ", tenant_id=" << tenant_id; + LOG(INFO) << "RetainGroups accepted " << accepted << "/" << group_ids.size() + << " groups for ttl_ms=" << ttl_ms << ", tenant_id=" << tenant_id; return result; } diff --git a/mooncake-store/tests/master_service_test.cpp b/mooncake-store/tests/master_service_test.cpp index d44d9c9f1b..d875b6a0e2 100644 --- a/mooncake-store/tests/master_service_test.cpp +++ b/mooncake-store/tests/master_service_test.cpp @@ -1447,9 +1447,8 @@ TEST_F(MasterServiceTest, RetainGroupsProtectsFutureGroupMembers) { constexpr size_t kSegmentSize = 4 * 1024 * 1024; constexpr size_t kObjectSize = 2 * 1024 * 1024; std::unique_ptr service_(new MasterService(service_config)); - [[maybe_unused]] const auto context = - PrepareSimpleSegment(*service_, "retained_group_segment", - kDefaultSegmentBase, kSegmentSize); + [[maybe_unused]] const auto context = PrepareSimpleSegment( + *service_, "retained_group_segment", kDefaultSegmentBase, kSegmentSize); const UUID client_id = generate_uuid(); const std::string key_a = "retained_group_key_a"; @@ -1459,8 +1458,8 @@ TEST_F(MasterServiceTest, RetainGroupsProtectsFutureGroupMembers) { config.replica_num = 1; config.group_ids = std::vector{group_id}; - auto retained = service_->RetainGroups({group_id, "missing_group"}, 1000, - "default"); + auto retained = + service_->RetainGroups({group_id, "missing_group"}, 1000, "default"); ASSERT_EQ(retained.size(), 2); ASSERT_TRUE(retained[0].has_value()); EXPECT_TRUE(retained[0].value()); @@ -1473,9 +1472,9 @@ TEST_F(MasterServiceTest, RetainGroupsProtectsFutureGroupMembers) { std::this_thread::sleep_for(std::chrono::milliseconds(200)); ReplicateConfig trigger_config; trigger_config.replica_num = 1; - auto trigger_result = service_->PutStart( - client_id, "trigger_retained_group_eviction", "default", kObjectSize, - trigger_config); + auto trigger_result = + service_->PutStart(client_id, "trigger_retained_group_eviction", + "default", kObjectSize, trigger_config); ASSERT_FALSE(trigger_result.has_value()); EXPECT_EQ(ErrorCode::NO_AVAILABLE_HANDLE, trigger_result.error()); EXPECT_TRUE(service_->GetReplicaList(key_a, "default").has_value());