rocksdb tune and async

Merge request reports

Loading
+1 −1
Changes for examples/gfind/pfind.sh: 1 added line, 1 removed line.
Original line number Diff line number Diff line
@@ -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
+2 −1
Changes for include/client/rpc/forward_metadata_proxy.hpp: 2 added lines, 1 removed line.
Original line number Diff line number Diff line
@@ -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

+7 −3
Changes for include/client/rpc/utils.hpp: 7 added lines, 3 removed lines.
Original line number Diff line number Diff line
@@ -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) {
@@ -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) {
@@ -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) {
@@ -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'.
+85 −137
Changes for src/client/rpc/forward_metadata.cpp: 85 added lines, 137 removed lines.
Original line number Diff line number Diff line
@@ -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);
@@ -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) {
@@ -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) {
@@ -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);
        }
    }

@@ -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;
@@ -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);
}
@@ -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;
@@ -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);
@@ -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;
@@ -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));
+3 −2
Changes for src/client/rpc/forward_metadata_proxy.cpp: 3 added lines, 2 removed lines.
Original line number Diff line number Diff line
@@ -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();
@@ -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