Loading include/common/rpc/cutshift_sorted.hpp +11 −10 Original line number Diff line number Diff line Loading @@ -31,18 +31,20 @@ namespace rpc { /// Collect gaps from shrinking nodes using CutShift+Sorted algorithm /// @param old_partitions Current partitions before expansion /// @param reductions Map of host_id -> capacity reduction amount /// @return Sorted list of gaps (intervals with host_id=0 indicating empty space) std::vector<Interval> collect_gaps_cutshift( const std::vector<Partition>& old_partitions, /// @return Sorted list of gaps (intervals with host_id=0 indicating empty /// space) std::vector<Interval> collect_gaps_cutshift(const std::vector<Partition>& old_partitions, const std::unordered_map<host_t, float>& reductions); /// Assign gaps to new nodes using greedy largest-first packing /// @param gaps List of available gaps /// @param new_hosts List of new host IDs needing intervals /// @param old_partitions Reference to old partitions (for capacity calculations) /// @param old_partitions Reference to old partitions (for capacity /// calculations) /// @return Updated partitions for new nodes std::vector<Partition> assign_gaps_to_new_nodes( std::vector<Interval> gaps, std::vector<Partition> assign_gaps_to_new_nodes(std::vector<Interval> gaps, const std::vector<host_t>& new_hosts, const std::vector<Partition>& old_partitions); Loading @@ -52,11 +54,10 @@ std::vector<Partition> assign_gaps_to_new_nodes( /// @param old_total_capacity Total capacity before expansion /// @param new_total_capacity Total capacity after expansion /// @return Updated partition table (old + new nodes) std::vector<Partition> expand_with_cutshift( std::vector<Partition> current_partitions, std::vector<Partition> expand_with_cutshift(std::vector<Partition> current_partitions, const std::vector<host_t>& new_hosts, float old_total_capacity, float new_total_capacity); float old_total_capacity, float new_total_capacity); } // namespace rpc } // namespace gkfs Loading include/common/rpc/data_migration_executor.hpp +15 −13 Original line number Diff line number Diff line Loading @@ -29,33 +29,34 @@ namespace gkfs { namespace rpc { /// Callback type for migration progress reporting using MigrationProgressCallback = std::function<void(size_t completed, size_t total)>; using MigrationProgressCallback = std::function<void(size_t completed, size_t total)>; /// Migration result status enum class MigrationStatus { Success, PartialFailure, AllFailed, Cancelled }; enum class MigrationStatus { Success, PartialFailure, AllFailed, Cancelled }; /// Execute a migration plan class DataMigrationExecutor { public: /// Execute all migration jobs with optional progress callback MigrationStatus execute(std::vector<MigrationJob>& jobs, MigrationStatus execute(std::vector<MigrationJob>& jobs, MigrationProgressCallback progress = nullptr); /// Execute migration jobs in batches of given size MigrationStatus execute_batched(std::vector<MigrationJob>& jobs, size_t batch_size, MigrationStatus execute_batched(std::vector<MigrationJob>& jobs, size_t batch_size, MigrationProgressCallback progress = nullptr); /// Get total bytes migrated (for accounting) size_t total_bytes_migrated() const { return total_bytes_.load(); } size_t total_bytes_migrated() const { return total_bytes_.load(); } /// Reset counters void reset() { void reset() { total_bytes_.store(0); success_count_.store(0); fail_count_.store(0); Loading @@ -67,7 +68,8 @@ public: size_t success_count; size_t fail_count; }; Stats get_stats() const { Stats get_stats() const { return {total_bytes_.load(), success_count_.load(), fail_count_.load()}; } Loading include/common/rpc/data_migrator.hpp +19 −13 Original line number Diff line number Diff line Loading @@ -38,43 +38,49 @@ struct MigrationJob { }; /// Helper to find which host owns a position in a partition table host_t find_host_for(const std::vector<Partition>& partitions, const std::string& path, chunkid_t chnk_id); host_t find_host_for(const std::vector<Partition>& partitions, const std::string& path, chunkid_t chnk_id); /// Compare old vs new partitions and compute migration jobs class DataMigrator { public: /// Compute which chunks need to move between old and new partitioning schemes. /// Compute which chunks need to move between old and new partitioning /// schemes. /// @param old_partitions The previous partition layout /// @param new_partitions The new partition layout /// @param chunk_sample_size Number of chunk IDs to sample per path for migration estimation /// @param chunk_sample_size Number of chunk IDs to sample per path for /// migration estimation /// @return List of migration jobs needed std::vector<MigrationJob> compute_migrations( const std::vector<Partition>& old_partitions, std::vector<MigrationJob> compute_migrations(const std::vector<Partition>& old_partitions, const std::vector<Partition>& new_partitions, int chunk_sample_size = 256); /// Get statistics about a migration plan struct MigrationStats { std::vector<MigrationJob> jobs; std::unordered_map<host_t, size_t> from_counts; // chunks leaving each host std::unordered_map<host_t, size_t> to_counts; // chunks entering each host std::unordered_map<host_t, size_t> from_counts; // chunks leaving each host std::unordered_map<host_t, size_t> to_counts; // chunks entering each host size_t total_chunks; size_t migrating_chunks; double migration_ratio; // migrating_chunks / total_chunks }; MigrationStats get_migration_stats( const std::vector<Partition>& old_partitions, MigrationStats get_migration_stats(const std::vector<Partition>& old_partitions, const std::vector<Partition>& new_partitions, int chunk_sample_size = 256); /// Print migration summary to stdout void print_migration_summary(const MigrationStats& stats); void print_migration_summary(const MigrationStats& stats); /// Check if migration is needed static bool needs_migration( const std::vector<Partition>& old_partitions, static bool needs_migration(const std::vector<Partition>& old_partitions, const std::vector<Partition>& new_partitions); }; Loading include/common/rpc/distribution_config.hpp +19 −11 Original line number Diff line number Diff line Loading @@ -31,30 +31,38 @@ namespace rpc { class DistributionConfig { public: /// Get the current distribution strategy DistributionStrategy get_strategy() const { return strategy_; } DistributionStrategy get_strategy() const { return strategy_; } /// Set strategy from string void set_strategy(const std::string& strategy) { void set_strategy(const std::string& strategy) { strategy_ = string_to_strategy(strategy); } /// Get strategy string for display std::string get_strategy_string() const { std::string get_strategy_string() const { return strategy_to_string(strategy_); } /// Check if random slicing is active bool is_random_slicing() const { bool is_random_slicing() const { return strategy_ == DistributionStrategy::RandomSlicing; } /// Check if simple hash is active bool is_simple_hash() const { bool is_simple_hash() const { return strategy_ == DistributionStrategy::SimpleHash; } /// Get the default strategy static DistributionStrategy default_strategy() { static DistributionStrategy default_strategy() { return DistributionStrategy::SimpleHash; } Loading @@ -63,14 +71,14 @@ private: }; /// Create a distributor from a DistributionConfig std::unique_ptr<Distributor> create_from_config(const DistributionConfig& config, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); std::unique_ptr<Distributor> create_from_config(const DistributionConfig& config, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); /// Read strategy from environment variable GKFS_DISTRIBUTION_STRATEGY /// Falls back to DistributionConfig::default_strategy() if not set DistributionStrategy read_strategy_from_env(); DistributionStrategy read_strategy_from_env(); } // namespace rpc } // namespace gkfs Loading include/common/rpc/distributor_factory.hpp +10 −12 Original line number Diff line number Diff line Loading @@ -37,10 +37,12 @@ enum class DistributionStrategy { }; /// Convert strategy enum to string const char* strategy_to_string(DistributionStrategy strategy); const char* strategy_to_string(DistributionStrategy strategy); /// Convert string to strategy enum DistributionStrategy string_to_strategy(const std::string& str); DistributionStrategy string_to_strategy(const std::string& str); /// Create a distributor based on the given strategy and parameters. /// @param strategy The distribution strategy to use Loading @@ -48,19 +50,15 @@ DistributionStrategy string_to_strategy(const std::string& str); /// @param hosts_size The number of hosts in the cluster /// @param fwd_host Optional forwarder host ID (used when strategy is Forwarder) /// @return A uniquely-owned distributor, or nullptr on error std::unique_ptr<Distributor> create_distributor( DistributionStrategy strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); std::unique_ptr<Distributor> create_distributor(DistributionStrategy strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); /// Convenience factory: creates a distributor from a string strategy name. /// Strings: "simple_hash", "random_slicing", "local_only", "forwarder" std::unique_ptr<Distributor> create_distributor_from_string( const std::string& strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); std::unique_ptr<Distributor> create_distributor_from_string(const std::string& strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); } // namespace rpc } // namespace gkfs Loading Loading
include/common/rpc/cutshift_sorted.hpp +11 −10 Original line number Diff line number Diff line Loading @@ -31,18 +31,20 @@ namespace rpc { /// Collect gaps from shrinking nodes using CutShift+Sorted algorithm /// @param old_partitions Current partitions before expansion /// @param reductions Map of host_id -> capacity reduction amount /// @return Sorted list of gaps (intervals with host_id=0 indicating empty space) std::vector<Interval> collect_gaps_cutshift( const std::vector<Partition>& old_partitions, /// @return Sorted list of gaps (intervals with host_id=0 indicating empty /// space) std::vector<Interval> collect_gaps_cutshift(const std::vector<Partition>& old_partitions, const std::unordered_map<host_t, float>& reductions); /// Assign gaps to new nodes using greedy largest-first packing /// @param gaps List of available gaps /// @param new_hosts List of new host IDs needing intervals /// @param old_partitions Reference to old partitions (for capacity calculations) /// @param old_partitions Reference to old partitions (for capacity /// calculations) /// @return Updated partitions for new nodes std::vector<Partition> assign_gaps_to_new_nodes( std::vector<Interval> gaps, std::vector<Partition> assign_gaps_to_new_nodes(std::vector<Interval> gaps, const std::vector<host_t>& new_hosts, const std::vector<Partition>& old_partitions); Loading @@ -52,11 +54,10 @@ std::vector<Partition> assign_gaps_to_new_nodes( /// @param old_total_capacity Total capacity before expansion /// @param new_total_capacity Total capacity after expansion /// @return Updated partition table (old + new nodes) std::vector<Partition> expand_with_cutshift( std::vector<Partition> current_partitions, std::vector<Partition> expand_with_cutshift(std::vector<Partition> current_partitions, const std::vector<host_t>& new_hosts, float old_total_capacity, float new_total_capacity); float old_total_capacity, float new_total_capacity); } // namespace rpc } // namespace gkfs Loading
include/common/rpc/data_migration_executor.hpp +15 −13 Original line number Diff line number Diff line Loading @@ -29,33 +29,34 @@ namespace gkfs { namespace rpc { /// Callback type for migration progress reporting using MigrationProgressCallback = std::function<void(size_t completed, size_t total)>; using MigrationProgressCallback = std::function<void(size_t completed, size_t total)>; /// Migration result status enum class MigrationStatus { Success, PartialFailure, AllFailed, Cancelled }; enum class MigrationStatus { Success, PartialFailure, AllFailed, Cancelled }; /// Execute a migration plan class DataMigrationExecutor { public: /// Execute all migration jobs with optional progress callback MigrationStatus execute(std::vector<MigrationJob>& jobs, MigrationStatus execute(std::vector<MigrationJob>& jobs, MigrationProgressCallback progress = nullptr); /// Execute migration jobs in batches of given size MigrationStatus execute_batched(std::vector<MigrationJob>& jobs, size_t batch_size, MigrationStatus execute_batched(std::vector<MigrationJob>& jobs, size_t batch_size, MigrationProgressCallback progress = nullptr); /// Get total bytes migrated (for accounting) size_t total_bytes_migrated() const { return total_bytes_.load(); } size_t total_bytes_migrated() const { return total_bytes_.load(); } /// Reset counters void reset() { void reset() { total_bytes_.store(0); success_count_.store(0); fail_count_.store(0); Loading @@ -67,7 +68,8 @@ public: size_t success_count; size_t fail_count; }; Stats get_stats() const { Stats get_stats() const { return {total_bytes_.load(), success_count_.load(), fail_count_.load()}; } Loading
include/common/rpc/data_migrator.hpp +19 −13 Original line number Diff line number Diff line Loading @@ -38,43 +38,49 @@ struct MigrationJob { }; /// Helper to find which host owns a position in a partition table host_t find_host_for(const std::vector<Partition>& partitions, const std::string& path, chunkid_t chnk_id); host_t find_host_for(const std::vector<Partition>& partitions, const std::string& path, chunkid_t chnk_id); /// Compare old vs new partitions and compute migration jobs class DataMigrator { public: /// Compute which chunks need to move between old and new partitioning schemes. /// Compute which chunks need to move between old and new partitioning /// schemes. /// @param old_partitions The previous partition layout /// @param new_partitions The new partition layout /// @param chunk_sample_size Number of chunk IDs to sample per path for migration estimation /// @param chunk_sample_size Number of chunk IDs to sample per path for /// migration estimation /// @return List of migration jobs needed std::vector<MigrationJob> compute_migrations( const std::vector<Partition>& old_partitions, std::vector<MigrationJob> compute_migrations(const std::vector<Partition>& old_partitions, const std::vector<Partition>& new_partitions, int chunk_sample_size = 256); /// Get statistics about a migration plan struct MigrationStats { std::vector<MigrationJob> jobs; std::unordered_map<host_t, size_t> from_counts; // chunks leaving each host std::unordered_map<host_t, size_t> to_counts; // chunks entering each host std::unordered_map<host_t, size_t> from_counts; // chunks leaving each host std::unordered_map<host_t, size_t> to_counts; // chunks entering each host size_t total_chunks; size_t migrating_chunks; double migration_ratio; // migrating_chunks / total_chunks }; MigrationStats get_migration_stats( const std::vector<Partition>& old_partitions, MigrationStats get_migration_stats(const std::vector<Partition>& old_partitions, const std::vector<Partition>& new_partitions, int chunk_sample_size = 256); /// Print migration summary to stdout void print_migration_summary(const MigrationStats& stats); void print_migration_summary(const MigrationStats& stats); /// Check if migration is needed static bool needs_migration( const std::vector<Partition>& old_partitions, static bool needs_migration(const std::vector<Partition>& old_partitions, const std::vector<Partition>& new_partitions); }; Loading
include/common/rpc/distribution_config.hpp +19 −11 Original line number Diff line number Diff line Loading @@ -31,30 +31,38 @@ namespace rpc { class DistributionConfig { public: /// Get the current distribution strategy DistributionStrategy get_strategy() const { return strategy_; } DistributionStrategy get_strategy() const { return strategy_; } /// Set strategy from string void set_strategy(const std::string& strategy) { void set_strategy(const std::string& strategy) { strategy_ = string_to_strategy(strategy); } /// Get strategy string for display std::string get_strategy_string() const { std::string get_strategy_string() const { return strategy_to_string(strategy_); } /// Check if random slicing is active bool is_random_slicing() const { bool is_random_slicing() const { return strategy_ == DistributionStrategy::RandomSlicing; } /// Check if simple hash is active bool is_simple_hash() const { bool is_simple_hash() const { return strategy_ == DistributionStrategy::SimpleHash; } /// Get the default strategy static DistributionStrategy default_strategy() { static DistributionStrategy default_strategy() { return DistributionStrategy::SimpleHash; } Loading @@ -63,14 +71,14 @@ private: }; /// Create a distributor from a DistributionConfig std::unique_ptr<Distributor> create_from_config(const DistributionConfig& config, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); std::unique_ptr<Distributor> create_from_config(const DistributionConfig& config, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); /// Read strategy from environment variable GKFS_DISTRIBUTION_STRATEGY /// Falls back to DistributionConfig::default_strategy() if not set DistributionStrategy read_strategy_from_env(); DistributionStrategy read_strategy_from_env(); } // namespace rpc } // namespace gkfs Loading
include/common/rpc/distributor_factory.hpp +10 −12 Original line number Diff line number Diff line Loading @@ -37,10 +37,12 @@ enum class DistributionStrategy { }; /// Convert strategy enum to string const char* strategy_to_string(DistributionStrategy strategy); const char* strategy_to_string(DistributionStrategy strategy); /// Convert string to strategy enum DistributionStrategy string_to_strategy(const std::string& str); DistributionStrategy string_to_strategy(const std::string& str); /// Create a distributor based on the given strategy and parameters. /// @param strategy The distribution strategy to use Loading @@ -48,19 +50,15 @@ DistributionStrategy string_to_strategy(const std::string& str); /// @param hosts_size The number of hosts in the cluster /// @param fwd_host Optional forwarder host ID (used when strategy is Forwarder) /// @return A uniquely-owned distributor, or nullptr on error std::unique_ptr<Distributor> create_distributor( DistributionStrategy strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); std::unique_ptr<Distributor> create_distributor(DistributionStrategy strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); /// Convenience factory: creates a distributor from a string strategy name. /// Strings: "simple_hash", "random_slicing", "local_only", "forwarder" std::unique_ptr<Distributor> create_distributor_from_string( const std::string& strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); std::unique_ptr<Distributor> create_distributor_from_string(const std::string& strategy, host_t localhost, unsigned int hosts_size, host_t fwd_host = 0); } // namespace rpc } // namespace gkfs Loading