diff --git a/src/core/partition_balancing/balancing_impl.cpp b/src/core/partition_balancing/balancing_impl.cpp index 6982f59..207a3fa 100644 --- a/src/core/partition_balancing/balancing_impl.cpp +++ b/src/core/partition_balancing/balancing_impl.cpp @@ -157,6 +157,7 @@ void RebalancePartitions( ExecuteRebalancingPhase(sortedHubs, hubPartitions, migrationContext, loadFactorPredictor, state, settings); } + assignedPartitions.clear(); for (auto it = hubPartitions.begin(); it != hubPartitions.end(); ++it) { auto& weightedPartitions = it->second; const auto& hub = it->first; diff --git a/src/core/partition_balancing/tests/rebalance_partitions_ut.cpp b/src/core/partition_balancing/tests/rebalance_partitions_ut.cpp new file mode 100644 index 0000000..5d334d8 --- /dev/null +++ b/src/core/partition_balancing/tests/rebalance_partitions_ut.cpp @@ -0,0 +1,219 @@ +#include + +#include "test_helpers_ut.hpp" +#include "load_factor_predictor_mock.hpp" + +#include + +namespace NCoordinator::NCore { + +//////////////////////////////////////////////////////////////////////////////// + +using namespace testing; +using namespace NDetail; +using namespace NDomain; + +class RebalancePartitionsTest + : public TBalancingTestBase +{ +}; + +TEST_F(RebalancePartitionsTest, DoesNothingWhenCVBelowThreshold) { + TCoordinationState::TClusterSnapshot snapshot; + snapshot.emplace(HUB("hub-1"), THubReport{EP(100), HUB("hub-1"), DC("myt"), LF(50), {}}); + snapshot.emplace(HUB("hub-2"), THubReport{EP(100), HUB("hub-2"), DC("myt"), LF(52), {}}); + + TPartitionMap map{ + .Partitions{}, + .Epoch{EP(100)}, + }; + + TCoordinationContext context; + TStateBuildingSettings stateSettings; + TCoordinationState state(map, snapshot, context, stateSettings); + + TSortedHubs sortedHubs{ + {LF(50), HUB("hub-1")}, + {LF(52), HUB("hub-2")}, + }; + + TAssignedPartitions assignedPartitions; + + TMigrationContext migrationContext; + TBalancingSettings settings; + settings.BalancingThresholdCV = 1.0; + + EXPECT_CALL(*Predictor_, PredictLoadFactor(_, _, _)).Times(0); + + RebalancePartitions( + sortedHubs, + assignedPartitions, + migrationContext, + Predictor_, + state, + settings); + + EXPECT_TRUE(assignedPartitions.empty()); + EXPECT_TRUE(migrationContext.MigratingPartitions.empty()); +} + +TEST_F(RebalancePartitionsTest, PerformsSingleRebalancingPhase) { + TCoordinationState::TClusterSnapshot snapshot; + snapshot.emplace(HUB("hub-max"), THubReport{ + EP(100), HUB("hub-max"), DC("myt"), LF(90), {{PID(1), PW(20)}} + }); + snapshot.emplace(HUB("hub-min"), THubReport{ + EP(100), HUB("hub-min"), DC("myt"), LF(10), {} + }); + + TPartitionMap map{ + .Partitions{{PID(1), HUB("hub-max")}}, + .Epoch{EP(100)}, + }; + + TCoordinationContext context{ + .PartitionWeights{{PID(1), PW(20)}}, + }; + + TStateBuildingSettings stateSettings; + TCoordinationState state(map, snapshot, context, stateSettings); + + TSortedHubs sortedHubs{ + {LF(90), HUB("hub-max")}, + {LF(10), HUB("hub-min")}, + }; + + TAssignedPartitions assignedPartitions{ + {PID(1), HUB("hub-max")}, + }; + + TMigrationContext migrationContext; + TBalancingSettings settings; + settings.BalancingThresholdCV = 0.1; + settings.BalancingTargetCV = 0.0; + settings.MinLoadFactorDelta = LF(5); + settings.MigratingWeightLimit = PW(1000); + settings.MaxRebalancePhases = 1; + + EXPECT_CALL(*Predictor_, PredictLoadFactor(LF(90), PW(20), _)) + .WillRepeatedly(Return(LF(70))); + + EXPECT_CALL(*Predictor_, PredictLoadFactor(LF(10), PW(20), _)) + .WillRepeatedly(Return(LF(30))); + + RebalancePartitions( + sortedHubs, + assignedPartitions, + migrationContext, + Predictor_, + state, + settings); + + + EXPECT_EQ(assignedPartitions.size(), 1u); + EXPECT_EQ(assignedPartitions.front().first, PID(1)); + EXPECT_EQ(assignedPartitions.front().second, HUB("hub-min")); + + EXPECT_EQ(migrationContext.TotalMigratingWeight, PW(20)); +} + +TEST_F(RebalancePartitionsTest, StopsWhenMigrationBudgetExceeded) { + TCoordinationState::TClusterSnapshot snapshot; + snapshot.emplace(HUB("hub-max"), THubReport{ + EP(100), HUB("hub-max"), DC("myt"), LF(90), {{PID(1), PW(50)}} + }); + snapshot.emplace(HUB("hub-min"), THubReport{ + EP(100), HUB("hub-min"), DC("myt"), LF(10), {} + }); + + TPartitionMap map{ + .Partitions{{PID(1), HUB("hub-max")}}, + .Epoch{EP(100)}, + }; + + TCoordinationContext context{ + .PartitionWeights{{PID(1), PW(50)}}, + }; + + TStateBuildingSettings stateSettings; + TCoordinationState state(map, snapshot, context, stateSettings); + + TSortedHubs sortedHubs{ + {LF(90), HUB("hub-max")}, + {LF(10), HUB("hub-min")}, + }; + + TAssignedPartitions assignedPartitions{ + {PID(1), HUB("hub-max")}, + }; + + TMigrationContext migrationContext; + migrationContext.TotalMigratingWeight = PW(45); + + TBalancingSettings settings; + settings.BalancingThresholdCV = 0.0; + settings.MigrationBudgetThreshold = PW(10); + settings.MigratingWeightLimit = PW(50); + settings.MaxRebalancePhases = 10; + + EXPECT_CALL(*Predictor_, PredictLoadFactor(_, _, _)).Times(0); + + RebalancePartitions( + sortedHubs, + assignedPartitions, + migrationContext, + Predictor_, + state, + settings); + + EXPECT_EQ(assignedPartitions.front().second, HUB("hub-max")); +} + +TEST_F(RebalancePartitionsTest, RespectsMaxRebalancePhases) { + TCoordinationState::TClusterSnapshot snapshot; + snapshot.emplace(HUB("hub-1"), THubReport{EP(100), HUB("hub-1"), DC("myt"), LF(80), {{PID(1), PW(10)}}}); + snapshot.emplace(HUB("hub-2"), THubReport{EP(100), HUB("hub-2"), DC("myt"), LF(20), {}}); + + TPartitionMap map{ + .Partitions{{PID(1), HUB("hub-1")}}, + .Epoch{EP(100)}, + }; + + TCoordinationContext context{ + .PartitionWeights{{PID(1), PW(10)}}, + }; + + TStateBuildingSettings stateSettings; + TCoordinationState state(map, snapshot, context, stateSettings); + + TSortedHubs sortedHubs{ + {LF(80), HUB("hub-1")}, + {LF(20), HUB("hub-2")}, + }; + + TAssignedPartitions assignedPartitions{ + {PID(1), HUB("hub-1")}, + }; + + TMigrationContext migrationContext; + TBalancingSettings settings; + settings.BalancingThresholdCV = 0.0; + settings.BalancingTargetCV = 0.0; + settings.MaxRebalancePhases = 0; + + EXPECT_CALL(*Predictor_, PredictLoadFactor(_, _, _)).Times(0); + + RebalancePartitions( + sortedHubs, + assignedPartitions, + migrationContext, + Predictor_, + state, + settings); + + EXPECT_EQ(assignedPartitions.front().second, HUB("hub-1")); +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace NCoordinator::NCore diff --git a/src/infra/serializer/serializer_ut.cpp b/src/infra/serializer/serializer_ut.cpp index 0ffdd02..5a9ff63 100644 --- a/src/infra/serializer/serializer_ut.cpp +++ b/src/infra/serializer/serializer_ut.cpp @@ -1 +1,126 @@ -// TODO \ No newline at end of file +#include "serializer.hpp" + +#include + +#include +#include + +namespace { + +using namespace NCoordinator::NCore::NDomain; +using namespace NCoordinator::NInfra; + +//////////////////////////////////////////////////////////////////////////////// + +TPartitionMap MakePartitionMap() +{ + TPartitionMap map; + map.Epoch = TEpoch{42}; + + map.Partitions.emplace_back( + TPartitionId{1}, + THubEndpoint{"hub-1"} + ); + map.Partitions.emplace_back( + TPartitionId{2}, + THubEndpoint{"hub-2"} + ); + + return map; +} + +//////////////////////////////////////////////////////////////////////////////// + +} // anonymous namespace + +TEST(PartitionMapSerializer, SerializeProducesValidJson) +{ + const auto map = MakePartitionMap(); + + const auto json = SerializePartitionMap(map); + + ASSERT_TRUE(json.HasMember("epoch")); + ASSERT_TRUE(json.HasMember("partitions")); + + EXPECT_EQ(json["epoch"].As(), 42); + + const auto& partitions = json["partitions"]; + ASSERT_EQ(partitions.GetSize(), 2); + + EXPECT_EQ(partitions[0]["id"].As(), 1); + EXPECT_EQ(partitions[0]["hub"].As(), "hub-1"); + + EXPECT_EQ(partitions[1]["id"].As(), 2); + EXPECT_EQ(partitions[1]["hub"].As(), "hub-2"); +} + +TEST(PartitionMapSerializer, DeserializeRestoresOriginalData) +{ + const auto original = MakePartitionMap(); + + const auto json = SerializePartitionMap(original); + const auto restored = DeserializePartitionMap(json); + + EXPECT_EQ(restored.Epoch, original.Epoch); + ASSERT_EQ(restored.Partitions.size(), original.Partitions.size()); + + for (size_t i = 0; i < original.Partitions.size(); ++i) { + EXPECT_EQ(restored.Partitions[i].first, original.Partitions[i].first); + EXPECT_EQ(restored.Partitions[i].second, original.Partitions[i].second); + } +} + +TEST(PartitionMapSerializer, HandlesEmptyPartitions) +{ + TPartitionMap map; + map.Epoch = TEpoch{123}; + + const auto json = SerializePartitionMap(map); + const auto restored = DeserializePartitionMap(json); + + EXPECT_EQ(restored.Epoch, TEpoch{123}); + EXPECT_TRUE(restored.Partitions.empty()); +} + +TEST(HubReportDeserializer, ParsesValidReport) +{ + userver::formats::json::ValueBuilder json; + + json["epoch"] = std::to_string(100); + json["endpoint"] = "hub-42"; + json["dc"] = "eu-west"; + json["load_factor"] = 73; + + userver::formats::json::ValueBuilder partitions; + partitions["1"] = std::to_string(10); + partitions["2"] = std::to_string(20); + + json["partitions"] = partitions.ExtractValue(); + + const auto report = DeserializeHubReport(json.ExtractValue()); + + EXPECT_EQ(report.Epoch, TEpoch{100}); + EXPECT_EQ(report.Endpoint, THubEndpoint{"hub-42"}); + EXPECT_EQ(report.DC, THubDC{"eu-west"}); + EXPECT_EQ(report.LoadFactor, TLoadFactor{73}); + + ASSERT_EQ(report.PartitionWeights.size(), 2); + + EXPECT_EQ(report.PartitionWeights.at(TPartitionId{1}), TPartitionWeight{10}); + EXPECT_EQ(report.PartitionWeights.at(TPartitionId{2}), TPartitionWeight{20}); +} + +TEST(HubReportDeserializer, HandlesEmptyPartitions) +{ + userver::formats::json::ValueBuilder json; + + json["epoch"] = "1"; + json["endpoint"] = "hub"; + json["dc"] = "dc"; + json["load_factor"] = 0; + json["partitions"] = userver::formats::json::ValueBuilder{}.ExtractValue(); + + const auto report = DeserializeHubReport(json.ExtractValue()); + + EXPECT_TRUE(report.PartitionWeights.empty()); +}