Commit 83b2f8c8 authored by Ramon Nou's avatar Ramon Nou
Browse files

Refactor replica helpers and client repair flow

parent 3ce312d6
Loading
Loading
Loading
Loading
Loading
+22 −4
Changes for include/client/preload_context.hpp: 22 added lines, 4 removed lines.
Original line number Diff line number Diff line
@@ -57,6 +57,7 @@
#include <client/rpc/repair_queue.hpp>
#include <client/rpc/repair_journal_lock.hpp>
#include <client/rpc/replica.hpp>
#include <client/rpc/replica_context.hpp>

#include <bitset>

@@ -159,7 +160,7 @@ private:
    bool internal_fds_must_relocate_;
    std::bitset<MAX_USER_FDS> protected_fds_;
    std::string hostname;
    int replicas_;
    gkfs::rpc::ReplicaContext replica_context_;

    bool protect_fds_{false};
    bool protect_files_generator_{false};
@@ -197,9 +198,6 @@ private:
    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_;
    std::vector<gkfs::rpc::repair_task> in_flight_repairs_;
    std::string repair_status_path_;
@@ -396,6 +394,21 @@ public:
    next_repair_generation(const std::string& path, gkfs::rpc::repair_kind kind,
                           uint64_t chunk_id = 0);

    void
    invalidate_repairs_for_write(const std::string& path, off64_t offset,
                                 size_t size);

    void
    invalidate_repairs_for_create(const std::string& path,
                                  bool includes_inline_data);

    void
    invalidate_repairs_for_remove(const std::string& path, uint64_t size);

    void
    invalidate_repairs_for_truncate(const std::string& path, off64_t old_size,
                                    off64_t new_size);

    uint64_t
    current_repair_generation(const std::string& path,
                              gkfs::rpc::repair_kind kind,
@@ -505,6 +518,11 @@ public:
    enqueue_async_write(const std::string& path, const void* buf,
                        off64_t offset, size_t count, int8_t num_copies);

    void
    enqueue_async_replica_write(const std::string& path, const void* buf,
                                off64_t offset, size_t count,
                                int8_t replica_count);

    void
    wait_async_writes();

+4 −0
Changes for include/client/rpc/forward_data.hpp: 4 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -62,6 +62,10 @@ forward_read(const std::string& path, void* buf, off64_t offset,
             size_t read_size, const int8_t num_copies,
             std::set<gkfs::rpc::host_t>& failed);

std::pair<int, ssize_t>
forward_read_with_replicas(const std::string& path, void* buf, off64_t offset,
                           size_t read_size, int8_t replica_count);

bool
repair_data_copy(const gkfs::rpc::repair_task& task);

+95 −0
Changes for include/client/rpc/replica_context.hpp: 95 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

  SPDX-License-Identifier: LGPL-3.0-or-later
*/

#ifndef GKFS_CLIENT_RPC_REPLICA_CONTEXT_HPP
#define GKFS_CLIENT_RPC_REPLICA_CONTEXT_HPP

#include <client/rpc/replica.hpp>

#include <mutex>
#include <string>
#include <unordered_map>

namespace gkfs::rpc {

class ReplicaContext {
public:
    void
    count(const int value) noexcept {
        replica_count_ = value;
    }

    int
    count() const noexcept {
        return replica_count_;
    }

    bool
    enabled() const noexcept {
        return replica_count_ > 0;
    }

    void
    mark_failure(const host_t host) {
        host_health_.mark_failure(host);
    }

    void
    mark_success(const host_t host) {
        host_health_.mark_success(host);
    }

    std::set<host_t>
    unavailable_hosts() {
        return host_health_.unavailable();
    }

    uint64_t
    next_generation(const std::string& path, const repair_kind kind,
                    const uint64_t chunk_id = 0) {
        std::lock_guard<std::mutex> lock(generation_mutex_);
        auto key = generation_key(path, kind, chunk_id);
        return ++generations_[key];
    }

    uint64_t
    current_generation(const std::string& path, const repair_kind kind,
                       const uint64_t chunk_id = 0) const {
        std::lock_guard<std::mutex> lock(generation_mutex_);
        const auto it = generations_.find(generation_key(path, kind, chunk_id));
        return it == generations_.end() ? 0 : it->second;
    }

    bool
    is_current(const repair_task& task) const {
        std::lock_guard<std::mutex> lock(generation_mutex_);
        const auto it = generations_.find(
                generation_key(task.path, task.kind, task.chunk_id));
        return it == generations_.end() || it->second == task.generation;
    }

private:
    static std::string
    generation_key(const std::string& path, const repair_kind kind,
                   const uint64_t chunk_id) {
        std::string key = path;
        key.push_back('\0');
        key += std::to_string(static_cast<uint8_t>(kind));
        key.push_back('\0');
        key += std::to_string(chunk_id);
        return key;
    }

    int replica_count_{0};
    host_health_table host_health_;
    mutable std::mutex generation_mutex_;
    std::unordered_map<std::string, uint64_t> generations_;
};

} // namespace gkfs::rpc

#endif
 No newline at end of file
+27 −0
Changes for include/client/rpc/replica_policy.hpp: 27 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

  SPDX-License-Identifier: LGPL-3.0-or-later
*/

#ifndef GKFS_CLIENT_RPC_REPLICA_POLICY_HPP
#define GKFS_CLIENT_RPC_REPLICA_POLICY_HPP

#include <client/rpc/replica.hpp>

namespace gkfs::rpc {

inline bool
replica_write_degraded(const replica_operation_result& result) noexcept {
    return result.degraded && result.successful_copies > 0;
}

inline bool
replica_write_usable(const replica_operation_result& result) noexcept {
    return result.successful_copies > 0 && result.error == 0;
}

} // namespace gkfs::rpc

#endif
 No newline at end of file
+34 −0
Changes for include/client/rpc/replica_repair.hpp: 34 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

  SPDX-License-Identifier: LGPL-3.0-or-later
*/

#ifndef GKFS_CLIENT_RPC_REPLICA_REPAIR_HPP
#define GKFS_CLIENT_RPC_REPLICA_REPAIR_HPP

#include <client/rpc/replica.hpp>

namespace gkfs::rpc {

inline std::vector<uint64_t>
chunk_ids_for_range(const off64_t offset, const size_t size,
                    const size_t chunk_size) {
    std::vector<uint64_t> ids;
    if(size == 0) {
        return ids;
    }
    const auto first = offset / static_cast<off64_t>(chunk_size);
    const auto last = (offset + static_cast<off64_t>(size) - 1) /
                      static_cast<off64_t>(chunk_size);
    ids.reserve(static_cast<std::size_t>(last - first + 1));
    for(auto id = first; id <= last; ++id) {
        ids.push_back(static_cast<uint64_t>(id));
    }
    return ids;
}

} // namespace gkfs::rpc

#endif
 No newline at end of file
Loading