Skip to content
Open
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions crates/datadog-serverless-compat/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
1 change: 1 addition & 0 deletions crates/datadog-trace-agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
173 changes: 173 additions & 0 deletions crates/datadog-trace-agent/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<f64>() {
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";

Expand Down Expand Up @@ -128,6 +174,10 @@ pub struct Config {
pub additional_metric_tags_cardinality_limit: Option<usize>,
/// 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 {
Expand Down Expand Up @@ -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(),
})
}
}
Expand All @@ -315,6 +366,7 @@ mod tests {
use std::collections::HashMap;

use crate::config;
use datadog_agent_trace_sampler::ErrorSamplerMode;

#[test]
#[serial]
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(),
}
}
}
11 changes: 6 additions & 5 deletions crates/datadog-trace-agent/src/stats_concentrator_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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(),
Expand Down
1 change: 1 addition & 0 deletions crates/datadog-trace-agent/src/stats_processor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
}
}

Expand Down
Loading
Loading