diff --git a/src/client/msgr_rdma.cpp b/src/client/msgr_rdma.cpp index 09fa497a..d614b1c9 100644 --- a/src/client/msgr_rdma.cpp +++ b/src/client/msgr_rdma.cpp @@ -558,7 +558,7 @@ static void try_send_rdma_wr(osd_client_t *cl, ibv_sge *sge, int op_sge) int osd_messenger_t::try_send_rdma_copy(osd_client_t *cl, uint8_t *dst, int dst_len) { int total_dst_len = dst_len; - while (dst_len > 0 && cl->write_ops.size()) + while (dst_len > 0 && (cl->write_op || cl->write_ops.size())) { if (!cl->write_op) { @@ -752,5 +752,4 @@ void osd_messenger_t::handle_rdma_events(msgr_rdma_context_t *rdma_context) } } } while (event_count > 0); - handle_immediate_ops(); } diff --git a/src/client/msgr_receive.cpp b/src/client/msgr_receive.cpp index 04dd1b67..23b47e70 100644 --- a/src/client/msgr_receive.cpp +++ b/src/client/msgr_receive.cpp @@ -75,6 +75,7 @@ void osd_messenger_t::read_requests() } } read_ready_clients.clear(); + handle_immediate_ops(); } void osd_messenger_t::handle_read(int result, osd_client_t *cl) @@ -112,7 +113,8 @@ void osd_messenger_t::handle_read(int result, osd_client_t *cl) if (!handle_read_buffer(cl, cl->in_buf, result)) { clear_immediate_ops(peer_fd); - handle_immediate_ops(); + if (set_immediate_ops.size()) + ringloop->wakeup(); return; } } @@ -155,7 +157,8 @@ void osd_messenger_t::handle_read(int result, osd_client_t *cl) { read_ready_clients.push_back(cl->peer_fd); } - handle_immediate_ops(); + if (set_immediate_ops.size()) + ringloop->wakeup(); } void osd_messenger_t::clear_immediate_ops(int peer_fd) diff --git a/src/client/msgr_send.cpp b/src/client/msgr_send.cpp index e8c462de..18a6cbd0 100644 --- a/src/client/msgr_send.cpp +++ b/src/client/msgr_send.cpp @@ -49,7 +49,7 @@ void osd_messenger_t::outbox_push(osd_op_t *cur_op) if (!ringloop) { // FIXME: It's worse because it doesn't allow batching - while (cl->write_ops.size()) + while (cl->write_op || cl->write_ops.size()) { try_send(cl); }