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


##########
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) {

Review Comment:
   Asserting strictly `ETIMEDOUT` can be environment-dependent on Linux. For 
example, if `net.ipv4.tcp_abort_on_overflow=1`, overflowing the accept queue 
may produce an immediate `ECONNREFUSED` (RST) instead of a timeout. To keep 
this test portable across CI hosts, consider accepting either `ETIMEDOUT` or 
`ECONNREFUSED` (or reading the sysctl and skipping/branching accordingly).



##########
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:
   Asserting strictly `ETIMEDOUT` can be environment-dependent on Linux. For 
example, if `net.ipv4.tcp_abort_on_overflow=1`, overflowing the accept queue 
may produce an immediate `ECONNREFUSED` (RST) instead of a timeout. To keep 
this test portable across CI hosts, consider accepting either `ETIMEDOUT` or 
`ECONNREFUSED` (or reading the sysctl and skipping/branching accordingly).



##########
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));

Review Comment:
   Asserting strictly `ETIMEDOUT` can be environment-dependent on Linux. For 
example, if `net.ipv4.tcp_abort_on_overflow=1`, overflowing the accept queue 
may produce an immediate `ECONNREFUSED` (RST) instead of a timeout. To keep 
this test portable across CI hosts, consider accepting either `ETIMEDOUT` or 
`ECONNREFUSED` (or reading the sysctl and skipping/branching accordingly).



##########
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:
   The test assumes all waiters will be registered and individually woken 
within a fixed 5s deadline, then hard-asserts `nwakeup == N`. On 
slow/oversubscribed CI, thread creation and waiter registration can 
legitimately exceed this budget, making the test flaky again. A more robust 
approach is to add an explicit readiness/registration counter (or barrier) in 
`dummy_waiter` and wait until all threads have entered the wait before starting 
the wake timing loop, or to change the post-condition to assert all threads 
eventually exit (and only report timing for those actually woken).



##########
test/bthread_timer_thread_unittest.cpp:
##########
@@ -85,11 +101,11 @@ class TimeKeeper {
     {
         ASSERT_TRUE(!_run_times.empty());
         long diff = timespec_diff_us(_run_times[0], expect_run_time);
-        EXPECT_LE(labs(diff), 50000);
+        EXPECT_GE(diff, 0);
     }

Review Comment:
   This assertion now only checks that the task did not run early, but it no 
longer bounds how late it may run. That can mask real timer regressions (e.g., 
tasks firing seconds/minutes late) while still passing. Consider keeping the 
anti-early check while also adding a generous upper bound (large enough for CI 
variance) so the test still detects stalls.



##########
test/bthread_timer_thread_unittest.cpp:
##########
@@ -98,59 +114,70 @@ class TimeKeeper {
         keeper->run();
     }
 
+    bool wait_finished() {
+        return WaitUntil([this] {
+            return _finished.load(butil::memory_order_acquire);
+        });
+    }
+
+    bool wait_started() {
+        return WaitUntil([this] {
+            return _started.load(butil::memory_order_acquire);
+        });
+    }
+
     timespec _expect_run_time;
     bthread::TimerThread::TaskId _task_id;
 
 private:
     const char* _name;
-    int _sleep_ms;
+    butil::atomic<int> _sleep_ms;
+    butil::atomic<bool> _started{false};
+    butil::atomic<bool> _finished{false};
     std::vector<timespec> _run_times;
 };
 
 TEST(TimerThreadTest, RunTasks) {
     bthread::TimerThread timer_thread;
     ASSERT_EQ(0, timer_thread.start(nullptr));
 
-    timespec _2s_later = butil::seconds_from_now(2);
+    timespec _2s_later = butil::milliseconds_from_now(20);

Review Comment:
   Variable names like `_2s_later`, `_1s_later`, and `_10s_later` no longer 
match their actual values (20ms/10ms/3600s). Renaming them to reflect the real 
delays (e.g., `_20ms_later`, `_10ms_later`, `_1h_later`) will prevent test 
logic from being misread and incorrectly modified later.



##########
test/bthread_timer_thread_unittest.cpp:
##########
@@ -98,59 +114,70 @@ class TimeKeeper {
         keeper->run();
     }
 
+    bool wait_finished() {
+        return WaitUntil([this] {
+            return _finished.load(butil::memory_order_acquire);
+        });
+    }
+
+    bool wait_started() {
+        return WaitUntil([this] {
+            return _started.load(butil::memory_order_acquire);
+        });
+    }
+
     timespec _expect_run_time;
     bthread::TimerThread::TaskId _task_id;
 
 private:
     const char* _name;
-    int _sleep_ms;
+    butil::atomic<int> _sleep_ms;
+    butil::atomic<bool> _started{false};
+    butil::atomic<bool> _finished{false};
     std::vector<timespec> _run_times;
 };
 
 TEST(TimerThreadTest, RunTasks) {
     bthread::TimerThread timer_thread;
     ASSERT_EQ(0, timer_thread.start(nullptr));
 
-    timespec _2s_later = butil::seconds_from_now(2);
+    timespec _2s_later = butil::milliseconds_from_now(20);
     TimeKeeper keeper1(_2s_later, "keeper1");
     keeper1.schedule(&timer_thread);
 
-    TimeKeeper keeper2(_2s_later, "keeper2");  // same time with keeper1
+    TimeKeeper keeper2(butil::seconds_from_now(3600), "keeper2");
     keeper2.schedule(&timer_thread);
     
-    timespec _1s_later = butil::seconds_from_now(1);
+    timespec _1s_later = butil::milliseconds_from_now(10);

Review Comment:
   Variable names like `_2s_later`, `_1s_later`, and `_10s_later` no longer 
match their actual values (20ms/10ms/3600s). Renaming them to reflect the real 
delays (e.g., `_20ms_later`, `_10ms_later`, `_1h_later`) will prevent test 
logic from being misread and incorrectly modified later.



##########
test/bthread_timer_thread_unittest.cpp:
##########
@@ -98,59 +114,70 @@ class TimeKeeper {
         keeper->run();
     }
 
+    bool wait_finished() {
+        return WaitUntil([this] {
+            return _finished.load(butil::memory_order_acquire);
+        });
+    }
+
+    bool wait_started() {
+        return WaitUntil([this] {
+            return _started.load(butil::memory_order_acquire);
+        });
+    }
+
     timespec _expect_run_time;
     bthread::TimerThread::TaskId _task_id;
 
 private:
     const char* _name;
-    int _sleep_ms;
+    butil::atomic<int> _sleep_ms;
+    butil::atomic<bool> _started{false};
+    butil::atomic<bool> _finished{false};
     std::vector<timespec> _run_times;
 };
 
 TEST(TimerThreadTest, RunTasks) {
     bthread::TimerThread timer_thread;
     ASSERT_EQ(0, timer_thread.start(nullptr));
 
-    timespec _2s_later = butil::seconds_from_now(2);
+    timespec _2s_later = butil::milliseconds_from_now(20);
     TimeKeeper keeper1(_2s_later, "keeper1");
     keeper1.schedule(&timer_thread);
 
-    TimeKeeper keeper2(_2s_later, "keeper2");  // same time with keeper1
+    TimeKeeper keeper2(butil::seconds_from_now(3600), "keeper2");
     keeper2.schedule(&timer_thread);
     
-    timespec _1s_later = butil::seconds_from_now(1);
+    timespec _1s_later = butil::milliseconds_from_now(10);
     TimeKeeper keeper3(_1s_later, "keeper3");
     keeper3.schedule(&timer_thread);
 
-    timespec _10s_later = butil::seconds_from_now(10);
+    timespec _10s_later = butil::seconds_from_now(3600);

Review Comment:
   Variable names like `_2s_later`, `_1s_later`, and `_10s_later` no longer 
match their actual values (20ms/10ms/3600s). Renaming them to reflect the real 
delays (e.g., `_20ms_later`, `_10ms_later`, `_1h_later`) will prevent test 
logic from being misread and incorrectly modified later.



##########
test/bthread_work_stealing_queue_unittest.cpp:
##########
@@ -85,6 +85,7 @@ void* pop_thread(void* arg) {
 
 
 TEST(WSQTest, sanity) {
+    g_stop = false;

Review Comment:
   Since `g_stop` is now an atomic, prefer using `g_stop.store(false, ...)` 
(and corresponding `load(...)` at read sites) for clarity and to make the 
intended memory ordering explicit. This also avoids relying on any non-standard 
implicit conversions/operators on the atomic type.



##########
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:
   The test assumes all waiters will be registered and individually woken 
within a fixed 5s deadline, then hard-asserts `nwakeup == N`. On 
slow/oversubscribed CI, thread creation and waiter registration can 
legitimately exceed this budget, making the test flaky again. A more robust 
approach is to add an explicit readiness/registration counter (or barrier) in 
`dummy_waiter` and wait until all threads have entered the wait before starting 
the wake timing loop, or to change the post-condition to assert all threads 
eventually exit (and only report timing for those actually woken).



-- 
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