diff --git a/components-rs/sidecar.rs b/components-rs/sidecar.rs index eae5a19357..af9cc6e411 100644 --- a/components-rs/sidecar.rs +++ b/components-rs/sidecar.rs @@ -3,11 +3,11 @@ use datadog_sidecar::service::blocking::{acquire_exception_hash_rate_limiter, Si use datadog_sidecar::service::exception_hash_rate_limiter::ExceptionHashRateLimiter; use datadog_sidecar::tracer::shm_limiter_path; use lazy_static::lazy_static; -use libdd_common::rate_limiter::{Limiter, LocalLimiter}; +use libdd_common::rate_limiter::LocalLimiter; use libdd_common::Endpoint; use libdd_common_ffi::slice::AsBytes; use libdd_common_ffi::{self as ffi, CharSlice, MaybeError}; -use libdd_ipc::rate_limiter::{AnyLimiter, ShmLimiterMemory}; +use libdd_ipc::rate_limiter::{ShmLimiter, ShmLimiterMemory}; use libdd_telemetry_ffi::try_c; #[cfg(windows)] use spawn_worker::{get_trampoline_target_data, LibDependency}; @@ -289,7 +289,14 @@ static SHM_LIMITER: LazyLock> = static EXCEPTION_HASH_LIMITER: LazyLock = LazyLock::new(ExceptionHashRateLimiter::new_reader); -pub struct MaybeShmLimiter(Option); +const SHM_LIMITER_GRANULARITY: Duration = Duration::from_secs(1); + +pub struct MaybeShmLimiter(Option); + +enum Limiter { + Local(LocalLimiter), + Shm(ShmLimiter<()>), +} impl MaybeShmLimiter { pub fn open(index: u32) -> Self { @@ -299,17 +306,17 @@ impl MaybeShmLimiter { Some( SHM_LIMITER .get(index) - .map(AnyLimiter::Shm) - .unwrap_or_else(|| AnyLimiter::Local(LocalLimiter::default())), + .map(Limiter::Shm) + .unwrap_or_else(|| Limiter::Local(LocalLimiter::default())), ) }) } pub fn inc(&self, limit: u32) -> bool { - if let Some(ref limiter) = self.0 { - limiter.inc(limit) - } else { - true + match &self.0 { + Some(Limiter::Local(limiter)) => limiter.inc(limit, SHM_LIMITER_GRANULARITY), + Some(Limiter::Shm(limiter)) => limiter.inc(limit, SHM_LIMITER_GRANULARITY), + None => true, } } } @@ -325,10 +332,10 @@ pub extern "C" fn ddog_exception_hash_limiter_inc( hash: u64, granularity_seconds: u32, ) -> bool { - if let Some(limiter) = EXCEPTION_HASH_LIMITER.find(hash) { + let granularity = Duration::from_secs(granularity_seconds as u64); + if let Some(limiter) = EXCEPTION_HASH_LIMITER.find(hash, granularity) { return limiter.inc(); } - let granularity = Duration::from_secs(granularity_seconds as u64); let _ = acquire_exception_hash_rate_limiter(connection, hash, granularity); true } diff --git a/ext/sidecar.c b/ext/sidecar.c index ee6f5745b5..e0082899a9 100644 --- a/ext/sidecar.c +++ b/ext/sidecar.c @@ -973,7 +973,9 @@ void datadog_sidecar_gshutdown(zend_datadog_globals *datadog_globals) { } bool datadog_alter_test_session_token(zval *old_value, zval *new_value, zend_string *new_str) { - UNUSED(old_value, new_str); + UNUSED(new_str); + bool token_changed = + Z_TYPE_P(old_value) != IS_STRING || !zend_string_equals(Z_STR_P(old_value), Z_STR_P(new_value)); if (datadog_endpoint) { ddog_endpoint_set_test_token_if_changed(datadog_endpoint, dd_zend_string_to_CharSlice(Z_STR_P(new_value))); } @@ -983,6 +985,13 @@ bool datadog_alter_test_session_token(zval *old_value, zval *new_value, zend_str } #if !defined(_WIN32) && defined(DDTRACE) ddtrace_coms_set_test_session_token(Z_STRVAL_P(new_value), Z_STRLEN_P(new_value)); +#endif +#ifdef DDTRACE + /* The test token is part of the named sampling-config shared-memory path. The reader keeps + * its own endpoint copy, so changing the sender endpoint does not retarget an existing reader. */ + if (token_changed) { + ddtrace_recreate_agent_config_reader(); + } #endif return true; } diff --git a/libdatadog b/libdatadog index d120c1090a..591bdb95e9 160000 --- a/libdatadog +++ b/libdatadog @@ -1 +1 @@ -Subproject commit d120c1090a03fa7274fed4121344d25660081895 +Subproject commit 591bdb95e9ec59d5cfa089418917f3bbb993d51d diff --git a/tests/Benchmarks/API/TraceSerializationBench.php b/tests/Benchmarks/API/TraceSerializationBench.php index 2e5ad35314..30c07c111f 100644 --- a/tests/Benchmarks/API/TraceSerializationBench.php +++ b/tests/Benchmarks/API/TraceSerializationBench.php @@ -30,22 +30,24 @@ public function warmUpSampling() throw new \RuntimeException('Trace serialization benchmark requires a ready agent /info response'); } - // Sampling rates are published after a trace response. Exercise automatic sampling - // and the real sender, without changing the priority of the measured trace. - $span = \DDTrace\start_trace_span(); - $span->name = 'bench.trace_serialization.warmup'; - \DDTrace\close_span(); - \DDTrace\flush(); - + // Sampling rates are published after a trace response. Retry the trace as well as the + // shared-memory read: a transient send failure would otherwise leave the worker polling + // data that cannot change. $deadline = microtime(true) + 5; do { - // Trace enqueueing is asynchronous, so the first flush can precede it. - \dd_trace_synchronous_flush(5000); + // Exercise automatic sampling and the real sender, without changing the priority of + // the measured trace. Trace enqueueing is asynchronous, so flush both layers. + $span = \DDTrace\start_trace_span(); + $span->name = 'bench.trace_serialization.warmup'; + \DDTrace\close_span(); + \DDTrace\flush(); + \dd_trace_synchronous_flush(1000); + $config = \dd_trace_internal_fn('get_agent_sampling_config'); if (isset($config['rate_by_service']) && is_array($config['rate_by_service'])) { return; } - usleep(10000); + usleep(100000); } while (microtime(true) < $deadline); throw new \RuntimeException('Trace serialization benchmark requires agent sampling rates'); diff --git a/tests/ext/background-sender/agent_sampling_sidecar_token_change.phpt b/tests/ext/background-sender/agent_sampling_sidecar_token_change.phpt new file mode 100644 index 0000000000..e143c79063 --- /dev/null +++ b/tests/ext/background-sender/agent_sampling_sidecar_token_change.phpt @@ -0,0 +1,49 @@ +--TEST-- +The sidecar sampling reader follows test session token changes +--SKIPIF-- + + +--ENV-- +DD_TRACE_LOG_LEVEL=warn +DD_AGENT_HOST=request-replayer +DD_TRACE_AGENT_PORT=80 +DD_TRACE_AGENT_FLUSH_INTERVAL=333 +DD_TRACE_GENERATE_ROOT_SPAN=0 +DD_INSTRUMENTATION_TELEMETRY_ENABLED=0 +DD_TRACE_SIDECAR_TRACE_SENDER=1 +DD_TRACE_SIDECAR_CONNECTION_MODE=thread +DD_TRACE_IGNORE_AGENT_SAMPLING_RATES=0 +--INI-- +datadog.trace.agent_test_session_token=background-sender/agent_sampling_sidecar_token_change-before +--FILE-- +replayRequest(); // cleanup possible leftover + +$marker = 'service:agent-sampling-sidecar-token-change-' . getmypid() . ',env:test'; +$rr->setResponse(['rate_by_service' => [$marker => 1]]); + +\DDTrace\start_span(); +\DDTrace\close_span(); +dd_trace_internal_fn('synchronous_flush'); +$rr->waitForDataAndReplay(); + +for ($i = 0; $i < 100; $i++) { + $sampling = dd_trace_internal_fn('get_agent_sampling_config'); + if (isset($sampling['rate_by_service'][$marker])) { + break; + } + usleep(100000); +} + +var_dump((float) ($sampling['rate_by_service'][$marker] ?? -1)); +?> +--EXPECT-- +float(1) diff --git a/tracer/ddtrace.c b/tracer/ddtrace.c index 4767d46aca..e450008dbd 100644 --- a/tracer/ddtrace.c +++ b/tracer/ddtrace.c @@ -440,6 +440,24 @@ void ddtrace_first_rinit(void) { dd_rinit_once_done = true; } +void ddtrace_recreate_agent_config_reader(void) { + if (!DDTRACE_G(agent_config_reader) || !get_global_DD_TRACE_SIDECAR_TRACE_SENDER()) { + return; + } + + ddog_agent_remote_config_reader_drop(DDTRACE_G(agent_config_reader)); + DDTRACE_G(agent_config_reader) = NULL; + + if (DDTRACE_G(agent_rate_by_service)) { + zai_json_release_persistent_array(DDTRACE_G(agent_rate_by_service)); + DDTRACE_G(agent_rate_by_service) = NULL; + } + + if (datadog_endpoint) { + DDTRACE_G(agent_config_reader) = ddog_agent_remote_config_reader_for_endpoint(datadog_endpoint); + } +} + static void dd_initialize_request(void) { DDTRACE_G(distributed_trace_id) = (datadog_trace_id){0}; DDTRACE_G(distributed_parent_trace_id) = 0; diff --git a/tracer/tracer_api.h b/tracer/tracer_api.h index 032bff78a5..178831f3ab 100644 --- a/tracer/tracer_api.h +++ b/tracer/tracer_api.h @@ -19,6 +19,7 @@ void ddtrace_rinit_early(void); void ddtrace_rinit(void); void ddtrace_rshutdown(bool fast_shutdown); void ddtrace_post_deactivate(void); +void ddtrace_recreate_agent_config_reader(void); // fork handling void ddtrace_internal_handle_fork(void);