From cecc06701dfeb0208cff97df5a369f3b2f46241c Mon Sep 17 00:00:00 2001 From: Raymond Zhao Date: Tue, 28 Jul 2026 15:33:50 -0400 Subject: [PATCH 1/2] fix(config): accept scalar additional endpoint API keys --- .../config/build/datadog_config_gen.rs | 70 ++++++++++++++----- .../src/generated/datadog_configuration.rs | 1 + lib/datadog-agent/config/src/lib.rs | 35 ++++++++++ lib/datadog-agent/config/src/list_de.rs | 60 +++++++++++++++- 4 files changed, 145 insertions(+), 21 deletions(-) diff --git a/lib/datadog-agent/config/build/datadog_config_gen.rs b/lib/datadog-agent/config/build/datadog_config_gen.rs index 353f70d973..ba9288cab8 100644 --- a/lib/datadog-agent/config/build/datadog_config_gen.rs +++ b/lib/datadog-agent/config/build/datadog_config_gen.rs @@ -15,11 +15,11 @@ //! cannot render it faithfully. Parsing the duration once, at the deserialization boundary, keeps a //! bare number from reaching the translator as an ambiguous unit. //! -//! Every `Vec` leaf also gets a shape-tolerant deserializer: a config file or the remote -//! Agent stream supplies a real sequence, but an environment variable supplies a single -//! space-separated string (`DD_DOGSTATSD_TAGS="env:prod team:core"`), and the field must accept -//! both. Handling that in the deserializer (rather than pre-splitting env values elsewhere) keeps -//! the concern in one place; see `stringlistize` and `crate::list_de`. +//! String-list fields also get shape-tolerant decoders. A standalone `Vec` may arrive +//! as a real sequence or a space-separated environment string, while a `HashMap>` may carry each map value as either one scalar string or a sequence. Handling those +//! shapes at the deserialization boundary keeps downstream types consistent; see `stringlistize` +//! and `crate::list_de`. use std::collections::{BTreeMap, HashMap, HashSet}; use std::path::Path; @@ -400,19 +400,20 @@ fn duration_defaults_module(durations: &BTreeMap) -> String { module } -/// Give every `Vec` leaf a shape-tolerant deserializer. +/// Give string-list leaves shape-tolerant decoders. /// /// String lists arrive as a real sequence from a file or the remote Agent stream, but as a single /// space-separated string from an environment variable (`DD_DOGSTATSD_TAGS="a b"`), so each field -/// must accept both. We push an extra `#[serde(deserialize_with = ...)]` attribute rather than -/// replacing the field's existing serde attributes: serde merges multiple `#[serde(...)]`, so the -/// field keeps its `default`/`skip_serializing_if`/`alias` and only gains the tolerant reader (the -/// same additive technique as `inject_serde_aliases`). +/// must accept both. Map values containing string lists may likewise arrive as either one scalar or +/// a sequence and are normalized into vectors. /// -/// Fields are matched by their Rust type, not by name: `Vec` is unambiguous, needs no -/// schema lookup, and correctly skips `HashMap>` (a map, not a list) and the -/// duration leaves (already retyped to `Duration`). Runs after `PathShortener`, so the type reads -/// as `Vec` rather than `::std::vec::Vec<::std::string::String>`. +/// We push an extra `#[serde(deserialize_with = ...)]` attribute rather than replacing the field's +/// existing serde attributes: serde merges multiple `#[serde(...)]`, so the field keeps its +/// `default`/`skip_serializing_if`/`alias` and only gains the tolerant reader (the same additive +/// technique as `inject_serde_aliases`). +/// +/// Fields are matched by their Rust type, not by name. Runs after `PathShortener`, so the types read +/// as `Vec` and `HashMap>` rather than fully qualified prelude paths. fn stringlistize(file: &mut syn::File) { for item in &mut file.items { let Item::Struct(s) = item else { continue }; @@ -420,12 +421,15 @@ fn stringlistize(file: &mut syn::File) { continue; }; for field in &mut fields.named { - if !is_vec_string(&field.ty) { - continue; + if is_vec_string(&field.ty) { + field.attrs.push(parse_quote!( + #[serde(deserialize_with = "crate::list_de::deserialize_space_separated_or_seq")] + )); + } else if is_string_map_vec_string(&field.ty) { + field.attrs.push(parse_quote!( + #[serde(deserialize_with = "crate::list_de::deserialize_string_map_scalar_or_seq")] + )); } - field.attrs.push(parse_quote!( - #[serde(deserialize_with = "crate::list_de::deserialize_space_separated_or_seq")] - )); } } } @@ -449,6 +453,34 @@ fn is_vec_string(ty: &syn::Type) -> bool { inner.path.segments.last().is_some_and(|seg| seg.ident == "String") } +/// Returns whether `ty` is exactly `HashMap>`. +fn is_string_map_vec_string(ty: &syn::Type) -> bool { + let syn::Type::Path(tp) = ty else { return false }; + let Some(last) = tp.path.segments.last() else { + return false; + }; + if last.ident != "HashMap" { + return false; + } + let syn::PathArguments::AngleBracketed(args) = &last.arguments else { + return false; + }; + let mut args = args.args.iter(); + let Some(syn::GenericArgument::Type(key)) = args.next() else { + return false; + }; + let Some(syn::GenericArgument::Type(value)) = args.next() else { + return false; + }; + is_string(key) && is_vec_string(value) +} + +/// Returns whether `ty` is exactly `String`. +fn is_string(ty: &syn::Type) -> bool { + let syn::Type::Path(tp) = ty else { return false }; + tp.path.segments.last().is_some_and(|seg| seg.ident == "String") +} + /// Make nested-section fields non-optional with `#[serde(default)]`. /// /// typify renders each non-required object property as `Option`. Every section diff --git a/lib/datadog-agent/config/src/generated/datadog_configuration.rs b/lib/datadog-agent/config/src/generated/datadog_configuration.rs index dc464ac607..d6b18c7ed8 100644 --- a/lib/datadog-agent/config/src/generated/datadog_configuration.rs +++ b/lib/datadog-agent/config/src/generated/datadog_configuration.rs @@ -42,6 +42,7 @@ pub mod error { #[derive(::serde::Deserialize, ::serde::Serialize, Clone, Debug)] pub struct DatadogConfiguration { #[serde(default, skip_serializing_if = ":: std :: collections :: HashMap::is_empty")] + #[serde(deserialize_with = "crate::list_de::deserialize_string_map_scalar_or_seq")] pub additional_endpoints: HashMap>, #[serde(default)] diff --git a/lib/datadog-agent/config/src/lib.rs b/lib/datadog-agent/config/src/lib.rs index e7f6431d09..fcf8678e0b 100644 --- a/lib/datadog-agent/config/src/lib.rs +++ b/lib/datadog-agent/config/src/lib.rs @@ -59,6 +59,41 @@ mod string_list_shape_tests { } } +#[cfg(test)] +mod string_map_list_shape_tests { + use super::DatadogConfiguration; + + #[test] + fn additional_endpoints_accept_scalar_values() { + let config: DatadogConfiguration = serde_json::from_value(serde_json::json!({ + "additional_endpoints": { + "https://agent.datadoghq.com.": "ENC[vault://api-key]" + } + })) + .expect("scalar additional endpoint API key deserializes"); + + assert_eq!( + config.additional_endpoints["https://agent.datadoghq.com."], + ["ENC[vault://api-key]"] + ); + } + + #[test] + fn additional_endpoints_accept_sequence_values() { + let config: DatadogConfiguration = serde_json::from_value(serde_json::json!({ + "additional_endpoints": { + "https://agent.datadoghq.com.": ["first", "second"] + } + })) + .expect("additional endpoint API key sequence deserializes"); + + assert_eq!( + config.additional_endpoints["https://agent.datadoghq.com."], + ["first", "second"] + ); + } +} + #[cfg(test)] mod byte_size_shape_tests { use super::DatadogConfiguration; diff --git a/lib/datadog-agent/config/src/list_de.rs b/lib/datadog-agent/config/src/list_de.rs index 362edce1b6..d0be35b67b 100644 --- a/lib/datadog-agent/config/src/list_de.rs +++ b/lib/datadog-agent/config/src/list_de.rs @@ -5,12 +5,15 @@ //! space-separated string (for example, `DD_DOGSTATSD_TAGS="env:prod team:core"`). A single field //! must accept both forms. //! -//! Deserializing here keeps that difference at the boundary: every string-list field accepts either -//! shape, while downstream consumers always receive a `Vec`. +//! Map values containing string lists have a similar compatibility shape: a single value can arrive +//! as a scalar string, while multiple values arrive as a sequence. Deserializing here keeps those +//! differences at the boundary, while downstream consumers always receive a `Vec`. +use std::collections::HashMap; use std::fmt; use serde::de::{self, Deserializer, SeqAccess, Visitor}; +use serde::Deserialize; /// Deserialize a `Vec` from either a sequence or a space-separated string. /// @@ -45,6 +48,36 @@ where deserializer.deserialize_any(SpaceSeparatedOrSeq) } +/// Deserialize string-list map values from either scalar strings or sequences. +/// +/// Scalar values are normalized into one-element vectors. Unlike standalone string-list fields, +/// scalar map values are not split on whitespace because each scalar represents one complete value. +pub(crate) fn deserialize_string_map_scalar_or_seq<'de, D>( + deserializer: D, +) -> Result>, D::Error> +where + D: Deserializer<'de>, +{ + #[derive(Deserialize)] + #[serde(untagged)] + enum ScalarOrSeq { + Scalar(String), + Seq(Vec), + } + + let values = HashMap::::deserialize(deserializer)?; + Ok(values + .into_iter() + .map(|(key, value)| { + let value = match value { + ScalarOrSeq::Scalar(value) => vec![value], + ScalarOrSeq::Seq(values) => values, + }; + (key, value) + }) + .collect()) +} + #[cfg(test)] mod tests { use super::*; @@ -84,4 +117,27 @@ mod tests { fn wrong_shape_is_rejected() { assert!(serde_json::from_str::(r#"{"list": 5}"#).is_err()); } + + #[derive(serde::Deserialize)] + struct MapHolder { + #[serde(deserialize_with = "deserialize_string_map_scalar_or_seq")] + map: HashMap>, + } + + fn parse_map(json: &str) -> HashMap> { + serde_json::from_str::(json).unwrap().map + } + + #[test] + fn string_map_accepts_scalar_and_sequence_values() { + let parsed = parse_map(r#"{"map":{"one":"api-key","many":["first","second"],"none":[]}}"#); + assert_eq!(parsed["one"], ["api-key"]); + assert_eq!(parsed["many"], ["first", "second"]); + assert!(parsed["none"].is_empty()); + } + + #[test] + fn string_map_rejects_non_string_values() { + assert!(serde_json::from_str::(r#"{"map":{"endpoint":5}}"#).is_err()); + } } From 6baeba29ed2d6795b73f8ab87b16c7f302aeb489 Mon Sep 17 00:00:00 2001 From: Raymond Zhao <35050708+rayz@users.noreply.github.com> Date: Tue, 28 Jul 2026 18:21:34 -0400 Subject: [PATCH 2/2] fix(metrics): skip v2 encoding when v3 is enabled for all endpoints (#2219) --- bin/agent-data-plane/src/cli/run.rs | 6 +- .../src/config_registry/forwarder.rs | 4 +- .../config/schema/schema_overlay.yaml | 4 +- .../src/common/datadog/config.rs | 6 + .../src/common/datadog/endpoints.rs | 6 +- .../src/encoders/datadog/metrics/mod.rs | 255 +++++++++++++++++- .../src/forwarders/cluster_agent/mod.rs | 22 ++ 7 files changed, 281 insertions(+), 22 deletions(-) diff --git a/bin/agent-data-plane/src/cli/run.rs b/bin/agent-data-plane/src/cli/run.rs index 697d5d274e..a215b3b87c 100644 --- a/bin/agent-data-plane/src/cli/run.rs +++ b/bin/agent-data-plane/src/cli/run.rs @@ -518,7 +518,8 @@ fn add_mrf_metrics_pipeline_to_blueprint( let mrf_gateway_config = MrfMetricsGatewayConfiguration::new(mrf_config.clone(), config.clone()); let mrf_metrics_config = DatadogMetricsConfiguration::from_configuration(config) - .error_context("Failed to configure Multi-Region Failover Datadog Metrics encoder.")?; + .error_context("Failed to configure Multi-Region Failover Datadog Metrics encoder.")? + .with_metrics_endpoint_override(mrf_dd_url.clone()); let mrf_forwarder_config = DatadogForwarderConfiguration::from_configuration(config) .map(|config| { @@ -571,7 +572,8 @@ fn add_autoscaling_failover_metrics_pipeline_to_blueprint( let af_gateway_config = AutoscalingFailoverGatewayConfiguration::new(af_config); let af_metrics_config = DatadogMetricsConfiguration::from_configuration(config) - .error_context("Failed to configure autoscaling failover metrics encoder.")?; + .error_context("Failed to configure autoscaling failover metrics encoder.")? + .with_v2_series_only(); let cluster_agent_forwarder_config = ClusterAgentForwarderConfiguration::from_configuration(config, ca_url, ca_token) .error_context("Failed to configure Cluster Agent forwarder.")?; diff --git a/lib/datadog-agent/config-testing/src/config_registry/forwarder.rs b/lib/datadog-agent/config-testing/src/config_registry/forwarder.rs index 516693b527..e48fad7abb 100644 --- a/lib/datadog-agent/config-testing/src/config_registry/forwarder.rs +++ b/lib/datadog-agent/config-testing/src/config_registry/forwarder.rs @@ -22,7 +22,7 @@ crate::declare_annotations! { support_level: SupportLevel::Full, additional_yaml_paths: &[], env_var_override: None, - used_by: &[structs::FORWARDER_CONFIGURATION], + used_by: &[structs::DATADOG_METRICS_CONFIGURATION, structs::FORWARDER_CONFIGURATION], value_type_override: None, test_json: None, pipeline_affinity: PipelineAffinity::CrossCutting, @@ -33,7 +33,7 @@ crate::declare_annotations! { support_level: SupportLevel::Full, additional_yaml_paths: &[], env_var_override: None, - used_by: &[structs::FORWARDER_CONFIGURATION], + used_by: &[structs::DATADOG_METRICS_CONFIGURATION, structs::FORWARDER_CONFIGURATION], value_type_override: None, test_json: None, pipeline_affinity: PipelineAffinity::CrossCutting, diff --git a/lib/datadog-agent/config/schema/schema_overlay.yaml b/lib/datadog-agent/config/schema/schema_overlay.yaml index 40c1191f87..501ef63c4d 100644 --- a/lib/datadog-agent/config/schema/schema_overlay.yaml +++ b/lib/datadog-agent/config/schema/schema_overlay.yaml @@ -713,7 +713,7 @@ inventory: pipelines: [cross_cutting] description: "Override intake endpoint URL" test_support: - used_by: [ForwarderConfiguration] + used_by: [DatadogMetricsConfiguration, ForwarderConfiguration] additional_attributes: config_registry_filename: forwarder.rs @@ -2388,7 +2388,7 @@ inventory: pipelines: [cross_cutting] description: "Datadog site domain" test_support: - used_by: [ForwarderConfiguration] + used_by: [DatadogMetricsConfiguration, ForwarderConfiguration] additional_attributes: config_registry_filename: forwarder.rs diff --git a/lib/saluki-components/src/common/datadog/config.rs b/lib/saluki-components/src/common/datadog/config.rs index 0b655e090c..78783b9d56 100644 --- a/lib/saluki-components/src/common/datadog/config.rs +++ b/lib/saluki-components/src/common/datadog/config.rs @@ -395,6 +395,12 @@ impl ForwarderConfiguration { self.opw_metrics = OpwMetricsConfiguration::default(); } + /// Forces series metrics routing to accept only V2 payloads. + pub(crate) fn force_v2_series(&mut self) { + self.data_plane_metrics_v3_series_enabled = false; + self.v3_api.series.shadow_sites.clear(); + } + /// Builds resolved endpoints with routing metadata. /// /// The normal primary and OPW metrics primary endpoints share the same dynamic API key source. diff --git a/lib/saluki-components/src/common/datadog/endpoints.rs b/lib/saluki-components/src/common/datadog/endpoints.rs index b0f3a455d9..82296061d3 100644 --- a/lib/saluki-components/src/common/datadog/endpoints.rs +++ b/lib/saluki-components/src/common/datadog/endpoints.rs @@ -32,7 +32,7 @@ pub const DEFAULT_SITE: &str = "datadoghq.com"; /// A `dd_url` equal to this constant carries no override intent and must not shadow `site`. const DEFAULT_PRIMARY_ENDPOINT: &str = "https://app.datadoghq.com"; -fn default_site() -> String { +pub(crate) fn default_site() -> String { DEFAULT_SITE.to_owned() } @@ -43,7 +43,7 @@ fn default_site() -> String { /// value equal to the default is treated as `None`, allowing `site` to determine the endpoint. This /// only affects the serde path; programmatic callers such as `set_dd_url` bypass serde and are /// unaffected. -fn deserialize_dd_url<'de, D>(deserializer: D) -> Result, D::Error> +pub(crate) fn deserialize_dd_url<'de, D>(deserializer: D) -> Result, D::Error> where D: serde::Deserializer<'de>, { @@ -895,7 +895,7 @@ fn add_data_plane_version_prefix(mut endpoint: Url) -> Result, site: &str, api_key: &str, ) -> Result { let raw_endpoint = match override_url { diff --git a/lib/saluki-components/src/encoders/datadog/metrics/mod.rs b/lib/saluki-components/src/encoders/datadog/metrics/mod.rs index eb83f3125f..f1786a0265 100644 --- a/lib/saluki-components/src/encoders/datadog/metrics/mod.rs +++ b/lib/saluki-components/src/encoders/datadog/metrics/mod.rs @@ -46,7 +46,10 @@ use self::{ use crate::{ common::datadog::{ clamp_payload_limits, default_serializer_compressor_kind, - endpoints::{series_v3_config_can_enable_v3, AdditionalEndpoints}, + endpoints::{ + calculate_resolved_endpoint, default_site, deserialize_dd_url, series_v3_config_can_enable_v3, + AdditionalEndpoints, EndpointV3Settings, ResolvedEndpoint, V3EndpointConfig, DEFAULT_SITE, + }, io::RB_BUFFER_CHUNK_SIZE, protocol::{MetricsPayloadInfo, UseV3ApiSeriesConfig, V3ApiConfig}, request_builder::{RequestBuilder, RequestBuilderError}, @@ -162,21 +165,33 @@ fn selected_metrics_primary_v3_override( opw_enabled: bool, opw_url: &str, opw_use_v3_series: bool, vector_enabled: bool, vector_url: &str, vector_use_v3_series: bool, ) -> Option { + selected_metrics_primary_endpoint( + opw_enabled, + opw_url, + opw_use_v3_series, + vector_enabled, + vector_url, + vector_use_v3_series, + ) + .map(|(_, use_v3_series)| use_v3_series) +} + +fn selected_metrics_primary_endpoint<'a>( + opw_enabled: bool, opw_url: &'a str, opw_use_v3_series: bool, vector_enabled: bool, vector_url: &'a str, + vector_use_v3_series: bool, +) -> Option<(&'a str, bool)> { if opw_enabled { - metrics_primary_v3_override_for_url(opw_url, opw_use_v3_series) + let opw_url = opw_url.trim(); + metrics_primary_url_can_resolve(opw_url).then_some((opw_url, opw_use_v3_series)) } else if vector_enabled { - metrics_primary_v3_override_for_url(vector_url, vector_use_v3_series) + let vector_url = vector_url.trim(); + metrics_primary_url_can_resolve(vector_url).then_some((vector_url, vector_use_v3_series)) } else { None } } -fn metrics_primary_v3_override_for_url(url: &str, use_v3_series: bool) -> Option { - metrics_primary_url_can_resolve(url).then_some(use_v3_series) -} - fn metrics_primary_url_can_resolve(url: &str) -> bool { - let url = url.trim(); if url.is_empty() { return false; } @@ -388,6 +403,18 @@ pub struct DatadogMetricsConfiguration { #[serde(default, rename = "vector_metrics_use_v3_api_series")] vector_metrics_use_v3_api_series: bool, + /// The Datadog site used to resolve the primary metrics endpoint. + /// + /// Defaults to `datadoghq.com`. + #[serde(default = "default_site")] + site: String, + + /// The optional explicit primary metrics endpoint. + /// + /// Defaults to unset, in which case `site` determines the endpoint. + #[serde(default, alias = "url", deserialize_with = "deserialize_dd_url")] + dd_url: Option, + /// Additional endpoints that metrics may be dual-shipped to. #[serde(default)] additional_endpoints: AdditionalEndpoints, @@ -405,6 +432,27 @@ impl DatadogMetricsConfiguration { self } + /// Restricts endpoint-aware protocol selection to a single overridden metrics endpoint. + /// + /// This mirrors a forwarder branch that replaces the normal primary endpoint and removes additional and + /// OPW/Vector endpoints, such as Multi-Region Failover. + pub fn with_metrics_endpoint_override(mut self, dd_url: String) -> Self { + self.dd_url = Some(dd_url); + self.additional_endpoints = AdditionalEndpoints::default(); + self.observability_pipelines_worker_metrics_enabled = false; + self.vector_metrics_enabled = false; + self + } + + /// Forces series metrics to use V2 without producing V3 shadow payloads. + /// + /// This is used for local destinations that only accept the V2 series protocol, such as the Cluster Agent. + pub fn with_v2_series_only(mut self) -> Self { + self.data_plane_metrics_v3_series_enabled = false; + self.v3_api.series.shadow_sample_rate = 0.0; + self + } + fn v3_payload_limits(&self) -> V3PayloadLimits { V3PayloadLimits::new( self.max_series_payload_size, @@ -413,6 +461,81 @@ impl DatadogMetricsConfiguration { self.max_series_points_per_payload, ) } + + fn endpoint_requires_v2_series( + &self, endpoint: &ResolvedEndpoint, metrics_primary_v3_override: Option, + serializer_v3_configured_endpoint: Option<&str>, + ) -> bool { + let settings = EndpointV3Settings::from_v3_config(V3EndpointConfig { + configured_endpoint: endpoint.configured_endpoint(), + resolved_endpoint: endpoint.endpoint(), + serializer_v3_configured_endpoint, + data_plane_v3_series_enabled: self.data_plane_metrics_v3_series_enabled, + series_config: &self.use_v3_api_series, + metrics_primary_v3_override, + serializer_v3_series_endpoints: &self.v3_api.series.endpoints, + serializer_v3_sketches_endpoints: &self.v3_api.sketches.endpoints, + series_validate: self.v3_api.series.validate, + sketches_validate: self.v3_api.sketches.validate, + series_shadow_sites: &self.v3_api.series.shadow_sites, + }); + + !settings.use_v3_series || settings.series_validation_mode + } + + fn configured_primary_endpoint(&self) -> String { + match self.dd_url.as_deref() { + Some(url) => url.to_string(), + None => { + let base_domain = if self.site.is_empty() { DEFAULT_SITE } else { &self.site }; + format!("https://app.{base_domain}") + } + } + } + + fn requires_v2_series(&self, metrics_v3_disabled_by_compressor: bool) -> Result { + if metrics_v3_disabled_by_compressor || !self.data_plane_metrics_v3_series_enabled { + return Ok(true); + } + + let configured_primary_endpoint = self.configured_primary_endpoint(); + if let Some((metrics_primary_url, metrics_primary_v3_override)) = selected_metrics_primary_endpoint( + self.observability_pipelines_worker_metrics_enabled, + &self.observability_pipelines_worker_metrics_url, + self.observability_pipelines_worker_metrics_use_v3_api_series, + self.vector_metrics_enabled, + &self.vector_metrics_url, + self.vector_metrics_use_v3_api_series, + ) { + let metrics_primary = calculate_resolved_endpoint(Some(metrics_primary_url), &self.site, "") + .error_context("Failed parsing/resolving the metrics primary destination endpoint.")?; + if self.endpoint_requires_v2_series( + &metrics_primary, + Some(metrics_primary_v3_override), + Some(&configured_primary_endpoint), + ) { + return Ok(true); + } + } else { + let primary = calculate_resolved_endpoint(self.dd_url.as_deref(), &self.site, "") + .error_context("Failed parsing/resolving the primary destination endpoint.")?; + if self.endpoint_requires_v2_series(&primary, None, None) { + return Ok(true); + } + } + + for endpoint in self + .additional_endpoints + .resolved_endpoints(None) + .error_context("Failed parsing/resolving the additional destination endpoints.")? + { + if self.endpoint_requires_v2_series(&endpoint, None, None) { + return Ok(true); + } + } + + Ok(false) + } } #[async_trait] @@ -524,11 +647,16 @@ impl EncoderBuilder for DatadogMetricsConfiguration { } else { generic_payload_limits }; - let mut v2_series_builder = v2::create_v2_request_builder(series_endpoint, &v2_endpoint_config) - .await - .error_context("Failed to create V2 series request builder.")?; - v2_series_builder.with_len_limits(series_uncompressed_limit, series_compressed_limit)?; - let v2_series_builder = Some(v2_series_builder); + let v2_series_builder = if self.requires_v2_series(metrics_v3_disabled_by_compressor)? { + let mut builder = v2::create_v2_request_builder(series_endpoint, &v2_endpoint_config) + .await + .error_context("Failed to create V2 series request builder.")?; + builder.with_len_limits(series_uncompressed_limit, series_compressed_limit)?; + Some(builder) + } else { + debug!("All metrics series endpoints use authoritative V3; disabling V2 series encoding."); + None + }; let (sketches_uncompressed_limit, sketches_compressed_limit) = generic_payload_limits; let mut v2_sketch_builder = v2::create_v2_request_builder(MetricsEndpoint::Sketches, &v2_endpoint_config) @@ -1900,6 +2028,107 @@ serializer_experimental_use_v3_api: ); } + fn v3_series_config(raw: &str) -> DatadogMetricsConfiguration { + serde_yaml::from_str(raw).expect("configuration should deserialize") + } + + #[test] + fn mixed_v2_and_v3_endpoints_require_v2_series() { + let config = v3_series_config( + r#" +data_plane_metrics_v3_series_enabled: true +use_v3_api_series_enabled: "datadog_only" +additional_endpoints: + https://custom.example.com: + - additional-api-key +"#, + ); + + assert!(config.requires_v2_series(false).expect("endpoints should resolve")); + } + + #[test] + fn validation_requires_v2_series() { + let config = v3_series_config( + r#" +data_plane_metrics_v3_series_enabled: true +use_v3_api_series_enabled: "true" +serializer_experimental_use_v3_api: + series: + validate: true +"#, + ); + + assert!(config.requires_v2_series(false).expect("endpoints should resolve")); + } + + #[test] + fn all_v3_serializer_endpoints_do_not_require_v2_series() { + let config = v3_series_config( + r#" +dd_url: https://agent.datad0g.com. +data_plane_metrics_v3_series_enabled: true +use_v3_api_series_enabled: "false" +additional_endpoints: + https://agent.datadoghq.com.: + - additional-api-key +serializer_experimental_use_v3_api: + series: + endpoints: + - https://agent.datad0g.com. + - https://agent.datadoghq.com. +"#, + ); + + assert!(!config.requires_v2_series(false).expect("endpoints should resolve")); + } + + #[test] + fn endpoint_override_uses_the_overridden_endpoint_protocol() { + let config = v3_series_config( + r#" +dd_url: https://primary.example.com +data_plane_metrics_v3_series_enabled: true +use_v3_api_series_enabled: "false" +serializer_experimental_use_v3_api: + series: + endpoints: + - https://primary.example.com + - https://v3-mrf.example.com +"#, + ); + + let v2_mrf_config = config + .clone() + .with_metrics_endpoint_override("https://v2-mrf.example.com".to_string()); + let v3_mrf_config = config.with_metrics_endpoint_override("https://v3-mrf.example.com".to_string()); + + assert!(v2_mrf_config + .requires_v2_series(false) + .expect("V2 MRF endpoint should resolve")); + assert!(!v3_mrf_config + .requires_v2_series(false) + .expect("V3 MRF endpoint should resolve")); + } + + #[test] + fn v2_series_only_override_keeps_v2_and_disables_shadowing() { + let config = v3_series_config( + r#" +data_plane_metrics_v3_series_enabled: true +use_v3_api_series_enabled: "true" +serializer_experimental_use_v3_api: + series: + shadow_sample_rate: 1.0 +"#, + ) + .with_v2_series_only(); + + assert!(!config.data_plane_metrics_v3_series_enabled); + assert_eq!(0.0, config.v3_api.series.shadow_sample_rate); + assert!(config.requires_v2_series(false).expect("endpoint should resolve")); + } + #[test] fn agent_default_v3_does_not_enable_opw_only_encoder_mode() { let series_config = UseV3ApiSeriesConfig::default(); diff --git a/lib/saluki-components/src/forwarders/cluster_agent/mod.rs b/lib/saluki-components/src/forwarders/cluster_agent/mod.rs index 58efdbdb08..22fef4b875 100644 --- a/lib/saluki-components/src/forwarders/cluster_agent/mod.rs +++ b/lib/saluki-components/src/forwarders/cluster_agent/mod.rs @@ -55,6 +55,7 @@ impl ClusterAgentForwarderConfiguration { endpoint.set_dd_url(endpoint_url); endpoint.set_api_key(auth_token); forwarder_config.clear_opw_metrics_endpoint(); + forwarder_config.force_v2_series(); Ok(Self { forwarder_config, @@ -250,6 +251,17 @@ mod tests { "enabled": true, "url": "https://opw.example.com" } + }, + "data_plane_metrics_v3_series_enabled": true, + "use_v3_api": { + "series": { + "enabled": "true" + } + }, + "serializer_experimental_use_v3_api": { + "series": { + "shadow_sites": ["example.com"] + } } })), None, @@ -257,6 +269,14 @@ mod tests { ) .await; + let unmodified_forwarder = + ForwarderConfiguration::from_configuration(&config).expect("forwarder configuration should parse"); + assert!(unmodified_forwarder.data_plane_metrics_v3_series_enabled()); + assert_eq!( + &["example.com".to_string()], + unmodified_forwarder.v3_api().series.shadow_sites.as_slice() + ); + let config = ClusterAgentForwarderConfiguration::from_configuration( &config, "https://cluster-agent.example.com".to_string(), @@ -275,5 +295,7 @@ mod tests { "https://cluster-agent.example.com/" ); assert_eq!(endpoints[0].endpoint().cached_api_key(), "secret-token"); + assert!(!config.forwarder_config.data_plane_metrics_v3_series_enabled()); + assert!(config.forwarder_config.v3_api().series.shadow_sites.is_empty()); } }