Loading include/common/rpc/replica_migrator.hpp +42 −2 Changes for include/common/rpc/replica_migrator.hpp: 42 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -13,6 +13,7 @@ #include <common/rpc/distributor.hpp> #include <algorithm> #include <cstdint> #include <stdexcept> #include <string> #include <vector> Loading @@ -26,12 +27,16 @@ struct replica_migration_job { int target_copy{}; host_t source_node{}; host_t target_node{}; uint64_t old_placement_generation{}; uint64_t new_placement_generation{}; }; struct logical_chunk_replica_plan { std::vector<host_t> old_hosts; std::vector<host_t> new_hosts; std::vector<replica_migration_job> jobs; uint64_t old_placement_generation{}; uint64_t new_placement_generation{}; }; /** Loading @@ -46,12 +51,16 @@ inline logical_chunk_replica_plan plan_replica_migration(const Distributor& old_distributor, const Distributor& new_distributor, const std::string& path, const chunkid_t chunk_id, const int replica_count) { const int replica_count, const uint64_t old_placement_generation = 0, const uint64_t new_placement_generation = 1) { if(replica_count < 0) { throw std::invalid_argument("replica count must not be negative"); } logical_chunk_replica_plan plan; plan.old_placement_generation = old_placement_generation; plan.new_placement_generation = new_placement_generation; plan.old_hosts.reserve(static_cast<std::size_t>(replica_count) + 1); plan.new_hosts.reserve(static_cast<std::size_t>(replica_count) + 1); for(int copy = 0; copy <= replica_count; ++copy) { Loading Loading @@ -93,11 +102,42 @@ plan_replica_migration(const Distributor& old_distributor, std::distance(plan.old_hosts.begin(), source_copy)); plan.jobs.push_back({path, chunk_id, source_copy_index, static_cast<int>(target_copy), *source_copy, target_node}); target_node, old_placement_generation, new_placement_generation}); } return plan; } inline bool validate_replica_migration_plan(const logical_chunk_replica_plan& plan, const uint64_t expected_old_generation, const uint64_t expected_new_generation) { if(plan.old_placement_generation != expected_old_generation || plan.new_placement_generation != expected_new_generation || expected_old_generation >= expected_new_generation || plan.old_hosts.empty() || plan.new_hosts.empty()) { return false; } std::vector<int> target_copies; target_copies.reserve(plan.jobs.size()); for(const auto& job : plan.jobs) { if(job.old_placement_generation != plan.old_placement_generation || job.new_placement_generation != plan.new_placement_generation || job.source_copy < 0 || job.target_copy < 0 || static_cast<std::size_t>(job.source_copy) >= plan.old_hosts.size() || static_cast<std::size_t>(job.target_copy) >= plan.new_hosts.size() || plan.old_hosts[job.source_copy] != job.source_node || plan.new_hosts[job.target_copy] != job.target_node || std::find(target_copies.begin(), target_copies.end(), job.target_copy) != target_copies.end()) { return false; } target_copies.push_back(job.target_copy); } return true; } } // namespace gkfs::rpc #endif // GKFS_RPC_REPLICA_MIGRATOR_HPP No newline at end of file tests/unit/test_replica_migrator.cpp +32 −1 Changes for tests/unit/test_replica_migrator.cpp: 32 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -68,16 +68,20 @@ TEST_CASE("replica migration plan preserves logical copy placement", const fixed_distributor new_distributor({4, 5, 6}); const auto plan = gkfs::rpc::plan_replica_migration( old_distributor, new_distributor, "/file", 7, 2); old_distributor, new_distributor, "/file", 7, 2, 4, 5); REQUIRE(plan.old_hosts == std::vector<gkfs::rpc::host_t>{1, 2, 3}); REQUIRE(plan.new_hosts == std::vector<gkfs::rpc::host_t>{4, 5, 6}); REQUIRE(plan.old_placement_generation == 4); REQUIRE(plan.new_placement_generation == 5); REQUIRE(plan.jobs.size() == 3); for(const auto& job : plan.jobs) { REQUIRE(job.path == "/file"); REQUIRE(job.chunk_id == 7); REQUIRE(job.source_copy == 0); REQUIRE(job.source_node == 1); REQUIRE(job.old_placement_generation == 4); REQUIRE(job.new_placement_generation == 5); } REQUIRE(plan.jobs[0].target_copy == 0); REQUIRE(plan.jobs[1].target_copy == 1); Loading @@ -101,6 +105,33 @@ TEST_CASE("replica migration plan reuses existing hosts and deduplicates targets 1); } TEST_CASE("replica migration plan validates placement generations and hosts", "[replication][mutation]") { const fixed_distributor old_distributor({1, 2, 3}); const fixed_distributor new_distributor({2, 4, 4}); const auto plan = gkfs::rpc::plan_replica_migration( old_distributor, new_distributor, "/file", 9, 2, 10, 11); REQUIRE(gkfs::rpc::validate_replica_migration_plan(plan, 10, 11)); REQUIRE_FALSE(gkfs::rpc::validate_replica_migration_plan(plan, 9, 11)); auto malformed = plan; malformed.jobs.front().target_node = 99; REQUIRE_FALSE( gkfs::rpc::validate_replica_migration_plan(malformed, 10, 11)); } TEST_CASE("replica migration plan rejects duplicate target copies", "[replication][mutation]") { const fixed_distributor old_distributor({1, 2, 3}); const fixed_distributor new_distributor({4, 5, 6}); auto plan = gkfs::rpc::plan_replica_migration( old_distributor, new_distributor, "/file", 9, 2, 1, 2); plan.jobs.push_back(plan.jobs.front()); REQUIRE_FALSE(gkfs::rpc::validate_replica_migration_plan(plan, 1, 2)); } TEST_CASE("replica migration plan rejects negative replica counts", "[replication][mutation]") { const fixed_distributor distributor({1}); Loading Loading
include/common/rpc/replica_migrator.hpp +42 −2 Changes for include/common/rpc/replica_migrator.hpp: 42 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -13,6 +13,7 @@ #include <common/rpc/distributor.hpp> #include <algorithm> #include <cstdint> #include <stdexcept> #include <string> #include <vector> Loading @@ -26,12 +27,16 @@ struct replica_migration_job { int target_copy{}; host_t source_node{}; host_t target_node{}; uint64_t old_placement_generation{}; uint64_t new_placement_generation{}; }; struct logical_chunk_replica_plan { std::vector<host_t> old_hosts; std::vector<host_t> new_hosts; std::vector<replica_migration_job> jobs; uint64_t old_placement_generation{}; uint64_t new_placement_generation{}; }; /** Loading @@ -46,12 +51,16 @@ inline logical_chunk_replica_plan plan_replica_migration(const Distributor& old_distributor, const Distributor& new_distributor, const std::string& path, const chunkid_t chunk_id, const int replica_count) { const int replica_count, const uint64_t old_placement_generation = 0, const uint64_t new_placement_generation = 1) { if(replica_count < 0) { throw std::invalid_argument("replica count must not be negative"); } logical_chunk_replica_plan plan; plan.old_placement_generation = old_placement_generation; plan.new_placement_generation = new_placement_generation; plan.old_hosts.reserve(static_cast<std::size_t>(replica_count) + 1); plan.new_hosts.reserve(static_cast<std::size_t>(replica_count) + 1); for(int copy = 0; copy <= replica_count; ++copy) { Loading Loading @@ -93,11 +102,42 @@ plan_replica_migration(const Distributor& old_distributor, std::distance(plan.old_hosts.begin(), source_copy)); plan.jobs.push_back({path, chunk_id, source_copy_index, static_cast<int>(target_copy), *source_copy, target_node}); target_node, old_placement_generation, new_placement_generation}); } return plan; } inline bool validate_replica_migration_plan(const logical_chunk_replica_plan& plan, const uint64_t expected_old_generation, const uint64_t expected_new_generation) { if(plan.old_placement_generation != expected_old_generation || plan.new_placement_generation != expected_new_generation || expected_old_generation >= expected_new_generation || plan.old_hosts.empty() || plan.new_hosts.empty()) { return false; } std::vector<int> target_copies; target_copies.reserve(plan.jobs.size()); for(const auto& job : plan.jobs) { if(job.old_placement_generation != plan.old_placement_generation || job.new_placement_generation != plan.new_placement_generation || job.source_copy < 0 || job.target_copy < 0 || static_cast<std::size_t>(job.source_copy) >= plan.old_hosts.size() || static_cast<std::size_t>(job.target_copy) >= plan.new_hosts.size() || plan.old_hosts[job.source_copy] != job.source_node || plan.new_hosts[job.target_copy] != job.target_node || std::find(target_copies.begin(), target_copies.end(), job.target_copy) != target_copies.end()) { return false; } target_copies.push_back(job.target_copy); } return true; } } // namespace gkfs::rpc #endif // GKFS_RPC_REPLICA_MIGRATOR_HPP No newline at end of file
tests/unit/test_replica_migrator.cpp +32 −1 Changes for tests/unit/test_replica_migrator.cpp: 32 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -68,16 +68,20 @@ TEST_CASE("replica migration plan preserves logical copy placement", const fixed_distributor new_distributor({4, 5, 6}); const auto plan = gkfs::rpc::plan_replica_migration( old_distributor, new_distributor, "/file", 7, 2); old_distributor, new_distributor, "/file", 7, 2, 4, 5); REQUIRE(plan.old_hosts == std::vector<gkfs::rpc::host_t>{1, 2, 3}); REQUIRE(plan.new_hosts == std::vector<gkfs::rpc::host_t>{4, 5, 6}); REQUIRE(plan.old_placement_generation == 4); REQUIRE(plan.new_placement_generation == 5); REQUIRE(plan.jobs.size() == 3); for(const auto& job : plan.jobs) { REQUIRE(job.path == "/file"); REQUIRE(job.chunk_id == 7); REQUIRE(job.source_copy == 0); REQUIRE(job.source_node == 1); REQUIRE(job.old_placement_generation == 4); REQUIRE(job.new_placement_generation == 5); } REQUIRE(plan.jobs[0].target_copy == 0); REQUIRE(plan.jobs[1].target_copy == 1); Loading @@ -101,6 +105,33 @@ TEST_CASE("replica migration plan reuses existing hosts and deduplicates targets 1); } TEST_CASE("replica migration plan validates placement generations and hosts", "[replication][mutation]") { const fixed_distributor old_distributor({1, 2, 3}); const fixed_distributor new_distributor({2, 4, 4}); const auto plan = gkfs::rpc::plan_replica_migration( old_distributor, new_distributor, "/file", 9, 2, 10, 11); REQUIRE(gkfs::rpc::validate_replica_migration_plan(plan, 10, 11)); REQUIRE_FALSE(gkfs::rpc::validate_replica_migration_plan(plan, 9, 11)); auto malformed = plan; malformed.jobs.front().target_node = 99; REQUIRE_FALSE( gkfs::rpc::validate_replica_migration_plan(malformed, 10, 11)); } TEST_CASE("replica migration plan rejects duplicate target copies", "[replication][mutation]") { const fixed_distributor old_distributor({1, 2, 3}); const fixed_distributor new_distributor({4, 5, 6}); auto plan = gkfs::rpc::plan_replica_migration( old_distributor, new_distributor, "/file", 9, 2, 1, 2); plan.jobs.push_back(plan.jobs.front()); REQUIRE_FALSE(gkfs::rpc::validate_replica_migration_plan(plan, 1, 2)); } TEST_CASE("replica migration plan rejects negative replica counts", "[replication][mutation]") { const fixed_distributor distributor({1}); Loading