diff --git a/src/client/messenger.h b/src/client/messenger.h index c0beb4d2..aa7b0d52 100644 --- a/src/client/messenger.h +++ b/src/client/messenger.h @@ -233,6 +233,7 @@ public: void parse_config(const json11::Json & config); void connect_peer(uint64_t osd_num, json11::Json peer_state); void stop_client(int peer_fd, bool force = false, bool force_delete = false); + void destroy_client(osd_client_t *cl); void outbox_push(osd_op_t *cur_op); std::function exec_op; std::function repeer_pgs; diff --git a/src/client/msgr_receive.cpp b/src/client/msgr_receive.cpp index 804cd778..bd45c4de 100644 --- a/src/client/msgr_receive.cpp +++ b/src/client/msgr_receive.cpp @@ -10,7 +10,7 @@ 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 || - cl_it->second->peer_state == PEER_RDMA || cl_it->second->peer_state == PEER_RDMA_CONNECTING) + cl_it->second->peer_state != PEER_CONNECTED) { continue; } @@ -83,7 +83,7 @@ bool osd_messenger_t::handle_read(int result, osd_client_t *cl) { if (cl->refs <= 0) { - delete cl; + destroy_client(cl); } return false; } diff --git a/src/client/msgr_send.cpp b/src/client/msgr_send.cpp index 5755d5d5..07ed11e2 100644 --- a/src/client/msgr_send.cpp +++ b/src/client/msgr_send.cpp @@ -184,7 +184,7 @@ void osd_messenger_t::measure_exec(osd_op_t *cur_op) bool osd_messenger_t::try_send(osd_client_t *cl) { int peer_fd = cl->peer_fd; - if (!cl->send_list.size() || cl->write_msg.msg_iovlen > 0) + if (!cl->send_list.size() || cl->write_msg.msg_iovlen > 0 || cl->peer_state == PEER_STOPPED) { return true; } @@ -274,7 +274,7 @@ void osd_messenger_t::handle_send(int result, bool prev, bool more, osd_client_t { if (cl->refs <= 0) { - delete cl; + destroy_client(cl); } return; } diff --git a/src/client/msgr_stop.cpp b/src/client/msgr_stop.cpp index c99880b3..e847c34a 100644 --- a/src/client/msgr_stop.cpp +++ b/src/client/msgr_stop.cpp @@ -43,6 +43,10 @@ void osd_op_t::cancel() } } +// force_delete means stop the client anyway, even if there are refs to it in the event loop. +// the flag should be used in the destructor. +// why? - because yes, we could close the FD first and let it fail all requests in the event loop, +// but in that case it can be quickly reopened and we can get old failed responses for the new FD. void osd_messenger_t::stop_client(int peer_fd, bool force, bool force_delete) { assert(peer_fd != 0); @@ -52,9 +56,16 @@ void osd_messenger_t::stop_client(int peer_fd, bool force, bool force_delete) return; } osd_client_t *cl = it->second; - // FIXME: This 'force' flag is probably an ugly reenterability hack - check its logic and maybe remove it + // FIXME "force" flag is required because otherwise a first failed operation + // may stop the client, make it start reconnecting, and then another failed + // operation may stop it again. The right fix would be to introduce unique peer ID + // and not use FDs for that. if (cl->peer_state == PEER_CONNECTING && !force || cl->peer_state == PEER_STOPPED) { + if (force_delete) + { + destroy_client(cl); + } return; } clear_immediate_ops(peer_fd); @@ -98,29 +109,11 @@ void osd_messenger_t::stop_client(int peer_fd, bool force, bool force_delete) } #endif #ifndef __MOCK__ - // Then remove FD from the eventloop so we don't accidentally read something - tfd->set_fd_handler(peer_fd, false, NULL); if (cl->connect_timeout_id >= 0) { tfd->clear_timer(cl->connect_timeout_id); cl->connect_timeout_id = -1; } - for (auto rit = read_ready_clients.begin(); rit != read_ready_clients.end(); rit++) - { - if (*rit == peer_fd) - { - read_ready_clients.erase(rit); - break; - } - } - for (auto wit = write_ready_clients.begin(); wit != write_ready_clients.end(); wit++) - { - if (*wit == peer_fd) - { - write_ready_clients.erase(wit); - break; - } - } #endif if (cl->in_osd_num && break_pg_locks) { @@ -135,17 +128,41 @@ void osd_messenger_t::stop_client(int peer_fd, bool force, bool force_delete) // so do not repeer on it. repeer_pgs(cl->osd_num); } + cl->refs--; + if (cl->refs <= 0 || force_delete) + { + destroy_client(cl); + } +} + +void osd_messenger_t::destroy_client(osd_client_t *cl) +{ // Find the item again because it can be invalidated at this point - it = clients.find(peer_fd); + auto it = clients.find(cl->peer_fd); if (it != clients.end()) { clients.erase(it); } - cl->refs--; - if (cl->refs <= 0 || force_delete) +#ifndef __MOCK__ + tfd->set_fd_handler(cl->peer_fd, false, NULL); + for (auto rit = read_ready_clients.begin(); rit != read_ready_clients.end(); rit++) { - delete cl; + if (*rit == cl->peer_fd) + { + read_ready_clients.erase(rit); + break; + } } + for (auto wit = write_ready_clients.begin(); wit != write_ready_clients.end(); wit++) + { + if (*wit == cl->peer_fd) + { + write_ready_clients.erase(wit); + break; + } + } +#endif + delete cl; } osd_client_t::~osd_client_t()