This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new a6431d7da1f [fix](be) Reap only the owned CDC client process (#68492)
a6431d7da1f is described below
commit a6431d7da1f4dc99db5d3e2e711a402f37e6c176
Author: Mingyu Chen (Rayner) <[email protected]>
AuthorDate: Mon Sep 28 16:16:16 2026 +0800
[fix](be) Reap only the owned CDC client process (#68492)
### What problem does this PR solve?
Issue Number: N/A
Problem Summary:
When BE starts the CDC client, `CdcClientMgr` installs a process-wide
SIGCHLD handler. The handler previously called `waitpid(-1, ...,
WNOHANG)`, which could reap any child of BE. The embedded JVM also
starts subprocesses and waits for their exit codes. If the CDC handler
reaps one first, the JVM can observe `ECHILD` and report exit code 0
even when the command failed. One observed result was a failed
libc-detection probe appearing successful, causing RocksDB JNI to select
the musl library on a glibc host.
The handler now waits only for the CDC child. A generation-qualified
identity and an exclusive operation claim coordinate it with startup,
health inspection, and shutdown, so those paths cannot act on a PID
after it has been reaped and reused. A pending-reap handoff preserves
SIGCHLD notifications arriving during inspection. The latest review
identified another handoff window: an old handler could read generation
N before claiming, then claim N after N+1 was published and hide N+1's
only exit notification. The handler now rechecks the published
generation after releasing that stale claim. A deterministic test pauses
at both sides of this window and verifies that N+1 is reaped.
The affected tests now fail if their process setup fails, observe
SIGTERM through child-reported markers and phase handshakes instead of
`kill(pid, 0)`, and reap an unrelated direct child even when a fatal
handler checkpoint fails. The forced-shutdown test uses a direct child
that ignores SIGTERM; the handler/shutdown test waits for an actual
failed ownership claim; and the stale-generation helper starts with a
controlled signal environment.
The ownership rule changes what a failed startup does to its forked
child. The health-check failure path now terminates and reaps it, where
the PID used to be dropped and the child left running unmanaged. That
abandoned child could finish warming up and be adopted by a later
`start_cdc_client` call. On a host whose cold `cdc-client.jar` JVM
consistently answers `/actuator/health` later than the unchanged startup
window, pre-start the CDC client so the adoption path finds it, or give
the JVM more room. The cleanup also removes the unmanaged JVM left
holding `cdc_client_port`.
The macOS BE test compatibility changes, including the Decimal64 literal
fix, were merged separately in #68520. This branch was rebased onto that
merge and its duplicate test-only commit was dropped.
---
be/src/runtime/cdc_client_mgr.cpp | 579 ++++++++++++++++++++++--
be/src/runtime/cdc_client_mgr.h | 34 +-
be/test/runtime/cdc_client_mgr_test.cpp | 754 +++++++++++++++++++++++++++++---
3 files changed, 1258 insertions(+), 109 deletions(-)
diff --git a/be/src/runtime/cdc_client_mgr.cpp
b/be/src/runtime/cdc_client_mgr.cpp
index b60c2c60bb1..acd742075a7 100644
--- a/be/src/runtime/cdc_client_mgr.cpp
+++ b/be/src/runtime/cdc_client_mgr.cpp
@@ -33,7 +33,10 @@
#endif
#include <atomic>
+#include <cerrno>
#include <chrono>
+#include <cstdint>
+#include <cstring>
#include <iterator>
#include <mutex>
#include <sstream>
@@ -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);
+ } while (wait_result < 0 && errno == EINTR);
}
-// Check CDC client health
+// Every normal-thread wait/signal is fenced by the exact published
generation. Clearing publication
+// while the operation claim is held transfers exclusive responsibility from
the signal handler to
+// this thread; the claim is released only after the child is reaped.
+bool terminate_owned_child(uint64_t identity) {
+ if (!claim_child_identity(identity)) {
+ return false;
+ }
+ uint64_t expected = identity;
+ if (!g_cdc_child_identity.compare_exchange_strong(expected, 0)) {
+ release_terminal_child_identity(identity);
+ return false;
+ }
+ terminate_and_reap_child(child_pid(identity));
+ release_terminal_child_identity(identity);
+ return true;
+}
+
+enum class OwnedChildState {
+ RUNNING,
+ EXITED,
+ NOT_OWNED,
+ WAIT_ERROR,
+};
+
+// Observes an owned child without leaving a raw pid usable after the
ownership claim. An exited child
+// is reaped and unpublished before the claim is released; a running child
remains published.
+[[maybe_unused]] OwnedChildState inspect_owned_child(uint64_t identity, int*
status,
+ int* wait_error) {
+ if (!claim_child_identity(identity)) {
+ return OwnedChildState::NOT_OWNED;
+ }
+
+ int local_status = 0;
+ int local_wait_error = 0;
+ OwnedChildState state = OwnedChildState::RUNNING;
+ while (true) {
+ pid_t wait_result;
+ do {
+ wait_result = waitpid(child_pid(identity), &local_status, WNOHANG);
+ } while (wait_result < 0 && errno == EINTR);
+ local_wait_error = wait_result < 0 ? errno : 0;
+
+ state = OwnedChildState::RUNNING;
+ if (wait_result == child_pid(identity) || (wait_result < 0 &&
local_wait_error == ECHILD)) {
+ uint64_t expected = identity;
+ g_cdc_child_identity.compare_exchange_strong(expected, 0);
+ state = OwnedChildState::EXITED;
+ } else if (wait_result < 0) {
+ state = OwnedChildState::WAIT_ERROR;
+ }
+#ifdef BE_TEST
+ if (wait_result == 0 &&
g_pause_cdc_child_inspection_after_running.load()) {
+ g_cdc_child_inspection_paused.store(true);
+ while (g_pause_cdc_child_inspection_after_running.load()) {
+ }
+ g_cdc_child_inspection_paused.store(false);
+ }
+#endif
+ if (state == OwnedChildState::EXITED) {
+ // The pid is reusable after waitpid/ECHILD. A pending handler
request refers to this
+ // now-unpublished generation and can be consumed without another
OS operation.
+ release_terminal_child_identity(identity);
+ break;
+ }
+ if (release_child_identity_if_quiescent(identity)) {
+ break;
+ }
+ // A SIGCHLD arrived after the WNOHANG observation while this claim
was still held. The
+ // pending-bit handoff makes this owner repeat the observation instead
of returning RUNNING.
+ }
+
+ if (status != nullptr) {
+ *status = local_status;
+ }
+ if (wait_error != nullptr) {
+ *wait_error = local_wait_error;
+ }
+ return state;
+}
+
+#ifdef BE_TEST
+// Simulates "old child reaped, numeric pid reused" without asking the kernel
to cycle its PID
+// allocator. This revokes one generation without touching the process, after
which a test can publish
+// the same number under a new generation and challenge cleanup with the stale
identity.
+bool revoke_owned_child_for_test(uint64_t identity) {
+ if (!claim_child_identity(identity)) {
+ return false;
+ }
+ uint64_t expected = identity;
+ const bool revoked =
g_cdc_child_identity.compare_exchange_strong(expected, 0);
+ release_terminal_child_identity(identity);
+ return revoked;
+}
+#endif
+
#ifndef BE_TEST
+std::string child_exit_description(int status) {
+ if (WIFEXITED(status)) {
+ return fmt::format("exit code {}", WEXITSTATUS(status));
+ }
+ if (WIFSIGNALED(status)) {
+ return fmt::format("signal {}", WTERMSIG(status));
+ }
+ return fmt::format("wait status {}", status);
+}
+
+// Check CDC client health
Status check_cdc_client_health(int retry_times, int sleep_time, std::string&
health_response) {
const std::string cdc_health_url =
"http://127.0.0.1:" +
std::to_string(doris::config::cdc_client_port) +
@@ -96,25 +449,105 @@ CdcClientMgr::~CdcClientMgr() {
stop();
}
+uint64_t CdcClientMgr::_publish_child_pid(pid_t pid) {
+ DORIS_CHECK(pid > 0);
+ // A handler may have unpublished its exited generation just before
releasing the operation
+ // claim. Do not publish a replacement into that short gap: the operation
token also carries
+ // the pending-reap bit for the currently published generation.
+ while (g_cdc_child_operation.load() != 0) {
+ std::this_thread::yield();
+ }
+ const uint64_t identity = new_child_identity(pid);
+ uint64_t empty = 0;
+ if (!g_cdc_child_identity.compare_exchange_strong(empty, identity)) {
+ return 0;
+ }
+ return identity;
+}
+
+uint64_t CdcClientMgr::_get_child_identity() const {
+ return g_cdc_child_identity.load();
+}
+
+pid_t CdcClientMgr::_get_child_pid() const {
+ return child_pid(_get_child_identity());
+}
+
+bool CdcClientMgr::_terminate_child_identity(uint64_t identity) {
+ return terminate_owned_child(identity);
+}
+
+#ifdef BE_TEST
+uint64_t CdcClientMgr::set_child_pid_for_test(pid_t pid) {
+ uint64_t current = _get_child_identity();
+ while (current != 0 && !revoke_owned_child_for_test(current)) {
+ current = _get_child_identity();
+ }
+ return _publish_child_pid(pid);
+}
+
+bool CdcClientMgr::terminate_child_identity_for_test(uint64_t identity) {
+ return _terminate_child_identity(identity);
+}
+
+bool CdcClientMgr::inspect_child_identity_for_test(uint64_t identity) {
+ return inspect_owned_child(identity, nullptr, nullptr) ==
OwnedChildState::RUNNING;
+}
+
+void CdcClientMgr::invoke_sigchld_handler_for_test() {
+ handle_sigchld(SIGCHLD);
+}
+
+void CdcClientMgr::pause_sigchld_handler_for_test(bool pause) {
+ g_pause_cdc_sigchld_handler.store(pause);
+}
+
+bool CdcClientMgr::sigchld_handler_paused_for_test() {
+ return g_cdc_sigchld_handler_paused.load();
+}
+
+void CdcClientMgr::reset_child_claim_failed_for_test() {
+ g_cdc_child_claim_failed.store(false);
+}
+
+bool CdcClientMgr::child_claim_failed_for_test() {
+ return g_cdc_child_claim_failed.load();
+}
+
+void CdcClientMgr::pause_child_inspection_after_running_for_test(bool pause) {
+ g_pause_cdc_child_inspection_after_running.store(pause);
+}
+
+bool CdcClientMgr::child_inspection_paused_for_test() {
+ return g_cdc_child_inspection_paused.load();
+}
+
+void CdcClientMgr::pause_sigchld_before_claim_for_test(bool pause) {
+ g_pause_cdc_sigchld_before_claim.store(pause);
+}
+
+bool CdcClientMgr::sigchld_before_claim_paused_for_test() {
+ return g_cdc_sigchld_before_claim_paused.load();
+}
+
+void CdcClientMgr::pause_stale_child_claim_for_test(bool pause) {
+ g_pause_cdc_stale_child_claim.store(pause);
+}
+
+bool CdcClientMgr::stale_child_claim_paused_for_test() {
+ return g_cdc_stale_child_claim_paused.load();
+}
+#endif
+
void CdcClientMgr::stop() {
- pid_t pid = _child_pid.load();
- if (pid > 0) {
- // Check if process is still alive
- if (kill(pid, 0) == 0) {
- LOG(INFO) << "Stopping CDC client process, pid=" << pid;
- // Send SIGTERM for graceful shutdown
- kill(pid, SIGTERM);
- // Wait a short time for graceful shutdown
- std::this_thread::sleep_for(std::chrono::milliseconds(200));
- // Force kill if still alive
- if (kill(pid, 0) == 0) {
- LOG(INFO) << "Force killing CDC client process, pid=" << pid;
- kill(pid, SIGKILL);
- int status = 0;
- waitpid(pid, &status, 0);
- }
+ std::lock_guard<std::mutex> lock(_start_mutex);
+ // Claim the exact generation before touching the OS process. If the
handler is operating it, wait;
+ // if the handler already reaped it, reload rather than retaining its
now-reusable numeric pid.
+ while (true) {
+ const uint64_t identity = _get_child_identity();
+ if (identity == 0 || _terminate_child_identity(identity)) {
+ break;
}
- _child_pid.store(0);
}
LOG(INFO) << "CdcClientMgr is stopped";
@@ -124,15 +557,18 @@ Status
CdcClientMgr::start_cdc_client(PRequestCdcClientResult* result) {
std::lock_guard<std::mutex> lock(_start_mutex);
Status st = Status::OK();
- pid_t exist_pid = _child_pid.load();
+ const uint64_t existing_identity = _get_child_identity();
+ const pid_t exist_pid = child_pid(existing_identity);
if (exist_pid > 0) {
#ifdef BE_TEST
// In test mode, directly return OK if PID exists
LOG(INFO) << "cdc client already started (BE_TEST mode), pid=" <<
exist_pid;
return Status::OK();
#else
- // Check if process is still alive
- if (kill(exist_pid, 0) == 0) {
+ int existing_wait_error = 0;
+ const OwnedChildState existing_state =
+ inspect_owned_child(existing_identity, nullptr,
&existing_wait_error);
+ if (existing_state == OwnedChildState::RUNNING) {
// Process exists, verify it's actually our CDC client by health
check
std::string check_response;
auto check_st = check_cdc_client_health(3, 1, check_response);
@@ -146,10 +582,13 @@ Status
CdcClientMgr::start_cdc_client(PRequestCdcClientResult* result) {
st.to_protobuf(result->mutable_status());
return st;
}
+ } else if (existing_state == OwnedChildState::WAIT_ERROR) {
+ st = Status::InternalError(fmt::format("Could not inspect CDC
client {}: {}", exist_pid,
+
strerror(existing_wait_error)));
+ st.to_protobuf(result->mutable_status());
+ return st;
} else {
LOG(INFO) << "CDC client is dead, pid=" << exist_pid;
- // Process is dead, reset PID and continue to start
- _child_pid.store(0);
}
#endif
} else if (!_adopted_external.load()) {
@@ -234,13 +673,26 @@ Status
CdcClientMgr::start_cdc_client(PRequestCdcClientResult* result) {
const std::string cdc_out_file = std::string(log_dir) + "/cdc-client.out";
- struct sigaction act;
- act.sa_flags = 0;
+ struct sigaction act {};
+ sigemptyset(&act.sa_mask);
+ // SA_RESTART: the handler runs on whichever thread the kernel picks, and
without it every
+ // blocking call in BE becomes interruptible whenever the cdc client
exits. SA_NOCLDSTOP: only
+ // the child's exit is interesting, not its stops.
+ act.sa_flags = SA_RESTART | SA_NOCLDSTOP;
act.sa_handler = handle_sigchld;
- sigaction(SIGCHLD, &act, NULL);
+ // Everything below depends on this disposition being installed: nothing
else consumes the
+ // published identity, and without the handler the pending-reap handoff is
dead.
+ DORIS_CHECK(sigaction(SIGCHLD, &act, nullptr) == 0);
LOG(INFO) << "Start to fork cdc client process with " << path;
#ifdef BE_TEST
- _child_pid.store(99999);
+ // Unit tests can construct several managers even though ExecEnv owns only
one in production.
+ // A concurrent test manager may win publication after our initial empty
check; in that case all
+ // managers observe the same process-wide test child and start is still
successful.
+ if (_publish_child_pid(99999) == 0 && _get_child_identity() == 0) {
+ st = Status::InternalError("Failed to publish test CDC child
identity");
+ st.to_protobuf(result->mutable_status());
+ return st;
+ }
st = Status::OK();
return st;
#else
@@ -266,30 +718,75 @@ Status
CdcClientMgr::start_cdc_client(PRequestCdcClientResult* result) {
perror("Cdc client child process error");
_exit(1);
} else {
- // Parent process: save PID and wait for startup
- _child_pid.store(pid);
+ // Parent process: publish a generation-qualified identity. The child
is not visible to the
+ // SIGCHLD handler before this succeeds, so a publication conflict
still leaves this thread as
+ // the only possible reaper of the just-forked pid.
+ const uint64_t forked_identity = _publish_child_pid(pid);
+ if (forked_identity == 0) {
+ terminate_and_reap_child(pid);
+ st = Status::InternalError("Another CDC child identity was
published during startup");
+ st.to_protobuf(result->mutable_status());
+ return st;
+ }
+ // A child that died between fork() returning and the store above
raised a SIGCHLD the
+ // handler saw with no identity to reap. Inspect the exact generation
here so the child is
+ // reaped before its pid can be reused, but decide whether that is
fatal only after the
+ // health check below: the port may be served by an external instance
again by then, and
+ // adoption is the outcome that used to be reached for a child that
died this early.
+ int forked_status = 0;
+ int forked_wait_error = 0;
+ const OwnedChildState initial_state =
+ inspect_owned_child(forked_identity, &forked_status,
&forked_wait_error);
+ if (initial_state == OwnedChildState::WAIT_ERROR) {
+ _terminate_child_identity(forked_identity);
+ st = Status::InternalError(
+ fmt::format("CDC client exited before startup or could not
be waited for: {}",
+ strerror(forked_wait_error)));
+ st.to_protobuf(result->mutable_status());
+ return st;
+ }
// Waiting for cdc to start, failed after more than 3 * 10 seconds
std::string health_response;
Status status = check_cdc_client_health(3, 10, health_response);
if (!status.ok()) {
- // Reset PID if startup failed
- _child_pid.store(0);
- st = Status::InternalError("Start cdc client failed.");
+ // Cleanup is conditional on the exact generation still being
ours. If the handler already
+ // reaped it, a same-parent child may now reuse the number and
must not be touched.
+ _terminate_child_identity(forked_identity);
+ if (initial_state == OwnedChildState::EXITED) {
+ st = Status::InternalError(fmt::format("CDC client exited
before startup with {}",
+
child_exit_description(forked_status)));
+ } else if (initial_state == OwnedChildState::NOT_OWNED) {
+ st = Status::InternalError("CDC client exited before startup");
+ } else {
+ st = Status::InternalError("Start cdc client failed.");
+ }
st.to_protobuf(result->mutable_status());
- } else if (kill(pid, 0) != 0) {
+ } else {
+ int final_wait_error = 0;
+ const OwnedChildState final_state =
+ inspect_owned_child(forked_identity, nullptr,
&final_wait_error);
+ if (final_state == OwnedChildState::WAIT_ERROR) {
+ _terminate_child_identity(forked_identity);
+ st = Status::InternalError(
+ fmt::format("Could not inspect started CDC client {}:
{}", pid,
+ strerror(final_wait_error)));
+ st.to_protobuf(result->mutable_status());
+ return st;
+ }
+ if (final_state == OwnedChildState::RUNNING) {
+ _adopted_external.store(false);
+ LOG(INFO) << "Start cdc client success, pid=" << pid
+ << ", status=" << status.to_string() << ",
response=" << health_response;
+ return st;
+ }
// Port healthy but our child has exited: an external process is
// answering. Treat as adoption instead of masking dead PID as
success.
- _child_pid.store(0);
if (!_adopted_external.exchange(true)) {
LOG(INFO) << "Forked cdc client " << pid << " exited but port "
<< doris::config::cdc_client_port
<< " is healthy, adopting external instance";
}
- } else {
- _adopted_external.store(false);
- LOG(INFO) << "Start cdc client success, pid=" << pid
- << ", status=" << status.to_string() << ", response=" <<
health_response;
}
}
#endif //BE_TEST
diff --git a/be/src/runtime/cdc_client_mgr.h b/be/src/runtime/cdc_client_mgr.h
index b3b350154b6..d054141993a 100644
--- a/be/src/runtime/cdc_client_mgr.h
+++ b/be/src/runtime/cdc_client_mgr.h
@@ -20,6 +20,7 @@
#include <gen_cpp/internal_service.pb.h>
#include <atomic>
+#include <cstdint>
#include <mutex>
#include <string>
@@ -50,17 +51,42 @@ public:
#ifdef BE_TEST
// For testing only: get current child PID
- pid_t get_child_pid() const { return _child_pid.load(); }
- // For testing only: set child PID directly
- void set_child_pid_for_test(pid_t pid) { _child_pid.store(pid); }
+ pid_t get_child_pid() const { return _get_child_pid(); }
+ // For testing only: publish a PID and return its generation-qualified
identity.
+ uint64_t set_child_pid_for_test(pid_t pid);
+ uint64_t get_child_identity_for_test() const { return
_get_child_identity(); }
+ // For testing only: run the production cleanup gate for an exact identity.
+ bool terminate_child_identity_for_test(uint64_t identity);
+ // For testing only: inspect the exact identity through the normal WNOHANG
path.
+ bool inspect_child_identity_for_test(uint64_t identity);
+ // For testing only: invoke the installed handler body deterministically.
+ static void invoke_sigchld_handler_for_test();
+ // For testing only: pause a deterministic handler after it has copied the
published identity.
+ static void pause_sigchld_handler_for_test(bool pause);
+ static bool sigchld_handler_paused_for_test();
+ // For testing only: observe a normal thread failing to claim an owned
child.
+ static void reset_child_claim_failed_for_test();
+ static bool child_claim_failed_for_test();
+ // For testing only: pause after inspect's WNOHANG=0 observation and
before claim release.
+ static void pause_child_inspection_after_running_for_test(bool pause);
+ static bool child_inspection_paused_for_test();
+ // Pause a handler before claiming a copied identity, and after claiming a
stale identity.
+ static void pause_sigchld_before_claim_for_test(bool pause);
+ static bool sigchld_before_claim_paused_for_test();
+ static void pause_stale_child_claim_for_test(bool pause);
+ static bool stale_child_claim_paused_for_test();
// For testing only: inspect / drive the adopt-external flag
bool get_adopted_external_for_test() const { return
_adopted_external.load(); }
void set_adopted_external_for_test(bool v) { _adopted_external.store(v); }
#endif
private:
+ uint64_t _get_child_identity() const;
+ pid_t _get_child_pid() const;
+ uint64_t _publish_child_pid(pid_t pid);
+ bool _terminate_child_identity(uint64_t identity);
+
std::mutex _start_mutex;
- std::atomic<pid_t> _child_pid {0};
std::atomic<bool> _adopted_external {false};
};
diff --git a/be/test/runtime/cdc_client_mgr_test.cpp
b/be/test/runtime/cdc_client_mgr_test.cpp
index c68e71956ad..fd35fc61c22 100644
--- a/be/test/runtime/cdc_client_mgr_test.cpp
+++ b/be/test/runtime/cdc_client_mgr_test.cpp
@@ -20,10 +20,13 @@
#include <gen_cpp/internal_service.pb.h>
#include <gtest/gtest.h>
#include <signal.h>
+#include <spawn.h>
#include <sys/stat.h>
#include <sys/wait.h>
#include <unistd.h>
+#include <atomic>
+#include <cerrno>
#include <chrono>
#include <cstdio>
#include <cstdlib>
@@ -34,9 +37,30 @@
#include "common/status.h"
#include "runtime/cluster_info.h"
#include "runtime/exec_env.h"
+#include "util/defer_op.h"
namespace doris {
+namespace {
+
+void reap_direct_child_if_present(pid_t pid) {
+ siginfo_t child_info {};
+ if (waitid(P_PID, pid, &child_info, WEXITED | WNOHANG | WNOWAIT) != 0) {
+ EXPECT_EQ(errno, ECHILD);
+ return;
+ }
+ if (child_info.si_pid == 0) {
+ EXPECT_EQ(kill(pid, SIGKILL), 0);
+ }
+ pid_t wait_result;
+ do {
+ wait_result = waitpid(pid, nullptr, 0);
+ } while (wait_result < 0 && errno == EINTR);
+ EXPECT_EQ(wait_result, pid);
+}
+
+} // namespace
+
class CdcClientMgrTest : public testing::Test {
public:
void SetUp() override {
@@ -123,7 +147,7 @@ TEST_F(CdcClientMgrTest, StopWithoutChild) {
mgr.stop();
}
-// Test stop when child process is already dead (covers lines 98-111:
kill(pid, 0) == 0 is false)
+// Test stop when the published process is already absent.
TEST_F(CdcClientMgrTest, StopWhenProcessDead) {
CdcClientMgr mgr;
@@ -133,85 +157,190 @@ TEST_F(CdcClientMgrTest, StopWhenProcessDead) {
EXPECT_TRUE(status.ok());
EXPECT_GT(mgr.get_child_pid(), 0);
- // Stop - since PID 99999 doesn't exist, kill(99999, 0) will fail
- // This should trigger the branch where kill(pid, 0) != 0 (process already
dead)
+ // The generation-qualified cleanup observes ECHILD and revokes ownership
without signalling.
mgr.stop();
// PID should be reset to 0
EXPECT_EQ(mgr.get_child_pid(), 0);
}
-// Test stop with real process that exits gracefully (covers lines 98-111:
graceful shutdown)
-TEST_F(CdcClientMgrTest, StopWithRealProcessGraceful) {
+// Test stop when the published pid cannot be reaped as this process's own
child. The
+// generation-qualified cleanup observes ECHILD there, which counts as
success, so stop() must
+// revoke ownership without signalling: a pid this process cannot reap is not
one it may operate on.
+// The signal and pipe checkpoints must stay in this process-lifecycle test.
+// NOLINTNEXTLINE(readability-function-cognitive-complexity)
+TEST_F(CdcClientMgrTest, StopDoesNotSignalANonChildPid) {
CdcClientMgr mgr;
- // Use popen to start a background sleep process and get its PID
- // This avoids fork() which conflicts with gcov/coverage tools
- FILE* pipe = popen("sleep 10 & echo $!", "r");
- if (pipe) {
- char buffer[128];
- if (fgets(buffer, sizeof(buffer), pipe) != nullptr) {
- pid_t real_pid = std::atoi(buffer);
- pclose(pipe);
-
- if (real_pid > 0) {
- // 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
- mgr.stop();
-
- // Verify PID is reset
- EXPECT_EQ(mgr.get_child_pid(), 0);
-
- // Clean up: make sure child is dead
- kill(real_pid, SIGKILL);
- waitpid(real_pid, nullptr, WNOHANG);
- }
- } else {
- pclose(pipe);
+ int pid_pipe[2];
+ ASSERT_EQ(pipe(pid_pipe), 0);
+ Defer close_pid_pipe {[&]() {
+ close(pid_pipe[0]);
+ if (pid_pipe[1] >= 0) {
+ close(pid_pipe[1]);
+ }
+ }};
+ int report_pipe[2];
+ ASSERT_EQ(pipe(report_pipe), 0);
+ Defer close_report_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_addclose(&actions, pid_pipe[0]), 0);
+ ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, report_pipe[0]), 0);
+ ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, phase_pipe[1]), 0);
+ ASSERT_EQ(posix_spawn_file_actions_adddup2(&actions, pid_pipe[1],
STDOUT_FILENO), 0);
+ ASSERT_EQ(posix_spawn_file_actions_adddup2(&actions, report_pipe[1],
STDERR_FILENO), 0);
+ ASSERT_EQ(posix_spawn_file_actions_adddup2(&actions, phase_pipe[0],
STDIN_FILENO), 0);
+ ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, pid_pipe[1]), 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);
+
+ 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 launcher reports its background helper's PID and exits. The helper
keeps the report and
+ // phase pipes open but is no longer this process's child when the
launcher has been reaped.
+ pid_t launcher = 0;
+ char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+ const_cast<char*>("sh -c 'trap \"printf S >&2\"
TERM; printf R >&2; "
+ "while IFS= read -r phase; do
printf A >&2; done' "
+ "<&0 >/dev/null & echo $!"),
+ nullptr};
+ char* const envp[] = {const_cast<char*>("PATH=/bin:/usr/bin"), nullptr};
+ ASSERT_EQ(posix_spawn(&launcher, "/bin/sh", &actions, &attr, argv, envp),
0);
+ ASSERT_GT(launcher, 0);
+ bool launcher_reaped = false;
+ Defer cleanup_launcher {[&]() {
+ if (!launcher_reaped) {
+ reap_direct_child_if_present(launcher);
+ }
+ }};
+ close(pid_pipe[1]);
+ pid_pipe[1] = -1;
+ close(report_pipe[1]);
+ report_pipe[1] = -1;
+ close(phase_pipe[0]);
+ phase_pipe[0] = -1;
+
+ char pid_text[32] {};
+ size_t pid_length = 0;
+ char digit = 0;
+ for (; pid_length < sizeof(pid_text) - 1; ++pid_length) {
+ ASSERT_EQ(read(pid_pipe[0], &digit, 1), 1);
+ if (digit == '\n') {
+ break;
+ }
+ pid_text[pid_length] = digit;
}
+ ASSERT_EQ(digit, '\n');
+ const pid_t helper_pid = std::atoi(pid_text);
+ ASSERT_GT(helper_pid, 0);
+ int launcher_status = 0;
+ ASSERT_EQ(waitpid(launcher, &launcher_status, 0), launcher);
+ launcher_reaped = true;
+ ASSERT_TRUE(WIFEXITED(launcher_status));
+ ASSERT_EQ(WEXITSTATUS(launcher_status), 0);
+ errno = 0;
+ ASSERT_EQ(waitpid(helper_pid, nullptr, WNOHANG), -1);
+ ASSERT_EQ(errno, ECHILD);
+
+ char ready = 0;
+ ASSERT_EQ(read(report_pipe[0], &ready, 1), 1);
+ ASSERT_EQ(ready, 'R');
+ mgr.set_child_pid_for_test(helper_pid);
+ mgr.stop();
+ EXPECT_EQ(mgr.get_child_pid(), 0);
+
+ 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') << "stop() sent SIGTERM to a process it could
not reap";
}
-// Test stop with real process that requires force kill (covers lines 98-111:
force kill path)
+// Test stop with a direct child that requires force kill.
TEST_F(CdcClientMgrTest, StopWithRealProcessForceKill) {
CdcClientMgr mgr;
- // Start a bash process that ignores SIGTERM by trapping it
- // This process will not exit on SIGTERM, requiring SIGKILL
- const char* script = "bash -c 'trap \"\" TERM; while true; do sleep 1;
done' & echo $!";
- FILE* pipe = popen(script, "r");
- if (pipe) {
- char buffer[128];
- if (fgets(buffer, sizeof(buffer), pipe) != nullptr) {
- pid_t real_pid = std::atoi(buffer);
- pclose(pipe);
-
- if (real_pid > 0) {
- // Give the process a moment to start
- std::this_thread::sleep_for(std::chrono::milliseconds(100));
-
- // Set the PID
- mgr.set_child_pid_for_test(real_pid);
-
- // Call stop - should try graceful shutdown first, then force
kill
- // Since process ignores SIGTERM, it will still be alive after
200ms
- // This should trigger the force kill path (lines 105-110)
- mgr.stop();
-
- // Verify PID is reset
- EXPECT_EQ(mgr.get_child_pid(), 0);
-
- // Clean up: make sure child is dead
- kill(real_pid, SIGKILL);
- waitpid(real_pid, nullptr, WNOHANG);
- }
- } else {
- pclose(pipe);
+ int ready_pipe[2];
+ ASSERT_EQ(pipe(ready_pipe), 0);
+ Defer close_pipe {[&]() {
+ close(ready_pipe[0]);
+ if (ready_pipe[1] >= 0) {
+ close(ready_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, ready_pipe[1],
STDOUT_FILENO), 0);
+ ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, ready_pipe[0]), 0);
+ ASSERT_EQ(posix_spawn_file_actions_addclose(&actions, ready_pipe[1]), 0);
+
+ // The shell reports readiness only after ignoring SIGTERM. exec preserves
ignored signals,
+ // leaving sleep as this test process's direct child with the same PID.
+ pid_t pid = 0;
+ char* const argv[] = {const_cast<char*>("sh"), const_cast<char*>("-c"),
+ const_cast<char*>("trap '' TERM; printf R; exec
sleep 3600"), 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);
+ bool child_reaped = false;
+ Defer cleanup_child {[&]() {
+ if (!child_reaped) {
+ kill(pid, SIGKILL);
+ waitpid(pid, nullptr, 0);
+ }
+ }};
+ close(ready_pipe[1]);
+ ready_pipe[1] = -1;
+ char ready = 0;
+ ASSERT_EQ(read(ready_pipe[0], &ready, 1), 1);
+ ASSERT_EQ(ready, 'R');
+
+ mgr.set_child_pid_for_test(pid);
+ mgr.stop();
+ EXPECT_EQ(mgr.get_child_pid(), 0);
+
+ errno = 0;
+ const pid_t wait_result = waitpid(pid, nullptr, WNOHANG);
+ const int wait_error = errno;
+ child_reaped = wait_result == pid || (wait_result < 0 && wait_error ==
ECHILD);
+ EXPECT_EQ(wait_result, -1) << "stop() did not collect the forced-kill
child";
+ EXPECT_EQ(wait_error, ECHILD);
}
// Test start_cdc_client with missing jar file
@@ -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();
+ EXPECT_EQ(mgr.get_child_identity_for_test(), successor_identity);
+ CdcClientMgr::pause_stale_child_claim_for_test(false);
+ delayed_handler.join();
+
+ 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_reaped = wait_result < 0 && wait_error == ECHILD;
+ EXPECT_EQ(wait_result, -1);
+ EXPECT_EQ(wait_error, ECHILD);
+}
+
+TEST_F(CdcClientMgrTest, ConcurrentHandlersHaveOneExclusiveProcessOperator) {
+ 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);
+ Defer resume_handler {[]() {
CdcClientMgr::pause_sigchld_handler_for_test(false); }};
+ std::thread first([]() { 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);
+ first.join();
+ kill(pid, SIGKILL);
+ waitpid(pid, nullptr, 0);
+ FAIL() << "the first handler did not acquire and pause its process
operation";
+ }
+
+ // The second handler must return instead of blocking on the claim the
first one holds, so
+ // join() returning is the proof; a flag written before the thread returns
would assert itself.
+ std::thread second([]() { CdcClientMgr::invoke_sigchld_handler_for_test();
});
+ second.join();
+ EXPECT_EQ(kill(pid, 0), 0);
+
+ CdcClientMgr::pause_sigchld_handler_for_test(false);
+ first.join();
+ mgr.stop();
+}
+
+TEST_F(CdcClientMgrTest, SigchldDuringRunningInspectionIsHandedBackForReap) {
+ 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); }};
+
+ // Keep delivery deterministic: the test invokes the production handler
only after waitid proves
+ // the child is waitable, while the inspecting thread still holds the
operation claim.
+ 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;
+ const uint64_t identity = mgr.set_child_pid_for_test(pid);
+ CdcClientMgr::pause_child_inspection_after_running_for_test(true);
+ std::atomic<bool> inspector_saw_running {true};
+ std::thread inspector(
+ [&]() {
inspector_saw_running.store(mgr.inspect_child_identity_for_test(identity)); });
+ Defer resume_and_join_inspector {[&]() {
+ CdcClientMgr::pause_child_inspection_after_running_for_test(false);
+ if (inspector.joinable()) {
+ inspector.join();
+ }
+ }};
+
+ for (int i = 0; i < 100 &&
!CdcClientMgr::child_inspection_paused_for_test(); ++i) {
+ std::this_thread::sleep_for(std::chrono::milliseconds(1));
+ }
+ if (!CdcClientMgr::child_inspection_paused_for_test()) {
+ FAIL() << "the inspector did not pause after observing WNOHANG=0";
+ }
+
+ ASSERT_EQ(kill(pid, SIGKILL), 0);
+ 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(1));
+ }
+ ASSERT_EQ(child_info.si_pid, pid);
+
+ // The handler cannot claim while inspect owns it. It must attach a
pending request rather than
+ // consume the only notification and return. No later signal and no
explicit stop() drive cleanup.
+ CdcClientMgr::invoke_sigchld_handler_for_test();
+ EXPECT_EQ(mgr.get_child_identity_for_test(), identity);
+ CdcClientMgr::pause_child_inspection_after_running_for_test(false);
+ inspector.join();
+
+ EXPECT_FALSE(inspector_saw_running.load());
+ EXPECT_EQ(mgr.get_child_identity_for_test(), 0);
+ errno = 0;
+ EXPECT_EQ(waitpid(pid, nullptr, WNOHANG), -1);
+ EXPECT_EQ(errno, ECHILD) << "the pending-reap handoff did not collect the
exited child";
+}
+
// Test start_cdc_client when environment is missing
TEST_F(CdcClientMgrTest, StartCdcClientMissingEnv) {
unsetenv("JAVA_HOME");
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]