Loading include/client/preload_context.hpp +92 −0 Changes for include/client/preload_context.hpp: 92 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -46,6 +46,7 @@ #include <future> #include <thread> #include <condition_variable> #include <functional> #include <thallium.hpp> #include <memory> #include <vector> Loading @@ -53,6 +54,8 @@ #include <mutex> #include <unordered_set> #include <config.hpp> #include <client/rpc/repair_queue.hpp> #include <client/rpc/replica.hpp> #include <bitset> Loading Loading @@ -186,6 +189,20 @@ private: std::thread async_write_thread_; bool async_write_stop_{false}; gkfs::rpc::repair_queue repair_queue_; mutable std::mutex repair_queue_mutex_; std::condition_variable repair_queue_cv_; std::vector<std::thread> repair_worker_threads_; bool repair_worker_stop_{false}; std::chrono::steady_clock::time_point repair_worker_deadline_{}; std::function<bool(const gkfs::rpc::repair_task&)> repair_executor_; gkfs::rpc::host_health_table host_health_; 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_; static constexpr uint32_t max_repairs_per_host_ = 1; static constexpr uint32_t repair_worker_count_ = 2; std::string ofi_interface_; std::chrono::milliseconds rpc_timeout_{0}; Loading Loading @@ -361,6 +378,27 @@ public: int get_replicas(); void mark_host_failure(gkfs::rpc::host_t host); void mark_host_success(gkfs::rpc::host_t host); std::set<gkfs::rpc::host_t> unavailable_hosts(); uint64_t next_repair_generation(const std::string& path, gkfs::rpc::repair_kind kind, uint64_t chunk_id = 0); uint64_t current_repair_generation(const std::string& path, gkfs::rpc::repair_kind kind, uint64_t chunk_id = 0) const; bool is_current_repair_generation(const gkfs::rpc::repair_task& task) const; bool protect_fds() const; Loading Loading @@ -464,6 +502,60 @@ public: void wait_async_writes(); void start_repair_worker(); void stop_repair_worker(); void repair_worker(); std::optional<gkfs::rpc::host_t> repair_target_host(const gkfs::rpc::repair_task& task) const; bool enqueue_repair_task(gkfs::rpc::repair_task task); void set_repair_executor( std::function<bool(const gkfs::rpc::repair_task&)> executor); void enqueue_repairs(const std::string& path, const std::vector<int8_t>& failed_copies, int8_t source_copy = 0, uint64_t chunk_id = 0); void enqueue_repairs(const std::string& path, const std::vector<int8_t>& failed_copies, uint64_t first_chunk, uint64_t last_chunk, int8_t source_copy = 0); void enqueue_metadata_create_repairs(const std::string& path, mode_t mode, const std::vector<int8_t>& failed_copies); void enqueue_metadata_create_inline_repairs( const std::string& path, mode_t mode, const std::string& data, const std::vector<int8_t>& failed_copies); void enqueue_metadata_size_repairs(const std::string& path, uint64_t size, int64_t offset, bool append, bool clear_inline, const std::vector<int8_t>& failed_copies); void enqueue_metadata_inline_repairs(const std::string& path, const std::string& data, uint64_t offset, bool append, const std::vector<int8_t>& failed_copies); gkfs::rpc::repair_queue::counters repair_queue_stats() const; }; } // namespace preload Loading include/client/rpc/forward_data.hpp +7 −1 Changes for include/client/rpc/forward_data.hpp: 7 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -42,9 +42,12 @@ #include <common/common_defs.hpp> #include <algorithm> #include <string> #include <memory> #include <set> #include <common/rpc/distributor.hpp> #include <client/rpc/repair_queue.hpp> namespace gkfs::rpc { // TODO once we have LEAF, remove all the error code returns and throw them as Loading @@ -57,7 +60,10 @@ forward_write(const std::string& path, const void* buf, off64_t offset, std::pair<int, ssize_t> forward_read(const std::string& path, void* buf, off64_t offset, size_t read_size, const int8_t num_copies, std::set<int8_t>& failed); std::set<gkfs::rpc::host_t>& failed); bool repair_data_copy(const gkfs::rpc::repair_task& task); int forward_truncate(const std::string& path, size_t current_size, size_t new_size, Loading include/client/rpc/forward_metadata.hpp +6 −2 Changes for include/client/rpc/forward_metadata.hpp: 6 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -44,6 +44,7 @@ #include <memory> #include <vector> #include <cstdint> #include <client/rpc/repair_queue.hpp> /* Forward declaration */ namespace gkfs { namespace filemap { Loading @@ -63,6 +64,9 @@ namespace rpc { int forward_create(const std::string& path, mode_t mode, const int copy); int repair_metadata_copy(const gkfs::rpc::repair_task& task); int forward_batch_create(uint64_t host_id, const std::vector<std::string>& paths, const std::vector<uint32_t>& modes); Loading Loading @@ -130,7 +134,7 @@ forward_mk_symlink(const std::string& path, const std::string& target_path); */ std::pair<int, off64_t> forward_write_inline(const std::string& path, const void* buf, off64_t offset, size_t write_size, bool append_flag); size_t write_size, bool append_flag, const int num_copies); /** * @brief Send an RPC request to read a small amount of data directly Loading @@ -144,7 +148,7 @@ forward_write_inline(const std::string& path, const void* buf, off64_t offset, */ std::pair<int, ssize_t> forward_read_inline(const std::string& path, void* buf, off64_t offset, size_t read_size); size_t read_size, const int num_copies); std::tuple<int, std::vector<std::tuple<const std::string, unsigned char, size_t, Loading include/client/rpc/repair_queue.hpp 0 → 100644 +204 −0 Changes for include/client/rpc/repair_queue.hpp: 204 added lines, 0 removed lines. Original line number Diff line number Diff line /* Copyright 2018-2025, Barcelona Supercomputing Center (BSC), Spain Copyright 2015-2025, Johannes Gutenberg Universitaet Mainz, Germany This file is part of GekkoFS. SPDX-License-Identifier: LGPL-3.0-or-later */ #ifndef GKFS_CLIENT_RPC_REPAIR_QUEUE_HPP #define GKFS_CLIENT_RPC_REPAIR_QUEUE_HPP #include <chrono> #include <cstdint> #include <algorithm> #include <atomic> #include <deque> #include <optional> #include <string> #include <unordered_map> namespace gkfs::rpc { enum class repair_kind : uint8_t { metadata_create, metadata_create_inline, metadata_size, metadata_inline, data_chunk, deferred, }; struct repair_task { std::string path; uint64_t chunk_id{}; int8_t source_copy{}; int8_t target_copy{}; uint64_t generation{}; uint32_t retry_count{}; int last_error{}; std::chrono::steady_clock::time_point next_attempt{}; bool data_chunk{false}; repair_kind kind{repair_kind::deferred}; uint32_t mode{}; uint64_t size{}; int64_t offset{}; bool append{}; bool clear_inline{}; std::string data; }; class repair_queue { public: static constexpr std::size_t default_capacity = 1024; static constexpr uint32_t default_retry_limit = 5; static constexpr uint64_t default_backoff_ms = 100; struct counters { uint64_t queued{}; uint64_t completed{}; uint64_t failed{}; uint64_t stale{}; }; explicit repair_queue(std::size_t capacity = default_capacity) : capacity_(capacity) {} bool enqueue(repair_task task) { const auto key = make_key(task); const auto existing = tasks_.find(key); if(existing != tasks_.end()) { const auto retry_count = std::max(existing->second.retry_count, task.retry_count); existing->second = std::move(task); existing->second.retry_count = retry_count; return true; } if(tasks_.size() >= capacity_) { return false; } if(task.next_attempt == std::chrono::steady_clock::time_point{}) { task.next_attempt = std::chrono::steady_clock::now(); } order_.push_back(key); tasks_.emplace(key, std::move(task)); queued_++; return true; } bool requeue(repair_task task, const int error, const std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now(), const uint32_t retry_limit = default_retry_limit, const uint64_t backoff_ms = default_backoff_ms) { task.last_error = error; if(task.retry_count >= retry_limit) { failed_++; return false; } ++task.retry_count; const auto delay = backoff_ms * (uint64_t{1} << std::min(task.retry_count - 1, 16u)); task.next_attempt = now + std::chrono::milliseconds(delay); return enqueue(std::move(task)); } void defer(repair_task task, const std::chrono::steady_clock::time_point next_attempt) { task.next_attempt = next_attempt; const auto key = make_key(task); if(tasks_.find(key) != tasks_.end()) { return; } order_.push_back(key); tasks_.emplace(key, std::move(task)); } void complete() { completed_++; } void discard_stale() { stale_++; } std::optional<repair_task> pop_ready(const std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now()) { return pop_ready([](const repair_task&) { return true; }, now); } template <typename Predicate> std::optional<repair_task> pop_ready(Predicate predicate, const std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now()) { const auto count = order_.size(); for(std::size_t i = 0; i < count; ++i) { auto key = std::move(order_.front()); order_.pop_front(); auto it = tasks_.find(key); if(it == tasks_.end()) { continue; } if(it->second.next_attempt <= now && predicate(it->second)) { auto task = std::move(it->second); tasks_.erase(it); return task; } order_.push_back(std::move(key)); } return std::nullopt; } bool empty() const { return tasks_.empty(); } std::size_t size() const { return tasks_.size(); } counters stats() const { return {queued_.load(), completed_.load(), failed_.load(), stale_.load()}; } private: using key_type = std::string; static key_type make_key(const repair_task& task) { key_type key = task.path; key.push_back('\0'); key += std::to_string(task.chunk_id); key.push_back('\0'); key += std::to_string(task.source_copy); key.push_back('\0'); key += std::to_string(task.target_copy); key.push_back('\0'); key += std::to_string(static_cast<uint8_t>(task.kind)); return key; } std::size_t capacity_; std::deque<key_type> order_; std::unordered_map<key_type, repair_task> tasks_; std::atomic<uint64_t> queued_{0}; std::atomic<uint64_t> completed_{0}; std::atomic<uint64_t> failed_{0}; std::atomic<uint64_t> stale_{0}; }; } // namespace gkfs::rpc #endif // GKFS_CLIENT_RPC_REPAIR_QUEUE_HPP No newline at end of file include/client/rpc/replica.hpp 0 → 100644 +331 −0 Changes for include/client/rpc/replica.hpp: 331 added lines, 0 removed lines. Original line number Diff line number Diff line /* Copyright 2018-2025, Barcelona Supercomputing Center (BSC), Spain Copyright 2015-2025, Johannes Gutenberg Universitaet Mainz, Germany This file is part of GekkoFS. SPDX-License-Identifier: LGPL-3.0-or-later */ #ifndef GKFS_CLIENT_RPC_REPLICA_HPP #define GKFS_CLIENT_RPC_REPLICA_HPP #include <common/rpc/distributor.hpp> #include <cstddef> #include <cstdint> #include <cerrno> #include <algorithm> #include <stdexcept> #include <string> #include <set> #include <optional> #include <atomic> #include <chrono> #include <mutex> #include <unordered_map> #include <vector> namespace gkfs::rpc { struct replica_target { int copy{}; host_t host{}; }; enum class host_health_state { healthy, suspect, down, recovering, }; struct host_health_status { host_health_state state{host_health_state::healthy}; uint32_t failure_count{}; std::chrono::steady_clock::time_point last_failure{}; std::chrono::steady_clock::time_point next_probe{}; }; class host_health_table { public: using clock = std::chrono::steady_clock; void mark_failure(const host_t host, const clock::time_point now = clock::now()) { std::lock_guard<std::mutex> lock(mutex_); auto& status = hosts_[host]; ++status.failure_count; status.last_failure = now; status.state = status.failure_count < failure_threshold ? host_health_state::suspect : host_health_state::down; if(status.state == host_health_state::down) { const auto exponent = std::min(status.failure_count - failure_threshold, max_backoff_shift); status.next_probe = now + base_backoff * (uint64_t{1} << exponent); } } void mark_success(const host_t host) { std::lock_guard<std::mutex> lock(mutex_); auto& status = hosts_[host]; status = {}; } std::set<host_t> unavailable(const clock::time_point now = clock::now()) { std::lock_guard<std::mutex> lock(mutex_); std::set<host_t> result; for(auto& [host, status] : hosts_) { if(status.state == host_health_state::down) { if(now < status.next_probe) { result.insert(host); } else { status.state = host_health_state::recovering; } } } return result; } host_health_status status(const host_t host) const { std::lock_guard<std::mutex> lock(mutex_); const auto it = hosts_.find(host); return it == hosts_.end() ? host_health_status{} : it->second; } private: static constexpr uint32_t failure_threshold = 2; static constexpr uint32_t max_backoff_shift = 6; static constexpr auto base_backoff = std::chrono::milliseconds(100); mutable std::mutex mutex_; std::unordered_map<host_t, host_health_status> hosts_; }; enum class replica_policy { read_any, write_any, write_all, strict, primary_required, }; struct replica_operation_result { int error{}; uint8_t successful_copies{}; uint8_t expected_copies{}; std::vector<int8_t> failed_copies; bool degraded{}; bool repair_pending{}; }; struct replica_counters { std::atomic<uint64_t> read_primary{0}; std::atomic<uint64_t> read_fallback{0}; std::atomic<uint64_t> read_failure{0}; std::atomic<uint64_t> write_success{0}; std::atomic<uint64_t> write_degraded{0}; std::atomic<uint64_t> write_failure{0}; std::atomic<uint64_t> repair_queued{0}; std::atomic<uint64_t> repair_completed{0}; std::atomic<uint64_t> repair_failed{0}; std::atomic<uint64_t> repair_stale{0}; }; inline replica_counters& replica_metrics() { static replica_counters counters; return counters; } inline bool mutation_supported_with_replicas(const int replica_count) { return replica_count == 0; } inline void record_replica_read(const int copy) { if(copy == 0) { replica_metrics().read_primary.fetch_add(1); } else { replica_metrics().read_fallback.fetch_add(1); } } inline void record_replica_read_failure() { replica_metrics().read_failure.fetch_add(1); } inline replica_operation_result classify_replica_operation(const std::vector<std::pair<int, int>>& copy_errors, const replica_policy policy) { replica_operation_result result; result.expected_copies = static_cast<uint8_t>(copy_errors.size()); for(const auto& [copy, error] : copy_errors) { if(error == 0) { ++result.successful_copies; } else { result.failed_copies.push_back(static_cast<int8_t>(copy)); if(result.error == 0) { result.error = error; } } } const auto all_succeeded = result.successful_copies == result.expected_copies; const auto any_succeeded = result.successful_copies > 0; switch(policy) { case replica_policy::read_any: case replica_policy::write_any: if(any_succeeded) { result.error = 0; } break; case replica_policy::write_all: case replica_policy::strict: if(all_succeeded) { result.error = 0; } else if(result.error == 0) { result.error = EIO; } break; case replica_policy::primary_required: { const auto primary = std::find_if( copy_errors.begin(), copy_errors.end(), [](const auto& entry) { return entry.first == 0; }); if(primary == copy_errors.end() || primary->second != 0) { result.error = primary == copy_errors.end() ? EIO : primary->second; } else { result.error = 0; } break; } } result.degraded = any_succeeded && !all_succeeded; result.repair_pending = result.degraded; if(!any_succeeded && result.error == 0) { result.error = EIO; } return result; } inline void record_replica_write(const replica_operation_result& result) { if(result.error != 0) { replica_metrics().write_failure.fetch_add(1); } else if(result.degraded) { replica_metrics().write_degraded.fetch_add(1); } else { replica_metrics().write_success.fetch_add(1); } } inline void record_replica_repair_queued() { replica_metrics().repair_queued.fetch_add(1); } inline void record_replica_repair_completed() { replica_metrics().repair_completed.fetch_add(1); } inline void record_replica_repair_failed() { replica_metrics().repair_failed.fetch_add(1); } inline void record_replica_repair_stale() { replica_metrics().repair_stale.fetch_add(1); } inline replica_operation_result classify_data_write(const int primary_error, const int replica_error, const int replica_count) { std::vector<std::pair<int, int>> copy_errors; copy_errors.emplace_back(0, primary_error); for(int copy = 1; copy <= replica_count; ++copy) { copy_errors.emplace_back(copy, replica_error); } return classify_replica_operation(copy_errors, replica_policy::write_any); } inline replica_target locate_metadata_replica(const Distributor& distributor, const std::string& path, const int copy) { if(copy < 0) { throw std::invalid_argument("replica copy must not be negative"); } return {copy, distributor.locate_file_metadata(path, copy)}; } /** * Return the primary and replica targets for a data chunk in copy order. * Copy zero is always the primary; copies one through replica_count are * replicas. */ inline std::vector<replica_target> locate_replicas(const Distributor& distributor, const std::string& path, const chunkid_t chunk_id, const int replica_count) { if(replica_count < 0) { throw std::invalid_argument("replica count must not be negative"); } std::vector<replica_target> targets; targets.reserve(static_cast<std::size_t>(replica_count) + 1); for(int copy = 0; copy <= replica_count; ++copy) { targets.push_back( {copy, distributor.locate_data(path, chunk_id, copy)}); } return targets; } /** * Return the primary and replica metadata targets in copy order. */ inline std::vector<replica_target> locate_metadata_replicas(const Distributor& distributor, const std::string& path, const int replica_count) { if(replica_count < 0) { throw std::invalid_argument("replica count must not be negative"); } std::vector<replica_target> targets; targets.reserve(static_cast<std::size_t>(replica_count) + 1); for(int copy = 0; copy <= replica_count; ++copy) { targets.push_back({copy, distributor.locate_file_metadata(path, copy)}); } return targets; } /** * Select the first healthy copy in deterministic primary-first order. */ inline std::optional<replica_target> locate_read_replica(const Distributor& distributor, const std::string& path, const chunkid_t chunk_id, const int replica_count, const std::set<host_t>& failed_hosts) { for(const auto& target : locate_replicas(distributor, path, chunk_id, replica_count)) { if(failed_hosts.find(target.host) == failed_hosts.end()) { return target; } } return std::nullopt; } } // namespace gkfs::rpc #endif // GKFS_CLIENT_RPC_REPLICA_HPP No newline at end of file Loading
include/client/preload_context.hpp +92 −0 Changes for include/client/preload_context.hpp: 92 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -46,6 +46,7 @@ #include <future> #include <thread> #include <condition_variable> #include <functional> #include <thallium.hpp> #include <memory> #include <vector> Loading @@ -53,6 +54,8 @@ #include <mutex> #include <unordered_set> #include <config.hpp> #include <client/rpc/repair_queue.hpp> #include <client/rpc/replica.hpp> #include <bitset> Loading Loading @@ -186,6 +189,20 @@ private: std::thread async_write_thread_; bool async_write_stop_{false}; gkfs::rpc::repair_queue repair_queue_; mutable std::mutex repair_queue_mutex_; std::condition_variable repair_queue_cv_; std::vector<std::thread> repair_worker_threads_; bool repair_worker_stop_{false}; std::chrono::steady_clock::time_point repair_worker_deadline_{}; std::function<bool(const gkfs::rpc::repair_task&)> repair_executor_; gkfs::rpc::host_health_table host_health_; 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_; static constexpr uint32_t max_repairs_per_host_ = 1; static constexpr uint32_t repair_worker_count_ = 2; std::string ofi_interface_; std::chrono::milliseconds rpc_timeout_{0}; Loading Loading @@ -361,6 +378,27 @@ public: int get_replicas(); void mark_host_failure(gkfs::rpc::host_t host); void mark_host_success(gkfs::rpc::host_t host); std::set<gkfs::rpc::host_t> unavailable_hosts(); uint64_t next_repair_generation(const std::string& path, gkfs::rpc::repair_kind kind, uint64_t chunk_id = 0); uint64_t current_repair_generation(const std::string& path, gkfs::rpc::repair_kind kind, uint64_t chunk_id = 0) const; bool is_current_repair_generation(const gkfs::rpc::repair_task& task) const; bool protect_fds() const; Loading Loading @@ -464,6 +502,60 @@ public: void wait_async_writes(); void start_repair_worker(); void stop_repair_worker(); void repair_worker(); std::optional<gkfs::rpc::host_t> repair_target_host(const gkfs::rpc::repair_task& task) const; bool enqueue_repair_task(gkfs::rpc::repair_task task); void set_repair_executor( std::function<bool(const gkfs::rpc::repair_task&)> executor); void enqueue_repairs(const std::string& path, const std::vector<int8_t>& failed_copies, int8_t source_copy = 0, uint64_t chunk_id = 0); void enqueue_repairs(const std::string& path, const std::vector<int8_t>& failed_copies, uint64_t first_chunk, uint64_t last_chunk, int8_t source_copy = 0); void enqueue_metadata_create_repairs(const std::string& path, mode_t mode, const std::vector<int8_t>& failed_copies); void enqueue_metadata_create_inline_repairs( const std::string& path, mode_t mode, const std::string& data, const std::vector<int8_t>& failed_copies); void enqueue_metadata_size_repairs(const std::string& path, uint64_t size, int64_t offset, bool append, bool clear_inline, const std::vector<int8_t>& failed_copies); void enqueue_metadata_inline_repairs(const std::string& path, const std::string& data, uint64_t offset, bool append, const std::vector<int8_t>& failed_copies); gkfs::rpc::repair_queue::counters repair_queue_stats() const; }; } // namespace preload Loading
include/client/rpc/forward_data.hpp +7 −1 Changes for include/client/rpc/forward_data.hpp: 7 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -42,9 +42,12 @@ #include <common/common_defs.hpp> #include <algorithm> #include <string> #include <memory> #include <set> #include <common/rpc/distributor.hpp> #include <client/rpc/repair_queue.hpp> namespace gkfs::rpc { // TODO once we have LEAF, remove all the error code returns and throw them as Loading @@ -57,7 +60,10 @@ forward_write(const std::string& path, const void* buf, off64_t offset, std::pair<int, ssize_t> forward_read(const std::string& path, void* buf, off64_t offset, size_t read_size, const int8_t num_copies, std::set<int8_t>& failed); std::set<gkfs::rpc::host_t>& failed); bool repair_data_copy(const gkfs::rpc::repair_task& task); int forward_truncate(const std::string& path, size_t current_size, size_t new_size, Loading
include/client/rpc/forward_metadata.hpp +6 −2 Changes for include/client/rpc/forward_metadata.hpp: 6 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -44,6 +44,7 @@ #include <memory> #include <vector> #include <cstdint> #include <client/rpc/repair_queue.hpp> /* Forward declaration */ namespace gkfs { namespace filemap { Loading @@ -63,6 +64,9 @@ namespace rpc { int forward_create(const std::string& path, mode_t mode, const int copy); int repair_metadata_copy(const gkfs::rpc::repair_task& task); int forward_batch_create(uint64_t host_id, const std::vector<std::string>& paths, const std::vector<uint32_t>& modes); Loading Loading @@ -130,7 +134,7 @@ forward_mk_symlink(const std::string& path, const std::string& target_path); */ std::pair<int, off64_t> forward_write_inline(const std::string& path, const void* buf, off64_t offset, size_t write_size, bool append_flag); size_t write_size, bool append_flag, const int num_copies); /** * @brief Send an RPC request to read a small amount of data directly Loading @@ -144,7 +148,7 @@ forward_write_inline(const std::string& path, const void* buf, off64_t offset, */ std::pair<int, ssize_t> forward_read_inline(const std::string& path, void* buf, off64_t offset, size_t read_size); size_t read_size, const int num_copies); std::tuple<int, std::vector<std::tuple<const std::string, unsigned char, size_t, Loading
include/client/rpc/repair_queue.hpp 0 → 100644 +204 −0 Changes for include/client/rpc/repair_queue.hpp: 204 added lines, 0 removed lines. Original line number Diff line number Diff line /* Copyright 2018-2025, Barcelona Supercomputing Center (BSC), Spain Copyright 2015-2025, Johannes Gutenberg Universitaet Mainz, Germany This file is part of GekkoFS. SPDX-License-Identifier: LGPL-3.0-or-later */ #ifndef GKFS_CLIENT_RPC_REPAIR_QUEUE_HPP #define GKFS_CLIENT_RPC_REPAIR_QUEUE_HPP #include <chrono> #include <cstdint> #include <algorithm> #include <atomic> #include <deque> #include <optional> #include <string> #include <unordered_map> namespace gkfs::rpc { enum class repair_kind : uint8_t { metadata_create, metadata_create_inline, metadata_size, metadata_inline, data_chunk, deferred, }; struct repair_task { std::string path; uint64_t chunk_id{}; int8_t source_copy{}; int8_t target_copy{}; uint64_t generation{}; uint32_t retry_count{}; int last_error{}; std::chrono::steady_clock::time_point next_attempt{}; bool data_chunk{false}; repair_kind kind{repair_kind::deferred}; uint32_t mode{}; uint64_t size{}; int64_t offset{}; bool append{}; bool clear_inline{}; std::string data; }; class repair_queue { public: static constexpr std::size_t default_capacity = 1024; static constexpr uint32_t default_retry_limit = 5; static constexpr uint64_t default_backoff_ms = 100; struct counters { uint64_t queued{}; uint64_t completed{}; uint64_t failed{}; uint64_t stale{}; }; explicit repair_queue(std::size_t capacity = default_capacity) : capacity_(capacity) {} bool enqueue(repair_task task) { const auto key = make_key(task); const auto existing = tasks_.find(key); if(existing != tasks_.end()) { const auto retry_count = std::max(existing->second.retry_count, task.retry_count); existing->second = std::move(task); existing->second.retry_count = retry_count; return true; } if(tasks_.size() >= capacity_) { return false; } if(task.next_attempt == std::chrono::steady_clock::time_point{}) { task.next_attempt = std::chrono::steady_clock::now(); } order_.push_back(key); tasks_.emplace(key, std::move(task)); queued_++; return true; } bool requeue(repair_task task, const int error, const std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now(), const uint32_t retry_limit = default_retry_limit, const uint64_t backoff_ms = default_backoff_ms) { task.last_error = error; if(task.retry_count >= retry_limit) { failed_++; return false; } ++task.retry_count; const auto delay = backoff_ms * (uint64_t{1} << std::min(task.retry_count - 1, 16u)); task.next_attempt = now + std::chrono::milliseconds(delay); return enqueue(std::move(task)); } void defer(repair_task task, const std::chrono::steady_clock::time_point next_attempt) { task.next_attempt = next_attempt; const auto key = make_key(task); if(tasks_.find(key) != tasks_.end()) { return; } order_.push_back(key); tasks_.emplace(key, std::move(task)); } void complete() { completed_++; } void discard_stale() { stale_++; } std::optional<repair_task> pop_ready(const std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now()) { return pop_ready([](const repair_task&) { return true; }, now); } template <typename Predicate> std::optional<repair_task> pop_ready(Predicate predicate, const std::chrono::steady_clock::time_point now = std::chrono::steady_clock::now()) { const auto count = order_.size(); for(std::size_t i = 0; i < count; ++i) { auto key = std::move(order_.front()); order_.pop_front(); auto it = tasks_.find(key); if(it == tasks_.end()) { continue; } if(it->second.next_attempt <= now && predicate(it->second)) { auto task = std::move(it->second); tasks_.erase(it); return task; } order_.push_back(std::move(key)); } return std::nullopt; } bool empty() const { return tasks_.empty(); } std::size_t size() const { return tasks_.size(); } counters stats() const { return {queued_.load(), completed_.load(), failed_.load(), stale_.load()}; } private: using key_type = std::string; static key_type make_key(const repair_task& task) { key_type key = task.path; key.push_back('\0'); key += std::to_string(task.chunk_id); key.push_back('\0'); key += std::to_string(task.source_copy); key.push_back('\0'); key += std::to_string(task.target_copy); key.push_back('\0'); key += std::to_string(static_cast<uint8_t>(task.kind)); return key; } std::size_t capacity_; std::deque<key_type> order_; std::unordered_map<key_type, repair_task> tasks_; std::atomic<uint64_t> queued_{0}; std::atomic<uint64_t> completed_{0}; std::atomic<uint64_t> failed_{0}; std::atomic<uint64_t> stale_{0}; }; } // namespace gkfs::rpc #endif // GKFS_CLIENT_RPC_REPAIR_QUEUE_HPP No newline at end of file
include/client/rpc/replica.hpp 0 → 100644 +331 −0 Changes for include/client/rpc/replica.hpp: 331 added lines, 0 removed lines. Original line number Diff line number Diff line /* Copyright 2018-2025, Barcelona Supercomputing Center (BSC), Spain Copyright 2015-2025, Johannes Gutenberg Universitaet Mainz, Germany This file is part of GekkoFS. SPDX-License-Identifier: LGPL-3.0-or-later */ #ifndef GKFS_CLIENT_RPC_REPLICA_HPP #define GKFS_CLIENT_RPC_REPLICA_HPP #include <common/rpc/distributor.hpp> #include <cstddef> #include <cstdint> #include <cerrno> #include <algorithm> #include <stdexcept> #include <string> #include <set> #include <optional> #include <atomic> #include <chrono> #include <mutex> #include <unordered_map> #include <vector> namespace gkfs::rpc { struct replica_target { int copy{}; host_t host{}; }; enum class host_health_state { healthy, suspect, down, recovering, }; struct host_health_status { host_health_state state{host_health_state::healthy}; uint32_t failure_count{}; std::chrono::steady_clock::time_point last_failure{}; std::chrono::steady_clock::time_point next_probe{}; }; class host_health_table { public: using clock = std::chrono::steady_clock; void mark_failure(const host_t host, const clock::time_point now = clock::now()) { std::lock_guard<std::mutex> lock(mutex_); auto& status = hosts_[host]; ++status.failure_count; status.last_failure = now; status.state = status.failure_count < failure_threshold ? host_health_state::suspect : host_health_state::down; if(status.state == host_health_state::down) { const auto exponent = std::min(status.failure_count - failure_threshold, max_backoff_shift); status.next_probe = now + base_backoff * (uint64_t{1} << exponent); } } void mark_success(const host_t host) { std::lock_guard<std::mutex> lock(mutex_); auto& status = hosts_[host]; status = {}; } std::set<host_t> unavailable(const clock::time_point now = clock::now()) { std::lock_guard<std::mutex> lock(mutex_); std::set<host_t> result; for(auto& [host, status] : hosts_) { if(status.state == host_health_state::down) { if(now < status.next_probe) { result.insert(host); } else { status.state = host_health_state::recovering; } } } return result; } host_health_status status(const host_t host) const { std::lock_guard<std::mutex> lock(mutex_); const auto it = hosts_.find(host); return it == hosts_.end() ? host_health_status{} : it->second; } private: static constexpr uint32_t failure_threshold = 2; static constexpr uint32_t max_backoff_shift = 6; static constexpr auto base_backoff = std::chrono::milliseconds(100); mutable std::mutex mutex_; std::unordered_map<host_t, host_health_status> hosts_; }; enum class replica_policy { read_any, write_any, write_all, strict, primary_required, }; struct replica_operation_result { int error{}; uint8_t successful_copies{}; uint8_t expected_copies{}; std::vector<int8_t> failed_copies; bool degraded{}; bool repair_pending{}; }; struct replica_counters { std::atomic<uint64_t> read_primary{0}; std::atomic<uint64_t> read_fallback{0}; std::atomic<uint64_t> read_failure{0}; std::atomic<uint64_t> write_success{0}; std::atomic<uint64_t> write_degraded{0}; std::atomic<uint64_t> write_failure{0}; std::atomic<uint64_t> repair_queued{0}; std::atomic<uint64_t> repair_completed{0}; std::atomic<uint64_t> repair_failed{0}; std::atomic<uint64_t> repair_stale{0}; }; inline replica_counters& replica_metrics() { static replica_counters counters; return counters; } inline bool mutation_supported_with_replicas(const int replica_count) { return replica_count == 0; } inline void record_replica_read(const int copy) { if(copy == 0) { replica_metrics().read_primary.fetch_add(1); } else { replica_metrics().read_fallback.fetch_add(1); } } inline void record_replica_read_failure() { replica_metrics().read_failure.fetch_add(1); } inline replica_operation_result classify_replica_operation(const std::vector<std::pair<int, int>>& copy_errors, const replica_policy policy) { replica_operation_result result; result.expected_copies = static_cast<uint8_t>(copy_errors.size()); for(const auto& [copy, error] : copy_errors) { if(error == 0) { ++result.successful_copies; } else { result.failed_copies.push_back(static_cast<int8_t>(copy)); if(result.error == 0) { result.error = error; } } } const auto all_succeeded = result.successful_copies == result.expected_copies; const auto any_succeeded = result.successful_copies > 0; switch(policy) { case replica_policy::read_any: case replica_policy::write_any: if(any_succeeded) { result.error = 0; } break; case replica_policy::write_all: case replica_policy::strict: if(all_succeeded) { result.error = 0; } else if(result.error == 0) { result.error = EIO; } break; case replica_policy::primary_required: { const auto primary = std::find_if( copy_errors.begin(), copy_errors.end(), [](const auto& entry) { return entry.first == 0; }); if(primary == copy_errors.end() || primary->second != 0) { result.error = primary == copy_errors.end() ? EIO : primary->second; } else { result.error = 0; } break; } } result.degraded = any_succeeded && !all_succeeded; result.repair_pending = result.degraded; if(!any_succeeded && result.error == 0) { result.error = EIO; } return result; } inline void record_replica_write(const replica_operation_result& result) { if(result.error != 0) { replica_metrics().write_failure.fetch_add(1); } else if(result.degraded) { replica_metrics().write_degraded.fetch_add(1); } else { replica_metrics().write_success.fetch_add(1); } } inline void record_replica_repair_queued() { replica_metrics().repair_queued.fetch_add(1); } inline void record_replica_repair_completed() { replica_metrics().repair_completed.fetch_add(1); } inline void record_replica_repair_failed() { replica_metrics().repair_failed.fetch_add(1); } inline void record_replica_repair_stale() { replica_metrics().repair_stale.fetch_add(1); } inline replica_operation_result classify_data_write(const int primary_error, const int replica_error, const int replica_count) { std::vector<std::pair<int, int>> copy_errors; copy_errors.emplace_back(0, primary_error); for(int copy = 1; copy <= replica_count; ++copy) { copy_errors.emplace_back(copy, replica_error); } return classify_replica_operation(copy_errors, replica_policy::write_any); } inline replica_target locate_metadata_replica(const Distributor& distributor, const std::string& path, const int copy) { if(copy < 0) { throw std::invalid_argument("replica copy must not be negative"); } return {copy, distributor.locate_file_metadata(path, copy)}; } /** * Return the primary and replica targets for a data chunk in copy order. * Copy zero is always the primary; copies one through replica_count are * replicas. */ inline std::vector<replica_target> locate_replicas(const Distributor& distributor, const std::string& path, const chunkid_t chunk_id, const int replica_count) { if(replica_count < 0) { throw std::invalid_argument("replica count must not be negative"); } std::vector<replica_target> targets; targets.reserve(static_cast<std::size_t>(replica_count) + 1); for(int copy = 0; copy <= replica_count; ++copy) { targets.push_back( {copy, distributor.locate_data(path, chunk_id, copy)}); } return targets; } /** * Return the primary and replica metadata targets in copy order. */ inline std::vector<replica_target> locate_metadata_replicas(const Distributor& distributor, const std::string& path, const int replica_count) { if(replica_count < 0) { throw std::invalid_argument("replica count must not be negative"); } std::vector<replica_target> targets; targets.reserve(static_cast<std::size_t>(replica_count) + 1); for(int copy = 0; copy <= replica_count; ++copy) { targets.push_back({copy, distributor.locate_file_metadata(path, copy)}); } return targets; } /** * Select the first healthy copy in deterministic primary-first order. */ inline std::optional<replica_target> locate_read_replica(const Distributor& distributor, const std::string& path, const chunkid_t chunk_id, const int replica_count, const std::set<host_t>& failed_hosts) { for(const auto& target : locate_replicas(distributor, path, chunk_id, replica_count)) { if(failed_hosts.find(target.host) == failed_hosts.end()) { return target; } } return std::nullopt; } } // namespace gkfs::rpc #endif // GKFS_CLIENT_RPC_REPLICA_HPP No newline at end of file