From feb3d45ca5b40eabce3309e8f75fa2b2879a7746 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gustavo=20Andr=C3=A9=20dos=20Santos=20Lopes?= Date: Mon, 21 Sep 2026 13:16:29 +0100 Subject: [PATCH 1/3] Make forced signal flush raw-clone safe Prepare an independent sidecar connection and serialized flush request before signal delivery. The SIGINT and SIGTERM handler starts a one-shot raw clone that uses audited direct syscalls while waiting for an ACK or timeout, avoiding libc and allocator state in the unrepaired clone. Track worker lifetime with lock-free the kernel clear_child_tid futex protocol. --- components-rs/datadog.h | 36 ++ components-rs/lib.rs | 2 + components-rs/signal_flush.rs | 707 ++++++++++++++++++++++++++++++++++ ext/sidecar.c | 24 ++ ext/signals.c | 270 ++++++++++--- ext/signals.h | 7 + ext/threads.c | 100 ++++- ext/threads.h | 19 +- 8 files changed, 1100 insertions(+), 65 deletions(-) create mode 100644 components-rs/signal_flush.rs diff --git a/components-rs/datadog.h b/components-rs/datadog.h index 4aefb5d8113..aa677ce8674 100644 --- a/components-rs/datadog.h +++ b/components-rs/datadog.h @@ -9,6 +9,14 @@ struct _zend_string; #include "telemetry.h" #include "sidecar.h" +#if defined(__linux__) +/** + * Prepared, independent, sessionless connection. Opaque to C; immutable after publication. + * Only normal initialized threads may construct or destroy this object. + */ +typedef struct ddog_SignalFlush ddog_SignalFlush; +#endif + extern void (*ddog_log_callback)(ddog_CharSlice); extern ddog_VecRemoteConfigProduct DATADOG_REMOTE_CONFIG_PRODUCTS; @@ -273,6 +281,34 @@ void datadog_sidecar_set_reconnect_fn(struct ddog_SidecarTransport **transport, void datadog_sidecar_clear_reconnect_fn(struct ddog_SidecarTransport **transport); +#if defined(__linux__) +/** + * Connect independently to the template's exact listener and prepare one Flush request. + * The template is borrowed only during this call: neither it nor its fd is retained. + * The returned object must outlive the raw worker and may only be dropped normally. + */ +ddog_MaybeError datadog_sidecar_prepare_signal_flush(struct ddog_SidecarTransport *template_, + struct ddog_SignalFlush **output); +#endif + +#if defined(__linux__) +/** + * Destroy an unpublished object, or one whose raw worker has exited. Normal context only. + */ +void datadog_sidecar_signal_flush_drop(struct ddog_SignalFlush *flush); +#endif + +#if defined(__linux__) +/** + * Execute one bounded exchange without calling libc, using TLS, allocating, or unwinding. + * If `terminate_process` is true, terminate the process after the flush completes, fails, or + * times out. Otherwise, return zero for the one-byte zero ACK or a negative Linux error number. + * The caller must keep the object alive and guarantee exclusive, one-shot use of its socket. + */ +int32_t datadog_sidecar_signal_flush_run(const struct ddog_SignalFlush *flush, + bool terminate_process); +#endif + bool ddog_shm_limiter_inc(const struct ddog_MaybeShmLimiter *limiter, uint32_t limit); bool ddog_exception_hash_limiter_inc(struct ddog_SidecarTransport *connection, diff --git a/components-rs/lib.rs b/components-rs/lib.rs index b7f5a035b18..16c8a10d2c1 100644 --- a/components-rs/lib.rs +++ b/components-rs/lib.rs @@ -45,6 +45,8 @@ pub mod log; pub mod remote_config; #[cfg(not(standalone_profiler))] pub mod sidecar; +#[cfg(all(not(standalone_profiler), target_os = "linux"))] +pub mod signal_flush; #[cfg(not(standalone_profiler))] pub mod stats; #[cfg(not(standalone_profiler))] diff --git a/components-rs/signal_flush.rs b/components-rs/signal_flush.rs new file mode 100644 index 00000000000..77a32ee45d3 --- /dev/null +++ b/components-rs/signal_flush.rs @@ -0,0 +1,707 @@ +//! The raw-clone signal worker owns no libc or Rust thread state. `prepare` and `drop` run on +//! ordinary initialized threads; only `run` may execute in the worker. `run` uses ordinary Rust +//! control flow, but every operation in its generated call graph must reduce to local memory +//! access or a direct Linux system call. Changes to it therefore require a disassembly audit in +//! both optimized and debug builds. + +use core::arch::asm; +use datadog_sidecar::service::blocking::SidecarTransport; +use datadog_sidecar::service::sidecar_interface::{SidecarFlushOptions, SidecarInterfaceRequest}; +use libdd_common_ffi::{Error, MaybeError}; +use libdd_ipc::SeqpacketConn; +use std::mem::{offset_of, size_of}; + +/// Prepared, independent, sessionless connection. Opaque to C; immutable after publication. +/// Only normal initialized threads may construct or destroy this object. +pub struct SignalFlush { + // Keep the raw view free of owning Rust values: the worker only copies these four fields. + raw: RawSignalFlush, + // These owners make the fd and request pointer in `raw` valid until normal-context drop. + _connection: SeqpacketConn, + _request: Box<[u8]>, +} + +#[repr(C)] +struct RawSignalFlush { + fd: i32, + owner_pid: i32, + request: *const u8, + request_len: usize, +} + +// Direct syscalls use the 64-bit Linux kernel ABI, not libc's platform types. +#[repr(C)] +struct KernelTimespec { + seconds: i64, + nanoseconds: i64, +} + +// Match the pollfd layout consumed directly by the kernel. +#[repr(C)] +struct KernelPollFd { + fd: i32, + events: i16, + revents: i16, +} + +/// Connect independently to the template's exact listener and prepare one Flush request. +/// The template is borrowed only during this call: neither it nor its fd is retained. +/// The returned object must outlive the raw worker and may only be dropped normally. +#[no_mangle] +pub unsafe extern "C" fn datadog_sidecar_prepare_signal_flush( + template: *mut SidecarTransport, + output: *mut *mut SignalFlush, +) -> MaybeError { + // Everything that may allocate, lock, inspect libc state, or use the sidecar codec belongs in + // this function. None of those operations can be deferred to the raw worker. + let result = (|| { + let output = output + .as_mut() + .ok_or_else(|| anyhow::anyhow!("signal flush output is null"))?; + // Never leave the caller holding an old value when a later preparation step fails. + *output = std::ptr::null_mut(); + let template = template + .as_mut() + .ok_or_else(|| anyhow::anyhow!("signal flush template is null"))?; + // This is a new socket, not a dup of the template fd. A dup would keep the template's + // connection alive and interfere with the sidecar's disconnect detection. + let connection = connect_to_template(template)?; + // Encoding now leaves `run` with one immutable packet and no serializer call graph. + let request = libdd_ipc::codec::encode(&SidecarInterfaceRequest::Flush { + // Preserve the previous forced-shutdown behavior: traces and stats are the only data + // whose delivery this signal path delays termination to await. + options: SidecarFlushOptions { + traces_and_stats: true, + flag_evaluations: false, + telemetry: false, + }, + }) + .into_boxed_slice(); + + *output = Box::into_raw(Box::new(SignalFlush { + raw: RawSignalFlush { + fd: connection.as_raw_fd(), + // A copied object must never send on the parent's connection after fork. + owner_pid: libc::getpid(), + request: request.as_ptr(), + request_len: request.len(), + }, + _connection: connection, + _request: request, + })); + Ok::<_, anyhow::Error>(()) + })(); + match result { + Ok(()) => MaybeError::None, + Err(error) => MaybeError::Some(Error::from(format!("{error:#}"))), + } +} + +/// Destroy an unpublished object, or one whose raw worker has exited. Normal context only. +#[no_mangle] +pub unsafe extern "C" fn datadog_sidecar_signal_flush_drop(flush: *mut SignalFlush) { + // An unpublished object has no worker borrower. For a published object, C joins the worker + // through clear_child_tid before calling this, so neither owner can be destroyed during `run`. + if !flush.is_null() { + drop(Box::from_raw(flush)); + } +} + +/// Execute one bounded exchange without calling libc, using TLS, allocating, or unwinding. +/// If `terminate_process` is true, terminate the process after the flush completes, fails, or +/// times out. Otherwise, return zero for the one-byte zero ACK or a negative Linux error number. +/// The caller must keep the object alive and guarantee exclusive, one-shot use of its socket. +#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))] +#[no_mangle] +#[inline(never)] +// An `extern "C"` definition gets a `panic_cannot_unwind` path in panic=unwind debug builds. +// `C-unwind` avoids that generated call. The body is deliberately written so it cannot unwind. +pub unsafe extern "C-unwind" fn datadog_sidecar_signal_flush_run( + flush: *const SignalFlush, + terminate_process: bool, +) -> i32 { + let result = signal_flush_exchange(flush); + if terminate_process { + // The raw clone has no usable libc thread state. Terminate through an inline syscall + // rather than returning to C or calling libc's `_exit`/`exit` machinery. + signal_safe_exit_group(); + } + result +} + +#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))] +#[inline(never)] +unsafe fn signal_flush_exchange(flush: *const SignalFlush) -> i32 { + // Raw Linux syscalls return `-errno` directly. They do not set libc's thread-local `errno`. + const EINTR: isize = 4; + const EAGAIN: isize = 11; + const CLOCK_MONOTONIC: usize = 1; + const MSG_DONTWAIT: usize = 0x40; + const MSG_NOSIGNAL: usize = 0x4000; + const MSG_TRUNC: usize = 0x20; + const POLLIN: i16 = 1; + const POLLOUT: i16 = 4; + + // Syscall numbers are architecture-specific, so keep them beside the audited inline syscall + // rather than routing through libc wrappers. + #[cfg(target_arch = "aarch64")] + const SYS_GETPID: usize = 172; + #[cfg(target_arch = "aarch64")] + const SYS_SENDTO: usize = 206; + #[cfg(target_arch = "aarch64")] + const SYS_RECVFROM: usize = 207; + #[cfg(target_arch = "aarch64")] + const SYS_CLOCK_GETTIME: usize = 113; + #[cfg(target_arch = "aarch64")] + const SYS_PPOLL: usize = 73; + + #[cfg(target_arch = "x86_64")] + const SYS_GETPID: usize = 39; + #[cfg(target_arch = "x86_64")] + const SYS_SENDTO: usize = 44; + #[cfg(target_arch = "x86_64")] + const SYS_RECVFROM: usize = 45; + #[cfg(target_arch = "x86_64")] + const SYS_CLOCK_GETTIME: usize = 228; + #[cfg(target_arch = "x86_64")] + const SYS_PPOLL: usize = 271; + + let (fd, owner_pid, request, request_len) = signal_safe_load(flush); + // `fork()` copies the object and fd but changes the process identity. Reject that copy even if + // a caller reaches this API before the C-side post-fork reset. + if signal_safe_syscall6(SYS_GETPID, 0, 0, 0, 0, 0, 0) != owner_pid as isize { + return -libc::ECHILD; + } + + let mut now = KernelTimespec { + seconds: 0, + nanoseconds: 0, + }; + let result = signal_safe_syscall6( + SYS_CLOCK_GETTIME, + CLOCK_MONOTONIC, + (&raw mut now).cast::() as usize, + 0, + 0, + 0, + 0, + ); + if result < 0 { + return result as i32; + } + // Use one absolute deadline for send, backpressure, and ACK. Restarting a relative timeout + // after EINTR or EAGAIN could otherwise keep process termination alive indefinitely. + // Wrapping arithmetic is intentional: debug overflow checks would introduce panic paths. + let deadline_seconds = now.seconds.wrapping_add(10); + let deadline_nanoseconds = now.nanoseconds; + // `receive` selects the protocol phase. `poll` records that the last nonblocking operation + // returned EAGAIN and must wait for readiness before retrying. + let mut receive = false; + let mut poll = false; + + loop { + let result = signal_safe_syscall6( + SYS_CLOCK_GETTIME, + CLOCK_MONOTONIC, + (&raw mut now).cast::() as usize, + 0, + 0, + 0, + 0, + ); + if result < 0 { + return result as i32; + } + let mut remaining_seconds = deadline_seconds.wrapping_sub(now.seconds); + let mut remaining_nanoseconds = deadline_nanoseconds.wrapping_sub(now.nanoseconds); + if remaining_nanoseconds < 0 { + remaining_nanoseconds = remaining_nanoseconds.wrapping_add(1_000_000_000); + remaining_seconds = remaining_seconds.wrapping_sub(1); + } + if remaining_seconds < 0 || (remaining_seconds == 0 && remaining_nanoseconds == 0) { + return -libc::ETIMEDOUT; + } + let mut remaining = KernelTimespec { + seconds: remaining_seconds, + nanoseconds: remaining_nanoseconds, + }; + + if poll { + let mut pollfd = KernelPollFd { + fd, + events: if receive { POLLIN } else { POLLOUT }, + revents: 0, + }; + let result = signal_safe_syscall6( + SYS_PPOLL, + (&raw mut pollfd).cast::() as usize, + 1, + (&raw mut remaining).cast::() as usize, + 0, + // The kernel's 64-bit signal-set size is eight bytes. No mask is supplied, but + // passing the ABI size keeps this a valid direct ppoll syscall on both targets. + 8, + 0, + ); + if result > 0 || result == -EINTR { + poll = false; + continue; + } + if result == 0 { + return -libc::ETIMEDOUT; + } + return result as i32; + } + + let result = if receive { + let mut ack = 1u8; + let result = signal_safe_syscall6( + SYS_RECVFROM, + fd as usize, + (&raw mut ack) as usize, + 1, + // MSG_TRUNC makes a packet larger than the one-byte buffer report its full size, + // so only an exact one-byte zero ACK can be accepted. + MSG_DONTWAIT | MSG_TRUNC, + 0, + 0, + ); + if result == 1 { + return if ack == 0 { 0 } else { -libc::EPROTO }; + } + result + } else { + let result = signal_safe_syscall6( + SYS_SENDTO, + fd as usize, + request as usize, + request_len, + // Nonblocking I/O lets the single ppoll deadline bound backpressure. MSG_NOSIGNAL + // prevents a closed sidecar socket from delivering SIGPIPE to the raw worker. + MSG_DONTWAIT | MSG_NOSIGNAL, + 0, + 0, + ); + // SOCK_SEQPACKET preserves message boundaries: a positive short send is a protocol + // failure rather than progress that can be resumed with a pointer offset. + if result == request_len as isize { + receive = true; + continue; + } + result + }; + + if result == -EINTR { + continue; + } + if result == -EAGAIN { + // Poll only after the kernel reports backpressure; a readiness wakeup retries the + // original send or receive operation and recomputes the remaining absolute deadline. + poll = true; + continue; + } + if result < 0 { + return result as i32; + } + // Any other positive result is an impossible short packet/send for this protocol. + return -libc::EPROTO; + } +} + +#[cfg(target_arch = "x86_64")] +#[inline(always)] +unsafe fn signal_safe_exit_group() -> ! { + // SYS_exit_group(0). The syscall cannot return, so no register clobbers remain live. + asm!( + "syscall", + in("rax") 231usize, + in("rdi") 0usize, + options(noreturn, nostack), + ); +} + +#[cfg(target_arch = "x86_64")] +#[inline(always)] +unsafe fn signal_safe_syscall6( + number: usize, + arg1: usize, + arg2: usize, + arg3: usize, + arg4: usize, + arg5: usize, + arg6: usize, +) -> isize { + let result: isize; + // Linux x86_64 uses r10, not the C ABI's rcx, for argument four. `syscall` itself clobbers rcx + // and r11; declaring both prevents the surrounding Rust from keeping live values there. + asm!( + "syscall", + inlateout("rax") number as isize => result, + in("rdi") arg1, + in("rsi") arg2, + in("rdx") arg3, + in("r10") arg4, + in("r8") arg5, + in("r9") arg6, + lateout("rcx") _, + lateout("r11") _, + options(nostack, preserves_flags), + ); + result +} + +#[cfg(target_arch = "x86_64")] +#[inline(always)] +unsafe fn signal_safe_load(flush: *const SignalFlush) -> (i32, i32, *const u8, usize) { + let fd: i32; + let owner_pid: i32; + let request: *const u8; + let request_len: usize; + // Do not replace this with `&(*flush).raw`. In panic=unwind debug builds that reference + // construction emitted null/alignment checks which call Rust panic handlers. Those handlers + // may allocate, use TLS, or otherwise depend on runtime state absent from the raw clone. + // These fixed-offset loads were audited to compile to four `mov` instructions and no calls. + asm!( + "movl {fd_offset}({base}), {fd:e}", + "movl {pid_offset}({base}), {owner_pid:e}", + "movq {request_offset}({base}), {request}", + "movq {length_offset}({base}), {request_len}", + base = in(reg) flush, + fd = out(reg) fd, + owner_pid = out(reg) owner_pid, + request = out(reg) request, + request_len = out(reg) request_len, + fd_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, fd), + pid_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, owner_pid), + request_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request), + length_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request_len), + options(att_syntax, nostack, readonly, preserves_flags), + ); + (fd, owner_pid, request, request_len) +} + +#[cfg(target_arch = "aarch64")] +#[inline(always)] +unsafe fn signal_safe_exit_group() -> ! { + // SYS_exit_group(0). The syscall cannot return, so no register clobbers remain live. + asm!( + "svc #0", + in("x8") 94usize, + in("x0") 0usize, + options(noreturn, nostack), + ); +} + +#[cfg(target_arch = "aarch64")] +#[inline(always)] +unsafe fn signal_safe_syscall6( + number: usize, + arg1: usize, + arg2: usize, + arg3: usize, + arg4: usize, + arg5: usize, + arg6: usize, +) -> isize { + let result: isize; + // Linux AArch64 takes the syscall number in x8, arguments in x0..x5, and returns in x0. + // Listing every register explicitly keeps the compiler from generating a helper call. + asm!( + "svc #0", + in("x8") number, + inlateout("x0") arg1 as isize => result, + in("x1") arg2, + in("x2") arg3, + in("x3") arg4, + in("x4") arg5, + in("x5") arg6, + options(nostack), + ); + result +} + +#[cfg(target_arch = "aarch64")] +#[inline(always)] +unsafe fn signal_safe_load(flush: *const SignalFlush) -> (i32, i32, *const u8, usize) { + let fd: i32; + let owner_pid: i32; + let request: *const u8; + let request_len: usize; + // As on x86_64, an ordinary Rust reference introduced panic-handler calls in debug builds. + // These fixed-offset loads were audited to compile to four `ldr` instructions and no calls. + asm!( + "ldr {fd:w}, [{base}, #{fd_offset}]", + "ldr {owner_pid:w}, [{base}, #{pid_offset}]", + "ldr {request}, [{base}, #{request_offset}]", + "ldr {request_len}, [{base}, #{length_offset}]", + base = in(reg) flush, + fd = out(reg) fd, + owner_pid = out(reg) owner_pid, + request = out(reg) request, + request_len = out(reg) request_len, + fd_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, fd), + pid_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, owner_pid), + request_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request), + length_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request_len), + options(nostack, readonly, preserves_flags), + ); + (fd, owner_pid, request, request_len) +} + +// The clone trampoline also rejects unsupported architectures. Never substitute libc here. +/// cbindgen:ignore +#[cfg(not(any(target_arch = "x86_64", target_arch = "aarch64")))] +#[no_mangle] +pub unsafe extern "C" fn datadog_sidecar_signal_flush_run( + _flush: *const SignalFlush, + _terminate_process: bool, +) -> i32 { + -libc::ENOTSUP +} + +fn connect_to_template(template: &mut SidecarTransport) -> anyhow::Result { + // SidecarTransport intentionally hides its endpoint. getpeername recovers the connected + // listener address without borrowing or retaining the template transport itself. + let mut address: libc::sockaddr_un = unsafe { std::mem::zeroed() }; + let mut length = size_of::() as libc::socklen_t; + if unsafe { + libc::getpeername( + template.as_raw_fd(), + (&mut address as *mut libc::sockaddr_un).cast(), + &mut length, + ) + } != 0 + { + return Err(std::io::Error::last_os_error().into()); + } + let start = offset_of!(libc::sockaddr_un, sun_path); + let length = length as usize; + anyhow::ensure!( + address.sun_family == libc::AF_UNIX as libc::sa_family_t + && length > start + 1 + && length <= size_of::() + && address.sun_path[0] == 0, + "signal flush template does not name an abstract Unix listener" + ); + // Linux abstract names are length-delimited and may contain embedded NULs. Use the sockaddr + // length returned by the kernel rather than treating sun_path as a C string. + let name = unsafe { + std::slice::from_raw_parts(address.sun_path.as_ptr().add(1).cast(), length - start - 1) + }; + Ok(SeqpacketConn::connect_abstract(name)?) +} + +#[cfg(test)] +mod tests { + use super::*; + use libdd_ipc::SeqpacketListener; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::time::{Duration, Instant}; + + fn pair() -> (SignalFlush, SeqpacketConn) { + let (connection, peer) = SeqpacketConn::socketpair().unwrap(); + let request = vec![42u8].into_boxed_slice(); + let flush = SignalFlush { + raw: RawSignalFlush { + fd: connection.as_raw_fd(), + owner_pid: unsafe { libc::getpid() }, + request: request.as_ptr(), + request_len: request.len(), + }, + _connection: connection, + _request: request, + }; + (flush, peer) + } + + fn send(peer: &SeqpacketConn, bytes: &[u8]) { + assert_eq!( + unsafe { libc::send(peer.as_raw_fd(), bytes.as_ptr().cast(), bytes.len(), 0) }, + bytes.len() as isize + ); + } + + #[test] + fn accepts_only_an_exact_zero_ack_packet() { + for (ack, expected) in [ + (&[0][..], 0), + (&[0, 0][..], -libc::EPROTO), + (&[1][..], -libc::EPROTO), + ] { + let (flush, peer) = pair(); + send(&peer, ack); + assert_eq!( + unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, + expected + ); + } + } + + #[test] + fn waits_for_the_ack_after_sending() { + let (flush, peer) = pair(); + let server = std::thread::spawn(move || { + let mut pollfd = libc::pollfd { + fd: peer.as_raw_fd(), + events: libc::POLLIN, + revents: 0, + }; + assert_eq!(unsafe { libc::poll(&mut pollfd, 1, 1000) }, 1); + let mut request = 0u8; + assert_eq!( + unsafe { libc::recv(peer.as_raw_fd(), (&mut request as *mut u8).cast(), 1, 0) }, + 1 + ); + assert_eq!(request, 42); + std::thread::sleep(Duration::from_millis(20)); + send(&peer, &[0]); + }); + assert_eq!( + unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, + 0 + ); + server.join().unwrap(); + } + + #[test] + fn retries_send_after_backpressure() { + let (flush, peer) = pair(); + let filler = [0u8; 512]; + loop { + let sent = unsafe { + libc::send( + flush.raw.fd, + filler.as_ptr().cast(), + filler.len(), + libc::MSG_DONTWAIT | libc::MSG_NOSIGNAL, + ) + }; + if sent < 0 { + assert_eq!( + std::io::Error::last_os_error().raw_os_error(), + Some(libc::EAGAIN) + ); + break; + } + } + let server = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(20)); + let mut bytes = [0u8; 512]; + loop { + let mut pollfd = libc::pollfd { + fd: peer.as_raw_fd(), + events: libc::POLLIN, + revents: 0, + }; + assert_eq!(unsafe { libc::poll(&mut pollfd, 1, 1000) }, 1); + let received = unsafe { + libc::recv(peer.as_raw_fd(), bytes.as_mut_ptr().cast(), bytes.len(), 0) + }; + assert!(received > 0); + if received == 1 && bytes[0] == 42 { + break; + } + assert_eq!(received, bytes.len() as isize); + } + send(&peer, &[0]); + }); + assert_eq!( + unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, + 0 + ); + server.join().unwrap(); + } + + #[test] + fn rejects_a_closed_peer_without_sigpipe() { + let (flush, peer) = pair(); + drop(peer); + assert!(unsafe { datadog_sidecar_signal_flush_run(&flush, false) } < 0); + } + + #[test] + fn refuses_inherited_objects() { + let (mut flush, _peer) = pair(); + flush.raw.owner_pid += 1; + assert_eq!( + unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, + -libc::ECHILD + ); + } + + #[test] + fn missing_ack_has_one_ten_second_deadline() { + let (flush, _peer) = pair(); + let start = Instant::now(); + assert_eq!( + unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, + -libc::ETIMEDOUT + ); + assert!(start.elapsed() >= Duration::from_secs(9)); + assert!(start.elapsed() < Duration::from_secs(15)); + } + + #[test] + fn preparation_creates_an_independent_connection_to_the_template_listener() { + static NEXT: AtomicUsize = AtomicUsize::new(0); + let name = format!( + "ddtrace-signal-flush-{}-{}", + std::process::id(), + NEXT.fetch_add(1, Ordering::Relaxed) + ); + let listener = SeqpacketListener::bind_abstract(name.as_bytes()).unwrap(); + let mut template = + SidecarTransport::from(SeqpacketConn::connect_abstract(name.as_bytes()).unwrap()); + let template_peer = listener.try_accept().unwrap(); + let mut output = std::ptr::null_mut(); + assert!(matches!( + unsafe { datadog_sidecar_prepare_signal_flush(&mut template, &mut output) }, + MaybeError::None + )); + let flush = unsafe { Box::from_raw(output) }; + let signal_peer = listener.try_accept().unwrap(); + assert_ne!(flush.raw.fd, template.as_raw_fd()); + + // Closing the normal connection must produce EOF despite the signal connection. + drop(template); + let mut bytes = [0u8; 256]; + assert_eq!( + unsafe { + libc::recv( + template_peer.as_raw_fd(), + bytes.as_mut_ptr().cast(), + bytes.len(), + 0, + ) + }, + 0 + ); + send(&signal_peer, &[0]); + assert_eq!( + unsafe { datadog_sidecar_signal_flush_run(&*flush, false) }, + 0 + ); + let received = unsafe { + libc::recv( + signal_peer.as_raw_fd(), + bytes.as_mut_ptr().cast(), + bytes.len(), + 0, + ) + }; + assert_eq!(received as usize, flush._request.len()); + assert_eq!(&bytes[..received as usize], &*flush._request); + drop(flush); + assert_eq!( + unsafe { + libc::recv( + signal_peer.as_raw_fd(), + bytes.as_mut_ptr().cast(), + bytes.len(), + 0, + ) + }, + 0 + ); + } +} diff --git a/ext/sidecar.c b/ext/sidecar.c index e5439e451cd..03f13867e06 100644 --- a/ext/sidecar.c +++ b/ext/sidecar.c @@ -13,6 +13,7 @@ #include "telemetry.h" #include "process_tags.h" #include "remote_config.h" +#include "signals.h" #include "string_utils.h" #include "target_metadata.h" #include "ffi_utils.h" @@ -204,6 +205,7 @@ void datadog_sidecar_refresh_user_service_defined(void) { } static void datadog_sidecar_setup_thread_mode(void); +static void dd_sidecar_setup_signal_transport(void); static void dd_sidecar_on_reconnect(ddog_SidecarTransport *transport) { if (!datadog_endpoint || !dogstatsd_endpoint) { @@ -308,6 +310,22 @@ static ddog_SidecarTransport *dd_sidecar_connect(bool as_worker, bool is_fork) { return sidecar_transport; } +static void dd_sidecar_setup_signal_transport(void) { +#ifdef __linux__ + if (datadog_signals_has_sidecar_flush() || !DATADOG_G(sidecar) || + (!get_global_DD_TRACE_FORCE_FLUSH_ON_SIGTERM() && !get_global_DD_TRACE_FORCE_FLUSH_ON_SIGINT())) { + return; + } + + ddog_SignalFlush *flush = NULL; + if (datadog_ffi_try("Failed preparing signal-only sidecar connection", + datadog_sidecar_prepare_signal_flush(DATADOG_G(sidecar), &flush))) { + // Takes ownership, including when another normal thread published first. + datadog_signals_set_sidecar_flush(flush); + } +#endif +} + static void datadog_sidecar_setup_thread_mode() { #ifndef _WIN32 int32_t current_pid = (int32_t)getpid(); @@ -470,6 +488,7 @@ void datadog_sidecar_setup(ddog_RemoteConfigFlags flags) { if (DATADOG_G(sidecar) && !datadog_sidecar_for_signal) { datadog_sidecar_for_signal = DATADOG_G(sidecar); } + dd_sidecar_setup_signal_transport(); } void datadog_sidecar_minit(void) { @@ -489,6 +508,9 @@ void datadog_sidecar_minit(void) { void datadog_sidecar_handle_fork(void) { #ifndef _WIN32 +#ifdef __linux__ + datadog_signals_reset_sidecar_flush_after_fork(); +#endif ddog_RemoteConfigFlags flags = {0}; bool enable_sidecar = datadog_sidecar_should_enable(&flags); @@ -546,6 +568,7 @@ void datadog_sidecar_handle_fork(void) { if (DATADOG_G(sidecar)) { datadog_sidecar_for_signal = DATADOG_G(sidecar); } + dd_sidecar_setup_signal_transport(); #endif } @@ -580,6 +603,7 @@ void datadog_sidecar_ensure_active(void) { datadog_sidecar_for_signal = DATADOG_G(sidecar); } } + dd_sidecar_setup_signal_transport(); } void datadog_sidecar_finalize(bool clear_id) { diff --git a/ext/signals.c b/ext/signals.c index 8926b4b442e..6025eb823d9 100644 --- a/ext/signals.c +++ b/ext/signals.c @@ -46,7 +46,10 @@ #endif #if __linux +#include #include +#include +#include #include #endif @@ -300,70 +303,234 @@ static struct sigaction dd_sigint_sigterm_sigaction; static struct sigaction dd_sigterm_prev_sigaction; static struct sigaction dd_sigint_prev_sigaction; -struct { - int sig; - siginfo_t si; - void *uc; -} dd_signal_data; +// The cleanup stack and prepared flush object are allocated in ordinary +// context. Once READY is published, the signal handler only reads them and +// exactly one raw worker may use them. +static char *dd_signal_cleanup_stack; +static size_t dd_signal_cleanup_stack_size; +static ddog_SignalFlush *dd_signal_flush; +// The process ID, rather than a thread ID, lets every PHP thread in this process +// pass while rejecting a copy inherited by a fork child. +static _Atomic(int) dd_signal_owner_pid; +// CLONE_PARENT_SETTID writes the worker TID here. CLONE_CHILD_CLEARTID clears it +// and performs a futex wake at thread exit, giving normal shutdown a join +// primitive that does not depend on pthread state. +static _Atomic(int) dd_signal_worker_tid; + +// The object is one-shot: +// +// DISABLED -> INSTALLING -> READY -> STARTING_* -> RUNNING_* -> STOPPED +// `-> FAILED -----------^ +// +// INSTALLING protects publication from teardown. STARTING protects clone setup. +// The DEFAULT and CUSTOM variants preserve who is responsible for termination. +enum { + DD_SIGNAL_DISABLED, // No flush object is installed. + DD_SIGNAL_INSTALLING, // A flush object is being published. + DD_SIGNAL_READY, // One signal may claim the published object. + DD_SIGNAL_STARTING_DEFAULT, // clone() is in progress; the worker will terminate the process. + DD_SIGNAL_STARTING_CUSTOM, // clone() is in progress; the worker will return after flushing. + DD_SIGNAL_RUNNING_DEFAULT, // The TID is joinable; the worker will terminate the process. + DD_SIGNAL_RUNNING_CUSTOM, // The TID is joinable; the worker will return after flushing. + DD_SIGNAL_FAILED, // Worker creation failed; another signal may not retry it. + DD_SIGNAL_STOPPED, // Module shutdown has disabled new workers. +}; +static _Atomic(int) dd_signal_state; +// Atomics used by a signal handler must compile to instructions; a libatomic +// fallback could lock or touch runtime state while the interrupted thread owns it. +// The TID object is also written directly by the kernel, whose clone/futex ABI +// requires the ordinary int representation. +_Static_assert(ATOMIC_INT_LOCK_FREE == 2, "signal atomics must be lock-free"); +_Static_assert(sizeof(dd_signal_worker_tid) == sizeof(int), "signal TID must have the kernel int layout"); + +static void dd_signals_init_cleanup_stack(void); +static void dd_signals_drop_sidecar_flush(void); +static void dd_call_prev_handler(int sig, siginfo_t *si, void *uc); + +bool datadog_signals_has_sidecar_flush(void) { + // Any state other than DISABLED owns or is publishing an object. In + // particular, FAILED and STOPPED must not allow a second publication. + return atomic_load_explicit(&dd_signal_state, memory_order_acquire) != DD_SIGNAL_DISABLED; +} -static int dd_call_prev_handler(bool flush) { - struct sigaction prev_sigaction = dd_signal_data.sig == SIGINT ? dd_sigint_prev_sigaction : dd_sigterm_prev_sigaction; - void *prev_handler = (prev_sigaction.sa_flags & SA_SIGINFO) ? (void *)prev_sigaction.sa_sigaction : (void *)prev_sigaction.sa_handler; - if (prev_handler == SIG_IGN) { - return 0; +void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush) { + // Consume `flush`. On success it becomes the process-wide signal flush + // object; if it cannot be installed, destroy it before returning. + int expected = DD_SIGNAL_DISABLED; + if (!flush) { + return; } + // prevent a previous handler from suspending this publication + sigset_t publication_signals, old_signals; + sigemptyset(&publication_signals); + sigaddset(&publication_signals, SIGTERM); + sigaddset(&publication_signals, SIGINT); + if (sigprocmask(SIG_BLOCK, &publication_signals, &old_signals) < 0) { + datadog_sidecar_signal_flush_drop(flush); + return; + } + if (!atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, DD_SIGNAL_INSTALLING, + memory_order_acq_rel, memory_order_acquire)) { + datadog_sidecar_signal_flush_drop(flush); + sigprocmask(SIG_SETMASK, &old_signals, NULL); + return; + } + dd_signal_flush = flush; + atomic_store_explicit(&dd_signal_owner_pid, getpid(), memory_order_relaxed); + // READY is the publication barrier for the pointer and owner PID. A handler + // can acquire READY only after both values have been initialized. + atomic_store_explicit(&dd_signal_state, DD_SIGNAL_READY, memory_order_release); + // A pending SIGTERM/SIGINT may run as soon as the old mask is restored; at + // that point it observes either the complete publication or a terminal state. + sigprocmask(SIG_SETMASK, &old_signals, NULL); +} - if (flush) { - ddog_sidecar_flush(&datadog_sidecar_for_signal, (ddog_SidecarFlushOptions){.traces_and_stats = true}); +void datadog_signals_reset_sidecar_flush_after_fork(void) { + atomic_store_explicit(&dd_signal_worker_tid, 0, memory_order_relaxed); + ddog_SignalFlush *flush = dd_signal_flush; + dd_signal_flush = NULL; + datadog_sidecar_signal_flush_drop(flush); + atomic_store_explicit(&dd_signal_state, DD_SIGNAL_DISABLED, memory_order_release); + // the owner pid is purposefully still the old one +} + +static void dd_signals_drop_sidecar_flush(void) { + // Called during module shutdown before freeing the flush object and stack. + // Wait if the signal handler is creating a worker, and wait for an existing + // worker to exit. Then set STOPPED so later signals use the previous handler. + for (;;) { + int state = atomic_load_explicit(&dd_signal_state, memory_order_acquire); + if (state == DD_SIGNAL_INSTALLING || state == DD_SIGNAL_STARTING_DEFAULT || + state == DD_SIGNAL_STARTING_CUSTOM) { + // Ordinary shutdown context. The signal handler performs only bounded + // startup. + sched_yield(); + continue; + } + if (state == DD_SIGNAL_RUNNING_DEFAULT || state == DD_SIGNAL_RUNNING_CUSTOM) { + int tid; + while ((tid = atomic_load_explicit(&dd_signal_worker_tid, memory_order_acquire)) != 0) { + // FUTEX_WAIT sleeps only if the word still equals tid, closing + // the exit-before-wait race. EAGAIN, EINTR, and spurious wakes + // are harmless because the loop reloads the word. The kernel's + // clear_child_tid wake uses a shared futex, not FUTEX_WAIT_PRIVATE. + datadog_raw_syscall6(SYS_futex, (long)(uintptr_t)&dd_signal_worker_tid, FUTEX_WAIT, tid, 0, 0, 0); + } + } + // A READY handler can race this transition. compare_exchange either + // stops it before clone, or reports its newer STARTING/RUNNING state so + // the loop waits for it above. + if (atomic_compare_exchange_strong_explicit(&dd_signal_state, &state, DD_SIGNAL_STOPPED, + memory_order_acq_rel, memory_order_acquire)) { + break; + } } + ddog_SignalFlush *flush = dd_signal_flush; + dd_signal_flush = NULL; + datadog_sidecar_signal_flush_drop(flush); +} + +static void dd_call_prev_handler(int sig, siginfo_t *si, void *uc) { + struct sigaction *prev_sigaction = sig == SIGINT ? &dd_sigint_prev_sigaction : &dd_sigterm_prev_sigaction; + void *prev_handler = (prev_sigaction->sa_flags & SA_SIGINFO) ? (void *)prev_sigaction->sa_sigaction + : (void *)prev_sigaction->sa_handler; + if (prev_handler == SIG_IGN) { + return; + } if (prev_handler == SIG_DFL) { _exit(0); } - - if (prev_sigaction.sa_flags & SA_SIGINFO) { - (*prev_sigaction.sa_sigaction)(dd_signal_data.sig, &dd_signal_data.si, dd_signal_data.uc); + if (prev_sigaction->sa_flags & SA_SIGINFO) { + (*prev_sigaction->sa_sigaction)(sig, si, uc); } else { - (*prev_sigaction.sa_handler)(dd_signal_data.sig); + (*prev_sigaction->sa_handler)(sig); } - - return 0; - } -static int dd_sigterm_cleanup_thread(void *arg) { - // Block all signals to prevent delivery to this thread - sigset_t set; - sigfillset(&set); - sigprocmask(SIG_BLOCK, &set, NULL); +static void dd_sigint_sigterm_handler(int sig, siginfo_t *si, void *uc) { + struct sigaction *prev_sigaction = sig == SIGINT ? &dd_sigint_prev_sigaction : &dd_sigterm_prev_sigaction; + void *prev_handler = (prev_sigaction->sa_flags & SA_SIGINFO) ? (void *)prev_sigaction->sa_sigaction + : (void *)prev_sigaction->sa_handler; + if (prev_handler == SIG_IGN) { + return; + } + bool terminate = prev_handler == SIG_DFL; + int expected = DD_SIGNAL_READY; + int starting = terminate ? DD_SIGNAL_STARTING_DEFAULT : DD_SIGNAL_STARTING_CUSTOM; + // READY is consumed exactly once. Besides preventing duplicate workers, the + // acquire half makes the prepared flush object and owner PID visible. + if (!atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, starting, memory_order_acq_rel, + memory_order_acquire)) { + // A worker handling a default-disposition signal will terminate the + // process after flushing. Do not let another default signal iterrupt + // the flush early by invoking the previous disposition. + if (terminate && (expected == DD_SIGNAL_STARTING_DEFAULT || expected == DD_SIGNAL_RUNNING_DEFAULT)) { + return; + } + dd_call_prev_handler(sig, si, uc); + return; + } - // Make the Go runtime believe, we are actually running on a signal stack - stack_t altstack; - altstack.ss_sp = dd_signal_async_stack; - if (altstack.ss_sp) { - altstack.ss_size = dd_signal_async_stack_size; - altstack.ss_flags = 0; - sigaltstack(&altstack, NULL); + // A fork child may inherit READY before its post-fork reset. Its owner PID + // is still the parent's, so reject it before accessing the inherited flush + // object. + if (datadog_raw_syscall6(SYS_getpid, 0, 0, 0, 0, 0, 0) != + atomic_load_explicit(&dd_signal_owner_pid, memory_order_relaxed)) { + atomic_store_explicit(&dd_signal_state, DD_SIGNAL_FAILED, memory_order_release); + dd_call_prev_handler(sig, si, uc); + return; + } + if (!dd_signal_cleanup_stack) { + atomic_store_explicit(&dd_signal_state, DD_SIGNAL_FAILED, memory_order_release); + dd_call_prev_handler(sig, si, uc); + return; } - return dd_call_prev_handler(true); + // The handler's sa_mask already blocks ordinary blockable signals. glibc and + // musl deliberately remove their internal signals from masks installed via + // public libc APIs, however. If one reached the raw clone, its libc handler + // could use the parent's inherited, unrepaired TLS. Bypass that filtering and + // let the clone inherit the complete kernel mask before it can run. Doing + // this inside the clone would leave a delivery window. SIGKILL and SIGSTOP + // remain unblockable, but neither executes a user-space handler. + uint64_t all_signals = UINT64_MAX, old_signals; + if (datadog_raw_syscall6(SYS_rt_sigprocmask, SIG_SETMASK, (long)(uintptr_t)&all_signals, + (long)(uintptr_t)&old_signals, sizeof(all_signals), 0, 0) < 0) { + atomic_store_explicit(&dd_signal_state, DD_SIGNAL_FAILED, memory_order_release); + dd_call_prev_handler(sig, si, uc); + return; + } + void *stack_top = dd_signal_cleanup_stack + dd_signal_cleanup_stack_size; + // These are pthread-like sharing flags without CLONE_SETTLS: the worker has + // no independent libc/Rust thread runtime and may execute only the audited + // raw call graph. PARENT_SETTID plus CHILD_CLEARTID make + // dd_signal_worker_tid a kernel-backed join word. + int flags = CLONE_VM | CLONE_FS | CLONE_FILES | CLONE_SIGHAND | CLONE_THREAD | CLONE_SYSVSEM | CLONE_PARENT_SETTID | + CLONE_CHILD_CLEARTID; + // The trampoline copies `terminate` onto the child stack before clone. Rust + // either returns for the custom-handler case or issues exit_group after the + // bounded exchange for the default disposition. + int result = datadog_clone_thread(datadog_sidecar_signal_flush_run, stack_top, flags, dd_signal_flush, terminate, + &dd_signal_worker_tid); + int running = terminate ? DD_SIGNAL_RUNNING_DEFAULT : DD_SIGNAL_RUNNING_CUSTOM; + // The child may finish before clone returns. That is safe: clear_child_tid + // will already be zero when normal teardown observes the RUNNING state. + atomic_store_explicit(&dd_signal_state, result < 0 ? DD_SIGNAL_FAILED : running, memory_order_release); + datadog_raw_syscall6(SYS_rt_sigprocmask, SIG_SETMASK, (long)(uintptr_t)&old_signals, 0, sizeof(old_signals), 0, 0); + + if (result < 0 || prev_handler != SIG_DFL) { + dd_call_prev_handler(sig, si, uc); + } // else the cleanup worker started and calling dd_call_prev_handler + // would exit the process } -static void dd_sigint_sigterm_handler(int sig, siginfo_t *si, void *uc) { - dd_signal_data.sig = sig; - memcpy(&dd_signal_data.si, si, sizeof(*si)); - dd_signal_data.uc = uc; - - if (datadog_sidecar_for_signal) { - // Spawn a thread using clone() to perform sidecar cleanup asynchronously to avoid async unsafeness in the signal handler - void *stack_top = dd_signal_async_stack + dd_signal_async_stack_size; - int flags = CLONE_VM | CLONE_FS | CLONE_FILES | CLONE_SIGHAND | CLONE_THREAD | CLONE_SYSVSEM; - if (datadog_clone_thread(dd_sigterm_cleanup_thread, stack_top, flags, NULL) < 0) { - // If the cleanup thread could not be started, we just do it ourselves. Will block, but that's okay then. - dd_call_prev_handler(true); - } - } else { - dd_call_prev_handler(false); +static void dd_signals_init_cleanup_stack(void) { + if (!dd_signal_cleanup_stack) { + // Allocate before signal delivery. The one-shot state gate ensures that + // no two workers ever share this stack. + dd_signal_cleanup_stack_size = MIN_STACKSZ; + dd_signal_cleanup_stack = malloc(dd_signal_cleanup_stack_size); } } #endif @@ -374,11 +541,11 @@ void datadog_signals_minit(void) { dd_sigint_sigterm_sigaction.sa_flags = SA_SIGINFO; sigemptyset(&dd_sigint_sigterm_sigaction.sa_mask); if (get_global_DD_TRACE_FORCE_FLUSH_ON_SIGTERM()) { - dd_signals_init_async_stack(); + dd_signals_init_cleanup_stack(); sigaction(SIGTERM, &dd_sigint_sigterm_sigaction, &dd_sigterm_prev_sigaction); } if (get_global_DD_TRACE_FORCE_FLUSH_ON_SIGINT()) { - dd_signals_init_async_stack(); + dd_signals_init_cleanup_stack(); sigaction(SIGINT, &dd_sigint_sigterm_sigaction, &dd_sigint_prev_sigaction); } #endif @@ -386,6 +553,8 @@ void datadog_signals_minit(void) { void datadog_signals_mshutdown(void) { #if __linux + // wait for the signal cleanup thread to exit + dd_signals_drop_sidecar_flush(); if (dd_sigint_sigterm_sigaction.sa_sigaction) { if (get_global_DD_TRACE_FORCE_FLUSH_ON_SIGTERM()) { sigaction(SIGTERM, &dd_sigterm_prev_sigaction, NULL); @@ -394,6 +563,9 @@ void datadog_signals_mshutdown(void) { sigaction(SIGINT, &dd_sigint_prev_sigaction, NULL); } } + + free(dd_signal_cleanup_stack); + dd_signal_cleanup_stack = NULL; #endif if (dd_signal_async_stack) { diff --git a/ext/signals.h b/ext/signals.h index 326abb9542c..bd29848ed6f 100644 --- a/ext/signals.h +++ b/ext/signals.h @@ -1,9 +1,16 @@ #ifndef DD_TRACE_SIGNALS_H #define DD_TRACE_SIGNALS_H +#include + +typedef struct ddog_SignalFlush ddog_SignalFlush; + void datadog_set_coredumpfilter(void); void datadog_signals_first_rinit(void); void datadog_signals_minit(void); void datadog_signals_mshutdown(void); +bool datadog_signals_has_sidecar_flush(void); +void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush); +void datadog_signals_reset_sidecar_flush_after_fork(void); #endif // DD_TRACE_SIGNALS_H diff --git a/ext/threads.c b/ext/threads.c index aee7d37476a..427d49cba01 100644 --- a/ext/threads.c +++ b/ext/threads.c @@ -156,23 +156,27 @@ __asm__( ".hidden datadog_clone_thread\n" ".type datadog_clone_thread,@function\n" "datadog_clone_thread:\n" - /* in: rdi = fn, rsi = stack_top, edx = flags, rcx = arg */ + /* in: rdi = fn, rsi = stack_top, edx = flags, rcx = arg, + * r8b = terminate_process, r9 = tid */ " andq $-16, %rsi\n" /* align the child stack */ - " subq $16, %rsi\n" /* hand fn and arg over on it */ + " subq $32, %rsi\n" /* hand fn, arg and flag over on it */ " movq %rdi, 0(%rsi)\n" " movq %rcx, 8(%rsi)\n" + " movb %r8b, 16(%rsi)\n" /* syscall: rdi = flags, rsi = newsp, rdx = parent_tid, r10 = child_tid, r8 = tls */ " movl %edx, %edi\n" - " xorl %edx, %edx\n" - " xorl %r10d, %r10d\n" + " movq %r9, %rdx\n" /* independent parent_tid / clear_child_tid word */ + " movq %r9, %r10\n" " xorl %r8d, %r8d\n" " movl $56, %eax\n" /* SYS_clone */ " syscall\n" " testq %rax, %rax\n" /* parent: tid or -errno; child: 0 */ " jnz 1f\n" " xorl %ebp, %ebp\n" /* end the frame pointer chain */ - " popq %rax\n" /* fn */ - " popq %rdi\n" /* arg */ + " movq 0(%rsp), %rax\n" /* fn */ + " movq 8(%rsp), %rdi\n" /* arg */ + " movzbl 16(%rsp), %esi\n" /* terminate_process */ + " addq $32, %rsp\n" /* restore alignment before call */ " callq *%rax\n" " movl %eax, %edi\n" /* fn's return value is the thread's exit status */ " movl $60, %eax\n" /* SYS_exit -- this thread only, not exit_group */ @@ -180,6 +184,24 @@ __asm__( " hlt\n" /* unreachable */ "1: ret\n" ".size datadog_clone_thread,.-datadog_clone_thread\n"); + +__asm__( + ".text\n" + ".globl datadog_raw_syscall6\n" + ".hidden datadog_raw_syscall6\n" + ".type datadog_raw_syscall6,@function\n" + "datadog_raw_syscall6:\n" + /* C ABI: rdi = number, rsi/rdi... = six syscall arguments. */ + " movq %rdi, %rax\n" + " movq %rsi, %rdi\n" + " movq %rdx, %rsi\n" + " movq %rcx, %rdx\n" + " movq %r8, %r10\n" + " movq %r9, %r8\n" + " movq 8(%rsp), %r9\n" + " syscall\n" + " ret\n" + ".size datadog_raw_syscall6,.-datadog_raw_syscall6\n"); #elif defined(__aarch64__) __asm__( ".text\n" @@ -187,29 +209,77 @@ __asm__( ".hidden datadog_clone_thread\n" ".type datadog_clone_thread,%function\n" "datadog_clone_thread:\n" - /* in: x0 = fn, x1 = stack_top, w2 = flags, x3 = arg */ + /* in: x0 = fn, x1 = stack_top, w2 = flags, x3 = arg, + * w4 = terminate_process, x5 = tid */ " and x1, x1, #-16\n" /* align the child stack */ - " stp x0, x3, [x1, #-16]!\n" /* hand fn and arg over on it; x1 becomes newsp */ + " sub x1, x1, #32\n" + " stp x0, x3, [x1]\n" /* hand fn and arg over on it; x1 is newsp */ + " strb w4, [x1, #16]\n" /* terminate_process */ /* syscall: x0 = flags, x1 = newsp, x2 = parent_tid, x3 = tls, x4 = child_tid */ " mov w0, w2\n" - " mov x2, #0\n" + " mov x2, x5\n" /* same independent parent_tid / clear_child_tid word */ " mov x3, #0\n" - " mov x4, #0\n" + " mov x4, x5\n" " mov x8, #220\n" /* SYS_clone */ " svc #0\n" " cbz x0, 1f\n" /* parent: tid or -errno; child: 0 */ " ret\n" - "1: ldp x1, x0, [sp], #16\n" /* x1 = fn, x0 = arg */ + "1: ldp x16, x0, [sp]\n" /* x16 = fn, x0 = arg */ + " ldrb w1, [sp, #16]\n" /* terminate_process */ + " add sp, sp, #32\n" " mov x29, #0\n" /* end the frame pointer chain */ - " blr x1\n" /* fn's return value is left in w0 */ + " blr x16\n" /* fn's return value is left in w0 */ " mov w8, #93\n" /* SYS_exit -- this thread only, not exit_group */ " svc #0\n" " brk #0\n" /* unreachable */ ".size datadog_clone_thread,.-datadog_clone_thread\n"); + +__asm__( + ".text\n" + ".globl datadog_raw_syscall6\n" + ".hidden datadog_raw_syscall6\n" + ".type datadog_raw_syscall6,%function\n" + "datadog_raw_syscall6:\n" + /* C ABI: x0 = number, x1..x6 = six syscall arguments. */ + " mov x8, x0\n" + " mov x0, x1\n" + " mov x1, x2\n" + " mov x2, x3\n" + " mov x3, x4\n" + " mov x4, x5\n" + " mov x5, x6\n" + " svc #0\n" + " ret\n" + ".size datadog_raw_syscall6,.-datadog_raw_syscall6\n"); #else -#include -int datadog_clone_thread(int (*fn)(void *), void *stack_top, int flags, void *arg) { - return clone(fn, stack_top, flags, arg); +int datadog_clone_thread(datadog_raw_clone_fn fn, void *stack_top, int flags, + const struct ddog_SignalFlush *arg, bool terminate_process, _Atomic(int) *tid) { + (void)fn; + (void)stack_top; + (void)flags; + (void)arg; + (void)terminate_process; + (void)tid; + return -1; +} + +long datadog_raw_syscall6( + long number, + long arg1, + long arg2, + long arg3, + long arg4, + long arg5, + long arg6 +) { + (void)number; + (void)arg1; + (void)arg2; + (void)arg3; + (void)arg4; + (void)arg5; + (void)arg6; + return -1; } #endif diff --git a/ext/threads.h b/ext/threads.h index a7678b05c1b..aace9961ee4 100644 --- a/ext/threads.h +++ b/ext/threads.h @@ -1,6 +1,8 @@ #ifndef DATADOG_THREADS_H #define DATADOG_THREADS_H +#include +#include #include #include @@ -27,7 +29,22 @@ TSRM_API int tsrm_mutex_unlock(MUTEX_T mutexp); #endif #ifdef __linux__ -int datadog_clone_thread(int (*fn)(void *), void *stack_top, int flags, void *arg); +#include + +struct ddog_SignalFlush; +typedef int32_t (*datadog_raw_clone_fn)(const struct ddog_SignalFlush *, bool); + +int datadog_clone_thread(datadog_raw_clone_fn fn, void *stack_top, int flags, + const struct ddog_SignalFlush *arg, bool terminate_process, _Atomic(int) *tid); +long datadog_raw_syscall6( + long number, + long arg1, + long arg2, + long arg3, + long arg4, + long arg5, + long arg6 +); #endif #endif // DATADOG_THREADS_H From 642ecc7edc9e4356817a2fe91d70fadb5e0e54f6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gustavo=20Andr=C3=A9=20dos=20Santos=20Lopes?= Date: Mon, 28 Sep 2026 14:45:23 +0100 Subject: [PATCH 2/3] Refresh signal flush transport after reconnect A sidecar restart left the prepared signal socket stale after the normal transport recovered, so forced shutdown could skip its flush. Refresh the socket from the replacement transport while it is still unclaimed. Publish before freeing the retired object so cross-thread signals cannot deadlock on allocator state. # Conflicts: # Makefile --- Cargo.lock | 199 +++++++----------- Makefile | 10 +- .../src/test/www/signal-flush/run.sh | 21 ++ .../www/signal-flush/signal_flush_worker.php | 49 +++++ ext/sidecar.c | 33 ++- ext/signals.c | 60 ++++-- ext/signals.h | 7 +- libdatadog | 2 +- 8 files changed, 226 insertions(+), 155 deletions(-) create mode 100755 appsec/tests/integration/src/test/www/signal-flush/run.sh create mode 100644 appsec/tests/integration/src/test/www/signal-flush/signal_flush_worker.php diff --git a/Cargo.lock b/Cargo.lock index 56177f8ccee..e3d3a28f563 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -154,12 +154,6 @@ version = "1.7.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69f7f8c3906b62b754cd5326047894316021dcfe5a194c8ea52bdd94934a3457" -[[package]] -name = "arrayref" -version = "0.3.9" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "76a2e8124351fda1ef8aaaa3bbd7ebbcb486bbcd4225aca0aa0d84bb2db8fecb" - [[package]] name = "assert-json-diff" version = "2.0.2" @@ -347,7 +341,7 @@ dependencies = [ "bitflags 2.13.0", "cexpr", "clang-sys", - "itertools 0.10.5", + "itertools 0.11.0", "log", "prettyplease", "proc-macro2", @@ -542,7 +536,6 @@ name = "build_common" version = "0.0.1" dependencies = [ "cbindgen 0.29.0", - "serde", "serde_json", ] @@ -639,7 +632,6 @@ version = "0.29.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "975982cdb7ad6a142be15bdf84aea7ec6a9e5d4d797c004d43185b24cfe4e684" dependencies = [ - "clap", "heck 0.5.0", "indexmap 2.12.1", "log", @@ -1295,7 +1287,6 @@ version = "0.0.1" dependencies = [ "anyhow", "arc-swap", - "arrayref", "base64 0.22.1", "bincode", "chrono", @@ -1304,7 +1295,6 @@ dependencies = [ "datadog-sidecar-macros", "futures", "http 1.4.2", - "http-body-util", "httpmock", "libc 0.2.186", "libdd-capabilities", @@ -1326,13 +1316,10 @@ dependencies = [ "libdd-trace-utils", "manual_future", "memory-stats", - "microseh", "nix 0.29.0", "prctl", "priority-queue", "rand 0.8.8", - "rmp-serde", - "sendfd", "serde", "serde_json", "serde_with", @@ -1347,8 +1334,6 @@ dependencies = [ "tracing-log", "tracing-subscriber", "winapi 0.3.9", - "windows 0.51.1", - "windows-sys 0.52.0", "zwohash", ] @@ -1371,10 +1356,8 @@ dependencies = [ "libdd-telemetry-ffi", "libdd-tinybytes", "libdd-trace-utils", - "paste", "rmp-serde", "serde_json", - "tempfile", "tracing", ] @@ -2302,13 +2285,14 @@ checksum = "9a3a5bfb195931eeb336b2a7b4d761daec841b97f947d34394601737a7bba5e4" [[package]] name = "hyper" -version = "1.6.0" +version = "1.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cc2b571658e38e0c01b1fdca3bbbe93c00d3d71693ff2770043f8c29bc7d6f80" +checksum = "27b501faa50e7a26c3d3560ca625132f4078a17771f4810baf70475ae48cbe43" dependencies = [ + "atomic-waker", "bytes", "futures-channel", - "futures-util", + "futures-core", "h2", "http 1.4.2", "http-body", @@ -2770,7 +2754,7 @@ dependencies = [ [[package]] name = "libdd-capabilities" -version = "3.0.1" +version = "4.0.0" dependencies = [ "anyhow", "bytes", @@ -2782,7 +2766,7 @@ dependencies = [ [[package]] name = "libdd-capabilities-impl" -version = "5.0.0" +version = "6.0.0" dependencies = [ "anyhow", "bytes", @@ -2796,7 +2780,7 @@ dependencies = [ [[package]] name = "libdd-common" -version = "6.0.0" +version = "7.0.0" dependencies = [ "anyhow", "bytes", @@ -2804,7 +2788,6 @@ dependencies = [ "const_format", "criterion", "futures", - "futures-core", "futures-util", "hex", "http 1.4.2", @@ -2827,7 +2810,6 @@ dependencies = [ "regex-lite", "reqwest 0.13.2", "rustls", - "rustls-native-certs", "rustls-platform-verifier", "rustls-webpki", "serde", @@ -2835,7 +2817,6 @@ dependencies = [ "tempfile", "thiserror 2.0.18", "tokio", - "tokio-rustls", "tower-service", "windows-sys 0.52.0", ] @@ -2884,7 +2865,7 @@ dependencies = [ "page_size", "portable-atomic", "rand 0.8.8", - "schemars", + "schemars 0.8.21", "serde", "serde_json", "symbolic-common", @@ -2907,7 +2888,6 @@ dependencies = [ "libdd-common", "libdd-common-ffi", "libdd-crashtracker", - "serde", "serde_json", "symbolic-common", "symbolic-demangle", @@ -2917,7 +2897,7 @@ dependencies = [ [[package]] name = "libdd-data-pipeline" -version = "10.0.0" +version = "11.0.0" dependencies = [ "anyhow", "arc-swap", @@ -2928,7 +2908,6 @@ dependencies = [ "duplicate", "either", "futures", - "getrandom 0.2.15", "h2", "http 1.4.2", "http-body-util", @@ -2945,7 +2924,6 @@ dependencies = [ "libdd-shared-runtime", "libdd-telemetry", "libdd-tinybytes", - "libdd-trace-normalization", "libdd-trace-obfuscation", "libdd-trace-protobuf", "libdd-trace-stats", @@ -2970,8 +2948,9 @@ dependencies = [ [[package]] name = "libdd-data-pipeline-core" -version = "1.0.0" +version = "2.0.0" dependencies = [ + "anyhow", "bytes", "futures", "http 1.4.2", @@ -2979,6 +2958,7 @@ dependencies = [ "libdd-common", "libdd-tinybytes", "libdd-trace-obfuscation", + "libdd-trace-stats", "libdd-trace-utils", "serde_json", "thiserror 2.0.18", @@ -2987,7 +2967,7 @@ dependencies = [ [[package]] name = "libdd-ddsketch" -version = "1.1.1" +version = "1.1.2" dependencies = [ "criterion", "prost", @@ -2999,7 +2979,7 @@ dependencies = [ [[package]] name = "libdd-dogstatsd-client" -version = "6.0.0" +version = "7.0.0" dependencies = [ "anyhow", "async-trait", @@ -3034,12 +3014,9 @@ dependencies = [ "pyo3", "semver", "serde", - "serde-bool", "serde_json", - "serde_with", "thiserror 2.0.18", "tokio", - "url", ] [[package]] @@ -3057,7 +3034,6 @@ dependencies = [ "anyhow", "bincode", "criterion", - "futures", "glibc_version", "io-lifetimes", "libc 0.2.186", @@ -3070,14 +3046,11 @@ dependencies = [ "memfd", "nix 0.29.0", "page_size", - "pretty_assertions", - "sendfd", "serde", "spawn_worker", "tempfile", "tokio", "tracing", - "tracing-subscriber", "winapi 0.3.9", "windows-sys 0.48.0", "zwohash", @@ -3095,7 +3068,7 @@ dependencies = [ [[package]] name = "libdd-library-config" -version = "4.0.0" +version = "4.1.0" dependencies = [ "anyhow", "libc 0.2.186", @@ -3103,7 +3076,6 @@ dependencies = [ "memfd", "prost", "rand 0.8.8", - "rmp", "rmp-serde", "serde", "serde_yaml", @@ -3121,7 +3093,6 @@ dependencies = [ "libdd-common", "libdd-common-ffi", "libdd-library-config", - "tempfile", ] [[package]] @@ -3173,7 +3144,6 @@ dependencies = [ "serde_json", "tokio", "tokio-util", - "uuid", ] [[package]] @@ -3196,24 +3166,20 @@ dependencies = [ "bitmaps", "bolero", "byteorder", - "bytes", "chrono", "criterion", "crossbeam-channel", "crossbeam-utils", "cxx", "cxx-build", - "futures", "hashbrown 0.17.1", "http 1.4.2", - "http-body-util", - "httparse", "indexmap 2.12.1", "libdd-alloc", "libdd-common", "libdd-profiling", "libdd-profiling-protobuf", - "mime", + "opentelemetry-proto", "parking_lot", "proptest", "prost", @@ -3245,7 +3211,7 @@ dependencies = [ [[package]] name = "libdd-remote-config" -version = "5.0.0" +version = "6.0.0" dependencies = [ "anyhow", "base64 0.22.1", @@ -3268,10 +3234,8 @@ dependencies = [ "manual_future", "prost", "rand 0.8.8", - "ring", "serde", "serde_json", - "serde_with", "sha2", "strum", "strum_macros", @@ -3285,7 +3249,7 @@ dependencies = [ [[package]] name = "libdd-shared-runtime" -version = "4.0.0" +version = "5.0.0" dependencies = [ "async-trait", "futures", @@ -3301,14 +3265,13 @@ dependencies = [ [[package]] name = "libdd-telemetry" -version = "8.0.0" +version = "9.0.0" dependencies = [ "anyhow", "async-trait", "base64 0.22.1", "bytes", "futures", - "getrandom 0.2.15", "hashbrown 0.17.1", "http 1.4.2", "httpmock", @@ -3338,7 +3301,6 @@ version = "0.0.1" dependencies = [ "build_common", "function_name", - "libc 0.2.186", "libdd-capabilities-impl", "libdd-common", "libdd-common-ffi", @@ -3350,7 +3312,7 @@ dependencies = [ [[package]] name = "libdd-tinybytes" -version = "1.1.3" +version = "1.1.4" dependencies = [ "libdd-tinybytes", "once_cell", @@ -3364,7 +3326,7 @@ dependencies = [ [[package]] name = "libdd-trace-normalization" -version = "4.0.0" +version = "4.1.0" dependencies = [ "anyhow", "arbitrary", @@ -3376,7 +3338,7 @@ dependencies = [ [[package]] name = "libdd-trace-obfuscation" -version = "8.0.0" +version = "9.0.0" dependencies = [ "anyhow", "criterion", @@ -3390,11 +3352,12 @@ dependencies = [ "percent-encoding", "serde", "serde_json", + "thiserror 2.0.18", ] [[package]] name = "libdd-trace-protobuf" -version = "5.0.0" +version = "5.0.1" dependencies = [ "bolero", "prost", @@ -3408,7 +3371,7 @@ dependencies = [ [[package]] name = "libdd-trace-stats" -version = "9.0.0" +version = "10.0.0" dependencies = [ "anyhow", "arc-swap", @@ -3432,14 +3395,13 @@ dependencies = [ "rmp-serde", "serde", "tokio", - "tokio-util", "tracing", "web-time", ] [[package]] name = "libdd-trace-utils" -version = "12.0.0" +version = "13.0.0" dependencies = [ "anyhow", "base64 0.22.1", @@ -3552,7 +3514,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc2f4eb4bc735547cfed7c0a4922cbd04a4655978c09b54f1f7b228750664c34" dependencies = [ "cfg-if", - "windows-targets 0.48.5", + "windows-targets 0.52.6", ] [[package]] @@ -3701,17 +3663,6 @@ dependencies = [ "windows-sys 0.52.0", ] -[[package]] -name = "microseh" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f26b2a7c5ccfb370edd57fda423f3a551516ee127e10bc22a6215e8c63b20a38" -dependencies = [ - "cc", - "libc 0.2.186", - "windows-sys 0.42.0", -] - [[package]] name = "mime" version = "0.3.17" @@ -4088,6 +4039,36 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff011a302c396a5197692431fc1948019154afc178baf7d8e37367442a4601cf" +[[package]] +name = "opentelemetry" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6cdb0b1b267eb9db3331b434ed9ddab10d50e280a9adf9d13e5233e2002b61b5" +dependencies = [ + "js-sys", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "25da1ac11a0aeccf38d7f77ee0348715adaf8340f65ad46c94a02c6b20e2f65d" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost", +] + +[[package]] +name = "opentelemetry_sdk" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb39533d9d1c912123efd7d41d7e0c29d16917b60ce15b4c8d87cb1af7f67520" +dependencies = [ + "opentelemetry", + "portable-atomic", +] + [[package]] name = "os_info" version = "3.14.0" @@ -4337,7 +4318,7 @@ version = "1.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "059a34f111a9dee2ce1ac2826a68b24601c4298cfeb1a587c3cb493d5ab46f52" dependencies = [ - "libc 0.1.12", + "libc 0.2.186", "nix 0.30.1", ] @@ -4491,7 +4472,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "343d3bd7056eda839b03204e68deff7d1b13aba7af2b2fd16890697274262ee7" dependencies = [ "heck 0.5.0", - "itertools 0.10.5", + "itertools 0.11.0", "log", "multimap", "petgraph", @@ -4510,7 +4491,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "27c6023962132f4b30eb4c172c91ce92d933da334c59c23cddee82358ddafb0b" dependencies = [ "anyhow", - "itertools 0.10.5", + "itertools 0.11.0", "proc-macro2", "quote", "syn 2.0.118", @@ -4584,9 +4565,9 @@ checksum = "7dc55d7dec32ecaf61e0bd90b3d2392d721a28b95cfd23c3e176eccefbeab2f2" [[package]] name = "pyo3" -version = "0.28.3" +version = "0.29.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "91fd8e38a3b50ed1167fb981cd6fd60147e091784c427b8f7183a7ee32c31c12" +checksum = "4688ddedf473e32662b9b067670129a8afb8c18e351482c70d62ba4a88171e8b" dependencies = [ "libc 0.2.186", "once_cell", @@ -4598,18 +4579,18 @@ dependencies = [ [[package]] name = "pyo3-build-config" -version = "0.28.3" +version = "0.29.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e368e7ddfdeb98c9bca7f8383be1648fd84ab466bf2bc015e94008db6d35611e" +checksum = "f41027e41b4bd03f6e60f9f417fe24a6341a6bb744edd62b6f709f2a52ea30e9" dependencies = [ "target-lexicon", ] [[package]] name = "pyo3-ffi" -version = "0.28.3" +version = "0.29.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7f29e10af80b1f7ccaf7f69eace800a03ecd13e883acfacc1e5d0988605f651e" +checksum = "e591a95526fead067432c3b3a33fc74770b87b1e04e73671090d9c2055a2b327" dependencies = [ "libc 0.2.186", "pyo3-build-config", @@ -4617,9 +4598,9 @@ dependencies = [ [[package]] name = "pyo3-macros" -version = "0.28.3" +version = "0.29.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df6e520eff47c45997d2fc7dd8214b25dd1310918bbb2642156ef66a67f29813" +checksum = "73225868fc1cd84eef2c3c230ddb91273bf1de46aeb8a4248da76d32a0924a1c" dependencies = [ "proc-macro2", "pyo3-macros-backend", @@ -4629,13 +4610,12 @@ dependencies = [ [[package]] name = "pyo3-macros-backend" -version = "0.28.3" +version = "0.29.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c4cdc218d835738f81c2338f822078af45b4afdf8b2e33cbb5916f108b813acb" +checksum = "571575aa3749fa6216757dd47d2a3e7ef360f329a40f0666a9fbd14889024952" dependencies = [ "heck 0.5.0", "proc-macro2", - "pyo3-build-config", "quote", "syn 2.0.118", ] @@ -5326,16 +5306,6 @@ dependencies = [ "serde", ] -[[package]] -name = "sendfd" -version = "0.4.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "604b71b8fc267e13bb3023a2c901126c8f349393666a6d98ac1ae5729b701798" -dependencies = [ - "libc 0.2.186", - "tokio", -] - [[package]] name = "serde" version = "1.0.228" @@ -5346,15 +5316,6 @@ dependencies = [ "serde_derive", ] -[[package]] -name = "serde-bool" -version = "0.1.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8fdd050c9c2ed5ae1fb29e71be0a6efdd9df43c7cb13ea5826528cfe10c51db0" -dependencies = [ - "serde", -] - [[package]] name = "serde-transcode" version = "1.1.1" @@ -5419,7 +5380,6 @@ version = "1.0.150" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9" dependencies = [ - "indexmap 2.12.1", "itoa 1.0.14", "memchr", "serde", @@ -6792,7 +6752,7 @@ version = "0.1.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cf221c93e13a30d793f7645a0e7762c55d169dbb0a49671918a2319d289b10bb" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.59.0", ] [[package]] @@ -6904,21 +6864,6 @@ dependencies = [ "windows-link 0.1.1", ] -[[package]] -name = "windows-sys" -version = "0.42.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5a3e1820f08b8513f676f7ab6c1f99ff312fb97b553d30ff4dd86f9f15728aa7" -dependencies = [ - "windows_aarch64_gnullvm 0.42.2", - "windows_aarch64_msvc 0.42.2", - "windows_i686_gnu 0.42.2", - "windows_i686_msvc 0.42.2", - "windows_x86_64_gnu 0.42.2", - "windows_x86_64_gnullvm 0.42.2", - "windows_x86_64_msvc 0.42.2", -] - [[package]] name = "windows-sys" version = "0.45.0" diff --git a/Makefile b/Makefile index ad3e24d8d8e..8bece99bc2f 100644 --- a/Makefile +++ b/Makefile @@ -221,6 +221,14 @@ test_c2php: $(SO_FILE) $(INIT_HOOK_TEST_FILES) $(BUILD_DIR)/run-tests.php test_with_init_hook: $(SO_FILE) $(INIT_HOOK_TEST_FILES) $(BUILD_DIR)/run-tests.php $(if $(ASAN), USE_ZEND_ALLOC=0 USE_TRACKED_ALLOC=1) $(RUN_TESTS_CMD) -d extension=$(SO_FILE) $(TRACER_SOURCES_INI) $(INIT_HOOK_TEST_FILES); +# Exercise the Linux-only raw-clone lifetime gate in the normal CI pass. +ifeq ($(shell uname -s),Linux) +test_extension_ci_normal: test_signal_flush +endif + +test_signal_flush: + bash tests/signal_flush/run.sh + # The .phpt suite runs twice: normally, and under valgrind for leak checking. # Separate targets so CI can parallelize them -- valgrind is far slower. # The PATH shim in tests/ext/valgrind adds the suppressions file. @@ -1706,6 +1714,6 @@ test_internal_api_randomized: $(SO_FILE) composer.lock: composer.json $(call run_composer_with_retry,,) -.PHONY: dev dist_clean clean cores all clang_format_check clang_format_fix install sudo_install test_c test_c_mem test_extension_ci test_extension_ci_normal test_extension_ci_valgrind test_zai test_zai_asan test install_ini install_all \ +.PHONY: dev dist_clean clean cores all clang_format_check clang_format_fix install sudo_install test_c test_c_mem test_extension_ci test_extension_ci_normal test_extension_ci_valgrind test_signal_flush test_zai test_zai_asan test install_ini install_all \ .apk .rpm .deb .tar.gz sudo debug prod strict run-tests.php verify_pecl_file_definitions verify_package_xml cbindgen cbindgen_binary \ compile_profiler install_profiler compile_profiler_asan diff --git a/appsec/tests/integration/src/test/www/signal-flush/run.sh b/appsec/tests/integration/src/test/www/signal-flush/run.sh new file mode 100755 index 00000000000..cb7e4ac5c57 --- /dev/null +++ b/appsec/tests/integration/src/test/www/signal-flush/run.sh @@ -0,0 +1,21 @@ +#!/bin/bash -e + +set -x + +mkdir -p /tmp/logs +logs=( + /tmp/logs/appsec.log + /tmp/logs/helper.log + /tmp/logs/php_error.log + /tmp/logs/php_server.log + /tmp/logs/sidecar.log +) +touch "${logs[@]}" + +enable_extensions.sh +echo datadog.trace.cli_enabled=true >> /etc/php/php.ini + +php -S 0.0.0.0:80 -t /var/www/public \ + >> /tmp/logs/php_server.log 2>&1 & + +exec tail -n +1 -F "${logs[@]}" diff --git a/appsec/tests/integration/src/test/www/signal-flush/signal_flush_worker.php b/appsec/tests/integration/src/test/www/signal-flush/signal_flush_worker.php new file mode 100644 index 00000000000..c9e1866a99d --- /dev/null +++ b/appsec/tests/integration/src/test/www/signal-flush/signal_flush_worker.php @@ -0,0 +1,49 @@ + INSTALLING -> READY -> STARTING_* -> RUNNING_* -> STOPPED -// `-> FAILED -----------^ +// Initial publication: DISABLED -> INSTALLING -> READY +// Successful refresh: READY -> INSTALLING -> READY +// Failed refresh: READY -> INSTALLING -> DISABLED +// Signal claim: READY -> STARTING_* -> RUNNING_* -> STOPPED +// `-> FAILED -----------^ // // INSTALLING protects publication from teardown. STARTING protects clone setup. +// Reconnection may move READY back to INSTALLING; consumed objects never rearm. // The DEFAULT and CUSTOM variants preserve who is responsible for termination. enum { DD_SIGNAL_DISABLED, // No flush object is installed. @@ -353,11 +358,10 @@ bool datadog_signals_has_sidecar_flush(void) { return atomic_load_explicit(&dd_signal_state, memory_order_acquire) != DD_SIGNAL_DISABLED; } -void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush) { +void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush, bool replace) { // Consume `flush`. On success it becomes the process-wide signal flush // object; if it cannot be installed, destroy it before returning. - int expected = DD_SIGNAL_DISABLED; - if (!flush) { + if (!flush && !replace) { return; } // prevent a previous handler from suspending this publication @@ -369,17 +373,32 @@ void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush) { datadog_sidecar_signal_flush_drop(flush); return; } - if (!atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, DD_SIGNAL_INSTALLING, + int owner_pid = getpid(); + int expected = atomic_load_explicit(&dd_signal_state, memory_order_acquire); + bool can_install = expected == DD_SIGNAL_DISABLED || (replace && expected == DD_SIGNAL_READY); + if (!can_install || + !atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, DD_SIGNAL_INSTALLING, memory_order_acq_rel, memory_order_acquire)) { datadog_sidecar_signal_flush_drop(flush); sigprocmask(SIG_SETMASK, &old_signals, NULL); return; } + // Winning READY -> INSTALLING excludes the handler's READY -> STARTING CAS. + // If the handler won instead, the new object was dropped above and the + // worker retains the old object through its existing clear_child_tid join. + ddog_SignalFlush *previous = dd_signal_flush; dd_signal_flush = flush; - atomic_store_explicit(&dd_signal_owner_pid, getpid(), memory_order_relaxed); - // READY is the publication barrier for the pointer and owner PID. A handler - // can acquire READY only after both values have been initialized. - atomic_store_explicit(&dd_signal_state, DD_SIGNAL_READY, memory_order_release); + atomic_store_explicit(&dd_signal_owner_pid, owner_pid, memory_order_relaxed); + // READY publishes the pointer and owner PID. A failed refresh leaves no + // object and returns to DISABLED so normal request setup can retry later. + atomic_store_explicit(&dd_signal_state, flush ? DD_SIGNAL_READY : DD_SIGNAL_DISABLED, memory_order_release); + // A handler on another thread waits while the state is INSTALLING. Publish + // the final state before freeing `previous`: the signal may have interrupted + // that thread while it held the allocator lock. Freeing the object first + // could then block the publisher on that lock while the handler waits for the + // publisher, causing a deadlock. No worker can claim `previous` after our + // successful READY -> INSTALLING transition. + datadog_sidecar_signal_flush_drop(previous); // A pending SIGTERM/SIGINT may run as soon as the old mask is restored; at // that point it observes either the complete publication or a terminal state. sigprocmask(SIG_SETMASK, &old_signals, NULL); @@ -460,8 +479,19 @@ static void dd_sigint_sigterm_handler(int sig, siginfo_t *si, void *uc) { int starting = terminate ? DD_SIGNAL_STARTING_DEFAULT : DD_SIGNAL_STARTING_CUSTOM; // READY is consumed exactly once. Besides preventing duplicate workers, the // acquire half makes the prepared flush object and owner PID visible. - if (!atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, starting, memory_order_acq_rel, - memory_order_acquire)) { + while (!atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, starting, memory_order_acq_rel, + memory_order_acquire)) { + if (expected == DD_SIGNAL_INSTALLING && + datadog_raw_syscall6(SYS_getpid, 0, 0, 0, 0, 0, 0) == + atomic_load_explicit(&dd_signal_owner_pid, memory_order_relaxed)) { + // The publisher blocks SIGTERM/SIGINT on its own thread and only + // performs lock-free stores while INSTALLING. Let it finish before + // claiming the replacement. Never wait for a publisher inherited + // across fork: that thread does not exist in the child. + datadog_raw_syscall6(SYS_sched_yield, 0, 0, 0, 0, 0, 0); + expected = DD_SIGNAL_READY; + continue; + } // A worker handling a default-disposition signal will terminate the // process after flushing. Do not let another default signal iterrupt // the flush early by invoking the previous disposition. diff --git a/ext/signals.h b/ext/signals.h index bd29848ed6f..94111b5e809 100644 --- a/ext/signals.h +++ b/ext/signals.h @@ -10,7 +10,12 @@ void datadog_signals_first_rinit(void); void datadog_signals_minit(void); void datadog_signals_mshutdown(void); bool datadog_signals_has_sidecar_flush(void); -void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush); +// Always takes ownership of `flush`, including when it cannot be published. +// With `replace` false, publishes only when no flush object is installed. With +// `replace` true, replaces or clears an installed object unless a signal handler +// has already claimed it. Passing NULL while replacing clears an unclaimed stale +// object so setup can retry later. +void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush, bool replace); void datadog_signals_reset_sidecar_flush_after_fork(void); #endif // DD_TRACE_SIGNALS_H diff --git a/libdatadog b/libdatadog index 836ff60ac46..588a120ba82 160000 --- a/libdatadog +++ b/libdatadog @@ -1 +1 @@ -Subproject commit 836ff60ac46244268c6b62b5246132a08cf09512 +Subproject commit 588a120ba82e41311496f057dc2340d731d43abe From cef0b2442e95fd01b851f0c4650909efbfdd69b5 Mon Sep 17 00:00:00 2001 From: Bob Weinand Date: Mon, 28 Sep 2026 23:06:46 +0200 Subject: [PATCH 3/3] Refactor out into libdatadog --- Cargo.lock | 1 + Makefile | 10 +- components-rs/datadog.h | 33 +- components-rs/sidecar.h | 46 +++ components-rs/signal_flush.rs | 705 +--------------------------------- config.m4 | 7 +- ext/sidecar.c | 5 +- ext/signals.c | 313 +++++---------- ext/threads.c | 52 --- ext/threads.h | 9 - libdatadog | 2 +- 11 files changed, 169 insertions(+), 1014 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index e3d3a28f563..a44107be962 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2783,6 +2783,7 @@ name = "libdd-common" version = "7.0.0" dependencies = [ "anyhow", + "arc-swap", "bytes", "cc", "const_format", diff --git a/Makefile b/Makefile index 8bece99bc2f..ad3e24d8d8e 100644 --- a/Makefile +++ b/Makefile @@ -221,14 +221,6 @@ test_c2php: $(SO_FILE) $(INIT_HOOK_TEST_FILES) $(BUILD_DIR)/run-tests.php test_with_init_hook: $(SO_FILE) $(INIT_HOOK_TEST_FILES) $(BUILD_DIR)/run-tests.php $(if $(ASAN), USE_ZEND_ALLOC=0 USE_TRACKED_ALLOC=1) $(RUN_TESTS_CMD) -d extension=$(SO_FILE) $(TRACER_SOURCES_INI) $(INIT_HOOK_TEST_FILES); -# Exercise the Linux-only raw-clone lifetime gate in the normal CI pass. -ifeq ($(shell uname -s),Linux) -test_extension_ci_normal: test_signal_flush -endif - -test_signal_flush: - bash tests/signal_flush/run.sh - # The .phpt suite runs twice: normally, and under valgrind for leak checking. # Separate targets so CI can parallelize them -- valgrind is far slower. # The PATH shim in tests/ext/valgrind adds the suppressions file. @@ -1714,6 +1706,6 @@ test_internal_api_randomized: $(SO_FILE) composer.lock: composer.json $(call run_composer_with_retry,,) -.PHONY: dev dist_clean clean cores all clang_format_check clang_format_fix install sudo_install test_c test_c_mem test_extension_ci test_extension_ci_normal test_extension_ci_valgrind test_signal_flush test_zai test_zai_asan test install_ini install_all \ +.PHONY: dev dist_clean clean cores all clang_format_check clang_format_fix install sudo_install test_c test_c_mem test_extension_ci test_extension_ci_normal test_extension_ci_valgrind test_zai test_zai_asan test install_ini install_all \ .apk .rpm .deb .tar.gz sudo debug prod strict run-tests.php verify_pecl_file_definitions verify_package_xml cbindgen cbindgen_binary \ compile_profiler install_profiler compile_profiler_asan diff --git a/components-rs/datadog.h b/components-rs/datadog.h index aa677ce8674..446d6fd51be 100644 --- a/components-rs/datadog.h +++ b/components-rs/datadog.h @@ -9,14 +9,6 @@ struct _zend_string; #include "telemetry.h" #include "sidecar.h" -#if defined(__linux__) -/** - * Prepared, independent, sessionless connection. Opaque to C; immutable after publication. - * Only normal initialized threads may construct or destroy this object. - */ -typedef struct ddog_SignalFlush ddog_SignalFlush; -#endif - extern void (*ddog_log_callback)(ddog_CharSlice); extern ddog_VecRemoteConfigProduct DATADOG_REMOTE_CONFIG_PRODUCTS; @@ -283,27 +275,10 @@ void datadog_sidecar_clear_reconnect_fn(struct ddog_SidecarTransport **transport #if defined(__linux__) /** - * Connect independently to the template's exact listener and prepare one Flush request. - * The template is borrowed only during this call: neither it nor its fd is retained. - * The returned object must outlive the raw worker and may only be dropped normally. - */ -ddog_MaybeError datadog_sidecar_prepare_signal_flush(struct ddog_SidecarTransport *template_, - struct ddog_SignalFlush **output); -#endif - -#if defined(__linux__) -/** - * Destroy an unpublished object, or one whose raw worker has exited. Normal context only. - */ -void datadog_sidecar_signal_flush_drop(struct ddog_SignalFlush *flush); -#endif - -#if defined(__linux__) -/** - * Execute one bounded exchange without calling libc, using TLS, allocating, or unwinding. - * If `terminate_process` is true, terminate the process after the flush completes, fails, or - * times out. Otherwise, return zero for the one-byte zero ACK or a negative Linux error number. - * The caller must keep the object alive and guarantee exclusive, one-shot use of its socket. + * Execute the prepared flush. For a default signal disposition, terminate the process afterward. + * The object must remain alive until the raw worker exits; custom handlers use normal shutdown + * to join that worker. `_exit` terminates the whole process without runtime cleanup. + * The signal handler must reject inherited state before starting a worker after fork. */ int32_t datadog_sidecar_signal_flush_run(const struct ddog_SignalFlush *flush, bool terminate_process); diff --git a/components-rs/sidecar.h b/components-rs/sidecar.h index a2413af792a..6ba2fc8511a 100644 --- a/components-rs/sidecar.h +++ b/components-rs/sidecar.h @@ -692,4 +692,50 @@ void ddog_add_event_attributes_float(ddog_SpanEventBytes *event, ddog_CharSlice */ ddog_CharSlice ddog_serialize_trace_into_charslice(ddog_TraceBytes *trace); +#if defined(__linux__) +/** + * One pre-encoded flush. The duplicate fd preserves the template connection until drop; it + * shares packet ordering with normal sends and never reads the normal client's replies. + * Construction and destruction require ordinary thread context. + */ +typedef struct ddog_SignalFlush ddog_SignalFlush; +#endif + +#if defined(__linux__) +/** + * Prepare a flush on this transport with a private completion pipe. + * Normal thread context only. The returned object owns a duplicate of the transport fd; + * refresh it after reconnect and drop it before normal connection shutdown. + * + * # Safety + * `transport` must be exclusively borrowed and `output` must be writable for this call. + */ +ddog_MaybeError ddog_sidecar_prepare_signal_flush(struct ddog_SidecarTransport *transport, + struct ddog_SidecarFlushOptions options, + struct ddog_SignalFlush **output); +#endif + +#if defined(__linux__) +/** + * Destroy a prepared flush in ordinary thread context. + * + * # Safety + * `flush` must be null or an owned pointer returned by prepare. Any raw worker must have exited. + */ +void ddog_sidecar_signal_flush_drop(struct ddog_SignalFlush *flush); +#endif + +#if defined(__linux__) +/** + * Run one bounded flush without TLS access, allocation, unwinding, or process termination. + * Returns zero when the sidecar closes the pipe (completion or exit), or a negative Linux errno. + * + * # Safety + * The object must remain alive through the call, with exclusive one-shot use of this object. + * The normal transport may continue sending and receiving concurrently. + * All worker signals must be blocked. Do not use an inherited object after fork. + */ +int32_t ddog_sidecar_signal_flush_run(const struct ddog_SignalFlush *flush); +#endif + #endif /* DDOG_SIDECAR_H */ diff --git a/components-rs/signal_flush.rs b/components-rs/signal_flush.rs index 77a32ee45d3..6290766b9fb 100644 --- a/components-rs/signal_flush.rs +++ b/components-rs/signal_flush.rs @@ -1,707 +1,20 @@ -//! The raw-clone signal worker owns no libc or Rust thread state. `prepare` and `drop` run on -//! ordinary initialized threads; only `run` may execute in the worker. `run` uses ordinary Rust -//! control flow, but every operation in its generated call graph must reduce to local memory -//! access or a direct Linux system call. Changes to it therefore require a disassembly audit in -//! both optimized and debug builds. +//! PHP process policy around libdatadog's raw-thread-safe flush exchange. -use core::arch::asm; -use datadog_sidecar::service::blocking::SidecarTransport; -use datadog_sidecar::service::sidecar_interface::{SidecarFlushOptions, SidecarInterfaceRequest}; -use libdd_common_ffi::{Error, MaybeError}; -use libdd_ipc::SeqpacketConn; -use std::mem::{offset_of, size_of}; +use datadog_sidecar::service::signal_flush::SignalFlush; -/// Prepared, independent, sessionless connection. Opaque to C; immutable after publication. -/// Only normal initialized threads may construct or destroy this object. -pub struct SignalFlush { - // Keep the raw view free of owning Rust values: the worker only copies these four fields. - raw: RawSignalFlush, - // These owners make the fd and request pointer in `raw` valid until normal-context drop. - _connection: SeqpacketConn, - _request: Box<[u8]>, -} - -#[repr(C)] -struct RawSignalFlush { - fd: i32, - owner_pid: i32, - request: *const u8, - request_len: usize, -} - -// Direct syscalls use the 64-bit Linux kernel ABI, not libc's platform types. -#[repr(C)] -struct KernelTimespec { - seconds: i64, - nanoseconds: i64, -} - -// Match the pollfd layout consumed directly by the kernel. -#[repr(C)] -struct KernelPollFd { - fd: i32, - events: i16, - revents: i16, -} - -/// Connect independently to the template's exact listener and prepare one Flush request. -/// The template is borrowed only during this call: neither it nor its fd is retained. -/// The returned object must outlive the raw worker and may only be dropped normally. -#[no_mangle] -pub unsafe extern "C" fn datadog_sidecar_prepare_signal_flush( - template: *mut SidecarTransport, - output: *mut *mut SignalFlush, -) -> MaybeError { - // Everything that may allocate, lock, inspect libc state, or use the sidecar codec belongs in - // this function. None of those operations can be deferred to the raw worker. - let result = (|| { - let output = output - .as_mut() - .ok_or_else(|| anyhow::anyhow!("signal flush output is null"))?; - // Never leave the caller holding an old value when a later preparation step fails. - *output = std::ptr::null_mut(); - let template = template - .as_mut() - .ok_or_else(|| anyhow::anyhow!("signal flush template is null"))?; - // This is a new socket, not a dup of the template fd. A dup would keep the template's - // connection alive and interfere with the sidecar's disconnect detection. - let connection = connect_to_template(template)?; - // Encoding now leaves `run` with one immutable packet and no serializer call graph. - let request = libdd_ipc::codec::encode(&SidecarInterfaceRequest::Flush { - // Preserve the previous forced-shutdown behavior: traces and stats are the only data - // whose delivery this signal path delays termination to await. - options: SidecarFlushOptions { - traces_and_stats: true, - flag_evaluations: false, - telemetry: false, - }, - }) - .into_boxed_slice(); - - *output = Box::into_raw(Box::new(SignalFlush { - raw: RawSignalFlush { - fd: connection.as_raw_fd(), - // A copied object must never send on the parent's connection after fork. - owner_pid: libc::getpid(), - request: request.as_ptr(), - request_len: request.len(), - }, - _connection: connection, - _request: request, - })); - Ok::<_, anyhow::Error>(()) - })(); - match result { - Ok(()) => MaybeError::None, - Err(error) => MaybeError::Some(Error::from(format!("{error:#}"))), - } -} - -/// Destroy an unpublished object, or one whose raw worker has exited. Normal context only. -#[no_mangle] -pub unsafe extern "C" fn datadog_sidecar_signal_flush_drop(flush: *mut SignalFlush) { - // An unpublished object has no worker borrower. For a published object, C joins the worker - // through clear_child_tid before calling this, so neither owner can be destroyed during `run`. - if !flush.is_null() { - drop(Box::from_raw(flush)); - } -} - -/// Execute one bounded exchange without calling libc, using TLS, allocating, or unwinding. -/// If `terminate_process` is true, terminate the process after the flush completes, fails, or -/// times out. Otherwise, return zero for the one-byte zero ACK or a negative Linux error number. -/// The caller must keep the object alive and guarantee exclusive, one-shot use of its socket. -#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))] +/// Execute the prepared flush. For a default signal disposition, terminate the process afterward. +/// The object must remain alive until the raw worker exits; custom handlers use normal shutdown +/// to join that worker. `_exit` terminates the whole process without runtime cleanup. +/// The signal handler must reject inherited state before starting a worker after fork. #[no_mangle] #[inline(never)] -// An `extern "C"` definition gets a `panic_cannot_unwind` path in panic=unwind debug builds. -// `C-unwind` avoids that generated call. The body is deliberately written so it cannot unwind. pub unsafe extern "C-unwind" fn datadog_sidecar_signal_flush_run( - flush: *const SignalFlush, + flush: &SignalFlush, terminate_process: bool, ) -> i32 { - let result = signal_flush_exchange(flush); + let result = flush.run(); if terminate_process { - // The raw clone has no usable libc thread state. Terminate through an inline syscall - // rather than returning to C or calling libc's `_exit`/`exit` machinery. - signal_safe_exit_group(); + libc::_exit(0); } result } - -#[cfg(any(target_arch = "x86_64", target_arch = "aarch64"))] -#[inline(never)] -unsafe fn signal_flush_exchange(flush: *const SignalFlush) -> i32 { - // Raw Linux syscalls return `-errno` directly. They do not set libc's thread-local `errno`. - const EINTR: isize = 4; - const EAGAIN: isize = 11; - const CLOCK_MONOTONIC: usize = 1; - const MSG_DONTWAIT: usize = 0x40; - const MSG_NOSIGNAL: usize = 0x4000; - const MSG_TRUNC: usize = 0x20; - const POLLIN: i16 = 1; - const POLLOUT: i16 = 4; - - // Syscall numbers are architecture-specific, so keep them beside the audited inline syscall - // rather than routing through libc wrappers. - #[cfg(target_arch = "aarch64")] - const SYS_GETPID: usize = 172; - #[cfg(target_arch = "aarch64")] - const SYS_SENDTO: usize = 206; - #[cfg(target_arch = "aarch64")] - const SYS_RECVFROM: usize = 207; - #[cfg(target_arch = "aarch64")] - const SYS_CLOCK_GETTIME: usize = 113; - #[cfg(target_arch = "aarch64")] - const SYS_PPOLL: usize = 73; - - #[cfg(target_arch = "x86_64")] - const SYS_GETPID: usize = 39; - #[cfg(target_arch = "x86_64")] - const SYS_SENDTO: usize = 44; - #[cfg(target_arch = "x86_64")] - const SYS_RECVFROM: usize = 45; - #[cfg(target_arch = "x86_64")] - const SYS_CLOCK_GETTIME: usize = 228; - #[cfg(target_arch = "x86_64")] - const SYS_PPOLL: usize = 271; - - let (fd, owner_pid, request, request_len) = signal_safe_load(flush); - // `fork()` copies the object and fd but changes the process identity. Reject that copy even if - // a caller reaches this API before the C-side post-fork reset. - if signal_safe_syscall6(SYS_GETPID, 0, 0, 0, 0, 0, 0) != owner_pid as isize { - return -libc::ECHILD; - } - - let mut now = KernelTimespec { - seconds: 0, - nanoseconds: 0, - }; - let result = signal_safe_syscall6( - SYS_CLOCK_GETTIME, - CLOCK_MONOTONIC, - (&raw mut now).cast::() as usize, - 0, - 0, - 0, - 0, - ); - if result < 0 { - return result as i32; - } - // Use one absolute deadline for send, backpressure, and ACK. Restarting a relative timeout - // after EINTR or EAGAIN could otherwise keep process termination alive indefinitely. - // Wrapping arithmetic is intentional: debug overflow checks would introduce panic paths. - let deadline_seconds = now.seconds.wrapping_add(10); - let deadline_nanoseconds = now.nanoseconds; - // `receive` selects the protocol phase. `poll` records that the last nonblocking operation - // returned EAGAIN and must wait for readiness before retrying. - let mut receive = false; - let mut poll = false; - - loop { - let result = signal_safe_syscall6( - SYS_CLOCK_GETTIME, - CLOCK_MONOTONIC, - (&raw mut now).cast::() as usize, - 0, - 0, - 0, - 0, - ); - if result < 0 { - return result as i32; - } - let mut remaining_seconds = deadline_seconds.wrapping_sub(now.seconds); - let mut remaining_nanoseconds = deadline_nanoseconds.wrapping_sub(now.nanoseconds); - if remaining_nanoseconds < 0 { - remaining_nanoseconds = remaining_nanoseconds.wrapping_add(1_000_000_000); - remaining_seconds = remaining_seconds.wrapping_sub(1); - } - if remaining_seconds < 0 || (remaining_seconds == 0 && remaining_nanoseconds == 0) { - return -libc::ETIMEDOUT; - } - let mut remaining = KernelTimespec { - seconds: remaining_seconds, - nanoseconds: remaining_nanoseconds, - }; - - if poll { - let mut pollfd = KernelPollFd { - fd, - events: if receive { POLLIN } else { POLLOUT }, - revents: 0, - }; - let result = signal_safe_syscall6( - SYS_PPOLL, - (&raw mut pollfd).cast::() as usize, - 1, - (&raw mut remaining).cast::() as usize, - 0, - // The kernel's 64-bit signal-set size is eight bytes. No mask is supplied, but - // passing the ABI size keeps this a valid direct ppoll syscall on both targets. - 8, - 0, - ); - if result > 0 || result == -EINTR { - poll = false; - continue; - } - if result == 0 { - return -libc::ETIMEDOUT; - } - return result as i32; - } - - let result = if receive { - let mut ack = 1u8; - let result = signal_safe_syscall6( - SYS_RECVFROM, - fd as usize, - (&raw mut ack) as usize, - 1, - // MSG_TRUNC makes a packet larger than the one-byte buffer report its full size, - // so only an exact one-byte zero ACK can be accepted. - MSG_DONTWAIT | MSG_TRUNC, - 0, - 0, - ); - if result == 1 { - return if ack == 0 { 0 } else { -libc::EPROTO }; - } - result - } else { - let result = signal_safe_syscall6( - SYS_SENDTO, - fd as usize, - request as usize, - request_len, - // Nonblocking I/O lets the single ppoll deadline bound backpressure. MSG_NOSIGNAL - // prevents a closed sidecar socket from delivering SIGPIPE to the raw worker. - MSG_DONTWAIT | MSG_NOSIGNAL, - 0, - 0, - ); - // SOCK_SEQPACKET preserves message boundaries: a positive short send is a protocol - // failure rather than progress that can be resumed with a pointer offset. - if result == request_len as isize { - receive = true; - continue; - } - result - }; - - if result == -EINTR { - continue; - } - if result == -EAGAIN { - // Poll only after the kernel reports backpressure; a readiness wakeup retries the - // original send or receive operation and recomputes the remaining absolute deadline. - poll = true; - continue; - } - if result < 0 { - return result as i32; - } - // Any other positive result is an impossible short packet/send for this protocol. - return -libc::EPROTO; - } -} - -#[cfg(target_arch = "x86_64")] -#[inline(always)] -unsafe fn signal_safe_exit_group() -> ! { - // SYS_exit_group(0). The syscall cannot return, so no register clobbers remain live. - asm!( - "syscall", - in("rax") 231usize, - in("rdi") 0usize, - options(noreturn, nostack), - ); -} - -#[cfg(target_arch = "x86_64")] -#[inline(always)] -unsafe fn signal_safe_syscall6( - number: usize, - arg1: usize, - arg2: usize, - arg3: usize, - arg4: usize, - arg5: usize, - arg6: usize, -) -> isize { - let result: isize; - // Linux x86_64 uses r10, not the C ABI's rcx, for argument four. `syscall` itself clobbers rcx - // and r11; declaring both prevents the surrounding Rust from keeping live values there. - asm!( - "syscall", - inlateout("rax") number as isize => result, - in("rdi") arg1, - in("rsi") arg2, - in("rdx") arg3, - in("r10") arg4, - in("r8") arg5, - in("r9") arg6, - lateout("rcx") _, - lateout("r11") _, - options(nostack, preserves_flags), - ); - result -} - -#[cfg(target_arch = "x86_64")] -#[inline(always)] -unsafe fn signal_safe_load(flush: *const SignalFlush) -> (i32, i32, *const u8, usize) { - let fd: i32; - let owner_pid: i32; - let request: *const u8; - let request_len: usize; - // Do not replace this with `&(*flush).raw`. In panic=unwind debug builds that reference - // construction emitted null/alignment checks which call Rust panic handlers. Those handlers - // may allocate, use TLS, or otherwise depend on runtime state absent from the raw clone. - // These fixed-offset loads were audited to compile to four `mov` instructions and no calls. - asm!( - "movl {fd_offset}({base}), {fd:e}", - "movl {pid_offset}({base}), {owner_pid:e}", - "movq {request_offset}({base}), {request}", - "movq {length_offset}({base}), {request_len}", - base = in(reg) flush, - fd = out(reg) fd, - owner_pid = out(reg) owner_pid, - request = out(reg) request, - request_len = out(reg) request_len, - fd_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, fd), - pid_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, owner_pid), - request_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request), - length_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request_len), - options(att_syntax, nostack, readonly, preserves_flags), - ); - (fd, owner_pid, request, request_len) -} - -#[cfg(target_arch = "aarch64")] -#[inline(always)] -unsafe fn signal_safe_exit_group() -> ! { - // SYS_exit_group(0). The syscall cannot return, so no register clobbers remain live. - asm!( - "svc #0", - in("x8") 94usize, - in("x0") 0usize, - options(noreturn, nostack), - ); -} - -#[cfg(target_arch = "aarch64")] -#[inline(always)] -unsafe fn signal_safe_syscall6( - number: usize, - arg1: usize, - arg2: usize, - arg3: usize, - arg4: usize, - arg5: usize, - arg6: usize, -) -> isize { - let result: isize; - // Linux AArch64 takes the syscall number in x8, arguments in x0..x5, and returns in x0. - // Listing every register explicitly keeps the compiler from generating a helper call. - asm!( - "svc #0", - in("x8") number, - inlateout("x0") arg1 as isize => result, - in("x1") arg2, - in("x2") arg3, - in("x3") arg4, - in("x4") arg5, - in("x5") arg6, - options(nostack), - ); - result -} - -#[cfg(target_arch = "aarch64")] -#[inline(always)] -unsafe fn signal_safe_load(flush: *const SignalFlush) -> (i32, i32, *const u8, usize) { - let fd: i32; - let owner_pid: i32; - let request: *const u8; - let request_len: usize; - // As on x86_64, an ordinary Rust reference introduced panic-handler calls in debug builds. - // These fixed-offset loads were audited to compile to four `ldr` instructions and no calls. - asm!( - "ldr {fd:w}, [{base}, #{fd_offset}]", - "ldr {owner_pid:w}, [{base}, #{pid_offset}]", - "ldr {request}, [{base}, #{request_offset}]", - "ldr {request_len}, [{base}, #{length_offset}]", - base = in(reg) flush, - fd = out(reg) fd, - owner_pid = out(reg) owner_pid, - request = out(reg) request, - request_len = out(reg) request_len, - fd_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, fd), - pid_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, owner_pid), - request_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request), - length_offset = const offset_of!(SignalFlush, raw) + offset_of!(RawSignalFlush, request_len), - options(nostack, readonly, preserves_flags), - ); - (fd, owner_pid, request, request_len) -} - -// The clone trampoline also rejects unsupported architectures. Never substitute libc here. -/// cbindgen:ignore -#[cfg(not(any(target_arch = "x86_64", target_arch = "aarch64")))] -#[no_mangle] -pub unsafe extern "C" fn datadog_sidecar_signal_flush_run( - _flush: *const SignalFlush, - _terminate_process: bool, -) -> i32 { - -libc::ENOTSUP -} - -fn connect_to_template(template: &mut SidecarTransport) -> anyhow::Result { - // SidecarTransport intentionally hides its endpoint. getpeername recovers the connected - // listener address without borrowing or retaining the template transport itself. - let mut address: libc::sockaddr_un = unsafe { std::mem::zeroed() }; - let mut length = size_of::() as libc::socklen_t; - if unsafe { - libc::getpeername( - template.as_raw_fd(), - (&mut address as *mut libc::sockaddr_un).cast(), - &mut length, - ) - } != 0 - { - return Err(std::io::Error::last_os_error().into()); - } - let start = offset_of!(libc::sockaddr_un, sun_path); - let length = length as usize; - anyhow::ensure!( - address.sun_family == libc::AF_UNIX as libc::sa_family_t - && length > start + 1 - && length <= size_of::() - && address.sun_path[0] == 0, - "signal flush template does not name an abstract Unix listener" - ); - // Linux abstract names are length-delimited and may contain embedded NULs. Use the sockaddr - // length returned by the kernel rather than treating sun_path as a C string. - let name = unsafe { - std::slice::from_raw_parts(address.sun_path.as_ptr().add(1).cast(), length - start - 1) - }; - Ok(SeqpacketConn::connect_abstract(name)?) -} - -#[cfg(test)] -mod tests { - use super::*; - use libdd_ipc::SeqpacketListener; - use std::sync::atomic::{AtomicUsize, Ordering}; - use std::time::{Duration, Instant}; - - fn pair() -> (SignalFlush, SeqpacketConn) { - let (connection, peer) = SeqpacketConn::socketpair().unwrap(); - let request = vec![42u8].into_boxed_slice(); - let flush = SignalFlush { - raw: RawSignalFlush { - fd: connection.as_raw_fd(), - owner_pid: unsafe { libc::getpid() }, - request: request.as_ptr(), - request_len: request.len(), - }, - _connection: connection, - _request: request, - }; - (flush, peer) - } - - fn send(peer: &SeqpacketConn, bytes: &[u8]) { - assert_eq!( - unsafe { libc::send(peer.as_raw_fd(), bytes.as_ptr().cast(), bytes.len(), 0) }, - bytes.len() as isize - ); - } - - #[test] - fn accepts_only_an_exact_zero_ack_packet() { - for (ack, expected) in [ - (&[0][..], 0), - (&[0, 0][..], -libc::EPROTO), - (&[1][..], -libc::EPROTO), - ] { - let (flush, peer) = pair(); - send(&peer, ack); - assert_eq!( - unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, - expected - ); - } - } - - #[test] - fn waits_for_the_ack_after_sending() { - let (flush, peer) = pair(); - let server = std::thread::spawn(move || { - let mut pollfd = libc::pollfd { - fd: peer.as_raw_fd(), - events: libc::POLLIN, - revents: 0, - }; - assert_eq!(unsafe { libc::poll(&mut pollfd, 1, 1000) }, 1); - let mut request = 0u8; - assert_eq!( - unsafe { libc::recv(peer.as_raw_fd(), (&mut request as *mut u8).cast(), 1, 0) }, - 1 - ); - assert_eq!(request, 42); - std::thread::sleep(Duration::from_millis(20)); - send(&peer, &[0]); - }); - assert_eq!( - unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, - 0 - ); - server.join().unwrap(); - } - - #[test] - fn retries_send_after_backpressure() { - let (flush, peer) = pair(); - let filler = [0u8; 512]; - loop { - let sent = unsafe { - libc::send( - flush.raw.fd, - filler.as_ptr().cast(), - filler.len(), - libc::MSG_DONTWAIT | libc::MSG_NOSIGNAL, - ) - }; - if sent < 0 { - assert_eq!( - std::io::Error::last_os_error().raw_os_error(), - Some(libc::EAGAIN) - ); - break; - } - } - let server = std::thread::spawn(move || { - std::thread::sleep(Duration::from_millis(20)); - let mut bytes = [0u8; 512]; - loop { - let mut pollfd = libc::pollfd { - fd: peer.as_raw_fd(), - events: libc::POLLIN, - revents: 0, - }; - assert_eq!(unsafe { libc::poll(&mut pollfd, 1, 1000) }, 1); - let received = unsafe { - libc::recv(peer.as_raw_fd(), bytes.as_mut_ptr().cast(), bytes.len(), 0) - }; - assert!(received > 0); - if received == 1 && bytes[0] == 42 { - break; - } - assert_eq!(received, bytes.len() as isize); - } - send(&peer, &[0]); - }); - assert_eq!( - unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, - 0 - ); - server.join().unwrap(); - } - - #[test] - fn rejects_a_closed_peer_without_sigpipe() { - let (flush, peer) = pair(); - drop(peer); - assert!(unsafe { datadog_sidecar_signal_flush_run(&flush, false) } < 0); - } - - #[test] - fn refuses_inherited_objects() { - let (mut flush, _peer) = pair(); - flush.raw.owner_pid += 1; - assert_eq!( - unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, - -libc::ECHILD - ); - } - - #[test] - fn missing_ack_has_one_ten_second_deadline() { - let (flush, _peer) = pair(); - let start = Instant::now(); - assert_eq!( - unsafe { datadog_sidecar_signal_flush_run(&flush, false) }, - -libc::ETIMEDOUT - ); - assert!(start.elapsed() >= Duration::from_secs(9)); - assert!(start.elapsed() < Duration::from_secs(15)); - } - - #[test] - fn preparation_creates_an_independent_connection_to_the_template_listener() { - static NEXT: AtomicUsize = AtomicUsize::new(0); - let name = format!( - "ddtrace-signal-flush-{}-{}", - std::process::id(), - NEXT.fetch_add(1, Ordering::Relaxed) - ); - let listener = SeqpacketListener::bind_abstract(name.as_bytes()).unwrap(); - let mut template = - SidecarTransport::from(SeqpacketConn::connect_abstract(name.as_bytes()).unwrap()); - let template_peer = listener.try_accept().unwrap(); - let mut output = std::ptr::null_mut(); - assert!(matches!( - unsafe { datadog_sidecar_prepare_signal_flush(&mut template, &mut output) }, - MaybeError::None - )); - let flush = unsafe { Box::from_raw(output) }; - let signal_peer = listener.try_accept().unwrap(); - assert_ne!(flush.raw.fd, template.as_raw_fd()); - - // Closing the normal connection must produce EOF despite the signal connection. - drop(template); - let mut bytes = [0u8; 256]; - assert_eq!( - unsafe { - libc::recv( - template_peer.as_raw_fd(), - bytes.as_mut_ptr().cast(), - bytes.len(), - 0, - ) - }, - 0 - ); - send(&signal_peer, &[0]); - assert_eq!( - unsafe { datadog_sidecar_signal_flush_run(&*flush, false) }, - 0 - ); - let received = unsafe { - libc::recv( - signal_peer.as_raw_fd(), - bytes.as_mut_ptr().cast(), - bytes.len(), - 0, - ) - }; - assert_eq!(received as usize, flush._request.len()); - assert_eq!(&bytes[..received as usize], &*flush._request); - drop(flush); - assert_eq!( - unsafe { - libc::recv( - signal_peer.as_raw_fd(), - bytes.as_mut_ptr().cast(), - bytes.len(), - 0, - ) - }, - 0 - ); - } -} diff --git a/config.m4 b/config.m4 index bbd6b8a03c6..5630a7c7d20 100644 --- a/config.m4 +++ b/config.m4 @@ -419,8 +419,13 @@ if test "$PHP_DDTRACE" != "no" && test "$PHP_DDTRACE_PROFILING" = "no"; then PHP_CHECK_LIBRARY(rt, shm_open, [EXTRA_LDFLAGS="$EXTRA_LDFLAGS -lrt"; DDTRACE_SHARED_LIBADD="${DDTRACE_SHARED_LIBADD:-} -lrt"]) - dnl rust imports these, so we need them to link + dnl Platform linker requirements for the Rust library case $host_os in + linux*) + dnl The signal worker calls _exit with shared TLS. Resolve libc symbols + dnl when loading the extension, including when Rust is linked as a static archive. + EXTRA_LDFLAGS="$EXTRA_LDFLAGS -Wl,-z,now" + ;; darwin*) EXTRA_LDFLAGS="$EXTRA_LDFLAGS -framework CoreFoundation -framework Security" PHP_ADD_FRAMEWORK([CoreFoundation]) diff --git a/ext/sidecar.c b/ext/sidecar.c index 9c7a712c423..37dc5e9398d 100644 --- a/ext/sidecar.c +++ b/ext/sidecar.c @@ -326,8 +326,9 @@ static void dd_sidecar_setup_signal_transport(ddog_SidecarTransport *transport, } ddog_SignalFlush *flush = NULL; - bool prepared = datadog_ffi_try("Failed preparing signal-only sidecar connection", - datadog_sidecar_prepare_signal_flush(transport, &flush)); + bool prepared = datadog_ffi_try("Failed preparing sidecar signal flush", + ddog_sidecar_prepare_signal_flush( + transport, (ddog_SidecarFlushOptions){.traces_and_stats = true}, &flush)); if (prepared || replace) { // Takes ownership, including when another normal thread published first. // A failed refresh clears the stale object so a later RINIT can retry. diff --git a/ext/signals.c b/ext/signals.c index 045bec91c7b..45bc57389f7 100644 --- a/ext/signals.c +++ b/ext/signals.c @@ -46,6 +46,7 @@ #endif #if __linux +#include #include #include #include @@ -303,279 +304,164 @@ static struct sigaction dd_sigint_sigterm_sigaction; static struct sigaction dd_sigterm_prev_sigaction; static struct sigaction dd_sigint_prev_sigaction; -// The cleanup stack and prepared flush object are allocated in ordinary -// context. A READY object may be replaced after reconnect, until exactly one -// raw worker claims it. The object is immutable for the worker's lifetime. -static char *dd_signal_cleanup_stack; -static size_t dd_signal_cleanup_stack_size; -static ddog_SignalFlush *dd_signal_flush; -// The process ID, rather than a thread ID, lets every PHP thread in this process -// pass while rejecting a copy inherited by a fork child. -static _Atomic(int) dd_signal_owner_pid; -// CLONE_PARENT_SETTID writes the worker TID here. CLONE_CHILD_CLEARTID clears it -// and performs a futex wake at thread exit, giving normal shutdown a join -// primitive that does not depend on pthread state. -static _Atomic(int) dd_signal_worker_tid; - -// A READY object may be replaced until a signal handler claims it. Claiming an -// object is one-shot: -// -// Initial publication: DISABLED -> INSTALLING -> READY -// Successful refresh: READY -> INSTALLING -> READY -// Failed refresh: READY -> INSTALLING -> DISABLED -// Signal claim: READY -> STARTING_* -> RUNNING_* -> STOPPED -// `-> FAILED -----------^ -// -// INSTALLING protects publication from teardown. STARTING protects clone setup. -// Reconnection may move READY back to INSTALLING; consumed objects never rearm. -// The DEFAULT and CUSTOM variants preserve who is responsible for termination. +// The request is prepared in ordinary context. Publication, signal delivery, +// and shutdown compete for this gate; a claimed request is never replaced or rearmed. enum { - DD_SIGNAL_DISABLED, // No flush object is installed. - DD_SIGNAL_INSTALLING, // A flush object is being published. - DD_SIGNAL_READY, // One signal may claim the published object. - DD_SIGNAL_STARTING_DEFAULT, // clone() is in progress; the worker will terminate the process. - DD_SIGNAL_STARTING_CUSTOM, // clone() is in progress; the worker will return after flushing. - DD_SIGNAL_RUNNING_DEFAULT, // The TID is joinable; the worker will terminate the process. - DD_SIGNAL_RUNNING_CUSTOM, // The TID is joinable; the worker will return after flushing. - DD_SIGNAL_FAILED, // Worker creation failed; another signal may not retry it. - DD_SIGNAL_STOPPED, // Module shutdown has disabled new workers. + DD_SIGNAL_DISABLED, + DD_SIGNAL_INSTALLING, + DD_SIGNAL_READY, + DD_SIGNAL_FLUSHING_DEFAULT, + DD_SIGNAL_FLUSHING_CUSTOM, + DD_SIGNAL_STOPPED, }; static _Atomic(int) dd_signal_state; -// Atomics used by a signal handler must compile to instructions; a libatomic -// fallback could lock or touch runtime state while the interrupted thread owns it. -// The TID object is also written directly by the kernel, whose clone/futex ABI -// requires the ordinary int representation. +static _Atomic(int) dd_signal_owner_pid; +static ddog_SignalFlush *dd_signal_flush; +// Only one worker can use this stack. Its storage lives as long as the extension. +static _Alignas(16) char dd_signal_cleanup_stack[MIN_STACKSZ]; + +// Before READY is published: -1 (not started). clone's PARENT_SETTID writes a positive +// TID; CHILD_CLEARTID clears it and wakes futex waiters when the worker has fully exited. +// A failed start clears it explicitly. This also covers a child exiting before clone returns. +static _Atomic(int) dd_signal_worker_tid; _Static_assert(ATOMIC_INT_LOCK_FREE == 2, "signal atomics must be lock-free"); _Static_assert(sizeof(dd_signal_worker_tid) == sizeof(int), "signal TID must have the kernel int layout"); -static void dd_signals_init_cleanup_stack(void); -static void dd_signals_drop_sidecar_flush(void); -static void dd_call_prev_handler(int sig, siginfo_t *si, void *uc); - bool datadog_signals_has_sidecar_flush(void) { - // Any state other than DISABLED owns or is publishing an object. In - // particular, FAILED and STOPPED must not allow a second publication. - return atomic_load_explicit(&dd_signal_state, memory_order_acquire) != DD_SIGNAL_DISABLED; + return atomic_load(&dd_signal_state) != DD_SIGNAL_DISABLED; } void datadog_signals_set_sidecar_flush(ddog_SignalFlush *flush, bool replace) { - // Consume `flush`. On success it becomes the process-wide signal flush - // object; if it cannot be installed, destroy it before returning. if (!flush && !replace) { return; } - // prevent a previous handler from suspending this publication + // A handler must not suspend its own publisher while INSTALLING. Other threads + // may spin on that state, so the protected section must contain no blocking calls. sigset_t publication_signals, old_signals; - sigemptyset(&publication_signals); - sigaddset(&publication_signals, SIGTERM); - sigaddset(&publication_signals, SIGINT); + sigfillset(&publication_signals); if (sigprocmask(SIG_BLOCK, &publication_signals, &old_signals) < 0) { - datadog_sidecar_signal_flush_drop(flush); + ddog_sidecar_signal_flush_drop(flush); return; } - int owner_pid = getpid(); - int expected = atomic_load_explicit(&dd_signal_state, memory_order_acquire); - bool can_install = expected == DD_SIGNAL_DISABLED || (replace && expected == DD_SIGNAL_READY); - if (!can_install || - !atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, DD_SIGNAL_INSTALLING, - memory_order_acq_rel, memory_order_acquire)) { - datadog_sidecar_signal_flush_drop(flush); - sigprocmask(SIG_SETMASK, &old_signals, NULL); - return; - } - // Winning READY -> INSTALLING excludes the handler's READY -> STARTING CAS. - // If the handler won instead, the new object was dropped above and the - // worker retains the old object through its existing clear_child_tid join. - ddog_SignalFlush *previous = dd_signal_flush; - dd_signal_flush = flush; - atomic_store_explicit(&dd_signal_owner_pid, owner_pid, memory_order_relaxed); - // READY publishes the pointer and owner PID. A failed refresh leaves no - // object and returns to DISABLED so normal request setup can retry later. - atomic_store_explicit(&dd_signal_state, flush ? DD_SIGNAL_READY : DD_SIGNAL_DISABLED, memory_order_release); - // A handler on another thread waits while the state is INSTALLING. Publish - // the final state before freeing `previous`: the signal may have interrupted - // that thread while it held the allocator lock. Freeing the object first - // could then block the publisher on that lock while the handler waits for the - // publisher, causing a deadlock. No worker can claim `previous` after our - // successful READY -> INSTALLING transition. - datadog_sidecar_signal_flush_drop(previous); - // A pending SIGTERM/SIGINT may run as soon as the old mask is restored; at - // that point it observes either the complete publication or a terminal state. + int state = atomic_load(&dd_signal_state); + if ((state == DD_SIGNAL_DISABLED || (replace && state == DD_SIGNAL_READY)) && + atomic_compare_exchange_strong(&dd_signal_state, &state, DD_SIGNAL_INSTALLING)) { + ddog_SignalFlush *previous = dd_signal_flush; + dd_signal_flush = flush; + atomic_store(&dd_signal_worker_tid, -1); + atomic_store(&dd_signal_state, flush ? DD_SIGNAL_READY : DD_SIGNAL_DISABLED); + // Drop AFTER publication: another thread's handler may have interrupted malloc. + flush = previous; + } + ddog_sidecar_signal_flush_drop(flush); sigprocmask(SIG_SETMASK, &old_signals, NULL); } void datadog_signals_reset_sidecar_flush_after_fork(void) { - atomic_store_explicit(&dd_signal_worker_tid, 0, memory_order_relaxed); - ddog_SignalFlush *flush = dd_signal_flush; + // The old PID makes handlers ignore inherited state until this reset is complete. + ddog_sidecar_signal_flush_drop(dd_signal_flush); dd_signal_flush = NULL; - datadog_sidecar_signal_flush_drop(flush); - atomic_store_explicit(&dd_signal_state, DD_SIGNAL_DISABLED, memory_order_release); - // the owner pid is purposefully still the old one + atomic_store(&dd_signal_worker_tid, 0); + atomic_store(&dd_signal_state, DD_SIGNAL_DISABLED); + atomic_store(&dd_signal_owner_pid, getpid()); } static void dd_signals_drop_sidecar_flush(void) { - // Called during module shutdown before freeing the flush object and stack. - // Wait if the signal handler is creating a worker, and wait for an existing - // worker to exit. Then set STOPPED so later signals use the previous handler. + // A signal may cause clean shutdown in a fork child before its PHP fork hook runs. + if (atomic_load(&dd_signal_owner_pid) != getpid()) { + datadog_signals_reset_sidecar_flush_after_fork(); + } for (;;) { - int state = atomic_load_explicit(&dd_signal_state, memory_order_acquire); - if (state == DD_SIGNAL_INSTALLING || state == DD_SIGNAL_STARTING_DEFAULT || - state == DD_SIGNAL_STARTING_CUSTOM) { - // Ordinary shutdown context. The signal handler performs only bounded - // startup. + int state = atomic_load(&dd_signal_state); + if (state == DD_SIGNAL_INSTALLING) { sched_yield(); continue; } - if (state == DD_SIGNAL_RUNNING_DEFAULT || state == DD_SIGNAL_RUNNING_CUSTOM) { + if (state == DD_SIGNAL_FLUSHING_DEFAULT || state == DD_SIGNAL_FLUSHING_CUSTOM) { int tid; - while ((tid = atomic_load_explicit(&dd_signal_worker_tid, memory_order_acquire)) != 0) { - // FUTEX_WAIT sleeps only if the word still equals tid, closing - // the exit-before-wait race. EAGAIN, EINTR, and spurious wakes - // are harmless because the loop reloads the word. The kernel's - // clear_child_tid wake uses a shared futex, not FUTEX_WAIT_PRIVATE. - datadog_raw_syscall6(SYS_futex, (long)(uintptr_t)&dd_signal_worker_tid, FUTEX_WAIT, tid, 0, 0, 0); + while ((tid = atomic_load(&dd_signal_worker_tid)) != 0) { + if (tid == -1) { + sched_yield(); // The handler has claimed the request but not finished clone. + } else { + // Kernel clear_child_tid uses a shared futex. Reload after EINTR/EAGAIN + // or a spurious wake; only zero proves the request and extension can be released. + syscall(SYS_futex, &dd_signal_worker_tid, FUTEX_WAIT, tid, NULL, NULL, 0); + } } } - // A READY handler can race this transition. compare_exchange either - // stops it before clone, or reports its newer STARTING/RUNNING state so - // the loop waits for it above. - if (atomic_compare_exchange_strong_explicit(&dd_signal_state, &state, DD_SIGNAL_STOPPED, - memory_order_acq_rel, memory_order_acquire)) { + if (atomic_compare_exchange_strong(&dd_signal_state, &state, DD_SIGNAL_STOPPED)) { break; } } - ddog_SignalFlush *flush = dd_signal_flush; + ddog_sidecar_signal_flush_drop(dd_signal_flush); dd_signal_flush = NULL; - datadog_sidecar_signal_flush_drop(flush); } -static void dd_call_prev_handler(int sig, siginfo_t *si, void *uc) { - struct sigaction *prev_sigaction = sig == SIGINT ? &dd_sigint_prev_sigaction : &dd_sigterm_prev_sigaction; - void *prev_handler = (prev_sigaction->sa_flags & SA_SIGINFO) ? (void *)prev_sigaction->sa_sigaction - : (void *)prev_sigaction->sa_handler; - - if (prev_handler == SIG_IGN) { - return; - } - if (prev_handler == SIG_DFL) { - _exit(0); +// True means a default-disposition worker owns process termination. False means chain +// the previous handler, including when startup fails or a custom handler owns shutdown. +static bool dd_signals_start_flush(bool terminate) { + if (getpid() != atomic_load(&dd_signal_owner_pid)) { + return false; + } + int state = DD_SIGNAL_READY; + int claimed = terminate ? DD_SIGNAL_FLUSHING_DEFAULT : DD_SIGNAL_FLUSHING_CUSTOM; + while (!atomic_compare_exchange_strong(&dd_signal_state, &state, claimed)) { + if (state != DD_SIGNAL_INSTALLING) { + return terminate && state == DD_SIGNAL_FLUSHING_DEFAULT; + } + sched_yield(); + state = DD_SIGNAL_READY; } - if (prev_sigaction->sa_flags & SA_SIGINFO) { - (*prev_sigaction->sa_sigaction)(sig, si, uc); - } else { - (*prev_sigaction->sa_handler)(sig); + + int flags = CLONE_VM | CLONE_FS | CLONE_FILES | CLONE_SIGHAND | CLONE_THREAD | CLONE_SYSVSEM | + CLONE_PARENT_SETTID | CLONE_CHILD_CLEARTID; + int result = datadog_clone_thread(datadog_sidecar_signal_flush_run, dd_signal_cleanup_stack + MIN_STACKSZ, + flags, dd_signal_flush, terminate, &dd_signal_worker_tid); + if (result < 0) { + atomic_store(&dd_signal_worker_tid, 0); } + return result >= 0 && terminate; } static void dd_sigint_sigterm_handler(int sig, siginfo_t *si, void *uc) { - struct sigaction *prev_sigaction = sig == SIGINT ? &dd_sigint_prev_sigaction : &dd_sigterm_prev_sigaction; - void *prev_handler = (prev_sigaction->sa_flags & SA_SIGINFO) ? (void *)prev_sigaction->sa_sigaction - : (void *)prev_sigaction->sa_handler; - if (prev_handler == SIG_IGN) { + struct sigaction *previous = sig == SIGINT ? &dd_sigint_prev_sigaction : &dd_sigterm_prev_sigaction; + if (previous->sa_handler == SIG_IGN) { return; } - bool terminate = prev_handler == SIG_DFL; - int expected = DD_SIGNAL_READY; - int starting = terminate ? DD_SIGNAL_STARTING_DEFAULT : DD_SIGNAL_STARTING_CUSTOM; - // READY is consumed exactly once. Besides preventing duplicate workers, the - // acquire half makes the prepared flush object and owner PID visible. - while (!atomic_compare_exchange_strong_explicit(&dd_signal_state, &expected, starting, memory_order_acq_rel, - memory_order_acquire)) { - if (expected == DD_SIGNAL_INSTALLING && - datadog_raw_syscall6(SYS_getpid, 0, 0, 0, 0, 0, 0) == - atomic_load_explicit(&dd_signal_owner_pid, memory_order_relaxed)) { - // The publisher blocks SIGTERM/SIGINT on its own thread and only - // performs lock-free stores while INSTALLING. Let it finish before - // claiming the replacement. Never wait for a publisher inherited - // across fork: that thread does not exist in the child. - datadog_raw_syscall6(SYS_sched_yield, 0, 0, 0, 0, 0, 0); - expected = DD_SIGNAL_READY; - continue; - } - // A worker handling a default-disposition signal will terminate the - // process after flushing. Do not let another default signal iterrupt - // the flush early by invoking the previous disposition. - if (terminate && (expected == DD_SIGNAL_STARTING_DEFAULT || expected == DD_SIGNAL_RUNNING_DEFAULT)) { - return; - } - dd_call_prev_handler(sig, si, uc); - return; - } - - // A fork child may inherit READY before its post-fork reset. Its owner PID - // is still the parent's, so reject it before accessing the inherited flush - // object. - if (datadog_raw_syscall6(SYS_getpid, 0, 0, 0, 0, 0, 0) != - atomic_load_explicit(&dd_signal_owner_pid, memory_order_relaxed)) { - atomic_store_explicit(&dd_signal_state, DD_SIGNAL_FAILED, memory_order_release); - dd_call_prev_handler(sig, si, uc); - return; - } - if (!dd_signal_cleanup_stack) { - atomic_store_explicit(&dd_signal_state, DD_SIGNAL_FAILED, memory_order_release); - dd_call_prev_handler(sig, si, uc); - return; - } - - // The handler's sa_mask already blocks ordinary blockable signals. glibc and - // musl deliberately remove their internal signals from masks installed via - // public libc APIs, however. If one reached the raw clone, its libc handler - // could use the parent's inherited, unrepaired TLS. Bypass that filtering and - // let the clone inherit the complete kernel mask before it can run. Doing - // this inside the clone would leave a delivery window. SIGKILL and SIGSTOP - // remain unblockable, but neither executes a user-space handler. + int saved_errno = errno; + // Block signals before claiming the request: a nested handler must not enter PHP + // shutdown while clone is starting. libc's mask APIs exclude reserved signals, so + // use the kernel mask for the worker's inherited TLS. The PHP thread can use syscall(). uint64_t all_signals = UINT64_MAX, old_signals; - if (datadog_raw_syscall6(SYS_rt_sigprocmask, SIG_SETMASK, (long)(uintptr_t)&all_signals, - (long)(uintptr_t)&old_signals, sizeof(all_signals), 0, 0) < 0) { - atomic_store_explicit(&dd_signal_state, DD_SIGNAL_FAILED, memory_order_release); - dd_call_prev_handler(sig, si, uc); + bool defer_termination = false; + if (syscall(SYS_rt_sigprocmask, SIG_SETMASK, &all_signals, &old_signals, sizeof(all_signals)) == 0) { + defer_termination = dd_signals_start_flush(previous->sa_handler == SIG_DFL); + syscall(SYS_rt_sigprocmask, SIG_SETMASK, &old_signals, NULL, sizeof(old_signals)); + } + errno = saved_errno; + if (defer_termination) { return; } - void *stack_top = dd_signal_cleanup_stack + dd_signal_cleanup_stack_size; - // These are pthread-like sharing flags without CLONE_SETTLS: the worker has - // no independent libc/Rust thread runtime and may execute only the audited - // raw call graph. PARENT_SETTID plus CHILD_CLEARTID make - // dd_signal_worker_tid a kernel-backed join word. - int flags = CLONE_VM | CLONE_FS | CLONE_FILES | CLONE_SIGHAND | CLONE_THREAD | CLONE_SYSVSEM | CLONE_PARENT_SETTID | - CLONE_CHILD_CLEARTID; - // The trampoline copies `terminate` onto the child stack before clone. Rust - // either returns for the custom-handler case or issues exit_group after the - // bounded exchange for the default disposition. - int result = datadog_clone_thread(datadog_sidecar_signal_flush_run, stack_top, flags, dd_signal_flush, terminate, - &dd_signal_worker_tid); - int running = terminate ? DD_SIGNAL_RUNNING_DEFAULT : DD_SIGNAL_RUNNING_CUSTOM; - // The child may finish before clone returns. That is safe: clear_child_tid - // will already be zero when normal teardown observes the RUNNING state. - atomic_store_explicit(&dd_signal_state, result < 0 ? DD_SIGNAL_FAILED : running, memory_order_release); - datadog_raw_syscall6(SYS_rt_sigprocmask, SIG_SETMASK, (long)(uintptr_t)&old_signals, 0, sizeof(old_signals), 0, 0); - - if (result < 0 || prev_handler != SIG_DFL) { - dd_call_prev_handler(sig, si, uc); - } // else the cleanup worker started and calling dd_call_prev_handler - // would exit the process -} - -static void dd_signals_init_cleanup_stack(void) { - if (!dd_signal_cleanup_stack) { - // Allocate before signal delivery. The one-shot state gate ensures that - // no two workers ever share this stack. - dd_signal_cleanup_stack_size = MIN_STACKSZ; - dd_signal_cleanup_stack = malloc(dd_signal_cleanup_stack_size); + if (previous->sa_handler == SIG_DFL) { + _exit(0); + } else if (previous->sa_flags & SA_SIGINFO) { + previous->sa_sigaction(sig, si, uc); + } else { + previous->sa_handler(sig); } } #endif void datadog_signals_minit(void) { #if __linux + atomic_store(&dd_signal_owner_pid, getpid()); dd_sigint_sigterm_sigaction.sa_sigaction = dd_sigint_sigterm_handler; dd_sigint_sigterm_sigaction.sa_flags = SA_SIGINFO; sigemptyset(&dd_sigint_sigterm_sigaction.sa_mask); if (get_global_DD_TRACE_FORCE_FLUSH_ON_SIGTERM()) { - dd_signals_init_cleanup_stack(); sigaction(SIGTERM, &dd_sigint_sigterm_sigaction, &dd_sigterm_prev_sigaction); } if (get_global_DD_TRACE_FORCE_FLUSH_ON_SIGINT()) { - dd_signals_init_cleanup_stack(); sigaction(SIGINT, &dd_sigint_sigterm_sigaction, &dd_sigint_prev_sigaction); } #endif @@ -593,9 +479,6 @@ void datadog_signals_mshutdown(void) { sigaction(SIGINT, &dd_sigint_prev_sigaction, NULL); } } - - free(dd_signal_cleanup_stack); - dd_signal_cleanup_stack = NULL; #endif if (dd_signal_async_stack) { diff --git a/ext/threads.c b/ext/threads.c index 427d49cba01..aa2fd532459 100644 --- a/ext/threads.c +++ b/ext/threads.c @@ -185,23 +185,6 @@ __asm__( "1: ret\n" ".size datadog_clone_thread,.-datadog_clone_thread\n"); -__asm__( - ".text\n" - ".globl datadog_raw_syscall6\n" - ".hidden datadog_raw_syscall6\n" - ".type datadog_raw_syscall6,@function\n" - "datadog_raw_syscall6:\n" - /* C ABI: rdi = number, rsi/rdi... = six syscall arguments. */ - " movq %rdi, %rax\n" - " movq %rsi, %rdi\n" - " movq %rdx, %rsi\n" - " movq %rcx, %rdx\n" - " movq %r8, %r10\n" - " movq %r9, %r8\n" - " movq 8(%rsp), %r9\n" - " syscall\n" - " ret\n" - ".size datadog_raw_syscall6,.-datadog_raw_syscall6\n"); #elif defined(__aarch64__) __asm__( ".text\n" @@ -234,23 +217,6 @@ __asm__( " brk #0\n" /* unreachable */ ".size datadog_clone_thread,.-datadog_clone_thread\n"); -__asm__( - ".text\n" - ".globl datadog_raw_syscall6\n" - ".hidden datadog_raw_syscall6\n" - ".type datadog_raw_syscall6,%function\n" - "datadog_raw_syscall6:\n" - /* C ABI: x0 = number, x1..x6 = six syscall arguments. */ - " mov x8, x0\n" - " mov x0, x1\n" - " mov x1, x2\n" - " mov x2, x3\n" - " mov x3, x4\n" - " mov x4, x5\n" - " mov x5, x6\n" - " svc #0\n" - " ret\n" - ".size datadog_raw_syscall6,.-datadog_raw_syscall6\n"); #else int datadog_clone_thread(datadog_raw_clone_fn fn, void *stack_top, int flags, const struct ddog_SignalFlush *arg, bool terminate_process, _Atomic(int) *tid) { @@ -263,24 +229,6 @@ int datadog_clone_thread(datadog_raw_clone_fn fn, void *stack_top, int flags, return -1; } -long datadog_raw_syscall6( - long number, - long arg1, - long arg2, - long arg3, - long arg4, - long arg5, - long arg6 -) { - (void)number; - (void)arg1; - (void)arg2; - (void)arg3; - (void)arg4; - (void)arg5; - (void)arg6; - return -1; -} #endif #endif diff --git a/ext/threads.h b/ext/threads.h index aace9961ee4..740b6b6c73f 100644 --- a/ext/threads.h +++ b/ext/threads.h @@ -36,15 +36,6 @@ typedef int32_t (*datadog_raw_clone_fn)(const struct ddog_SignalFlush *, bool); int datadog_clone_thread(datadog_raw_clone_fn fn, void *stack_top, int flags, const struct ddog_SignalFlush *arg, bool terminate_process, _Atomic(int) *tid); -long datadog_raw_syscall6( - long number, - long arg1, - long arg2, - long arg3, - long arg4, - long arg5, - long arg6 -); #endif #endif // DATADOG_THREADS_H diff --git a/libdatadog b/libdatadog index 588a120ba82..89ee2a8b042 160000 --- a/libdatadog +++ b/libdatadog @@ -1 +1 @@ -Subproject commit 588a120ba82e41311496f057dc2340d731d43abe +Subproject commit 89ee2a8b042db71e99017d5156facf5dd897c36b