saluki_components/encoders/datadog/metrics/
mod.rs

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