Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 18 additions & 11 deletions components-rs/sidecar.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -289,7 +289,14 @@ static SHM_LIMITER: LazyLock<ShmLimiterMemory<()>> =
static EXCEPTION_HASH_LIMITER: LazyLock<ExceptionHashRateLimiter> =
LazyLock::new(ExceptionHashRateLimiter::new_reader);

pub struct MaybeShmLimiter(Option<AnyLimiter>);
const SHM_LIMITER_GRANULARITY: Duration = Duration::from_secs(1);

pub struct MaybeShmLimiter(Option<Limiter>);

enum Limiter {
Local(LocalLimiter),
Shm(ShmLimiter<()>),
}

impl MaybeShmLimiter {
pub fn open(index: u32) -> Self {
Expand All @@ -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,
}
}
}
Expand All @@ -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
}
11 changes: 10 additions & 1 deletion ext/sidecar.c
Original file line number Diff line number Diff line change
Expand Up @@ -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)));
}
Expand All @@ -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;
}
Expand Down
22 changes: 12 additions & 10 deletions tests/Benchmarks/API/TraceSerializationBench.php
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
--TEST--
The sidecar sampling reader follows test session token changes
--SKIPIF--
<?php include __DIR__ . '/../includes/skipif_no_dev_env.inc'; ?>
<?php if (getenv('USE_ZEND_ALLOC') === '0' && !getenv('SKIP_ASAN')) die('skip timing sensitive test - valgrind is too slow'); ?>
--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--
<?php
include __DIR__ . '/../includes/request_replayer.inc';

ini_set(
'datadog.trace.agent_test_session_token',
'background-sender/agent_sampling_sidecar_token_change-after'
);

$rr = new RequestReplayer();
$rr->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)
18 changes: 18 additions & 0 deletions tracer/ddtrace.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
1 change: 1 addition & 0 deletions tracer/tracer_api.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading