From baec68ddf8ff3b04bcf382edafe35956e1f932d9 Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Fri, 18 Sep 2026 14:23:19 -0400 Subject: [PATCH 1/9] Rescue errored P0 trace chunks from backend drop Wire the shared datadog-agent-trace-sampler into the SCL trace processor so error chunks with automatic-drop priority (0) get a second look before being forwarded. On a keep, the sampler's rate is stamped as a positive `_dd.errors_sr` metric on the chunk's root span, which the APM backend treats as an `error` ingestion reason and retains without any priority promotion. Chunks are never dropped locally: unrescued chunks are forwarded unchanged for the backend to discard. - Default to rate-limited rescue at 10 error TPS per process, with a shared budget across requests and processor clones - Configure via DD_APM_ERROR_SAMPLER_MODE (rate_limited | always_keep) and DD_APM_ERROR_TPS (<= 0 disables rescue in either mode); invalid values warn and fall back to defaults - Gate rescue on agent stats computation and submit all chunks to the stats concentrator before any rescue stamping --- Cargo.lock | 1 + crates/datadog-serverless-compat/src/main.rs | 1 + crates/datadog-trace-agent/Cargo.toml | 1 + crates/datadog-trace-agent/src/config.rs | 173 ++++ .../src/stats_processor.rs | 1 + .../src/trace_processor.rs | 952 +++++++++++++++++- .../tests/integration_test.rs | 225 ++++- 7 files changed, 1340 insertions(+), 14 deletions(-) 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_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..410a0c0 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -4,10 +4,13 @@ use std::sync::Arc; 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 +30,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 +133,75 @@ 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>, } 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, + } + } + + /// 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. + /// + /// `now_unix_secs` is passed in explicitly so tests can exercise the + /// rolling window without sleeping. + fn apply_error_rescue( + &self, + payload: &mut TracerPayloadCollection, + config: &Config, + now_unix_secs: i64, + ) { + let TracerPayloadCollection::V07(tracer_payloads) = payload else { + return; + }; + 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; + } + 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(&mut sampler, chunk, env, now_unix_secs); + } } } @@ -179,6 +249,86 @@ impl ServerlessTraceProcessor { } } +/// Chooses the env used to key the error sampler's per-signature rates: the +/// tracer payload's own env when nonempty, otherwise the agent's configured +/// env. This matches how stats prefer the payload env and fall back to the +/// agent config env. +fn resolve_payload_env<'a>(tracer_payload_env: &'a str, config_env: &'a str) -> &'a str { + if tracer_payload_env.is_empty() { + config_env + } else { + tracer_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 or rootless chunks are left unchanged. 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 +432,17 @@ 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 { + let now_unix_secs = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |d| d.as_secs() as i64); + self.apply_error_rescue(&mut payload, &config, now_unix_secs); + } + let pieces: Vec<(TracerPayloadCollection, usize)> = match payload { TracerPayloadCollection::V07(payloads) => { let split_budget = @@ -360,6 +521,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 +588,16 @@ 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> { + Arc::new(std::sync::Mutex::new(ErrorsSampler::new( + ErrorSamplerConfig::default(), + ))) + } + fn create_test_metadata() -> MiniAgentMetadata { MiniAgentMetadata { azure_spring_app_hostname: Default::default(), @@ -542,7 +713,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 +787,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 +874,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 +933,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 +1002,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 +1044,768 @@ 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(&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["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(|s| SpanView { + service: &s.service, + name: &s.name, + resource: &s.resource, + error: s.error != 0, + http_status_code: s.meta.get("http.status_code").map(String::as_str), + error_type: s.meta.get("error.type").map(String::as_str), + }) + .collect(); + let trace = equivalent_trace_view(chunk, "test-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: "test-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_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["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" + ); + } +} From ca14e598f01909802672843a72d892d097a47c49 Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Fri, 18 Sep 2026 14:46:44 -0400 Subject: [PATCH 2/9] Exercise error rescue in stats regression tests --- crates/datadog-trace-agent/src/trace_processor.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 410a0c0..2d12fc6 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -1361,6 +1361,7 @@ mod tests { 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() @@ -1768,6 +1769,7 @@ mod tests { 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() From bfbe16481bfd558481fa28939eaf59da64c5540e Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Fri, 18 Sep 2026 15:07:37 -0400 Subject: [PATCH 3/9] Read sampler clock while holding the sampler lock --- .../src/trace_processor.rs | 53 ++++++++++++++----- 1 file changed, 39 insertions(+), 14 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 2d12fc6..c087c9c 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -177,30 +177,58 @@ impl ServerlessTraceProcessor { /// Chunks are never removed or reordered: unrescued chunks stay in the /// payload and the backend discards them. /// - /// `now_unix_secs` is passed in explicitly so tests can exercise the - /// rolling window without sleeping. - fn apply_error_rescue( + 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, which is + // monotonic. Reading the clock before acquiring the lock could hand the + // sampler an older time after a newer one, moving the rolling window + // backwards and clearing counts from the current bucket, which would + // undercount TPS and rescue too many error chunks. + let now_unix_secs = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_or(0, |d| d.as_secs() as i64); + 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 TracerPayloadCollection::V07(tracer_payloads) = payload else { - return; - }; 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; } + self.rescue_with_sampler(payload, config, now_unix_secs, &mut sampler); + } + + 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(&mut sampler, chunk, env, now_unix_secs); + sample_and_stamp(sampler, chunk, env, now_unix_secs); } } } @@ -437,10 +465,7 @@ impl TraceProcessor for ServerlessTraceProcessor { // 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 { - let now_unix_secs = SystemTime::now() - .duration_since(UNIX_EPOCH) - .map_or(0, |d| d.as_secs() as i64); - self.apply_error_rescue(&mut payload, &config, now_unix_secs); + self.apply_error_rescue(&mut payload, &config); } let pieces: Vec<(TracerPayloadCollection, usize)> = match payload { @@ -1113,7 +1138,7 @@ mod tests { chunks, ..Default::default() }]); - processor.apply_error_rescue(&mut payload, config, now_unix_secs); + processor.apply_error_rescue_at(&mut payload, config, now_unix_secs); match payload { TracerPayloadCollection::V07(mut payloads) => payloads.remove(0).chunks, _ => unreachable!(), From 1cf32c8dbe7e3ead1605b1c3f66da5936c99a38b Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Thu, 24 Sep 2026 13:08:31 -0400 Subject: [PATCH 4/9] Clamp sampler timestamps against backward clock steps --- .../src/trace_processor.rs | 44 ++++++++++++++++--- 1 file changed, 39 insertions(+), 5 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index c087c9c..1dab5fb 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -2,6 +2,7 @@ // 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}; @@ -137,6 +138,10 @@ pub struct ServerlessTraceProcessor { /// 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 { @@ -149,6 +154,7 @@ impl 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)), } } @@ -185,14 +191,13 @@ impl ServerlessTraceProcessor { return; } // Read the clock while holding the sampler lock so that concurrent - // requests deliver timestamps in lock-acquisition order, which is - // monotonic. Reading the clock before acquiring the lock could hand the - // sampler an older time after a newer one, moving the rolling window - // backwards and clearing counts from the current bucket, which would - // undercount TPS and rescue too many error chunks. + // 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); } @@ -212,6 +217,17 @@ impl ServerlessTraceProcessor { 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, @@ -1641,6 +1657,24 @@ mod tests { ); } + #[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!( From 37a17c5c23563df33f1af304ce7b134f6c594ce6 Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Thu, 24 Sep 2026 13:12:17 -0400 Subject: [PATCH 5/9] Match test oracle env with processor fallback env --- crates/datadog-trace-agent/src/trace_processor.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index 1dab5fb..c19b3c7 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -1547,7 +1547,7 @@ mod tests { error_type: s.meta.get("error.type").map(String::as_str), }) .collect(); - let trace = equivalent_trace_view(chunk, "test-env", &views); + 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 }, From f99948f7d79bab2b1332c309d104f38d9011c85a Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Thu, 24 Sep 2026 13:30:40 -0400 Subject: [PATCH 6/9] Match clone-budget oracle env with processor fallback env --- crates/datadog-trace-agent/src/trace_processor.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index c19b3c7..a28b5cd 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -1632,7 +1632,7 @@ mod tests { }]; for (i, id) in ids.iter().enumerate() { let trace = TraceView { - env: "test-env", + env: &config.env, trace_id: *id, root_index: 0, root_global_sample_rate: 1.0, From 171a77015493376c07d47c346271be1d3c35fcde Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Thu, 24 Sep 2026 13:35:09 -0400 Subject: [PATCH 7/9] Document root-resolution fallback in sampler doc comment --- crates/datadog-trace-agent/src/trace_processor.rs | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index a28b5cd..fda8dc3 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -326,9 +326,11 @@ fn span_view(span: &pb::Span) -> SpanView<'_> { /// 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 or rootless chunks are left unchanged. 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. +/// `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, From 99686d37dad872534c2fe09e440cbd929203b4c6 Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Thu, 24 Sep 2026 14:18:58 -0400 Subject: [PATCH 8/9] Share the payload-env fallback rule between stats and error rescue Extracts the "prefer tracer payload env, fall back to agent config env" rule into one function reused by stats flushing and error-rescue sampling, and removes duplicated sampler lock/disabled-check and SpanView-construction code in the error sampler and its tests. --- .../src/stats_concentrator_service.rs | 11 +++-- .../src/trace_processor.rs | 47 ++++++------------- 2 files changed, 21 insertions(+), 37 deletions(-) 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/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index fda8dc3..e6882cd 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -184,12 +184,6 @@ impl ServerlessTraceProcessor { /// 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 @@ -198,12 +192,14 @@ impl ServerlessTraceProcessor { .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); + self.apply_error_rescue_at(payload, config, now_unix_secs); } - /// Test-only variant that injects a synthetic timestamp so tests can - /// exercise the rolling window without sleeping. - #[cfg(test)] + /// Applies the rescue pass for an explicit timestamp, taking the sampler + /// lock and skipping all view construction when the sampler is disabled + /// by config (target_tps <= 0). Used by `apply_error_rescue` with the + /// clamped wall clock, and directly by tests with a synthetic timestamp + /// so they can exercise the rolling window without sleeping. fn apply_error_rescue_at( &self, payload: &mut TracerPayloadCollection, @@ -293,15 +289,15 @@ impl ServerlessTraceProcessor { } } -/// Chooses the env used to key the error sampler's per-signature rates: the -/// tracer payload's own env when nonempty, otherwise the agent's configured -/// env. This matches how stats prefer the payload env and fall back to the -/// agent config env. -fn resolve_payload_env<'a>(tracer_payload_env: &'a str, config_env: &'a str) -> &'a str { - if tracer_payload_env.is_empty() { +/// 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 { - tracer_payload_env + payload_env } } @@ -636,9 +632,7 @@ mod tests { } fn default_error_sampler() -> Arc> { - Arc::new(std::sync::Mutex::new(ErrorsSampler::new( - ErrorSamplerConfig::default(), - ))) + sampler_for(&ErrorSamplerConfig::default()) } fn create_test_metadata() -> MiniAgentMetadata { @@ -1537,18 +1531,7 @@ mod tests { // 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(|s| SpanView { - service: &s.service, - name: &s.name, - resource: &s.resource, - error: s.error != 0, - http_status_code: s.meta.get("http.status_code").map(String::as_str), - error_type: s.meta.get("error.type").map(String::as_str), - }) - .collect(); + 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)) { From 69fc9771f62b94cdc1d2cb776b19cbf21f2754f6 Mon Sep 17 00:00:00 2001 From: Lucas Pimentel Date: Thu, 24 Sep 2026 14:40:53 -0400 Subject: [PATCH 9/9] Fix error-rescue clock read to happen under the sampler lock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The clock read and timestamp clamp were happening before the sampler lock was acquired, letting concurrent requests deliver timestamps out of lock-acquisition order and undercount TPS in the rolling window. 🤖 --- .../datadog-trace-agent/src/trace_processor.rs | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/crates/datadog-trace-agent/src/trace_processor.rs b/crates/datadog-trace-agent/src/trace_processor.rs index e6882cd..c150d5e 100644 --- a/crates/datadog-trace-agent/src/trace_processor.rs +++ b/crates/datadog-trace-agent/src/trace_processor.rs @@ -184,6 +184,12 @@ impl ServerlessTraceProcessor { /// 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 @@ -192,14 +198,12 @@ impl ServerlessTraceProcessor { .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.apply_error_rescue_at(payload, config, now_unix_secs); + self.rescue_with_sampler(payload, config, now_unix_secs, &mut sampler); } - /// Applies the rescue pass for an explicit timestamp, taking the sampler - /// lock and skipping all view construction when the sampler is disabled - /// by config (target_tps <= 0). Used by `apply_error_rescue` with the - /// clamped wall clock, and directly by tests with a synthetic timestamp - /// so they can exercise the rolling window without sleeping. + /// 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,