chenBright commented on code in PR #3545:
URL: https://github.com/apache/brpc/pull/3545#discussion_r4089275249


##########
test/bthread_fd_unittest.cpp:
##########
@@ -566,73 +621,74 @@ TEST(FDTest, double_close) {
     ASSERT_EQ(ec, errno);
 }
 
-const char* g_hostname1 = "github.com";
-const char* g_hostname2 = "baidu.com";
-TEST(FDTest, bthread_connect) {
-    butil::EndPoint ep1;
-    butil::EndPoint ep2;
-    ASSERT_EQ(0, butil::hostname2endpoint(g_hostname1, 80, &ep1));
-    ASSERT_EQ(0, butil::hostname2endpoint(g_hostname2, 80, &ep2));
-
-    {
-        struct sockaddr_storage serv_addr{};
-        socklen_t serv_addr_size = 0;
-        ASSERT_EQ(0, endpoint2sockaddr(ep1, &serv_addr, &serv_addr_size));
-        butil::fd_guard sockfd(socket(serv_addr.ss_family, SOCK_STREAM, 0));
-        ASSERT_LE(0, sockfd);
-        bool is_blocking = butil::is_blocking(sockfd);
-        ASSERT_LE(0, sockfd);
-        ASSERT_EQ(0, bthread_connect(sockfd, (struct sockaddr*) &serv_addr, 
serv_addr_size));
-        ASSERT_EQ(is_blocking, butil::is_blocking(sockfd));
+// Local listeners keep connect tests independent of DNS, Internet latency,
+// and the assumption that a connection cannot complete within one millisecond.
+void TestLocalConnect(bool timed) {
+    butil::EndPoint endpoint;
+    ASSERT_EQ(0, butil::str2endpoint("127.0.0.1:0", &endpoint));
+    butil::fd_guard listener(butil::tcp_listen(endpoint));
+    ASSERT_GE(listener, 0);
+    ASSERT_EQ(0, butil::get_local_side(listener, &endpoint));
+    struct sockaddr_storage address{};
+    socklen_t length = 0;
+    ASSERT_EQ(0, endpoint2sockaddr(endpoint, &address, &length));
+    butil::fd_guard client(socket(address.ss_family, SOCK_STREAM, 0));
+    ASSERT_GE(client, 0);
+    bool was_blocking = butil::is_blocking(client);
+    timespec deadline = butil::seconds_from_now(10);
+    int rc = bthread_timed_connect(
+        client, reinterpret_cast<sockaddr*>(&address), length,
+        timed ? &deadline : nullptr);
+    ASSERT_EQ(0, rc) << "errno=" << errno;
+    ASSERT_EQ(was_blocking, butil::is_blocking(client));
+    ASSERT_EQ(0, butil::is_connected(client));
+    // The handshake does not require a concurrent accept thread.
+    butil::fd_guard accepted(accept(listener, nullptr, nullptr));
+    ASSERT_GE(accepted, 0);
+}
 
-    }
+TEST(FDTest, bthread_connect) {
+    TestLocalConnect(false);
+    TestLocalConnect(true);
+}
 
-    {
-        struct sockaddr_storage serv_addr{};
-        socklen_t serv_addr_size = 0;
-        ASSERT_EQ(0, endpoint2sockaddr(ep2, &serv_addr, &serv_addr_size));
-        butil::fd_guard sockfd(socket(serv_addr.ss_family, SOCK_STREAM, 0));
-        ASSERT_LE(0, sockfd);
-        bool is_blocking = butil::is_blocking(sockfd);
-        // In most cases, 1 millisecond will result in a connection timeout.
-        timespec abstime = butil::milliseconds_from_now(1);
-        const int rc = bthread_timed_connect(
-            sockfd, (struct sockaddr*) &serv_addr,
-            serv_addr_size, &abstime);
-        ASSERT_EQ(-1, rc);
-        ASSERT_EQ(ETIMEDOUT, errno);
-        ASSERT_EQ(is_blocking, butil::is_blocking(sockfd));
-    }
+#if defined(OS_LINUX)
+TEST(FDTest, connect_timeout_with_full_accept_queue) {
+    butil::EndPoint endpoint;
+    ASSERT_EQ(0, butil::str2endpoint("127.0.0.1:0", &endpoint));
+    butil::fd_guard listener(butil::tcp_listen(endpoint));
+    ASSERT_GE(listener, 0);
+    // Linux permits one queued connection for backlog=0. Leave it unaccepted
+    // so the next handshake cannot complete; no external network is needed.
+    ASSERT_EQ(0, listen(listener, 0));
+    ASSERT_EQ(0, butil::get_local_side(listener, &endpoint));
+    sockaddr_storage address{};
+    socklen_t length = 0;
+    ASSERT_EQ(0, endpoint2sockaddr(endpoint, &address, &length));
+    butil::fd_guard queued(socket(AF_INET, SOCK_STREAM, 0));
+    ASSERT_GE(queued, 0);
+    timespec setup_deadline = butil::seconds_from_now(5);
+    ASSERT_EQ(0, bthread_timed_connect(queued,
+        reinterpret_cast<sockaddr*>(&address), length, &setup_deadline));
+    butil::fd_guard client(socket(AF_INET, SOCK_STREAM, 0));
+    ASSERT_GE(client, 0);
+    timespec deadline = butil::milliseconds_from_now(50);
+    int rc = bthread_timed_connect(client,
+        reinterpret_cast<sockaddr*>(&address), length, &deadline);
+    int error = errno;
+    ASSERT_EQ(-1, rc);
+    ASSERT_EQ(ETIMEDOUT, error);

Review Comment:
   Same as above.



##########
test/bthread_futex_unittest.cpp:
##########
@@ -113,30 +116,55 @@ TEST(FutexTest, futex_wake_before_wait) {
 }
 
 void* dummy_waiter(void* lock) {
-    bthread::futex_wait_private(lock, 0, nullptr);
+    timespec timeout = butil::seconds_to_timespec(10);
+    int rc;
+    do {
+        rc = bthread::futex_wait_private(lock, 0, &timeout);
+    } while (rc != 0 && errno == EINTR);
+    EXPECT_EQ(0, rc);
     return nullptr;
 }
 
 TEST(FutexTest, futex_wake_many_waiters_perf) {
     
-    int lock1 = 0;
-    size_t N = 0;
-    pthread_t th;
-    for (; N < 1000 && !pthread_create(&th, nullptr, dummy_waiter, &lock1); 
++N) {}
-    
-    sleep(1);
+    butil::atomic<int> lock1(0);
+    std::vector<pthread_t> threads;
+    for (size_t i = 0; i < 1000; ++i) {
+        pthread_t th;
+        if (pthread_create(&th, nullptr, dummy_waiter, &lock1) != 0) {
+            break;
+        }
+        threads.push_back(th);
+    }
+    ASSERT_FALSE(threads.empty());
+    size_t N = threads.size();
     int nwakeup = 0;
+    int64_t wake_ns = 0;
+    int64_t deadline = butil::cpuwide_time_us() + 5000000L;
     butil::Timer tm;
-    tm.start();
-    for (size_t i = 0; i < N; ++i) {
-        nwakeup += bthread::futex_wake_private(&lock1, 1);
+    while (static_cast<size_t>(nwakeup) < N &&
+           butil::cpuwide_time_us() < deadline) {
+        tm.start();
+        int rc = bthread::futex_wake_private(&lock1, 1);
+        tm.stop();
+        EXPECT_GE(rc, 0);
+        if (rc > 0) {
+            nwakeup += rc;
+            wake_ns += tm.n_elapsed();
+        } else {
+            usleep(1000);
+        }
     }
-    tm.stop();
-    printf("N=%lu, futex_wake a thread = %" PRId64 "ns\n", N, tm.n_elapsed() / 
N);
-    ASSERT_EQ(N, (size_t)nwakeup);
+    // Also release late waiters on failure; a wake alone is not persistent.
+    lock1.store(1);
+    bthread::futex_wake_private(&lock1, INT_MAX);
+    for (pthread_t th : threads) {
+        EXPECT_EQ(0, pthread_join(th, nullptr));
+    }
+    ASSERT_EQ(N, static_cast<size_t>(nwakeup));
+    printf("N=%lu, futex_wake a thread = %" PRId64 "ns\n", N, wake_ns / N);

Review Comment:
   Fixed in 53373a5d52bed4a1f06c00bcc108002c31181818 . The waiter now accepts 
EWOULDBLOCK only after the controller has published cleanup and changed the 
futex value. Normal execution still requires a successful wake, so a late 
waiter no longer adds misleading failures while the registration failure 
remains reported.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to