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
8 changes: 5 additions & 3 deletions datadog-ipc/src/shm_stats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,8 @@ use zwohash::ZwoHasher;
use libdd_ddsketch::DDSketch;
use libdd_trace_protobuf::pb;
use libdd_trace_stats::span_concentrator::{
FixedAggregationKey, FlushResult, FlushableConcentrator,
cardinality_limit_telemetry::CollapsedFieldsMetrics, FixedAggregationKey, FlushResult,
FlushableConcentrator,
};

use crate::platform::{FileBackedHandle, MappedMem, NamedShmHandle};
Expand Down Expand Up @@ -822,12 +823,13 @@ impl ShmSpanConcentrator {
}

impl FlushableConcentrator for ShmSpanConcentrator {
fn flush_buckets(&mut self, force: bool) -> FlushResult {
// The SHM concentrator does not perform client-side obfuscation.
fn flush_buckets(&mut self, force: bool) -> FlushResult<pb::ClientStatsBucket> {
// The SHM concentrator does not perform client-side obfuscation nor emits telemetry.
Comment thread
Eldolfin marked this conversation as resolved.
FlushResult {
obfuscated_buckets: vec![],
unobfuscated_buckets: self.drain_buckets(force),
collapsed_spans: 0,
collapsed_fields_metrics: CollapsedFieldsMetrics::zero(),
}
}
}
Expand Down
23 changes: 20 additions & 3 deletions libdd-trace-stats/src/span_concentrator/aggregation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,12 @@ use std::{
};
use tracing::warn;

use crate::span_concentrator::StatSpan;
use crate::span_concentrator::{cardinality_limit_telemetry::CollapsedFieldSet, StatSpan};

use super::CardinalityLimitConfig;
use super::{
cardinality_limit_telemetry::{self, CollapsedFieldsMetrics},
CardinalityLimitConfig,
};

/// Sentinel value used for cardinality limiting.
pub const TRACER_BLOCKED_VALUE: &str = "tracer_blocked_value";
Expand Down Expand Up @@ -527,8 +530,9 @@ pub(super) struct StatsBucket {
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.
/// Number of spans collapsed into the overflow bucket due to whole-key cardinality limiting.
collapsed_count: u64,
collapsed_fields_metrics: CollapsedFieldsMetrics,
/// Indicates if stats obfuscated in this bucket. This is set once at creation and stays
/// constant per bucket
#[cfg(feature = "stats-obfuscation")]
Expand All @@ -555,9 +559,15 @@ impl StatsBucket {
distinct_http_endpoints: HashSet::new(),
distinct_peer_tags: HashSet::new(),
distinct_additional_tags: HashSet::new(),
collapsed_fields_metrics: cardinality_limit_telemetry::CollapsedFieldsMetrics::zero(),
}
}

/// Returns metrics on spans field collapse with reasons.
pub fn collapsed_fields_metrics(&self) -> cardinality_limit_telemetry::CollapsedFieldsMetrics {
self.collapsed_fields_metrics
}

/// Return the number of spans collapsed into the overflow bucket.
pub(super) fn collapsed_count(&self) -> u64 {
self.collapsed_count
Expand Down Expand Up @@ -614,11 +624,14 @@ impl StatsBucket {
hasher.finish()
}

let mut collapsed_fields = CollapsedFieldSet::empty();

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;
collapsed_fields.add(CollapsedFieldSet::RESOURCE_NAME);
} else {
slot.insert();
}
Expand All @@ -629,6 +642,7 @@ impl StatsBucket {
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;
collapsed_fields.add(CollapsedFieldSet::HTTP_ENDPOINT);
} else {
slot.insert();
}
Expand All @@ -639,6 +653,7 @@ impl StatsBucket {
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(""))];
collapsed_fields.add(CollapsedFieldSet::PEER_TAGS);
} else {
slot.insert();
}
Expand All @@ -649,10 +664,12 @@ impl StatsBucket {
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, "")];
collapsed_fields.add(CollapsedFieldSet::ADDITIONAL_TAGS);
} else {
slot.insert();
}
}
self.collapsed_fields_metrics.increment(collapsed_fields);
}

/// Consume the bucket and return a ClientStatsBucket containing the bucket stats.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
// Copyright 2024-Present Datadog, Inc. https://www.datadoghq.com/
// SPDX-License-Identifier: Apache-2.0

//! This module implement the logic for sending telemetry/dogstatsd related to cardinality limits
// Because the cardinality limit RFC requires one point of telemetry with a set of collapsed fields
// and not one point for each collapsed field, this solution was found as an alternative to adding
// telemetry and dogstatsd client in the SpanConcentrator to keep it as "pure computation logic"

#[cfg(feature = "telemetry")]
use libdd_capabilities::{HttpClientCapability, MaybeSend, SleepCapability};
use libdd_common::tag::const_assert;

/// Bitset of collapsed stats key fields
pub struct CollapsedFieldSet(usize);
impl CollapsedFieldSet {
pub const RESOURCE_NAME: usize = 1 << 0;
pub const HTTP_ENDPOINT: usize = 1 << 1;
pub const PEER_TAGS: usize = 1 << 2;
pub const ADDITIONAL_TAGS: usize = 1 << 3;

const FIELDS: [usize; 4] = [
Self::RESOURCE_NAME,
Self::HTTP_ENDPOINT,
Self::PEER_TAGS,
Self::ADDITIONAL_TAGS,
];

/// Default value with none of the fields set (no collapsed fields)
pub fn empty() -> CollapsedFieldSet {
Self(0)
}

/// Enable the bit corresponding to `field`
pub fn add(&mut self, field: usize) {
debug_assert!(Self::FIELDS.contains(&field));
self.0 |= field;
}
}

const COLLAPSED_FIELD_METRIC_SIZE: usize = 1 << CollapsedFieldSet::FIELDS.len();
// Verify the array is of a reasonable size
const_assert!(COLLAPSED_FIELD_METRIC_SIZE <= 16);
/// Counter of combination of collapsed stats key fields
// Note: slot 0 is a counter for non_collapsed spans. It's not used for emitting telemetry
#[derive(Debug, Clone, Default, Copy)]
pub struct CollapsedFieldsMetrics([usize; COLLAPSED_FIELD_METRIC_SIZE]);

impl CollapsedFieldsMetrics {
/// Default value every combination of collapsed fields set to 0
pub fn zero() -> Self {
Self::default()
}

#[cfg(feature = "dogstatsd")]
pub fn emit_dogstatsd(&self, dogstatsd: &libdd_dogstatsd_client::DogStatsDClient) {
// skip the first slot that is used to count span which have no collapsed fields
for (mask, &count) in self.0.iter().enumerate().skip(1) {
if count > 0 {
let tags = Self::fields_mask_to_list(mask);
dogstatsd.send(vec![libdd_dogstatsd_client::DogStatsDAction::Count(
"datadog.tracer.stats.collapsed_spans",
count as i64,
tags.iter(),
)]);
}
}
}

#[cfg(feature = "telemetry")]
pub fn emit_telemetry<
Cap: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static,
>(
&self,
handle: &libdd_telemetry::worker::TelemetryWorkerHandle<Cap>,
context_key: &libdd_telemetry::metrics::ContextKey,
) {
// skip the first slot that is used to count span which have no collapsed fields
for (mask, &count) in self.0.iter().enumerate().skip(1) {
if count > 0 {
let tags = Self::fields_mask_to_list(mask);
let _ = handle.add_point(count as f64, context_key, tags);
}
}
}

/// Given a bitmask of collapsed fields, returns the list of tags to attach to
/// telemetry/dogstatsd
#[cfg(any(feature = "telemetry", feature = "dogstatsd"))]
fn fields_mask_to_list(mask: usize) -> Vec<libdd_common::tag::Tag> {
let mut tags = Vec::new();
for field_pow in 0..CollapsedFieldSet::FIELDS.len() {
let field_value = 1 << field_pow;
debug_assert!(
CollapsedFieldSet::FIELDS.contains(&field_value),
"{field_value} is an invalid value for a CollapsedFieldSet"
);
let has_field = (mask & field_value) != 0;
if !has_field {
continue;
}
let field_tag = match field_value {
CollapsedFieldSet::RESOURCE_NAME => {
libdd_common::tag!("collapsed_spans", "resource")
}
CollapsedFieldSet::HTTP_ENDPOINT => {
libdd_common::tag!("collapsed_spans", "http_endpoint")
}
CollapsedFieldSet::PEER_TAGS => {
libdd_common::tag!("collapsed_spans", "peer_tags")
}
CollapsedFieldSet::ADDITIONAL_TAGS => {
libdd_common::tag!("collapsed_spans", "additional_metric_tags")
}
// Should be unreachable, but don't fail in prod if provided with an invalid field
// set
_ => continue,
};
tags.push(field_tag);
}
debug_assert!(!tags.is_empty());
tags
}

/// Increment the telemetry counter corresponding to this collapsed field combination
pub fn increment(&mut self, field_set: CollapsedFieldSet) {
self.0[field_set.0] += 1;
}
}

impl std::ops::AddAssign for CollapsedFieldsMetrics {
fn add_assign(&mut self, rhs: Self) {
for i in 0..self.0.len() {
self.0[i] += rhs.0[i];
}
}
}
Loading
Loading