Loading include/daemon/classes/fs_data.hpp +36 −0 Changes for include/daemon/classes/fs_data.hpp: 36 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -51,6 +51,17 @@ #include <thallium.hpp> /* Forward declarations */ namespace gkfs::daemon { enum class mutation_phase { idle, planned, running, failed, ready_to_finalize, }; } // namespace gkfs::daemon namespace gkfs { namespace metadata { class MetadataDB; Loading Loading @@ -146,6 +157,10 @@ private: int mutate_old_server_conf_ = 0; int mutate_new_server_conf_ = 0; std::string mutate_hosts_file_; mutation_phase mutation_phase_{mutation_phase::idle}; uint64_t placement_generation_{0}; uint64_t mutate_old_placement_generation_{0}; uint64_t mutate_new_placement_generation_{0}; // redist_running_ indicates to client that redistribution is running std::atomic<bool> redist_running_{false}; // redist_failed_ indicates that the last redistribution did not complete Loading Loading @@ -441,6 +456,27 @@ public: std::string mutate_hosts_file() const; mutation_phase mutation_state() const; uint64_t placement_generation() const; uint64_t mutate_old_placement_generation() const; uint64_t mutate_new_placement_generation() const; void mutation_running(); void mutation_completed(bool failed); void mutation_committed(); bool redist_running() const; Loading src/daemon/classes/fs_data.cpp +75 −0 Changes for src/daemon/classes/fs_data.cpp: 75 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -506,6 +506,9 @@ FsData::begin_mutate_start(int old_server_conf, int new_server_conf, mutate_old_server_conf_ = old_server_conf; mutate_new_server_conf_ = new_server_conf; mutate_hosts_file_ = hosts_file; mutation_phase_ = mutation_phase::planned; mutate_old_placement_generation_ = placement_generation_; mutate_new_placement_generation_ = placement_generation_ + 1; ABT_mutex_unlock(maintenance_mode_mutex_); return false; } Loading @@ -515,6 +518,7 @@ FsData::end_mutate_start() { ABT_mutex_lock(maintenance_mode_mutex_); mutate_start_active_ = false; mutate_hosts_file_.clear(); mutation_phase_ = mutation_phase::idle; ABT_mutex_unlock(maintenance_mode_mutex_); } Loading Loading @@ -554,6 +558,77 @@ FsData::mutate_hosts_file() const { return value; } mutation_phase FsData::mutation_state() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = mutation_phase_; ABT_mutex_unlock(*mutex); return value; } uint64_t FsData::placement_generation() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = placement_generation_; ABT_mutex_unlock(*mutex); return value; } uint64_t FsData::mutate_old_placement_generation() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = mutate_old_placement_generation_; ABT_mutex_unlock(*mutex); return value; } uint64_t FsData::mutate_new_placement_generation() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = mutate_new_placement_generation_; ABT_mutex_unlock(*mutex); return value; } void FsData::mutation_running() { ABT_mutex_lock(maintenance_mode_mutex_); if(mutation_phase_ != mutation_phase::planned) { ABT_mutex_unlock(maintenance_mode_mutex_); throw std::runtime_error("Mutation is not in planned state"); } mutation_phase_ = mutation_phase::running; ABT_mutex_unlock(maintenance_mode_mutex_); } void FsData::mutation_completed(const bool failed) { ABT_mutex_lock(maintenance_mode_mutex_); if(mutation_phase_ != mutation_phase::running) { ABT_mutex_unlock(maintenance_mode_mutex_); throw std::runtime_error("Mutation is not running"); } mutation_phase_ = failed ? mutation_phase::failed : mutation_phase::ready_to_finalize; ABT_mutex_unlock(maintenance_mode_mutex_); } void FsData::mutation_committed() { ABT_mutex_lock(maintenance_mode_mutex_); if(mutation_phase_ != mutation_phase::ready_to_finalize) { ABT_mutex_unlock(maintenance_mode_mutex_); throw std::runtime_error("Mutation is not ready to finalize"); } placement_generation_ = mutate_new_placement_generation_; mutation_phase_ = mutation_phase::idle; ABT_mutex_unlock(maintenance_mode_mutex_); } bool FsData::redist_running() const { return redist_running_.load(); Loading src/daemon/handler/srv_malleability.cpp +10 −0 Changes for src/daemon/handler/srv_malleability.cpp: 10 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -195,6 +195,15 @@ rpc_srv_mutate_finalize(const tl::request& req, gkfs::utils::safe_respond(req, out, __func__); return; } if(GKFS_DATA->mutation_state() != gkfs::daemon::mutation_phase::ready_to_finalize) { GKFS_DATA->spdlogger()->warn( "{}() Refusing finalize before successful mutation completion", __func__); out.err = EBUSY; gkfs::utils::safe_respond(req, out, __func__); return; } if(!GKFS_DATA->mutate_start_active()) { GKFS_DATA->spdlogger()->warn( "{}() Refusing finalize without an active mutation", Loading @@ -205,6 +214,7 @@ rpc_srv_mutate_finalize(const tl::request& req, } GKFS_DATA->maintenance_mode(false); GKFS_DATA->keep_hosts_file(true); GKFS_DATA->mutation_committed(); GKFS_DATA->end_mutate_start(); out.err = 0; } catch(const std::exception& e) { Loading src/daemon/malleability/malleable_manager.cpp +5 −0 Changes for src/daemon/malleability/malleable_manager.cpp: 5 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -850,6 +850,7 @@ MalleableManager::expand_abt_v2(void* _arg) { } GKFS_DATA->spdlogger()->info("{}() Starting v2 expansion process...", __func__); GKFS_DATA->mutation_running(); GKFS_DATA->redist_running(true); GKFS_DATA->redist_failed(false); auto failed = GKFS_DATA->malleable_manager()->redistribute_metadata() != 0; Loading @@ -868,6 +869,7 @@ MalleableManager::expand_abt_v2(void* _arg) { } GKFS_DATA->redist_failed(failed); GKFS_DATA->redist_running(false); GKFS_DATA->mutation_completed(failed); GKFS_DATA->spdlogger()->info( "{}() V2 expansion process successfully finished.", __func__); } Loading @@ -883,6 +885,7 @@ MalleableManager::expand_on_demand_abt(void* _arg) { GKFS_DATA->spdlogger()->info( "{}() Starting expand-on-demand process: metadata eager, data lazy.", __func__); GKFS_DATA->mutation_running(); GKFS_DATA->redist_running(true); bool failed = false; try { Loading @@ -900,6 +903,7 @@ MalleableManager::expand_on_demand_abt(void* _arg) { } GKFS_DATA->redist_failed(failed); GKFS_DATA->redist_running(false); GKFS_DATA->mutation_completed(failed); GKFS_DATA->spdlogger()->info( "{}() Expand-on-demand metadata redistribution finished. Data chunks remain lazy.", __func__); Loading Loading @@ -1206,6 +1210,7 @@ MalleableManager::mutate_start(int old_server_conf, int new_server_conf, } // Use v2 pipeline for mutate GKFS_DATA->mutation_running(); GKFS_DATA->redist_running(true); auto abt_err = ABT_thread_create(RPC_DATA->io_pool(), expand_abt_v2, this, ABT_THREAD_ATTR_NULL, &redist_thread_); Loading tests/unit/test_fs_data_faults.cpp +41 −0 Changes for tests/unit/test_fs_data_faults.cpp: 41 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -38,3 +38,44 @@ TEST_CASE("Conflicting mutation requests are rejected", data->end_mutate_start(); } TEST_CASE("Mutation state tracks generations and rejects interrupted finalize", "[resilience][mutation]") { auto* data = gkfs::daemon::FsData::getInstance(); data->end_mutate_start(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::idle); const auto old_generation = data->placement_generation(); REQUIRE_FALSE(data->begin_mutate_start(1, 2, "/tmp/hosts-generation")); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::planned); REQUIRE(data->mutate_old_placement_generation() == old_generation); REQUIRE(data->mutate_new_placement_generation() == old_generation + 1); REQUIRE_THROWS_WITH(data->mutation_committed(), "Mutation is not ready to finalize"); data->mutation_running(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::running); data->mutation_completed(true); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::failed); REQUIRE_THROWS_WITH(data->mutation_committed(), "Mutation is not ready to finalize"); data->end_mutate_start(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::idle); } TEST_CASE("Successful mutation commits the placement generation", "[resilience][mutation]") { auto* data = gkfs::daemon::FsData::getInstance(); data->end_mutate_start(); const auto old_generation = data->placement_generation(); REQUIRE_FALSE(data->begin_mutate_start(1, 2, "/tmp/hosts-commit")); data->mutation_running(); data->mutation_completed(false); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::ready_to_finalize); data->mutation_committed(); REQUIRE(data->placement_generation() == old_generation + 1); data->end_mutate_start(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::idle); } No newline at end of file Loading
include/daemon/classes/fs_data.hpp +36 −0 Changes for include/daemon/classes/fs_data.hpp: 36 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -51,6 +51,17 @@ #include <thallium.hpp> /* Forward declarations */ namespace gkfs::daemon { enum class mutation_phase { idle, planned, running, failed, ready_to_finalize, }; } // namespace gkfs::daemon namespace gkfs { namespace metadata { class MetadataDB; Loading Loading @@ -146,6 +157,10 @@ private: int mutate_old_server_conf_ = 0; int mutate_new_server_conf_ = 0; std::string mutate_hosts_file_; mutation_phase mutation_phase_{mutation_phase::idle}; uint64_t placement_generation_{0}; uint64_t mutate_old_placement_generation_{0}; uint64_t mutate_new_placement_generation_{0}; // redist_running_ indicates to client that redistribution is running std::atomic<bool> redist_running_{false}; // redist_failed_ indicates that the last redistribution did not complete Loading Loading @@ -441,6 +456,27 @@ public: std::string mutate_hosts_file() const; mutation_phase mutation_state() const; uint64_t placement_generation() const; uint64_t mutate_old_placement_generation() const; uint64_t mutate_new_placement_generation() const; void mutation_running(); void mutation_completed(bool failed); void mutation_committed(); bool redist_running() const; Loading
src/daemon/classes/fs_data.cpp +75 −0 Changes for src/daemon/classes/fs_data.cpp: 75 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -506,6 +506,9 @@ FsData::begin_mutate_start(int old_server_conf, int new_server_conf, mutate_old_server_conf_ = old_server_conf; mutate_new_server_conf_ = new_server_conf; mutate_hosts_file_ = hosts_file; mutation_phase_ = mutation_phase::planned; mutate_old_placement_generation_ = placement_generation_; mutate_new_placement_generation_ = placement_generation_ + 1; ABT_mutex_unlock(maintenance_mode_mutex_); return false; } Loading @@ -515,6 +518,7 @@ FsData::end_mutate_start() { ABT_mutex_lock(maintenance_mode_mutex_); mutate_start_active_ = false; mutate_hosts_file_.clear(); mutation_phase_ = mutation_phase::idle; ABT_mutex_unlock(maintenance_mode_mutex_); } Loading Loading @@ -554,6 +558,77 @@ FsData::mutate_hosts_file() const { return value; } mutation_phase FsData::mutation_state() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = mutation_phase_; ABT_mutex_unlock(*mutex); return value; } uint64_t FsData::placement_generation() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = placement_generation_; ABT_mutex_unlock(*mutex); return value; } uint64_t FsData::mutate_old_placement_generation() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = mutate_old_placement_generation_; ABT_mutex_unlock(*mutex); return value; } uint64_t FsData::mutate_new_placement_generation() const { auto* mutex = const_cast<ABT_mutex*>(&maintenance_mode_mutex_); ABT_mutex_lock(*mutex); const auto value = mutate_new_placement_generation_; ABT_mutex_unlock(*mutex); return value; } void FsData::mutation_running() { ABT_mutex_lock(maintenance_mode_mutex_); if(mutation_phase_ != mutation_phase::planned) { ABT_mutex_unlock(maintenance_mode_mutex_); throw std::runtime_error("Mutation is not in planned state"); } mutation_phase_ = mutation_phase::running; ABT_mutex_unlock(maintenance_mode_mutex_); } void FsData::mutation_completed(const bool failed) { ABT_mutex_lock(maintenance_mode_mutex_); if(mutation_phase_ != mutation_phase::running) { ABT_mutex_unlock(maintenance_mode_mutex_); throw std::runtime_error("Mutation is not running"); } mutation_phase_ = failed ? mutation_phase::failed : mutation_phase::ready_to_finalize; ABT_mutex_unlock(maintenance_mode_mutex_); } void FsData::mutation_committed() { ABT_mutex_lock(maintenance_mode_mutex_); if(mutation_phase_ != mutation_phase::ready_to_finalize) { ABT_mutex_unlock(maintenance_mode_mutex_); throw std::runtime_error("Mutation is not ready to finalize"); } placement_generation_ = mutate_new_placement_generation_; mutation_phase_ = mutation_phase::idle; ABT_mutex_unlock(maintenance_mode_mutex_); } bool FsData::redist_running() const { return redist_running_.load(); Loading
src/daemon/handler/srv_malleability.cpp +10 −0 Changes for src/daemon/handler/srv_malleability.cpp: 10 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -195,6 +195,15 @@ rpc_srv_mutate_finalize(const tl::request& req, gkfs::utils::safe_respond(req, out, __func__); return; } if(GKFS_DATA->mutation_state() != gkfs::daemon::mutation_phase::ready_to_finalize) { GKFS_DATA->spdlogger()->warn( "{}() Refusing finalize before successful mutation completion", __func__); out.err = EBUSY; gkfs::utils::safe_respond(req, out, __func__); return; } if(!GKFS_DATA->mutate_start_active()) { GKFS_DATA->spdlogger()->warn( "{}() Refusing finalize without an active mutation", Loading @@ -205,6 +214,7 @@ rpc_srv_mutate_finalize(const tl::request& req, } GKFS_DATA->maintenance_mode(false); GKFS_DATA->keep_hosts_file(true); GKFS_DATA->mutation_committed(); GKFS_DATA->end_mutate_start(); out.err = 0; } catch(const std::exception& e) { Loading
src/daemon/malleability/malleable_manager.cpp +5 −0 Changes for src/daemon/malleability/malleable_manager.cpp: 5 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -850,6 +850,7 @@ MalleableManager::expand_abt_v2(void* _arg) { } GKFS_DATA->spdlogger()->info("{}() Starting v2 expansion process...", __func__); GKFS_DATA->mutation_running(); GKFS_DATA->redist_running(true); GKFS_DATA->redist_failed(false); auto failed = GKFS_DATA->malleable_manager()->redistribute_metadata() != 0; Loading @@ -868,6 +869,7 @@ MalleableManager::expand_abt_v2(void* _arg) { } GKFS_DATA->redist_failed(failed); GKFS_DATA->redist_running(false); GKFS_DATA->mutation_completed(failed); GKFS_DATA->spdlogger()->info( "{}() V2 expansion process successfully finished.", __func__); } Loading @@ -883,6 +885,7 @@ MalleableManager::expand_on_demand_abt(void* _arg) { GKFS_DATA->spdlogger()->info( "{}() Starting expand-on-demand process: metadata eager, data lazy.", __func__); GKFS_DATA->mutation_running(); GKFS_DATA->redist_running(true); bool failed = false; try { Loading @@ -900,6 +903,7 @@ MalleableManager::expand_on_demand_abt(void* _arg) { } GKFS_DATA->redist_failed(failed); GKFS_DATA->redist_running(false); GKFS_DATA->mutation_completed(failed); GKFS_DATA->spdlogger()->info( "{}() Expand-on-demand metadata redistribution finished. Data chunks remain lazy.", __func__); Loading Loading @@ -1206,6 +1210,7 @@ MalleableManager::mutate_start(int old_server_conf, int new_server_conf, } // Use v2 pipeline for mutate GKFS_DATA->mutation_running(); GKFS_DATA->redist_running(true); auto abt_err = ABT_thread_create(RPC_DATA->io_pool(), expand_abt_v2, this, ABT_THREAD_ATTR_NULL, &redist_thread_); Loading
tests/unit/test_fs_data_faults.cpp +41 −0 Changes for tests/unit/test_fs_data_faults.cpp: 41 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -38,3 +38,44 @@ TEST_CASE("Conflicting mutation requests are rejected", data->end_mutate_start(); } TEST_CASE("Mutation state tracks generations and rejects interrupted finalize", "[resilience][mutation]") { auto* data = gkfs::daemon::FsData::getInstance(); data->end_mutate_start(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::idle); const auto old_generation = data->placement_generation(); REQUIRE_FALSE(data->begin_mutate_start(1, 2, "/tmp/hosts-generation")); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::planned); REQUIRE(data->mutate_old_placement_generation() == old_generation); REQUIRE(data->mutate_new_placement_generation() == old_generation + 1); REQUIRE_THROWS_WITH(data->mutation_committed(), "Mutation is not ready to finalize"); data->mutation_running(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::running); data->mutation_completed(true); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::failed); REQUIRE_THROWS_WITH(data->mutation_committed(), "Mutation is not ready to finalize"); data->end_mutate_start(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::idle); } TEST_CASE("Successful mutation commits the placement generation", "[resilience][mutation]") { auto* data = gkfs::daemon::FsData::getInstance(); data->end_mutate_start(); const auto old_generation = data->placement_generation(); REQUIRE_FALSE(data->begin_mutate_start(1, 2, "/tmp/hosts-commit")); data->mutation_running(); data->mutation_completed(false); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::ready_to_finalize); data->mutation_committed(); REQUIRE(data->placement_generation() == old_generation + 1); data->end_mutate_start(); REQUIRE(data->mutation_state() == gkfs::daemon::mutation_phase::idle); } No newline at end of file