Commit 17622406 authored by Ramon Nou's avatar Ramon Nou
Browse files

contain daemon RPC response failures

parent 4e643235
Loading
Loading
Loading
Loading
+33 −11
Original line number Diff line number Diff line
#pragma once
#include <thallium.hpp>
#include <spdlog/spdlog.h>
#include <atomic>
#include <cstdint>

namespace gkfs::utils {

@@ -13,31 +15,51 @@ namespace gkfs::utils {
 * @endinternal
 */
template <typename RequestType, typename ResponseType>
void
safe_respond(RequestType& req, const ResponseType& resp) {
bool
safe_respond(RequestType& req, const ResponseType& resp,
             const char* context = "unknown_rpc") noexcept {
    try {
        req.respond(resp);
        return true;
    } catch(const thallium::margo_exception& e) {
        // Client vanished — log and silently discard.
        // This is a normal part of malleable workloads, not an error.
        static std::atomic<uint64_t> transport_failures{0};
        const auto failure = transport_failures.fetch_add(1) + 1;
        if(failure == 1 || failure % 100 == 0) {
            try {
                auto logger = spdlog::get("daemon");
                if(logger) {
            logger->debug(
                    "handler: client vanished mid-RPC, respond failed: {}",
                    e.what());
                    logger->warn(
                            "rpc_response_failed=1 rpc_context={} "
                            "failure_kind=client_disconnect failure_count={} "
                            "cause='{}'",
                            context, failure, e.what());
                }
            } catch(...) {
            }
        }
        return false;
    } catch(const std::exception& e) {
        // Unknown error — log but do not abort
        try {
            auto logger = spdlog::get("daemon");
            if(logger) {
            logger->error("handler: unexpected respond error: {}", e.what());
                logger->error("rpc_response_failed=1 rpc_context={} "
                              "failure_kind=unexpected cause='{}'",
                              context, e.what());
            }
        } catch(...) {
        // Keep transport-specific failures from escaping the RPC handler.
        }
        return false;
    } catch(...) {
        try {
            auto logger = spdlog::get("daemon");
            if(logger) {
            logger->error("handler: unknown respond error");
                logger->error("rpc_response_failed=1 rpc_context={} "
                              "failure_kind=unknown",
                              context);
            }
        } catch(...) {
        }
        return false;
    }
}

+6 −4
Original line number Diff line number Diff line
@@ -67,7 +67,8 @@ namespace gkfs::rpc {
 */
template <typename InputType, typename OutputType, typename Func>
void
run_rpc_handler(const tl::request& req, const InputType& in, Func func) {
run_rpc_handler(const tl::request& req, const InputType& in, Func func,
                const char* context = "unknown_rpc") {
    OutputType out{};
    const auto started = std::chrono::steady_clock::now();
    try {
@@ -122,7 +123,7 @@ run_rpc_handler(const tl::request& req, const InputType& in, Func func) {
        GKFS_DATA->stats()->record_rpc(static_cast<uint64_t>(elapsed.count()),
                                       out.err != 0);
    }
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, context);
}

/**
@@ -135,7 +136,8 @@ run_rpc_handler(const tl::request& req, const InputType& in, Func func) {
 */
template <typename OutputType, typename Func>
void
run_rpc_handler(const tl::request& req, Func func) {
run_rpc_handler(const tl::request& req, Func func,
                const char* context = "unknown_rpc") {
    OutputType out{};
    const auto started = std::chrono::steady_clock::now();
    try {
@@ -173,7 +175,7 @@ run_rpc_handler(const tl::request& req, Func func) {
        GKFS_DATA->stats()->record_rpc(static_cast<uint64_t>(elapsed.count()),
                                       out.err != 0);
    }
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, context);
}

} // namespace gkfs::rpc
+13 −13
Original line number Diff line number Diff line
@@ -71,7 +71,7 @@ rpc_srv_mutate_start(const tl::request& req,
                in.old_server_conf, in.new_server_conf, in.new_hosts_file);
        if(duplicate_request) {
            out.err = 0;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return;
        }
        registered_mutate = true;
@@ -106,7 +106,7 @@ rpc_srv_mutate_start(const tl::request& req,

    GKFS_DATA->spdlogger()->debug("{}() Sending output err '{}'", __func__,
                                  out.err);
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

void
@@ -130,7 +130,7 @@ rpc_srv_mutate_status(const tl::request& req) {
    }
    GKFS_DATA->spdlogger()->debug("{}() Sending output err '{}'", __func__,
                                  out.err);
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

void
@@ -155,7 +155,7 @@ rpc_srv_mutate_detailed_status(const tl::request& req) {
    } catch(...) {
        out.err = -1;
    }
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

void
@@ -168,7 +168,7 @@ rpc_srv_mutate_finalize(const tl::request& req) {
                    "{}() Refusing finalize while redistribution is running",
                    __func__);
            out.err = EBUSY;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return;
        }
        if(GKFS_DATA->redist_failed()) {
@@ -176,7 +176,7 @@ rpc_srv_mutate_finalize(const tl::request& req) {
                    "{}() Refusing finalize after failed redistribution",
                    __func__);
            out.err = EIO;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return;
        }
        if(!GKFS_DATA->mutate_start_active()) {
@@ -184,7 +184,7 @@ rpc_srv_mutate_finalize(const tl::request& req) {
                    "{}() Refusing finalize without an active mutation",
                    __func__);
            out.err = EINVAL;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return;
        }
        GKFS_DATA->maintenance_mode(false);
@@ -204,7 +204,7 @@ rpc_srv_mutate_finalize(const tl::request& req) {

    GKFS_DATA->spdlogger()->debug("{}() Sending output err '{}'", __func__,
                                  out.err);
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

void
@@ -218,7 +218,7 @@ rpc_srv_mutate_reload(const tl::request& req,
            GKFS_DATA->spdlogger()->warn(
                    "{}() Refusing reload while mutation is active", __func__);
            out.err = EBUSY;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return;
        }
        GKFS_DATA->malleable_manager()->reload_hosts_file(in.hosts_file);
@@ -233,7 +233,7 @@ rpc_srv_mutate_reload(const tl::request& req,
                __func__);
        out.err = -1;
    }
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

void
@@ -249,7 +249,7 @@ rpc_srv_mutate_shutdown(const tl::request& req) {
                    "{}() Refusing shutdown while mutation is incomplete",
                    __func__);
            out.err = EBUSY;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return;
        }
        GKFS_DATA->maintenance_mode(false);
@@ -265,7 +265,7 @@ rpc_srv_mutate_shutdown(const tl::request& req) {
        out.err = -1;
    }

    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);

    if(out.err == 0) {
        std::thread([]() {
@@ -287,7 +287,7 @@ rpc_srv_migrate_metadata(const tl::request& req,
                                      __func__, e.what());
        out.err = -1;
    }
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

// } // namespace
+19 −18
Original line number Diff line number Diff line
@@ -102,7 +102,8 @@ rpc_srv_create(const tl::request& req, const gkfs::rpc::rpc_mk_node_in_t& in) {
                gkfs::metadata::Metadata md(in.mode);
                gkfs::metadata::create(in.path, md);
                out.err = 0;
            });
            },
            __func__);
    if(GKFS_DATA->enable_stats()) {
        GKFS_DATA->stats()->add_value_iops(
                gkfs::utils::Stats::IopsOp::iops_create);
@@ -546,7 +547,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,

    if(entries.empty()) {
        out.err = 0;
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return 0;
    }

@@ -613,7 +614,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,
                        if(entries_serialized == 0) {
                            out.err = ENOBUFS;
                            out.dirents_size = compressed_bound;
                            gkfs::utils::safe_respond(req, out);
                            gkfs::utils::safe_respond(req, out, __func__);
                            return 0;
                        } else {
                            break; // Buffer full
@@ -626,7 +627,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,
                            out.err = ENOBUFS;
                            out.dirents_size =
                                    uncompressed_data.size() + entry_size;
                            gkfs::utils::safe_respond(req, out);
                            gkfs::utils::safe_respond(req, out, __func__);
                            return 0;
                        } else {
                            break; // Buffer full
@@ -682,7 +683,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,
                                          __func__,
                                          ZSTD_getErrorName(compressed_size));
            out.err = EIO;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return 0;
        }
        // Double check fits (should match bound check roughly)
@@ -692,7 +693,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,
                    __func__, compressed_size, client_bulk_size);
            out.err = ENOBUFS;
            out.dirents_size = compressed_size;
            gkfs::utils::safe_respond(req, out);
            gkfs::utils::safe_respond(req, out, __func__);
            return 0;
        }

@@ -723,7 +724,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,
        GKFS_DATA->spdlogger()->error("{}() Failed to push data to client: {}",
                                      __func__, e.what());
        out.err = EBUSY;
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return 0;
    }

@@ -732,7 +733,7 @@ get_dirents_helper(const std::shared_ptr<tl::engine>& engine,
    GKFS_DATA->spdlogger()->debug(
            "{}() Sending output response: err='{}', size='{}'. DONE", __func__,
            out.err, out.dirents_size);
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
    return entries_serialized;
}

@@ -755,7 +756,7 @@ rpc_srv_get_dirents(const std::shared_ptr<tl::engine>& engine,
    } catch(const std::exception& e) {
        GKFS_DATA->spdlogger()->error("{}() Error during get_dirents(): '{}'",
                                      __func__, e.what());
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return;
    }

@@ -795,7 +796,7 @@ rpc_srv_get_dirents_extended(const std::shared_ptr<tl::engine>& engine,
        GKFS_DATA->spdlogger()->error(
                "{}() Error during get_dirents_extended(): '{}'", __func__,
                e.what());
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return;
    }

@@ -842,7 +843,7 @@ rpc_srv_get_dirents_filtered(
        GKFS_DATA->spdlogger()->error(
                "{}() Error during get_dirents_filtered(): '{}'", __func__,
                e.what());
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return;
    }

@@ -887,7 +888,7 @@ rpc_srv_mk_symlink(const tl::request& req,

    if(!gkfs::config::metadata::symlink_support) {
        out.err = ENOTSUP;
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return;
    }

@@ -916,7 +917,7 @@ rpc_srv_mk_symlink(const tl::request& req,

    GKFS_DATA->spdlogger()->debug("{}() Sending output err '{}'", __func__,
                                  out.err);
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

/**
@@ -936,7 +937,7 @@ rpc_srv_rename(const tl::request& req, const gkfs::rpc::rpc_rename_in_t& in) {

    if(!gkfs::config::metadata::rename_support) {
        out.err = ENOTSUP;
        gkfs::utils::safe_respond(req, out);
        gkfs::utils::safe_respond(req, out, __func__);
        return;
    }

@@ -970,7 +971,7 @@ rpc_srv_rename(const tl::request& req, const gkfs::rpc::rpc_rename_in_t& in) {

    GKFS_DATA->spdlogger()->debug("{}() Sending output err '{}'", __func__,
                                  out.err);
    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

/**
@@ -1053,7 +1054,7 @@ rpc_srv_write_data_inline(const tl::request& req,
        out.err = EIO;
    }

    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

/**
@@ -1101,7 +1102,7 @@ rpc_srv_create_write_inline(const tl::request& req,
        out.err = -1;
    }

    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}

/**
@@ -1163,7 +1164,7 @@ rpc_srv_read_data_inline(const tl::request& req,
        out.err = EIO;
    }

    gkfs::utils::safe_respond(req, out);
    gkfs::utils::safe_respond(req, out, __func__);
}