1use std::future::Future;
5use std::sync::{Arc, PoisonError, RwLock};
6use std::time::Duration;
7
8use agent_data_plane_config::{Live, SalukiConfiguration};
9use arc_swap::ArcSwap;
10use datadog_agent_config::{DatadogConfiguration, TranslateErrors};
11use saluki_config::dynamic::ConfigUpdate;
12use serde::Deserialize;
13use serde_json::Value;
14use snafu::Snafu;
15use tokio::sync::{mpsc, watch, Mutex};
16use tokio::time::timeout;
17use tracing::{debug, warn};
18
19use crate::saluki_only::SalukiOnly;
20use crate::source::SourceTree;
21use crate::translators::DatadogTranslator;
22
23const INITIAL_CONFIG_SNAPSHOT_TIMEOUT_SECS: u64 = 15;
28const INITIAL_CONFIG_SNAPSHOT_TIMEOUT: Duration = Duration::from_secs(INITIAL_CONFIG_SNAPSHOT_TIMEOUT_SECS);
29
30#[derive(Debug, Snafu)]
32pub enum Error {
33 #[snafu(context(false), display("{source}"))]
35 Deserialize {
36 source: serde_json::Error,
38 },
39
40 #[snafu(display("configuration stream closed before the initial snapshot"))]
42 StreamClosed,
43
44 #[snafu(display("configuration stream closed; no further configuration updates can be applied"))]
47 UpdateStreamClosed,
48
49 #[snafu(display("timed out waiting for the initial snapshot ({INITIAL_CONFIG_SNAPSHOT_TIMEOUT_SECS} seconds)"))]
51 SnapshotTimeOut,
52
53 #[snafu(display("failed to build the configuration base: {message}"))]
55 Base {
56 message: String,
58 },
59
60 #[snafu(display("{source}"))]
62 Translate {
63 source: TranslateErrors,
65 },
66
67 #[snafu(display(
69 "no Datadog API key is configured: set `api_key` in the Datadog Agent's configuration, or \
70 `DD_API_KEY` in the environment. Every payload is authenticated with this key, so nothing can \
71 be submitted without one"
72 ))]
73 MissingApiKey,
74}
75
76type Result<T> = std::result::Result<T, Error>;
77
78pub struct ConfigurationSystem {
85 current: Arc<ArcSwap<SalukiConfiguration>>,
86 sources: Arc<RwLock<Arc<SourceTree>>>,
90 tick: Arc<watch::Sender<()>>,
94}
95
96impl ConfigurationSystem {
97 pub(crate) async fn connected(
111 mut agent_rx: mpsc::Receiver<ConfigUpdate>, base: SourceTree,
112 ) -> Result<(Self, ConfigurationUpdates)> {
113 saluki_antithesis::reachable!("config readiness wait entered");
116 let first = match timeout(INITIAL_CONFIG_SNAPSHOT_TIMEOUT, agent_rx.recv()).await {
117 Ok(None) => {
118 saluki_antithesis::unreachable!("config stream closed before the initial snapshot");
119 return Err(Error::StreamClosed);
120 }
121 Ok(Some(first)) => first,
122 Err(_) => return Err(Error::SnapshotTimeOut),
123 };
124 saluki_antithesis::sometimes!(true, "config readiness signal received");
125
126 let mut agent = SourceTree::empty();
127 fold(&mut agent, &first);
128
129 let merged = base.overlay(&agent);
134 let config = translate_authoritative(&merged)?;
135
136 let current = Arc::new(ArcSwap::from_pointee(config));
137 let sources = Arc::new(RwLock::new(Arc::new(merged)));
138 let (tick, _) = watch::channel(());
141 let tick = Arc::new(tick);
142
143 let updates = ConfigurationUpdates {
144 state: Arc::new(Mutex::new(UpdateState {
145 agent_rx,
146 base,
147 agent,
148 current: Arc::clone(¤t),
149 sources: Arc::clone(&sources),
150 tick: Arc::clone(&tick),
151 })),
152 };
153
154 Ok((Self { current, sources, tick }, updates))
155 }
156
157 pub(crate) fn standalone(config: SalukiConfiguration, sources: SourceTree) -> Self {
161 let current = Arc::new(ArcSwap::from_pointee(config));
162 let (tick, _) = watch::channel(());
163 Self {
164 current,
165 sources: Arc::new(RwLock::new(Arc::new(sources))),
166 tick: Arc::new(tick),
167 }
168 }
169
170 pub fn live<T>(&self, project: impl for<'a> Fn(&'a SalukiConfiguration) -> &'a T + Send + Sync + 'static) -> Live<T>
173 where
174 T: Clone + PartialEq + 'static,
175 {
176 Live::new_dynamic(Arc::clone(&self.current), self.tick.subscribe(), project)
177 }
178
179 pub fn config(&self) -> arc_swap::Guard<Arc<SalukiConfiguration>> {
183 self.current.load()
184 }
185
186 pub fn current_handle(&self) -> Arc<ArcSwap<SalukiConfiguration>> {
189 Arc::clone(&self.current)
190 }
191
192 pub fn raw_snapshot(&self) -> Arc<dyn Fn() -> Value + Send + Sync> {
198 let sources = Arc::clone(&self.sources);
199 Arc::new(move || load_sources(&sources).to_value())
200 }
201
202 pub(crate) fn sources(&self) -> Arc<SourceTree> {
204 load_sources(&self.sources)
205 }
206}
207
208fn load_sources(sources: &RwLock<Arc<SourceTree>>) -> Arc<SourceTree> {
209 Arc::clone(&sources.read().unwrap_or_else(PoisonError::into_inner))
210}
211
212#[must_use = "configuration updates from the Datadog Agent are never applied unless the updates are run"]
218pub struct ConfigurationUpdates {
219 state: Arc<Mutex<UpdateState>>,
220}
221
222impl ConfigurationUpdates {
223 pub fn run(&self) -> impl Future<Output = std::result::Result<(), Error>> + Send + 'static {
234 let state = Arc::clone(&self.state);
235 async move {
236 let mut state = state.lock_owned().await;
238 state.apply_updates().await
239 }
240 }
241}
242
243struct UpdateState {
245 agent_rx: mpsc::Receiver<ConfigUpdate>,
246 base: SourceTree,
247 agent: SourceTree,
249 current: Arc<ArcSwap<SalukiConfiguration>>,
250 sources: Arc<RwLock<Arc<SourceTree>>>,
251 tick: Arc<watch::Sender<()>>,
252}
253
254impl UpdateState {
255 async fn apply_updates(&mut self) -> Result<()> {
261 while let Some(update) = self.agent_rx.recv().await {
262 let mut tentative = self.agent.clone();
266 fold(&mut tentative, &update);
267 let merged = self.base.overlay(&tentative);
268 match translate_authoritative(&merged) {
269 Ok(config) => {
270 self.agent = tentative;
271 {
272 let mut sources = self.sources.write().unwrap_or_else(PoisonError::into_inner);
273 *sources = Arc::new(merged);
274 self.current.store(Arc::new(config));
275 }
276 self.tick.send_replace(());
277 debug!("Applied configuration update.");
278 }
279 Err(e) => warn!(
280 error = %e,
281 "Rejected configuration update; keeping the last-known-good typed configuration."
282 ),
283 }
284 }
285
286 Err(Error::UpdateStreamClosed)
287 }
288}
289
290fn fold(agent: &mut SourceTree, update: &ConfigUpdate) {
297 match update {
298 ConfigUpdate::Snapshot(settings) => *agent = SourceTree::from_settings(settings),
299 ConfigUpdate::Partial(setting) => agent.set(setting),
300 }
301}
302
303pub(crate) fn translate_strict(merged: &SourceTree) -> Result<SalukiConfiguration> {
309 let Sources { datadog, saluki } = deserialize_sources(&merged.to_value())?;
310 let (config, errors) = translate(&datadog, &saluki, merged);
311 if let Some(errors) = errors {
312 return Err(Error::Translate { source: errors });
313 }
314 Ok(config)
315}
316
317pub(crate) fn translate_authoritative(merged: &SourceTree) -> Result<SalukiConfiguration> {
327 let config = translate_strict(merged)?;
328 validate(&config)?;
329 Ok(config)
330}
331
332pub(crate) fn validate(config: &SalukiConfiguration) -> Result<()> {
351 if config.shared.endpoints.api_key.trim().is_empty() {
354 return Err(Error::MissingApiKey);
355 }
356
357 Ok(())
358}
359
360struct Sources {
367 datadog: DatadogConfiguration,
368 saluki: SalukiOnly,
369}
370
371fn deserialize_sources(merged: &Value) -> Result<Sources> {
378 let saluki = SalukiOnly::deserialize(merged)?;
379 let datadog = DatadogConfiguration::deserialize(merged)?;
380 Ok(Sources { datadog, saluki })
381}
382
383fn translate(
395 datadog: &DatadogConfiguration, saluki: &SalukiOnly, sources: &SourceTree,
396) -> (SalukiConfiguration, Option<TranslateErrors>) {
397 let (mut config, errors) = DatadogTranslator::new(datadog, sources).translate();
398 saluki.seed(&mut config);
399 (config, errors)
400}
401
402#[cfg(test)]
403mod tests {
404 use std::collections::HashSet;
405 use std::sync::Arc;
406 use std::time::Duration;
407
408 use agent_data_plane_config::domains::dogstatsd::OriginTagCardinality;
409 use agent_data_plane_config::shared::V3SeriesMode;
410 use agent_data_plane_config::Provenance;
411 use agent_data_plane_config::{Live, SalukiConfiguration};
412 use datadog_agent_config::DatadogConfiguration;
413 use saluki_config::dynamic::{ConfigSetting, ConfigUpdate, Provenance as StreamProvenance};
414 use serde_json::{json, Value};
415 use tokio::sync::mpsc;
416
417 use super::{
418 translate, translate_authoritative, translate_strict, ConfigurationSystem, Error, SalukiOnly, SourceTree,
419 };
420
421 const TEST_API_KEY: &str = "test-api-key";
427
428 fn standalone_system(file: Value) -> Result<ConfigurationSystem, Error> {
433 let base = SourceTree::all_explicit(file);
434 let config = translate_strict(&base)?;
435 Ok(ConfigurationSystem::standalone(config, base))
436 }
437
438 async fn connected_system(mut base: Value) -> (ConfigurationSystem, mpsc::Sender<ConfigUpdate>) {
445 let (agent_tx, agent_rx) = mpsc::channel(100);
446 agent_tx.send(ConfigUpdate::snapshot([])).await.unwrap();
447 if let Some(base) = base.as_object_mut() {
448 base.entry("api_key").or_insert(json!(TEST_API_KEY));
449 }
450 let base = SourceTree::all_explicit(base);
451 let (system, updates) = ConfigurationSystem::connected(agent_rx, base)
452 .await
453 .expect("system builds");
454 tokio::spawn(updates.run());
455 (system, agent_tx)
456 }
457
458 async fn await_config(system: &ConfigurationSystem, what: &str, predicate: impl Fn(&SalukiConfiguration) -> bool) {
460 tokio::time::timeout(Duration::from_secs(2), async {
461 while !predicate(&system.config()) {
462 tokio::time::sleep(Duration::from_millis(5)).await;
463 }
464 })
465 .await
466 .unwrap_or_else(|_| panic!("timed out waiting for {what}"));
467 }
468
469 #[tokio::test]
470 async fn startup_current_reflects_translation() {
471 let system = standalone_system(json!({ "log_level": "warn", "dogstatsd_port": 9125 })).expect("system builds");
472 let config = system.config();
473
474 assert_eq!(config.control.logging.level, "warn");
475 assert_eq!(config.domains.dogstatsd.listeners.port, 9125);
476 }
477
478 #[test]
479 fn boolean_use_v3_api_series_enabled_is_normalized() {
480 let sources = SourceTree::all_explicit(json!({
481 "use_v3_api": {
482 "series": {
483 "enabled": true
484 }
485 }
486 }));
487
488 let config = translate_strict(&sources).expect("a boolean V3 series mode should translate");
489
490 assert_eq!(config.shared.metrics_encoding.v3_series_mode, V3SeriesMode::Enabled);
491 }
492
493 #[test]
494 fn compound_v3_series_endpoint_modes_are_rejected_at_startup() {
495 for mode in [json!(["true"]), json!({ "enabled": true })] {
496 let sources = SourceTree::all_explicit(json!({
497 "use_v3_api": { "series": { "endpoints": { "https://app.datadoghq.com": mode } } }
498 }));
499
500 assert!(translate_strict(&sources).is_err());
501 }
502 }
503
504 #[tokio::test]
505 async fn malformed_v3_endpoint_update_keeps_last_known_good_and_recovers() {
506 let (system, agent_tx) = connected_system(json!({
507 "dogstatsd_port": 9125,
508 "use_v3_api": { "series": { "endpoints": { "https://app.datadoghq.com": true } } }
509 }))
510 .await;
511
512 agent_tx
513 .send(ConfigUpdate::snapshot([
514 ConfigSetting::explicit("dogstatsd_port", json!(9999)),
515 ConfigSetting::explicit(
516 "use_v3_api.series.endpoints",
517 json!({ "https://app.datadoghq.com": ["false"] }),
518 ),
519 ]))
520 .await
521 .unwrap();
522 agent_tx
523 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
524 "log_level",
525 json!("error"),
526 )))
527 .await
528 .unwrap();
529 await_config(&system, "the update following the rejected snapshot", |config| {
530 config.control.logging.level == "error"
531 })
532 .await;
533
534 assert_eq!(system.config().domains.dogstatsd.listeners.port, 9125);
535 assert_eq!(
536 system.config().shared.metrics_encoding.v3_series_endpoint_modes["https://app.datadoghq.com"],
537 V3SeriesMode::Enabled
538 );
539
540 agent_tx
541 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
542 "use_v3_api.series.endpoints",
543 json!({ "https://app.datadoghq.com": false }),
544 )))
545 .await
546 .unwrap();
547 await_config(&system, "the corrected endpoint mode", |config| {
548 config.shared.metrics_encoding.v3_series_endpoint_modes["https://app.datadoghq.com"]
549 == V3SeriesMode::Disabled
550 })
551 .await;
552 }
553
554 #[tokio::test]
557 async fn connected_startup_accepts_json_encoded_map_settings() {
558 let (agent_tx, agent_rx) = mpsc::channel(1);
559 agent_tx
560 .send(ConfigUpdate::snapshot([
561 ConfigSetting::explicit(
562 "additional_endpoints",
563 json!(r#"{"https://app.datadoghq.com": ["second-org-key"]}"#),
564 ),
565 ConfigSetting::explicit(
566 "use_v3_api.series.endpoints",
567 json!(r#"{"https://app.datadoghq.com": "true"}"#),
568 ),
569 ]))
570 .await
571 .unwrap();
572 let base = SourceTree::all_explicit(json!({ "api_key": TEST_API_KEY }));
573
574 let (system, _updates) = ConfigurationSystem::connected(agent_rx, base)
576 .await
577 .expect("startup accepts the streamed JSON strings");
578
579 let config = system.config();
580 assert_eq!(
581 config.shared.endpoints.additional_endpoints["https://app.datadoghq.com"],
582 ["second-org-key"]
583 );
584 assert_eq!(
585 config.shared.metrics_encoding.v3_series_endpoint_modes["https://app.datadoghq.com"],
586 V3SeriesMode::Enabled
587 );
588 }
589
590 #[tokio::test]
593 async fn connected_startup_reads_null_collections_as_empty() {
594 let (agent_tx, agent_rx) = mpsc::channel(1);
595 agent_tx
596 .send(ConfigUpdate::snapshot([
597 ConfigSetting::explicit("histogram_aggregates", Value::Null),
598 ConfigSetting::explicit("proxy.no_proxy", Value::Null),
599 ConfigSetting::explicit("additional_endpoints", Value::Null),
600 ]))
601 .await
602 .unwrap();
603 let base = SourceTree::all_explicit(json!({ "api_key": TEST_API_KEY }));
604
605 let (system, _updates) = ConfigurationSystem::connected(agent_rx, base)
606 .await
607 .expect("startup accepts streamed null collections");
608
609 let config = system.config();
610 assert!(config.shared.metrics_encoding.histogram.aggregates.is_empty());
611 assert!(config.shared.endpoints.proxy.no_proxy.is_empty());
612 assert!(config.shared.endpoints.additional_endpoints.is_empty());
613 }
614
615 #[tokio::test]
616 async fn null_update_clears_a_collection_rather_than_restoring_its_default() {
617 let (system, agent_tx) = connected_system(json!({
618 "additional_endpoints": { "https://app.datadoghq.com": ["second-org-key"] },
619 "histogram_aggregates": ["max"]
620 }))
621 .await;
622 assert!(!system.config().shared.endpoints.additional_endpoints.is_empty());
623
624 for key in ["additional_endpoints", "histogram_aggregates"] {
625 agent_tx
626 .send(ConfigUpdate::Partial(ConfigSetting::explicit(key, Value::Null)))
627 .await
628 .unwrap();
629 }
630 await_config(&system, "the cleared collections", |config| {
631 config.shared.endpoints.additional_endpoints.is_empty()
632 && config.shared.metrics_encoding.histogram.aggregates.is_empty()
633 })
634 .await;
635 }
636
637 #[tokio::test]
638 async fn connected_stream_translates_metrics_v3_routing_configuration() {
639 let (system, agent_tx) = connected_system(json!({
640 "data_plane": {
641 "metrics": {
642 "v3": {
643 "series": {
644 "enabled": true
645 }
646 }
647 }
648 }
649 }))
650 .await;
651
652 assert_eq!(
653 system.config().shared.metrics_encoding.v3_series_mode,
654 V3SeriesMode::DatadogOnly
655 );
656
657 agent_tx
658 .send(ConfigUpdate::snapshot([
659 ConfigSetting::explicit("serializer_compressor_kind", json!("zstd")),
660 ConfigSetting::explicit("serializer_experimental_use_v3_api.compression_level", json!(7)),
661 ConfigSetting::explicit("use_v2_api.series", json!(false)),
662 ConfigSetting::explicit("use_v3_api.series.enabled", json!("false")),
663 ConfigSetting::explicit(
665 "use_v3_api.series.endpoints",
666 json!({ "https://app.datadoghq.com": "true" }),
667 ),
668 ConfigSetting::explicit("observability_pipelines_worker.metrics.enabled", json!(true)),
669 ConfigSetting::explicit(
670 "observability_pipelines_worker.metrics.url",
671 json!("https://opw.example.com"),
672 ),
673 ConfigSetting::explicit("observability_pipelines_worker.metrics.use_v3_api.series", json!(true)),
674 ]))
675 .await
676 .unwrap();
677
678 await_config(&system, "the streamed metrics V3 routing configuration", |config| {
679 config.shared.metrics_encoding.v3_series_mode == V3SeriesMode::Disabled
680 && config
681 .shared
682 .metrics_encoding
683 .v3_series_endpoint_modes
684 .get("https://app.datadoghq.com")
685 == Some(&V3SeriesMode::Enabled)
686 })
687 .await;
688
689 let config = system.config();
690 let metrics = &config.shared.metrics_encoding;
691 assert!(!metrics.use_v2_series_api);
692 assert_eq!(metrics.v3_api.compression_level, 7);
693 let opw = &config.shared.endpoints.opw_intake;
694 assert!(opw.enabled);
695 assert_eq!(opw.url, "https://opw.example.com");
696 assert!(opw.use_v3_series);
697 }
698
699 #[tokio::test]
700 async fn raw_snapshot_reflects_streamed_updates() {
701 let (system, agent_tx) = connected_system(json!({})).await;
702 let snapshot = system.raw_snapshot();
703 let cloned = Arc::clone(&snapshot);
704
705 assert_eq!(snapshot().pointer("/dogstatsd_port"), None);
706
707 agent_tx
708 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
709 "dogstatsd_port",
710 json!(9125),
711 )))
712 .await
713 .unwrap();
714
715 tokio::time::timeout(Duration::from_secs(2), async {
716 while cloned().pointer("/dogstatsd_port") != Some(&json!(9125)) {
717 tokio::time::sleep(Duration::from_millis(5)).await;
718 }
719 })
720 .await
721 .expect("timed out waiting for the raw snapshot to reflect the update");
722 }
723
724 #[tokio::test]
725 async fn raw_snapshot_keeps_keys_the_typed_model_does_not_carry() {
726 let (system, agent_tx) = connected_system(json!({ "dogstatsd_stats_buffer": 42 })).await;
730 let snapshot = system.raw_snapshot();
731
732 assert_eq!(snapshot().pointer("/dogstatsd_stats_buffer"), Some(&json!(42)));
733
734 agent_tx
735 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
736 "config_id",
737 json!("abc"),
738 )))
739 .await
740 .unwrap();
741
742 tokio::time::timeout(Duration::from_secs(2), async {
743 while snapshot().pointer("/config_id") != Some(&json!("abc")) {
744 tokio::time::sleep(Duration::from_millis(5)).await;
745 }
746 })
747 .await
748 .expect("timed out waiting for the streamed unsupported key");
749 }
750
751 #[tokio::test]
752 async fn nested_datadog_key_reaches_the_model() {
753 let system = standalone_system(json!({
757 "autoscaling": {
758 "failover": {
759 "enabled": true,
760 "metrics": "container.memory.usage container.cpu.usage",
761 }
762 }
763 }))
764 .expect("system builds");
765 let config = system.config();
766
767 assert!(config.shared.autoscaling_failover.enabled);
768 assert_eq!(
769 config.shared.autoscaling_failover.metrics,
770 vec!["container.memory.usage".to_string(), "container.cpu.usage".to_string()]
771 );
772 }
773
774 #[tokio::test]
775 async fn unset_autoscaling_failover_keeps_its_schema_defaults() {
776 let system = standalone_system(json!({})).expect("system builds");
779 let config = system.config();
780
781 assert!(!config.shared.autoscaling_failover.enabled);
782 assert_eq!(
783 config.shared.autoscaling_failover.metrics,
784 vec!["container.memory.usage".to_string(), "container.cpu.usage".to_string()]
785 );
786 }
787
788 #[tokio::test]
789 async fn nested_saluki_only_key_seeds_the_model() {
790 let system = standalone_system(json!({ "data_plane": { "standalone_mode": true } })).expect("system builds");
791
792 assert!(system.config().control.standalone_mode);
793 }
794
795 #[tokio::test]
796 async fn experimental_metric_allowlists_seed_the_model() {
797 let policies = json!({
798 "https://primary.example.com": ["allowed.metric"],
799 "https://secondary.example.com": []
800 });
801 let system = standalone_system(json!({
802 "experimental": {
803 "metrics_endpoint_routing": {
804 "metric_allowlist": policies.clone(),
805 "metric_prefix_allowlist": { "https://primary.example.com": ["billing."] }
806 }
807 }
808 }))
809 .expect("system builds");
810
811 assert_eq!(
812 serde_json::to_value(&system.config().domains.metrics_endpoint_routing.metric_allowlists).unwrap(),
813 policies
814 );
815 assert_eq!(
816 system
817 .config()
818 .domains
819 .metrics_endpoint_routing
820 .metric_prefix_allowlists,
821 std::collections::HashMap::from([(
822 "https://primary.example.com".to_string(),
823 vec!["billing.".to_string()]
824 )])
825 );
826 }
827
828 #[tokio::test]
829 async fn empty_experimental_sections_preserve_default_routing() {
830 for source in [
831 json!({}),
832 json!({ "experimental": {} }),
833 json!({ "experimental": { "metrics_endpoint_routing": {} } }),
834 json!({ "experimental": { "metrics_endpoint_routing": { "metric_allowlist": {} } } }),
835 json!({ "experimental": { "metrics_endpoint_routing": { "metric_prefix_allowlist": {} } } }),
836 ] {
837 let system = standalone_system(source).expect("system builds");
838 assert!(system
839 .config()
840 .domains
841 .metrics_endpoint_routing
842 .metric_allowlists
843 .is_empty());
844 assert!(system
845 .config()
846 .domains
847 .metrics_endpoint_routing
848 .metric_prefix_allowlists
849 .is_empty());
850 }
851 }
852
853 #[tokio::test]
854 async fn a_flattened_spelling_of_a_nested_key_is_not_read() {
855 let system = standalone_system(json!({ "autoscaling_failover_enabled": true })).expect("system builds");
860
861 assert!(!system.config().shared.autoscaling_failover.enabled);
862 }
863
864 #[tokio::test]
865 async fn load_fails_on_translation_invalid_startup_config() {
866 let result = standalone_system(json!({ "dogstatsd_tag_cardinality": "bogus" }));
869
870 assert!(matches!(result, Err(Error::Translate { .. })));
871 }
872
873 #[tokio::test]
874 async fn negative_dogstatsd_workers_count_is_rejected_at_startup() {
875 let result = standalone_system(json!({ "dogstatsd_workers_count": -1 }));
876
877 let Err(error) = result else {
878 panic!("negative worker count should fail the startup translation gate");
879 };
880 assert!(matches!(error, Error::Translate { .. }));
881 assert!(error.to_string().contains("dogstatsd_workers_count"));
882 assert!(error.to_string().contains("greater than or equal to 0"));
883 }
884
885 #[test]
886 fn zero_otlp_trace_interner_size_is_rejected() {
887 let sources = SourceTree::all_explicit(json!({ "otlp_config": { "traces": { "string_interner_size": 0 } } }));
890 let error = translate_strict(&sources).expect_err("zero trace interner size should fail translation");
891
892 assert!(matches!(error, Error::Deserialize { .. }));
893 assert!(error.to_string().contains("value of bytes must be greater than zero"));
894 }
895
896 #[test]
897 fn oversized_otlp_trace_interner_size_is_rejected() {
898 let sources =
899 SourceTree::all_explicit(json!({ "otlp_config": { "traces": { "string_interner_size": "2GiB" } } }));
900 let error = translate_strict(&sources).expect_err("oversized trace interner should fail translation");
901
902 assert!(matches!(error, Error::Deserialize { .. }));
903 assert!(error.to_string().contains("must not exceed 1073741824 bytes"));
904 }
905
906 #[test]
907 fn positive_otlp_trace_interner_size_is_accepted() {
908 let sources =
909 SourceTree::all_explicit(json!({ "otlp_config": { "traces": { "string_interner_size": "512KiB" } } }));
910 let config = translate_strict(&sources).expect("positive trace interner size should translate");
911
912 assert_eq!(config.domains.otlp.traces.string_interner_size.get(), 512 * 1024);
913 }
914
915 #[test]
916 fn invalid_metric_tag_value_allowlist_is_rejected_before_publication() {
917 let sources = SourceTree::all_explicit(json!({
918 "metric_tag_value_allowlist": [
919 { "metric_prefix": "requests.", "tag_name": "customer_id" },
920 { "metric_prefix": "requests.api.", "tag_name": "customer_id" }
921 ]
922 }));
923 let error = translate_strict(&sources).expect_err("overlapping allow-list prefixes should fail translation");
924
925 assert!(matches!(error, Error::Deserialize { .. }));
926 assert!(error.to_string().contains("overlapping metric prefixes"));
927 }
928
929 #[tokio::test]
930 async fn standalone_loads_numeric_byte_size() {
931 let system =
935 standalone_system(json!({ "dogstatsd_log_file_max_size": 10485760 })).expect("numeric byte size boots");
936
937 assert_eq!(system.config().domains.dogstatsd.debug_log.log_file_max_size, 10485760);
938 }
939
940 #[test]
941 fn a_configuration_without_an_api_key_is_rejected() {
942 let sources = SourceTree::all_explicit(json!({}));
946 translate_strict(&sources).expect("an absent API key still translates");
947
948 let error = translate_authoritative(&sources).expect_err("an absent API key should fail the gate");
949
950 assert!(matches!(error, Error::MissingApiKey));
951 assert!(error.to_string().contains("api_key"));
952 }
953
954 #[test]
955 fn a_blank_api_key_is_rejected() {
956 for key in ["", " ", "\t\n"] {
959 let sources = SourceTree::all_explicit(json!({ "api_key": key }));
960
961 assert!(
962 matches!(translate_authoritative(&sources), Err(Error::MissingApiKey)),
963 "{key:?} should not count as an API key"
964 );
965 }
966
967 let sources = SourceTree::all_explicit(json!({ "api_key": TEST_API_KEY }));
968 assert_eq!(
969 TEST_API_KEY,
970 translate_authoritative(&sources)
971 .expect("a real key passes the gate")
972 .shared
973 .endpoints
974 .api_key
975 );
976 }
977
978 #[tokio::test]
979 async fn an_update_that_blanks_the_api_key_is_rejected_keeping_last_known_good() {
980 let (system, agent_tx) = connected_system(json!({ "log_level": "warn" })).await;
983 assert_eq!(TEST_API_KEY, system.config().shared.endpoints.api_key);
984
985 agent_tx
986 .send(ConfigUpdate::Partial(ConfigSetting::explicit("api_key", json!(""))))
987 .await
988 .unwrap();
989 agent_tx
990 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
991 "log_level",
992 json!("error"),
993 )))
994 .await
995 .unwrap();
996
997 await_config(&system, "the later valid update to take effect", |c| {
998 c.control.logging.level == "error"
999 })
1000 .await;
1001 assert_eq!(TEST_API_KEY, system.config().shared.endpoints.api_key);
1002 }
1003
1004 #[tokio::test]
1005 async fn standalone_loads_scalars_written_in_any_form_the_agent_casts() {
1006 let system = standalone_system(json!({
1011 "use_v3_api": { "series": { "enabled": true } },
1012 "dogstatsd_port": "8126",
1013 }))
1014 .expect("scalars in Agent-castable forms boot");
1015
1016 assert_eq!(
1017 system.config().shared.metrics_encoding.v3_series_mode,
1018 V3SeriesMode::Enabled
1019 );
1020 assert_eq!(system.config().domains.dogstatsd.listeners.port, 8126);
1021 }
1022
1023 #[tokio::test]
1024 async fn translation_invalid_update_is_rejected_keeping_last_known_good() {
1025 let (system, agent_tx) =
1026 connected_system(json!({ "log_level": "warn", "dogstatsd_tag_cardinality": "high" })).await;
1027 assert_eq!(
1028 system.config().domains.dogstatsd.origin.tag_cardinality,
1029 OriginTagCardinality::High
1030 );
1031
1032 agent_tx
1035 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1036 "dogstatsd_tag_cardinality",
1037 json!("bogus"),
1038 )))
1039 .await
1040 .unwrap();
1041 agent_tx
1042 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1043 "log_level",
1044 json!("error"),
1045 )))
1046 .await
1047 .unwrap();
1048
1049 await_config(&system, "the later valid update to take effect", |c| {
1050 c.control.logging.level == "error"
1051 })
1052 .await;
1053 assert_eq!(
1056 system.config().domains.dogstatsd.origin.tag_cardinality,
1057 OriginTagCardinality::High
1058 );
1059 }
1060
1061 #[tokio::test]
1062 async fn invalid_metric_tag_value_allowlist_update_keeps_last_known_good() {
1063 let initial_allowlist = json!([{
1064 "metric_prefix": "requests.",
1065 "tag_name": "customer_id",
1066 "values": ["customer-1"]
1067 }]);
1068 let (system, agent_tx) = connected_system(json!({
1069 "log_level": "warn",
1070 "metric_tag_value_allowlist": initial_allowlist
1071 }))
1072 .await;
1073
1074 agent_tx
1075 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1076 "metric_tag_value_allowlist",
1077 json!([
1078 { "metric_prefix": "requests.", "tag_name": "customer_id" },
1079 { "metric_prefix": "requests.api.", "tag_name": "customer_id" }
1080 ]),
1081 )))
1082 .await
1083 .unwrap();
1084 agent_tx
1085 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1086 "log_level",
1087 json!("error"),
1088 )))
1089 .await
1090 .unwrap();
1091
1092 await_config(&system, "the later valid update to take effect", |config| {
1093 config.control.logging.level == "error"
1094 })
1095 .await;
1096 let config = system.config();
1097 assert_eq!(config.domains.dogstatsd.tag_value_allowlist.len(), 1);
1098 assert_eq!(
1099 config.domains.dogstatsd.tag_value_allowlist[0].metric_prefix,
1100 "requests."
1101 );
1102 assert_eq!(config.domains.dogstatsd.tag_value_allowlist[0].values, ["customer-1"]);
1103 }
1104
1105 #[tokio::test]
1106 async fn converges_to_latest_value_under_burst() {
1107 let (system, agent_tx) = connected_system(json!({ "log_level": "info" })).await;
1108
1109 let burst = [
1110 "warn", "error", "debug", "trace", "info", "warn", "error", "debug", "trace", "info", "warn", "error",
1111 "debug", "trace", "info", "warn", "error", "debug", "trace",
1112 ];
1113 for (i, level) in burst.iter().enumerate() {
1114 agent_tx
1115 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1116 "log_level",
1117 json!(level),
1118 )))
1119 .await
1120 .unwrap();
1121 if i == burst.len() / 2 {
1125 agent_tx
1126 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1127 "dogstatsd_tag_cardinality",
1128 json!("bogus"),
1129 )))
1130 .await
1131 .unwrap();
1132 agent_tx
1133 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1134 "dogstatsd_tag_cardinality",
1135 json!("high"),
1136 )))
1137 .await
1138 .unwrap();
1139 }
1140 }
1141 let final_level = "error";
1142 agent_tx
1143 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1144 "log_level",
1145 json!(final_level),
1146 )))
1147 .await
1148 .unwrap();
1149
1150 await_config(
1151 &system,
1152 "the current configuration to converge to the final value",
1153 |c| c.control.logging.level == final_level,
1154 )
1155 .await;
1156 assert_eq!(system.config().control.logging.level, final_level);
1157 }
1158
1159 #[tokio::test]
1162 async fn a_defaulted_agent_value_does_not_erase_a_local_one() {
1163 let (system, agent_tx) = connected_system(json!({ "dd_url": "https://vector.example.com" })).await;
1164
1165 agent_tx
1166 .send(ConfigUpdate::snapshot([ConfigSetting::new(
1167 "dd_url",
1168 json!("https://app.datadoghq.com"),
1169 StreamProvenance::Default,
1170 )]))
1171 .await
1172 .unwrap();
1173 agent_tx
1175 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1176 "log_level",
1177 json!("error"),
1178 )))
1179 .await
1180 .unwrap();
1181 await_config(&system, "the trailing update to take effect", |c| {
1182 c.control.logging.level == "error"
1183 })
1184 .await;
1185
1186 let dd_url = &system.config().shared.endpoints.dd_url;
1187 assert!(dd_url.is_explicit());
1188 assert_eq!(dd_url.value, "https://vector.example.com");
1189 }
1190
1191 #[tokio::test]
1192 async fn demoting_an_agent_value_to_a_default_reveals_the_local_value() {
1193 let (system, agent_tx) = connected_system(json!({ "dd_url": "https://vector.example.com" })).await;
1194
1195 agent_tx
1197 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1198 "dd_url",
1199 json!("https://app.datadoghq.eu"),
1200 )))
1201 .await
1202 .unwrap();
1203 await_config(&system, "the Agent override to take effect", |c| {
1204 c.shared.endpoints.dd_url.value == "https://app.datadoghq.eu"
1205 })
1206 .await;
1207
1208 agent_tx
1212 .send(ConfigUpdate::Partial(ConfigSetting::new(
1213 "dd_url",
1214 json!("https://app.datadoghq.com"),
1215 StreamProvenance::Default,
1216 )))
1217 .await
1218 .unwrap();
1219
1220 await_config(&system, "the local value to be revealed again", |c| {
1221 c.shared.endpoints.dd_url.value == "https://vector.example.com"
1222 })
1223 .await;
1224 assert!(system.config().shared.endpoints.dd_url.is_explicit());
1225 }
1226
1227 #[tokio::test]
1228 async fn the_retained_sources_follow_accepted_updates_only() {
1229 let (system, agent_tx) = connected_system(json!({})).await;
1233 let all_pipelines = HashSet::new();
1234 system
1235 .check_compatibility(&all_pipelines)
1236 .expect("nothing unsupported is set yet");
1237
1238 agent_tx
1241 .send(ConfigUpdate::snapshot([
1242 ConfigSetting::explicit("heroku_dyno", json!(true)),
1243 ConfigSetting::explicit("dogstatsd_tag_cardinality", json!("bogus")),
1244 ]))
1245 .await
1246 .unwrap();
1247 agent_tx
1248 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1249 "log_level",
1250 json!("error"),
1251 )))
1252 .await
1253 .unwrap();
1254 await_config(&system, "the update following the rejected snapshot", |config| {
1255 config.control.logging.level == "error"
1256 })
1257 .await;
1258 system
1259 .check_compatibility(&all_pipelines)
1260 .expect("a rejected update does not reach the sources");
1261
1262 agent_tx
1263 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1264 "heroku_dyno",
1265 json!(true),
1266 )))
1267 .await
1268 .unwrap();
1269 agent_tx
1270 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1271 "log_level",
1272 json!("warn"),
1273 )))
1274 .await
1275 .unwrap();
1276 await_config(&system, "the accepted unsupported key", |config| {
1277 config.control.logging.level == "warn"
1278 })
1279 .await;
1280 system
1281 .check_compatibility(&all_pipelines)
1282 .expect_err("an accepted update reaches the sources");
1283 }
1284
1285 #[tokio::test]
1286 async fn a_snapshot_replaces_the_agent_layer() {
1287 let (system, agent_tx) = connected_system(json!({})).await;
1288
1289 agent_tx
1290 .send(ConfigUpdate::snapshot([ConfigSetting::explicit(
1291 "dogstatsd_port",
1292 json!(9125),
1293 )]))
1294 .await
1295 .unwrap();
1296 await_config(&system, "the first snapshot to take effect", |c| {
1297 c.domains.dogstatsd.listeners.port == 9125
1298 })
1299 .await;
1300
1301 agent_tx
1304 .send(ConfigUpdate::snapshot([ConfigSetting::explicit(
1305 "log_level",
1306 json!("error"),
1307 )]))
1308 .await
1309 .unwrap();
1310
1311 await_config(&system, "the replacing snapshot to take effect", |c| {
1312 c.control.logging.level == "error"
1313 })
1314 .await;
1315 assert_eq!(system.config().domains.dogstatsd.listeners.port, 8125);
1316 }
1317
1318 #[tokio::test]
1319 async fn a_live_view_wakes_when_only_provenance_changes() {
1320 let (system, agent_tx) = connected_system(json!({})).await;
1321 let mut dd_url = system.live(|c| &c.shared.endpoints.dd_url);
1322 assert_eq!(dd_url.provenance, Provenance::Default);
1324 assert_eq!(dd_url.value, "https://app.datadoghq.com");
1325
1326 agent_tx
1329 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1330 "dd_url",
1331 json!("https://app.datadoghq.com"),
1332 )))
1333 .await
1334 .unwrap();
1335
1336 let updated = tokio::time::timeout(Duration::from_secs(2), dd_url.changed())
1337 .await
1338 .expect("the view observes a provenance-only change");
1339 assert_eq!(updated.provenance, Provenance::Explicit);
1340 assert_eq!(updated.value, "https://app.datadoghq.com");
1341 }
1342
1343 #[tokio::test]
1344 async fn live_view_observes_debug_log_update() {
1345 let (system, agent_tx) = connected_system(json!({ "dogstatsd_metrics_stats_enable": false })).await;
1346 let mut view = system.live(|c| &c.domains.dogstatsd.debug_log);
1347 assert!(!view.metrics_stats_enable);
1348
1349 agent_tx
1350 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1351 "dogstatsd_metrics_stats_enable",
1352 json!(true),
1353 )))
1354 .await
1355 .unwrap();
1356
1357 let updated = tokio::time::timeout(Duration::from_secs(2), view.changed())
1358 .await
1359 .expect("view observes the debug-log update");
1360 assert!(updated.metrics_stats_enable);
1361 assert!(view.metrics_stats_enable);
1363 }
1364
1365 #[tokio::test]
1366 async fn field_view_wakes_on_its_field() {
1367 let (system, agent_tx) = connected_system(json!({ "dogstatsd_metrics_stats_enable": false })).await;
1370 let mut stats = system.live(|c| &c.domains.dogstatsd.debug_log.metrics_stats_enable);
1371 assert!(!*stats);
1372
1373 agent_tx
1374 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1375 "dogstatsd_metrics_stats_enable",
1376 json!(true),
1377 )))
1378 .await
1379 .unwrap();
1380
1381 let updated = tokio::time::timeout(Duration::from_secs(2), stats.changed())
1382 .await
1383 .expect("field view observes its field's update");
1384 assert!(updated);
1385 assert!(*stats);
1386 }
1387
1388 #[tokio::test]
1389 async fn live_metric_filter_follows_current_and_legacy_precedence() {
1390 let (system, agent_tx) = connected_system(json!({})).await;
1392 let mut metric_filter = system.live(|c| &c.domains.dogstatsd.metric_filter);
1393 assert!(metric_filter.values.is_empty());
1394
1395 agent_tx
1396 .send(ConfigUpdate::snapshot([
1397 ConfigSetting::explicit("metric_filterlist", json!(["current.duration.max"])),
1398 ConfigSetting::explicit("metric_filterlist_match_prefix", json!(false)),
1399 ConfigSetting::explicit("statsd_metric_blocklist", json!(["legacy.duration"])),
1400 ConfigSetting::explicit("statsd_metric_blocklist_match_prefix", json!(true)),
1401 ]))
1402 .await
1403 .unwrap();
1404
1405 let current = tokio::time::timeout(Duration::from_secs(2), metric_filter.changed())
1406 .await
1407 .expect("view observes the current filterlist taking precedence");
1408 assert_eq!(current.values, vec!["current.duration.max".to_string()]);
1409 assert!(!current.match_prefix);
1410
1411 agent_tx
1412 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1413 "metric_filterlist",
1414 json!([]),
1415 )))
1416 .await
1417 .unwrap();
1418
1419 let legacy = tokio::time::timeout(Duration::from_secs(2), metric_filter.changed())
1420 .await
1421 .expect("view observes the fallback to the legacy blocklist");
1422 assert_eq!(legacy.values, vec!["legacy.duration".to_string()]);
1423 assert!(legacy.match_prefix);
1424
1425 agent_tx
1426 .send(ConfigUpdate::Partial(ConfigSetting::explicit(
1427 "metric_filterlist",
1428 json!(["current.duration.avg"]),
1429 )))
1430 .await
1431 .unwrap();
1432
1433 let restored = tokio::time::timeout(Duration::from_secs(2), metric_filter.changed())
1434 .await
1435 .expect("view observes the current filterlist shadowing the legacy blocklist again");
1436 assert_eq!(restored.values, vec!["current.duration.avg".to_string()]);
1437 assert!(!restored.match_prefix);
1438 }
1439
1440 #[tokio::test]
1441 async fn fixed_view_never_changes() {
1442 let mut view: Live<bool> = Live::new_fixed(true);
1443 assert!(*view);
1444 assert!(tokio::time::timeout(Duration::from_millis(100), view.changed())
1446 .await
1447 .is_err());
1448 }
1449
1450 #[tokio::test]
1451 async fn live_views_reflect_startup_configuration() {
1452 let system = standalone_system(json!({ "dogstatsd_metrics_stats_enable": true })).expect("system builds");
1453 let config = system.config();
1454
1455 let debug_log = system.live(|c| &c.domains.dogstatsd.debug_log);
1456 assert_eq!(&*debug_log, &config.domains.dogstatsd.debug_log);
1457
1458 let prefix_filter = system.live(|c| &c.domains.dogstatsd.prefix_filter);
1459 assert_eq!(&*prefix_filter, &config.domains.dogstatsd.prefix_filter);
1460
1461 let multi_region_failover = system.live(|c| &c.domains.multi_region_failover);
1462 assert_eq!(&*multi_region_failover, &config.domains.multi_region_failover);
1463 }
1464
1465 #[test]
1466 fn translate_small_map_through_witness_and_seed() {
1467 let sources = SourceTree::all_explicit(json!({
1470 "api_key": "abc",
1471 "dd_url": "https://custom.example.com",
1472 "dogstatsd_port": 9125,
1473 "dogstatsd_tag_cardinality": "high",
1474 "expected_tags_duration": "15s",
1475 "provider_kind": "gke-autopilot",
1476 "eks_fargate": true,
1477 "kubernetes_kubelet_nodename": "fargate-node",
1478 "cluster_name": "fargate-cluster",
1479 "telemetry": { "dogstatsd_origin": true },
1480 "dogstatsd_tcp_port": 8126,
1481 }));
1482 let value = sources.to_value();
1483 let datadog: DatadogConfiguration = serde_json::from_value(value.clone()).expect("datadog source deserializes");
1484 let saluki: SalukiOnly = serde_json::from_value(value).expect("saluki-only source deserializes");
1485
1486 let (config, errors) = translate(&datadog, &saluki, &sources);
1487 assert!(errors.is_none(), "translation of a valid map records no error");
1488
1489 assert_eq!(config.domains.dogstatsd.listeners.port, 9125);
1491 assert_eq!(
1493 config.domains.dogstatsd.origin.tag_cardinality,
1494 OriginTagCardinality::High
1495 );
1496 assert_eq!(config.shared.tags.expected_tags_duration, Duration::from_secs(15));
1498 assert_eq!(config.shared.static_tags.provider_kind, "gke-autopilot");
1500 assert!(config.shared.static_tags.eks_fargate);
1501 assert_eq!(config.shared.static_tags.kubernetes_kubelet_nodename, "fargate-node");
1502 assert_eq!(config.shared.static_tags.cluster_name, "fargate-cluster");
1503 assert!(config.domains.dogstatsd.telemetry.origin_breakdown);
1505 assert_eq!(config.shared.endpoints.api_key, "abc");
1507 assert!(config.shared.endpoints.dd_url.is_explicit());
1508 assert_eq!(config.shared.endpoints.dd_url.value, "https://custom.example.com");
1509 assert_eq!(config.domains.dogstatsd.listeners.tcp_port, 8126);
1511 }
1512
1513 #[test]
1515 fn datadog_aggregation_keys_reach_the_model() {
1516 let sources = SourceTree::all_explicit(json!({
1517 "dogstatsd_expiry_seconds": 60,
1518 "dogstatsd_flush_incomplete_buckets": true,
1519 }));
1520 let value = sources.to_value();
1521 let datadog: DatadogConfiguration = serde_json::from_value(value.clone()).expect("datadog source deserializes");
1522 let saluki: SalukiOnly = serde_json::from_value(value).expect("saluki-only source deserializes");
1523
1524 let (config, errors) = translate(&datadog, &saluki, &sources);
1525 assert!(errors.is_none(), "translation of a valid map records no error");
1526
1527 let aggregation = &config.domains.dogstatsd.aggregation;
1528 assert_eq!(aggregation.counter_expiry_seconds, Some(60));
1529 assert!(aggregation.flush_open_windows);
1530 }
1531}