agent_data_plane_config_system/
system.rs

1//! [`ConfigurationSystem`]: the runtime configuration, translated from the raw sources and kept
2//! current as the Datadog Agent streams updates.
3
4use 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
23// Timeout for waiting for the initial configuration snapshot from the Datadog Agent.
24//
25// This is separately from any timeout related to _connected_ to the Datadog Agent and establishing the configuration
26// update stream in the first place.
27const INITIAL_CONFIG_SNAPSHOT_TIMEOUT_SECS: u64 = 15;
28const INITIAL_CONFIG_SNAPSHOT_TIMEOUT: Duration = Duration::from_secs(INITIAL_CONFIG_SNAPSHOT_TIMEOUT_SECS);
29
30/// An error building the translated configuration from the merged sources.
31#[derive(Debug, Snafu)]
32pub enum Error {
33    /// A source model could not be deserialized from the merged configuration value.
34    #[snafu(context(false), display("{source}"))]
35    Deserialize {
36        /// The underlying deserialization error.
37        source: serde_json::Error,
38    },
39
40    /// The Datadog Agent closed the configuration stream before sending the initial snapshot.
41    #[snafu(display("configuration stream closed before the initial snapshot"))]
42    StreamClosed,
43
44    /// The Datadog Agent configuration stream closed after the initial snapshot, so no further
45    /// configuration updates can be applied.
46    #[snafu(display("configuration stream closed; no further configuration updates can be applied"))]
47    UpdateStreamClosed,
48
49    /// Timed out waiting for the initial snapshot from the Datadog Agent configuration stream.
50    #[snafu(display("timed out waiting for the initial snapshot ({INITIAL_CONFIG_SNAPSHOT_TIMEOUT_SECS} seconds)"))]
51    SnapshotTimeOut,
52
53    /// The typed base could not be built from the file and environment.
54    #[snafu(display("failed to build the configuration base: {message}"))]
55    Base {
56        /// What went wrong reading the file, parsing YAML, or decoding an environment variable.
57        message: String,
58    },
59
60    /// Translating the sources into the model failed on one or more keys.
61    #[snafu(display("{source}"))]
62    Translate {
63        /// Every translation error recorded.
64        source: TranslateErrors,
65    },
66
67    /// The fully merged configuration resolved no usable Datadog API key.
68    #[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
78/// The runtime configuration, translated from the merged sources and kept current.
79///
80/// The configuration system is the single owner of the Datadog Agent's `ConfigUpdate` stream. It
81/// folds each update onto the local source base and builds the typed [`SalukiConfiguration`] from
82/// the result. The current configuration lives in an [`ArcSwap`] cell so readers load a whole,
83/// self-consistent version with no lock, while the update task replaces it in one atomic store.
84pub struct ConfigurationSystem {
85    current: Arc<ArcSwap<SalukiConfiguration>>,
86    // The merged sources the current configuration was translated from, kept so consumers that work
87    // a key at a time have a by-key view carrying provenance. The update task holds the write lock
88    // while it replaces both, so a reader never sees sources that `current` does not match.
89    sources: Arc<RwLock<Arc<SourceTree>>>,
90    // Fired once after each accepted update so live views wake and re-project. Shared with the
91    // update task via `Arc` because `watch::Sender` is not `Clone` and both the system (to mint
92    // views) and the task (to notify) need it.
93    tick: Arc<watch::Sender<()>>,
94}
95
96impl ConfigurationSystem {
97    /// Connected authority: takes ownership of the Datadog Agent's config stream and builds the typed
98    /// model from the stream folded onto the local `base` (file + environment).
99    ///
100    /// Blocks for the first authoritative snapshot and is the strict startup gate: a snapshot that
101    /// never arrives, cannot be deserialized, or fails translation aborts the boot.
102    ///
103    /// Returns the system together with the [`ConfigurationUpdates`] that apply each later update.
104    /// No update after the first snapshot is applied until the caller runs them.
105    ///
106    /// # Errors
107    ///
108    /// Returns an error if the stream closes before the first snapshot, or the initial configuration
109    /// cannot be deserialized or translated.
110    pub(crate) async fn connected(
111        mut agent_rx: mpsc::Receiver<ConfigUpdate>, base: SourceTree,
112    ) -> Result<(Self, ConfigurationUpdates)> {
113        // The first stream message is the authoritative initial snapshot. There is no timeout on this
114        // wait by design: if the Datadog Agent never sends the snapshot, startup blocks here forever.
115        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        // Startup is the strict gate: this is the first, authoritative Agent snapshot, so any error
130        // fails the boot and we never run on bad config. At runtime (see `ConfigurationUpdates`) the
131        // same check instead rejects the offending update and keeps the last-known-good
132        // configuration, because a runtime update must never take the system down.
133        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        // The initial receiver is dropped immediately; `send_replace` works with zero receivers, and
139        // each live view subscribes its own receiver from the sender.
140        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(&current),
149                sources: Arc::clone(&sources),
150                tick: Arc::clone(&tick),
151            })),
152        };
153
154        Ok((Self { current, sources, tick }, updates))
155    }
156
157    /// Installs a static configuration without an update task.
158    ///
159    /// Live views retain their initial values because this system sends no update notifications.
160    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    /// Returns a live view of the given projection of the current configuration. Narrow further with
171    /// [`Live::project`]. This is the only way a consumer subscribes to runtime updates.
172    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    /// Loads the current translated configuration.
180    ///
181    /// The returned guard pins one whole version; a concurrent refresh never tears the read.
182    pub fn config(&self) -> arc_swap::Guard<Arc<SalukiConfiguration>> {
183        self.current.load()
184    }
185
186    /// Returns a shared handle to the current-configuration cell for readers that load it
187    /// independently.
188    pub fn current_handle(&self) -> Arc<ArcSwap<SalukiConfiguration>> {
189        Arc::clone(&self.current)
190    }
191
192    /// Returns a callback that reads the current merged sources as JSON.
193    ///
194    /// This is the raw source view: the values every input supplied, untranslated, including keys the
195    /// typed model does not carry. It reflects the last update the typed path accepted, so it is
196    /// consistent with the running configuration.
197    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    /// Loads the merged sources the current configuration was translated from.
203    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/// Applies each configuration update from the Datadog Agent after the initial snapshot.
213///
214/// [`LoadedConfiguration::run`][crate::LoadedConfiguration::run] returns this once the system holds
215/// the initial snapshot. No later update is applied until the caller calls [`run`][Self::run].
216/// Until then, updates stay queued in the Agent configuration stream.
217#[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    /// Applies each update from the Datadog Agent until the configuration stream closes.
224    ///
225    /// Each call continues from where the previous call stopped, because the stream and the
226    /// accepted configuration persist between calls. Thus, a caller can run this again after a call
227    /// fails or panics. If a call panicked while it applied an update, that update is lost.
228    ///
229    /// # Errors
230    ///
231    /// Returns [`Error::UpdateStreamClosed`] when the configuration stream closes. The stream does
232    /// not reopen, so each later call fails in the same way at once.
233    pub fn run(&self) -> impl Future<Output = std::result::Result<(), Error>> + Send + 'static {
234        let state = Arc::clone(&self.state);
235        async move {
236            // A panic does not poison this lock, so a later call takes over the state as the panic left it.
237            let mut state = state.lock_owned().await;
238            state.apply_updates().await
239        }
240    }
241}
242
243/// The state that [`ConfigurationUpdates`] keeps between runs.
244struct UpdateState {
245    agent_rx: mpsc::Receiver<ConfigUpdate>,
246    base: SourceTree,
247    // The accumulated Agent layer of the last accepted update.
248    agent: SourceTree,
249    current: Arc<ArcSwap<SalukiConfiguration>>,
250    sources: Arc<RwLock<Arc<SourceTree>>>,
251    tick: Arc<watch::Sender<()>>,
252}
253
254impl UpdateState {
255    /// Validates each update from the Datadog Agent config stream against the typed model and commits
256    /// it on success. Ends when the stream closes.
257    ///
258    /// Each update is processed individually (no burst collapse) so a rejection can be attributed to
259    /// the exact update that caused it. Updates are infrequent, so re-translating per update is cheap.
260    async fn apply_updates(&mut self) -> Result<()> {
261        while let Some(update) = self.agent_rx.recv().await {
262            // Validate-then-commit: fold onto a tentative copy of the Agent layer and drive the typed
263            // model from it. Only a fully successful update advances the committed layer, so a rejected
264            // value never lingers to re-poison a later merge.
265            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
290/// Folds one update into the accumulating Agent layer.
291///
292/// `Snapshot` replaces the layer; `Partial` applies one (possibly dotted) key.
293///
294/// Each setting's provenance is retained, which is what lets a later update that demotes a value to
295/// an Agent default stop shadowing the local value it had been overriding.
296fn 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
303/// Deserializes and translates merged source values, rejecting partially translated configuration.
304///
305/// # Errors
306///
307/// Returns an error if either source model cannot be deserialized or any key fails translation.
308pub(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
317/// Translates merged sources that are authoritative for the running process, rejecting a
318/// configuration ADP cannot run on.
319///
320/// This is [`translate_strict`] plus [`validate`]. Use it where the merged sources are complete: the
321/// Datadog Agent's snapshot layered over the local base, or the local base alone in standalone mode.
322///
323/// # Errors
324///
325/// Returns an error if translation fails, or if the translated configuration fails validation.
326pub(crate) fn translate_authoritative(merged: &SourceTree) -> Result<SalukiConfiguration> {
327    let config = translate_strict(merged)?;
328    validate(&config)?;
329    Ok(config)
330}
331
332/// Checks the invariants a configuration must satisfy for this process to do useful work.
333///
334/// Translation alone cannot make these checks. It converts one key at a time and every schema key
335/// has a default, so a setting the operator never supplied is indistinguishable from one they did
336/// until the whole merged configuration is in hand.
337///
338/// Apply this only to an authoritative configuration. The local snapshot
339/// [`LoadedConfiguration::load`][crate::LoadedConfiguration::load] produces is incomplete by design:
340/// under the Datadog Agent the API key arrives over the configuration stream, so a local-only
341/// snapshot legitimately has none, and CLI subcommands read that snapshot without ever submitting a
342/// payload.
343///
344/// # Errors
345///
346/// Returns [`Error::MissingApiKey`] if no usable API key resolved. Every payload ADP submits is
347/// authenticated with this key, so an empty one turns each flush into a rejected request that the
348/// forwarder then retries. Failing here names the cause once instead of leaving an operator to infer
349/// it from a stream of authentication failures.
350pub(crate) fn validate(config: &SalukiConfiguration) -> Result<()> {
351    // A blank key is as unusable as an absent one, and a padded key is a typo we should name rather
352    // than send.
353    if config.shared.endpoints.api_key.trim().is_empty() {
354        return Err(Error::MissingApiKey);
355    }
356
357    Ok(())
358}
359
360// TODO: A map/array-valued schema leaf is replaced wholesale when any source (file, environment, or
361// the Agent config stream) supplies it. Verify this is the intended semantic for the remote Agent
362// config stream: ADP is that stream's first consumer, so the correct behavior for a stream update to
363// a map-shaped setting may not have been defined yet.
364
365/// The sources deserialized from the merged configuration value, separated by source authority.
366struct Sources {
367    datadog: DatadogConfiguration,
368    saluki: SalukiOnly,
369}
370
371/// Deserializes both source models from the merged configuration value.
372///
373/// The source models use ordinary serde-compatible field types, so deserializing from
374/// `serde_json::Value` preserves the values. Both read the canonical nested shape: the local base
375/// is built that way by the schema-driven environment readers, and the Datadog Agent's stream
376/// delivers dotted keys that are nested on arrival.
377fn 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
383/// Translates the Datadog and Saluki-only sources into one [`SalukiConfiguration`], returning every
384/// error recorded while converting an individual Datadog value.
385///
386/// The Datadog `drive` feeds every supported key to a `DatadogTranslator`; a value that cannot be
387/// converted leaves its field at the model default and records an error. The Saluki-only values
388/// then seed their disjoint destinations, which cannot fail. The returned configuration is always
389/// complete: every valid value is present, and every invalid one holds its default.
390///
391/// `sources` is the same merged layer the models were deserialized from. The translator consults it
392/// for provenance, which a deserialized source model cannot supply: a schema key with a default is
393/// always present, so its value alone cannot say whether an input set it explicitly.
394fn 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    /// API key the connected-system fixture puts in the local base.
422    ///
423    /// An authoritative configuration must resolve one (see [`super::validate`]), and putting it in
424    /// the base rather than in a streamed snapshot keeps it in place across the snapshot replacements
425    /// these tests exercise.
426    const TEST_API_KEY: &str = "test-api-key";
427
428    /// Builds a standalone system whose authority is the local sources.
429    ///
430    /// Translates without validating so a test can state only the setting it is exercising. The
431    /// production standalone path validates; `loaded.rs` covers that.
432    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    /// Builds a connected system whose base is `base` and whose authority is the returned Agent
439    /// stream. The initial (empty) snapshot is queued before the system blocks on it, so the caller
440    /// gets back a stream ready for `Partial`/`Snapshot` updates.
441    ///
442    /// `base` is given an [`api_key`][TEST_API_KEY] unless it states its own, so a caller varying
443    /// some unrelated setting need not restate what validation requires.
444    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    /// Polls the current configuration until `predicate` holds, failing if it never does.
459    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    // Regression: the Agent streams `DD_ADDITIONAL_ENDPOINTS` as the raw JSON string it stores, and the
555    // map deserializer rejected it, so the startup gate failed and ADP restarted in a loop.
556    #[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        // The updates only include later updates. This test does not need them, so it never runs them.
575        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    // The Agent streams `null` for an empty list, including one produced by ordinary loading, such as
591    // `histogram_aggregates: []` in YAML. Its accessors read the null as empty, so startup must too.
592    #[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                // The Agent sends an object-valued setting whole, and these entry keys contain dots.
664                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        // The raw view is for diagnosing a configuration, so it must show the unsupported keys an
727        // input supplied. They are pruned from the typed model by construction, which is why the view
728        // reads the sources rather than the model.
729        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        // Sources deliver the Agent's canonical nested shape, which is what the Datadog
754        // deserializer reads. A string list supplied as one space-separated string (the form an
755        // environment variable carries) is still split on whitespace at the leaf.
756        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        // Nothing is set here, so both fields must come back as the schema defaults; the component
777        // layer no longer supplies fallbacks of its own.
778        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        // Nothing translates `autoscaling_failover_enabled` into the nested slot: resolving an
856        // environment variable to its canonical path is the environment readers' job, and they do it
857        // before a value ever reaches this point. A flattened key arriving from any other source is
858        // simply not a key the model knows.
859        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        // Startup is the strict gate: a value the sources carry but the model rejects fails the load,
867        // so the process never boots on bad config.
868        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        // Component builders used to discover this after translation. Reject zero before publishing
888        // an invalid typed model.
889        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        // A byte-size setting documented as accepting a bare integer (`10485760`) rather than a
932        // string (`"10MB"`) must not abort the strict startup gate. The typed model normalizes it,
933        // and the translator resolves it to the same byte count.
934        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        // Nothing ADP submits is accepted without a key, so the authoritative gate names the cause
943        // rather than letting every flush fail authentication. Translation alone accepts this: the
944        // schema default for `api_key` is the empty string, so nothing is missing to translate.
945        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        // An explicitly blank key is as unusable as an absent one, and whitespace is a typo worth
957        // reporting rather than submitting.
958        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        // Validation covers runtime updates too: an update that leaves the process unable to submit
981        // anything is rejected like any other invalid one, and the working key stays in place.
982        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        // The Agent reads a setting by casting whatever its configuration holds to the accessor's
1007        // type, so a boolean written where the schema declares a string, or a quoted integer, is a
1008        // configuration it accepts. Each must reach the typed model instead of aborting the strict
1009        // startup gate.
1010        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        // Send a translation-invalid update, then a valid update to a different field. Updates are
1033        // processed in order, so once the second is observed the first has already been handled.
1034        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        // The invalid update was rejected whole: the field keeps its last-known-good value rather
1054        // than falling back to a default, and the later valid update still applied.
1055        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            // Interleave a translation-invalid update mid-burst, then correct it. The invalid value
1122            // is rejected whole (last-known-good retained) rather than wedging the task, so the
1123            // baseline keeps converging on the latest valid value regardless of the transient bad one.
1124            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    // Issue #1965: the Core Agent streams every setting it knows about, including the ones nobody
1160    // configured, so its schema defaults must not overwrite the local file.
1161    #[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        // A later unrelated update is observable, so once it lands the snapshot above has been handled.
1174        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        // The Agent takes over the setting.
1196        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        // The operator removes it from the Agent's configuration, so the Agent now reports its own
1209        // default. The Agent layer must stop shadowing the local value rather than pinning the value
1210        // it last held.
1211        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        // The retained sources are the by-key view of the same configuration `current` holds, so an
1230        // update the typed path rejected must not appear in them. `check_compatibility` is the reader
1231        // that shows it: it classifies keys the typed model does not carry.
1232        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        // The cardinality value fails translation, so the whole snapshot is rejected and the
1239        // unsupported key it also carries never reaches the sources.
1240        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        // A snapshot is the producer's complete state, so a setting it omits is no longer set and the
1302        // port returns to the schema default rather than lingering at 9125.
1303        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        // Nothing has set the URL, so it is the schema default the Agent supplies.
1323        assert_eq!(dd_url.provenance, Provenance::Default);
1324        assert_eq!(dd_url.value, "https://app.datadoghq.com");
1325
1326        // The same URL, now deliberately chosen. The value is unchanged, so only provenance can carry
1327        // the fact that it became an override.
1328        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        // `Deref` reflects the value returned by the last `changed`.
1362        assert!(view.metrics_stats_enable);
1363    }
1364
1365    #[tokio::test]
1366    async fn field_view_wakes_on_its_field() {
1367        // Projecting straight to a single field needs no schema change and no central registration:
1368        // the granularity is chosen at the call site.
1369        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        // Regression: clearing `metric_filterlist` must restore the legacy list and match mode.
1391        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        // A fixed view never resolves, so this bound is deterministic rather than timing-dependent.
1445        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        // A small raw source map exercising a scalar conversion, an enum parse, a duration parse, the
1468        // raw endpoint inputs, and one seeded Saluki-only field.
1469        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        // Driven scalar conversion: i64 -> u16.
1490        assert_eq!(config.domains.dogstatsd.listeners.port, 9125);
1491        // Driven enum parse.
1492        assert_eq!(
1493            config.domains.dogstatsd.origin.tag_cardinality,
1494            OriginTagCardinality::High
1495        );
1496        // Driven `format: duration` parse: a Go duration string becomes a `Duration`.
1497        assert_eq!(config.shared.tags.expected_tags_duration, Duration::from_secs(15));
1498        // Shared deployment inputs used to derive static tags for both DogStatsD and OTLP.
1499        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        // Driven bool in a nested Datadog section.
1504        assert!(config.domains.dogstatsd.telemetry.origin_breakdown);
1505        // Raw endpoint inputs: carried through without selecting a primary endpoint here.
1506        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        // Seeded Saluki-only field.
1510        assert_eq!(config.domains.dogstatsd.listeners.tcp_port, 8126);
1511    }
1512
1513    /// Datadog-defined aggregation keys reach their typed model fields through the witness translator.
1514    #[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}