diff --git a/Cargo.lock b/Cargo.lock index a8d1b46..b36f51c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -663,6 +663,7 @@ dependencies = [ "anyhow", "async-trait", "bytes", + "datadog-agent-trace-sampler", "datadog-fips", "duplicate", "flate2", diff --git a/crates/datadog-serverless-compat/src/main.rs b/crates/datadog-serverless-compat/src/main.rs index cdcdcab..b7e22e9 100644 --- a/crates/datadog-serverless-compat/src/main.rs +++ b/crates/datadog-serverless-compat/src/main.rs @@ -186,6 +186,7 @@ pub async fn main() { let trace_processor = Arc::new(trace_processor::ServerlessTraceProcessor::new( stats_concentrator.as_ref().map(|c| c.handle.clone()), + trace_processor::ServerlessTraceProcessor::new_error_sampler(&config), )); let stats_flusher = Arc::new(stats_flusher::ServerlessStatsFlusher { diff --git a/crates/datadog-trace-agent/Cargo.toml b/crates/datadog-trace-agent/Cargo.toml index 291ff87..17734f9 100644 --- a/crates/datadog-trace-agent/Cargo.toml +++ b/crates/datadog-trace-agent/Cargo.toml @@ -33,6 +33,7 @@ libdd-trace-stats = { version = "7.0.0", features = ["stats-obfuscation"] } libdd-common = { workspace = true, features = ["https"] } libdd-trace-obfuscation = { workspace = true, features = ["https"] } libdd-trace-utils = { workspace = true, features = ["https", "mini_agent"] } +datadog-agent-trace-sampler = { path = "../datadog-agent-trace-sampler" } datadog-fips = { path = "../datadog-fips" } reqwest = { version = "0.12.23", features = [ "json", diff --git a/crates/datadog-trace-agent/src/config.rs b/crates/datadog-trace-agent/src/config.rs index 186ce0a..53976bd 100644 --- a/crates/datadog-trace-agent/src/config.rs +++ b/crates/datadog-trace-agent/src/config.rs @@ -9,14 +9,60 @@ use std::env; use std::str::FromStr; use std::sync::OnceLock; +use datadog_agent_trace_sampler::{ErrorSamplerConfig, ErrorSamplerMode}; use libdd_trace_obfuscation::obfuscation_config; use libdd_trace_utils::config_utils::{ read_cloud_env, trace_intake_url, trace_intake_url_prefixed, trace_stats_url, trace_stats_url_prefixed, }; use libdd_trace_utils::trace_utils; +use tracing::warn; const DEFAULT_APM_RECEIVER_PORT: u16 = 8126; +/// Default error rescue budget: error traces per second, matching the Go agent's `ErrorTPS`. +const DEFAULT_ERROR_SAMPLER_TPS: f64 = 10.0; + +/// Parses the error rescue sampler settings from the environment. +/// +/// - `DD_APM_ERROR_SAMPLER_MODE`: `rate_limited` (default) or `always_keep`. +/// Surrounding whitespace and casing are normalized; any other value warns +/// and falls back to `rate_limited`. +/// - `DD_APM_ERROR_TPS`: target error traces per second, default 10. Parsed as +/// `f64`; non-finite or unparsable values warn and fall back to the default. +/// A finite value <= 0 disables rescue in either mode. +/// +/// `extra_sample_rate` is fixed at 1.0 and intentionally not configurable. +/// Startup stays operational when either optional setting is malformed. +fn parse_error_sampler_config() -> ErrorSamplerConfig { + let mut config = ErrorSamplerConfig { + mode: ErrorSamplerMode::RateLimited, + target_tps: DEFAULT_ERROR_SAMPLER_TPS, + extra_sample_rate: 1.0, + }; + + if let Ok(raw) = env::var("DD_APM_ERROR_SAMPLER_MODE") { + match raw.trim().to_lowercase().as_str() { + "rate_limited" => config.mode = ErrorSamplerMode::RateLimited, + "always_keep" => config.mode = ErrorSamplerMode::AlwaysKeep, + _ => { + warn!("Invalid DD_APM_ERROR_SAMPLER_MODE {raw:?}; using default mode rate_limited"); + } + } + } + + if let Ok(raw) = env::var("DD_APM_ERROR_TPS") { + match raw.trim().parse::() { + Ok(tps) if tps.is_finite() => config.target_tps = tps, + _ => { + warn!( + "Invalid DD_APM_ERROR_TPS {raw:?}; using default {DEFAULT_ERROR_SAMPLER_TPS}" + ); + } + } + } + + config +} const DEFAULT_DOGSTATSD_PORT: u16 = 8125; const DSM_PIPELINE_STATS_ROUTE: &str = "/api/v0.1/pipeline_stats"; @@ -128,6 +174,10 @@ pub struct Config { pub additional_metric_tags_cardinality_limit: Option, /// Whether the agent should compute trace stats pub agent_stats_computation_enabled: bool, + /// Error rescue sampler settings, parsed from `DD_APM_ERROR_SAMPLER_MODE` + /// and `DD_APM_ERROR_TPS`. `extra_sample_rate` is fixed at 1.0 and not + /// configurable. See `parse_error_sampler_config` for the defaults. + pub error_sampler: ErrorSamplerConfig, } impl Config { @@ -304,6 +354,7 @@ impl Config { agent_stats_computation_enabled: env::var("DD_AGENT_STATS_COMPUTATION_ENABLED") .map(|val| val.to_lowercase() == "true") .unwrap_or(true), + error_sampler: parse_error_sampler_config(), }) } } @@ -315,6 +366,7 @@ mod tests { use std::collections::HashMap; use crate::config; + use datadog_agent_trace_sampler::ErrorSamplerMode; #[test] #[serial] @@ -888,6 +940,126 @@ mod tests { }, ); } + + fn assert_rate_limited_default(error_sampler: &config::ErrorSamplerConfig) { + assert!( + matches!(error_sampler.mode, ErrorSamplerMode::RateLimited), + "expected default mode RateLimited" + ); + assert_eq!(error_sampler.target_tps, 10.0); + assert_eq!(error_sampler.extra_sample_rate, 1.0); + } + + #[test] + #[serial] + fn test_error_sampler_defaults() { + temp_env::with_vars( + [ + ("DD_API_KEY", Some("_not_a_real_key_")), + ("FUNCTIONS_EXTENSION_VERSION", Some("~4")), + ("FUNCTIONS_WORKER_RUNTIME", Some("dotnet")), + ("WEBSITE_SITE_NAME", Some("my-azure-function")), + ], + || { + let config = config::Config::new().unwrap(); + assert_rate_limited_default(&config.error_sampler); + }, + ); + } + + #[test] + #[serial] + fn test_error_sampler_mode_normalization() { + for (raw, expected) in [ + ("always_keep", ErrorSamplerMode::AlwaysKeep), + ("ALWAYS_KEEP", ErrorSamplerMode::AlwaysKeep), + (" Always_Keep ", ErrorSamplerMode::AlwaysKeep), + ("rate_limited", ErrorSamplerMode::RateLimited), + ("RATE_LIMITED", ErrorSamplerMode::RateLimited), + (" Rate_Limited ", ErrorSamplerMode::RateLimited), + ] { + temp_env::with_vars( + [ + ("DD_API_KEY", Some("_not_a_real_key_")), + ("FUNCTIONS_EXTENSION_VERSION", Some("~4")), + ("FUNCTIONS_WORKER_RUNTIME", Some("dotnet")), + ("WEBSITE_SITE_NAME", Some("my-azure-function")), + ("DD_APM_ERROR_SAMPLER_MODE", Some(raw)), + ], + || { + let config = config::Config::new().unwrap(); + assert!( + matches!(config.error_sampler.mode, m if std::mem::discriminant(&m) == std::mem::discriminant(&expected)), + "mode {raw:?} should parse to {expected:?}, got {:?}", + config.error_sampler.mode + ); + }, + ); + } + } + + #[test] + #[serial] + fn test_error_sampler_invalid_mode_falls_back_to_default() { + for raw in ["bogus", "", "always-keep"] { + temp_env::with_vars( + [ + ("DD_API_KEY", Some("_not_a_real_key_")), + ("FUNCTIONS_EXTENSION_VERSION", Some("~4")), + ("FUNCTIONS_WORKER_RUNTIME", Some("dotnet")), + ("WEBSITE_SITE_NAME", Some("my-azure-function")), + ("DD_APM_ERROR_SAMPLER_MODE", Some(raw)), + ], + || { + let config = config::Config::new().unwrap(); + assert_rate_limited_default(&config.error_sampler); + }, + ); + } + } + + #[test] + #[serial] + fn test_error_sampler_tps_parsing() { + for (raw, expected_tps) in [("5.5", 5.5), ("0", 0.0), ("-1", -1.0), (" 10 ", 10.0)] { + temp_env::with_vars( + [ + ("DD_API_KEY", Some("_not_a_real_key_")), + ("FUNCTIONS_EXTENSION_VERSION", Some("~4")), + ("FUNCTIONS_WORKER_RUNTIME", Some("dotnet")), + ("WEBSITE_SITE_NAME", Some("my-azure-function")), + ("DD_APM_ERROR_TPS", Some(raw)), + ], + || { + let config = config::Config::new().unwrap(); + assert_eq!( + config.error_sampler.target_tps, expected_tps, + "DD_APM_ERROR_TPS {raw:?} should parse to {expected_tps}" + ); + }, + ); + } + } + + #[test] + #[serial] + fn test_error_sampler_invalid_tps_falls_back_to_default() { + for raw in ["not_a_number", "inf", "-inf", "NaN", "1e400", ""] { + temp_env::with_vars( + [ + ("DD_API_KEY", Some("_not_a_real_key_")), + ("FUNCTIONS_EXTENSION_VERSION", Some("~4")), + ("FUNCTIONS_WORKER_RUNTIME", Some("dotnet")), + ("WEBSITE_SITE_NAME", Some("my-azure-function")), + ("DD_APM_ERROR_TPS", Some(raw)), + ], + || { + let config = config::Config::new().unwrap(); + assert_rate_limited_default(&config.error_sampler); + }, + ); + } + } } /// Test helpers for creating Config instances in tests @@ -930,6 +1102,7 @@ pub mod test_helpers { additional_metric_tags: vec![], additional_metric_tags_cardinality_limit: None, agent_stats_computation_enabled: true, + error_sampler: ErrorSamplerConfig::default(), } } } diff --git a/crates/datadog-trace-agent/src/stats_concentrator_service.rs b/crates/datadog-trace-agent/src/stats_concentrator_service.rs index 01d5164..4689774 100644 --- a/crates/datadog-trace-agent/src/stats_concentrator_service.rs +++ b/crates/datadog-trace-agent/src/stats_concentrator_service.rs @@ -7,6 +7,7 @@ use std::sync::Arc; use tokio::sync::{mpsc, oneshot}; use crate::config::Config; +use crate::trace_processor::resolve_payload_env; use libdd_library_config::tracer_metadata::TracerMetadata; use libdd_trace_protobuf::pb::{ClientStatsPayload, TraceChunk}; use libdd_trace_stats::span_concentrator::{CardinalityLimitConfig, SpanConcentrator}; @@ -191,11 +192,11 @@ impl StatsConcentratorService { // Do not set hostname so the trace stats backend can aggregate stats properly hostname: String::new(), // Prefer env from the tracer payload, fall back to agent config - env: metadata - .service_env - .clone() - .filter(|s| !s.is_empty()) - .unwrap_or_else(|| self.config.env.clone()), + env: resolve_payload_env( + metadata.service_env.as_deref().unwrap_or(""), + &self.config.env, + ) + .to_string(), version: metadata.service_version.clone().unwrap_or_default(), lang: metadata.tracer_language.clone(), tracer_version: metadata.tracer_version.clone(), diff --git a/crates/datadog-trace-agent/src/stats_processor.rs b/crates/datadog-trace-agent/src/stats_processor.rs index 990d706..a50e59d 100644 --- a/crates/datadog-trace-agent/src/stats_processor.rs +++ b/crates/datadog-trace-agent/src/stats_processor.rs @@ -164,6 +164,7 @@ mod tests { additional_metric_tags: vec![], additional_metric_tags_cardinality_limit: None, agent_stats_computation_enabled, + error_sampler: datadog_agent_trace_sampler::ErrorSamplerConfig::default(), } } diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index e26465c..c150d5e 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -2,12 +2,16 @@ // SPDX-License-Identifier: Apache-2.0 use std::sync::Arc; +use std::sync::atomic::{AtomicI64, Ordering}; use async_trait::async_trait; +use datadog_agent_trace_sampler::{ErrorsSampler, SampleDecision, SpanView, TraceView}; use http_body_util::BodyExt; use hyper::{StatusCode, http}; use libdd_common::http_common; use libdd_library_config::tracer_metadata::TracerMetadata; +use std::sync::{Mutex, MutexGuard, PoisonError}; +use std::time::{SystemTime, UNIX_EPOCH}; use tokio::sync::mpsc::Sender; use tracing::{debug, error, warn}; @@ -27,6 +31,12 @@ use crate::{ const TRACER_PAYLOAD_FUNCTION_TAGS_TAG_KEY: &str = "_dd.tags.function"; +/// The root-span metric the backend uses to rescue traces the ordinary P0 drop +/// would discard: a positive `_dd.errors_sr` resolves the chunk's ingestion +/// reason to `error`, which is retained at low priority without any priority +/// promotion. +const ERRORS_SR_METRIC_KEY: &str = "_dd.errors_sr"; + /// Rough upper bound on the protobuf framing overhead added when a V07 `TracerPayload` is /// wrapped in the outer `AgentPayload` envelope before being sent const V07_ENVELOPE_OVERHEAD_BYTES: usize = 64; @@ -124,14 +134,118 @@ const MAX_IN_FLIGHT_ENQUEUES: usize = 10; pub struct ServerlessTraceProcessor { pub stats_concentrator: Option, enqueue_permits: Arc, + /// Shared error rescue sampler. The `Arc` means processor clones (one per + /// connection) share a single sampler state and TPS budget for the whole + /// process lifetime. + error_sampler: Arc>, + /// Last timestamp handed to the sampler, shared across clones. Clamps the + /// clock so a backward wall-clock step cannot move sampler bucket IDs + /// backwards. Updated only under the sampler lock. + last_sampler_timestamp: Arc, } impl ServerlessTraceProcessor { #[allow(clippy::must_use_candidate)] - pub fn new(stats_concentrator: Option) -> Self { + pub fn new( + stats_concentrator: Option, + error_sampler: Arc>, + ) -> Self { ServerlessTraceProcessor { stats_concentrator, enqueue_permits: Arc::new(tokio::sync::Semaphore::new(MAX_IN_FLIGHT_ENQUEUES)), + error_sampler, + last_sampler_timestamp: Arc::new(AtomicI64::new(i64::MIN)), + } + } + + /// Builds the error rescue sampler from parsed config settings. + #[must_use] + pub fn new_error_sampler(config: &Config) -> Arc> { + Arc::new(Mutex::new(ErrorsSampler::new(config.error_sampler))) + } + + /// Locks the shared sampler, recovering from a poisoned guard (a panic in + /// another thread while holding the lock) instead of panicking on unlock. + fn lock_sampler(&self) -> MutexGuard<'_, ErrorsSampler> { + self.error_sampler + .lock() + .unwrap_or_else(PoisonError::into_inner) + } + + /// Applies the error rescue pass to a payload collection. + /// + /// For every V07 chunk that would be dropped by the backend's ordinary P0 + /// drop (chunk priority is exactly 0, i.e. an automatic drop, not an + /// explicit user drop or the no-priority sentinel) and contains at least + /// one span with a non-zero error flag, consults the shared error sampler. + /// On a keep, stamps the sampler's `_dd.errors_sr` on the chunk's root + /// span; the backend keeps such chunks without any priority promotion. + /// Chunks are never removed or reordered: unrescued chunks stay in the + /// payload and the backend discards them. + /// + fn apply_error_rescue(&self, payload: &mut TracerPayloadCollection, config: &Config) { + let mut sampler = self.lock_sampler(); + // Skip all view construction when the sampler is disabled by config + // (target_tps <= 0): nothing can be rescued. + if sampler.is_disabled() { + return; + } + // Read the clock while holding the sampler lock so that concurrent + // requests deliver timestamps in lock-acquisition order, and clamp so a + // backward wall-clock step cannot move the sampler's rolling window + // backwards, which would undercount TPS and rescue too many chunks. + let now_unix_secs = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |d| d.as_secs() as i64); + let now_unix_secs = self.clamp_sampler_timestamp(now_unix_secs); + self.rescue_with_sampler(payload, config, now_unix_secs, &mut sampler); + } + + /// Test-only variant that injects a synthetic timestamp so tests can + /// exercise the rolling window without sleeping. + #[cfg(test)] + fn apply_error_rescue_at( + &self, + payload: &mut TracerPayloadCollection, + config: &Config, + now_unix_secs: i64, + ) { + let mut sampler = self.lock_sampler(); + if sampler.is_disabled() { + return; + } + self.rescue_with_sampler(payload, config, now_unix_secs, &mut sampler); + } + + /// Clamps the timestamp so it never moves backwards relative to the last + /// one handed to the shared sampler, and records it as the new floor. Must + /// be called while holding the sampler lock, alongside the clock read, so + /// concurrent requests cannot interleave a read with the clamp. + fn clamp_sampler_timestamp(&self, now_unix_secs: i64) -> i64 { + let previous = self + .last_sampler_timestamp + .fetch_max(now_unix_secs, Ordering::Relaxed); + now_unix_secs.max(previous) + } + + fn rescue_with_sampler( + &self, + payload: &mut TracerPayloadCollection, + config: &Config, + now_unix_secs: i64, + sampler: &mut ErrorsSampler, + ) { + let TracerPayloadCollection::V07(tracer_payloads) = payload else { + return; + }; + for tracer_payload in tracer_payloads.iter_mut() { + // The sampler keys its per-signature rate limits on the env the + // tracer reported for this payload, falling back to the agent's + // configured env, consistent with how stats are flushed. + let env: &str = resolve_payload_env(&tracer_payload.env, &config.env); + for chunk in tracer_payload.chunks.iter_mut() { + sample_and_stamp(sampler, chunk, env, now_unix_secs); + } } } @@ -179,6 +293,88 @@ impl ServerlessTraceProcessor { } } +/// Chooses the env used to key per-signature state: the payload-reported env +/// when nonempty, otherwise the agent's configured env. Shared by stats +/// flushing and error-rescue sampling so the fallback rule can't drift +/// between the two. +pub(crate) fn resolve_payload_env<'a>(payload_env: &'a str, config_env: &'a str) -> &'a str { + if payload_env.is_empty() { + config_env + } else { + payload_env + } +} + +/// Builds the sampler's read-only view of a span from its already enriched and +/// obfuscated fields. +fn span_view(span: &pb::Span) -> SpanView<'_> { + SpanView { + service: &span.service, + name: &span.name, + resource: &span.resource, + error: span.error != 0, + http_status_code: span.meta.get("http.status_code").map(String::as_str), + error_type: span.meta.get("error.type").map(String::as_str), + } +} + +/// Consults the error sampler for one chunk and, on a keep, stamps the +/// sampler's `_dd.errors_sr` value on the chunk's root span. +/// +/// Only automatic-drop chunks are candidates: chunk priority must be exactly 0. +/// Explicit user drops (-1), other negative priorities, positive priorities, +/// and the no-priority sentinel (`i8::MIN`) are all left untouched. The chunk +/// must contain at least one span with a non-zero error flag; HTTP status or +/// error metadata alone does not qualify. The root is resolved with +/// `get_root_span_index`: empty chunks are left unchanged, and non-empty +/// chunks always resolve a root, falling back to the last span when no span +/// has `parent_id 0` (e.g. a cyclic chunk). Chunks are never removed here: on +/// a Drop decision the unrescued chunk is forwarded as-is and the backend's +/// ordinary P0 drop handles it. +fn sample_and_stamp( + sampler: &mut ErrorsSampler, + chunk: &mut pb::TraceChunk, + env: &str, + now_unix_secs: i64, +) { + if chunk.priority != 0 { + return; + } + // An error anywhere in the chunk makes it a rescue candidate, not just an + // error on the root span. Checked before the root-span search because + // non-errored chunks are the common case on this path. + if !chunk.spans.iter().any(|span| span.error != 0) { + return; + } + let Ok(root_index) = trace_utils::get_root_span_index(&chunk.spans) else { + // No identifiable root span: leave the chunk unchanged rather than + // guessing an index. + return; + }; + + let Some(root) = chunk.spans.get(root_index) else { + return; + }; + let views: Vec = chunk.spans.iter().map(span_view).collect(); + let trace = TraceView { + env, + trace_id: root.trace_id, + root_index, + // The raw `_sample_rate` wire value is passed through: the shared + // sampler sanitizes non-finite or out-of-range rates to 1.0 itself. + root_global_sample_rate: root.metrics.get("_sample_rate").copied().unwrap_or(1.0), + spans: &views, + }; + let decision = sampler.sample(now_unix_secs, &trace); + + if let SampleDecision::Keep { errors_sr } = decision + && let Some(root) = chunk.spans.get_mut(root_index) + { + root.metrics + .insert(ERRORS_SR_METRIC_KEY.to_string(), errors_sr); + } +} + #[async_trait] impl TraceProcessor for ServerlessTraceProcessor { async fn process_traces( @@ -282,6 +478,14 @@ impl TraceProcessor for ServerlessTraceProcessor { Self::send_to_concentrator(concentrator, &payload); } + // Error rescue runs after stats submission so the concentrator observes every + // submitted chunk exactly as the tracer sent it, and before payload splitting so + // the newly inserted metric is included in the recomputed outbound size. It is + // gated on agent stats computation: without it, P0 chunks are not expected here. + if config.agent_stats_computation_enabled { + self.apply_error_rescue(&mut payload, &config); + } + let pieces: Vec<(TracerPayloadCollection, usize)> = match payload { TracerPayloadCollection::V07(payloads) => { let split_budget = @@ -360,6 +564,9 @@ mod tests { encoded_size, split_oversized_payloads, }, }; + use datadog_agent_trace_sampler::{ + ErrorSamplerConfig, ErrorSamplerMode, ErrorsSampler, SampleDecision, SpanView, TraceView, + }; use libdd_common::{Endpoint, http_common}; use libdd_trace_protobuf::pb; use libdd_trace_utils::test_utils::{create_test_gcp_json_span, create_test_gcp_span}; @@ -424,9 +631,14 @@ mod tests { additional_metric_tags: vec![], additional_metric_tags_cardinality_limit: None, agent_stats_computation_enabled: false, + error_sampler: ErrorSamplerConfig::default(), } } + fn default_error_sampler() -> Arc> { + sampler_for(&ErrorSamplerConfig::default()) + } + fn create_test_metadata() -> MiniAgentMetadata { MiniAgentMetadata { azure_spring_app_hostname: Default::default(), @@ -542,7 +754,8 @@ mod tests { .body(http_common::Body::from(bytes)) .unwrap(); - let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + let trace_processor = + trace_processor::ServerlessTraceProcessor::new(None, default_error_sampler()); let res = trace_processor .process_traces( Arc::new(create_test_config()), @@ -615,7 +828,8 @@ mod tests { .body(http_common::Body::from(bytes)) .unwrap(); - let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + let trace_processor = + trace_processor::ServerlessTraceProcessor::new(None, default_error_sampler()); let res = trace_processor .process_traces( Arc::new(create_test_config()), @@ -701,7 +915,8 @@ mod tests { .body(http_common::Body::from(bytes)) .unwrap(); - let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + let trace_processor = + trace_processor::ServerlessTraceProcessor::new(None, default_error_sampler()); let res = trace_processor .process_traces( Arc::new(create_test_config()), @@ -759,7 +974,8 @@ mod tests { Receiver, ) = mpsc::channel(1); - let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + let trace_processor = + trace_processor::ServerlessTraceProcessor::new(None, default_error_sampler()); let config = Arc::new(Config { enqueue_permit_timeout_secs: 0, ..create_test_config() @@ -827,7 +1043,8 @@ mod tests { Receiver, ) = mpsc::channel(MAX_IN_FLIGHT_ENQUEUES + 1); - let trace_processor = trace_processor::ServerlessTraceProcessor::new(None); + let trace_processor = + trace_processor::ServerlessTraceProcessor::new(None, default_error_sampler()); // Uses the default (non-zero) enqueue_permit_timeout_secs, so the 11th request's permit // acquire has real time to succeed once an earlier request's send completes and // releases its permit, rather than racing a 0-second timeout against task scheduling. @@ -868,4 +1085,777 @@ mod tests { allowing all of them to succeed rather than shedding load" ); } + + // ---- Error rescue tests ---- + + const RESCUE_NOW: i64 = 1_700_000_000; + + fn test_span(trace_id: u64, span_id: u64, parent_id: u64, error: i32) -> pb::Span { + pb::Span { + service: "test-service".to_string(), + name: "test-operation".to_string(), + resource: "GET /test".to_string(), + trace_id, + span_id, + parent_id, + error, + ..Default::default() + } + } + + fn test_chunk(spans: Vec, priority: i32) -> pb::TraceChunk { + pb::TraceChunk { + spans, + priority, + ..Default::default() + } + } + + fn errored_root_chunk(trace_id: u64, priority: i32) -> pb::TraceChunk { + test_chunk(vec![test_span(trace_id, trace_id + 1, 0, 1)], priority) + } + + fn rescue_config(error_sampler: ErrorSamplerConfig) -> Config { + Config { + agent_stats_computation_enabled: true, + env: "agent-env".to_string(), + error_sampler, + ..create_test_config() + } + } + + fn always_keep_sampler_config() -> ErrorSamplerConfig { + ErrorSamplerConfig { + mode: ErrorSamplerMode::AlwaysKeep, + target_tps: 10.0, + extra_sample_rate: 1.0, + } + } + + fn rate_limited_sampler_config(target_tps: f64) -> ErrorSamplerConfig { + ErrorSamplerConfig { + mode: ErrorSamplerMode::RateLimited, + target_tps, + extra_sample_rate: 1.0, + } + } + + fn sampler_for(config: &ErrorSamplerConfig) -> Arc> { + Arc::new(std::sync::Mutex::new(ErrorsSampler::new(*config))) + } + + fn run_rescue( + processor: &trace_processor::ServerlessTraceProcessor, + config: &Config, + chunks: Vec, + now_unix_secs: i64, + ) -> Vec { + let mut payload = TracerPayloadCollection::V07(vec![pb::TracerPayload { + chunks, + ..Default::default() + }]); + processor.apply_error_rescue_at(&mut payload, config, now_unix_secs); + match payload { + TracerPayloadCollection::V07(mut payloads) => payloads.remove(0).chunks, + _ => unreachable!(), + } + } + + fn root(chunk: &pb::TraceChunk) -> &pb::Span { + // All fixtures in this section place the root first or resolve it the + // same way the processor does. + &chunk.spans[0] + } + + fn errors_sr(span: &pb::Span) -> Option { + span.metrics.get("_dd.errors_sr").copied() + } + + /// Builds the equivalent standalone `TraceView` for a fixture chunk, used to + /// compare adapter decisions with direct shared-sampler usage. + fn equivalent_trace_view<'a>( + chunk: &pb::TraceChunk, + env: &'a str, + views: &'a [SpanView<'a>], + ) -> TraceView<'a> { + let root_index = trace_utils::get_root_span_index(&chunk.spans).unwrap(); + let root = &chunk.spans[root_index]; + TraceView { + env, + trace_id: root.trace_id, + root_index, + root_global_sample_rate: root.metrics.get("_sample_rate").copied().unwrap_or(1.0), + spans: views, + } + } + + #[test] + fn test_rescue_stamps_errored_p0_root_and_keeps_priority() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let chunks = run_rescue( + &processor, + &config, + vec![errored_root_chunk(0xdead_beef, 0)], + RESCUE_NOW, + ); + + assert_eq!(chunks.len(), 1, "rescued chunk must still be forwarded"); + assert_eq!(chunks[0].priority, 0, "priority must not be promoted"); + assert_eq!(errors_sr(root(&chunks[0])), Some(1.0)); + } + + #[test] + fn test_rescue_stamps_actual_root_when_error_is_on_child() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + // Healthy root first, errored child second: only the root may be stamped. + let spans = vec![ + test_span(0xfeed, 0x101, 0, 0), + test_span(0xfeed, 0x102, 0x101, 1), + ]; + let chunks = run_rescue(&processor, &config, vec![test_chunk(spans, 0)], RESCUE_NOW); + + assert_eq!(chunks.len(), 1); + assert_eq!(errors_sr(&chunks[0].spans[0]), Some(1.0), "root stamped"); + assert_eq!( + errors_sr(&chunks[0].spans[1]), + None, + "errored child must not be stamped" + ); + } + + #[test] + fn test_rescue_uses_resolved_root_not_first_span() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + // Root (parent_id 0) appears last; the first span is a child. + let spans = vec![ + test_span(0xbeef, 0x201, 0x203, 0), + test_span(0xbeef, 0x202, 0x203, 1), + test_span(0xbeef, 0x203, 0, 0), + ]; + let chunks = run_rescue(&processor, &config, vec![test_chunk(spans, 0)], RESCUE_NOW); + + assert_eq!(chunks.len(), 1); + assert_eq!( + errors_sr(&chunks[0].spans[2]), + Some(1.0), + "the resolved root (last span) must be stamped" + ); + assert_eq!(errors_sr(&chunks[0].spans[0]), None); + assert_eq!(errors_sr(&chunks[0].spans[1]), None); + } + + #[test] + fn test_no_error_flags_leaves_chunk_unchanged_and_spends_no_budget() { + // With RateLimited, counting is global across signatures via + // `all_sigs_seen`, so if the healthy chunks below were wrongly fed to + // the sampler, the final errored chunk's rate would drop below 1.0 and + // it would not be stamped with 1.0. Correct behavior: only errored + // chunks reach the sampler, so the errored chunk is the first count and + // is kept at rate 1.0. + let config = rescue_config(rate_limited_sampler_config(1.0)); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let mut chunks = Vec::new(); + for i in 0..9_u64 { + chunks.push(errored_root_chunk(0xa000 + i, 0)); + // Strip the error flag: healthy P0 chunk with a distinct signature. + chunks.last_mut().unwrap().spans[0].error = 0; + } + chunks.push(errored_root_chunk(0xb000, 0)); + + let chunks = run_rescue(&processor, &config, chunks, RESCUE_NOW); + + for (i, chunk) in chunks.iter().enumerate() { + let is_last = i == chunks.len() - 1; + assert_eq!( + errors_sr(root(chunk)), + if is_last { Some(1.0) } else { None }, + "chunk {i}: only the errored chunk may be rescued" + ); + } + } + + #[test] + fn test_http_500_metadata_alone_is_not_an_error() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let mut span = test_span(0xc0de, 0xc0de + 1, 0, 0); + span.meta + .insert("http.status_code".to_string(), "500".to_string()); + span.meta + .insert("error.type".to_string(), "Error".to_string()); + + let chunks = run_rescue( + &processor, + &config, + vec![test_chunk(vec![span], 0)], + RESCUE_NOW, + ); + + assert_eq!(chunks.len(), 1); + assert_eq!(errors_sr(root(&chunks[0])), None); + } + + #[test] + fn test_non_automatic_drop_priorities_are_never_rescued() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + // -1 (explicit user drop), other negatives, positive priorities, and + // the no-priority sentinel (i8::MIN) are all out of scope. + let priorities = [-1_i32, -5, 1, 2, i8::MIN as i32]; + let chunks = run_rescue( + &processor, + &config, + priorities + .iter() + .map(|p| errored_root_chunk(0x1000_u64 + (*p).unsigned_abs() as u64, *p)) + .collect(), + RESCUE_NOW, + ); + + for chunk in &chunks { + assert_eq!( + errors_sr(root(chunk)), + None, + "priority {} must not be rescued", + chunk.priority + ); + } + } + + #[test] + fn test_empty_chunk_does_not_panic_and_stays_unchanged() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + // An empty chunk cannot be scored: root resolution fails and it must be + // left unchanged without panicking. For non-empty chunks, + // `get_root_span_index` always resolves (falling back to the last span + // when no span has parent_id 0), so the cyclic chunk below exercises + // the fallback path and gets stamped on the resolved fallback root. + let cyclic = test_chunk( + vec![ + test_span(0xd00d, 0xd01, 0xd02, 1), + test_span(0xd00d, 0xd02, 0xd01, 0), + ], + 0, + ); + + let chunks = run_rescue( + &processor, + &config, + vec![test_chunk(vec![], 0), cyclic], + RESCUE_NOW, + ); + + assert_eq!(chunks.len(), 2); + assert!(chunks[0].spans.is_empty(), "empty chunk unchanged"); + assert_eq!( + errors_sr(&chunks[1].spans[1]), + Some(1.0), + "fallback root (last span) stamped" + ); + assert_eq!(errors_sr(&chunks[1].spans[0]), None); + } + + #[tokio::test] + async fn test_rescue_is_noop_when_agent_stats_disabled() { + let config = rescue_config(always_keep_sampler_config()); + let config = Config { + agent_stats_computation_enabled: false, + ..config + }; + let (tx, mut rx): ( + Sender, + tokio::sync::mpsc::Receiver, + ) = mpsc::channel(1); + + let start = get_current_timestamp_nanos(); + let mut json_span = create_test_json_span(11, 222, 333, start, true); + // Root span with an error and an automatic-drop priority: eligible. + json_span["error"] = serde_json::json!(1); + json_span["metrics"]["_sampling_priority_v1"] = serde_json::json!(0.0); + let bytes = rmp_serde::to_vec(&vec![vec![json_span]]).unwrap(); + let request = Request::builder() + .header("datadog-meta-tracer-version", "4.0.0") + .header("datadog-meta-lang", "nodejs") + .header("datadog-meta-lang-version", "v19.7.0") + .header("datadog-meta-lang-interpreter", "v8") + .header("content-length", "100") + .body(http_common::Body::from(bytes)) + .unwrap(); + + let trace_processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + let res = trace_processor + .process_traces( + Arc::new(config), + request, + tx, + Arc::new(create_test_metadata()), + ) + .await; + assert!(res.is_ok()); + + let send_data = rx.recv().await.expect("payload forwarded"); + let payloads = send_data.get_payloads(); + let TracerPayloadCollection::V07(tracer_payloads) = payloads else { + panic!("expected V07 payload"); + }; + let chunk = &tracer_payloads[0].chunks[0]; + assert_eq!(chunk.priority, 0); + for span in &chunk.spans { + assert_eq!( + errors_sr(span), + None, + "rescue must be a no-op when agent stats computation is disabled" + ); + } + } + + #[test] + fn test_disabled_tps_disables_rescue_in_both_modes() { + for mode in [ErrorSamplerMode::RateLimited, ErrorSamplerMode::AlwaysKeep] { + for tps in [0.0_f64, -3.0] { + let sampler_config = ErrorSamplerConfig { + mode, + target_tps: tps, + extra_sample_rate: 1.0, + }; + let config = rescue_config(sampler_config); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let chunks = run_rescue( + &processor, + &config, + vec![errored_root_chunk(0xe000, 0)], + RESCUE_NOW, + ); + + assert_eq!(chunks.len(), 1, "chunk still forwarded"); + assert_eq!( + errors_sr(root(&chunks[0])), + None, + "mode {mode:?} tps {tps}: disabled sampler must not rescue" + ); + } + } + } + + #[test] + fn test_always_keep_rescues_every_eligible_chunk() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let chunks = run_rescue( + &processor, + &config, + (0..10_u64) + .map(|i| errored_root_chunk(0xf000 + i, 0)) + .collect(), + RESCUE_NOW, + ); + + assert_eq!(chunks.len(), 10); + for chunk in &chunks { + assert_eq!(errors_sr(root(chunk)), Some(1.0)); + assert_eq!(chunk.priority, 0); + } + } + + #[test] + fn test_rate_limited_under_sustained_load_keeps_and_rejects() { + let sampler_config = rate_limited_sampler_config(10.0); + let config = rescue_config(sampler_config); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + // 100 chunks with the same signature in one 5-second bucket: the + // default rate is 10 / (100 / 5) = 0.5, so a deterministic mix of keeps + // and rejections is expected. No exact keep count is asserted: the + // sampler is an adaptive rolling-window sampler, not a hard token + // bucket. + let ids: Vec = (0..100_u64).map(|i| 0x11_0000 + i).collect(); + let input: Vec = ids.iter().map(|id| errored_root_chunk(*id, 0)).collect(); + let chunks = run_rescue(&processor, &config, input, RESCUE_NOW); + + let keeps = chunks + .iter() + .filter(|c| errors_sr(root(c)).is_some()) + .count(); + let drops = chunks.len() - keeps; + assert!(keeps > 0, "expected some keeps, got none"); + assert!(drops > 0, "expected some drops, got none"); + assert_eq!(chunks.len(), 100, "every chunk is still forwarded"); + for chunk in &chunks { + assert_eq!(chunk.priority, 0, "priority unchanged on rescue decision"); + } + + // The same input through a standalone shared sampler must produce the + // identical decision sequence, validating the adapter wiring (env, + // sample rate, views) against direct crate usage. + let mut standalone = ErrorsSampler::new(sampler_config); + for (chunk, id) in chunks.iter().zip(&ids) { + let views: Vec = chunk.spans.iter().map(super::span_view).collect(); + let trace = equivalent_trace_view(chunk, &config.env, &views); + let expected = standalone.sample(RESCUE_NOW, &trace); + let actual = match errors_sr(root(chunk)) { + Some(errors_sr) => SampleDecision::Keep { errors_sr }, + None => SampleDecision::Drop, + }; + assert_eq!(actual, expected, "decision mismatch for trace id {id}"); + } + } + + #[test] + fn test_rate_limited_bucket_transitions_and_steady_state() { + let sampler_config = rate_limited_sampler_config(10.0); + let config = rescue_config(sampler_config); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + // First bucket: 100 distinct signatures push the default rate to 0.5. + let ids: Vec = (0..100_u64).map(|i| 0x22_0000 + i).collect(); + let first: Vec = ids.iter().map(|id| errored_root_chunk(*id, 0)).collect(); + let first_chunks = run_rescue(&processor, &config, first, RESCUE_NOW); + let first_keeps = first_chunks + .iter() + .filter(|c| errors_sr(root(c)).is_some()) + .count(); + assert!( + first_keeps > 0 && first_keeps < 100, + "mixed decisions in the first bucket" + ); + + // A full window later (the rolling window is 6 buckets of 5 seconds), + // the same IDs in a fresh bucket still produce mixed decisions, and + // every chunk is still forwarded. + let second: Vec = ids.iter().map(|id| errored_root_chunk(*id, 0)).collect(); + let second_chunks = run_rescue(&processor, &config, second, RESCUE_NOW + 40); + assert_eq!(second_chunks.len(), 100); + let second_keeps = second_chunks + .iter() + .filter(|c| errors_sr(root(c)).is_some()) + .count(); + assert!( + second_keeps > 0 && second_keeps < 100, + "mixed decisions after window rotation" + ); + } + + #[test] + fn test_processor_clones_share_one_budget() { + let sampler_config = rate_limited_sampler_config(1.0); + let config = rescue_config(sampler_config); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + let processor_clone = processor.clone(); + + // Ten errored chunks with the same signature but distinct IDs: five + // through the original, five through the clone. All must draw from the + // same budget, matching a standalone sampler fed the same sequence. + let ids: Vec = (0..10_u64).map(|i| 0x33_0000 + i).collect(); + let mut observed = Vec::new(); + for (i, id) in ids.iter().enumerate() { + let target = if i < 5 { &processor } else { &processor_clone }; + let chunks = run_rescue( + target, + &config, + vec![errored_root_chunk(*id, 0)], + RESCUE_NOW, + ); + observed.push(errors_sr(root(&chunks[0])).is_some()); + } + + let mut standalone = ErrorsSampler::new(sampler_config); + let spans = [SpanView { + service: "test-service", + name: "test-operation", + resource: "GET /test", + error: true, + http_status_code: None, + error_type: None, + }]; + for (i, id) in ids.iter().enumerate() { + let trace = TraceView { + env: &config.env, + trace_id: *id, + root_index: 0, + root_global_sample_rate: 1.0, + spans: &spans, + }; + let expected = matches!( + standalone.sample(RESCUE_NOW, &trace), + SampleDecision::Keep { .. } + ); + assert_eq!( + observed[i], expected, + "clone {i} decision diverged from the shared-budget standalone sampler" + ); + } + assert!( + observed.iter().any(|kept| *kept), + "expected at least one keep" + ); + assert!( + !observed.iter().all(|kept| *kept), + "expected at least one drop" + ); + } + + #[test] + fn test_sampler_timestamp_clamp_never_moves_backwards() { + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&rate_limited_sampler_config(1.0)), + ); + // A clone shares the clamp floor with the original. + let clone = processor.clone(); + + assert_eq!(processor.clamp_sampler_timestamp(100), 100); + // Backward timestamps, through either instance, are clamped to the + // last value seen; forward ones update the floor. + assert_eq!(clone.clamp_sampler_timestamp(50), 100); + assert_eq!(processor.clamp_sampler_timestamp(100), 100); + assert_eq!(clone.clamp_sampler_timestamp(200), 200); + assert_eq!(processor.clamp_sampler_timestamp(199), 200); + } + + #[test] + fn test_rescue_uses_payload_env_with_config_fallback() { + assert_eq!( + super::resolve_payload_env("tracer-env", "agent-env"), + "tracer-env", + "nonempty payload env wins" + ); + assert_eq!( + super::resolve_payload_env("", "agent-env"), + "agent-env", + "empty payload env falls back to the agent config env" + ); + assert_eq!( + super::resolve_payload_env("", ""), + "", + "both empty stays empty" + ); + } + + #[test] + fn test_rescue_preserves_chunk_metadata_and_span_order() { + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let mut root_span = test_span(0x501d, 0x502, 0, 0); + root_span + .metrics + .insert("_sampling_priority_v1".to_string(), 0.0); + root_span.metrics.insert("_sample_rate".to_string(), 0.25); + root_span + .metrics + .insert("_dd.span_sampling.rule".to_string(), 1.0); + let mut child = test_span(0x501d, 0x503, 0x502, 1); + child.meta.insert("keep".to_string(), "me".to_string()); + + let mut chunk = test_chunk(vec![root_span, child], 0); + chunk.tags.insert("_dd.p.dm".to_string(), "-4".to_string()); + chunk + .tags + .insert("origin".to_string(), "synthetics".to_string()); + chunk.origin = "synthetics".to_string(); + + let chunks = run_rescue(&processor, &config, vec![chunk], RESCUE_NOW); + + assert_eq!(chunks.len(), 1); + let rescued = &chunks[0]; + assert_eq!(rescued.priority, 0, "priority preserved"); + assert_eq!(rescued.origin, "synthetics", "origin preserved"); + assert_eq!( + rescued.tags.get("_dd.p.dm").map(String::as_str), + Some("-4"), + "decision maker preserved" + ); + assert_eq!( + rescued.tags.get("origin").map(String::as_str), + Some("synthetics") + ); + assert_eq!( + rescued.spans[0] + .metrics + .get("_dd.span_sampling.rule") + .copied(), + Some(1.0), + "single-span sampling metrics preserved" + ); + assert_eq!( + rescued.spans[0] + .metrics + .get("_sampling_priority_v1") + .copied(), + Some(0.0), + "sampling priority metric preserved" + ); + assert_eq!( + errors_sr(&rescued.spans[0]), + Some(1.0), + "root stamped despite existing metrics" + ); + assert_eq!( + rescued.spans[1].meta.get("keep").map(String::as_str), + Some("me"), + "span metadata and order preserved" + ); + assert_eq!(rescued.spans.len(), 2, "no spans added or removed"); + } + + #[test] + fn test_probabilistic_decision_maker_chunk_gets_no_workaround() { + // A priority-0 chunk carrying the chunk-level probabilistic decision + // maker (`_dd.p.dm = "-9"`) is rescued like any other automatic-drop + // chunk: priority and decision maker are preserved untouched. The + // backend resolves such chunks to the `probabilistic` ingestion reason + // before checking `_dd.errors_sr`, so it may still drop them; SCL does + // not work around that limitation. + let config = rescue_config(always_keep_sampler_config()); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let mut chunk = errored_root_chunk(0x600d, 0); + chunk.tags.insert("_dd.p.dm".to_string(), "-9".to_string()); + + let chunks = run_rescue(&processor, &config, vec![chunk], RESCUE_NOW); + + assert_eq!(chunks.len(), 1); + assert_eq!(chunks[0].priority, 0, "priority preserved, no promotion"); + assert_eq!( + chunks[0].tags.get("_dd.p.dm").map(String::as_str), + Some("-9"), + "probabilistic decision maker preserved, no rewrite" + ); + assert_eq!(errors_sr(root(&chunks[0])), Some(1.0)); + } + + #[test] + fn test_rescue_passes_raw_sample_rate_to_sampler() { + // The raw `_sample_rate` wire value is passed to the shared sampler, + // which sanitizes it. A value outside (0, 1] falls back to 1.0, so a + // bogus rate does not change the stamped rescue rate. + let sampler_config = rate_limited_sampler_config(10.0); + let config = rescue_config(sampler_config); + let processor = trace_processor::ServerlessTraceProcessor::new( + None, + sampler_for(&config.error_sampler), + ); + + let mut chunk = errored_root_chunk(0x77_00, 0); + chunk.spans[0] + .metrics + .insert("_sample_rate".to_string(), f64::NAN); + + let chunks = run_rescue(&processor, &config, vec![chunk], RESCUE_NOW); + + // Single signature, well under budget: sanitized rate 1.0 keeps it. + assert_eq!(errors_sr(root(&chunks[0])), Some(1.0)); + } + + #[tokio::test] + async fn test_stats_concentrator_observes_pre_rescue_chunks() { + let config = rescue_config(always_keep_sampler_config()); + let (stats_tx, mut stats_rx) = tokio::sync::mpsc::unbounded_channel(); + let concentrator = super::StatsConcentratorHandle::new(stats_tx); + let processor = trace_processor::ServerlessTraceProcessor::new( + Some(concentrator), + sampler_for(&config.error_sampler), + ); + + let start = get_current_timestamp_nanos(); + let mut json_span = create_test_json_span(11, 222, 333, start, true); + json_span["error"] = serde_json::json!(1); + json_span["metrics"]["_sampling_priority_v1"] = serde_json::json!(0.0); + let bytes = rmp_serde::to_vec(&vec![vec![json_span]]).unwrap(); + let request = Request::builder() + .header("datadog-meta-tracer-version", "4.0.0") + .header("datadog-meta-lang", "nodejs") + .header("datadog-meta-lang-version", "v19.7.0") + .header("datadog-meta-lang-interpreter", "v8") + .header("content-length", "100") + .body(http_common::Body::from(bytes)) + .unwrap(); + + let res = processor + .process_traces( + Arc::new(config), + request, + mpsc::channel(1).0, + Arc::new(create_test_metadata()), + ) + .await; + assert!(res.is_ok()); + + // The concentrator must receive the chunk exactly as the tracer sent + // it: before rescue stamping, with no `_dd.errors_sr` anywhere. + let (chunk, _metadata) = match stats_rx.try_recv() { + Ok(crate::stats_concentrator_service::ConcentratorCommand::AddChunk( + chunk, + metadata, + )) => (*chunk, metadata), + Ok(_) => panic!("expected an AddChunk command"), + Err(err) => panic!("expected an AddChunk command, got {err}"), + }; + assert_eq!(chunk.priority, 0); + for span in &chunk.spans { + assert_eq!( + errors_sr(span), + None, + "stats must observe the pre-rescue chunk" + ); + } + } } diff --git a/crates/datadog-trace-agent/tests/integration_test.rs b/crates/datadog-trace-agent/tests/integration_test.rs index cb0117b..800cde8 100644 --- a/crates/datadog-trace-agent/tests/integration_test.rs +++ b/crates/datadog-trace-agent/tests/integration_test.rs @@ -21,10 +21,11 @@ use datadog_trace_agent::{ }; use http_body_util::BodyExt; use hyper::StatusCode; -use serde_json::Value; +use libdd_trace_utils::test_utils::create_test_json_span; +use serde_json::{Value, json}; use serial_test::serial; use std::sync::Arc; -use std::time::Duration; +use std::time::{Duration, UNIX_EPOCH}; #[cfg(all(windows, feature = "windows-pipes"))] use common::helpers::send_named_pipe_request; @@ -97,9 +98,10 @@ pub fn create_mini_agent_with_real_flushers( let aggregator = Arc::new(tokio::sync::Mutex::new(TraceAggregator::default())); let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor::new(Some( - stats_concentrator_handle.clone(), - ))), + trace_processor: Arc::new(ServerlessTraceProcessor::new( + Some(stats_concentrator_handle.clone()), + ServerlessTraceProcessor::new_error_sampler(&config), + )), trace_flusher: Arc::new(ServerlessTraceFlusher::new( aggregator.clone(), config.clone(), @@ -249,7 +251,10 @@ async fn test_mini_agent_tcp_handles_requests() { let test_port = config.dd_apm_receiver_port; let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor::new(None)), + trace_processor: Arc::new(ServerlessTraceProcessor::new( + None, + ServerlessTraceProcessor::new_error_sampler(&config), + )), trace_flusher: Arc::new(MockTraceFlusher), stats_processor: Arc::new(MockStatsProcessor), stats_flusher: Arc::new(MockStatsFlusher), @@ -363,7 +368,10 @@ async fn test_mini_agent_named_pipe_handles_requests() { let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor::new(None)), + trace_processor: Arc::new(ServerlessTraceProcessor::new( + None, + ServerlessTraceProcessor::new_error_sampler(&config), + )), trace_flusher: Arc::new(MockTraceFlusher), stats_processor: Arc::new(MockStatsProcessor), stats_flusher: Arc::new(MockStatsFlusher), @@ -535,7 +543,10 @@ async fn test_mini_agent_tcp_proxies_dsm_requests() { let mini_agent = MiniAgent { config: config.clone(), - trace_processor: Arc::new(ServerlessTraceProcessor::new(None)), + trace_processor: Arc::new(ServerlessTraceProcessor::new( + None, + ServerlessTraceProcessor::new_error_sampler(&config), + )), trace_flusher: Arc::new(MockTraceFlusher), stats_processor: Arc::new(MockStatsProcessor), stats_flusher: Arc::new(MockStatsFlusher), @@ -1247,3 +1258,201 @@ async fn test_mini_agent_dual_transport_with_real_flushers() { let _ = agent_handle.await; verify_stats_request(&mock_server).await; } + +/// Builds a mixed four-trace payload for the error rescue wire test: an +/// eligible errored automatic-drop chunk, a healthy automatic-drop chunk, an +/// explicit user drop, and a positive-priority chunk. +fn create_error_rescue_test_payload() -> Vec { + let start = UNIX_EPOCH.elapsed().unwrap().as_nanos() as i64; + let span = |trace_id: u64, name: &str, error: i64, priority: f64| { + let mut span = create_test_json_span(trace_id, trace_id + 1, 0, start, false); + span["name"] = json!(name); + span["error"] = json!(error); + span["metrics"]["_sampling_priority_v1"] = json!(priority); + span + }; + let traces = vec![ + vec![span(700, "rescued_error_p0", 1, 0.0)], + vec![span(701, "healthy_p0", 0, 0.0)], + vec![span(702, "user_drop_p0", 1, -1.0)], + vec![span(703, "priority_1", 1, 1.0)], + ]; + rmp_serde::to_vec(&traces).expect("Failed to serialize error rescue test trace") +} + +/// End-to-end error rescue verification against the outbound wire payload and +/// agent-computed stats: an eligible errored P0 chunk is rescued (positive +/// `_dd.errors_sr` on the root, priority still 0) while all other chunks are +/// forwarded unchanged, and every submitted span still contributes to stats. +/// This proves the wire contract only; backend retention of rescued chunks is +/// established by the backend implementation, not by this fake intake. +#[cfg(test)] +#[tokio::test] +#[serial] +async fn test_error_rescue_outbound_payload_and_stats() { + use libdd_trace_protobuf::pb::AgentPayload; + use prost::Message as _; + + let mock_server: MockServer = MockServer::start().await; + tokio::time::sleep(Duration::from_millis(50)).await; + + let mut config = create_tcp_test_config(8137); // use different port to avoid race condition with other tests + configure_mock_endpoints(&mut config, &mock_server.url()); + config.agent_stats_computation_enabled = true; + // Deterministic rescue settings: every eligible chunk is kept at 1.0. + config.error_sampler = datadog_agent_trace_sampler::ErrorSamplerConfig { + mode: datadog_agent_trace_sampler::ErrorSamplerMode::AlwaysKeep, + target_tps: 1.0, + extra_sample_rate: 1.0, + }; + let config = Arc::new(config); + let test_port = config.dd_apm_receiver_port; + + let (mini_agent, stats_concentrator_service_handle) = + create_mini_agent_with_real_flushers(config); + + let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(false); + let agent_handle = tokio::spawn(async move { + let _ = mini_agent + .start_mini_agent(shutdown_rx, Some(stats_concentrator_service_handle)) + .await; + }); + + let mut server_ready = false; + for _ in 0..20 { + tokio::time::sleep(Duration::from_millis(50)).await; + if let Ok(response) = send_tcp_request(test_port, "/info", "GET", None, &[]).await + && response.status().is_success() + { + server_ready = true; + break; + } + } + assert!( + server_ready, + "Mini agent server failed to start within timeout" + ); + + let trace_response = send_tcp_request( + test_port, + "/v0.4/traces", + "POST", + Some(create_error_rescue_test_payload()), + &[], + ) + .await + .expect("Failed to send /v0.4/traces request"); + assert_eq!(trace_response.status(), StatusCode::OK); + + verify_trace_request(&mock_server).await; + + // Decode the actual outbound protobuf payload(s) and locate each submitted + // chunk by its root span name. + let trace_reqs = mock_server.get_requests_for_path("/api/v0.2/traces"); + let mut chunks_by_root_name: std::collections::HashMap< + String, + libdd_trace_protobuf::pb::TraceChunk, + > = std::collections::HashMap::new(); + for req in &trace_reqs { + let agent_payload = AgentPayload::decode(&req.body[..]) + .expect("Failed to decode outbound AgentPayload protobuf"); + for tracer_payload in agent_payload.tracer_payloads { + for chunk in tracer_payload.chunks { + let root_name = chunk + .spans + .iter() + .find(|s| s.parent_id == 0) + .map(|s| s.name.clone()) + .expect("chunk has a root span"); + chunks_by_root_name.insert(root_name, chunk); + } + } + } + + let names = [ + "rescued_error_p0", + "healthy_p0", + "user_drop_p0", + "priority_1", + ]; + for name in names { + assert!( + chunks_by_root_name.contains_key(name), + "expected chunk {name} to be forwarded to the backend, got: {:?}", + chunks_by_root_name.keys().collect::>() + ); + } + + let errors_sr_of = |name: &str| -> Option { + chunks_by_root_name[name] + .spans + .iter() + .find(|s| s.parent_id == 0) + .and_then(|root| root.metrics.get("_dd.errors_sr").copied()) + }; + + // The eligible errored automatic-drop chunk is rescued: positive + // `_dd.errors_sr` on the root, chunk priority still 0. + assert_eq!( + errors_sr_of("rescued_error_p0"), + Some(1.0), + "rescued root must carry a positive _dd.errors_sr" + ); + assert_eq!( + chunks_by_root_name["rescued_error_p0"].priority, 0, + "rescued chunk priority must remain 0, no promotion" + ); + + // Non-candidate chunks receive no rescue metric and keep their priority. + assert_eq!(errors_sr_of("healthy_p0"), None, "healthy P0 not rescued"); + assert_eq!(chunks_by_root_name["healthy_p0"].priority, 0); + assert_eq!( + errors_sr_of("user_drop_p0"), + None, + "explicit user drop never rescued" + ); + assert_eq!(chunks_by_root_name["user_drop_p0"].priority, -1); + assert_eq!( + errors_sr_of("priority_1"), + None, + "positive priority not rescued" + ); + assert_eq!(chunks_by_root_name["priority_1"].priority, 1); + + // Wait for the stats flush, then assert agent-computed stats include all + // four submitted spans, independent of rescue decisions. + tokio::time::sleep(FLUSH_WAIT_DURATION).await; + let _ = shutdown_tx.send(true); + let _ = agent_handle.await; + + let stats_reqs = mock_server.get_requests_for_path("/api/v0.2/stats"); + assert!( + !stats_reqs.is_empty(), + "Expected at least one stats request" + ); + + let all_groups: Vec<_> = stats_reqs + .iter() + .map(|req| decode_stats_payload(&req.body)) + .flat_map(|payload| { + payload + .stats + .into_iter() + .flat_map(|csp| csp.stats.into_iter()) + .flat_map(|bucket| bucket.stats.into_iter()) + }) + .collect(); + + for name in names { + let group = all_groups.iter().find(|g| g.name == name); + assert!( + group.is_some(), + "expected span {name} to contribute to agent-computed stats, got: {:?}", + all_groups.iter().map(|g| &g.name).collect::>() + ); + assert!( + group.unwrap().hits > 0, + "expected span {name} to have a positive hit count in stats" + ); + } +}