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