Loading include/client/env.hpp +2 −0 Changes for include/client/env.hpp: 2 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -86,6 +86,8 @@ static constexpr auto RANGE_FD = ADD_PREFIX("RANGE_FD"); static constexpr auto DIRENTS_BUFF_SIZE = ADD_PREFIX("DIRENTS_BUFF_SIZE"); static constexpr auto NUM_REPL = ADD_PREFIX("NUM_REPL"); static constexpr auto REPAIR_STATUS_PATH = ADD_PREFIX("REPAIR_STATUS_PATH"); static constexpr auto REPAIR_JOURNAL_PATH = ADD_PREFIX("REPAIR_JOURNAL_PATH"); static constexpr auto PROXY_PID_FILE = ADD_PREFIX("PROXY_PID_FILE"); namespace cache { static constexpr auto DENTRY = ADD_PREFIX("DENTRY_CACHE"); Loading include/client/preload_context.hpp +15 −0 Changes for include/client/preload_context.hpp: 15 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -200,6 +200,9 @@ private: std::unordered_map<std::string, uint64_t> repair_generations_; mutable std::mutex repair_generation_mutex_; std::unordered_map<gkfs::rpc::host_t, uint32_t> active_repairs_; std::vector<gkfs::rpc::repair_task> in_flight_repairs_; std::string repair_status_path_; std::string repair_journal_path_; static constexpr uint32_t max_repairs_per_host_ = 1; static constexpr uint32_t repair_worker_count_ = 2; Loading Loading @@ -518,6 +521,18 @@ public: bool enqueue_repair_task(gkfs::rpc::repair_task task); void publish_repair_status_locked() const; void load_repair_journal(); void persist_repair_journal_locked() const; void remove_in_flight_repair(const gkfs::rpc::repair_task& task); void set_repair_executor( std::function<bool(const gkfs::rpc::repair_task&)> executor); Loading include/client/rpc/repair_queue.hpp +12 −0 Changes for include/client/rpc/repair_queue.hpp: 12 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -18,6 +18,7 @@ #include <optional> #include <string> #include <unordered_map> #include <vector> namespace gkfs::rpc { Loading Loading @@ -177,6 +178,17 @@ public: stale_.load(), coalesced_.load(), capacity_dropped_.load()}; } std::vector<repair_task> snapshot() const { std::vector<repair_task> result; result.reserve(tasks_.size()); for(const auto& [key, task] : tasks_) { (void) key; result.push_back(task); } return result; } private: using key_type = std::string; Loading include/common/rpc/distribution_config.hpp +8 −0 Changes for include/common/rpc/distribution_config.hpp: 8 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -22,6 +22,7 @@ #define GKFS_RPC_DISTRIBUTION_CONFIG_H #include "common/rpc/distributor_factory.hpp" #include <stdexcept> #include <string> namespace gkfs { Loading Loading @@ -60,6 +61,13 @@ public: return strategy_ == DistributionStrategy::SimpleHash; } /// Random Slicing does not yet define independent replica placement. bool supports_replication(const int replica_count) const { return replica_count == 0 || strategy_ != DistributionStrategy::RandomSlicing; } /// Get the default strategy static DistributionStrategy default_strategy() { Loading src/client/preload.cpp +6 −0 Changes for src/client/preload.cpp: 6 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -369,6 +369,12 @@ init_environment() { if(env_val != nullptr && env_val[0] != '\0') { config.set_strategy(env_val); } if(!config.supports_replication(CTX->get_replicas())) { exit_error_msg( EXIT_FAILURE, "LIBGKFS_NUM_REPL is not supported with random_slicing " "distribution"); } LOG(INFO, "{}() Distribution strategy: '{}'", __func__, config.get_strategy_string()); Loading Loading
include/client/env.hpp +2 −0 Changes for include/client/env.hpp: 2 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -86,6 +86,8 @@ static constexpr auto RANGE_FD = ADD_PREFIX("RANGE_FD"); static constexpr auto DIRENTS_BUFF_SIZE = ADD_PREFIX("DIRENTS_BUFF_SIZE"); static constexpr auto NUM_REPL = ADD_PREFIX("NUM_REPL"); static constexpr auto REPAIR_STATUS_PATH = ADD_PREFIX("REPAIR_STATUS_PATH"); static constexpr auto REPAIR_JOURNAL_PATH = ADD_PREFIX("REPAIR_JOURNAL_PATH"); static constexpr auto PROXY_PID_FILE = ADD_PREFIX("PROXY_PID_FILE"); namespace cache { static constexpr auto DENTRY = ADD_PREFIX("DENTRY_CACHE"); Loading
include/client/preload_context.hpp +15 −0 Changes for include/client/preload_context.hpp: 15 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -200,6 +200,9 @@ private: std::unordered_map<std::string, uint64_t> repair_generations_; mutable std::mutex repair_generation_mutex_; std::unordered_map<gkfs::rpc::host_t, uint32_t> active_repairs_; std::vector<gkfs::rpc::repair_task> in_flight_repairs_; std::string repair_status_path_; std::string repair_journal_path_; static constexpr uint32_t max_repairs_per_host_ = 1; static constexpr uint32_t repair_worker_count_ = 2; Loading Loading @@ -518,6 +521,18 @@ public: bool enqueue_repair_task(gkfs::rpc::repair_task task); void publish_repair_status_locked() const; void load_repair_journal(); void persist_repair_journal_locked() const; void remove_in_flight_repair(const gkfs::rpc::repair_task& task); void set_repair_executor( std::function<bool(const gkfs::rpc::repair_task&)> executor); Loading
include/client/rpc/repair_queue.hpp +12 −0 Changes for include/client/rpc/repair_queue.hpp: 12 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -18,6 +18,7 @@ #include <optional> #include <string> #include <unordered_map> #include <vector> namespace gkfs::rpc { Loading Loading @@ -177,6 +178,17 @@ public: stale_.load(), coalesced_.load(), capacity_dropped_.load()}; } std::vector<repair_task> snapshot() const { std::vector<repair_task> result; result.reserve(tasks_.size()); for(const auto& [key, task] : tasks_) { (void) key; result.push_back(task); } return result; } private: using key_type = std::string; Loading
include/common/rpc/distribution_config.hpp +8 −0 Changes for include/common/rpc/distribution_config.hpp: 8 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -22,6 +22,7 @@ #define GKFS_RPC_DISTRIBUTION_CONFIG_H #include "common/rpc/distributor_factory.hpp" #include <stdexcept> #include <string> namespace gkfs { Loading Loading @@ -60,6 +61,13 @@ public: return strategy_ == DistributionStrategy::SimpleHash; } /// Random Slicing does not yet define independent replica placement. bool supports_replication(const int replica_count) const { return replica_count == 0 || strategy_ != DistributionStrategy::RandomSlicing; } /// Get the default strategy static DistributionStrategy default_strategy() { Loading
src/client/preload.cpp +6 −0 Changes for src/client/preload.cpp: 6 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -369,6 +369,12 @@ init_environment() { if(env_val != nullptr && env_val[0] != '\0') { config.set_strategy(env_val); } if(!config.supports_replication(CTX->get_replicas())) { exit_error_msg( EXIT_FAILURE, "LIBGKFS_NUM_REPL is not supported with random_slicing " "distribution"); } LOG(INFO, "{}() Distribution strategy: '{}'", __func__, config.get_strategy_string()); Loading