saluki_components/encoders/datadog/metrics/
mod.rs

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/// Encoding mode for a metrics endpoint.
66#[derive(Clone, Copy, Debug, PartialEq, Eq)]
67enum MetricsEncoderMode {
68    /// Send V2 payloads only.
69    V2Only,
70    /// V3 is enabled for at least one endpoint; generate tagged V2 and V3 payloads so each endpoint
71    /// receives the protocol version configured for it.
72    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/// Datadog Metrics encoder.
129///
130/// Generates Datadog metrics payloads for the Datadog platform.
131#[derive(Clone)]
132#[cfg_attr(test, derive(Debug, PartialEq))]
133pub struct DatadogMetricsConfiguration {
134    /// Maximum number of input metrics to encode into a single request payload.
135    ///
136    /// This applies both to the series and sketches endpoints.
137    max_metrics_per_payload: usize,
138
139    /// Maximum compressed size, in bytes, of generic payloads.
140    ///
141    /// This applies to V1 JSON series payloads and sketch payloads, matching the Datadog Agent's generic payload
142    /// builder. V2 series payloads use the series-specific limit instead. The effective value is clamped to
143    /// the Agent's default intake-safe limit of 2,621,440 bytes, so larger configured values do not allow payloads that
144    /// intake may reject. If set to `0`, every non-empty compressed payload exceeds the limit and is dropped during
145    /// flush.
146    max_payload_size: usize,
147
148    /// Maximum uncompressed size, in bytes, of generic payloads.
149    ///
150    /// This applies to V1 JSON series payloads and sketch payloads, matching the Datadog Agent's generic payload
151    /// builder. V2 series payloads use the series-specific limit instead. The effective value
152    /// is clamped to the Agent's default intake-safe limit of 4,194,304 bytes, so larger configured values do not allow
153    /// payloads that intake may reject. Values smaller than the minimum endpoint framing size prevent the request
154    /// builder from starting.
155    max_uncompressed_payload_size: usize,
156
157    /// Maximum compressed size, in bytes, of a V2 series payload.
158    ///
159    /// This applies only when the V2 series API is in use. V1 series and sketches use the generic payload limit
160    /// instead. The effective value is clamped to the V2 series API limit of 512,000 bytes, so larger configured values
161    /// do not allow payloads that intake would reject. If set to `0`, every non-empty compressed payload exceeds the
162    /// limit and is dropped during flush.
163    max_series_payload_size: usize,
164
165    /// Maximum uncompressed size, in bytes, of a V2 series payload.
166    ///
167    /// This applies only when the V2 series API is in use. V1 series and sketches use the generic uncompressed payload
168    /// limit instead. The effective value is clamped to the V2 series API limit of 5,242,880 bytes, so larger configured
169    /// values do not allow payloads that intake would reject. This limit protects the encoder before compression, so
170    /// compressed payload size may still force a separate flush. Values smaller than the minimum endpoint framing size
171    /// prevent the request builder from starting.
172    max_series_uncompressed_payload_size: usize,
173
174    /// Maximum number of data points, across all series, to encode into a single series request payload.
175    ///
176    /// This applies only to series metrics (counters, gauges, rates, sets) and not to sketch metrics (histograms,
177    /// distributions). A single metric series may contribute multiple data points when it carries more than one
178    /// timestamp/value pair. When encoding an input would cause the running data point total to exceed this limit, the
179    /// current payload is flushed first and the input is placed in the next payload.
180    max_series_points_per_payload: usize,
181
182    /// Flush timeout for pending requests.
183    ///
184    /// When the destination has written metrics to the in-flight request payload, but it has not yet reached the
185    /// payload size limits that would force the payload to be flushed, the destination will wait for a period of time
186    /// before flushing the in-flight request payload. This allows for the possibility of other events to be processed
187    /// and written into the request payload, thereby maximizing the payload size and reducing the number of requests
188    /// generated and sent overall.
189    flush_timeout: Duration,
190
191    /// Compression kind to use for the request payloads.
192    compressor_kind: String,
193
194    /// Effective zstd compression level for the request payloads.
195    zstd_compressor_level: i32,
196
197    /// Whether to use the V2 API for series metrics.
198    ///
199    /// When `true`, series metrics are sent to the V2 protobuf endpoint (`/api/v2/series`). When `false`, series
200    /// metrics are sent to the legacy V1 JSON endpoint (`/api/v1/series`). Sketch metrics always use the V2 endpoint
201    /// (`/api/beta/sketches`) regardless of this setting.
202    use_v2_series_api: bool,
203
204    /// Whether to log metric payload contents before encoding.
205    ///
206    /// This logs decoded metric objects, not the encoded JSON/protobuf HTTP body.
207    log_payloads: bool,
208
209    /// Additional tags to apply to all forwarded metrics.
210    additional_tags: Option<SharedTagSet>,
211
212    /// V3 API configuration for per-endpoint V3 support.
213    ///
214    /// Configures which endpoints receive V3 payloads.
215    v3_api: V3ApiConfig,
216
217    /// Agent-compatible V3 API configuration.
218    use_v3_api: UseV3ApiConfig,
219
220    /// Metrics routing to an alternate intake.
221    opw_metrics: OpwMetricsConfiguration,
222
223    /// The primary metrics endpoint, as configured and not altered in any way.
224    primary_endpoint: String,
225
226    /// Additional endpoints that metrics may be dual-shipped to, keyed by endpoint URL with their API keys.
227    additional_endpoints: HashMap<String, Vec<String>>,
228}
229
230impl DatadogMetricsConfiguration {
231    /// Creates a new `DatadogMetricsConfiguration` from the resolved shared configuration.
232    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    /// Sets additional tags to be applied uniformly to all metrics forwarded by this destination.
258    pub fn with_additional_tags(mut self, additional_tags: SharedTagSet) -> Self {
259        self.additional_tags = Some(additional_tags);
260        self
261    }
262
263    /// Restricts endpoint-aware protocol selection to a single overridden metrics endpoint.
264    ///
265    /// This mirrors a forwarder branch that replaces the normal primary endpoint and removes additional and
266    /// alternate-intake endpoints, such as Multi-Region Failover.
267    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    /// Forces series metrics to use V2.
275    ///
276    /// This is used for local destinations that only accept the V2 series protocol, such as the Cluster Agent.
277    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            // Actually enforced by V3PayloadLimits, required for the
398            // constructor shared between V1/V2/V3.
399            usize::MAX,
400            self.additional_tags.clone(),
401        );
402
403        // Derive the encoding mode for each metric type from the configuration.
404        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            // We always give ourselves a minimum flush timeout of 10ms to allow for some very minimal amount of
456            // batching, while still practically flushing things almost immediately.
457            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        // TODO: How do we properly represent the requests we can generate that may be sitting around in-flight?
487        //
488        // Theoretically, we'll end up being limited by the size of the downstream forwarder's interconnect, and however
489        // many payloads it will buffer internally... so realistically the firm limit boils down to the forwarder itself
490        // but we'll have a hard time in the forwarder knowing the maximum size of any given payload being sent in, which
491        // then makes it hard to calculate a proper firm bound even though we know the rest of the values required to
492        // calculate the firm bound.
493        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            // Capture the size of the "split re-encode" buffers in the request builders, which is where we keep owned
502            // versions of metrics that we encode in case we need to actually re-encode them during a split operation.
503            .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        // Run our request builder task on the worker pool.
546        //
547        // The request builder task ignores the shutdown signal on purpose: it drains its incoming event buffer channel
548        // until the channel closes, which is what guarantees every buffered metric is encoded and dispatched.
549        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                        // Both channels between this task and the request builder task are bounded, and each task is
584                        // the producer on one and the consumer on the other. Blocking outright on `events_tx` would
585                        // stop us draining `payloads_rx`, which can deadlock: the builder task blocks sending a
586                        // payload once `payloads_tx` is full, so it never returns to receive from `events_rx`, so
587                        // capacity never frees up here. Keep draining payloads while we wait for capacity so that at
588                        // least one of the two channels is always making progress.
589                        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                                    // Our payloads channel is gone, which means our request builder task went away unexpectedly.
601                                    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 the events sender, which signals the request builder task to stop.
614        drop(events_tx);
615
616        // Continue draining the payloads receiver until it is closed.
617        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        // Draining `payloads_rx` to completion already implies the request builder finished: it owns the only sender,
624        // so the channel only closes once that child's future has run to completion (or been dropped).
625        debug!("Datadog Metrics encoder stopped.");
626
627        Ok(())
628    }
629}
630
631/// Logs the decoded contents of a metric prior to encoding.
632///
633/// This logs the metric object itself, not the encoded JSON/protobuf HTTP body.
634fn 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                    // A series metric whose points are all non-finite would encode to a series with no points, which
688                    // intake rejects as an empty value set. Drop it whole rather than emit an empty series.
689                    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                    // Figure out which endpoint the metric belongs to, and grab the relevant V2 builder/V3 storage.
696                    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                    // Store a copy of the metric in `maybe_v3_metrics` if it's present.
715                    //
716                    // We have to do this before encoding because `RequestBuilder::encode` consumes the metric. This also means we'll
717                    // need to _remove_ the metric if encoding fails.
718                    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                    // Attempt encoding the metric for V2 if configured.
726                    //
727                    // If the metric could not be encoded, remove its buffered V3 copy when this endpoint also requires
728                    // a V2 encoding.
729                    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 we're not already pending a flush, we'll start the countdown.
816                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                // Flush any pending series metrics.
827                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                // Flush any pending sketch metrics.
859                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            // Event buffers channel has been closed, and we have no pending flushing, so we're all done.
894            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    // Encode the metric. If we get it back, that means the current request is full, and we need to
925    // flush it before we can try to encode the metric again... so we'll hold on to it in that case
926    // before flushing and trying to encode it again.
927    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    // Now try to encode the metric again. If it fails again, we'll just log it because it shouldn't
945    // be possible to fail at this point, otherwise we would have already caught that the first
946    // time.
947    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            // TODO: Increment a counter here that metrics were dropped due to a flush failure.
976            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
1054/// Smallest compression ratio used when translating the compressed size limit into an uncompressed byte target.
1055///
1056/// Guards against dividing by a pathologically small ratio and producing an enormous target.
1057const V3_MIN_COMPRESSION_RATIO: f64 = 0.02;
1058
1059/// Fraction of the derived byte target that a batch is actually cut at.
1060///
1061/// Leaves headroom so that ordinary variation between batches does not push the finalized payload over the real limit
1062/// and force the expensive re-encode path.
1063const V3_BATCH_TARGET_MARGIN: f64 = 0.9;
1064
1065/// Weight given to the newest observation when updating the compression ratio.
1066const V3_RATIO_SMOOTHING: f64 = 0.25;
1067
1068/// Tracks the compression ratio observed on emitted V3 payloads, so batches can be sized against the compressed limit.
1069///
1070/// V3 payloads are columnar and only compressed once the whole batch is finalized, so unlike the V2 request builder
1071/// there is no live compressor to consult with a [`CompressionEstimator`] while batching. Instead we remember how well
1072/// recent payloads compressed and use that to turn the compressed size limit (which is the limit that actually binds
1073/// for metrics intake) into an uncompressed byte target for the next batch.
1074#[derive(Default)]
1075struct V3CompressionRatio {
1076    ratio: Option<f64>,
1077}
1078
1079impl V3CompressionRatio {
1080    /// Returns the uncompressed byte target that a batch should be cut at.
1081    fn uncompressed_target(&self, limits: V3PayloadLimits) -> usize {
1082        // Until a payload has been observed, assume the data will not compress at all. That errs towards batches that
1083        // are too small, which costs an extra payload, rather than too large, which costs a full re-encode.
1084        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        // Always leave room for at least one metric, however small the limits are.
1089        (((target as f64) * V3_BATCH_TARGET_MARGIN) as usize).max(1)
1090    }
1091
1092    /// Records the uncompressed and compressed sizes of an emitted payload.
1093    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
1106/// A contiguous run of metrics encoded into a single V3 payload.
1107struct V3Batch {
1108    encoded: V3EncodedMetrics,
1109    event_count: usize,
1110    data_point_count: usize,
1111}
1112
1113/// Encodes metrics starting at `*idx` into a single V3 payload, stopping once a payload limit or
1114/// `target_uncompressed_len` is reached, and advances `*idx` past everything consumed.
1115///
1116/// Sizing the batch while encoding it is what keeps this to a single encode pass. The alternative (encoding a large
1117/// batch and halving it whenever the result is too big) re-encodes and re-compresses the same metrics once per level of
1118/// splitting, which is dominated by dictionary rebuilding on high-cardinality data
1119fn 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            // This metric exceeds the point limit on its own, so it cannot fit in any payload.
1142            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        // Limits that have to be honoured before the metric goes in. The first metric of a batch is always admitted,
1158        // since it has already been checked to fit by itself.
1159        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        // The size contribution of a metric is only knowable once it is written, since dictionary deduplication makes
1180        // it depend on what came before. Cutting after the fact can overshoot by one metric, which the target's
1181        // headroom absorbs; the finalized payload is checked against the real limits regardless.
1182        if writer.estimated_uncompressed_len() >= target_uncompressed_len {
1183            cut_for_size = true;
1184            break;
1185        }
1186    }
1187
1188    // Only count this as a payload boundary if metrics actually remain: a batch that ends on the last metric was not
1189    // split by anything.
1190    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
1215/// Turns an encoded batch into a payload request, returning `None` if it does not fit the payload limits.
1216async 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    // Feed the observation back even when the request is too big: an oversized batch is exactly the case where the
1237    // target needs correcting.
1238    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    // Per-column compressed sizes are measured only for requests we actually emit. Measuring them eagerly in
1245    // `create_v3_request` would compress every column of every oversized attempt that gets discarded, which is pure
1246    // waste.
1247    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
1263/// Re-encodes an oversized range as progressively smaller ranges until each one fits.
1264///
1265/// This is the fallback for a batch that was sized against the byte target but still exceeded the real payload limits.
1266/// It re-encodes and re-compresses at every level of splitting, so it is deliberately only reached when the target
1267/// guessed wrong.
1268async 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    // `build_v3_batch` advances past metrics it drops. Preserve that decision when the accepted batch has to be
1273    // re-encoded: retrying the raw source range would otherwise resurrect those metrics.
1274    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            // The encoded request is too large and this range cannot be split any further.
1318            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        // Retry this oversized range as two smaller ranges, preserving the original metric order.
1333        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            // `build_v3_batch` always consumes at least one metric when one is available, but guard against spinning
1356            // on a batch that somehow produced nothing.
1357            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            // The batch we sized against the byte target still came out over the real limits, so fall back to
1374            // re-encoding it as progressively smaller ranges. The target self-corrects from the ratio we record on
1375            // every emitted payload, so this should be rare.
1376            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
1420/// Measures the compressed size of each V3 column, for telemetry purposes.
1421///
1422/// This is deliberately separate from building the request: it is only worth paying for columns belonging to a request
1423/// that will actually be sent, since oversized requests are discarded and re-encoded as smaller ranges.
1424async 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
1478// Encodes a batch of metrics to V3 columnar format.
1479fn 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
1494/// Writes a single metric to the V3 writer.
1495pub(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    // Tags - chain instrumented + additional + origin tags
1532    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    // Resources - extract host and, for series, promoted resource tags.
1545    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    // Origin metadata
1576    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    // Points based on metric type
1598    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                // Scale by interval as done in V2. A zero interval is an unnormalized rate.
1612                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            // Set values are already converted to count in the iterator
1626            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                // Convert histogram to DDSketch
1656                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
1693/// Creates a V3 HTTP request from encoded payload data.
1694async fn create_v3_request(
1695    endpoint_uri: &str, encoded: V3EncodedMetrics, compression_scheme: CompressionScheme,
1696) -> Result<V3EncodedRequest, GenericError> {
1697    // Keep the wire payload as one continuous compressed stream. Per-column compressed sizes are measured
1698    // independently for telemetry.
1699    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    /// Returns an encoder built from shared configuration shaped like a translated Agent configuration.
1792    fn metrics_config_from(shared: &SharedConfiguration) -> DatadogMetricsConfiguration {
1793        DatadogMetricsConfiguration::from_configuration(shared)
1794    }
1795
1796    /// Returns the Agent-compatible V3 series settings a default configuration resolves to.
1797    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        // A selected endpoint means alternate-intake routing is enabled; it carries the URL and V3 override.
1844        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        // The primary Datadog endpoint is V3-authoritative under `datadog_only`, while the additional
1864        // endpoint is not a Datadog URL and stays on V2.
1865        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        // An alternate intake whose URL cannot resolve does not fall through to the Vector intake.
1986        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        // Values are 1, 2, 3, 4, 3 (rate scaled by 10s), and 1 (set cardinality), encoded as sint64.
2177        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    /// Collected form of a payload emitted by `encode_and_flush_v3_payload_requests`.
2287    struct CollectedV3Payload {
2288        event_count: usize,
2289        data_point_count: usize,
2290        request: Request<FrozenChunkedBytesBuffer>,
2291    }
2292
2293    /// Drives `encode_and_flush_v3_payload_requests` and collects the payloads it emits.
2294    ///
2295    /// The channel is sized generously because nothing drains it concurrently here; production drains it from the
2296    /// encoder's main task.
2297    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    /// Drives `build_v3_batch` to exhaustion, returning the metric range and event count of each batch it cut.
2472    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        // A target of one byte forces a cut after every metric, since the estimate is non-zero once anything is
2542        // written. This is the mechanism that keeps high-cardinality batches from being encoded and then re-split.
2543        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        // Dropping the oversized metric no longer forces a payload boundary: the metrics either side of it share a
2593        // batch, since together they still fit the point limit.
2594        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        // With nothing observed yet, assume no compression: the target tracks the compressed limit directly so the
2622        // first batch cannot wildly overshoot.
2623        let initial_target = ratio.uncompressed_target(limits);
2624        assert_eq!((500_000.0 * V3_BATCH_TARGET_MARGIN) as usize, initial_target);
2625
2626        // After observing 4:1 compression, the target grows towards the uncompressed budget those bytes imply.
2627        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        // It is still capped by the uncompressed limit, however well the data compresses.
2635        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, // field 25, length-delimited.
2660            0x0c, // field payload length: varint string length + string bytes.
2661            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, // field 26, length-delimited.
2673            0x02, // packed field payload length.
2674            0x02, 0x00, // sparse unit refs for metrics 1 and 3 only: refs [1, 1] -> deltas [1, 0].
2675        ];
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, // field 4, length-delimited.
2730            0x25, // field payload length.
2731            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, // field 4, length-delimited.
2763            0x34, // field payload length.
2764            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, // field 4, length-delimited.
2792            0x0c, // field payload length.
2793            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    /// Distinct inputs for a [`run_request_builder`] integration test.
2799    struct RequestBuilderScenario {
2800        series_mode: MetricsEncoderMode,
2801        sketches_mode: MetricsEncoderMode,
2802        payload_limits: V3PayloadLimits,
2803        flush_timeout: Duration,
2804    }
2805
2806    impl RequestBuilderScenario {
2807        /// Creates a scenario with unbounded V3 size/point limits and a 10 ms flush timeout.
2808        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        /// Sets the V3 per-payload point limit that drives point-count split flushes.
2818        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        /// Overrides the flush timeout used to bound pending flushes.
2824        fn with_flush_timeout(mut self, flush_timeout: Duration) -> Self {
2825            self.flush_timeout = flush_timeout;
2826            self
2827        }
2828
2829        /// Spawns [`run_request_builder`] for this scenario, returning a harness for pushing metrics and draining
2830        /// flushed payloads.
2831        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    /// A spawned [`run_request_builder`] task plus its input/output channels.
2866    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        /// Sends a single event buffer carrying `metrics` to the request builder.
2874        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        /// Receives the next flushed payload, unwrapping it into its HTTP metadata and request.
2889        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        /// Asserts that no payload is flushed within `window`.
2901        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        /// Drops the events sender and waits for the request builder task to stop cleanly.
2909        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        // Authoritative V3 batches independently of V2. With a 3-point limit, the first two-point metric fits but
2921        // the second would exceed it, so the first flushes as a point-limit split and the second flushes on the
2922        // pending-flush timeout.
2923        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        // Point-count split flushes the first metric before the second exceeds the limit.
2940        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        // The carried-over metric flushes when the pending-flush timeout fires.
2945        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        // The same authoritative point-limit split applies to sketches, which always target the V3 sketches
2958        // endpoint: the first distribution flushes as a point-limit split, the second on the pending-flush timeout.
2959        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    // Regression test to ensure the V2 series request builder enforces `max_series_points_per_payload`.
3009    //
3010    // The test encodes more total points than the configured limit and asserts the builder splits them across
3011    // multiple payloads without any payload exceeding the limit.
3012    #[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}