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]