github-actions[bot] commented on code in PR #68492:
URL: https://github.com/apache/doris/pull/68492#discussion_r4092566007


##########
be/test/runtime/cdc_client_mgr_test.cpp:
##########
@@ -333,6 +336,310 @@ TEST_F(CdcClientMgrTest, StartCdcClientWithResult) {
     EXPECT_GT(mgr.get_child_pid(), 0); // PID should be set
 }
 
+// Scenario: starting the cdc client installs a process-wide SIGCHLD handler, 
and a process-wide
+// handler sees every child of the BE, not just the cdc client. BE also runs 
an embedded JVM, which
+// forks children of its own for Runtime.exec() and reads their exit status 
from its process-reaper
+// thread. Reaping one of those here makes that thread find the child already 
gone, and
+// java.lang.ProcessHandleImpl turns the resulting ECHILD into exit code 0 
whatever the child
+// really returned - Java code inside BE that branches on an exit status then 
takes the wrong
+// branch silently. The handler must wait on the cdc client's pid alone.
+TEST_F(CdcClientMgrTest, SigchldHandlerDoesNotReapOtherChildren) {
+    CdcClientMgr mgr;
+    PRequestCdcClientResult result;
+    ASSERT_TRUE(mgr.start_cdc_client(&result).ok());
+
+    pid_t pid = 0;
+    char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+                          const_cast<char*>("sleep 0.2; exit 7"), nullptr};
+    char* const envp[] = {const_cast<char*>("PATH=/bin:/usr/bin"), nullptr};
+    ASSERT_EQ(posix_spawn(&pid, "/bin/sh", nullptr, nullptr, argv, envp), 0);
+    ASSERT_GT(pid, 0);
+
+    // Be somewhere other than waitpid() when the child exits, so its SIGCHLD 
reaches the handler
+    // rather than a waiter that is already blocked on this pid.
+    std::this_thread::sleep_for(std::chrono::milliseconds(600));
+
+    int child_status = 0;
+    const pid_t reaped = waitpid(pid, &child_status, 0);
+    ASSERT_EQ(reaped, pid) << "the cdc SIGCHLD handler consumed a child that 
is not the cdc client";
+    ASSERT_TRUE(WIFEXITED(child_status));
+    EXPECT_EQ(WEXITSTATUS(child_status), 7);
+
+    mgr.stop();
+}
+
+TEST_F(CdcClientMgrTest, SigchldHandlerReapsOwnedChildAndPreservesErrno) {
+    sigset_t blocked;
+    sigset_t old_mask;
+    sigemptyset(&blocked);
+    sigaddset(&blocked, SIGCHLD);
+    ASSERT_EQ(pthread_sigmask(SIG_BLOCK, &blocked, &old_mask), 0);
+    Defer restore_mask {[&]() { pthread_sigmask(SIG_SETMASK, &old_mask, 
nullptr); }};
+
+    // Earlier cases install the production SIGCHLD handler process-wide. 
Temporarily restore the
+    // default disposition so only the deterministic direct invocation below 
can collect this child;
+    // blocking SIGCHLD on this thread alone cannot stop another test/runtime 
thread receiving it.
+    struct sigaction old_action {};
+    ASSERT_EQ(sigaction(SIGCHLD, nullptr, &old_action), 0);
+    struct sigaction default_action {};
+    default_action.sa_handler = SIG_DFL;
+    sigemptyset(&default_action.sa_mask);
+    ASSERT_EQ(sigaction(SIGCHLD, &default_action, nullptr), 0);
+    Defer restore_action {[&]() { sigaction(SIGCHLD, &old_action, nullptr); }};
+
+    pid_t pid = 0;
+    char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+                          const_cast<char*>("exit 7"), nullptr};
+    char* const envp[] = {const_cast<char*>("PATH=/bin:/usr/bin"), nullptr};
+    ASSERT_EQ(posix_spawn(&pid, "/bin/sh", nullptr, nullptr, argv, envp), 0);
+    ASSERT_GT(pid, 0);
+
+    CdcClientMgr mgr;
+    mgr.set_child_pid_for_test(pid);
+
+    siginfo_t child_info {};
+    for (int i = 0; i < 100; ++i) {
+        ASSERT_EQ(waitid(P_PID, pid, &child_info, WEXITED | WNOHANG | 
WNOWAIT), 0);
+        if (child_info.si_pid == pid) {
+            break;
+        }
+        std::this_thread::sleep_for(std::chrono::milliseconds(10));
+    }
+    ASSERT_EQ(child_info.si_pid, pid);
+
+    errno = EBUSY;
+    CdcClientMgr::invoke_sigchld_handler_for_test();
+    EXPECT_EQ(errno, EBUSY);
+    EXPECT_EQ(mgr.get_child_pid(), 0)
+            << "reaping the owned child must also revoke the manager's 
ownership";
+
+    int status = 0;
+    errno = 0;
+    EXPECT_EQ(waitpid(pid, &status, WNOHANG), -1);
+    EXPECT_EQ(errno, ECHILD) << "the handler must collect the cdc child 
itself";
+
+    // Once ownership is revoked, stop() must not act on another live child of 
the same BE. This
+    // covers the dangerous same-parent case: waitpid() would accept that 
child, unlike a reused PID
+    // owned by another process.
+    pid_t unrelated_pid = 0;
+    char* const unrelated_argv[] = {const_cast<char*>("sh"), 
const_cast<char*>("-c"),
+                                    const_cast<char*>("sleep 10"), nullptr};
+    ASSERT_EQ(posix_spawn(&unrelated_pid, "/bin/sh", nullptr, nullptr, 
unrelated_argv, envp), 0);
+    ASSERT_GT(unrelated_pid, 0);
+    Defer cleanup_unrelated {[&]() {
+        kill(unrelated_pid, SIGKILL);
+        waitpid(unrelated_pid, nullptr, 0);
+    }};
+
+    mgr.stop();
+    EXPECT_EQ(kill(unrelated_pid, 0), 0)
+            << "stop() signalled a child after the CDC ownership had been 
revoked";
+}
+
+TEST_F(CdcClientMgrTest, StopWaitsForAHandlerHoldingTheOldChildIdentity) {
+    sigset_t blocked;
+    sigset_t old_mask;
+    sigemptyset(&blocked);
+    sigaddset(&blocked, SIGCHLD);
+    ASSERT_EQ(pthread_sigmask(SIG_BLOCK, &blocked, &old_mask), 0);
+    Defer restore_mask {[&]() { pthread_sigmask(SIG_SETMASK, &old_mask, 
nullptr); }};
+
+    struct sigaction old_action {};
+    ASSERT_EQ(sigaction(SIGCHLD, nullptr, &old_action), 0);
+    struct sigaction default_action {};
+    default_action.sa_handler = SIG_DFL;
+    sigemptyset(&default_action.sa_mask);
+    ASSERT_EQ(sigaction(SIGCHLD, &default_action, nullptr), 0);
+    Defer restore_action {[&]() { sigaction(SIGCHLD, &old_action, nullptr); }};
+
+    pid_t pid = 0;
+    char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+                          const_cast<char*>("sleep 10"), nullptr};
+    char* const envp[] = {const_cast<char*>("PATH=/bin:/usr/bin"), nullptr};
+    ASSERT_EQ(posix_spawn(&pid, "/bin/sh", nullptr, nullptr, argv, envp), 0);
+    ASSERT_GT(pid, 0);
+
+    CdcClientMgr mgr;
+    mgr.set_child_pid_for_test(pid);
+    CdcClientMgr::pause_sigchld_handler_for_test(true);
+    std::thread handler([]() { 
CdcClientMgr::invoke_sigchld_handler_for_test(); });
+
+    for (int i = 0; i < 100 && 
!CdcClientMgr::sigchld_handler_paused_for_test(); ++i) {
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    if (!CdcClientMgr::sigchld_handler_paused_for_test()) {
+        CdcClientMgr::pause_sigchld_handler_for_test(false);
+        handler.join();
+        kill(pid, SIGKILL);
+        waitpid(pid, nullptr, 0);
+        FAIL() << "the deterministic handler did not reach its pause point";
+    }
+
+    std::atomic<bool> stop_finished {false};

Review Comment:
   [P2] Prove the stopper reached the ownership gate
   
   The 20 ms sleep does not establish that the stopper has entered `stop()` 
while the handler still owns the operation token. On a loaded runner the 
stopper can remain unscheduled, `stop_finished` is still false, and the test 
then unpauses the handler before the stopper ever attempts the claim; the test 
passes without exercising the exclusion it names. Please add a deterministic 
test checkpoint/barrier showing the stopper reached the claim attempt while the 
handler is paused, then release the handler and verify the stopper crosses only 
afterward.



##########
be/src/runtime/cdc_client_mgr.cpp:
##########
@@ -50,16 +53,330 @@
 namespace doris {
 
 namespace {
-// Handle SIGCHLD signal to prevent zombie processes
+// The identity of the cdc client this process forked, published for 
handle_sigchld(). A signal
+// handler may only touch lock-free atomics, so the pid and its generation 
live in one 64-bit word
+// rather than behind CdcClientMgr's mutex. The generation prevents a delayed 
handler from clearing
+// ownership after the kernel has already reused the same numeric pid for a 
replacement child.
+// ExecEnv owns a single CdcClientMgr, so there is a single published identity.
+static_assert(sizeof(pid_t) <= sizeof(uint32_t));
+static_assert(std::atomic<uint64_t>::is_always_lock_free);
+static_assert(std::atomic<uint32_t>::is_always_lock_free);
+std::atomic<uint64_t> g_cdc_child_identity {0};
+std::atomic<uint32_t> g_cdc_child_generation {0};
+// The identity whose OS process operations are currently owned by exactly one 
actor. A handler or
+// normal thread must claim the published identity here BEFORE waitpid/kill 
and keep the claim until
+// its last syscall. While claimed, the original child is either running or 
remains a zombie, so its
+// numeric pid cannot be reused for another same-parent child.
+std::atomic<uint64_t> g_cdc_child_operation {0};
+// pid_t is signed and a published child pid is positive, so bit 31 of the pid 
half is free for a
+// pending-reap handoff. Keeping the request in the SAME atomic word as the 
operation claim closes
+// the release/request race: either the owner observes the bit before 
releasing, or the handler wins
+// the release CAS and claims the now-idle identity itself.
+constexpr uint64_t CDC_CHILD_REAP_PENDING = uint64_t {1} << 31;
+static_assert(std::atomic<uint64_t>::is_always_lock_free);
+
+#ifdef BE_TEST
+static_assert(std::atomic<bool>::is_always_lock_free);
+std::atomic<bool> g_pause_cdc_sigchld_handler {false};
+std::atomic<bool> g_cdc_sigchld_handler_paused {false};
+std::atomic<bool> g_pause_cdc_child_inspection_after_running {false};
+std::atomic<bool> g_cdc_child_inspection_paused {false};
+#endif
+
+pid_t child_pid(uint64_t identity) {
+    return static_cast<pid_t>(static_cast<uint32_t>(identity));
+}
+
+uint64_t new_child_identity(pid_t pid) {
+    const uint64_t generation = g_cdc_child_generation.fetch_add(1, 
std::memory_order_relaxed) + 1;
+    return (generation << 32) | static_cast<uint32_t>(pid);
+}
+
+// Signal-safe, non-blocking acquisition. Revalidate after publishing the 
claim: a normal thread may
+// have revoked the identity between the first load and this CAS.
+bool try_claim_child_identity(uint64_t identity) {
+    if (identity == 0) {
+        return false;
+    }
+    uint64_t unclaimed = 0;
+    if (!g_cdc_child_operation.compare_exchange_strong(unclaimed, identity)) {
+        return false;
+    }
+    if (g_cdc_child_identity.load() == identity) {
+        return true;
+    }
+    // The publication was revoked while the CAS was in flight. No actor can 
acquire the operation
+    // until this store, and a pending request for the revoked generation no 
longer needs service.
+    g_cdc_child_operation.store(0);
+    return false;
+}
+
+// Signal-safe acquisition with a lossless handoff when a normal thread 
already owns the claim.
+// The pending bit and claim are changed by one CAS, so a concurrent release 
cannot pass between
+// "request reap" and "observe idle" and lose the only SIGCHLD notification.
+bool try_claim_or_request_child_reap(uint64_t identity) {
+    if (identity == 0) {
+        return false;
+    }
+    uint64_t operation = 0;
+    while (true) {
+        if (operation == 0) {
+            if (g_cdc_child_operation.compare_exchange_weak(operation, 
identity)) {
+                if (g_cdc_child_identity.load() == identity) {
+                    return true;
+                }
+                g_cdc_child_operation.store(0);
+                return false;
+            }
+            continue;
+        }
+        if ((operation & ~CDC_CHILD_REAP_PENDING) != identity ||
+            (operation & CDC_CHILD_REAP_PENDING) != 0) {
+            return false;
+        }
+        if (g_cdc_child_operation.compare_exchange_weak(operation,
+                                                        operation | 
CDC_CHILD_REAP_PENDING)) {
+            return false;
+        }
+    }
+}
+
+// Normal threads may wait for a handler's short WNOHANG operation. Returning 
false means the exact
+// generation was revoked; the caller must not operate on its numeric pid.
+bool claim_child_identity(uint64_t identity) {
+    while (g_cdc_child_identity.load() == identity) {
+        if (try_claim_child_identity(identity)) {
+            return true;
+        }
+        std::this_thread::yield();
+    }
+    return false;
+}
+
+// Returns true when the claim was released. A false return means a handler 
attached a reap request
+// to this exact claim; the owner consumed it and must repeat waitpid before 
trying to release again.
+// The final CAS races atomically with the handler's pending-bit CAS, which is 
the missed-wakeup fence.
+bool release_child_identity_if_quiescent(uint64_t identity) {
+    while (true) {
+        uint64_t operation = identity;
+        if (g_cdc_child_operation.compare_exchange_strong(operation, 0)) {
+            return true;
+        }
+        if (operation != (identity | CDC_CHILD_REAP_PENDING)) {
+            // Only this owner can change a non-pending claim. Treat an 
already released/revoked
+            // token as quiescent rather than touching another generation.
+            return true;
+        }
+        if (g_cdc_child_operation.compare_exchange_strong(operation, 
identity)) {
+            return false;
+        }
+    }
+}
+
+void release_terminal_child_identity(uint64_t identity) {
+    // The child is already unpublished and reaped. Consume any handler 
request that was based on a
+    // pre-revocation identity, but no further waitpid is necessary.
+    while (!release_child_identity_if_quiescent(identity)) {
+    }
+}
+
+// Reap the cdc client so it does not linger as a zombie.
+//
+// waitpid(-1) here would reap ANY child of this process, including the ones 
the embedded JVM forks
+// for Runtime.exec(): the JVM's process-reaper thread would then find its own 
child already gone,
+// and java.lang.ProcessHandleImpl turns that ECHILD into exit code 0 no 
matter what the child
+// really returned. Java code inside BE that branches on an exit status would 
silently take the
+// wrong branch - which is how frocksdbjni's `ldd /usr/bin/env | grep -q musl` 
probe answered "yes"
+// on a glibc host and loaded the musl build of librocksdbjni.so. Wait for our 
own pid only.
 void handle_sigchld(int sig_no) {
+    const int saved_errno = errno;
+    const uint64_t cdc_identity = g_cdc_child_identity.load();
+    const pid_t cdc_pid = child_pid(cdc_identity);
+    // Never retain a raw pid without the operation claim. If another actor 
owns this identity, the
+    // helper attaches a pending-reap handoff to that claim before this 
handler returns.
+    if (cdc_pid <= 0 || !try_claim_or_request_child_reap(cdc_identity)) {
+        errno = saved_errno;
+        return;
+    }
+#ifdef BE_TEST
+    if (g_pause_cdc_sigchld_handler.load()) {
+        g_cdc_sigchld_handler_paused.store(true);
+        while (g_pause_cdc_sigchld_handler.load()) {
+        }
+        g_cdc_sigchld_handler_paused.store(false);
+    }
+#endif
+    while (true) {
+        int status = 0;
+        pid_t wait_result;
+        do {
+            wait_result = waitpid(cdc_pid, &status, WNOHANG);
+        } while (wait_result < 0 && errno == EINTR);
+        if (wait_result == cdc_pid || (wait_result < 0 && errno == ECHILD)) {
+            uint64_t expected = cdc_identity;
+            g_cdc_child_identity.compare_exchange_strong(expected, 0);
+            // Reaping makes the numeric pid reusable. Consume stale handoffs 
without issuing
+            // another syscall against a pid that may already belong to a 
different child.
+            release_terminal_child_identity(cdc_identity);
+            break;
+        }
+        // No syscall may use cdc_pid after a successful release. If a 
concurrent handler marked
+        // this claim pending, consume the request and repeat waitpid before 
releasing instead.
+        if (release_child_identity_if_quiescent(cdc_identity)) {
+            break;
+        }
+    }
+    errno = saved_errno;
+}
+
+// Terminate and collect one child owned by this process. The SIGCHLD handler 
may have won the
+// waitpid race already; ECHILD is therefore success, not a reason to send a 
signal to a reused pid.
+void terminate_and_reap_child(pid_t pid) {
+    if (pid <= 0) {
+        return;
+    }
+
     int status = 0;
-    pid_t pid;
-    while ((pid = waitpid(-1, &status, WNOHANG)) > 0) {
+    pid_t wait_result;
+    do {
+        wait_result = waitpid(pid, &status, WNOHANG);
+    } while (wait_result < 0 && errno == EINTR);
+    if (wait_result == pid || (wait_result < 0 && errno == ECHILD)) {

Review Comment:
   [P2] Keep the forced-kill path under test
   
   This new `ECHILD` return makes `StopWithRealProcessForceKill` pass before it 
exercises the fallback below. That test gets a background PID from a shell 
launched with `popen`; once `pclose()` reaps the shell, the background process 
is not this test process's child, so `waitpid(real_pid, ..., WNOHANG)` returns 
`ECHILD` here and neither `SIGTERM` nor `SIGKILL` is sent by `stop()`. The new 
`posix_spawn` cases use children that accept `SIGTERM`, so they cover graceful 
termination but not lines 263-270. Please make the ignore-`SIGTERM` case a 
direct child (and retain failure cleanup) so the changed forced-kill/reap path 
is actually verified.



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