Loading include/client/env.hpp +4 −0 Viewed Changes for include/client/env.hpp: 4 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -69,6 +69,10 @@ static constexpr auto METRICS_IP_PORT = ADD_PREFIX("METRICS_IP_PORT"); #endif static constexpr auto PROTECT_FD = ADD_PREFIX("PROTECT_FD"); static constexpr auto PROTECT_FILES_GENERATOR = ADD_PREFIX("PROTECT_FILES_GENERATOR"); static constexpr auto PROTECT_FILES_CONSUMER = ADD_PREFIX("PROTECT_FILES_CONSUMER"); static constexpr auto NUM_REPL = ADD_PREFIX("NUM_REPL"); static constexpr auto PROXY_PID_FILE = ADD_PREFIX("PROXY_PID_FILE"); namespace cache { Loading include/client/preload_context.hpp +14 −0 Viewed Changes for include/client/preload_context.hpp: 14 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -140,6 +140,8 @@ private: int replicas_; bool protect_fds_{false}; bool protect_files_generator_{false}; bool protect_files_consumer_{false}; std::shared_ptr<gkfs::messagepack::ClientMetrics> write_metrics_; std::shared_ptr<gkfs::messagepack::ClientMetrics> read_metrics_; Loading Loading @@ -322,6 +324,18 @@ public: void protect_fds(bool protect); bool protect_files_generator() const; void protect_files_generator(bool protect); bool protect_files_consumer() const; void protect_files_consumer(bool protect); const std::shared_ptr<gkfs::messagepack::ClientMetrics> write_metrics(); Loading src/client/gkfs_functions.cpp +107 −3 Viewed Changes for src/client/gkfs_functions.cpp: 107 added lines, 3 removed lines. Original line number Diff line number Diff line Loading @@ -48,6 +48,8 @@ #include <client/rpc/forward_data_proxy.hpp> #include <client/open_dir.hpp> #include <client/cache.hpp> #include <string> #include <string_view> #include <common/path_util.hpp> #ifdef GKFS_ENABLE_CLIENT_METRICS Loading Loading @@ -145,6 +147,65 @@ check_parent_dir(const std::string& path) { namespace gkfs::syscall { /** * @brief generate_lock_file * @param path * @param increase * * Creates, if it does not exist, a lock file, path+".lockgekko", empty * If increase is true, increase the size of the file +1 * if increase is false, decrease the size of the file -1 * If size == 0, delete the file * Using calls : forward_create, forward_stat, forward_remove, forward_decr_size * and forward_update_metadentry_size Proxy not supported */ void generate_lock_file(const std::string& path, bool increase) { auto lock_path = path + ".lockgekko"; if(increase) { auto md = gkfs::utils::get_metadata(lock_path); if(!md) { gkfs::rpc::forward_create(lock_path, 0777 | S_IFREG, 0); } gkfs::rpc::forward_update_metadentry_size(lock_path, 1, 0, false, 0); } else { auto md = gkfs::utils::get_metadata(lock_path); if(md) { if(md.value().size() == 1) { LOG(DEBUG, "Deleting Lock file {}", lock_path); gkfs::rpc::forward_remove(lock_path, false, 0); } else { gkfs::rpc::forward_decr_size(lock_path, md.value().size() - 1, 0); } } } } /** * @brief generate_lock_file * @param path * * Test if the lock file exists, if it exists, wait 0.5 second and check again * (max 80 times) Using calls : forward_stat */ void test_lock_file(const std::string& path) { auto lock_path = path + ".lockgekko"; auto md = gkfs::utils::get_metadata(lock_path); if(md) { LOG(DEBUG, "Lock file exists {} --> {}", lock_path, md->size()); for(int i = 0; i < 80; i++) { if(!md) { break; } std::this_thread::sleep_for(std::chrono::milliseconds(500)); md = gkfs::utils::get_metadata(lock_path); } } } /** * gkfs wrapper for open() system calls * errno may be set Loading Loading @@ -197,9 +258,14 @@ gkfs_open(const std::string& path, mode_t mode, int flags) { return -1; } } else { // file was successfully created. Add to filemap return CTX->file_map()->add( auto fd = CTX->file_map()->add( std::make_shared<gkfs::filemap::OpenFile>(path, flags)); // CREATE_MODE if(CTX->protect_files_generator()) { generate_lock_file(path, true); } // file was successfully created. Add to filemap return fd; } } else { auto md_ = gkfs::utils::get_metadata(path); Loading Loading @@ -277,8 +343,18 @@ gkfs_open(const std::string& path, mode_t mode, int flags) { } } return CTX->file_map()->add( auto fd = CTX->file_map()->add( std::make_shared<gkfs::filemap::OpenFile>(path, flags)); if(CTX->protect_files_consumer()) { test_lock_file(path); } if(CTX->protect_files_generator()) { generate_lock_file(path, true); } return fd; } /** Loading Loading @@ -1480,6 +1556,14 @@ gkfs_getdents(unsigned int fd, struct linux_dirent* dirp, unsigned int count) { while(pos < open_dir->size()) { // get dentry fir current position auto de = open_dir->getdent(pos); if(CTX->protect_files_consumer() or CTX->protect_files_generator()) { // if de.name ends with lockgekko jump to the next file if(de.name().size() >= 10 && de.name().substr(de.name().size() - 10) == ".lockgekko") { pos++; continue; } } /* * Calculate the total dentry size within the kernel struct * `linux_dirent` depending on the file name size. The size is then Loading Loading @@ -1549,6 +1633,14 @@ gkfs_getdents64(unsigned int fd, struct linux_dirent64* dirp, struct linux_dirent64* current_dirp = nullptr; while(pos < open_dir->size()) { auto de = open_dir->getdent(pos); if(CTX->protect_files_consumer() or CTX->protect_files_generator()) { // if de.name ends with lockgekko jump to the next file if(de.name().size() >= 10 && de.name().substr(de.name().size() - 10) == ".lockgekko") { pos++; continue; } } /* * Calculate the total dentry size within the kernel struct * `linux_dirent` depending on the file name size. The size is then Loading Loading @@ -1644,6 +1736,11 @@ gkfs_close(unsigned int fd) { CTX->file_map()->get(fd)->path()); } } if(CTX->protect_files_generator()) { auto path = CTX->file_map()->get(fd)->path(); generate_lock_file(path, false); } // No call to the daemon is required CTX->file_map()->remove(fd); return 0; Loading Loading @@ -1764,6 +1861,13 @@ gkfs_get_file_list(const std::string& path) { while(pos < open_dir->size()) { auto de = open_dir->getdent(pos++); if(CTX->protect_files_consumer() or CTX->protect_files_generator()) { // if de.name ends with lockgekko jump to the next file if(de.name().size() >= 10 && de.name().substr(de.name().size() - 10) == ".lockgekko") { continue; } } file_list.push_back(de.name()); } return file_list; Loading src/client/preload.cpp +9 −0 Viewed Changes for src/client/preload.cpp: 9 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -239,6 +239,15 @@ init_environment() { "Failed to connect to hosts: "s + e.what()); } CTX->protect_files_generator( gkfs::env::get_var(gkfs::env::PROTECT_FILES_GENERATOR, 0)); CTX->protect_files_consumer( gkfs::env::get_var(gkfs::env::PROTECT_FILES_CONSUMER, 0)); LOG(INFO, "Lock-Files : Generator = {} / Consumer = {}", CTX->protect_files_generator(), CTX->protect_files_consumer()); /* Setup distributor */ auto forwarding_map_file = gkfs::env::get_var( gkfs::env::FORWARDING_MAP_FILE, gkfs::config::forwarding_file_path); Loading src/client/preload_context.cpp +21 −0 Viewed Changes for src/client/preload_context.cpp: 21 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -659,6 +659,27 @@ PreloadContext::get_replicas() { return replicas_; } bool PreloadContext::protect_files_generator() const { return protect_files_generator_; } void PreloadContext::protect_files_generator(bool protect) { protect_files_generator_ = protect; } bool PreloadContext::protect_files_consumer() const { return protect_files_consumer_; } void PreloadContext::protect_files_consumer(bool protect) { protect_files_consumer_ = protect; } const std::shared_ptr<messagepack::ClientMetrics> PreloadContext::write_metrics() { return write_metrics_; Loading Loading
include/client/env.hpp +4 −0 Viewed Changes for include/client/env.hpp: 4 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -69,6 +69,10 @@ static constexpr auto METRICS_IP_PORT = ADD_PREFIX("METRICS_IP_PORT"); #endif static constexpr auto PROTECT_FD = ADD_PREFIX("PROTECT_FD"); static constexpr auto PROTECT_FILES_GENERATOR = ADD_PREFIX("PROTECT_FILES_GENERATOR"); static constexpr auto PROTECT_FILES_CONSUMER = ADD_PREFIX("PROTECT_FILES_CONSUMER"); static constexpr auto NUM_REPL = ADD_PREFIX("NUM_REPL"); static constexpr auto PROXY_PID_FILE = ADD_PREFIX("PROXY_PID_FILE"); namespace cache { Loading
include/client/preload_context.hpp +14 −0 Viewed Changes for include/client/preload_context.hpp: 14 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -140,6 +140,8 @@ private: int replicas_; bool protect_fds_{false}; bool protect_files_generator_{false}; bool protect_files_consumer_{false}; std::shared_ptr<gkfs::messagepack::ClientMetrics> write_metrics_; std::shared_ptr<gkfs::messagepack::ClientMetrics> read_metrics_; Loading Loading @@ -322,6 +324,18 @@ public: void protect_fds(bool protect); bool protect_files_generator() const; void protect_files_generator(bool protect); bool protect_files_consumer() const; void protect_files_consumer(bool protect); const std::shared_ptr<gkfs::messagepack::ClientMetrics> write_metrics(); Loading
src/client/gkfs_functions.cpp +107 −3 Viewed Changes for src/client/gkfs_functions.cpp: 107 added lines, 3 removed lines. Original line number Diff line number Diff line Loading @@ -48,6 +48,8 @@ #include <client/rpc/forward_data_proxy.hpp> #include <client/open_dir.hpp> #include <client/cache.hpp> #include <string> #include <string_view> #include <common/path_util.hpp> #ifdef GKFS_ENABLE_CLIENT_METRICS Loading Loading @@ -145,6 +147,65 @@ check_parent_dir(const std::string& path) { namespace gkfs::syscall { /** * @brief generate_lock_file * @param path * @param increase * * Creates, if it does not exist, a lock file, path+".lockgekko", empty * If increase is true, increase the size of the file +1 * if increase is false, decrease the size of the file -1 * If size == 0, delete the file * Using calls : forward_create, forward_stat, forward_remove, forward_decr_size * and forward_update_metadentry_size Proxy not supported */ void generate_lock_file(const std::string& path, bool increase) { auto lock_path = path + ".lockgekko"; if(increase) { auto md = gkfs::utils::get_metadata(lock_path); if(!md) { gkfs::rpc::forward_create(lock_path, 0777 | S_IFREG, 0); } gkfs::rpc::forward_update_metadentry_size(lock_path, 1, 0, false, 0); } else { auto md = gkfs::utils::get_metadata(lock_path); if(md) { if(md.value().size() == 1) { LOG(DEBUG, "Deleting Lock file {}", lock_path); gkfs::rpc::forward_remove(lock_path, false, 0); } else { gkfs::rpc::forward_decr_size(lock_path, md.value().size() - 1, 0); } } } } /** * @brief generate_lock_file * @param path * * Test if the lock file exists, if it exists, wait 0.5 second and check again * (max 80 times) Using calls : forward_stat */ void test_lock_file(const std::string& path) { auto lock_path = path + ".lockgekko"; auto md = gkfs::utils::get_metadata(lock_path); if(md) { LOG(DEBUG, "Lock file exists {} --> {}", lock_path, md->size()); for(int i = 0; i < 80; i++) { if(!md) { break; } std::this_thread::sleep_for(std::chrono::milliseconds(500)); md = gkfs::utils::get_metadata(lock_path); } } } /** * gkfs wrapper for open() system calls * errno may be set Loading Loading @@ -197,9 +258,14 @@ gkfs_open(const std::string& path, mode_t mode, int flags) { return -1; } } else { // file was successfully created. Add to filemap return CTX->file_map()->add( auto fd = CTX->file_map()->add( std::make_shared<gkfs::filemap::OpenFile>(path, flags)); // CREATE_MODE if(CTX->protect_files_generator()) { generate_lock_file(path, true); } // file was successfully created. Add to filemap return fd; } } else { auto md_ = gkfs::utils::get_metadata(path); Loading Loading @@ -277,8 +343,18 @@ gkfs_open(const std::string& path, mode_t mode, int flags) { } } return CTX->file_map()->add( auto fd = CTX->file_map()->add( std::make_shared<gkfs::filemap::OpenFile>(path, flags)); if(CTX->protect_files_consumer()) { test_lock_file(path); } if(CTX->protect_files_generator()) { generate_lock_file(path, true); } return fd; } /** Loading Loading @@ -1480,6 +1556,14 @@ gkfs_getdents(unsigned int fd, struct linux_dirent* dirp, unsigned int count) { while(pos < open_dir->size()) { // get dentry fir current position auto de = open_dir->getdent(pos); if(CTX->protect_files_consumer() or CTX->protect_files_generator()) { // if de.name ends with lockgekko jump to the next file if(de.name().size() >= 10 && de.name().substr(de.name().size() - 10) == ".lockgekko") { pos++; continue; } } /* * Calculate the total dentry size within the kernel struct * `linux_dirent` depending on the file name size. The size is then Loading Loading @@ -1549,6 +1633,14 @@ gkfs_getdents64(unsigned int fd, struct linux_dirent64* dirp, struct linux_dirent64* current_dirp = nullptr; while(pos < open_dir->size()) { auto de = open_dir->getdent(pos); if(CTX->protect_files_consumer() or CTX->protect_files_generator()) { // if de.name ends with lockgekko jump to the next file if(de.name().size() >= 10 && de.name().substr(de.name().size() - 10) == ".lockgekko") { pos++; continue; } } /* * Calculate the total dentry size within the kernel struct * `linux_dirent` depending on the file name size. The size is then Loading Loading @@ -1644,6 +1736,11 @@ gkfs_close(unsigned int fd) { CTX->file_map()->get(fd)->path()); } } if(CTX->protect_files_generator()) { auto path = CTX->file_map()->get(fd)->path(); generate_lock_file(path, false); } // No call to the daemon is required CTX->file_map()->remove(fd); return 0; Loading Loading @@ -1764,6 +1861,13 @@ gkfs_get_file_list(const std::string& path) { while(pos < open_dir->size()) { auto de = open_dir->getdent(pos++); if(CTX->protect_files_consumer() or CTX->protect_files_generator()) { // if de.name ends with lockgekko jump to the next file if(de.name().size() >= 10 && de.name().substr(de.name().size() - 10) == ".lockgekko") { continue; } } file_list.push_back(de.name()); } return file_list; Loading
src/client/preload.cpp +9 −0 Viewed Changes for src/client/preload.cpp: 9 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -239,6 +239,15 @@ init_environment() { "Failed to connect to hosts: "s + e.what()); } CTX->protect_files_generator( gkfs::env::get_var(gkfs::env::PROTECT_FILES_GENERATOR, 0)); CTX->protect_files_consumer( gkfs::env::get_var(gkfs::env::PROTECT_FILES_CONSUMER, 0)); LOG(INFO, "Lock-Files : Generator = {} / Consumer = {}", CTX->protect_files_generator(), CTX->protect_files_consumer()); /* Setup distributor */ auto forwarding_map_file = gkfs::env::get_var( gkfs::env::FORWARDING_MAP_FILE, gkfs::config::forwarding_file_path); Loading
src/client/preload_context.cpp +21 −0 Viewed Changes for src/client/preload_context.cpp: 21 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -659,6 +659,27 @@ PreloadContext::get_replicas() { return replicas_; } bool PreloadContext::protect_files_generator() const { return protect_files_generator_; } void PreloadContext::protect_files_generator(bool protect) { protect_files_generator_ = protect; } bool PreloadContext::protect_files_consumer() const { return protect_files_consumer_; } void PreloadContext::protect_files_consumer(bool protect) { protect_files_consumer_ = protect; } const std::shared_ptr<messagepack::ClientMetrics> PreloadContext::write_metrics() { return write_metrics_; Loading