From 3d302515b927b17050db91310a452327fa9652b0 Mon Sep 17 00:00:00 2001 From: Ray Yan Date: Thu, 23 Jul 2026 10:44:40 +0800 Subject: [PATCH 1/3] This is an automated cherry-pick of #10997 Signed-off-by: ti-chi-bot --- .../LocalAdmissionController.cpp | 64 +++- .../LocalAdmissionController.h | 5 + dbms/src/Flash/ResourceControl/TokenBucket.h | 2 + .../gtest_local_admission_controller.cpp | 279 ++++++++++++++++++ 4 files changed, 334 insertions(+), 16 deletions(-) create mode 100644 dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp diff --git a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp index 12520da17df..79a3a96c3b4 100644 --- a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp +++ b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp @@ -73,7 +73,7 @@ std::optional ResourceGroup::buildRequestInfoIfNecessary(const S { acquire_tokens = getAcquireRUNumWithoutLock( consumption_delta_info.speed, - LocalAdmissionController::DEFAULT_TARGET_PERIOD.count(), + REFILL_TOKEN_INTERVAL.count(), LocalAdmissionController::ACQUIRE_RU_AMPLIFICATION); assert(acquire_tokens >= 0.0); @@ -125,26 +125,50 @@ bool ResourceGroup::shouldReportRUConsumption(const SteadyClock::time_point & no return false; } -double ResourceGroup::getAcquireRUNumWithoutLock(double speed, uint32_t n_sec, double amplification) const +bool ResourceGroup::shouldRefillToken(const SteadyClock::time_point & now) const { - assert(amplification > 1.0); + std::lock_guard lock(mu); + if (burstable || bucket_mode != normal_mode || request_in_progress) + return false; + + const auto elapsed = now - last_request_gac_timepoint; + RUNTIME_CHECK(elapsed.count() >= 0, elapsed.count()); + if (elapsed < REFILL_TOKEN_INTERVAL) + return false; - double remaining_ru = 0.0; - remaining_ru = bucket->peek(); + const auto refill_threshold = getTokenHighWatermarkWithoutLock() * REFILL_TOKEN_THRESHOLD_RATE; + return bucket->peek() <= refill_threshold; +} - // Appropriate amplification is necessary to prevent situation that GAC has sufficient RU, - // but user query speed is limited due to LAC requests too few RU. - double acquire_num = speed * n_sec * amplification; +double ResourceGroup::getTokenHighWatermarkWithoutLock() const +{ + // The resource group definition contains the global burst limit. Only use capacity as a local high watermark after + // GAC has returned the capacity assigned to this client. Before that, keep the startup fill rate as the watermark. + const auto high_watermark = has_gac_capacity ? bucket->getCapacity() : static_cast(user_ru_per_sec); + if unlikely (high_watermark <= 0.0 && !burstable) + return DEFAULT_BUFFER_TOKENS; + return high_watermark; +} - // This should not happen, but still add this to avoid stuck. - if unlikely (acquire_num == 0.0 && remaining_ru == 0.0) - acquire_num = DEFAULT_BUFFER_TOKENS; +double ResourceGroup::getAcquireRUNumWithoutLock(double speed, uint32_t n_sec, double amplification) const +{ + assert(amplification > 1.0); - // The purpose of subtracting remaining_ru is try to ensure that the number of local tokens - // always stays same with the amount consumed. - acquire_num -= remaining_ru; - acquire_num = (acquire_num > 0.0 ? acquire_num : 0.0); - return acquire_num; + const auto remaining_ru = bucket->peek(); + const auto high_watermark = getTokenHighWatermarkWithoutLock(); + auto acquire_num = high_watermark - remaining_ru; + if (acquire_num <= 0.0) + return 0.0; + + if (bucket->lowToken()) + return acquire_num; + + // Refill at most one second of predicted consumption in normal mode. The fallback batch gradually raises an idle + // bucket without transferring the whole capacity from GAC in one request. Low-token mode bypasses this limit above. + const auto refill_window = high_watermark * (1.0 - REFILL_TOKEN_THRESHOLD_RATE); + const auto fallback_batch = std::min(static_cast(DEFAULT_BUFFER_TOKENS), refill_window); + const auto incremental_batch = std::max(speed * n_sec * amplification, fallback_batch); + return std::min(acquire_num, incremental_batch); } void ResourceGroup::updateNormalMode(double add_tokens, double new_capacity, const SteadyClock::time_point & now) @@ -160,6 +184,7 @@ void ResourceGroup::updateNormalMode(double add_tokens, double new_capacity, con burstable = true; return; } + has_gac_capacity = true; auto config = bucket->getConfig(); std::string ori_bucket_info = bucket->toString(); @@ -195,6 +220,7 @@ void ResourceGroup::updateTrickleMode( burstable = true; return; } + has_gac_capacity = true; bucket_mode = TokenBucketMode::trickle_mode; double new_fill_rate = add_tokens / (static_cast(trickle_ms) / 1000); @@ -431,9 +457,15 @@ std::optional LocalAdmissionController::b for (const auto & iter : local_resource_groups) { +<<<<<<< HEAD const auto rg_name = iter.first; const bool need_fetch_token = local_low_token_resource_groups.contains(rg_name); const bool need_report = iter.second->shouldReportRUConsumption(current_tick); +======= + const bool need_fetch_token = local_keyspace_low_token_resource_groups.contains(ele.first) + || ele.second->shouldRefillToken(current_tick); + const bool need_report = ele.second->shouldReportRUConsumption(current_tick); +>>>>>>> 555bb7cc2a (LAC: smoothing tokens request and keep in high level (#10997)) if (need_fetch_token || need_report) { diff --git a/dbms/src/Flash/ResourceControl/LocalAdmissionController.h b/dbms/src/Flash/ResourceControl/LocalAdmissionController.h index 965acb42764..9d2decacc58 100644 --- a/dbms/src/Flash/ResourceControl/LocalAdmissionController.h +++ b/dbms/src/Flash/ResourceControl/LocalAdmissionController.h @@ -133,6 +133,8 @@ class ResourceGroup final : private boost::noncopyable static constexpr auto REPORT_RU_CONSUMPTION_DELTA_THRESHOLD = 100; static constexpr auto EXTENDING_REPORT_RU_CONSUMPTION_FACTOR = 4; static constexpr auto DEFAULT_BUFFER_TOKENS = 5000; + static constexpr auto REFILL_TOKEN_INTERVAL = std::chrono::seconds(1); + static constexpr double REFILL_TOKEN_THRESHOLD_RATE = 0.8; // Indicate the round trip time of gac request. static constexpr auto GAC_RTT_ANTICIPATION = std::chrono::seconds(1); @@ -180,9 +182,11 @@ class ResourceGroup final : private boost::noncopyable endRequestWithoutLock(); } bool shouldReportRUConsumption(const SteadyClock::time_point & now) const; + bool shouldRefillToken(const SteadyClock::time_point & now) const; std::optional buildRequestInfoIfNecessary(const SteadyClock::time_point & now); LACRUConsumptionDeltaInfo updateRUConsumptionDeltaInfoWithoutLock(); double getAcquireRUNumWithoutLock(double speed, uint32_t n_sec, double amplification) const; + double getTokenHighWatermarkWithoutLock() const; void updateRUConsumptionSpeedIfNecessary(const SteadyClock::time_point & now); // Called when user change config of resource group. @@ -280,6 +284,7 @@ class ResourceGroup final : private boost::noncopyable // Local token bucket. TokenBucketPtr bucket; TokenBucketMode bucket_mode = TokenBucketMode::normal_mode; + bool has_gac_capacity = false; // For compute priority. uint64_t cpu_time_in_ns = 0; diff --git a/dbms/src/Flash/ResourceControl/TokenBucket.h b/dbms/src/Flash/ResourceControl/TokenBucket.h index 61c052e0a63..eaa27376b0a 100644 --- a/dbms/src/Flash/ResourceControl/TokenBucket.h +++ b/dbms/src/Flash/ResourceControl/TokenBucket.h @@ -94,6 +94,8 @@ class TokenBucket final bool isStatic() const { return fill_rate == 0.0; } + double getCapacity() const { return capacity; } + std::string toString() const { return fmt::format("tokens: {}, fill_rate: {}, capacity: {}", tokens, fill_rate, capacity); diff --git a/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp b/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp new file mode 100644 index 00000000000..12e913c0a7d --- /dev/null +++ b/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp @@ -0,0 +1,279 @@ +// Copyright 2026 PingCAP, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#include +#include +#include + +#include +#include + +namespace DB::tests +{ +namespace +{ +class TestPDClient final : public pingcap::pd::MockPDClient +{ +public: + explicit TestPDClient( + std::function + get_resource_group_) + : get_resource_group(std::move(get_resource_group_)) + {} + + resource_manager::GetResourceGroupResponse getResourceGroup( + const resource_manager::GetResourceGroupRequest & req) override + { + return get_resource_group(req); + } + + resource_manager::TokenBucketsResponse acquireTokenBuckets(const resource_manager::TokenBucketsRequest &) override + { + return {}; + } + +private: + std::function + get_resource_group; +}; + +resource_manager::GetResourceGroupResponse buildResourceGroupResp( + const resource_manager::GetResourceGroupRequest & req, + uint64_t fill_rate, + bool with_keyspace_id, + uint32_t priority = ResourceGroup::UserLowPriority) +{ + resource_manager::GetResourceGroupResponse resp; + auto * group = resp.mutable_group(); + group->set_name(req.resource_group_name()); + group->set_mode(resource_manager::GroupMode::RUMode); + group->set_priority(priority); + if (with_keyspace_id) + group->mutable_keyspace_id()->set_value(req.keyspace_id().value()); + auto * settings = group->mutable_r_u_settings()->mutable_r_u()->mutable_settings(); + settings->set_fill_rate(fill_rate); + settings->set_burst_limit(fill_rate); + return resp; +} +} // namespace + +TEST(LocalAdmissionControllerTest, LegacyBackendCachesSharedGroupAndOmitsKeyspaceAfterDetection) +{ + constexpr KeyspaceID keyspace_id = 23; + std::vector request_has_keyspace; + pingcap::kv::Cluster cluster; + cluster.pd_client = std::make_shared( + [&request_has_keyspace](const resource_manager::GetResourceGroupRequest & req) { + request_has_keyspace.push_back(req.has_keyspace_id()); + const auto fill_rate = req.resource_group_name() == "default" ? 1024 : 2048; + return buildResourceGroupResp(req, fill_rate, false); + }); + + LocalAdmissionController lac(&cluster, nullptr, true, false); + + ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default")); + ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id + 1, "analytics")); + + ASSERT_EQ(request_has_keyspace.size(), 2); + EXPECT_TRUE(request_has_keyspace[0]); + EXPECT_FALSE(request_has_keyspace[1]); + + EXPECT_EQ(lac.getCachedResourceGroupForTest(keyspace_id, "default"), nullptr); + auto default_group = lac.getCachedResourceGroupForTest(NullspaceID, "default"); + ASSERT_NE(default_group, nullptr); + EXPECT_TRUE(lac.getPriority(keyspace_id, "default").has_value()); + + default_group->consumeResource(1, 0); + auto request_info = default_group->buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_EQ(request_info->keyspace_id, NullspaceID); + + auto analytics_group = lac.getCachedResourceGroupForTest(NullspaceID, "analytics"); + ASSERT_NE(analytics_group, nullptr); + EXPECT_EQ(analytics_group->user_ru_per_sec, 2048); +} + +TEST(LocalAdmissionControllerTest, KeyspaceScopedBackendCachesExactGroupAndKeepsKeyspaceRequests) +{ + constexpr KeyspaceID keyspace_id = 29; + std::vector request_has_keyspace; + pingcap::kv::Cluster cluster; + cluster.pd_client = std::make_shared( + [&request_has_keyspace](const resource_manager::GetResourceGroupRequest & req) { + request_has_keyspace.push_back(req.has_keyspace_id()); + return buildResourceGroupResp(req, 4096, true); + }); + + LocalAdmissionController lac(&cluster, nullptr, true, false); + + ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default")); + ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id + 1, "analytics")); + + ASSERT_EQ(request_has_keyspace.size(), 2); + EXPECT_TRUE(request_has_keyspace[0]); + EXPECT_TRUE(request_has_keyspace[1]); + + EXPECT_EQ(lac.getCachedResourceGroupForTest(NullspaceID, "default"), nullptr); + auto default_group = lac.getCachedResourceGroupForTest(keyspace_id, "default"); + ASSERT_NE(default_group, nullptr); + default_group->consumeResource(1, 0); + auto request_info = default_group->buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_EQ(request_info->keyspace_id, keyspace_id); + + auto analytics_group = lac.getCachedResourceGroupForTest(keyspace_id + 1, "analytics"); + ASSERT_NE(analytics_group, nullptr); + EXPECT_EQ(analytics_group->group_pb.keyspace_id().value(), keyspace_id + 1); +} + +TEST(LocalAdmissionControllerTest, WarmupErrorPropagatesWithoutCompatibilityFallback) +{ + constexpr KeyspaceID keyspace_id = 31; + size_t request_count = 0; + pingcap::kv::Cluster cluster; + cluster.pd_client + = std::make_shared([&request_count](const resource_manager::GetResourceGroupRequest &) { + ++request_count; + resource_manager::GetResourceGroupResponse resp; + resp.mutable_error()->set_message("resource group default not found"); + return resp; + }); + + LocalAdmissionController lac(&cluster, nullptr, true, false); + + EXPECT_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default"), DB::Exception); + EXPECT_EQ(request_count, 1); + EXPECT_EQ(lac.cachedResourceGroupCount(), 0); +} + +TEST(LocalAdmissionControllerTest, UnexpectedPriorityFallsBackToMedium) +{ + constexpr KeyspaceID keyspace_id = 37; + pingcap::kv::Cluster cluster; + cluster.pd_client = std::make_shared([](const resource_manager::GetResourceGroupRequest & req) { + const auto priority = req.resource_group_name() == "default" ? 0U : 9U; + return buildResourceGroupResp(req, 2048, false, priority); + }); + + LocalAdmissionController lac(&cluster, nullptr, true, false); + + ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default")); + ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "analytics")); + + auto default_group = lac.getCachedResourceGroupForTest(NullspaceID, "default"); + ASSERT_NE(default_group, nullptr); + EXPECT_EQ(default_group->user_priority_val, ResourceGroup::MediumPriorityValue); + + auto analytics_group = lac.getCachedResourceGroupForTest(NullspaceID, "analytics"); + ASSERT_NE(analytics_group, nullptr); + EXPECT_EQ(analytics_group->user_priority_val, ResourceGroup::MediumPriorityValue); +} + +TEST(LocalAdmissionControllerTest, StartupRefillDoesNotUseGlobalBurstLimitAsHighWatermark) +{ + constexpr double fill_rate = 1000; + constexpr int64_t global_burst_limit = 10000; + constexpr double consumed_tokens = 500; + const auto start_time = SteadyClock::now() - 2 * ResourceGroup::REFILL_TOKEN_INTERVAL; + + resource_manager::ResourceGroup group_pb; + group_pb.set_name("startup"); + group_pb.set_mode(resource_manager::GroupMode::RUMode); + group_pb.set_priority(ResourceGroup::UserMediumPriority); + auto * settings = group_pb.mutable_r_u_settings()->mutable_r_u()->mutable_settings(); + settings->set_fill_rate(fill_rate); + settings->set_burst_limit(global_burst_limit); + + ResourceGroup group(NullspaceID, group_pb, start_time); + group.smooth_ru_consumption_speed = 0; + group.consumeResource(consumed_tokens, 0); + + EXPECT_FALSE(group.shouldRefillToken(start_time + ResourceGroup::REFILL_TOKEN_INTERVAL / 2)); + EXPECT_TRUE(group.shouldRefillToken(start_time + ResourceGroup::REFILL_TOKEN_INTERVAL)); + + const auto request_info = group.buildRequestInfoIfNecessary(start_time + ResourceGroup::REFILL_TOKEN_INTERVAL); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, fill_rate * (1 - ResourceGroup::REFILL_TOKEN_THRESHOLD_RATE)); + EXPECT_DOUBLE_EQ(request_info->ru_consumption_delta, consumed_tokens); + + const auto response_time = start_time + ResourceGroup::REFILL_TOKEN_INTERVAL; + group.updateNormalMode(request_info->acquire_tokens, global_burst_limit, response_time); + EXPECT_FALSE(group.lowToken()); + EXPECT_TRUE(group.shouldRefillToken(response_time + ResourceGroup::REFILL_TOKEN_INTERVAL)); + + const auto next_request_info + = group.buildRequestInfoIfNecessary(response_time + ResourceGroup::REFILL_TOKEN_INTERVAL); + ASSERT_TRUE(next_request_info.has_value()); + EXPECT_DOUBLE_EQ( + next_request_info->acquire_tokens, + global_burst_limit * (1 - ResourceGroup::REFILL_TOKEN_THRESHOLD_RATE)); +} + +TEST(LocalAdmissionControllerTest, RefillTokensIncrementallyAboveLowWatermark) +{ + constexpr double capacity = 10000; + constexpr double consumed_tokens = 3000; + + ResourceGroup group( + "normal_refill", + ResourceGroup::UserMediumPriority, + capacity, + /*burstable_=*/false); + group.smooth_ru_consumption_speed = 0; + group.consumeResource(consumed_tokens, 0); + + ASSERT_TRUE(group.shouldRefillToken(SteadyClock::now())); + auto request_info = group.buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, capacity * (1 - ResourceGroup::REFILL_TOKEN_THRESHOLD_RATE)); + EXPECT_DOUBLE_EQ(request_info->ru_consumption_delta, consumed_tokens); +} + +TEST(LocalAdmissionControllerTest, RefillTokensUsesPredictedConsumptionWhenHigher) +{ + constexpr double capacity = 10000; + constexpr double consumed_tokens = 3000; + + ResourceGroup group( + "high_speed", + ResourceGroup::UserMediumPriority, + capacity, + /*burstable_=*/false); + group.smooth_ru_consumption_speed = 4000; + group.consumeResource(consumed_tokens, 0); + + auto request_info = group.buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, consumed_tokens); +} + +TEST(LocalAdmissionControllerTest, LowTokenRefillBypassesIncrementalLimit) +{ + constexpr double capacity = 10000; + constexpr double consumed_tokens = 7500; + + ResourceGroup group( + "emergency_refill", + ResourceGroup::UserMediumPriority, + capacity, + /*burstable_=*/false); + group.smooth_ru_consumption_speed = 0; + group.consumeResource(consumed_tokens, 0); + + auto request_info = group.buildRequestInfoIfNecessary(SteadyClock::now()); + ASSERT_TRUE(request_info.has_value()); + EXPECT_DOUBLE_EQ(request_info->acquire_tokens, consumed_tokens); +} + +} // namespace DB::tests From 5c92b727261067ead7668faddb4b96e7e74ca34d Mon Sep 17 00:00:00 2001 From: JaySon-Huang Date: Wed, 5 Aug 2026 13:46:10 +0800 Subject: [PATCH 2/3] Resolve cherry-pick conflicts for LAC token refill on release-8.5. Keep the release-8.5 LAC API while retaining shouldRefillToken, and adapt the new unit tests away from master keyspace-only helpers. --- .../LocalAdmissionController.cpp | 9 +- .../gtest_local_admission_controller.cpp | 165 +----------------- 2 files changed, 3 insertions(+), 171 deletions(-) diff --git a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp index 79a3a96c3b4..78c63773b79 100644 --- a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp +++ b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp @@ -457,15 +457,10 @@ std::optional LocalAdmissionController::b for (const auto & iter : local_resource_groups) { -<<<<<<< HEAD const auto rg_name = iter.first; - const bool need_fetch_token = local_low_token_resource_groups.contains(rg_name); + const bool need_fetch_token = local_low_token_resource_groups.contains(rg_name) + || iter.second->shouldRefillToken(current_tick); const bool need_report = iter.second->shouldReportRUConsumption(current_tick); -======= - const bool need_fetch_token = local_keyspace_low_token_resource_groups.contains(ele.first) - || ele.second->shouldRefillToken(current_tick); - const bool need_report = ele.second->shouldReportRUConsumption(current_tick); ->>>>>>> 555bb7cc2a (LAC: smoothing tokens request and keep in high level (#10997)) if (need_fetch_token || need_report) { diff --git a/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp b/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp index 12e913c0a7d..cebe1387a0b 100644 --- a/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp +++ b/dbms/src/Flash/ResourceControl/tests/gtest_local_admission_controller.cpp @@ -14,172 +14,9 @@ #include #include -#include - -#include -#include namespace DB::tests { -namespace -{ -class TestPDClient final : public pingcap::pd::MockPDClient -{ -public: - explicit TestPDClient( - std::function - get_resource_group_) - : get_resource_group(std::move(get_resource_group_)) - {} - - resource_manager::GetResourceGroupResponse getResourceGroup( - const resource_manager::GetResourceGroupRequest & req) override - { - return get_resource_group(req); - } - - resource_manager::TokenBucketsResponse acquireTokenBuckets(const resource_manager::TokenBucketsRequest &) override - { - return {}; - } - -private: - std::function - get_resource_group; -}; - -resource_manager::GetResourceGroupResponse buildResourceGroupResp( - const resource_manager::GetResourceGroupRequest & req, - uint64_t fill_rate, - bool with_keyspace_id, - uint32_t priority = ResourceGroup::UserLowPriority) -{ - resource_manager::GetResourceGroupResponse resp; - auto * group = resp.mutable_group(); - group->set_name(req.resource_group_name()); - group->set_mode(resource_manager::GroupMode::RUMode); - group->set_priority(priority); - if (with_keyspace_id) - group->mutable_keyspace_id()->set_value(req.keyspace_id().value()); - auto * settings = group->mutable_r_u_settings()->mutable_r_u()->mutable_settings(); - settings->set_fill_rate(fill_rate); - settings->set_burst_limit(fill_rate); - return resp; -} -} // namespace - -TEST(LocalAdmissionControllerTest, LegacyBackendCachesSharedGroupAndOmitsKeyspaceAfterDetection) -{ - constexpr KeyspaceID keyspace_id = 23; - std::vector request_has_keyspace; - pingcap::kv::Cluster cluster; - cluster.pd_client = std::make_shared( - [&request_has_keyspace](const resource_manager::GetResourceGroupRequest & req) { - request_has_keyspace.push_back(req.has_keyspace_id()); - const auto fill_rate = req.resource_group_name() == "default" ? 1024 : 2048; - return buildResourceGroupResp(req, fill_rate, false); - }); - - LocalAdmissionController lac(&cluster, nullptr, true, false); - - ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default")); - ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id + 1, "analytics")); - - ASSERT_EQ(request_has_keyspace.size(), 2); - EXPECT_TRUE(request_has_keyspace[0]); - EXPECT_FALSE(request_has_keyspace[1]); - - EXPECT_EQ(lac.getCachedResourceGroupForTest(keyspace_id, "default"), nullptr); - auto default_group = lac.getCachedResourceGroupForTest(NullspaceID, "default"); - ASSERT_NE(default_group, nullptr); - EXPECT_TRUE(lac.getPriority(keyspace_id, "default").has_value()); - - default_group->consumeResource(1, 0); - auto request_info = default_group->buildRequestInfoIfNecessary(SteadyClock::now()); - ASSERT_TRUE(request_info.has_value()); - EXPECT_EQ(request_info->keyspace_id, NullspaceID); - - auto analytics_group = lac.getCachedResourceGroupForTest(NullspaceID, "analytics"); - ASSERT_NE(analytics_group, nullptr); - EXPECT_EQ(analytics_group->user_ru_per_sec, 2048); -} - -TEST(LocalAdmissionControllerTest, KeyspaceScopedBackendCachesExactGroupAndKeepsKeyspaceRequests) -{ - constexpr KeyspaceID keyspace_id = 29; - std::vector request_has_keyspace; - pingcap::kv::Cluster cluster; - cluster.pd_client = std::make_shared( - [&request_has_keyspace](const resource_manager::GetResourceGroupRequest & req) { - request_has_keyspace.push_back(req.has_keyspace_id()); - return buildResourceGroupResp(req, 4096, true); - }); - - LocalAdmissionController lac(&cluster, nullptr, true, false); - - ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default")); - ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id + 1, "analytics")); - - ASSERT_EQ(request_has_keyspace.size(), 2); - EXPECT_TRUE(request_has_keyspace[0]); - EXPECT_TRUE(request_has_keyspace[1]); - - EXPECT_EQ(lac.getCachedResourceGroupForTest(NullspaceID, "default"), nullptr); - auto default_group = lac.getCachedResourceGroupForTest(keyspace_id, "default"); - ASSERT_NE(default_group, nullptr); - default_group->consumeResource(1, 0); - auto request_info = default_group->buildRequestInfoIfNecessary(SteadyClock::now()); - ASSERT_TRUE(request_info.has_value()); - EXPECT_EQ(request_info->keyspace_id, keyspace_id); - - auto analytics_group = lac.getCachedResourceGroupForTest(keyspace_id + 1, "analytics"); - ASSERT_NE(analytics_group, nullptr); - EXPECT_EQ(analytics_group->group_pb.keyspace_id().value(), keyspace_id + 1); -} - -TEST(LocalAdmissionControllerTest, WarmupErrorPropagatesWithoutCompatibilityFallback) -{ - constexpr KeyspaceID keyspace_id = 31; - size_t request_count = 0; - pingcap::kv::Cluster cluster; - cluster.pd_client - = std::make_shared([&request_count](const resource_manager::GetResourceGroupRequest &) { - ++request_count; - resource_manager::GetResourceGroupResponse resp; - resp.mutable_error()->set_message("resource group default not found"); - return resp; - }); - - LocalAdmissionController lac(&cluster, nullptr, true, false); - - EXPECT_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default"), DB::Exception); - EXPECT_EQ(request_count, 1); - EXPECT_EQ(lac.cachedResourceGroupCount(), 0); -} - -TEST(LocalAdmissionControllerTest, UnexpectedPriorityFallsBackToMedium) -{ - constexpr KeyspaceID keyspace_id = 37; - pingcap::kv::Cluster cluster; - cluster.pd_client = std::make_shared([](const resource_manager::GetResourceGroupRequest & req) { - const auto priority = req.resource_group_name() == "default" ? 0U : 9U; - return buildResourceGroupResp(req, 2048, false, priority); - }); - - LocalAdmissionController lac(&cluster, nullptr, true, false); - - ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "default")); - ASSERT_NO_THROW(lac.warmupResourceGroupInfoCache(keyspace_id, "analytics")); - - auto default_group = lac.getCachedResourceGroupForTest(NullspaceID, "default"); - ASSERT_NE(default_group, nullptr); - EXPECT_EQ(default_group->user_priority_val, ResourceGroup::MediumPriorityValue); - - auto analytics_group = lac.getCachedResourceGroupForTest(NullspaceID, "analytics"); - ASSERT_NE(analytics_group, nullptr); - EXPECT_EQ(analytics_group->user_priority_val, ResourceGroup::MediumPriorityValue); -} - TEST(LocalAdmissionControllerTest, StartupRefillDoesNotUseGlobalBurstLimitAsHighWatermark) { constexpr double fill_rate = 1000; @@ -195,7 +32,7 @@ TEST(LocalAdmissionControllerTest, StartupRefillDoesNotUseGlobalBurstLimitAsHigh settings->set_fill_rate(fill_rate); settings->set_burst_limit(global_burst_limit); - ResourceGroup group(NullspaceID, group_pb, start_time); + ResourceGroup group(group_pb, start_time); group.smooth_ru_consumption_speed = 0; group.consumeResource(consumed_tokens, 0); From cc6532cd796207f035cf6cedb26cf83a29925265 Mon Sep 17 00:00:00 2001 From: JaySon-Huang Date: Wed, 5 Aug 2026 14:07:17 +0800 Subject: [PATCH 3/3] Format codes Signed-off-by: JaySon-Huang --- dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp index 78c63773b79..07618045e88 100644 --- a/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp +++ b/dbms/src/Flash/ResourceControl/LocalAdmissionController.cpp @@ -458,8 +458,8 @@ std::optional LocalAdmissionController::b for (const auto & iter : local_resource_groups) { const auto rg_name = iter.first; - const bool need_fetch_token = local_low_token_resource_groups.contains(rg_name) - || iter.second->shouldRefillToken(current_tick); + const bool need_fetch_token + = local_low_token_resource_groups.contains(rg_name) || iter.second->shouldRefillToken(current_tick); const bool need_report = iter.second->shouldReportRUConsumption(current_tick); if (need_fetch_token || need_report)