use std::{
sync::{
atomic::{AtomicU64, Ordering},
Arc, Mutex,
},
time,
};
use crate::span_concentrator::{FlushableConcentrator, SpanConcentrator};
use async_trait::async_trait;
use futures::stream::FuturesUnordered;
use futures::StreamExt as _;
use libdd_capabilities::{HttpClientCapability, MaybeSend, SleepCapability};
use libdd_common::{Endpoint, MutexExt};
use libdd_shared_runtime::Worker;
use libdd_trace_protobuf::pb;
use libdd_trace_utils::send_with_retry::{
send_with_retry, CompressionStrategy, RetryBackoffType, RetryStrategy,
};
use libdd_trace_utils::stats_payload_encoder::{
build_stats_payload, encode_stats_payload_msgpack, split_stats_buckets,
MAX_GROUPED_STATS_PER_PAYLOAD,
};
use libdd_trace_utils::trace_utils::TracerHeaderTags;
use libdd_trace_utils::tracer_metadata::TracerMetadata;
use std::fmt::Debug;
use tracing::error;
pub const STATS_ENDPOINT_PATH: &str = "/v0.6/stats";
pub struct StatsRequest {
pub body: Vec<u8>,
pub headers: http::HeaderMap,
pub compression: CompressionStrategy,
pub endpoint: Endpoint,
pub retry: RetryStrategy,
}
#[derive(Debug)]
pub enum StatsDestination {
Agent { endpoint: Endpoint },
Agentless(AgentlessStatsTarget),
}
#[derive(Debug)]
pub struct AgentlessStatsTarget {
pub endpoint: Endpoint,
pub version: String,
}
pub const COLLAPSED_SPANS_HEALTH_METRIC: &str = "datadog.tracer.stats.collapsed_spans";
pub const COLLAPSED_SPANS_TELEMETRY_METRIC: &str = "tracers.stats_collapsed_spans";
#[derive(Clone, Default, Debug)]
pub struct StatsMetadata {
pub hostname: String,
pub env: String,
pub app_version: String,
pub runtime_id: String,
pub language: String,
pub lang_version: String,
pub lang_interpreter: String,
pub lang_vendor: String,
pub tracer_version: String,
pub git_commit_sha: String,
pub process_tags: String,
pub service: String,
pub container_id: String,
}
impl<'a> From<&'a StatsMetadata> for TracerHeaderTags<'a> {
fn from(m: &'a StatsMetadata) -> TracerHeaderTags<'a> {
TracerHeaderTags {
lang: &m.language,
lang_version: &m.lang_version,
lang_interpreter: &m.lang_interpreter,
lang_vendor: &m.lang_vendor,
tracer_version: &m.tracer_version,
..Default::default()
}
}
}
impl From<TracerMetadata> for StatsMetadata {
fn from(m: TracerMetadata) -> StatsMetadata {
StatsMetadata {
hostname: m.hostname,
env: m.env,
app_version: m.app_version,
runtime_id: m.runtime_id,
language: m.language,
lang_version: m.language_version,
lang_interpreter: m.language_interpreter,
lang_vendor: m.language_interpreter_vendor,
tracer_version: m.tracer_version,
git_commit_sha: m.git_commit_sha,
process_tags: m.process_tags,
service: m.service,
container_id: String::new(),
}
}
}
#[derive(Debug)]
pub struct StatsExporter<
Cap: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static,
Con: FlushableConcentrator = SpanConcentrator,
> {
flush_interval: time::Duration,
concentrator: Arc<Mutex<Con>>,
destination: StatsDestination,
meta: StatsMetadata,
sequence_id: AtomicU64,
capabilities: Cap,
#[cfg(feature = "stats-obfuscation")]
supported_obfuscation_version: &'static str,
#[cfg(feature = "telemetry")]
telemetry: Option<(
libdd_telemetry::worker::TelemetryWorkerHandle<Cap>,
libdd_telemetry::metrics::ContextKey,
)>,
#[cfg(feature = "dogstatsd")]
dogstatsd: Option<libdd_dogstatsd_client::DogStatsDClient>,
}
impl<
Cap: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static,
Con: FlushableConcentrator,
> StatsExporter<Cap, Con>
{
#[allow(clippy::too_many_arguments)]
pub fn new(
flush_interval: time::Duration,
concentrator: Arc<Mutex<Con>>,
meta: StatsMetadata,
endpoint: Endpoint,
capabilities: Cap,
#[cfg(feature = "stats-obfuscation")] supported_obfuscation_version: &'static str,
#[cfg(feature = "telemetry")] telemetry: Option<
libdd_telemetry::worker::TelemetryWorkerHandle<Cap>,
>,
#[cfg(feature = "dogstatsd")] dogstatsd: Option<libdd_dogstatsd_client::DogStatsDClient>,
) -> Self {
Self::from_parts(
flush_interval,
concentrator,
meta,
StatsDestination::Agent { endpoint },
capabilities,
#[cfg(feature = "stats-obfuscation")]
supported_obfuscation_version,
#[cfg(feature = "telemetry")]
telemetry,
#[cfg(feature = "dogstatsd")]
dogstatsd,
)
}
#[allow(clippy::too_many_arguments)]
pub fn new_agentless(
flush_interval: time::Duration,
concentrator: Arc<Mutex<Con>>,
meta: StatsMetadata,
target: AgentlessStatsTarget,
capabilities: Cap,
#[cfg(feature = "telemetry")] telemetry: Option<
libdd_telemetry::worker::TelemetryWorkerHandle<Cap>,
>,
#[cfg(feature = "dogstatsd")] dogstatsd: Option<libdd_dogstatsd_client::DogStatsDClient>,
) -> Self {
Self::from_parts(
flush_interval,
concentrator,
meta,
StatsDestination::Agentless(target),
capabilities,
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
telemetry,
#[cfg(feature = "dogstatsd")]
dogstatsd,
)
}
#[allow(clippy::too_many_arguments)]
fn from_parts(
flush_interval: time::Duration,
concentrator: Arc<Mutex<Con>>,
meta: StatsMetadata,
destination: StatsDestination,
capabilities: Cap,
#[cfg(feature = "stats-obfuscation")] supported_obfuscation_version: &'static str,
#[cfg(feature = "telemetry")] telemetry: Option<
libdd_telemetry::worker::TelemetryWorkerHandle<Cap>,
>,
#[cfg(feature = "dogstatsd")] dogstatsd: Option<libdd_dogstatsd_client::DogStatsDClient>,
) -> Self {
#[cfg(feature = "telemetry")]
let telemetry = telemetry.map(|handle| {
let key = handle.register_metric_context(
COLLAPSED_SPANS_TELEMETRY_METRIC.to_string(),
vec![],
libdd_telemetry::data::metrics::MetricType::Count,
true,
libdd_telemetry::data::metrics::MetricNamespace::Tracers,
);
(handle, key)
});
Self {
flush_interval,
concentrator,
destination,
meta,
sequence_id: AtomicU64::new(0),
capabilities,
#[cfg(feature = "stats-obfuscation")]
supported_obfuscation_version,
#[cfg(feature = "telemetry")]
telemetry,
#[cfg(feature = "dogstatsd")]
dogstatsd,
}
}
pub async fn send(&self, force_flush: bool) -> anyhow::Result<bool> {
let flush = {
let mut concentrator = self.concentrator.lock_or_panic();
concentrator.flush_buckets(force_flush)
};
#[cfg(feature = "telemetry")]
if let Some((handle, key)) = &self.telemetry {
if flush.collapsed_spans > 0 {
let _ = handle.add_point(
flush.collapsed_spans as f64,
key,
vec![libdd_common::tag!("collapsed_spans", "whole_key")],
);
}
flush.collapsed_fields_metrics.emit_telemetry(handle, key);
}
#[cfg(feature = "dogstatsd")]
if let Some(client) = &self.dogstatsd {
if flush.collapsed_spans > 0 {
client.send(vec![libdd_dogstatsd_client::DogStatsDAction::Count(
COLLAPSED_SPANS_HEALTH_METRIC,
flush.collapsed_spans as i64,
[libdd_common::tag!("collapsed_spans", "whole_key")].iter(),
)]);
}
flush.collapsed_fields_metrics.emit_dogstatsd(client);
}
let futures = FuturesUnordered::new();
if !flush.obfuscated_buckets.is_empty() {
futures.push(self.send_payload(flush.obfuscated_buckets, true));
}
if !flush.unobfuscated_buckets.is_empty() {
futures.push(self.send_payload(flush.unobfuscated_buckets, false));
}
let sent_stats = !futures.is_empty();
futures
.collect::<Vec<anyhow::Result<()>>>()
.await
.into_iter()
.collect::<anyhow::Result<()>>()?;
Ok(sent_stats)
}
async fn send_payload(
&self,
buckets: Vec<pb::ClientStatsBucket>,
obfuscated: bool,
) -> anyhow::Result<()> {
let groups = split_stats_buckets(buckets, MAX_GROUPED_STATS_PER_PAYLOAD);
let split = groups.len() > 1;
let sequence = self.sequence_id.fetch_add(1, Ordering::Relaxed);
let mut errors = Vec::new();
for group in groups {
if let Err(e) = self
.send_single_payload(group, obfuscated, split, sequence)
.await
{
errors.push(e);
}
}
if let Some(last_err) = errors.pop() {
if !errors.is_empty() {
struct AdditionalErrors(Vec<anyhow::Error>);
impl std::fmt::Display for AdditionalErrors {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
writeln!(f, "with {} additional errors:", self.0.len())?;
for e in &self.0 {
writeln!(f, "{e}")?;
}
Ok(())
}
}
return Err(last_err.context(AdditionalErrors(errors)));
}
return Err(last_err);
}
Ok(())
}
async fn send_single_payload(
&self,
buckets: Vec<pb::ClientStatsBucket>,
obfuscated: bool,
split: bool,
sequence: u64,
) -> anyhow::Result<()> {
let request = match &self.destination {
StatsDestination::Agent { endpoint } => {
self.build_agent_request(endpoint.clone(), sequence, buckets, obfuscated)?
}
StatsDestination::Agentless(target) => {
build_agentless_request(&self.meta, sequence, buckets, target, split)?
}
};
let result = send_with_retry(
&self.capabilities,
&request.endpoint,
request.body,
&request.headers,
&request.retry,
request.compression,
)
.await;
match result {
Ok(_) => Ok(()),
Err(err) => {
error!(?err, "Error with the StatsExporter when sending stats");
anyhow::bail!("Failed to send stats: {err}");
}
}
}
fn build_agent_request(
&self,
endpoint: Endpoint,
sequence: u64,
buckets: Vec<pb::ClientStatsBucket>,
#[cfg_attr(not(feature = "stats-obfuscation"), allow(unused))] obfuscated: bool,
) -> anyhow::Result<StatsRequest> {
let payload = encode_stats_payload(&self.meta, sequence, buckets);
let body = rmp_serde::encode::to_vec_named(&payload)?;
let mut headers: http::HeaderMap = TracerHeaderTags::from(&self.meta).into();
headers.insert(
http::header::CONTENT_TYPE,
libdd_common::header::APPLICATION_MSGPACK,
);
#[cfg(feature = "stats-obfuscation")]
if obfuscated {
headers.insert(
http::HeaderName::from_static("datadog-obfuscation-version"),
http::HeaderValue::from_static(self.supported_obfuscation_version),
);
}
Ok(StatsRequest {
body,
headers,
compression: CompressionStrategy::None,
endpoint,
retry: RetryStrategy::new(0, 0, RetryBackoffType::Constant, None),
})
}
}
const AGENTLESS_STATS_MAX_RETRIES: u32 = 2;
const AGENTLESS_STATS_RETRY_DELAY_MS: u64 = 1000;
fn build_agentless_request(
meta: &StatsMetadata,
sequence: u64,
buckets: Vec<pb::ClientStatsBucket>,
target: &AgentlessStatsTarget,
split: bool,
) -> anyhow::Result<StatsRequest> {
let mut client_payload = encode_stats_payload(meta, sequence, buckets);
client_payload.lang = meta.language.clone();
client_payload.tracer_version = meta.tracer_version.clone();
client_payload.container_id = meta.container_id.clone();
let payload = build_stats_payload(
client_payload,
meta.hostname.clone(),
meta.env.clone(),
target.version.clone(),
split,
);
let body = encode_stats_payload_msgpack(&payload)?;
let mut headers: http::HeaderMap = TracerHeaderTags::from(meta).into();
headers.insert(
http::header::CONTENT_TYPE,
libdd_common::header::APPLICATION_MSGPACK,
);
#[cfg(feature = "compression")]
let compression = CompressionStrategy::Zstd { level: 1 };
#[cfg(not(feature = "compression"))]
let compression = CompressionStrategy::None;
Ok(StatsRequest {
body,
headers,
compression,
endpoint: target.endpoint.clone(),
retry: RetryStrategy::new(
AGENTLESS_STATS_MAX_RETRIES,
AGENTLESS_STATS_RETRY_DELAY_MS,
RetryBackoffType::Exponential,
None,
),
})
}
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
impl<
Cap: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static,
Con: FlushableConcentrator + Send + Debug,
> Worker for StatsExporter<Cap, Con>
{
async fn trigger(&mut self) {
self.capabilities.sleep(self.flush_interval).await;
}
async fn run(&mut self) {
let _ = self.send(false).await; }
fn reset(&mut self) {
let _ = self.concentrator.lock_or_panic().flush_buckets(true);
self.sequence_id.store(0, Ordering::Relaxed);
}
async fn shutdown(&mut self) {
let _ = self.send(true).await;
}
}
fn encode_stats_payload(
meta: &StatsMetadata,
sequence: u64,
buckets: Vec<pb::ClientStatsBucket>,
) -> pb::ClientStatsPayload {
pb::ClientStatsPayload {
hostname: meta.hostname.clone(),
env: if meta.env.is_empty() {
"unknown-env".to_string()
} else {
meta.env.clone()
},
version: meta.app_version.clone(),
runtime_id: meta.runtime_id.clone(),
sequence,
service: meta.service.clone(),
stats: buckets,
git_commit_sha: meta.git_commit_sha.clone(),
process_tags: meta.process_tags.clone(),
container_id: String::new(),
tags: Vec::new(),
agent_aggregation: String::new(),
image_tag: String::new(),
process_tags_hash: 0,
lang: String::new(),
tracer_version: String::new(),
}
}
pub fn stats_url_from_agent_url(agent_url: &str) -> anyhow::Result<http::Uri> {
let mut parts = agent_url.parse::<http::Uri>()?.into_parts();
parts.path_and_query = Some(http::uri::PathAndQuery::from_static(STATS_ENDPOINT_PATH));
Ok(http::Uri::from_parts(parts)?)
}
#[cfg(test)]
mod tests {
use super::*;
use httpmock::prelude::*;
use httpmock::MockServer;
use libdd_capabilities_impl::NativeCapabilities;
use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime, SharedRuntime};
use libdd_trace_utils::span::{trace_utils, v04::SpanSlice};
use libdd_trace_utils::test_utils::{poll_for_mock_hit, poll_for_mock_hits};
use std::borrow::Cow;
use time::Duration;
use time::SystemTime;
fn is_send<T: Send>() {}
fn is_sync<T: Sync>() {}
const BUCKETS_DURATION: Duration = Duration::from_secs(10);
#[test]
fn test_stats_exporter_sync_send() {
let _ = is_send::<StatsExporter<NativeCapabilities>>;
let _ = is_sync::<StatsExporter<NativeCapabilities>>;
}
fn get_test_metadata() -> StatsMetadata {
StatsMetadata {
hostname: "libdatadog-test".into(),
env: "test".into(),
app_version: "0.0.0".into(),
language: "rust".into(),
tracer_version: "0.0.0".into(),
runtime_id: "e39d6d12-0752-489f-b488-cf80006c0378".into(),
process_tags: "key1:value1,key2:value2".into(),
..Default::default()
}
}
fn get_test_concentrator() -> SpanConcentrator {
get_test_concentrator_with_obfuscation_config(
#[cfg(feature = "stats-obfuscation")]
None,
)
}
fn get_test_concentrator_with_obfuscation_config(
#[cfg(feature = "stats-obfuscation")] obfuscation_config: Option<
crate::span_concentrator::SharedStatsComputationObfuscationConfig,
>,
) -> SpanConcentrator {
let mut concentrator = SpanConcentrator::new(
BUCKETS_DURATION,
SystemTime::now() - BUCKETS_DURATION * 3,
vec![],
vec![],
None,
vec![],
#[cfg(feature = "stats-obfuscation")]
obfuscation_config,
);
let mut trace = vec![];
for i in 1..100 {
trace.push(SpanSlice {
service: Cow::Borrowed("libdatadog-test"),
duration: i,
..Default::default()
})
}
trace_utils::compute_top_level_span(trace.as_mut_slice());
for span in trace.iter() {
concentrator.add_span(span);
}
concentrator
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_send_stats() {
let server = MockServer::start_async().await;
let mock = server
.mock_async(|when, then| {
when.method(POST)
.header("Content-type", "application/msgpack")
.path("/v0.6/stats")
.body_includes("libdatadog-test")
.body_includes("key1:value1,key2:value2");
then.status(200).body("");
})
.await;
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let send_status = stats_exporter.send(true).await;
send_status.unwrap();
mock.assert_async().await;
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_send_agentless_stats() {
use super::AgentlessStatsTarget;
let server = MockServer::start_async().await;
let mock = server
.mock_async(|when, then| {
let w = when
.method(POST)
.header("Content-type", "application/msgpack")
.header("dd-api-key", "test-api-key")
.path("/api/v0.2/stats");
#[cfg(feature = "compression")]
let w = w.header("Content-Encoding", "zstd");
#[cfg(not(feature = "compression"))]
let w = w.body_includes("libdatadog-test").body_includes("rust");
let _ = w;
then.status(202).body("");
})
.await;
let target = AgentlessStatsTarget {
endpoint: Endpoint {
api_key: Some("test-api-key".into()),
..Endpoint::from_slice(&server.url("/api/v0.2/stats"))
},
version: "1.2.3-libdatadog".to_string(),
};
let stats_exporter = StatsExporter::<NativeCapabilities>::new_agentless(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
target,
NativeCapabilities::new_client(),
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let send_status = stats_exporter.send(true).await;
send_status.unwrap();
mock.assert_async().await;
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_send_agentless_stats_fail_retries() {
use super::AgentlessStatsTarget;
let server = MockServer::start_async().await;
let mut mock = server
.mock_async(|when, then| {
when.method(POST)
.header("Content-type", "application/msgpack")
.header("dd-api-key", "test-api-key")
.path("/api/v0.2/stats");
then.status(503)
.header("content-type", "application/json")
.body(r#"{"status":"error"}"#);
})
.await;
let target = AgentlessStatsTarget {
endpoint: Endpoint {
api_key: Some("test-api-key".into()),
..Endpoint::from_slice(&server.url("/api/v0.2/stats"))
},
version: "1.2.3-libdatadog".to_string(),
};
let stats_exporter = StatsExporter::<NativeCapabilities>::new_agentless(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
target,
NativeCapabilities::new_client(),
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let send_status = stats_exporter.send(true).await;
send_status.expect_err("agentless stats send should fail after exhausting retries");
assert!(
poll_for_mock_hits(
&mut mock,
80,
100,
(AGENTLESS_STATS_MAX_RETRIES + 1) as usize
)
.await,
"Expected {} attempts (initial + retries) for the agentless intake",
AGENTLESS_STATS_MAX_RETRIES + 1
);
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_send_stats_fail() {
let server = MockServer::start_async().await;
let mut mock = server
.mock_async(|_when, then| {
then.status(503)
.header("content-type", "application/json")
.body(r#"{"status":"error"}"#);
})
.await;
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let send_status = stats_exporter.send(true).await;
send_status.unwrap_err();
assert!(
poll_for_mock_hit(&mut mock, 10, 100, 1, true).await,
"Expected a single attempt with no retries"
);
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_run() {
let shared_runtime = ForkSafeRuntime::new().expect("Failed to create runtime");
let server = MockServer::start();
let mut mock = server.mock(|when, then| {
when.method(POST)
.header("Content-type", "application/msgpack")
.path("/v0.6/stats")
.body_includes("libdatadog-test")
.body_includes("key1:value1,key2:value2");
then.status(200).body("");
});
let caps = NativeCapabilities::new();
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
Duration::from_secs(1),
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
caps.clone(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let _handle = shared_runtime
.spawn_worker(stats_exporter, true)
.expect("Failed to spawn worker");
std::thread::sleep(Duration::from_secs(1));
assert!(
shared_runtime
.block_on(poll_for_mock_hit(&mut mock, 10, 100, 1, false))
.expect("Failed to use runtime"),
"Expected max retry attempts"
);
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_worker_shutdown() {
let shared_runtime = ForkSafeRuntime::new().expect("Failed to create runtime");
let server = MockServer::start();
let mut mock = server.mock(|when, then| {
when.method(POST)
.header("Content-type", "application/msgpack")
.path("/v0.6/stats")
.body_includes("libdatadog-test")
.body_includes("key1:value1,key2:value2");
then.status(200).body("");
});
let buckets_duration = Duration::from_secs(10);
let caps = NativeCapabilities::new();
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
buckets_duration,
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
caps.clone(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let _handle = shared_runtime
.spawn_worker(stats_exporter, true)
.expect("Failed to spawn worker");
shared_runtime.shutdown(None).unwrap();
assert!(
shared_runtime
.block_on(poll_for_mock_hit(&mut mock, 10, 100, 1, false))
.expect("Failed to get runtime"),
"Expected max retry attempts"
);
}
#[test]
fn test_encode_stats_payload_defaults_empty_env() {
let mut meta_with_empty_env = get_test_metadata();
meta_with_empty_env.env = "".to_string();
let buckets = vec![];
let payload = encode_stats_payload(&meta_with_empty_env, 1, buckets.clone());
assert_eq!(
payload.env, "unknown-env",
"Empty env should default to 'unknown-env'"
);
let meta_with_env = get_test_metadata();
let payload_with_env = encode_stats_payload(&meta_with_env, 2, buckets);
assert_eq!(
payload_with_env.env, "test",
"Non-empty env should be preserved"
);
}
#[cfg(feature = "stats-obfuscation")]
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_send_stats_with_obfuscation_header() {
use crate::span_concentrator::StatsComputationObfuscationConfig;
use arc_swap::ArcSwap;
let server = MockServer::start_async().await;
let mock = server
.mock_async(|when, then| {
when.method(POST)
.header("Content-type", "application/msgpack")
.header("datadog-obfuscation-version", "1")
.path("/v0.6/stats")
.body_includes("libdatadog-test");
then.status(200).body("");
})
.await;
let concentrator = get_test_concentrator_with_obfuscation_config(Some(Arc::new(
ArcSwap::from_pointee(StatsComputationObfuscationConfig {
enabled: true,
..Default::default()
}),
)));
let stats_exporter = StatsExporter::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(concentrator)),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
#[cfg(feature = "dogstatsd")]
None,
);
let send_status = stats_exporter.send(true).await;
send_status.unwrap();
mock.assert_async().await;
}
#[cfg(any(feature = "telemetry", feature = "dogstatsd"))]
fn get_collapsed_concentrator(per_key_collapsed: bool) -> SpanConcentrator {
use crate::span_concentrator::CardinalityLimitConfig;
use libdd_trace_utils::span::{
trace_utils,
v04::{SpanSlice, VecMap},
};
let mut cardinality_limit_config = CardinalityLimitConfig {
whole_key_limit: 2, ..Default::default()
};
if per_key_collapsed {
cardinality_limit_config.resource_limit = 1;
cardinality_limit_config.http_endpoint_limit = 1;
}
let mut concentrator = SpanConcentrator::new(
BUCKETS_DURATION,
SystemTime::now(),
vec![],
vec![],
Some(cardinality_limit_config),
vec![],
#[cfg(feature = "stats-obfuscation")]
None,
);
let mut trace = vec![
SpanSlice {
service: Cow::Borrowed("svc-a"),
resource: Cow::Borrowed("resource-a"),
duration: 10,
meta: VecMap::from_iter([(Cow::Borrowed("http.endpoint"), Cow::Borrowed("/"))]),
..Default::default()
},
SpanSlice {
service: Cow::Borrowed("svc-a"),
resource: Cow::Borrowed("resource-b"),
duration: 20,
meta: VecMap::from_iter([(Cow::Borrowed("http.endpoint"), Cow::Borrowed("/"))]),
..Default::default()
},
SpanSlice {
service: Cow::Borrowed("svc-b"),
resource: Cow::Borrowed("resource-c"),
duration: 20,
meta: VecMap::from_iter([(
Cow::Borrowed("http.endpoint"),
Cow::Borrowed("/hello.txt"),
)]),
..Default::default()
},
SpanSlice {
service: Cow::Borrowed("svc-b"),
resource: Cow::Borrowed("resource-b"),
duration: 20,
..Default::default()
},
];
trace_utils::compute_top_level_span(trace.as_mut_slice());
for span in &trace {
concentrator.add_span(span);
}
concentrator
}
#[cfg(feature = "dogstatsd")]
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_no_emission_when_zero() {
use std::net;
let server = MockServer::start_async().await;
server
.mock_async(|_when, then| {
then.status(200).body("");
})
.await;
let socket = net::UdpSocket::bind("127.0.0.1:0").expect("failed to bind UDP socket");
socket
.set_read_timeout(Some(std::time::Duration::from_millis(200)))
.unwrap();
let addr = socket.local_addr().unwrap().to_string();
let dogstatsd_client =
libdd_dogstatsd_client::DogStatsDClient::new(libdd_common::Endpoint::from_slice(&addr))
.expect("failed to create dogstatsd client");
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_test_concentrator())),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
Some(dogstatsd_client),
);
stats_exporter.send(true).await.unwrap();
let mut buf = [0u8; 256];
let result = socket.recv(&mut buf);
assert!(
result.is_err(),
"No DogStatsD datagram expected when collapsed_spans == 0. Got {}",
std::str::from_utf8(&buf[..result.unwrap()]).unwrap()
);
}
#[cfg(feature = "dogstatsd")]
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_collapsed_spans_dogstatsd() {
use std::net;
let server = MockServer::start_async().await;
server
.mock_async(|_when, then| {
then.status(200).body("");
})
.await;
let socket = net::UdpSocket::bind("127.0.0.1:0").expect("failed to bind UDP socket");
socket
.set_read_timeout(Some(std::time::Duration::from_millis(500)))
.unwrap();
let addr = socket.local_addr().unwrap().to_string();
let dogstatsd_client =
libdd_dogstatsd_client::DogStatsDClient::new(libdd_common::Endpoint::from_slice(&addr))
.expect("failed to create dogstatsd client");
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_collapsed_concentrator(false))),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
Some(dogstatsd_client),
);
stats_exporter.send(true).await.unwrap();
let mut buf = [0u8; 256];
let n = socket
.recv(&mut buf)
.expect("expected a DogStatsD datagram");
let datagram = std::str::from_utf8(&buf[..n]).expect("valid utf-8");
assert_eq!(
datagram, "datadog.tracer.stats.collapsed_spans:2|c|#collapsed_spans:whole_key",
"DogStatsD datagram must match the expected format"
);
}
#[cfg(feature = "dogstatsd")]
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_collapsed_spans_per_key_dogstatsd() {
use std::net;
let server = MockServer::start_async().await;
server
.mock_async(|_when, then| {
then.status(200).body("");
})
.await;
let socket = net::UdpSocket::bind("127.0.0.1:0").expect("failed to bind UDP socket");
socket
.set_read_timeout(Some(std::time::Duration::from_millis(500)))
.unwrap();
let addr = socket.local_addr().unwrap().to_string();
let dogstatsd_client =
libdd_dogstatsd_client::DogStatsDClient::new(libdd_common::Endpoint::from_slice(&addr))
.expect("failed to create dogstatsd client");
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_collapsed_concentrator(true))),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
None,
Some(dogstatsd_client),
);
stats_exporter.send(true).await.unwrap();
let mut buf = [0u8; 256];
let n = socket
.recv(&mut buf)
.expect("expected a DogStatsD datagram");
let datagram = std::str::from_utf8(&buf[..n]).expect("valid utf-8");
assert_eq!(
datagram, "datadog.tracer.stats.collapsed_spans:2|c|#collapsed_spans:whole_key",
"DogStatsD datagram must match the expected format"
);
let n = socket
.recv(&mut buf)
.expect("expected a DogStatsD datagram");
let datagram = std::str::from_utf8(&buf[..n]).expect("valid utf-8");
assert_eq!(
datagram, "datadog.tracer.stats.collapsed_spans:1|c|#collapsed_spans:resource",
"DogStatsD datagram must match the expected format"
);
let n = socket
.recv(&mut buf)
.expect("expected a DogStatsD datagram");
let datagram = std::str::from_utf8(&buf[..n]).expect("valid utf-8");
assert_eq!(
datagram, "datadog.tracer.stats.collapsed_spans:2|c|#collapsed_spans:resource,collapsed_spans:http_endpoint",
"DogStatsD datagram must match the expected format"
);
}
#[cfg(feature = "telemetry")]
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn test_collapsed_spans_telemetry() {
use libdd_telemetry::worker::TelemetryWorkerBuilder;
let server = MockServer::start_async().await;
server
.mock_async(|_when, then| {
then.status(200).body("");
})
.await;
let (handle, _join_handle) = TelemetryWorkerBuilder::new(
"test-host".to_string(),
"test-service".to_string(),
"rust".to_string(),
"1.0".to_string(),
"0.0.0".to_string(),
)
.spawn();
let stats_exporter = StatsExporter::<NativeCapabilities>::new(
BUCKETS_DURATION,
Arc::new(Mutex::new(get_collapsed_concentrator(true))),
get_test_metadata(),
Endpoint::from_url(stats_url_from_agent_url(&server.url("/")).unwrap()),
NativeCapabilities::new_client(),
#[cfg(feature = "stats-obfuscation")]
"1",
#[cfg(feature = "telemetry")]
Some(handle),
#[cfg(feature = "dogstatsd")]
None,
);
stats_exporter.send(true).await.unwrap();
let stats_exporter_ref = &stats_exporter;
let (handle_ref, _key) = stats_exporter_ref
.telemetry
.as_ref()
.expect("telemetry must be set");
let receiver = handle_ref.stats().expect("failed to request stats");
let stats = receiver.await.expect("failed to receive stats");
assert_eq!(
stats.metric_contexts, 1,
"exactly one metric context (COLLAPSED_SPANS_METRIC) should be registered"
);
assert_eq!(
stats.metric_buckets.buckets, 3,
"exactly 3 metric bucket expected after one whole-key collapsed-spans, one resource key collapsed-span and one resource key+http_endpoint key emissions"
);
}
}