From 0edcb1bdfaf9578491e02c0ba8d9ca789cd795bb Mon Sep 17 00:00:00 2001 From: Ketor Date: Fri, 7 Aug 2026 06:31:10 +0800 Subject: [PATCH] =?UTF-8?q?fix(rdma):=20v2.6.1=20hardening=20=E2=80=94=20r?= =?UTF-8?q?evert=20unilateral,=20fix=20MR=20leak,=20drop=20per-request=20l?= =?UTF-8?q?ogs?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Release-blocking fixes from review: 1. Revert unilateral PUT (PR#263 RoundTrip path) to dual-sided WRITE_WITH_IMM. Unilateral never entered the standard client path: KVClient::Put/BatchPut both route through CacheMany (pipelined), so the advertised performance path did not match the call chain. Removes PostSendNotify + nbuf_/nmr_ entirely, which also eliminates the per-QP notification MR leak (Close() cleared nmr_/nbuf_ vectors BEFORE the dereg loops, so both cleanup loops were always empty). 2. Fix Close() ordering: dereg/free loops run before vector clear (smr_/rmr_/dmr_ kept; nbuf_/nmr_ gone with the unilateral revert). 3. Remove per-request INFO logging in server decode (every RDMA completion logged with global-mutex fprintf), the duplicate unilateral decoders (three copies existed), and the type-punned reinterpret_cast wire load. 4. Remove dead reaper: StartReaper/StopReaper had no callers, reaper_enabled_ unused, Close() never called StopReaper(). 5. Default idle back to 10 min (DFKV_RDMA_IDLE_MS=600000); K8s launcher exports 30000 explicitly. Verified on 0064 (B200, mlx5_0, 32GB recv segment, 32GB cap, 64KB slab granularity, credits=256): t32 PUT+GET 10k @ batch 1/2/8/32, depth 1/4/8, 1MB size — all 0 fails, 0 evictions. Release claimed t32/b1 ~83% fails: root cause was double-PostRecv RQ misalignment (fixed in #266) plus test-env slab eviction, not shared-QP WR interleaving (Acquire is mutex-guarded). --- src/cache/rdma_server.cc | 72 ++------------------------------- src/transport/rdma_transport.cc | 8 ++-- src/transport/rdma_transport.h | 1 - src/transport/rdma_verbs.cc | 56 +------------------------ src/transport/rdma_verbs.h | 12 ------ 5 files changed, 7 insertions(+), 142 deletions(-) diff --git a/src/cache/rdma_server.cc b/src/cache/rdma_server.cc index 83f5851..a88f942 100644 --- a/src/cache/rdma_server.cc +++ b/src/cache/rdma_server.cc @@ -302,7 +302,7 @@ int ServerIdleMs() { // the client re-dials a stale pooled connection via RdmaTransport's 2-attempt // retry. Default 10 min keeps active/recently-used pooled conns alive; set // DFKV_RDMA_IDLE_MS=0 disables the reaper and waits indefinitely. - int out = 30000; // 30 seconds (was 10 min; Mooncake has no idle timeout but uses SIEVE eviction) + int out = 600000; // 10 min; K8s launcher exports DFKV_RDMA_IDLE_MS=30000 explicitly const char* e = std::getenv("DFKV_RDMA_IDLE_MS"); if (e && *e) { long v = std::strtol(e, nullptr, 10); @@ -543,7 +543,6 @@ void RdmaServer::Serve(int boot_fd) { request->recv_slot = recv_slot; request->data_slot = recv_slot; request->recv_bytes = completion.byte_len; - DFKV_LOG_INFO("rdma: recv slot=" + std::to_string(recv_slot) + " byte_len=" + std::to_string(completion.byte_len) + " opcode=" + std::to_string(completion.opcode) + " has_imm=" + std::to_string(has_immediate)); if (write_imm_recv) { const size_t data_slot = static_cast(ntohl(completion.imm_data)); if (data_slot >= K || completion.byte_len < kReqPrefix) return false; @@ -578,73 +577,8 @@ void RdmaServer::Serve(int boot_fd) { } return request->get.targets.size() <= rdma::kV2MaxGetTargets; } - // Unilateral PUT path: the notification is a 50B request prefix sent via - // SEND. The offset field carries the data_slot index where [header|payload] - // was written via plain RDMA WRITE to the shared receive segment. - if (request->fields.op == static_cast(WireOp::kCache)) { - if (completion.byte_len < kReqPrefix) return false; - const size_t data_slot = static_cast(request->fields.offset); - if (data_slot >= K) return false; - const char* data_frame = recv_lease.data() + data_slot * slot_size + - rdma::kV2PutPrefixOffset; - DFKV_LOG_INFO("rdma: unilateral PUT data_slot=" + std::to_string(data_slot) + - " slot_size=" + std::to_string(slot_size) + - " data_frame_ver=" + std::to_string(static_cast(static_cast(data_frame[0]))) + - " data_frame_op=" + std::to_string(static_cast(static_cast(data_frame[1]))) + - " payload_len=" + std::to_string(*reinterpret_cast(data_frame + 42))); - if (!DecodeReqVersion( - data_frame, kNativeProtoRdmaV2, &request->fields, - static_cast(logical_data_cap)) || - request->fields.op != static_cast(WireOp::kCache) || - request->fields.payload_len > static_cast(logical_data_cap)) { - return false; - } - request->data_slot = data_slot; - request->contiguous_payload = data_frame + kReqPrefix; - v2_put_writes_.fetch_add(1, std::memory_order_relaxed); - return true; - } - // Unilateral PUT path: the notification is a 50B request prefix sent via - // SEND on the control QP. The offset field carries the data_slot index - // where [header|payload] was written via plain RDMA WRITE to the shared - // receive segment. nbuf_ is separate from sbuf_ so no race. - if (request->fields.op == static_cast(WireOp::kCache)) { - if (completion.byte_len < kReqPrefix) return false; - const size_t data_slot = static_cast(request->fields.offset); - if (data_slot >= K) return false; - const char* data_frame = recv_lease.data() + data_slot * slot_size + - rdma::kV2PutPrefixOffset; - if (!DecodeReqVersion( - data_frame, kNativeProtoRdmaV2, &request->fields, - static_cast(logical_data_cap)) || - request->fields.op != static_cast(WireOp::kCache) || - request->fields.payload_len > static_cast(logical_data_cap)) { - return false; - } - request->data_slot = data_slot; - request->contiguous_payload = data_frame + kReqPrefix; - v2_put_writes_.fetch_add(1, std::memory_order_relaxed); - return true; - } - // Unilateral PUT: notification SEND carries offset=data_slot. - // Data was written via unsignaled WRITE to recv segment slot. - if (request->fields.op == static_cast(WireOp::kCache)) { - const size_t data_slot = static_cast(request->fields.offset); - if (data_slot >= K) return false; - const char* data_frame = recv_lease.data() + data_slot * slot_size + - rdma::kV2PutPrefixOffset; - if (!DecodeReqVersion( - data_frame, kNativeProtoRdmaV2, &request->fields, - static_cast(logical_data_cap)) || - request->fields.op != static_cast(WireOp::kCache) || - request->fields.payload_len > static_cast(logical_data_cap)) { - return false; - } - request->data_slot = data_slot; - request->contiguous_payload = data_frame + kReqPrefix; - v2_put_writes_.fetch_add(1, std::memory_order_relaxed); - return true; - } + + // Other control ops (Exist/Remove/Members/Lookup) with inline payload if (completion.byte_len < kReqPrefix + request->fields.payload_len) { diff --git a/src/transport/rdma_transport.cc b/src/transport/rdma_transport.cc index a60ab68..73ad41c 100644 --- a/src/transport/rdma_transport.cc +++ b/src/transport/rdma_transport.cc @@ -907,12 +907,10 @@ Status RdmaTransport::RoundTrip(const std::string& node, WireOp op, return Status::kIOError; } conn->Encode(ep.sbuf(0), op, key, offset, length, payload_len); - std::memcpy(ep.nbuf(0), ep.sbuf(0), kReqPrefix); - ok = ep.PostWriteScatter( + ok = ep.PostRecv(0) && + ep.PostWriteImmScatter( 0, kReqPrefix, payload, static_cast(payload_len), - payload_mr, conn->put_addr(0), conn->recv_segment.rkey); - if (ok) - ok = ep.PostRecv(0) && ep.PostSendNotify(0, kReqPrefix); + payload_mr, conn->put_addr(0), conn->recv_segment.rkey, 0); if (ok) v2_put_writes_.fetch_add(1, std::memory_order_relaxed); } else if (op == WireOp::kRange) { diff --git a/src/transport/rdma_transport.h b/src/transport/rdma_transport.h index 3fdc5d5..8837b93 100644 --- a/src/transport/rdma_transport.h +++ b/src/transport/rdma_transport.h @@ -209,7 +209,6 @@ class RdmaTransport : public Transport { std::atomic numa_caller_unknown_fallbacks_{0}; std::atomic numa_no_local_fallbacks_{0}; // #1: async CQ reaper for high-concurrency batch operations - bool reaper_enabled_ = false; }; } // namespace dfkv diff --git a/src/transport/rdma_verbs.cc b/src/transport/rdma_verbs.cc index 6667664..d113906 100644 --- a/src/transport/rdma_verbs.cc +++ b/src/transport/rdma_verbs.cc @@ -273,10 +273,8 @@ void RcEndpoint::Close() { // only after the last endpoint that can still return it has gone idle/closed. for (const auto& pool : pool_mr_) SharedReleasePoolMr(ctx_, pool.mr); pool_mr_.clear(); - smr_.clear(); rmr_.clear(); dmr_.clear(); nmr_.clear(); nbuf_.clear(); + smr_.clear(); rmr_.clear(); dmr_.clear(); for (auto* b : sbuf_) delete[] b; - for (auto* m : nmr_) if (m) ibv_dereg_mr(m); - for (auto* b : nbuf_) delete[] b; for (auto* b : rbuf_) delete[] b; for (auto* b : dbuf_) std::free(b); sbuf_.clear(); rbuf_.clear(); dbuf_.clear(); @@ -414,7 +412,6 @@ bool RcEndpoint::Open(const char* dev_name, size_t cap, size_t depth, // sizes (>=128 KiB) new[] is mmap-backed = page-aligned, which mbind needs. numa_node_ = numa::DeviceNode(dev_name); sbuf_.resize(depth_, nullptr); rbuf_.resize(depth_, nullptr); - nbuf_.resize(depth_, nullptr); nmr_.resize(depth_, nullptr); smr_.resize(depth_, nullptr); rmr_.resize(depth_, nullptr); if (direct_io_buffers) { const size_t dio_cap = direct_io_cap ? direct_io_cap : cap_; @@ -439,9 +436,6 @@ bool RcEndpoint::Open(const char* dev_name, size_t cap, size_t depth, } } for (size_t i = 0; i < depth_; ++i) { - nbuf_[i] = new char[64]; - nmr_[i] = ibv_reg_mr(pd_, nbuf_[i], 64, IBV_ACCESS_LOCAL_WRITE); - if (!nmr_[i]) { Close(); return false; } } // local addressing info @@ -897,21 +891,6 @@ bool RcEndpoint::PostWriteScatter( return ibv_post_send(qp_, &wr, &bad) == 0; } -bool RcEndpoint::PostSendNotify(size_t slot, size_t len) { - if (slot >= depth_ || len > 64 || !nbuf_[slot] || !nmr_[slot]) return false; - ibv_sge sge{}; - sge.addr = reinterpret_cast(nbuf_[slot]); - sge.length = static_cast(len); - sge.lkey = nmr_[slot]->lkey; - ibv_send_wr wr{}, *bad = nullptr; - wr.wr_id = slot; - wr.sg_list = &sge; - wr.num_sge = 1; - wr.opcode = IBV_WR_SEND; - wr.send_flags = IBV_SEND_SIGNALED; - return ibv_post_send(qp_, &wr, &bad) == 0; -} - bool RcEndpoint::PostWriteImmScatterMulti( size_t slot, size_t header_len, const std::vector>& segs, @@ -1005,38 +984,5 @@ void RcEndpoint::Wake() { } -void RcEndpoint::StartReaper(std::vector*>* slots) { - reap_slots_ = *slots; - reap_stop_.store(false, std::memory_order_relaxed); - reap_thread_ = new std::thread([this] { - std::vector wcs(64); - while (!reap_stop_.load(std::memory_order_relaxed)) { - int got = ibv_poll_cq(cq_, static_cast(wcs.size()), wcs.data()); - if (got > 0) { - for (int i = 0; i < got; ++i) { - size_t slot = static_cast(wcs[i].wr_id); - if (slot < reap_slots_.size() && reap_slots_[slot]) { - // Store byte_len for RECV, mark done for SEND - uint32_t val = (wcs[i].opcode == IBV_WC_RECV) ? wcs[i].byte_len : 0xFFFFFFFF; - reap_slots_[slot]->store(val, std::memory_order_release); - } - } - } else if (got == 0) { - // No completion: brief yield to avoid 100% CPU - std::this_thread::yield(); - } - } - }); -} - -void RcEndpoint::StopReaper() { - if (!reap_thread_) return; - reap_stop_.store(true, std::memory_order_relaxed); - reap_thread_->join(); - delete reap_thread_; - reap_thread_ = nullptr; - reap_slots_.clear(); -} - } // namespace rdma } // namespace dfkv diff --git a/src/transport/rdma_verbs.h b/src/transport/rdma_verbs.h index 54700cd..3545e31 100644 --- a/src/transport/rdma_verbs.h +++ b/src/transport/rdma_verbs.h @@ -129,7 +129,6 @@ class RcEndpoint { int numa_node() const { return numa_node_; } // device's NUMA node, or -1 char* sbuf(size_t slot) { return sbuf_[slot]; } char* rbuf(size_t slot) { return rbuf_[slot]; } - char* nbuf(size_t slot) { return nbuf_[slot]; } char* dbuf(size_t slot) { return dbuf_.empty() ? nullptr : dbuf_[slot]; } size_t dbuf_cap() const { return dbuf_cap_; } ibv_mr* dmr(size_t slot) { return dmr_.empty() ? nullptr : dmr_[slot]; } @@ -248,8 +247,6 @@ class RcEndpoint { size_t slot, size_t header_len, const void* payload, size_t payload_len, ibv_mr* payload_mr, uint64_t remote_addr, uint32_t remote_rkey); - // Signaled SEND from nbuf_[slot] (notification buffer). - bool PostSendNotify(size_t slot, size_t len); // QP scatter-gather capability negotiated in Open() = min(kMaxSge, device cap). // The SG datapath caps raw-payload segments at max_sge()-1 (SGE0 is the wire prefix). @@ -271,8 +268,6 @@ class RcEndpoint { // #1: Start a background CQ reaper that polls ibv_poll_cq in a tight // loop and fills per-slot atomic flags. Eliminates CQ channel // syscall and CQ contention at high thread counts. - void StartReaper(std::vector*>* slots); - void StopReaper(); private: void Close(); @@ -290,19 +285,12 @@ class RcEndpoint { unsigned cq_armed_unacked_ = 0; bool busy_poll_ = false; size_t num_qp_ = 1; - // #1: async CQ reaper support. Per-slot atomic completion flags - // filled by a background reaper thread, checked by ReapPosted. - std::vector*> reap_slots_; - std::atomic reap_stop_{false}; - std::thread* reap_thread_ = nullptr; size_t cap_ = 0, depth_ = 0, dbuf_cap_ = 0; uint16_t remote_depth_ = 0; // Set from the required v2 QP advertisement. size_t max_sge_ = 2; // QP max_send_sge/max_recv_sge = min(kMaxSge, device cap) std::vector sbuf_, rbuf_, dbuf_; std::vector smr_, rmr_, dmr_; - std::vector nbuf_; - std::vector nmr_; // Big pre-registered caller regions (the host KV pool). Each cache entry owns // one shared-registry reference. Successful growth replaces the idle // endpoint's older same-base generation; endpoints with an in-flight user