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]