// Copyright (c) Vitaliy Filippov, 2019+ // License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) #define _XOPEN_SOURCE #include #include "messenger.h" void osd_messenger_t::read_requests() { for (int i = 0; i < read_ready_clients.size(); i++) { uint64_t client_id = read_ready_clients[i]; auto cl_it = clients.find(client_id); if (cl_it == clients.end() || !cl_it->second || cl_it->second->read_msg.msg_iovlen || cl_it->second->peer_state != PEER_CONNECTED) { continue; } auto cl = cl_it->second; if (cl->read_op && cl->read_op_size-(cl->read_op_pos-OSD_PACKET_SIZE) >= receive_buffer_size) { op_get_read_buffers(cl, cl->recv_list); } if (!cl->recv_list.size()) { cl->read_iov.iov_base = cl->in_buf; cl->read_iov.iov_len = receive_buffer_size; cl->read_msg.msg_iov = &cl->read_iov; cl->read_msg.msg_iovlen = 1; } else { cl->read_iov.iov_base = 0; cl->read_iov.iov_len = 0; cl->read_msg.msg_iov = cl->recv_list.data(); cl->read_msg.msg_iovlen = cl->recv_list.size(); } assert(!cl->read_op || cl->read_op_pos < OSD_PACKET_SIZE || cl->read_op_size >= (cl->read_op_pos-OSD_PACKET_SIZE)); cl->refs++; if (ringloop && !use_sync_send_recv) { auto iothread = iothreads.size() ? iothreads[cl->peer_fd % iothreads.size()] : NULL; io_uring_sqe sqe_local; ring_data_t data_local; io_uring_sqe* sqe = (iothread ? &sqe_local : ringloop->get_sqe()); if (iothread) { sqe_local = { .user_data = (uint64_t)&data_local }; data_local = {}; } if (!sqe) { cl->refs--; cl->read_msg.msg_iovlen = 0; read_ready_clients.erase(read_ready_clients.begin(), read_ready_clients.begin() + i); return; } ring_data_t* data = ((ring_data_t*)sqe->user_data); data->callback = [this, cl](ring_data_t *data) { handle_read(data->res, cl); }; io_uring_prep_recvmsg(sqe, cl->peer_fd, &cl->read_msg, cl->recv_list.size() ? MSG_WAITALL : 0); if (iothread) { iothread->add_sqe(sqe_local); } } else { int result = recvmsg(cl->peer_fd, &cl->read_msg, 0); if (result < 0) { result = -errno; } // like set_immediate tfd->set_timer_us(0, false, [this, result, cl](int){ handle_read(result, cl); }); } } read_ready_clients.clear(); handle_immediate_ops(); } void osd_messenger_t::handle_read(int result, osd_client_t *cl) { cl->refs--; if (cl->peer_state == PEER_RDMA) { return; } if (cl->peer_state == PEER_STOPPED) { if (cl->refs <= 0) { destroy_client(cl); } return; } if (result <= 0 && result != -EAGAIN && result != -EINTR) { // this is a client socket, so don't panic on error. just disconnect it if (result != 0) { fprintf(stderr, "Client %ju socket read error: %d (%s). Disconnecting client\n", cl->client_id, -result, strerror(-result)); } stop_client(cl->client_id); return; } bool full_read = false; if (result > 0) { if (cl->read_iov.iov_base == cl->in_buf) { full_read = result >= cl->read_iov.iov_len; if (!handle_read_buffer(cl, cl->in_buf, result)) { if (set_immediate_ops.size()) ringloop->wakeup(); return; } } else { // Reset OSD ping state cl->ping_time_remaining = 0; cl->idle_time_remaining = osd_idle_timeout; // Long data size_t i = 0; while (i < cl->recv_list.size() && result >= cl->recv_list[i].iov_len) { result -= cl->recv_list[i].iov_len; i++; } if (i < cl->recv_list.size()) { cl->recv_list[i].iov_base += result; cl->recv_list[i].iov_len -= result; } else { full_read = true; } cl->recv_list.erase(cl->recv_list.begin(), cl->recv_list.begin()+i); if (!cl->recv_list.size()) { handle_finished_op(cl); } } } cl->read_msg.msg_iovlen = 0; if (result == -EAGAIN || result == -EINTR || !full_read) { cl->read_ready--; if (cl->read_ready > 0) read_ready_clients.push_back(cl->client_id); } else { read_ready_clients.push_back(cl->client_id); } if (set_immediate_ops.size()) ringloop->wakeup(); } void osd_messenger_t::handle_immediate_ops() { while (set_immediate_ops.size()) { auto op = set_immediate_ops.front(); set_immediate_ops.pop_front(); if (op->op_type == OSD_OP_IN) { auto cl_it = clients.find(op->client_id); if (cl_it != clients.end() && cl_it->second->peer_state != PEER_STOPPED) exec_op(op); else delete op; } else { // Copy lambda to be unaffected by `delete op` std::function(op->callback)(op); } } } bool osd_messenger_t::handle_read_buffer(osd_client_t *cl, uint8_t *curbuf, size_t bufsize) { // Reset OSD ping state cl->ping_time_remaining = 0; cl->idle_time_remaining = osd_idle_timeout; // Compose operation(s) from the buffer size_t done = 0; while (done < bufsize) { if (!cl->read_op) { cl->read_op = new osd_op_t; cl->read_op->client_id = cl->client_id; cl->read_op->op_type = OSD_OP_IN; cl->read_op_pos = 0; cl->read_op_size = 0; cl->read_op_inline_decrypt_in = 0; cl->read_op_inline_decrypt_pos = (size_t)-1; } if (cl->read_op_pos < OSD_PACKET_SIZE) { int len = OSD_PACKET_SIZE - cl->read_op_pos; if (len > bufsize-done) len = bufsize-done; memcpy(cl->read_op->req.buf + cl->read_op_pos, curbuf+done, len); done += len; cl->read_op_pos += len; if (cl->read_op_pos < OSD_PACKET_SIZE) return true; if (!handle_hdr(cl)) { stop_client(cl->client_id); return false; } } op_copy_from(cl, curbuf, bufsize, done); } return true; } bool osd_messenger_t::handle_hdr(osd_client_t *cl) { if (cl->read_op->req.hdr.magic == SECONDARY_OSD_REPLY_MAGIC) { auto req_it = cl->sent_ops.find(cl->read_op->req.hdr.id); if (req_it == cl->sent_ops.end()) { // Command out of sync. Drop connection fprintf(stderr, "Client %ju command out of sync: id %ju\n", cl->client_id, cl->read_op->req.hdr.id); return false; } osd_op_t *op = req_it->second; memcpy(op->reply.buf, cl->read_op->req.buf, OSD_PACKET_SIZE); if (!allocate_reply_buffers(cl, op)) { return false; } cl->sent_ops.erase(req_it); delete cl->read_op; cl->read_op = op; } else if (cl->read_op->req.hdr.magic == SECONDARY_OSD_OP_MAGIC) { if (cl->check_sequencing) { if (cl->read_op->req.hdr.id != cl->read_op_id) { fprintf(stderr, "Warning: operation sequencing is broken on client %d: expected num %ju, got %ju, stopping client\n", cl->peer_fd, cl->read_op_id, cl->read_op->req.hdr.id); return false; } cl->read_op_id++; } if (!allocate_op_buffers(cl)) { return false; } } else { fprintf(stderr, "Received garbage: magic=%jx id=%ju opcode=%jx from client %ju\n", cl->read_op->req.hdr.magic, cl->read_op->req.hdr.id, cl->read_op->req.hdr.opcode, cl->client_id); return false; } return true; } bool osd_messenger_t::allocate_op_buffers(osd_client_t *cl) { osd_op_t *cur_op = cl->read_op; cl->read_op_size = 0; if (cur_op->req.hdr.opcode == OSD_OP_SEC_WRITE || cur_op->req.hdr.opcode == OSD_OP_SEC_WRITE_STABLE) { if (cur_op->req.sec_rw.attr_len > 0) { if (cur_op->req.sec_rw.attr_len > sizeof(unsigned)) cur_op->bitmap = cur_op->rmw_buf = malloc_or_die(cur_op->req.sec_rw.attr_len); else cur_op->bitmap = &cur_op->bmp_data; } if (cur_op->req.sec_rw.len > 0) { cur_op->buf = memalign_or_die(MEM_ALIGNMENT, cur_op->req.sec_rw.len); } cl->read_op_size = cur_op->req.sec_rw.len + cur_op->req.sec_rw.attr_len; } else if (cur_op->req.hdr.opcode == OSD_OP_SEC_STABILIZE || cur_op->req.hdr.opcode == OSD_OP_SEC_ROLLBACK) { if (cur_op->req.sec_stab.len > 0) { cur_op->buf = memalign_or_die(MEM_ALIGNMENT, cur_op->req.sec_stab.len); } cl->read_op_size = cur_op->req.sec_stab.len; } else if (cur_op->req.hdr.opcode == OSD_OP_SEC_READ_BMP) { if (cur_op->req.sec_read_bmp.len > 0) { cur_op->buf = memalign_or_die(MEM_ALIGNMENT, cur_op->req.sec_read_bmp.len); } cl->read_op_size = cur_op->req.sec_read_bmp.len; } else if (cur_op->req.hdr.opcode == OSD_OP_WRITE) { if (cur_op->req.rw.len > 0) { cur_op->buf = memalign_or_die(MEM_ALIGNMENT, cur_op->req.rw.len); } cl->read_op_size = cur_op->req.rw.len; } else if (cur_op->req.hdr.opcode == OSD_OP_SHOW_CONFIG) { if (cur_op->req.show_conf.json_len > 0) { cur_op->buf = malloc_or_die(cur_op->req.show_conf.json_len+1); ((uint8_t*)cur_op->buf)[cur_op->req.show_conf.json_len] = 0; } cl->read_op_size = cur_op->req.show_conf.json_len; } return true; } bool osd_messenger_t::allocate_reply_buffers(osd_client_t *cl, osd_op_t *op) { cl->read_op_size = 0; if (op->reply.hdr.opcode == OSD_OP_SEC_READ || op->reply.hdr.opcode == OSD_OP_READ) { // Read data. In this case we assume that the buffer is preallocated by the caller (!) unsigned bmp_len = (op->reply.hdr.opcode == OSD_OP_SEC_READ ? op->reply.sec_rw.attr_len : op->reply.rw.bitmap_len); unsigned expected_size = (op->reply.hdr.opcode == OSD_OP_SEC_READ ? op->req.sec_rw.len : op->req.rw.len); if (op->reply.hdr.retval >= 0 && (op->reply.hdr.retval != expected_size || bmp_len > op->bitmap_len)) { // Check reply length to not overflow the buffer fprintf(stderr, "Client %ju read reply of different length: expected %u+%u, got %jd+%u\n", cl->client_id, expected_size, op->bitmap_len, op->reply.hdr.retval, bmp_len); return false; } if (bmp_len > 0) { assert(op->bitmap); cl->read_op_size += bmp_len; } if (op->reply.hdr.retval > 0) { assert(op->iov.count > 0); cl->read_op_size += op->reply.hdr.retval; } } else if (op->reply.hdr.opcode == OSD_OP_SEC_LIST && op->reply.hdr.retval > 0) { assert(!op->iov.count); cl->read_op_size = sizeof(obj_ver_id) * op->reply.hdr.retval; op->buf = memalign_or_die(MEM_ALIGNMENT, cl->read_op_size); } else if (op->reply.hdr.opcode == OSD_OP_SEC_READ_BMP && op->reply.hdr.retval > 0) { assert(!op->iov.count); cl->read_op_size = op->reply.hdr.retval; free(op->buf); op->buf = memalign_or_die(MEM_ALIGNMENT, cl->read_op_size); } else if (op->reply.hdr.opcode == OSD_OP_SHOW_CONFIG && op->reply.hdr.retval > 0) { cl->read_op_size = op->reply.hdr.retval; free(op->buf); op->buf = malloc_or_die(op->reply.hdr.retval); } else if (op->reply.hdr.opcode == OSD_OP_DESCRIBE && op->reply.describe.result_bytes > 0) { cl->read_op_size = op->reply.describe.result_bytes; free(op->buf); op->buf = malloc_or_die(op->reply.describe.result_bytes); } return true; } size_t osd_messenger_t::op_copy_from(osd_client_t *cl, uint8_t *src, size_t src_len, size_t & done) { osd_op_t *op = cl->read_op; size_t from = cl->read_op_pos-OSD_PACKET_SIZE; auto op_read_buf = [&](uint8_t *dst, size_t dst_len) { if (from < dst_len) { size_t n = dst_len-from; if (n > src_len-done) n = src_len-done; memcpy(dst+from, src+done, n); done += n; cl->read_op_pos += n; from += n; if (from < dst_len) return false; from = 0; } else from -= dst_len; return true; }; if (op->op_type == OSD_OP_IN) { if (op->req.hdr.opcode == OSD_OP_SEC_WRITE || op->req.hdr.opcode == OSD_OP_SEC_WRITE_STABLE) { if (!op_read_buf((uint8_t*)op->bitmap, op->req.sec_rw.attr_len)) return done; if (!op_read_buf((uint8_t*)op->buf, op->req.sec_rw.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_SEC_STABILIZE || op->req.hdr.opcode == OSD_OP_SEC_ROLLBACK) { if (!op_read_buf((uint8_t*)op->buf, op->req.sec_stab.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_SEC_READ_BMP) { if (!op_read_buf((uint8_t*)op->buf, op->req.sec_read_bmp.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_WRITE) { if (!op_read_buf((uint8_t*)op->buf, op->req.rw.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_SHOW_CONFIG) { if (!op_read_buf((uint8_t*)op->buf, op->req.show_conf.json_len)) return done; } } else { if (op->reply.hdr.opcode == OSD_OP_SEC_READ) { if (op->reply.sec_rw.attr_len > 0) { if (!op_read_buf((uint8_t*)op->bitmap, op->reply.sec_rw.attr_len)) return done; } if (op->reply.hdr.retval > 0) { for (int i = 0; i < op->iov.count; i++) if (!op_read_buf((uint8_t*)op->iov.buf[i].iov_base, op->iov.buf[i].iov_len)) return done; } } else if (op->reply.hdr.opcode == OSD_OP_READ) { if (op->reply.rw.bitmap_len > 0) { if (!op_read_buf((uint8_t*)op->bitmap, op->reply.rw.bitmap_len)) return done; } if (op->reply.hdr.retval > 0) { if (op->enc) { if (!op_decrypted_copy_data_from(cl, src, src_len, from, done)) return done; } else { for (int i = 0; i < op->iov.count; i++) if (!op_read_buf((uint8_t*)op->iov.buf[i].iov_base, op->iov.buf[i].iov_len)) return done; } } } else if (op->reply.hdr.opcode == OSD_OP_SEC_LIST && op->reply.hdr.retval > 0) { if (!op_read_buf((uint8_t*)op->buf, sizeof(obj_ver_id) * op->reply.hdr.retval)) return done; } else if ((op->reply.hdr.opcode == OSD_OP_SEC_READ_BMP || op->reply.hdr.opcode == OSD_OP_SHOW_CONFIG) && op->reply.hdr.retval > 0) { if (!op_read_buf((uint8_t*)op->buf, op->reply.hdr.retval)) return done; } else if (op->reply.hdr.opcode == OSD_OP_DESCRIBE && op->reply.describe.result_bytes > 0) { if (!op_read_buf((uint8_t*)op->buf, op->reply.describe.result_bytes)) return done; } } handle_finished_op(cl); return done; } size_t osd_messenger_t::op_get_read_buffers(osd_client_t *cl, std::vector & lst) { osd_op_t *op = cl->read_op; size_t from = cl->read_op_pos-OSD_PACKET_SIZE; size_t done = 0; auto op_read_buf = [&](uint8_t *dst, size_t dst_len) { if (lst.size() >= IOV_MAX) return false; if (from < dst_len) { lst.push_back((iovec){ .iov_base = dst+from, .iov_len = dst_len-from }); cl->read_op_pos += dst_len-from; done += dst_len-from; from = 0; } else from -= dst_len; return true; }; if (op->op_type == OSD_OP_IN) { if (op->req.hdr.opcode == OSD_OP_SEC_WRITE || op->req.hdr.opcode == OSD_OP_SEC_WRITE_STABLE) { if (!op_read_buf((uint8_t*)op->bitmap, op->req.sec_rw.attr_len)) return done; if (!op_read_buf((uint8_t*)op->buf, op->req.sec_rw.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_SEC_STABILIZE || op->req.hdr.opcode == OSD_OP_SEC_ROLLBACK) { if (!op_read_buf((uint8_t*)op->buf, op->req.sec_stab.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_SEC_READ_BMP) { if (!op_read_buf((uint8_t*)op->buf, op->req.sec_read_bmp.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_WRITE) { if (!op_read_buf((uint8_t*)op->buf, op->req.rw.len)) return done; } else if (op->req.hdr.opcode == OSD_OP_SHOW_CONFIG) { if (!op_read_buf((uint8_t*)op->buf, op->req.show_conf.json_len)) return done; } } else { if (op->reply.hdr.opcode == OSD_OP_SEC_READ) { if (op->reply.sec_rw.attr_len > 0) { if (!op_read_buf((uint8_t*)op->bitmap, op->reply.sec_rw.attr_len)) return done; } if (op->reply.hdr.retval > 0) { for (int i = 0; i < op->iov.count; i++) if (!op_read_buf((uint8_t*)op->iov.buf[i].iov_base, op->iov.buf[i].iov_len)) return done; } } else if (op->reply.hdr.opcode == OSD_OP_READ) { if (op->reply.rw.bitmap_len > 0) { if (!op_read_buf((uint8_t*)op->bitmap, op->reply.rw.bitmap_len)) return done; } if (op->reply.hdr.retval > 0) { if (op->enc) { cl->read_op_inline_decrypt_pos = cl->read_op_pos; cl->read_op_pos = cl->read_op_inline_decrypt_in + OSD_PACKET_SIZE + op->reply.rw.bitmap_len; from = cl->read_op_inline_decrypt_in; } for (int i = 0; i < op->iov.count; i++) if (!op_read_buf((uint8_t*)op->iov.buf[i].iov_base, op->iov.buf[i].iov_len)) return done; } } else if (op->reply.hdr.opcode == OSD_OP_SEC_LIST && op->reply.hdr.retval > 0) { if (!op_read_buf((uint8_t*)op->buf, sizeof(obj_ver_id) * op->reply.hdr.retval)) return done; } else if ((op->reply.hdr.opcode == OSD_OP_SEC_READ_BMP || op->reply.hdr.opcode == OSD_OP_SHOW_CONFIG) && op->reply.hdr.retval > 0) { if (!op_read_buf((uint8_t*)op->buf, op->reply.hdr.retval)) return done; } else if (op->reply.hdr.opcode == OSD_OP_DESCRIBE && op->reply.describe.result_bytes > 0) { if (!op_read_buf((uint8_t*)op->buf, op->reply.describe.result_bytes)) return done; } } return done; } void osd_messenger_t::handle_finished_op(osd_client_t *cl) { osd_op_t *op = cl->read_op; if (op->op_type == OSD_OP_IN) { // Operation is ready cl->received_ops.push_back(op); } else { // Inline decryption if (cl->read_op_inline_decrypt_pos != (size_t)-1) { op_decrypt_inline(cl); cl->read_op_inline_decrypt_pos = (size_t)-1; } // Measure subop (outbound op) latency timespec tv_end; clock_gettime(CLOCK_REALTIME, &tv_end); stats.subop_stat_count[op->req.hdr.opcode]++; if (!stats.subop_stat_count[op->req.hdr.opcode]) { stats.subop_stat_count[op->req.hdr.opcode]++; stats.subop_stat_sum[op->req.hdr.opcode] = 0; } stats.subop_stat_sum[op->req.hdr.opcode] += ( (tv_end.tv_sec - op->tv_begin.tv_sec)*1000000 + (tv_end.tv_nsec - op->tv_begin.tv_nsec)/1000 ); } set_immediate_ops.push_back(op); cl->read_op = NULL; }