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 serde_json::Value;
15use tokio::sync::mpsc;
16
17use crate::saluki_env_overlay;
18use crate::source::SourceTree;
19use crate::system::{translate_strict, validate, ConfigurationSystem, ConfigurationUpdates, Error};
20
21#[derive(Clone, Copy, Debug, Eq, PartialEq)]
23pub enum EnvPrecedence {
24 BeforeFile,
28 AfterFile,
31 Disabled,
33}
34
35pub struct LoadedConfiguration {
37 base: SourceTree,
39 local: SalukiConfiguration,
42}
43
44impl LoadedConfiguration {
45 pub async fn load(path: impl AsRef<Path>, env: EnvPrecedence) -> Result<Self, Error> {
56 let base = SourceTree::all_explicit(build_base(path.as_ref(), env)?);
59 let local = translate_strict(&base)?;
60 Ok(Self { base, local })
61 }
62
63 pub fn local(&self) -> &SalukiConfiguration {
67 &self.local
68 }
69
70 pub async fn run(
81 self, config_stream: mpsc::Receiver<ConfigUpdate>,
82 ) -> Result<(ConfigurationSystem, ConfigurationUpdates), Error> {
83 ConfigurationSystem::connected(config_stream, self.base).await
84 }
85
86 pub async fn standalone(self) -> Result<ConfigurationSystem, Error> {
96 validate(&self.local)?;
97 Ok(ConfigurationSystem::standalone(self.local, self.base))
98 }
99}
100
101fn build_base(path: &Path, env: EnvPrecedence) -> Result<Value, Error> {
109 let text = std::fs::read_to_string(path).map_err(|e| Error::Base {
110 message: format!("read `{}`: {e}", path.display()),
111 })?;
112 let mut base: Value = serde_yaml::from_str(&text).map_err(|e| Error::Base {
113 message: format!("parse `{}`: {e}", path.display()),
114 })?;
115 drop_nulls(&mut base);
116 if base.is_null() {
117 base = Value::Object(serde_json::Map::new());
118 }
119
120 let overwrite = match env {
121 EnvPrecedence::Disabled => return Ok(base),
122 EnvPrecedence::AfterFile => true,
123 EnvPrecedence::BeforeFile => false,
124 };
125 apply_datadog_env(&mut base, overwrite).map_err(|message| Error::Base { message })?;
126 saluki_env_overlay::apply_env(&mut base, overwrite).map_err(|message| Error::Base { message })?;
127 Ok(base)
128}
129
130fn drop_nulls(value: &mut Value) {
133 if let Value::Object(map) = value {
134 map.retain(|_, v| !v.is_null());
135 for v in map.values_mut() {
136 drop_nulls(v);
137 }
138 }
139}
140
141#[cfg(test)]
142mod tests {
143 use bytesize::ByteSize;
144 use saluki_config::test_env_lock;
145 use serde_json::json;
146
147 use super::*;
148
149 #[test]
150 fn build_base_composes_file_and_environment_by_precedence() {
151 let _guard = test_env_lock();
152 let path = std::env::temp_dir().join(format!("adp_build_base_{}.yaml", std::process::id()));
153 std::fs::write(
154 &path,
155 "dogstatsd_port: 8125\ndogstatsd_non_local_traffic: false\nempty_key:\n",
156 )
157 .unwrap();
158 std::env::set_var("DD_DOGSTATSD_PORT", "9125");
159
160 let base = build_base(&path, EnvPrecedence::AfterFile).expect("base builds");
163 assert_eq!(base.get("dogstatsd_port"), Some(&json!(9125)));
164 assert_eq!(base.get("dogstatsd_non_local_traffic"), Some(&json!(false)));
165 assert!(base.get("empty_key").is_none());
166
167 let base = build_base(&path, EnvPrecedence::BeforeFile).expect("base builds");
169 assert_eq!(base.get("dogstatsd_port"), Some(&json!(8125)));
170
171 let base = build_base(&path, EnvPrecedence::Disabled).expect("base builds");
173 assert_eq!(base.get("dogstatsd_port"), Some(&json!(8125)));
174
175 std::env::remove_var("DD_DOGSTATSD_PORT");
176 std::fs::remove_file(&path).ok();
177 }
178
179 #[tokio::test]
180 async fn local_exposes_translated_configuration() {
181 let path = std::env::temp_dir().join(format!("adp_local_{}.yaml", std::process::id()));
183 std::fs::write(&path, "log_level: warn\ndogstatsd_port: 9125\n").unwrap();
184
185 let loaded = LoadedConfiguration::load(&path, EnvPrecedence::Disabled)
186 .await
187 .expect("local sources load");
188 let config = loaded.local();
189
190 assert_eq!(config.control.logging.level, "warn");
191 assert_eq!(config.domains.dogstatsd.listeners.port, 9125);
192
193 std::fs::remove_file(&path).ok();
194 }
195
196 #[tokio::test]
197 async fn load_rejects_translation_invalid_local_sources() {
198 let path = std::env::temp_dir().join(format!("adp_local_bad_{}.yaml", std::process::id()));
199 std::fs::write(&path, "dogstatsd_tag_cardinality: bogus\n").unwrap();
201
202 let result = LoadedConfiguration::load(&path, EnvPrecedence::Disabled).await;
203
204 std::fs::remove_file(&path).ok();
205 assert!(matches!(result, Err(Error::Translate { .. })));
206 }
207
208 #[tokio::test]
209 async fn load_leaves_an_incomplete_local_snapshot_to_the_selected_authority() {
210 let path = std::env::temp_dir().join(format!("adp_no_api_key_{}.yaml", std::process::id()));
214 std::fs::write(&path, "log_level: warn\n").unwrap();
215
216 let loaded = LoadedConfiguration::load(&path, EnvPrecedence::Disabled)
217 .await
218 .expect("a local snapshot without an API key loads");
219 assert_eq!("", loaded.local().shared.endpoints.api_key);
220
221 let result = loaded.standalone().await;
223
224 std::fs::remove_file(&path).ok();
225 assert!(matches!(result, Err(Error::MissingApiKey)));
226 }
227
228 async fn load_from_file(test: &str, contents: &str) -> LoadedConfiguration {
232 let path = std::env::temp_dir().join(format!("adp_{test}_{}.yaml", std::process::id()));
233 std::fs::write(&path, contents).unwrap();
234 let loaded = LoadedConfiguration::load(&path, EnvPrecedence::Disabled)
235 .await
236 .expect("local sources load");
237 std::fs::remove_file(&path).ok();
238 loaded
239 }
240
241 async fn await_log_level(system: &ConfigurationSystem, level: &str) {
243 tokio::time::timeout(std::time::Duration::from_secs(2), async {
244 while system.config().control.logging.level != level {
245 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
246 }
247 })
248 .await
249 .unwrap_or_else(|_| panic!("the configuration should report the log level `{level}`"));
250 }
251
252 fn log_level_update(level: &str) -> ConfigUpdate {
253 use saluki_config::dynamic::ConfigSetting;
254
255 ConfigUpdate::Partial(ConfigSetting::explicit("log_level", json!(level)))
256 }
257
258 #[tokio::test]
259 async fn run_applies_later_updates_only_once_the_updates_run() {
260 let loaded = load_from_file("run_updates", "api_key: test-api-key\nlog_level: info\n").await;
261
262 let (agent_tx, agent_rx) = mpsc::channel(8);
264 agent_tx.send(ConfigUpdate::snapshot([])).await.unwrap();
265 let (system, updates) = loaded.run(agent_rx).await.expect("the initial snapshot is accepted");
266 assert_eq!(system.config().control.logging.level, "info");
267
268 agent_tx.send(log_level_update("warn")).await.unwrap();
271 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
272 assert_eq!(system.config().control.logging.level, "info");
273
274 tokio::spawn(updates.run());
275 await_log_level(&system, "warn").await;
276 }
277
278 #[tokio::test]
279 async fn updates_continue_where_a_stopped_run_left_off() {
280 let loaded = load_from_file("run_restart", "api_key: test-api-key\nlog_level: info\n").await;
281
282 let (agent_tx, agent_rx) = mpsc::channel(8);
283 agent_tx.send(ConfigUpdate::snapshot([])).await.unwrap();
284 let (system, updates) = loaded.run(agent_rx).await.expect("the initial snapshot is accepted");
285
286 let first_run = tokio::spawn(updates.run());
287 agent_tx.send(log_level_update("warn")).await.unwrap();
288 await_log_level(&system, "warn").await;
289
290 first_run.abort();
293 let _ = first_run.await;
294 agent_tx.send(log_level_update("error")).await.unwrap();
295
296 tokio::spawn(updates.run());
297 await_log_level(&system, "error").await;
298 }
299
300 #[tokio::test]
301 async fn updates_fail_once_the_stream_closes() {
302 let loaded = load_from_file("run_ended", "api_key: test-api-key\n").await;
303
304 let (agent_tx, agent_rx) = mpsc::channel(8);
305 agent_tx.send(ConfigUpdate::snapshot([])).await.unwrap();
306 let (_system, updates) = loaded.run(agent_rx).await.expect("the initial snapshot is accepted");
307 drop(agent_tx);
308
309 assert!(matches!(updates.run().await, Err(Error::UpdateStreamClosed)));
311 assert!(matches!(updates.run().await, Err(Error::UpdateStreamClosed)));
312 }
313
314 #[tokio::test]
315 async fn run_fails_if_the_stream_closes_before_the_initial_snapshot() {
316 let loaded = load_from_file("run_closed", "api_key: test-api-key\n").await;
317
318 let (agent_tx, agent_rx) = mpsc::channel::<ConfigUpdate>(1);
319 drop(agent_tx);
320
321 let result = loaded.run(agent_rx).await;
322 assert!(matches!(result, Err(Error::StreamClosed)));
323 }
324
325 fn block_on<F: std::future::Future>(future: F) -> F::Output {
329 tokio::runtime::Builder::new_current_thread()
330 .build()
331 .expect("runtime builds")
332 .block_on(future)
333 }
334
335 #[test]
336 fn a_saluki_only_environment_variable_reaches_the_model() {
337 let _guard = test_env_lock();
340 let path = std::env::temp_dir().join(format!("adp_saluki_env_{}.yaml", std::process::id()));
341 std::fs::write(&path, "{}\n").unwrap();
342 std::env::set_var("DD_DATA_PLANE_STANDALONE_MODE", "true");
343
344 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
345
346 std::env::remove_var("DD_DATA_PLANE_STANDALONE_MODE");
347 std::fs::remove_file(&path).ok();
348 assert!(loaded.local().control.standalone_mode);
349 }
350
351 #[test]
352 fn ottl_filter_error_mode_environment_variable_reaches_the_model() {
353 use agent_data_plane_config::domains::traces::OttlErrorMode;
357
358 let _guard = test_env_lock();
359 let path = std::env::temp_dir().join(format!("adp_ottl_env_{}.yaml", std::process::id()));
360 std::fs::write(&path, "{}\n").unwrap();
361 std::env::set_var("DD_OTTL_FILTER_CONFIG_ERROR_MODE", "silent");
362
363 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
364
365 std::env::remove_var("DD_OTTL_FILTER_CONFIG_ERROR_MODE");
366 std::fs::remove_file(&path).ok();
367 assert_eq!(
368 loaded.local().domains.traces.ottl_filter.error_mode,
369 OttlErrorMode::Silent
370 );
371 }
372
373 #[test]
374 fn a_structured_saluki_only_environment_variable_reaches_the_model() {
375 let _guard = test_env_lock();
376 let path = std::env::temp_dir().join(format!("adp_saluki_structured_env_{}.yaml", std::process::id()));
377 std::fs::write(&path, "{}\n").unwrap();
378 std::env::set_var(
379 "DD_METRIC_TAG_VALUE_ALLOWLIST",
380 r#"[{"metric_prefix":"requests.","tag_name":"customer_id","values":["customer-1"],"on_miss":"replace","replacement":"other"}]"#,
381 );
382
383 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
384
385 std::env::remove_var("DD_METRIC_TAG_VALUE_ALLOWLIST");
386 std::fs::remove_file(&path).ok();
387 let entries = &loaded.local().domains.dogstatsd.tag_value_allowlist;
388 assert_eq!(entries.len(), 1);
389 assert_eq!(entries[0].metric_prefix, "requests.");
390 assert_eq!(entries[0].tag_name, "customer_id");
391 assert_eq!(entries[0].values, ["customer-1"]);
392 assert_eq!(
393 entries[0].on_miss,
394 agent_data_plane_config::domains::dogstatsd::TagValueMismatchAction::Replace
395 );
396 assert_eq!(entries[0].replacement, "other");
397 }
398
399 #[test]
400 fn metric_tag_value_allowlist_from_file_preserves_whitespace() {
401 let path = std::env::temp_dir().join(format!("adp_tag_value_allowlist_{}.yaml", std::process::id()));
402 std::fs::write(
403 &path,
404 "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",
405 )
406 .unwrap();
407
408 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::Disabled))
409 .expect("allow-list strings containing whitespace should load");
410
411 std::fs::remove_file(&path).ok();
412 let entries = &loaded.local().domains.dogstatsd.tag_value_allowlist;
413 assert_eq!(entries.len(), 1);
414 assert_eq!(entries[0].metric_prefix, " requests. ");
415 assert_eq!(entries[0].tag_name, " customer_id ");
416 assert_eq!(entries[0].values, [" customer-1 "]);
417 assert_eq!(entries[0].replacement, " other ");
418 }
419
420 #[test]
421 fn invalid_metric_tag_value_allowlist_from_environment_fails_load() {
422 let _guard = test_env_lock();
423 let path = std::env::temp_dir().join(format!("adp_bad_tag_value_allowlist_env_{}.yaml", std::process::id()));
424 std::fs::write(&path, "{}\n").unwrap();
425 std::env::set_var(
426 "DD_METRIC_TAG_VALUE_ALLOWLIST",
427 r#"[{"metric_prefix":"requests.","tag_name":"customer_id"},{"metric_prefix":"requests.api.","tag_name":"customer_id"}]"#,
428 );
429
430 let result = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile));
431
432 std::env::remove_var("DD_METRIC_TAG_VALUE_ALLOWLIST");
433 std::fs::remove_file(&path).ok();
434 let error = match result {
435 Err(error) => error,
436 Ok(_) => panic!("overlapping environment allow-list should fail loading"),
437 };
438 assert!(matches!(error, Error::Deserialize { .. }));
439 assert!(error.to_string().contains("overlapping metric prefixes"));
440 }
441
442 #[test]
443 fn the_adp_zstd_override_reaches_the_model_from_the_environment() {
444 let _guard = test_env_lock();
447 let path = std::env::temp_dir().join(format!("adp_zstd_env_{}.yaml", std::process::id()));
448 std::fs::write(&path, "{}\n").unwrap();
449 std::env::set_var("DD_DATA_PLANE_SERIALIZER_ZSTD_COMPRESSOR_LEVEL", "7");
450
451 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
452 let level = loaded.local().shared.endpoints.compression.effective_zstd_level();
453
454 std::env::remove_var("DD_DATA_PLANE_SERIALIZER_ZSTD_COMPRESSOR_LEVEL");
455 std::fs::remove_file(&path).ok();
456 assert_eq!(level, 7);
457 }
458
459 #[test]
460 fn the_adp_stop_timeout_override_reaches_the_model_from_the_environment() {
461 let _guard = test_env_lock();
462 let path = std::env::temp_dir().join(format!("adp_stop_timeout_env_{}.yaml", std::process::id()));
463 std::fs::write(&path, "{}\n").unwrap();
464 std::env::set_var("DD_DATA_PLANE_STOP_TIMEOUT", "45");
465
466 let loaded = block_on(LoadedConfiguration::load(&path, EnvPrecedence::AfterFile)).expect("local sources load");
467 let stop_timeout = loaded.local().control.stop_timeout;
468
469 std::env::remove_var("DD_DATA_PLANE_STOP_TIMEOUT");
470 std::fs::remove_file(&path).ok();
471 assert_eq!(stop_timeout, Some(std::time::Duration::from_secs(45)));
472 }
473
474 #[test]
475 fn build_base_rejects_a_malformed_environment_value() {
476 let _guard = test_env_lock();
477 let path = std::env::temp_dir().join(format!("adp_build_base_bad_{}.yaml", std::process::id()));
478 std::fs::write(&path, "dogstatsd_port: 8125\n").unwrap();
479 std::env::set_var("DD_DOGSTATSD_PORT", "not-a-number");
480
481 let result = build_base(&path, EnvPrecedence::AfterFile);
482
483 std::env::remove_var("DD_DOGSTATSD_PORT");
484 std::fs::remove_file(&path).ok();
485 assert!(matches!(result, Err(Error::Base { .. })));
486 }
487
488 #[test]
489 fn build_base_accepts_a_human_readable_dogstatsd_interner_size() {
490 let _guard = test_env_lock();
491 let path = std::env::temp_dir().join(format!("adp_build_base_interner_{}.yaml", std::process::id()));
492 std::fs::write(&path, "{}\n").unwrap();
493 std::env::set_var("DD_DOGSTATSD_STRING_INTERNER_SIZE_BYTES", "12MiB");
494
495 let base = SourceTree::all_explicit(build_base(&path, EnvPrecedence::AfterFile).expect("base builds"));
496 let config = translate_strict(&base).expect("human-readable byte size translates");
497
498 std::env::remove_var("DD_DOGSTATSD_STRING_INTERNER_SIZE_BYTES");
499 std::fs::remove_file(&path).ok();
500 assert_eq!(
501 config.domains.dogstatsd.contexts.string_interner_size_bytes,
502 Some(ByteSize::mib(12).as_u64())
503 );
504 }
505}