Loading examples/gfind/pfind.sh +1 −1 Viewed Changes for examples/gfind/pfind.sh: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -34,7 +34,7 @@ GKFS_FIND_PROCESS=10 GKFS_SERVERS=$SLURM_JOB_NUM_NODES GKFS_FIND=sfind srun -N $NUM_NODES -n $GKFS_FIND_PROCESS --overlap --overcommit --mem=0 --oversubscribe --export=ALL,LD_PRELOAD=${GKFS} $GKFS_FIND $@ -M $GKFS_MNT -S $GKFS_SERVERS srun -N $NUM_NODES -n $GKFS_FIND_PROCESS --overlap --overcommit --mem=0 --oversubscribe --export=ALL,LD_PRELOAD=${GKFS} $GKFS_FIND $@ -M $GKFS_MNT -S $GKFS_SERVERS --server-side -C #!/bin/bash # scripts/aggregate_sfind_results.sh Loading include/client/rpc/forward_metadata_proxy.hpp +2 −1 Viewed Changes for include/client/rpc/forward_metadata_proxy.hpp: 2 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -54,7 +54,8 @@ forward_get_metadentry_size_proxy(const std::string& path); std::pair<int, std::unique_ptr<std::vector<std::tuple< const std::string, unsigned char, size_t, time_t>>>> forward_get_dirents_single_proxy_v2(const std::string& path, int server); forward_get_dirents_single_proxy_v2(const std::string& path, int server, const std::string& start_key = ""); } // namespace gkfs::rpc Loading include/client/rpc/utils.hpp +7 −3 Viewed Changes for include/client/rpc/utils.hpp: 7 added lines, 3 removed lines. Original line number Diff line number Diff line Loading @@ -105,6 +105,7 @@ decompress_and_parse_entries(const OutputOrErr& out, std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> entries; entries.reserve(static_cast<size_t>((end - p) / 32) + 1); while(p < end) { if(p + sizeof(unsigned char) > end) { Loading @@ -112,7 +113,8 @@ decompress_and_parse_entries(const OutputOrErr& out, __func__); break; } unsigned char type = *reinterpret_cast<const unsigned char*>(p); unsigned char type = 0; std::memcpy(&type, p, sizeof(unsigned char)); p += sizeof(unsigned char); if(p + sizeof(size_t) > end) { Loading @@ -120,7 +122,8 @@ decompress_and_parse_entries(const OutputOrErr& out, __func__); break; } size_t file_size = *reinterpret_cast<const size_t*>(p); size_t file_size = 0; std::memcpy(&file_size, p, sizeof(size_t)); p += sizeof(size_t); if(p + sizeof(time_t) > end) { Loading @@ -128,7 +131,8 @@ decompress_and_parse_entries(const OutputOrErr& out, __func__); break; } time_t ctime = *reinterpret_cast<const time_t*>(p); time_t ctime = 0; std::memcpy(&ctime, p, sizeof(time_t)); p += sizeof(time_t); // Name is null-terminated. We must check we don't go past 'end'. Loading src/client/rpc/forward_metadata.cpp +85 −137 Viewed Changes for src/client/rpc/forward_metadata.cpp: 85 added lines, 137 removed lines. Original line number Diff line number Diff line Loading @@ -238,24 +238,18 @@ forward_mk_symlink(const std::string& path, const std::string& target_path) { } // function matches the standard rpc_get_dirents_out_t inline std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> decompress_and_parse_entries_standard( const gkfs::rpc::rpc_get_dirents_out_t& out, const void* compressed_buffer) { // Duplicated: 'out' type differs (proxy vs daemon) template <typename OutputType> std::pair<const char*, std::size_t> decompress_dirents_payload(const OutputType& out, const void* compressed_buffer, std::vector<char>& decompressed_data) { if(out.err != 0) { throw std::runtime_error("Server returned an error: " + std::to_string(out.err)); } if(out.dirents_size == 0) { return {}; // No entries, return empty vector return {nullptr, 0}; } const char* p = nullptr; const char* end = nullptr; std::vector<char> decompressed_data; if(gkfs::config::rpc::use_dirents_compression) { const unsigned long long uncompressed_size = ZSTD_getFrameContentSize(compressed_buffer, out.dirents_size); Loading Loading @@ -283,100 +277,29 @@ decompress_and_parse_entries_standard( throw std::runtime_error("Decompression size mismatch."); } p = decompressed_data.data(); end = p + uncompressed_size; return {decompressed_data.data(), static_cast<std::size_t>(uncompressed_size)}; } else { p = static_cast<const char*>(compressed_buffer); end = p + out.dirents_size; } std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> entries; entries.reserve(out.dirents_size); // Approx while(p < end) { if(p + sizeof(unsigned char) > end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing type", __func__); break; } unsigned char type = *reinterpret_cast<const unsigned char*>(p); p += sizeof(unsigned char); // Name is null-terminated. We must check we don't go past 'end'. size_t name_len = strnlen(p, end - p); if(p + name_len >= end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing name", __func__); break; } std::string name(p, name_len); p += name_len + 1; if(!name.empty()) { entries.emplace_back(name, type, 0, 0); return {static_cast<const char*>(compressed_buffer), out.dirents_size}; } } return entries; } // Helper for filtered entries which include size and ctime inline std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> decompress_and_parse_entries_filtered( const gkfs::rpc::rpc_get_dirents_filtered_out_t& out, decompress_and_parse_entries_standard( const gkfs::rpc::rpc_get_dirents_out_t& out, const void* compressed_buffer) { if(out.err != 0) { throw std::runtime_error("Server returned an error: " + std::to_string(out.err)); } if(out.dirents_size == 0) { return {}; // No entries, return empty vector } const char* p = nullptr; const char* end = nullptr; std::vector<char> decompressed_data; if(gkfs::config::rpc::use_dirents_compression) { const unsigned long long uncompressed_size = ZSTD_getFrameContentSize(compressed_buffer, out.dirents_size); if(uncompressed_size == ZSTD_CONTENTSIZE_ERROR) { throw std::runtime_error( "Received data is not a valid Zstd frame."); } if(uncompressed_size == ZSTD_CONTENTSIZE_UNKNOWN) { throw std::runtime_error( "Zstd frame content size is unknown and was not written in the frame."); } decompressed_data.resize(uncompressed_size); const size_t result_size = ZSTD_decompress(decompressed_data.data(), uncompressed_size, compressed_buffer, out.dirents_size); if(ZSTD_isError(result_size)) { throw std::runtime_error( "Zstd decompression failed: " + std::string(ZSTD_getErrorName(result_size))); } if(result_size != uncompressed_size) { throw std::runtime_error("Decompression size mismatch."); auto [payload, payload_size] = decompress_dirents_payload( out, compressed_buffer, decompressed_data); if(payload_size == 0) { return {}; } const char* p = payload; const char* end = payload + payload_size; p = decompressed_data.data(); end = p + uncompressed_size; } else { p = static_cast<const char*>(compressed_buffer); end = p + out.dirents_size; } std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> entries; // We don't know exact count, but can optimize reserve? // entries.reserve(out.dirents_size / 32); entries.reserve((payload_size / 32U) + 1U); while(p < end) { if(p + sizeof(unsigned char) > end) { Loading @@ -387,22 +310,6 @@ decompress_and_parse_entries_filtered( unsigned char type = *reinterpret_cast<const unsigned char*>(p); p += sizeof(unsigned char); if(p + sizeof(size_t) > end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing size", __func__); break; } size_t size = *reinterpret_cast<const size_t*>(p); p += sizeof(size_t); if(p + sizeof(time_t) > end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing ctime", __func__); break; } time_t ctime = *reinterpret_cast<const time_t*>(p); p += sizeof(time_t); // Name is null-terminated. We must check we don't go past 'end'. size_t name_len = strnlen(p, end - p); if(p + name_len >= end) { Loading @@ -415,7 +322,7 @@ decompress_and_parse_entries_filtered( p += name_len + 1; if(!name.empty()) { entries.emplace_back(name, type, size, ctime); entries.emplace_back(name, type, 0, 0); } } Loading Loading @@ -682,12 +589,23 @@ forward_update_metadentry_size(const string& path, const size_t size, int err = 0; off64_t ret_offset = 0; bool valid = false; bool primary_valid = false; off64_t primary_offset = 0; struct pending_update { int copy; thallium::async_response response; }; std::vector<pending_update> waiters; waiters.reserve(num_copies + 1); auto update_rpc = CTX->rpc_engine()->define(gkfs::rpc::tag::update_metadentry_size); for(auto copy = 0; copy < num_copies + 1; copy++) { auto endp = CTX->hosts().at( CTX->distributor()->locate_file_metadata(path, copy)); try { gkfs::rpc::rpc_update_metadentry_size_in_t in; in.path = path; in.size = size; Loading @@ -695,27 +613,45 @@ forward_update_metadentry_size(const string& path, const size_t size, in.append = append_flag; in.clear_inline = clear_inline_flag; auto update_rpc = CTX->rpc_engine()->define( gkfs::rpc::tag::update_metadentry_size); gkfs::rpc::rpc_update_metadentry_size_out_t out = update_rpc.on(endp)(in); try { waiters.push_back({copy, update_rpc.on(endp).async(in)}); } catch(const std::exception& ex) { LOG(ERROR, "{}() posting rpc for path '{}' replica {} failed: {}", __func__, path, copy, ex.what()); err = EBUSY; } } for(auto& waiter : waiters) { try { gkfs::rpc::rpc_update_metadentry_size_out_t out = waiter.response.wait(); if(out.err == 0) { valid = true; if(waiter.copy == 0) { primary_valid = true; primary_offset = out.ret_offset; } if(!primary_valid) { ret_offset = out.ret_offset; } } else { err = out.err; } } catch(const std::exception& ex) { LOG(ERROR, "{}() getting rpc output for path '{}' replica {} failed: {}", __func__, path, copy, ex.what()); // Continue to other replicas __func__, path, waiter.copy, ex.what()); err = EBUSY; } } if(!valid) return make_pair(err, 0); return make_pair(err ? err : EIO, 0); if(primary_valid) { ret_offset = primary_offset; } return make_pair(err, ret_offset); } Loading @@ -738,13 +674,9 @@ forward_get_dirents_single(const string& path, int server, vector<tuple<const std::string, unsigned char, size_t, time_t>>>(); int err = 0; string start_key = start_key_arg; // Chunking loop while(true) { size_t buffer_size = CTX->dirents_buff_size(); auto large_buffer = std::unique_ptr<char[]>(new char[buffer_size]); const int max_retries = 2; bool chunk_success = false; auto endp = CTX->hosts().at(targets[server]); auto get_dirents_rpc = CTX->rpc_engine()->define(gkfs::rpc::tag::get_dirents_extended); // Ensure path ends with / for get_dirents prefix check std::string root_path = path; Loading @@ -752,29 +684,39 @@ forward_get_dirents_single(const string& path, int server, root_path += '/'; } thread_local std::vector<char> large_buffer; size_t buffer_size = std::max<std::size_t>(std::size_t{1}, CTX->dirents_buff_size()); if(large_buffer.size() < buffer_size) { large_buffer.resize(buffer_size); } thallium::bulk exposed_buffer; bool expose_needed = true; // Chunking loop while(true) { const int max_retries = 2; bool chunk_success = false; for(int attempt = 0; attempt < max_retries; ++attempt) { // Expose buffer if(expose_needed) { std::vector<std::pair<void*, std::size_t>> segments = { std::make_pair(large_buffer.get(), buffer_size)}; thallium::bulk exposed_buffer; std::make_pair(large_buffer.data(), buffer_size)}; try { exposed_buffer = CTX->rpc_engine()->expose( segments, thallium::bulk_mode::read_write); expose_needed = false; } catch(const std::exception& e) { LOG(ERROR, "Failed to expose buffer: {}", e.what()); return make_pair(EBUSY, nullptr); } auto endp = CTX->hosts().at(targets[server]); } gkfs::rpc::rpc_get_dirents_in_t in; in.path = root_path; in.start_key = start_key; in.bulk_handle = exposed_buffer; auto get_dirents_rpc = CTX->rpc_engine()->define( gkfs::rpc::tag::get_dirents_extended); try { gkfs::rpc::rpc_get_dirents_out_t out = get_dirents_rpc.on(endp)(in); Loading @@ -784,9 +726,15 @@ forward_get_dirents_single(const string& path, int server, LOG(WARNING, "{}() Buffer too small. Server requested {} bytes. Retrying.", __func__, required_size); if(required_size <= buffer_size) { err = ENOBUFS; break; } buffer_size = required_size; large_buffer = std::unique_ptr<char[]>(new char[buffer_size]); if(large_buffer.size() < buffer_size) { large_buffer.resize(buffer_size); } expose_needed = true; continue; // Retry with new buffer size } else if(out.err != 0) { err = out.err; Loading @@ -794,7 +742,7 @@ forward_get_dirents_single(const string& path, int server, } auto current_entries = gkfs::rpc::decompress_and_parse_entries( out, large_buffer.get()); out, large_buffer.data()); if(current_entries.empty()) { return make_pair(0, std::move(all_entries)); Loading src/client/rpc/forward_metadata_proxy.cpp +3 −2 Viewed Changes for src/client/rpc/forward_metadata_proxy.cpp: 3 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -143,7 +143,8 @@ forward_get_metadentry_size_proxy(const std::string& path) { pair<int, unique_ptr<vector< tuple<const std::string, unsigned char, size_t, time_t>>>> forward_get_dirents_single_proxy_v2(const string& path, int server) { forward_get_dirents_single_proxy_v2(const string& path, int server, const std::string& start_key_arg) { LOG(DEBUG, "{}() enter for path '{}', server '{}'", __func__, path, server); auto endp = CTX->proxy_host(); Loading @@ -151,7 +152,7 @@ forward_get_dirents_single_proxy_v2(const string& path, int server) { auto all_entries = make_unique< vector<tuple<const std::string, unsigned char, size_t, time_t>>>(); int err = 0; string start_key = ""; string start_key = start_key_arg; // Chunking loop: keep fetching until no more entries are returned while(true) { Loading Loading
examples/gfind/pfind.sh +1 −1 Viewed Changes for examples/gfind/pfind.sh: 1 added line, 1 removed line. Original line number Diff line number Diff line Loading @@ -34,7 +34,7 @@ GKFS_FIND_PROCESS=10 GKFS_SERVERS=$SLURM_JOB_NUM_NODES GKFS_FIND=sfind srun -N $NUM_NODES -n $GKFS_FIND_PROCESS --overlap --overcommit --mem=0 --oversubscribe --export=ALL,LD_PRELOAD=${GKFS} $GKFS_FIND $@ -M $GKFS_MNT -S $GKFS_SERVERS srun -N $NUM_NODES -n $GKFS_FIND_PROCESS --overlap --overcommit --mem=0 --oversubscribe --export=ALL,LD_PRELOAD=${GKFS} $GKFS_FIND $@ -M $GKFS_MNT -S $GKFS_SERVERS --server-side -C #!/bin/bash # scripts/aggregate_sfind_results.sh Loading
include/client/rpc/forward_metadata_proxy.hpp +2 −1 Viewed Changes for include/client/rpc/forward_metadata_proxy.hpp: 2 added lines, 1 removed line. Original line number Diff line number Diff line Loading @@ -54,7 +54,8 @@ forward_get_metadentry_size_proxy(const std::string& path); std::pair<int, std::unique_ptr<std::vector<std::tuple< const std::string, unsigned char, size_t, time_t>>>> forward_get_dirents_single_proxy_v2(const std::string& path, int server); forward_get_dirents_single_proxy_v2(const std::string& path, int server, const std::string& start_key = ""); } // namespace gkfs::rpc Loading
include/client/rpc/utils.hpp +7 −3 Viewed Changes for include/client/rpc/utils.hpp: 7 added lines, 3 removed lines. Original line number Diff line number Diff line Loading @@ -105,6 +105,7 @@ decompress_and_parse_entries(const OutputOrErr& out, std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> entries; entries.reserve(static_cast<size_t>((end - p) / 32) + 1); while(p < end) { if(p + sizeof(unsigned char) > end) { Loading @@ -112,7 +113,8 @@ decompress_and_parse_entries(const OutputOrErr& out, __func__); break; } unsigned char type = *reinterpret_cast<const unsigned char*>(p); unsigned char type = 0; std::memcpy(&type, p, sizeof(unsigned char)); p += sizeof(unsigned char); if(p + sizeof(size_t) > end) { Loading @@ -120,7 +122,8 @@ decompress_and_parse_entries(const OutputOrErr& out, __func__); break; } size_t file_size = *reinterpret_cast<const size_t*>(p); size_t file_size = 0; std::memcpy(&file_size, p, sizeof(size_t)); p += sizeof(size_t); if(p + sizeof(time_t) > end) { Loading @@ -128,7 +131,8 @@ decompress_and_parse_entries(const OutputOrErr& out, __func__); break; } time_t ctime = *reinterpret_cast<const time_t*>(p); time_t ctime = 0; std::memcpy(&ctime, p, sizeof(time_t)); p += sizeof(time_t); // Name is null-terminated. We must check we don't go past 'end'. Loading
src/client/rpc/forward_metadata.cpp +85 −137 Viewed Changes for src/client/rpc/forward_metadata.cpp: 85 added lines, 137 removed lines. Original line number Diff line number Diff line Loading @@ -238,24 +238,18 @@ forward_mk_symlink(const std::string& path, const std::string& target_path) { } // function matches the standard rpc_get_dirents_out_t inline std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> decompress_and_parse_entries_standard( const gkfs::rpc::rpc_get_dirents_out_t& out, const void* compressed_buffer) { // Duplicated: 'out' type differs (proxy vs daemon) template <typename OutputType> std::pair<const char*, std::size_t> decompress_dirents_payload(const OutputType& out, const void* compressed_buffer, std::vector<char>& decompressed_data) { if(out.err != 0) { throw std::runtime_error("Server returned an error: " + std::to_string(out.err)); } if(out.dirents_size == 0) { return {}; // No entries, return empty vector return {nullptr, 0}; } const char* p = nullptr; const char* end = nullptr; std::vector<char> decompressed_data; if(gkfs::config::rpc::use_dirents_compression) { const unsigned long long uncompressed_size = ZSTD_getFrameContentSize(compressed_buffer, out.dirents_size); Loading Loading @@ -283,100 +277,29 @@ decompress_and_parse_entries_standard( throw std::runtime_error("Decompression size mismatch."); } p = decompressed_data.data(); end = p + uncompressed_size; return {decompressed_data.data(), static_cast<std::size_t>(uncompressed_size)}; } else { p = static_cast<const char*>(compressed_buffer); end = p + out.dirents_size; } std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> entries; entries.reserve(out.dirents_size); // Approx while(p < end) { if(p + sizeof(unsigned char) > end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing type", __func__); break; } unsigned char type = *reinterpret_cast<const unsigned char*>(p); p += sizeof(unsigned char); // Name is null-terminated. We must check we don't go past 'end'. size_t name_len = strnlen(p, end - p); if(p + name_len >= end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing name", __func__); break; } std::string name(p, name_len); p += name_len + 1; if(!name.empty()) { entries.emplace_back(name, type, 0, 0); return {static_cast<const char*>(compressed_buffer), out.dirents_size}; } } return entries; } // Helper for filtered entries which include size and ctime inline std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> decompress_and_parse_entries_filtered( const gkfs::rpc::rpc_get_dirents_filtered_out_t& out, decompress_and_parse_entries_standard( const gkfs::rpc::rpc_get_dirents_out_t& out, const void* compressed_buffer) { if(out.err != 0) { throw std::runtime_error("Server returned an error: " + std::to_string(out.err)); } if(out.dirents_size == 0) { return {}; // No entries, return empty vector } const char* p = nullptr; const char* end = nullptr; std::vector<char> decompressed_data; if(gkfs::config::rpc::use_dirents_compression) { const unsigned long long uncompressed_size = ZSTD_getFrameContentSize(compressed_buffer, out.dirents_size); if(uncompressed_size == ZSTD_CONTENTSIZE_ERROR) { throw std::runtime_error( "Received data is not a valid Zstd frame."); } if(uncompressed_size == ZSTD_CONTENTSIZE_UNKNOWN) { throw std::runtime_error( "Zstd frame content size is unknown and was not written in the frame."); } decompressed_data.resize(uncompressed_size); const size_t result_size = ZSTD_decompress(decompressed_data.data(), uncompressed_size, compressed_buffer, out.dirents_size); if(ZSTD_isError(result_size)) { throw std::runtime_error( "Zstd decompression failed: " + std::string(ZSTD_getErrorName(result_size))); } if(result_size != uncompressed_size) { throw std::runtime_error("Decompression size mismatch."); auto [payload, payload_size] = decompress_dirents_payload( out, compressed_buffer, decompressed_data); if(payload_size == 0) { return {}; } const char* p = payload; const char* end = payload + payload_size; p = decompressed_data.data(); end = p + uncompressed_size; } else { p = static_cast<const char*>(compressed_buffer); end = p + out.dirents_size; } std::vector<std::tuple<const std::string, unsigned char, size_t, time_t>> entries; // We don't know exact count, but can optimize reserve? // entries.reserve(out.dirents_size / 32); entries.reserve((payload_size / 32U) + 1U); while(p < end) { if(p + sizeof(unsigned char) > end) { Loading @@ -387,22 +310,6 @@ decompress_and_parse_entries_filtered( unsigned char type = *reinterpret_cast<const unsigned char*>(p); p += sizeof(unsigned char); if(p + sizeof(size_t) > end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing size", __func__); break; } size_t size = *reinterpret_cast<const size_t*>(p); p += sizeof(size_t); if(p + sizeof(time_t) > end) { LOG(ERROR, "{}() Unexpected end of buffer while parsing ctime", __func__); break; } time_t ctime = *reinterpret_cast<const time_t*>(p); p += sizeof(time_t); // Name is null-terminated. We must check we don't go past 'end'. size_t name_len = strnlen(p, end - p); if(p + name_len >= end) { Loading @@ -415,7 +322,7 @@ decompress_and_parse_entries_filtered( p += name_len + 1; if(!name.empty()) { entries.emplace_back(name, type, size, ctime); entries.emplace_back(name, type, 0, 0); } } Loading Loading @@ -682,12 +589,23 @@ forward_update_metadentry_size(const string& path, const size_t size, int err = 0; off64_t ret_offset = 0; bool valid = false; bool primary_valid = false; off64_t primary_offset = 0; struct pending_update { int copy; thallium::async_response response; }; std::vector<pending_update> waiters; waiters.reserve(num_copies + 1); auto update_rpc = CTX->rpc_engine()->define(gkfs::rpc::tag::update_metadentry_size); for(auto copy = 0; copy < num_copies + 1; copy++) { auto endp = CTX->hosts().at( CTX->distributor()->locate_file_metadata(path, copy)); try { gkfs::rpc::rpc_update_metadentry_size_in_t in; in.path = path; in.size = size; Loading @@ -695,27 +613,45 @@ forward_update_metadentry_size(const string& path, const size_t size, in.append = append_flag; in.clear_inline = clear_inline_flag; auto update_rpc = CTX->rpc_engine()->define( gkfs::rpc::tag::update_metadentry_size); gkfs::rpc::rpc_update_metadentry_size_out_t out = update_rpc.on(endp)(in); try { waiters.push_back({copy, update_rpc.on(endp).async(in)}); } catch(const std::exception& ex) { LOG(ERROR, "{}() posting rpc for path '{}' replica {} failed: {}", __func__, path, copy, ex.what()); err = EBUSY; } } for(auto& waiter : waiters) { try { gkfs::rpc::rpc_update_metadentry_size_out_t out = waiter.response.wait(); if(out.err == 0) { valid = true; if(waiter.copy == 0) { primary_valid = true; primary_offset = out.ret_offset; } if(!primary_valid) { ret_offset = out.ret_offset; } } else { err = out.err; } } catch(const std::exception& ex) { LOG(ERROR, "{}() getting rpc output for path '{}' replica {} failed: {}", __func__, path, copy, ex.what()); // Continue to other replicas __func__, path, waiter.copy, ex.what()); err = EBUSY; } } if(!valid) return make_pair(err, 0); return make_pair(err ? err : EIO, 0); if(primary_valid) { ret_offset = primary_offset; } return make_pair(err, ret_offset); } Loading @@ -738,13 +674,9 @@ forward_get_dirents_single(const string& path, int server, vector<tuple<const std::string, unsigned char, size_t, time_t>>>(); int err = 0; string start_key = start_key_arg; // Chunking loop while(true) { size_t buffer_size = CTX->dirents_buff_size(); auto large_buffer = std::unique_ptr<char[]>(new char[buffer_size]); const int max_retries = 2; bool chunk_success = false; auto endp = CTX->hosts().at(targets[server]); auto get_dirents_rpc = CTX->rpc_engine()->define(gkfs::rpc::tag::get_dirents_extended); // Ensure path ends with / for get_dirents prefix check std::string root_path = path; Loading @@ -752,29 +684,39 @@ forward_get_dirents_single(const string& path, int server, root_path += '/'; } thread_local std::vector<char> large_buffer; size_t buffer_size = std::max<std::size_t>(std::size_t{1}, CTX->dirents_buff_size()); if(large_buffer.size() < buffer_size) { large_buffer.resize(buffer_size); } thallium::bulk exposed_buffer; bool expose_needed = true; // Chunking loop while(true) { const int max_retries = 2; bool chunk_success = false; for(int attempt = 0; attempt < max_retries; ++attempt) { // Expose buffer if(expose_needed) { std::vector<std::pair<void*, std::size_t>> segments = { std::make_pair(large_buffer.get(), buffer_size)}; thallium::bulk exposed_buffer; std::make_pair(large_buffer.data(), buffer_size)}; try { exposed_buffer = CTX->rpc_engine()->expose( segments, thallium::bulk_mode::read_write); expose_needed = false; } catch(const std::exception& e) { LOG(ERROR, "Failed to expose buffer: {}", e.what()); return make_pair(EBUSY, nullptr); } auto endp = CTX->hosts().at(targets[server]); } gkfs::rpc::rpc_get_dirents_in_t in; in.path = root_path; in.start_key = start_key; in.bulk_handle = exposed_buffer; auto get_dirents_rpc = CTX->rpc_engine()->define( gkfs::rpc::tag::get_dirents_extended); try { gkfs::rpc::rpc_get_dirents_out_t out = get_dirents_rpc.on(endp)(in); Loading @@ -784,9 +726,15 @@ forward_get_dirents_single(const string& path, int server, LOG(WARNING, "{}() Buffer too small. Server requested {} bytes. Retrying.", __func__, required_size); if(required_size <= buffer_size) { err = ENOBUFS; break; } buffer_size = required_size; large_buffer = std::unique_ptr<char[]>(new char[buffer_size]); if(large_buffer.size() < buffer_size) { large_buffer.resize(buffer_size); } expose_needed = true; continue; // Retry with new buffer size } else if(out.err != 0) { err = out.err; Loading @@ -794,7 +742,7 @@ forward_get_dirents_single(const string& path, int server, } auto current_entries = gkfs::rpc::decompress_and_parse_entries( out, large_buffer.get()); out, large_buffer.data()); if(current_entries.empty()) { return make_pair(0, std::move(all_entries)); Loading
src/client/rpc/forward_metadata_proxy.cpp +3 −2 Viewed Changes for src/client/rpc/forward_metadata_proxy.cpp: 3 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -143,7 +143,8 @@ forward_get_metadentry_size_proxy(const std::string& path) { pair<int, unique_ptr<vector< tuple<const std::string, unsigned char, size_t, time_t>>>> forward_get_dirents_single_proxy_v2(const string& path, int server) { forward_get_dirents_single_proxy_v2(const string& path, int server, const std::string& start_key_arg) { LOG(DEBUG, "{}() enter for path '{}', server '{}'", __func__, path, server); auto endp = CTX->proxy_host(); Loading @@ -151,7 +152,7 @@ forward_get_dirents_single_proxy_v2(const string& path, int server) { auto all_entries = make_unique< vector<tuple<const std::string, unsigned char, size_t, time_t>>>(); int err = 0; string start_key = ""; string start_key = start_key_arg; // Chunking loop: keep fetching until no more entries are returned while(true) { Loading