From d07e07221219f1236fe7892563353e041ef11678 Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Mon, 24 Jun 2024 00:54:39 +0300 Subject: [PATCH] Change bool wr to event mask in epoll_manager --- src/client/http_client.cpp | 6 +++--- src/client/messenger.cpp | 10 +++++----- src/client/msgr_send.cpp | 2 +- src/client/msgr_stop.cpp | 2 +- src/client/nbd_proxy.cpp | 2 +- src/kv/kv_cli.cpp | 4 ++-- src/nfs/nfs_proxy.cpp | 8 ++++---- src/osd/osd.cpp | 2 +- src/test/stub_uring_osd.cpp | 2 +- src/util/epoll_manager.cpp | 6 +++--- src/util/epoll_manager.h | 4 +++- src/util/timerfd_manager.cpp | 6 +++--- src/util/timerfd_manager.h | 4 ++-- 13 files changed, 30 insertions(+), 28 deletions(-) diff --git a/src/client/http_client.cpp b/src/client/http_client.cpp index d0ec9a49..de20b69e 100644 --- a/src/client/http_client.cpp +++ b/src/client/http_client.cpp @@ -271,7 +271,7 @@ void http_co_t::close_connection() } if (peer_fd >= 0) { - tfd->set_fd_handler(peer_fd, false, NULL); + tfd->set_fd_handler(peer_fd, 0, NULL); close(peer_fd); peer_fd = -1; } @@ -314,7 +314,7 @@ void http_co_t::start_connection() stackout(); return; } - tfd->set_fd_handler(peer_fd, true, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(peer_fd, EPOLLIN|EPOLLOUT, [this](int peer_fd, int epoll_events) { this->epoll_events |= epoll_events; handle_events(); @@ -372,7 +372,7 @@ void http_co_t::handle_connect_result() } int one = 1; setsockopt(peer_fd, SOL_TCP, TCP_NODELAY, &one, sizeof(one)); - tfd->set_fd_handler(peer_fd, false, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(peer_fd, EPOLLIN, [this](int peer_fd, int epoll_events) { this->epoll_events |= epoll_events; handle_events(); diff --git a/src/client/messenger.cpp b/src/client/messenger.cpp index 74d0d2d2..abe241be 100644 --- a/src/client/messenger.cpp +++ b/src/client/messenger.cpp @@ -35,7 +35,7 @@ void osd_messenger_t::init() ? rdma_max_sge : rdma_context->attrx.orig_attr.max_sge; fprintf(stderr, "[OSD %ju] RDMA initialized successfully\n", osd_num); fcntl(rdma_context->channel->fd, F_SETFL, fcntl(rdma_context->channel->fd, F_GETFL, 0) | O_NONBLOCK); - tfd->set_fd_handler(rdma_context->channel->fd, false, [this](int notify_fd, int epoll_events) + tfd->set_fd_handler(rdma_context->channel->fd, EPOLLIN, [this](int notify_fd, int epoll_events) { handle_rdma_events(); }); @@ -262,7 +262,7 @@ void osd_messenger_t::try_connect_peer_addr(osd_num_t peer_osd, const char *peer clients[peer_fd]->connect_timeout_id = -1; clients[peer_fd]->osd_num = peer_osd; clients[peer_fd]->in_buf = malloc_or_die(receive_buffer_size); - tfd->set_fd_handler(peer_fd, true, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(peer_fd, EPOLLIN|EPOLLOUT, [this](int peer_fd, int epoll_events) { // Either OUT (connected) or HUP handle_connect_epoll(peer_fd); @@ -303,7 +303,7 @@ void osd_messenger_t::handle_connect_epoll(int peer_fd) int one = 1; setsockopt(peer_fd, SOL_TCP, TCP_NODELAY, &one, sizeof(one)); cl->peer_state = PEER_CONNECTED; - tfd->set_fd_handler(peer_fd, false, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(peer_fd, EPOLLIN, [this](int peer_fd, int epoll_events) { handle_peer_epoll(peer_fd, epoll_events); }); @@ -487,7 +487,7 @@ void osd_messenger_t::check_peer_config(osd_client_t *cl) fprintf(stderr, "Connected to OSD %ju using RDMA\n", cl->osd_num); } cl->peer_state = PEER_RDMA; - tfd->set_fd_handler(cl->peer_fd, false, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(cl->peer_fd, 0, [this](int peer_fd, int epoll_events) { // Do not miss the disconnection! if (epoll_events & EPOLLRDHUP) @@ -528,7 +528,7 @@ void osd_messenger_t::accept_connections(int listen_fd) clients[peer_fd]->peer_state = PEER_CONNECTED; clients[peer_fd]->in_buf = malloc_or_die(receive_buffer_size); // Add FD to epoll - tfd->set_fd_handler(peer_fd, false, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(peer_fd, EPOLLIN, [this](int peer_fd, int epoll_events) { handle_peer_epoll(peer_fd, epoll_events); }); diff --git a/src/client/msgr_send.cpp b/src/client/msgr_send.cpp index 7374a1e7..9cb8dfe7 100644 --- a/src/client/msgr_send.cpp +++ b/src/client/msgr_send.cpp @@ -296,7 +296,7 @@ void osd_messenger_t::handle_send(int result, osd_client_t *cl) fprintf(stderr, "Successfully connected with client %d using RDMA\n", cl->peer_fd); } cl->peer_state = PEER_RDMA; - tfd->set_fd_handler(cl->peer_fd, false, [this](int peer_fd, int epoll_events) + tfd->set_fd_handler(cl->peer_fd, 0, [this](int peer_fd, int epoll_events) { // Do not miss the disconnection! if (epoll_events & EPOLLRDHUP) diff --git a/src/client/msgr_stop.cpp b/src/client/msgr_stop.cpp index 26394761..68051ad6 100644 --- a/src/client/msgr_stop.cpp +++ b/src/client/msgr_stop.cpp @@ -78,7 +78,7 @@ void osd_messenger_t::stop_client(int peer_fd, bool force, bool force_delete) } #ifndef __MOCK__ // Then remove FD from the eventloop so we don't accidentally read something - tfd->set_fd_handler(peer_fd, false, NULL); + tfd->set_fd_handler(peer_fd, 0, NULL); if (cl->connect_timeout_id >= 0) { tfd->clear_timer(cl->connect_timeout_id); diff --git a/src/client/nbd_proxy.cpp b/src/client/nbd_proxy.cpp index e82a7892..0ddd3fb4 100644 --- a/src/client/nbd_proxy.cpp +++ b/src/client/nbd_proxy.cpp @@ -655,7 +655,7 @@ help: ringloop->register_consumer(&consumer); // Add FD to epoll bool stop = false; - epmgr->tfd->set_fd_handler(sockfd[0], false, [this, &stop](int peer_fd, int epoll_events) + epmgr->tfd->set_fd_handler(sockfd[0], EPOLLIN, [this, &stop](int peer_fd, int epoll_events) { if (epoll_events & EPOLLRDHUP) { diff --git a/src/kv/kv_cli.cpp b/src/kv/kv_cli.cpp index fd3d7d0a..1f194dad 100644 --- a/src/kv/kv_cli.cpp +++ b/src/kv/kv_cli.cpp @@ -185,7 +185,7 @@ void kv_cli_t::run() fcntl(0, F_SETFL, fcntl(0, F_GETFL, 0) | O_NONBLOCK); try { - epmgr->tfd->set_fd_handler(0, false, [this](int fd, int events) + epmgr->tfd->set_fd_handler(0, EPOLLIN, [this](int fd, int events) { if (events & EPOLLIN) { @@ -193,7 +193,7 @@ void kv_cli_t::run() } if (events & EPOLLRDHUP) { - epmgr->tfd->set_fd_handler(0, false, NULL); + epmgr->tfd->set_fd_handler(0, 0, NULL); finished = true; } }); diff --git a/src/nfs/nfs_proxy.cpp b/src/nfs/nfs_proxy.cpp index 65496d45..0ed3475a 100644 --- a/src/nfs/nfs_proxy.cpp +++ b/src/nfs/nfs_proxy.cpp @@ -243,7 +243,7 @@ void nfs_proxy_t::run(json11::Json cfg) // Create NFS socket and add it to epoll int nfs_socket = create_and_bind_socket(bind_address, nfs_port, 128, &listening_port); fcntl(nfs_socket, F_SETFL, fcntl(nfs_socket, F_GETFL, 0) | O_NONBLOCK); - epmgr->tfd->set_fd_handler(nfs_socket, false, [this](int nfs_socket, int epoll_events) + epmgr->tfd->set_fd_handler(nfs_socket, EPOLLIN, [this](int nfs_socket, int epoll_events) { if (epoll_events & EPOLLRDHUP) { @@ -260,7 +260,7 @@ void nfs_proxy_t::run(json11::Json cfg) // Create portmap socket and add it to epoll int portmap_socket = create_and_bind_socket(bind_address, 111, 128, NULL); fcntl(portmap_socket, F_SETFL, fcntl(portmap_socket, F_GETFL, 0) | O_NONBLOCK); - epmgr->tfd->set_fd_handler(portmap_socket, false, [this](int portmap_socket, int epoll_events) + epmgr->tfd->set_fd_handler(portmap_socket, EPOLLIN, [this](int portmap_socket, int epoll_events) { if (epoll_events & EPOLLRDHUP) { @@ -466,7 +466,7 @@ void nfs_proxy_t::do_accept(int listen_fd) { cli->proc_table.insert(fn); } - epmgr->tfd->set_fd_handler(nfs_fd, true, [cli](int nfs_fd, int epoll_events) + epmgr->tfd->set_fd_handler(nfs_fd, EPOLLIN|EPOLLOUT, [cli](int nfs_fd, int epoll_events) { // Handle incoming event if (epoll_events & EPOLLRDHUP) @@ -723,7 +723,7 @@ void nfs_client_t::stop() stopped = true; if (refs <= 0) { - parent->epmgr->tfd->set_fd_handler(nfs_fd, true, NULL); + parent->epmgr->tfd->set_fd_handler(nfs_fd, 0, NULL); close(nfs_fd); delete this; } diff --git a/src/osd/osd.cpp b/src/osd/osd.cpp index b8af5f15..208c4e20 100644 --- a/src/osd/osd.cpp +++ b/src/osd/osd.cpp @@ -361,7 +361,7 @@ void osd_t::bind_socket() listen_fd = create_and_bind_socket(bind_address, bind_port, listen_backlog, &listening_port); fcntl(listen_fd, F_SETFL, fcntl(listen_fd, F_GETFL, 0) | O_NONBLOCK); - epmgr->set_fd_handler(listen_fd, false, [this](int fd, int events) + epmgr->set_fd_handler(listen_fd, EPOLLIN, [this](int fd, int events) { msgr.accept_connections(listen_fd); }); diff --git a/src/test/stub_uring_osd.cpp b/src/test/stub_uring_osd.cpp index cc5c7188..40643297 100644 --- a/src/test/stub_uring_osd.cpp +++ b/src/test/stub_uring_osd.cpp @@ -43,7 +43,7 @@ int main(int narg, char *args[]) // Accept new connections int listen_fd = create_and_bind_socket("0.0.0.0", 11203, 128, NULL); fcntl(listen_fd, F_SETFL, fcntl(listen_fd, F_GETFL, 0) | O_NONBLOCK); - epmgr->set_fd_handler(listen_fd, false, [listen_fd, msgr](int fd, int events) + epmgr->set_fd_handler(listen_fd, EPOLLIN, [listen_fd, msgr](int fd, int events) { msgr->accept_connections(listen_fd); }); diff --git a/src/util/epoll_manager.cpp b/src/util/epoll_manager.cpp index 7dc127ae..d380dd4f 100644 --- a/src/util/epoll_manager.cpp +++ b/src/util/epoll_manager.cpp @@ -21,7 +21,7 @@ epoll_manager_t::epoll_manager_t(ring_loop_t *ringloop) throw std::runtime_error(std::string("epoll_create: ") + strerror(errno)); } - tfd = new timerfd_manager_t([this](int fd, bool wr, std::function handler) { set_fd_handler(fd, wr, handler); }); + tfd = new timerfd_manager_t([this](int fd, int events, std::function handler) { set_fd_handler(fd, events, handler); }); if (ringloop) { @@ -54,14 +54,14 @@ int epoll_manager_t::get_fd() return epoll_fd; } -void epoll_manager_t::set_fd_handler(int fd, bool wr, std::function handler) +void epoll_manager_t::set_fd_handler(int fd, int events, std::function handler) { if (handler != NULL) { bool exists = epoll_handlers.find(fd) != epoll_handlers.end(); epoll_event ev; ev.data.fd = fd; - ev.events = (wr ? EPOLLOUT : 0) | EPOLLIN | EPOLLRDHUP | EPOLLET; + ev.events = events | EPOLLRDHUP | EPOLLET; if (epoll_ctl(epoll_fd, exists ? EPOLL_CTL_MOD : EPOLL_CTL_ADD, fd, &ev) < 0) { if (errno == ENOENT) diff --git a/src/util/epoll_manager.h b/src/util/epoll_manager.h index 2316d329..2ae4d55c 100644 --- a/src/util/epoll_manager.h +++ b/src/util/epoll_manager.h @@ -3,6 +3,8 @@ #pragma once +#include + #include #include "ringloop.h" @@ -21,7 +23,7 @@ public: epoll_manager_t(ring_loop_t *ringloop); ~epoll_manager_t(); int get_fd(); - void set_fd_handler(int fd, bool wr, std::function handler); + void set_fd_handler(int fd, int events, std::function handler); void handle_events(int timeout); timerfd_manager_t *tfd; diff --git a/src/util/timerfd_manager.cpp b/src/util/timerfd_manager.cpp index c46f3f5f..c71bc5ee 100644 --- a/src/util/timerfd_manager.cpp +++ b/src/util/timerfd_manager.cpp @@ -11,7 +11,7 @@ #include #include "timerfd_manager.h" -timerfd_manager_t::timerfd_manager_t(std::function)> set_fd_handler) +timerfd_manager_t::timerfd_manager_t(std::function)> set_fd_handler) { this->set_fd_handler = set_fd_handler; wait_state = 0; @@ -20,7 +20,7 @@ timerfd_manager_t::timerfd_manager_t(std::function)> set_fd_handler; + std::function)> set_fd_handler; - timerfd_manager_t(std::function)> set_fd_handler); + timerfd_manager_t(std::function)> set_fd_handler); ~timerfd_manager_t(); int set_timer(uint64_t millis, bool repeat, std::function callback); int set_timer_us(uint64_t micros, bool repeat, std::function callback);