This is an automated email from the ASF dual-hosted git repository.

chenBright pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/brpc.git


The following commit(s) were added to refs/heads/master by this push:
     new 0c8aede1 Fix rdma handshake failing the socket instead of falling back 
to TCP (#3424)
0c8aede1 is described below

commit 0c8aede17187363a0b03584388fb26d57f438db0
Author: Bright Chen <[email protected]>
AuthorDate: Fri Aug 7 17:38:10 2026 +0800

    Fix rdma handshake failing the socket instead of falling back to TCP (#3424)
---
 src/brpc/rdma/rdma_endpoint.cpp |  75 +++++++++++++++++++++-------
 src/brpc/rdma/rdma_endpoint.h   |  10 +++-
 test/brpc_rdma_unittest.cpp     | 107 ++++++++++++++++++++++++++++++++++++++++
 3 files changed, 173 insertions(+), 19 deletions(-)

diff --git a/src/brpc/rdma/rdma_endpoint.cpp b/src/brpc/rdma/rdma_endpoint.cpp
index 8660d8d9..e2ce2e0c 100644
--- a/src/brpc/rdma/rdma_endpoint.cpp
+++ b/src/brpc/rdma/rdma_endpoint.cpp
@@ -53,6 +53,10 @@ extern int (*IbvQueryEce)(ibv_qp*, ibv_ece*);
 extern int (*IbvSetEce)(ibv_qp*, ibv_ece*);
 extern bool g_skip_rdma_init;
 
+// Only for UT: force AllocateResources() to fail, so that the "fallback to 
TCP" path
+// of the handshake can be tested without a real RDMA device.
+bool g_fail_resource_alloc_for_test = false;
+
 DEFINE_int32(rdma_sq_size, 128, "SQ size for RDMA");
 DEFINE_int32(rdma_rq_size, 128, "RQ size for RDMA");
 DEFINE_bool(rdma_recv_zerocopy, true, "Enable zerocopy for receive side");
@@ -430,7 +434,9 @@ void* RdmaEndpoint::ProcessHandshakeAtClient(void* arg) {
     // First initialize CQ and QP resources.
     ep->_state.store(C_ALLOC_QPCQ, butil::memory_order_relaxed);
     if (ep->AllocateResources() < 0) {
-        LOG(WARNING) << "Fallback to tcp:" << s->description();
+        PLOG(WARNING) << "Fail to allocate rdma resources, fallback to tcp:"
+                      << s->description();
+        errno = 0;
         rdma_transport->_rdma_state = RdmaTransport::RDMA_OFF;
         ep->_state.store(FALLBACK_TCP, butil::memory_order_release);
         return NULL;
@@ -563,8 +569,8 @@ ParseResult 
RdmaEndpoint::ExecuteServerHandshake(butil::IOBuf* source, Socket* s
             ep->ApplyRemoteHello(remote);
             ep->_state.store(S_ALLOC_QPCQ, butil::memory_order_relaxed);
             if (ep->AllocateResources() < 0) {
-                LOG(WARNING) << "Fail to allocate rdma resources, fallback to 
tcp:"
-                             << s->description();
+                PLOG(WARNING) << "Fail to allocate rdma resources, fallback to 
tcp:"
+                              << s->description();
                 negotiated = false;
             } else {
                 ep->_state.store(S_BRINGUP_QP, butil::memory_order_relaxed);
@@ -1072,8 +1078,26 @@ static RdmaResource* AllocateQpCq(uint16_t sq_size, 
uint16_t rq_size) {
 }
 
 int RdmaEndpoint::AllocateResources() {
+    if (DoAllocateResources() == 0) {
+        return 0;
+    }
+
+    const int saved_errno = errno;
+    DeallocateResources();
+    _sbuf.clear();
+    _rbuf.clear();
+    _rbuf_data.clear();
+    errno = saved_errno;
+    return -1;
+}
+
+int RdmaEndpoint::DoAllocateResources() {
     if (BAIDU_UNLIKELY(g_skip_rdma_init)) {
         // For UT
+        if (BAIDU_UNLIKELY(g_fail_resource_alloc_for_test)) {
+            errno = EINVAL;
+            return -1;
+        }
         return 0;
     }
 
@@ -1097,10 +1121,10 @@ int RdmaEndpoint::AllocateResources() {
     }
 
     if (!FLAGS_rdma_use_polling) {
-        if (0 != ReqNotifyCq(true)) {
+        if (0 != ReqNotifyCq(true, false)) {
             return -1;
         }
-        if (0 != ReqNotifyCq(false)) {
+        if (0 != ReqNotifyCq(false, false)) {
             return -1;
         }
 
@@ -1364,10 +1388,20 @@ _reclaim:
             goto _reclaim;
         }
 
-        BAIDU_SCOPED_LOCK(*g_rdma_resource_mutex);
-        _resource->next = g_rdma_resource_list;
-        g_rdma_resource_list = _resource;
+        {
+            BAIDU_SCOPED_LOCK(*g_rdma_resource_mutex);
+            _resource->next = g_rdma_resource_list;
+            g_rdma_resource_list = _resource;
+        }
+        _resource = NULL;
     }
+
+    // Detach everything from this endpoint so that the function is
+    // idempotent: it is called both when the endpoint is reset/destroyed
+    // and when AllocateResources() fails halfway.
+    _cq_sid = INVALID_SOCKET_ID;
+    _send_cq_events = 0;
+    _recv_cq_events = 0;
 }
 
 static const int MAX_CQ_EVENTS = 128;
@@ -1411,17 +1445,21 @@ int RdmaEndpoint::GetAndAckEvents(SocketUniquePtr& s) {
     return 0;
 }
 
-int RdmaEndpoint::ReqNotifyCq(bool send_cq) {
-    errno = ibv_req_notify_cq(
+int RdmaEndpoint::ReqNotifyCq(bool send_cq, bool fatal_on_error) {
+    const int err = ibv_req_notify_cq(
         send_cq ? _resource->send_cq : _resource->recv_cq,
         send_cq ? 0 : 1);
-    if (0 != errno) {
-        const int saved_errno = errno;
+    if (0 != err) {
+        errno = err;
         PLOG(WARNING) << "Fail to arm " << (send_cq ? "send" : "recv")
                       << " CQ comp channel from " << _socket->description();
-        _socket->SetFailed(saved_errno, "Fail to arm %s CQ channel from %s: 
%s",
-                           send_cq ? "send" : "recv", 
_socket->description().c_str(),
-                           berror(saved_errno));
+        if (fatal_on_error) {
+            _socket->SetFailed(err, "Fail to arm %s CQ channel from %s: %s",
+                               send_cq ? "send" : "recv", 
_socket->description().c_str(),
+                               berror(err));
+        }
+        // The logging and SetFailed() above may clobber errno.
+        errno = err;
         return -1;
     }
 
@@ -1485,10 +1523,13 @@ void RdmaEndpoint::PollCq(Socket* m) {
                 // that the event arrives after the poll but before the notify,
                 // we should re-poll the CQ once after the notify to check if
                 // there is an available CQE.
-                if (0 != ep->ReqNotifyCq(true)) {
+                // The connection is already working in RDMA mode here, a
+                // failed re-arm means no more CQ event will be reported,
+                // which is fatal for this connection.
+                if (0 != ep->ReqNotifyCq(true, true)) {
                     return;
                 }
-                if (0 != ep->ReqNotifyCq(false)) {
+                if (0 != ep->ReqNotifyCq(false, true)) {
                     return;
                 }
                 notified = true;
diff --git a/src/brpc/rdma/rdma_endpoint.h b/src/brpc/rdma/rdma_endpoint.h
index 36e22ad2..03bec814 100644
--- a/src/brpc/rdma/rdma_endpoint.h
+++ b/src/brpc/rdma/rdma_endpoint.h
@@ -164,10 +164,16 @@ private:
     // Process handshake at the client
     static void* ProcessHandshakeAtClient(void* arg);
 
-    // Allocate resources
+    // Allocate resources. On failure the endpoint is left with no RDMA
+    // resource attached, so that the handshake can safely fall back to TCP.
     // Return 0 if success, -1 if failed and errno set
     int AllocateResources();
 
+    // The real implementation of AllocateResources(), which may return
+    // in the middle with resources partially allocated.
+    // Return 0 if success, -1 if failed and errno set
+    int DoAllocateResources();
+
     // Release resources
     void DeallocateResources();
 
@@ -244,7 +250,7 @@ private:
     int GetAndAckEvents(SocketUniquePtr& s);
 
     // Request completion notification on a send/recv CQ.
-    int ReqNotifyCq(bool send_cq);
+    int ReqNotifyCq(bool send_cq, bool fatal_on_error);
 
     // Poll CQ and get the work completion
     static void PollCq(Socket* m);
diff --git a/test/brpc_rdma_unittest.cpp b/test/brpc_rdma_unittest.cpp
index e30ae09f..2ecd1f3c 100644
--- a/test/brpc_rdma_unittest.cpp
+++ b/test/brpc_rdma_unittest.cpp
@@ -72,6 +72,7 @@ extern int (*IbvQueryQp)(ibv_qp*, ibv_qp_attr*, 
ibv_qp_attr_mask, ibv_qp_init_at
 extern int (*IbvDestroyQp)(ibv_qp*);
 extern butil::atomic<bool> g_rdma_available;
 extern bool g_skip_rdma_init;
+extern bool g_fail_resource_alloc_for_test;
 } // namespace rdma
 } // namespace brpc
 
@@ -1919,6 +1920,112 @@ TEST_F(RdmaTest, 
v3_server_reply_has_no_ece_without_hw_negotiation) {
     StopServer();
 }
 
+class ResourceAllocFailGuard {
+public:
+    explicit ResourceAllocFailGuard(bool v)
+        : _saved(rdma::g_fail_resource_alloc_for_test) {
+        rdma::g_fail_resource_alloc_for_test = v;
+    }
+    ~ResourceAllocFailGuard() {
+        rdma::g_fail_resource_alloc_for_test = _saved;
+    }
+private:
+    bool _saved;
+};
+
+TEST_F(RdmaTest, client_alloc_resource_fail_fallback_tcp) {
+    StartServer();
+    ResourceAllocFailGuard alloc_fail_guard(true);
+
+    Channel channel;
+    ChannelOptions chan_options;
+    chan_options.socket_mode = SOCKET_MODE_RDMA;
+    chan_options.connect_timeout_ms = 500;
+    chan_options.timeout_ms = 500;
+    chan_options.max_retry = 0;
+    ASSERT_EQ(0, channel.Init(g_ep, &chan_options));
+
+    Controller cntl;
+    test::EchoRequest req;
+    test::EchoResponse res;
+    req.set_message(__FUNCTION__);
+    req.set_sleep_us(200000);
+    google::protobuf::Closure* done = DoNothing();
+    ::test::EchoService::Stub(&channel).Echo(&cntl, &req, &res, done);
+    usleep(100000);
+
+    SocketUniquePtr s;
+    ASSERT_EQ(0, Socket::Address(cntl._single_server_id, &s));
+    ASSERT_EQ(rdma::RdmaEndpoint::FALLBACK_TCP,
+              
static_cast<RdmaTransport*>(s->_transport.get())->_rdma_ep->_state);
+    ASSERT_EQ(RdmaTransport::RDMA_OFF,
+              static_cast<RdmaTransport*>(s->_transport.get())->_rdma_state);
+    // The socket must not be failed, otherwise it can no longer carry TCP.
+    ASSERT_FALSE(s->Failed());
+
+    // The RPC still completes over TCP.
+    bthread_id_join(cntl.call_id());
+    ASSERT_EQ(0, cntl.ErrorCode()) << cntl.ErrorText();
+
+    StopServer();
+}
+
+TEST_F(RdmaTest, server_alloc_resource_fail_fallback_tcp) {
+    StartServer();
+    ResourceAllocFailGuard alloc_fail_guard(true);
+
+    sockaddr_in addr;
+    bzero((char*)&addr, sizeof(addr));
+    addr.sin_family = AF_INET;
+    addr.sin_port = htons(PORT);
+    butil::fd_guard sockfd(socket(AF_INET, SOCK_STREAM, 0));
+    ASSERT_TRUE(sockfd >= 0);
+    ASSERT_EQ(0, connect(sockfd, (sockaddr*)&addr, sizeof(sockaddr)));
+    usleep(100000);  // wait for server to handle the msg
+    Socket* s = GetSocketFromServer(0);
+    ASSERT_TRUE(s != NULL);
+    ASSERT_EQ(rdma::RdmaEndpoint::UNINIT,
+              
static_cast<RdmaTransport*>(s->_transport.get())->_rdma_ep->_state);
+
+    // Send a well-formed v2 hello: the negotiation succeeds
+    // but the resource allocation does not.
+    rdma::v2_wire::HelloMessage msg{};
+    msg.msg_len = rdma::HELLO_V2_MSG_LEN_MIN;
+    msg.hello_ver = rdma::HELLO_V2_VERSION;
+    msg.impl_ver = rdma::IMPL_V2_VERSION;
+    msg.sq_size = 16;
+    msg.rq_size = 16;
+    msg.block_size = 8192;
+    msg.qp_num = 0;
+    msg.gid = rdma::GetRdmaGid();
+
+    uint8_t data[rdma::HELLO_V2_MSG_LEN_MIN];
+    memcpy(data, "RDMA", 4);
+    msg.Serialize(data + 4);
+    ASSERT_EQ(rdma::HELLO_V2_MSG_LEN_MIN,
+              write(sockfd, data, rdma::HELLO_V2_MSG_LEN_MIN));
+    usleep(100000);
+    ASSERT_EQ(rdma::RdmaEndpoint::S_ACK_WAIT,
+              
static_cast<RdmaTransport*>(s->_transport.get())->_rdma_ep->_state);
+    ASSERT_EQ(RdmaTransport::RDMA_OFF,
+              static_cast<RdmaTransport*>(s->_transport.get())->_rdma_state);
+    ASSERT_FALSE(s->Failed());
+
+    // Ack without RDMA so that the server finishes the handshake in TCP mode.
+    uint32_t flags = butil::HostToNet32(0);
+    ASSERT_EQ(sizeof(flags), write(sockfd, &flags, sizeof(flags)));
+    usleep(100000);
+    ASSERT_EQ(rdma::RdmaEndpoint::FALLBACK_TCP,
+              
static_cast<RdmaTransport*>(s->_transport.get())->_rdma_ep->_state);
+    ASSERT_FALSE(s->Failed());
+
+    sockfd.reset(-1);
+    usleep(100000);
+    ASSERT_EQ(NULL, GetSocketFromServer(0));
+
+    StopServer();
+}
+
 TEST_F(RdmaTest, try_global_disable_rdma) {
     StartServer();
     rdma::g_rdma_available.store(false, butil::memory_order_relaxed);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to