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]