diff --git a/src/client/messenger.cpp b/src/client/messenger.cpp index 55b3a42c..20f04df8 100644 --- a/src/client/messenger.cpp +++ b/src/client/messenger.cpp @@ -739,14 +739,6 @@ 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) - { - // Do not miss the disconnection! - if (epoll_events & EPOLLRDHUP) - { - handle_peer_epoll(peer_fd, epoll_events); - } - }); // Add the initial receive request init_recv_rdma(cl); } diff --git a/src/client/msgr_rdma.cpp b/src/client/msgr_rdma.cpp index 68e09464..13eb2464 100644 --- a/src/client/msgr_rdma.cpp +++ b/src/client/msgr_rdma.cpp @@ -696,6 +696,10 @@ void osd_messenger_t::handle_rdma_events(msgr_rdma_context_t *rdma_context) continue; } osd_client_t *cl = cl_it->second; + if (cl->peer_state == PEER_STOPPED) + { + continue; + } auto rc = cl->rdma_conn; if (wc[i].status != IBV_WC_SUCCESS) { diff --git a/src/client/msgr_receive.cpp b/src/client/msgr_receive.cpp index 60e9b9f0..0f2ef30b 100644 --- a/src/client/msgr_receive.cpp +++ b/src/client/msgr_receive.cpp @@ -9,7 +9,8 @@ void osd_messenger_t::read_requests() { int peer_fd = read_ready_clients[i]; auto cl_it = clients.find(peer_fd); - if (cl_it == clients.end() || !cl_it->second || cl_it->second->read_msg.msg_iovlen) + if (cl_it == clients.end() || !cl_it->second || cl_it->second->read_msg.msg_iovlen || + cl_it->second->peer_state == PEER_RDMA || cl_it->second->peer_state == PEER_RDMA_CONNECTING) { continue; } @@ -75,6 +76,10 @@ bool osd_messenger_t::handle_read(int result, osd_client_t *cl) int peer_fd = cl->peer_fd; cl->read_msg.msg_iovlen = 0; cl->refs--; + if (cl->peer_state == PEER_RDMA) + { + return true; + } if (cl->peer_state == PEER_STOPPED) { if (cl->refs <= 0) diff --git a/src/client/msgr_send.cpp b/src/client/msgr_send.cpp index 99ca9bc8..5755d5d5 100644 --- a/src/client/msgr_send.cpp +++ b/src/client/msgr_send.cpp @@ -251,7 +251,7 @@ void osd_messenger_t::send_replies() { int peer_fd = write_ready_clients[i]; auto cl_it = clients.find(peer_fd); - if (cl_it != clients.end() && !try_send(cl_it->second)) + if (cl_it != clients.end() && cl_it->second->peer_state != PEER_RDMA && !try_send(cl_it->second)) { write_ready_clients.erase(write_ready_clients.begin(), write_ready_clients.begin() + i); return; @@ -349,21 +349,12 @@ void osd_messenger_t::handle_send(int result, bool prev, bool more, osd_client_t #ifdef WITH_RDMA if (cl->rdma_conn && !cl->outbox.size() && cl->peer_state == PEER_RDMA_CONNECTING) { - // FIXME: Do something better than just forgetting the FD // FIXME: Ignore pings during RDMA state transition if (log_level > 0) { 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) - { - // Do not miss the disconnection! - if (epoll_events & EPOLLRDHUP) - { - handle_peer_epoll(peer_fd, epoll_events); - } - }); // Add the initial receive request init_recv_rdma(cl); }