Loading .gitlab-ci.yml +5 −0 Changes for .gitlab-ci.yml: 5 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -114,8 +114,13 @@ gkfs:integration-1: - export PATH=${PATH}:/usr/local/bin - mkdir -p ${BUILD_PATH}/tests/run - cd ${BUILD_PATH}/tests/integration - LIBGKFS_RPC_TIMEOUT_MS=5000 ${PYTEST} -v -n $(nproc) --dist=loadscope --timeout=900 --timeout-method=thread ${INTEGRATION_TESTS_BIN_PATH}/data/test_replication.py --basetemp=${BUILD_PATH}/tests/run/replication --junit-xml=report-1-replication.xml - ${PYTEST} -v -n $(nproc) --dist=loadscope --timeout=900 --timeout-method=thread ${INTEGRATION_TESTS_BIN_PATH}/data --ignore=${INTEGRATION_TESTS_BIN_PATH}/data/test_replication.py ${INTEGRATION_TESTS_BIN_PATH}/fuse ${INTEGRATION_TESTS_BIN_PATH}/shell ${INTEGRATION_TESTS_BIN_PATH}/compatibility Loading include/client/rpc/repair_journal_lock.hpp +18 −2 Changes for include/client/rpc/repair_journal_lock.hpp: 18 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -15,6 +15,7 @@ #include <fcntl.h> #include <unistd.h> #include <algorithm> #include <cerrno> #include <chrono> #include <cstdint> Loading @@ -26,6 +27,7 @@ #include <stdexcept> #include <string> #include <sys/types.h> #include <vector> namespace gkfs::rpc { Loading @@ -45,12 +47,26 @@ repair_journal_topology_identity(const std::string& hostfile_path) { constexpr uint64_t fnv_offset = 14695981039346656037ULL; constexpr uint64_t fnv_prime = 1099511628211ULL; std::vector<std::string> logical_hosts; std::string line; while(std::getline(input, line)) { std::istringstream fields(line); std::string host; if(fields >> host) { logical_hosts.push_back(std::move(host)); } } std::sort(logical_hosts.begin(), logical_hosts.end()); uint64_t hash = fnv_offset; char byte{}; while(input.get(byte)) { for(const auto& host : logical_hosts) { for(const auto byte : host) { hash ^= static_cast<uint64_t>(static_cast<unsigned char>(byte)); hash *= fnv_prime; } hash ^= static_cast<uint64_t>('\n'); hash *= fnv_prime; } std::ostringstream result; result << "fnv1a-64:" << std::hex << std::setw(16) << std::setfill('0') Loading scripts/dev/coverage.py +4 −4 Changes for scripts/dev/coverage.py: 4 added lines, 4 removed lines. Original line number Diff line number Diff line Loading @@ -175,7 +175,7 @@ def capture(args): pipeline.append([ "lcov", "--rc", "branch_coverage=1", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--capture", "--quiet" if args.verbosity > 3 else "", "--initial" if args.initial else "", Loading @@ -187,7 +187,7 @@ def capture(args): pipeline.append([ "lcov", "--rc", "branch_coverage=1", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--quiet" if args.verbosity > 3 else "", "--remove=$INPUT", *args.exclusion_patterns, Loading Loading @@ -226,7 +226,7 @@ def merge(args): pipeline.append([ "lcov", "--rc", "branch_coverage=1", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--quiet" if args.verbosity > 3 else "", *(f"--add-tracefile={t}" for t in tracefiles), "--output-file=$OUTPUT"]) Loading Loading @@ -262,7 +262,7 @@ def html_report(args): "--quiet" if args.verbosity > 3 else "", "--show-details", "--filter", "missing", "--ignore-errors", "inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors", "inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--legend", "--demangle-cpp", "--frames", Loading src/client/gkfs_data.cpp +5 −4 Changes for src/client/gkfs_data.cpp: 5 added lines, 4 removed lines. Original line number Diff line number Diff line Loading @@ -610,13 +610,14 @@ gkfs_do_read(const gkfs::filemap::OpenFile& file, char* buf, size_t count, } else { std::set<gkfs::rpc::host_t> failed; // set with failed targets. if(CTX->get_replicas() != 0) { ret = gkfs::rpc::forward_read(file.path(), buf, offset, count, CTX->get_replicas(), failed); while(ret.first == EIO) { const auto max_attempts = CTX->get_replicas() + 1; for(int attempt = 0; attempt < max_attempts; ++attempt) { ret = gkfs::rpc::forward_read(file.path(), buf, offset, count, CTX->get_replicas(), failed); if(ret.first != EIO) { break; } LOG(WARNING, "gkfs::rpc::forward_read() failed with ret '{}'", ret.first); Loading src/client/rpc/forward_data.cpp +20 −24 Changes for src/client/rpc/forward_data.cpp: 20 added lines, 24 removed lines. Original line number Diff line number Diff line Loading @@ -184,7 +184,7 @@ forward_write(const string& path, const void* buf, const off64_t offset, std::set<uint64_t> chnk_end_target{}; std::unordered_map<uint64_t, std::vector<uint8_t>> write_ops_vect; const auto unavailable_hosts = num_copies > 0 const auto unavailable_hosts = CTX->get_replicas() > 0 ? CTX->unavailable_hosts() : std::set<gkfs::rpc::host_t>{}; Loading Loading @@ -289,11 +289,11 @@ forward_write(const string& path, const void* buf, const off64_t offset, in.bulk_handle = bulk_handle; try { if(num_copies > 0 && CTX->rpc_timeout().count() == 0) { if((num_copies > 0 || CTX->get_replicas() > 0) && CTX->rpc_timeout().count() == 0) { waiters.push_back( write_rpc.on(CTX->hosts().at(target)) .timed_async(std::chrono::milliseconds(100), in)); .timed_async(std::chrono::seconds(5), in)); } else { waiters.push_back(forward_async_with_timeout( write_rpc, CTX->hosts().at(target), in)); Loading Loading @@ -351,7 +351,9 @@ forward_write(const string& path, const void* buf, const off64_t offset, } } if(submission_failed && num_copies == 0) { if(waiters.empty() && num_copies == 0) { err = EHOSTUNREACH; } else if(submission_failed && num_copies == 0) { err = EIO; } Loading Loading @@ -483,26 +485,23 @@ forward_read(const string& path, void* buf, const off64_t offset, std::vector<thallium::async_response> waiters; waiters.reserve(targets.size()); std::vector<uint64_t> waiter_targets; // track targets for error reporting std::vector<uint64_t> waiter_targets; waiter_targets.reserve(targets.size()); auto read_rpc = CTX->rpc_engine()->define(gkfs::rpc::tag::read); // Issue non-blocking RPC requests and wait for the result later for(std::size_t i = 0; i < targets.size(); ++i) { auto target = targets[i]; auto err = 0; ssize_t out_size = 0; // total chunk_size for target for(std::size_t i = 0; i < targets.size(); ++i) { const auto target = targets[i]; auto total_chunk_size = target_chnks[target].size() * gkfs::config::rpc::chunksize; // receiver of first chunk must subtract the offset from first chunk if(target == chnk_start_target) { total_chunk_size -= block_overrun(offset, gkfs::config::rpc::chunksize); } // receiver of last chunk must subtract if(target == chnk_end_target && !is_aligned(offset + read_size, gkfs::config::rpc::chunksize)) { total_chunk_size -= block_underrun(offset + read_size, Loading @@ -528,13 +527,18 @@ forward_read(const string& path, void* buf, const off64_t offset, in.bulk_handle = bulk_handle; try { if(num_copies > 0 && CTX->rpc_timeout().count() == 0) { waiters.push_back( read_rpc.on(CTX->hosts().at(target)) .timed_async(std::chrono::seconds(5), in)); } else { waiters.push_back(forward_async_with_timeout( read_rpc, CTX->hosts().at(target), in)); } waiter_targets.push_back(target); } catch(const std::exception& ex) { LOG(ERROR, "Failed to send RPC to host {}: {}", target, ex.what()); // Treat submission failures like failed responses so the caller // can retry the affected read against another replica. LOG(ERROR, "RPC read failed for host {}: {}", target, ex.what()); err = EIO; failed.insert(target); CTX->mark_host_failure(target); } Loading @@ -545,10 +549,6 @@ forward_read(const string& path, void* buf, const off64_t offset, in.offset); } // Wait for RPC responses and then get response auto err = 0; ssize_t out_size = 0; for(std::size_t i = 0; i < waiters.size(); ++i) { try { gkfs::rpc::rpc_data_out_t out = waiters[i].wait(); Loading @@ -569,10 +569,6 @@ forward_read(const string& path, void* buf, const off64_t offset, } } if(waiters.size() != targets.size()) { return make_pair(EIO, 0); } /* * Typically file systems return the size even if only a part of it was * read. In our case, we do not keep track which daemon fully read its Loading Loading
.gitlab-ci.yml +5 −0 Changes for .gitlab-ci.yml: 5 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -114,8 +114,13 @@ gkfs:integration-1: - export PATH=${PATH}:/usr/local/bin - mkdir -p ${BUILD_PATH}/tests/run - cd ${BUILD_PATH}/tests/integration - LIBGKFS_RPC_TIMEOUT_MS=5000 ${PYTEST} -v -n $(nproc) --dist=loadscope --timeout=900 --timeout-method=thread ${INTEGRATION_TESTS_BIN_PATH}/data/test_replication.py --basetemp=${BUILD_PATH}/tests/run/replication --junit-xml=report-1-replication.xml - ${PYTEST} -v -n $(nproc) --dist=loadscope --timeout=900 --timeout-method=thread ${INTEGRATION_TESTS_BIN_PATH}/data --ignore=${INTEGRATION_TESTS_BIN_PATH}/data/test_replication.py ${INTEGRATION_TESTS_BIN_PATH}/fuse ${INTEGRATION_TESTS_BIN_PATH}/shell ${INTEGRATION_TESTS_BIN_PATH}/compatibility Loading
include/client/rpc/repair_journal_lock.hpp +18 −2 Changes for include/client/rpc/repair_journal_lock.hpp: 18 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -15,6 +15,7 @@ #include <fcntl.h> #include <unistd.h> #include <algorithm> #include <cerrno> #include <chrono> #include <cstdint> Loading @@ -26,6 +27,7 @@ #include <stdexcept> #include <string> #include <sys/types.h> #include <vector> namespace gkfs::rpc { Loading @@ -45,12 +47,26 @@ repair_journal_topology_identity(const std::string& hostfile_path) { constexpr uint64_t fnv_offset = 14695981039346656037ULL; constexpr uint64_t fnv_prime = 1099511628211ULL; std::vector<std::string> logical_hosts; std::string line; while(std::getline(input, line)) { std::istringstream fields(line); std::string host; if(fields >> host) { logical_hosts.push_back(std::move(host)); } } std::sort(logical_hosts.begin(), logical_hosts.end()); uint64_t hash = fnv_offset; char byte{}; while(input.get(byte)) { for(const auto& host : logical_hosts) { for(const auto byte : host) { hash ^= static_cast<uint64_t>(static_cast<unsigned char>(byte)); hash *= fnv_prime; } hash ^= static_cast<uint64_t>('\n'); hash *= fnv_prime; } std::ostringstream result; result << "fnv1a-64:" << std::hex << std::setw(16) << std::setfill('0') Loading
scripts/dev/coverage.py +4 −4 Changes for scripts/dev/coverage.py: 4 added lines, 4 removed lines. Original line number Diff line number Diff line Loading @@ -175,7 +175,7 @@ def capture(args): pipeline.append([ "lcov", "--rc", "branch_coverage=1", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--capture", "--quiet" if args.verbosity > 3 else "", "--initial" if args.initial else "", Loading @@ -187,7 +187,7 @@ def capture(args): pipeline.append([ "lcov", "--rc", "branch_coverage=1", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--quiet" if args.verbosity > 3 else "", "--remove=$INPUT", *args.exclusion_patterns, Loading Loading @@ -226,7 +226,7 @@ def merge(args): pipeline.append([ "lcov", "--rc", "branch_coverage=1", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors","inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--quiet" if args.verbosity > 3 else "", *(f"--add-tracefile={t}" for t in tracefiles), "--output-file=$OUTPUT"]) Loading Loading @@ -262,7 +262,7 @@ def html_report(args): "--quiet" if args.verbosity > 3 else "", "--show-details", "--filter", "missing", "--ignore-errors", "inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov", "--ignore-errors", "inconsistent,inconsistent,corrupt,corrupt,mismatch,mismatch,negative,negative,version,version,empty,empty,unused,unused,source,source,missing,missing,count,count,gcov,gcov,range,range", "--legend", "--demangle-cpp", "--frames", Loading
src/client/gkfs_data.cpp +5 −4 Changes for src/client/gkfs_data.cpp: 5 added lines, 4 removed lines. Original line number Diff line number Diff line Loading @@ -610,13 +610,14 @@ gkfs_do_read(const gkfs::filemap::OpenFile& file, char* buf, size_t count, } else { std::set<gkfs::rpc::host_t> failed; // set with failed targets. if(CTX->get_replicas() != 0) { ret = gkfs::rpc::forward_read(file.path(), buf, offset, count, CTX->get_replicas(), failed); while(ret.first == EIO) { const auto max_attempts = CTX->get_replicas() + 1; for(int attempt = 0; attempt < max_attempts; ++attempt) { ret = gkfs::rpc::forward_read(file.path(), buf, offset, count, CTX->get_replicas(), failed); if(ret.first != EIO) { break; } LOG(WARNING, "gkfs::rpc::forward_read() failed with ret '{}'", ret.first); Loading
src/client/rpc/forward_data.cpp +20 −24 Changes for src/client/rpc/forward_data.cpp: 20 added lines, 24 removed lines. Original line number Diff line number Diff line Loading @@ -184,7 +184,7 @@ forward_write(const string& path, const void* buf, const off64_t offset, std::set<uint64_t> chnk_end_target{}; std::unordered_map<uint64_t, std::vector<uint8_t>> write_ops_vect; const auto unavailable_hosts = num_copies > 0 const auto unavailable_hosts = CTX->get_replicas() > 0 ? CTX->unavailable_hosts() : std::set<gkfs::rpc::host_t>{}; Loading Loading @@ -289,11 +289,11 @@ forward_write(const string& path, const void* buf, const off64_t offset, in.bulk_handle = bulk_handle; try { if(num_copies > 0 && CTX->rpc_timeout().count() == 0) { if((num_copies > 0 || CTX->get_replicas() > 0) && CTX->rpc_timeout().count() == 0) { waiters.push_back( write_rpc.on(CTX->hosts().at(target)) .timed_async(std::chrono::milliseconds(100), in)); .timed_async(std::chrono::seconds(5), in)); } else { waiters.push_back(forward_async_with_timeout( write_rpc, CTX->hosts().at(target), in)); Loading Loading @@ -351,7 +351,9 @@ forward_write(const string& path, const void* buf, const off64_t offset, } } if(submission_failed && num_copies == 0) { if(waiters.empty() && num_copies == 0) { err = EHOSTUNREACH; } else if(submission_failed && num_copies == 0) { err = EIO; } Loading Loading @@ -483,26 +485,23 @@ forward_read(const string& path, void* buf, const off64_t offset, std::vector<thallium::async_response> waiters; waiters.reserve(targets.size()); std::vector<uint64_t> waiter_targets; // track targets for error reporting std::vector<uint64_t> waiter_targets; waiter_targets.reserve(targets.size()); auto read_rpc = CTX->rpc_engine()->define(gkfs::rpc::tag::read); // Issue non-blocking RPC requests and wait for the result later for(std::size_t i = 0; i < targets.size(); ++i) { auto target = targets[i]; auto err = 0; ssize_t out_size = 0; // total chunk_size for target for(std::size_t i = 0; i < targets.size(); ++i) { const auto target = targets[i]; auto total_chunk_size = target_chnks[target].size() * gkfs::config::rpc::chunksize; // receiver of first chunk must subtract the offset from first chunk if(target == chnk_start_target) { total_chunk_size -= block_overrun(offset, gkfs::config::rpc::chunksize); } // receiver of last chunk must subtract if(target == chnk_end_target && !is_aligned(offset + read_size, gkfs::config::rpc::chunksize)) { total_chunk_size -= block_underrun(offset + read_size, Loading @@ -528,13 +527,18 @@ forward_read(const string& path, void* buf, const off64_t offset, in.bulk_handle = bulk_handle; try { if(num_copies > 0 && CTX->rpc_timeout().count() == 0) { waiters.push_back( read_rpc.on(CTX->hosts().at(target)) .timed_async(std::chrono::seconds(5), in)); } else { waiters.push_back(forward_async_with_timeout( read_rpc, CTX->hosts().at(target), in)); } waiter_targets.push_back(target); } catch(const std::exception& ex) { LOG(ERROR, "Failed to send RPC to host {}: {}", target, ex.what()); // Treat submission failures like failed responses so the caller // can retry the affected read against another replica. LOG(ERROR, "RPC read failed for host {}: {}", target, ex.what()); err = EIO; failed.insert(target); CTX->mark_host_failure(target); } Loading @@ -545,10 +549,6 @@ forward_read(const string& path, void* buf, const off64_t offset, in.offset); } // Wait for RPC responses and then get response auto err = 0; ssize_t out_size = 0; for(std::size_t i = 0; i < waiters.size(); ++i) { try { gkfs::rpc::rpc_data_out_t out = waiters[i].wait(); Loading @@ -569,10 +569,6 @@ forward_read(const string& path, void* buf, const off64_t offset, } } if(waiters.size() != targets.size()) { return make_pair(EIO, 0); } /* * Typically file systems return the size even if only a part of it was * read. In our case, we do not keep track which daemon fully read its Loading