From 1dbbb0c3f8eec67f297e562ebfd56df44d0d6bf6 Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Mon, 4 Nov 2024 18:58:46 +0300 Subject: [PATCH] Implement NFS RDMA support --- src/CMakeLists.txt | 4 + src/nfs/CMakeLists.txt | 4 + src/nfs/nfs_block.cpp | 3 +- src/nfs/nfs_kv_lookup.cpp | 6 +- src/nfs/nfs_kv_read.cpp | 11 +- src/nfs/nfs_proxy.cpp | 287 ++++++--- src/nfs/nfs_proxy.h | 39 +- src/nfs/nfs_proxy_rdma.cpp | 1072 +++++++++++++++++++++++++++++++ src/nfs/proto/nfs.x | 8 +- src/nfs/proto/nfs_xdr.cpp | 12 +- src/nfs/proto/nfs_xdr.cpp.diff | 53 ++ src/nfs/proto/rpc_impl.h | 7 +- src/nfs/proto/rpc_rdma.h | 144 +++++ src/nfs/proto/rpc_rdma.x | 166 +++++ src/nfs/proto/rpc_rdma_xdr.cpp | 200 ++++++ src/nfs/proto/run-rpcgen.sh | 2 + src/nfs/proto/xdr_impl.cpp | 45 ++ src/nfs/proto/xdr_impl.h | 15 + src/nfs/proto/xdr_impl_inline.h | 43 +- src/nfs/rdma_alloc.cpp | 216 +++++++ src/nfs/rdma_alloc.h | 17 + 21 files changed, 2247 insertions(+), 107 deletions(-) create mode 100644 src/nfs/nfs_proxy_rdma.cpp create mode 100644 src/nfs/proto/nfs_xdr.cpp.diff create mode 100644 src/nfs/proto/rpc_rdma.h create mode 100644 src/nfs/proto/rpc_rdma.x create mode 100644 src/nfs/proto/rpc_rdma_xdr.cpp create mode 100644 src/nfs/rdma_alloc.cpp create mode 100644 src/nfs/rdma_alloc.h diff --git a/src/CMakeLists.txt b/src/CMakeLists.txt index 329dedf7..310cad57 100644 --- a/src/CMakeLists.txt +++ b/src/CMakeLists.txt @@ -61,6 +61,10 @@ pkg_check_modules(ISAL libisal) if (ISAL_LIBRARIES) add_definitions(-DWITH_ISAL) endif (ISAL_LIBRARIES) +pkg_check_modules(RDMACM librdmacm) +if (RDMACM_LIBRARIES) + add_definitions(-DWITH_RDMACM) +endif (RDMACM_LIBRARIES) add_custom_target(build_tests) add_custom_target(test diff --git a/src/nfs/CMakeLists.txt b/src/nfs/CMakeLists.txt index 42891c39..df1b4ba3 100644 --- a/src/nfs/CMakeLists.txt +++ b/src/nfs/CMakeLists.txt @@ -5,6 +5,7 @@ project(vitastor) # vitastor-nfs add_executable(vitastor-nfs nfs_proxy.cpp + nfs_proxy_rdma.cpp nfs_block.cpp nfs_kv.cpp nfs_kv_create.cpp @@ -21,8 +22,10 @@ add_executable(vitastor-nfs nfs_fsstat.cpp nfs_mount.cpp nfs_portmap.cpp + rdma_alloc.cpp ../util/sha256.c proto/xdr_impl.cpp + proto/rpc_rdma_xdr.cpp proto/rpc_xdr.cpp proto/portmap_xdr.cpp proto/nfs_xdr.cpp @@ -30,4 +33,5 @@ add_executable(vitastor-nfs target_link_libraries(vitastor-nfs vitastor_client vitastor_kv + ${RDMACM_LIBRARIES} ) diff --git a/src/nfs/nfs_block.cpp b/src/nfs/nfs_block.cpp index dae9aaca..245b5604 100644 --- a/src/nfs/nfs_block.cpp +++ b/src/nfs/nfs_block.cpp @@ -315,8 +315,7 @@ static int block_nfs3_read_proc(void *opaque, rpc_op_t *rop) if (aligned_count % alignment) aligned_count = aligned_count + alignment - (aligned_count % alignment); aligned_count -= aligned_offset; - void *buf = malloc_or_die(aligned_count); - xdr_add_malloc(rop->xdrs, buf); + void *buf = self->malloc_or_rdma(rop, aligned_count); cluster_op_t *op = new cluster_op_t; op->opcode = OSD_OP_READ; op->inode = ino_it->second; diff --git a/src/nfs/nfs_kv_lookup.cpp b/src/nfs/nfs_kv_lookup.cpp index fab097e8..9636840b 100644 --- a/src/nfs/nfs_kv_lookup.cpp +++ b/src/nfs/nfs_kv_lookup.cpp @@ -91,10 +91,14 @@ int kv_nfs3_readlink_proc(void *opaque, rpc_op_t *rop) } else { + std::string link_target = attrs["symlink"].string_value(); + char *cp = (char*)self->malloc_or_rdma(rop, link_target.size()+1); + memcpy(cp, link_target.data(), link_target.size()); + cp[link_target.size()] = 0; *reply = (READLINK3res){ .status = NFS3_OK, .resok = (READLINK3resok){ - .data = xdr_copy_string(rop->xdrs, attrs["symlink"].string_value()), + .data = (xdr_string_t){ link_target.size(), cp }, }, }; } diff --git a/src/nfs/nfs_kv_read.cpp b/src/nfs/nfs_kv_read.cpp index c0f3199f..78eae9fe 100644 --- a/src/nfs/nfs_kv_read.cpp +++ b/src/nfs/nfs_kv_read.cpp @@ -96,7 +96,7 @@ resume_1: } read_size += sizeof(shared_file_header_t); assert(!st->aligned_buf); - st->aligned_buf = (uint8_t*)malloc_or_die(read_size); + st->aligned_buf = (uint8_t*)st->self->malloc_or_rdma(st->rop, read_size); st->buf = st->aligned_buf + sizeof(shared_file_header_t) + st->offset; st->op->iov.push_back(st->aligned_buf, read_size); st->op->len = align_up(read_offset+read_size) - st->op->offset; @@ -117,7 +117,7 @@ resume_1: resume_2: if (st->res < 0) { - free(st->aligned_buf); + st->self->free_or_rdma(st->rop, st->aligned_buf); st->aligned_buf = NULL; auto cb = std::move(st->cb); cb(st->res); @@ -131,7 +131,7 @@ resume_2: " 0x%jx offset 0x%jx: probably a read/write conflict, retrying\n", st->ino, st->ientry["shared_ino"].uint64_value(), st->ientry["shared_offset"].uint64_value()); st->retry++; - free(st->aligned_buf); + st->self->free_or_rdma(st->rop, st->aligned_buf); st->aligned_buf = NULL; st->allow_cache = false; goto resume_0; @@ -144,7 +144,7 @@ resume_2: st->aligned_offset = align_down(st->offset); st->aligned_size = align_up(st->offset+st->size) - st->aligned_offset; assert(!st->aligned_buf); - st->aligned_buf = (uint8_t*)malloc_or_die(st->aligned_size); + st->aligned_buf = (uint8_t*)st->self->malloc_or_rdma(st->rop, st->aligned_size); st->buf = st->aligned_buf + st->offset - st->aligned_offset; st->op = new cluster_op_t; st->op->opcode = OSD_OP_READ; @@ -163,7 +163,7 @@ resume_2: resume_3: if (st->res < 0) { - free(st->aligned_buf); + st->self->free_or_rdma(st->rop, st->aligned_buf); st->aligned_buf = NULL; } auto cb = std::move(st->cb); @@ -194,7 +194,6 @@ int kv_nfs3_read_proc(void *opaque, rpc_op_t *rop) *reply = (READ3res){ .status = vitastor_nfs_map_err(res) }; if (res == 0) { - xdr_add_malloc(st->rop->xdrs, st->aligned_buf); reply->resok.data.data = (char*)st->buf; reply->resok.data.size = st->size; reply->resok.count = st->size; diff --git a/src/nfs/nfs_proxy.cpp b/src/nfs/nfs_proxy.cpp index 25d99e61..16130375 100644 --- a/src/nfs/nfs_proxy.cpp +++ b/src/nfs/nfs_proxy.cpp @@ -34,6 +34,9 @@ const char *exe_name = NULL; nfs_proxy_t::~nfs_proxy_t() { +#ifdef WITH_RDMACM + destroy_rdma(); +#endif if (kvfs) delete kvfs; if (blockfs) @@ -65,9 +68,14 @@ static const char* help_text = "\n" "vitastor-nfs (--fs | --block) start\n" " Start network NFS server. Options:\n" - " --bind bind service to address (default 0.0.0.0)\n" - " --port use port for NFS services (default is 2049)\n" - " --portmap 0 do not listen on port 111 (portmap/rpcbind, requires root)\n" + " --bind bind service to address (default 0.0.0.0)\n" + " --port use port for NFS services (default is 2049)\n" + " --portmap 0 do not listen on port 111 (portmap/rpcbind, requires root)\n" + " --nfs_rdma enable NFS-RDMA at RDMA-CM port (you can try 20049)\n" + " --nfs_rdma_credit 16 maximum operation credit for RDMA clients (max iodepth)\n" + " --nfs_rdma_send 1024 maximum RDMA send operation count (should be larger than iodepth)\n" + " --nfs_rdma_alloc 1M RDMA memory allocation rounding\n" + " --nfs_rdma_gc 500M maximum unused RDMA buffers\n" "\n" "vitastor-nfs --fs upgrade\n" " Upgrade FS metadata. Can be run online, but server(s) should be restarted\n" @@ -184,6 +192,7 @@ void nfs_proxy_t::run(json11::Json cfg) srand48(tv.tv_sec*1000000000 + tv.tv_nsec); server_id = (uint64_t)lrand48() | ((uint64_t)lrand48() << 31) | ((uint64_t)lrand48() << 62); // Parse options + mountpoint = cfg["mount"].string_value(); if (cfg["logfile"].string_value() != "") logfile = cfg["logfile"].string_value(); pidfile = cfg["pidfile"].string_value(); @@ -194,8 +203,22 @@ void nfs_proxy_t::run(json11::Json cfg) default_pool = cfg["pool"].as_string(); portmap_enabled = !json_is_false(cfg["portmap"]); nfs_port = cfg["port"].uint64_value() & 0xffff; + nfs_rdma_port = cfg["nfs_rdma"].uint64_value() & 0xffff; + // Allow RDMA-only mode if port is explicitly set to 0 if (!nfs_port) - nfs_port = 2049; + nfs_port = !cfg["port"].is_null() && nfs_rdma_port ? -1 : 2049; + nfs_rdma_credit = cfg["nfs_rdma_credit"].uint64_value(); + if (!nfs_rdma_credit) + nfs_rdma_credit = 16; + nfs_rdma_max_send = cfg["nfs_rdma_send"].uint64_value(); + if (!nfs_rdma_max_send) + nfs_rdma_max_send = 1024; + nfs_rdma_alloc = cfg["nfs_rdma_alloc"].uint64_value(); + if (!nfs_rdma_alloc) + nfs_rdma_alloc = 1048576; + nfs_rdma_gc = cfg["nfs_rdma_gc"].uint64_value(); + if (!nfs_rdma_gc) + nfs_rdma_gc = 500*1048576; export_root = cfg["nfspath"].string_value(); if (!export_root.size()) export_root = "/"; @@ -207,7 +230,6 @@ void nfs_proxy_t::run(json11::Json cfg) obj["client_writeback_allowed"] = true; cfg = obj; } - mountpoint = cfg["mount"].string_value(); if (mountpoint != "") { bind_address = "127.0.0.1"; @@ -292,49 +314,56 @@ void nfs_proxy_t::run(json11::Json cfg) void nfs_proxy_t::run_server(json11::Json cfg) { + if (nfs_port != -1) + { + // Create NFS socket and add it to epoll + int nfs_socket = create_and_bind_socket(bind_address, nfs_port, 128, &listening_port); + fcntl(nfs_socket, F_SETFL, fcntl(nfs_socket, F_GETFL, 0) | O_NONBLOCK); + epmgr->tfd->set_fd_handler(nfs_socket, false, [this](int nfs_socket, int epoll_events) + { + if (epoll_events & EPOLLRDHUP) + { + fprintf(stderr, "Listening portmap socket disconnected, exiting\n"); + exit(1); + } + else + { + do_accept(nfs_socket); + } + }); + } + else + { + listening_port = nfs_rdma_port; + } // Self-register portmap and NFS pmap.reg_ports.insert((portmap_id_t){ .prog = PMAP_PROGRAM, .vers = PMAP_V2, - .port = portmap_enabled ? 111 : nfs_port, + .port = (unsigned)(portmap_enabled ? 111 : listening_port), .owner = "portmapper-service", - .addr = portmap_enabled ? "0.0.0.0.0.111" : ("0.0.0.0.0."+std::to_string(nfs_port)), + .addr = portmap_enabled ? "0.0.0.0.0.111" : ("0.0.0.0.0."+std::to_string(listening_port)), }); pmap.reg_ports.insert((portmap_id_t){ .prog = PMAP_PROGRAM, .vers = PMAP_V3, - .port = portmap_enabled ? 111 : nfs_port, + .port = (unsigned)(portmap_enabled ? 111 : listening_port), .owner = "portmapper-service", - .addr = portmap_enabled ? "0.0.0.0.0.111" : ("0.0.0.0.0."+std::to_string(nfs_port)), + .addr = portmap_enabled ? "0.0.0.0.0.111" : ("0.0.0.0.0."+std::to_string(listening_port)), }); pmap.reg_ports.insert((portmap_id_t){ .prog = NFS_PROGRAM, .vers = NFS_V3, - .port = nfs_port, + .port = (unsigned)listening_port, .owner = "nfs-server", - .addr = "0.0.0.0.0."+std::to_string(nfs_port), + .addr = "0.0.0.0.0."+std::to_string(listening_port), }); pmap.reg_ports.insert((portmap_id_t){ .prog = MOUNT_PROGRAM, .vers = MOUNT_V3, - .port = nfs_port, + .port = (unsigned)listening_port, .owner = "rpc.mountd", - .addr = "0.0.0.0.0."+std::to_string(nfs_port), - }); - // Create NFS socket and add it to epoll - int nfs_socket = create_and_bind_socket(bind_address, nfs_port, 128, &listening_port); - fcntl(nfs_socket, F_SETFL, fcntl(nfs_socket, F_GETFL, 0) | O_NONBLOCK); - epmgr->tfd->set_fd_handler(nfs_socket, false, [this](int nfs_socket, int epoll_events) - { - if (epoll_events & EPOLLRDHUP) - { - fprintf(stderr, "Listening portmap socket disconnected, exiting\n"); - exit(1); - } - else - { - do_accept(nfs_socket); - } + .addr = "0.0.0.0.0."+std::to_string(listening_port), }); if (portmap_enabled) { @@ -354,6 +383,10 @@ void nfs_proxy_t::run_server(json11::Json cfg) } }); } + if (nfs_rdma_port) + { + rdma_context = create_rdma(bind_address, nfs_rdma_port, nfs_rdma_credit, nfs_rdma_max_send, nfs_rdma_alloc, nfs_rdma_gc); + } if (mountpoint != "") { mount_fs(); @@ -499,6 +532,20 @@ void nfs_proxy_t::check_default_pool() } } +nfs_client_t *nfs_proxy_t::create_client() +{ + auto cli = new nfs_client_t(); + cli->parent = this; + if (kvfs) + nfs_kv_procs(cli); + else + nfs_block_procs(cli); + for (auto & fn: pmap.proc_table) + cli->proc_table.insert(fn); + rpc_clients.insert(cli); + return cli; +} + void nfs_proxy_t::do_accept(int listen_fd) { struct sockaddr_storage addr; @@ -512,18 +559,8 @@ void nfs_proxy_t::do_accept(int listen_fd) fcntl(nfs_fd, F_SETFL, fcntl(nfs_fd, F_GETFL, 0) | O_NONBLOCK); int one = 1; setsockopt(nfs_fd, SOL_TCP, TCP_NODELAY, &one, sizeof(one)); - auto cli = new nfs_client_t(); - if (kvfs) - nfs_kv_procs(cli); - else - nfs_block_procs(cli); - cli->parent = this; + auto cli = this->create_client(); cli->nfs_fd = nfs_fd; - for (auto & fn: pmap.proc_table) - { - cli->proc_table.insert(fn); - } - rpc_clients[nfs_fd] = cli; epmgr->tfd->set_fd_handler(nfs_fd, true, [cli](int nfs_fd, int epoll_events) { // Handle incoming event @@ -780,11 +817,17 @@ void nfs_client_t::stop() stopped = true; if (refs <= 0) { +#ifdef WITH_RDMACM + destroy_rdma_conn(); +#endif auto parent = this->parent; - parent->rpc_clients.erase(nfs_fd); + parent->rpc_clients.erase(this); parent->active_connections--; - parent->epmgr->tfd->set_fd_handler(nfs_fd, true, NULL); - close(nfs_fd); + if (nfs_fd >= 0) + { + parent->epmgr->tfd->set_fd_handler(nfs_fd, true, NULL); + close(nfs_fd); + } delete this; parent->check_exit(); } @@ -813,8 +856,7 @@ void nfs_client_t::handle_send(int result) if (rop) { // Reply fully sent - xdr_reset(rop->xdrs); - parent->xdr_pool.push_back(rop->xdrs); + parent->free_xdr(rop->xdrs); if (rop->buffer && rop->referenced) { // Dereference the buffer @@ -831,7 +873,7 @@ void nfs_client_t::handle_send(int result) { // FIXME Maybe put free_buffers into parent free_buffers.push_back((rpc_free_buffer_t){ - .buf = rop->buffer, + .buf = (uint8_t*)rop->buffer, .size = ub.size, }); used_buffers.erase(rop->buffer); @@ -876,8 +918,6 @@ void nfs_client_t::handle_send(int result) void rpc_queue_reply(rpc_op_t *rop) { nfs_client_t *self = (nfs_client_t*)rop->client; - iovec *iov_list = NULL; - unsigned iov_count = 0; int r = xdr_encode(rop->xdrs, (xdrproc_t)xdr_rpc_msg, &rop->out_msg); assert(r); if (rop->reply_fn != NULL) @@ -885,55 +925,78 @@ void rpc_queue_reply(rpc_op_t *rop) r = xdr_encode(rop->xdrs, rop->reply_fn, rop->reply); assert(r); } - xdr_encode_finish(rop->xdrs, &iov_list, &iov_count); - assert(iov_count > 0); - rop->reply_marker = 0; - for (unsigned i = 0; i < iov_count; i++) +#ifdef WITH_RDMACM + if (!self->rdma_conn) +#endif { - rop->reply_marker += iov_list[i].iov_len; - } - rop->reply_marker = htobe32(rop->reply_marker | 0x80000000); - auto & to_send_list = self->write_msg.msg_iovlen ? self->next_send_list : self->send_list; - auto & to_outbox = self->write_msg.msg_iovlen ? self->next_outbox : self->outbox; - to_send_list.push_back((iovec){ .iov_base = &rop->reply_marker, .iov_len = 4 }); - to_outbox.push_back(NULL); - for (unsigned i = 0; i < iov_count; i++) - { - to_send_list.push_back(iov_list[i]); + iovec *iov_list = NULL; + unsigned iov_count = 0; + xdr_encode_finish(rop->xdrs, &iov_list, &iov_count); + assert(iov_count > 0); + rop->reply_marker = 0; + for (unsigned i = 0; i < iov_count; i++) + { + rop->reply_marker += iov_list[i].iov_len; + } + rop->reply_marker = htobe32(rop->reply_marker | 0x80000000); + auto & to_send_list = self->write_msg.msg_iovlen ? self->next_send_list : self->send_list; + auto & to_outbox = self->write_msg.msg_iovlen ? self->next_outbox : self->outbox; + to_send_list.push_back((iovec){ .iov_base = &rop->reply_marker, .iov_len = 4 }); to_outbox.push_back(NULL); + for (unsigned i = 0; i < iov_count; i++) + { + to_send_list.push_back(iov_list[i]); + to_outbox.push_back(NULL); + } + to_outbox[to_outbox.size()-1] = rop; + self->submit_send(); } - to_outbox[to_outbox.size()-1] = rop; - self->submit_send(); +#ifdef WITH_RDMACM + else + { + self->rdma_queue_reply(rop); + } +#endif } -int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg_len) +XDR *nfs_proxy_t::get_xdr() { // Take an XDR object from the pool XDR *xdrs; - if (parent->xdr_pool.size()) + if (xdr_pool.size()) { - xdrs = parent->xdr_pool.back(); - parent->xdr_pool.pop_back(); + xdrs = xdr_pool.back(); + xdr_pool.pop_back(); } else { xdrs = xdr_create(); } + return xdrs; +} + +void nfs_proxy_t::free_xdr(XDR *xdrs) +{ + xdr_reset(xdrs); + xdr_pool.push_back(xdrs); +} + +int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg_len) +{ + XDR *xdrs = parent->get_xdr(); // Decode the RPC header char inmsg_data[sizeof(rpc_msg)]; rpc_msg *inmsg = (rpc_msg*)&inmsg_data; if (!xdr_decode(xdrs, msg_buf, msg_len, (xdrproc_t)xdr_rpc_msg, inmsg)) { // Invalid message, ignore it - xdr_reset(xdrs); - parent->xdr_pool.push_back(xdrs); + parent->free_xdr(xdrs); return 0; } if (inmsg->body.dir != RPC_CALL) { // Reply sent to the server? Strange thing. Also ignore it - xdr_reset(xdrs); - parent->xdr_pool.push_back(xdrs); + parent->free_xdr(xdrs); return 0; } if (inmsg->body.cbody.rpcvers != RPC_MSG_VERSION) @@ -968,6 +1031,17 @@ int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg // Incoming buffer isn't needed to handle request, so return 0 return 0; } + auto rop = create_rpc_op(xdrs, base_buf, inmsg, NULL); + if (!rop) + { + // No such procedure + return 0; + } + return handle_rpc_op(rop); +} + +rpc_op_t *nfs_client_t::create_rpc_op(XDR *xdrs, void *buffer, rpc_msg *inmsg, rdma_msg *rmsg) +{ // Find decoder for the request auto proc_it = proc_table.find((rpc_service_proc_t){ .prog = inmsg->body.cbody.prog, @@ -995,6 +1069,7 @@ int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg rpc_op_t *rop = (rpc_op_t*)malloc_or_die(sizeof(rpc_op_t)); *rop = (rpc_op_t){ .client = this, + .buffer = buffer, .xdrs = xdrs, .out_msg = (rpc_msg){ .xid = inmsg->xid, @@ -1017,9 +1092,15 @@ int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg }, }, }; + // FIXME: malloc and avoid copy? + memcpy(&rop->in_msg, inmsg, sizeof(rpc_msg)); + if (rmsg) + { + memcpy(&rop->in_rdma_msg, rmsg, sizeof(rdma_msg)); + } rpc_queue_reply(rop); // Incoming buffer isn't needed to handle request, so return 0 - return 0; + return NULL; } // Allocate memory rpc_op_t *rop = (rpc_op_t*)malloc_or_die( @@ -1028,7 +1109,7 @@ int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg rpc_reply_stat x = RPC_MSG_ACCEPTED; *rop = (rpc_op_t){ .client = this, - .buffer = (uint8_t*)base_buf, + .buffer = buffer, .xdrs = xdrs, .out_msg = (rpc_msg){ .xid = inmsg->xid, @@ -1045,10 +1126,25 @@ int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg .request = ((uint8_t*)rop) + sizeof(rpc_op_t), .reply = ((uint8_t*)rop) + sizeof(rpc_op_t) + proc_it->req_size, }; + // FIXME: malloc and avoid copy? memcpy(&rop->in_msg, inmsg, sizeof(rpc_msg)); + if (rmsg) + { + memcpy(&rop->in_rdma_msg, rmsg, sizeof(rdma_msg)); + } + return rop; +} + +int nfs_client_t::handle_rpc_op(rpc_op_t *rop) +{ // Try to decode the request // req_fn may be NULL, that means function has no arguments - if (proc_it->req_fn && !proc_it->req_fn(xdrs, rop->request)) + auto proc_it = proc_table.find((rpc_service_proc_t){ + .prog = rop->in_msg.body.cbody.prog, + .vers = rop->in_msg.body.cbody.vers, + .proc = rop->in_msg.body.cbody.proc, + }); + if (proc_it == proc_table.end() || proc_it->req_fn && !proc_it->req_fn(rop->xdrs, rop->request)) { // Invalid request rop->out_msg.body.rbody.areply.reply_data.stat = RPC_GARBAGE_ARGS; @@ -1058,18 +1154,55 @@ int nfs_client_t::handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg } rop->out_msg.body.rbody.areply.reply_data.stat = RPC_SUCCESS; rop->reply_fn = proc_it->resp_fn; + rop->referenced = 0; int ref = proc_it->handler_fn(proc_it->opaque, rop); - rop->referenced = ref ? 1 : 0; + if (ref) + rop->referenced = 1; return ref; } +void *nfs_client_t::malloc_or_rdma(rpc_op_t *rop, size_t size) +{ +#ifdef WITH_RDMACM + if (!rdma_conn) + { +#endif + void *buf = malloc_or_die(size); + xdr_add_malloc(rop->xdrs, buf); + return buf; +#ifdef WITH_RDMACM + } + void *buf = rdma_malloc(size); + xdr_set_rdma_chunk(rop->xdrs, buf); + return buf; +#endif +} + +void nfs_client_t::free_or_rdma(rpc_op_t *rop, void *buf) +{ +#ifdef WITH_RDMACM + if (!rdma_conn) + { +#endif + xdr_del_malloc(rop->xdrs, buf); + free(buf); +#ifdef WITH_RDMACM + } + else + { + xdr_set_rdma_chunk(rop->xdrs, NULL); + rdma_free(buf); + } +#endif +} + void nfs_proxy_t::daemonize() { // Stop all clients because client I/O sometimes breaks during daemonize // I.e. the new process stops receiving events on the old FD // It doesn't happen if we call sleep(1) here, but we don't want to call sleep(1)... - for (auto & clp: rpc_clients) - clp.second->stop(); + for (auto & cli: rpc_clients) + cli->stop(); if (fork()) exit(0); setsid(); diff --git a/src/nfs/nfs_proxy.h b/src/nfs/nfs_proxy.h index 29f12019..ad60c4a6 100644 --- a/src/nfs/nfs_proxy.h +++ b/src/nfs/nfs_proxy.h @@ -22,6 +22,7 @@ class cli_tool_t; struct kv_fs_state_t; struct block_fs_state_t; class nfs_client_t; +struct nfs_rdma_context_t; class nfs_proxy_t { @@ -33,7 +34,12 @@ public: std::string default_pool; std::string export_root; bool portmap_enabled; - unsigned nfs_port; + unsigned nfs_port = 0; + unsigned nfs_rdma_port = 0; + uint32_t nfs_rdma_credit = 16; + uint32_t nfs_rdma_max_send = 1024; + uint64_t nfs_rdma_alloc = 1048576; + uint64_t nfs_rdma_gc = 500*1048576; int trace = 0; std::string logfile = "/dev/null"; std::string pidfile; @@ -55,7 +61,8 @@ public: vitastorkv_dbw_t *db = NULL; kv_fs_state_t *kvfs = NULL; block_fs_state_t *blockfs = NULL; - std::map rpc_clients; + nfs_rdma_context_t* rdma_context = NULL; + std::set rpc_clients; std::vector xdr_pool; @@ -72,12 +79,20 @@ public: void watch_stats(); void parse_stats(etcd_kv_t & kv); void check_default_pool(); + nfs_client_t* create_client(); void do_accept(int listen_fd); void daemonize(); void write_pid(); void mount_fs(); void check_already_mounted(); void check_exit(); + + nfs_rdma_context_t* create_rdma(const std::string & bind_address, int rdmacm_port, + uint32_t max_iodepth, uint32_t max_send_wr, uint64_t rdma_malloc_round_to, uint64_t rdma_max_unused_buffers); + void destroy_rdma(); + + XDR *get_xdr(); + void free_xdr(XDR *xdrs); }; struct rpc_cur_buffer_t @@ -101,19 +116,24 @@ struct rpc_free_buffer_t unsigned size; }; +struct nfs_rdma_conn_t; + class nfs_client_t { public: nfs_proxy_t *parent = NULL; - int nfs_fd; - int epoll_events = 0; int refs = 0; bool stopped = false; std::set proc_table; + nfs_rdma_conn_t *rdma_conn = NULL; + + // + int nfs_fd = -1; + int epoll_events = 0; // Read state rpc_cur_buffer_t cur_buffer = { 0 }; - std::map used_buffers; + std::map used_buffers; std::vector free_buffers; iovec read_iov; @@ -130,7 +150,16 @@ public: void submit_send(); void handle_send(int result); int handle_rpc_message(void *base_buf, void *msg_buf, uint32_t msg_len); + // + rpc_op_t *create_rpc_op(XDR *xdrs, void *buffer, rpc_msg *inmsg, rdma_msg *rmsg); + int handle_rpc_op(rpc_op_t *rop); bool deref(); void stop(); + void *malloc_or_rdma(rpc_op_t *rop, size_t size); + void free_or_rdma(rpc_op_t *rop, void *buf); + void *rdma_malloc(size_t size); + void rdma_free(void *buf); + void rdma_queue_reply(rpc_op_t *rop); + void destroy_rdma_conn(); }; diff --git a/src/nfs/nfs_proxy_rdma.cpp b/src/nfs/nfs_proxy_rdma.cpp new file mode 100644 index 00000000..449266cf --- /dev/null +++ b/src/nfs/nfs_proxy_rdma.cpp @@ -0,0 +1,1072 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 (see README.md for details) +// +// NFS RDMA support + +#ifdef WITH_RDMACM + +#define _XOPEN_SOURCE + +#include +#include + +#include "addr_util.h" + +#include "proto/nfs.h" +#include "proto/rpc.h" +#include "proto/rpc_rdma.h" + +#include "nfs_proxy.h" + +#include "rdma_alloc.h" + +#define NFS_RDMACM_PRIVATE_DATA_MAGIC_LE 0x180eabf6 + +struct __attribute__((__packed__)) nfs_rdmacm_private +{ + uint32_t format_identifier; // magic, should be 0xf6ab0e18 in big endian + uint8_t version; // version, 1 + uint8_t remote_invalidate; // remote invalidation flag (1 or 0) + uint8_t max_send_size; // maximum RDMA Send operation size / 1024 - 1 (i.e. 0 is 1 KB, 255 is 256 KB) + uint8_t max_recv_size; // maximum RDMA Receive operation size / 1024 - 1 (i.e. 0 is 1 KB, 255 is 256 KB) +}; + +struct nfs_rdma_buf_t +{ + void *buf = NULL; + size_t len = 0; + ibv_mr *mr = NULL; +}; + +struct nfs_rdma_conn_t; + +struct nfs_rdma_context_t +{ + std::string bind_address; + int rdmacm_port = 0; + uint32_t max_iodepth = 16, max_send_wr = 1024; + uint64_t rdma_malloc_round_to = 1048576, rdma_max_unused_buffers = 500*1048576; + uint64_t max_send_size = 256*1024, max_recv_size = 256*1024; + + nfs_proxy_t *proxy = NULL; + epoll_manager_t *epmgr = NULL; + + int max_cqe = 0, used_max_cqe = 0; + rdma_event_channel *rdmacm_evch = NULL; + rdma_cm_id *listener_id = NULL; + ibv_device_attr listener_dev_attr; // present only if bound to a specific device + ibv_comp_channel *channel = NULL; + ibv_cq *cq = NULL; + rdma_allocator_t *alloc = NULL; + std::map rdma_connections; + std::map rdma_connections_by_qp; + + ~nfs_rdma_context_t(); + void handle_io(); + void handle_rdmacm_events(); + void rdmacm_accept(rdma_cm_event *ev); + void rdmacm_established(rdma_cm_event *ev); +}; + +struct nfs_rdma_conn_t +{ + nfs_rdma_context_t *ctx = NULL; + nfs_client_t *client = NULL; + rdma_cm_id *id = NULL; + int max_send_size = 256*1024, max_recv_size = 256*1024; + int max_buf_size = 256*1024; + int remote_max_send_size = 1024, remote_max_recv_size = 1024; + int max_rdma_reads = 16; + bool remote_invalidate = true; + bool established = false; + uint32_t cur_credit = 16; + uint32_t cur_send = 0; + uint32_t cur_rdma_reads = 0; + std::vector recv_buffers; + std::map used_buffers; + int next_recv_buf = 0; + std::vector outbox; + std::vector outbox_wrs; + std::vector chunk_inbox, chunk_read_postponed; + int outbox_pos = 0; + + void post_initial_receives(); + ~nfs_rdma_conn_t(); + nfs_rdma_buf_t create_buf(size_t len); + void post_recv(nfs_rdma_buf_t b); + void post_send(); + int post_chunk_reads(rpc_op_t *rop, bool push); + bool handle_recv(void *buf, size_t len); + void free_rdma_rpc_op(rpc_op_t *rop); + void reuse_buffer(void *buf); + void rdma_encode_header(XDR *xdrs, rpc_op_t *rop, bool nomsg); +}; + +nfs_rdma_context_t* nfs_proxy_t::create_rdma(const std::string & bind_address, int rdmacm_port, + uint32_t max_iodepth, uint32_t max_send_wr, uint64_t rdma_malloc_round_to, uint64_t rdma_max_unused_buffers) +{ + nfs_rdma_context_t* self = new nfs_rdma_context_t; + self->proxy = this; + self->epmgr = epmgr; + self->bind_address = bind_address; + self->rdmacm_port = rdmacm_port; + self->max_iodepth = max_iodepth; + self->max_send_wr = max_send_wr ? max_send_wr : 1024; + self->rdma_malloc_round_to = rdma_malloc_round_to; + self->rdma_max_unused_buffers = rdma_max_unused_buffers; + self->rdmacm_evch = rdma_create_event_channel(); + if (!self->rdmacm_evch) + { + fprintf(stderr, "Failed to initialize RDMA-CM event channel: %s (code %d)\n", strerror(errno), errno); + delete self; + return NULL; + } + fcntl(self->rdmacm_evch->fd, F_SETFL, fcntl(self->rdmacm_evch->fd, F_GETFL, 0) | O_NONBLOCK); + epmgr->tfd->set_fd_handler(self->rdmacm_evch->fd, false, [self](int rdmacm_eventfd, int epoll_events) + { + self->handle_rdmacm_events(); + }); + int r = rdma_create_id(self->rdmacm_evch, &self->listener_id, NULL, RDMA_PS_TCP); + if (r != 0) + { + fprintf(stderr, "Failed to create RDMA-CM ID: %s (code %d)\n", strerror(errno), errno); + delete self; + return NULL; + } + sockaddr_storage addr; + if (!string_to_addr(bind_address, 0, rdmacm_port, &addr)) + { + fprintf(stderr, "Server address: %s is not valid\n", bind_address.c_str()); + delete self; + return NULL; + } + r = rdma_bind_addr(self->listener_id, (sockaddr*)&addr); + if (r != 0) + { + fprintf(stderr, "Failed to bind RDMA-CM to %s:%d: %s (code %d)\n", bind_address.c_str(), rdmacm_port, strerror(errno), errno); + delete self; + return NULL; + } + r = rdma_listen(self->listener_id, 128); + if (r != 0) + { + fprintf(stderr, "Failed to listen RDMA-CM: %s (code %d)\n", strerror(errno), errno); + delete self; + return NULL; + } + if (self->listener_id->verbs) + { + r = ibv_query_device(self->listener_id->verbs, &self->listener_dev_attr); + if (r != 0) + { + fprintf(stderr, "Failed to query listerning RDMA device: %s (code %d)\n", strerror(r), r); + delete self; + return NULL; + } + } + // FIXME: What if self->listener_id->verbs is empty... + self->channel = ibv_create_comp_channel(self->listener_id->verbs); + if (!self->channel) + { + fprintf(stderr, "Couldn't create RDMA completion channel\n"); + delete self; + return NULL; + } + self->max_cqe = 4096; + self->cq = ibv_create_cq(self->listener_id->verbs, self->max_cqe, NULL, self->channel, 0); + if (!self->cq) + { + fprintf(stderr, "Couldn't create RDMA completion queue\n"); + delete self; + return NULL; + } + self->alloc = rdma_malloc_create(self->listener_id->pd, self->rdma_malloc_round_to, self->rdma_max_unused_buffers, IBV_ACCESS_LOCAL_WRITE); + fcntl(self->channel->fd, F_SETFL, fcntl(self->channel->fd, F_GETFL, 0) | O_NONBLOCK); + epmgr->tfd->set_fd_handler(self->channel->fd, false, [self](int channel_eventfd, int epoll_events) + { + self->handle_io(); + }); + // run handle_io() once to reset poll state + self->handle_io(); + return self; +} + +void nfs_proxy_t::destroy_rdma() +{ + if (rdma_context) + { + delete rdma_context; + rdma_context = NULL; + } +} + +nfs_rdma_context_t::~nfs_rdma_context_t() +{ + if (listener_id) + { + int r = rdma_destroy_id(listener_id); + if (r != 0) + fprintf(stderr, "Failed to destroy RDMA-CM ID: %s (code %d)\n", strerror(errno), errno); + else + listener_id = NULL; + } + if (rdmacm_evch) + { + epmgr->tfd->set_fd_handler(rdmacm_evch->fd, false, NULL); + rdma_destroy_event_channel(rdmacm_evch); + rdmacm_evch = NULL; + } + if (cq) + { + ibv_destroy_cq(cq); + cq = NULL; + } + if (channel) + { + epmgr->tfd->set_fd_handler(channel->fd, false, NULL); + ibv_destroy_comp_channel(channel); + channel = NULL; + } + if (alloc) + { + rdma_malloc_destroy(alloc); + alloc = NULL; + } +} + +void nfs_rdma_context_t::handle_rdmacm_events() +{ + rdma_cm_event *ev = NULL; + std::vector stop_clients; + while (1) + { + int r = rdma_get_cm_event(rdmacm_evch, &ev); + if (r != 0) + { + if (errno == EAGAIN || errno == EINTR) + break; + fprintf(stderr, "Failed to get RDMA-CM event: %s (code %d)\n", strerror(errno), errno); + exit(1); + } + if (ev->event == RDMA_CM_EVENT_CONNECT_REQUEST) + { + rdmacm_accept(ev); + } + else if (ev->event == RDMA_CM_EVENT_CONNECT_ERROR || + ev->event == RDMA_CM_EVENT_REJECTED || + ev->event == RDMA_CM_EVENT_DISCONNECTED || + ev->event == RDMA_CM_EVENT_DEVICE_REMOVAL) + { + auto event_type_name = ev->event == RDMA_CM_EVENT_CONNECT_ERROR ? "RDMA_CM_EVENT_CONNECT_ERROR" : ( + ev->event == RDMA_CM_EVENT_REJECTED ? "RDMA_CM_EVENT_REJECTED" : ( + ev->event == RDMA_CM_EVENT_DISCONNECTED ? "RDMA_CM_EVENT_DISCONNECTED" : "RDMA_CM_EVENT_DEVICE_REMOVAL")); + auto conn_it = rdma_connections.find(ev->id); + if (conn_it == rdma_connections.end()) + { + fprintf(stderr, "Received %s event for an unknown connection 0x%jx - ignoring\n", + event_type_name, (uint64_t)ev->id); + } + else + { + fprintf(stderr, "Received %s event for connection 0x%jx - closing it\n", + event_type_name, (uint64_t)ev->id); + auto conn = conn_it->second; + stop_clients.push_back(conn->client); + } + } + else if (ev->event == RDMA_CM_EVENT_ESTABLISHED) + { + rdmacm_established(ev); + } + else if (ev->event == RDMA_CM_EVENT_ADDR_CHANGE || ev->event == RDMA_CM_EVENT_TIMEWAIT_EXIT) + { + // Do nothing + } + else + { + // Other events are unexpected + fprintf(stderr, "Unexpected RDMA-CM event type: %d\n", ev->event); + } + r = rdma_ack_cm_event(ev); + if (r != 0) + { + fprintf(stderr, "Failed to ack (free) RDMA-CM event: %s (code %d)\n", strerror(errno), errno); + exit(1); + } + } + // Stop only after flushing all events, otherwise rdma_destroy_id infinitely waits for pthread_cond + for (auto cli: stop_clients) + { + cli->stop(); + } +} + +void nfs_rdma_context_t::rdmacm_accept(rdma_cm_event *ev) +{ + this->used_max_cqe += max_iodepth*2; + if (this->used_max_cqe > this->max_cqe) + { + // Resize CQ + int new_max_cqe = this->max_cqe; + while (this->used_max_cqe > new_max_cqe) + { + new_max_cqe *= 2; + } + if (ibv_resize_cq(this->cq, new_max_cqe) != 0) + { + fprintf(stderr, "Couldn't resize RDMA completion queue to %d entries\n", new_max_cqe); + rdma_destroy_id(ev->id); + return; + } + this->max_cqe = new_max_cqe; + } + ibv_qp_init_attr init_attr = { + .send_cq = this->cq, + .recv_cq = this->cq, + .cap = { + // each op at each moment takes 1 RDMA_RECV or 1 RDMA_READ or 1 RDMA_WRITE + 1 RDMA_SEND + .max_send_wr = max_send_wr, + .max_recv_wr = max_iodepth, + .max_send_sge = 1, + .max_recv_sge = 1, // we don't need S/G currently + }, + .qp_type = IBV_QPT_RC, + }; + int r = rdma_create_qp(ev->id, NULL, &init_attr); + if (r != 0) + { + fprintf(stderr, "Failed to create a queue pair via RDMA-CM: %s (code %d)\n", strerror(errno), errno); + rdma_destroy_id(ev->id); + return; + } + assert(ev->id->qp->send_cq == this->cq); + assert(ev->id->qp->recv_cq == this->cq); + nfs_rdmacm_private private_data = { + .format_identifier = NFS_RDMACM_PRIVATE_DATA_MAGIC_LE, + .version = 1, + .remote_invalidate = 1, + .max_send_size = (uint8_t)(max_send_size <= 256*1024 ? max_send_size/1024 - 1 : 255), + .max_recv_size = (uint8_t)(max_recv_size <= 256*1024 ? max_recv_size/1024 - 1 : 255), + }; + // should be responder_resources, but it's 0 for Linux NFS client + // so probably RDMA-CM swaps these values in _accept for the "convenience" :) + int max_rdma_reads = ev->param.conn.initiator_depth; + // just 16 on ConnectX-4... + if (!max_rdma_reads || max_rdma_reads > listener_dev_attr.max_qp_rd_atom) + max_rdma_reads = listener_dev_attr.max_qp_rd_atom; + if (max_rdma_reads > listener_dev_attr.max_qp_init_rd_atom) + max_rdma_reads = listener_dev_attr.max_qp_init_rd_atom; + if (max_rdma_reads > 255) + max_rdma_reads = 255; + rdma_conn_param conn_params = { + .private_data = &private_data, + .private_data_len = sizeof(private_data), + .responder_resources = (uint8_t)max_rdma_reads, // min(max_qp_rd_atom of the local device, max_qp_init_rd_atom of the remote device) + .initiator_depth = (uint8_t)max_rdma_reads, // max_qp_init_rd_atom of the local device + .rnr_retry_count = 7, + }; + assert(ev->id->verbs == listener_id->verbs); + r = rdma_accept(ev->id, &conn_params); + if (r != 0) + { + fprintf(stderr, "Failed to accept RDMA-CM connection: %s (code %d)\n", strerror(errno), errno); + rdma_destroy_qp(ev->id); + rdma_destroy_id(ev->id); + } + else + { + auto conn = new nfs_rdma_conn_t(); + conn->ctx = this; + conn->id = ev->id; + conn->cur_credit = max_iodepth; + conn->max_send_size = max_send_size; + conn->max_recv_size = max_recv_size; + conn->max_rdma_reads = max_rdma_reads; + rdma_connections[ev->id] = conn; + rdma_connections_by_qp[conn->id->qp->qp_num] = conn; + // Handle NFS private_data + if (ev->param.conn.private_data_len >= sizeof(nfs_rdmacm_private)) + { + nfs_rdmacm_private *private_data = (nfs_rdmacm_private *)ev->param.conn.private_data; + if (private_data->format_identifier == NFS_RDMACM_PRIVATE_DATA_MAGIC_LE && + private_data->version == 1) + { + conn->remote_invalidate = private_data->remote_invalidate; + conn->remote_max_send_size = (private_data->max_send_size+1) * 1024; + conn->remote_max_recv_size = (private_data->max_recv_size+1) * 1024; + if (conn->remote_max_recv_size < conn->max_send_size) + conn->max_send_size = conn->remote_max_recv_size; + } + } + auto cli = this->proxy->create_client(); + conn->client = cli; + cli->rdma_conn = conn; + // Post initial receive requests + conn->post_initial_receives(); + } +} + +nfs_rdma_conn_t::~nfs_rdma_conn_t() +{ + assert(!outbox.size()); + assert(!chunk_inbox.size()); + assert(!chunk_read_postponed.size()); + for (auto & b: recv_buffers) + { + ibv_dereg_mr(b.mr); + free(b.buf); + } + recv_buffers.clear(); + if (id) + { + ctx->rdma_connections.erase(id); + if (id->qp) + { + ctx->rdma_connections_by_qp.erase(id->qp->qp_num); + rdma_destroy_qp(id); + } + rdma_destroy_id(id); + } +} + +void nfs_rdma_context_t::rdmacm_established(rdma_cm_event *ev) +{ + auto conn_it = rdma_connections.find(ev->id); + if (conn_it == rdma_connections.end()) + { + fprintf(stderr, "Received RDMA_CM_EVENT_ESTABLISHED event for an unknown connection 0x%jx - ignoring\n", (uint64_t)ev->id); + return; + } + fprintf(stderr, "Received RDMA_CM_EVENT_ESTABLISHED event for connection 0x%jx - connection established\n", (uint64_t)ev->id); + auto conn = conn_it->second; + conn->established = true; +} + +void nfs_rdma_conn_t::post_initial_receives() +{ + for (int i = 0; i < cur_credit; i++) + { + auto b = create_buf(max_recv_size); + recv_buffers.push_back(b); + used_buffers[b.buf] = b; + post_recv(b); + } +} + +nfs_rdma_buf_t nfs_rdma_conn_t::create_buf(size_t len) +{ + nfs_rdma_buf_t b; + b.buf = malloc_or_die(len); + b.len = len; + b.mr = ibv_reg_mr(id->pd, b.buf, len, IBV_ACCESS_LOCAL_WRITE); + if (!b.mr) + { + fprintf(stderr, "Failed to register RDMA memory region: %s\n", strerror(errno)); + exit(1); + } + return b; +} + +void nfs_rdma_conn_t::post_recv(nfs_rdma_buf_t b) +{ + if (client->stopped) + { + return; + } + ibv_sge sge = { + .addr = (uintptr_t)b.buf, + .length = (uint32_t)b.len, + .lkey = b.mr->lkey, + }; + ibv_recv_wr *bad_wr = NULL; + ibv_recv_wr wr = { + .wr_id = 1, // 1 is any read, 2 is any write :) + .sg_list = &sge, + .num_sge = 1, + }; + int err = ibv_post_recv(id->qp, &wr, &bad_wr); + if (err || bad_wr) + { + fprintf(stderr, "RDMA receive failed: %s\n", strerror(err)); + exit(1); + } +} + +void nfs_rdma_conn_t::rdma_encode_header(XDR *xdrs, rpc_op_t *rop, bool nomsg) +{ + rdma_msg outrmsg = { + .rdma_xid = rop->in_rdma_msg.rdma_xid, + .rdma_vers = rop->in_rdma_msg.rdma_vers, + .rdma_credit = cur_credit, + .rdma_body = { + .proc = rop->rdma_error ? RDMA_ERROR : (nomsg ? RDMA_NOMSG : RDMA_MSG), + }, + }; + if (rop->rdma_error) + { + outrmsg.rdma_body.rdma_error.err = rop->rdma_error; + if (rop->rdma_error == ERR_VERS) + outrmsg.rdma_body.rdma_error.range = (rpc_rdma_errvers){ 1, 1 }; + } + else + { + // Copy chunks... it's a real shit + outrmsg.rdma_body.rdma_msg = { + .rdma_writes = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes, + .rdma_reply = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_reply, + }; + } + int r = xdr_encode(xdrs, (xdrproc_t)xdr_rdma_msg, &outrmsg); + assert(r); +} + +void nfs_client_t::rdma_queue_reply(rpc_op_t *rop) +{ + rdma_conn->outbox.push_back(rop); + rdma_conn->post_send(); +} + +void nfs_rdma_conn_t::post_send() +{ + while (chunk_read_postponed.size() > 0) + { + int posted = post_chunk_reads(chunk_read_postponed[0], false); + if (posted) + chunk_read_postponed.erase(chunk_read_postponed.begin(), chunk_read_postponed.begin()+1); + else + break; + } + while (outbox.size() > outbox_pos) + { + auto rop = outbox[outbox_pos]; + if (rop->buffer) + { + reuse_buffer(rop->buffer); + rop->buffer = NULL; + } + XDR *hdr_xdr = ctx->proxy->get_xdr(); +send_again: + rdma_encode_header(hdr_xdr, rop, false); + size_t hdr_size = xdr_encode_get_size(hdr_xdr); + iovec *chunk_iov = NULL; + iovec *iov_list = NULL; + unsigned iov_count = 0; + if (!rop->rdma_error) + { + xdr_encode_finish(rop->xdrs, &iov_list, &iov_count); + assert(iov_count > 0); + // READ3resok and READLINK3resok - extract last byte buffer from iovecs and send it in a "write chunk" + if (rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes && + rop->in_msg.body.cbody.prog == NFS_PROGRAM && + (rop->in_msg.body.cbody.proc == NFS3_READ && ((READ3res*)rop->reply)->status == NFS3_OK || + rop->in_msg.body.cbody.proc == NFS3_READLINK && ((READLINK3res*)rop->reply)->status == NFS3_OK)) + { + assert(iov_count > 1); + iov_count--; + chunk_iov = &iov_list[iov_count]; + auto expected_size = (rop->in_msg.body.cbody.proc == NFS3_READ + ? ((READ3res*)rop->reply)->resok.count + : ((READLINK3res*)rop->reply)->resok.data.size); + assert(chunk_iov->iov_len == expected_size); + } + } + size_t msg_size = 0; + for (unsigned i = 0; i < iov_count; i++) + { + msg_size += iov_list[i].iov_len; + } + // Estimate reply WR count, create WR and SGE arrays + xdr_write_chunk *reply_chunk = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_reply; + int reply_chunk_wr_count = (reply_chunk ? reply_chunk->target.target_len : 0); + uint32_t wr_count = 1 + (chunk_iov ? 1 : 0) + (reply_chunk ? reply_chunk_wr_count : 0); + if (wr_count > ctx->max_send_wr) + { + fprintf(stderr, "Reply fragmentation (%u) exceeds max_send_wr (%u), sending ERR_CHUNK\n", wr_count, ctx->max_send_wr); +chunk_error: + xdr_reset(hdr_xdr); + rop->rdma_error = ERR_CHUNK; + goto send_again; + } + if (msg_size+hdr_size > max_send_size && !reply_chunk) + { + fprintf(stderr, "RPC message size (%zu) exceeds client's max_send_size (%u)" + " and reply chunk is not provided, sending ERR_CHUNK\n", msg_size+hdr_size, max_send_size); + goto chunk_error; + } + if (chunk_iov && chunk_iov->iov_len > rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes->entry.target.target_val[0].length) + { + fprintf(stderr, "%s write chunk size (%u) is smaller than read size (%zu), sending ERR_CHUNK\n", + rop->in_msg.body.cbody.proc == NFS3_READ ? "READ3" : "READLINK3", + rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes->entry.target.target_val[0].length, + chunk_iov->iov_len); + goto chunk_error; + } + if (cur_send + wr_count > ctx->max_send_wr) + { + // Retry later + ctx->proxy->free_xdr(hdr_xdr); + return; + } + // Check reply chunk size and set actual reply segment sizes + bool reencode = false; + if (reply_chunk) + { + size_t reply_chunk_len = 0; + size_t left = msg_size; + for (uint32_t i = 0; i < reply_chunk->target.target_len; i++) + { + reply_chunk_len += reply_chunk->target.target_val[i].length; + if (reply_chunk->target.target_val[i].length > left) + reply_chunk->target.target_val[i].length = left; + left -= reply_chunk->target.target_val[i].length; + } + if (left > 0) + { + // Message doesn't fit even in the reply chunk + fprintf(stderr, "RPC message payload size (%zu) exceeds reply chunk size (%zu), sending ERR_CHUNK\n", + msg_size, reply_chunk_len); + goto chunk_error; + } + reencode = true; + } + if (chunk_iov && chunk_iov->iov_len < rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes->entry.target.target_val[0].length) + { + // Set actual write chunk size + reencode = true; + rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes->entry.target.target_val[0].length = chunk_iov->iov_len; + } + if (reencode) + { + // ...and re-encode the header :-( shitty protocol + xdr_reset(hdr_xdr); + rdma_encode_header(hdr_xdr, rop, reply_chunk != NULL); + hdr_size = xdr_encode_get_size(hdr_xdr); + } + ibv_sge sges[wr_count]; + ibv_send_wr wrs[wr_count]; + int wr_pos = 0; + // Use a buffer from rdma_malloc for the reply + assert(!rop->buffer); + rop->buffer = rdma_malloc_alloc(ctx->alloc, hdr_size+msg_size); + auto buf_lkey = rdma_malloc_get_lkey(ctx->alloc, rop->buffer); + size_t pos = 0; + { + // Copy and free the RDMA-RPC header + iovec *hdr_iov_list = NULL; + unsigned hdr_iov_count = 0; + xdr_encode_finish(hdr_xdr, &hdr_iov_list, &hdr_iov_count); + assert(hdr_iov_count > 0); + for (unsigned i = 0; i < hdr_iov_count; i++) + { + memcpy(rop->buffer + pos, hdr_iov_list[i].iov_base, hdr_iov_list[i].iov_len); + pos += hdr_iov_list[i].iov_len; + } + assert(pos == hdr_size); + ctx->proxy->free_xdr(hdr_xdr); + } + for (unsigned i = 0; i < iov_count; i++) + { + memcpy(rop->buffer + pos, iov_list[i].iov_base, iov_list[i].iov_len); + pos += iov_list[i].iov_len; + } + // Include header size in msg_size now and then + msg_size += hdr_size; + assert(msg_size == pos); + // Check if we have to use the reply chunk + if (reply_chunk) + { + size_t pos = hdr_size; + for (uint32_t i = 0; i < reply_chunk->target.target_len && pos < msg_size; i++) + { + uint32_t len = (reply_chunk->target.target_val[i].length < msg_size-pos + ? reply_chunk->target.target_val[i].length : msg_size-pos); + sges[wr_pos] = { + .addr = (uintptr_t)(rop->buffer + pos), + .length = len, + .lkey = buf_lkey, + }; + wrs[wr_pos] = { + .wr_id = 4, // 4 is chunk write + .opcode = IBV_WR_RDMA_WRITE, + .wr = { + .rdma = { + .remote_addr = reply_chunk->target.target_val[i].offset, + .rkey = reply_chunk->target.target_val[i].handle, + }, + }, + }; + wr_pos++; + pos += len; + } + } + // Check if we have a DDP chunk + xdr_rdma_segment *wr_chunk = NULL; + if (chunk_iov != NULL) + { + wr_chunk = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_writes->entry.target.target_val; + sges[wr_pos] = { + .addr = (uintptr_t)chunk_iov->iov_base, + .length = (uint32_t)chunk_iov->iov_len, + .lkey = rdma_malloc_get_lkey(ctx->alloc, chunk_iov->iov_base), + }; + wrs[wr_pos] = { + .wr_id = 4, // 4 is chunk write + .opcode = IBV_WR_RDMA_WRITE, + .wr = { + .rdma = { + .remote_addr = wr_chunk->offset, + .rkey = wr_chunk->handle, + }, + }, + }; + wr_pos++; + } + // Add the normal send + sges[wr_pos] = { + .addr = (uintptr_t)rop->buffer, + .length = (uint32_t)(reply_chunk ? hdr_size : msg_size), + .lkey = buf_lkey, + }; + wrs[wr_pos] = { + .wr_id = 2, // 2 is send + .opcode = remote_invalidate && !reply_chunk && wr_chunk ? IBV_WR_SEND_WITH_INV : IBV_WR_SEND, + .send_flags = IBV_SEND_SIGNALED, + .invalidate_rkey = remote_invalidate && !reply_chunk && wr_chunk ? wr_chunk->handle : 0, + }; + wr_pos++; + // Send it all + for (int i = 0; i < wr_pos; i++) + { + wrs[i].next = &wrs[i+1]; + wrs[i].sg_list = &sges[i]; + wrs[i].num_sge = 1; + } + wrs[wr_pos-1].next = NULL; + ibv_send_wr *bad_wr = NULL; + int err = ibv_post_send(id->qp, &wrs[0], &bad_wr); + if (err || bad_wr) + { + fprintf(stderr, "Posting RDMA send failed: %s\n", strerror(err)); + exit(1); + } + cur_send += wr_pos; + outbox_wrs.push_back(wr_pos); + outbox_pos++; + } +} + +#define RDMA_EVENTS_AT_ONCE 32 + +void nfs_rdma_context_t::handle_io() +{ + // Request next notification + ibv_cq *ev_cq; + void *ev_ctx; + // This is inefficient as it calls read()... (but there is no other way) + if (ibv_get_cq_event(channel, &ev_cq, &ev_ctx) == 0) + { + ibv_ack_cq_events(cq, 1); + } + if (ibv_req_notify_cq(cq, 0) != 0) + { + fprintf(stderr, "Failed to request RDMA completion notification, exiting\n"); + exit(1); + } + ibv_wc wc[RDMA_EVENTS_AT_ONCE]; + int event_count; + do + { + event_count = ibv_poll_cq(cq, RDMA_EVENTS_AT_ONCE, wc); + for (int i = 0; i < event_count; i++) + { + auto conn_it = rdma_connections_by_qp.find(wc[i].qp_num); + if (conn_it == rdma_connections_by_qp.end()) + { + continue; + } + auto conn = conn_it->second; + if (wc[i].status != IBV_WC_SUCCESS) + { + fprintf(stderr, "RDMA work request failed for queue %d with status: %s, stopping client\n", wc[i].qp_num, ibv_wc_status_str(wc[i].status)); + conn->client->stop(); + // but continue to handle events to purge the queue + } + if (wc[i].wr_id == 1) + { + // 1 = receive + auto b = conn->recv_buffers[conn->next_recv_buf]; + // Due to the credit-based flow control in RPC-RDMA, we can just remove that buffer and reuse it later + conn->recv_buffers.erase(conn->recv_buffers.begin()+conn->next_recv_buf, conn->recv_buffers.begin()+conn->next_recv_buf+1); + if (conn->cur_credit > 0 && !conn->client->stopped) + { + // Increase client refcount while the RPC call is being processed + conn->client->refs++; + conn->cur_credit--; + conn->handle_recv(b.buf, wc[i].byte_len); + } + else + { + fprintf(stderr, "Warning: NFS client credit exceeded for queue %d, stopping client\n", wc[i].qp_num); + conn->client->stop(); + } + } + else if (wc[i].wr_id == 2) + { + // 2 = send + auto rop = conn->outbox[0]; + conn->cur_send -= conn->outbox_wrs[0]; + conn->outbox.erase(conn->outbox.begin(), conn->outbox.begin()+1); + conn->outbox_wrs.erase(conn->outbox_wrs.begin(), conn->outbox_wrs.begin()+1); + conn->outbox_pos--; + // Retry send for corner cases of exceeded max_send_wr + conn->post_send(); + // Free rpc_op + conn->free_rdma_rpc_op(rop); + } + else if (wc[i].wr_id == 3) + { + // 3 = chunk read + auto rop = conn->chunk_inbox[0]; + conn->chunk_inbox.erase(conn->chunk_inbox.begin(), conn->chunk_inbox.begin()+1); + size_t read_chunk_count = 0; + for (auto cur = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_reads; cur; cur = cur->next) + read_chunk_count++; + conn->cur_send -= read_chunk_count; + conn->cur_rdma_reads -= read_chunk_count; + conn->client->handle_rpc_op(rop); + // Retry send for corner cases of exceeded max_send_wr + conn->post_send(); + } + else + { + fprintf(stderr, "BUG: unknown RDMA work completion with wr_id=%ju\n", wc[i].wr_id); + } + } + } while (event_count > 0); +} + +void nfs_rdma_conn_t::free_rdma_rpc_op(rpc_op_t *rop) +{ + if (rop->buffer) + { + rdma_malloc_free(ctx->alloc, rop->buffer); + rop->buffer = NULL; + } + auto rdma_chunk = xdr_get_rdma_chunk(rop->xdrs); + if (rdma_chunk) + { + rdma_malloc_free(ctx->alloc, rdma_chunk); + xdr_set_rdma_chunk(rop->xdrs, NULL); + } + ctx->proxy->free_xdr(rop->xdrs); + free(rop); + cur_credit++; + client->deref(); +} + +void nfs_rdma_conn_t::reuse_buffer(void *buf) +{ + auto ub = used_buffers.at(buf); + recv_buffers.push_back(ub); + post_recv(ub); +} + +// returns false if handling is done, returns true if handling is continued asynchronously +bool nfs_rdma_conn_t::handle_recv(void *buf, size_t len) +{ + XDR *xdrs = ctx->proxy->get_xdr(); + xdr_set_rdma(xdrs); + // Decode the RDMA-RPC header + rdma_msg rmsg; + if (!xdr_decode(xdrs, buf, len, (xdrproc_t)xdr_rdma_msg, &rmsg)) + { + // Invalid message, ignore it + fprintf(stderr, "Invalid RDMA-RPC header on connection 0x%jx, ignoring message\n", (uint64_t)this); +ignore_msg: + reuse_buffer(buf); + ctx->proxy->free_xdr(xdrs); + client->deref(); + return 0; + } + if (rmsg.rdma_vers != 1 || rmsg.rdma_body.proc != RDMA_MSG) + { + // Bad RDMA-RPC version or message type + fprintf( + stderr, "Unsupported RDMA-RPC version (%d) or message type (%d), sending %s\n", + rmsg.rdma_vers, rmsg.rdma_body.proc, + rmsg.rdma_vers != 1 ? "ERR_VERS" : "ERR_CHUNK" + ); + rpc_op_t *rop = (rpc_op_t*)malloc_or_die(sizeof(rpc_op_t)); + *rop = (rpc_op_t){ + .client = this, + .xdrs = xdrs, + .in_msg = { + .xid = rmsg.rdma_xid, + }, + .in_rdma_msg = rmsg, + .rdma_error = rmsg.rdma_vers != 1 ? ERR_VERS : ERR_CHUNK, + }; + rpc_queue_reply(rop); + // Incoming buffer isn't needed to handle request, so return 0 + return 0; + } + rpc_msg inmsg; + if (!xdr_rpc_msg(xdrs, &inmsg)) + { + // Invalid message, ignore it + fprintf(stderr, "Invalid RDMA-RPC body on connection 0x%jx, ignoring message\n", (uint64_t)this); + goto ignore_msg; + } + if (inmsg.xid != rmsg.rdma_xid) + { + fprintf(stderr, "RPC XID (%u) and RDMA XID (%u) mismatch on connection 0x%jx, ignoring message\n", inmsg.xid, rmsg.rdma_xid, (uint64_t)this); + goto ignore_msg; + } + if (inmsg.body.dir != RPC_CALL) + { + fprintf(stderr, "Non-RPC_CALL message direction (%d) on connection 0x%jx, ignoring message\n", inmsg.body.dir, (uint64_t)this); + goto ignore_msg; + } + rpc_op_t *rop = client->create_rpc_op(xdrs, buf, &inmsg, &rmsg); + if (!rop) + { + // No such procedure + return 0; + } + // Read chunks may only be provided for WRITE3 and SYMLINK3 + if (inmsg.body.cbody.prog == NFS_PROGRAM && + rmsg.rdma_body.rdma_msg.rdma_reads && + inmsg.body.cbody.proc != NFS3_WRITE && inmsg.body.cbody.proc != NFS3_SYMLINK) + { + fprintf(stderr, "Read chunk(s) are provided for non-WRITE or SYMLINK operation (procedure %u), sending ERR_CHUNK\n", inmsg.body.cbody.proc); + rop->rdma_error = ERR_CHUNK; + rpc_queue_reply(rop); + return 0; + } + // Write chunk may be provided only for READ3 and READLINK3 + if (inmsg.body.cbody.prog == NFS_PROGRAM && + rmsg.rdma_body.rdma_msg.rdma_writes && + inmsg.body.cbody.proc != NFS3_READ && inmsg.body.cbody.proc != NFS3_READLINK) + { + fprintf(stderr, "Write chunk(s) are provided for non-READ or READLINK operation (procedure %u), sending ERR_CHUNK\n", inmsg.body.cbody.proc); + rop->rdma_error = ERR_CHUNK; + rpc_queue_reply(rop); + return 0; + } + // Only 1 write chunk is supported + if (inmsg.body.cbody.prog == NFS_PROGRAM && + rmsg.rdma_body.rdma_msg.rdma_writes && ( + rmsg.rdma_body.rdma_msg.rdma_writes->next || + rmsg.rdma_body.rdma_msg.rdma_writes->entry.target.target_len != 1)) + { + int n = 0, seg = 0; + for (auto cur = rmsg.rdma_body.rdma_msg.rdma_writes; cur; cur = cur->next) + { + n++; + seg += cur->entry.target.target_len; + } + fprintf(stderr, "Only 1 write chunk with 1 segment may be provided, but we got %d chunks with %d segments, sending ERR_CHUNK\n", n, seg); + rop->rdma_error = ERR_CHUNK; + rpc_queue_reply(rop); + return 0; + } + // Process read chunk(s) + if (rmsg.rdma_body.rdma_msg.rdma_reads) + { + int r = post_chunk_reads(rop, true); + return (r >= 0 ? 1 : 0); + } + return client->handle_rpc_op(rop); +} + +int nfs_rdma_conn_t::post_chunk_reads(rpc_op_t *rop, bool push) +{ + rop->referenced = 1; + size_t read_chunk_size = 0; + size_t read_chunk_count = 0; + for (auto cur = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_reads; cur; cur = cur->next) + { + read_chunk_size += cur->entry.target.length; + read_chunk_count++; + } + if (read_chunk_count > max_rdma_reads) + { + fprintf(stderr, "Read chunk count (%zu) exceeds max_rdma_reads (%d), sending ERR_CHUNK\n", read_chunk_count, max_rdma_reads); + rop->rdma_error = ERR_CHUNK; + rpc_queue_reply(rop); + return -1; + } + if (cur_send+read_chunk_count > ctx->max_send_wr || + cur_rdma_reads+read_chunk_count > max_rdma_reads) + { + // Send it later + if (push) + chunk_read_postponed.push_back(rop); + return 0; + } + void *buf = rdma_malloc_alloc(ctx->alloc, read_chunk_size); + auto buf_lkey = rdma_malloc_get_lkey(ctx->alloc, buf); + ibv_sge chunk_sge[read_chunk_count]; + ibv_send_wr chunk_wr[read_chunk_count]; + size_t i = 0; + size_t pos = 0; + for (auto cur = rop->in_rdma_msg.rdma_body.rdma_msg.rdma_reads; cur; cur = cur->next) + { + chunk_sge[i] = { + .addr = (uintptr_t)((uint8_t*)buf + pos), + .length = cur->entry.target.length, + .lkey = buf_lkey, + }; + chunk_wr[i] = { + .wr_id = 3, // 3 is chunk read + .next = (i == read_chunk_count-1 ? NULL : &chunk_wr[i+1]), + .sg_list = &chunk_sge[i], + .num_sge = 1, + .opcode = IBV_WR_RDMA_READ, + .send_flags = (unsigned)(i == read_chunk_count-1 ? IBV_SEND_SIGNALED : 0), + .wr = { + .rdma = { + .remote_addr = cur->entry.target.offset, + .rkey = cur->entry.target.handle, + }, + }, + }; + pos += cur->entry.target.length; + i++; + } + assert(i == read_chunk_count); + assert(pos == read_chunk_size); + ibv_send_wr *bad_wr = NULL; + int err = ibv_post_send(id->qp, &chunk_wr[0], &bad_wr); + if (err || bad_wr) + { + fprintf(stderr, "Posting RDMA read failed: %s\n", strerror(err)); + exit(1); + } + xdr_set_rdma_chunk(rop->xdrs, buf); + chunk_inbox.push_back(rop); + cur_send += read_chunk_count; + cur_rdma_reads += read_chunk_count; + return 1; +} + +void *nfs_client_t::rdma_malloc(size_t size) +{ + return rdma_malloc_alloc(rdma_conn->ctx->alloc, size); +} + +void nfs_client_t::rdma_free(void *buf) +{ + rdma_malloc_free(rdma_conn->conn_dev->alloc, buf); +} + +void nfs_client_t::destroy_rdma_conn() +{ + if (rdma_conn) + { + delete rdma_conn; + rdma_conn = NULL; + } +} + +#endif diff --git a/src/nfs/proto/nfs.x b/src/nfs/proto/nfs.x index 738516a0..c02c9278 100644 --- a/src/nfs/proto/nfs.x +++ b/src/nfs/proto/nfs.x @@ -168,7 +168,7 @@ struct WRITE3args { offset3 offset; count3 count; stable_how stable; - opaque data<>; + opaque data<>; /* RDMA DDP-eligible */ }; typedef opaque writeverf3[NFS3_WRITEVERFSIZE]; @@ -409,7 +409,7 @@ struct READ3resok { post_op_attr file_attributes; count3 count; bool eof; - opaque data<>; + opaque data<>; /* RDMA DDP-eligible */ }; struct READ3resfail { @@ -514,7 +514,7 @@ typedef string nfspath3<>; struct symlinkdata3 { sattr3 symlink_attributes; - nfspath3 symlink_data; + nfspath3 symlink_data; /* RDMA DDP-eligible */ }; struct SYMLINK3args { @@ -546,7 +546,7 @@ struct READLINK3args { struct READLINK3resok { post_op_attr symlink_attributes; - nfspath3 data; + nfspath3 data; /* RDMA DDP-eligible */ }; struct READLINK3resfail { diff --git a/src/nfs/proto/nfs_xdr.cpp b/src/nfs/proto/nfs_xdr.cpp index 87451293..5897e6ad 100644 --- a/src/nfs/proto/nfs_xdr.cpp +++ b/src/nfs/proto/nfs_xdr.cpp @@ -272,7 +272,7 @@ xdr_WRITE3args (XDR *xdrs, WRITE3args *objp) return FALSE; if (!xdr_stable_how (xdrs, &objp->stable)) return FALSE; - if (!xdr_bytes(xdrs, &objp->data, ~0)) + if (!xdr_bytes(xdrs, &objp->data, ~0, true)) return FALSE; return TRUE; } @@ -829,7 +829,7 @@ xdr_READ3resok (XDR *xdrs, READ3resok *objp) return FALSE; if (!xdr_bool (xdrs, &objp->eof)) return FALSE; - if (!xdr_bytes(xdrs, &objp->data, ~0)) + if (!xdr_bytes(xdrs, &objp->data, ~0, true)) return FALSE; return TRUE; } @@ -1173,10 +1173,10 @@ xdr_PATHCONF3res (XDR *xdrs, PATHCONF3res *objp) } bool_t -xdr_nfspath3 (XDR *xdrs, nfspath3 *objp) +xdr_nfspath3 (XDR *xdrs, nfspath3 *objp, bool rdma_chunk) { - if (!xdr_string (xdrs, objp, ~0)) + if (!xdr_string (xdrs, objp, ~0, rdma_chunk)) return FALSE; return TRUE; } @@ -1187,7 +1187,7 @@ xdr_symlinkdata3 (XDR *xdrs, symlinkdata3 *objp) if (!xdr_sattr3 (xdrs, &objp->symlink_attributes)) return FALSE; - if (!xdr_nfspath3 (xdrs, &objp->symlink_data)) + if (!xdr_nfspath3 (xdrs, &objp->symlink_data, true)) return FALSE; return TRUE; } @@ -1259,7 +1259,7 @@ xdr_READLINK3resok (XDR *xdrs, READLINK3resok *objp) if (!xdr_post_op_attr (xdrs, &objp->symlink_attributes)) return FALSE; - if (!xdr_nfspath3 (xdrs, &objp->data)) + if (!xdr_nfspath3 (xdrs, &objp->data, true)) return FALSE; return TRUE; } diff --git a/src/nfs/proto/nfs_xdr.cpp.diff b/src/nfs/proto/nfs_xdr.cpp.diff new file mode 100644 index 00000000..fe22424e --- /dev/null +++ b/src/nfs/proto/nfs_xdr.cpp.diff @@ -0,0 +1,53 @@ +diff --git a/src/nfs/proto/nfs_xdr.cpp b/src/nfs/proto/nfs_xdr.cpp +index 87451293..5897e6ad 100644 +--- a/src/nfs/proto/nfs_xdr.cpp ++++ b/src/nfs/proto/nfs_xdr.cpp +@@ -272,7 +272,7 @@ xdr_WRITE3args (XDR *xdrs, WRITE3args *objp) + return FALSE; + if (!xdr_stable_how (xdrs, &objp->stable)) + return FALSE; +- if (!xdr_bytes(xdrs, &objp->data, ~0)) ++ if (!xdr_bytes(xdrs, &objp->data, ~0, true)) + return FALSE; + return TRUE; + } +@@ -829,7 +829,7 @@ xdr_READ3resok (XDR *xdrs, READ3resok *objp) + return FALSE; + if (!xdr_bool (xdrs, &objp->eof)) + return FALSE; +- if (!xdr_bytes(xdrs, &objp->data, ~0)) ++ if (!xdr_bytes(xdrs, &objp->data, ~0, true)) + return FALSE; + return TRUE; + } +@@ -1173,10 +1173,10 @@ xdr_PATHCONF3res (XDR *xdrs, PATHCONF3res *objp) + } + + bool_t +-xdr_nfspath3 (XDR *xdrs, nfspath3 *objp) ++xdr_nfspath3 (XDR *xdrs, nfspath3 *objp, bool rdma_chunk) + { + +- if (!xdr_string (xdrs, objp, ~0)) ++ if (!xdr_string (xdrs, objp, ~0, rdma_chunk)) + return FALSE; + return TRUE; + } +@@ -1187,7 +1187,7 @@ xdr_symlinkdata3 (XDR *xdrs, symlinkdata3 *objp) + + if (!xdr_sattr3 (xdrs, &objp->symlink_attributes)) + return FALSE; +- if (!xdr_nfspath3 (xdrs, &objp->symlink_data)) ++ if (!xdr_nfspath3 (xdrs, &objp->symlink_data, true)) + return FALSE; + return TRUE; + } +@@ -1259,7 +1259,7 @@ xdr_READLINK3resok (XDR *xdrs, READLINK3resok *objp) + + if (!xdr_post_op_attr (xdrs, &objp->symlink_attributes)) + return FALSE; +- if (!xdr_nfspath3 (xdrs, &objp->data)) ++ if (!xdr_nfspath3 (xdrs, &objp->data, true)) + return FALSE; + return TRUE; + } diff --git a/src/nfs/proto/rpc_impl.h b/src/nfs/proto/rpc_impl.h index 4acfdf87..965beb0f 100644 --- a/src/nfs/proto/rpc_impl.h +++ b/src/nfs/proto/rpc_impl.h @@ -1,6 +1,7 @@ #pragma once #include "rpc.h" +#include "rpc_rdma.h" struct rpc_op_t; @@ -27,12 +28,16 @@ inline bool operator < (const rpc_service_proc_t & a, const rpc_service_proc_t & return a.prog < b.prog || a.prog == b.prog && (a.vers < b.vers || a.vers == b.vers && a.proc < b.proc); } +struct rdma_msg; + struct rpc_op_t { void *client; - uint8_t *buffer; + void *buffer; XDR *xdrs; rpc_msg in_msg, out_msg; + rdma_msg in_rdma_msg; + rpc_rdma_errcode rdma_error; void *request; void *reply; xdrproc_t reply_fn; diff --git a/src/nfs/proto/rpc_rdma.h b/src/nfs/proto/rpc_rdma.h new file mode 100644 index 00000000..8801dd11 --- /dev/null +++ b/src/nfs/proto/rpc_rdma.h @@ -0,0 +1,144 @@ +/* + * Please do not edit this file. + * It was generated using rpcgen. + */ + +#ifndef _RPC_RDMA_H_RPCGEN +#define _RPC_RDMA_H_RPCGEN + +#include "xdr_impl.h" + + +#ifdef __cplusplus +extern "C" { +#endif + + +struct xdr_rdma_segment { + uint32_t handle; + uint32_t length; + uint64_t offset; +}; +typedef struct xdr_rdma_segment xdr_rdma_segment; + +struct xdr_read_chunk { + uint32_t position; + struct xdr_rdma_segment target; +}; +typedef struct xdr_read_chunk xdr_read_chunk; + +struct xdr_read_list { + struct xdr_read_chunk entry; + struct xdr_read_list *next; +}; +typedef struct xdr_read_list xdr_read_list; + +struct xdr_write_chunk { + struct { + u_int target_len; + struct xdr_rdma_segment *target_val; + } target; +}; +typedef struct xdr_write_chunk xdr_write_chunk; + +struct xdr_write_list { + struct xdr_write_chunk entry; + struct xdr_write_list *next; +}; +typedef struct xdr_write_list xdr_write_list; + +struct rpc_rdma_header { + struct xdr_read_list *rdma_reads; + struct xdr_write_list *rdma_writes; + struct xdr_write_chunk *rdma_reply; +}; +typedef struct rpc_rdma_header rpc_rdma_header; + +struct rpc_rdma_header_nomsg { + struct xdr_read_list *rdma_reads; + struct xdr_write_list *rdma_writes; + struct xdr_write_chunk *rdma_reply; +}; +typedef struct rpc_rdma_header_nomsg rpc_rdma_header_nomsg; + +struct rpc_rdma_header_padded { + uint32_t rdma_align; + uint32_t rdma_thresh; + struct xdr_read_list *rdma_reads; + struct xdr_write_list *rdma_writes; + struct xdr_write_chunk *rdma_reply; +}; +typedef struct rpc_rdma_header_padded rpc_rdma_header_padded; + +enum rpc_rdma_errcode { + ERR_VERS = 1, + ERR_CHUNK = 2, +}; +typedef enum rpc_rdma_errcode rpc_rdma_errcode; + +struct rpc_rdma_errvers { + uint32_t rdma_vers_low; + uint32_t rdma_vers_high; +}; +typedef struct rpc_rdma_errvers rpc_rdma_errvers; + +struct rpc_rdma_error { + rpc_rdma_errcode err; + union { + rpc_rdma_errvers range; + }; +}; +typedef struct rpc_rdma_error rpc_rdma_error; + +enum rdma_proc { + RDMA_MSG = 0, + RDMA_NOMSG = 1, + RDMA_MSGP = 2, + RDMA_DONE = 3, + RDMA_ERROR = 4, +}; +typedef enum rdma_proc rdma_proc; + +struct rdma_body { + rdma_proc proc; + union { + rpc_rdma_header rdma_msg; + rpc_rdma_header_nomsg rdma_nomsg; + rpc_rdma_header_padded rdma_msgp; + rpc_rdma_error rdma_error; + }; +}; +typedef struct rdma_body rdma_body; + +struct rdma_msg { + uint32_t rdma_xid; + uint32_t rdma_vers; + uint32_t rdma_credit; + rdma_body rdma_body; +}; +typedef struct rdma_msg rdma_msg; + +/* the xdr functions */ + + +extern bool_t xdr_xdr_rdma_segment (XDR *, xdr_rdma_segment*); +extern bool_t xdr_xdr_read_chunk (XDR *, xdr_read_chunk*); +extern bool_t xdr_xdr_read_list (XDR *, xdr_read_list*); +extern bool_t xdr_xdr_write_chunk (XDR *, xdr_write_chunk*); +extern bool_t xdr_xdr_write_list (XDR *, xdr_write_list*); +extern bool_t xdr_rpc_rdma_header (XDR *, rpc_rdma_header*); +extern bool_t xdr_rpc_rdma_header_nomsg (XDR *, rpc_rdma_header_nomsg*); +extern bool_t xdr_rpc_rdma_header_padded (XDR *, rpc_rdma_header_padded*); +extern bool_t xdr_rpc_rdma_errcode (XDR *, rpc_rdma_errcode*); +extern bool_t xdr_rpc_rdma_errvers (XDR *, rpc_rdma_errvers*); +extern bool_t xdr_rpc_rdma_error (XDR *, rpc_rdma_error*); +extern bool_t xdr_rdma_proc (XDR *, rdma_proc*); +extern bool_t xdr_rdma_body (XDR *, rdma_body*); +extern bool_t xdr_rdma_msg (XDR *, rdma_msg*); + + +#ifdef __cplusplus +} +#endif + +#endif /* !_RPC_RDMA_H_RPCGEN */ diff --git a/src/nfs/proto/rpc_rdma.x b/src/nfs/proto/rpc_rdma.x new file mode 100644 index 00000000..03d26a8d --- /dev/null +++ b/src/nfs/proto/rpc_rdma.x @@ -0,0 +1,166 @@ +/* RFC 8166 - Remote Direct Memory Access Transport for Remote Procedure Call Version 1 */ + +/* + * Copyright (c) 2010-2017 IETF Trust and the persons + * identified as authors of the code. All rights reserved. + * + * The authors of the code are: + * B. Callaghan, T. Talpey, and C. Lever + * + * Redistribution and use in source and binary forms, with + * or without modification, are permitted provided that the + * following conditions are met: + * + * - Redistributions of source code must retain the above + * copyright notice, this list of conditions and the + * following disclaimer. + * + * - Redistributions in binary form must reproduce the above + * copyright notice, this list of conditions and the + * following disclaimer in the documentation and/or other + * materials provided with the distribution. + * + * - Neither the name of Internet Society, IETF or IETF + * Trust, nor the names of specific contributors, may be + * used to endorse or promote products derived from this + * software without specific prior written permission. + * + * THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS + * AND CONTRIBUTORS "AS IS" AND ANY EXPRESS OR IMPLIED + * WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE + * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS + * FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO + * EVENT SHALL THE COPYRIGHT OWNER OR CONTRIBUTORS BE + * LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, + * EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT + * NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR + * SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS + * INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF + * LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, + * OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING + * IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF + * ADVISED OF THE POSSIBILITY OF SUCH DAMAGE. + */ + +/* + * Plain RDMA segment (Section 3.4.3) + */ +struct xdr_rdma_segment { + uint32_t handle; /* Registered memory handle */ + uint32_t length; /* Length of the chunk in bytes */ + uint64_t offset; /* Chunk virtual address or offset */ +}; + +/* + * RDMA read segment (Section 3.4.5) + */ +struct xdr_read_chunk { + uint32_t position; /* Position in XDR stream */ + struct xdr_rdma_segment target; +}; + +/* + * Read list (Section 4.3.1) + */ +struct xdr_read_list { + struct xdr_read_chunk entry; + struct xdr_read_list *next; +}; + +/* + * Write chunk (Section 3.4.6) + */ +struct xdr_write_chunk { + struct xdr_rdma_segment target<>; +}; + +/* + * Write list (Section 4.3.2) + */ +struct xdr_write_list { + struct xdr_write_chunk entry; + struct xdr_write_list *next; +}; + +/* + * Chunk lists (Section 4.3) + */ +struct rpc_rdma_header { + struct xdr_read_list *rdma_reads; + struct xdr_write_list *rdma_writes; + struct xdr_write_chunk *rdma_reply; + /* rpc body follows */ +}; + +struct rpc_rdma_header_nomsg { + struct xdr_read_list *rdma_reads; + struct xdr_write_list *rdma_writes; + struct xdr_write_chunk *rdma_reply; +}; + +/* Not to be used */ +struct rpc_rdma_header_padded { + uint32_t rdma_align; + uint32_t rdma_thresh; + struct xdr_read_list *rdma_reads; + struct xdr_write_list *rdma_writes; + struct xdr_write_chunk *rdma_reply; + /* rpc body follows */ +}; + +/* + * Error handling (Section 4.5) + */ +enum rpc_rdma_errcode { + ERR_VERS = 1, /* Value fixed for all versions */ + ERR_CHUNK = 2 +}; + +/* Structure fixed for all versions */ +struct rpc_rdma_errvers { + uint32_t rdma_vers_low; + uint32_t rdma_vers_high; +}; + +union rpc_rdma_error switch (rpc_rdma_errcode err) { + case ERR_VERS: + rpc_rdma_errvers range; + case ERR_CHUNK: + void; +}; + +/* + * Procedures (Section 4.2.4) + */ +enum rdma_proc { + RDMA_MSG = 0, /* Value fixed for all versions */ + RDMA_NOMSG = 1, /* Value fixed for all versions */ + RDMA_MSGP = 2, /* Not to be used */ + RDMA_DONE = 3, /* Not to be used */ + RDMA_ERROR = 4 /* Value fixed for all versions */ +}; + +/* The position of the proc discriminator field is + * fixed for all versions */ +union rdma_body switch (rdma_proc proc) { + case RDMA_MSG: + rpc_rdma_header rdma_msg; + case RDMA_NOMSG: + rpc_rdma_header_nomsg rdma_nomsg; + case RDMA_MSGP: /* Not to be used */ + rpc_rdma_header_padded rdma_msgp; + case RDMA_DONE: /* Not to be used */ + void; + case RDMA_ERROR: + rpc_rdma_error rdma_error; +}; + +/* + * Fixed header fields (Section 4.2) + */ +struct rdma_msg { + uint32_t rdma_xid; /* Position fixed for all versions */ + uint32_t rdma_vers; /* Position fixed for all versions */ + uint32_t rdma_credit; /* Position fixed for all versions */ + rdma_body rdma_body; +}; diff --git a/src/nfs/proto/rpc_rdma_xdr.cpp b/src/nfs/proto/rpc_rdma_xdr.cpp new file mode 100644 index 00000000..ba1cb466 --- /dev/null +++ b/src/nfs/proto/rpc_rdma_xdr.cpp @@ -0,0 +1,200 @@ +/* + * Please do not edit this file. + * It was generated using rpcgen. + */ + +#include "rpc_rdma.h" +#include "xdr_impl_inline.h" + +bool_t +xdr_xdr_rdma_segment (XDR *xdrs, xdr_rdma_segment *objp) +{ + + if (!xdr_uint32_t (xdrs, &objp->handle)) + return FALSE; + if (!xdr_uint32_t (xdrs, &objp->length)) + return FALSE; + if (!xdr_uint64_t (xdrs, &objp->offset)) + return FALSE; + return TRUE; +} + +bool_t +xdr_xdr_read_chunk (XDR *xdrs, xdr_read_chunk *objp) +{ + + if (!xdr_uint32_t (xdrs, &objp->position)) + return FALSE; + if (!xdr_xdr_rdma_segment (xdrs, &objp->target)) + return FALSE; + return TRUE; +} + +bool_t +xdr_xdr_read_list (XDR *xdrs, xdr_read_list *objp) +{ + + if (!xdr_xdr_read_chunk (xdrs, &objp->entry)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->next, sizeof (xdr_read_list), (xdrproc_t) xdr_xdr_read_list)) + return FALSE; + return TRUE; +} + +bool_t +xdr_xdr_write_chunk (XDR *xdrs, xdr_write_chunk *objp) +{ + + if (!xdr_array (xdrs, (char **)&objp->target.target_val, (u_int *) &objp->target.target_len, ~0, + sizeof (xdr_rdma_segment), (xdrproc_t) xdr_xdr_rdma_segment)) + return FALSE; + return TRUE; +} + +bool_t +xdr_xdr_write_list (XDR *xdrs, xdr_write_list *objp) +{ + + if (!xdr_xdr_write_chunk (xdrs, &objp->entry)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->next, sizeof (xdr_write_list), (xdrproc_t) xdr_xdr_write_list)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rpc_rdma_header (XDR *xdrs, rpc_rdma_header *objp) +{ + + if (!xdr_pointer (xdrs, (char **)&objp->rdma_reads, sizeof (xdr_read_list), (xdrproc_t) xdr_xdr_read_list)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_writes, sizeof (xdr_write_list), (xdrproc_t) xdr_xdr_write_list)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_reply, sizeof (xdr_write_chunk), (xdrproc_t) xdr_xdr_write_chunk)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rpc_rdma_header_nomsg (XDR *xdrs, rpc_rdma_header_nomsg *objp) +{ + + if (!xdr_pointer (xdrs, (char **)&objp->rdma_reads, sizeof (xdr_read_list), (xdrproc_t) xdr_xdr_read_list)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_writes, sizeof (xdr_write_list), (xdrproc_t) xdr_xdr_write_list)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_reply, sizeof (xdr_write_chunk), (xdrproc_t) xdr_xdr_write_chunk)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rpc_rdma_header_padded (XDR *xdrs, rpc_rdma_header_padded *objp) +{ + + if (!xdr_uint32_t (xdrs, &objp->rdma_align)) + return FALSE; + if (!xdr_uint32_t (xdrs, &objp->rdma_thresh)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_reads, sizeof (xdr_read_list), (xdrproc_t) xdr_xdr_read_list)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_writes, sizeof (xdr_write_list), (xdrproc_t) xdr_xdr_write_list)) + return FALSE; + if (!xdr_pointer (xdrs, (char **)&objp->rdma_reply, sizeof (xdr_write_chunk), (xdrproc_t) xdr_xdr_write_chunk)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rpc_rdma_errcode (XDR *xdrs, rpc_rdma_errcode *objp) +{ + + if (!xdr_enum (xdrs, (enum_t *) objp)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rpc_rdma_errvers (XDR *xdrs, rpc_rdma_errvers *objp) +{ + + if (!xdr_uint32_t (xdrs, &objp->rdma_vers_low)) + return FALSE; + if (!xdr_uint32_t (xdrs, &objp->rdma_vers_high)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rpc_rdma_error (XDR *xdrs, rpc_rdma_error *objp) +{ + + if (!xdr_rpc_rdma_errcode (xdrs, &objp->err)) + return FALSE; + switch (objp->err) { + case ERR_VERS: + if (!xdr_rpc_rdma_errvers (xdrs, &objp->range)) + return FALSE; + break; + case ERR_CHUNK: + break; + default: + return FALSE; + } + return TRUE; +} + +bool_t +xdr_rdma_proc (XDR *xdrs, rdma_proc *objp) +{ + + if (!xdr_enum (xdrs, (enum_t *) objp)) + return FALSE; + return TRUE; +} + +bool_t +xdr_rdma_body (XDR *xdrs, rdma_body *objp) +{ + + if (!xdr_rdma_proc (xdrs, &objp->proc)) + return FALSE; + switch (objp->proc) { + case RDMA_MSG: + if (!xdr_rpc_rdma_header (xdrs, &objp->rdma_msg)) + return FALSE; + break; + case RDMA_NOMSG: + if (!xdr_rpc_rdma_header_nomsg (xdrs, &objp->rdma_nomsg)) + return FALSE; + break; + case RDMA_MSGP: + if (!xdr_rpc_rdma_header_padded (xdrs, &objp->rdma_msgp)) + return FALSE; + break; + case RDMA_DONE: + break; + case RDMA_ERROR: + if (!xdr_rpc_rdma_error (xdrs, &objp->rdma_error)) + return FALSE; + break; + default: + return FALSE; + } + return TRUE; +} + +bool_t +xdr_rdma_msg (XDR *xdrs, rdma_msg *objp) +{ + + if (!xdr_uint32_t (xdrs, &objp->rdma_xid)) + return FALSE; + if (!xdr_uint32_t (xdrs, &objp->rdma_vers)) + return FALSE; + if (!xdr_uint32_t (xdrs, &objp->rdma_credit)) + return FALSE; + if (!xdr_rdma_body (xdrs, &objp->rdma_body)) + return FALSE; + return TRUE; +} diff --git a/src/nfs/proto/run-rpcgen.sh b/src/nfs/proto/run-rpcgen.sh index 9224d174..840bfb84 100755 --- a/src/nfs/proto/run-rpcgen.sh +++ b/src/nfs/proto/run-rpcgen.sh @@ -46,3 +46,5 @@ run_rpcgen() { run_rpcgen nfs run_rpcgen rpc run_rpcgen portmap +run_rpcgen rpc_rdma +patch nfs_xdr.cpp < nfs_xdr.cpp.diff diff --git a/src/nfs/proto/xdr_impl.cpp b/src/nfs/proto/xdr_impl.cpp index 1c264898..61486982 100644 --- a/src/nfs/proto/xdr_impl.cpp +++ b/src/nfs/proto/xdr_impl.cpp @@ -16,6 +16,22 @@ void xdr_destroy(XDR* xdrs) delete xdrs; } +void xdr_set_rdma(XDR *xdrs) +{ + xdrs->rdma = true; +} + +void xdr_set_rdma_chunk(XDR *xdrs, void *chunk) +{ + assert(!xdrs->rdma_chunk || !chunk); + xdrs->rdma_chunk = chunk; +} + +void* xdr_get_rdma_chunk(XDR *xdrs) +{ + return xdrs->rdma_chunk; +} + void xdr_reset(XDR *xdrs) { for (auto buf: xdrs->allocs) @@ -23,6 +39,9 @@ void xdr_reset(XDR *xdrs) free(buf); } xdrs->buf = NULL; + xdrs->rdma = false; + xdrs->rdma_chunk = NULL; + xdrs->rdma_chunk_used = false; xdrs->avail = 0; xdrs->allocs.resize(0); xdrs->in_linked_list.resize(0); @@ -45,6 +64,20 @@ int xdr_encode(XDR *xdrs, xdrproc_t fn, void *data) return fn(xdrs, data); } +size_t xdr_encode_get_size(XDR *xdrs) +{ + size_t len = 0; + for (auto & buf: xdrs->buf_list) + { + len += buf.iov_len; + } + if (xdrs->last_end < xdrs->cur_out.size()) + { + len += xdrs->cur_out.size() - xdrs->last_end; + } + return len; +} + void xdr_encode_finish(XDR *xdrs, iovec **iov_list, unsigned *iov_count) { if (xdrs->last_end < xdrs->cur_out.size()) @@ -83,6 +116,18 @@ void xdr_add_malloc(XDR *xdrs, void *buf) xdrs->allocs.push_back(buf); } +void xdr_del_malloc(XDR *xdrs, void *buf) +{ + for (int i = 0; i < xdrs->allocs.size(); i++) + { + if (xdrs->allocs[i] == buf) + { + xdrs->allocs.erase(xdrs->allocs.begin()+i); + break; + } + } +} + xdr_string_t xdr_copy_string(XDR *xdrs, const std::string & str) { char *cp = (char*)malloc_or_die(str.size()+1); diff --git a/src/nfs/proto/xdr_impl.h b/src/nfs/proto/xdr_impl.h index badfe75f..3e79e326 100644 --- a/src/nfs/proto/xdr_impl.h +++ b/src/nfs/proto/xdr_impl.h @@ -55,6 +55,15 @@ void xdr_destroy(XDR* xdrs); // Free resources from any previous xdr_decode/xdr_encode calls void xdr_reset(XDR *xdrs); +// Mark XDR as used for RDMA +void xdr_set_rdma(XDR *xdrs); + +// Set (single) RDMA chunk buffer for this xdr before decoding an RDMA message +void xdr_set_rdma_chunk(XDR *xdrs, void *chunk); + +// Get the current RDMA chunk buffer +void* xdr_get_rdma_chunk(XDR *xdrs); + // Try to decode bytes from buffer using // Result may contain memory allocations that will be valid until the next call to xdr_{reset,destroy,decode,encode} int xdr_decode(XDR *xdrs, void *buf, unsigned size, xdrproc_t fn, void *data); @@ -64,6 +73,9 @@ int xdr_decode(XDR *xdrs, void *buf, unsigned size, xdrproc_t fn, void *data); // May be called multiple times to encode multiple parts of the same message int xdr_encode(XDR *xdrs, xdrproc_t fn, void *data); +// Get current size of encoded data in +size_t xdr_encode_get_size(XDR *xdrs); + // Get the result of previous xdr_encodes as a list of 's // in (start) and (count). // The resulting iov_list is valid until the next call to xdr_{reset,destroy}. @@ -74,6 +86,9 @@ void xdr_encode_finish(XDR *xdrs, iovec **iov_list, unsigned *iov_count); // Remember an allocated buffer to free it later on xdr_reset() or xdr_destroy() void xdr_add_malloc(XDR *xdrs, void *buf); +// Remove an allocated buffer from XDR +void xdr_del_malloc(XDR *xdrs, void *buf); + xdr_string_t xdr_copy_string(XDR *xdrs, const std::string & str); xdr_string_t xdr_copy_string(XDR *xdrs, const char *str); diff --git a/src/nfs/proto/xdr_impl_inline.h b/src/nfs/proto/xdr_impl_inline.h index 63e1a4bd..9e78d70a 100644 --- a/src/nfs/proto/xdr_impl_inline.h +++ b/src/nfs/proto/xdr_impl_inline.h @@ -28,6 +28,19 @@ // RPC over TCP: // // BE 32bit length, then rpc_msg, then the procedure message itself +// +// RPC over RDMA: +// RFC 8166 - Remote Direct Memory Access Transport for Remote Procedure Call Version 1 +// RFC 8267 - Network File System (NFS) Upper-Layer Binding to RPC-over-RDMA Version 1 +// RFC 8797 - Remote Direct Memory Access - Connection Manager (RDMA-CM) Private Data for RPC-over-RDMA Version 1 +// message is received in an RDMA Receive operation +// message: list of read chunks, list of write chunks, optional reply write chunk, then actual RPC body if present +// read chunk: BE 32bit position, BE 32bit registered memory key, BE 32bit length, BE 64bit offset +// write chunk: BE 32bit registered memory key, BE 32bit length, BE 64bit offset +// in reality for NFS 3.0: only 1 read chunk in write3 and symlink3, only 1 write chunk in read3 and readlink3 +// read chunk is read by the server using RDMA Read from the client memory after receiving RPC request +// write chunk is pushed by the server using RDMA Write to the client memory before sending RPC reply +// connection is established using RDMA-CM at default port 20049 #pragma once @@ -35,6 +48,7 @@ #include #include +#include #include #include "malloc_or_die.h" @@ -61,6 +75,9 @@ struct xdr_linked_list_t struct XDR { int x_op; + bool rdma = false; + void *rdma_chunk = NULL; + bool rdma_chunk_used = false; // For decoding: uint8_t *buf = NULL; @@ -106,13 +123,22 @@ inline int xdr_opaque(XDR *xdrs, void *data, uint32_t len) return 1; } -inline int xdr_bytes(XDR *xdrs, xdr_string_t *data, uint32_t maxlen) +inline int xdr_bytes(XDR *xdrs, xdr_string_t *data, uint32_t maxlen, bool rdma_chunk = false) { if (xdrs->x_op == XDR_DECODE) { if (xdrs->avail < 4) return 0; uint32_t len = be32toh(*((uint32_t*)xdrs->buf)); + if (rdma_chunk && xdrs->rdma && xdrs->rdma_chunk) + { + // Take (only a single) RDMA chunk from xdrs->rdma_chunk while decoding + assert(!xdrs->rdma_chunk_used); + xdrs->rdma_chunk_used = true; + data->data = (char*)xdrs->rdma_chunk; + data->size = len; + return 1; + } uint32_t padded = len_pad4(len); if (xdrs->avail < 4+padded) return 0; @@ -123,7 +149,8 @@ inline int xdr_bytes(XDR *xdrs, xdr_string_t *data, uint32_t maxlen) } else { - if (data->size < XDR_COPY_LENGTH) + // Always encode RDMA chunks as separate iovecs + if (data->size < XDR_COPY_LENGTH && (!rdma_chunk || !xdrs->rdma)) { unsigned old = xdrs->cur_out.size(); xdrs->cur_out.resize(old + 4+data->size); @@ -146,8 +173,9 @@ inline int xdr_bytes(XDR *xdrs, xdr_string_t *data, uint32_t maxlen) .iov_len = data->size, }); } - if (data->size & 3) + if ((data->size & 3) && (!rdma_chunk || !xdrs->rdma)) { + // No padding for RDMA chunks int pad = 4-(data->size & 3); unsigned old = xdrs->cur_out.size(); xdrs->cur_out.resize(old+pad); @@ -158,9 +186,9 @@ inline int xdr_bytes(XDR *xdrs, xdr_string_t *data, uint32_t maxlen) return 1; } -inline int xdr_string(XDR *xdrs, xdr_string_t *data, uint32_t maxlen) +inline int xdr_string(XDR *xdrs, xdr_string_t *data, uint32_t maxlen, bool rdma_chunk = false) { - return xdr_bytes(xdrs, data, maxlen); + return xdr_bytes(xdrs, data, maxlen, rdma_chunk); } inline int xdr_u_int(XDR *xdrs, void *data) @@ -182,6 +210,11 @@ inline int xdr_u_int(XDR *xdrs, void *data) return 1; } +inline int xdr_uint32_t(XDR *xdrs, void *data) +{ + return xdr_u_int(xdrs, data); +} + inline int xdr_enum(XDR *xdrs, void *data) { return xdr_u_int(xdrs, data); diff --git a/src/nfs/rdma_alloc.cpp b/src/nfs/rdma_alloc.cpp new file mode 100644 index 00000000..a2e027f2 --- /dev/null +++ b/src/nfs/rdma_alloc.cpp @@ -0,0 +1,216 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 (see README.md for details) +// +// Simple & stupid RDMA-enabled memory allocator (allocates buffers within ibv_mr's) + +#include +#include +#include +#include +#include "rdma_alloc.h" +#include "malloc_or_die.h" + +struct rdma_region_t +{ + void *buf = NULL; + size_t len = 0; + ibv_mr *mr = NULL; +}; + +struct rdma_frag_t +{ + rdma_region_t *rgn = NULL; + size_t len = 0; + bool is_free = false; +}; + +struct rdma_free_t +{ + size_t len = 0; + void *buf = NULL; +}; + +inline bool operator < (const rdma_free_t &a, const rdma_free_t &b) +{ + return a.len < b.len || a.len == b.len && a.buf < b.buf; +} + +struct rdma_allocator_t +{ + size_t rdma_alloc_size = 1048576; + size_t rdma_max_unused = 500*1048576; + int rdma_access = IBV_ACCESS_LOCAL_WRITE; + ibv_pd *pd = NULL; + + std::set regions; + std::map frags; + std::set freelist; + size_t freebuffers = 0; +}; + +rdma_allocator_t *rdma_malloc_create(ibv_pd *pd, size_t rdma_alloc_size, size_t rdma_max_unused, int rdma_access) +{ + rdma_allocator_t *self = new rdma_allocator_t(); + self->pd = pd; + self->rdma_alloc_size = rdma_alloc_size ? rdma_alloc_size : 1048576; + self->rdma_max_unused = rdma_max_unused ? rdma_max_unused : 500*1048576; + self->rdma_access = rdma_access; + return self; +} + +static void rdma_malloc_free_unused_buffers(rdma_allocator_t *self, size_t max_unused, bool force) +{ + auto free_it = self->freelist.end(); + if (free_it == self->freelist.begin()) + return; + free_it--; + do + { + auto frag_it = self->frags.find(free_it->buf); + assert(frag_it != self->frags.end()); + if (frag_it->second.len != frag_it->second.rgn->len) + { + if (force) + { + fprintf(stderr, "BUG: Attempt to destroy RDMA allocator while buffers are not freed yet\n"); + abort(); + } + break; + } + self->freebuffers -= frag_it->second.rgn->len; + ibv_dereg_mr(frag_it->second.rgn->mr); + free(frag_it->second.rgn); + self->regions.erase(frag_it->second.rgn); + self->frags.erase(frag_it); + if (free_it == self->freelist.begin()) + { + self->freelist.erase(free_it); + break; + } + self->freelist.erase(free_it--); + } while (self->freebuffers > max_unused); +} + +void rdma_malloc_destroy(rdma_allocator_t *self) +{ + rdma_malloc_free_unused_buffers(self, 0, true); + assert(!self->freebuffers); + assert(!self->regions.size()); + assert(!self->frags.size()); + assert(!self->freelist.size()); + delete self; +} + +void *rdma_malloc_alloc(rdma_allocator_t *self, size_t size) +{ + auto it = self->freelist.lower_bound((rdma_free_t){ .len = size }); + if (it == self->freelist.end()) + { + // round size up to rdma_malloc_size (1 MB) + size_t alloc_size = ((size + self->rdma_alloc_size - 1) / self->rdma_alloc_size) * self->rdma_alloc_size; + rdma_region_t *r = (rdma_region_t*)malloc_or_die(alloc_size + sizeof(rdma_region_t)); + r->buf = r+1; + r->len = alloc_size; + r->mr = ibv_reg_mr(self->pd, r->buf, r->len, self->rdma_access); + if (!r->mr) + { + fprintf(stderr, "Failed to register RDMA memory region: %s\n", strerror(errno)); + exit(1); + } + self->regions.insert(r); + self->frags[r->buf] = (rdma_frag_t){ .rgn = r, .len = alloc_size, .is_free = true }; + it = self->freelist.insert((rdma_free_t){ .len = alloc_size, .buf = r->buf }).first; + self->freebuffers += alloc_size; + } + void *ptr = it->buf; + auto & frag = self->frags.at(ptr); + self->freelist.erase(it); + assert(frag.len >= size && frag.is_free); + if (frag.len == frag.rgn->len) + { + self->freebuffers -= frag.rgn->len; + } + if (frag.len == size) + { + frag.is_free = false; + } + else + { + frag.len -= size; + ptr = (uint8_t*)ptr + frag.len; + self->freelist.insert((rdma_free_t){ .len = frag.len, .buf = frag.rgn->buf }); + self->frags[ptr] = (rdma_frag_t){ .rgn = frag.rgn, .len = size, .is_free = false }; + } + return ptr; +} + +void rdma_malloc_free(rdma_allocator_t *self, void *buf) +{ + auto frag_it = self->frags.find(buf); + if (frag_it == self->frags.end()) + { + fprintf(stderr, "BUG: Attempt to double-free RDMA buffer fragment 0x%jx\n", (size_t)buf); + return; + } + auto prev_it = frag_it, next_it = frag_it; + if (frag_it != self->frags.begin()) + prev_it--; + next_it++; + bool merge_back = prev_it != frag_it && + prev_it->second.is_free && + prev_it->second.rgn == frag_it->second.rgn && + (uint8_t*)prev_it->first+prev_it->second.len == frag_it->first; + bool merge_next = next_it != self->frags.end() && + next_it->second.is_free && + next_it->second.rgn == frag_it->second.rgn && + next_it->first == (uint8_t*)frag_it->first+frag_it->second.len; + if (merge_back && merge_next) + { + prev_it->second.len += frag_it->second.len + next_it->second.len; + self->freelist.erase((rdma_free_t){ .len = next_it->second.len, .buf = next_it->first }); + self->frags.erase(next_it); + self->frags.erase(frag_it); + frag_it = prev_it; + } + else if (merge_back) + { + prev_it->second.len += frag_it->second.len; + self->frags.erase(frag_it); + frag_it = prev_it; + } + else if (merge_next) + { + frag_it->second.is_free = true; + frag_it->second.len += next_it->second.len; + self->freelist.erase((rdma_free_t){ .len = next_it->second.len, .buf = next_it->first }); + self->frags.erase(next_it); + } + else + { + frag_it->second.is_free = true; + self->freelist.insert((rdma_free_t){ .len = frag_it->second.len, .buf = frag_it->first }); + } + assert(frag_it->second.len <= frag_it->second.rgn->len); + if (frag_it->second.len == frag_it->second.rgn->len) + { + // The whole buffer is freed + self->freebuffers += frag_it->second.rgn->len; + if (self->freebuffers > self->rdma_max_unused) + { + rdma_malloc_free_unused_buffers(self, self->rdma_max_unused, false); + } + } +} + +uint32_t rdma_malloc_get_lkey(rdma_allocator_t *self, void *buf) +{ + auto frag_it = self->frags.upper_bound(buf); + if (frag_it != self->frags.begin()) + { + frag_it--; + if ((uint8_t*)frag_it->first + frag_it->second.len > buf) + return frag_it->second.rgn->mr->lkey; + } + fprintf(stderr, "BUG: Attempt to use an unknown RDMA buffer fragment 0x%zx\n", (size_t)buf); + abort(); +} diff --git a/src/nfs/rdma_alloc.h b/src/nfs/rdma_alloc.h new file mode 100644 index 00000000..83ed6a2c --- /dev/null +++ b/src/nfs/rdma_alloc.h @@ -0,0 +1,17 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 (see README.md for details) +// +// Simple & stupid RDMA-enabled memory allocator (allocates buffers within ibv_mr's) + +#pragma once + +#include +#include + +struct rdma_allocator_t; + +rdma_allocator_t *rdma_malloc_create(ibv_pd *pd, size_t rdma_alloc_size, size_t rdma_max_unused, int rdma_access); +void rdma_malloc_destroy(rdma_allocator_t *self); +void *rdma_malloc_alloc(rdma_allocator_t *self, size_t size); +void rdma_malloc_free(rdma_allocator_t *self, void *buf); +uint32_t rdma_malloc_get_lkey(rdma_allocator_t *self, void *buf);