Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
39 commits
Select commit Hold shift + click to select a range
52710cf
fix: attach obfuscation flag to stats buckets, fix obfuscation versio…
Eldolfin Jul 6, 2026
a62a150
fix: fmt
Eldolfin Jul 7, 2026
8f65b1d
fix: flush buckets which have the same obfuscation indicator together
Eldolfin Jul 7, 2026
a0bfe8c
fix(sidecar): remove useless stats-obfuscation feature
Eldolfin Jul 7, 2026
6c6f058
Revert "fix(sidecar): remove useless stats-obfuscation feature"
Eldolfin Jul 7, 2026
a44467a
fix: outdated comment, remove metadata return from otlp flush
Eldolfin Jul 7, 2026
3113a7e
fix: tests lints
Eldolfin Jul 7, 2026
909d654
fix: local review suggestions
Eldolfin Jul 7, 2026
85eaf0c
fix: fmt
Eldolfin Jul 7, 2026
05d2364
fix(test): put back obfuscation enabled in SpanConcentrator instead o…
Eldolfin Jul 7, 2026
ed0277a
fix: remove unused file results.xml
Eldolfin Jul 7, 2026
888e56c
fix error message StateExporter -> StatsExporter
Eldolfin Jul 7, 2026
9d91a3c
fix nit comment
Eldolfin Jul 7, 2026
2342ce1
refacto: cleaner flush result struct, simplify span obfuscation path
Eldolfin Jul 8, 2026
2d4eebc
feat: send obfuscated/unobfuscated stats concurrently
Eldolfin Jul 9, 2026
9f7474e
Merge branch 'main' into oscarld/fix-css-obfuscation-logic
Eldolfin Jul 10, 2026
35d265b
fix: FuturesUnordered over tokio::join!
Eldolfin Jul 10, 2026
e082005
feat!(stats): per-key cardinality limits
Eldolfin Jul 8, 2026
1c9dca9
feat: store distinct hashes instead of values
Eldolfin Jul 8, 2026
43a1b0b
fix: libdd-data-pipeline(-ffi) build
Eldolfin Jul 8, 2026
586e6d7
fix: wrong peer_tags collapse value
Eldolfin Jul 8, 2026
cc5bf3b
fix: fmt
Eldolfin Jul 8, 2026
726c57a
fix: collapse key fields before applying whole-key cardinality limits
Eldolfin Jul 8, 2026
bd6e3bc
feat: add a warning when whole-key limit is lower than per-field limits
Eldolfin Jul 8, 2026
1c46a80
fix: fmt
Eldolfin Jul 8, 2026
04befd7
fix(ffi): add libdd-trace-stats to cbindgen config to generate Cardin…
Eldolfin Jul 8, 2026
74f105a
fix: inline DEFAULT_MAX_ENTRIES_PER_BUCKET
Eldolfin Jul 10, 2026
edf1e0c
fix: link span derived primary tags ticket and pr
Eldolfin Jul 10, 2026
5cedcf4
fix: avoid useless re-hashes
Eldolfin Jul 10, 2026
83043a4
fix: warn on any cardinality limit = 0
Eldolfin Jul 10, 2026
39427ce
fix: fmt
Eldolfin Jul 10, 2026
a14a2b6
feat: add function to get default cardinality limit config
Eldolfin Jul 20, 2026
45e3b60
Merge remote-tracking branch 'origin/main' into oscarld/stats-per-key…
Eldolfin Jul 21, 2026
a1c321a
fix: bad merge lints
Eldolfin Jul 21, 2026
4b3981f
fix: apply configured limits to OTLP stats
Eldolfin Jul 22, 2026
909fb50
Merge remote-tracking branch 'origin/main' into oscarld/stats-per-key…
Eldolfin Jul 22, 2026
854ef55
explain HashSet of hashes reasoning
Eldolfin Jul 22, 2026
2906b45
implement additional_metrics_tags field cardinality limit
Eldolfin Jul 22, 2026
85c73e9
fix: missing resource cardinality limit = 0 warn
Eldolfin Jul 22, 2026
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
2 changes: 1 addition & 1 deletion libdd-data-pipeline-ffi/cbindgen.toml
Original file line number Diff line number Diff line change
Expand Up @@ -47,4 +47,4 @@ must_use = "DDOG_CHECK_RETURN"

[parse]
parse_deps = true
include = ["libdd-common", "libdd-common-ffi", "libdd-shared-runtime", "libdd-data-pipeline", "libdd-trace-utils"]
include = ["libdd-common", "libdd-common-ffi", "libdd-shared-runtime", "libdd-data-pipeline", "libdd-trace-utils", "libdd-trace-stats"]
44 changes: 35 additions & 9 deletions libdd-data-pipeline-ffi/src/trace_exporter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use libdd_common_ffi::{
CharSlice,
{slice::AsBytes, slice::ByteSlice},
};
use libdd_data_pipeline::trace_exporter::stats::CardinalityLimitConfig;
use libdd_data_pipeline::trace_exporter::{
TelemetryConfig, TelemetryInstrumentationSessions, TraceExporter as GenericTraceExporter,
TraceExporterInputFormat, TraceExporterOutputFormat,
Expand Down Expand Up @@ -92,7 +93,7 @@ pub struct TraceExporterConfig {
otlp_instrumentation_scope_version: Option<String>,
output_to_log: bool,
log_max_line_size: Option<usize>,
stats_cardinality_limit: Option<usize>,
stats_cardinality_limits: Option<CardinalityLimitConfig>,
}

#[no_mangle]
Expand Down Expand Up @@ -588,11 +589,11 @@ pub unsafe extern "C" fn ddog_trace_exporter_config_set_otlp_instrumentation_sco
#[no_mangle]
pub unsafe extern "C" fn ddog_trace_exporter_config_set_stats_cardinality_limit(
config: Option<&mut TraceExporterConfig>,
limit: usize,
limits: CardinalityLimitConfig,
) -> Option<Box<ExporterError>> {
catch_panic!(
if let Some(handle) = config {
handle.stats_cardinality_limit = Some(limit);
handle.stats_cardinality_limits = Some(limits);
Comment thread
Eldolfin marked this conversation as resolved.
None
} else {
gen_error!(ErrorCode::InvalidArgument)
Expand All @@ -601,6 +602,13 @@ pub unsafe extern "C" fn ddog_trace_exporter_config_set_stats_cardinality_limit(
)
}

/// Returns a `CardinalityLimitConfig` with default values
#[no_mangle]
pub unsafe extern "C" fn ddog_trace_exporter_config_default_stats_cardinality_limit(
) -> CardinalityLimitConfig {
CardinalityLimitConfig::default()
}

/// Configure the exporter to write traces as newline-delimited JSON to stdout (the Datadog
/// Forwarder "log exporter" path) instead of sending them to a Datadog agent. Used in serverless
/// environments (e.g. AWS Lambda) when no agent is reachable.
Expand Down Expand Up @@ -678,7 +686,7 @@ pub unsafe extern "C" fn ddog_trace_exporter_new(
.set_output_format(config.output_format)
.set_connection_timeout(config.connection_timeout);

if let Some(limit) = config.stats_cardinality_limit {
if let Some(limit) = config.stats_cardinality_limits {
builder.set_stats_cardinality_limit(limit);
}

Expand Down Expand Up @@ -825,7 +833,7 @@ mod tests {
assert!(cfg.connection_timeout.is_none());
assert!(!cfg.output_to_log);
assert_eq!(cfg.log_max_line_size, None);
assert_eq!(cfg.stats_cardinality_limit, None);
assert_eq!(cfg.stats_cardinality_limits, None);
assert!(cfg.otlp_instrumentation_scope_name.is_none());
assert!(cfg.otlp_instrumentation_scope_version.is_none());

Expand Down Expand Up @@ -1596,16 +1604,34 @@ mod tests {
fn config_stats_cardinality_limit_test() {
unsafe {
// Null config → InvalidArgument
let error = ddog_trace_exporter_config_set_stats_cardinality_limit(None, 100);
let error = ddog_trace_exporter_config_set_stats_cardinality_limit(
None,
CardinalityLimitConfig {
whole_key_limit: 100,
..Default::default()
},
);
assert_eq!(error.as_ref().unwrap().code, ErrorCode::InvalidArgument);
ddog_trace_exporter_error_free(error);

// Valid config → value stored
let mut config = Some(TraceExporterConfig::default());
let error =
ddog_trace_exporter_config_set_stats_cardinality_limit(config.as_mut(), 500);
let error = ddog_trace_exporter_config_set_stats_cardinality_limit(
config.as_mut(),
CardinalityLimitConfig {
whole_key_limit: 500,
..Default::default()
},
);
assert_eq!(error, None);
assert_eq!(config.unwrap().stats_cardinality_limit, Some(500));
assert_eq!(
config
.unwrap()
.stats_cardinality_limits
.unwrap()
.whole_key_limit,
500
);
}
}

Expand Down
16 changes: 10 additions & 6 deletions libdd-data-pipeline/src/trace_exporter/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ use libdd_dogstatsd_client::DogStatsDClient;
use libdd_shared_runtime::SharedRuntime;
#[cfg(not(target_arch = "wasm32"))]
use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime};
use libdd_trace_stats::span_concentrator::CardinalityLimitConfig;
use libdd_trace_utils::trace_filter::TraceFilterer;
use std::sync::Arc;
use std::time::Duration;
Expand Down Expand Up @@ -75,7 +76,7 @@ pub struct TraceExporterBuilder<R: SharedRuntime> {
/// A Some value enables stats-computation, None if it is disabled
stats_bucket_size: Option<Duration>,
peer_tags: Vec<String>,
stats_cardinality_limit: Option<usize>,
stats_cardinality_limits: Option<CardinalityLimitConfig>,
#[cfg(feature = "stats-obfuscation")]
client_side_stats_obfuscation_enabled: bool,
#[cfg(feature = "telemetry")]
Expand Down Expand Up @@ -144,7 +145,7 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
client_computed_top_level: false,
stats_bucket_size: None,
peer_tags: Vec::new(),
stats_cardinality_limit: None,
stats_cardinality_limits: None,
#[cfg(feature = "stats-obfuscation")]
client_side_stats_obfuscation_enabled: false,
#[cfg(feature = "telemetry")]
Expand Down Expand Up @@ -339,8 +340,11 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
/// This bounds memory usage when the trace population has very high cardinality.
///
/// Has no effect unless stats computation is enabled.
pub fn set_stats_cardinality_limit(&mut self, cardinality_limit: usize) -> &mut Self {
self.stats_cardinality_limit = Some(cardinality_limit);
pub fn set_stats_cardinality_limit(
&mut self,
cardinality_limits: CardinalityLimitConfig,
) -> &mut Self {
self.stats_cardinality_limits = Some(cardinality_limits);
self
}

Expand Down Expand Up @@ -753,7 +757,7 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
web_time::SystemTime::now(),
span_kinds,
self.peer_tags.clone(),
None,
self.stats_cardinality_limits,
vec![],
#[cfg(feature = "stats-obfuscation")]
None,
Expand Down Expand Up @@ -826,7 +830,7 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
common_stats_tags: vec![libdatadog_version],
client_side_stats: StatsComputationConfig {
status: ArcSwap::new(stats.into()),
stats_cardinality_limit: self.stats_cardinality_limit,
stats_cardinality_limits: self.stats_cardinality_limits,
Comment thread
Eldolfin marked this conversation as resolved.
#[cfg(feature = "stats-obfuscation")]
obfuscation_config: Arc::new(ArcSwap::from_pointee(
StatsComputationObfuscationConfig::default(),
Expand Down
6 changes: 3 additions & 3 deletions libdd-data-pipeline/src/trace_exporter/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -459,7 +459,7 @@ impl<
metadata: &self.metadata,
endpoint_url: &self.endpoint.url,
shared_runtime: &*self.shared_runtime,
stats_cardinality_limit: self.client_side_stats.stats_cardinality_limit,
stats_cardinality_limits: self.client_side_stats.stats_cardinality_limits,
dogstatsd: if self.health_metrics_enabled {
self.dogstatsd.clone()
} else {
Expand Down Expand Up @@ -862,8 +862,8 @@ impl<
&& self.v1_active.swap(false, Ordering::Relaxed)
{
warn!(
"V1 trace send returned 404; agent no longer advertises {V1_TRACES_ENDPOINT} — falling back to V0.4"
);
"V1 trace send returned 404; agent no longer advertises {V1_TRACES_ENDPOINT} — falling back to V0.4"
);
self.info_response_observer.manual_trigger();
}
}
Expand Down
8 changes: 5 additions & 3 deletions libdd-data-pipeline/src/trace_exporter/stats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@
//! including starting/stopping stats workers, managing the span concentrator,
//! and processing traces for stats collection.

pub use libdd_trace_stats::span_concentrator::CardinalityLimitConfig;

use super::add_path;
use super::TracerMetadata;
use crate::agent_info::schema::AgentInfo;
Expand Down Expand Up @@ -47,7 +49,7 @@ pub(crate) struct StatsContext<
pub metadata: &'a TracerMetadata,
pub endpoint_url: &'a http::Uri,
pub shared_runtime: &'a R,
pub stats_cardinality_limit: Option<usize>,
pub stats_cardinality_limits: Option<CardinalityLimitConfig>,
/// Optional DogStatsD client forwarded to the [`StatsExporter`].
pub dogstatsd: Option<libdd_dogstatsd_client::DogStatsDClient>,
/// Optional telemetry handle forwarded to the [`StatsExporter`].
Expand Down Expand Up @@ -75,7 +77,7 @@ pub(crate) enum StatsComputationStatus {
#[derive(Debug)]
pub(crate) struct StatsComputationConfig {
pub(crate) status: ArcSwap<StatsComputationStatus>,
pub(crate) stats_cardinality_limit: Option<usize>,
pub(crate) stats_cardinality_limits: Option<CardinalityLimitConfig>,
#[cfg(feature = "stats-obfuscation")]
pub(crate) obfuscation_config: SharedStatsComputationObfuscationConfig,
/// Builder-level opt-in. When false, stats obfuscation stays off
Expand Down Expand Up @@ -137,7 +139,7 @@ pub(crate) fn start_stats_computation<
SystemTime::now(),
span_kinds,
peer_tags,
ctx.stats_cardinality_limit,
ctx.stats_cardinality_limits,
vec![],
#[cfg(feature = "stats-obfuscation")]
Some(client_side_stats.obfuscation_config.clone()),
Expand Down
99 changes: 87 additions & 12 deletions libdd-trace-stats/src/span_concentrator/aggregation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,15 +5,20 @@
//! This includes the aggregation key to group spans together and the computation of stats from a
//! span.

use hashbrown::HashMap;
use hashbrown::{HashMap, HashSet};
use libdd_trace_obfuscation::ip_address::quantize_peer_ip_addresses;
use libdd_trace_protobuf::pb;
use libdd_trace_utils::span::SpanText;
use std::borrow::{Borrow, Cow};
use std::{
borrow::{Borrow, Cow},
hash::{DefaultHasher, Hash, Hasher as _},
};
use tracing::warn;

use crate::span_concentrator::StatSpan;

use super::CardinalityLimitConfig;

/// Sentinel value used for cardinality limiting.
pub const TRACER_BLOCKED_VALUE: &str = "tracer_blocked_value";

Expand Down Expand Up @@ -349,8 +354,8 @@ impl OwnedAggregationKey {
is_synthetics_request: false,
is_trace_root: pb::Trilean::NotSet,
},
peer_tags: vec![],
additional_metric_tags: vec![],
peer_tags: vec![(TRACER_BLOCKED_VALUE.to_owned(), "".to_owned())],
additional_metric_tags: vec![(TRACER_BLOCKED_VALUE.to_owned(), "".to_owned())],
}
}
}
Expand Down Expand Up @@ -487,7 +492,15 @@ pub(super) struct StatsBucket {
start: u64,
/// Maximum number of distinct aggregation keys this bucket will hold before collapsing new
/// ones into the overflow sentinel key.
max_entries: usize,
cardinality_limits: CardinalityLimitConfig,
// HashSet of hashes of field values so we save memory
// This is not 100% accurate but the probability of getting collision is close to 0
// In the very rare case we get a collision, we would get one extra bucket which is totally
// fine
distinct_resources: HashSet<u64>,
distinct_http_endpoints: HashSet<u64>,
distinct_peer_tags: HashSet<u64>,
distinct_additional_tags: HashSet<u64>,
/// Number of spans collapsed into the overflow bucket due to cardinality limiting.
collapsed_count: u64,
/// Indicates if stats obfuscated in this bucket. This is set once at creation and stays
Expand All @@ -499,20 +512,23 @@ pub(super) struct StatsBucket {
impl StatsBucket {
/// Return a new StatsBucket starting at `start_timestamp`.
///
/// `max_entries` is the maximum number of distinct aggregation keys the bucket will hold.
/// Once the limit is reached, new distinct keys are collapsed into the overflow sentinel key.
/// `cardinality_limits` are the values for whole-key and per-field cardinality limits
pub(super) fn new(
start_timestamp: u64,
max_entries: usize,
cardinality_limits: CardinalityLimitConfig,
#[cfg(feature = "stats-obfuscation")] obfuscation_enabled: bool,
) -> Self {
Self {
data: HashMap::new(),
start: start_timestamp,
max_entries,
cardinality_limits,
collapsed_count: 0,
#[cfg(feature = "stats-obfuscation")]
obfuscated: obfuscation_enabled,
distinct_resources: HashSet::new(),
distinct_http_endpoints: HashSet::new(),
distinct_peer_tags: HashSet::new(),
distinct_additional_tags: HashSet::new(),
}
}

Expand All @@ -528,11 +544,14 @@ impl StatsBucket {
/// `max_entries` limit, which collapses it into the overflow sentinel key.
pub(super) fn insert(
&mut self,
key: BorrowedAggregationKey<'_>,
mut key: BorrowedAggregationKey<'_>,
duration: i64,
is_error: bool,
is_top_level: bool,
) {
// Per field cardinality limiting
self.collapse_key_fields_cardinality(&mut key);

// The map can't change size before the entry below is resolved, so this single read
// covers the `max_entries` check in the vacant branch without a further lookup.
let len_before_insert = self.data.len();
Expand All @@ -545,7 +564,7 @@ impl StatsBucket {
hashbrown::hash_map::EntryRef::Vacant(e) => {
// New key over the max entry limit, collapse into the overflow
// sentinel.
if len_before_insert >= self.max_entries {
if len_before_insert >= self.cardinality_limits.whole_key_limit {
self.collapsed_count += 1;
self.data
.entry(OwnedAggregationKey::overflow_key())
Expand All @@ -560,6 +579,56 @@ impl StatsBucket {
}
}

/// Collapse an aggregation key fields following the bucket's `CardinalityLimitConfig`
fn collapse_key_fields_cardinality(&mut self, key: &mut BorrowedAggregationKey<'_>) {
use hashbrown::hash_set::Entry;
fn hash(input: &impl Hash) -> u64 {
let mut hasher = DefaultHasher::new();
input.hash(&mut hasher);
hasher.finish()
}

let resource_hash = hash(&key.fixed.resource_name);
let resources_count = self.distinct_resources.len();
if let Entry::Vacant(slot) = self.distinct_resources.entry(resource_hash) {
if resources_count >= self.cardinality_limits.resource_limit {
key.fixed.resource_name = TRACER_BLOCKED_VALUE;
} else {
slot.insert();
}
}

let http_endpoint_hash = hash(&key.fixed.http_endpoint);
let http_endpoints_count = self.distinct_http_endpoints.len();
if let Entry::Vacant(slot) = self.distinct_http_endpoints.entry(http_endpoint_hash) {
if http_endpoints_count >= self.cardinality_limits.http_endpoint_limit {
key.fixed.http_endpoint = TRACER_BLOCKED_VALUE;
} else {
slot.insert();
}
}

let peer_tags_hash = hash(&key.peer_tags);
let peer_tags_count = self.distinct_peer_tags.len();
if let Entry::Vacant(slot) = self.distinct_peer_tags.entry(peer_tags_hash) {
if peer_tags_count >= self.cardinality_limits.peer_tags_limit {
key.peer_tags = vec![(TRACER_BLOCKED_VALUE, Cow::Borrowed(""))];
} else {
slot.insert();
}
}

let additional_tags_hash = hash(&key.additional_metric_tags);
let additional_tags_count = self.distinct_additional_tags.len();
if let Entry::Vacant(slot) = self.distinct_additional_tags.entry(additional_tags_hash) {
if additional_tags_count >= self.cardinality_limits.additional_tags_limit {
key.additional_metric_tags = vec![(TRACER_BLOCKED_VALUE, "")];
} else {
slot.insert();
}
}
}

/// Consume the bucket and return a ClientStatsBucket containing the bucket stats.
/// `bucket_duration` is the size of buckets for the concentrator containing the bucket.
pub(super) fn flush(self, bucket_duration: u64) -> pb::ClientStatsBucket {
Expand Down Expand Up @@ -624,7 +693,13 @@ fn encode_grouped_stats(key: OwnedAggregationKey, group: GroupedStats) -> pb::Cl
peer_tags: key
.peer_tags
.into_iter()
.map(|(k, v)| format!("{k}:{v}"))
.map(|(k, v)| {
if v.is_empty() {
k.to_string()
Comment thread
Eldolfin marked this conversation as resolved.
} else {
format!("{k}:{v}")
}
})
.collect(),
is_trace_root: f.is_trace_root.into(),
http_method: f.http_method,
Expand Down
Loading
Loading