saluki_components/encoders/datadog/metrics/
mod.rs

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