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


##########
be/test/runtime/cdc_client_mgr_test.cpp:
##########
@@ -333,6 +357,358 @@ 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());
+    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);
+
+    // 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

Review Comment:
   [P2] Reap the unrelated child when this fatal checkpoint fails
   
   If this assertion times out because the handler did not run—the behavior the 
test is checking—the `exit 7` child is already waitable, but the fatal return 
skips the first `waitpid(pid, ...)` below. `mgr` only owns synthetic PID 99999, 
so its destructor cannot collect this unrelated child and the zombie remains in 
the shared BE test process for later signal tests. Please install a 
direct-child cleanup guard immediately after the successful spawn (kill only if 
still running, then reap) and disarm it after the explicit successful `waitpid`.



##########
be/src/runtime/cdc_client_mgr.cpp:
##########
@@ -50,16 +53,335 @@
 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};
+#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 ||

Review Comment:
   [P2] Preserve the current child's SIGCHLD behind a stale claim
   
   A handler can load generation N at the start of `handle_sigchld` and be 
descheduled before it claims the operation word. Another thread can then reap 
N, publish N+1, complete both startup inspections, and return; 
`_publish_child_pid()` cannot see the delayed handler because it has not 
claimed yet. If that old handler resumes and installs token N, then N+1 exits 
before the old handler clears it, the N+1 handler takes this mismatch branch 
and returns without recording a pending reap. The old handler then clears N and 
also returns without inspecting N+1, so the current child can remain published 
as a zombie until a later start/stop call (standard SIGCHLD instances are not 
guaranteed to queue another delivery). Please preserve/recheck the current 
generation when releasing a stale claim, or otherwise let the current handler 
record its pending reap, and add a deterministic cross-generation test for this 
window.



##########
be/test/runtime/cdc_client_mgr_test.cpp:
##########
@@ -333,6 +357,358 @@ 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());
+    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);
+
+    // 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);
+    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);
+}
+
+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]);
+        }
+    }};
+
+    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_addclose(&actions, report_pipe[0]), 0);
+    ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, report_pipe[1]), 0);
+
+    // The shell reports readiness and then reports that SIGTERM reached it, 
so the graceful half of
+    // the reap sequence is observable: a child taken out by the forced-kill 
fallback reports
+    // nothing, and then the pipe is at EOF rather than blocking, because the 
background sleep is
+    // detached from it. `wait` keeps the trap prompt - a shell blocked on a 
foreground child defers
+    // traps until that child exits, which is longer than the grace window.
+    pid_t pid = 0;
+    char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+                          const_cast<char*>("trap 'printf T; exit 0' TERM; 
printf R; "
+                                            "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, nullptr, argv, envp), 0);
+    ASSERT_GT(pid, 0);
+    close(report_pipe[1]);
+    report_pipe[1] = -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) {
+            kill(pid, SIGKILL);
+            waitpid(pid, nullptr, 0);
+        }
+    }};
+
+    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);

Review Comment:
   [P2] Distinguish a stale SIGTERM from the intentional stop
   
   This liveness probe does not prove that the stale identity refrained from 
signalling the replacement: it can return zero before the target processes 
SIGTERM, and it also returns zero while a signalled child is a zombie. If the 
stale cleanup wrongly sends TERM, the trap can write `T` and exit; `mgr.stop()` 
then merely reaps that zombie, and the final single-byte read consumes the 
earlier `T`, so every assertion still passes. Please synchronize with the 
helper after the stale call and assert that no TERM marker was emitted, then 
use a separate phase or marker to verify the later intentional `mgr.stop()` 
signal.



##########
be/test/runtime/cdc_client_mgr_test.cpp:
##########
@@ -158,13 +163,18 @@ TEST_F(CdcClientMgrTest, StopWithRealProcessGraceful) {
                 // Set the PID in the manager
                 mgr.set_child_pid_for_test(real_pid);
 
-                // Call stop - process will respond to SIGTERM and exit
-                // This covers the graceful shutdown path
+                // pclose() already reaped the shell, so the background sleep 
is not this process's
+                // child and stop() reaches terminate_and_reap_child only to 
find ECHILD.
                 mgr.stop();
 
                 // Verify PID is reset
                 EXPECT_EQ(mgr.get_child_pid(), 0);
 

Review Comment:
   [P2] Observe whether the non-child actually received SIGTERM
   
   This immediate `kill(pid, 0)` is not proof that `stop()` refrained from 
signalling the process: if a regression sends SIGTERM, `kill` can return before 
the detached target is scheduled and exits, so this existence probe can still 
succeed and the cleanup below then hides the signal by killing the target 
itself. The surrounding `if (pipe)`, `if (fgets(...))`, and positive-PID checks 
also let setup failures report a passing test without reaching any assertion. 
Please make setup fail closed and use a non-child helper that reports SIGTERM 
reception (for example through a pipe plus a post-`stop()` handshake) so the 
test deterministically proves no signal was delivered before cleanup.



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