From c312557acebec0d248aa5e3c11d17137a5d80dee Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Sun, 10 Nov 2024 16:44:13 +0300 Subject: [PATCH] Do not execute remaining operations if the client is stopped during read --- src/client/messenger.h | 4 ++- src/client/msgr_rdma.cpp | 8 ++--- src/client/msgr_receive.cpp | 61 ++++++++++++++++++++++++++++--------- 3 files changed, 52 insertions(+), 21 deletions(-) diff --git a/src/client/messenger.h b/src/client/messenger.h index 3d447cdb..1de228d2 100644 --- a/src/client/messenger.h +++ b/src/client/messenger.h @@ -177,7 +177,7 @@ protected: std::vector read_ready_clients; std::vector write_ready_clients; // We don't use ringloop->set_immediate here because we may have no ringloop in client :) - std::vector> set_immediate; + std::vector set_immediate_ops; public: timerfd_manager_t *tfd; @@ -237,6 +237,8 @@ protected: void handle_op_hdr(osd_client_t *cl); bool handle_reply_hdr(osd_client_t *cl); void handle_reply_ready(osd_op_t *op); + void handle_immediate_ops(); + void clear_immediate_ops(int peer_fd); #ifdef WITH_RDMA void try_send_rdma(osd_client_t *cl); diff --git a/src/client/msgr_rdma.cpp b/src/client/msgr_rdma.cpp index a31905ea..10d2c408 100644 --- a/src/client/msgr_rdma.cpp +++ b/src/client/msgr_rdma.cpp @@ -598,6 +598,7 @@ void osd_messenger_t::handle_rdma_events() } fprintf(stderr, " with status: %s, stopping client\n", ibv_wc_status_str(wc[i].status)); stop_client(client_id); + clear_immediate_ops(client_id); continue; } if (!is_send) @@ -606,6 +607,7 @@ void osd_messenger_t::handle_rdma_events() if (!handle_read_buffer(cl, rc->recv_buffers[rc->next_recv_buf].buf, wc[i].byte_len)) { // handle_read_buffer may stop the client + clear_immediate_ops(client_id); continue; } try_recv_rdma_wr(cl, rc->recv_buffers[rc->next_recv_buf]); @@ -666,9 +668,5 @@ void osd_messenger_t::handle_rdma_events() } } } while (event_count > 0); - for (auto cb: set_immediate) - { - cb(); - } - set_immediate.clear(); + handle_immediate_ops(); } diff --git a/src/client/msgr_receive.cpp b/src/client/msgr_receive.cpp index a43729cd..45807921 100644 --- a/src/client/msgr_receive.cpp +++ b/src/client/msgr_receive.cpp @@ -65,6 +65,7 @@ void osd_messenger_t::read_requests() bool osd_messenger_t::handle_read(int result, osd_client_t *cl) { bool ret = false; + int peer_fd = cl->peer_fd; cl->read_msg.msg_iovlen = 0; cl->refs--; if (cl->peer_state == PEER_STOPPED) @@ -101,7 +102,8 @@ bool osd_messenger_t::handle_read(int result, osd_client_t *cl) { if (!handle_read_buffer(cl, cl->in_buf, result)) { - goto fin; + clear_immediate_ops(peer_fd); + return false; } } else @@ -113,7 +115,8 @@ bool osd_messenger_t::handle_read(int result, osd_client_t *cl) { if (!handle_finished_read(cl)) { - goto fin; + clear_immediate_ops(peer_fd); + return false; } } } @@ -122,15 +125,47 @@ bool osd_messenger_t::handle_read(int result, osd_client_t *cl) ret = true; } } -fin: - for (auto cb: set_immediate) - { - cb(); - } - set_immediate.clear(); + handle_immediate_ops(); return ret; } +void osd_messenger_t::clear_immediate_ops(int peer_fd) +{ + size_t i = 0, j = 0; + while (i < set_immediate_ops.size()) + { + if (set_immediate_ops[i]->peer_fd == peer_fd) + { + delete set_immediate_ops[i]; + } + else + { + if (i != j) + set_immediate_ops[j] = set_immediate_ops[i]; + j++; + } + i++; + } + set_immediate_ops.resize(j); +} + +void osd_messenger_t::handle_immediate_ops() +{ + for (auto op: set_immediate_ops) + { + if (op->op_type == OSD_OP_IN) + { + exec_op(op); + } + else + { + // Copy lambda to be unaffected by `delete op` + std::function(op->callback)(op); + } + } + set_immediate_ops.clear(); +} + bool osd_messenger_t::handle_read_buffer(osd_client_t *cl, void *curbuf, int remain) { // Compose operation(s) from the buffer @@ -199,7 +234,7 @@ bool osd_messenger_t::handle_finished_read(osd_client_t *cl) { // Operation is ready cl->received_ops.push_back(cl->read_op); - set_immediate.push_back([this, op = cl->read_op]() { exec_op(op); }); + set_immediate_ops.push_back(cl->read_op); cl->read_op = NULL; cl->read_state = 0; } @@ -295,7 +330,7 @@ void osd_messenger_t::handle_op_hdr(osd_client_t *cl) { // Operation is ready cl->received_ops.push_back(cur_op); - set_immediate.push_back([this, cur_op]() { exec_op(cur_op); }); + set_immediate_ops.push_back(cur_op); cl->read_op = NULL; cl->read_state = 0; } @@ -416,9 +451,5 @@ void osd_messenger_t::handle_reply_ready(osd_op_t *op) (tv_end.tv_sec - op->tv_begin.tv_sec)*1000000 + (tv_end.tv_nsec - op->tv_begin.tv_nsec)/1000 ); - set_immediate.push_back([op]() - { - // Copy lambda to be unaffected by `delete op` - std::function(op->callback)(op); - }); + set_immediate_ops.push_back(op); }