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