agent_data_plane_config_system/
loaded.rs1use std::path::Path;
10
11use agent_data_plane_config::SalukiConfiguration;
12use datadog_agent_config::apply_datadog_env;
13use saluki_config::dynamic::ConfigUpdate;
14use saluki_config::{ConfigurationLoader, GenericConfiguration};
15use serde_json::Value;
16use tokio::sync::mpsc;
17
18use crate::env_provider::EnvironmentProvider;
19use crate::saluki_env_overlay;
20use crate::source::SourceTree;
21use crate::system::{translate_strict, validate, ConfigurationSystem, Error};
22
23const ENV_VAR_PREFIX: &str = "DD";
27
28const COMPAT_FORWARD_CHANNEL_SIZE: usize = 100;
31
32#[derive(Clone, Copy, Debug, Eq, PartialEq)]
37pub enum EnvPrecedence {
38 BeforeFile,
42 AfterFile,
45 Disabled,
47}
48
49pub struct LoadedConfiguration {
54 loader: ConfigurationLoader,
55 base: SourceTree,
57 local: SalukiConfiguration,
60}
61
62impl LoadedConfiguration {
63 pub async fn load(path: impl AsRef<Path>, env: EnvPrecedence) -> Result<Self, Error> {
74 let loader = build_loader(path.as_ref(), env)?;
75 let base = SourceTree::all_explicit(build_base(path.as_ref(), env)?);
78 let local = translate_strict(&base)?;
79 Ok(Self { loader, base, local })
80 }
81
82 pub fn local(&self) -> &SalukiConfiguration {
86 &self.local
87 }
88
89 pub fn raw_config(&self) -> GenericConfiguration {
92 self.loader.bootstrap_generic()
93 }
94
95 pub async fn run(self, config_stream: mpsc::Receiver<ConfigUpdate>) -> Result<ConfigurationSystem, Error> {
105 let (compat_tx, compat_rx) = mpsc::channel(COMPAT_FORWARD_CHANNEL_SIZE);
108 let compat_map = self.loader.with_dynamic_configuration(compat_rx).into_generic().await?;
109
110 ConfigurationSystem::connected(config_stream, compat_tx, compat_map, self.base).await
111 }
112
113 pub async fn standalone(self) -> Result<ConfigurationSystem, Error> {
124 validate(&self.local)?;
125 let compat_map = self.loader.into_generic().await?;
126 Ok(ConfigurationSystem::standalone(compat_map, self.local))
127 }
128}
129
130fn build_base(path: &Path, env: EnvPrecedence) -> Result<Value, Error> {
139 let text = std::fs::read_to_string(path).map_err(|e| Error::Base {
140 message: format!("read `{}`: {e}", path.display()),
141 })?;
142 let mut base: Value = serde_yaml::from_str(&text).map_err(|e| Error::Base {
143 message: format!("parse `{}`: {e}", path.display()),
144 })?;
145 drop_nulls(&mut base);
146 if base.is_null() {
147 base = Value::Object(serde_json::Map::new());
148 }
149
150 let overwrite = match env {
151 EnvPrecedence::Disabled => return Ok(base),
152 EnvPrecedence::AfterFile => true,
153 EnvPrecedence::BeforeFile => false,
154 };
155 apply_datadog_env(&mut base, overwrite).map_err(|message| Error::Base { message })?;
156 saluki_env_overlay::apply_env(&mut base, overwrite).map_err(|message| Error::Base { message })?;
157 Ok(base)
158}
159
160fn drop_nulls(value: &mut Value) {
163 if let Value::Object(map) = value {
164 map.retain(|_, v| !v.is_null());
165 for v in map.values_mut() {
166 drop_nulls(v);
167 }
168 }
169}
170
171fn build_loader(path: &Path, env: EnvPrecedence) -> Result<ConfigurationLoader, Error> {
181 let loader = ConfigurationLoader::default();
182 let loader = match env {
183 EnvPrecedence::AfterFile => loader
184 .from_yaml(path)?
185 .from_environment(ENV_VAR_PREFIX)?
186 .add_providers([schema_env_provider()?]),
187 EnvPrecedence::BeforeFile => loader
188 .from_environment(ENV_VAR_PREFIX)?
189 .add_providers([schema_env_provider()?])
190 .from_yaml(path)?,
191 EnvPrecedence::Disabled => loader.from_yaml(path)?,
192 };
193 Ok(loader)
194}
195
196fn schema_env_provider() -> Result<EnvironmentProvider, Error> {
199 EnvironmentProvider::new().map_err(|message| Error::Base { message })
200}
201
202#[cfg(test)]
203mod tests {
204 use bytesize::ByteSize;
205 use saluki_config::test_env_lock;
206 use serde_json::json;
207
208 use super::*;
209
210 #[test]
211 fn build_base_composes_file_and_environment_by_precedence() {
212 let _guard = test_env_lock();
213 let path = std::env::temp_dir().join(format!("adp_build_base_{}.yaml", std::process::id()));
214 std::fs::write(
215 &path,
216 "dogstatsd_port: 8125\ndogstatsd_non_local_traffic: false\nempty_key:\n",
217 )
218 .unwrap();
219 std::env::set_var("DD_DOGSTATSD_PORT", "9125");
220
221 let base = build_base(&path, EnvPrecedence::AfterFile).expect("base builds");
224 assert_eq!(base.get("dogstatsd_port"), Some(&json!(9125)));
225 assert_eq!(base.get("dogstatsd_non_local_traffic"), Some(&json!(false)));
226 assert!(base.get("empty_key").is_none());
227
228 let base = build_base(&path, EnvPrecedence::BeforeFile).expect("base builds");
230 assert_eq!(base.get("dogstatsd_port"), Some(&json!(8125)));
231
232 let base = build_base(&path, EnvPrecedence::Disabled).expect("base builds");
234 assert_eq!(base.get("dogstatsd_port"), Some(&json!(8125)));
235
236 std::env::remove_var("DD_DOGSTATSD_PORT");
237 std::fs::remove_file(&path).ok();
238 }
239
240 #[tokio::test]
241 async fn local_exposes_translated_configuration() {
242 let path = std::env::temp_dir().join(format!("adp_local_{}.yaml", std::process::id()));
244 std::fs::write(&path, "log_level: warn\ndogstatsd_port: 9125\n").unwrap();
245
246 let loaded = LoadedConfiguration::load(&path, EnvPrecedence::Disabled)
247 .await
248 .expect("local sources load");
249 let config = loaded.local();
250
251 assert_eq!(config.control.logging.level, "warn");
252 assert_eq!(config.domains.dogstatsd.listeners.port, 9125);
253
254 std::fs::remove_file(&path).ok();
255 }
256
257 #[tokio::test]
258 async fn load_rejects_translation_invalid_local_sources() {
259 let path = std::env::temp_dir().join(format!("adp_local_bad_{}.yaml", std::process::id()));
260 std::fs::write(&path, "dogstatsd_tag_cardinality: bogus\n").unwrap();
262
263 let result = LoadedConfiguration::load(&path, EnvPrecedence::Disabled).await;
264
265 std::fs::remove_file(&path).ok();
266 assert!(matches!(result, Err(Error::Translate { .. })));
267 }
268
269 #[tokio::test]
270 async fn load_leaves_an_incomplete_local_snapshot_to_the_selected_authority() {
271 let path = std::env::temp_dir().join(format!("adp_no_api_key_{}.yaml", std::process::id()));
275 std::fs::write(&path, "log_level: warn\n").unwrap();
276
277 let loaded = LoadedConfiguration::load(&path, EnvPrecedence::Disabled)
278 .await
279 .expect("a local snapshot without an API key loads");
280 assert_eq!("", loaded.local().shared.endpoints.api_key);
281
282 let result = loaded.standalone().await;
284
285 std::fs::remove_file(&path).ok();
286 assert!(matches!(result, Err(Error::MissingApiKey)));
287 }
288
289 fn block_on<F: std::future::Future>(future: F) -> F::Output {
293 tokio::runtime::Builder::new_current_thread()
294 .build()
295 .expect("runtime builds")
296 .block_on(future)
297 }
298
299 #[test]
300 fn a_saluki_only_environment_variable_reaches_the_model() {
301 let _guard = test_env_lock();
304 let path = std::env::temp_dir().join(format!("adp_saluki_env_{}.yaml", std::process::id()));
305 std::fs::write(&path, "{}\n").unwrap();
306 std::env::set_var("DD_DATA_PLANE_STANDALONE_MODE", "true");
307
308 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
309
310 std::env::remove_var("DD_DATA_PLANE_STANDALONE_MODE");
311 std::fs::remove_file(&path).ok();
312 assert!(loaded.local().control.standalone_mode);
313 }
314
315 #[test]
316 fn ottl_filter_error_mode_environment_variable_reaches_the_model() {
317 use agent_data_plane_config::domains::traces::OttlErrorMode;
321
322 let _guard = test_env_lock();
323 let path = std::env::temp_dir().join(format!("adp_ottl_env_{}.yaml", std::process::id()));
324 std::fs::write(&path, "{}\n").unwrap();
325 std::env::set_var("DD_OTTL_FILTER_CONFIG_ERROR_MODE", "silent");
326
327 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
328
329 std::env::remove_var("DD_OTTL_FILTER_CONFIG_ERROR_MODE");
330 std::fs::remove_file(&path).ok();
331 assert_eq!(
332 loaded.local().domains.traces.ottl_filter.error_mode,
333 OttlErrorMode::Silent
334 );
335 }
336
337 #[test]
338 fn a_structured_saluki_only_environment_variable_reaches_the_model() {
339 let _guard = test_env_lock();
340 let path = std::env::temp_dir().join(format!("adp_saluki_structured_env_{}.yaml", std::process::id()));
341 std::fs::write(&path, "{}\n").unwrap();
342 std::env::set_var(
343 "DD_METRIC_TAG_VALUE_ALLOWLIST",
344 r#"[{"metric_prefix":"requests.","tag_name":"customer_id","values":["customer-1"],"on_miss":"replace","replacement":"other"}]"#,
345 );
346
347 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
348
349 std::env::remove_var("DD_METRIC_TAG_VALUE_ALLOWLIST");
350 std::fs::remove_file(&path).ok();
351 let entries = &loaded.local().domains.dogstatsd.tag_value_allowlist;
352 assert_eq!(entries.len(), 1);
353 assert_eq!(entries[0].metric_prefix, "requests.");
354 assert_eq!(entries[0].tag_name, "customer_id");
355 assert_eq!(entries[0].values, ["customer-1"]);
356 assert_eq!(
357 entries[0].on_miss,
358 agent_data_plane_config::domains::dogstatsd::TagValueMismatchAction::Replace
359 );
360 assert_eq!(entries[0].replacement, "other");
361 }
362
363 #[test]
364 fn metric_tag_value_allowlist_from_file_preserves_whitespace() {
365 let path = std::env::temp_dir().join(format!("adp_tag_value_allowlist_{}.yaml", std::process::id()));
366 std::fs::write(
367 &path,
368 "metric_tag_value_allowlist:\n - metric_prefix: ' requests. '\n tag_name: ' customer_id '\n values: [' customer-1 ']\n on_miss: replace\n replacement: ' other '\n",
369 )
370 .unwrap();
371
372 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::Disabled))
373 .expect("allow-list strings containing whitespace should load");
374
375 std::fs::remove_file(&path).ok();
376 let entries = &loaded.local().domains.dogstatsd.tag_value_allowlist;
377 assert_eq!(entries.len(), 1);
378 assert_eq!(entries[0].metric_prefix, " requests. ");
379 assert_eq!(entries[0].tag_name, " customer_id ");
380 assert_eq!(entries[0].values, [" customer-1 "]);
381 assert_eq!(entries[0].replacement, " other ");
382 }
383
384 #[test]
385 fn invalid_metric_tag_value_allowlist_from_environment_fails_load() {
386 let _guard = test_env_lock();
387 let path = std::env::temp_dir().join(format!("adp_bad_tag_value_allowlist_env_{}.yaml", std::process::id()));
388 std::fs::write(&path, "{}\n").unwrap();
389 std::env::set_var(
390 "DD_METRIC_TAG_VALUE_ALLOWLIST",
391 r#"[{"metric_prefix":"requests.","tag_name":"customer_id"},{"metric_prefix":"requests.api.","tag_name":"customer_id"}]"#,
392 );
393
394 let result = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile));
395
396 std::env::remove_var("DD_METRIC_TAG_VALUE_ALLOWLIST");
397 std::fs::remove_file(&path).ok();
398 let error = match result {
399 Err(error) => error,
400 Ok(_) => panic!("overlapping environment allow-list should fail loading"),
401 };
402 assert!(matches!(error, Error::Deserialize { .. }));
403 assert!(error.to_string().contains("overlapping metric prefixes"));
404 }
405
406 #[test]
407 fn a_nested_datadog_environment_variable_reaches_the_by_key_view() {
408 let _guard = test_env_lock();
411 let path = std::env::temp_dir().join(format!("adp_bykey_env_{}.yaml", std::process::id()));
412 std::fs::write(&path, "{}\n").unwrap();
413 std::env::set_var("DD_PROXY_HTTP", "http://proxy.example.com");
414
415 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
416 let raw = loaded.raw_config();
417
418 std::env::remove_var("DD_PROXY_HTTP");
419 std::fs::remove_file(&path).ok();
420 assert_eq!(
421 raw.try_get_typed::<String>("proxy.http").expect("key reads"),
422 Some("http://proxy.example.com".to_string())
423 );
424 }
425
426 #[test]
427 fn the_adp_zstd_override_reaches_both_views_from_the_environment() {
428 let _guard = test_env_lock();
431 let path = std::env::temp_dir().join(format!("adp_zstd_env_{}.yaml", std::process::id()));
432 std::fs::write(&path, "{}\n").unwrap();
433 std::env::set_var("DD_DATA_PLANE_SERIALIZER_ZSTD_COMPRESSOR_LEVEL", "7");
434
435 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
436 let from_by_key = loaded
437 .raw_config()
438 .try_get_typed::<i32>("data_plane.serializer_zstd_compressor_level")
439 .expect("key reads");
440 let from_typed = loaded.local().shared.endpoints.compression.effective_zstd_level();
441
442 std::env::remove_var("DD_DATA_PLANE_SERIALIZER_ZSTD_COMPRESSOR_LEVEL");
443 std::fs::remove_file(&path).ok();
444 assert_eq!(from_by_key, Some(7));
445 assert_eq!(from_typed, 7);
446 }
447
448 #[test]
449 fn build_base_rejects_a_malformed_environment_value() {
450 let _guard = test_env_lock();
451 let path = std::env::temp_dir().join(format!("adp_build_base_bad_{}.yaml", std::process::id()));
452 std::fs::write(&path, "dogstatsd_port: 8125\n").unwrap();
453 std::env::set_var("DD_DOGSTATSD_PORT", "not-a-number");
454
455 let result = build_base(&path, EnvPrecedence::AfterFile);
456
457 std::env::remove_var("DD_DOGSTATSD_PORT");
458 std::fs::remove_file(&path).ok();
459 assert!(matches!(result, Err(Error::Base { .. })));
460 }
461
462 #[test]
463 fn build_base_accepts_a_human_readable_dogstatsd_interner_size() {
464 let _guard = test_env_lock();
465 let path = std::env::temp_dir().join(format!("adp_build_base_interner_{}.yaml", std::process::id()));
466 std::fs::write(&path, "{}\n").unwrap();
467 std::env::set_var("DD_DOGSTATSD_STRING_INTERNER_SIZE_BYTES", "12MiB");
468
469 let base = SourceTree::all_explicit(build_base(&path, EnvPrecedence::AfterFile).expect("base builds"));
470 let config = translate_strict(&base).expect("human-readable byte size translates");
471
472 std::env::remove_var("DD_DOGSTATSD_STRING_INTERNER_SIZE_BYTES");
473 std::fs::remove_file(&path).ok();
474 assert_eq!(
475 config.domains.dogstatsd.contexts.string_interner_size_bytes,
476 Some(ByteSize::mib(12).as_u64())
477 );
478 }
479}