1use std::{collections::VecDeque, ops::Range, time::Duration};
2
3use agent_data_plane_config::shared::{Endpoints, MetricsEncoding};
4use async_trait::async_trait;
5use ddsketch::DDSketch;
6use facet::Facet;
7use http::{HeaderValue, Method, Request};
8use protobuf::{rt::WireType, CodedOutputStream};
9use saluki_common::{
10 buf::{ChunkedBytesBuffer, FrozenChunkedBytesBuffer},
11 iter::ReusableDeduplicator,
12 task::HandleExt as _,
13};
14use saluki_config::GenericConfiguration;
15use saluki_context::tags::{SharedTagSet, Tag};
16use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
17use saluki_core::{
18 components::{encoders::*, ComponentContext},
19 data_model::{
20 event::{
21 metric::{Metric, MetricOrigin, MetricValues},
22 EventType,
23 },
24 payload::{HttpPayload, Payload, PayloadMetadata, PayloadType},
25 },
26 observability::ComponentMetricsExt as _,
27 topology::{EventsBuffer, PayloadsBuffer},
28};
29use saluki_error::{generic_error, ErrorContext as _, GenericError};
30use saluki_io::compression::{CompressionScheme, Compressor};
31use saluki_metrics::MetricsBuilder;
32use serde::Deserialize;
33use tokio::{io::AsyncWriteExt as _, select, sync::mpsc, time::sleep};
34use tracing::{debug, error, warn};
35use url::Url;
36use uuid::Uuid;
37
38#[cfg(test)]
39use self::shadow::shadow_sample_matches;
40use self::{
41 shadow::{SeriesShadowConfig, SeriesShadowState},
42 v3::{
43 V3EncodedMetrics, V3EncodedRequest, V3EncoderStats, V3MetricType, V3PayloadLimits, V3PayloadRequest,
44 V3PayloadSplitReason, V3SerializerTelemetry, V3Writer,
45 },
46};
47use crate::{
48 common::datadog::{
49 clamp_payload_limits, default_serializer_compressor_kind,
50 endpoints::{
51 calculate_resolved_endpoint, default_site, deserialize_dd_url, series_v3_config_can_enable_v3,
52 AdditionalEndpoints, EndpointV3Settings, ResolvedEndpoint, V3EndpointConfig, DEFAULT_SITE,
53 },
54 io::RB_BUFFER_CHUNK_SIZE,
55 protocol::{MetricsPayloadInfo, UseV3ApiSeriesConfig, V3ApiConfig},
56 request_builder::{RequestBuilder, RequestBuilderError},
57 resolve_zstd_compressor_level,
58 telemetry::ComponentTelemetry,
59 DEFAULT_SERIALIZER_COMPRESSED_SIZE_LIMIT, DEFAULT_SERIALIZER_UNCOMPRESSED_SIZE_LIMIT, METRICS_SERIES_V3_PATH,
60 METRICS_SKETCHES_V3_PATH,
61 },
62 encoders::datadog::metrics::v2::MetricsEndpointEncoder,
63};
64
65mod endpoint;
66use self::endpoint::{EndpointConfiguration, MetricsEndpoint};
67
68mod shadow;
69mod v1;
70mod v2;
71mod v3;
72
73const V3_SERIES_ENDPOINT_URI: &str = METRICS_SERIES_V3_PATH;
74const V3_SKETCHES_ENDPOINT_URI: &str = METRICS_SKETCHES_V3_PATH;
75
76const fn default_max_metrics_per_payload() -> usize {
77 10_000
78}
79
80const fn default_max_payload_size() -> usize {
81 DEFAULT_SERIALIZER_COMPRESSED_SIZE_LIMIT
82}
83
84const fn default_max_uncompressed_payload_size() -> usize {
85 DEFAULT_SERIALIZER_UNCOMPRESSED_SIZE_LIMIT
86}
87
88const fn default_max_series_payload_size() -> usize {
89 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT
90}
91
92const fn default_max_series_uncompressed_payload_size() -> usize {
93 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT
94}
95
96const fn default_max_series_points_per_payload() -> usize {
97 10_000
98}
99
100const fn default_flush_timeout_secs() -> u64 {
101 2
102}
103
104const fn default_use_v2_api_series() -> bool {
105 true
106}
107
108const fn default_log_payloads() -> bool {
109 false
110}
111
112fn series_shadow_config_for_endpoint(
113 series_endpoint: MetricsEndpoint, sample_rate: f64, metrics_v3_disabled_by_compressor: bool,
114) -> SeriesShadowConfig {
115 SeriesShadowConfig::new(
116 if !metrics_v3_disabled_by_compressor && series_endpoint == MetricsEndpoint::SeriesV2 {
117 sample_rate
118 } else {
119 0.0
120 },
121 )
122}
123
124#[derive(Clone, Copy, Debug, PartialEq, Eq)]
126enum MetricsEncoderMode {
127 V2Only,
129 V3Enabled,
132 Validation,
134}
135
136impl MetricsEncoderMode {
137 fn from_config(use_v3: bool, validate: bool) -> Self {
138 match (use_v3, validate) {
139 (false, _) => Self::V2Only,
140 (true, false) => Self::V3Enabled,
141 (true, true) => Self::Validation,
142 }
143 }
144
145 fn needs_v3(self) -> bool {
146 matches!(self, Self::V3Enabled | Self::Validation)
147 }
148
149 fn needs_batch_id(self) -> bool {
150 matches!(self, Self::Validation)
151 }
152
153 fn needs_tagging(self) -> bool {
154 matches!(self, Self::V3Enabled | Self::Validation)
155 }
156}
157
158fn metrics_encoder_mode_for_config(
159 use_v3: bool, validate: bool, metrics_v3_disabled_by_compressor: bool,
160) -> MetricsEncoderMode {
161 let use_v3 = use_v3 && !metrics_v3_disabled_by_compressor;
162 MetricsEncoderMode::from_config(use_v3, use_v3 && validate)
163}
164
165fn selected_metrics_primary_v3_override(
166 opw_enabled: bool, opw_url: &str, opw_use_v3_series: bool, vector_enabled: bool, vector_url: &str,
167 vector_use_v3_series: bool,
168) -> Option<bool> {
169 selected_metrics_primary_endpoint(
170 opw_enabled,
171 opw_url,
172 opw_use_v3_series,
173 vector_enabled,
174 vector_url,
175 vector_use_v3_series,
176 )
177 .map(|(_, use_v3_series)| use_v3_series)
178}
179
180fn selected_metrics_primary_endpoint<'a>(
181 opw_enabled: bool, opw_url: &'a str, opw_use_v3_series: bool, vector_enabled: bool, vector_url: &'a str,
182 vector_use_v3_series: bool,
183) -> Option<(&'a str, bool)> {
184 if opw_enabled {
185 let opw_url = opw_url.trim();
186 metrics_primary_url_can_resolve(opw_url).then_some((opw_url, opw_use_v3_series))
187 } else if vector_enabled {
188 let vector_url = vector_url.trim();
189 metrics_primary_url_can_resolve(vector_url).then_some((vector_url, vector_use_v3_series))
190 } else {
191 None
192 }
193}
194
195fn metrics_primary_url_can_resolve(url: &str) -> bool {
196 if url.is_empty() {
197 return false;
198 }
199 if url.starts_with("http://") || url.starts_with("https://") {
200 Url::parse(url).is_ok_and(|url| url.host_str().is_some())
201 } else {
202 Url::parse(&format!("https://{url}")).is_ok_and(|url| url.host_str().is_some())
203 }
204}
205
206fn series_v3_can_be_enabled_for_config(
207 use_v2_api_series: bool, serializer_use_v3_series: bool, metrics_primary_v3_override: Option<bool>,
208 has_additional_endpoints: bool, series_config: &UseV3ApiSeriesConfig,
209) -> bool {
210 use_v2_api_series
211 && (serializer_use_v3_series
212 || metrics_primary_v3_override == Some(true)
213 || ((metrics_primary_v3_override != Some(false) || has_additional_endpoints)
214 && series_v3_config_can_enable_v3(series_config)))
215}
216
217#[derive(Clone, Deserialize, Facet)]
221#[cfg_attr(test, derive(Debug, PartialEq, serde::Serialize))]
222pub struct DatadogMetricsConfiguration {
223 #[serde(
229 rename = "serializer_max_metrics_per_payload",
230 default = "default_max_metrics_per_payload"
231 )]
232 max_metrics_per_payload: usize,
233
234 #[serde(rename = "serializer_max_payload_size", default = "default_max_payload_size")]
244 max_payload_size: usize,
245
246 #[serde(
256 rename = "serializer_max_uncompressed_payload_size",
257 default = "default_max_uncompressed_payload_size"
258 )]
259 max_uncompressed_payload_size: usize,
260
261 #[serde(
271 rename = "serializer_max_series_payload_size",
272 default = "default_max_series_payload_size"
273 )]
274 max_series_payload_size: usize,
275
276 #[serde(
286 rename = "serializer_max_series_uncompressed_payload_size",
287 default = "default_max_series_uncompressed_payload_size"
288 )]
289 max_series_uncompressed_payload_size: usize,
290
291 #[serde(
300 rename = "serializer_max_series_points_per_payload",
301 default = "default_max_series_points_per_payload"
302 )]
303 max_series_points_per_payload: usize,
304
305 #[serde(default = "default_flush_timeout_secs")]
315 flush_timeout_secs: u64,
316
317 #[serde(
321 rename = "serializer_compressor_kind",
322 default = "default_serializer_compressor_kind"
323 )]
324 compressor_kind: String,
325
326 #[serde(rename = "data_plane_serializer_zstd_compressor_level", default)]
330 data_plane_zstd_compressor_level: Option<i32>,
331
332 #[serde(rename = "serializer_zstd_compressor_level", default)]
337 serializer_zstd_compressor_level: Option<i32>,
338
339 #[serde(default = "default_use_v2_api_series")]
347 use_v2_api_series: bool,
348
349 #[serde(default = "default_log_payloads")]
355 log_payloads: bool,
356
357 #[serde(default, skip)]
359 #[facet(opaque)]
360 additional_tags: Option<SharedTagSet>,
361
362 #[serde(rename = "serializer_experimental_use_v3_api", default)]
366 v3_api: V3ApiConfig,
367
368 #[serde(flatten)]
370 use_v3_api_series: UseV3ApiSeriesConfig,
371
372 #[serde(default, rename = "observability_pipelines_worker_metrics_enabled")]
374 observability_pipelines_worker_metrics_enabled: bool,
375
376 #[serde(default, rename = "observability_pipelines_worker_metrics_url")]
378 observability_pipelines_worker_metrics_url: String,
379
380 #[serde(default, rename = "observability_pipelines_worker_metrics_use_v3_api_series")]
384 observability_pipelines_worker_metrics_use_v3_api_series: bool,
385
386 #[serde(default, rename = "vector_metrics_enabled")]
388 vector_metrics_enabled: bool,
389
390 #[serde(default, rename = "vector_metrics_url")]
392 vector_metrics_url: String,
393
394 #[serde(default, rename = "vector_metrics_use_v3_api_series")]
400 vector_metrics_use_v3_api_series: bool,
401
402 #[serde(default = "default_site")]
406 site: String,
407
408 #[serde(default, alias = "url", deserialize_with = "deserialize_dd_url")]
412 dd_url: Option<String>,
413
414 #[serde(default)]
416 additional_endpoints: AdditionalEndpoints,
417}
418
419impl DatadogMetricsConfiguration {
420 pub fn from_configuration(config: &GenericConfiguration) -> Result<Self, GenericError> {
422 Ok(config.as_typed()?)
423 }
424
425 pub fn from_configuration_with_metrics_routing(
427 config: &GenericConfiguration, metrics: &MetricsEncoding, endpoints: &Endpoints,
428 ) -> Result<Self, GenericError> {
429 let mut metrics_config = Self::from_configuration(config)?;
430
431 metrics_config.compressor_kind = endpoints.compression.compressor_kind.clone();
432 metrics_config.use_v2_api_series = metrics.use_v2_series_api;
433 metrics_config.v3_api = (&metrics.v3_api).into();
434 metrics_config.use_v3_api_series = (&metrics.v3_series_mode).into();
435 metrics_config.observability_pipelines_worker_metrics_enabled = endpoints.opw_intake.enabled;
436 metrics_config.observability_pipelines_worker_metrics_url = endpoints.opw_intake.url.clone();
437 metrics_config.observability_pipelines_worker_metrics_use_v3_api_series = endpoints.opw_intake.use_v3_series;
438 metrics_config.vector_metrics_enabled = endpoints.vector_intake.enabled;
439 metrics_config.vector_metrics_url = endpoints.vector_intake.url.clone();
440 metrics_config.vector_metrics_use_v3_api_series = endpoints.vector_intake.use_v3_series;
441
442 Ok(metrics_config)
443 }
444
445 pub fn with_additional_tags(mut self, additional_tags: SharedTagSet) -> Self {
447 self.additional_tags = Some(additional_tags);
448 self
449 }
450
451 pub fn with_metrics_endpoint_override(mut self, dd_url: String) -> Self {
456 self.dd_url = Some(dd_url);
457 self.additional_endpoints = AdditionalEndpoints::default();
458 self.observability_pipelines_worker_metrics_enabled = false;
459 self.vector_metrics_enabled = false;
460 self
461 }
462
463 pub fn with_v2_series_only(mut self) -> Self {
467 self.use_v3_api_series.enabled = "false".to_string();
468 self.use_v3_api_series.endpoints.clear();
469 self.v3_api.series.endpoints.clear();
470 self.v3_api.series.shadow_sample_rate = 0.0;
471 self.observability_pipelines_worker_metrics_use_v3_api_series = false;
472 self.vector_metrics_use_v3_api_series = false;
473 self
474 }
475
476 fn v3_payload_limits(&self) -> V3PayloadLimits {
477 V3PayloadLimits::new(
478 self.max_series_payload_size,
479 self.max_series_uncompressed_payload_size,
480 self.max_metrics_per_payload,
481 self.max_series_points_per_payload,
482 )
483 }
484
485 fn endpoint_v3_settings(
486 &self, endpoint: &ResolvedEndpoint, metrics_primary_v3_override: Option<bool>,
487 serializer_v3_configured_endpoint: Option<&str>,
488 ) -> EndpointV3Settings {
489 EndpointV3Settings::from_v3_config(V3EndpointConfig {
490 configured_endpoint: endpoint.configured_endpoint(),
491 resolved_endpoint: endpoint.endpoint(),
492 serializer_v3_configured_endpoint,
493 series_config: &self.use_v3_api_series,
494 metrics_primary_v3_override,
495 serializer_v3_series_endpoints: &self.v3_api.series.endpoints,
496 serializer_v3_sketches_endpoints: &self.v3_api.sketches.endpoints,
497 series_validate: self.v3_api.series.validate,
498 sketches_validate: self.v3_api.sketches.validate,
499 series_shadow_sites: &self.v3_api.series.shadow_sites,
500 })
501 }
502
503 fn configured_primary_endpoint(&self) -> String {
504 match self.dd_url.as_deref() {
505 Some(url) => url.to_string(),
506 None => {
507 let base_domain = if self.site.is_empty() { DEFAULT_SITE } else { &self.site };
508 format!("https://app.{base_domain}")
509 }
510 }
511 }
512
513 fn any_series_endpoint_matches(
514 &self, mut predicate: impl FnMut(&EndpointV3Settings) -> bool,
515 ) -> Result<bool, GenericError> {
516 let configured_primary_endpoint = self.configured_primary_endpoint();
517 if let Some((metrics_primary_url, metrics_primary_v3_override)) = selected_metrics_primary_endpoint(
518 self.observability_pipelines_worker_metrics_enabled,
519 &self.observability_pipelines_worker_metrics_url,
520 self.observability_pipelines_worker_metrics_use_v3_api_series,
521 self.vector_metrics_enabled,
522 &self.vector_metrics_url,
523 self.vector_metrics_use_v3_api_series,
524 ) {
525 let metrics_primary = calculate_resolved_endpoint(Some(metrics_primary_url), &self.site, "")
526 .error_context("Failed parsing/resolving the metrics primary destination endpoint.")?;
527 let settings = self.endpoint_v3_settings(
528 &metrics_primary,
529 Some(metrics_primary_v3_override),
530 Some(&configured_primary_endpoint),
531 );
532 if predicate(&settings) {
533 return Ok(true);
534 }
535 } else {
536 let primary = calculate_resolved_endpoint(self.dd_url.as_deref(), &self.site, "")
537 .error_context("Failed parsing/resolving the primary destination endpoint.")?;
538 let settings = self.endpoint_v3_settings(&primary, None, None);
539 if predicate(&settings) {
540 return Ok(true);
541 }
542 }
543
544 for endpoint in self
545 .additional_endpoints
546 .resolved_endpoints(None)
547 .error_context("Failed parsing/resolving the additional destination endpoints.")?
548 {
549 let settings = self.endpoint_v3_settings(&endpoint, None, None);
550 if predicate(&settings) {
551 return Ok(true);
552 }
553 }
554
555 Ok(false)
556 }
557
558 fn requires_v2_series(&self, metrics_v3_disabled_by_compressor: bool) -> Result<bool, GenericError> {
559 if !self.use_v2_api_series || metrics_v3_disabled_by_compressor {
560 return Ok(true);
561 }
562
563 self.any_series_endpoint_matches(|settings| !settings.use_v3_series || settings.series_validation_mode)
564 }
565
566 fn requires_v3_series(&self, metrics_v3_disabled_by_compressor: bool) -> Result<bool, GenericError> {
567 if metrics_v3_disabled_by_compressor {
568 return Ok(false);
569 }
570
571 let metrics_primary_v3_override = selected_metrics_primary_v3_override(
572 self.observability_pipelines_worker_metrics_enabled,
573 &self.observability_pipelines_worker_metrics_url,
574 self.observability_pipelines_worker_metrics_use_v3_api_series,
575 self.vector_metrics_enabled,
576 &self.vector_metrics_url,
577 self.vector_metrics_use_v3_api_series,
578 );
579 if !series_v3_can_be_enabled_for_config(
580 self.use_v2_api_series,
581 self.v3_api.use_v3_series(),
582 metrics_primary_v3_override,
583 !self.additional_endpoints.is_empty(),
584 &self.use_v3_api_series,
585 ) {
586 return Ok(false);
587 }
588
589 self.any_series_endpoint_matches(|settings| settings.use_v3_series)
590 }
591}
592
593#[async_trait]
594impl EncoderBuilder for DatadogMetricsConfiguration {
595 fn input_event_type(&self) -> EventType {
596 EventType::Metric
597 }
598
599 fn output_payload_type(&self) -> PayloadType {
600 PayloadType::Http
601 }
602
603 async fn build(&self, context: ComponentContext) -> Result<Box<dyn Encoder + Send>, GenericError> {
604 let metrics_builder = MetricsBuilder::from_component_context(&context);
605 let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
606 let v3_serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
607
608 let zstd_compressor_level = resolve_zstd_compressor_level(
609 self.data_plane_zstd_compressor_level,
610 self.serializer_zstd_compressor_level,
611 );
612 let v2_compression_scheme = CompressionScheme::new(&self.compressor_kind, zstd_compressor_level);
613 let v3_compression_scheme = if self.v3_api.compression_level > 0 {
614 CompressionScheme::new(&self.compressor_kind, self.v3_api.compression_level)
615 } else {
616 v2_compression_scheme
617 };
618 let series_endpoint_uri = if self.v3_api.series.use_beta {
619 self.v3_api.series.beta_route.clone()
620 } else {
621 V3_SERIES_ENDPOINT_URI.to_string()
622 };
623 let shadow_series_endpoint_uri = self.v3_api.series.beta_route.clone();
624 let payload_limits = self.v3_payload_limits();
625
626 let v2_endpoint_config = EndpointConfiguration::new(
627 v2_compression_scheme,
628 self.max_metrics_per_payload,
629 self.max_series_points_per_payload,
630 self.additional_tags.clone(),
631 );
632 let endpoint_config = EndpointConfiguration::new(
633 v3_compression_scheme,
634 self.max_metrics_per_payload,
635 usize::MAX,
638 self.additional_tags.clone(),
639 );
640
641 let metrics_v3_disabled_by_compressor = matches!(v3_compression_scheme, CompressionScheme::Zlib(_));
643 let use_v3_series = self.requires_v3_series(metrics_v3_disabled_by_compressor)?;
644 let series_mode = metrics_encoder_mode_for_config(
645 use_v3_series,
646 self.v3_api.series.validate,
647 metrics_v3_disabled_by_compressor,
648 );
649 let sketches_mode = metrics_encoder_mode_for_config(
650 self.v3_api.use_v3_sketches(),
651 self.v3_api.sketches.validate,
652 metrics_v3_disabled_by_compressor,
653 );
654 let series_endpoint = if self.use_v2_api_series {
655 MetricsEndpoint::SeriesV2
656 } else {
657 MetricsEndpoint::SeriesV1
658 };
659 let series_shadow_config = series_shadow_config_for_endpoint(
660 series_endpoint,
661 self.v3_api.series.shadow_sample_rate,
662 metrics_v3_disabled_by_compressor,
663 );
664 let v3_runtime_config = V3RuntimeConfig {
665 endpoint_config,
666 payload_limits,
667 series_endpoint_uri,
668 shadow_series_endpoint_uri,
669 series_shadow_config,
670 serializer_telemetry: v3_serializer_telemetry,
671 };
672 let generic_payload_limits = clamp_payload_limits(
673 self.max_uncompressed_payload_size,
674 self.max_payload_size,
675 DEFAULT_SERIALIZER_UNCOMPRESSED_SIZE_LIMIT,
676 DEFAULT_SERIALIZER_COMPRESSED_SIZE_LIMIT,
677 );
678 let (series_uncompressed_limit, series_compressed_limit) = if series_endpoint == MetricsEndpoint::SeriesV2 {
679 clamp_payload_limits(
680 self.max_series_uncompressed_payload_size,
681 self.max_series_payload_size,
682 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT,
683 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT,
684 )
685 } else {
686 generic_payload_limits
687 };
688 let v2_series_builder = if self.requires_v2_series(metrics_v3_disabled_by_compressor)? {
689 let mut builder = v2::create_v2_request_builder(series_endpoint, &v2_endpoint_config)
690 .await
691 .error_context("Failed to create V2 series request builder.")?;
692 builder.with_len_limits(series_uncompressed_limit, series_compressed_limit)?;
693 Some(builder)
694 } else {
695 debug!("All metrics series endpoints use authoritative V3; disabling V2 series encoding.");
696 None
697 };
698
699 let (sketches_uncompressed_limit, sketches_compressed_limit) = generic_payload_limits;
700 let mut v2_sketch_builder = v2::create_v2_request_builder(MetricsEndpoint::Sketches, &v2_endpoint_config)
701 .await
702 .error_context("Failed to create V2 sketches request builder.")?;
703 v2_sketch_builder.with_len_limits(sketches_uncompressed_limit, sketches_compressed_limit)?;
704 let v2_sketch_builder = Some(v2_sketch_builder);
705
706 let flush_timeout = match self.flush_timeout_secs {
707 0 => Duration::from_millis(10),
710 secs => Duration::from_secs(secs),
711 };
712
713 if series_mode.needs_v3() || sketches_mode.needs_v3() {
714 debug!(
715 ?series_mode,
716 ?sketches_mode,
717 v3_series_endpoints = ?self.v3_api.series.endpoints,
718 v3_sketches_endpoints = ?self.v3_api.sketches.endpoints,
719 "V3 encoding support is enabled."
720 );
721 }
722
723 Ok(Box::new(DatadogMetrics {
724 v2_series_builder,
725 v2_sketch_builder,
726 series_mode,
727 sketches_mode,
728 v3_runtime_config,
729 telemetry,
730 flush_timeout,
731 log_payloads: self.log_payloads,
732 }))
733 }
734}
735
736impl MemoryBounds for DatadogMetricsConfiguration {
737 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
738 builder
746 .minimum()
747 .with_single_value::<DatadogMetrics>("component struct")
748 .with_array::<EventsBuffer>("request builder events channel", 8)
749 .with_array::<PayloadsBuffer>("request builder payloads channel", 8);
750
751 builder
752 .firm()
753 .with_array::<Metric>("series metrics split re-encode buffer", self.max_metrics_per_payload)
756 .with_array::<Metric>("sketch metrics split re-encode buffer", self.max_metrics_per_payload);
757 }
758}
759
760pub struct DatadogMetrics {
761 v2_series_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>,
762 v2_sketch_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>,
763 series_mode: MetricsEncoderMode,
764 sketches_mode: MetricsEncoderMode,
765 v3_runtime_config: V3RuntimeConfig,
766 telemetry: ComponentTelemetry,
767 flush_timeout: Duration,
768 log_payloads: bool,
769}
770
771struct V3RuntimeConfig {
772 endpoint_config: EndpointConfiguration,
773 payload_limits: V3PayloadLimits,
774 series_endpoint_uri: String,
775 shadow_series_endpoint_uri: String,
776 series_shadow_config: SeriesShadowConfig,
777 serializer_telemetry: V3SerializerTelemetry,
778}
779
780#[async_trait]
781impl Encoder for DatadogMetrics {
782 async fn run(mut self: Box<Self>, mut context: EncoderContext) -> Result<(), GenericError> {
783 let Self {
784 v2_series_builder,
785 v2_sketch_builder,
786 series_mode,
787 sketches_mode,
788 v3_runtime_config,
789 telemetry,
790 flush_timeout,
791 log_payloads,
792 } = *self;
793
794 let mut health = context.take_health_handle();
795
796 let (events_tx, events_rx) = mpsc::channel(8);
798 let (payloads_tx, mut payloads_rx) = mpsc::channel(8);
799 let request_builder_fut = run_request_builder(
800 v2_series_builder,
801 v2_sketch_builder,
802 series_mode,
803 sketches_mode,
804 v3_runtime_config,
805 telemetry,
806 events_rx,
807 payloads_tx,
808 flush_timeout,
809 log_payloads,
810 );
811 let request_builder_handle = context
812 .topology_context()
813 .global_thread_pool()
814 .spawn_traced_named("dd-metrics-request-builder", request_builder_fut);
815
816 health.mark_ready();
817 debug!("Datadog Metrics encoder started.");
818
819 loop {
820 select! {
821 biased;
822
823 _ = health.live() => continue,
824 maybe_payload = payloads_rx.recv() => match maybe_payload {
825 Some(payload) => {
826 if let Err(e) = context.dispatcher().dispatch(payload).await {
827 error!("Failed to dispatch payload: {}", e);
828 }
829 }
830 None => break,
831 },
832 maybe_event_buffer = context.events().next() => match maybe_event_buffer {
833 Some(event_buffer) => {
834 let permit = loop {
841 select! {
842 biased;
843
844 permit = events_tx.reserve() => break permit
845 .error_context("Failed to reserve capacity for event buffer.")?,
846 maybe_payload = payloads_rx.recv() => match maybe_payload {
847 Some(payload) => if let Err(e) = context.dispatcher().dispatch(payload).await {
848 error!("Failed to dispatch payload: {}", e);
849 },
850
851 None => return Err(generic_error!("Request builder task stopped before accepting event buffer.")),
853 }
854 }
855 };
856
857 permit.send(event_buffer);
858 }
859 None => break,
860 },
861 }
862 }
863
864 drop(events_tx);
866
867 while let Some(payload) = payloads_rx.recv().await {
869 if let Err(e) = context.dispatcher().dispatch(payload).await {
870 error!("Failed to dispatch payload: {}", e);
871 }
872 }
873
874 match request_builder_handle.await {
876 Ok(Ok(())) => debug!("Request builder task stopped."),
877 Ok(Err(e)) => error!(error = %e, "Request builder task failed."),
878 Err(e) => error!(error = %e, "Request builder task panicked."),
879 }
880
881 debug!("Datadog Metrics encoder stopped.");
882
883 Ok(())
884 }
885}
886
887fn log_metric_payload(metric: &Metric) {
891 match metric.values() {
892 MetricValues::Counter(..) | MetricValues::Rate(..) | MetricValues::Gauge(..) | MetricValues::Set(..) => {
893 debug!(?metric, "Flushing series metric.")
894 }
895 MetricValues::Histogram(..) | MetricValues::Distribution(..) => {
896 debug!(?metric, "Flushing sketch metric.")
897 }
898 }
899}
900
901#[allow(clippy::too_many_arguments)]
902async fn run_request_builder(
903 mut v2_series_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>,
904 mut v2_sketch_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>, series_mode: MetricsEncoderMode,
905 sketches_mode: MetricsEncoderMode, v3_runtime_config: V3RuntimeConfig, telemetry: ComponentTelemetry,
906 mut events_rx: mpsc::Receiver<EventsBuffer>, mut payloads_tx: mpsc::Sender<PayloadsBuffer>,
907 flush_timeout: Duration, log_payloads: bool,
908) -> Result<(), GenericError> {
909 let mut pending_flush = false;
910 let pending_flush_timeout = sleep(flush_timeout);
911 tokio::pin!(pending_flush_timeout);
912
913 let mut v3_series_metrics =
914 (series_mode.needs_v3() || v3_runtime_config.series_shadow_config.is_enabled()).then(Vec::<Metric>::new);
915 let mut v3_sketch_metrics = sketches_mode.needs_v3().then(Vec::<Metric>::new);
916 let mut v3_series_points = 0usize;
917 let mut v3_sketch_points = 0usize;
918 let mut v3_series_ratio = V3CompressionRatio::default();
919 let mut v3_sketch_ratio = V3CompressionRatio::default();
920
921 let mut series_batch_id = None;
922 let mut sketches_batch_id = None;
923 let mut series_shadow_state = SeriesShadowState::default();
924 let series_shadow_config = v3_runtime_config.series_shadow_config;
925
926 let tag_series = series_mode.needs_tagging();
927 let tag_sketches = sketches_mode.needs_tagging();
928 let v3_flush_context = V3FlushContext {
929 endpoint_config: &v3_runtime_config.endpoint_config,
930 payload_limits: v3_runtime_config.payload_limits,
931 series_endpoint_uri: &v3_runtime_config.series_endpoint_uri,
932 serializer_telemetry: &v3_runtime_config.serializer_telemetry,
933 telemetry: &telemetry,
934 };
935 let v3_shadow_flush_context = V3FlushContext {
936 endpoint_config: &v3_runtime_config.endpoint_config,
937 payload_limits: v3_runtime_config.payload_limits,
938 series_endpoint_uri: &v3_runtime_config.shadow_series_endpoint_uri,
939 serializer_telemetry: &v3_runtime_config.serializer_telemetry,
940 telemetry: &telemetry,
941 };
942
943 loop {
944 select! {
945 Some(event_buffer) = events_rx.recv() => {
946 for event in event_buffer {
947 let metric = match event.try_into_metric() {
948 Some(metric) => metric,
949 None => continue,
950 };
951
952 if log_payloads {
953 log_metric_payload(&metric);
954 }
955
956 if !v1::has_emittable_point(&metric) {
959 debug!(metric = %metric.context().name(), "Dropping series metric with no finite points.");
960 telemetry.events_dropped_encoder().increment(1);
961 continue;
962 }
963
964 let endpoint = MetricsEndpoint::from_metric(&metric);
966 let (endpoint_mode, maybe_v2_builder, maybe_v3_metrics, v3_points, batch_id) = match endpoint {
967 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => (
968 series_mode,
969 &mut v2_series_builder,
970 &mut v3_series_metrics,
971 &mut v3_series_points,
972 &mut series_batch_id,
973 ),
974 MetricsEndpoint::Sketches => (
975 sketches_mode,
976 &mut v2_sketch_builder,
977 &mut v3_sketch_metrics,
978 &mut v3_sketch_points,
979 &mut sketches_batch_id,
980 ),
981 };
982 let metric_point_count = metric.values().len();
983 let is_series = matches!(endpoint, MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2);
984 let series_shadow_active = is_series
985 && matches!(endpoint_mode, MetricsEncoderMode::V2Only)
986 && series_shadow_state.ensure_decision(series_shadow_config);
987 let needs_batch_id = endpoint_mode.needs_batch_id() || series_shadow_active;
988 if needs_batch_id && batch_id.is_none() {
989 *batch_id = Some(Uuid::now_v7());
990 }
991 let active_batch_id = needs_batch_id.then_some(batch_id.as_ref()).flatten();
992 let should_buffer_v3 = endpoint_mode.needs_v3() || series_shadow_active;
993
994 if should_buffer_v3 {
999 if let Some(metrics) = maybe_v3_metrics {
1000 metrics.push(metric.clone());
1001 *v3_points += metric_point_count;
1002 }
1003 }
1004
1005 let v2_payload_info = match endpoint {
1012 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => {
1013 if series_shadow_active {
1014 Some(MetricsPayloadInfo::v2_shadow_series())
1015 } else {
1016 tag_series.then(MetricsPayloadInfo::v2_series)
1017 }
1018 }
1019 MetricsEndpoint::Sketches => tag_sketches.then(MetricsPayloadInfo::v2_sketches),
1020 };
1021 let metric_for_next_shadow_batch = (is_series
1024 && matches!(endpoint_mode, MetricsEncoderMode::V2Only)
1025 && series_shadow_config.is_enabled()
1026 && !series_shadow_active)
1027 .then(|| metric.clone());
1028 let mut v2_encoded = false;
1029 let v2_flushed = if let Some(builder) = maybe_v2_builder {
1030 let result =
1031 encode_v2_metrics(builder, metric, &telemetry, &mut payloads_tx, active_batch_id, v2_payload_info).await?;
1032 v2_encoded = result.encoded();
1033 if should_buffer_v3
1034 && !result.encoded()
1035 && !matches!(endpoint_mode, MetricsEncoderMode::V3Enabled)
1036 {
1037 if let Some(metrics) = maybe_v3_metrics {
1038 let _ = metrics.pop();
1039 *v3_points = v3_points.saturating_sub(metric_point_count);
1040 }
1041 }
1042
1043 result.flushed()
1044 } else {
1045 false
1046 };
1047
1048 let v3_payload_info = match endpoint {
1051 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => {
1052 if series_shadow_active {
1053 Some(MetricsPayloadInfo::v3_shadow_series())
1054 } else {
1055 tag_series.then(MetricsPayloadInfo::v3_series)
1056 }
1057 }
1058 MetricsEndpoint::Sketches => tag_sketches.then(MetricsPayloadInfo::v3_sketches),
1059 };
1060 let mut split_metric = None;
1061 let v3_flushed = if let Some(v3_metrics) = maybe_v3_metrics {
1062 let should_flush_v3 = match endpoint_mode {
1063 MetricsEncoderMode::V2Only => series_shadow_active && v2_flushed,
1064 MetricsEncoderMode::V3Enabled => {
1065 if v3_flush_context.payload_limits.point_count_fits(metric_point_count)
1066 && v3_flush_context.payload_limits.point_count_exceeds_limit(*v3_points)
1067 && v3_metrics.len() > 1
1068 {
1069 v3_flush_context
1070 .serializer_telemetry
1071 .record_split_reason(V3PayloadSplitReason::MaxPoints);
1072 if let Some(metric) = v3_metrics.pop() {
1073 *v3_points = v3_points.saturating_sub(metric.values().len());
1074 split_metric = Some(metric);
1075 }
1076 true
1077 } else {
1078 v3_flush_context.payload_limits.should_flush_point_count_limit(*v3_points)
1079 || v3_flush_context
1080 .payload_limits
1081 .should_flush_metric_count_limit(v3_metrics)
1082 }
1083 }
1084 MetricsEncoderMode::Validation => v2_flushed,
1085 };
1086 if should_flush_v3 {
1087 if !matches!(endpoint_mode, MetricsEncoderMode::V3Enabled) {
1088 if v2_flushed {
1091 if let Some(metric) = v3_metrics.pop() {
1092 *v3_points = v3_points.saturating_sub(metric.values().len());
1093 split_metric = Some(metric);
1094 }
1095 }
1096 }
1097 let flush_context = if series_shadow_active {
1098 v3_shadow_flush_context
1099 } else {
1100 v3_flush_context
1101 };
1102 encode_and_flush_v3_metrics(
1103 endpoint,
1104 flush_context,
1105 v3_metrics,
1106 &mut v3_series_ratio,
1107 &mut v3_sketch_ratio,
1108 &mut payloads_tx,
1109 active_batch_id,
1110 v3_payload_info,
1111 )
1112 .await?;
1113 *v3_points = 0;
1114 true
1115 } else {
1116 false
1117 }
1118 } else {
1119 false
1120 };
1121
1122 if matches!(endpoint_mode, MetricsEncoderMode::V3Enabled) {
1123 if let Some(m) = split_metric.take() {
1126 let point_count = m.values().len();
1127 if let Some(metrics) = maybe_v3_metrics {
1128 metrics.push(m);
1129 *v3_points += point_count;
1130 }
1131 }
1132 } else if endpoint_mode.needs_batch_id() && (v2_flushed || v3_flushed) {
1133 *batch_id = if let Some(m) = split_metric.take() {
1136 let point_count = m.values().len();
1137 if let Some(metrics) = maybe_v3_metrics {
1138 metrics.push(m);
1139 *v3_points += point_count;
1140 }
1141 Some(Uuid::now_v7())
1142 } else {
1143 None
1144 };
1145 } else if is_series && series_shadow_active && (v2_flushed || v3_flushed) {
1146 series_shadow_state.reset();
1149 *batch_id = if let Some(m) = split_metric.take() {
1150 if series_shadow_state.ensure_decision(series_shadow_config) {
1151 let point_count = m.values().len();
1152 if let Some(metrics) = maybe_v3_metrics {
1153 metrics.push(m);
1154 *v3_points += point_count;
1155 }
1156 Some(Uuid::now_v7())
1157 } else {
1158 None
1159 }
1160 } else {
1161 None
1162 };
1163 } else if is_series
1164 && matches!(endpoint_mode, MetricsEncoderMode::V2Only)
1165 && series_shadow_config.is_enabled()
1166 && v2_flushed
1167 {
1168 series_shadow_state.reset();
1171 *batch_id = None;
1172 if v2_encoded {
1173 if let Some(m) = metric_for_next_shadow_batch {
1174 if series_shadow_state.ensure_decision(series_shadow_config) {
1175 let point_count = m.values().len();
1176 if let Some(metrics) = maybe_v3_metrics {
1177 metrics.push(m);
1178 *v3_points += point_count;
1179 }
1180 *batch_id = Some(Uuid::now_v7());
1181 }
1182 }
1183 }
1184 }
1185 }
1186
1187 debug!("Processed event buffer.");
1188
1189 if !pending_flush {
1191 pending_flush_timeout.as_mut().reset(tokio::time::Instant::now() + flush_timeout);
1192 pending_flush = true;
1193 }
1194 },
1195 _ = &mut pending_flush_timeout, if pending_flush => {
1196 debug!("Flushing pending request(s).");
1197
1198 pending_flush = false;
1199
1200 let series_shadow_active = matches!(series_mode, MetricsEncoderMode::V2Only)
1204 && series_shadow_state.active.unwrap_or(false);
1205 let v2_series_payload_info = if series_shadow_active {
1206 Some(MetricsPayloadInfo::v2_shadow_series())
1207 } else {
1208 tag_series.then(MetricsPayloadInfo::v2_series)
1209 };
1210 let series_active_batch_id = (series_mode.needs_batch_id() || series_shadow_active)
1211 .then_some(series_batch_id.as_ref())
1212 .flatten();
1213 let mut v2_series_flush_succeeded = true;
1214 if let Some(builder) = &mut v2_series_builder {
1215 if let Err(e) = flush_v2_metrics(builder, &mut payloads_tx, series_active_batch_id, v2_series_payload_info).await {
1216 error!(error = %e, "Failed to flush V2 series metrics: {}", e);
1217 v2_series_flush_succeeded = false;
1218 }
1219 }
1220
1221 let v3_series_payload_info = if series_shadow_active {
1222 Some(MetricsPayloadInfo::v3_shadow_series())
1223 } else {
1224 tag_series.then(MetricsPayloadInfo::v3_series)
1225 };
1226 if let Some(metrics) = &mut v3_series_metrics {
1227 if v2_series_flush_succeeded || matches!(series_mode, MetricsEncoderMode::V3Enabled) {
1228 let flush_context = if series_shadow_active {
1230 v3_shadow_flush_context
1231 } else {
1232 v3_flush_context
1233 };
1234 if let Err(e) = encode_and_flush_v3_series_metrics(
1235 flush_context,
1236 metrics,
1237 &mut v3_series_ratio,
1238 &mut payloads_tx,
1239 series_active_batch_id,
1240 v3_series_payload_info,
1241 )
1242 .await
1243 {
1244 error!(error = %e, "Failed to flush V3 series metrics: {}", e);
1245 }
1246 v3_series_points = 0;
1247 } else {
1248 warn!("Failed to flush V2 series metrics, skipping V3 series flush.");
1250 metrics.clear();
1251 v3_series_points = 0;
1252 }
1253 }
1254 if series_mode.needs_batch_id() {
1255 series_batch_id = None;
1256 }
1257 if matches!(series_mode, MetricsEncoderMode::V2Only) && series_shadow_config.is_enabled() {
1258 series_shadow_state.reset();
1259 series_batch_id = None;
1260 }
1261
1262 let v2_sketches_payload_info = tag_sketches.then(MetricsPayloadInfo::v2_sketches);
1264 let sketches_active_batch_id = sketches_mode.needs_batch_id().then_some(sketches_batch_id.as_ref()).flatten();
1265 let mut v2_sketches_flush_succeeded = true;
1266 if let Some(builder) = &mut v2_sketch_builder {
1267 if let Err(e) = flush_v2_metrics(builder, &mut payloads_tx, sketches_active_batch_id, v2_sketches_payload_info).await {
1268 error!(error = %e, "Failed to flush V2 sketch metrics: {}", e);
1269 v2_sketches_flush_succeeded = false;
1270 }
1271 }
1272
1273 let v3_sketches_payload_info = tag_sketches.then(MetricsPayloadInfo::v3_sketches);
1274 if let Some(metrics) = &mut v3_sketch_metrics {
1275 if v2_sketches_flush_succeeded || matches!(sketches_mode, MetricsEncoderMode::V3Enabled) {
1276 if let Err(e) = encode_and_flush_v3_sketch_metrics(
1277 v3_flush_context,
1278 metrics,
1279 &mut v3_sketch_ratio,
1280 &mut payloads_tx,
1281 sketches_active_batch_id,
1282 v3_sketches_payload_info,
1283 )
1284 .await
1285 {
1286 error!(error = %e, "Failed to flush V3 sketch metrics: {}", e);
1287 }
1288 v3_sketch_points = 0;
1289 } else {
1290 warn!("Failed to flush V2 sketch metrics, skipping V3 sketch flush.");
1291 metrics.clear();
1292 v3_sketch_points = 0;
1293 }
1294 }
1295 if sketches_mode.needs_batch_id() {
1296 sketches_batch_id = None;
1297 }
1298
1299 debug!("All flushed requests sent to I/O task. Waiting for next event buffer...");
1300 },
1301
1302 else => break,
1304 }
1305 }
1306
1307 Ok(())
1308}
1309
1310struct EncodeResult {
1311 encoded: bool,
1312 flushed: bool,
1313}
1314
1315impl EncodeResult {
1316 pub const fn new(encoded: bool, flushed: bool) -> Self {
1317 Self { encoded, flushed }
1318 }
1319
1320 pub const fn encoded(&self) -> bool {
1321 self.encoded
1322 }
1323
1324 pub const fn flushed(&self) -> bool {
1325 self.flushed
1326 }
1327}
1328
1329async fn encode_v2_metrics(
1330 request_builder: &mut RequestBuilder<v2::MetricsEndpointEncoder>, metric: Metric, telemetry: &ComponentTelemetry,
1331 payloads_tx: &mut mpsc::Sender<Payload>, batch_id: Option<&Uuid>, payload_info: Option<MetricsPayloadInfo>,
1332) -> Result<EncodeResult, GenericError> {
1333 let metric_to_retry = match request_builder.encode(metric).await {
1337 Ok(None) => return Ok(EncodeResult::new(true, false)),
1338 Ok(Some(metric)) => metric,
1339 Err(RequestBuilderError::InvalidInput { input }) => {
1340 debug!(metric_name = %input.context().name(), "Dropping metric with no emittable values.");
1341 telemetry.events_dropped_encoder().increment(1);
1342 return Ok(EncodeResult::new(false, false));
1343 }
1344 Err(e) => {
1345 error!(error = %e, "Failed to encode metric.");
1346 telemetry.events_dropped_encoder().increment(1);
1347 return Ok(EncodeResult::new(false, false));
1348 }
1349 };
1350
1351 flush_v2_metrics(request_builder, payloads_tx, batch_id, payload_info).await?;
1352
1353 match request_builder.encode(metric_to_retry).await {
1357 Ok(None) => Ok(EncodeResult::new(true, true)),
1358 Ok(Some(_)) => unreachable!(
1359 "failure to encode due to size should never occur after flush for metrics which aren't unencodable"
1360 ),
1361 Err(e) => {
1362 error!(error = %e, "Failed to encode metric.");
1363 telemetry.events_dropped_encoder().increment(1);
1364 Ok(EncodeResult::new(false, true))
1365 }
1366 }
1367}
1368
1369async fn flush_v2_metrics(
1370 request_builder: &mut RequestBuilder<MetricsEndpointEncoder>, payloads_tx: &mut mpsc::Sender<Payload>,
1371 batch_id: Option<&Uuid>, payload_info: Option<MetricsPayloadInfo>,
1372) -> Result<usize, GenericError> {
1373 let mut requests_flushed = 0;
1374
1375 let maybe_requests = request_builder.flush().await;
1376 let batch_len = maybe_requests.len();
1377 for (batch_seq, maybe_request) in maybe_requests.into_iter().enumerate() {
1378 match maybe_request {
1379 Ok((events, data_points, request)) => {
1380 requests_flushed += 1;
1381
1382 flush_payload(
1383 request,
1384 events,
1385 data_points,
1386 payloads_tx,
1387 batch_id,
1388 batch_seq,
1389 batch_len,
1390 payload_info,
1391 )
1392 .await?;
1393 }
1394
1395 Err(e) => {
1397 if !e.is_recoverable() {
1398 return Err(GenericError::from(e).context("Failed to flush request."));
1399 }
1400 }
1401 }
1402 }
1403
1404 Ok(requests_flushed)
1405}
1406
1407#[derive(Clone, Copy)]
1408struct V3FlushContext<'a> {
1409 endpoint_config: &'a EndpointConfiguration,
1410 payload_limits: V3PayloadLimits,
1411 series_endpoint_uri: &'a str,
1412 serializer_telemetry: &'a V3SerializerTelemetry,
1413 telemetry: &'a ComponentTelemetry,
1414}
1415
1416#[allow(clippy::too_many_arguments)]
1417async fn encode_and_flush_v3_metrics(
1418 endpoint: MetricsEndpoint, context: V3FlushContext<'_>, metrics: &mut Vec<Metric>,
1419 series_ratio: &mut V3CompressionRatio, sketches_ratio: &mut V3CompressionRatio,
1420 payloads_tx: &mut mpsc::Sender<Payload>, batch_id: Option<&Uuid>, payload_info: Option<MetricsPayloadInfo>,
1421) -> Result<(), GenericError> {
1422 match endpoint {
1423 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => {
1424 encode_and_flush_v3_series_metrics(context, metrics, series_ratio, payloads_tx, batch_id, payload_info)
1425 .await
1426 }
1427 MetricsEndpoint::Sketches => {
1428 encode_and_flush_v3_sketch_metrics(context, metrics, sketches_ratio, payloads_tx, batch_id, payload_info)
1429 .await
1430 }
1431 }
1432}
1433
1434async fn encode_and_flush_v3_series_metrics(
1435 context: V3FlushContext<'_>, metrics: &mut Vec<Metric>, ratio: &mut V3CompressionRatio,
1436 payloads_tx: &mut mpsc::Sender<Payload>, batch_id: Option<&Uuid>, payload_info: Option<MetricsPayloadInfo>,
1437) -> Result<(), GenericError> {
1438 if metrics.is_empty() {
1439 return Ok(());
1440 }
1441 let metrics_to_flush = std::mem::take(metrics);
1442
1443 encode_and_flush_v3_payload_requests(
1444 context.series_endpoint_uri,
1445 &metrics_to_flush,
1446 context,
1447 "series",
1448 ratio,
1449 payloads_tx,
1450 batch_id,
1451 payload_info,
1452 )
1453 .await
1454}
1455
1456async fn encode_and_flush_v3_sketch_metrics(
1457 context: V3FlushContext<'_>, metrics: &mut Vec<Metric>, ratio: &mut V3CompressionRatio,
1458 payloads_tx: &mut mpsc::Sender<Payload>, batch_id: Option<&Uuid>, payload_info: Option<MetricsPayloadInfo>,
1459) -> Result<(), GenericError> {
1460 if metrics.is_empty() {
1461 return Ok(());
1462 }
1463 let metrics_to_flush = std::mem::take(metrics);
1464
1465 encode_and_flush_v3_payload_requests(
1466 V3_SKETCHES_ENDPOINT_URI,
1467 &metrics_to_flush,
1468 context,
1469 "sketches",
1470 ratio,
1471 payloads_tx,
1472 batch_id,
1473 payload_info,
1474 )
1475 .await
1476}
1477
1478#[allow(clippy::too_many_arguments)]
1485const V3_MIN_COMPRESSION_RATIO: f64 = 0.02;
1489
1490const V3_BATCH_TARGET_MARGIN: f64 = 0.9;
1495
1496const V3_RATIO_SMOOTHING: f64 = 0.25;
1498
1499#[derive(Default)]
1506struct V3CompressionRatio {
1507 ratio: Option<f64>,
1508}
1509
1510impl V3CompressionRatio {
1511 fn uncompressed_target(&self, limits: V3PayloadLimits) -> usize {
1513 let ratio = self.ratio.unwrap_or(1.0).clamp(V3_MIN_COMPRESSION_RATIO, 1.0);
1516 let target_from_compressed_limit = (limits.max_compressed_size as f64 / ratio) as usize;
1517 let target = target_from_compressed_limit.min(limits.max_uncompressed_size);
1518
1519 (((target as f64) * V3_BATCH_TARGET_MARGIN) as usize).max(1)
1521 }
1522
1523 fn record(&mut self, uncompressed_len: usize, compressed_len: usize) {
1525 if uncompressed_len == 0 {
1526 return;
1527 }
1528
1529 let observed = compressed_len as f64 / uncompressed_len as f64;
1530 self.ratio = Some(match self.ratio {
1531 Some(current) => (current * (1.0 - V3_RATIO_SMOOTHING)) + (observed * V3_RATIO_SMOOTHING),
1532 None => observed,
1533 });
1534 }
1535}
1536
1537struct V3Batch {
1539 encoded: V3EncodedMetrics,
1540 event_count: usize,
1541 data_point_count: usize,
1542}
1543
1544fn build_v3_batch(
1551 metrics: &[Metric], idx: &mut usize, context: V3FlushContext<'_>, payload_kind: &'static str,
1552 target_uncompressed_len: usize,
1553) -> Option<V3Batch> {
1554 let mut writer = V3Writer::new();
1555 let mut tags_deduplicator = ReusableDeduplicator::new();
1556 let additional_tags = context.endpoint_config.additional_tags();
1557 let mut data_point_count = 0usize;
1558 let mut cut_for_size = false;
1559
1560 while *idx < metrics.len() {
1561 let metric = &metrics[*idx];
1562
1563 if !metric_has_emittable_values(metric) {
1564 debug!(metric_name = %metric.context().name(), "Dropping metric with no emittable values.");
1565 context.telemetry.events_dropped_encoder().increment(1);
1566 *idx += 1;
1567 continue;
1568 }
1569
1570 let metric_points = metric.values().len();
1571 if !context.payload_limits.point_count_fits(metric_points) {
1572 context.serializer_telemetry.record_item_too_big();
1574 context
1575 .serializer_telemetry
1576 .record_split_reason(V3PayloadSplitReason::ItemTooBig);
1577 warn!(
1578 payload_kind,
1579 data_points = metric_points,
1580 point_limit = context.payload_limits.max_points_per_payload,
1581 "Dropping oversized V3 metric that exceeds the point-count limit."
1582 );
1583 context.telemetry.events_dropped_encoder().increment(1);
1584 *idx += 1;
1585 continue;
1586 }
1587
1588 if writer.metric_count() > 0 {
1591 if !context
1592 .payload_limits
1593 .point_count_fits(data_point_count + metric_points)
1594 {
1595 context
1596 .serializer_telemetry
1597 .record_split_reason(V3PayloadSplitReason::MaxPoints);
1598 break;
1599 }
1600
1601 if context.payload_limits.metric_count_reached(writer.metric_count()) {
1602 break;
1603 }
1604 }
1605
1606 write_metric_to_v3(&mut writer, metric, additional_tags, &mut tags_deduplicator);
1607 data_point_count += metric_points;
1608 *idx += 1;
1609
1610 if writer.estimated_uncompressed_len() >= target_uncompressed_len {
1614 cut_for_size = true;
1615 break;
1616 }
1617 }
1618
1619 if cut_for_size && *idx < metrics.len() {
1622 context
1623 .serializer_telemetry
1624 .record_split_reason(V3PayloadSplitReason::PayloadFull);
1625 }
1626
1627 let event_count = writer.metric_count();
1628 if event_count == 0 {
1629 return None;
1630 }
1631
1632 match writer.finalize() {
1633 Ok(encoded) => Some(V3Batch {
1634 encoded,
1635 event_count,
1636 data_point_count,
1637 }),
1638 Err(e) => {
1639 error!(error = %e, payload_kind, events = event_count, "Failed to encode V3 metrics payload request.");
1640 context.telemetry.events_dropped_encoder().increment(event_count as u64);
1641 None
1642 }
1643 }
1644}
1645
1646async fn create_v3_batch_request(
1648 endpoint_uri: &str, batch: V3Batch, context: V3FlushContext<'_>, payload_kind: &'static str,
1649 ratio: &mut V3CompressionRatio,
1650) -> Option<V3PayloadRequest> {
1651 let V3Batch {
1652 encoded,
1653 event_count,
1654 data_point_count,
1655 } = batch;
1656
1657 let mut encoded_request =
1658 match create_v3_request(endpoint_uri, encoded, context.endpoint_config.compression_scheme()).await {
1659 Ok(request) => request,
1660 Err(e) => {
1661 error!(error = %e, payload_kind, events = event_count, "Failed to create V3 metrics request.");
1662 context.telemetry.events_dropped_encoder().increment(event_count as u64);
1663 return None;
1664 }
1665 };
1666
1667 ratio.record(encoded_request.uncompressed_len, encoded_request.compressed_len);
1670
1671 if !context.payload_limits.request_fits(&encoded_request) {
1672 return None;
1673 }
1674
1675 if let Err(e) =
1679 measure_v3_column_compressed_sizes(&mut encoded_request.stats, context.endpoint_config.compression_scheme())
1680 .await
1681 {
1682 error!(error = %e, payload_kind, "Failed to measure V3 column compressed sizes.");
1683 } else {
1684 record_v3_serializer_stats(context.serializer_telemetry, &encoded_request.stats);
1685 }
1686
1687 Some(V3PayloadRequest {
1688 request: encoded_request.request,
1689 event_count,
1690 data_point_count,
1691 })
1692}
1693
1694async fn split_and_encode_oversized_v3_range(
1700 endpoint_uri: &str, metrics: &[Metric], range: Range<usize>, context: V3FlushContext<'_>,
1701 payload_kind: &'static str, ratio: &mut V3CompressionRatio,
1702) -> Vec<V3PayloadRequest> {
1703 let accepted_metrics = metrics[range]
1706 .iter()
1707 .filter(|metric| {
1708 metric_has_emittable_values(metric) && context.payload_limits.point_count_fits(metric.values().len())
1709 })
1710 .collect::<Vec<_>>();
1711
1712 let mut requests = Vec::new();
1713 let mut pending_ranges = VecDeque::new();
1714 pending_ranges.push_back(0..accepted_metrics.len());
1715
1716 while let Some(range) = pending_ranges.pop_front() {
1717 if range.is_empty() {
1718 continue;
1719 }
1720
1721 let metrics_in_range = &accepted_metrics[range.clone()];
1722 let event_count = metrics_in_range.len();
1723 let data_point_count = metrics_in_range.iter().map(|metric| metric.values().len()).sum();
1724
1725 let encoded = match encode_v3_metrics_batch(
1726 metrics_in_range.iter().copied(),
1727 context.endpoint_config.additional_tags(),
1728 ) {
1729 Ok(encoded) => encoded,
1730 Err(e) => {
1731 error!(error = %e, payload_kind, events = event_count, "Failed to encode V3 metrics payload request.");
1732 context.telemetry.events_dropped_encoder().increment(event_count as u64);
1733 continue;
1734 }
1735 };
1736 let batch = V3Batch {
1737 encoded,
1738 event_count,
1739 data_point_count,
1740 };
1741
1742 if let Some(request) = create_v3_batch_request(endpoint_uri, batch, context, payload_kind, ratio).await {
1743 requests.push(request);
1744 continue;
1745 }
1746
1747 if range.len() == 1 {
1748 context.serializer_telemetry.record_item_too_big();
1750 context
1751 .serializer_telemetry
1752 .record_split_reason(V3PayloadSplitReason::ItemTooBig);
1753 warn!(
1754 payload_kind,
1755 compressed_limit = context.payload_limits.max_compressed_size,
1756 uncompressed_limit = context.payload_limits.max_uncompressed_size,
1757 "Dropping oversized V3 metric that cannot be split further."
1758 );
1759 context.telemetry.events_dropped_encoder().increment(1);
1760 continue;
1761 }
1762
1763 context
1765 .serializer_telemetry
1766 .record_split_reason(V3PayloadSplitReason::PayloadFull);
1767 let pivot = range.start + range.len() / 2;
1768 pending_ranges.push_front(pivot..range.end);
1769 pending_ranges.push_front(range.start..pivot);
1770 }
1771
1772 requests
1773}
1774
1775#[allow(clippy::too_many_arguments)]
1776async fn encode_and_flush_v3_payload_requests(
1777 endpoint_uri: &str, metrics: &[Metric], context: V3FlushContext<'_>, payload_kind: &'static str,
1778 ratio: &mut V3CompressionRatio, payloads_tx: &mut mpsc::Sender<Payload>, batch_id: Option<&Uuid>,
1779 payload_info: Option<MetricsPayloadInfo>,
1780) -> Result<(), GenericError> {
1781 let stream_payloads = batch_id.is_none();
1783 let mut pending_requests = Vec::new();
1784 let target_uncompressed_len = ratio.uncompressed_target(context.payload_limits);
1785 let mut idx = 0;
1786
1787 while idx < metrics.len() {
1788 let batch_start = idx;
1789 let Some(batch) = build_v3_batch(metrics, &mut idx, context, payload_kind, target_uncompressed_len) else {
1790 if idx == batch_start {
1793 error!(
1794 payload_kind,
1795 "V3 batching made no progress; dropping remaining metrics."
1796 );
1797 context
1798 .telemetry
1799 .events_dropped_encoder()
1800 .increment((metrics.len() - idx) as u64);
1801 break;
1802 }
1803 continue;
1804 };
1805
1806 let requests = match create_v3_batch_request(endpoint_uri, batch, context, payload_kind, ratio).await {
1807 Some(request) => vec![request],
1808 None => {
1812 split_and_encode_oversized_v3_range(
1813 endpoint_uri,
1814 metrics,
1815 batch_start..idx,
1816 context,
1817 payload_kind,
1818 ratio,
1819 )
1820 .await
1821 }
1822 };
1823
1824 for payload_request in requests {
1825 if stream_payloads {
1826 flush_payload(
1827 payload_request.request,
1828 payload_request.event_count,
1829 payload_request.data_point_count,
1830 payloads_tx,
1831 None,
1832 0,
1833 0,
1834 payload_info,
1835 )
1836 .await?;
1837 debug!(
1838 payload_kind,
1839 events = payload_request.event_count,
1840 data_points = payload_request.data_point_count,
1841 "Sent V3 payload."
1842 );
1843 } else {
1844 pending_requests.push(payload_request);
1845 }
1846 }
1847 }
1848
1849 let batch_len = pending_requests.len();
1850 for (batch_seq, payload_request) in pending_requests.into_iter().enumerate() {
1851 flush_payload(
1852 payload_request.request,
1853 payload_request.event_count,
1854 payload_request.data_point_count,
1855 payloads_tx,
1856 batch_id,
1857 batch_seq,
1858 batch_len,
1859 payload_info,
1860 )
1861 .await?;
1862 debug!(
1863 payload_kind,
1864 events = payload_request.event_count,
1865 data_points = payload_request.data_point_count,
1866 "Sent V3 payload."
1867 );
1868 }
1869
1870 Ok(())
1871}
1872
1873fn record_v3_serializer_stats(telemetry: &V3SerializerTelemetry, stats: &V3EncoderStats) {
1874 telemetry.record_values_count(stats.value_encoding_stats);
1875
1876 for column in &stats.columns {
1877 let uncompressed_size = column.bytes.len() as u64;
1878 let compressed_size = column.compressed_len as u64;
1879 telemetry.record_column_size(column.field_number, uncompressed_size, compressed_size);
1880 }
1881}
1882
1883async fn measure_v3_column_compressed_sizes(
1888 stats: &mut V3EncoderStats, compression_scheme: CompressionScheme,
1889) -> Result<(), GenericError> {
1890 for column in &mut stats.columns {
1891 column.compressed_len = compressed_v3_len(&column.bytes, compression_scheme)
1892 .await
1893 .error_context("Failed to measure V3 column compressed size.")?;
1894 }
1895
1896 Ok(())
1897}
1898
1899async fn compressed_v3_len(bytes: &[u8], compression_scheme: CompressionScheme) -> Result<usize, GenericError> {
1900 if matches!(compression_scheme, CompressionScheme::Noop) {
1901 return Ok(bytes.len());
1902 }
1903
1904 let buffer = ChunkedBytesBuffer::new(RB_BUFFER_CHUNK_SIZE);
1905 let mut compressor = Compressor::from_scheme(compression_scheme, buffer);
1906 compressor
1907 .write_all(bytes)
1908 .await
1909 .error_context("Failed to compress V3 bytes.")?;
1910 compressor
1911 .flush()
1912 .await
1913 .error_context("Failed to flush V3 compressor.")?;
1914 compressor
1915 .shutdown()
1916 .await
1917 .error_context("Failed to shutdown V3 compressor.")?;
1918
1919 Ok(compressor.into_inner().freeze().len())
1920}
1921
1922fn uuid_to_header_value(uuid: &Uuid) -> HeaderValue {
1924 let s = uuid.as_hyphenated().to_string();
1925 HeaderValue::try_from(s).expect("hyphenated UUID should be a valid header value")
1926}
1927
1928fn usize_to_header_value(value: usize) -> HeaderValue {
1930 let s = value.to_string();
1931 HeaderValue::try_from(s).expect("usize should be a valid header value")
1932}
1933
1934async fn flush_payload(
1935 mut request: Request<FrozenChunkedBytesBuffer>, event_count: usize, data_point_count: usize,
1936 payloads_tx: &mut mpsc::Sender<Payload>, batch_id: Option<&Uuid>, batch_seq: usize, batch_len: usize,
1937 payload_info: Option<MetricsPayloadInfo>,
1938) -> Result<(), GenericError> {
1939 if let Some(batch_id) = batch_id {
1941 let headers = request.headers_mut();
1942 headers.insert("X-Metrics-Request-ID", uuid_to_header_value(batch_id));
1943 headers.insert("X-Metrics-Request-Seq", usize_to_header_value(batch_seq));
1944 headers.insert("X-Metrics-Request-Len", usize_to_header_value(batch_len));
1945 }
1946
1947 let mut payload_meta = PayloadMetadata::from_event_and_data_point_count(event_count, data_point_count);
1948 if let Some(info) = payload_info {
1949 payload_meta = payload_meta.with(info);
1950 }
1951 let http_payload = HttpPayload::new(payload_meta, request);
1952 let payload = Payload::Http(http_payload);
1953
1954 payloads_tx
1955 .send(payload)
1956 .await
1957 .error_context("Failed to send payload.")?;
1958
1959 Ok(())
1960}
1961
1962fn encode_v3_metrics_batch<'a>(
1964 metrics: impl IntoIterator<Item = &'a Metric>, additional_tags: &SharedTagSet,
1965) -> Result<V3EncodedMetrics, GenericError> {
1966 let mut writer = V3Writer::new();
1967 let mut tags_deduplicator = ReusableDeduplicator::new();
1968
1969 for metric in metrics {
1970 write_metric_to_v3(&mut writer, metric, additional_tags, &mut tags_deduplicator);
1971 }
1972
1973 writer
1974 .finalize()
1975 .map_err(|e| generic_error!("Failed to serialize V3 payload: {}", e))
1976}
1977
1978pub(super) fn sketch_has_emittable_values(sketch: &DDSketch) -> bool {
1980 !sketch.is_empty() || !sketch.bins().is_empty()
1981}
1982
1983pub(super) fn metric_has_emittable_values(metric: &Metric) -> bool {
1984 match metric.values() {
1985 MetricValues::Counter(points) | MetricValues::Rate(points, _) | MetricValues::Gauge(points) => {
1986 points.into_iter().any(|(_, value)| emittable_scalar_point(value))
1987 }
1988 MetricValues::Set(points) => points.into_iter().any(|(_, value)| emittable_scalar_point(value)),
1989 MetricValues::Distribution(sketches) => sketches
1990 .into_iter()
1991 .any(|(_, sketch)| sketch_has_emittable_values(sketch)),
1992 MetricValues::Histogram(points) => points.into_iter().any(|(_, histogram)| !histogram.samples().is_empty()),
1993 }
1994}
1995
1996fn write_metric_to_v3(
1997 writer: &mut V3Writer, metric: &Metric, additional_tags: &SharedTagSet,
1998 tags_deduplicator: &mut ReusableDeduplicator<Tag>,
1999) {
2000 if !metric_has_emittable_values(metric) {
2001 debug!(metric_name = %metric.context().name(), "Dropping metric with no emittable values.");
2002 return;
2003 }
2004
2005 let metric_type = match metric.values() {
2006 MetricValues::Counter(..) => V3MetricType::Count,
2007 MetricValues::Rate(..) => V3MetricType::Rate,
2008 MetricValues::Gauge(..) | MetricValues::Set(..) => V3MetricType::Gauge,
2009 MetricValues::Histogram(..) | MetricValues::Distribution(..) => V3MetricType::Sketch,
2010 };
2011 let is_sketch = metric_type == V3MetricType::Sketch;
2012
2013 let mut builder = writer.write(metric_type, metric.context().name());
2014
2015 let chained_tags = metric
2017 .context()
2018 .tags()
2019 .into_iter()
2020 .chain(additional_tags)
2021 .chain(metric.context().origin_tags());
2022 let all_tags = tags_deduplicator
2023 .deduplicated(chained_tags)
2024 .filter(|t| is_sketch || !is_v3_series_resource_tag(t) && !is_v3_series_device_tag(t))
2025 .map(|t| t.as_str());
2026 builder.set_tags(all_tags);
2027
2028 let mut resources = Vec::new();
2030 if let Some(host) = metric.context().host().filter(|host| !host.is_empty()) {
2031 resources.push(("host", host));
2032 }
2033 if !is_sketch {
2034 let mut device_resource = None;
2035 let chained_tags = metric
2036 .context()
2037 .origin_tags()
2038 .into_iter()
2039 .chain(metric.context().tags())
2040 .chain(additional_tags);
2041 for tag in tags_deduplicator.deduplicated(chained_tags) {
2042 if is_v3_series_device_tag(tag) {
2043 device_resource = tag.value().filter(|device| !device.is_empty());
2044 } else if is_v3_series_resource_tag(tag) {
2045 if let Some((rtype, rname)) = tag.value().and_then(|value| value.split_once(':')) {
2046 if !rtype.is_empty() && !rname.is_empty() {
2047 resources.push((rtype, rname));
2048 }
2049 }
2050 }
2051 }
2052 if let Some(device) = device_resource {
2053 let device_idx = usize::from(metric.context().host().is_some_and(|host| !host.is_empty()));
2054 resources.insert(device_idx, ("device", device));
2055 }
2056 }
2057 builder.set_resources(&resources);
2058
2059 if let Some(origin) = metric.metadata().origin() {
2061 match origin {
2062 MetricOrigin::SourceType(source_type) => {
2063 builder.set_source_type(source_type.as_ref());
2064 }
2065 MetricOrigin::OriginMetadata {
2066 product,
2067 subproduct,
2068 product_detail,
2069 } => {
2070 builder.set_origin(*product, *subproduct, *product_detail, false);
2071 }
2072 }
2073 }
2074
2075 if metric_type != V3MetricType::Sketch {
2076 if let Some(unit) = metric.metadata().unit() {
2077 builder.set_unit(unit);
2078 }
2079 }
2080
2081 match metric.values() {
2083 MetricValues::Counter(points) | MetricValues::Gauge(points) => {
2084 for (ts, val) in points {
2085 if !emittable_scalar_point(val) {
2086 continue;
2087 }
2088 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
2089 builder.add_point(timestamp, val);
2090 }
2091 }
2092 MetricValues::Rate(points, interval) => {
2093 builder.set_interval(interval.as_secs());
2094 for (ts, val) in points {
2095 let scaled = val / interval.as_secs_f64();
2097 if !emittable_scalar_point(scaled) {
2098 continue;
2099 }
2100 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
2101 builder.add_point(timestamp, scaled);
2102 }
2103 }
2104 MetricValues::Set(points) => {
2105 for (ts, count) in points {
2107 if !emittable_scalar_point(count) {
2108 continue;
2109 }
2110 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
2111 builder.add_point(timestamp, count);
2112 }
2113 }
2114 MetricValues::Distribution(sketches) => {
2115 for (ts, sketch) in sketches {
2116 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
2117 if sketch_has_emittable_values(sketch) {
2118 let bin_keys: Vec<i32> = sketch.bins().iter().map(|b| b.key()).collect();
2119 let bin_counts: Vec<u32> = sketch.bins().iter().map(|b| b.count()).collect();
2120 builder.add_sketch(
2121 timestamp,
2122 sketch.count() as i64,
2123 sketch.stored_sum(),
2124 sketch.stored_min(),
2125 sketch.stored_max(),
2126 &bin_keys,
2127 &bin_counts,
2128 );
2129 }
2130 }
2131 }
2132 MetricValues::Histogram(histograms) => {
2133 for (ts, histogram) in histograms {
2134 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
2135 let mut sketch = DDSketch::default();
2137 for sample in histogram.samples() {
2138 sketch.insert_n(sample.value.into_inner(), sample.weight.0 as u64);
2139 }
2140 if !sketch.is_empty() {
2141 let bin_keys: Vec<i32> = sketch.bins().iter().map(|b| b.key()).collect();
2142 let bin_counts: Vec<u32> = sketch.bins().iter().map(|b| b.count()).collect();
2143 builder.add_sketch(
2144 timestamp,
2145 sketch.count() as i64,
2146 sketch.sum().unwrap_or(0.0),
2147 sketch.min().unwrap_or(0.0),
2148 sketch.max().unwrap_or(0.0),
2149 &bin_keys,
2150 &bin_counts,
2151 );
2152 }
2153 }
2154 }
2155 }
2156
2157 builder.close();
2158}
2159
2160#[inline]
2161fn emittable_scalar_point(point: f64) -> bool {
2162 point.is_finite()
2163}
2164
2165fn is_v3_series_device_tag(tag: &Tag) -> bool {
2166 tag.name() == "device" && tag.value().is_some()
2167}
2168
2169fn is_v3_series_resource_tag(tag: &Tag) -> bool {
2170 tag.name() == "dd.internal.resource" && tag.value().is_some()
2171}
2172
2173async fn create_v3_request(
2175 endpoint_uri: &str, encoded: V3EncodedMetrics, compression_scheme: CompressionScheme,
2176) -> Result<V3EncodedRequest, GenericError> {
2177 let mut header_buf = [0; 16];
2180 let header_len = {
2181 let mut header_writer = CodedOutputStream::bytes(&mut header_buf);
2182 header_writer.write_tag(3, WireType::LengthDelimited)?;
2183 header_writer.write_uint64_no_tag(encoded.payload.len() as u64)?;
2184 header_writer.flush()?;
2185 header_writer.total_bytes_written() as usize
2186 };
2187
2188 let uncompressed_len = header_len + encoded.payload.len();
2189
2190 let buffer = ChunkedBytesBuffer::new(RB_BUFFER_CHUNK_SIZE);
2191 let mut compressor = Compressor::from_scheme(compression_scheme, buffer);
2192 compressor
2193 .write_all(&header_buf[..header_len])
2194 .await
2195 .error_context("Failed to compress V3 payload.")?;
2196 compressor
2197 .write_all(&encoded.payload)
2198 .await
2199 .error_context("Failed to compress V3 payload.")?;
2200 compressor
2201 .flush()
2202 .await
2203 .error_context("Failed to flush V3 compressor.")?;
2204 compressor
2205 .shutdown()
2206 .await
2207 .error_context("Failed to shutdown V3 compressor.")?;
2208
2209 let compressed_buf = compressor.into_inner().freeze();
2210 let compressed_len = compressed_buf.len();
2211
2212 let mut builder = Request::builder()
2213 .method(Method::POST)
2214 .uri(endpoint_uri)
2215 .header(http::header::CONTENT_TYPE, "application/x-protobuf");
2216
2217 if let Some(encoding) = content_encoding_for_scheme(compression_scheme) {
2218 builder = builder.header(http::header::CONTENT_ENCODING, encoding);
2219 }
2220
2221 let request = builder
2222 .body(compressed_buf)
2223 .map_err(|e| generic_error!("Failed to build V3 request: {}", e))?;
2224
2225 Ok(V3EncodedRequest {
2226 request,
2227 compressed_len,
2228 uncompressed_len,
2229 stats: encoded.stats,
2230 })
2231}
2232
2233fn content_encoding_for_scheme(compression_scheme: CompressionScheme) -> Option<HeaderValue> {
2234 match compression_scheme {
2235 CompressionScheme::Noop => None,
2236 CompressionScheme::Gzip(_) => Some(HeaderValue::from_static("gzip")),
2237 CompressionScheme::Zlib(_) => Some(HeaderValue::from_static("deflate")),
2238 CompressionScheme::Zstd(_) => Some(HeaderValue::from_static("zstd")),
2239 }
2240}
2241
2242#[cfg(test)]
2243mod tests {
2244 use std::{collections::HashMap, io::Cursor};
2245
2246 use bytes::Bytes;
2247 use datadog_protos::metrics::v3::MetricData as V3MetricData;
2248 use protobuf::Message as _;
2249 use saluki_config::ConfigurationLoader;
2250 use saluki_context::{
2251 tags::{Tag, TagSet},
2252 Context,
2253 };
2254 use saluki_core::data_model::{
2255 event::{metric::MetricMetadata, Event},
2256 payload::Payload,
2257 };
2258 use saluki_metrics::test::TestRecorder;
2259 use stringtheory::MetaString;
2260 use tokio::time::timeout;
2261
2262 use super::*;
2263
2264 #[test]
2265 fn deser_agent_v3_api_nested_settings() {
2266 let raw = r#"
2267serializer_experimental_use_v3_api:
2268 compression_level: 7
2269 series:
2270 endpoints:
2271 - https://app.datadoghq.com
2272 validate: true
2273 use_beta: true
2274 beta_route: /api/intake/metrics/custom/series
2275 shadow_sample_rate: 0.25
2276 shadow_sites:
2277 - datadoghq.eu
2278 sketches:
2279 endpoints:
2280 - https://app.datadoghq.eu
2281"#;
2282
2283 let config =
2284 serde_yaml::from_str::<DatadogMetricsConfiguration>(raw).expect("configuration should deserialize");
2285
2286 assert_eq!(7, config.v3_api.compression_level);
2287 assert_eq!(
2288 Some("https://app.datadoghq.com"),
2289 config.v3_api.series.endpoints.first().map(String::as_str)
2290 );
2291 assert!(config.v3_api.series.validate);
2292 assert!(config.v3_api.series.use_beta);
2293 assert_eq!("/api/intake/metrics/custom/series", config.v3_api.series.beta_route);
2294 assert_eq!(0.25, config.v3_api.series.shadow_sample_rate);
2295 assert_eq!(vec!["datadoghq.eu"], config.v3_api.series.shadow_sites);
2296 assert_eq!(
2297 Some("https://app.datadoghq.eu"),
2298 config.v3_api.sketches.endpoints.first().map(String::as_str)
2299 );
2300 }
2301
2302 #[tokio::test]
2303 async fn typed_metrics_routing_is_authoritative() {
2304 let (raw, _) = ConfigurationLoader::for_tests(
2305 Some(serde_json::json!({
2306 "serializer_compressor_kind": "zlib",
2307 "use_v2_api_series": true,
2308 "use_v3_api_series_enabled": "true",
2309 "observability_pipelines_worker_metrics_enabled": false,
2310 })),
2311 None,
2312 false,
2313 )
2314 .await;
2315
2316 let mut endpoints = Endpoints::default();
2317 endpoints.compression.compressor_kind = "zstd".to_string();
2318 endpoints.opw_intake.enabled = true;
2319 endpoints.opw_intake.url = "https://opw.example.com".to_string();
2320 endpoints.opw_intake.use_v3_series = true;
2321 let mut metrics = MetricsEncoding {
2322 use_v2_series_api: false,
2323 ..Default::default()
2324 };
2325 metrics.v3_api.compression_level = 7;
2326 metrics.v3_api.series.validate = true;
2327 metrics.v3_series_mode.mode = "false".to_string();
2328 metrics.v3_series_mode.endpoint_modes =
2329 HashMap::from([("https://app.datadoghq.com".to_string(), "true".to_string())]);
2330
2331 let config = DatadogMetricsConfiguration::from_configuration_with_metrics_routing(&raw, &metrics, &endpoints)
2332 .expect("configuration should deserialize");
2333
2334 assert_eq!(config.compressor_kind, "zstd");
2335 assert!(!config.use_v2_api_series);
2336 assert_eq!(config.v3_api.compression_level, 7);
2337 assert!(config.v3_api.series.validate);
2338 assert_eq!(config.use_v3_api_series.enabled, "false");
2339 assert_eq!(
2340 config.use_v3_api_series.endpoints.get("https://app.datadoghq.com"),
2341 Some(&"true".to_string())
2342 );
2343 assert!(config.observability_pipelines_worker_metrics_enabled);
2344 assert_eq!(
2345 config.observability_pipelines_worker_metrics_url,
2346 "https://opw.example.com"
2347 );
2348 assert!(config.observability_pipelines_worker_metrics_use_v3_api_series);
2349 }
2350
2351 #[test]
2352 fn agent_v3_api_shadow_defaults_match_agent() {
2353 let config = serde_yaml::from_str::<DatadogMetricsConfiguration>("").expect("configuration should deserialize");
2354
2355 assert_eq!(0.0, config.v3_api.series.shadow_sample_rate);
2356 assert_eq!(vec!["datadoghq.com"], config.v3_api.series.shadow_sites);
2357 }
2358
2359 #[test]
2360 fn shadow_sample_matches_agent_threshold_behavior() {
2361 assert!(!shadow_sample_matches(0.0, 0.0));
2362 assert!(shadow_sample_matches(0.5, 0.4));
2363 assert!(!shadow_sample_matches(0.5, 0.5));
2364 assert!(!shadow_sample_matches(0.5, 0.6));
2365 }
2366
2367 #[test]
2368 fn shadow_sampling_is_disabled_for_v1_series_baseline_or_v3_incompatible_compressor() {
2369 assert!(series_shadow_config_for_endpoint(MetricsEndpoint::SeriesV2, 1.0, false).is_enabled());
2370 assert!(!series_shadow_config_for_endpoint(MetricsEndpoint::SeriesV1, 1.0, false).is_enabled());
2371 assert!(!series_shadow_config_for_endpoint(MetricsEndpoint::SeriesV2, 1.0, true).is_enabled());
2372 }
2373
2374 #[test]
2375 fn metrics_v3_disabled_by_compressor_uses_v2_only() {
2376 assert_eq!(
2377 MetricsEncoderMode::V2Only,
2378 metrics_encoder_mode_for_config(true, false, true)
2379 );
2380 assert_eq!(
2381 MetricsEncoderMode::V2Only,
2382 metrics_encoder_mode_for_config(true, true, true)
2383 );
2384 assert_eq!(
2385 MetricsEncoderMode::V3Enabled,
2386 metrics_encoder_mode_for_config(true, false, false)
2387 );
2388 assert_eq!(
2389 MetricsEncoderMode::Validation,
2390 metrics_encoder_mode_for_config(true, true, false)
2391 );
2392 }
2393
2394 fn v3_series_config(raw: &str) -> DatadogMetricsConfiguration {
2395 serde_yaml::from_str(raw).expect("configuration should deserialize")
2396 }
2397
2398 #[test]
2399 fn mixed_v2_and_v3_endpoints_require_both_series_encoders() {
2400 let config = v3_series_config(
2401 r#"
2402use_v3_api_series_enabled: "datadog_only"
2403additional_endpoints:
2404 https://custom.example.com:
2405 - additional-api-key
2406"#,
2407 );
2408
2409 assert!(config.requires_v2_series(false).expect("endpoints should resolve"));
2410 assert!(config.requires_v3_series(false).expect("endpoints should resolve"));
2411 }
2412
2413 #[test]
2414 fn all_v2_endpoints_do_not_require_v3_series() {
2415 let config = v3_series_config(
2416 r#"
2417dd_url: http://127.0.0.1:9091
2418use_v3_api_series_enabled: "datadog_only"
2419"#,
2420 );
2421
2422 assert!(config.requires_v2_series(false).expect("endpoints should resolve"));
2423 assert!(!config.requires_v3_series(false).expect("endpoints should resolve"));
2424 }
2425
2426 #[test]
2427 fn validation_requires_v2_series() {
2428 let config = v3_series_config(
2429 r#"
2430use_v3_api_series_enabled: "true"
2431serializer_experimental_use_v3_api:
2432 series:
2433 validate: true
2434"#,
2435 );
2436
2437 assert!(config.requires_v2_series(false).expect("endpoints should resolve"));
2438 assert!(config.requires_v3_series(false).expect("endpoints should resolve"));
2439 }
2440
2441 #[test]
2442 fn v1_series_configuration_takes_precedence_over_v3() {
2443 let config = v3_series_config(
2444 r#"
2445use_v2_api_series: false
2446use_v3_api_series_enabled: "true"
2447"#,
2448 );
2449
2450 let series_v3_can_be_enabled = series_v3_can_be_enabled_for_config(
2451 config.use_v2_api_series,
2452 false,
2453 None,
2454 false,
2455 &config.use_v3_api_series,
2456 );
2457
2458 assert!(!series_v3_can_be_enabled);
2459
2460 assert_eq!(
2461 MetricsEncoderMode::V2Only,
2462 metrics_encoder_mode_for_config(series_v3_can_be_enabled, false, false)
2463 );
2464
2465 assert!(
2466 config.requires_v2_series(false).expect("endpoint should resolve"),
2467 "V1 configuration must retain the legacy series builder even when V3 is enabled"
2468 );
2469 }
2470
2471 #[test]
2472 fn all_v3_serializer_endpoints_require_only_v3_series() {
2473 let config = v3_series_config(
2474 r#"
2475dd_url: https://agent.datad0g.com.
2476use_v3_api_series_enabled: "false"
2477additional_endpoints:
2478 https://agent.datadoghq.com.:
2479 - additional-api-key
2480serializer_experimental_use_v3_api:
2481 series:
2482 endpoints:
2483 - https://agent.datad0g.com.
2484 - https://agent.datadoghq.com.
2485"#,
2486 );
2487
2488 assert!(!config.requires_v2_series(false).expect("endpoints should resolve"));
2489 assert!(config.requires_v3_series(false).expect("endpoints should resolve"));
2490 }
2491
2492 #[test]
2493 fn endpoint_override_uses_the_overridden_endpoint_protocol() {
2494 let config = v3_series_config(
2495 r#"
2496dd_url: https://primary.example.com
2497use_v3_api_series_enabled: "false"
2498serializer_experimental_use_v3_api:
2499 series:
2500 endpoints:
2501 - https://primary.example.com
2502 - https://v3-mrf.example.com
2503"#,
2504 );
2505
2506 let v2_mrf_config = config
2507 .clone()
2508 .with_metrics_endpoint_override("https://v2-mrf.example.com".to_string());
2509 let v3_mrf_config = config.with_metrics_endpoint_override("https://v3-mrf.example.com".to_string());
2510
2511 assert!(v2_mrf_config
2512 .requires_v2_series(false)
2513 .expect("V2 MRF endpoint should resolve"));
2514 assert!(!v2_mrf_config
2515 .requires_v3_series(false)
2516 .expect("V2 MRF endpoint should resolve"));
2517 assert!(!v3_mrf_config
2518 .requires_v2_series(false)
2519 .expect("V3 MRF endpoint should resolve"));
2520 assert!(v3_mrf_config
2521 .requires_v3_series(false)
2522 .expect("V3 MRF endpoint should resolve"));
2523 }
2524
2525 #[test]
2526 fn v2_series_only_override_keeps_v2_and_disables_shadowing() {
2527 let config = v3_series_config(
2528 r#"
2529use_v3_api_series_enabled: "true"
2530serializer_experimental_use_v3_api:
2531 series:
2532 shadow_sample_rate: 1.0
2533"#,
2534 )
2535 .with_v2_series_only();
2536
2537 assert_eq!("false", config.use_v3_api_series.enabled);
2538 assert!(config.use_v3_api_series.endpoints.is_empty());
2539 assert!(config.v3_api.series.endpoints.is_empty());
2540 assert_eq!(0.0, config.v3_api.series.shadow_sample_rate);
2541 assert!(config.requires_v2_series(false).expect("endpoint should resolve"));
2542 assert!(!config.requires_v3_series(false).expect("endpoint should resolve"));
2543 }
2544
2545 #[test]
2546 fn agent_default_v3_does_not_enable_opw_only_encoder_mode() {
2547 let series_config = UseV3ApiSeriesConfig::default();
2548 let invalid_metrics_primary_override =
2549 selected_metrics_primary_v3_override(true, "http://[::1", false, true, "http://vector.example.com", false);
2550
2551 assert!(!series_v3_can_be_enabled_for_config(
2552 true,
2553 false,
2554 Some(false),
2555 false,
2556 &series_config
2557 ));
2558 assert!(series_v3_can_be_enabled_for_config(
2559 true,
2560 false,
2561 Some(true),
2562 false,
2563 &series_config
2564 ));
2565 assert!(series_v3_can_be_enabled_for_config(
2566 true,
2567 false,
2568 Some(false),
2569 true,
2570 &series_config
2571 ));
2572 assert!(series_v3_can_be_enabled_for_config(
2573 true,
2574 true,
2575 Some(false),
2576 false,
2577 &series_config
2578 ));
2579 assert_eq!(None, invalid_metrics_primary_override);
2580 assert!(series_v3_can_be_enabled_for_config(
2581 true,
2582 false,
2583 invalid_metrics_primary_override,
2584 false,
2585 &series_config
2586 ));
2587 }
2588
2589 #[test]
2590 fn v3_encodes_zero_count_distribution_with_bins() {
2591 let mut sketch = DDSketch::default();
2592 sketch.insert_n(42.0, 10);
2593 sketch.set_count(0);
2594 sketch.set_min(40.0);
2595 sketch.set_max(45.0);
2596 sketch.set_sum(123.0);
2597 sketch.set_avg(f64::NAN);
2598
2599 let metric = Metric::from_parts(
2600 Context::from_static_parts("zero.count.distribution", &[]),
2601 MetricValues::distribution((123_u64, sketch)),
2602 MetricMetadata::default(),
2603 );
2604
2605 let encoded = encode_v3_metrics_batch(&[metric], &SharedTagSet::default())
2606 .expect("zero-count distribution should encode to V3");
2607 let metric_data = V3MetricData::parse_from_bytes(&encoded.payload).expect("V3 metric data should decode");
2608
2609 assert_eq!(metric_data.numPoints, vec![1]);
2610 }
2611
2612 #[test]
2613 fn v3_skips_empty_distribution() {
2614 let metric = Metric::from_parts(
2615 Context::from_static_parts("empty.distribution", &[]),
2616 MetricValues::distribution((123_u64, DDSketch::default())),
2617 MetricMetadata::default(),
2618 );
2619
2620 let encoded =
2621 encode_v3_metrics_batch(&[metric], &SharedTagSet::default()).expect("empty distribution should encode");
2622 let metric_data = V3MetricData::parse_from_bytes(&encoded.payload).expect("V3 metric data should decode");
2623
2624 assert!(metric_data.numPoints.is_empty());
2625 }
2626
2627 #[tokio::test]
2628 async fn create_v3_request_uses_configured_endpoint_uri() {
2629 let encoded = V3Writer::new().finalize().expect("empty V3 payload should encode");
2630 let request = create_v3_request("/api/intake/metrics/custom/series", encoded, CompressionScheme::noop())
2631 .await
2632 .expect("request should be created");
2633
2634 assert_eq!("/api/intake/metrics/custom/series", request.request.uri());
2635 }
2636
2637 #[tokio::test]
2638 async fn create_v3_request_uses_single_stream_body_and_column_telemetry() {
2639 let metrics = vec![Metric::counter("v3.single.stream", 42.0)];
2640 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2641 let expected_payload = encoded.payload.clone();
2642
2643 let mut request = create_v3_request(V3_SERIES_ENDPOINT_URI, encoded, CompressionScheme::noop())
2644 .await
2645 .expect("request should be created");
2646 measure_v3_column_compressed_sizes(&mut request.stats, CompressionScheme::noop())
2647 .await
2648 .expect("column compressed sizes should be measured");
2649
2650 for column in &request.stats.columns {
2651 assert_eq!(column.compressed_len, column.bytes.len());
2652 }
2653
2654 let mut expected_body = Vec::new();
2655 {
2656 let mut os = CodedOutputStream::vec(&mut expected_body);
2657 os.write_tag(3, WireType::LengthDelimited).unwrap();
2658 os.write_uint64_no_tag(expected_payload.len() as u64).unwrap();
2659 os.flush().unwrap();
2660 }
2661 expected_body.extend_from_slice(&expected_payload);
2662
2663 assert_eq!(request.request.into_body().into_bytes(), Bytes::from(expected_body));
2664 }
2665
2666 #[tokio::test]
2667 async fn create_v3_request_zstd_body_decodes_to_metric_payload() {
2668 let metrics = vec![Metric::counter("v3.single.stream.zstd", 42.0)];
2669 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2670 let expected_payload = encoded.payload.clone();
2671
2672 let mut request = create_v3_request(V3_SERIES_ENDPOINT_URI, encoded, CompressionScheme::zstd_default())
2673 .await
2674 .expect("request should be created");
2675 measure_v3_column_compressed_sizes(&mut request.stats, CompressionScheme::zstd_default())
2676 .await
2677 .expect("column compressed sizes should be measured");
2678
2679 for column in &request.stats.columns {
2680 let expected_compressed_len = compressed_v3_len(&column.bytes, CompressionScheme::zstd_default())
2681 .await
2682 .expect("column compressed size should be measured");
2683 assert_eq!(column.compressed_len, expected_compressed_len);
2684 }
2685
2686 let mut expected_body = Vec::new();
2687 {
2688 let mut os = CodedOutputStream::vec(&mut expected_body);
2689 os.write_tag(3, WireType::LengthDelimited).unwrap();
2690 os.write_uint64_no_tag(expected_payload.len() as u64).unwrap();
2691 os.flush().unwrap();
2692 }
2693 expected_body.extend_from_slice(&expected_payload);
2694
2695 let decoded_body = zstd::stream::decode_all(Cursor::new(request.request.into_body().into_bytes()))
2696 .expect("compressed V3 body should decode");
2697 assert_eq!(decoded_body, expected_body);
2698 }
2699
2700 #[test]
2701 fn v3_drops_non_finite_scalar_points() {
2702 const NUM_POINTS_FIELD_NUMBER: u32 = 15;
2703 const VALUE_SINT64_FIELD_NUMBER: u32 = 17;
2704
2705 let metrics = vec![
2706 Metric::from_parts(
2707 Context::from_static_parts("v3.finite.counter", &[]),
2708 MetricValues::counter([(1, 1.0_f64), (2, f64::NAN), (3, f64::INFINITY), (4, 2.0)]),
2709 MetricMetadata::default(),
2710 ),
2711 Metric::from_parts(
2712 Context::from_static_parts("v3.finite.gauge", &[]),
2713 MetricValues::gauge([(1, 3.0_f64), (2, f64::NAN), (3, 4.0)]),
2714 MetricMetadata::default(),
2715 ),
2716 Metric::from_parts(
2717 Context::from_static_parts("v3.finite.rate", &[]),
2718 MetricValues::rate(
2719 [(1, 30.0_f64), (2, f64::INFINITY), (3, f64::NAN)],
2720 Duration::from_secs(10),
2721 ),
2722 MetricMetadata::default(),
2723 ),
2724 Metric::set("v3.finite.set", "alpha"),
2725 ];
2726
2727 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2728 let column = |field_number| {
2729 encoded
2730 .stats
2731 .columns
2732 .iter()
2733 .find(|column| column.field_number == field_number)
2734 .map(|column| column.bytes.as_slice())
2735 };
2736
2737 assert_eq!(column(NUM_POINTS_FIELD_NUMBER), Some(&[2, 2, 1, 1][..]));
2738 assert_eq!(column(VALUE_SINT64_FIELD_NUMBER), Some(&[2, 4, 6, 8, 6, 2][..]));
2740 assert_eq!(encoded.stats.value_encoding_stats.sint64, 6);
2741 assert_eq!(encoded.stats.value_encoding_stats.float64, 0);
2742 }
2743
2744 #[tokio::test]
2745 async fn v3_serializer_stats_record_agent_style_value_and_column_counters() {
2746 const VALUE_SINT64_FIELD_NUMBER: u32 = 17;
2747
2748 let recorder = TestRecorder::default();
2749 let _local = metrics::set_default_local_recorder(&recorder);
2750 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2751 let metrics = vec![
2752 Metric::gauge("v3.telemetry.zero", [(123, 0.0), (124, 0.0)]),
2753 Metric::counter("v3.telemetry.sint64", [(123, 100.0), (124, 200.0)]),
2754 Metric::gauge("v3.telemetry.float32", [(123, 1.5), (124, 2.25)]),
2755 Metric::gauge("v3.telemetry.float64", [(123, (1i64 << 30) as f64), (124, 1.5)]),
2756 ];
2757 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2758 let mut request = create_v3_request(V3_SERIES_ENDPOINT_URI, encoded, CompressionScheme::noop())
2759 .await
2760 .expect("request should be created");
2761 measure_v3_column_compressed_sizes(&mut request.stats, CompressionScheme::noop())
2762 .await
2763 .expect("column compressed sizes should be measured");
2764
2765 let value_sint64_column = request
2766 .stats
2767 .columns
2768 .iter()
2769 .find(|column| column.field_number == VALUE_SINT64_FIELD_NUMBER)
2770 .expect("sint64 value column should be present");
2771 let value_sint64_len = value_sint64_column.bytes.len() as u64;
2772 assert_eq!(value_sint64_column.compressed_len as u64, value_sint64_len);
2773
2774 record_v3_serializer_stats(&serializer_telemetry, &request.stats);
2775
2776 assert_eq!(
2777 recorder.counter(("serializer.v3_values_count", &[("type", "zero")])),
2778 Some(2)
2779 );
2780 assert_eq!(
2781 recorder.counter(("serializer.v3_values_count", &[("type", "sint64")])),
2782 Some(2)
2783 );
2784 assert_eq!(
2785 recorder.counter(("serializer.v3_values_count", &[("type", "float32")])),
2786 Some(2)
2787 );
2788 assert_eq!(
2789 recorder.counter(("serializer.v3_values_count", &[("type", "float64")])),
2790 Some(2)
2791 );
2792 assert_eq!(
2793 recorder.counter((
2794 "serializer.v3_column_size",
2795 &[("column", "ValueSint64"), ("compressed", "uncompressed")]
2796 )),
2797 Some(value_sint64_len)
2798 );
2799 assert_eq!(
2800 recorder.counter((
2801 "serializer.v3_column_size",
2802 &[("column", "ValueSint64"), ("compressed", "compressed")]
2803 )),
2804 Some(value_sint64_len)
2805 );
2806 }
2807
2808 async fn create_v3_test_request(metrics: &[Metric]) -> V3EncodedRequest {
2809 let encoded = encode_v3_metrics_batch(metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2810 create_v3_request(V3_SERIES_ENDPOINT_URI, encoded, CompressionScheme::noop())
2811 .await
2812 .expect("request should be created")
2813 }
2814
2815 fn test_v3_flush_context<'a>(
2816 ep_config: &'a EndpointConfiguration, payload_limits: V3PayloadLimits,
2817 serializer_telemetry: &'a V3SerializerTelemetry, telemetry: &'a ComponentTelemetry,
2818 ) -> V3FlushContext<'a> {
2819 V3FlushContext {
2820 endpoint_config: ep_config,
2821 payload_limits,
2822 series_endpoint_uri: V3_SERIES_ENDPOINT_URI,
2823 serializer_telemetry,
2824 telemetry,
2825 }
2826 }
2827
2828 struct CollectedV3Payload {
2830 event_count: usize,
2831 data_point_count: usize,
2832 request: Request<FrozenChunkedBytesBuffer>,
2833 }
2834
2835 async fn collect_v3_payload_requests(metrics: &[Metric], context: V3FlushContext<'_>) -> Vec<CollectedV3Payload> {
2840 let (mut payloads_tx, mut payloads_rx) = tokio::sync::mpsc::channel(64);
2841 let mut ratio = V3CompressionRatio::default();
2842 encode_and_flush_v3_payload_requests(
2843 V3_SERIES_ENDPOINT_URI,
2844 metrics,
2845 context,
2846 "series",
2847 &mut ratio,
2848 &mut payloads_tx,
2849 None,
2850 None,
2851 )
2852 .await
2853 .expect("payload requests should encode and flush");
2854 drop(payloads_tx);
2855
2856 let mut collected = Vec::new();
2857 while let Some(payload) = payloads_rx.recv().await {
2858 let Payload::Http(http_payload) = payload else {
2859 panic!("expected HTTP payload");
2860 };
2861 let (metadata, request) = http_payload.into_parts();
2862 collected.push(CollectedV3Payload {
2863 event_count: metadata.event_count(),
2864 data_point_count: metadata.data_point_count(),
2865 request,
2866 });
2867 }
2868
2869 collected
2870 }
2871
2872 #[tokio::test]
2873 async fn v3_payload_requests_split_by_compressed_size_limit() {
2874 let metrics = vec![
2875 Metric::counter("v3.compressed.split.one", 1.0),
2876 Metric::counter("v3.compressed.split.two", 2.0),
2877 ];
2878 let single_request = create_v3_test_request(&metrics[..1]).await;
2879 let combined_request = create_v3_test_request(&metrics).await;
2880 assert!(combined_request.compressed_len > single_request.compressed_len);
2881
2882 let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2883 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2884 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2885 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2886
2887 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2888 let requests = collect_v3_payload_requests(&metrics, context).await;
2889
2890 assert_eq!(2, requests.len());
2891 assert_eq!(
2892 vec![1, 1],
2893 requests.iter().map(|request| request.event_count).collect::<Vec<_>>()
2894 );
2895 assert!(requests
2896 .iter()
2897 .all(|request| request.request.body().len() <= limits.max_compressed_size));
2898 }
2899
2900 #[tokio::test]
2901 async fn v3_byte_fallback_does_not_resurrect_metric_over_point_limit() {
2902 let point_oversized_metric = Metric::counter("v3.fallback.over.point.limit", [(123, 1.0), (124, 2.0)]);
2903
2904 let byte_oversized_metric = Metric::counter(
2905 concat!(
2906 "v3.fallback.over.byte.limit.",
2907 "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
2908 "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
2909 "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc",
2910 "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd",
2911 ),
2912 3.0,
2913 );
2914
2915 let point_oversized_request = create_v3_test_request(std::slice::from_ref(&point_oversized_metric)).await;
2916 let byte_oversized_request = create_v3_test_request(std::slice::from_ref(&byte_oversized_metric)).await;
2917
2918 assert!(byte_oversized_request.compressed_len > point_oversized_request.compressed_len);
2919
2920 let limits = V3PayloadLimits::new(point_oversized_request.compressed_len, usize::MAX, 10_000, 1);
2921
2922 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2923 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2924 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2925 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2926
2927 let requests = collect_v3_payload_requests(&[point_oversized_metric, byte_oversized_metric], context).await;
2928
2929 let emitted_point_counts = requests
2930 .iter()
2931 .map(|request| request.data_point_count)
2932 .collect::<Vec<_>>();
2933
2934 assert!(
2935 emitted_point_counts
2936 .iter()
2937 .all(|point_count| *point_count <= limits.max_points_per_payload),
2938 "fallback emitted request point counts {emitted_point_counts:?} above limit {}",
2939 limits.max_points_per_payload
2940 );
2941 }
2942
2943 #[tokio::test]
2944 async fn v3_serializer_stats_record_payload_full_split_reason() {
2945 let metrics = vec![
2946 Metric::counter("v3.telemetry.payload_full.one", 1.0),
2947 Metric::counter("v3.telemetry.payload_full.two", 2.0),
2948 ];
2949 let single_request = create_v3_test_request(&metrics[..1]).await;
2950 let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2951 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2952 let recorder = TestRecorder::default();
2953 let _local = metrics::set_default_local_recorder(&recorder);
2954 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2955 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2956
2957 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2958 let requests = collect_v3_payload_requests(&metrics, context).await;
2959
2960 assert_eq!(2, requests.len());
2961 assert_eq!(
2962 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "payload_full")])),
2963 Some(1)
2964 );
2965 }
2966
2967 #[tokio::test]
2968 async fn v3_serializer_stats_record_item_too_big_split_reason() {
2969 let metrics = vec![Metric::counter("v3.telemetry.item_too_big", 1.0)];
2970 let request = create_v3_test_request(&metrics).await;
2971 let limits = V3PayloadLimits::new(request.compressed_len - 1, usize::MAX, 10_000, 10_000);
2972 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2973 let recorder = TestRecorder::default();
2974 let _local = metrics::set_default_local_recorder(&recorder);
2975 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2976 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2977
2978 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2979 let requests = collect_v3_payload_requests(&metrics, context).await;
2980
2981 assert!(requests.is_empty());
2982 assert_eq!(recorder.counter("serializer.v3_item_too_big"), Some(1));
2983 assert_eq!(
2984 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "item_too_big")])),
2985 Some(1)
2986 );
2987 }
2988
2989 #[tokio::test]
2990 async fn v3_payload_requests_split_by_uncompressed_size_limit() {
2991 let metrics = vec![
2992 Metric::counter("v3.uncompressed.split.one", 1.0),
2993 Metric::counter("v3.uncompressed.split.two", 2.0),
2994 ];
2995 let single_request = create_v3_test_request(&metrics[..1]).await;
2996 let combined_request = create_v3_test_request(&metrics).await;
2997 assert!(combined_request.uncompressed_len > single_request.uncompressed_len);
2998
2999 let limits = V3PayloadLimits::new(usize::MAX, single_request.uncompressed_len, 10_000, 10_000);
3000 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3001 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3002 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3003
3004 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3005 let requests = collect_v3_payload_requests(&metrics, context).await;
3006
3007 assert_eq!(2, requests.len());
3008 assert_eq!(
3009 vec![1, 1],
3010 requests.iter().map(|request| request.event_count).collect::<Vec<_>>()
3011 );
3012 }
3013
3014 fn collect_v3_batches(
3016 metrics: &[Metric], context: V3FlushContext<'_>, target_uncompressed_len: usize,
3017 ) -> Vec<(Range<usize>, usize)> {
3018 let mut batches = Vec::new();
3019 let mut idx = 0;
3020 while idx < metrics.len() {
3021 let start = idx;
3022 match build_v3_batch(metrics, &mut idx, context, "series", target_uncompressed_len) {
3023 Some(batch) => batches.push((start..idx, batch.event_count)),
3024 None if idx == start => break,
3025 None => {}
3026 }
3027 }
3028
3029 batches
3030 }
3031
3032 #[test]
3033 fn v3_batches_are_cut_at_the_point_limit() {
3034 let metrics = vec![
3035 Metric::counter("v3.points.split.one", [(123, 1.0), (124, 2.0)]),
3036 Metric::counter("v3.points.split.two", [(123, 3.0), (124, 4.0)]),
3037 Metric::counter("v3.points.split.three", 5.0),
3038 ];
3039 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 3);
3040 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3041 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3042 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3043 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3044
3045 let batches = collect_v3_batches(&metrics, context, usize::MAX);
3046
3047 assert_eq!(vec![(0..1, 1), (1..3, 2)], batches);
3048 }
3049
3050 #[test]
3051 fn v3_batches_are_cut_at_the_metric_limit() {
3052 let metrics = vec![
3053 Metric::counter("v3.metric.limit.a", 1.0),
3054 Metric::counter("v3.metric.limit.b", 2.0),
3055 Metric::counter("v3.metric.limit.c", 3.0),
3056 Metric::counter("v3.metric.limit.d", 4.0),
3057 Metric::counter("v3.metric.limit.e", 5.0),
3058 ];
3059 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 2, 10_000);
3060 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 2, usize::MAX, None);
3061 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3062 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3063 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3064
3065 let batches = collect_v3_batches(&metrics, context, usize::MAX);
3066
3067 assert_eq!(vec![(0..2, 2), (2..4, 2), (4..5, 1)], batches);
3068 }
3069
3070 #[test]
3071 fn v3_batches_are_cut_at_the_uncompressed_byte_target() {
3072 let metrics = vec![
3073 Metric::counter("v3.byte.target.a", 1.0),
3074 Metric::counter("v3.byte.target.b", 2.0),
3075 Metric::counter("v3.byte.target.c", 3.0),
3076 Metric::counter("v3.byte.target.d", 4.0),
3077 ];
3078 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000);
3079 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3080 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3081 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3082 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3083
3084 let batches = collect_v3_batches(&metrics, context, 1);
3087
3088 assert_eq!(4, batches.len());
3089 assert!(batches.iter().all(|(_, event_count)| *event_count == 1));
3090 }
3091
3092 #[test]
3093 fn v3_serializer_stats_record_max_points_split_reason() {
3094 let metrics = vec![
3095 Metric::counter("v3.telemetry.max_points.one", [(123, 1.0), (124, 2.0)]),
3096 Metric::counter("v3.telemetry.max_points.two", [(123, 3.0), (124, 4.0)]),
3097 ];
3098 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 2);
3099 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3100 let recorder = TestRecorder::default();
3101 let _local = metrics::set_default_local_recorder(&recorder);
3102 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3103 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3104 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3105
3106 let batches = collect_v3_batches(&metrics, context, usize::MAX);
3107
3108 assert_eq!(vec![(0..1, 1), (1..2, 1)], batches);
3109 assert_eq!(
3110 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "max_points")])),
3111 Some(1)
3112 );
3113 }
3114
3115 #[test]
3116 fn v3_batches_drop_oversized_metric_and_keep_batching() {
3117 let metrics = vec![
3118 Metric::counter("v3.points.oversized.before", [(123, 1.0), (124, 2.0)]),
3119 Metric::counter(
3120 "v3.points.oversized.too_big",
3121 [(123, 3.0), (124, 4.0), (125, 5.0), (126, 6.0)],
3122 ),
3123 Metric::counter("v3.points.oversized.after", 7.0),
3124 ];
3125 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 3);
3126 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3127 let recorder = TestRecorder::default();
3128 let _local = metrics::set_default_local_recorder(&recorder);
3129 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3130 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3131 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3132
3133 let batches = collect_v3_batches(&metrics, context, usize::MAX);
3134
3135 assert_eq!(vec![(0..3, 2)], batches);
3138 assert_eq!(recorder.counter("serializer.v3_item_too_big"), Some(1));
3139 }
3140
3141 #[test]
3142 fn v3_batches_skip_zero_point_metrics_without_splitting() {
3143 let metrics = vec![
3144 Metric::counter("v3.points.zero.before", 1.0),
3145 Metric::counter("v3.points.zero.empty", &[] as &[f64]),
3146 Metric::counter("v3.points.zero.after", 2.0),
3147 ];
3148 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000);
3149 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3150 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3151 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3152 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3153
3154 let batches = collect_v3_batches(&metrics, context, usize::MAX);
3155
3156 assert_eq!(vec![(0..3, 2)], batches);
3157 }
3158
3159 #[test]
3160 fn v3_compression_ratio_targets_the_compressed_limit() {
3161 let limits = V3PayloadLimits::new(500_000, 5_000_000, 10_000, 10_000);
3162 let mut ratio = V3CompressionRatio::default();
3163
3164 let initial_target = ratio.uncompressed_target(limits);
3167 assert_eq!((500_000.0 * V3_BATCH_TARGET_MARGIN) as usize, initial_target);
3168
3169 ratio.record(4_000_000, 1_000_000);
3171 let observed_target = ratio.uncompressed_target(limits);
3172 assert!(
3173 observed_target > initial_target,
3174 "target should grow once compression is observed: {observed_target} vs {initial_target}"
3175 );
3176
3177 ratio.record(4_000_000, 40_000);
3179 ratio.record(4_000_000, 40_000);
3180 ratio.record(4_000_000, 40_000);
3181 ratio.record(4_000_000, 40_000);
3182 ratio.record(4_000_000, 40_000);
3183 assert!(ratio.uncompressed_target(limits) <= limits.max_uncompressed_size);
3184 }
3185
3186 #[tokio::test]
3187 async fn v3_split_flush_uses_payload_request_batch_headers() {
3188 let mut metrics = vec![
3189 Metric::counter("v3.headers.split.one", 1.0),
3190 Metric::counter("v3.headers.split.two", 2.0),
3191 ];
3192 let single_request = create_v3_test_request(&metrics[..1]).await;
3193 let combined_request = create_v3_test_request(&metrics).await;
3194 assert!(combined_request.compressed_len > single_request.compressed_len);
3195
3196 let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
3197 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3198 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3199 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3200 let batch_id = Uuid::now_v7();
3201 let (mut payloads_tx, mut payloads_rx) = tokio::sync::mpsc::channel(8);
3202
3203 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3204 let mut ratio = V3CompressionRatio::default();
3205 encode_and_flush_v3_series_metrics(
3206 context,
3207 &mut metrics,
3208 &mut ratio,
3209 &mut payloads_tx,
3210 Some(&batch_id),
3211 Some(MetricsPayloadInfo::v3_series()),
3212 )
3213 .await
3214 .expect("V3 metrics should flush");
3215
3216 for expected_seq in 0..2 {
3217 let payload = payloads_rx.recv().await.expect("payload should be emitted");
3218 let Payload::Http(http_payload) = payload else {
3219 panic!("expected HTTP payload");
3220 };
3221 let (_, request) = http_payload.into_parts();
3222 assert_eq!(
3223 batch_id.as_hyphenated().to_string(),
3224 request
3225 .headers()
3226 .get("X-Metrics-Request-ID")
3227 .expect("batch ID header should be present")
3228 .to_str()
3229 .expect("batch ID header should be valid")
3230 );
3231 assert_eq!(
3232 expected_seq.to_string(),
3233 request
3234 .headers()
3235 .get("X-Metrics-Request-Seq")
3236 .expect("batch sequence header should be present")
3237 .to_str()
3238 .expect("batch sequence header should be valid")
3239 );
3240 assert_eq!(
3241 "2",
3242 request
3243 .headers()
3244 .get("X-Metrics-Request-Len")
3245 .expect("batch length header should be present")
3246 .to_str()
3247 .expect("batch length header should be valid")
3248 );
3249 }
3250
3251 assert!(metrics.is_empty());
3252 }
3253
3254 #[tokio::test]
3255 async fn v3_sketch_flush_uses_split_payload_requests() {
3256 let mut metrics = vec![
3257 Metric::distribution("v3.sketch.split.one", [1.0, 2.0, 3.0]),
3258 Metric::distribution("v3.sketch.split.two", [4.0, 5.0, 6.0]),
3259 ];
3260 let single_request = create_v3_test_request(&metrics[..1]).await;
3261 let combined_request = create_v3_test_request(&metrics).await;
3262 assert!(combined_request.compressed_len > single_request.compressed_len);
3263
3264 let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
3265 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3266 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
3267 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
3268 let batch_id = Uuid::now_v7();
3269 let (mut payloads_tx, mut payloads_rx) = tokio::sync::mpsc::channel(8);
3270
3271 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
3272 let mut ratio = V3CompressionRatio::default();
3273 encode_and_flush_v3_sketch_metrics(
3274 context,
3275 &mut metrics,
3276 &mut ratio,
3277 &mut payloads_tx,
3278 Some(&batch_id),
3279 Some(MetricsPayloadInfo::v3_sketches()),
3280 )
3281 .await
3282 .expect("V3 sketches should flush");
3283
3284 for expected_seq in 0..2 {
3285 let payload = payloads_rx.recv().await.expect("payload should be emitted");
3286 let Payload::Http(http_payload) = payload else {
3287 panic!("expected HTTP payload");
3288 };
3289 let (_, request) = http_payload.into_parts();
3290 assert_eq!(V3_SKETCHES_ENDPOINT_URI, request.uri());
3291 assert_eq!(
3292 expected_seq.to_string(),
3293 request
3294 .headers()
3295 .get("X-Metrics-Request-Seq")
3296 .expect("batch sequence header should be present")
3297 .to_str()
3298 .expect("batch sequence header should be valid")
3299 );
3300 assert_eq!(
3301 "2",
3302 request
3303 .headers()
3304 .get("X-Metrics-Request-Len")
3305 .expect("batch length header should be present")
3306 .to_str()
3307 .expect("batch length header should be valid")
3308 );
3309 }
3310
3311 assert!(metrics.is_empty());
3312 }
3313
3314 #[test]
3315 fn v3_series_metric_unit_refs_are_encoded_sparsely() {
3316 let context = Context::from_static_parts("my.timer.avg", &[]);
3317 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
3318 let gauge = Metric::from_parts(context, MetricValues::gauge([1.0_f64]), metadata);
3319 let context = Context::from_static_parts("my.counter", &[]);
3320 let no_unit = Metric::from_parts(context, MetricValues::gauge([2.0_f64]), MetricMetadata::default());
3321 let context = Context::from_static_parts("my.timer.max", &[]);
3322 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
3323 let same_unit = Metric::from_parts(context, MetricValues::gauge([3.0_f64]), metadata);
3324
3325 let payload = encode_v3_metrics_batch(&[gauge, no_unit, same_unit], &SharedTagSet::default())
3326 .expect("V3 metric should encode successfully")
3327 .payload;
3328
3329 let expected_unit_dict = [
3330 0xca, 0x01, 0x0c, 0x0b, b'm', b'i', b'l', b'l', b'i', b's', b'e', b'c', b'o', b'n', b'd',
3333 ];
3334 assert!(
3335 payload
3336 .windows(expected_unit_dict.len())
3337 .any(|window| window == expected_unit_dict),
3338 "V3 payload should contain DictUnitStr field for 'millisecond', got bytes: {:?}",
3339 payload
3340 );
3341
3342 let expected_unit_ref = [
3343 0xd2, 0x01, 0x02, 0x02, 0x00, ];
3347 assert!(
3348 payload
3349 .windows(expected_unit_ref.len())
3350 .any(|window| window == expected_unit_ref),
3351 "V3 payload should contain UnitRef field for 'millisecond', got bytes: {:?}",
3352 payload
3353 );
3354 }
3355
3356 #[test]
3357 fn v3_sketch_metric_unit_not_encoded() {
3358 let context = Context::from_static_parts("my.histogram", &[]);
3359 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
3360 let histogram = Metric::from_parts(context, MetricValues::histogram([1.0_f64]), metadata);
3361
3362 let payload = encode_v3_metrics_batch(&[histogram], &SharedTagSet::default())
3363 .expect("V3 sketch metric should encode successfully")
3364 .payload;
3365
3366 assert!(
3367 !payload
3368 .windows(b"millisecond".len())
3369 .any(|window| window == b"millisecond"),
3370 "V3 sketch payload should not contain unit bytes, matching the Agent V3 sketch builder: {:?}",
3371 payload
3372 );
3373 }
3374
3375 #[test]
3376 fn v3_series_promotes_device_and_internal_resource_tags_to_resources() {
3377 let context = Context::from_static_parts(
3378 "series.resources",
3379 &[
3380 "env:prod",
3381 "device:switch1",
3382 "dd.internal.resource:pod:pod-a",
3383 "dd.internal.resource:malformed",
3384 ],
3385 );
3386 let context = context.with_host(Some(MetaString::from_static("host-a")));
3387 let metadata = MetricMetadata::default();
3388 let metric = Metric::from_parts(context, MetricValues::gauge([1.0_f64]), metadata);
3389
3390 let payload = encode_v3_metrics_batch(&[metric], &SharedTagSet::default())
3391 .expect("V3 series should encode successfully")
3392 .payload;
3393
3394 assert_contains_bytes(&payload, b"env:prod");
3395 assert!(!contains_bytes(&payload, b"device:switch1"));
3396 assert!(!contains_bytes(&payload, b"dd.internal.resource:pod:pod-a"));
3397 assert!(!contains_bytes(&payload, b"dd.internal.resource:malformed"));
3398
3399 let expected_resource_dict = [
3400 0x22, 0x25, 0x04, b'h', b'o', b's', b't', 0x06, b'h', b'o', b's', b't', b'-', b'a', 0x06, b'd', b'e', b'v', b'i', b'c',
3403 b'e', 0x07, b's', b'w', b'i', b't', b'c', b'h', b'1', 0x03, b'p', b'o', b'd', 0x05, b'p', b'o', b'd', b'-',
3404 b'a',
3405 ];
3406 assert_contains_bytes(&payload, &expected_resource_dict);
3407 }
3408
3409 #[test]
3410 fn v3_series_promotes_additional_and_origin_resource_tags_without_empty_host() {
3411 let context = Context::from_static_parts("series.additional_origin_resources", &["env:prod"])
3412 .with_origin_tags(tag_set(["dd.internal.resource:pod:pod-origin"]));
3413 let additional_tags = SharedTagSet::from(tag_set([
3414 "team:core",
3415 "device:switch1",
3416 "dd.internal.resource:container:container-a",
3417 ]));
3418 let context = context.with_host(Some(MetaString::empty()));
3419 let metadata = MetricMetadata::default();
3420 let metric = Metric::from_parts(context, MetricValues::gauge([1.0_f64]), metadata);
3421
3422 let payload = encode_v3_metrics_batch(&[metric], &additional_tags)
3423 .expect("V3 series should encode successfully")
3424 .payload;
3425
3426 assert_contains_bytes(&payload, b"env:prod");
3427 assert_contains_bytes(&payload, b"team:core");
3428 assert!(!contains_bytes(&payload, b"device:switch1"));
3429 assert!(!contains_bytes(&payload, b"dd.internal.resource:container:container-a"));
3430 assert!(!contains_bytes(&payload, b"dd.internal.resource:pod:pod-origin"));
3431
3432 let expected_resource_dict = [
3433 0x22, 0x34, 0x06, b'd', b'e', b'v', b'i', b'c', b'e', 0x07, b's', b'w', b'i', b't', b'c', b'h', b'1', 0x03, b'p', b'o',
3436 b'd', 0x0a, b'p', b'o', b'd', b'-', b'o', b'r', b'i', b'g', b'i', b'n', 0x09, b'c', b'o', b'n', b't', b'a',
3437 b'i', b'n', b'e', b'r', 0x0b, b'c', b'o', b'n', b't', b'a', b'i', b'n', b'e', b'r', b'-', b'a',
3438 ];
3439 assert_contains_bytes(&payload, &expected_resource_dict);
3440 assert!(!contains_bytes(&payload, b"host"));
3441 }
3442
3443 #[test]
3444 fn v3_sketch_keeps_device_and_internal_resource_tags_as_tags() {
3445 let context = Context::from_static_parts(
3446 "sketch.resources",
3447 &["env:prod", "device:switch1", "dd.internal.resource:pod:pod-a"],
3448 );
3449 let context = context.with_host(Some(MetaString::from_static("host-a")));
3450 let metadata = MetricMetadata::default();
3451 let metric = Metric::from_parts(context, MetricValues::histogram([1.0_f64]), metadata);
3452
3453 let payload = encode_v3_metrics_batch(&[metric], &SharedTagSet::default())
3454 .expect("V3 sketch should encode successfully")
3455 .payload;
3456
3457 assert_contains_bytes(&payload, b"env:prod");
3458 assert_contains_bytes(&payload, b"device:switch1");
3459 assert_contains_bytes(&payload, b"dd.internal.resource:pod:pod-a");
3460
3461 let expected_resource_dict = [
3462 0x22, 0x0c, 0x04, b'h', b'o', b's', b't', 0x06, b'h', b'o', b's', b't', b'-', b'a',
3465 ];
3466 assert_contains_bytes(&payload, &expected_resource_dict);
3467 }
3468
3469 struct RequestBuilderScenario {
3471 v2_series_max_metrics: Option<usize>,
3472 series_mode: MetricsEncoderMode,
3473 sketches_mode: MetricsEncoderMode,
3474 payload_limits: V3PayloadLimits,
3475 shadow_series_endpoint_uri: String,
3476 series_shadow_config: SeriesShadowConfig,
3477 flush_timeout: Duration,
3478 }
3479
3480 impl RequestBuilderScenario {
3481 fn new(series_mode: MetricsEncoderMode, sketches_mode: MetricsEncoderMode) -> Self {
3484 Self {
3485 v2_series_max_metrics: None,
3486 series_mode,
3487 sketches_mode,
3488 payload_limits: V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000),
3489 shadow_series_endpoint_uri: "/api/intake/metrics/v3beta/series".to_string(),
3490 series_shadow_config: SeriesShadowConfig::new(0.0),
3491 flush_timeout: Duration::from_millis(10),
3492 }
3493 }
3494
3495 fn with_v2_series_builder(mut self, max_metrics_per_payload: usize) -> Self {
3497 self.v2_series_max_metrics = Some(max_metrics_per_payload);
3498 self
3499 }
3500
3501 fn with_max_points_per_payload(mut self, max_points_per_payload: usize) -> Self {
3503 self.payload_limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, max_points_per_payload);
3504 self
3505 }
3506
3507 fn with_series_shadow(mut self, sample_rate: f64, shadow_series_endpoint_uri: &str) -> Self {
3509 self.series_shadow_config = SeriesShadowConfig::new(sample_rate);
3510 self.shadow_series_endpoint_uri = shadow_series_endpoint_uri.to_string();
3511 self
3512 }
3513
3514 fn with_flush_timeout(mut self, flush_timeout: Duration) -> Self {
3516 self.flush_timeout = flush_timeout;
3517 self
3518 }
3519
3520 async fn spawn(self) -> RequestBuilderHarness {
3523 let v2_series_builder = match self.v2_series_max_metrics {
3524 Some(max_metrics_per_payload) => {
3525 let v2_endpoint_config = EndpointConfiguration::new(
3526 CompressionScheme::noop(),
3527 max_metrics_per_payload,
3528 usize::MAX,
3529 None,
3530 );
3531 Some(
3532 v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
3533 .await
3534 .expect("V2 request builder should be created"),
3535 )
3536 }
3537 None => None,
3538 };
3539 let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
3540 let metrics_builder = MetricsBuilder::default();
3541 let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
3542 let serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
3543 let (events_tx, events_rx) = tokio::sync::mpsc::channel(1);
3544 let (payloads_tx, payloads_rx) = tokio::sync::mpsc::channel(8);
3545
3546 let handle = tokio::spawn(run_request_builder(
3547 v2_series_builder,
3548 None,
3549 self.series_mode,
3550 self.sketches_mode,
3551 V3RuntimeConfig {
3552 endpoint_config: v3_endpoint_config,
3553 payload_limits: self.payload_limits,
3554 series_endpoint_uri: V3_SERIES_ENDPOINT_URI.to_string(),
3555 shadow_series_endpoint_uri: self.shadow_series_endpoint_uri,
3556 series_shadow_config: self.series_shadow_config,
3557 serializer_telemetry,
3558 },
3559 telemetry,
3560 events_rx,
3561 payloads_tx,
3562 self.flush_timeout,
3563 false,
3564 ));
3565
3566 RequestBuilderHarness {
3567 events_tx,
3568 payloads_rx,
3569 handle,
3570 }
3571 }
3572 }
3573
3574 struct RequestBuilderHarness {
3576 events_tx: tokio::sync::mpsc::Sender<EventsBuffer>,
3577 payloads_rx: tokio::sync::mpsc::Receiver<PayloadsBuffer>,
3578 handle: tokio::task::JoinHandle<Result<(), GenericError>>,
3579 }
3580
3581 impl RequestBuilderHarness {
3582 async fn push_metrics(&self, metrics: impl IntoIterator<Item = Metric>) {
3584 let mut events = EventsBuffer::default();
3585 for metric in metrics {
3586 assert!(
3587 events.try_push(Event::Metric(metric)).is_none(),
3588 "event buffer should accept metric"
3589 );
3590 }
3591 self.events_tx
3592 .send(events)
3593 .await
3594 .expect("events should be sent to request builder");
3595 }
3596
3597 async fn next_http_request(&mut self) -> (PayloadMetadata, Request<FrozenChunkedBytesBuffer>) {
3599 let payload = timeout(Duration::from_secs(1), self.payloads_rx.recv())
3600 .await
3601 .expect("payload should arrive before timeout")
3602 .expect("payload channel should remain open");
3603 match payload {
3604 Payload::Http(http_payload) => http_payload.into_parts(),
3605 _ => panic!("expected HTTP payload"),
3606 }
3607 }
3608
3609 async fn assert_no_payload_within(&mut self, window: Duration) {
3611 assert!(
3612 timeout(window, self.payloads_rx.recv()).await.is_err(),
3613 "no payload should be flushed within {window:?}"
3614 );
3615 }
3616
3617 async fn shutdown(self) {
3619 drop(self.events_tx);
3620 self.handle
3621 .await
3622 .expect("request builder task should complete")
3623 .expect("request builder should stop cleanly");
3624 }
3625 }
3626
3627 #[tokio::test]
3628 async fn validation_split_flush_assigns_batch_id_to_carried_metric() {
3629 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::Validation, MetricsEncoderMode::V2Only)
3633 .with_v2_series_builder(1)
3634 .spawn()
3635 .await;
3636
3637 harness
3638 .push_metrics([
3639 Metric::counter("validation.split.one", 1.0),
3640 Metric::counter("validation.split.two", 2.0),
3641 ])
3642 .await;
3643
3644 let mut flushed_requests = Vec::new();
3645 for _ in 0..4 {
3646 let (_, request) = harness.next_http_request().await;
3647 let batch_id = request
3648 .headers()
3649 .get("X-Metrics-Request-ID")
3650 .expect("validation batch ID header should be present")
3651 .to_str()
3652 .expect("validation batch ID should be valid header text")
3653 .to_string();
3654 flushed_requests.push((request.uri().to_string(), batch_id));
3655 }
3656
3657 assert_eq!("/api/v2/series", flushed_requests[0].0);
3658 assert_eq!(V3_SERIES_ENDPOINT_URI, flushed_requests[1].0);
3659 assert_eq!("/api/v2/series", flushed_requests[2].0);
3660 assert_eq!(V3_SERIES_ENDPOINT_URI, flushed_requests[3].0);
3661
3662 assert_eq!(flushed_requests[0].1, flushed_requests[1].1);
3663 assert_eq!(flushed_requests[2].1, flushed_requests[3].1);
3664 assert_ne!(flushed_requests[0].1, flushed_requests[2].1);
3665
3666 harness.shutdown().await;
3667 }
3668
3669 #[tokio::test]
3670 async fn authoritative_v3_flushes_previous_point_limit_batch() {
3671 let recorder = TestRecorder::default();
3675 let _local = metrics::set_default_local_recorder(&recorder);
3676
3677 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V3Enabled, MetricsEncoderMode::V2Only)
3678 .with_max_points_per_payload(3)
3679 .with_flush_timeout(Duration::from_millis(250))
3680 .spawn()
3681 .await;
3682
3683 harness
3684 .push_metrics([
3685 Metric::counter("authoritative.v3.points.one", [(123, 1.0), (124, 2.0)]),
3686 Metric::counter("authoritative.v3.points.two", [(123, 3.0), (124, 4.0)]),
3687 ])
3688 .await;
3689
3690 let (_, request) = harness.next_http_request().await;
3692 assert_eq!(V3_SERIES_ENDPOINT_URI, request.uri());
3693 harness.assert_no_payload_within(Duration::from_millis(50)).await;
3694
3695 let (_, request) = harness.next_http_request().await;
3697 assert_eq!(V3_SERIES_ENDPOINT_URI, request.uri());
3698 assert_eq!(
3699 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "max_points")])),
3700 Some(1)
3701 );
3702
3703 harness.shutdown().await;
3704 }
3705
3706 #[tokio::test]
3707 async fn authoritative_v3_sketches_flush_previous_point_limit_batch() {
3708 let recorder = TestRecorder::default();
3711 let _local = metrics::set_default_local_recorder(&recorder);
3712
3713 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V2Only, MetricsEncoderMode::V3Enabled)
3714 .with_max_points_per_payload(3)
3715 .with_flush_timeout(Duration::from_millis(250))
3716 .spawn()
3717 .await;
3718
3719 harness
3720 .push_metrics([
3721 Metric::distribution("authoritative.v3.sketch.points.one", [(123, 1.0), (124, 2.0)]),
3722 Metric::distribution("authoritative.v3.sketch.points.two", [(123, 3.0), (124, 4.0)]),
3723 ])
3724 .await;
3725
3726 for stage in ["point-limit split", "timeout"] {
3727 let (_, request) = harness.next_http_request().await;
3728 assert_eq!(
3729 V3_SKETCHES_ENDPOINT_URI,
3730 request.uri(),
3731 "{stage} sketches payload should target the V3 sketches endpoint"
3732 );
3733 }
3734 assert_eq!(
3735 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "max_points")])),
3736 Some(1)
3737 );
3738
3739 harness.shutdown().await;
3740 }
3741
3742 #[tokio::test]
3743 async fn authoritative_v3_does_not_flush_on_v2_boundary() {
3744 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V3Enabled, MetricsEncoderMode::V2Only)
3747 .with_v2_series_builder(1)
3748 .with_flush_timeout(Duration::from_millis(250))
3749 .spawn()
3750 .await;
3751
3752 harness
3753 .push_metrics([
3754 Metric::counter("authoritative.v3.decouple.one", 1.0),
3755 Metric::counter("authoritative.v3.decouple.two", 2.0),
3756 ])
3757 .await;
3758
3759 let (_, request) = harness.next_http_request().await;
3760 assert_eq!("/api/v2/series", request.uri());
3761 harness.assert_no_payload_within(Duration::from_millis(50)).await;
3762
3763 let mut timeout_flush_uris = Vec::new();
3764 for _ in 0..2 {
3765 let (_, request) = harness.next_http_request().await;
3766 timeout_flush_uris.push(request.uri().to_string());
3767 }
3768 assert_eq!(2, timeout_flush_uris.len());
3769 assert!(timeout_flush_uris.iter().any(|uri| uri == "/api/v2/series"));
3770 assert!(timeout_flush_uris.iter().any(|uri| uri == V3_SERIES_ENDPOINT_URI));
3771
3772 harness.shutdown().await;
3773 }
3774
3775 #[tokio::test]
3776 async fn shadow_sampled_series_flush_sends_v2_and_v3_beta_with_same_batch_id() {
3777 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V2Only, MetricsEncoderMode::V2Only)
3780 .with_v2_series_builder(10_000)
3781 .with_series_shadow(1.0, "/api/intake/metrics/v3beta/custom")
3782 .spawn()
3783 .await;
3784
3785 harness.push_metrics([Metric::counter("shadow.sampled", 1.0)]).await;
3786
3787 let mut flushed_requests = Vec::new();
3788 for _ in 0..2 {
3789 let (metadata, request) = harness.next_http_request().await;
3790 let payload_info = *metadata
3791 .get::<MetricsPayloadInfo>()
3792 .expect("metrics payload info should be present");
3793 let batch_id = request
3794 .headers()
3795 .get("X-Metrics-Request-ID")
3796 .expect("shadow batch ID header should be present")
3797 .to_str()
3798 .expect("shadow batch ID should be valid header text")
3799 .to_string();
3800 flushed_requests.push((request.uri().to_string(), payload_info, batch_id));
3801 }
3802
3803 assert_eq!("/api/v2/series", flushed_requests[0].0);
3804 assert_eq!(MetricsPayloadInfo::v2_shadow_series(), flushed_requests[0].1);
3805 assert_eq!("/api/intake/metrics/v3beta/custom", flushed_requests[1].0);
3806 assert_eq!(MetricsPayloadInfo::v3_shadow_series(), flushed_requests[1].1);
3807 assert_eq!(flushed_requests[0].2, flushed_requests[1].2);
3808
3809 harness.shutdown().await;
3810 }
3811
3812 #[tokio::test]
3813 async fn shadow_sample_rate_zero_sends_only_v2_without_validation_headers() {
3814 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V2Only, MetricsEncoderMode::V2Only)
3817 .with_v2_series_builder(10_000)
3818 .spawn()
3819 .await;
3820
3821 harness.push_metrics([Metric::counter("shadow.disabled", 1.0)]).await;
3822
3823 let (_, request) = harness.next_http_request().await;
3824 assert_eq!("/api/v2/series", request.uri());
3825 assert!(!request.headers().contains_key("X-Metrics-Request-ID"));
3826 harness.assert_no_payload_within(Duration::from_millis(50)).await;
3827
3828 harness.shutdown().await;
3829 }
3830
3831 fn contains_bytes(haystack: &[u8], needle: &[u8]) -> bool {
3832 haystack.windows(needle.len()).any(|window| window == needle)
3833 }
3834
3835 fn assert_contains_bytes(haystack: &[u8], needle: &[u8]) {
3836 assert!(
3837 contains_bytes(haystack, needle),
3838 "expected payload to contain bytes {:?}, got {:?}",
3839 needle,
3840 haystack
3841 );
3842 }
3843
3844 fn tag_set<const N: usize>(tags: [&'static str; N]) -> TagSet {
3845 tags.into_iter().map(Tag::from_static).collect()
3846 }
3847
3848 #[tokio::test]
3853 async fn v2_series_builder_enforces_max_series_points_per_payload() {
3854 let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, 10_000, None);
3855 let mut builder = v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
3856 .await
3857 .expect("V2 request builder should be created");
3858 builder
3859 .with_len_limits(usize::MAX, usize::MAX)
3860 .expect("byte limits should be accepted");
3861
3862 let mut total_points = 0;
3863 for i in 0..6_000u32 {
3864 let metric = Metric::gauge(
3865 Context::from_parts(MetaString::from(format!("g{i}")), TagSet::default()),
3866 [(1u64, 1.0), (2u64, 2.0)],
3867 );
3868 let mut pending = Some(metric);
3869 while let Some(metric) = pending.take() {
3870 if let Some(returned) = builder.encode(metric).await.expect("encode should not error") {
3871 for request in builder.flush().await {
3872 let (_, points, _) = request.expect("request should build");
3873 assert!(
3874 points <= 10_000,
3875 "payload carried {} points, over the 10000 limit",
3876 points
3877 );
3878 total_points += points;
3879 }
3880 pending = Some(returned);
3881 }
3882 }
3883 }
3884 for request in builder.flush().await {
3885 let (_, points, _) = request.expect("request should build");
3886 assert!(
3887 points <= 10_000,
3888 "final payload carried {} points, over the 10000 limit",
3889 points
3890 );
3891 total_points += points;
3892 }
3893
3894 assert_eq!(
3895 total_points, 12_000,
3896 "all points should be emitted across the split payloads"
3897 );
3898 }
3899}
3900
3901#[cfg(test)]
3902mod config_smoke {
3903 use datadog_agent_config_testing::config_registry::structs;
3904 use datadog_agent_config_testing::run_config_smoke_tests;
3905 use serde_json::json;
3906
3907 use super::DatadogMetricsConfiguration;
3908 use crate::config::{DatadogRemapper, KEY_ALIASES};
3909
3910 #[tokio::test]
3911 async fn smoke_test() {
3912 run_config_smoke_tests(
3913 structs::DATADOG_METRICS_CONFIGURATION,
3914 &[
3915 "serializer_experimental_use_v3_api.sketches.beta_route",
3916 "serializer_experimental_use_v3_api.sketches.shadow_sample_rate",
3917 "serializer_experimental_use_v3_api.sketches.shadow_sites",
3918 "serializer_experimental_use_v3_api.sketches.use_beta",
3919 ],
3920 json!({}),
3921 |cfg| {
3922 cfg.as_typed::<DatadogMetricsConfiguration>()
3923 .expect("DatadogMetricsConfiguration should deserialize")
3924 },
3925 KEY_ALIASES,
3926 DatadogRemapper::from_env_vars,
3927 )
3928 .await
3929 }
3930}
3931
3932#[cfg(test)]
3933mod use_v2_api_series_default {
3934 use saluki_config::ConfigurationLoader;
3935 use serde_json::json;
3936
3937 use super::{v2, DatadogMetricsConfiguration};
3938 use crate::{common::datadog::clamp_payload_limits, config::KEY_ALIASES};
3939
3940 #[tokio::test]
3942 async fn defaults_to_true_when_absent() {
3943 let cfg = ConfigurationLoader::default()
3944 .with_key_aliases(KEY_ALIASES)
3945 .add_providers([figment::providers::Serialized::defaults(json!({}))])
3946 .into_generic()
3947 .await
3948 .expect("config should load");
3949 let parsed: DatadogMetricsConfiguration = cfg.as_typed().expect("should deserialize");
3950 assert!(parsed.use_v2_api_series);
3951 }
3952
3953 #[tokio::test]
3954 async fn deserializes_payload_limit_keys() {
3955 let cfg = ConfigurationLoader::default()
3956 .with_key_aliases(KEY_ALIASES)
3957 .add_providers([figment::providers::Serialized::defaults(json!({
3958 "serializer_max_payload_size": 4321,
3959 "serializer_max_uncompressed_payload_size": 8765,
3960 "serializer_max_series_payload_size": 1234,
3961 "serializer_max_series_uncompressed_payload_size": 5678,
3962 }))])
3963 .into_generic()
3964 .await
3965 .expect("config should load");
3966 let parsed: DatadogMetricsConfiguration = cfg.as_typed().expect("should deserialize");
3967
3968 assert_eq!(parsed.max_payload_size, 4321);
3969 assert_eq!(parsed.max_uncompressed_payload_size, 8765);
3970 assert_eq!(parsed.max_series_payload_size, 1234);
3971 assert_eq!(parsed.max_series_uncompressed_payload_size, 5678);
3972 }
3973
3974 #[tokio::test]
3975 async fn deserializes_max_series_points_per_payload() {
3976 let cfg = ConfigurationLoader::default()
3978 .with_key_aliases(KEY_ALIASES)
3979 .add_providers([figment::providers::Serialized::defaults(json!({}))])
3980 .into_generic()
3981 .await
3982 .expect("config should load");
3983 let parsed: DatadogMetricsConfiguration = cfg.as_typed().expect("should deserialize");
3984 assert_eq!(parsed.max_series_points_per_payload, 10_000);
3985 assert_eq!(parsed.v3_payload_limits().max_points_per_payload, 10_000);
3986
3987 let cfg = ConfigurationLoader::default()
3989 .with_key_aliases(KEY_ALIASES)
3990 .add_providers([figment::providers::Serialized::defaults(json!({
3991 "serializer_max_series_points_per_payload": 500,
3992 }))])
3993 .into_generic()
3994 .await
3995 .expect("config should load");
3996 let parsed: DatadogMetricsConfiguration = cfg.as_typed().expect("should deserialize");
3997 assert_eq!(parsed.max_series_points_per_payload, 500);
3998 assert_eq!(parsed.v3_payload_limits().max_points_per_payload, 500);
3999 }
4000
4001 #[test]
4002 fn clamps_series_payload_limit_keys_to_api_limits() {
4003 let (uncompressed_limit, compressed_limit) = clamp_payload_limits(
4004 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT + 1,
4005 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT + 1,
4006 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT,
4007 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT,
4008 );
4009 assert_eq!(uncompressed_limit, v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT);
4010 assert_eq!(compressed_limit, v2::SERIES_V2_COMPRESSED_SIZE_LIMIT);
4011
4012 let (uncompressed_limit, compressed_limit) = clamp_payload_limits(
4013 5678,
4014 1234,
4015 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT,
4016 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT,
4017 );
4018 assert_eq!(uncompressed_limit, 5678);
4019 assert_eq!(compressed_limit, 1234);
4020 }
4021}