1use std::{collections::HashMap, collections::VecDeque, ops::Range, sync::Arc, time::Duration};
2
3use agent_data_plane_config::shared::{SharedConfiguration, V3SeriesMode};
4use async_trait::async_trait;
5use ddsketch::DDSketch;
6use http::{HeaderValue, Method, Request};
7use protobuf::{rt::WireType, CodedOutputStream};
8use saluki_common::{
9 buf::{ChunkedBytesBuffer, FrozenChunkedBytesBuffer},
10 iter::ReusableDeduplicator,
11};
12use saluki_core::{
13 accounting::{MemoryBounds, MemoryBoundsBuilder},
14 components::{encoders::*, BuildContext},
15 data_model::{
16 event::{
17 metric::{Metric, MetricOrigin, MetricValues},
18 EventType,
19 },
20 payload::{HttpPayload, Payload, PayloadMetadata, PayloadType},
21 tags::{SharedTagSet, Tag},
22 },
23 observability::ComponentMetricsExt as _,
24 runtime,
25 topology::{EventsBuffer, PayloadsBuffer},
26};
27use saluki_error::{generic_error, ErrorContext as _, GenericError};
28use saluki_io::compression::{CompressionScheme, Compressor};
29use saluki_metrics::MetricsBuilder;
30use tokio::{io::AsyncWriteExt as _, select, sync::mpsc, time::sleep};
31use tracing::{debug, error, warn};
32use url::Url;
33
34use self::v3::{
35 V3EncodedMetrics, V3EncodedRequest, V3EncoderStats, V3MetricType, V3PayloadLimits, V3PayloadRequest,
36 V3PayloadSplitReason, V3SerializerTelemetry, V3Writer,
37};
38use crate::{
39 common::datadog::{
40 clamp_payload_limits,
41 config::OpwMetricsConfiguration,
42 endpoints::{
43 resolve_additional_endpoints, series_v3_config_can_enable_v3, EndpointV3Settings, ResolvedEndpoint,
44 V3EndpointConfig,
45 },
46 io::RB_BUFFER_CHUNK_SIZE,
47 protocol::{MetricsEndpointRouting, MetricsPayloadInfo, UseV3ApiConfig, UseV3ApiSeriesConfig, V3ApiConfig},
48 request_builder::{RequestBuilder, RequestBuilderError},
49 telemetry::ComponentTelemetry,
50 DEFAULT_SERIALIZER_COMPRESSED_SIZE_LIMIT, DEFAULT_SERIALIZER_UNCOMPRESSED_SIZE_LIMIT, METRICS_SERIES_V3_PATH,
51 METRICS_SKETCHES_V3_PATH,
52 },
53 encoders::datadog::metrics::v2::MetricsEndpointEncoder,
54};
55
56mod endpoint;
57use self::endpoint::{EndpointConfiguration, MetricsEndpoint};
58
59mod v1;
60mod v2;
61mod v3;
62
63const V3_SKETCHES_ENDPOINT_URI: &str = METRICS_SKETCHES_V3_PATH;
64
65#[derive(Clone, Copy, Debug, PartialEq, Eq)]
67enum MetricsEncoderMode {
68 V2Only,
70 V3Enabled,
73}
74
75impl MetricsEncoderMode {
76 fn from_config(use_v3: bool) -> Self {
77 match use_v3 {
78 false => Self::V2Only,
79 true => Self::V3Enabled,
80 }
81 }
82
83 fn needs_v3(self) -> bool {
84 matches!(self, Self::V3Enabled)
85 }
86
87 fn needs_tagging(self) -> bool {
88 matches!(self, Self::V3Enabled)
89 }
90}
91
92fn metrics_encoder_mode_for_config(use_v3: bool, metrics_v3_disabled_by_compressor: bool) -> MetricsEncoderMode {
93 let use_v3 = use_v3 && !metrics_v3_disabled_by_compressor;
94 MetricsEncoderMode::from_config(use_v3)
95}
96
97fn selected_metrics_primary_v3_override(opw_metrics: &OpwMetricsConfiguration) -> Option<bool> {
98 selected_metrics_primary_endpoint(opw_metrics).map(|(_, use_v3_series)| use_v3_series)
99}
100
101fn selected_metrics_primary_endpoint(opw_metrics: &OpwMetricsConfiguration) -> Option<(&str, bool)> {
102 let selected = opw_metrics.selected_endpoint()?;
103 let url = selected.url.trim();
104 metrics_primary_url_can_resolve(url).then_some((url, selected.use_v3_series))
105}
106
107fn metrics_primary_url_can_resolve(url: &str) -> bool {
108 if url.is_empty() {
109 return false;
110 }
111 if url.starts_with("http://") || url.starts_with("https://") {
112 Url::parse(url).is_ok_and(|url| url.host_str().is_some())
113 } else {
114 Url::parse(&format!("https://{url}")).is_ok_and(|url| url.host_str().is_some())
115 }
116}
117
118fn series_v3_can_be_enabled_for_config(
119 use_v2_api_series: bool, metrics_primary_v3_override: Option<bool>, has_additional_endpoints: bool,
120 series_config: &UseV3ApiSeriesConfig,
121) -> bool {
122 use_v2_api_series
123 && (metrics_primary_v3_override == Some(true)
124 || ((metrics_primary_v3_override != Some(false) || has_additional_endpoints)
125 && series_v3_config_can_enable_v3(series_config)))
126}
127
128#[derive(Clone)]
132#[cfg_attr(test, derive(Debug, PartialEq))]
133pub struct DatadogMetricsConfiguration {
134 max_metrics_per_payload: usize,
138
139 max_payload_size: usize,
147
148 max_uncompressed_payload_size: usize,
156
157 max_series_payload_size: usize,
164
165 max_series_uncompressed_payload_size: usize,
173
174 max_series_points_per_payload: usize,
181
182 flush_timeout: Duration,
190
191 compressor_kind: String,
193
194 zstd_compressor_level: i32,
196
197 use_v2_series_api: bool,
203
204 log_payloads: bool,
208
209 additional_tags: Option<SharedTagSet>,
211
212 v3_api: V3ApiConfig,
216
217 use_v3_api: UseV3ApiConfig,
219
220 opw_metrics: OpwMetricsConfiguration,
222
223 primary_endpoint: String,
225
226 additional_endpoints: HashMap<String, Vec<String>>,
228
229 endpoint_routing: Option<Arc<MetricsEndpointRouting>>,
231}
232
233impl DatadogMetricsConfiguration {
234 pub fn from_configuration(shared: &SharedConfiguration) -> Self {
236 let endpoints = &shared.endpoints;
237 let metrics = &shared.metrics_encoding;
238
239 Self {
240 max_metrics_per_payload: metrics.max_metrics_per_payload,
241 max_payload_size: metrics.max_payload_size,
242 max_uncompressed_payload_size: metrics.max_uncompressed_payload_size,
243 max_series_payload_size: metrics.max_series_payload_size,
244 max_series_uncompressed_payload_size: metrics.max_series_uncompressed_payload_size,
245 max_series_points_per_payload: metrics.max_series_points_per_payload,
246 flush_timeout: metrics.flush_timeout,
247 compressor_kind: endpoints.compression.compressor_kind.clone(),
248 zstd_compressor_level: endpoints.compression.effective_zstd_level(),
249 use_v2_series_api: metrics.use_v2_series_api,
250 log_payloads: metrics.log_payloads,
251 additional_tags: None,
252 v3_api: (&metrics.v3_api).into(),
253 use_v3_api: UseV3ApiConfig { series: metrics.into() },
254 opw_metrics: OpwMetricsConfiguration::from_configuration(endpoints),
255 primary_endpoint: endpoints.primary_endpoint(),
256 additional_endpoints: endpoints.additional_endpoints.clone(),
257 endpoint_routing: None,
258 }
259 }
260
261 pub fn with_additional_tags(mut self, additional_tags: SharedTagSet) -> Self {
263 self.additional_tags = Some(additional_tags);
264 self
265 }
266
267 pub fn with_endpoint_routing(mut self, endpoint_routing: MetricsEndpointRouting) -> Self {
269 self.endpoint_routing = Some(Arc::new(endpoint_routing));
270 self
271 }
272
273 pub fn with_metrics_endpoint_override(mut self, dd_url: String) -> Self {
278 self.primary_endpoint = dd_url;
279 self.additional_endpoints = HashMap::new();
280 self.opw_metrics.disable();
281 self
282 }
283
284 pub fn with_v2_series_only(mut self) -> Self {
288 self.use_v3_api.series.enabled = V3SeriesMode::Disabled;
289 self.use_v3_api.series.endpoints.clear();
290 self.opw_metrics.clear_v3_series_overrides();
291 self
292 }
293
294 fn v3_payload_limits(&self) -> V3PayloadLimits {
295 V3PayloadLimits::new(
296 self.max_series_payload_size,
297 self.max_series_uncompressed_payload_size,
298 self.max_metrics_per_payload,
299 self.max_series_points_per_payload,
300 )
301 }
302
303 fn endpoint_v3_settings(
304 &self, endpoint: &ResolvedEndpoint, metrics_primary_v3_override: Option<bool>,
305 ) -> EndpointV3Settings {
306 EndpointV3Settings::from_v3_config(V3EndpointConfig {
307 configured_endpoint: endpoint.configured_endpoint(),
308 series_config: &self.use_v3_api.series,
309 metrics_primary_v3_override,
310 serializer_v3_sketches_endpoints: &self.v3_api.sketches.endpoints,
311 })
312 }
313
314 fn any_series_endpoint_matches(
315 &self, mut predicate: impl FnMut(&EndpointV3Settings) -> bool,
316 ) -> Result<bool, GenericError> {
317 if let Some((metrics_primary_url, metrics_primary_v3_override)) =
318 selected_metrics_primary_endpoint(&self.opw_metrics)
319 {
320 let metrics_primary = ResolvedEndpoint::from_raw_endpoint(metrics_primary_url, "")
321 .error_context("Failed parsing/resolving the metrics primary destination endpoint.")?;
322 let settings = self.endpoint_v3_settings(&metrics_primary, Some(metrics_primary_v3_override));
323 if self.routes_to_endpoint(&self.primary_endpoint) && predicate(&settings) {
325 return Ok(true);
326 }
327 } else {
328 let primary = ResolvedEndpoint::from_raw_endpoint(&self.primary_endpoint, "")
329 .error_context("Failed parsing/resolving the primary destination endpoint.")?;
330 let settings = self.endpoint_v3_settings(&primary, None);
331 if self.routes_to_endpoint(&self.primary_endpoint) && predicate(&settings) {
332 return Ok(true);
333 }
334 }
335
336 for endpoint in resolve_additional_endpoints(&self.additional_endpoints)
337 .error_context("Failed parsing/resolving the additional destination endpoints.")?
338 {
339 let settings = self.endpoint_v3_settings(&endpoint, None);
340 if self.routes_to_endpoint(endpoint.configured_endpoint()) && predicate(&settings) {
341 return Ok(true);
342 }
343 }
344
345 Ok(false)
346 }
347
348 fn routes_to_endpoint(&self, policy_endpoint: &str) -> bool {
349 self.endpoint_routing
350 .as_ref()
351 .is_none_or(|routing| routing.should_route_to(policy_endpoint))
352 }
353
354 fn requires_v2_series(&self, metrics_v3_disabled_by_compressor: bool) -> Result<bool, GenericError> {
355 if !self.use_v2_series_api || metrics_v3_disabled_by_compressor {
356 return Ok(true);
357 }
358
359 self.any_series_endpoint_matches(|settings| !settings.use_v3_series)
360 }
361
362 fn requires_v3_series(&self, metrics_v3_disabled_by_compressor: bool) -> Result<bool, GenericError> {
363 if metrics_v3_disabled_by_compressor {
364 return Ok(false);
365 }
366
367 let metrics_primary_v3_override = selected_metrics_primary_v3_override(&self.opw_metrics);
368 if !series_v3_can_be_enabled_for_config(
369 self.use_v2_series_api,
370 metrics_primary_v3_override,
371 !self.additional_endpoints.is_empty(),
372 &self.use_v3_api.series,
373 ) {
374 return Ok(false);
375 }
376
377 self.any_series_endpoint_matches(|settings| settings.use_v3_series)
378 }
379}
380
381#[async_trait]
382impl EncoderBuilder for DatadogMetricsConfiguration {
383 fn input_event_type(&self) -> EventType {
384 EventType::Metric
385 }
386
387 fn output_payload_type(&self) -> PayloadType {
388 PayloadType::Http
389 }
390
391 async fn build(&self, context: BuildContext) -> Result<Box<dyn Encoder + Send>, GenericError> {
392 let metrics_builder = MetricsBuilder::from_component_context(context.component_context());
393 let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
394 let v3_serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
395
396 let v2_compression_scheme = CompressionScheme::new(&self.compressor_kind, self.zstd_compressor_level);
397 let v3_compression_scheme = if self.v3_api.compression_level > 0 {
398 CompressionScheme::new(&self.compressor_kind, self.v3_api.compression_level)
399 } else {
400 v2_compression_scheme
401 };
402 let series_endpoint_uri = METRICS_SERIES_V3_PATH.to_string();
403 let payload_limits = self.v3_payload_limits();
404
405 let v2_endpoint_config = EndpointConfiguration::new(
406 v2_compression_scheme,
407 self.max_metrics_per_payload,
408 self.max_series_points_per_payload,
409 self.additional_tags.clone(),
410 );
411 let endpoint_config = EndpointConfiguration::new(
412 v3_compression_scheme,
413 self.max_metrics_per_payload,
414 usize::MAX,
417 self.additional_tags.clone(),
418 );
419
420 let metrics_v3_disabled_by_compressor = matches!(v3_compression_scheme, CompressionScheme::Zlib(_));
422 let use_v3_series = self.requires_v3_series(metrics_v3_disabled_by_compressor)?;
423 let series_mode = metrics_encoder_mode_for_config(use_v3_series, metrics_v3_disabled_by_compressor);
424 let sketches_mode =
425 metrics_encoder_mode_for_config(self.v3_api.use_v3_sketches(), metrics_v3_disabled_by_compressor);
426 let series_endpoint = if self.use_v2_series_api {
427 MetricsEndpoint::SeriesV2
428 } else {
429 MetricsEndpoint::SeriesV1
430 };
431 let v3_runtime_config = V3RuntimeConfig {
432 endpoint_config,
433 payload_limits,
434 series_endpoint_uri,
435 serializer_telemetry: v3_serializer_telemetry,
436 };
437 let generic_payload_limits = clamp_payload_limits(
438 self.max_uncompressed_payload_size,
439 self.max_payload_size,
440 DEFAULT_SERIALIZER_UNCOMPRESSED_SIZE_LIMIT,
441 DEFAULT_SERIALIZER_COMPRESSED_SIZE_LIMIT,
442 );
443 let (series_uncompressed_limit, series_compressed_limit) = if series_endpoint == MetricsEndpoint::SeriesV2 {
444 clamp_payload_limits(
445 self.max_series_uncompressed_payload_size,
446 self.max_series_payload_size,
447 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT,
448 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT,
449 )
450 } else {
451 generic_payload_limits
452 };
453 let v2_series_builder = if self.requires_v2_series(metrics_v3_disabled_by_compressor)? {
454 let mut builder = v2::create_v2_request_builder(series_endpoint, &v2_endpoint_config)
455 .await
456 .error_context("Failed to create V2 series request builder.")?;
457 builder.with_len_limits(series_uncompressed_limit, series_compressed_limit)?;
458 Some(builder)
459 } else {
460 debug!("All metrics series endpoints use authoritative V3; disabling V2 series encoding.");
461 None
462 };
463
464 let (sketches_uncompressed_limit, sketches_compressed_limit) = generic_payload_limits;
465 let mut v2_sketch_builder = v2::create_v2_request_builder(MetricsEndpoint::Sketches, &v2_endpoint_config)
466 .await
467 .error_context("Failed to create V2 sketches request builder.")?;
468 v2_sketch_builder.with_len_limits(sketches_uncompressed_limit, sketches_compressed_limit)?;
469 let v2_sketch_builder = Some(v2_sketch_builder);
470
471 let flush_timeout = if self.flush_timeout.is_zero() {
472 Duration::from_millis(10)
475 } else {
476 self.flush_timeout
477 };
478
479 if series_mode.needs_v3() || sketches_mode.needs_v3() {
480 debug!(
481 ?series_mode,
482 ?sketches_mode,
483 v3_sketches_endpoints = ?self.v3_api.sketches.endpoints,
484 "V3 encoding support is enabled."
485 );
486 }
487
488 Ok(Box::new(DatadogMetrics {
489 v2_series_builder,
490 v2_sketch_builder,
491 series_mode,
492 sketches_mode,
493 v3_runtime_config,
494 telemetry,
495 flush_timeout,
496 log_payloads: self.log_payloads,
497 endpoint_routing: self.endpoint_routing.clone(),
498 }))
499 }
500}
501
502impl MemoryBounds for DatadogMetricsConfiguration {
503 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
504 builder
512 .minimum()
513 .with_single_value::<DatadogMetrics>("component struct")
514 .with_array::<EventsBuffer>("request builder events channel", 8)
515 .with_array::<PayloadsBuffer>("request builder payloads channel", 8);
516
517 builder
518 .firm()
519 .with_array::<Metric>("series metrics split re-encode buffer", self.max_metrics_per_payload)
522 .with_array::<Metric>("sketch metrics split re-encode buffer", self.max_metrics_per_payload);
523 }
524}
525
526pub struct DatadogMetrics {
527 v2_series_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>,
528 v2_sketch_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>,
529 series_mode: MetricsEncoderMode,
530 sketches_mode: MetricsEncoderMode,
531 v3_runtime_config: V3RuntimeConfig,
532 telemetry: ComponentTelemetry,
533 flush_timeout: Duration,
534 log_payloads: bool,
535 endpoint_routing: Option<Arc<MetricsEndpointRouting>>,
536}
537
538struct V3RuntimeConfig {
539 endpoint_config: EndpointConfiguration,
540 payload_limits: V3PayloadLimits,
541 series_endpoint_uri: String,
542 serializer_telemetry: V3SerializerTelemetry,
543}
544
545#[async_trait]
546impl Encoder for DatadogMetrics {
547 async fn run(mut self: Box<Self>, mut context: EncoderContext) -> Result<(), GenericError> {
548 let Self {
549 v2_series_builder,
550 v2_sketch_builder,
551 series_mode,
552 sketches_mode,
553 v3_runtime_config,
554 telemetry,
555 flush_timeout,
556 log_payloads,
557 endpoint_routing,
558 } = *self;
559
560 let mut health = context.take_health_handle();
561
562 let (events_tx, events_rx) = mpsc::channel(8);
563 let (payloads_tx, mut payloads_rx) = mpsc::channel(8);
564
565 let request_builder_fut = run_request_builder(
570 v2_series_builder,
571 v2_sketch_builder,
572 series_mode,
573 sketches_mode,
574 v3_runtime_config,
575 telemetry,
576 events_rx,
577 payloads_tx,
578 flush_timeout,
579 log_payloads,
580 );
581 runtime::worker("request_builder", request_builder_fut)
582 .on_runtime(context.topology_context().global_thread_pool().clone())
583 .spawn();
584
585 health.mark_ready();
586 debug!("Datadog Metrics encoder started.");
587
588 loop {
589 select! {
590 biased;
591
592 _ = health.live() => continue,
593 maybe_payload = payloads_rx.recv() => match maybe_payload {
594 Some(payload) => {
595 let payload = apply_endpoint_routing(payload, endpoint_routing.as_ref());
596 if let Err(e) = context.dispatcher().dispatch(payload).await {
597 error!("Failed to dispatch payload: {}", e);
598 }
599 }
600 None => break,
601 },
602 maybe_event_buffer = context.events().next() => match maybe_event_buffer {
603 Some(event_buffer) => {
604 let permit = loop {
611 select! {
612 biased;
613
614 permit = events_tx.reserve() => break permit
615 .error_context("Failed to reserve capacity for event buffer.")?,
616 maybe_payload = payloads_rx.recv() => match maybe_payload {
617 Some(payload) => {
618 let payload = apply_endpoint_routing(payload, endpoint_routing.as_ref());
619 if let Err(e) = context.dispatcher().dispatch(payload).await {
620 error!("Failed to dispatch payload: {}", e);
621 }
622 },
623
624 None => return Err(generic_error!("Request builder task stopped before accepting event buffer.")),
626 }
627 }
628 };
629
630 permit.send(event_buffer);
631 }
632 None => break,
633 },
634 }
635 }
636
637 drop(events_tx);
639
640 while let Some(payload) = payloads_rx.recv().await {
642 let payload = apply_endpoint_routing(payload, endpoint_routing.as_ref());
643 if let Err(e) = context.dispatcher().dispatch(payload).await {
644 error!("Failed to dispatch payload: {}", e);
645 }
646 }
647
648 debug!("Datadog Metrics encoder stopped.");
651
652 Ok(())
653 }
654}
655
656fn apply_endpoint_routing(payload: Payload, endpoint_routing: Option<&Arc<MetricsEndpointRouting>>) -> Payload {
657 let Some(endpoint_routing) = endpoint_routing else {
658 return payload;
659 };
660
661 match payload {
662 Payload::Http(http_payload) => {
663 let (mut metadata, request) = http_payload.into_parts();
664 metadata.set(Arc::clone(endpoint_routing));
665 Payload::Http(HttpPayload::new(metadata, request))
666 }
667 payload => payload,
668 }
669}
670
671fn log_metric_payload(metric: &Metric) {
675 match metric.values() {
676 MetricValues::Counter(..) | MetricValues::Rate(..) | MetricValues::Gauge(..) | MetricValues::Set(..) => {
677 debug!(?metric, "Flushing series metric.")
678 }
679 MetricValues::Histogram(..) | MetricValues::Distribution(..) => {
680 debug!(?metric, "Flushing sketch metric.")
681 }
682 }
683}
684
685#[allow(clippy::too_many_arguments)]
686async fn run_request_builder(
687 mut v2_series_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>,
688 mut v2_sketch_builder: Option<RequestBuilder<v2::MetricsEndpointEncoder>>, series_mode: MetricsEncoderMode,
689 sketches_mode: MetricsEncoderMode, v3_runtime_config: V3RuntimeConfig, telemetry: ComponentTelemetry,
690 mut events_rx: mpsc::Receiver<EventsBuffer>, mut payloads_tx: mpsc::Sender<PayloadsBuffer>,
691 flush_timeout: Duration, log_payloads: bool,
692) -> Result<(), GenericError> {
693 let mut pending_flush = false;
694 let pending_flush_timeout = sleep(flush_timeout);
695 tokio::pin!(pending_flush_timeout);
696
697 let mut v3_series_metrics = series_mode.needs_v3().then(Vec::<Metric>::new);
698 let mut v3_sketch_metrics = sketches_mode.needs_v3().then(Vec::<Metric>::new);
699 let mut v3_series_points = 0usize;
700 let mut v3_sketch_points = 0usize;
701 let mut v3_series_ratio = V3CompressionRatio::default();
702 let mut v3_sketch_ratio = V3CompressionRatio::default();
703
704 let tag_series = series_mode.needs_tagging();
705 let tag_sketches = sketches_mode.needs_tagging();
706 let v3_flush_context = V3FlushContext {
707 endpoint_config: &v3_runtime_config.endpoint_config,
708 payload_limits: v3_runtime_config.payload_limits,
709 series_endpoint_uri: &v3_runtime_config.series_endpoint_uri,
710 serializer_telemetry: &v3_runtime_config.serializer_telemetry,
711 telemetry: &telemetry,
712 };
713
714 loop {
715 select! {
716 Some(event_buffer) = events_rx.recv() => {
717 for event in event_buffer {
718 let metric = match event.try_into_metric() {
719 Some(metric) => metric,
720 None => continue,
721 };
722
723 if log_payloads {
724 log_metric_payload(&metric);
725 }
726
727 if !v1::has_emittable_point(&metric) {
730 debug!(metric = %metric.context().name(), "Dropping series metric with no finite points.");
731 telemetry.events_dropped_encoder().increment(1);
732 continue;
733 }
734
735 let endpoint = MetricsEndpoint::from_metric(&metric);
737 let (endpoint_mode, maybe_v2_builder, maybe_v3_metrics, v3_points) = match endpoint {
738 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => (
739 series_mode,
740 &mut v2_series_builder,
741 &mut v3_series_metrics,
742 &mut v3_series_points,
743 ),
744 MetricsEndpoint::Sketches => (
745 sketches_mode,
746 &mut v2_sketch_builder,
747 &mut v3_sketch_metrics,
748 &mut v3_sketch_points,
749 ),
750 };
751 let metric_point_count = metric.values().len();
752 let should_buffer_v3 = endpoint_mode.needs_v3();
753
754 if should_buffer_v3 {
759 if let Some(metrics) = maybe_v3_metrics {
760 metrics.push(metric.clone());
761 *v3_points += metric_point_count;
762 }
763 }
764
765 let v2_payload_info = match endpoint {
770 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => {
771 tag_series.then(MetricsPayloadInfo::v2_series)
772 }
773 MetricsEndpoint::Sketches => tag_sketches.then(MetricsPayloadInfo::v2_sketches),
774 };
775 let _v2_flushed = if let Some(builder) = maybe_v2_builder {
776 let result = encode_v2_metrics(builder, metric, &telemetry, &mut payloads_tx, v2_payload_info).await?;
777 if should_buffer_v3
778 && !result.encoded()
779 && !matches!(endpoint_mode, MetricsEncoderMode::V3Enabled)
780 {
781 if let Some(metrics) = maybe_v3_metrics {
782 let _ = metrics.pop();
783 *v3_points = v3_points.saturating_sub(metric_point_count);
784 }
785 }
786 result.flushed()
787 } else {
788 false
789 };
790
791 let v3_payload_info = match endpoint {
792 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => {
793 tag_series.then(MetricsPayloadInfo::v3_series)
794 }
795 MetricsEndpoint::Sketches => tag_sketches.then(MetricsPayloadInfo::v3_sketches),
796 };
797 let mut split_metric = None;
798 let _v3_flushed = if let Some(v3_metrics) = maybe_v3_metrics {
799 let should_flush_v3 = match endpoint_mode {
800 MetricsEncoderMode::V2Only => false,
801 MetricsEncoderMode::V3Enabled => {
802 if v3_flush_context.payload_limits.point_count_fits(metric_point_count)
803 && v3_flush_context.payload_limits.point_count_exceeds_limit(*v3_points)
804 && v3_metrics.len() > 1
805 {
806 v3_flush_context
807 .serializer_telemetry
808 .record_split_reason(V3PayloadSplitReason::MaxPoints);
809 if let Some(metric) = v3_metrics.pop() {
810 *v3_points = v3_points.saturating_sub(metric.values().len());
811 split_metric = Some(metric);
812 }
813 true
814 } else {
815 v3_flush_context.payload_limits.should_flush_point_count_limit(*v3_points)
816 || v3_flush_context
817 .payload_limits
818 .should_flush_metric_count_limit(v3_metrics)
819 }
820 }
821 };
822 if should_flush_v3 {
823 encode_and_flush_v3_metrics(
824 endpoint,
825 v3_flush_context,
826 v3_metrics,
827 &mut v3_series_ratio,
828 &mut v3_sketch_ratio,
829 &mut payloads_tx,
830 v3_payload_info,
831 )
832 .await?;
833 *v3_points = 0;
834 true
835 } else {
836 false
837 }
838 } else {
839 false
840 };
841
842 if matches!(endpoint_mode, MetricsEncoderMode::V3Enabled) {
843 if let Some(m) = split_metric.take() {
844 let point_count = m.values().len();
845 if let Some(metrics) = maybe_v3_metrics {
846 metrics.push(m);
847 *v3_points += point_count;
848 }
849 }
850 }
851 }
852
853 debug!("Processed event buffer.");
854
855 if !pending_flush {
857 pending_flush_timeout.as_mut().reset(tokio::time::Instant::now() + flush_timeout);
858 pending_flush = true;
859 }
860 },
861 _ = &mut pending_flush_timeout, if pending_flush => {
862 debug!("Flushing pending request(s).");
863
864 pending_flush = false;
865
866 let v2_series_payload_info = tag_series.then(MetricsPayloadInfo::v2_series);
868 let mut v2_series_flush_succeeded = true;
869 if let Some(builder) = &mut v2_series_builder {
870 if let Err(e) = flush_v2_metrics(builder, &mut payloads_tx, v2_series_payload_info).await {
871 error!(error = %e, "Failed to flush V2 series metrics: {}", e);
872 v2_series_flush_succeeded = false;
873 }
874 }
875
876 let v3_series_payload_info = tag_series.then(MetricsPayloadInfo::v3_series);
877 if let Some(metrics) = &mut v3_series_metrics {
878 if v2_series_flush_succeeded || matches!(series_mode, MetricsEncoderMode::V3Enabled) {
879 if let Err(e) = encode_and_flush_v3_series_metrics(
880 v3_flush_context,
881 metrics,
882 &mut v3_series_ratio,
883 &mut payloads_tx,
884 v3_series_payload_info,
885 )
886 .await
887 {
888 error!(error = %e, "Failed to flush V3 series metrics: {}", e);
889 }
890 v3_series_points = 0;
891 } else {
892 warn!("Failed to flush V2 series metrics, skipping V3 series flush.");
893 metrics.clear();
894 v3_series_points = 0;
895 }
896 }
897
898 let v2_sketches_payload_info = tag_sketches.then(MetricsPayloadInfo::v2_sketches);
900 let mut v2_sketches_flush_succeeded = true;
901 if let Some(builder) = &mut v2_sketch_builder {
902 if let Err(e) = flush_v2_metrics(builder, &mut payloads_tx, v2_sketches_payload_info).await {
903 error!(error = %e, "Failed to flush V2 sketch metrics: {}", e);
904 v2_sketches_flush_succeeded = false;
905 }
906 }
907
908 let v3_sketches_payload_info = tag_sketches.then(MetricsPayloadInfo::v3_sketches);
909 if let Some(metrics) = &mut v3_sketch_metrics {
910 if v2_sketches_flush_succeeded || matches!(sketches_mode, MetricsEncoderMode::V3Enabled) {
911 if let Err(e) = encode_and_flush_v3_sketch_metrics(
912 v3_flush_context,
913 metrics,
914 &mut v3_sketch_ratio,
915 &mut payloads_tx,
916 v3_sketches_payload_info,
917 )
918 .await
919 {
920 error!(error = %e, "Failed to flush V3 sketch metrics: {}", e);
921 }
922 v3_sketch_points = 0;
923 } else {
924 warn!("Failed to flush V2 sketch metrics, skipping V3 sketch flush.");
925 metrics.clear();
926 v3_sketch_points = 0;
927 }
928 }
929
930 debug!("All flushed requests sent to I/O task. Waiting for next event buffer...");
931 },
932
933 else => break,
935 }
936 }
937
938 Ok(())
939}
940
941struct EncodeResult {
942 encoded: bool,
943 flushed: bool,
944}
945
946impl EncodeResult {
947 pub const fn new(encoded: bool, flushed: bool) -> Self {
948 Self { encoded, flushed }
949 }
950
951 pub const fn encoded(&self) -> bool {
952 self.encoded
953 }
954
955 pub const fn flushed(&self) -> bool {
956 self.flushed
957 }
958}
959
960async fn encode_v2_metrics(
961 request_builder: &mut RequestBuilder<v2::MetricsEndpointEncoder>, metric: Metric, telemetry: &ComponentTelemetry,
962 payloads_tx: &mut mpsc::Sender<Payload>, payload_info: Option<MetricsPayloadInfo>,
963) -> Result<EncodeResult, GenericError> {
964 let metric_to_retry = match request_builder.encode(metric).await {
968 Ok(None) => return Ok(EncodeResult::new(true, false)),
969 Ok(Some(metric)) => metric,
970 Err(RequestBuilderError::InvalidInput { input }) => {
971 debug!(metric_name = %input.context().name(), "Dropping metric with no emittable values.");
972 telemetry.events_dropped_encoder().increment(1);
973 return Ok(EncodeResult::new(false, false));
974 }
975 Err(e) => {
976 error!(error = %e, "Failed to encode metric.");
977 telemetry.events_dropped_encoder().increment(1);
978 return Ok(EncodeResult::new(false, false));
979 }
980 };
981
982 flush_v2_metrics(request_builder, payloads_tx, payload_info).await?;
983
984 match request_builder.encode(metric_to_retry).await {
988 Ok(None) => Ok(EncodeResult::new(true, true)),
989 Ok(Some(_)) => unreachable!(
990 "failure to encode due to size should never occur after flush for metrics which aren't unencodable"
991 ),
992 Err(e) => {
993 error!(error = %e, "Failed to encode metric.");
994 telemetry.events_dropped_encoder().increment(1);
995 Ok(EncodeResult::new(false, true))
996 }
997 }
998}
999
1000async fn flush_v2_metrics(
1001 request_builder: &mut RequestBuilder<MetricsEndpointEncoder>, payloads_tx: &mut mpsc::Sender<Payload>,
1002 payload_info: Option<MetricsPayloadInfo>,
1003) -> Result<usize, GenericError> {
1004 let mut requests_flushed = 0;
1005
1006 let maybe_requests = request_builder.flush().await;
1007 for maybe_request in maybe_requests {
1008 match maybe_request {
1009 Ok((events, data_points, request)) => {
1010 requests_flushed += 1;
1011
1012 flush_payload(request, events, data_points, payloads_tx, payload_info).await?;
1013 }
1014
1015 Err(e) => {
1017 if !e.is_recoverable() {
1018 return Err(GenericError::from(e).context("Failed to flush request."));
1019 }
1020 }
1021 }
1022 }
1023
1024 Ok(requests_flushed)
1025}
1026
1027#[derive(Clone, Copy)]
1028struct V3FlushContext<'a> {
1029 endpoint_config: &'a EndpointConfiguration,
1030 payload_limits: V3PayloadLimits,
1031 series_endpoint_uri: &'a str,
1032 serializer_telemetry: &'a V3SerializerTelemetry,
1033 telemetry: &'a ComponentTelemetry,
1034}
1035
1036#[allow(clippy::too_many_arguments)]
1037async fn encode_and_flush_v3_metrics(
1038 endpoint: MetricsEndpoint, context: V3FlushContext<'_>, metrics: &mut Vec<Metric>,
1039 series_ratio: &mut V3CompressionRatio, sketches_ratio: &mut V3CompressionRatio,
1040 payloads_tx: &mut mpsc::Sender<Payload>, payload_info: Option<MetricsPayloadInfo>,
1041) -> Result<(), GenericError> {
1042 match endpoint {
1043 MetricsEndpoint::SeriesV1 | MetricsEndpoint::SeriesV2 => {
1044 encode_and_flush_v3_series_metrics(context, metrics, series_ratio, payloads_tx, payload_info).await
1045 }
1046 MetricsEndpoint::Sketches => {
1047 encode_and_flush_v3_sketch_metrics(context, metrics, sketches_ratio, payloads_tx, payload_info).await
1048 }
1049 }
1050}
1051
1052async fn encode_and_flush_v3_series_metrics(
1053 context: V3FlushContext<'_>, metrics: &mut Vec<Metric>, ratio: &mut V3CompressionRatio,
1054 payloads_tx: &mut mpsc::Sender<Payload>, payload_info: Option<MetricsPayloadInfo>,
1055) -> Result<(), GenericError> {
1056 if metrics.is_empty() {
1057 return Ok(());
1058 }
1059 let metrics_to_flush = std::mem::take(metrics);
1060
1061 encode_and_flush_v3_payload_requests(
1062 context.series_endpoint_uri,
1063 &metrics_to_flush,
1064 context,
1065 "series",
1066 ratio,
1067 payloads_tx,
1068 payload_info,
1069 )
1070 .await
1071}
1072
1073async fn encode_and_flush_v3_sketch_metrics(
1074 context: V3FlushContext<'_>, metrics: &mut Vec<Metric>, ratio: &mut V3CompressionRatio,
1075 payloads_tx: &mut mpsc::Sender<Payload>, payload_info: Option<MetricsPayloadInfo>,
1076) -> Result<(), GenericError> {
1077 if metrics.is_empty() {
1078 return Ok(());
1079 }
1080 let metrics_to_flush = std::mem::take(metrics);
1081
1082 encode_and_flush_v3_payload_requests(
1083 V3_SKETCHES_ENDPOINT_URI,
1084 &metrics_to_flush,
1085 context,
1086 "sketches",
1087 ratio,
1088 payloads_tx,
1089 payload_info,
1090 )
1091 .await
1092}
1093
1094const V3_MIN_COMPRESSION_RATIO: f64 = 0.02;
1098
1099const V3_BATCH_TARGET_MARGIN: f64 = 0.9;
1104
1105const V3_RATIO_SMOOTHING: f64 = 0.25;
1107
1108#[derive(Default)]
1115struct V3CompressionRatio {
1116 ratio: Option<f64>,
1117}
1118
1119impl V3CompressionRatio {
1120 fn uncompressed_target(&self, limits: V3PayloadLimits) -> usize {
1122 let ratio = self.ratio.unwrap_or(1.0).clamp(V3_MIN_COMPRESSION_RATIO, 1.0);
1125 let target_from_compressed_limit = (limits.max_compressed_size as f64 / ratio) as usize;
1126 let target = target_from_compressed_limit.min(limits.max_uncompressed_size);
1127
1128 (((target as f64) * V3_BATCH_TARGET_MARGIN) as usize).max(1)
1130 }
1131
1132 fn record(&mut self, uncompressed_len: usize, compressed_len: usize) {
1134 if uncompressed_len == 0 {
1135 return;
1136 }
1137
1138 let observed = compressed_len as f64 / uncompressed_len as f64;
1139 self.ratio = Some(match self.ratio {
1140 Some(current) => (current * (1.0 - V3_RATIO_SMOOTHING)) + (observed * V3_RATIO_SMOOTHING),
1141 None => observed,
1142 });
1143 }
1144}
1145
1146struct V3Batch {
1148 encoded: V3EncodedMetrics,
1149 event_count: usize,
1150 data_point_count: usize,
1151}
1152
1153fn build_v3_batch(
1160 metrics: &[Metric], idx: &mut usize, context: V3FlushContext<'_>, payload_kind: &'static str,
1161 target_uncompressed_len: usize,
1162) -> Option<V3Batch> {
1163 let mut writer = V3Writer::new();
1164 let mut tags_deduplicator = ReusableDeduplicator::new();
1165 let additional_tags = context.endpoint_config.additional_tags();
1166 let mut data_point_count = 0usize;
1167 let mut cut_for_size = false;
1168
1169 while *idx < metrics.len() {
1170 let metric = &metrics[*idx];
1171
1172 if !metric_has_emittable_values(metric) {
1173 debug!(metric_name = %metric.context().name(), "Dropping metric with no emittable values.");
1174 context.telemetry.events_dropped_encoder().increment(1);
1175 *idx += 1;
1176 continue;
1177 }
1178
1179 let metric_points = metric.values().len();
1180 if !context.payload_limits.point_count_fits(metric_points) {
1181 context.serializer_telemetry.record_item_too_big();
1183 context
1184 .serializer_telemetry
1185 .record_split_reason(V3PayloadSplitReason::ItemTooBig);
1186 warn!(
1187 payload_kind,
1188 data_points = metric_points,
1189 point_limit = context.payload_limits.max_points_per_payload,
1190 "Dropping oversized V3 metric that exceeds the point-count limit."
1191 );
1192 context.telemetry.events_dropped_encoder().increment(1);
1193 *idx += 1;
1194 continue;
1195 }
1196
1197 if writer.metric_count() > 0 {
1200 if !context
1201 .payload_limits
1202 .point_count_fits(data_point_count + metric_points)
1203 {
1204 context
1205 .serializer_telemetry
1206 .record_split_reason(V3PayloadSplitReason::MaxPoints);
1207 break;
1208 }
1209
1210 if context.payload_limits.metric_count_reached(writer.metric_count()) {
1211 break;
1212 }
1213 }
1214
1215 write_metric_to_v3(&mut writer, metric, additional_tags, &mut tags_deduplicator);
1216 data_point_count += metric_points;
1217 *idx += 1;
1218
1219 if writer.estimated_uncompressed_len() >= target_uncompressed_len {
1223 cut_for_size = true;
1224 break;
1225 }
1226 }
1227
1228 if cut_for_size && *idx < metrics.len() {
1231 context
1232 .serializer_telemetry
1233 .record_split_reason(V3PayloadSplitReason::PayloadFull);
1234 }
1235
1236 let event_count = writer.metric_count();
1237 if event_count == 0 {
1238 return None;
1239 }
1240
1241 match writer.finalize() {
1242 Ok(encoded) => Some(V3Batch {
1243 encoded,
1244 event_count,
1245 data_point_count,
1246 }),
1247 Err(e) => {
1248 error!(error = %e, payload_kind, events = event_count, "Failed to encode V3 metrics payload request.");
1249 context.telemetry.events_dropped_encoder().increment(event_count as u64);
1250 None
1251 }
1252 }
1253}
1254
1255async fn create_v3_batch_request(
1257 endpoint_uri: &str, batch: V3Batch, context: V3FlushContext<'_>, payload_kind: &'static str,
1258 ratio: &mut V3CompressionRatio,
1259) -> Option<V3PayloadRequest> {
1260 let V3Batch {
1261 encoded,
1262 event_count,
1263 data_point_count,
1264 } = batch;
1265
1266 let mut encoded_request =
1267 match create_v3_request(endpoint_uri, encoded, context.endpoint_config.compression_scheme()).await {
1268 Ok(request) => request,
1269 Err(e) => {
1270 error!(error = %e, payload_kind, events = event_count, "Failed to create V3 metrics request.");
1271 context.telemetry.events_dropped_encoder().increment(event_count as u64);
1272 return None;
1273 }
1274 };
1275
1276 ratio.record(encoded_request.uncompressed_len, encoded_request.compressed_len);
1279
1280 if !context.payload_limits.request_fits(&encoded_request) {
1281 return None;
1282 }
1283
1284 if let Err(e) =
1288 measure_v3_column_compressed_sizes(&mut encoded_request.stats, context.endpoint_config.compression_scheme())
1289 .await
1290 {
1291 error!(error = %e, payload_kind, "Failed to measure V3 column compressed sizes.");
1292 } else {
1293 record_v3_serializer_stats(context.serializer_telemetry, &encoded_request.stats);
1294 }
1295
1296 Some(V3PayloadRequest {
1297 request: encoded_request.request,
1298 event_count,
1299 data_point_count,
1300 })
1301}
1302
1303async fn split_and_encode_oversized_v3_range(
1309 endpoint_uri: &str, metrics: &[Metric], range: Range<usize>, context: V3FlushContext<'_>,
1310 payload_kind: &'static str, ratio: &mut V3CompressionRatio,
1311) -> Vec<V3PayloadRequest> {
1312 let accepted_metrics = metrics[range]
1315 .iter()
1316 .filter(|metric| {
1317 metric_has_emittable_values(metric) && context.payload_limits.point_count_fits(metric.values().len())
1318 })
1319 .collect::<Vec<_>>();
1320
1321 let mut requests = Vec::new();
1322 let mut pending_ranges = VecDeque::new();
1323 pending_ranges.push_back(0..accepted_metrics.len());
1324
1325 while let Some(range) = pending_ranges.pop_front() {
1326 if range.is_empty() {
1327 continue;
1328 }
1329
1330 let metrics_in_range = &accepted_metrics[range.clone()];
1331 let event_count = metrics_in_range.len();
1332 let data_point_count = metrics_in_range.iter().map(|metric| metric.values().len()).sum();
1333
1334 let encoded = match encode_v3_metrics_batch(
1335 metrics_in_range.iter().copied(),
1336 context.endpoint_config.additional_tags(),
1337 ) {
1338 Ok(encoded) => encoded,
1339 Err(e) => {
1340 error!(error = %e, payload_kind, events = event_count, "Failed to encode V3 metrics payload request.");
1341 context.telemetry.events_dropped_encoder().increment(event_count as u64);
1342 continue;
1343 }
1344 };
1345 let batch = V3Batch {
1346 encoded,
1347 event_count,
1348 data_point_count,
1349 };
1350
1351 if let Some(request) = create_v3_batch_request(endpoint_uri, batch, context, payload_kind, ratio).await {
1352 requests.push(request);
1353 continue;
1354 }
1355
1356 if range.len() == 1 {
1357 context.serializer_telemetry.record_item_too_big();
1359 context
1360 .serializer_telemetry
1361 .record_split_reason(V3PayloadSplitReason::ItemTooBig);
1362 warn!(
1363 payload_kind,
1364 compressed_limit = context.payload_limits.max_compressed_size,
1365 uncompressed_limit = context.payload_limits.max_uncompressed_size,
1366 "Dropping oversized V3 metric that cannot be split further."
1367 );
1368 context.telemetry.events_dropped_encoder().increment(1);
1369 continue;
1370 }
1371
1372 context
1374 .serializer_telemetry
1375 .record_split_reason(V3PayloadSplitReason::PayloadFull);
1376 let pivot = range.start + range.len() / 2;
1377 pending_ranges.push_front(pivot..range.end);
1378 pending_ranges.push_front(range.start..pivot);
1379 }
1380
1381 requests
1382}
1383
1384#[allow(clippy::too_many_arguments)]
1385async fn encode_and_flush_v3_payload_requests(
1386 endpoint_uri: &str, metrics: &[Metric], context: V3FlushContext<'_>, payload_kind: &'static str,
1387 ratio: &mut V3CompressionRatio, payloads_tx: &mut mpsc::Sender<Payload>, payload_info: Option<MetricsPayloadInfo>,
1388) -> Result<(), GenericError> {
1389 let target_uncompressed_len = ratio.uncompressed_target(context.payload_limits);
1390 let mut idx = 0;
1391
1392 while idx < metrics.len() {
1393 let batch_start = idx;
1394 let Some(batch) = build_v3_batch(metrics, &mut idx, context, payload_kind, target_uncompressed_len) else {
1395 if idx == batch_start {
1398 error!(
1399 payload_kind,
1400 "V3 batching made no progress; dropping remaining metrics."
1401 );
1402 context
1403 .telemetry
1404 .events_dropped_encoder()
1405 .increment((metrics.len() - idx) as u64);
1406 break;
1407 }
1408 continue;
1409 };
1410
1411 let requests = match create_v3_batch_request(endpoint_uri, batch, context, payload_kind, ratio).await {
1412 Some(request) => vec![request],
1413 None => {
1417 split_and_encode_oversized_v3_range(
1418 endpoint_uri,
1419 metrics,
1420 batch_start..idx,
1421 context,
1422 payload_kind,
1423 ratio,
1424 )
1425 .await
1426 }
1427 };
1428
1429 for payload_request in requests {
1430 flush_payload(
1431 payload_request.request,
1432 payload_request.event_count,
1433 payload_request.data_point_count,
1434 payloads_tx,
1435 payload_info,
1436 )
1437 .await?;
1438 debug!(
1439 payload_kind,
1440 events = payload_request.event_count,
1441 data_points = payload_request.data_point_count,
1442 "Sent V3 payload."
1443 );
1444 }
1445 }
1446
1447 Ok(())
1448}
1449
1450fn record_v3_serializer_stats(telemetry: &V3SerializerTelemetry, stats: &V3EncoderStats) {
1451 telemetry.record_values_count(stats.value_encoding_stats);
1452
1453 for column in &stats.columns {
1454 let uncompressed_size = column.bytes.len() as u64;
1455 let compressed_size = column.compressed_len as u64;
1456 telemetry.record_column_size(column.field_number, uncompressed_size, compressed_size);
1457 }
1458}
1459
1460async fn measure_v3_column_compressed_sizes(
1465 stats: &mut V3EncoderStats, compression_scheme: CompressionScheme,
1466) -> Result<(), GenericError> {
1467 for column in &mut stats.columns {
1468 column.compressed_len = compressed_v3_len(&column.bytes, compression_scheme)
1469 .await
1470 .error_context("Failed to measure V3 column compressed size.")?;
1471 }
1472
1473 Ok(())
1474}
1475
1476async fn compressed_v3_len(bytes: &[u8], compression_scheme: CompressionScheme) -> Result<usize, GenericError> {
1477 if matches!(compression_scheme, CompressionScheme::Noop) {
1478 return Ok(bytes.len());
1479 }
1480
1481 let buffer = ChunkedBytesBuffer::new(RB_BUFFER_CHUNK_SIZE);
1482 let mut compressor = Compressor::from_scheme(compression_scheme, buffer);
1483 compressor
1484 .write_all(bytes)
1485 .await
1486 .error_context("Failed to compress V3 bytes.")?;
1487 compressor
1488 .flush()
1489 .await
1490 .error_context("Failed to flush V3 compressor.")?;
1491 compressor
1492 .shutdown()
1493 .await
1494 .error_context("Failed to shutdown V3 compressor.")?;
1495
1496 Ok(compressor.into_inner().freeze().len())
1497}
1498
1499async fn flush_payload(
1500 request: Request<FrozenChunkedBytesBuffer>, event_count: usize, data_point_count: usize,
1501 payloads_tx: &mut mpsc::Sender<Payload>, payload_info: Option<MetricsPayloadInfo>,
1502) -> Result<(), GenericError> {
1503 let mut payload_meta = PayloadMetadata::from_event_and_data_point_count(event_count, data_point_count);
1504 if let Some(info) = payload_info {
1505 payload_meta = payload_meta.with(info);
1506 }
1507 let http_payload = HttpPayload::new(payload_meta, request);
1508 let payload = Payload::Http(http_payload);
1509
1510 payloads_tx
1511 .send(payload)
1512 .await
1513 .error_context("Failed to send payload.")?;
1514
1515 Ok(())
1516}
1517
1518fn encode_v3_metrics_batch<'a>(
1520 metrics: impl IntoIterator<Item = &'a Metric>, additional_tags: &SharedTagSet,
1521) -> Result<V3EncodedMetrics, GenericError> {
1522 let mut writer = V3Writer::new();
1523 let mut tags_deduplicator = ReusableDeduplicator::new();
1524
1525 for metric in metrics {
1526 write_metric_to_v3(&mut writer, metric, additional_tags, &mut tags_deduplicator);
1527 }
1528
1529 writer
1530 .finalize()
1531 .map_err(|e| generic_error!("Failed to serialize V3 payload: {}", e))
1532}
1533
1534pub(super) fn sketch_has_emittable_values(sketch: &DDSketch) -> bool {
1536 !sketch.is_empty() || !sketch.bins().is_empty()
1537}
1538
1539pub(super) fn metric_has_emittable_values(metric: &Metric) -> bool {
1540 match metric.values() {
1541 MetricValues::Counter(points) | MetricValues::Rate(points, _) | MetricValues::Gauge(points) => {
1542 points.into_iter().any(|(_, value)| emittable_scalar_point(value))
1543 }
1544 MetricValues::Set(points) => points.into_iter().any(|(_, value)| emittable_scalar_point(value)),
1545 MetricValues::Distribution(sketches) => sketches
1546 .into_iter()
1547 .any(|(_, sketch)| sketch_has_emittable_values(sketch)),
1548 MetricValues::Histogram(points) => points.into_iter().any(|(_, histogram)| !histogram.samples().is_empty()),
1549 }
1550}
1551
1552fn write_metric_to_v3(
1553 writer: &mut V3Writer, metric: &Metric, additional_tags: &SharedTagSet,
1554 tags_deduplicator: &mut ReusableDeduplicator<Tag>,
1555) {
1556 if !metric_has_emittable_values(metric) {
1557 debug!(metric_name = %metric.context().name(), "Dropping metric with no emittable values.");
1558 return;
1559 }
1560
1561 let metric_type = match metric.values() {
1562 MetricValues::Counter(..) => V3MetricType::Count,
1563 MetricValues::Rate(..) => V3MetricType::Rate,
1564 MetricValues::Gauge(..) | MetricValues::Set(..) => V3MetricType::Gauge,
1565 MetricValues::Histogram(..) | MetricValues::Distribution(..) => V3MetricType::Sketch,
1566 };
1567 let is_sketch = metric_type == V3MetricType::Sketch;
1568
1569 let mut builder = writer.write(metric_type, metric.context().name());
1570
1571 let chained_tags = metric
1573 .context()
1574 .tags()
1575 .into_iter()
1576 .chain(additional_tags)
1577 .chain(metric.context().origin_tags());
1578 let all_tags = tags_deduplicator
1579 .deduplicated(chained_tags)
1580 .filter(|t| is_sketch || !is_v3_series_resource_tag(t) && !is_v3_series_device_tag(t))
1581 .map(|t| t.as_str());
1582 builder.set_tags(all_tags);
1583
1584 let mut resources = Vec::new();
1586 if let Some(host) = metric.context().host().filter(|host| !host.is_empty()) {
1587 resources.push(("host", host));
1588 }
1589 if !is_sketch {
1590 let mut device_resource = None;
1591 let chained_tags = metric
1592 .context()
1593 .origin_tags()
1594 .into_iter()
1595 .chain(metric.context().tags())
1596 .chain(additional_tags);
1597 for tag in tags_deduplicator.deduplicated(chained_tags) {
1598 if is_v3_series_device_tag(tag) {
1599 device_resource = tag.value().filter(|device| !device.is_empty());
1600 } else if is_v3_series_resource_tag(tag) {
1601 if let Some((rtype, rname)) = tag.value().and_then(|value| value.split_once(':')) {
1602 if !rtype.is_empty() && !rname.is_empty() {
1603 resources.push((rtype, rname));
1604 }
1605 }
1606 }
1607 }
1608 if let Some(device) = device_resource {
1609 let device_idx = usize::from(metric.context().host().is_some_and(|host| !host.is_empty()));
1610 resources.insert(device_idx, ("device", device));
1611 }
1612 }
1613 builder.set_resources(&resources);
1614
1615 if let Some(origin) = metric.metadata().origin() {
1617 match origin {
1618 MetricOrigin::SourceType(source_type) => {
1619 builder.set_source_type(source_type.as_ref());
1620 }
1621 MetricOrigin::OriginMetadata {
1622 product,
1623 subproduct,
1624 product_detail,
1625 } => {
1626 builder.set_origin(*product, *subproduct, *product_detail, false);
1627 }
1628 }
1629 }
1630
1631 if metric_type != V3MetricType::Sketch {
1632 if let Some(unit) = metric.metadata().unit() {
1633 builder.set_unit(unit);
1634 }
1635 }
1636
1637 match metric.values() {
1639 MetricValues::Counter(points) | MetricValues::Gauge(points) => {
1640 for (ts, val) in points {
1641 if !emittable_scalar_point(val) {
1642 continue;
1643 }
1644 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
1645 builder.add_point(timestamp, val);
1646 }
1647 }
1648 MetricValues::Rate(points, interval) => {
1649 builder.set_interval(interval.as_secs());
1650 for (ts, val) in points {
1651 let scaled = if interval.is_zero() {
1653 val
1654 } else {
1655 val / interval.as_secs_f64()
1656 };
1657 if !emittable_scalar_point(scaled) {
1658 continue;
1659 }
1660 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
1661 builder.add_point(timestamp, scaled);
1662 }
1663 }
1664 MetricValues::Set(points) => {
1665 for (ts, count) in points {
1667 if !emittable_scalar_point(count) {
1668 continue;
1669 }
1670 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
1671 builder.add_point(timestamp, count);
1672 }
1673 }
1674 MetricValues::Distribution(sketches) => {
1675 for (ts, sketch) in sketches {
1676 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
1677 if sketch_has_emittable_values(sketch) {
1678 let bin_keys: Vec<i32> = sketch.bins().iter().map(|b| b.key()).collect();
1679 let bin_counts: Vec<u32> = sketch.bins().iter().map(|b| b.count()).collect();
1680 builder.add_sketch(
1681 timestamp,
1682 sketch.count() as i64,
1683 sketch.stored_sum(),
1684 sketch.stored_min(),
1685 sketch.stored_max(),
1686 &bin_keys,
1687 &bin_counts,
1688 );
1689 }
1690 }
1691 }
1692 MetricValues::Histogram(histograms) => {
1693 for (ts, histogram) in histograms {
1694 let timestamp = ts.map(|t| t.get() as i64).unwrap_or(0);
1695 let mut sketch = DDSketch::default();
1697 for sample in histogram.samples() {
1698 sketch.insert_n(sample.value.into_inner(), sample.weight.0 as u64);
1699 }
1700 if !sketch.is_empty() {
1701 let bin_keys: Vec<i32> = sketch.bins().iter().map(|b| b.key()).collect();
1702 let bin_counts: Vec<u32> = sketch.bins().iter().map(|b| b.count()).collect();
1703 builder.add_sketch(
1704 timestamp,
1705 sketch.count() as i64,
1706 sketch.sum().unwrap_or(0.0),
1707 sketch.min().unwrap_or(0.0),
1708 sketch.max().unwrap_or(0.0),
1709 &bin_keys,
1710 &bin_counts,
1711 );
1712 }
1713 }
1714 }
1715 }
1716
1717 builder.close();
1718}
1719
1720#[inline]
1721fn emittable_scalar_point(point: f64) -> bool {
1722 point.is_finite()
1723}
1724
1725fn is_v3_series_device_tag(tag: &Tag) -> bool {
1726 tag.name() == "device" && tag.value().is_some()
1727}
1728
1729fn is_v3_series_resource_tag(tag: &Tag) -> bool {
1730 tag.name() == "dd.internal.resource" && tag.value().is_some()
1731}
1732
1733async fn create_v3_request(
1735 endpoint_uri: &str, encoded: V3EncodedMetrics, compression_scheme: CompressionScheme,
1736) -> Result<V3EncodedRequest, GenericError> {
1737 let mut header_buf = [0; 16];
1740 let header_len = {
1741 let mut header_writer = CodedOutputStream::bytes(&mut header_buf);
1742 header_writer.write_tag(3, WireType::LengthDelimited)?;
1743 header_writer.write_uint64_no_tag(encoded.payload.len() as u64)?;
1744 header_writer.flush()?;
1745 header_writer.total_bytes_written() as usize
1746 };
1747
1748 let uncompressed_len = header_len + encoded.payload.len();
1749
1750 let buffer = ChunkedBytesBuffer::new(RB_BUFFER_CHUNK_SIZE);
1751 let mut compressor = Compressor::from_scheme(compression_scheme, buffer);
1752 compressor
1753 .write_all(&header_buf[..header_len])
1754 .await
1755 .error_context("Failed to compress V3 payload.")?;
1756 compressor
1757 .write_all(&encoded.payload)
1758 .await
1759 .error_context("Failed to compress V3 payload.")?;
1760 compressor
1761 .flush()
1762 .await
1763 .error_context("Failed to flush V3 compressor.")?;
1764 compressor
1765 .shutdown()
1766 .await
1767 .error_context("Failed to shutdown V3 compressor.")?;
1768
1769 let compressed_buf = compressor.into_inner().freeze();
1770 let compressed_len = compressed_buf.len();
1771
1772 let mut builder = Request::builder()
1773 .method(Method::POST)
1774 .uri(endpoint_uri)
1775 .header(http::header::CONTENT_TYPE, "application/x-protobuf");
1776
1777 if let Some(encoding) = content_encoding_for_scheme(compression_scheme) {
1778 builder = builder.header(http::header::CONTENT_ENCODING, encoding);
1779 }
1780
1781 let request = builder
1782 .body(compressed_buf)
1783 .map_err(|e| generic_error!("Failed to build V3 request: {}", e))?;
1784
1785 Ok(V3EncodedRequest {
1786 request,
1787 compressed_len,
1788 uncompressed_len,
1789 stats: encoded.stats,
1790 })
1791}
1792
1793fn content_encoding_for_scheme(compression_scheme: CompressionScheme) -> Option<HeaderValue> {
1794 match compression_scheme {
1795 CompressionScheme::Noop => None,
1796 CompressionScheme::Gzip(_) => Some(HeaderValue::from_static("gzip")),
1797 CompressionScheme::Zlib(_) => Some(HeaderValue::from_static("deflate")),
1798 CompressionScheme::Zstd(_) => Some(HeaderValue::from_static("zstd")),
1799 }
1800}
1801
1802#[cfg(test)]
1803mod tests {
1804 use std::{collections::HashMap, io::Cursor};
1805
1806 use agent_data_plane_config::{
1807 shared::{
1808 AltMetricsIntake, MetricsEncoding as TypedMetricsEncoding, V3ApiSettings as TypedV3ApiSettings,
1809 V3SeriesMode,
1810 },
1811 ConfigValue,
1812 };
1813 use bytes::Bytes;
1814 use datadog_protos::metrics::v3::MetricData as V3MetricData;
1815 use protobuf::Message as _;
1816 use saluki_core::data_model::{
1817 event::{
1818 metric::{context::Context, MetricMetadata},
1819 Event,
1820 },
1821 payload::Payload,
1822 tags::{Tag, TagSet},
1823 };
1824 use saluki_metrics::test::TestRecorder;
1825 use stringtheory::MetaString;
1826 use tokio::time::timeout;
1827
1828 use super::*;
1829 use crate::common::datadog::test_util::shared_configuration;
1830
1831 fn metrics_config_from(shared: &SharedConfiguration) -> DatadogMetricsConfiguration {
1833 DatadogMetricsConfiguration::from_configuration(shared)
1834 }
1835
1836 fn agent_series_config() -> UseV3ApiSeriesConfig {
1838 (&TypedMetricsEncoding::default()).into()
1839 }
1840
1841 #[test]
1842 fn endpoint_routing_restricts_series_and_sketch_payloads() {
1843 use crate::common::datadog::{
1844 METRICS_SERIES_V1_PATH, METRICS_SERIES_V2_PATH, METRICS_SERIES_V3_BETA_PATH, METRICS_SKETCHES_PATH,
1845 };
1846
1847 for routing in [
1848 MetricsEndpointRouting::AllExcept(["https://selected.example.com".to_string()].into()),
1849 MetricsEndpointRouting::Only(["https://selected.example.com".to_string()].into()),
1850 ] {
1851 let routing = Arc::new(routing);
1852 for (path, info) in [
1853 (METRICS_SERIES_V1_PATH, MetricsPayloadInfo::v2_series()),
1854 (METRICS_SERIES_V2_PATH, MetricsPayloadInfo::v2_series()),
1855 (METRICS_SERIES_V3_BETA_PATH, MetricsPayloadInfo::v3_series()),
1856 (METRICS_SERIES_V3_PATH, MetricsPayloadInfo::v3_series()),
1857 (METRICS_SKETCHES_PATH, MetricsPayloadInfo::v2_sketches()),
1858 (METRICS_SKETCHES_V3_PATH, MetricsPayloadInfo::v3_sketches()),
1859 ] {
1860 for tagged in [false, true] {
1861 let mut metadata = PayloadMetadata::from_event_count(1);
1862 if tagged {
1863 metadata.set(info);
1864 }
1865 let body = ChunkedBytesBuffer::new(RB_BUFFER_CHUNK_SIZE).freeze();
1866 let request = Request::builder().uri(path).body(body).unwrap();
1867 let payload = Payload::Http(HttpPayload::new(metadata, request));
1868 let routed = apply_endpoint_routing(payload, Some(&routing))
1869 .try_into_http_payload()
1870 .unwrap();
1871 let (metadata, _) = routed.into_parts();
1872 let payload_routing = metadata.get::<Arc<MetricsEndpointRouting>>();
1873 assert!(Arc::ptr_eq(payload_routing.unwrap(), &routing));
1874 assert_eq!(metadata.get::<MetricsPayloadInfo>(), tagged.then_some(&info));
1875 }
1876 }
1877 }
1878 }
1879
1880 #[test]
1881 fn v3_api_settings_come_from_resolved_configuration() {
1882 let mut shared = shared_configuration();
1883 shared.metrics_encoding.v3_api.compression_level = 7;
1884 shared.metrics_encoding.v3_api.sketches = TypedV3ApiSettings {
1885 endpoints: vec!["https://app.datadoghq.eu".to_string()],
1886 };
1887
1888 let config = metrics_config_from(&shared);
1889
1890 assert_eq!(7, config.v3_api.compression_level);
1891 assert_eq!(
1892 Some("https://app.datadoghq.eu"),
1893 config.v3_api.sketches.endpoints.first().map(String::as_str)
1894 );
1895 }
1896
1897 #[test]
1898 fn metrics_routing_comes_from_resolved_configuration() {
1899 let mut shared = shared_configuration();
1900 shared.endpoints.compression.compressor_kind = "zstd".to_string();
1901 shared.endpoints.opw_intake = AltMetricsIntake {
1902 enabled: true,
1903 url: "https://opw.example.com".to_string(),
1904 use_v3_series: true,
1905 };
1906 shared.metrics_encoding.use_v2_series_api = false;
1907 shared.metrics_encoding.v3_api.compression_level = 7;
1908 shared.metrics_encoding.v3_series_mode = V3SeriesMode::Disabled;
1909 shared.metrics_encoding.v3_series_endpoint_modes =
1910 HashMap::from([("https://app.datadoghq.com".to_string(), V3SeriesMode::Enabled)]);
1911
1912 let config = metrics_config_from(&shared);
1913
1914 assert_eq!("zstd", config.compressor_kind);
1915 assert!(!config.use_v2_series_api);
1916 assert_eq!(7, config.v3_api.compression_level);
1917 assert_eq!(V3SeriesMode::Disabled, config.use_v3_api.series.enabled);
1918 assert_eq!(
1919 Some(&V3SeriesMode::Enabled),
1920 config.use_v3_api.series.endpoints.get("https://app.datadoghq.com")
1921 );
1922 let opw = config
1924 .opw_metrics
1925 .selected_endpoint()
1926 .expect("alternate intake routing should be selected");
1927 assert_eq!("https://opw.example.com", opw.url);
1928 assert!(opw.use_v3_series);
1929 }
1930
1931 #[test]
1932 fn metrics_v3_disabled_by_compressor_uses_v2_only() {
1933 assert_eq!(MetricsEncoderMode::V2Only, metrics_encoder_mode_for_config(true, true));
1934 assert_eq!(
1935 MetricsEncoderMode::V3Enabled,
1936 metrics_encoder_mode_for_config(true, false)
1937 );
1938 }
1939
1940 #[test]
1941 fn mixed_v2_and_v3_endpoints_require_both_series_encoders() {
1942 let mut shared = shared_configuration();
1945 shared.metrics_encoding.v3_series_mode = V3SeriesMode::DatadogOnly;
1946 shared.endpoints.additional_endpoints = HashMap::from([(
1947 "https://custom.example.com".to_string(),
1948 vec!["additional-api-key".to_string()],
1949 )]);
1950 let config = metrics_config_from(&shared);
1951
1952 assert!(config.requires_v2_series(false).expect("endpoints should resolve"));
1953 assert!(config.requires_v3_series(false).expect("endpoints should resolve"));
1954 }
1955
1956 #[test]
1957 fn endpoint_routing_limits_protocol_selection_to_payload_targets() {
1958 let mut shared = shared_configuration();
1962 shared.metrics_encoding.v3_series_mode = V3SeriesMode::DatadogOnly;
1963 shared.endpoints.additional_endpoints = HashMap::from([(
1964 "https://custom.example.com".to_string(),
1965 vec!["additional-api-key".to_string()],
1966 )]);
1967
1968 let primary = shared.endpoints.primary_endpoint();
1969 let additional = "https://custom.example.com".to_string();
1970 for (routing, expects_v2) in [
1971 (MetricsEndpointRouting::AllExcept([additional.clone()].into()), false),
1972 (MetricsEndpointRouting::Only([additional].into()), true),
1973 (MetricsEndpointRouting::AllExcept([primary.clone()].into()), true),
1974 (MetricsEndpointRouting::Only([primary].into()), false),
1975 ] {
1976 let config = metrics_config_from(&shared).with_endpoint_routing(routing);
1977 assert_eq!(config.requires_v2_series(false).unwrap(), expects_v2);
1978 assert_eq!(config.requires_v3_series(false).unwrap(), !expects_v2);
1979 }
1980 }
1981
1982 #[test]
1983 fn alternate_metrics_intakes_inherit_primary_routing_but_keep_their_protocol() {
1984 for use_vector in [false, true] {
1985 for use_v3_series in [false, true] {
1986 let mut shared = shared_configuration();
1987 shared.metrics_encoding.v3_series_mode = if use_v3_series {
1989 V3SeriesMode::Disabled
1990 } else {
1991 V3SeriesMode::Enabled
1992 };
1993 let alternate = AltMetricsIntake {
1994 enabled: true,
1995 url: "https://alternate.example.com".to_string(),
1996 use_v3_series,
1997 };
1998 if use_vector {
1999 shared.endpoints.vector_intake = alternate;
2000 } else {
2001 shared.endpoints.opw_intake = alternate;
2002 }
2003 shared.endpoints.additional_endpoints = HashMap::from([(
2004 "https://alternate.example.com".to_string(),
2005 vec!["additional-key".to_string()],
2006 )]);
2007 let primary = shared.endpoints.primary_endpoint();
2008 for (routing, expects_v3) in [
2009 (MetricsEndpointRouting::Only([primary.clone()].into()), use_v3_series),
2010 (MetricsEndpointRouting::AllExcept([primary].into()), !use_v3_series),
2011 ] {
2012 let config = metrics_config_from(&shared).with_endpoint_routing(routing);
2013 assert_eq!(config.requires_v2_series(false).unwrap(), !expects_v3);
2014 assert_eq!(config.requires_v3_series(false).unwrap(), expects_v3);
2015 }
2016 }
2017 }
2018 }
2019
2020 #[test]
2021 fn all_v2_endpoints_do_not_require_v3_series() {
2022 let mut shared = shared_configuration();
2023 shared.metrics_encoding.v3_series_mode = V3SeriesMode::DatadogOnly;
2024 shared.endpoints.dd_url = ConfigValue::explicit("http://127.0.0.1:9091".to_string());
2025 let config = metrics_config_from(&shared);
2026
2027 assert!(config.requires_v2_series(false).expect("endpoints should resolve"));
2028 assert!(!config.requires_v3_series(false).expect("endpoints should resolve"));
2029 }
2030
2031 #[test]
2032 fn enabled_v3_series_mode_requires_only_v3_series() {
2033 let mut shared = shared_configuration();
2034 shared.metrics_encoding.v3_series_mode = V3SeriesMode::Enabled;
2035 let config = metrics_config_from(&shared);
2036
2037 assert!(!config.requires_v2_series(false).expect("endpoints should resolve"));
2038 assert!(config.requires_v3_series(false).expect("endpoints should resolve"));
2039 }
2040
2041 #[test]
2042 fn v1_series_configuration_takes_precedence_over_v3() {
2043 let mut shared = shared_configuration();
2044 shared.metrics_encoding.use_v2_series_api = false;
2045 shared.metrics_encoding.v3_series_mode = V3SeriesMode::Enabled;
2046 let config = metrics_config_from(&shared);
2047
2048 let series_v3_can_be_enabled =
2049 series_v3_can_be_enabled_for_config(config.use_v2_series_api, None, false, &config.use_v3_api.series);
2050
2051 assert!(!series_v3_can_be_enabled);
2052
2053 assert_eq!(
2054 MetricsEncoderMode::V2Only,
2055 metrics_encoder_mode_for_config(series_v3_can_be_enabled, false)
2056 );
2057
2058 assert!(
2059 config.requires_v2_series(false).expect("endpoint should resolve"),
2060 "V1 configuration must retain the legacy series builder even when V3 is enabled"
2061 );
2062 }
2063
2064 #[test]
2065 fn endpoint_override_uses_the_overridden_endpoint_protocol() {
2066 let mut shared = shared_configuration();
2067 shared.metrics_encoding.v3_series_mode = V3SeriesMode::Disabled;
2068 shared.endpoints.dd_url = ConfigValue::explicit("https://primary.example.com".to_string());
2069 shared.metrics_encoding.v3_series_endpoint_modes =
2070 HashMap::from([("https://v3-mrf.example.com".to_string(), V3SeriesMode::Enabled)]);
2071 let config = metrics_config_from(&shared);
2072
2073 let v2_mrf_config = config
2074 .clone()
2075 .with_metrics_endpoint_override("https://v2-mrf.example.com".to_string());
2076 let v3_mrf_config = config.with_metrics_endpoint_override("https://v3-mrf.example.com".to_string());
2077
2078 assert!(v2_mrf_config
2079 .requires_v2_series(false)
2080 .expect("V2 MRF endpoint should resolve"));
2081 assert!(!v2_mrf_config
2082 .requires_v3_series(false)
2083 .expect("V2 MRF endpoint should resolve"));
2084 assert!(!v3_mrf_config
2085 .requires_v2_series(false)
2086 .expect("V3 MRF endpoint should resolve"));
2087 assert!(v3_mrf_config
2088 .requires_v3_series(false)
2089 .expect("V3 MRF endpoint should resolve"));
2090 }
2091
2092 #[test]
2093 fn an_endpoint_override_drops_the_configured_additional_and_alternate_endpoints() {
2094 let mut shared = shared_configuration();
2095 shared.endpoints.dd_url = ConfigValue::explicit("https://primary.example.com".to_string());
2096 shared.endpoints.additional_endpoints = HashMap::from([(
2097 "https://additional.example.com".to_string(),
2098 vec!["additional-api-key".to_string()],
2099 )]);
2100 shared.endpoints.opw_intake = AltMetricsIntake {
2101 enabled: true,
2102 url: "https://opw.example.com".to_string(),
2103 use_v3_series: true,
2104 };
2105
2106 let config = metrics_config_from(&shared).with_metrics_endpoint_override("https://mrf.example.com".to_string());
2107
2108 assert_eq!("https://mrf.example.com", config.primary_endpoint);
2109 assert!(config.additional_endpoints.is_empty());
2110 assert_eq!(None, config.opw_metrics.selected_endpoint().map(|opw| opw.url));
2111 }
2112
2113 #[test]
2114 fn v2_series_only_override_requires_only_v2_series() {
2115 let mut shared = shared_configuration();
2116 shared.metrics_encoding.v3_series_mode = V3SeriesMode::Enabled;
2117 let config = metrics_config_from(&shared).with_v2_series_only();
2118
2119 assert_eq!(V3SeriesMode::Disabled, config.use_v3_api.series.enabled);
2120 assert!(config.use_v3_api.series.endpoints.is_empty());
2121 assert!(config.requires_v2_series(false).expect("endpoint should resolve"));
2122 assert!(!config.requires_v3_series(false).expect("endpoint should resolve"));
2123 }
2124
2125 #[test]
2126 fn agent_default_v3_does_not_enable_opw_only_encoder_mode() {
2127 let series_config = agent_series_config();
2128 let mut shared = shared_configuration();
2130 shared.endpoints.opw_intake = AltMetricsIntake {
2131 enabled: true,
2132 url: "http://[::1".to_string(),
2133 use_v3_series: false,
2134 };
2135 shared.endpoints.vector_intake = AltMetricsIntake {
2136 enabled: true,
2137 url: "http://vector.example.com".to_string(),
2138 use_v3_series: false,
2139 };
2140 let opw_metrics = OpwMetricsConfiguration::from_configuration(&shared.endpoints);
2141 let invalid_metrics_primary_override = selected_metrics_primary_v3_override(&opw_metrics);
2142
2143 assert!(!series_v3_can_be_enabled_for_config(
2144 true,
2145 Some(false),
2146 false,
2147 &series_config
2148 ));
2149 assert!(series_v3_can_be_enabled_for_config(
2150 true,
2151 Some(true),
2152 false,
2153 &series_config
2154 ));
2155 assert!(series_v3_can_be_enabled_for_config(
2156 true,
2157 Some(false),
2158 true,
2159 &series_config
2160 ));
2161 assert_eq!(None, invalid_metrics_primary_override);
2162 assert!(series_v3_can_be_enabled_for_config(
2163 true,
2164 invalid_metrics_primary_override,
2165 false,
2166 &series_config
2167 ));
2168 }
2169
2170 #[test]
2171 fn v3_encodes_zero_count_distribution_with_bins() {
2172 let mut sketch = DDSketch::default();
2173 sketch.insert_n(42.0, 10);
2174 sketch.set_count(0);
2175 sketch.set_min(40.0);
2176 sketch.set_max(45.0);
2177 sketch.set_sum(123.0);
2178 sketch.set_avg(f64::NAN);
2179
2180 let metric = Metric::from_parts(
2181 Context::from_static_parts("zero.count.distribution", &[]),
2182 MetricValues::distribution((123_u64, sketch)),
2183 MetricMetadata::default(),
2184 );
2185
2186 let encoded = encode_v3_metrics_batch(&[metric], &SharedTagSet::default())
2187 .expect("zero-count distribution should encode to V3");
2188 let metric_data = V3MetricData::parse_from_bytes(&encoded.payload).expect("V3 metric data should decode");
2189
2190 assert_eq!(metric_data.numPoints, vec![1]);
2191 }
2192
2193 #[test]
2194 fn v3_skips_empty_distribution() {
2195 let metric = Metric::from_parts(
2196 Context::from_static_parts("empty.distribution", &[]),
2197 MetricValues::distribution((123_u64, DDSketch::default())),
2198 MetricMetadata::default(),
2199 );
2200
2201 let encoded =
2202 encode_v3_metrics_batch(&[metric], &SharedTagSet::default()).expect("empty distribution should encode");
2203 let metric_data = V3MetricData::parse_from_bytes(&encoded.payload).expect("V3 metric data should decode");
2204
2205 assert!(metric_data.numPoints.is_empty());
2206 }
2207
2208 #[tokio::test]
2209 async fn create_v3_request_uses_configured_endpoint_uri() {
2210 let encoded = V3Writer::new().finalize().expect("empty V3 payload should encode");
2211 let request = create_v3_request("/api/intake/metrics/custom/series", encoded, CompressionScheme::noop())
2212 .await
2213 .expect("request should be created");
2214
2215 assert_eq!("/api/intake/metrics/custom/series", request.request.uri());
2216 }
2217
2218 #[tokio::test]
2219 async fn create_v3_request_uses_single_stream_body_and_column_telemetry() {
2220 let metrics = vec![Metric::counter("v3.single.stream", 42.0)];
2221 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2222 let expected_payload = encoded.payload.clone();
2223
2224 let mut request = create_v3_request(METRICS_SERIES_V3_PATH, encoded, CompressionScheme::noop())
2225 .await
2226 .expect("request should be created");
2227 measure_v3_column_compressed_sizes(&mut request.stats, CompressionScheme::noop())
2228 .await
2229 .expect("column compressed sizes should be measured");
2230
2231 for column in &request.stats.columns {
2232 assert_eq!(column.compressed_len, column.bytes.len());
2233 }
2234
2235 let mut expected_body = Vec::new();
2236 {
2237 let mut os = CodedOutputStream::vec(&mut expected_body);
2238 os.write_tag(3, WireType::LengthDelimited).unwrap();
2239 os.write_uint64_no_tag(expected_payload.len() as u64).unwrap();
2240 os.flush().unwrap();
2241 }
2242 expected_body.extend_from_slice(&expected_payload);
2243
2244 assert_eq!(request.request.into_body().into_bytes(), Bytes::from(expected_body));
2245 }
2246
2247 #[tokio::test]
2248 async fn create_v3_request_zstd_body_decodes_to_metric_payload() {
2249 let metrics = vec![Metric::counter("v3.single.stream.zstd", 42.0)];
2250 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2251 let expected_payload = encoded.payload.clone();
2252
2253 let mut request = create_v3_request(METRICS_SERIES_V3_PATH, encoded, CompressionScheme::zstd_default())
2254 .await
2255 .expect("request should be created");
2256 measure_v3_column_compressed_sizes(&mut request.stats, CompressionScheme::zstd_default())
2257 .await
2258 .expect("column compressed sizes should be measured");
2259
2260 for column in &request.stats.columns {
2261 let expected_compressed_len = compressed_v3_len(&column.bytes, CompressionScheme::zstd_default())
2262 .await
2263 .expect("column compressed size should be measured");
2264 assert_eq!(column.compressed_len, expected_compressed_len);
2265 }
2266
2267 let mut expected_body = Vec::new();
2268 {
2269 let mut os = CodedOutputStream::vec(&mut expected_body);
2270 os.write_tag(3, WireType::LengthDelimited).unwrap();
2271 os.write_uint64_no_tag(expected_payload.len() as u64).unwrap();
2272 os.flush().unwrap();
2273 }
2274 expected_body.extend_from_slice(&expected_payload);
2275
2276 let decoded_body = zstd::stream::decode_all(Cursor::new(request.request.into_body().into_bytes()))
2277 .expect("compressed V3 body should decode");
2278 assert_eq!(decoded_body, expected_body);
2279 }
2280
2281 #[test]
2282 fn v3_drops_non_finite_scalar_points() {
2283 const NUM_POINTS_FIELD_NUMBER: u32 = 15;
2284 const VALUE_SINT64_FIELD_NUMBER: u32 = 17;
2285
2286 let metrics = vec![
2287 Metric::from_parts(
2288 Context::from_static_parts("v3.finite.counter", &[]),
2289 MetricValues::counter([(1, 1.0_f64), (2, f64::NAN), (3, f64::INFINITY), (4, 2.0)]),
2290 MetricMetadata::default(),
2291 ),
2292 Metric::from_parts(
2293 Context::from_static_parts("v3.finite.gauge", &[]),
2294 MetricValues::gauge([(1, 3.0_f64), (2, f64::NAN), (3, 4.0)]),
2295 MetricMetadata::default(),
2296 ),
2297 Metric::from_parts(
2298 Context::from_static_parts("v3.finite.rate", &[]),
2299 MetricValues::rate(
2300 [(1, 30.0_f64), (2, f64::INFINITY), (3, f64::NAN)],
2301 Duration::from_secs(10),
2302 ),
2303 MetricMetadata::default(),
2304 ),
2305 Metric::set("v3.finite.set", "alpha"),
2306 ];
2307
2308 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2309 let column = |field_number| {
2310 encoded
2311 .stats
2312 .columns
2313 .iter()
2314 .find(|column| column.field_number == field_number)
2315 .map(|column| column.bytes.as_slice())
2316 };
2317
2318 assert_eq!(column(NUM_POINTS_FIELD_NUMBER), Some(&[2, 2, 1, 1][..]));
2319 assert_eq!(column(VALUE_SINT64_FIELD_NUMBER), Some(&[2, 4, 6, 8, 6, 2][..]));
2321 assert_eq!(encoded.stats.value_encoding_stats.sint64, 6);
2322 assert_eq!(encoded.stats.value_encoding_stats.float64, 0);
2323 }
2324
2325 #[test]
2326 fn v3_encodes_zero_interval_rate_without_scaling() {
2327 const NUM_POINTS_FIELD_NUMBER: u32 = 15;
2328 const VALUE_SINT64_FIELD_NUMBER: u32 = 17;
2329
2330 let metrics = vec![Metric::rate("v3.unnormalized.rate", 42.0, Duration::ZERO)];
2331 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("rate should encode to V3");
2332 let column = |field_number| {
2333 encoded
2334 .stats
2335 .columns
2336 .iter()
2337 .find(|column| column.field_number == field_number)
2338 .map(|column| column.bytes.as_slice())
2339 };
2340
2341 assert_eq!(column(NUM_POINTS_FIELD_NUMBER), Some(&[1][..]));
2342 assert_eq!(column(VALUE_SINT64_FIELD_NUMBER), Some(&[84][..]));
2343 }
2344
2345 #[tokio::test]
2346 async fn v3_serializer_stats_record_agent_style_value_and_column_counters() {
2347 const VALUE_SINT64_FIELD_NUMBER: u32 = 17;
2348
2349 let recorder = TestRecorder::default();
2350 let _local = metrics::set_default_local_recorder(&recorder);
2351 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2352 let metrics = vec![
2353 Metric::gauge("v3.telemetry.zero", [(123, 0.0), (124, 0.0)]),
2354 Metric::counter("v3.telemetry.sint64", [(123, 100.0), (124, 200.0)]),
2355 Metric::gauge("v3.telemetry.float32", [(123, 1.5), (124, 2.25)]),
2356 Metric::gauge("v3.telemetry.float64", [(123, (1i64 << 30) as f64), (124, 1.5)]),
2357 ];
2358 let encoded = encode_v3_metrics_batch(&metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2359 let mut request = create_v3_request(METRICS_SERIES_V3_PATH, encoded, CompressionScheme::noop())
2360 .await
2361 .expect("request should be created");
2362 measure_v3_column_compressed_sizes(&mut request.stats, CompressionScheme::noop())
2363 .await
2364 .expect("column compressed sizes should be measured");
2365
2366 let value_sint64_column = request
2367 .stats
2368 .columns
2369 .iter()
2370 .find(|column| column.field_number == VALUE_SINT64_FIELD_NUMBER)
2371 .expect("sint64 value column should be present");
2372 let value_sint64_len = value_sint64_column.bytes.len() as u64;
2373 assert_eq!(value_sint64_column.compressed_len as u64, value_sint64_len);
2374
2375 record_v3_serializer_stats(&serializer_telemetry, &request.stats);
2376
2377 assert_eq!(
2378 recorder.counter(("serializer.v3_values_count", &[("type", "zero")])),
2379 Some(2)
2380 );
2381 assert_eq!(
2382 recorder.counter(("serializer.v3_values_count", &[("type", "sint64")])),
2383 Some(2)
2384 );
2385 assert_eq!(
2386 recorder.counter(("serializer.v3_values_count", &[("type", "float32")])),
2387 Some(2)
2388 );
2389 assert_eq!(
2390 recorder.counter(("serializer.v3_values_count", &[("type", "float64")])),
2391 Some(2)
2392 );
2393 assert_eq!(
2394 recorder.counter((
2395 "serializer.v3_column_size",
2396 &[("column", "ValueSint64"), ("compressed", "uncompressed")]
2397 )),
2398 Some(value_sint64_len)
2399 );
2400 assert_eq!(
2401 recorder.counter((
2402 "serializer.v3_column_size",
2403 &[("column", "ValueSint64"), ("compressed", "compressed")]
2404 )),
2405 Some(value_sint64_len)
2406 );
2407 }
2408
2409 async fn create_v3_test_request(metrics: &[Metric]) -> V3EncodedRequest {
2410 let encoded = encode_v3_metrics_batch(metrics, &SharedTagSet::default()).expect("metrics should encode to V3");
2411 create_v3_request(METRICS_SERIES_V3_PATH, encoded, CompressionScheme::noop())
2412 .await
2413 .expect("request should be created")
2414 }
2415
2416 fn test_v3_flush_context<'a>(
2417 ep_config: &'a EndpointConfiguration, payload_limits: V3PayloadLimits,
2418 serializer_telemetry: &'a V3SerializerTelemetry, telemetry: &'a ComponentTelemetry,
2419 ) -> V3FlushContext<'a> {
2420 V3FlushContext {
2421 endpoint_config: ep_config,
2422 payload_limits,
2423 series_endpoint_uri: METRICS_SERIES_V3_PATH,
2424 serializer_telemetry,
2425 telemetry,
2426 }
2427 }
2428
2429 struct CollectedV3Payload {
2431 event_count: usize,
2432 data_point_count: usize,
2433 request: Request<FrozenChunkedBytesBuffer>,
2434 }
2435
2436 async fn collect_v3_payload_requests(metrics: &[Metric], context: V3FlushContext<'_>) -> Vec<CollectedV3Payload> {
2441 let (mut payloads_tx, mut payloads_rx) = tokio::sync::mpsc::channel(64);
2442 let mut ratio = V3CompressionRatio::default();
2443 encode_and_flush_v3_payload_requests(
2444 METRICS_SERIES_V3_PATH,
2445 metrics,
2446 context,
2447 "series",
2448 &mut ratio,
2449 &mut payloads_tx,
2450 None,
2451 )
2452 .await
2453 .expect("payload requests should encode and flush");
2454 drop(payloads_tx);
2455
2456 let mut collected = Vec::new();
2457 while let Some(payload) = payloads_rx.recv().await {
2458 let Payload::Http(http_payload) = payload else {
2459 panic!("expected HTTP payload");
2460 };
2461 let (metadata, request) = http_payload.into_parts();
2462 collected.push(CollectedV3Payload {
2463 event_count: metadata.event_count(),
2464 data_point_count: metadata.data_point_count(),
2465 request,
2466 });
2467 }
2468
2469 collected
2470 }
2471
2472 #[tokio::test]
2473 async fn v3_payload_requests_split_by_compressed_size_limit() {
2474 let metrics = vec![
2475 Metric::counter("v3.compressed.split.one", 1.0),
2476 Metric::counter("v3.compressed.split.two", 2.0),
2477 ];
2478 let single_request = create_v3_test_request(&metrics[..1]).await;
2479 let combined_request = create_v3_test_request(&metrics).await;
2480 assert!(combined_request.compressed_len > single_request.compressed_len);
2481
2482 let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2483 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2484 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2485 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2486
2487 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2488 let requests = collect_v3_payload_requests(&metrics, context).await;
2489
2490 assert_eq!(2, requests.len());
2491 assert_eq!(
2492 vec![1, 1],
2493 requests.iter().map(|request| request.event_count).collect::<Vec<_>>()
2494 );
2495 assert!(requests
2496 .iter()
2497 .all(|request| request.request.body().len() <= limits.max_compressed_size));
2498 }
2499
2500 #[tokio::test]
2501 async fn v3_byte_fallback_does_not_resurrect_metric_over_point_limit() {
2502 let point_oversized_metric = Metric::counter("v3.fallback.over.point.limit", [(123, 1.0), (124, 2.0)]);
2503
2504 let byte_oversized_metric = Metric::counter(
2505 concat!(
2506 "v3.fallback.over.byte.limit.",
2507 "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
2508 "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
2509 "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc",
2510 "dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd",
2511 ),
2512 3.0,
2513 );
2514
2515 let point_oversized_request = create_v3_test_request(std::slice::from_ref(&point_oversized_metric)).await;
2516 let byte_oversized_request = create_v3_test_request(std::slice::from_ref(&byte_oversized_metric)).await;
2517
2518 assert!(byte_oversized_request.compressed_len > point_oversized_request.compressed_len);
2519
2520 let limits = V3PayloadLimits::new(point_oversized_request.compressed_len, usize::MAX, 10_000, 1);
2521
2522 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2523 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2524 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2525 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2526
2527 let requests = collect_v3_payload_requests(&[point_oversized_metric, byte_oversized_metric], context).await;
2528
2529 let emitted_point_counts = requests
2530 .iter()
2531 .map(|request| request.data_point_count)
2532 .collect::<Vec<_>>();
2533
2534 assert!(
2535 emitted_point_counts
2536 .iter()
2537 .all(|point_count| *point_count <= limits.max_points_per_payload),
2538 "fallback emitted request point counts {emitted_point_counts:?} above limit {}",
2539 limits.max_points_per_payload
2540 );
2541 }
2542
2543 #[tokio::test]
2544 async fn v3_serializer_stats_record_payload_full_split_reason() {
2545 let metrics = vec![
2546 Metric::counter("v3.telemetry.payload_full.one", 1.0),
2547 Metric::counter("v3.telemetry.payload_full.two", 2.0),
2548 ];
2549 let single_request = create_v3_test_request(&metrics[..1]).await;
2550 let limits = V3PayloadLimits::new(single_request.compressed_len, usize::MAX, 10_000, 10_000);
2551 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2552 let recorder = TestRecorder::default();
2553 let _local = metrics::set_default_local_recorder(&recorder);
2554 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2555 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2556
2557 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2558 let requests = collect_v3_payload_requests(&metrics, context).await;
2559
2560 assert_eq!(2, requests.len());
2561 assert_eq!(
2562 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "payload_full")])),
2563 Some(1)
2564 );
2565 }
2566
2567 #[tokio::test]
2568 async fn v3_serializer_stats_record_item_too_big_split_reason() {
2569 let metrics = vec![Metric::counter("v3.telemetry.item_too_big", 1.0)];
2570 let request = create_v3_test_request(&metrics).await;
2571 let limits = V3PayloadLimits::new(request.compressed_len - 1, usize::MAX, 10_000, 10_000);
2572 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2573 let recorder = TestRecorder::default();
2574 let _local = metrics::set_default_local_recorder(&recorder);
2575 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2576 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2577
2578 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2579 let requests = collect_v3_payload_requests(&metrics, context).await;
2580
2581 assert!(requests.is_empty());
2582 assert_eq!(recorder.counter("serializer.v3_item_too_big"), Some(1));
2583 assert_eq!(
2584 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "item_too_big")])),
2585 Some(1)
2586 );
2587 }
2588
2589 #[tokio::test]
2590 async fn v3_payload_requests_split_by_uncompressed_size_limit() {
2591 let metrics = vec![
2592 Metric::counter("v3.uncompressed.split.one", 1.0),
2593 Metric::counter("v3.uncompressed.split.two", 2.0),
2594 ];
2595 let single_request = create_v3_test_request(&metrics[..1]).await;
2596 let combined_request = create_v3_test_request(&metrics).await;
2597 assert!(combined_request.uncompressed_len > single_request.uncompressed_len);
2598
2599 let limits = V3PayloadLimits::new(usize::MAX, single_request.uncompressed_len, 10_000, 10_000);
2600 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2601 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2602 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2603
2604 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2605 let requests = collect_v3_payload_requests(&metrics, context).await;
2606
2607 assert_eq!(2, requests.len());
2608 assert_eq!(
2609 vec![1, 1],
2610 requests.iter().map(|request| request.event_count).collect::<Vec<_>>()
2611 );
2612 }
2613
2614 fn collect_v3_batches(
2616 metrics: &[Metric], context: V3FlushContext<'_>, target_uncompressed_len: usize,
2617 ) -> Vec<(Range<usize>, usize)> {
2618 let mut batches = Vec::new();
2619 let mut idx = 0;
2620 while idx < metrics.len() {
2621 let start = idx;
2622 match build_v3_batch(metrics, &mut idx, context, "series", target_uncompressed_len) {
2623 Some(batch) => batches.push((start..idx, batch.event_count)),
2624 None if idx == start => break,
2625 None => {}
2626 }
2627 }
2628
2629 batches
2630 }
2631
2632 #[test]
2633 fn v3_batches_are_cut_at_the_point_limit() {
2634 let metrics = vec![
2635 Metric::counter("v3.points.split.one", [(123, 1.0), (124, 2.0)]),
2636 Metric::counter("v3.points.split.two", [(123, 3.0), (124, 4.0)]),
2637 Metric::counter("v3.points.split.three", 5.0),
2638 ];
2639 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 3);
2640 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2641 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2642 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2643 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2644
2645 let batches = collect_v3_batches(&metrics, context, usize::MAX);
2646
2647 assert_eq!(vec![(0..1, 1), (1..3, 2)], batches);
2648 }
2649
2650 #[test]
2651 fn v3_batches_are_cut_at_the_metric_limit() {
2652 let metrics = vec![
2653 Metric::counter("v3.metric.limit.a", 1.0),
2654 Metric::counter("v3.metric.limit.b", 2.0),
2655 Metric::counter("v3.metric.limit.c", 3.0),
2656 Metric::counter("v3.metric.limit.d", 4.0),
2657 Metric::counter("v3.metric.limit.e", 5.0),
2658 ];
2659 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 2, 10_000);
2660 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 2, usize::MAX, None);
2661 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2662 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2663 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2664
2665 let batches = collect_v3_batches(&metrics, context, usize::MAX);
2666
2667 assert_eq!(vec![(0..2, 2), (2..4, 2), (4..5, 1)], batches);
2668 }
2669
2670 #[test]
2671 fn v3_batches_are_cut_at_the_uncompressed_byte_target() {
2672 let metrics = vec![
2673 Metric::counter("v3.byte.target.a", 1.0),
2674 Metric::counter("v3.byte.target.b", 2.0),
2675 Metric::counter("v3.byte.target.c", 3.0),
2676 Metric::counter("v3.byte.target.d", 4.0),
2677 ];
2678 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000);
2679 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2680 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2681 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2682 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2683
2684 let batches = collect_v3_batches(&metrics, context, 1);
2687
2688 assert_eq!(4, batches.len());
2689 assert!(batches.iter().all(|(_, event_count)| *event_count == 1));
2690 }
2691
2692 #[test]
2693 fn v3_serializer_stats_record_max_points_split_reason() {
2694 let metrics = vec![
2695 Metric::counter("v3.telemetry.max_points.one", [(123, 1.0), (124, 2.0)]),
2696 Metric::counter("v3.telemetry.max_points.two", [(123, 3.0), (124, 4.0)]),
2697 ];
2698 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 2);
2699 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2700 let recorder = TestRecorder::default();
2701 let _local = metrics::set_default_local_recorder(&recorder);
2702 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2703 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2704 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2705
2706 let batches = collect_v3_batches(&metrics, context, usize::MAX);
2707
2708 assert_eq!(vec![(0..1, 1), (1..2, 1)], batches);
2709 assert_eq!(
2710 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "max_points")])),
2711 Some(1)
2712 );
2713 }
2714
2715 #[test]
2716 fn v3_batches_drop_oversized_metric_and_keep_batching() {
2717 let metrics = vec![
2718 Metric::counter("v3.points.oversized.before", [(123, 1.0), (124, 2.0)]),
2719 Metric::counter(
2720 "v3.points.oversized.too_big",
2721 [(123, 3.0), (124, 4.0), (125, 5.0), (126, 6.0)],
2722 ),
2723 Metric::counter("v3.points.oversized.after", 7.0),
2724 ];
2725 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 3);
2726 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2727 let recorder = TestRecorder::default();
2728 let _local = metrics::set_default_local_recorder(&recorder);
2729 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2730 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2731 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2732
2733 let batches = collect_v3_batches(&metrics, context, usize::MAX);
2734
2735 assert_eq!(vec![(0..3, 2)], batches);
2738 assert_eq!(recorder.counter("serializer.v3_item_too_big"), Some(1));
2739 }
2740
2741 #[test]
2742 fn v3_batches_skip_zero_point_metrics_without_splitting() {
2743 let metrics = vec![
2744 Metric::counter("v3.points.zero.before", 1.0),
2745 Metric::counter("v3.points.zero.empty", &[] as &[f64]),
2746 Metric::counter("v3.points.zero.after", 2.0),
2747 ];
2748 let limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000);
2749 let ep_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2750 let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());
2751 let serializer_telemetry = V3SerializerTelemetry::from_builder(&MetricsBuilder::default());
2752 let context = test_v3_flush_context(&ep_config, limits, &serializer_telemetry, &telemetry);
2753
2754 let batches = collect_v3_batches(&metrics, context, usize::MAX);
2755
2756 assert_eq!(vec![(0..3, 2)], batches);
2757 }
2758
2759 #[test]
2760 fn v3_compression_ratio_targets_the_compressed_limit() {
2761 let limits = V3PayloadLimits::new(500_000, 5_000_000, 10_000, 10_000);
2762 let mut ratio = V3CompressionRatio::default();
2763
2764 let initial_target = ratio.uncompressed_target(limits);
2767 assert_eq!((500_000.0 * V3_BATCH_TARGET_MARGIN) as usize, initial_target);
2768
2769 ratio.record(4_000_000, 1_000_000);
2771 let observed_target = ratio.uncompressed_target(limits);
2772 assert!(
2773 observed_target > initial_target,
2774 "target should grow once compression is observed: {observed_target} vs {initial_target}"
2775 );
2776
2777 ratio.record(4_000_000, 40_000);
2779 ratio.record(4_000_000, 40_000);
2780 ratio.record(4_000_000, 40_000);
2781 ratio.record(4_000_000, 40_000);
2782 ratio.record(4_000_000, 40_000);
2783 assert!(ratio.uncompressed_target(limits) <= limits.max_uncompressed_size);
2784 }
2785
2786 #[test]
2787 fn v3_series_metric_unit_refs_are_encoded_sparsely() {
2788 let context = Context::from_static_parts("my.timer.avg", &[]);
2789 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
2790 let gauge = Metric::from_parts(context, MetricValues::gauge([1.0_f64]), metadata);
2791 let context = Context::from_static_parts("my.counter", &[]);
2792 let no_unit = Metric::from_parts(context, MetricValues::gauge([2.0_f64]), MetricMetadata::default());
2793 let context = Context::from_static_parts("my.timer.max", &[]);
2794 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
2795 let same_unit = Metric::from_parts(context, MetricValues::gauge([3.0_f64]), metadata);
2796
2797 let payload = encode_v3_metrics_batch(&[gauge, no_unit, same_unit], &SharedTagSet::default())
2798 .expect("V3 metric should encode successfully")
2799 .payload;
2800
2801 let expected_unit_dict = [
2802 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',
2805 ];
2806 assert!(
2807 payload
2808 .windows(expected_unit_dict.len())
2809 .any(|window| window == expected_unit_dict),
2810 "V3 payload should contain DictUnitStr field for 'millisecond', got bytes: {:?}",
2811 payload
2812 );
2813
2814 let expected_unit_ref = [
2815 0xd2, 0x01, 0x02, 0x02, 0x00, ];
2819 assert!(
2820 payload
2821 .windows(expected_unit_ref.len())
2822 .any(|window| window == expected_unit_ref),
2823 "V3 payload should contain UnitRef field for 'millisecond', got bytes: {:?}",
2824 payload
2825 );
2826 }
2827
2828 #[test]
2829 fn v3_sketch_metric_unit_not_encoded() {
2830 let context = Context::from_static_parts("my.histogram", &[]);
2831 let metadata = MetricMetadata::default().with_unit(MetaString::from_static("millisecond"));
2832 let histogram = Metric::from_parts(context, MetricValues::histogram([1.0_f64]), metadata);
2833
2834 let payload = encode_v3_metrics_batch(&[histogram], &SharedTagSet::default())
2835 .expect("V3 sketch metric should encode successfully")
2836 .payload;
2837
2838 assert!(
2839 !payload
2840 .windows(b"millisecond".len())
2841 .any(|window| window == b"millisecond"),
2842 "V3 sketch payload should not contain unit bytes, matching the Agent V3 sketch builder: {:?}",
2843 payload
2844 );
2845 }
2846
2847 #[test]
2848 fn v3_series_promotes_device_and_internal_resource_tags_to_resources() {
2849 let context = Context::from_static_parts(
2850 "series.resources",
2851 &[
2852 "env:prod",
2853 "device:switch1",
2854 "dd.internal.resource:pod:pod-a",
2855 "dd.internal.resource:malformed",
2856 ],
2857 );
2858 let context = context.with_host(Some(MetaString::from_static("host-a")));
2859 let metadata = MetricMetadata::default();
2860 let metric = Metric::from_parts(context, MetricValues::gauge([1.0_f64]), metadata);
2861
2862 let payload = encode_v3_metrics_batch(&[metric], &SharedTagSet::default())
2863 .expect("V3 series should encode successfully")
2864 .payload;
2865
2866 assert_contains_bytes(&payload, b"env:prod");
2867 assert!(!contains_bytes(&payload, b"device:switch1"));
2868 assert!(!contains_bytes(&payload, b"dd.internal.resource:pod:pod-a"));
2869 assert!(!contains_bytes(&payload, b"dd.internal.resource:malformed"));
2870
2871 let expected_resource_dict = [
2872 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',
2875 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'-',
2876 b'a',
2877 ];
2878 assert_contains_bytes(&payload, &expected_resource_dict);
2879 }
2880
2881 #[test]
2882 fn v3_series_promotes_additional_and_origin_resource_tags_without_empty_host() {
2883 let context = Context::from_static_parts("series.additional_origin_resources", &["env:prod"])
2884 .with_origin_tags(tag_set(["dd.internal.resource:pod:pod-origin"]));
2885 let additional_tags = SharedTagSet::from(tag_set([
2886 "team:core",
2887 "device:switch1",
2888 "dd.internal.resource:container:container-a",
2889 ]));
2890 let context = context.with_host(Some(MetaString::empty()));
2891 let metadata = MetricMetadata::default();
2892 let metric = Metric::from_parts(context, MetricValues::gauge([1.0_f64]), metadata);
2893
2894 let payload = encode_v3_metrics_batch(&[metric], &additional_tags)
2895 .expect("V3 series should encode successfully")
2896 .payload;
2897
2898 assert_contains_bytes(&payload, b"env:prod");
2899 assert_contains_bytes(&payload, b"team:core");
2900 assert!(!contains_bytes(&payload, b"device:switch1"));
2901 assert!(!contains_bytes(&payload, b"dd.internal.resource:container:container-a"));
2902 assert!(!contains_bytes(&payload, b"dd.internal.resource:pod:pod-origin"));
2903
2904 let expected_resource_dict = [
2905 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',
2908 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',
2909 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',
2910 ];
2911 assert_contains_bytes(&payload, &expected_resource_dict);
2912 assert!(!contains_bytes(&payload, b"host"));
2913 }
2914
2915 #[test]
2916 fn v3_sketch_keeps_device_and_internal_resource_tags_as_tags() {
2917 let context = Context::from_static_parts(
2918 "sketch.resources",
2919 &["env:prod", "device:switch1", "dd.internal.resource:pod:pod-a"],
2920 );
2921 let context = context.with_host(Some(MetaString::from_static("host-a")));
2922 let metadata = MetricMetadata::default();
2923 let metric = Metric::from_parts(context, MetricValues::histogram([1.0_f64]), metadata);
2924
2925 let payload = encode_v3_metrics_batch(&[metric], &SharedTagSet::default())
2926 .expect("V3 sketch should encode successfully")
2927 .payload;
2928
2929 assert_contains_bytes(&payload, b"env:prod");
2930 assert_contains_bytes(&payload, b"device:switch1");
2931 assert_contains_bytes(&payload, b"dd.internal.resource:pod:pod-a");
2932
2933 let expected_resource_dict = [
2934 0x22, 0x0c, 0x04, b'h', b'o', b's', b't', 0x06, b'h', b'o', b's', b't', b'-', b'a',
2937 ];
2938 assert_contains_bytes(&payload, &expected_resource_dict);
2939 }
2940
2941 struct RequestBuilderScenario {
2943 series_mode: MetricsEncoderMode,
2944 sketches_mode: MetricsEncoderMode,
2945 payload_limits: V3PayloadLimits,
2946 flush_timeout: Duration,
2947 }
2948
2949 impl RequestBuilderScenario {
2950 fn new(series_mode: MetricsEncoderMode, sketches_mode: MetricsEncoderMode) -> Self {
2952 Self {
2953 series_mode,
2954 sketches_mode,
2955 payload_limits: V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, 10_000),
2956 flush_timeout: Duration::from_millis(10),
2957 }
2958 }
2959
2960 fn with_max_points_per_payload(mut self, max_points_per_payload: usize) -> Self {
2962 self.payload_limits = V3PayloadLimits::new(usize::MAX, usize::MAX, 10_000, max_points_per_payload);
2963 self
2964 }
2965
2966 fn with_flush_timeout(mut self, flush_timeout: Duration) -> Self {
2968 self.flush_timeout = flush_timeout;
2969 self
2970 }
2971
2972 async fn spawn(self) -> RequestBuilderHarness {
2975 let v3_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, usize::MAX, None);
2976 let metrics_builder = MetricsBuilder::default();
2977 let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
2978 let serializer_telemetry = V3SerializerTelemetry::from_builder(&metrics_builder);
2979 let (events_tx, events_rx) = tokio::sync::mpsc::channel(1);
2980 let (payloads_tx, payloads_rx) = tokio::sync::mpsc::channel(8);
2981
2982 let handle = tokio::spawn(run_request_builder(
2983 None,
2984 None,
2985 self.series_mode,
2986 self.sketches_mode,
2987 V3RuntimeConfig {
2988 endpoint_config: v3_endpoint_config,
2989 payload_limits: self.payload_limits,
2990 series_endpoint_uri: METRICS_SERIES_V3_PATH.to_string(),
2991 serializer_telemetry,
2992 },
2993 telemetry,
2994 events_rx,
2995 payloads_tx,
2996 self.flush_timeout,
2997 false,
2998 ));
2999
3000 RequestBuilderHarness {
3001 events_tx,
3002 payloads_rx,
3003 handle,
3004 }
3005 }
3006 }
3007
3008 struct RequestBuilderHarness {
3010 events_tx: tokio::sync::mpsc::Sender<EventsBuffer>,
3011 payloads_rx: tokio::sync::mpsc::Receiver<PayloadsBuffer>,
3012 handle: tokio::task::JoinHandle<Result<(), GenericError>>,
3013 }
3014
3015 impl RequestBuilderHarness {
3016 async fn push_metrics(&self, metrics: impl IntoIterator<Item = Metric>) {
3018 let mut events = EventsBuffer::default();
3019 for metric in metrics {
3020 assert!(
3021 events.try_push(Event::Metric(metric)).is_none(),
3022 "event buffer should accept metric"
3023 );
3024 }
3025 self.events_tx
3026 .send(events)
3027 .await
3028 .expect("events should be sent to request builder");
3029 }
3030
3031 async fn next_http_request(&mut self) -> (PayloadMetadata, Request<FrozenChunkedBytesBuffer>) {
3033 let payload = timeout(Duration::from_secs(1), self.payloads_rx.recv())
3034 .await
3035 .expect("payload should arrive before timeout")
3036 .expect("payload channel should remain open");
3037 match payload {
3038 Payload::Http(http_payload) => http_payload.into_parts(),
3039 _ => panic!("expected HTTP payload"),
3040 }
3041 }
3042
3043 async fn assert_no_payload_within(&mut self, window: Duration) {
3045 assert!(
3046 timeout(window, self.payloads_rx.recv()).await.is_err(),
3047 "no payload should be flushed within {window:?}"
3048 );
3049 }
3050
3051 async fn shutdown(self) {
3053 drop(self.events_tx);
3054 self.handle
3055 .await
3056 .expect("request builder task should complete")
3057 .expect("request builder should stop cleanly");
3058 }
3059 }
3060
3061 #[tokio::test]
3062 async fn authoritative_v3_flushes_previous_point_limit_batch() {
3063 let recorder = TestRecorder::default();
3067 let _local = metrics::set_default_local_recorder(&recorder);
3068
3069 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V3Enabled, MetricsEncoderMode::V2Only)
3070 .with_max_points_per_payload(3)
3071 .with_flush_timeout(Duration::from_millis(250))
3072 .spawn()
3073 .await;
3074
3075 harness
3076 .push_metrics([
3077 Metric::counter("authoritative.v3.points.one", [(123, 1.0), (124, 2.0)]),
3078 Metric::counter("authoritative.v3.points.two", [(123, 3.0), (124, 4.0)]),
3079 ])
3080 .await;
3081
3082 let (_, request) = harness.next_http_request().await;
3084 assert_eq!(METRICS_SERIES_V3_PATH, request.uri());
3085 harness.assert_no_payload_within(Duration::from_millis(50)).await;
3086
3087 let (_, request) = harness.next_http_request().await;
3089 assert_eq!(METRICS_SERIES_V3_PATH, request.uri());
3090 assert_eq!(
3091 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "max_points")])),
3092 Some(1)
3093 );
3094
3095 harness.shutdown().await;
3096 }
3097
3098 #[tokio::test]
3099 async fn authoritative_v3_sketches_flush_previous_point_limit_batch() {
3100 let recorder = TestRecorder::default();
3103 let _local = metrics::set_default_local_recorder(&recorder);
3104
3105 let mut harness = RequestBuilderScenario::new(MetricsEncoderMode::V2Only, MetricsEncoderMode::V3Enabled)
3106 .with_max_points_per_payload(3)
3107 .with_flush_timeout(Duration::from_millis(250))
3108 .spawn()
3109 .await;
3110
3111 harness
3112 .push_metrics([
3113 Metric::distribution("authoritative.v3.sketch.points.one", [(123, 1.0), (124, 2.0)]),
3114 Metric::distribution("authoritative.v3.sketch.points.two", [(123, 3.0), (124, 4.0)]),
3115 ])
3116 .await;
3117
3118 for stage in ["point-limit split", "timeout"] {
3119 let (_, request) = harness.next_http_request().await;
3120 assert_eq!(
3121 V3_SKETCHES_ENDPOINT_URI,
3122 request.uri(),
3123 "{stage} sketches payload should target the V3 sketches endpoint"
3124 );
3125 }
3126 assert_eq!(
3127 recorder.counter(("serializer.v3_payload_split_reason", &[("reason", "max_points")])),
3128 Some(1)
3129 );
3130
3131 harness.shutdown().await;
3132 }
3133
3134 fn contains_bytes(haystack: &[u8], needle: &[u8]) -> bool {
3135 haystack.windows(needle.len()).any(|window| window == needle)
3136 }
3137
3138 fn assert_contains_bytes(haystack: &[u8], needle: &[u8]) {
3139 assert!(
3140 contains_bytes(haystack, needle),
3141 "expected payload to contain bytes {:?}, got {:?}",
3142 needle,
3143 haystack
3144 );
3145 }
3146
3147 fn tag_set<const N: usize>(tags: [&'static str; N]) -> TagSet {
3148 tags.into_iter().map(Tag::from_static).collect()
3149 }
3150
3151 #[tokio::test]
3156 async fn v2_series_builder_enforces_max_series_points_per_payload() {
3157 let v2_endpoint_config = EndpointConfiguration::new(CompressionScheme::noop(), 10_000, 10_000, None);
3158 let mut builder = v2::create_v2_request_builder(MetricsEndpoint::SeriesV2, &v2_endpoint_config)
3159 .await
3160 .expect("V2 request builder should be created");
3161 builder
3162 .with_len_limits(usize::MAX, usize::MAX)
3163 .expect("byte limits should be accepted");
3164
3165 let mut total_points = 0;
3166 for i in 0..6_000u32 {
3167 let metric = Metric::gauge(
3168 Context::from_parts(MetaString::from(format!("g{i}")), TagSet::default()),
3169 [(1u64, 1.0), (2u64, 2.0)],
3170 );
3171 let mut pending = Some(metric);
3172 while let Some(metric) = pending.take() {
3173 if let Some(returned) = builder.encode(metric).await.expect("encode should not error") {
3174 for request in builder.flush().await {
3175 let (_, points, _) = request.expect("request should build");
3176 assert!(
3177 points <= 10_000,
3178 "payload carried {} points, over the 10000 limit",
3179 points
3180 );
3181 total_points += points;
3182 }
3183 pending = Some(returned);
3184 }
3185 }
3186 }
3187 for request in builder.flush().await {
3188 let (_, points, _) = request.expect("request should build");
3189 assert!(
3190 points <= 10_000,
3191 "final payload carried {} points, over the 10000 limit",
3192 points
3193 );
3194 total_points += points;
3195 }
3196
3197 assert_eq!(
3198 total_points, 12_000,
3199 "all points should be emitted across the split payloads"
3200 );
3201 }
3202}
3203
3204#[cfg(test)]
3205mod payload_limits {
3206 use super::{v2, DatadogMetricsConfiguration};
3207 use crate::common::datadog::{clamp_payload_limits, test_util::shared_configuration};
3208
3209 #[test]
3210 fn payload_limits_come_from_resolved_configuration() {
3211 let mut shared = shared_configuration();
3212 shared.metrics_encoding.max_payload_size = 4321;
3213 shared.metrics_encoding.max_uncompressed_payload_size = 8765;
3214 shared.metrics_encoding.max_series_payload_size = 1234;
3215 shared.metrics_encoding.max_series_uncompressed_payload_size = 5678;
3216 shared.metrics_encoding.max_series_points_per_payload = 500;
3217 shared.metrics_encoding.max_metrics_per_payload = 42;
3218
3219 let config = DatadogMetricsConfiguration::from_configuration(&shared);
3220
3221 assert_eq!(4321, config.max_payload_size);
3222 assert_eq!(8765, config.max_uncompressed_payload_size);
3223 assert_eq!(1234, config.max_series_payload_size);
3224 assert_eq!(5678, config.max_series_uncompressed_payload_size);
3225 assert_eq!(500, config.max_series_points_per_payload);
3226 assert_eq!(500, config.v3_payload_limits().max_points_per_payload);
3227 assert_eq!(42, config.max_metrics_per_payload);
3228 }
3229
3230 #[test]
3231 fn clamps_series_payload_limit_keys_to_api_limits() {
3232 let (uncompressed_limit, compressed_limit) = clamp_payload_limits(
3233 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT + 1,
3234 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT + 1,
3235 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT,
3236 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT,
3237 );
3238 assert_eq!(uncompressed_limit, v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT);
3239 assert_eq!(compressed_limit, v2::SERIES_V2_COMPRESSED_SIZE_LIMIT);
3240
3241 let (uncompressed_limit, compressed_limit) = clamp_payload_limits(
3242 5678,
3243 1234,
3244 v2::SERIES_V2_UNCOMPRESSED_SIZE_LIMIT,
3245 v2::SERIES_V2_COMPRESSED_SIZE_LIMIT,
3246 );
3247 assert_eq!(uncompressed_limit, 5678);
3248 assert_eq!(compressed_limit, 1234);
3249 }
3250}