libdd_trace_stats/span_concentrator/
mod.rs1mod aggregation;
5pub mod cardinality_limit_telemetry;
6pub mod stat_span;
7
8use std::collections::HashMap;
9use std::time::Duration;
10use tracing::{debug, warn};
11use web_time::{SystemTime, UNIX_EPOCH};
13
14use libdd_trace_protobuf::pb;
15
16use aggregation::StatsBucket;
17
18use aggregation::BorrowedAggregationKey;
19pub use aggregation::{FixedAggregationKey, OtlpExactCell, OtlpExactGroup, OtlpStatsBucket};
20use cardinality_limit_telemetry::CollapsedFieldsMetrics;
21
22pub use stat_span::{ChunkSpanView, StatSpan};
23
24const ADDITIONAL_METRIC_TAGS_MAX_KEYS: usize = 4;
25
26fn normalize_additional_metric_tag_keys(mut keys: Vec<String>) -> Vec<String> {
29 keys.sort_unstable();
30 keys.dedup();
31 if keys.len() > ADDITIONAL_METRIC_TAGS_MAX_KEYS {
32 let dropped = keys.split_off(ADDITIONAL_METRIC_TAGS_MAX_KEYS);
33 warn!(
34 "additional_metric_tag_keys: {} additional metric tag keys exceed the cap of {}; dropping: {:?}",
35 dropped.len() + ADDITIONAL_METRIC_TAGS_MAX_KEYS,
36 ADDITIONAL_METRIC_TAGS_MAX_KEYS,
37 dropped,
38 );
39 }
40 keys
41}
42
43pub struct FlushResult<T> {
48 pub obfuscated_buckets: Vec<T>,
50 pub unobfuscated_buckets: Vec<T>,
52 pub collapsed_spans: u64,
55 pub collapsed_fields_metrics: CollapsedFieldsMetrics,
56}
57
58impl<T> FlushResult<T> {
59 pub fn all_buckets(self) -> Vec<T> {
61 let mut buckets = self.obfuscated_buckets;
62 buckets.extend(self.unobfuscated_buckets);
63 buckets
64 }
65}
66
67pub trait FlushableConcentrator {
72 fn flush_buckets(&mut self, force: bool) -> FlushResult<pb::ClientStatsBucket>;
75}
76
77impl FlushableConcentrator for SpanConcentrator {
78 fn flush_buckets(&mut self, force: bool) -> FlushResult<pb::ClientStatsBucket> {
79 self.flush(SystemTime::now(), force)
80 }
81}
82
83fn system_time_to_unix_duration(t: SystemTime) -> Duration {
86 t.duration_since(UNIX_EPOCH)
87 .unwrap_or(Duration::from_nanos(0))
88}
89
90#[inline]
92fn align_timestamp(t: u64, bucket_size: u64) -> u64 {
93 t - (t % bucket_size)
94}
95
96pub fn is_span_eligible<'a, T>(span: &'a T, span_kinds_stats_computed: &[String]) -> bool
98where
99 T: StatSpan<'a>,
100{
101 (span.has_top_level() || span.is_measured() || {
102 span.get_meta("span.kind")
103 .is_some_and(|span_kind| span_kinds_stats_computed.contains(&span_kind.to_lowercase()))
104 }) && !span.is_partial_snapshot()
105}
106
107#[cfg(feature = "stats-obfuscation")]
108#[derive(Clone, Debug, Default)]
109#[cfg_attr(target_arch = "wasm32", allow(dead_code))]
110pub struct StatsComputationObfuscationConfig {
111 pub enabled: bool,
112 pub sql_obfuscation_mode: libdd_trace_obfuscation::sql::SqlObfuscationMode,
113}
114
115#[cfg(feature = "stats-obfuscation")]
116pub type SharedStatsComputationObfuscationConfig =
117 std::sync::Arc<arc_swap::ArcSwap<StatsComputationObfuscationConfig>>;
118
119#[derive(Debug, Clone, Copy, PartialEq, Eq)]
121#[repr(C)]
122pub struct CardinalityLimitConfig {
123 pub whole_key_limit: usize,
125 pub resource_limit: usize,
127 pub http_endpoint_limit: usize,
129 pub peer_tags_limit: usize,
131 pub additional_tags_limit: usize,
133}
134
135impl Default for CardinalityLimitConfig {
136 fn default() -> Self {
137 Self {
138 whole_key_limit: 7_000,
146 resource_limit: 1_024,
148 http_endpoint_limit: 512,
149 peer_tags_limit: 512,
150 additional_tags_limit: 100,
151 }
152 }
153}
154
155#[derive(Debug, Clone)]
176pub struct SpanConcentrator {
177 bucket_size: u64,
179 buckets: HashMap<u64, StatsBucket>,
180 oldest_timestamp: u64,
183 buffer_len: usize,
185 cardinality_limits: CardinalityLimitConfig,
187 span_kinds_stats_computed: Vec<String>,
189 peer_tag_keys: Vec<String>,
191 additional_metric_tag_keys: Vec<String>,
193 big_resource: bool,
195 #[cfg(feature = "stats-obfuscation")]
196 obfuscation_config: SharedStatsComputationObfuscationConfig,
197}
198
199impl SpanConcentrator {
200 pub fn new(
210 bucket_size: Duration,
211 now: SystemTime,
212 span_kinds_stats_computed: Vec<String>,
213 peer_tag_keys: Vec<String>,
214 override_cardinality_limits: Option<CardinalityLimitConfig>,
215 additional_metric_tag_keys: Vec<String>,
216 #[cfg(feature = "stats-obfuscation")] obfuscation_config: Option<
217 SharedStatsComputationObfuscationConfig,
218 >,
219 ) -> SpanConcentrator {
220 if let Some(cardinality_limit_config) = override_cardinality_limits.as_ref() {
221 if cardinality_limit_config.whole_key_limit == 0
222 || cardinality_limit_config.resource_limit == 0
223 || cardinality_limit_config.http_endpoint_limit == 0
224 || cardinality_limit_config.peer_tags_limit == 0
225 || cardinality_limit_config.additional_tags_limit == 0
226 {
227 warn!(
228 ?cardinality_limit_config,
229 "Stats cardinality limit is misconfigured: cardinality limits must not be 0 otherwise all the stats get collapsed!"
230 );
231 }
232 if cardinality_limit_config.whole_key_limit <= cardinality_limit_config.resource_limit
233 || cardinality_limit_config.whole_key_limit
234 <= cardinality_limit_config.http_endpoint_limit
235 || cardinality_limit_config.whole_key_limit
236 <= cardinality_limit_config.peer_tags_limit
237 || cardinality_limit_config.whole_key_limit
238 <= cardinality_limit_config.additional_tags_limit
239 {
240 warn!(
241 ?cardinality_limit_config,
242 "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!"
243 );
244 }
245 }
246 SpanConcentrator {
247 bucket_size: bucket_size.as_nanos() as u64,
248 buckets: HashMap::new(),
249 oldest_timestamp: align_timestamp(
250 system_time_to_unix_duration(now).as_nanos() as u64,
251 bucket_size.as_nanos() as u64,
252 ),
253 buffer_len: 2,
254 cardinality_limits: override_cardinality_limits.unwrap_or_default(),
255 span_kinds_stats_computed,
256 peer_tag_keys,
257 additional_metric_tag_keys: normalize_additional_metric_tag_keys(
258 additional_metric_tag_keys,
259 ),
260 big_resource: false,
261 #[cfg(feature = "stats-obfuscation")]
262 obfuscation_config: obfuscation_config.unwrap_or_default(),
263 }
264 }
265
266 pub fn span_kinds(&self) -> &[String] {
268 &self.span_kinds_stats_computed
269 }
270
271 pub fn set_span_kinds(&mut self, span_kinds: Vec<String>) {
273 self.span_kinds_stats_computed = span_kinds;
274 }
275
276 pub fn peer_tag_keys(&self) -> &[String] {
278 &self.peer_tag_keys
279 }
280
281 pub fn set_peer_tags(&mut self, peer_tags: Vec<String>) {
283 self.peer_tag_keys = peer_tags;
284 }
285
286 pub fn set_big_resource(&mut self, big_resource: bool) {
288 self.big_resource = big_resource;
289 }
290
291 pub fn additional_metric_tag_keys(&self) -> &[String] {
293 &self.additional_metric_tag_keys
294 }
295
296 pub fn set_additional_metric_tag_keys(&mut self, tag_keys: Vec<String>) {
298 self.additional_metric_tag_keys = normalize_additional_metric_tag_keys(tag_keys);
299 }
300
301 pub fn get_bucket_size(&self) -> Duration {
303 Duration::from_nanos(self.bucket_size)
304 }
305
306 pub fn add_span<'a>(&'a mut self, span: &'a impl StatSpan<'a>) {
309 if !is_span_eligible(span, self.span_kinds_stats_computed.as_slice()) {
310 return;
311 }
312 let mut bucket_timestamp =
313 align_timestamp((span.start() + span.duration()) as u64, self.bucket_size);
314 if bucket_timestamp < self.oldest_timestamp {
317 bucket_timestamp = self.oldest_timestamp;
318 }
319
320 let target_bucket = self.buckets.entry(bucket_timestamp).or_insert_with(|| {
321 StatsBucket::new(
322 bucket_timestamp,
323 self.cardinality_limits,
324 #[cfg(feature = "stats-obfuscation")]
325 self.obfuscation_config.load().enabled,
326 )
327 });
328 #[cfg(feature = "stats-obfuscation")]
329 let obfuscated_resource = if target_bucket.obfuscated {
330 Self::compute_obfuscated_span(self.obfuscation_config.load().sql_obfuscation_mode, span)
331 } else {
332 None
333 };
334 #[cfg(not(feature = "stats-obfuscation"))]
335 let obfuscated_resource: Option<String> = None;
336 let agg_key = match obfuscated_resource.as_deref() {
337 Some(res) => BorrowedAggregationKey::from_obfuscated_span(
338 res,
339 span,
340 self.peer_tag_keys.as_slice(),
341 self.additional_metric_tag_keys.as_slice(),
342 ),
343 None => BorrowedAggregationKey::from_span(
344 span,
345 self.peer_tag_keys.as_slice(),
346 self.additional_metric_tag_keys.as_slice(),
347 ),
348 };
349 #[cfg(feature = "stats-obfuscation")]
351 let mut agg_key = agg_key;
352 #[cfg(feature = "stats-obfuscation")]
353 if target_bucket.obfuscated {
354 agg_key.truncate(self.big_resource);
355 }
356 target_bucket.insert(
357 agg_key,
358 span.duration(),
359 span.is_error(),
360 span.has_top_level(),
361 );
362 }
363
364 #[cfg(feature = "stats-obfuscation")]
365 fn compute_obfuscated_span<'a>(
366 sql_obfuscation_mode: libdd_trace_obfuscation::sql::SqlObfuscationMode,
367 span: &'a impl StatSpan<'a>,
368 ) -> Option<String> {
369 let dbms_hint: Option<&str> = span.get_meta("db.type");
370 libdd_trace_obfuscation::obfuscate::obfuscate_resource_for_stats(
371 span.r#type(),
372 span.resource(),
373 dbms_hint,
374 sql_obfuscation_mode,
375 )
376 }
377
378 pub fn flush(&mut self, now: SystemTime, force: bool) -> FlushResult<pb::ClientStatsBucket> {
383 self.drain_due_buckets(now, force, StatsBucket::flush)
384 }
385
386 pub fn flush_with_otlp_exact(&mut self, now: SystemTime, force: bool) -> Vec<OtlpStatsBucket> {
390 self.drain_due_buckets(now, force, StatsBucket::flush_with_otlp_exact)
391 .all_buckets()
392 }
393
394 fn drain_due_buckets<T>(
401 &mut self,
402 now: SystemTime,
403 force: bool,
404 encode: impl Fn(StatsBucket, u64) -> T,
405 ) -> FlushResult<T> {
406 let now_timestamp = system_time_to_unix_duration(now).as_nanos() as u64;
408 let buckets: Vec<(u64, StatsBucket)> = self.buckets.drain().collect();
409 self.oldest_timestamp = if force {
410 align_timestamp(now_timestamp, self.bucket_size)
411 } else {
412 align_timestamp(now_timestamp, self.bucket_size)
413 - (self.buffer_len as u64 - 1) * self.bucket_size
414 };
415 let mut collapsed_spans = 0;
416 let mut collapsed_fields_metrics = CollapsedFieldsMetrics::zero();
417 let buckets_pb: Vec<(bool, T)> = buckets
418 .into_iter()
419 .filter_map(|(timestamp, bucket)| {
420 let keep = !force
428 && timestamp > (now_timestamp - self.buffer_len as u64 * self.bucket_size);
429 if keep {
430 self.buckets.insert(timestamp, bucket);
431 return None;
432 }
433 collapsed_spans += bucket.collapsed_count();
434 collapsed_fields_metrics += bucket.collapsed_fields_metrics();
435 #[cfg(feature = "stats-obfuscation")]
436 let obfuscated = bucket.obfuscated;
437 #[cfg(not(feature = "stats-obfuscation"))]
438 let obfuscated = false;
439 Some((obfuscated, encode(bucket, self.bucket_size)))
440 })
441 .collect();
442 if collapsed_spans > 0 {
443 debug!(
444 max_entries_per_bucket = self.cardinality_limits.whole_key_limit,
445 collapsed_spans,
446 "Client-side stats values have been collapsed to 'tracer_blocked_value'. This is due to the cardinality exceeding DD_TRACE_STATS_CARDINALITY_LIMIT"
447 );
448 }
449
450 let mut obfuscated_buckets = Vec::new();
451 let mut unobfuscated_buckets = Vec::new();
452 for (obfuscated, bucket) in buckets_pb {
453 if obfuscated {
454 obfuscated_buckets.push(bucket);
455 } else {
456 unobfuscated_buckets.push(bucket);
457 }
458 }
459
460 FlushResult {
461 obfuscated_buckets,
462 unobfuscated_buckets,
463 collapsed_spans,
464 collapsed_fields_metrics,
465 }
466 }
467}
468
469#[cfg(feature = "stats-obfuscation")]
470impl StatsComputationObfuscationConfig {
471 pub fn disabled() -> SharedStatsComputationObfuscationConfig {
472 use arc_swap::ArcSwap;
473 use std::sync::Arc;
474
475 Arc::new(ArcSwap::from_pointee(
476 StatsComputationObfuscationConfig::default(),
477 ))
478 }
479}
480
481#[cfg(test)]
482mod tests;