saluki_components/encoders/datadog/metrics/
mod.rs

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