mod aggregation;
pub mod cardinality_limit_telemetry;
pub mod stat_span;
use std::collections::HashMap;
use std::time::Duration;
use tracing::{debug, warn};
use web_time::{SystemTime, UNIX_EPOCH};
use libdd_trace_protobuf::pb;
use aggregation::StatsBucket;
use aggregation::BorrowedAggregationKey;
pub use aggregation::{FixedAggregationKey, OtlpExactCell, OtlpExactGroup, OtlpStatsBucket};
use cardinality_limit_telemetry::CollapsedFieldsMetrics;
pub use stat_span::{ChunkSpanView, StatSpan};
const ADDITIONAL_METRIC_TAGS_MAX_KEYS: usize = 4;
fn normalize_additional_metric_tag_keys(mut keys: Vec<String>) -> Vec<String> {
keys.sort_unstable();
keys.dedup();
if keys.len() > ADDITIONAL_METRIC_TAGS_MAX_KEYS {
let dropped = keys.split_off(ADDITIONAL_METRIC_TAGS_MAX_KEYS);
warn!(
"additional_metric_tag_keys: {} additional metric tag keys exceed the cap of {}; dropping: {:?}",
dropped.len() + ADDITIONAL_METRIC_TAGS_MAX_KEYS,
ADDITIONAL_METRIC_TAGS_MAX_KEYS,
dropped,
);
}
keys
}
pub struct FlushResult<T> {
pub obfuscated_buckets: Vec<T>,
pub unobfuscated_buckets: Vec<T>,
pub collapsed_spans: u64,
pub collapsed_fields_metrics: CollapsedFieldsMetrics,
}
impl<T> FlushResult<T> {
pub fn all_buckets(self) -> Vec<T> {
let mut buckets = self.obfuscated_buckets;
buckets.extend(self.unobfuscated_buckets);
buckets
}
}
pub trait FlushableConcentrator {
fn flush_buckets(&mut self, force: bool) -> FlushResult<pb::ClientStatsBucket>;
}
impl FlushableConcentrator for SpanConcentrator {
fn flush_buckets(&mut self, force: bool) -> FlushResult<pb::ClientStatsBucket> {
self.flush(SystemTime::now(), force)
}
}
fn system_time_to_unix_duration(t: SystemTime) -> Duration {
t.duration_since(UNIX_EPOCH)
.unwrap_or(Duration::from_nanos(0))
}
#[inline]
fn align_timestamp(t: u64, bucket_size: u64) -> u64 {
t - (t % bucket_size)
}
pub fn is_span_eligible<'a, T>(span: &'a T, span_kinds_stats_computed: &[String]) -> bool
where
T: StatSpan<'a>,
{
(span.has_top_level() || span.is_measured() || {
span.get_meta("span.kind")
.is_some_and(|span_kind| span_kinds_stats_computed.contains(&span_kind.to_lowercase()))
}) && !span.is_partial_snapshot()
}
#[cfg(feature = "stats-obfuscation")]
#[derive(Clone, Debug, Default)]
#[cfg_attr(target_arch = "wasm32", allow(dead_code))]
pub struct StatsComputationObfuscationConfig {
pub enabled: bool,
pub sql_obfuscation_mode: libdd_trace_obfuscation::sql::SqlObfuscationMode,
}
#[cfg(feature = "stats-obfuscation")]
pub type SharedStatsComputationObfuscationConfig =
std::sync::Arc<arc_swap::ArcSwap<StatsComputationObfuscationConfig>>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(C)]
pub struct CardinalityLimitConfig {
pub whole_key_limit: usize,
pub resource_limit: usize,
pub http_endpoint_limit: usize,
pub peer_tags_limit: usize,
pub additional_tags_limit: usize,
}
impl Default for CardinalityLimitConfig {
fn default() -> Self {
Self {
whole_key_limit: 7_000,
resource_limit: 1_024,
http_endpoint_limit: 512,
peer_tags_limit: 512,
additional_tags_limit: 100,
}
}
}
#[derive(Debug, Clone)]
pub struct SpanConcentrator {
bucket_size: u64,
buckets: HashMap<u64, StatsBucket>,
oldest_timestamp: u64,
buffer_len: usize,
cardinality_limits: CardinalityLimitConfig,
span_kinds_stats_computed: Vec<String>,
peer_tag_keys: Vec<String>,
additional_metric_tag_keys: Vec<String>,
big_resource: bool,
#[cfg(feature = "stats-obfuscation")]
obfuscation_config: SharedStatsComputationObfuscationConfig,
}
impl SpanConcentrator {
pub fn new(
bucket_size: Duration,
now: SystemTime,
span_kinds_stats_computed: Vec<String>,
peer_tag_keys: Vec<String>,
override_cardinality_limits: Option<CardinalityLimitConfig>,
additional_metric_tag_keys: Vec<String>,
#[cfg(feature = "stats-obfuscation")] obfuscation_config: Option<
SharedStatsComputationObfuscationConfig,
>,
) -> SpanConcentrator {
if let Some(cardinality_limit_config) = override_cardinality_limits.as_ref() {
if cardinality_limit_config.whole_key_limit == 0
|| cardinality_limit_config.resource_limit == 0
|| cardinality_limit_config.http_endpoint_limit == 0
|| cardinality_limit_config.peer_tags_limit == 0
|| cardinality_limit_config.additional_tags_limit == 0
{
warn!(
?cardinality_limit_config,
"Stats cardinality limit is misconfigured: cardinality limits must not be 0 otherwise all the stats get collapsed!"
);
}
if cardinality_limit_config.whole_key_limit <= cardinality_limit_config.resource_limit
|| cardinality_limit_config.whole_key_limit
<= cardinality_limit_config.http_endpoint_limit
|| cardinality_limit_config.whole_key_limit
<= cardinality_limit_config.peer_tags_limit
|| cardinality_limit_config.whole_key_limit
<= cardinality_limit_config.additional_tags_limit
{
warn!(
?cardinality_limit_config,
"Stats cardinality limit is misconfigured: per-field limits must be lower than whole-key limit otherwise they have no effect and you will get over-collapsed stats!"
);
}
}
SpanConcentrator {
bucket_size: bucket_size.as_nanos() as u64,
buckets: HashMap::new(),
oldest_timestamp: align_timestamp(
system_time_to_unix_duration(now).as_nanos() as u64,
bucket_size.as_nanos() as u64,
),
buffer_len: 2,
cardinality_limits: override_cardinality_limits.unwrap_or_default(),
span_kinds_stats_computed,
peer_tag_keys,
additional_metric_tag_keys: normalize_additional_metric_tag_keys(
additional_metric_tag_keys,
),
big_resource: false,
#[cfg(feature = "stats-obfuscation")]
obfuscation_config: obfuscation_config.unwrap_or_default(),
}
}
pub fn span_kinds(&self) -> &[String] {
&self.span_kinds_stats_computed
}
pub fn set_span_kinds(&mut self, span_kinds: Vec<String>) {
self.span_kinds_stats_computed = span_kinds;
}
pub fn peer_tag_keys(&self) -> &[String] {
&self.peer_tag_keys
}
pub fn set_peer_tags(&mut self, peer_tags: Vec<String>) {
self.peer_tag_keys = peer_tags;
}
pub fn set_big_resource(&mut self, big_resource: bool) {
self.big_resource = big_resource;
}
pub fn additional_metric_tag_keys(&self) -> &[String] {
&self.additional_metric_tag_keys
}
pub fn set_additional_metric_tag_keys(&mut self, tag_keys: Vec<String>) {
self.additional_metric_tag_keys = normalize_additional_metric_tag_keys(tag_keys);
}
pub fn get_bucket_size(&self) -> Duration {
Duration::from_nanos(self.bucket_size)
}
pub fn add_span<'a>(&'a mut self, span: &'a impl StatSpan<'a>) {
if !is_span_eligible(span, self.span_kinds_stats_computed.as_slice()) {
return;
}
let mut bucket_timestamp =
align_timestamp((span.start() + span.duration()) as u64, self.bucket_size);
if bucket_timestamp < self.oldest_timestamp {
bucket_timestamp = self.oldest_timestamp;
}
let target_bucket = self.buckets.entry(bucket_timestamp).or_insert_with(|| {
StatsBucket::new(
bucket_timestamp,
self.cardinality_limits,
#[cfg(feature = "stats-obfuscation")]
self.obfuscation_config.load().enabled,
)
});
#[cfg(feature = "stats-obfuscation")]
let obfuscated_resource = if target_bucket.obfuscated {
Self::compute_obfuscated_span(self.obfuscation_config.load().sql_obfuscation_mode, span)
} else {
None
};
#[cfg(not(feature = "stats-obfuscation"))]
let obfuscated_resource: Option<String> = None;
let agg_key = match obfuscated_resource.as_deref() {
Some(res) => BorrowedAggregationKey::from_obfuscated_span(
res,
span,
self.peer_tag_keys.as_slice(),
self.additional_metric_tag_keys.as_slice(),
),
None => BorrowedAggregationKey::from_span(
span,
self.peer_tag_keys.as_slice(),
self.additional_metric_tag_keys.as_slice(),
),
};
#[cfg(feature = "stats-obfuscation")]
let mut agg_key = agg_key;
#[cfg(feature = "stats-obfuscation")]
if target_bucket.obfuscated {
agg_key.truncate(self.big_resource);
}
target_bucket.insert(
agg_key,
span.duration(),
span.is_error(),
span.has_top_level(),
);
}
#[cfg(feature = "stats-obfuscation")]
fn compute_obfuscated_span<'a>(
sql_obfuscation_mode: libdd_trace_obfuscation::sql::SqlObfuscationMode,
span: &'a impl StatSpan<'a>,
) -> Option<String> {
let dbms_hint: Option<&str> = span.get_meta("db.type");
libdd_trace_obfuscation::obfuscate::obfuscate_resource_for_stats(
span.r#type(),
span.resource(),
dbms_hint,
sql_obfuscation_mode,
)
}
pub fn flush(&mut self, now: SystemTime, force: bool) -> FlushResult<pb::ClientStatsBucket> {
self.drain_due_buckets(now, force, StatsBucket::flush)
}
pub fn flush_with_otlp_exact(&mut self, now: SystemTime, force: bool) -> Vec<OtlpStatsBucket> {
self.drain_due_buckets(now, force, StatsBucket::flush_with_otlp_exact)
.all_buckets()
}
fn drain_due_buckets<T>(
&mut self,
now: SystemTime,
force: bool,
encode: impl Fn(StatsBucket, u64) -> T,
) -> FlushResult<T> {
let now_timestamp = system_time_to_unix_duration(now).as_nanos() as u64;
let buckets: Vec<(u64, StatsBucket)> = self.buckets.drain().collect();
self.oldest_timestamp = if force {
align_timestamp(now_timestamp, self.bucket_size)
} else {
align_timestamp(now_timestamp, self.bucket_size)
- (self.buffer_len as u64 - 1) * self.bucket_size
};
let mut collapsed_spans = 0;
let mut collapsed_fields_metrics = CollapsedFieldsMetrics::zero();
let buckets_pb: Vec<(bool, T)> = buckets
.into_iter()
.filter_map(|(timestamp, bucket)| {
let keep = !force
&& timestamp > (now_timestamp - self.buffer_len as u64 * self.bucket_size);
if keep {
self.buckets.insert(timestamp, bucket);
return None;
}
collapsed_spans += bucket.collapsed_count();
collapsed_fields_metrics += bucket.collapsed_fields_metrics();
#[cfg(feature = "stats-obfuscation")]
let obfuscated = bucket.obfuscated;
#[cfg(not(feature = "stats-obfuscation"))]
let obfuscated = false;
Some((obfuscated, encode(bucket, self.bucket_size)))
})
.collect();
if collapsed_spans > 0 {
debug!(
max_entries_per_bucket = self.cardinality_limits.whole_key_limit,
collapsed_spans,
"Client-side stats values have been collapsed to 'tracer_blocked_value'. This is due to the cardinality exceeding DD_TRACE_STATS_CARDINALITY_LIMIT"
);
}
let mut obfuscated_buckets = Vec::new();
let mut unobfuscated_buckets = Vec::new();
for (obfuscated, bucket) in buckets_pb {
if obfuscated {
obfuscated_buckets.push(bucket);
} else {
unobfuscated_buckets.push(bucket);
}
}
FlushResult {
obfuscated_buckets,
unobfuscated_buckets,
collapsed_spans,
collapsed_fields_metrics,
}
}
}
#[cfg(feature = "stats-obfuscation")]
impl StatsComputationObfuscationConfig {
pub fn disabled() -> SharedStatsComputationObfuscationConfig {
use arc_swap::ArcSwap;
use std::sync::Arc;
Arc::new(ArcSwap::from_pointee(
StatsComputationObfuscationConfig::default(),
))
}
}
#[cfg(test)]
mod tests;