Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/core/partition_balancing/balancing_impl.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
219 changes: 219 additions & 0 deletions src/core/partition_balancing/tests/rebalance_partitions_ut.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,219 @@
#include <core/partition_balancing/balancing_impl.hpp>

#include "test_helpers_ut.hpp"
#include "load_factor_predictor_mock.hpp"

#include <gtest/gtest.h>

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
127 changes: 126 additions & 1 deletion src/infra/serializer/serializer_ut.cpp
Original file line number Diff line number Diff line change
@@ -1 +1,126 @@
// TODO
#include "serializer.hpp"

#include <gtest/gtest.h>

#include <userver/formats/json.hpp>
#include <userver/formats/json/value_builder.hpp>

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<uint64_t>(), 42);

const auto& partitions = json["partitions"];
ASSERT_EQ(partitions.GetSize(), 2);

EXPECT_EQ(partitions[0]["id"].As<uint64_t>(), 1);
EXPECT_EQ(partitions[0]["hub"].As<std::string>(), "hub-1");

EXPECT_EQ(partitions[1]["id"].As<uint64_t>(), 2);
EXPECT_EQ(partitions[1]["hub"].As<std::string>(), "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());
}