diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index 15ca9bac56..7ed6e9d3bb 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 current and future objects in each group for the requested " + "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/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_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 42fd27efa3..04bb49c1f0 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 @@ -1218,11 +1222,14 @@ 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(); } }; @@ -1397,6 +1404,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, @@ -1538,6 +1546,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/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/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/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..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), @@ -1606,11 +1608,33 @@ 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); + retained_group_count_.fetch_sub(1, std::memory_order_relaxed); + 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 +1975,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"; } @@ -2176,6 +2203,119 @@ 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 || 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); + 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_); + 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); + retained_group_count_.fetch_sub(1, std::memory_order_relaxed); + } else { + ++it; + } + } + 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 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() || + group_it->second.empty()) { + results[i] = true; + continue; + } + 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.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); + retained_group_count_.fetch_sub(1, + std::memory_order_relaxed); + } 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; @@ -7761,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/src/real_client.cpp b/mooncake-store/src/real_client.cpp index fe00d17224..c803107ee2 100644 --- a/mooncake-store/src/real_client.cpp +++ b/mooncake-store/src/real_client.cpp @@ -2053,6 +2053,18 @@ 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 +4710,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..3272860faa 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 @@ -76,6 +78,18 @@ std::vector> WrappedMasterService::BatchExistKey( return result; } +std::vector> WrappedMasterService::RetainGroups( + const std::vector& group_ids, uint64_t ttl_ms, + const std::string& 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< std::unordered_map, boost::hash>, ErrorCode> @@ -1349,6 +1363,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..d875b6a0e2 100644 --- a/mooncake-store/tests/master_service_test.cpp +++ b/mooncake-store/tests/master_service_test.cpp @@ -1441,6 +1441,148 @@ TEST_F(MasterServiceTest, } } +TEST_F(MasterServiceTest, RetainGroupsProtectsFutureGroupMembers) { + 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}; + + 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_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; + 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, 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;