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


##########
be/test/runtime/cdc_client_mgr_test.cpp:
##########
@@ -333,6 +462,503 @@ 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.
+// The failure cleanup and handler checkpoint are part of one 
process-lifecycle test.
+// NOLINTNEXTLINE(readability-function-cognitive-complexity)
+TEST_F(CdcClientMgrTest, SigchldHandlerDoesNotReapOtherChildren) {
+    CdcClientMgr mgr;
+    PRequestCdcClientResult result;
+    ASSERT_TRUE(mgr.start_cdc_client(&result).ok());
+    ASSERT_GT(mgr.get_child_pid(), 0);
+
+    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);
+    bool child_reaped = false;
+    Defer cleanup_child {[&]() {
+        if (!child_reaped) {
+            reap_direct_child_if_present(pid);
+        }
+    }};
+
+    // Wait for the handler itself to run before reaping anything. The handler 
unpublishes the
+    // identity it was handed once its own waitpid has returned, so the 
published test child
+    // disappearing is the observable proof that the handler already executed 
for this child's exit.
+    // A fixed sleep proves nothing: on a delayed delivery the blocking 
waitpid below would win the
+    // race and reap the child itself, passing the case without exercising the 
replacement for
+    // waitpid(-1).
+    for (int i = 0; i < 5000 && mgr.get_child_pid() != 0; ++i) {
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_EQ(mgr.get_child_pid(), 0)
+            << "the SIGCHLD handler did not run to completion for the 
unrelated child's exit";
+
+    int child_status = 0;
+    const pid_t reaped = waitpid(pid, &child_status, 0);
+    child_reaped = reaped == pid;
+    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";
+    }
+
+    CdcClientMgr::reset_child_claim_failed_for_test();
+    std::atomic<bool> stop_finished {false};
+    std::thread stopper([&]() {
+        mgr.stop();
+        stop_finished.store(true);
+    });
+    // The handler keeps the identity published while it owns the 
process-operation claim. That
+    // prevents a replacement generation from publishing the same numeric pid 
until the handler's
+    // final syscall has completed.
+    EXPECT_EQ(mgr.get_child_pid(), pid);
+    for (int i = 0; i < 5000 && !CdcClientMgr::child_claim_failed_for_test(); 
++i) {
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    EXPECT_TRUE(CdcClientMgr::child_claim_failed_for_test())
+            << "stop did not attempt to claim the child while the handler was 
paused";
+    EXPECT_FALSE(stop_finished.load())
+            << "stop returned while a signal handler could still operate the 
old numeric pid";
+
+    CdcClientMgr::pause_sigchld_handler_for_test(false);
+    handler.join();
+    stopper.join();
+    EXPECT_EQ(mgr.get_child_pid(), 0);
+
+    errno = 0;
+    EXPECT_EQ(waitpid(pid, nullptr, WNOHANG), -1);
+    EXPECT_EQ(errno, ECHILD);
+}
+
+// Both signal phases and their cleanup must remain visible in one test.
+// NOLINTNEXTLINE(readability-function-cognitive-complexity)
+TEST_F(CdcClientMgrTest, StaleGenerationCannotTerminateAReusedNumericPid) {
+    int report_pipe[2];
+    ASSERT_EQ(pipe(report_pipe), 0);
+    Defer close_pipe {[&]() {
+        close(report_pipe[0]);
+        if (report_pipe[1] >= 0) {
+            close(report_pipe[1]);
+        }
+    }};
+    int phase_pipe[2];
+    ASSERT_EQ(pipe(phase_pipe), 0);
+    Defer close_phase_pipe {[&]() {
+        if (phase_pipe[0] >= 0) {
+            close(phase_pipe[0]);
+        }
+        close(phase_pipe[1]);
+    }};
+
+    struct sigaction old_pipe_action {};
+    ASSERT_EQ(sigaction(SIGPIPE, nullptr, &old_pipe_action), 0);
+    struct sigaction ignored_pipe_action {};
+    ignored_pipe_action.sa_handler = SIG_IGN;
+    sigemptyset(&ignored_pipe_action.sa_mask);
+    ASSERT_EQ(sigaction(SIGPIPE, &ignored_pipe_action, nullptr), 0);
+    Defer restore_pipe_action {[&]() { sigaction(SIGPIPE, &old_pipe_action, 
nullptr); }};
+
+    posix_spawn_file_actions_t actions;
+    ASSERT_EQ(posix_spawn_file_actions_init(&actions), 0);
+    Defer destroy_actions {[&]() { posix_spawn_file_actions_destroy(&actions); 
}};
+    ASSERT_EQ(posix_spawn_file_actions_adddup2(&actions, report_pipe[1], 
STDOUT_FILENO), 0);
+    ASSERT_EQ(posix_spawn_file_actions_adddup2(&actions, phase_pipe[0], 
STDIN_FILENO), 0);
+    ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, report_pipe[0]), 0);
+    ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, report_pipe[1]), 0);
+    ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, phase_pipe[0]), 0);
+    ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, phase_pipe[1]), 0);
+
+    // Give the child its own signal environment instead of the inherited one: 
all of the suites run
+    // in one process, so any case that leaves SIGTERM ignored or blocked 
behind it decides whether
+    // this child can act on the signal at all. A non-interactive shell 
silently refuses to trap a
+    // signal that was ignored on entry, and a blocked one is never delivered; 
both leave the report
+    // below empty, which is indistinguishable from the forced kill. SIGTERM 
therefore starts at its
+    // default disposition and unblocked, and SIGCHLD the same so that the 
child's own `wait` works.
+    posix_spawnattr_t attr;
+    ASSERT_EQ(posix_spawnattr_init(&attr), 0);
+    Defer destroy_attr {[&]() { posix_spawnattr_destroy(&attr); }};
+    sigset_t no_signals;
+    sigemptyset(&no_signals);
+    sigset_t default_signals;
+    sigemptyset(&default_signals);
+    sigaddset(&default_signals, SIGTERM);
+    sigaddset(&default_signals, SIGCHLD);
+    ASSERT_EQ(posix_spawnattr_setsigmask(&attr, &no_signals), 0);
+    ASSERT_EQ(posix_spawnattr_setsigdefault(&attr, &default_signals), 0);
+    ASSERT_EQ(posix_spawnattr_setflags(&attr, POSIX_SPAWN_SETSIGMASK | 
POSIX_SPAWN_SETSIGDEF), 0);
+
+    // The shell reports S if the stale identity signals it. It acknowledges 
the phase change with
+    // A only after processing the first phase, then reports T for the 
intentional stop().
+    pid_t pid = 0;
+    char* const argv[] = {
+            const_cast<char*>("sh"), const_cast<char*>("-c"),
+            const_cast<char*>("trap 'printf S' TERM; printf R; "
+                              "IFS= read -r phase; trap 'printf T; exit 0' 
TERM; "
+                              "printf A; sleep 10 </dev/null >/dev/null 2>&1 & 
wait $!"),
+            nullptr};
+    char* const envp[] = {const_cast<char*>("PATH=/bin:/usr/bin"), nullptr};
+    ASSERT_EQ(posix_spawn(&pid, "/bin/sh", &actions, &attr, argv, envp), 0);
+    ASSERT_GT(pid, 0);
+    close(report_pipe[1]);
+    report_pipe[1] = -1;
+    close(phase_pipe[0]);
+    phase_pipe[0] = -1;
+
+    // Declare the manager before the cleanup guard so that a fatal assertion 
destroys the guard
+    // first: it kills and reaps this direct child, and the manager's stop() 
then observes ECHILD
+    // instead of signalling a numeric pid whose ownership it has already 
given up.
+    CdcClientMgr mgr;
+    bool child_needs_cleanup = true;
+    Defer cleanup {[&]() {
+        if (child_needs_cleanup) {
+            reap_direct_child_if_present(pid);
+        }
+    }};
+
+    char ready = 0;
+    ASSERT_EQ(read(report_pipe[0], &ready, 1), 1);
+    ASSERT_EQ(ready, 'R');
+
+    const uint64_t old_identity = mgr.set_child_pid_for_test(pid);
+    // Republish the same numeric pid under a new generation. This 
deterministically models the
+    // kernel reusing a reaped CDC pid for another same-parent child without 
depending on PID churn.
+    const uint64_t replacement_identity = mgr.set_child_pid_for_test(pid);
+    ASSERT_NE(old_identity, replacement_identity);
+    ASSERT_EQ(mgr.get_child_identity_for_test(), replacement_identity);
+
+    ASSERT_FALSE(mgr.terminate_child_identity_for_test(old_identity));
+    ASSERT_EQ(write(phase_pipe[1], "\n", 1), 1);
+    char phase_ack = 0;
+    ASSERT_EQ(read(report_pipe[0], &phase_ack, 1), 1);
+    ASSERT_EQ(phase_ack, 'A') << "the stale generation delivered SIGTERM 
before stop()";
+
+    mgr.stop();
+    EXPECT_EQ(mgr.get_child_identity_for_test(), 0);
+    errno = 0;
+    const pid_t wait_result = waitpid(pid, nullptr, WNOHANG);
+    const int wait_error = errno;
+    child_needs_cleanup = wait_result < 0 && wait_error == ECHILD;
+    EXPECT_EQ(wait_result, -1);
+    EXPECT_EQ(wait_error, ECHILD);
+
+    char handled = 0;
+    EXPECT_EQ(read(report_pipe[0], &handled, 1), 1)
+            << "the published child was collected without reporting the 
SIGTERM trap: it either "
+               "never got to run it, or was not scheduled inside the grace 
window";
+    EXPECT_EQ(handled, 'T') << "stop() fell through the grace window to the 
forced kill";
+}
+
+// The controlled interleaving and cleanup must remain visible in one test.
+// NOLINTNEXTLINE(readability-function-cognitive-complexity)
+TEST_F(CdcClientMgrTest, DelayedHandlerReapsExitedSuccessorGeneration) {
+    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); }};
+
+    CdcClientMgr mgr;
+    const uint64_t old_identity = mgr.set_child_pid_for_test(99999);
+    ASSERT_NE(old_identity, 0);
+    CdcClientMgr::pause_sigchld_before_claim_for_test(true);
+    Defer resume_handlers {[]() {
+        CdcClientMgr::pause_sigchld_before_claim_for_test(false);
+        CdcClientMgr::pause_stale_child_claim_for_test(false);
+    }};
+    std::thread delayed_handler([]() { 
CdcClientMgr::invoke_sigchld_handler_for_test(); });
+    Defer resume_and_join_delayed {[&]() {
+        CdcClientMgr::pause_sigchld_before_claim_for_test(false);
+        CdcClientMgr::pause_stale_child_claim_for_test(false);
+        if (delayed_handler.joinable()) {
+            delayed_handler.join();
+        }
+    }};
+    for (int i = 0; i < 5000 && 
!CdcClientMgr::sigchld_before_claim_paused_for_test(); ++i) {
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_TRUE(CdcClientMgr::sigchld_before_claim_paused_for_test());
+
+    pid_t pid = 0;
+    char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+                          const_cast<char*>("exec sleep 3600"), 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);
+    bool child_reaped = false;
+    Defer cleanup_child {[&]() {
+        if (!child_reaped) {
+            reap_direct_child_if_present(pid);
+        }
+    }};
+
+    const uint64_t successor_identity = mgr.set_child_pid_for_test(pid);
+    ASSERT_NE(successor_identity, old_identity);
+    ASSERT_TRUE(mgr.inspect_child_identity_for_test(successor_identity));
+    ASSERT_TRUE(mgr.inspect_child_identity_for_test(successor_identity));
+
+    CdcClientMgr::pause_stale_child_claim_for_test(true);
+    CdcClientMgr::pause_sigchld_before_claim_for_test(false);
+    for (int i = 0; i < 5000 && 
!CdcClientMgr::stale_child_claim_paused_for_test(); ++i) {
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_TRUE(CdcClientMgr::stale_child_claim_paused_for_test());
+
+    ASSERT_EQ(kill(pid, SIGKILL), 0);
+    siginfo_t child_info {};
+    for (int i = 0; i < 5000 && child_info.si_pid != pid; ++i) {
+        ASSERT_EQ(waitid(P_PID, pid, &child_info, WEXITED | WNOHANG | 
WNOWAIT), 0);
+        std::this_thread::sleep_for(std::chrono::milliseconds(1));
+    }
+    ASSERT_EQ(child_info.si_pid, pid);
+
+    // This handler sees the successor, but the stale claim prevents it from 
recording a pending
+    // reap. The delayed handler must take responsibility after releasing that 
stale claim.
+    CdcClientMgr::invoke_sigchld_handler_for_test();

Review Comment:
   [P2] Bound the handler-contention checks
   
   This synchronous invocation runs while the delayed handler deliberately 
holds the stale generation's operation token, and that owner is not released 
until line 844. If the foreign-token path regresses to waiting for the token, 
this call never returns and its only release is unreachable, hanging the whole 
shared BE test binary instead of failing the case. The same circular harness 
dependency occurs at `second.join()` around line 884 and the synchronous 
invocation around line 950. Please run each challenged handler asynchronously, 
observe completion with a bounded deadline, then always release/join the owner 
before failing (or isolate these liveness checks in a subprocess with a 
timeout).



##########
be/src/runtime/cdc_client_mgr.cpp:
##########
@@ -50,16 +53,366 @@
 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_cdc_child_claim_failed {false};
+std::atomic<bool> g_pause_cdc_child_inspection_after_running {false};
+std::atomic<bool> g_cdc_child_inspection_paused {false};
+std::atomic<bool> g_pause_cdc_sigchld_before_claim {false};
+std::atomic<bool> g_cdc_sigchld_before_claim_paused {false};
+std::atomic<bool> g_pause_cdc_stale_child_claim {false};
+std::atomic<bool> g_cdc_stale_child_claim_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)) {
+#ifdef BE_TEST
+                if (g_pause_cdc_stale_child_claim.load() &&
+                    g_cdc_child_identity.load() != identity) {
+                    g_cdc_stale_child_claim_paused.store(true);
+                    while (g_pause_cdc_stale_child_claim.load()) {
+                    }
+                    g_cdc_stale_child_claim_paused.store(false);
+                }
+#endif
+                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) {
+    // Identity 0 names no generation: the loop below could never observe it 
change and would spin
+    // forever while stop() or start_cdc_client() holds _start_mutex.
+    DORIS_CHECK(identity != 0);
+    while (g_cdc_child_identity.load() == identity) {
+        if (try_claim_child_identity(identity)) {
+            return true;
+        }
+#ifdef BE_TEST
+        g_cdc_child_claim_failed.store(true);
+#endif
+        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;
+    uint64_t cdc_identity = g_cdc_child_identity.load();
+    while (child_pid(cdc_identity) > 0) {
+#ifdef BE_TEST
+        if (g_pause_cdc_sigchld_before_claim.load()) {
+            g_cdc_sigchld_before_claim_paused.store(true);
+            while (g_pause_cdc_sigchld_before_claim.load()) {
+            }
+            g_cdc_sigchld_before_claim_paused.store(false);
+        }
+#endif
+        // 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 (!try_claim_or_request_child_reap(cdc_identity)) {
+            const uint64_t current_identity = g_cdc_child_identity.load();
+            if (current_identity == cdc_identity) {
+                break;
+            }
+            // A delayed handler can claim a revoked generation after its 
successor was published.
+            // Its stale claim briefly blocks the successor's handler. Recheck 
the current generation
+            // after releasing that claim so the successor's only SIGCHLD is 
not lost.
+            cdc_identity = current_identity;
+            continue;
+        }
+#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
+        const pid_t cdc_pid = child_pid(cdc_identity);
+        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;
+            }
+        }
+        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) {
+    DORIS_CHECK(pid > 0);
+
     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)) {
+        return;
+    }
+
+    LOG(INFO) << "Stopping CDC client process, pid=" << pid;
+    if (kill(pid, SIGTERM) != 0 && errno != ESRCH) {
+        LOG(WARNING) << "Failed to terminate CDC client process, pid=" << pid
+                     << ", error=" << strerror(errno);
+    }
+    for (int i = 0; i < 20; ++i) {
+        do {
+            wait_result = waitpid(pid, &status, WNOHANG);
+        } while (wait_result < 0 && errno == EINTR);
+        if (wait_result == pid || (wait_result < 0 && errno == ECHILD)) {
+            return;
+        }
+        std::this_thread::sleep_for(std::chrono::milliseconds(10));
+    }
+
+    LOG(INFO) << "Force killing CDC client process, pid=" << pid;
+    if (kill(pid, SIGKILL) != 0 && errno != ESRCH) {
+        LOG(WARNING) << "Failed to kill CDC client process, pid=" << pid
+                     << ", error=" << strerror(errno);
     }
+    do {
+        wait_result = waitpid(pid, &status, 0);

Review Comment:
   [P2] Keep failed-start cleanup bounded after SIGKILL
   
   The health-check failure path now reaches this helper while 
`start_cdc_client()` still holds `_start_mutex`, but after `SIGKILL` this 
blocking `waitpid(..., 0)` has no deadline. A task stuck in uninterruptible I/O 
can remain un-waitable despite the pending kill, leaving the request, the 
operation token, every later start/stop, and eventual manager destruction 
blocked indefinitely. Doris already handles this condition in 
`PythonUDFProcess::shutdown()` with a bounded post-kill wait and a background 
reaper. Please preserve the exact-generation ownership guarantee by 
transferring an unreaped child to a unique per-PID asynchronous owner before 
releasing the foreground claim, so failed startup can return without blocking 
successor publication.



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