blob: 05fbb2ce35ec034f52be8f53a6a9bf47a84edde1 [file]
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you 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 <bvar/bvar.h>
#include <cpp/token_bucket_rate_limiter.h>
#include <gtest/gtest.h>
#include <atomic>
#include <chrono>
#include <thread>
#include <vector>
#include "common/configbase.h"
#include "common/logging.h"
using namespace doris::cloud;
namespace doris {
extern bvar::Adder<int64_t> s3_put_rate_limit_rejected_count;
} // namespace doris
int main(int argc, char** argv) {
auto conf_file = "doris_cloud.conf";
if (!doris::cloud::config::init(conf_file, true)) {
std::cerr << "failed to init config file, conf=" << conf_file << std::endl;
return -1;
}
if (!doris::cloud::init_glog("util")) {
std::cerr << "failed to init glog" << std::endl;
return -1;
}
::testing::InitGoogleTest(&argc, argv);
return RUN_ALL_TESTS();
}
// #define S3RateLimiterTest DISABLED_S3RateLimiterTest
TEST(S3RateLimiterTest, normal) {
auto rate_limiter = doris::S3RateLimiter(1, 5, 10);
std::atomic_int64_t failed;
std::atomic_int64_t succ;
std::atomic_int64_t sleep_thread_num;
std::atomic_int64_t sleep;
auto request_thread = [&]() {
std::this_thread::sleep_for(std::chrono::milliseconds(100));
auto ms = rate_limiter.add(1);
if (ms < 0) {
failed++;
} else if (ms == 0) {
succ++;
} else {
sleep += ms;
sleep_thread_num++;
}
};
{
std::vector<std::thread> threads;
for (size_t i = 0; i < 20; i++) {
threads.emplace_back(request_thread);
}
for (auto& t : threads) {
if (t.joinable()) {
t.join();
}
}
}
EXPECT_EQ(failed, 10);
EXPECT_EQ(succ, 5);
EXPECT_EQ(sleep_thread_num, 5);
}
TEST(S3RateLimiterTest, BasicRateLimit) {
auto rate_limiter = doris::S3RateLimiter(100, 100, 200);
int64_t sleep_time = rate_limiter.add(50);
EXPECT_EQ(sleep_time, 0);
sleep_time = rate_limiter.add(60);
EXPECT_GT(sleep_time, 0);
}
TEST(S3RateLimiterTest, BurstCapacity) {
auto rate_limiter = doris::S3RateLimiter(125, 250, 1000);
int64_t sleep_time = rate_limiter.add(250);
EXPECT_EQ(sleep_time, 0);
sleep_time = rate_limiter.add(1);
EXPECT_GT(sleep_time, 0);
}
TEST(S3RateLimiterTest, ExceedLimit) {
auto rate_limiter = doris::S3RateLimiter(125, 250, 250);
int64_t sleep_time = rate_limiter.add(250);
EXPECT_EQ(sleep_time, 0);
sleep_time = rate_limiter.add(1);
EXPECT_EQ(sleep_time, -1);
}
TEST(S3RateLimiterHolderTest, BvarMetric) {
bvar::Adder<int64_t> rate_limit_sleep_ns("rate_limit_sleep_ns");
bvar::Adder<int64_t> rate_limit_sleep_count("rate_limit_sleep_count");
bvar::Adder<int64_t> rate_limit_rejected_count("rate_limit_rejected_count");
auto rate_limiter_holder = doris::S3RateLimiterHolder(
125, 250, 251,
doris::metric_func_factory(rate_limit_sleep_ns, rate_limit_sleep_count,
&rate_limit_rejected_count));
int64_t sleep_time = rate_limiter_holder.add(250);
EXPECT_EQ(sleep_time, 0);
EXPECT_EQ(rate_limit_sleep_ns.get_value(), 0);
EXPECT_EQ(rate_limit_sleep_count.get_value(), 0);
EXPECT_EQ(rate_limit_rejected_count.get_value(), 0);
sleep_time = rate_limiter_holder.add(1);
EXPECT_GT(sleep_time, 0);
EXPECT_GT(rate_limit_sleep_ns.get_value(), 0);
EXPECT_EQ(rate_limit_sleep_count.get_value(), 1);
EXPECT_EQ(rate_limit_rejected_count.get_value(), 0);
sleep_time = rate_limiter_holder.add(1);
EXPECT_EQ(sleep_time, -1);
EXPECT_EQ(rate_limit_rejected_count.get_value(), 1);
}
TEST(S3RateLimiterHolderTest, ApplyS3RateLimitRecordsRejectedMetric) {
auto rate_limiter_holder = doris::S3RateLimiterHolder(
125, 250, 1, doris::s3_rate_limiter_metric_func(doris::S3RateLimitType::PUT));
auto rejected_count = doris::s3_put_rate_limit_rejected_count.get_value();
int64_t sleep_time =
doris::apply_s3_rate_limit(doris::S3RateLimitType::PUT, &rate_limiter_holder, 1);
EXPECT_EQ(sleep_time, 0);
EXPECT_EQ(doris::s3_put_rate_limit_rejected_count.get_value(), rejected_count);
sleep_time = doris::apply_s3_rate_limit(doris::S3RateLimitType::PUT, &rate_limiter_holder, 1);
EXPECT_EQ(sleep_time, -1);
EXPECT_EQ(doris::s3_put_rate_limit_rejected_count.get_value(), rejected_count + 1);
}
TEST(S3RateLimiterHolderTest, ConcurrentResetReturnsConsistentConfig) {
constexpr size_t config_a_speed = 101;
constexpr size_t config_a_burst = 102;
constexpr size_t config_a_limit = 103;
constexpr size_t config_b_speed = 201;
constexpr size_t config_b_burst = 202;
constexpr size_t config_b_limit = 203;
doris::S3RateLimiterHolder rate_limiter_holder(config_a_speed, config_a_burst, config_a_limit,
[](int64_t) {});
std::atomic<bool> start {false};
std::thread reset_thread([&]() {
while (!start.load(std::memory_order_acquire)) {
std::this_thread::yield();
}
for (size_t i = 0; i < 10000; ++i) {
if (i % 2 == 0) {
rate_limiter_holder.reset(config_b_speed, config_b_burst, config_b_limit);
} else {
rate_limiter_holder.reset(config_a_speed, config_a_burst, config_a_limit);
}
}
});
start.store(true, std::memory_order_release);
for (size_t i = 0; i < 10000; ++i) {
auto result = rate_limiter_holder.add_with_config(0);
bool is_config_a = result.max_speed == config_a_speed &&
result.max_burst == config_a_burst && result.limit == config_a_limit;
bool is_config_b = result.max_speed == config_b_speed &&
result.max_burst == config_b_burst && result.limit == config_b_limit;
EXPECT_TRUE(is_config_a || is_config_b);
}
reset_thread.join();
}