Implement NFS RDMA support

This commit is contained in:
Vitaliy Filippov
2024-12-11 21:09:36 +03:00
parent 64db31ec10
commit 1dbbb0c3f8
21 changed files with 2247 additions and 107 deletions
+4
View File
@@ -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
+4
View File
@@ -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}
)
+1 -2
View File
@@ -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;
+5 -1
View File
@@ -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 },
},
};
}
+5 -6
View File
@@ -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;
+210 -77
View File
@@ -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 <NAME> | --block) start\n"
" Start network NFS server. Options:\n"
" --bind <IP> bind service to <IP> address (default 0.0.0.0)\n"
" --port <PORT> use port <PORT> for NFS services (default is 2049)\n"
" --portmap 0 do not listen on port 111 (portmap/rpcbind, requires root)\n"
" --bind <IP> bind service to <IP> address (default 0.0.0.0)\n"
" --port <PORT> use port <PORT> for NFS services (default is 2049)\n"
" --portmap 0 do not listen on port 111 (portmap/rpcbind, requires root)\n"
" --nfs_rdma <PORT> enable NFS-RDMA at RDMA-CM port <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 <NAME> 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();
+34 -5
View File
@@ -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<int, nfs_client_t*> rpc_clients;
nfs_rdma_context_t* rdma_context = NULL;
std::set<nfs_client_t*> rpc_clients;
std::vector<XDR*> 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<rpc_service_proc_t> proc_table;
nfs_rdma_conn_t *rdma_conn = NULL;
// <TCP>
int nfs_fd = -1;
int epoll_events = 0;
// Read state
rpc_cur_buffer_t cur_buffer = { 0 };
std::map<uint8_t*, rpc_used_buffer_t> used_buffers;
std::map<void*, rpc_used_buffer_t> used_buffers;
std::vector<rpc_free_buffer_t> 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);
// </TCP>
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();
};
File diff suppressed because it is too large Load Diff
+4 -4
View File
@@ -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 {
+6 -6
View File
@@ -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;
}
+53
View File
@@ -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;
}
+6 -1
View File
@@ -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;
+144
View File
@@ -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 */
+166
View File
@@ -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;
};
+200
View File
@@ -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;
}
+2
View File
@@ -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
+45
View File
@@ -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);
+15
View File
@@ -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 <size> bytes from buffer <buf> using <fn>
// 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 <xdrs>
size_t xdr_encode_get_size(XDR *xdrs);
// Get the result of previous xdr_encodes as a list of <struct iovec>'s
// in <iov_list> (start) and <iov_count> (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);
+38 -5
View File
@@ -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 <string.h>
#include <endian.h>
#include <assert.h>
#include <vector>
#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);
+216
View File
@@ -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 <stdio.h>
#include <assert.h>
#include <map>
#include <set>
#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<rdma_region_t*> regions;
std::map<void*, rdma_frag_t> frags;
std::set<rdma_free_t> 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();
}
+17
View File
@@ -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 <infiniband/verbs.h>
#include <stdint.h>
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);