use std::collections::HashMap;
use std::time::{self, Duration, SystemTime};
use libdd_trace_protobuf::pb;
use aggregation::StatsBucket;
mod aggregation;
use aggregation::BorrowedAggregationKey;
pub use aggregation::FixedAggregationKey;
pub mod stat_span;
pub use stat_span::StatSpan;
pub trait FlushableConcentrator {
fn flush_buckets(&mut self, force: bool) -> Vec<pb::ClientStatsBucket>;
}
impl FlushableConcentrator for SpanConcentrator {
fn flush_buckets(&mut self, force: bool) -> Vec<pb::ClientStatsBucket> {
self.flush(SystemTime::now(), force)
}
}
fn system_time_to_unix_duration(t: SystemTime) -> Duration {
t.duration_since(time::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)]
pub struct SpanConcentrator {
bucket_size: u64,
buckets: HashMap<u64, StatsBucket>,
oldest_timestamp: u64,
buffer_len: usize,
span_kinds_stats_computed: Vec<String>,
peer_tag_keys: Vec<String>,
#[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>,
#[cfg(feature = "stats-obfuscation")] obfuscation_config: Option<
SharedStatsComputationObfuscationConfig,
>,
) -> SpanConcentrator {
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,
span_kinds_stats_computed,
peer_tag_keys,
#[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 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 obfuscated_resource = self.compute_obfuscated_span(span);
let agg_key = match obfuscated_resource.as_deref() {
Some(res) => BorrowedAggregationKey::from_obfuscated_span(
res,
span,
self.peer_tag_keys.as_slice(),
),
None => BorrowedAggregationKey::from_span(span, self.peer_tag_keys.as_slice()),
};
self.buckets
.entry(bucket_timestamp)
.or_insert(StatsBucket::new(bucket_timestamp))
.insert(
agg_key,
span.duration(),
span.is_error(),
span.has_top_level(),
);
}
fn compute_obfuscated_span<'a>(
&self,
#[allow(unused)] span: &'a impl StatSpan<'a>,
) -> Option<String> {
#[cfg(feature = "stats-obfuscation")]
if self.obfuscation_config.load().enabled {
let dbms_hint: Option<&str> = span.get_meta("db.type");
return libdd_trace_obfuscation::obfuscate::obfuscate_resource_for_stats(
span.r#type(),
span.resource(),
dbms_hint,
self.obfuscation_config.load().sql_obfuscation_mode,
);
}
None
}
pub fn flush(&mut self, now: SystemTime, force: bool) -> Vec<pb::ClientStatsBucket> {
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
};
buckets
.into_iter()
.filter_map(|(timestamp, bucket)| {
if !force && timestamp > (now_timestamp - self.buffer_len as u64 * self.bucket_size)
{
self.buckets.insert(timestamp, bucket);
return None;
}
Some(bucket.flush(self.bucket_size))
})
.collect()
}
}
#[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;