Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,19 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
per-stream response field-count and name/value-byte limits across final
headers and trailers, resetting an offending stream before the field is
copied (#998).
- **Non-blocking RDMA endpoint teardown**: `rdma_ibverbs::endpoint` destruction
no longer sleeps on scheduler workers, and the new awaitable `shutdown()` joins
the CQ pump before releasing verbs resources. Provider teardown failures are
reported without discarding ownership, allowing a serialized retry; checked
QP destruction no longer loses provider failures through `rdma_destroy_qp()`.
Shutdown is terminal for data-path use and requires the pump scheduler to
remain live until completion. Fatal undrainable eager operations can call
`abandon_outstanding()` while their state and buffers remain alive; only a
`true` result confirms QP destruction before the remaining verbs resources
are intentionally relinquished.
Post-shutdown calls to `conn()`, `cas()`, `fetch_add()`, and `dispatcher_ref()`
now throw `std::logic_error` instead of retaining their former `noexcept`
specification (#975).
- **Bounded raw WebSocket parsing by default**: `websocket::frame_parser` now
applies the same 16 MiB aggregate message limit as the high-level client and
server configurations. Direct parser users must explicitly set a limit of
Expand Down
152 changes: 81 additions & 71 deletions examples/rdma_gpu_bw.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -192,38 +192,42 @@ coro::task<int> async_main(int argc, char* argv[]) {
ep.start_cq_pump(*sched);
std::printf("Client connected.\n");

// Allocate GPU buffer and register
elio::rdma_cuda::gpu_memory_region gpu_mr{
ep.pd(), cfg.size,
IBV_ACCESS_LOCAL_WRITE | IBV_ACCESS_REMOTE_WRITE |
IBV_ACCESS_REMOTE_READ};

// Exchange remote info via SEND/RECV
auto remote = gpu_mr.remote();
xchg_msg msg_out{remote.addr, remote.length, remote.rkey};

auto send_mr = ep.register_buffer(&msg_out, sizeof(msg_out),
IBV_ACCESS_LOCAL_WRITE);
auto send_wc = co_await ep.conn().send(send_mr.view());
if (!send_wc.ok()) {
std::fprintf(stderr, "server send failed status=%d\n",
static_cast<int>(send_wc.status));
co_return 1;
int result = 0;
{
// This scope owns every MR; it ends before endpoint shutdown.
elio::rdma_cuda::gpu_memory_region gpu_mr{
ep.pd(), cfg.size,
IBV_ACCESS_LOCAL_WRITE | IBV_ACCESS_REMOTE_WRITE |
IBV_ACCESS_REMOTE_READ};

auto remote = gpu_mr.remote();
xchg_msg msg_out{remote.addr, remote.length, remote.rkey};

auto send_mr = ep.register_buffer(&msg_out, sizeof(msg_out),
IBV_ACCESS_LOCAL_WRITE);
auto send_wc = co_await ep.conn().send(send_mr.view());
if (!send_wc.ok()) {
std::fprintf(stderr, "server send failed status=%d\n",
static_cast<int>(send_wc.status));
result = 1;
} else {
xchg_msg msg_in{};
auto recv_mr = ep.register_buffer(&msg_in, sizeof(msg_in),
IBV_ACCESS_LOCAL_WRITE);
auto recv_wc = co_await ep.conn().recv(recv_mr.view());
if (!recv_wc.ok()) {
std::fprintf(stderr, "server recv failed status=%d\n",
static_cast<int>(recv_wc.status));
result = 1;
}
}
}

// Wait for client "done" signal
xchg_msg msg_in{};
auto recv_mr = ep.register_buffer(&msg_in, sizeof(msg_in),
IBV_ACCESS_LOCAL_WRITE);
auto recv_wc = co_await ep.conn().recv(recv_mr.view());
if (!recv_wc.ok()) {
std::fprintf(stderr, "server recv failed status=%d\n",
static_cast<int>(recv_wc.status));
co_return 1;
co_await ep.shutdown();
if (result == 0) {
std::printf("Benchmark complete.\n");
}

std::printf("Benchmark complete.\n");
co_return 0;
co_return result;
}

// Client
Expand All @@ -235,54 +239,60 @@ coro::task<int> async_main(int argc, char* argv[]) {
cm_ch, reinterpret_cast<sockaddr*>(&addr), sizeof(addr), ep_cfg);
ep.start_cq_pump(*sched);

// Receive server's remote buffer info
xchg_msg srv_info{};
auto recv_mr = ep.register_buffer(&srv_info, sizeof(srv_info),
IBV_ACCESS_LOCAL_WRITE);
auto recv_wc = co_await ep.conn().recv(recv_mr.view());
if (!recv_wc.ok()) {
std::fprintf(stderr, "client recv failed status=%d\n",
static_cast<int>(recv_wc.status));
co_return 1;
}

rdma::remote_buffer server_buf{srv_info.addr, srv_info.length, srv_info.rkey};
std::printf("Connected. Server buffer: %zu bytes\n", cfg.size);
std::printf("Running RDMA WRITE bandwidth test (%d iters, %zu B):\n\n",
int result = 0;
{
// Receive server's remote buffer info. All MRs in this scope are
// released before endpoint shutdown below.
xchg_msg srv_info{};
auto recv_mr = ep.register_buffer(&srv_info, sizeof(srv_info),
IBV_ACCESS_LOCAL_WRITE);
auto recv_wc = co_await ep.conn().recv(recv_mr.view());
if (!recv_wc.ok()) {
std::fprintf(stderr, "client recv failed status=%d\n",
static_cast<int>(recv_wc.status));
result = 1;
} else {
rdma::remote_buffer server_buf{
srv_info.addr, srv_info.length, srv_info.rkey};
std::printf("Connected. Server buffer: %zu bytes\n", cfg.size);
std::printf(
"Running RDMA WRITE bandwidth test (%d iters, %zu B):\n\n",
cfg.iters, cfg.size);

// GPU path
if (cfg.buf_mode == "gpu" || cfg.buf_mode == "both") {
elio::rdma_cuda::gpu_memory_region gpu_mr{
ep.pd(), cfg.size, IBV_ACCESS_LOCAL_WRITE};
auto r = co_await run_write_bw(ep, gpu_mr.view(), server_buf,
DEFAULT_WARMUP, cfg.iters);
print_result("GPU", r);
}
if (cfg.buf_mode == "gpu" || cfg.buf_mode == "both") {
elio::rdma_cuda::gpu_memory_region gpu_mr{
ep.pd(), cfg.size, IBV_ACCESS_LOCAL_WRITE};
auto r = co_await run_write_bw(
ep, gpu_mr.view(), server_buf, DEFAULT_WARMUP, cfg.iters);
print_result("GPU", r);
}

// CPU path
if (cfg.buf_mode == "cpu" || cfg.buf_mode == "both") {
std::vector<char> cpu_buf(cfg.size, 'A');
auto cpu_mr = ep.register_buffer(cpu_buf.data(), cfg.size,
IBV_ACCESS_LOCAL_WRITE);
auto r = co_await run_write_bw(ep, cpu_mr.view(), server_buf,
DEFAULT_WARMUP, cfg.iters);
print_result("CPU", r);
}
if (cfg.buf_mode == "cpu" || cfg.buf_mode == "both") {
std::vector<char> cpu_buf(cfg.size, 'A');
auto cpu_mr = ep.register_buffer(
cpu_buf.data(), cfg.size, IBV_ACCESS_LOCAL_WRITE);
auto r = co_await run_write_bw(
ep, cpu_mr.view(), server_buf, DEFAULT_WARMUP, cfg.iters);
print_result("CPU", r);
}

// Signal server we're done
xchg_msg done_msg{};
auto done_mr = ep.register_buffer(&done_msg, sizeof(done_msg),
IBV_ACCESS_LOCAL_WRITE);
auto done_wc = co_await ep.conn().send(done_mr.view());
if (!done_wc.ok()) {
std::fprintf(stderr, "client done send failed status=%d\n",
static_cast<int>(done_wc.status));
co_return 1;
xchg_msg done_msg{};
auto done_mr = ep.register_buffer(
&done_msg, sizeof(done_msg), IBV_ACCESS_LOCAL_WRITE);
auto done_wc = co_await ep.conn().send(done_mr.view());
if (!done_wc.ok()) {
std::fprintf(stderr, "client done send failed status=%d\n",
static_cast<int>(done_wc.status));
result = 1;
}
}
}

std::printf("\nDone.\n");
co_return 0;
co_await ep.shutdown();
if (result == 0) {
std::printf("\nDone.\n");
}
co_return result;
}

ELIO_ASYNC_MAIN(async_main)
Expand Down
Loading