From befd1d0541a6d0d3801ff8a171f2fede72db90fd Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Mon, 10 Aug 2026 14:53:41 -0400 Subject: [PATCH 1/7] feat(data-pipeline): emit trace export telemetry metrics --- libdd-data-pipeline/src/agentless/exporter.rs | 15 +- libdd-data-pipeline/src/agentless/mod.rs | 1 - libdd-data-pipeline/src/otlp/exporter.rs | 50 +++- libdd-data-pipeline/src/otlp/mod.rs | 1 - libdd-data-pipeline/src/telemetry/metrics.rs | 26 ++ libdd-data-pipeline/src/telemetry/mod.rs | 104 ++++++- libdd-data-pipeline/src/trace_exporter/mod.rs | 265 ++++++++++++++++-- 7 files changed, 413 insertions(+), 49 deletions(-) diff --git a/libdd-data-pipeline/src/agentless/exporter.rs b/libdd-data-pipeline/src/agentless/exporter.rs index 5ac9d05209..a0ca3e44a3 100644 --- a/libdd-data-pipeline/src/agentless/exporter.rs +++ b/libdd-data-pipeline/src/agentless/exporter.rs @@ -10,6 +10,7 @@ use libdd_capabilities::{HttpClientCapability, SleepCapability}; use libdd_common::Endpoint; use libdd_trace_utils::send_with_retry::{ send_with_retry, CompressionStrategy, RetryBackoffType, RetryStrategy, SendWithRetryError, + SendWithRetryResult, }; use tracing::error; @@ -21,12 +22,12 @@ const AGENTLESS_RETRY_DELAY_MS: u64 = 1000; /// `headers` should already contain all required headers (api key, content-type, meta-*, /// entity, trace-count, etc.). `test_token` is forwarded as `X-Datadog-Test-Session-Token` /// when set, enabling snapshot tests against a local mock. -pub async fn send_agentless_traces_http( +pub(crate) async fn send_agentless_traces_http( capabilities: &C, config: &AgentlessTraceConfig, headers: HeaderMap, json_body: Vec, -) -> Result<(), TraceExporterError> { +) -> Result { let url = libdd_common::parse_uri(&config.endpoint_url).map_err(|e| { TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "Invalid agentless endpoint URL: {e}" @@ -51,7 +52,7 @@ pub async fn send_agentless_traces_http Ok(()), - Err(e) => Err(map_send_error(e)), - } + .await) } -fn map_send_error(err: SendWithRetryError) -> TraceExporterError { +pub(crate) fn map_send_error(err: SendWithRetryError) -> TraceExporterError { match err { SendWithRetryError::Http(response, _) => { let status = response.status(); diff --git a/libdd-data-pipeline/src/agentless/mod.rs b/libdd-data-pipeline/src/agentless/mod.rs index 50263ef59c..8cbd9f34a5 100644 --- a/libdd-data-pipeline/src/agentless/mod.rs +++ b/libdd-data-pipeline/src/agentless/mod.rs @@ -27,4 +27,3 @@ pub(crate) mod config; pub(crate) mod exporter; pub use config::AgentlessTraceConfig; -pub use exporter::send_agentless_traces_http; diff --git a/libdd-data-pipeline/src/otlp/exporter.rs b/libdd-data-pipeline/src/otlp/exporter.rs index 1297a7a25e..9ff78b0048 100644 --- a/libdd-data-pipeline/src/otlp/exporter.rs +++ b/libdd-data-pipeline/src/otlp/exporter.rs @@ -10,6 +10,7 @@ use libdd_capabilities::{HttpClientCapability, SleepCapability}; use libdd_common::Endpoint; use libdd_trace_utils::send_with_retry::{ send_with_retry, CompressionStrategy, RetryBackoffType, RetryStrategy, SendWithRetryError, + SendWithRetryResult, }; use std::time::Duration; @@ -36,6 +37,34 @@ pub(crate) async fn send_otlp_http( body: Vec, max_retries: u32, ) -> Result<(), TraceExporterError> { + match send_otlp_http_with_result( + capabilities, + endpoint_url, + config_headers, + timeout, + test_token, + content_type, + body, + max_retries, + ) + .await? + { + Ok(_) => Ok(()), + Err(e) => Err(map_send_error(e)), + } +} + +#[allow(clippy::too_many_arguments)] +async fn send_otlp_http_with_result( + capabilities: &C, + endpoint_url: &str, + config_headers: &HeaderMap, + timeout: Duration, + test_token: Option<&str>, + content_type: http::HeaderValue, + body: Vec, + max_retries: u32, +) -> Result { let url = libdd_common::parse_uri(endpoint_url).map_err(|e| { TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "Invalid OTLP endpoint URL: {}", @@ -67,7 +96,7 @@ pub(crate) async fn send_otlp_http( None, ); - match send_with_retry( + Ok(send_with_retry( capabilities, &target, body, @@ -75,25 +104,16 @@ pub(crate) async fn send_otlp_http( &retry_strategy, CompressionStrategy::None, ) - .await - { - Ok(_) => Ok(()), - Err(e) => Err(map_send_error(e).await), - } + .await) } -/// Send an OTLP trace payload to the configured endpoint with retries. The `Content-Type` is -/// derived from `config.protocol`, which also selected the body encoding (JSON or protobuf). -/// -/// `test_token` is forwarded as `X-Datadog-Test-Session-Token` when set, enabling snapshot tests -/// against the Datadog test agent's OTLP endpoint. -pub async fn send_otlp_traces_http( +pub(crate) async fn send_otlp_traces_http( capabilities: &C, config: &OtlpTraceConfig, test_token: Option<&str>, body: Vec, -) -> Result<(), TraceExporterError> { - send_otlp_http( +) -> Result { + send_otlp_http_with_result( capabilities, &config.endpoint_url, &config.headers, @@ -106,7 +126,7 @@ pub async fn send_otlp_traces_http( .await } -async fn map_send_error(err: SendWithRetryError) -> TraceExporterError { +pub(crate) fn map_send_error(err: SendWithRetryError) -> TraceExporterError { match err { SendWithRetryError::Http(response, _) => { let status = response.status(); diff --git a/libdd-data-pipeline/src/otlp/mod.rs b/libdd-data-pipeline/src/otlp/mod.rs index 0bda6b1b7e..c6edf66e28 100644 --- a/libdd-data-pipeline/src/otlp/mod.rs +++ b/libdd-data-pipeline/src/otlp/mod.rs @@ -32,6 +32,5 @@ pub mod exporter; pub mod metrics; pub use config::{OtlpMetricsConfig, OtlpProtocol, OtlpTraceConfig}; -pub use exporter::send_otlp_traces_http; pub use libdd_trace_utils::otlp_encoder::{map_traces_to_otlp, OtlpResourceInfo}; pub use metrics::OtlpStatsExporter; diff --git a/libdd-data-pipeline/src/telemetry/metrics.rs b/libdd-data-pipeline/src/telemetry/metrics.rs index c003dfbecd..2e11a185ec 100644 --- a/libdd-data-pipeline/src/telemetry/metrics.rs +++ b/libdd-data-pipeline/src/telemetry/metrics.rs @@ -34,6 +34,12 @@ pub enum MetricKind { ChunksDroppedSerializationError, /// trace_chunks_dropped metric (reason: send_failure) ChunksDroppedSendFailure, + /// spans_enqueued_for_serialization metric + SpansEnqueuedForSerialization, + /// spans_dropped metric (reason: serialization_error) + SpansDroppedSerializationError, + /// spans_dropped metric (reason: api_error) + SpansDroppedApiError, } /// Constants for metric names @@ -44,6 +50,8 @@ const API_BYTES_STR: &str = "trace_api.bytes"; const API_RESPONSES_STR: &str = "trace_api.responses"; const CHUNKS_SENT_STR: &str = "trace_chunks_sent"; const CHUNKS_DROPPED_STR: &str = "trace_chunks_dropped"; +const SPANS_ENQUEUED_FOR_SERIALIZATION_STR: &str = "spans_enqueued_for_serialization"; +const SPANS_DROPPED_STR: &str = "spans_dropped"; #[derive(Debug)] struct Metric { @@ -132,6 +140,24 @@ const METRICS: &[Metric] = &[ tag!["reason", "send_failure"], ], }, + Metric { + name: SPANS_ENQUEUED_FOR_SERIALIZATION_STR, + metric_type: MetricType::Count, + namespace: MetricNamespace::Tracers, + tags: &[], + }, + Metric { + name: SPANS_DROPPED_STR, + metric_type: MetricType::Count, + namespace: MetricNamespace::Tracers, + tags: &[tag!["reason", "serialization_error"]], + }, + Metric { + name: SPANS_DROPPED_STR, + metric_type: MetricType::Count, + namespace: MetricNamespace::Tracers, + tags: &[tag!["reason", "api_error"]], + }, ]; /// Structure to accumulate partial results coming from sending traces to the agent. diff --git a/libdd-data-pipeline/src/telemetry/mod.rs b/libdd-data-pipeline/src/telemetry/mod.rs index c371436996..0696397cf8 100644 --- a/libdd-data-pipeline/src/telemetry/mod.rs +++ b/libdd-data-pipeline/src/telemetry/mod.rs @@ -212,6 +212,9 @@ pub struct SendPayloadTelemetry { chunks_sent: u64, chunks_dropped_serialization_error: u64, chunks_dropped_send_failure: u64, + spans_enqueued_for_serialization: u64, + spans_dropped_serialization_error: u64, + spans_dropped_api_error: u64, responses_count_per_code: HashMap, } @@ -232,13 +235,35 @@ impl From<&SendDataResult> for SendPayloadTelemetry { } impl SendPayloadTelemetry { - /// Create a [`SendPayloadTelemetry`] from a [`SendWithRetryResult`]. + /// Create telemetry for spans accepted by the trace-export serialization pipeline. + pub fn spans_enqueued(span_count: u64) -> Self { + Self { + spans_enqueued_for_serialization: span_count, + ..Default::default() + } + } + + /// Create telemetry for spans dropped before an HTTP request could be built. + pub fn spans_dropped_serialization_error(span_count: u64) -> Self { + Self { + spans_dropped_serialization_error: span_count, + ..Default::default() + } + } + + /// Create send telemetry and attribute a definitive failure to every span in the payload. /// /// # Arguments /// * `value` - The result of sending traces with retry /// * `bytes_sent` - The number of bytes in the payload /// * `chunks` - The number of trace chunks in the payload - pub fn from_retry_result(value: &SendWithRetryResult, bytes_sent: u64, chunks: u64) -> Self { + /// * `spans` - The number of spans in the payload + pub fn from_retry_result_with_spans( + value: &SendWithRetryResult, + bytes_sent: u64, + chunks: u64, + spans: u64, + ) -> Self { let mut telemetry = Self::default(); match value { Ok((response, attempts)) => { @@ -252,6 +277,7 @@ impl SendPayloadTelemetry { Err(err) => match err { SendWithRetryError::Http(response, attempts) => { telemetry.chunks_dropped_send_failure = chunks; + telemetry.spans_dropped_api_error = spans; telemetry.errors_status_code = 1; telemetry .responses_count_per_code @@ -260,21 +286,25 @@ impl SendPayloadTelemetry { } SendWithRetryError::Timeout(attempts) => { telemetry.chunks_dropped_send_failure = chunks; + telemetry.spans_dropped_api_error = spans; telemetry.errors_timeout = 1; telemetry.requests_count = *attempts as u64; } SendWithRetryError::Network(_, attempts) => { telemetry.chunks_dropped_send_failure = chunks; + telemetry.spans_dropped_api_error = spans; telemetry.errors_network = 1; telemetry.requests_count = *attempts as u64; } SendWithRetryError::ResponseBody(attempts) => { telemetry.chunks_dropped_send_failure = chunks; + telemetry.spans_dropped_api_error = spans; telemetry.errors_network = 1; telemetry.requests_count = *attempts as u64; } SendWithRetryError::Build(attempts) => { telemetry.chunks_dropped_serialization_error = chunks; + telemetry.spans_dropped_serialization_error = spans; telemetry.requests_count = *attempts as u64; } }, @@ -333,6 +363,25 @@ impl Tel self.worker .add_point(data.chunks_dropped_send_failure as f64, key, vec![])?; } + if data.spans_enqueued_for_serialization > 0 { + let key = self + .metrics + .get(metrics::MetricKind::SpansEnqueuedForSerialization); + self.worker + .add_point(data.spans_enqueued_for_serialization as f64, key, vec![])?; + } + if data.spans_dropped_serialization_error > 0 { + let key = self + .metrics + .get(metrics::MetricKind::SpansDroppedSerializationError); + self.worker + .add_point(data.spans_dropped_serialization_error as f64, key, vec![])?; + } + if data.spans_dropped_api_error > 0 { + let key = self.metrics.get(metrics::MetricKind::SpansDroppedApiError); + self.worker + .add_point(data.spans_dropped_api_error as f64, key, vec![])?; + } if !data.responses_count_per_code.is_empty() { let key = self.metrics.get(metrics::MetricKind::ApiResponses); for (status_code, count) in &data.responses_count_per_code { @@ -699,6 +748,43 @@ mod tests { .expect("Failed to get runtime"); } + #[cfg_attr(miri, ignore)] + #[test] + fn span_writer_metrics_test() { + let spans_enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,2\.0\]\],"tags":\[\],"common":true,"type":"count"#).unwrap(); + let serialization_error = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,3\.0\]\],"tags":\["reason:serialization_error"\],"common":true,"type":"count"#).unwrap(); + let api_error = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,4\.0\]\],"tags":\["reason:api_error"\],"common":true,"type":"count"#).unwrap(); + let shared_runtime = ForkSafeRuntime::new().expect("Failed to create runtime"); + let server = MockServer::start(); + let mut telemetry_srv = server.mock(|when, then| { + when.method(POST) + .body_matches(spans_enqueued) + .body_matches(serialization_error) + .body_matches(api_error); + then.status(200).body(""); + }); + let data = SendPayloadTelemetry { + spans_enqueued_for_serialization: 2, + spans_dropped_serialization_error: 3, + spans_dropped_api_error: 4, + ..Default::default() + }; + let (client, handle) = get_test_client(&server.url("/"), &shared_runtime); + shared_runtime + .block_on(async { + let _ = client.start(); + let _ = client.send(&data); + sleep(Duration::from_millis(100)).await; + + handle.stop().await.expect("Failed to stop worker"); + assert!( + poll_for_mock_hits(&mut telemetry_srv, 1000, 10, 1).await, + "telemetry server did not receive calls within timeout" + ); + }) + .expect("Failed to get runtime"); + } + #[cfg_attr(miri, ignore)] #[test] fn send_client_side_stats_drops_test() { @@ -769,7 +855,7 @@ mod tests { .unwrap(), 3, )); - let telemetry = SendPayloadTelemetry::from_retry_result(&result, 4, 5); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 4, 5, 6); assert_eq!( telemetry, SendPayloadTelemetry { @@ -789,11 +875,12 @@ mod tests { .body(Bytes::new()) .unwrap(); let result = Err(SendWithRetryError::Http(error_response, 5)); - let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, + spans_dropped_api_error: 3, requests_count: 5, errors_status_code: 1, responses_count_per_code: HashMap::from([(400, 1)]), @@ -808,11 +895,12 @@ mod tests { HttpError::Network(anyhow::anyhow!("connection refused")), 5, )); - let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, + spans_dropped_api_error: 3, requests_count: 5, errors_network: 1, ..Default::default() @@ -823,11 +911,12 @@ mod tests { #[test] fn telemetry_from_timeout_error_test() { let result = Err(SendWithRetryError::Timeout(5)); - let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, + spans_dropped_api_error: 3, requests_count: 5, errors_timeout: 1, ..Default::default() @@ -839,11 +928,12 @@ mod tests { #[test] fn telemetry_from_build_error_test() { let result = Err(SendWithRetryError::Build(5)); - let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_serialization_error: 2, + spans_dropped_serialization_error: 3, requests_count: 5, ..Default::default() } diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index a29ffc6a55..7ee5fb790e 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -18,8 +18,12 @@ use self::metrics::MetricsEmitter; use self::stats::StatsComputationStatus; use self::trace_serializer::TraceSerializer; use crate::agent_info::ResponseObserver; -use crate::agentless::{send_agentless_traces_http, AgentlessTraceConfig}; -use crate::otlp::{map_traces_to_otlp, send_otlp_traces_http, OtlpResourceInfo, OtlpTraceConfig}; +use crate::agentless::exporter::{ + map_send_error as map_agentless_send_error, send_agentless_traces_http, +}; +use crate::agentless::AgentlessTraceConfig; +use crate::otlp::exporter::{map_send_error as map_otlp_send_error, send_otlp_traces_http}; +use crate::otlp::{map_traces_to_otlp, OtlpResourceInfo, OtlpTraceConfig}; #[cfg(feature = "telemetry")] use crate::telemetry::{SendPayloadTelemetry, TelemetryClient}; use crate::trace_exporter::agent_response::{ @@ -308,6 +312,15 @@ impl< .store(handle.map(|h| Arc::new(TelemetryClient::with_handle(h)))); } + #[cfg(feature = "telemetry")] + fn emit_trace_telemetry(&self, data: &SendPayloadTelemetry) { + if let Some(telemetry) = self.telemetry.load_full().as_deref() { + if let Err(e) = telemetry.send(data) { + error!(?e, "Error sending telemetry"); + } + } + } + /// Stop the background workers owned by this exporter. /// /// Sync facade over [`Self::shutdown_async`]; panics inside an existing tokio context. @@ -658,18 +671,36 @@ impl< config: &AgentlessTraceConfig, ) -> Result { let trace_count = traces.len(); + #[cfg(feature = "telemetry")] + let span_count = traces.iter().map(Vec::len).sum::(); let json_body = libdd_trace_utils::agentless_encoder::encode_payload( &traces, &self.metadata, ) .map_err(|e| { error!("Agentless JSON serialization error: {e}"); + #[cfg(feature = "telemetry")] + self.emit_trace_telemetry(&SendPayloadTelemetry::spans_dropped_serialization_error( + span_count as u64, + )); TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(e.to_string())) })?; let headers = build_agentless_headers(&self.metadata, &config.api_key, trace_count)?; + #[cfg(feature = "telemetry")] + let payload_len = json_body.len(); + let result = + send_agentless_traces_http(&self.capabilities, config, headers, json_body).await?; + + #[cfg(feature = "telemetry")] + self.emit_trace_telemetry(&SendPayloadTelemetry::from_retry_result_with_spans( + &result, + payload_len as u64, + trace_count as u64, + span_count as u64, + )); - send_agentless_traces_http(&self.capabilities, config, headers, json_body).await?; + result.map_err(map_agentless_send_error)?; Ok(AgentResponse::Unchanged) } @@ -679,6 +710,10 @@ impl< traces: Vec>>, config: &OtlpTraceConfig, ) -> Result { + #[cfg(feature = "telemetry")] + let trace_count = traces.len(); + #[cfg(feature = "telemetry")] + let span_count = traces.iter().map(Vec::len).sum::(); let resource_info = { let mut r = OtlpResourceInfo::default(); r.service = self.metadata.service.clone(); @@ -699,6 +734,10 @@ impl< map_traces_to_otlp(traces, &resource_info, config.otel_trace_semantics_enabled); let body = config.protocol.encode(&request).map_err(|e| { error!("OTLP serialization error: {e}"); + #[cfg(feature = "telemetry")] + self.emit_trace_telemetry(&SendPayloadTelemetry::spans_dropped_serialization_error( + span_count as u64, + )); TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "failed to encode OTLP request: {e}" ))) @@ -718,13 +757,25 @@ impl< } else { config }; - send_otlp_traces_http( + #[cfg(feature = "telemetry")] + let payload_len = body.len(); + let result = send_otlp_traces_http( &self.capabilities, config_to_use, self.endpoint.test_token.as_deref(), body, ) .await?; + + #[cfg(feature = "telemetry")] + self.emit_trace_telemetry(&SendPayloadTelemetry::from_retry_result_with_spans( + &result, + payload_len as u64, + trace_count as u64, + span_count as u64, + )); + + result.map_err(map_otlp_send_error)?; Ok(AgentResponse::Unchanged) } @@ -735,7 +786,10 @@ impl< mp_payload: Vec, headers: HeaderMap, chunks: usize, + spans: usize, ) -> Result { + #[cfg(not(feature = "telemetry"))] + let _ = spans; let strategy = RetryStrategy::default(); let payload_len = mp_payload.len(); @@ -751,15 +805,12 @@ impl< .await; #[cfg(feature = "telemetry")] - if let Some(telemetry) = self.telemetry.load_full().as_deref() { - if let Err(e) = telemetry.send(&SendPayloadTelemetry::from_retry_result( - &result, - payload_len as u64, - chunks as u64, - )) { - error!(?e, "Error sending telemetry"); - } - } + self.emit_trace_telemetry(&SendPayloadTelemetry::from_retry_result_with_spans( + &result, + payload_len as u64, + chunks as u64, + spans as u64, + )); self.handle_send_result(result, chunks, payload_len).await } @@ -799,6 +850,11 @@ impl< return self.send_trace_chunks_to_log(&traces, max_line_size); } + #[cfg(feature = "telemetry")] + self.emit_trace_telemetry(&SendPayloadTelemetry::spans_enqueued( + traces.iter().map(Vec::len).sum::() as u64, + )); + let mut header_tags: TracerHeaderTags = self.metadata.borrow().into(); if let Some(ref config) = self.agentless_config { @@ -845,6 +901,7 @@ impl< // Snapshot the effective format once so the serializer and the URL agree even if // `v1_active` flips mid-send (the background `/info` fetcher can race us otherwise). let effective_format = self.effective_output_format(); + let span_count = traces.iter().map(Vec::len).sum::(); let prepared = match self.serializer.prepare_traces_payload( traces, @@ -860,6 +917,10 @@ impl< HealthMetric::Count(health_metrics::SERIALIZE_TRACES_ERRORS, 1), None, ); + #[cfg(feature = "telemetry")] + self.emit_trace_telemetry( + &SendPayloadTelemetry::spans_dropped_serialization_error(span_count as u64), + ); return Err(e); } }; @@ -875,6 +936,7 @@ impl< prepared.data, prepared.headers, prepared.chunk_count, + span_count, ) .await; @@ -2354,13 +2416,15 @@ mod tests { #[cfg(feature = "telemetry")] mod telemetry_metrics_tests { use super::*; + use crate::telemetry::TelemetryClientBuilder; use crate::trace_exporter::tests::build_test_exporter; use httpmock::prelude::*; use httpmock::MockServer; use libdd_capabilities_impl::NativeCapabilities; use libdd_shared_runtime::ForkSafeRuntime; use libdd_tinybytes::BytesString; - use libdd_trace_utils::span::v05; + use libdd_trace_utils::msgpack_encoder; + use libdd_trace_utils::span::{v04::SpanBytes, v05}; // v05 messagepack empty payload -> [[""], []] const V5_EMPTY: [u8; 4] = [0x92, 0x91, 0xA0, 0x90]; @@ -2385,6 +2449,7 @@ mod telemetry_metrics_tests { let metrics_endpoint = server.mock(|when, then| { when.method(POST) .body_includes("\"metric\":\"trace_api.bytes\"") + .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") .path("/telemetry/proxy/api/v2/apmtelemetry"); then.status(200) .header("content-type", "application/json") @@ -2406,7 +2471,7 @@ mod telemetry_metrics_tests { }); let exporter = builder.build::().unwrap(); - let traces = vec![0x90]; + let traces = msgpack_encoder::v04::to_vec_from_v04(&[vec![SpanBytes::default()]]); let result = exporter.send(traces.as_ref()).unwrap(); let AgentResponse::Changed { body } = result else { panic!("Expected Changed response"); @@ -2440,6 +2505,7 @@ mod telemetry_metrics_tests { let metrics_endpoint = server.mock(|when, then| { when.method(POST) .body_includes("\"metric\":\"trace_api.bytes\"") + .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") .path("/telemetry/proxy/api/v2/apmtelemetry"); then.status(200) .header("content-type", "application/json") @@ -2455,7 +2521,10 @@ mod telemetry_metrics_tests { true, ); - let v5: (Vec, Vec>) = (vec![], vec![]); + let v5: (Vec, Vec>) = ( + vec![BytesString::from_static("")], + vec![vec![v05::Span::default()]], + ); let traces = rmp_serde::to_vec(&v5).unwrap(); let result = exporter.send(traces.as_ref()).unwrap(); let AgentResponse::Changed { body } = result else { @@ -2470,6 +2539,170 @@ mod telemetry_metrics_tests { metrics_endpoint.assert_calls(1); } + #[test] + #[cfg_attr(miri, ignore)] + fn test_exporter_metrics_v1() { + let server = MockServer::start(); + let traces_endpoint = server.mock(|when, then| { + when.method(POST).path(V1_TRACES_ENDPOINT); + then.status(200).body("{}"); + }); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .body_includes("\"metric\":\"trace_api.requests\"") + .body_includes("\"metric\":\"trace_api.responses\"") + .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") + .path("/telemetry/proxy/api/v2/apmtelemetry"); + then.status(200).body(""); + }); + + let exporter = build_test_exporter( + server.url("/"), + None, + TraceExporterInputFormat::V04, + TraceExporterOutputFormat::V1, + true, + true, + ); + exporter.v1_active.store(true, Ordering::Relaxed); + + let result = exporter + .shared_runtime + .block_on(exporter.send_trace_chunks_inner(vec![vec![SpanBytes::default()]])) + .unwrap() + .unwrap(); + assert!(matches!(result, AgentResponse::Changed { .. })); + + traces_endpoint.assert_calls(1); + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + } + + #[test] + #[cfg_attr(miri, ignore)] + fn test_exporter_metrics_otlp() { + let server = MockServer::start(); + let traces_endpoint = server.mock(|when, then| { + when.method(POST) + .path("/v1/traces") + .header("content-type", "application/json"); + then.status(200).body(""); + }); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .body_includes("\"metric\":\"trace_api.requests\"") + .body_includes("\"metric\":\"trace_api.responses\"") + .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") + .path("/telemetry/proxy/api/v2/apmtelemetry"); + then.status(200).body(""); + }); + + let otlp_endpoint = format!("{}/v1/traces", server.url("/").trim_end_matches('/')); + let mut builder = TraceExporter::::builder(); + builder + .set_url(&server.url("/")) + .set_service("foo") + .set_env("foo-env") + .set_tracer_version("v0.1") + .set_language("nodejs") + .set_language_version("1.0") + .set_language_interpreter("v8") + .set_otlp_endpoint(&otlp_endpoint) + .enable_telemetry(TelemetryConfig { + heartbeat: 100, + ..Default::default() + }); + let exporter = builder.build::().unwrap(); + + let result = exporter + .shared_runtime + .block_on(exporter.send_trace_chunks_inner(vec![vec![SpanBytes::default()]])) + .unwrap() + .unwrap(); + assert_eq!(result, AgentResponse::Unchanged); + + traces_endpoint.assert_calls(1); + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + } + + #[test] + #[cfg_attr(miri, ignore)] + fn test_exporter_metrics_agentless_with_shared_telemetry() { + let server = MockServer::start(); + let traces_endpoint = server.mock(|when, then| { + when.method(POST) + .path("/v1/input") + .header("dd-api-key", "test-api-key"); + then.status(200).body(""); + }); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .body_includes("\"metric\":\"trace_api.requests\"") + .body_includes("\"metric\":\"trace_api.responses\"") + .body_includes("\"metric\":\"spans_enqueued_for_serialization\""); + then.status(200).body(""); + }); + + let telemetry_runtime = ForkSafeRuntime::new().unwrap(); + let (telemetry_client, telemetry_worker) = TelemetryClientBuilder::default() + .set_service_name("foo") + .set_service_version("1.0") + .set_env("foo-env") + .set_language("nodejs") + .set_language_version("1.0") + .set_tracer_version("v0.1") + .set_url(&server.url("/")) + .set_heartbeat(100) + .build::() + .unwrap(); + let telemetry_worker = telemetry_runtime + .spawn_worker(telemetry_worker, true) + .unwrap(); + telemetry_client.start().unwrap(); + + let intake_endpoint = format!("{}/v1/input", server.url("/").trim_end_matches('/')); + let mut builder = TraceExporter::::builder(); + builder + .set_service("foo") + .set_env("foo-env") + .set_tracer_version("v0.1") + .set_language("nodejs") + .set_language_version("1.0") + .set_language_interpreter("v8") + .set_agentless_endpoint(&intake_endpoint, "test-api-key"); + let exporter = builder.build::().unwrap(); + exporter.set_telemetry_handle(Some(telemetry_client.clone_handle())); + + let result = exporter + .send_trace_chunks(vec![vec![SpanBytes::default()]], None) + .unwrap(); + assert_eq!(result, AgentResponse::Unchanged); + + traces_endpoint.assert_calls(1); + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + telemetry_runtime + .block_on(telemetry_worker.stop()) + .unwrap() + .unwrap(); + } + #[test] #[cfg_attr(miri, ignore)] fn test_exporter_metrics_v4_to_v5() { From a4070e10eed1939eb7398713be20084fc66602c9 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Mon, 10 Aug 2026 16:55:08 -0400 Subject: [PATCH 2/7] refactor(data-pipeline): preserve exporter visibility --- libdd-data-pipeline/src/agentless/exporter.rs | 2 +- libdd-data-pipeline/src/agentless/mod.rs | 1 + libdd-data-pipeline/src/otlp/exporter.rs | 7 ++++++- libdd-data-pipeline/src/otlp/mod.rs | 1 + libdd-data-pipeline/src/trace_exporter/mod.rs | 10 ++++------ 5 files changed, 13 insertions(+), 8 deletions(-) diff --git a/libdd-data-pipeline/src/agentless/exporter.rs b/libdd-data-pipeline/src/agentless/exporter.rs index a0ca3e44a3..227fa4dc0a 100644 --- a/libdd-data-pipeline/src/agentless/exporter.rs +++ b/libdd-data-pipeline/src/agentless/exporter.rs @@ -22,7 +22,7 @@ const AGENTLESS_RETRY_DELAY_MS: u64 = 1000; /// `headers` should already contain all required headers (api key, content-type, meta-*, /// entity, trace-count, etc.). `test_token` is forwarded as `X-Datadog-Test-Session-Token` /// when set, enabling snapshot tests against a local mock. -pub(crate) async fn send_agentless_traces_http( +pub async fn send_agentless_traces_http( capabilities: &C, config: &AgentlessTraceConfig, headers: HeaderMap, diff --git a/libdd-data-pipeline/src/agentless/mod.rs b/libdd-data-pipeline/src/agentless/mod.rs index 8cbd9f34a5..50263ef59c 100644 --- a/libdd-data-pipeline/src/agentless/mod.rs +++ b/libdd-data-pipeline/src/agentless/mod.rs @@ -27,3 +27,4 @@ pub(crate) mod config; pub(crate) mod exporter; pub use config::AgentlessTraceConfig; +pub use exporter::send_agentless_traces_http; diff --git a/libdd-data-pipeline/src/otlp/exporter.rs b/libdd-data-pipeline/src/otlp/exporter.rs index 9ff78b0048..505f0995af 100644 --- a/libdd-data-pipeline/src/otlp/exporter.rs +++ b/libdd-data-pipeline/src/otlp/exporter.rs @@ -107,7 +107,12 @@ async fn send_otlp_http_with_result( .await) } -pub(crate) async fn send_otlp_traces_http( +/// Send an OTLP trace payload to the configured endpoint with retries. The `Content-Type` is +/// derived from `config.protocol`, which also selected the body encoding (JSON or protobuf). +/// +/// `test_token` is forwarded as `X-Datadog-Test-Session-Token` when set, enabling snapshot tests +/// against the Datadog test agent's OTLP endpoint. +pub async fn send_otlp_traces_http( capabilities: &C, config: &OtlpTraceConfig, test_token: Option<&str>, diff --git a/libdd-data-pipeline/src/otlp/mod.rs b/libdd-data-pipeline/src/otlp/mod.rs index c6edf66e28..0bda6b1b7e 100644 --- a/libdd-data-pipeline/src/otlp/mod.rs +++ b/libdd-data-pipeline/src/otlp/mod.rs @@ -32,5 +32,6 @@ pub mod exporter; pub mod metrics; pub use config::{OtlpMetricsConfig, OtlpProtocol, OtlpTraceConfig}; +pub use exporter::send_otlp_traces_http; pub use libdd_trace_utils::otlp_encoder::{map_traces_to_otlp, OtlpResourceInfo}; pub use metrics::OtlpStatsExporter; diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index 7ee5fb790e..7db77c1254 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -18,12 +18,10 @@ use self::metrics::MetricsEmitter; use self::stats::StatsComputationStatus; use self::trace_serializer::TraceSerializer; use crate::agent_info::ResponseObserver; -use crate::agentless::exporter::{ - map_send_error as map_agentless_send_error, send_agentless_traces_http, -}; -use crate::agentless::AgentlessTraceConfig; -use crate::otlp::exporter::{map_send_error as map_otlp_send_error, send_otlp_traces_http}; -use crate::otlp::{map_traces_to_otlp, OtlpResourceInfo, OtlpTraceConfig}; +use crate::agentless::exporter::map_send_error as map_agentless_send_error; +use crate::agentless::{send_agentless_traces_http, AgentlessTraceConfig}; +use crate::otlp::exporter::map_send_error as map_otlp_send_error; +use crate::otlp::{map_traces_to_otlp, send_otlp_traces_http, OtlpResourceInfo, OtlpTraceConfig}; #[cfg(feature = "telemetry")] use crate::telemetry::{SendPayloadTelemetry, TelemetryClient}; use crate::trace_exporter::agent_response::{ From a07af01a22f1eeac8c202954a216af214c59a975 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Mon, 10 Aug 2026 16:57:05 -0400 Subject: [PATCH 3/7] refactor(data-pipeline): preserve mapper signature --- libdd-data-pipeline/src/otlp/exporter.rs | 4 ++-- libdd-data-pipeline/src/trace_exporter/mod.rs | 4 +++- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/libdd-data-pipeline/src/otlp/exporter.rs b/libdd-data-pipeline/src/otlp/exporter.rs index 505f0995af..b556e4a056 100644 --- a/libdd-data-pipeline/src/otlp/exporter.rs +++ b/libdd-data-pipeline/src/otlp/exporter.rs @@ -50,7 +50,7 @@ pub(crate) async fn send_otlp_http( .await? { Ok(_) => Ok(()), - Err(e) => Err(map_send_error(e)), + Err(e) => Err(map_send_error(e).await), } } @@ -131,7 +131,7 @@ pub async fn send_otlp_traces_http( .await } -pub(crate) fn map_send_error(err: SendWithRetryError) -> TraceExporterError { +pub(crate) async fn map_send_error(err: SendWithRetryError) -> TraceExporterError { match err { SendWithRetryError::Http(response, _) => { let status = response.status(); diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index 7db77c1254..46ba5c9154 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -773,7 +773,9 @@ impl< span_count as u64, )); - result.map_err(map_otlp_send_error)?; + if let Err(err) = result { + return Err(map_otlp_send_error(err).await); + } Ok(AgentResponse::Unchanged) } From 3b00525bd8a6439279f216ad0718b465678e94f6 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Thu, 20 Aug 2026 09:41:49 -0400 Subject: [PATCH 4/7] fix(data-pipeline): preserve trace transport behavior --- libdd-data-pipeline/src/agentless/exporter.rs | 29 +- libdd-data-pipeline/src/agentless/mod.rs | 1 + libdd-data-pipeline/src/otlp/exporter.rs | 33 +- libdd-data-pipeline/src/otlp/mod.rs | 1 + libdd-data-pipeline/src/telemetry/mod.rs | 116 ++----- libdd-data-pipeline/src/trace_exporter/mod.rs | 297 +++++------------- libdd-trace-utils/src/send_with_retry/mod.rs | 41 ++- 7 files changed, 173 insertions(+), 345 deletions(-) diff --git a/libdd-data-pipeline/src/agentless/exporter.rs b/libdd-data-pipeline/src/agentless/exporter.rs index 227fa4dc0a..6b4b2fd615 100644 --- a/libdd-data-pipeline/src/agentless/exporter.rs +++ b/libdd-data-pipeline/src/agentless/exporter.rs @@ -9,8 +9,8 @@ use http::HeaderMap; use libdd_capabilities::{HttpClientCapability, SleepCapability}; use libdd_common::Endpoint; use libdd_trace_utils::send_with_retry::{ - send_with_retry, CompressionStrategy, RetryBackoffType, RetryStrategy, SendWithRetryError, - SendWithRetryResult, + send_with_retry_and_size, CompressionStrategy, RetryBackoffType, RetryStrategy, + SendWithRetryError, SendWithRetryResult, }; use tracing::error; @@ -22,12 +22,27 @@ const AGENTLESS_RETRY_DELAY_MS: u64 = 1000; /// `headers` should already contain all required headers (api key, content-type, meta-*, /// entity, trace-count, etc.). `test_token` is forwarded as `X-Datadog-Test-Session-Token` /// when set, enabling snapshot tests against a local mock. +#[allow(dead_code)] // Retains the existing contract while telemetry uses the observer variant. pub async fn send_agentless_traces_http( capabilities: &C, config: &AgentlessTraceConfig, headers: HeaderMap, json_body: Vec, -) -> Result { +) -> Result<(), TraceExporterError> { + send_agentless_traces_http_with_observer(capabilities, config, headers, json_body, |_, _| {}) + .await +} + +pub(crate) async fn send_agentless_traces_http_with_observer< + C: HttpClientCapability + SleepCapability, + F: FnOnce(&SendWithRetryResult, usize), +>( + capabilities: &C, + config: &AgentlessTraceConfig, + headers: HeaderMap, + json_body: Vec, + observer: F, +) -> Result<(), TraceExporterError> { let url = libdd_common::parse_uri(&config.endpoint_url).map_err(|e| { TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "Invalid agentless endpoint URL: {e}" @@ -52,7 +67,7 @@ pub async fn send_agentless_traces_http TraceExporterError { +fn map_send_error(err: SendWithRetryError) -> TraceExporterError { match err { SendWithRetryError::Http(response, _) => { let status = response.status(); diff --git a/libdd-data-pipeline/src/agentless/mod.rs b/libdd-data-pipeline/src/agentless/mod.rs index 50263ef59c..a4a869c146 100644 --- a/libdd-data-pipeline/src/agentless/mod.rs +++ b/libdd-data-pipeline/src/agentless/mod.rs @@ -27,4 +27,5 @@ pub(crate) mod config; pub(crate) mod exporter; pub use config::AgentlessTraceConfig; +#[allow(unused_imports)] pub use exporter::send_agentless_traces_http; diff --git a/libdd-data-pipeline/src/otlp/exporter.rs b/libdd-data-pipeline/src/otlp/exporter.rs index b556e4a056..bf4aa8c9a8 100644 --- a/libdd-data-pipeline/src/otlp/exporter.rs +++ b/libdd-data-pipeline/src/otlp/exporter.rs @@ -37,7 +37,7 @@ pub(crate) async fn send_otlp_http( body: Vec, max_retries: u32, ) -> Result<(), TraceExporterError> { - match send_otlp_http_with_result( + send_otlp_http_with_observer( capabilities, endpoint_url, config_headers, @@ -46,16 +46,16 @@ pub(crate) async fn send_otlp_http( content_type, body, max_retries, + |_| {}, ) - .await? - { - Ok(_) => Ok(()), - Err(e) => Err(map_send_error(e).await), - } + .await } #[allow(clippy::too_many_arguments)] -async fn send_otlp_http_with_result( +pub(crate) async fn send_otlp_http_with_observer< + C: HttpClientCapability + SleepCapability, + F: FnOnce(&SendWithRetryResult), +>( capabilities: &C, endpoint_url: &str, config_headers: &HeaderMap, @@ -64,7 +64,8 @@ async fn send_otlp_http_with_result( content_type: http::HeaderValue, body: Vec, max_retries: u32, -) -> Result { + observer: F, +) -> Result<(), TraceExporterError> { let url = libdd_common::parse_uri(endpoint_url).map_err(|e| { TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "Invalid OTLP endpoint URL: {}", @@ -96,7 +97,7 @@ async fn send_otlp_http_with_result( None, ); - Ok(send_with_retry( + let result = send_with_retry( capabilities, &target, body, @@ -104,7 +105,12 @@ async fn send_otlp_http_with_result( &retry_strategy, CompressionStrategy::None, ) - .await) + .await; + observer(&result); + match result { + Ok(_) => Ok(()), + Err(e) => Err(map_send_error(e).await), + } } /// Send an OTLP trace payload to the configured endpoint with retries. The `Content-Type` is @@ -112,13 +118,14 @@ async fn send_otlp_http_with_result( /// /// `test_token` is forwarded as `X-Datadog-Test-Session-Token` when set, enabling snapshot tests /// against the Datadog test agent's OTLP endpoint. +#[allow(dead_code)] // Retains the existing contract while telemetry uses the observer variant. pub async fn send_otlp_traces_http( capabilities: &C, config: &OtlpTraceConfig, test_token: Option<&str>, body: Vec, -) -> Result { - send_otlp_http_with_result( +) -> Result<(), TraceExporterError> { + send_otlp_http( capabilities, &config.endpoint_url, &config.headers, @@ -131,7 +138,7 @@ pub async fn send_otlp_traces_http( .await } -pub(crate) async fn map_send_error(err: SendWithRetryError) -> TraceExporterError { +async fn map_send_error(err: SendWithRetryError) -> TraceExporterError { match err { SendWithRetryError::Http(response, _) => { let status = response.status(); diff --git a/libdd-data-pipeline/src/otlp/mod.rs b/libdd-data-pipeline/src/otlp/mod.rs index 0bda6b1b7e..073372ba02 100644 --- a/libdd-data-pipeline/src/otlp/mod.rs +++ b/libdd-data-pipeline/src/otlp/mod.rs @@ -32,6 +32,7 @@ pub mod exporter; pub mod metrics; pub use config::{OtlpMetricsConfig, OtlpProtocol, OtlpTraceConfig}; +#[allow(unused_imports)] pub use exporter::send_otlp_traces_http; pub use libdd_trace_utils::otlp_encoder::{map_traces_to_otlp, OtlpResourceInfo}; pub use metrics::OtlpStatsExporter; diff --git a/libdd-data-pipeline/src/telemetry/mod.rs b/libdd-data-pipeline/src/telemetry/mod.rs index 0696397cf8..a69f8c752a 100644 --- a/libdd-data-pipeline/src/telemetry/mod.rs +++ b/libdd-data-pipeline/src/telemetry/mod.rs @@ -198,6 +198,18 @@ impl Tel worker: handle, } } + + pub(crate) fn send_count_metric( + &self, + metric: metrics::MetricKind, + count: usize, + ) -> Result<(), TelemetryError> { + if count > 0 { + self.worker + .add_point(count as f64, self.metrics.get(metric), vec![])?; + } + Ok(()) + } } /// Telemetry describing the sending of a trace payload @@ -212,9 +224,6 @@ pub struct SendPayloadTelemetry { chunks_sent: u64, chunks_dropped_serialization_error: u64, chunks_dropped_send_failure: u64, - spans_enqueued_for_serialization: u64, - spans_dropped_serialization_error: u64, - spans_dropped_api_error: u64, responses_count_per_code: HashMap, } @@ -235,35 +244,13 @@ impl From<&SendDataResult> for SendPayloadTelemetry { } impl SendPayloadTelemetry { - /// Create telemetry for spans accepted by the trace-export serialization pipeline. - pub fn spans_enqueued(span_count: u64) -> Self { - Self { - spans_enqueued_for_serialization: span_count, - ..Default::default() - } - } - - /// Create telemetry for spans dropped before an HTTP request could be built. - pub fn spans_dropped_serialization_error(span_count: u64) -> Self { - Self { - spans_dropped_serialization_error: span_count, - ..Default::default() - } - } - - /// Create send telemetry and attribute a definitive failure to every span in the payload. + /// Create a [`SendPayloadTelemetry`] from a [`SendWithRetryResult`]. /// /// # Arguments /// * `value` - The result of sending traces with retry /// * `bytes_sent` - The number of bytes in the payload /// * `chunks` - The number of trace chunks in the payload - /// * `spans` - The number of spans in the payload - pub fn from_retry_result_with_spans( - value: &SendWithRetryResult, - bytes_sent: u64, - chunks: u64, - spans: u64, - ) -> Self { + pub fn from_retry_result(value: &SendWithRetryResult, bytes_sent: u64, chunks: u64) -> Self { let mut telemetry = Self::default(); match value { Ok((response, attempts)) => { @@ -277,7 +264,6 @@ impl SendPayloadTelemetry { Err(err) => match err { SendWithRetryError::Http(response, attempts) => { telemetry.chunks_dropped_send_failure = chunks; - telemetry.spans_dropped_api_error = spans; telemetry.errors_status_code = 1; telemetry .responses_count_per_code @@ -286,25 +272,21 @@ impl SendPayloadTelemetry { } SendWithRetryError::Timeout(attempts) => { telemetry.chunks_dropped_send_failure = chunks; - telemetry.spans_dropped_api_error = spans; telemetry.errors_timeout = 1; telemetry.requests_count = *attempts as u64; } SendWithRetryError::Network(_, attempts) => { telemetry.chunks_dropped_send_failure = chunks; - telemetry.spans_dropped_api_error = spans; telemetry.errors_network = 1; telemetry.requests_count = *attempts as u64; } SendWithRetryError::ResponseBody(attempts) => { telemetry.chunks_dropped_send_failure = chunks; - telemetry.spans_dropped_api_error = spans; telemetry.errors_network = 1; telemetry.requests_count = *attempts as u64; } SendWithRetryError::Build(attempts) => { telemetry.chunks_dropped_serialization_error = chunks; - telemetry.spans_dropped_serialization_error = spans; telemetry.requests_count = *attempts as u64; } }, @@ -363,25 +345,6 @@ impl Tel self.worker .add_point(data.chunks_dropped_send_failure as f64, key, vec![])?; } - if data.spans_enqueued_for_serialization > 0 { - let key = self - .metrics - .get(metrics::MetricKind::SpansEnqueuedForSerialization); - self.worker - .add_point(data.spans_enqueued_for_serialization as f64, key, vec![])?; - } - if data.spans_dropped_serialization_error > 0 { - let key = self - .metrics - .get(metrics::MetricKind::SpansDroppedSerializationError); - self.worker - .add_point(data.spans_dropped_serialization_error as f64, key, vec![])?; - } - if data.spans_dropped_api_error > 0 { - let key = self.metrics.get(metrics::MetricKind::SpansDroppedApiError); - self.worker - .add_point(data.spans_dropped_api_error as f64, key, vec![])?; - } if !data.responses_count_per_code.is_empty() { let key = self.metrics.get(metrics::MetricKind::ApiResponses); for (status_code, count) in &data.responses_count_per_code { @@ -748,43 +711,6 @@ mod tests { .expect("Failed to get runtime"); } - #[cfg_attr(miri, ignore)] - #[test] - fn span_writer_metrics_test() { - let spans_enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,2\.0\]\],"tags":\[\],"common":true,"type":"count"#).unwrap(); - let serialization_error = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,3\.0\]\],"tags":\["reason:serialization_error"\],"common":true,"type":"count"#).unwrap(); - let api_error = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,4\.0\]\],"tags":\["reason:api_error"\],"common":true,"type":"count"#).unwrap(); - let shared_runtime = ForkSafeRuntime::new().expect("Failed to create runtime"); - let server = MockServer::start(); - let mut telemetry_srv = server.mock(|when, then| { - when.method(POST) - .body_matches(spans_enqueued) - .body_matches(serialization_error) - .body_matches(api_error); - then.status(200).body(""); - }); - let data = SendPayloadTelemetry { - spans_enqueued_for_serialization: 2, - spans_dropped_serialization_error: 3, - spans_dropped_api_error: 4, - ..Default::default() - }; - let (client, handle) = get_test_client(&server.url("/"), &shared_runtime); - shared_runtime - .block_on(async { - let _ = client.start(); - let _ = client.send(&data); - sleep(Duration::from_millis(100)).await; - - handle.stop().await.expect("Failed to stop worker"); - assert!( - poll_for_mock_hits(&mut telemetry_srv, 1000, 10, 1).await, - "telemetry server did not receive calls within timeout" - ); - }) - .expect("Failed to get runtime"); - } - #[cfg_attr(miri, ignore)] #[test] fn send_client_side_stats_drops_test() { @@ -855,7 +781,7 @@ mod tests { .unwrap(), 3, )); - let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 4, 5, 6); + let telemetry = SendPayloadTelemetry::from_retry_result(&result, 4, 5); assert_eq!( telemetry, SendPayloadTelemetry { @@ -875,12 +801,11 @@ mod tests { .body(Bytes::new()) .unwrap(); let result = Err(SendWithRetryError::Http(error_response, 5)); - let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); + let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, - spans_dropped_api_error: 3, requests_count: 5, errors_status_code: 1, responses_count_per_code: HashMap::from([(400, 1)]), @@ -895,12 +820,11 @@ mod tests { HttpError::Network(anyhow::anyhow!("connection refused")), 5, )); - let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); + let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, - spans_dropped_api_error: 3, requests_count: 5, errors_network: 1, ..Default::default() @@ -911,12 +835,11 @@ mod tests { #[test] fn telemetry_from_timeout_error_test() { let result = Err(SendWithRetryError::Timeout(5)); - let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); + let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, - spans_dropped_api_error: 3, requests_count: 5, errors_timeout: 1, ..Default::default() @@ -928,12 +851,11 @@ mod tests { #[test] fn telemetry_from_build_error_test() { let result = Err(SendWithRetryError::Build(5)); - let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 3); + let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_serialization_error: 2, - spans_dropped_serialization_error: 3, requests_count: 5, ..Default::default() } diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index 46ba5c9154..01b2486596 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -18,10 +18,12 @@ use self::metrics::MetricsEmitter; use self::stats::StatsComputationStatus; use self::trace_serializer::TraceSerializer; use crate::agent_info::ResponseObserver; -use crate::agentless::exporter::map_send_error as map_agentless_send_error; -use crate::agentless::{send_agentless_traces_http, AgentlessTraceConfig}; -use crate::otlp::exporter::map_send_error as map_otlp_send_error; -use crate::otlp::{map_traces_to_otlp, send_otlp_traces_http, OtlpResourceInfo, OtlpTraceConfig}; +use crate::agentless::exporter::send_agentless_traces_http_with_observer; +use crate::agentless::AgentlessTraceConfig; +use crate::otlp::exporter::{send_otlp_http_with_observer, OTLP_MAX_RETRIES}; +use crate::otlp::{map_traces_to_otlp, OtlpResourceInfo, OtlpTraceConfig}; +#[cfg(feature = "telemetry")] +use crate::telemetry::metrics::MetricKind; #[cfg(feature = "telemetry")] use crate::telemetry::{SendPayloadTelemetry, TelemetryClient}; use crate::trace_exporter::agent_response::{ @@ -311,9 +313,44 @@ impl< } #[cfg(feature = "telemetry")] - fn emit_trace_telemetry(&self, data: &SendPayloadTelemetry) { + fn emit_trace_counts(&self, counts: &[(MetricKind, usize)]) { if let Some(telemetry) = self.telemetry.load_full().as_deref() { - if let Err(e) = telemetry.send(data) { + for &(metric, count) in counts { + if let Err(e) = telemetry.send_count_metric(metric, count) { + error!(?e, "Error sending telemetry"); + } + } + } + } + + #[cfg(feature = "telemetry")] + fn emit_serialization_drop(&self, chunks: usize, spans: usize) { + self.emit_trace_counts(&[ + (MetricKind::ChunksDroppedSerializationError, chunks), + (MetricKind::SpansDroppedSerializationError, spans), + ]); + } + + #[cfg(feature = "telemetry")] + fn emit_send_telemetry( + &self, + result: &SendWithRetryResult, + bytes: usize, + chunks: usize, + spans: usize, + ) { + if let Some(telemetry) = self.telemetry.load_full().as_deref() { + let payload = + SendPayloadTelemetry::from_retry_result(result, bytes as u64, chunks as u64); + if let Err(e) = telemetry.send(&payload) { + error!(?e, "Error sending telemetry"); + } + let metric = match result { + Err(SendWithRetryError::Build(_)) => MetricKind::SpansDroppedSerializationError, + Err(_) => MetricKind::SpansDroppedApiError, + Ok(_) => return, + }; + if let Err(e) = telemetry.send_count_metric(metric, spans) { error!(?e, "Error sending telemetry"); } } @@ -678,27 +715,22 @@ impl< .map_err(|e| { error!("Agentless JSON serialization error: {e}"); #[cfg(feature = "telemetry")] - self.emit_trace_telemetry(&SendPayloadTelemetry::spans_dropped_serialization_error( - span_count as u64, - )); + self.emit_serialization_drop(trace_count, span_count); TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(e.to_string())) })?; let headers = build_agentless_headers(&self.metadata, &config.api_key, trace_count)?; - #[cfg(feature = "telemetry")] - let payload_len = json_body.len(); - let result = - send_agentless_traces_http(&self.capabilities, config, headers, json_body).await?; - - #[cfg(feature = "telemetry")] - self.emit_trace_telemetry(&SendPayloadTelemetry::from_retry_result_with_spans( - &result, - payload_len as u64, - trace_count as u64, - span_count as u64, - )); - - result.map_err(map_agentless_send_error)?; + send_agentless_traces_http_with_observer( + &self.capabilities, + config, + headers, + json_body, + |_result, _payload_len| { + #[cfg(feature = "telemetry")] + self.emit_send_telemetry(_result, _payload_len, trace_count, span_count); + }, + ) + .await?; Ok(AgentResponse::Unchanged) } @@ -733,9 +765,7 @@ impl< let body = config.protocol.encode(&request).map_err(|e| { error!("OTLP serialization error: {e}"); #[cfg(feature = "telemetry")] - self.emit_trace_telemetry(&SendPayloadTelemetry::spans_dropped_serialization_error( - span_count as u64, - )); + self.emit_serialization_drop(trace_count, span_count); TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "failed to encode OTLP request: {e}" ))) @@ -757,25 +787,21 @@ impl< }; #[cfg(feature = "telemetry")] let payload_len = body.len(); - let result = send_otlp_traces_http( + send_otlp_http_with_observer( &self.capabilities, - config_to_use, + &config_to_use.endpoint_url, + &config_to_use.headers, + config_to_use.timeout, self.endpoint.test_token.as_deref(), + config_to_use.protocol.content_type(), body, + OTLP_MAX_RETRIES, + |_result| { + #[cfg(feature = "telemetry")] + self.emit_send_telemetry(_result, payload_len, trace_count, span_count); + }, ) .await?; - - #[cfg(feature = "telemetry")] - self.emit_trace_telemetry(&SendPayloadTelemetry::from_retry_result_with_spans( - &result, - payload_len as u64, - trace_count as u64, - span_count as u64, - )); - - if let Err(err) = result { - return Err(map_otlp_send_error(err).await); - } Ok(AgentResponse::Unchanged) } @@ -805,12 +831,7 @@ impl< .await; #[cfg(feature = "telemetry")] - self.emit_trace_telemetry(&SendPayloadTelemetry::from_retry_result_with_spans( - &result, - payload_len as u64, - chunks as u64, - spans as u64, - )); + self.emit_send_telemetry(&result, payload_len, chunks, spans); self.handle_send_result(result, chunks, payload_len).await } @@ -851,9 +872,10 @@ impl< } #[cfg(feature = "telemetry")] - self.emit_trace_telemetry(&SendPayloadTelemetry::spans_enqueued( - traces.iter().map(Vec::len).sum::() as u64, - )); + self.emit_trace_counts(&[( + MetricKind::SpansEnqueuedForSerialization, + traces.iter().map(Vec::len).sum(), + )]); let mut header_tags: TracerHeaderTags = self.metadata.borrow().into(); @@ -901,6 +923,8 @@ impl< // Snapshot the effective format once so the serializer and the URL agree even if // `v1_active` flips mid-send (the background `/info` fetcher can race us otherwise). let effective_format = self.effective_output_format(); + #[cfg(feature = "telemetry")] + let trace_count = traces.len(); let span_count = traces.iter().map(Vec::len).sum::(); let prepared = match self.serializer.prepare_traces_payload( @@ -918,9 +942,7 @@ impl< None, ); #[cfg(feature = "telemetry")] - self.emit_trace_telemetry( - &SendPayloadTelemetry::spans_dropped_serialization_error(span_count as u64), - ); + self.emit_serialization_drop(trace_count, span_count); return Err(e); } }; @@ -2416,7 +2438,6 @@ mod tests { #[cfg(feature = "telemetry")] mod telemetry_metrics_tests { use super::*; - use crate::telemetry::TelemetryClientBuilder; use crate::trace_exporter::tests::build_test_exporter; use httpmock::prelude::*; use httpmock::MockServer; @@ -2505,7 +2526,6 @@ mod telemetry_metrics_tests { let metrics_endpoint = server.mock(|when, then| { when.method(POST) .body_includes("\"metric\":\"trace_api.bytes\"") - .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") .path("/telemetry/proxy/api/v2/apmtelemetry"); then.status(200) .header("content-type", "application/json") @@ -2521,10 +2541,7 @@ mod telemetry_metrics_tests { true, ); - let v5: (Vec, Vec>) = ( - vec![BytesString::from_static("")], - vec![vec![v05::Span::default()]], - ); + let v5: (Vec, Vec>) = (vec![], vec![]); let traces = rmp_serde::to_vec(&v5).unwrap(); let result = exporter.send(traces.as_ref()).unwrap(); let AgentResponse::Changed { body } = result else { @@ -2539,170 +2556,6 @@ mod telemetry_metrics_tests { metrics_endpoint.assert_calls(1); } - #[test] - #[cfg_attr(miri, ignore)] - fn test_exporter_metrics_v1() { - let server = MockServer::start(); - let traces_endpoint = server.mock(|when, then| { - when.method(POST).path(V1_TRACES_ENDPOINT); - then.status(200).body("{}"); - }); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .body_includes("\"metric\":\"trace_api.requests\"") - .body_includes("\"metric\":\"trace_api.responses\"") - .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") - .path("/telemetry/proxy/api/v2/apmtelemetry"); - then.status(200).body(""); - }); - - let exporter = build_test_exporter( - server.url("/"), - None, - TraceExporterInputFormat::V04, - TraceExporterOutputFormat::V1, - true, - true, - ); - exporter.v1_active.store(true, Ordering::Relaxed); - - let result = exporter - .shared_runtime - .block_on(exporter.send_trace_chunks_inner(vec![vec![SpanBytes::default()]])) - .unwrap() - .unwrap(); - assert!(matches!(result, AgentResponse::Changed { .. })); - - traces_endpoint.assert_calls(1); - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - } - - #[test] - #[cfg_attr(miri, ignore)] - fn test_exporter_metrics_otlp() { - let server = MockServer::start(); - let traces_endpoint = server.mock(|when, then| { - when.method(POST) - .path("/v1/traces") - .header("content-type", "application/json"); - then.status(200).body(""); - }); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .body_includes("\"metric\":\"trace_api.requests\"") - .body_includes("\"metric\":\"trace_api.responses\"") - .body_includes("\"metric\":\"spans_enqueued_for_serialization\"") - .path("/telemetry/proxy/api/v2/apmtelemetry"); - then.status(200).body(""); - }); - - let otlp_endpoint = format!("{}/v1/traces", server.url("/").trim_end_matches('/')); - let mut builder = TraceExporter::::builder(); - builder - .set_url(&server.url("/")) - .set_service("foo") - .set_env("foo-env") - .set_tracer_version("v0.1") - .set_language("nodejs") - .set_language_version("1.0") - .set_language_interpreter("v8") - .set_otlp_endpoint(&otlp_endpoint) - .enable_telemetry(TelemetryConfig { - heartbeat: 100, - ..Default::default() - }); - let exporter = builder.build::().unwrap(); - - let result = exporter - .shared_runtime - .block_on(exporter.send_trace_chunks_inner(vec![vec![SpanBytes::default()]])) - .unwrap() - .unwrap(); - assert_eq!(result, AgentResponse::Unchanged); - - traces_endpoint.assert_calls(1); - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - } - - #[test] - #[cfg_attr(miri, ignore)] - fn test_exporter_metrics_agentless_with_shared_telemetry() { - let server = MockServer::start(); - let traces_endpoint = server.mock(|when, then| { - when.method(POST) - .path("/v1/input") - .header("dd-api-key", "test-api-key"); - then.status(200).body(""); - }); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .body_includes("\"metric\":\"trace_api.requests\"") - .body_includes("\"metric\":\"trace_api.responses\"") - .body_includes("\"metric\":\"spans_enqueued_for_serialization\""); - then.status(200).body(""); - }); - - let telemetry_runtime = ForkSafeRuntime::new().unwrap(); - let (telemetry_client, telemetry_worker) = TelemetryClientBuilder::default() - .set_service_name("foo") - .set_service_version("1.0") - .set_env("foo-env") - .set_language("nodejs") - .set_language_version("1.0") - .set_tracer_version("v0.1") - .set_url(&server.url("/")) - .set_heartbeat(100) - .build::() - .unwrap(); - let telemetry_worker = telemetry_runtime - .spawn_worker(telemetry_worker, true) - .unwrap(); - telemetry_client.start().unwrap(); - - let intake_endpoint = format!("{}/v1/input", server.url("/").trim_end_matches('/')); - let mut builder = TraceExporter::::builder(); - builder - .set_service("foo") - .set_env("foo-env") - .set_tracer_version("v0.1") - .set_language("nodejs") - .set_language_version("1.0") - .set_language_interpreter("v8") - .set_agentless_endpoint(&intake_endpoint, "test-api-key"); - let exporter = builder.build::().unwrap(); - exporter.set_telemetry_handle(Some(telemetry_client.clone_handle())); - - let result = exporter - .send_trace_chunks(vec![vec![SpanBytes::default()]], None) - .unwrap(); - assert_eq!(result, AgentResponse::Unchanged); - - traces_endpoint.assert_calls(1); - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - telemetry_runtime - .block_on(telemetry_worker.stop()) - .unwrap() - .unwrap(); - } - #[test] #[cfg_attr(miri, ignore)] fn test_exporter_metrics_v4_to_v5() { diff --git a/libdd-trace-utils/src/send_with_retry/mod.rs b/libdd-trace-utils/src/send_with_retry/mod.rs index 41129f4d9d..50f4382aca 100644 --- a/libdd-trace-utils/src/send_with_retry/mod.rs +++ b/libdd-trace-utils/src/send_with_retry/mod.rs @@ -106,6 +106,31 @@ pub async fn send_with_retry( retry_strategy: &RetryStrategy, compression_strategy: CompressionStrategy, ) -> SendWithRetryResult { + send_with_retry_and_size( + capabilities, + target, + payload, + headers, + retry_strategy, + compression_strategy, + ) + .await + .0 +} + +/// Send a payload with retries and return its post-compression size. +/// +/// This is equivalent to [`send_with_retry`], with the payload size exposed for transport +/// telemetry. The result contains the final outcome and the total number of request attempts. +#[allow(clippy::result_large_err)] +pub async fn send_with_retry_and_size( + capabilities: &C, + target: &Endpoint, + payload: Vec, + headers: &HeaderMap, + retry_strategy: &RetryStrategy, + compression_strategy: CompressionStrategy, +) -> (SendWithRetryResult, usize) { let mut request_attempt = 0; let timeout = Duration::from_millis(target.timeout_ms); @@ -118,8 +143,9 @@ pub async fn send_with_retry( let (compressed, compression_strategy) = compression::compress(payload, compression_strategy); let payload = Bytes::from(compressed); + let payload_size = payload.len(); - loop { + let result = loop { request_attempt += 1; debug!( @@ -144,7 +170,7 @@ pub async fn send_with_retry( let req = match builder.body(payload.clone()) { Ok(r) => r, Err(_) => { - return Err(SendWithRetryError::Build(request_attempt)); + break Err(SendWithRetryError::Build(request_attempt)); } }; @@ -186,7 +212,7 @@ pub async fn send_with_retry( attempts = request_attempt, "Max retries exceeded, returning HTTP error" ); - return Err(SendWithRetryError::Http(response, request_attempt)); + break Err(SendWithRetryError::Http(response, request_attempt)); } } else { debug!( @@ -194,7 +220,7 @@ pub async fn send_with_retry( attempts = request_attempt, "Request succeeded" ); - return Ok((response, request_attempt)); + break Ok((response, request_attempt)); } } Ok(Err(e)) => { @@ -228,7 +254,7 @@ pub async fn send_with_retry( attempts = request_attempt, "Max retries exceeded, returning request error" ); - return Err(classified_error); + break Err(classified_error); } } Err(_) => { @@ -252,11 +278,12 @@ pub async fn send_with_retry( attempts = request_attempt, "Max retries exceeded, returning timeout error" ); - return Err(SendWithRetryError::Timeout(request_attempt)); + break Err(SendWithRetryError::Timeout(request_attempt)); } } } - } + }; + (result, payload_size) } #[cfg(test)] From 6c0dae8c8d58a72c5aa361e0382d8b5f0d3af385 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Tue, 25 Aug 2026 10:03:22 -0400 Subject: [PATCH 5/7] fix(data-pipeline): count spans after export filtering --- libdd-data-pipeline/src/trace_exporter/mod.rs | 77 +++++++++++++++++-- 1 file changed, 71 insertions(+), 6 deletions(-) diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index 01b2486596..fe010022a4 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -708,6 +708,8 @@ impl< let trace_count = traces.len(); #[cfg(feature = "telemetry")] let span_count = traces.iter().map(Vec::len).sum::(); + #[cfg(feature = "telemetry")] + self.emit_trace_counts(&[(MetricKind::SpansEnqueuedForSerialization, span_count)]); let json_body = libdd_trace_utils::agentless_encoder::encode_payload( &traces, &self.metadata, @@ -744,6 +746,8 @@ impl< let trace_count = traces.len(); #[cfg(feature = "telemetry")] let span_count = traces.iter().map(Vec::len).sum::(); + #[cfg(feature = "telemetry")] + self.emit_trace_counts(&[(MetricKind::SpansEnqueuedForSerialization, span_count)]); let resource_info = { let mut r = OtlpResourceInfo::default(); r.service = self.metadata.service.clone(); @@ -871,12 +875,6 @@ impl< return self.send_trace_chunks_to_log(&traces, max_line_size); } - #[cfg(feature = "telemetry")] - self.emit_trace_counts(&[( - MetricKind::SpansEnqueuedForSerialization, - traces.iter().map(Vec::len).sum(), - )]); - let mut header_tags: TracerHeaderTags = self.metadata.borrow().into(); if let Some(ref config) = self.agentless_config { @@ -926,6 +924,8 @@ impl< #[cfg(feature = "telemetry")] let trace_count = traces.len(); let span_count = traces.iter().map(Vec::len).sum::(); + #[cfg(feature = "telemetry")] + self.emit_trace_counts(&[(MetricKind::SpansEnqueuedForSerialization, span_count)]); let prepared = match self.serializer.prepare_traces_payload( traces, @@ -2446,6 +2446,7 @@ mod telemetry_metrics_tests { use libdd_tinybytes::BytesString; use libdd_trace_utils::msgpack_encoder; use libdd_trace_utils::span::{v04::SpanBytes, v05}; + use regex::Regex; // v05 messagepack empty payload -> [[""], []] const V5_EMPTY: [u8; 4] = [0x92, 0x91, 0xA0, 0x90]; @@ -2506,6 +2507,70 @@ mod telemetry_metrics_tests { metrics_endpoint.assert_calls(1); } + #[test] + #[cfg_attr(miri, ignore)] + fn test_otlp_enqueued_metric_excludes_dropped_p0_spans() { + let server = MockServer::start(); + let traces_endpoint = server.mock(|when, then| { + when.method(POST) + .path("/v1/traces") + .header("Content-Type", "application/json"); + then.status(200).body(""); + }); + let enqueued_metric = Regex::new( + r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,1\.0\]\],"tags":\[\],"common":true,"type":"count""#, + ) + .unwrap(); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .body_matches(enqueued_metric) + .path("/telemetry/proxy/api/v2/apmtelemetry"); + then.status(200) + .header("content-type", "application/json") + .body(""); + }); + + let otlp_endpoint = format!("{}/v1/traces", server.url("/").trim_end_matches('/')); + let mut builder = TraceExporter::::builder(); + builder + .set_url(&server.url("/")) + .set_service("foo") + .set_env("foo-env") + .set_tracer_version("v0.1") + .set_language("nodejs") + .set_language_version("1.0") + .set_language_interpreter("v8") + .set_otlp_endpoint(&otlp_endpoint) + .enable_telemetry(TelemetryConfig { + heartbeat: 100, + ..Default::default() + }); + let exporter = builder.build::().unwrap(); + + let kept_span = SpanBytes { + trace_id: 1, + span_id: 1, + ..Default::default() + }; + let dropped_span = SpanBytes { + trace_id: 2, + span_id: 2, + metrics: vec![(BytesString::from_static("_sampling_priority_v1"), -1.0)].into(), + ..Default::default() + }; + let traces = msgpack_encoder::v04::to_vec_from_v04(&[vec![kept_span], vec![dropped_span]]); + + assert_eq!( + exporter.send(traces.as_ref()).unwrap(), + AgentResponse::Unchanged + ); + traces_endpoint.assert_calls(1); + while metrics_endpoint.calls() == 0 { + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + } + #[test] #[cfg_attr(miri, ignore)] fn test_exporter_metrics_v5() { From 45359c6565cf30ad8f3ea17053595d8ee4feadb4 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Tue, 25 Aug 2026 10:42:15 -0400 Subject: [PATCH 6/7] refactor(data-pipeline): centralize payload telemetry counts --- libdd-data-pipeline/src/telemetry/mod.rs | 150 ++++- libdd-data-pipeline/src/trace_exporter/mod.rs | 538 +++++++++++++++--- libdd-trace-utils/src/send_with_retry/mod.rs | 32 ++ 3 files changed, 632 insertions(+), 88 deletions(-) diff --git a/libdd-data-pipeline/src/telemetry/mod.rs b/libdd-data-pipeline/src/telemetry/mod.rs index a69f8c752a..d402a32876 100644 --- a/libdd-data-pipeline/src/telemetry/mod.rs +++ b/libdd-data-pipeline/src/telemetry/mod.rs @@ -198,18 +198,6 @@ impl Tel worker: handle, } } - - pub(crate) fn send_count_metric( - &self, - metric: metrics::MetricKind, - count: usize, - ) -> Result<(), TelemetryError> { - if count > 0 { - self.worker - .add_point(count as f64, self.metrics.get(metric), vec![])?; - } - Ok(()) - } } /// Telemetry describing the sending of a trace payload @@ -224,6 +212,9 @@ pub struct SendPayloadTelemetry { chunks_sent: u64, chunks_dropped_serialization_error: u64, chunks_dropped_send_failure: u64, + spans_enqueued_for_serialization: u64, + spans_dropped_serialization_error: u64, + spans_dropped_api_error: u64, responses_count_per_code: HashMap, } @@ -293,6 +284,42 @@ impl SendPayloadTelemetry { }; telemetry } + + pub(crate) fn from_retry_result_with_spans( + value: &SendWithRetryResult, + bytes_sent: u64, + chunks: u64, + spans: u64, + ) -> Self { + let mut telemetry = Self::from_retry_result(value, bytes_sent, chunks); + telemetry.spans_enqueued_for_serialization = spans; + match value { + Err(SendWithRetryError::Build(_)) => { + telemetry.spans_dropped_serialization_error = spans; + } + Err(_) => telemetry.spans_dropped_api_error = spans, + Ok(_) => {} + } + telemetry + } + + pub(crate) fn from_serialization_error(chunks: u64, spans: u64) -> Self { + Self { + chunks_dropped_serialization_error: chunks, + spans_enqueued_for_serialization: spans, + spans_dropped_serialization_error: spans, + ..Default::default() + } + } + + pub(crate) fn from_canceled_send(chunks: u64, spans: u64) -> Self { + Self { + chunks_dropped_send_failure: chunks, + spans_enqueued_for_serialization: spans, + spans_dropped_api_error: spans, + ..Default::default() + } + } } impl TelemetryClient { @@ -345,6 +372,25 @@ impl Tel self.worker .add_point(data.chunks_dropped_send_failure as f64, key, vec![])?; } + if data.spans_enqueued_for_serialization > 0 { + let key = self + .metrics + .get(metrics::MetricKind::SpansEnqueuedForSerialization); + self.worker + .add_point(data.spans_enqueued_for_serialization as f64, key, vec![])?; + } + if data.spans_dropped_serialization_error > 0 { + let key = self + .metrics + .get(metrics::MetricKind::SpansDroppedSerializationError); + self.worker + .add_point(data.spans_dropped_serialization_error as f64, key, vec![])?; + } + if data.spans_dropped_api_error > 0 { + let key = self.metrics.get(metrics::MetricKind::SpansDroppedApiError); + self.worker + .add_point(data.spans_dropped_api_error as f64, key, vec![])?; + } if !data.responses_count_per_code.is_empty() { let key = self.metrics.get(metrics::MetricKind::ApiResponses); for (status_code, count) in &data.responses_count_per_code { @@ -772,6 +818,43 @@ mod tests { .expect("Failed to get runtime"); } + #[cfg_attr(miri, ignore)] + #[test] + fn span_counts_test() { + let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,7\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); + let serialization = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,3\.0\]\],"tags":\["reason:serialization_error"\],"common":true,"type":"count""#).unwrap(); + let api = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,4\.0\]\],"tags":\["reason:api_error"\],"common":true,"type":"count""#).unwrap(); + let shared_runtime = ForkSafeRuntime::new().expect("Failed to create runtime"); + let server = MockServer::start(); + let mut telemetry_srv = server.mock(|when, then| { + when.method(POST) + .body_matches(enqueued) + .body_matches(serialization) + .body_matches(api); + then.status(200).body(""); + }); + let data = SendPayloadTelemetry { + spans_enqueued_for_serialization: 7, + spans_dropped_serialization_error: 3, + spans_dropped_api_error: 4, + ..Default::default() + }; + let (client, handle) = get_test_client(&server.url("/"), &shared_runtime); + shared_runtime + .block_on(async { + let _ = client.start(); + let _ = client.send(&data); + sleep(Duration::from_millis(100)).await; + + handle.stop().await.expect("Failed to stop worker"); + assert!( + poll_for_mock_hits(&mut telemetry_srv, 1000, 10, 1).await, + "telemetry server did not receive calls within timeout" + ); + }) + .expect("Failed to get runtime"); + } + #[test] fn telemetry_from_ok_response_test() { let result = Ok(( @@ -862,6 +945,49 @@ mod tests { ) } + #[test] + fn telemetry_from_retry_result_with_spans_test() { + let result = Err(SendWithRetryError::Timeout(5)); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 7); + assert_eq!( + telemetry, + SendPayloadTelemetry { + chunks_dropped_send_failure: 2, + spans_enqueued_for_serialization: 7, + spans_dropped_api_error: 7, + requests_count: 5, + errors_timeout: 1, + ..Default::default() + } + ) + } + + #[test] + fn telemetry_from_serialization_error_test() { + assert_eq!( + SendPayloadTelemetry::from_serialization_error(2, 7), + SendPayloadTelemetry { + chunks_dropped_serialization_error: 2, + spans_enqueued_for_serialization: 7, + spans_dropped_serialization_error: 7, + ..Default::default() + } + ) + } + + #[test] + fn telemetry_from_canceled_send_test() { + assert_eq!( + SendPayloadTelemetry::from_canceled_send(2, 7), + SendPayloadTelemetry { + chunks_dropped_send_failure: 2, + spans_enqueued_for_serialization: 7, + spans_dropped_api_error: 7, + ..Default::default() + } + ) + } + #[test] fn telemetry_from_send_data_result_test() { let result = SendDataResult { diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index fe010022a4..840112c13e 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -23,8 +23,6 @@ use crate::agentless::AgentlessTraceConfig; use crate::otlp::exporter::{send_otlp_http_with_observer, OTLP_MAX_RETRIES}; use crate::otlp::{map_traces_to_otlp, OtlpResourceInfo, OtlpTraceConfig}; #[cfg(feature = "telemetry")] -use crate::telemetry::metrics::MetricKind; -#[cfg(feature = "telemetry")] use crate::telemetry::{SendPayloadTelemetry, TelemetryClient}; use crate::trace_exporter::agent_response::{ AgentResponsePayloadVersion, DATADOG_RATES_PAYLOAD_VERSION, @@ -72,6 +70,92 @@ const V04_TRACES_ENDPOINT: &str = "/v0.4/traces"; const V05_TRACES_ENDPOINT: &str = "/v0.5/traces"; const V1_TRACES_ENDPOINT: &str = "/v1.0/traces"; +#[derive(Clone, Copy)] +struct PayloadCounts { + chunks: usize, + #[cfg(feature = "telemetry")] + spans: usize, +} + +impl PayloadCounts { + fn from_traces(traces: &[Vec>]) -> Self { + Self { + chunks: traces.len(), + #[cfg(feature = "telemetry")] + spans: traces.iter().map(Vec::len).sum(), + } + } +} + +#[cfg(feature = "telemetry")] +struct PayloadTelemetryGuard< + 'a, + C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static, +> { + telemetry: &'a ArcSwapOption>, + counts: PayloadCounts, + pending: bool, +} + +#[cfg(feature = "telemetry")] +impl<'a, C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> + PayloadTelemetryGuard<'a, C> +{ + fn new(telemetry: &'a ArcSwapOption>, counts: PayloadCounts) -> Self { + Self { + telemetry, + counts, + pending: true, + } + } + + fn emit(&mut self, payload: &SendPayloadTelemetry) { + self.pending = false; + if let Some(telemetry) = self.telemetry.load_full().as_deref() { + if let Err(e) = telemetry.send(payload) { + error!(?e, "Error sending telemetry"); + } + } + } + + fn emit_retry_result(&mut self, result: &SendWithRetryResult, bytes: usize) { + let payload = SendPayloadTelemetry::from_retry_result_with_spans( + result, + bytes as u64, + self.counts.chunks as u64, + self.counts.spans as u64, + ); + self.emit(&payload); + } + + fn emit_serialization_error(&mut self) { + let payload = SendPayloadTelemetry::from_serialization_error( + self.counts.chunks as u64, + self.counts.spans as u64, + ); + self.emit(&payload); + } + + fn is_pending(&self) -> bool { + self.pending + } +} + +#[cfg(feature = "telemetry")] +impl Drop + for PayloadTelemetryGuard<'_, C> +{ + fn drop(&mut self) { + if self.pending { + let payload = SendPayloadTelemetry::from_canceled_send( + self.counts.chunks as u64, + self.counts.spans as u64, + ); + self.emit(&payload); + } + } +} + /// Build the HTTP headers required by the agentless intake. /// /// Includes the API key, content-type, trace count, `Datadog-Meta-*` tracer headers, @@ -313,47 +397,20 @@ impl< } #[cfg(feature = "telemetry")] - fn emit_trace_counts(&self, counts: &[(MetricKind, usize)]) { + fn emit_payload_telemetry(&self, payload: &SendPayloadTelemetry) { if let Some(telemetry) = self.telemetry.load_full().as_deref() { - for &(metric, count) in counts { - if let Err(e) = telemetry.send_count_metric(metric, count) { - error!(?e, "Error sending telemetry"); - } + if let Err(e) = telemetry.send(payload) { + error!(?e, "Error sending telemetry"); } } } #[cfg(feature = "telemetry")] - fn emit_serialization_drop(&self, chunks: usize, spans: usize) { - self.emit_trace_counts(&[ - (MetricKind::ChunksDroppedSerializationError, chunks), - (MetricKind::SpansDroppedSerializationError, spans), - ]); - } - - #[cfg(feature = "telemetry")] - fn emit_send_telemetry( - &self, - result: &SendWithRetryResult, - bytes: usize, - chunks: usize, - spans: usize, - ) { - if let Some(telemetry) = self.telemetry.load_full().as_deref() { - let payload = - SendPayloadTelemetry::from_retry_result(result, bytes as u64, chunks as u64); - if let Err(e) = telemetry.send(&payload) { - error!(?e, "Error sending telemetry"); - } - let metric = match result { - Err(SendWithRetryError::Build(_)) => MetricKind::SpansDroppedSerializationError, - Err(_) => MetricKind::SpansDroppedApiError, - Ok(_) => return, - }; - if let Err(e) = telemetry.send_count_metric(metric, spans) { - error!(?e, "Error sending telemetry"); - } - } + fn emit_serialization_drop(&self, counts: PayloadCounts) { + self.emit_payload_telemetry(&SendPayloadTelemetry::from_serialization_error( + counts.chunks as u64, + counts.spans as u64, + )); } /// Stop the background workers owned by this exporter. @@ -705,11 +762,7 @@ impl< traces: Vec>>, config: &AgentlessTraceConfig, ) -> Result { - let trace_count = traces.len(); - #[cfg(feature = "telemetry")] - let span_count = traces.iter().map(Vec::len).sum::(); - #[cfg(feature = "telemetry")] - self.emit_trace_counts(&[(MetricKind::SpansEnqueuedForSerialization, span_count)]); + let counts = PayloadCounts::from_traces(&traces); let json_body = libdd_trace_utils::agentless_encoder::encode_payload( &traces, &self.metadata, @@ -717,22 +770,37 @@ impl< .map_err(|e| { error!("Agentless JSON serialization error: {e}"); #[cfg(feature = "telemetry")] - self.emit_serialization_drop(trace_count, span_count); + self.emit_serialization_drop(counts); TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(e.to_string())) })?; - let headers = build_agentless_headers(&self.metadata, &config.api_key, trace_count)?; - send_agentless_traces_http_with_observer( + #[cfg(feature = "telemetry")] + let mut telemetry = PayloadTelemetryGuard::new(&self.telemetry, counts); + let headers = match build_agentless_headers(&self.metadata, &config.api_key, counts.chunks) + { + Ok(headers) => headers, + Err(e) => { + #[cfg(feature = "telemetry")] + telemetry.emit_serialization_error(); + return Err(e); + } + }; + let result = send_agentless_traces_http_with_observer( &self.capabilities, config, headers, json_body, |_result, _payload_len| { #[cfg(feature = "telemetry")] - self.emit_send_telemetry(_result, _payload_len, trace_count, span_count); + telemetry.emit_retry_result(_result, _payload_len); }, ) - .await?; + .await; + #[cfg(feature = "telemetry")] + if result.is_err() && telemetry.is_pending() { + telemetry.emit_serialization_error(); + } + result?; Ok(AgentResponse::Unchanged) } @@ -743,11 +811,7 @@ impl< config: &OtlpTraceConfig, ) -> Result { #[cfg(feature = "telemetry")] - let trace_count = traces.len(); - #[cfg(feature = "telemetry")] - let span_count = traces.iter().map(Vec::len).sum::(); - #[cfg(feature = "telemetry")] - self.emit_trace_counts(&[(MetricKind::SpansEnqueuedForSerialization, span_count)]); + let chunk_count = traces.len(); let resource_info = { let mut r = OtlpResourceInfo::default(); r.service = self.metadata.service.clone(); @@ -766,14 +830,26 @@ impl< // the mapper. let request = map_traces_to_otlp(traces, &resource_info, config.otel_trace_semantics_enabled); + #[cfg(feature = "telemetry")] + let counts = PayloadCounts { + chunks: chunk_count, + spans: request + .resource_spans + .iter() + .flat_map(|resource| &resource.scope_spans) + .map(|scope| scope.spans.len()) + .sum(), + }; let body = config.protocol.encode(&request).map_err(|e| { error!("OTLP serialization error: {e}"); #[cfg(feature = "telemetry")] - self.emit_serialization_drop(trace_count, span_count); + self.emit_serialization_drop(counts); TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(format!( "failed to encode OTLP request: {e}" ))) })?; + #[cfg(feature = "telemetry")] + let mut telemetry = PayloadTelemetryGuard::new(&self.telemetry, counts); // Also set the header: resource attributes survive Collector hops, headers don't. let effective_config; let config_to_use = if self.otlp_stats_enabled { @@ -791,7 +867,7 @@ impl< }; #[cfg(feature = "telemetry")] let payload_len = body.len(); - send_otlp_http_with_observer( + let result = send_otlp_http_with_observer( &self.capabilities, &config_to_use.endpoint_url, &config_to_use.headers, @@ -802,10 +878,15 @@ impl< OTLP_MAX_RETRIES, |_result| { #[cfg(feature = "telemetry")] - self.emit_send_telemetry(_result, payload_len, trace_count, span_count); + telemetry.emit_retry_result(_result, payload_len); }, ) - .await?; + .await; + #[cfg(feature = "telemetry")] + if result.is_err() && telemetry.is_pending() { + telemetry.emit_serialization_error(); + } + result?; Ok(AgentResponse::Unchanged) } @@ -815,13 +896,12 @@ impl< endpoint: &Endpoint, mp_payload: Vec, headers: HeaderMap, - chunks: usize, - spans: usize, + counts: PayloadCounts, ) -> Result { - #[cfg(not(feature = "telemetry"))] - let _ = spans; let strategy = RetryStrategy::default(); let payload_len = mp_payload.len(); + #[cfg(feature = "telemetry")] + let mut telemetry = PayloadTelemetryGuard::new(&self.telemetry, counts); // Send traces to the agent let result = send_with_retry( @@ -835,9 +915,10 @@ impl< .await; #[cfg(feature = "telemetry")] - self.emit_send_telemetry(&result, payload_len, chunks, spans); + telemetry.emit_retry_result(&result, payload_len); - self.handle_send_result(result, chunks, payload_len).await + self.handle_send_result(result, counts.chunks, payload_len) + .await } /// Synchronous log-export path: encode every span to newline-delimited @@ -921,11 +1002,7 @@ impl< // Snapshot the effective format once so the serializer and the URL agree even if // `v1_active` flips mid-send (the background `/info` fetcher can race us otherwise). let effective_format = self.effective_output_format(); - #[cfg(feature = "telemetry")] - let trace_count = traces.len(); - let span_count = traces.iter().map(Vec::len).sum::(); - #[cfg(feature = "telemetry")] - self.emit_trace_counts(&[(MetricKind::SpansEnqueuedForSerialization, span_count)]); + let counts = PayloadCounts::from_traces(&traces); let prepared = match self.serializer.prepare_traces_payload( traces, @@ -942,10 +1019,14 @@ impl< None, ); #[cfg(feature = "telemetry")] - self.emit_serialization_drop(trace_count, span_count); + self.emit_serialization_drop(counts); return Err(e); } }; + let counts = PayloadCounts { + chunks: prepared.chunk_count, + ..counts + }; let endpoint = Endpoint { url: effective_format.add_path(&self.endpoint.url), @@ -953,13 +1034,7 @@ impl< }; let result = self - .send_traces_with_telemetry( - &endpoint, - prepared.data, - prepared.headers, - prepared.chunk_count, - span_count, - ) + .send_traces_with_telemetry(&endpoint, prepared.data, prepared.headers, counts) .await; // State-hash trap mitigation: the agent does not return a `Datadog-Agent-State` @@ -1218,6 +1293,76 @@ mod tests { use libdd_trace_utils::span::v04::SpanBytes; use std::net; + #[derive(Clone, Debug)] + struct RejectingCapabilities; + + impl HttpClientCapability for RejectingCapabilities { + fn new_client() -> Self { + Self + } + + fn new_without_connection_pooling() -> Self { + Self + } + + fn request( + &self, + _req: http::Request, + ) -> impl std::future::Future< + Output = Result, libdd_capabilities::HttpError>, + > + MaybeSend { + std::future::ready(Ok(http::Response::builder() + .status(http::StatusCode::INTERNAL_SERVER_ERROR) + .body(bytes::Bytes::new()) + .unwrap())) + } + } + + impl SleepCapability for RejectingCapabilities { + fn new() -> Self { + Self + } + + fn sleep(&self, _duration: Duration) -> impl std::future::Future + MaybeSend { + std::future::ready(()) + } + } + + #[test] + fn public_transport_senders_preserve_error_contract() { + let runtime = ForkSafeRuntime::new().unwrap(); + let agentless = AgentlessTraceConfig { + endpoint_url: "http://localhost/v1/input".to_string(), + api_key: "test-key".to_string(), + timeout: Duration::from_secs(1), + }; + let otlp = OtlpTraceConfig { + endpoint_url: "http://localhost/v1/traces".to_string(), + headers: HeaderMap::new(), + timeout: Duration::from_secs(1), + protocol: crate::otlp::OtlpProtocol::HttpJson, + instrumentation_scope_name: String::new(), + instrumentation_scope_version: String::new(), + otel_trace_semantics_enabled: false, + }; + + runtime + .block_on(async { + crate::agentless::send_agentless_traces_http( + &RejectingCapabilities, + &agentless, + HeaderMap::new(), + Vec::new(), + ) + .await + .unwrap_err(); + crate::otlp::send_otlp_traces_http(&RejectingCapabilities, &otlp, None, Vec::new()) + .await + .unwrap_err(); + }) + .unwrap(); + } + #[test] fn test_from_tracer_tags_to_tracer_header_tags() { let tracer_tags = TracerMetadata { @@ -2438,6 +2583,7 @@ mod tests { #[cfg(feature = "telemetry")] mod telemetry_metrics_tests { use super::*; + use crate::telemetry::TelemetryClientBuilder; use crate::trace_exporter::tests::build_test_exporter; use httpmock::prelude::*; use httpmock::MockServer; @@ -2451,6 +2597,124 @@ mod telemetry_metrics_tests { // v05 messagepack empty payload -> [[""], []] const V5_EMPTY: [u8; 4] = [0x92, 0x91, 0xA0, 0x90]; + fn start_test_telemetry( + server: &MockServer, + ) -> ( + ForkSafeRuntime, + TelemetryClient, + WorkerHandle, + ) { + let runtime = ForkSafeRuntime::new().unwrap(); + let (client, worker) = TelemetryClientBuilder::default() + .set_service_name("foo") + .set_service_version("1.0") + .set_env("foo-env") + .set_language("nodejs") + .set_language_version("1.0") + .set_tracer_version("v0.1") + .set_url(&server.url("/")) + .set_heartbeat(100) + .build::() + .unwrap(); + let worker = runtime.spawn_worker(worker, true).unwrap(); + client.start().unwrap(); + (runtime, client, worker) + } + + #[test] + #[cfg_attr(miri, ignore)] + fn test_canceled_payload_emits_terminal_metrics_once() { + let server = MockServer::start(); + let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,3\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); + let spans_dropped = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,3\.0\]\],"tags":\["reason:api_error"\],"common":true,"type":"count""#).unwrap(); + let chunks_dropped = Regex::new(r#""metric":"trace_chunks_dropped","points":\[\[\d+,2\.0\]\],"tags":\["src_library:libdatadog","reason:send_failure"\],"common":true,"type":"count""#).unwrap(); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .body_matches(enqueued) + .body_matches(spans_dropped) + .body_matches(chunks_dropped); + then.status(200).body(""); + }); + let (runtime, client, worker) = start_test_telemetry(&server); + let telemetry = ArcSwapOption::new(None); + let guard = PayloadTelemetryGuard::::new( + &telemetry, + PayloadCounts { + chunks: 2, + spans: 3, + }, + ); + telemetry.store(Some(Arc::new(TelemetryClient::with_handle( + client.clone_handle(), + )))); + + drop(guard); + + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + runtime.block_on(worker.stop()).unwrap().unwrap(); + } + + #[test] + #[cfg_attr(miri, ignore)] + fn test_exporter_metrics_agentless_with_shared_telemetry() { + let server = MockServer::start(); + let traces_endpoint = server.mock(|when, then| { + when.method(POST) + .path("/v1/input") + .header("dd-api-key", "test-api-key"); + then.status(200).body(""); + }); + let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,3\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); + let chunks_sent = Regex::new(r#""metric":"trace_chunks_sent","points":\[\[\d+,2\.0\]\],"tags":\["src_library:libdatadog"\],"common":true,"type":"count""#).unwrap(); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .body_matches(enqueued) + .body_matches(chunks_sent); + then.status(200).body(""); + }); + let (runtime, client, worker) = start_test_telemetry(&server); + + let intake_endpoint = format!("{}/v1/input", server.url("/").trim_end_matches('/')); + let mut builder = TraceExporter::::builder(); + builder + .set_service("foo") + .set_env("foo-env") + .set_tracer_version("v0.1") + .set_language("nodejs") + .set_language_version("1.0") + .set_language_interpreter("v8") + .set_agentless_endpoint(&intake_endpoint, "test-api-key"); + let exporter = builder.build::().unwrap(); + exporter.set_telemetry_handle(Some(client.clone_handle())); + + let result = exporter + .send_trace_chunks( + vec![ + vec![SpanBytes::default(), SpanBytes::default()], + vec![SpanBytes::default()], + ], + None, + ) + .unwrap(); + assert_eq!(result, AgentResponse::Unchanged); + + traces_endpoint.assert_calls(1); + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + runtime.block_on(worker.stop()).unwrap().unwrap(); + } + #[test] #[cfg_attr(miri, ignore)] fn test_exporter_metrics_v4() { @@ -2565,7 +2829,10 @@ mod telemetry_metrics_tests { AgentResponse::Unchanged ); traces_endpoint.assert_calls(1); - while metrics_endpoint.calls() == 0 { + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } std::thread::sleep(Duration::from_millis(100)); } metrics_endpoint.assert_calls(1); @@ -2691,8 +2958,127 @@ mod single_threaded_tests { use httpmock::prelude::*; use libdd_capabilities_impl::NativeCapabilities; use libdd_shared_runtime::ForkSafeRuntime; + use libdd_tinybytes::BytesString; use libdd_trace_utils::msgpack_encoder; use libdd_trace_utils::span::v04::SpanBytes; + use regex::Regex; + + #[cfg_attr(miri, ignore)] + #[test] + fn test_client_side_stats_preserve_post_filter_payload_counts() { + agent_info::clear_cache_for_test(); + + let server = MockServer::start(); + let mock_traces = server.mock(|when, then| { + when.method(POST) + .path(V04_TRACES_ENDPOINT) + .header("x-datadog-trace-count", "1") + .header("datadog-client-computed-stats", "true") + .header("datadog-client-dropped-p0-traces", "1") + .header("datadog-client-dropped-p0-spans", "1"); + then.status(200).body(r#"{"rate_by_service":{}}"#); + }); + let mock_stats = server.mock(|when, then| { + when.method(POST).path(STATS_ENDPOINT); + then.status(200).body(""); + }); + let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,2\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); + let metrics_endpoint = server.mock(|when, then| { + when.method(POST) + .path("/telemetry/proxy/api/v2/apmtelemetry") + .body_matches(enqueued); + then.status(200).body(""); + }); + let _mock_info = server.mock(|when, then| { + when.method(GET).path(INFO_ENDPOINT); + then.status(200) + .header("content-type", "application/json") + .header("datadog-agent-state", "css-counts") + .body(format!( + r#"{{"version":"1","client_drop_p0s":true,"endpoints":["{V04_TRACES_ENDPOINT}","{STATS_ENDPOINT}"],"filter_tags":{{"reject":["drop:true"]}}}}"# + )); + }); + + let runtime = Arc::new(ForkSafeRuntime::new().unwrap()); + let mut builder = TraceExporter::::builder(); + builder + .set_url(&server.url("/")) + .set_service("test") + .set_env("staging") + .set_tracer_version("v0.1") + .set_language("nodejs") + .set_language_version("1.0") + .set_language_interpreter("v8") + .set_shared_runtime(runtime.clone()) + .enable_stats(Duration::from_secs(10)) + .enable_telemetry(TelemetryConfig { + heartbeat: 100, + ..Default::default() + }); + let exporter = builder.build::().unwrap(); + + while agent_info::get_agent_info().is_none() { + std::thread::sleep(Duration::from_millis(10)); + } + + let kept = vec![ + SpanBytes { + trace_id: 1, + span_id: 1, + duration: 10, + ..Default::default() + }, + SpanBytes { + trace_id: 1, + span_id: 2, + parent_id: 1, + duration: 10, + ..Default::default() + }, + ]; + let dropped_p0 = vec![SpanBytes { + trace_id: 2, + span_id: 3, + duration: 10, + metrics: vec![(BytesString::from_static("_sampling_priority_v1"), -1.0)].into(), + ..Default::default() + }]; + let filtered = vec![SpanBytes { + trace_id: 3, + span_id: 4, + duration: 10, + meta: vec![( + BytesString::from_static("drop"), + BytesString::from_static("true"), + )] + .into(), + ..Default::default() + }]; + let data = msgpack_encoder::v04::to_vec_from_v04(&[kept, dropped_p0, filtered]); + + assert!(matches!( + exporter.send(data.as_ref()).unwrap(), + AgentResponse::Changed { .. } + )); + mock_traces.assert_calls(1); + for _ in 0..50 { + if metrics_endpoint.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + metrics_endpoint.assert_calls(1); + + runtime.shutdown(None).unwrap(); + for _ in 0..100 { + if mock_stats.calls() > 0 { + break; + } + std::thread::sleep(Duration::from_millis(10)); + } + mock_stats.assert(); + agent_info::clear_cache_for_test(); + } #[cfg_attr(miri, ignore)] #[test] diff --git a/libdd-trace-utils/src/send_with_retry/mod.rs b/libdd-trace-utils/src/send_with_retry/mod.rs index 50f4382aca..05379f081e 100644 --- a/libdd-trace-utils/src/send_with_retry/mod.rs +++ b/libdd-trace-utils/src/send_with_retry/mod.rs @@ -491,4 +491,36 @@ mod tests { "Expected only one request attempt" ); } + + #[cfg(feature = "compression")] + #[cfg_attr(miri, ignore)] + #[tokio::test] + async fn test_reports_compressed_payload_size() { + let server = MockServer::start(); + let mock_202 = server + .mock_async(|_when, then| { + then.status(202); + }) + .await; + let endpoint = Endpoint { + url: server.url("").parse().unwrap(), + ..Default::default() + }; + let payload = vec![0; 1024]; + let expected_size = zstd::encode_all(payload.as_slice(), 1).unwrap().len(); + + let (result, payload_size) = send_with_retry_and_size( + &NativeCapabilities::new_client(), + &endpoint, + payload, + &HeaderMap::new(), + &RetryStrategy::new(0, 0, RetryBackoffType::Constant, None), + CompressionStrategy::Zstd { level: 1 }, + ) + .await; + + assert!(result.is_ok()); + assert_eq!(payload_size, expected_size); + mock_202.assert_calls_async(1).await; + } } From 8c9d3fc3102f273e6d3b9512bd2163dc1906d3b3 Mon Sep 17 00:00:00 2001 From: Munir Abdinur Date: Tue, 25 Aug 2026 11:28:29 -0400 Subject: [PATCH 7/7] refactor(data-pipeline): reduce telemetry accounting diff --- libdd-data-pipeline/src/agentless/exporter.rs | 2 +- libdd-data-pipeline/src/otlp/exporter.rs | 2 +- libdd-data-pipeline/src/telemetry/mod.rs | 102 +--- libdd-data-pipeline/src/trace_exporter/mod.rs | 511 +----------------- libdd-trace-utils/src/send_with_retry/mod.rs | 35 -- 5 files changed, 23 insertions(+), 629 deletions(-) diff --git a/libdd-data-pipeline/src/agentless/exporter.rs b/libdd-data-pipeline/src/agentless/exporter.rs index 6b4b2fd615..61c0ef4c7b 100644 --- a/libdd-data-pipeline/src/agentless/exporter.rs +++ b/libdd-data-pipeline/src/agentless/exporter.rs @@ -22,7 +22,7 @@ const AGENTLESS_RETRY_DELAY_MS: u64 = 1000; /// `headers` should already contain all required headers (api key, content-type, meta-*, /// entity, trace-count, etc.). `test_token` is forwarded as `X-Datadog-Test-Session-Token` /// when set, enabling snapshot tests against a local mock. -#[allow(dead_code)] // Retains the existing contract while telemetry uses the observer variant. +#[allow(dead_code)] pub async fn send_agentless_traces_http( capabilities: &C, config: &AgentlessTraceConfig, diff --git a/libdd-data-pipeline/src/otlp/exporter.rs b/libdd-data-pipeline/src/otlp/exporter.rs index bf4aa8c9a8..a05a04cc6a 100644 --- a/libdd-data-pipeline/src/otlp/exporter.rs +++ b/libdd-data-pipeline/src/otlp/exporter.rs @@ -118,7 +118,7 @@ pub(crate) async fn send_otlp_http_with_observer< /// /// `test_token` is forwarded as `X-Datadog-Test-Session-Token` when set, enabling snapshot tests /// against the Datadog test agent's OTLP endpoint. -#[allow(dead_code)] // Retains the existing contract while telemetry uses the observer variant. +#[allow(dead_code)] pub async fn send_otlp_traces_http( capabilities: &C, config: &OtlpTraceConfig, diff --git a/libdd-data-pipeline/src/telemetry/mod.rs b/libdd-data-pipeline/src/telemetry/mod.rs index d402a32876..94a081b407 100644 --- a/libdd-data-pipeline/src/telemetry/mod.rs +++ b/libdd-data-pipeline/src/telemetry/mod.rs @@ -302,24 +302,6 @@ impl SendPayloadTelemetry { } telemetry } - - pub(crate) fn from_serialization_error(chunks: u64, spans: u64) -> Self { - Self { - chunks_dropped_serialization_error: chunks, - spans_enqueued_for_serialization: spans, - spans_dropped_serialization_error: spans, - ..Default::default() - } - } - - pub(crate) fn from_canceled_send(chunks: u64, spans: u64) -> Self { - Self { - chunks_dropped_send_failure: chunks, - spans_enqueued_for_serialization: spans, - spans_dropped_api_error: spans, - ..Default::default() - } - } } impl TelemetryClient { @@ -818,43 +800,6 @@ mod tests { .expect("Failed to get runtime"); } - #[cfg_attr(miri, ignore)] - #[test] - fn span_counts_test() { - let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,7\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); - let serialization = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,3\.0\]\],"tags":\["reason:serialization_error"\],"common":true,"type":"count""#).unwrap(); - let api = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,4\.0\]\],"tags":\["reason:api_error"\],"common":true,"type":"count""#).unwrap(); - let shared_runtime = ForkSafeRuntime::new().expect("Failed to create runtime"); - let server = MockServer::start(); - let mut telemetry_srv = server.mock(|when, then| { - when.method(POST) - .body_matches(enqueued) - .body_matches(serialization) - .body_matches(api); - then.status(200).body(""); - }); - let data = SendPayloadTelemetry { - spans_enqueued_for_serialization: 7, - spans_dropped_serialization_error: 3, - spans_dropped_api_error: 4, - ..Default::default() - }; - let (client, handle) = get_test_client(&server.url("/"), &shared_runtime); - shared_runtime - .block_on(async { - let _ = client.start(); - let _ = client.send(&data); - sleep(Duration::from_millis(100)).await; - - handle.stop().await.expect("Failed to stop worker"); - assert!( - poll_for_mock_hits(&mut telemetry_srv, 1000, 10, 1).await, - "telemetry server did not receive calls within timeout" - ); - }) - .expect("Failed to get runtime"); - } - #[test] fn telemetry_from_ok_response_test() { let result = Ok(( @@ -884,11 +829,13 @@ mod tests { .body(Bytes::new()) .unwrap(); let result = Err(SendWithRetryError::Http(error_response, 5)); - let telemetry = SendPayloadTelemetry::from_retry_result(&result, 1, 2); + let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 7); assert_eq!( telemetry, SendPayloadTelemetry { chunks_dropped_send_failure: 2, + spans_enqueued_for_serialization: 7, + spans_dropped_api_error: 7, requests_count: 5, errors_status_code: 1, responses_count_per_code: HashMap::from([(400, 1)]), @@ -945,49 +892,6 @@ mod tests { ) } - #[test] - fn telemetry_from_retry_result_with_spans_test() { - let result = Err(SendWithRetryError::Timeout(5)); - let telemetry = SendPayloadTelemetry::from_retry_result_with_spans(&result, 1, 2, 7); - assert_eq!( - telemetry, - SendPayloadTelemetry { - chunks_dropped_send_failure: 2, - spans_enqueued_for_serialization: 7, - spans_dropped_api_error: 7, - requests_count: 5, - errors_timeout: 1, - ..Default::default() - } - ) - } - - #[test] - fn telemetry_from_serialization_error_test() { - assert_eq!( - SendPayloadTelemetry::from_serialization_error(2, 7), - SendPayloadTelemetry { - chunks_dropped_serialization_error: 2, - spans_enqueued_for_serialization: 7, - spans_dropped_serialization_error: 7, - ..Default::default() - } - ) - } - - #[test] - fn telemetry_from_canceled_send_test() { - assert_eq!( - SendPayloadTelemetry::from_canceled_send(2, 7), - SendPayloadTelemetry { - chunks_dropped_send_failure: 2, - spans_enqueued_for_serialization: 7, - spans_dropped_api_error: 7, - ..Default::default() - } - ) - } - #[test] fn telemetry_from_send_data_result_test() { let result = SendDataResult { diff --git a/libdd-data-pipeline/src/trace_exporter/mod.rs b/libdd-data-pipeline/src/trace_exporter/mod.rs index 840112c13e..994ee8ca8f 100644 --- a/libdd-data-pipeline/src/trace_exporter/mod.rs +++ b/libdd-data-pipeline/src/trace_exporter/mod.rs @@ -87,75 +87,6 @@ impl PayloadCounts { } } -#[cfg(feature = "telemetry")] -struct PayloadTelemetryGuard< - 'a, - C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static, -> { - telemetry: &'a ArcSwapOption>, - counts: PayloadCounts, - pending: bool, -} - -#[cfg(feature = "telemetry")] -impl<'a, C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> - PayloadTelemetryGuard<'a, C> -{ - fn new(telemetry: &'a ArcSwapOption>, counts: PayloadCounts) -> Self { - Self { - telemetry, - counts, - pending: true, - } - } - - fn emit(&mut self, payload: &SendPayloadTelemetry) { - self.pending = false; - if let Some(telemetry) = self.telemetry.load_full().as_deref() { - if let Err(e) = telemetry.send(payload) { - error!(?e, "Error sending telemetry"); - } - } - } - - fn emit_retry_result(&mut self, result: &SendWithRetryResult, bytes: usize) { - let payload = SendPayloadTelemetry::from_retry_result_with_spans( - result, - bytes as u64, - self.counts.chunks as u64, - self.counts.spans as u64, - ); - self.emit(&payload); - } - - fn emit_serialization_error(&mut self) { - let payload = SendPayloadTelemetry::from_serialization_error( - self.counts.chunks as u64, - self.counts.spans as u64, - ); - self.emit(&payload); - } - - fn is_pending(&self) -> bool { - self.pending - } -} - -#[cfg(feature = "telemetry")] -impl Drop - for PayloadTelemetryGuard<'_, C> -{ - fn drop(&mut self) { - if self.pending { - let payload = SendPayloadTelemetry::from_canceled_send( - self.counts.chunks as u64, - self.counts.spans as u64, - ); - self.emit(&payload); - } - } -} - /// Build the HTTP headers required by the agentless intake. /// /// Includes the API key, content-type, trace count, `Datadog-Meta-*` tracer headers, @@ -397,22 +328,25 @@ impl< } #[cfg(feature = "telemetry")] - fn emit_payload_telemetry(&self, payload: &SendPayloadTelemetry) { + fn emit_serialization_drop(&self, counts: PayloadCounts) { + self.emit_retry_result(&Err(SendWithRetryError::Build(0)), 0, counts); + } + + #[cfg(feature = "telemetry")] + fn emit_retry_result(&self, result: &SendWithRetryResult, bytes: usize, counts: PayloadCounts) { if let Some(telemetry) = self.telemetry.load_full().as_deref() { - if let Err(e) = telemetry.send(payload) { + let payload = SendPayloadTelemetry::from_retry_result_with_spans( + result, + bytes as u64, + counts.chunks as u64, + counts.spans as u64, + ); + if let Err(e) = telemetry.send(&payload) { error!(?e, "Error sending telemetry"); } } } - #[cfg(feature = "telemetry")] - fn emit_serialization_drop(&self, counts: PayloadCounts) { - self.emit_payload_telemetry(&SendPayloadTelemetry::from_serialization_error( - counts.chunks as u64, - counts.spans as u64, - )); - } - /// Stop the background workers owned by this exporter. /// /// Sync facade over [`Self::shutdown_async`]; panics inside an existing tokio context. @@ -774,17 +708,7 @@ impl< TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(e.to_string())) })?; - #[cfg(feature = "telemetry")] - let mut telemetry = PayloadTelemetryGuard::new(&self.telemetry, counts); - let headers = match build_agentless_headers(&self.metadata, &config.api_key, counts.chunks) - { - Ok(headers) => headers, - Err(e) => { - #[cfg(feature = "telemetry")] - telemetry.emit_serialization_error(); - return Err(e); - } - }; + let headers = build_agentless_headers(&self.metadata, &config.api_key, counts.chunks)?; let result = send_agentless_traces_http_with_observer( &self.capabilities, config, @@ -792,14 +716,10 @@ impl< json_body, |_result, _payload_len| { #[cfg(feature = "telemetry")] - telemetry.emit_retry_result(_result, _payload_len); + self.emit_retry_result(_result, _payload_len, counts); }, ) .await; - #[cfg(feature = "telemetry")] - if result.is_err() && telemetry.is_pending() { - telemetry.emit_serialization_error(); - } result?; Ok(AgentResponse::Unchanged) } @@ -811,7 +731,7 @@ impl< config: &OtlpTraceConfig, ) -> Result { #[cfg(feature = "telemetry")] - let chunk_count = traces.len(); + let counts = PayloadCounts::from_traces(&traces); let resource_info = { let mut r = OtlpResourceInfo::default(); r.service = self.metadata.service.clone(); @@ -830,16 +750,6 @@ impl< // the mapper. let request = map_traces_to_otlp(traces, &resource_info, config.otel_trace_semantics_enabled); - #[cfg(feature = "telemetry")] - let counts = PayloadCounts { - chunks: chunk_count, - spans: request - .resource_spans - .iter() - .flat_map(|resource| &resource.scope_spans) - .map(|scope| scope.spans.len()) - .sum(), - }; let body = config.protocol.encode(&request).map_err(|e| { error!("OTLP serialization error: {e}"); #[cfg(feature = "telemetry")] @@ -848,8 +758,6 @@ impl< "failed to encode OTLP request: {e}" ))) })?; - #[cfg(feature = "telemetry")] - let mut telemetry = PayloadTelemetryGuard::new(&self.telemetry, counts); // Also set the header: resource attributes survive Collector hops, headers don't. let effective_config; let config_to_use = if self.otlp_stats_enabled { @@ -878,14 +786,10 @@ impl< OTLP_MAX_RETRIES, |_result| { #[cfg(feature = "telemetry")] - telemetry.emit_retry_result(_result, payload_len); + self.emit_retry_result(_result, payload_len, counts); }, ) .await; - #[cfg(feature = "telemetry")] - if result.is_err() && telemetry.is_pending() { - telemetry.emit_serialization_error(); - } result?; Ok(AgentResponse::Unchanged) } @@ -900,9 +804,6 @@ impl< ) -> Result { let strategy = RetryStrategy::default(); let payload_len = mp_payload.len(); - #[cfg(feature = "telemetry")] - let mut telemetry = PayloadTelemetryGuard::new(&self.telemetry, counts); - // Send traces to the agent let result = send_with_retry( &self.capabilities, @@ -915,7 +816,7 @@ impl< .await; #[cfg(feature = "telemetry")] - telemetry.emit_retry_result(&result, payload_len); + self.emit_retry_result(&result, payload_len, counts); self.handle_send_result(result, counts.chunks, payload_len) .await @@ -1293,76 +1194,6 @@ mod tests { use libdd_trace_utils::span::v04::SpanBytes; use std::net; - #[derive(Clone, Debug)] - struct RejectingCapabilities; - - impl HttpClientCapability for RejectingCapabilities { - fn new_client() -> Self { - Self - } - - fn new_without_connection_pooling() -> Self { - Self - } - - fn request( - &self, - _req: http::Request, - ) -> impl std::future::Future< - Output = Result, libdd_capabilities::HttpError>, - > + MaybeSend { - std::future::ready(Ok(http::Response::builder() - .status(http::StatusCode::INTERNAL_SERVER_ERROR) - .body(bytes::Bytes::new()) - .unwrap())) - } - } - - impl SleepCapability for RejectingCapabilities { - fn new() -> Self { - Self - } - - fn sleep(&self, _duration: Duration) -> impl std::future::Future + MaybeSend { - std::future::ready(()) - } - } - - #[test] - fn public_transport_senders_preserve_error_contract() { - let runtime = ForkSafeRuntime::new().unwrap(); - let agentless = AgentlessTraceConfig { - endpoint_url: "http://localhost/v1/input".to_string(), - api_key: "test-key".to_string(), - timeout: Duration::from_secs(1), - }; - let otlp = OtlpTraceConfig { - endpoint_url: "http://localhost/v1/traces".to_string(), - headers: HeaderMap::new(), - timeout: Duration::from_secs(1), - protocol: crate::otlp::OtlpProtocol::HttpJson, - instrumentation_scope_name: String::new(), - instrumentation_scope_version: String::new(), - otel_trace_semantics_enabled: false, - }; - - runtime - .block_on(async { - crate::agentless::send_agentless_traces_http( - &RejectingCapabilities, - &agentless, - HeaderMap::new(), - Vec::new(), - ) - .await - .unwrap_err(); - crate::otlp::send_otlp_traces_http(&RejectingCapabilities, &otlp, None, Vec::new()) - .await - .unwrap_err(); - }) - .unwrap(); - } - #[test] fn test_from_tracer_tags_to_tracer_header_tags() { let tracer_tags = TracerMetadata { @@ -2583,7 +2414,6 @@ mod tests { #[cfg(feature = "telemetry")] mod telemetry_metrics_tests { use super::*; - use crate::telemetry::TelemetryClientBuilder; use crate::trace_exporter::tests::build_test_exporter; use httpmock::prelude::*; use httpmock::MockServer; @@ -2592,129 +2422,10 @@ mod telemetry_metrics_tests { use libdd_tinybytes::BytesString; use libdd_trace_utils::msgpack_encoder; use libdd_trace_utils::span::{v04::SpanBytes, v05}; - use regex::Regex; // v05 messagepack empty payload -> [[""], []] const V5_EMPTY: [u8; 4] = [0x92, 0x91, 0xA0, 0x90]; - fn start_test_telemetry( - server: &MockServer, - ) -> ( - ForkSafeRuntime, - TelemetryClient, - WorkerHandle, - ) { - let runtime = ForkSafeRuntime::new().unwrap(); - let (client, worker) = TelemetryClientBuilder::default() - .set_service_name("foo") - .set_service_version("1.0") - .set_env("foo-env") - .set_language("nodejs") - .set_language_version("1.0") - .set_tracer_version("v0.1") - .set_url(&server.url("/")) - .set_heartbeat(100) - .build::() - .unwrap(); - let worker = runtime.spawn_worker(worker, true).unwrap(); - client.start().unwrap(); - (runtime, client, worker) - } - - #[test] - #[cfg_attr(miri, ignore)] - fn test_canceled_payload_emits_terminal_metrics_once() { - let server = MockServer::start(); - let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,3\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); - let spans_dropped = Regex::new(r#""metric":"spans_dropped","points":\[\[\d+,3\.0\]\],"tags":\["reason:api_error"\],"common":true,"type":"count""#).unwrap(); - let chunks_dropped = Regex::new(r#""metric":"trace_chunks_dropped","points":\[\[\d+,2\.0\]\],"tags":\["src_library:libdatadog","reason:send_failure"\],"common":true,"type":"count""#).unwrap(); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .body_matches(enqueued) - .body_matches(spans_dropped) - .body_matches(chunks_dropped); - then.status(200).body(""); - }); - let (runtime, client, worker) = start_test_telemetry(&server); - let telemetry = ArcSwapOption::new(None); - let guard = PayloadTelemetryGuard::::new( - &telemetry, - PayloadCounts { - chunks: 2, - spans: 3, - }, - ); - telemetry.store(Some(Arc::new(TelemetryClient::with_handle( - client.clone_handle(), - )))); - - drop(guard); - - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - runtime.block_on(worker.stop()).unwrap().unwrap(); - } - - #[test] - #[cfg_attr(miri, ignore)] - fn test_exporter_metrics_agentless_with_shared_telemetry() { - let server = MockServer::start(); - let traces_endpoint = server.mock(|when, then| { - when.method(POST) - .path("/v1/input") - .header("dd-api-key", "test-api-key"); - then.status(200).body(""); - }); - let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,3\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); - let chunks_sent = Regex::new(r#""metric":"trace_chunks_sent","points":\[\[\d+,2\.0\]\],"tags":\["src_library:libdatadog"\],"common":true,"type":"count""#).unwrap(); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .body_matches(enqueued) - .body_matches(chunks_sent); - then.status(200).body(""); - }); - let (runtime, client, worker) = start_test_telemetry(&server); - - let intake_endpoint = format!("{}/v1/input", server.url("/").trim_end_matches('/')); - let mut builder = TraceExporter::::builder(); - builder - .set_service("foo") - .set_env("foo-env") - .set_tracer_version("v0.1") - .set_language("nodejs") - .set_language_version("1.0") - .set_language_interpreter("v8") - .set_agentless_endpoint(&intake_endpoint, "test-api-key"); - let exporter = builder.build::().unwrap(); - exporter.set_telemetry_handle(Some(client.clone_handle())); - - let result = exporter - .send_trace_chunks( - vec![ - vec![SpanBytes::default(), SpanBytes::default()], - vec![SpanBytes::default()], - ], - None, - ) - .unwrap(); - assert_eq!(result, AgentResponse::Unchanged); - - traces_endpoint.assert_calls(1); - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - runtime.block_on(worker.stop()).unwrap().unwrap(); - } - #[test] #[cfg_attr(miri, ignore)] fn test_exporter_metrics_v4() { @@ -2771,73 +2482,6 @@ mod telemetry_metrics_tests { metrics_endpoint.assert_calls(1); } - #[test] - #[cfg_attr(miri, ignore)] - fn test_otlp_enqueued_metric_excludes_dropped_p0_spans() { - let server = MockServer::start(); - let traces_endpoint = server.mock(|when, then| { - when.method(POST) - .path("/v1/traces") - .header("Content-Type", "application/json"); - then.status(200).body(""); - }); - let enqueued_metric = Regex::new( - r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,1\.0\]\],"tags":\[\],"common":true,"type":"count""#, - ) - .unwrap(); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .body_matches(enqueued_metric) - .path("/telemetry/proxy/api/v2/apmtelemetry"); - then.status(200) - .header("content-type", "application/json") - .body(""); - }); - - let otlp_endpoint = format!("{}/v1/traces", server.url("/").trim_end_matches('/')); - let mut builder = TraceExporter::::builder(); - builder - .set_url(&server.url("/")) - .set_service("foo") - .set_env("foo-env") - .set_tracer_version("v0.1") - .set_language("nodejs") - .set_language_version("1.0") - .set_language_interpreter("v8") - .set_otlp_endpoint(&otlp_endpoint) - .enable_telemetry(TelemetryConfig { - heartbeat: 100, - ..Default::default() - }); - let exporter = builder.build::().unwrap(); - - let kept_span = SpanBytes { - trace_id: 1, - span_id: 1, - ..Default::default() - }; - let dropped_span = SpanBytes { - trace_id: 2, - span_id: 2, - metrics: vec![(BytesString::from_static("_sampling_priority_v1"), -1.0)].into(), - ..Default::default() - }; - let traces = msgpack_encoder::v04::to_vec_from_v04(&[vec![kept_span], vec![dropped_span]]); - - assert_eq!( - exporter.send(traces.as_ref()).unwrap(), - AgentResponse::Unchanged - ); - traces_endpoint.assert_calls(1); - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - } - #[test] #[cfg_attr(miri, ignore)] fn test_exporter_metrics_v5() { @@ -2958,127 +2602,8 @@ mod single_threaded_tests { use httpmock::prelude::*; use libdd_capabilities_impl::NativeCapabilities; use libdd_shared_runtime::ForkSafeRuntime; - use libdd_tinybytes::BytesString; use libdd_trace_utils::msgpack_encoder; use libdd_trace_utils::span::v04::SpanBytes; - use regex::Regex; - - #[cfg_attr(miri, ignore)] - #[test] - fn test_client_side_stats_preserve_post_filter_payload_counts() { - agent_info::clear_cache_for_test(); - - let server = MockServer::start(); - let mock_traces = server.mock(|when, then| { - when.method(POST) - .path(V04_TRACES_ENDPOINT) - .header("x-datadog-trace-count", "1") - .header("datadog-client-computed-stats", "true") - .header("datadog-client-dropped-p0-traces", "1") - .header("datadog-client-dropped-p0-spans", "1"); - then.status(200).body(r#"{"rate_by_service":{}}"#); - }); - let mock_stats = server.mock(|when, then| { - when.method(POST).path(STATS_ENDPOINT); - then.status(200).body(""); - }); - let enqueued = Regex::new(r#""metric":"spans_enqueued_for_serialization","points":\[\[\d+,2\.0\]\],"tags":\[\],"common":true,"type":"count""#).unwrap(); - let metrics_endpoint = server.mock(|when, then| { - when.method(POST) - .path("/telemetry/proxy/api/v2/apmtelemetry") - .body_matches(enqueued); - then.status(200).body(""); - }); - let _mock_info = server.mock(|when, then| { - when.method(GET).path(INFO_ENDPOINT); - then.status(200) - .header("content-type", "application/json") - .header("datadog-agent-state", "css-counts") - .body(format!( - r#"{{"version":"1","client_drop_p0s":true,"endpoints":["{V04_TRACES_ENDPOINT}","{STATS_ENDPOINT}"],"filter_tags":{{"reject":["drop:true"]}}}}"# - )); - }); - - let runtime = Arc::new(ForkSafeRuntime::new().unwrap()); - let mut builder = TraceExporter::::builder(); - builder - .set_url(&server.url("/")) - .set_service("test") - .set_env("staging") - .set_tracer_version("v0.1") - .set_language("nodejs") - .set_language_version("1.0") - .set_language_interpreter("v8") - .set_shared_runtime(runtime.clone()) - .enable_stats(Duration::from_secs(10)) - .enable_telemetry(TelemetryConfig { - heartbeat: 100, - ..Default::default() - }); - let exporter = builder.build::().unwrap(); - - while agent_info::get_agent_info().is_none() { - std::thread::sleep(Duration::from_millis(10)); - } - - let kept = vec![ - SpanBytes { - trace_id: 1, - span_id: 1, - duration: 10, - ..Default::default() - }, - SpanBytes { - trace_id: 1, - span_id: 2, - parent_id: 1, - duration: 10, - ..Default::default() - }, - ]; - let dropped_p0 = vec![SpanBytes { - trace_id: 2, - span_id: 3, - duration: 10, - metrics: vec![(BytesString::from_static("_sampling_priority_v1"), -1.0)].into(), - ..Default::default() - }]; - let filtered = vec![SpanBytes { - trace_id: 3, - span_id: 4, - duration: 10, - meta: vec![( - BytesString::from_static("drop"), - BytesString::from_static("true"), - )] - .into(), - ..Default::default() - }]; - let data = msgpack_encoder::v04::to_vec_from_v04(&[kept, dropped_p0, filtered]); - - assert!(matches!( - exporter.send(data.as_ref()).unwrap(), - AgentResponse::Changed { .. } - )); - mock_traces.assert_calls(1); - for _ in 0..50 { - if metrics_endpoint.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(100)); - } - metrics_endpoint.assert_calls(1); - - runtime.shutdown(None).unwrap(); - for _ in 0..100 { - if mock_stats.calls() > 0 { - break; - } - std::thread::sleep(Duration::from_millis(10)); - } - mock_stats.assert(); - agent_info::clear_cache_for_test(); - } #[cfg_attr(miri, ignore)] #[test] diff --git a/libdd-trace-utils/src/send_with_retry/mod.rs b/libdd-trace-utils/src/send_with_retry/mod.rs index 05379f081e..0e932c6f45 100644 --- a/libdd-trace-utils/src/send_with_retry/mod.rs +++ b/libdd-trace-utils/src/send_with_retry/mod.rs @@ -119,9 +119,6 @@ pub async fn send_with_retry( } /// Send a payload with retries and return its post-compression size. -/// -/// This is equivalent to [`send_with_retry`], with the payload size exposed for transport -/// telemetry. The result contains the final outcome and the total number of request attempts. #[allow(clippy::result_large_err)] pub async fn send_with_retry_and_size( capabilities: &C, @@ -491,36 +488,4 @@ mod tests { "Expected only one request attempt" ); } - - #[cfg(feature = "compression")] - #[cfg_attr(miri, ignore)] - #[tokio::test] - async fn test_reports_compressed_payload_size() { - let server = MockServer::start(); - let mock_202 = server - .mock_async(|_when, then| { - then.status(202); - }) - .await; - let endpoint = Endpoint { - url: server.url("").parse().unwrap(), - ..Default::default() - }; - let payload = vec![0; 1024]; - let expected_size = zstd::encode_all(payload.as_slice(), 1).unwrap().len(); - - let (result, payload_size) = send_with_retry_and_size( - &NativeCapabilities::new_client(), - &endpoint, - payload, - &HeaderMap::new(), - &RetryStrategy::new(0, 0, RetryBackoffType::Constant, None), - CompressionStrategy::Zstd { level: 1 }, - ) - .await; - - assert!(result.is_ok()); - assert_eq!(payload_size, expected_size); - mock_202.assert_calls_async(1).await; - } }