Skip to main content

servers/otlp/
metrics.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use ahash::HashSet;
16use api::greptime_proto::io::prometheus::write::v2::histogram::{
17    Count as PromCount, ResetHint, ZeroCount as PromZeroCount,
18};
19use api::greptime_proto::io::prometheus::write::v2::{BucketSpan, Histogram as PromHistogram};
20use api::v1::value::ValueData;
21use api::v1::{RowInsertRequests, SemanticType, Value};
22use common_grpc::precision::Precision;
23use common_query::native_histogram::{
24    encode_native_histogram, native_histogram_column_schema, native_histogram_value_type,
25};
26use common_query::prelude::{
27    GREPTIME_COUNT, GREPTIME_TEMPORALITY_DELTA, OTLP_AGGREGATION_TEMPORALITY_LABEL,
28    greptime_timestamp, greptime_value,
29};
30use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS;
31use common_telemetry::warn;
32use lazy_static::lazy_static;
33use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest;
34use otel_arrow_rust::proto::opentelemetry::common::v1::{AnyValue, KeyValue, any_value};
35use otel_arrow_rust::proto::opentelemetry::metrics::v1::{metric, number_data_point, *};
36use session::protocol_ctx::{MetricType, OtlpMetricCtx};
37use table::requests::{
38    METADATA_QUALITY_DECLARED, METRIC_TEMPORALITY_CUMULATIVE, METRIC_TEMPORALITY_DELTA,
39    SEMANTIC_METRIC_METADATA_QUALITY, SEMANTIC_METRIC_ORIGINAL_NAME, SEMANTIC_METRIC_TEMPORALITY,
40    SEMANTIC_METRIC_TYPE, SEMANTIC_METRIC_UNIT,
41};
42
43use crate::error::{self, Result};
44use crate::metrics::OTLP_EXPONENTIAL_HISTOGRAM_REJECTED_DATA_POINTS;
45use crate::otlp::trace::{KEY_SERVICE_INSTANCE_ID, KEY_SERVICE_NAME, KEY_SERVICE_NAMESPACE};
46use crate::query_handler::MetricsIngestOutcome;
47use crate::row_writer::{self, MultiTableData, TableData};
48pub use crate::semantic::SemanticIndex;
49use crate::semantic::{
50    METRIC_TYPE_COUNTER, METRIC_TYPE_GAUGE, METRIC_TYPE_HISTOGRAM, METRIC_TYPE_SUMMARY,
51    METRIC_TYPE_UPDOWN_COUNTER,
52};
53
54mod resource_info;
55mod translator;
56
57pub use resource_info::OTEL_RESOURCE_INFO_TABLE_NAME;
58use resource_info::ResourceInfoData;
59pub use translator::legacy_normalize_otlp_name;
60pub(crate) use translator::ucum_to_openmetrics_unit;
61use translator::{translate_label_name, translate_metric_name};
62
63/// the default column count for table writer
64const APPROXIMATE_COLUMN_COUNT: usize = 8;
65
66const COUNT_TABLE_SUFFIX: &str = "_count";
67const SUM_TABLE_SUFFIX: &str = "_sum";
68const BUCKET_TABLE_SUFFIX: &str = "_bucket";
69
70const JOB_KEY: &str = "job";
71const INSTANCE_KEY: &str = "instance";
72
73// see: https://prometheus.io/docs/guides/opentelemetry/#promoting-resource-attributes
74const DEFAULT_PROMOTE_ATTRS: [&str; 19] = [
75    "service.instance.id",
76    "service.name",
77    "service.namespace",
78    "service.version",
79    "cloud.availability_zone",
80    "cloud.region",
81    "container.name",
82    "deployment.environment",
83    "deployment.environment.name",
84    "k8s.cluster.name",
85    "k8s.container.name",
86    "k8s.cronjob.name",
87    "k8s.daemonset.name",
88    "k8s.deployment.name",
89    "k8s.job.name",
90    "k8s.namespace.name",
91    "k8s.pod.name",
92    "k8s.replicaset.name",
93    "k8s.statefulset.name",
94];
95
96lazy_static! {
97    static ref DEFAULT_PROMOTE_ATTRS_SET: HashSet<String> =
98        HashSet::from_iter(DEFAULT_PROMOTE_ATTRS.iter().map(|s| s.to_string()));
99}
100
101const OTEL_SCOPE_NAME: &str = "name";
102const OTEL_SCOPE_VERSION: &str = "version";
103const OTEL_SCOPE_SCHEMA_URL: &str = "schema_url";
104const MIN_EXPONENTIAL_HISTOGRAM_SCALE: i32 = -4;
105const MAX_EXPONENTIAL_HISTOGRAM_SCALE: i32 = 8;
106const MAX_REJECTION_MESSAGE_BYTES: usize = 512;
107
108/// Result of converting one OTLP metrics request.
109#[derive(Debug)]
110pub struct MetricsConversion {
111    pub requests: RowInsertRequests,
112    /// Row count of `requests`; the resource descriptor is not counted.
113    pub rows: usize,
114    /// Per-table semantic index for the auto-create path to stamp as table
115    /// options; covers the descriptor table too.
116    pub semantic_index: SemanticIndex,
117    /// The synthesized resource descriptor, written separately from the
118    /// metric data. See [`resource_info`].
119    pub resource_info: Option<RowInsertRequests>,
120    pub outcome: MetricsIngestOutcome,
121}
122
123/// Convert OpenTelemetry metrics to GreptimeDB insert requests
124///
125/// See
126/// <https://github.com/open-telemetry/opentelemetry-proto/blob/main/opentelemetry/proto/metrics/v1/metrics.proto>
127/// for data structure of OTLP metrics.
128pub fn to_grpc_insert_requests(
129    request: ExportMetricsServiceRequest,
130    metric_ctx: &mut OtlpMetricCtx,
131) -> Result<MetricsConversion> {
132    let mut table_writer = MultiTableData::default();
133    let mut semantic_index = SemanticIndex::default();
134    let mut outcome = MetricsIngestOutcome::default();
135    let mut resource_info = ResourceInfoData::default();
136
137    for resource in &request.resource_metrics {
138        if metric_ctx.resource_info
139            && !metric_ctx.is_legacy
140            && let Some(r) = resource.resource.as_ref()
141        {
142            resource_info.observe(&r.attributes, resource);
143        }
144
145        let resource_attrs = resource.resource.as_ref().map(|r| {
146            let mut attrs = r.attributes.clone();
147            process_resource_attrs(&mut attrs, metric_ctx);
148            attrs
149        });
150
151        for scope in &resource.scope_metrics {
152            let scope_attrs = process_scope_attrs(scope, metric_ctx);
153
154            for metric in &scope.metrics {
155                if metric.data.is_none() {
156                    continue;
157                }
158                if let Some(t) = metric.data.as_ref().map(from_metric_type) {
159                    metric_ctx.set_metric_type(t);
160                }
161
162                encode_metrics(
163                    &mut table_writer,
164                    metric,
165                    resource_attrs.as_ref(),
166                    scope_attrs.as_ref(),
167                    metric_ctx,
168                    &mut semantic_index,
169                    &mut outcome,
170                )?;
171            }
172        }
173    }
174
175    let (requests, rows) = table_writer.into_row_insert_requests();
176
177    validate_sample_kinds(&requests)?;
178
179    // The metric and the descriptor would fight over the same table across
180    // two engines. The user's metric wins.
181    let resource_info = if !metric_ctx.resource_info {
182        None
183    } else if requests
184        .inserts
185        .iter()
186        .any(|r| r.table_name == OTEL_RESOURCE_INFO_TABLE_NAME)
187    {
188        warn!(
189            "Skipping OTLP resource descriptor synthesis: the request writes \
190             a metric named `{OTEL_RESOURCE_INFO_TABLE_NAME}`"
191        );
192        None
193    } else {
194        resource_info.into_row_insert_requests()?
195    };
196    if resource_info.is_some() {
197        semantic_index.record_scalar(
198            OTEL_RESOURCE_INFO_TABLE_NAME,
199            SEMANTIC_METRIC_TYPE,
200            crate::semantic::METRIC_TYPE_INFO,
201        );
202        semantic_index.record_scalar(
203            OTEL_RESOURCE_INFO_TABLE_NAME,
204            SEMANTIC_METRIC_METADATA_QUALITY,
205            METADATA_QUALITY_DECLARED,
206        );
207    }
208
209    Ok(MetricsConversion {
210        requests,
211        rows,
212        semantic_index,
213        outcome,
214        resource_info,
215    })
216}
217
218fn validate_sample_kinds(requests: &RowInsertRequests) -> Result<()> {
219    for request in &requests.inserts {
220        let Some(rows) = &request.rows else {
221            continue;
222        };
223        let field_count = rows
224            .schema
225            .iter()
226            .filter(|column| column.semantic_type == SemanticType::Field as i32)
227            .count();
228        let has_native_histogram = rows.schema.iter().any(|column| {
229            column.semantic_type == SemanticType::Field as i32
230                && api::helper::is_column_type_value_eq(
231                    column.datatype,
232                    column.datatype_extension.clone(),
233                    native_histogram_value_type(),
234                )
235        });
236        if has_native_histogram && field_count != 1 {
237            return Err(error::InvalidParameterSnafu {
238                reason: format!(
239                    "OTLP metric `{}` cannot mix native histogram and float sample fields",
240                    request.table_name
241                ),
242            }
243            .build());
244        }
245    }
246    Ok(())
247}
248
249/// The tables a metric emits and their per-table `metric.type`. Histogram fans
250/// out into `_bucket` (the histogram) plus `_sum`/`_count` counters; summary
251/// fans out into the quantile table plus `_count`/`_sum` counters (legacy
252/// summary stays a single table).
253fn emitted_semantic_tables(
254    metric_type: &MetricType,
255    is_legacy: bool,
256    base: &str,
257) -> Vec<(String, &'static str)> {
258    match metric_type {
259        MetricType::Gauge => vec![(base.to_string(), METRIC_TYPE_GAUGE)],
260        MetricType::MonotonicSum => vec![(base.to_string(), METRIC_TYPE_COUNTER)],
261        MetricType::NonMonotonicSum => vec![(base.to_string(), METRIC_TYPE_UPDOWN_COUNTER)],
262        MetricType::Histogram => vec![
263            (
264                format!("{base}{BUCKET_TABLE_SUFFIX}"),
265                METRIC_TYPE_HISTOGRAM,
266            ),
267            (format!("{base}{SUM_TABLE_SUFFIX}"), METRIC_TYPE_COUNTER),
268            (format!("{base}{COUNT_TABLE_SUFFIX}"), METRIC_TYPE_COUNTER),
269        ],
270        MetricType::Summary if is_legacy => vec![(base.to_string(), METRIC_TYPE_SUMMARY)],
271        MetricType::Summary => vec![
272            (base.to_string(), METRIC_TYPE_SUMMARY),
273            (format!("{base}{COUNT_TABLE_SUFFIX}"), METRIC_TYPE_COUNTER),
274            (format!("{base}{SUM_TABLE_SUFFIX}"), METRIC_TYPE_COUNTER),
275        ],
276        MetricType::ExponentialHistogram => {
277            vec![(base.to_string(), METRIC_TYPE_HISTOGRAM)]
278        }
279        MetricType::Init => vec![],
280    }
281}
282
283/// Maps OTLP `aggregation_temporality` to the semantic value, or `None` when the
284/// instrument has no temporality (gauge/summary) or it is unspecified.
285fn temporality_value(data: &metric::Data) -> Option<&'static str> {
286    let raw = match data {
287        metric::Data::Sum(sum) => sum.aggregation_temporality,
288        metric::Data::Histogram(hist) => hist.aggregation_temporality,
289        metric::Data::ExponentialHistogram(hist) => hist.aggregation_temporality,
290        _ => return None,
291    };
292    match AggregationTemporality::try_from(raw) {
293        Ok(AggregationTemporality::Delta) => Some(METRIC_TEMPORALITY_DELTA),
294        Ok(AggregationTemporality::Cumulative) => Some(METRIC_TEMPORALITY_CUMULATIVE),
295        _ => None,
296    }
297}
298
299/// Records the declared metric-level semantic keys for every table this metric
300/// emits.
301fn record_metric_semantics(
302    index: &mut SemanticIndex,
303    metric: &Metric,
304    name: &str,
305    metric_ctx: &OtlpMetricCtx,
306) {
307    let emitted = emitted_semantic_tables(&metric_ctx.metric_type, metric_ctx.is_legacy, name);
308    if emitted.is_empty() {
309        return;
310    }
311
312    let temporality = metric.data.as_ref().and_then(temporality_value);
313    let unit = metric.unit.trim();
314    // `original_name` is meaningful only when translation renamed the metric.
315    let original_name = (name != metric.name.as_str()).then_some(metric.name.as_str());
316
317    for (table, metric_type) in &emitted {
318        index.record_scalar(table, SEMANTIC_METRIC_TYPE, metric_type);
319        index.record_scalar(
320            table,
321            SEMANTIC_METRIC_METADATA_QUALITY,
322            METADATA_QUALITY_DECLARED,
323        );
324        if let Some(temporality) = temporality {
325            index.record_scalar(table, SEMANTIC_METRIC_TEMPORALITY, temporality);
326        }
327        if !unit.is_empty() {
328            index.record_scalar(table, SEMANTIC_METRIC_UNIT, unit);
329        }
330        if let Some(original_name) = original_name {
331            index.record_scalar(table, SEMANTIC_METRIC_ORIGINAL_NAME, original_name);
332        }
333    }
334}
335
336fn from_metric_type(data: &metric::Data) -> MetricType {
337    match data {
338        metric::Data::Gauge(_) => MetricType::Gauge,
339        metric::Data::Sum(s) => {
340            if s.is_monotonic {
341                MetricType::MonotonicSum
342            } else {
343                MetricType::NonMonotonicSum
344            }
345        }
346        metric::Data::Histogram(_) => MetricType::Histogram,
347        metric::Data::ExponentialHistogram(_) => MetricType::ExponentialHistogram,
348        metric::Data::Summary(_) => MetricType::Summary,
349    }
350}
351
352/// Non-scalar values (bool, arrays, maps, bytes) are not representable as tags.
353fn scalar_value_string(value: Option<&AnyValue>) -> Option<String> {
354    match value.and_then(|v| v.value.as_ref())? {
355        any_value::Value::StringValue(s) => Some(s.clone()),
356        any_value::Value::IntValue(v) => Some(v.to_string()),
357        any_value::Value::DoubleValue(v) => Some(v.to_string()),
358        _ => None,
359    }
360}
361
362/// Prometheus-style `(job, instance)` identity. Per the OTel Prometheus
363/// compatibility spec `job` folds in `service.namespace` when present; a
364/// resource without `service.name` gets no job rather than a fabricated one.
365/// Both fields are optional strings, so a named type keeps a swapped
366/// destructuring from type-checking at the call sites.
367pub(crate) struct ServiceIdentity {
368    pub job: Option<String>,
369    pub instance: Option<String>,
370}
371
372pub(crate) fn service_identity(attrs: &[KeyValue]) -> ServiceIdentity {
373    let mut name = None;
374    let mut namespace = None;
375    let mut instance = None;
376    for kv in attrs {
377        match kv.key.as_str() {
378            KEY_SERVICE_NAME => name = scalar_value_string(kv.value.as_ref()),
379            KEY_SERVICE_NAMESPACE => namespace = scalar_value_string(kv.value.as_ref()),
380            KEY_SERVICE_INSTANCE_ID => instance = scalar_value_string(kv.value.as_ref()),
381            _ => {}
382        }
383    }
384    let job = name.map(|name| match namespace {
385        Some(ns) if !ns.is_empty() => format!("{ns}/{name}"),
386        _ => name,
387    });
388    ServiceIdentity { job, instance }
389}
390
391fn string_key_value(key: &str, value: String) -> KeyValue {
392    KeyValue {
393        key: key.to_string(),
394        value: Some(AnyValue {
395            value: Some(any_value::Value::StringValue(value)),
396        }),
397    }
398}
399
400fn process_resource_attrs(attrs: &mut Vec<KeyValue>, metric_ctx: &OtlpMetricCtx) {
401    if metric_ctx.is_legacy {
402        return;
403    }
404
405    // remap the service identity attributes to job and instance
406    let ServiceIdentity { job, instance } = service_identity(attrs);
407
408    // if promote all, then exclude the list, else, include the list
409    if metric_ctx.promote_all_resource_attrs {
410        attrs.retain(|kv| !metric_ctx.resource_attrs.contains(&kv.key));
411    } else {
412        attrs.retain(|kv| {
413            metric_ctx.resource_attrs.contains(&kv.key)
414                || DEFAULT_PROMOTE_ATTRS_SET.contains(&kv.key)
415        });
416    }
417
418    if let Some(job) = job {
419        attrs.push(string_key_value(JOB_KEY, job));
420    }
421    if let Some(instance) = instance {
422        attrs.push(string_key_value(INSTANCE_KEY, instance));
423    }
424}
425
426fn process_scope_attrs(scope: &ScopeMetrics, metric_ctx: &OtlpMetricCtx) -> Option<Vec<KeyValue>> {
427    if metric_ctx.is_legacy {
428        return scope.scope.as_ref().map(|s| s.attributes.clone());
429    };
430
431    if !metric_ctx.promote_scope_attrs {
432        return None;
433    }
434
435    // persist scope attrs with name, version and schema_url
436    scope.scope.as_ref().map(|s| {
437        let mut attrs = s.attributes.clone();
438        attrs.push(KeyValue {
439            key: OTEL_SCOPE_NAME.to_string(),
440            value: Some(AnyValue {
441                value: Some(any_value::Value::StringValue(s.name.clone())),
442            }),
443        });
444        attrs.push(KeyValue {
445            key: OTEL_SCOPE_VERSION.to_string(),
446            value: Some(AnyValue {
447                value: Some(any_value::Value::StringValue(s.version.clone())),
448            }),
449        });
450        attrs.push(KeyValue {
451            key: OTEL_SCOPE_SCHEMA_URL.to_string(),
452            value: Some(AnyValue {
453                value: Some(any_value::Value::StringValue(scope.schema_url.clone())),
454            }),
455        });
456        attrs
457    })
458}
459
460fn encode_metrics(
461    table_writer: &mut MultiTableData,
462    metric: &Metric,
463    resource_attrs: Option<&Vec<KeyValue>>,
464    scope_attrs: Option<&Vec<KeyValue>>,
465    metric_ctx: &OtlpMetricCtx,
466    semantic_index: &mut SemanticIndex,
467    outcome: &mut MetricsIngestOutcome,
468) -> Result<()> {
469    let name = if metric_ctx.is_legacy {
470        legacy_normalize_otlp_name(&metric.name)
471    } else {
472        translate_metric_name(
473            metric,
474            &metric_ctx.metric_type,
475            metric_ctx.metric_translation_strategy,
476        )
477    };
478
479    let emitted = if let Some(data) = &metric.data {
480        match data {
481            metric::Data::Gauge(gauge) => {
482                encode_gauge(
483                    table_writer,
484                    &name,
485                    gauge,
486                    resource_attrs,
487                    scope_attrs,
488                    metric_ctx,
489                )?;
490                add_accepted_data_points(outcome, gauge.data_points.len())?;
491                !gauge.data_points.is_empty()
492            }
493            metric::Data::Sum(sum) => {
494                encode_sum(
495                    table_writer,
496                    &name,
497                    sum,
498                    resource_attrs,
499                    scope_attrs,
500                    metric_ctx,
501                )?;
502                add_accepted_data_points(outcome, sum.data_points.len())?;
503                !sum.data_points.is_empty()
504            }
505            metric::Data::Summary(summary) => {
506                encode_summary(
507                    table_writer,
508                    &name,
509                    summary,
510                    resource_attrs,
511                    scope_attrs,
512                    metric_ctx,
513                )?;
514                add_accepted_data_points(outcome, summary.data_points.len())?;
515                !summary.data_points.is_empty()
516            }
517            metric::Data::Histogram(hist) => encode_histogram(
518                table_writer,
519                &name,
520                hist,
521                resource_attrs,
522                scope_attrs,
523                metric_ctx,
524                outcome,
525            )?,
526            metric::Data::ExponentialHistogram(hist) => encode_exponential_histogram(
527                table_writer,
528                &name,
529                hist,
530                resource_attrs,
531                scope_attrs,
532                metric_ctx,
533                outcome,
534            )?,
535        }
536    } else {
537        false
538    };
539
540    if emitted {
541        // Stamp semantic metadata only after at least one row was accepted.
542        record_metric_semantics(semantic_index, metric, &name, metric_ctx);
543    }
544
545    Ok(())
546}
547
548fn add_accepted_data_points(outcome: &mut MetricsIngestOutcome, count: usize) -> Result<()> {
549    let count = i64::try_from(count).map_err(|_| {
550        error::InvalidParameterSnafu {
551            reason: "OTLP metrics data-point count exceeds i64",
552        }
553        .build()
554    })?;
555    outcome.accepted_data_points =
556        outcome
557            .accepted_data_points
558            .checked_add(count)
559            .ok_or_else(|| {
560                error::InvalidParameterSnafu {
561                    reason: "OTLP accepted data-point count overflows i64",
562                }
563                .build()
564            })?;
565    Ok(())
566}
567
568fn reject_data_points(
569    outcome: &mut MetricsIngestOutcome,
570    count: usize,
571    reason: impl FnOnce() -> String,
572) -> Result<()> {
573    if count == 0 {
574        return Ok(());
575    }
576    let count = i64::try_from(count).map_err(|_| {
577        error::InvalidParameterSnafu {
578            reason: "OTLP rejected data-point count exceeds i64",
579        }
580        .build()
581    })?;
582    outcome.rejected_data_points =
583        outcome
584            .rejected_data_points
585            .checked_add(count)
586            .ok_or_else(|| {
587                error::InvalidParameterSnafu {
588                    reason: "OTLP rejected data-point count overflows i64",
589                }
590                .build()
591            })?;
592    append_rejection_message(&mut outcome.error_message, reason);
593    Ok(())
594}
595
596fn append_rejection_message(message: &mut Option<String>, reason: impl FnOnce() -> String) {
597    let message = message.get_or_insert_with(String::new);
598    let separator = if message.is_empty() { "" } else { "; " };
599    let Some(available) = MAX_REJECTION_MESSAGE_BYTES.checked_sub(message.len()) else {
600        return;
601    };
602    if available <= separator.len() {
603        return;
604    }
605    let reason = reason();
606    message.push_str(separator);
607
608    let available = MAX_REJECTION_MESSAGE_BYTES - message.len();
609    if reason.len() <= available {
610        message.push_str(&reason);
611        return;
612    }
613
614    const ELLIPSIS: &str = "...";
615    let mut end = available.saturating_sub(ELLIPSIS.len());
616    while !reason.is_char_boundary(end) {
617        end -= 1;
618    }
619    message.push_str(&reason[..end]);
620    if available >= ELLIPSIS.len() {
621        message.push_str(ELLIPSIS);
622    }
623}
624
625fn encode_exponential_histogram(
626    table_writer: &mut MultiTableData,
627    name: &str,
628    histogram: &ExponentialHistogram,
629    resource_attrs: Option<&Vec<KeyValue>>,
630    scope_attrs: Option<&Vec<KeyValue>>,
631    metric_ctx: &OtlpMetricCtx,
632    outcome: &mut MetricsIngestOutcome,
633) -> Result<bool> {
634    if let Err(rejection) = exponential_histogram_gate(histogram) {
635        reject_data_points(outcome, histogram.data_points.len(), || {
636            rejection.message(name)
637        })?;
638        OTLP_EXPONENTIAL_HISTOGRAM_REJECTED_DATA_POINTS
639            .with_label_values(&[rejection.reason_label()])
640            .inc_by(histogram.data_points.len() as u64);
641        return Ok(false);
642    }
643
644    let column_schema = native_histogram_column_schema().map_err(|error| {
645        error::InternalSnafu {
646            err_msg: error.to_string(),
647        }
648        .build()
649    })?;
650    let mut emitted = false;
651    for (index, data_point) in histogram.data_points.iter().enumerate() {
652        let (value, timestamp_nanos) = match exponential_histogram_value(data_point) {
653            Ok(value) => value,
654            Err(reason) => {
655                reject_data_points(outcome, 1, || {
656                    format!("metric `{name}` data point {index}: {reason}")
657                })?;
658                OTLP_EXPONENTIAL_HISTOGRAM_REJECTED_DATA_POINTS
659                    .with_label_values(&["invalid_data_point"])
660                    .inc();
661                continue;
662            }
663        };
664
665        let table = table_writer.get_or_default_table_data(
666            name,
667            APPROXIMATE_COLUMN_COUNT,
668            histogram.data_points.len(),
669        );
670        let mut row = table.alloc_one_row();
671        write_tags_and_timestamp(
672            table,
673            &mut row,
674            resource_attrs,
675            scope_attrs,
676            Some(data_point.attributes.as_ref()),
677            timestamp_nanos,
678            metric_ctx,
679        )?;
680        row_writer::write_by_schema(
681            table,
682            std::iter::once((column_schema.clone(), Some(value))),
683            &mut row,
684        )?;
685        table.add_row(row);
686        add_accepted_data_points(outcome, 1)?;
687        emitted = true;
688    }
689
690    Ok(emitted)
691}
692
693pub(crate) enum ExponentialHistogramRejection {
694    DeltaTemporality,
695    UnspecifiedTemporality,
696}
697
698impl ExponentialHistogramRejection {
699    fn reason_label(&self) -> &'static str {
700        match self {
701            Self::DeltaTemporality => "delta_temporality",
702            Self::UnspecifiedTemporality => "unspecified_temporality",
703        }
704    }
705
706    fn message(&self, name: &str) -> String {
707        match self {
708            Self::DeltaTemporality => format!(
709                "metric `{name}` uses delta OTLP exponential histograms; only cumulative temporality is supported"
710            ),
711            Self::UnspecifiedTemporality => format!(
712                "metric `{name}` has unspecified OTLP exponential histogram temporality; cumulative temporality is required"
713            ),
714        }
715    }
716}
717
718/// Whole-metric acceptance, decided once for the encoder and for the resource
719/// descriptor, which must not describe a resource whose data was rejected.
720/// Individual points can still fail [`exponential_histogram_value`].
721/// Raw-delta storage currently supports sums and explicit histograms; exponential
722/// histograms still require cumulative input. Delta-to-cumulative conversion is out of scope.
723pub(crate) fn exponential_histogram_gate(
724    histogram: &ExponentialHistogram,
725) -> std::result::Result<(), ExponentialHistogramRejection> {
726    match AggregationTemporality::try_from(histogram.aggregation_temporality) {
727        Ok(AggregationTemporality::Cumulative) => Ok(()),
728        Ok(AggregationTemporality::Delta) => Err(ExponentialHistogramRejection::DeltaTemporality),
729        _ => Err(ExponentialHistogramRejection::UnspecifiedTemporality),
730    }
731}
732
733pub(crate) fn exponential_histogram_value(
734    data_point: &ExponentialHistogramDataPoint,
735) -> std::result::Result<(ValueData, i64), String> {
736    if data_point.start_time_unix_nano > data_point.time_unix_nano {
737        return Err(format!(
738            "start_time_unix_nano {} exceeds time_unix_nano {}",
739            data_point.start_time_unix_nano, data_point.time_unix_nano
740        ));
741    }
742
743    let timestamp_nanos = i64::try_from(data_point.time_unix_nano)
744        .map_err(|_| format!("time_unix_nano {} overflows i64", data_point.time_unix_nano))?;
745    let timestamp = timestamp_millis(data_point.time_unix_nano, "time_unix_nano")?;
746    let start_timestamp =
747        timestamp_millis(data_point.start_time_unix_nano, "start_time_unix_nano")?;
748
749    let no_recorded_value = data_point.flags & DataPointFlags::NoRecordedValueMask as u32 != 0;
750    let histogram = if no_recorded_value {
751        PromHistogram {
752            sum: f64::from_bits(PROMETHEUS_STALE_NAN_BITS),
753            schema: 0,
754            zero_threshold: 0.0,
755            reset_hint: ResetHint::Unspecified as i32,
756            timestamp,
757            start_timestamp,
758            count: Some(PromCount::CountInt(0)),
759            zero_count: Some(PromZeroCount::ZeroCountInt(0)),
760            ..Default::default()
761        }
762    } else {
763        if data_point.scale < MIN_EXPONENTIAL_HISTOGRAM_SCALE {
764            return Err(format!(
765                "scale {} is unsupported; minimum supported scale is {}",
766                data_point.scale, MIN_EXPONENTIAL_HISTOGRAM_SCALE
767            ));
768        }
769        if !data_point.zero_threshold.is_finite() || data_point.zero_threshold < 0.0 {
770            return Err(format!(
771                "zero_threshold {} must be finite and non-negative",
772                data_point.zero_threshold
773            ));
774        }
775
776        let downscale_shift = if data_point.scale > MAX_EXPONENTIAL_HISTOGRAM_SCALE {
777            let shift = data_point
778                .scale
779                .checked_sub(MAX_EXPONENTIAL_HISTOGRAM_SCALE)
780                .ok_or_else(|| "downscale shift overflows i32".to_string())?;
781            u32::try_from(shift).map_err(|_| "downscale shift exceeds u32".to_string())?
782        } else {
783            0
784        };
785        let (positive_spans, positive_deltas, positive_count) =
786            convert_bucket_range("positive", data_point.positive.as_ref(), downscale_shift)?;
787        let (negative_spans, negative_deltas, negative_count) =
788            convert_bucket_range("negative", data_point.negative.as_ref(), downscale_shift)?;
789        let bucket_count = data_point
790            .zero_count
791            .checked_add(positive_count)
792            .and_then(|count| count.checked_add(negative_count))
793            .ok_or_else(|| "bucket observation total overflows u64".to_string())?;
794        if bucket_count != data_point.count {
795            return Err(format!(
796                "buckets contain {bucket_count} observations, declared count is {}",
797                data_point.count
798            ));
799        }
800
801        i64::try_from(data_point.count)
802            .map_err(|_| format!("count {} overflows i64", data_point.count))?;
803        i64::try_from(data_point.zero_count)
804            .map_err(|_| format!("zero_count {} overflows i64", data_point.zero_count))?;
805
806        let sum = match data_point.sum {
807            Some(sum) if sum.is_nan() => f64::NAN,
808            Some(sum) => sum,
809            None => f64::NAN,
810        };
811        PromHistogram {
812            sum,
813            schema: data_point.scale.min(MAX_EXPONENTIAL_HISTOGRAM_SCALE),
814            zero_threshold: data_point.zero_threshold,
815            negative_spans,
816            negative_deltas,
817            positive_spans,
818            positive_deltas,
819            reset_hint: ResetHint::Unspecified as i32,
820            timestamp,
821            start_timestamp,
822            count: Some(PromCount::CountInt(data_point.count)),
823            zero_count: Some(PromZeroCount::ZeroCountInt(data_point.zero_count)),
824            ..Default::default()
825        }
826    };
827
828    encode_native_histogram(&histogram)
829        .map(|value| (value, timestamp_nanos))
830        .map_err(|error| format!("OTLP exponential histogram cannot be encoded: {error}"))
831}
832
833fn timestamp_millis(timestamp_nanos: u64, name: &str) -> std::result::Result<i64, String> {
834    i64::try_from(timestamp_nanos / 1_000_000)
835        .map_err(|_| format!("{name} {timestamp_nanos} milliseconds overflow i64"))
836}
837
838fn convert_bucket_range(
839    name: &str,
840    buckets: Option<&exponential_histogram_data_point::Buckets>,
841    downscale_shift: u32,
842) -> std::result::Result<(Vec<BucketSpan>, Vec<i64>, u64), String> {
843    let Some(buckets) = buckets else {
844        return Ok((Vec::new(), Vec::new(), 0));
845    };
846    if buckets.bucket_counts.is_empty() {
847        return Ok((Vec::new(), Vec::new(), 0));
848    }
849
850    let mut merged = Vec::<(i32, u64)>::with_capacity(buckets.bucket_counts.len());
851    let mut total = 0u64;
852    for (position, count) in buckets.bucket_counts.iter().copied().enumerate() {
853        let position =
854            i32::try_from(position).map_err(|_| format!("{name} bucket length exceeds i32"))?;
855        let source_index = buckets.offset.checked_add(position).ok_or_else(|| {
856            format!(
857                "{name} bucket index overflows i32 at offset {} and position {position}",
858                buckets.offset
859            )
860        })?;
861        let target_index = downscale_bucket_index(source_index, downscale_shift)?
862            .checked_add(1)
863            .ok_or_else(|| format!("{name} shifted bucket index overflows i32"))?;
864        if let Some((last_index, last_count)) = merged.last_mut() {
865            if *last_index == target_index {
866                *last_count = last_count
867                    .checked_add(count)
868                    .ok_or_else(|| format!("{name} merged bucket count overflows u64"))?;
869                total = total
870                    .checked_add(count)
871                    .ok_or_else(|| format!("{name} bucket count total overflows u64"))?;
872                continue;
873            }
874            let next_index = last_index
875                .checked_add(1)
876                .ok_or_else(|| format!("{name} target bucket index overflows i32"))?;
877            if target_index != next_index {
878                return Err(format!(
879                    "{name} bucket indexes are not contiguous after downscaling"
880                ));
881            }
882        }
883        total = total
884            .checked_add(count)
885            .ok_or_else(|| format!("{name} bucket count total overflows u64"))?;
886        merged.push((target_index, count));
887    }
888
889    let mut spans = Vec::<BucketSpan>::new();
890    let mut deltas = Vec::with_capacity(merged.len());
891    let mut previous_index = None::<i32>;
892    let mut previous = 0i64;
893    for (index, count) in merged {
894        if count == 0 {
895            continue;
896        }
897        match (spans.last_mut(), previous_index) {
898            (Some(span), Some(prev_index)) if prev_index.checked_add(1) == Some(index) => {
899                span.length = span
900                    .length
901                    .checked_add(1)
902                    .ok_or_else(|| format!("{name} bucket span length exceeds u32"))?;
903            }
904            (_, prev_index) => {
905                let offset = match prev_index {
906                    Some(prev_index) => index
907                        .checked_sub(prev_index)
908                        .and_then(|gap| gap.checked_sub(1))
909                        .ok_or_else(|| format!("{name} bucket span offset overflows i32"))?,
910                    None => index,
911                };
912                spans.push(BucketSpan { offset, length: 1 });
913            }
914        }
915        let count = i64::try_from(count)
916            .map_err(|_| format!("{name} bucket count {count} overflows i64"))?;
917        let delta = count
918            .checked_sub(previous)
919            .ok_or_else(|| format!("{name} bucket delta overflows i64"))?;
920        deltas.push(delta);
921        previous_index = Some(index);
922        // Bucket deltas continue across span gaps; omitted zero buckets do not reset them.
923        previous = count;
924    }
925
926    Ok((spans, deltas, total))
927}
928
929fn downscale_bucket_index(index: i32, downscale_shift: u32) -> std::result::Result<i32, String> {
930    if downscale_shift == 0 {
931        return Ok(index);
932    }
933    if downscale_shift >= i32::BITS {
934        return Ok(if index < 0 { -1 } else { 0 });
935    }
936
937    let divisor = 1i64
938        .checked_shl(downscale_shift)
939        .ok_or_else(|| format!("downscale shift {downscale_shift} is invalid"))?;
940    i32::try_from(i64::from(index).div_euclid(divisor))
941        .map_err(|_| format!("downscaled bucket index {index} overflows i32"))
942}
943
944#[derive(Debug, Clone, Copy, PartialEq, Eq)]
945enum AttributeType {
946    Resource,
947    Scope,
948    DataPoint,
949    Legacy,
950}
951
952fn write_attributes(
953    writer: &mut TableData,
954    row: &mut Vec<Value>,
955    attrs: Option<&Vec<KeyValue>>,
956    attribute_type: AttributeType,
957    metric_ctx: &OtlpMetricCtx,
958) -> Result<()> {
959    let Some(attrs) = attrs else {
960        return Ok(());
961    };
962
963    let mut tags = Vec::with_capacity(attrs.len());
964    for attr in attrs {
965        // TODO(sunng87): allow different type of values
966        let Some(value) = scalar_value_string(attr.value.as_ref()) else {
967            continue;
968        };
969        let key = match attribute_type {
970            AttributeType::Resource | AttributeType::DataPoint => {
971                translate_label_name(&attr.key, metric_ctx.metric_translation_strategy)
972            }
973            AttributeType::Scope => {
974                format!(
975                    "otel_scope_{}",
976                    translate_label_name(&attr.key, metric_ctx.metric_translation_strategy)
977                )
978            }
979            AttributeType::Legacy => legacy_normalize_otlp_name(&attr.key),
980        };
981        if key == OTLP_AGGREGATION_TEMPORALITY_LABEL {
982            return Err(error::InvalidOtlpMetricInputSnafu {
983                reason: format!(
984                    "OTLP attribute `{}` resolves to reserved label `{}`",
985                    attr.key, OTLP_AGGREGATION_TEMPORALITY_LABEL
986                ),
987            }
988            .build());
989        }
990        tags.push((key, value));
991    }
992    row_writer::write_tags(writer, tags.into_iter(), row)?;
993
994    Ok(())
995}
996
997fn write_timestamp(
998    table: &mut TableData,
999    row: &mut Vec<Value>,
1000    time_nano: i64,
1001    metric_ctx: &OtlpMetricCtx,
1002) -> Result<()> {
1003    // Keep the full nanosecond precision whenever the request is headed for
1004    // the metric engine: `Inserter::handle_metric_row_inserts` converts the
1005    // timestamps to the physical table's time index unit (which may be
1006    // micro/nanosecond). Only the non-metric prometheus-compatible path is
1007    // fixed to milliseconds, to keep auto-created mito tables on the
1008    // millisecond time index.
1009    if metric_ctx.is_legacy || metric_ctx.with_metric_engine {
1010        row_writer::write_ts_to_nanos(
1011            table,
1012            greptime_timestamp(),
1013            Some(time_nano),
1014            Precision::Nanosecond,
1015            row,
1016        )
1017    } else {
1018        row_writer::write_ts_to_millis(
1019            table,
1020            greptime_timestamp(),
1021            Some(time_nano / 1000000),
1022            Precision::Millisecond,
1023            row,
1024        )
1025    }
1026}
1027
1028fn write_data_point_value(
1029    table: &mut TableData,
1030    row: &mut Vec<Value>,
1031    field: &str,
1032    value: &Option<number_data_point::Value>,
1033) -> Result<()> {
1034    match value {
1035        Some(number_data_point::Value::AsInt(val)) => {
1036            // we coerce all values to f64
1037            row_writer::write_f64(table, field, *val as f64, row)?;
1038        }
1039        Some(number_data_point::Value::AsDouble(val)) => {
1040            row_writer::write_f64(table, field, *val, row)?;
1041        }
1042        _ => {}
1043    }
1044    Ok(())
1045}
1046
1047fn write_temporality_tag(
1048    table: &mut TableData,
1049    row: &mut Vec<Value>,
1050    is_delta: bool,
1051) -> Result<()> {
1052    if is_delta {
1053        row_writer::write_tag(
1054            table,
1055            OTLP_AGGREGATION_TEMPORALITY_LABEL,
1056            GREPTIME_TEMPORALITY_DELTA,
1057            row,
1058        )?;
1059    }
1060    Ok(())
1061}
1062
1063fn has_no_recorded_value(flags: u32) -> bool {
1064    flags & DataPointFlags::NoRecordedValueMask as u32 != 0
1065}
1066
1067fn write_tags_and_timestamp(
1068    table: &mut TableData,
1069    row: &mut Vec<Value>,
1070    resource_attrs: Option<&Vec<KeyValue>>,
1071    scope_attrs: Option<&Vec<KeyValue>>,
1072    data_point_attrs: Option<&Vec<KeyValue>>,
1073    timestamp_nanos: i64,
1074    metric_ctx: &OtlpMetricCtx,
1075) -> Result<()> {
1076    if metric_ctx.is_legacy {
1077        write_attributes(
1078            table,
1079            row,
1080            resource_attrs,
1081            AttributeType::Legacy,
1082            metric_ctx,
1083        )?;
1084        write_attributes(table, row, scope_attrs, AttributeType::Legacy, metric_ctx)?;
1085        write_attributes(
1086            table,
1087            row,
1088            data_point_attrs,
1089            AttributeType::Legacy,
1090            metric_ctx,
1091        )?;
1092    } else {
1093        // TODO(shuiyisong): check `__type__` and `__unit__` tags in prometheus
1094        write_attributes(
1095            table,
1096            row,
1097            resource_attrs,
1098            AttributeType::Resource,
1099            metric_ctx,
1100        )?;
1101        write_attributes(table, row, scope_attrs, AttributeType::Scope, metric_ctx)?;
1102        write_attributes(
1103            table,
1104            row,
1105            data_point_attrs,
1106            AttributeType::DataPoint,
1107            metric_ctx,
1108        )?;
1109    }
1110
1111    write_timestamp(table, row, timestamp_nanos, metric_ctx)?;
1112
1113    Ok(())
1114}
1115
1116/// encode this gauge metric
1117///
1118/// note that there can be multiple data points in the request, it's going to be
1119/// stored as multiple rows
1120fn encode_gauge(
1121    table_writer: &mut MultiTableData,
1122    name: &str,
1123    gauge: &Gauge,
1124    resource_attrs: Option<&Vec<KeyValue>>,
1125    scope_attrs: Option<&Vec<KeyValue>>,
1126    metric_ctx: &OtlpMetricCtx,
1127) -> Result<()> {
1128    let table = table_writer.get_or_default_table_data(
1129        name,
1130        APPROXIMATE_COLUMN_COUNT,
1131        gauge.data_points.len(),
1132    );
1133
1134    for data_point in &gauge.data_points {
1135        let mut row = table.alloc_one_row();
1136        write_tags_and_timestamp(
1137            table,
1138            &mut row,
1139            resource_attrs,
1140            scope_attrs,
1141            Some(data_point.attributes.as_ref()),
1142            data_point.time_unix_nano as i64,
1143            metric_ctx,
1144        )?;
1145
1146        write_data_point_value(table, &mut row, greptime_value(), &data_point.value)?;
1147        table.add_row(row);
1148    }
1149
1150    Ok(())
1151}
1152
1153/// Encodes sums, preserving delta points as interval values with a temporality tag.
1154/// Ingestion is stateless: values are never accumulated across timestamps.
1155fn encode_sum(
1156    table_writer: &mut MultiTableData,
1157    name: &str,
1158    sum: &Sum,
1159    resource_attrs: Option<&Vec<KeyValue>>,
1160    scope_attrs: Option<&Vec<KeyValue>>,
1161    metric_ctx: &OtlpMetricCtx,
1162) -> Result<()> {
1163    let is_delta = matches!(
1164        AggregationTemporality::try_from(sum.aggregation_temporality),
1165        Ok(AggregationTemporality::Delta)
1166    );
1167    let table = table_writer.get_or_default_table_data(
1168        name,
1169        APPROXIMATE_COLUMN_COUNT,
1170        sum.data_points.len(),
1171    );
1172
1173    for data_point in &sum.data_points {
1174        let mut row = table.alloc_one_row();
1175        write_tags_and_timestamp(
1176            table,
1177            &mut row,
1178            resource_attrs,
1179            scope_attrs,
1180            Some(data_point.attributes.as_ref()),
1181            data_point.time_unix_nano as i64,
1182            metric_ctx,
1183        )?;
1184        write_temporality_tag(table, &mut row, is_delta)?;
1185        if has_no_recorded_value(data_point.flags) {
1186            row_writer::write_f64(
1187                table,
1188                greptime_value(),
1189                f64::from_bits(PROMETHEUS_STALE_NAN_BITS),
1190                &mut row,
1191            )?;
1192        } else {
1193            write_data_point_value(table, &mut row, greptime_value(), &data_point.value)?;
1194        }
1195        table.add_row(row);
1196    }
1197
1198    Ok(())
1199}
1200
1201const HISTOGRAM_LE_COLUMN: &str = "le";
1202
1203/// Encode histogram data. This function returns 3 insert requests for 3 tables.
1204///
1205/// The implementation has been following Prometheus histogram table format:
1206///
1207/// - A `%metric%_bucket` table including `greptime_le` tag that stores bucket upper
1208///   limit, and `greptime_value` for bucket count
1209/// - A `%metric%_sum` table storing sum of samples
1210/// -  A `%metric%_count` table storing count of samples.
1211///
1212/// By its Prometheus compatibility, we hope to be able to use prometheus
1213/// quantile functions on this table.
1214/// Delta points retain their temporality tag and interval values. Bucket counts
1215/// are prefix-summed within each point, never accumulated across timestamps.
1216fn encode_histogram(
1217    table_writer: &mut MultiTableData,
1218    name: &str,
1219    hist: &Histogram,
1220    resource_attrs: Option<&Vec<KeyValue>>,
1221    scope_attrs: Option<&Vec<KeyValue>>,
1222    metric_ctx: &OtlpMetricCtx,
1223    outcome: &mut MetricsIngestOutcome,
1224) -> Result<bool> {
1225    let normalized_name = name;
1226
1227    let bucket_table_name = format!("{}{}", normalized_name, BUCKET_TABLE_SUFFIX);
1228    let sum_table_name = format!("{}{}", normalized_name, SUM_TABLE_SUFFIX);
1229    let count_table_name = format!("{}{}", normalized_name, COUNT_TABLE_SUFFIX);
1230
1231    let is_delta = matches!(
1232        AggregationTemporality::try_from(hist.aggregation_temporality),
1233        Ok(AggregationTemporality::Delta)
1234    );
1235    let stale_value = f64::from_bits(PROMETHEUS_STALE_NAN_BITS);
1236    let mut emitted = false;
1237    for (index, data_point) in hist.data_points.iter().enumerate() {
1238        if let Some(reason) = histogram_data_point_rejection(data_point, is_delta) {
1239            reject_data_points(outcome, 1, || {
1240                format!("metric `{name}` data point {index}: {reason}")
1241            })?;
1242            continue;
1243        }
1244
1245        let bucket_table =
1246            table_writer.get_or_default_table_data(&bucket_table_name, APPROXIMATE_COLUMN_COUNT, 0);
1247        let no_recorded_value = has_no_recorded_value(data_point.flags);
1248        if no_recorded_value {
1249            bucket_table.reserve_rows(data_point.explicit_bounds.len());
1250            bucket_table.reserve_rows(1);
1251        } else {
1252            bucket_table.reserve_rows(data_point.bucket_counts.len().max(1));
1253        }
1254        let bucket_values = if no_recorded_value {
1255            data_point
1256                .explicit_bounds
1257                .iter()
1258                .copied()
1259                .chain(std::iter::once(f64::INFINITY))
1260                .map(|bound| (bound, stale_value))
1261                .collect::<Vec<_>>()
1262        } else if data_point.bucket_counts.is_empty() && data_point.explicit_bounds.is_empty() {
1263            // OTLP count-only histograms map to one implicit Prometheus infinity bucket.
1264            vec![(f64::INFINITY, data_point.count as f64)]
1265        } else {
1266            let mut accumulated_count = 0u64;
1267            let mut values = Vec::with_capacity(data_point.bucket_counts.len());
1268            for (idx, count) in data_point.bucket_counts.iter().enumerate() {
1269                accumulated_count = accumulated_count.checked_add(*count).ok_or_else(|| {
1270                    error::InvalidParameterSnafu {
1271                        reason: format!(
1272                            "metric `{name}` data point {index}: bucket prefix overflows u64"
1273                        ),
1274                    }
1275                    .build()
1276                })?;
1277                let bound =
1278                    data_point.explicit_bounds.get(idx).copied().or_else(|| {
1279                        (idx == data_point.explicit_bounds.len()).then_some(f64::INFINITY)
1280                    });
1281                if let Some(bound) = bound {
1282                    values.push((bound, accumulated_count as f64));
1283                }
1284            }
1285            values
1286        };
1287        for (bound, value) in bucket_values {
1288            let mut bucket_row = bucket_table.alloc_one_row();
1289            write_tags_and_timestamp(
1290                bucket_table,
1291                &mut bucket_row,
1292                resource_attrs,
1293                scope_attrs,
1294                Some(data_point.attributes.as_ref()),
1295                data_point.time_unix_nano as i64,
1296                metric_ctx,
1297            )?;
1298            write_temporality_tag(bucket_table, &mut bucket_row, is_delta)?;
1299            row_writer::write_tag(bucket_table, HISTOGRAM_LE_COLUMN, bound, &mut bucket_row)?;
1300            row_writer::write_f64(bucket_table, greptime_value(), value, &mut bucket_row)?;
1301
1302            bucket_table.add_row(bucket_row);
1303        }
1304
1305        if let Some(sum) = data_point.sum {
1306            let sum_table = table_writer.get_or_default_table_data(
1307                &sum_table_name,
1308                APPROXIMATE_COLUMN_COUNT,
1309                hist.data_points.len(),
1310            );
1311            let mut sum_row = sum_table.alloc_one_row();
1312            write_tags_and_timestamp(
1313                sum_table,
1314                &mut sum_row,
1315                resource_attrs,
1316                scope_attrs,
1317                Some(data_point.attributes.as_ref()),
1318                data_point.time_unix_nano as i64,
1319                metric_ctx,
1320            )?;
1321            write_temporality_tag(sum_table, &mut sum_row, is_delta)?;
1322            row_writer::write_f64(
1323                sum_table,
1324                greptime_value(),
1325                if no_recorded_value { stale_value } else { sum },
1326                &mut sum_row,
1327            )?;
1328            sum_table.add_row(sum_row);
1329        }
1330
1331        let count_table = table_writer.get_or_default_table_data(
1332            &count_table_name,
1333            APPROXIMATE_COLUMN_COUNT,
1334            hist.data_points.len(),
1335        );
1336        let mut count_row = count_table.alloc_one_row();
1337        write_tags_and_timestamp(
1338            count_table,
1339            &mut count_row,
1340            resource_attrs,
1341            scope_attrs,
1342            Some(data_point.attributes.as_ref()),
1343            data_point.time_unix_nano as i64,
1344            metric_ctx,
1345        )?;
1346        write_temporality_tag(count_table, &mut count_row, is_delta)?;
1347        row_writer::write_f64(
1348            count_table,
1349            greptime_value(),
1350            if no_recorded_value {
1351                stale_value
1352            } else {
1353                data_point.count as f64
1354            },
1355            &mut count_row,
1356        )?;
1357        count_table.add_row(count_row);
1358        add_accepted_data_points(outcome, 1)?;
1359        emitted = true;
1360    }
1361
1362    Ok(emitted)
1363}
1364
1365pub(crate) fn histogram_data_point_rejection(
1366    data_point: &HistogramDataPoint,
1367    is_delta: bool,
1368) -> Option<String> {
1369    if has_no_recorded_value(data_point.flags) {
1370        return None;
1371    }
1372
1373    if is_delta {
1374        let valid_empty_layout =
1375            data_point.bucket_counts.is_empty() && data_point.explicit_bounds.is_empty();
1376        let expected_buckets = data_point.explicit_bounds.len().checked_add(1);
1377        if !valid_empty_layout && expected_buckets != Some(data_point.bucket_counts.len()) {
1378            return Some(format!(
1379                "bucket_counts length {} must equal explicit_bounds length {} plus one",
1380                data_point.bucket_counts.len(),
1381                data_point.explicit_bounds.len()
1382            ));
1383        }
1384        if data_point
1385            .explicit_bounds
1386            .iter()
1387            .any(|bound| !bound.is_finite())
1388            || data_point
1389                .explicit_bounds
1390                .windows(2)
1391                .any(|bounds| bounds[0] >= bounds[1])
1392        {
1393            return Some("explicit_bounds must be finite and strictly increasing".to_string());
1394        }
1395        if data_point.count == 0 && data_point.sum.is_some_and(|sum| sum != 0.0) {
1396            return Some("sum must be absent or zero when count is zero".to_string());
1397        }
1398    }
1399
1400    let bucket_total = data_point
1401        .bucket_counts
1402        .iter()
1403        .try_fold(0u64, |total, count| total.checked_add(*count));
1404    let Some(bucket_total) = bucket_total else {
1405        return Some("bucket prefix overflows u64".to_string());
1406    };
1407    if is_delta && !data_point.bucket_counts.is_empty() && bucket_total != data_point.count {
1408        return Some(format!(
1409            "buckets contain {bucket_total} observations, declared count is {}",
1410            data_point.count
1411        ));
1412    }
1413
1414    None
1415}
1416
1417fn encode_summary(
1418    table_writer: &mut MultiTableData,
1419    name: &str,
1420    summary: &Summary,
1421    resource_attrs: Option<&Vec<KeyValue>>,
1422    scope_attrs: Option<&Vec<KeyValue>>,
1423    metric_ctx: &OtlpMetricCtx,
1424) -> Result<()> {
1425    if metric_ctx.is_legacy {
1426        let table = table_writer.get_or_default_table_data(
1427            name,
1428            APPROXIMATE_COLUMN_COUNT,
1429            summary.data_points.len(),
1430        );
1431
1432        for data_point in &summary.data_points {
1433            let mut row = table.alloc_one_row();
1434            write_tags_and_timestamp(
1435                table,
1436                &mut row,
1437                resource_attrs,
1438                scope_attrs,
1439                Some(data_point.attributes.as_ref()),
1440                data_point.time_unix_nano as i64,
1441                metric_ctx,
1442            )?;
1443
1444            for quantile in &data_point.quantile_values {
1445                row_writer::write_f64(
1446                    table,
1447                    format!("greptime_p{:02}", quantile.quantile * 100f64),
1448                    quantile.value,
1449                    &mut row,
1450                )?;
1451            }
1452
1453            row_writer::write_f64(table, GREPTIME_COUNT, data_point.count as f64, &mut row)?;
1454            table.add_row(row);
1455        }
1456    } else {
1457        // 1. quantile table
1458        // 2. count table
1459        // 3. sum table
1460
1461        let metric_name = name;
1462        let count_name = format!("{}{}", metric_name, COUNT_TABLE_SUFFIX);
1463        let sum_name = format!("{}{}", metric_name, SUM_TABLE_SUFFIX);
1464
1465        for data_point in &summary.data_points {
1466            {
1467                let quantile_table = table_writer.get_or_default_table_data(
1468                    metric_name,
1469                    APPROXIMATE_COLUMN_COUNT,
1470                    summary.data_points.len(),
1471                );
1472
1473                for quantile in &data_point.quantile_values {
1474                    let mut row = quantile_table.alloc_one_row();
1475                    write_tags_and_timestamp(
1476                        quantile_table,
1477                        &mut row,
1478                        resource_attrs,
1479                        scope_attrs,
1480                        Some(data_point.attributes.as_ref()),
1481                        data_point.time_unix_nano as i64,
1482                        metric_ctx,
1483                    )?;
1484                    row_writer::write_tag(quantile_table, "quantile", quantile.quantile, &mut row)?;
1485                    row_writer::write_f64(
1486                        quantile_table,
1487                        greptime_value(),
1488                        quantile.value,
1489                        &mut row,
1490                    )?;
1491                    quantile_table.add_row(row);
1492                }
1493            }
1494            {
1495                let count_table = table_writer.get_or_default_table_data(
1496                    &count_name,
1497                    APPROXIMATE_COLUMN_COUNT,
1498                    summary.data_points.len(),
1499                );
1500                let mut row = count_table.alloc_one_row();
1501                write_tags_and_timestamp(
1502                    count_table,
1503                    &mut row,
1504                    resource_attrs,
1505                    scope_attrs,
1506                    Some(data_point.attributes.as_ref()),
1507                    data_point.time_unix_nano as i64,
1508                    metric_ctx,
1509                )?;
1510
1511                row_writer::write_f64(
1512                    count_table,
1513                    greptime_value(),
1514                    data_point.count as f64,
1515                    &mut row,
1516                )?;
1517
1518                count_table.add_row(row);
1519            }
1520            {
1521                let sum_table = table_writer.get_or_default_table_data(
1522                    &sum_name,
1523                    APPROXIMATE_COLUMN_COUNT,
1524                    summary.data_points.len(),
1525                );
1526
1527                let mut row = sum_table.alloc_one_row();
1528                write_tags_and_timestamp(
1529                    sum_table,
1530                    &mut row,
1531                    resource_attrs,
1532                    scope_attrs,
1533                    Some(data_point.attributes.as_ref()),
1534                    data_point.time_unix_nano as i64,
1535                    metric_ctx,
1536                )?;
1537
1538                row_writer::write_f64(sum_table, greptime_value(), data_point.sum, &mut row)?;
1539
1540                sum_table.add_row(row);
1541            }
1542        }
1543    }
1544
1545    Ok(())
1546}
1547
1548#[cfg(test)]
1549mod tests {
1550    use api::v1::ColumnDataType;
1551    use common_query::prelude::set_default_prefix;
1552    use otel_arrow_rust::proto::opentelemetry::common::v1::AnyValue;
1553    use otel_arrow_rust::proto::opentelemetry::common::v1::any_value::Value as Val;
1554    use otel_arrow_rust::proto::opentelemetry::metrics::v1::number_data_point::Value;
1555    use otel_arrow_rust::proto::opentelemetry::metrics::v1::summary_data_point::ValueAtQuantile;
1556    use otel_arrow_rust::proto::opentelemetry::metrics::v1::{
1557        AggregationTemporality, HistogramDataPoint, NumberDataPoint, SummaryDataPoint,
1558    };
1559    use otel_arrow_rust::proto::opentelemetry::resource::v1::Resource;
1560
1561    use super::*;
1562
1563    mod delta;
1564
1565    fn keyvalue(key: &str, value: &str) -> KeyValue {
1566        KeyValue {
1567            key: key.into(),
1568            value: Some(AnyValue {
1569                value: Some(Val::StringValue(value.into())),
1570            }),
1571        }
1572    }
1573
1574    fn descriptor_ctx() -> OtlpMetricCtx {
1575        OtlpMetricCtx {
1576            resource_info: true,
1577            ..Default::default()
1578        }
1579    }
1580
1581    fn attr_value(attrs: &[KeyValue], key: &str) -> Option<String> {
1582        attrs
1583            .iter()
1584            .find(|kv| kv.key == key)
1585            .and_then(|kv| scalar_value_string(kv.value.as_ref()))
1586    }
1587
1588    fn gauge_request(
1589        resource_attrs: Vec<KeyValue>,
1590        metric_name: &str,
1591    ) -> ExportMetricsServiceRequest {
1592        use otel_arrow_rust::proto::opentelemetry::resource::v1::Resource;
1593        ExportMetricsServiceRequest {
1594            resource_metrics: vec![ResourceMetrics {
1595                resource: Some(Resource {
1596                    attributes: resource_attrs,
1597                    ..Default::default()
1598                }),
1599                scope_metrics: vec![ScopeMetrics {
1600                    metrics: vec![Metric {
1601                        name: metric_name.to_string(),
1602                        data: Some(metric::Data::Gauge(Gauge {
1603                            data_points: vec![NumberDataPoint {
1604                                time_unix_nano: 1_000_000,
1605                                value: Some(Value::AsInt(1)),
1606                                ..Default::default()
1607                            }],
1608                        })),
1609                        ..Default::default()
1610                    }],
1611                    ..Default::default()
1612                }],
1613                ..Default::default()
1614            }],
1615        }
1616    }
1617
1618    fn column_names(request: &RowInsertRequests, table: &str) -> Vec<String> {
1619        request
1620            .inserts
1621            .iter()
1622            .find(|r| r.table_name == table)
1623            .unwrap_or_else(|| panic!("missing table {table}"))
1624            .rows
1625            .as_ref()
1626            .unwrap()
1627            .schema
1628            .iter()
1629            .map(|c| c.column_name.clone())
1630            .collect()
1631    }
1632
1633    #[test]
1634    fn test_conversion_synthesizes_resource_descriptor() {
1635        set_default_prefix(None).unwrap();
1636        let request = gauge_request(
1637            vec![keyvalue("service.name", "api"), keyvalue("host.id", "h-1")],
1638            "my_gauge",
1639        );
1640        let conversion = to_grpc_insert_requests(request, &mut descriptor_ctx()).unwrap();
1641
1642        // descriptor keeps raw OTel keys while the metric table's labels went
1643        // through the (default underscore-escaping) translation strategy
1644        let resource_info = conversion.resource_info.expect("descriptor synthesized");
1645        let descriptor_cols = column_names(&resource_info, OTEL_RESOURCE_INFO_TABLE_NAME);
1646        assert!(descriptor_cols.contains(&"host.id".to_string()));
1647        assert!(descriptor_cols.contains(&"service.name".to_string()));
1648        assert!(descriptor_cols.contains(&"job".to_string()));
1649        let metric_cols = column_names(&conversion.requests, "my_gauge");
1650        assert!(metric_cols.contains(&"service_name".to_string()));
1651        assert!(!metric_cols.contains(&"service.name".to_string()));
1652        // host.id is not in the promote list: only the descriptor keeps it
1653        assert!(!metric_cols.contains(&"host_id".to_string()));
1654
1655        let decoded = decode(&conversion.semantic_index);
1656        let t = &decoded[OTEL_RESOURCE_INFO_TABLE_NAME];
1657        assert_eq!(
1658            t.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
1659            Some("info")
1660        );
1661        assert_eq!(
1662            t.get(SEMANTIC_METRIC_METADATA_QUALITY).map(String::as_str),
1663            Some("declared")
1664        );
1665    }
1666
1667    #[test]
1668    fn test_conversion_skips_descriptor_for_legacy_mode() {
1669        set_default_prefix(None).unwrap();
1670        let request = gauge_request(
1671            vec![keyvalue("service.name", "api"), keyvalue("host.id", "h-1")],
1672            "my_gauge",
1673        );
1674        let mut ctx = OtlpMetricCtx {
1675            is_legacy: true,
1676            ..descriptor_ctx()
1677        };
1678        let conversion = to_grpc_insert_requests(request, &mut ctx).unwrap();
1679        assert!(conversion.resource_info.is_none());
1680
1681        // legacy tables predate job/instance and the promote filter; adding
1682        // either would alter the schema of tables already in use
1683        let cols = column_names(&conversion.requests, "my_gauge");
1684        assert!(!cols.contains(&"job".to_string()));
1685        assert!(cols.contains(&"host_id".to_string()));
1686    }
1687
1688    #[test]
1689    fn test_conversion_skips_descriptor_on_metric_name_collision() {
1690        set_default_prefix(None).unwrap();
1691        let request = gauge_request(
1692            vec![keyvalue("service.name", "api")],
1693            OTEL_RESOURCE_INFO_TABLE_NAME,
1694        );
1695        let conversion = to_grpc_insert_requests(request, &mut descriptor_ctx()).unwrap();
1696        assert!(conversion.resource_info.is_none());
1697        // the metric itself still goes through the main path
1698        assert!(
1699            conversion
1700                .requests
1701                .inserts
1702                .iter()
1703                .any(|r| r.table_name == OTEL_RESOURCE_INFO_TABLE_NAME)
1704        );
1705    }
1706
1707    #[test]
1708    fn test_job_composition_follows_service_namespace() {
1709        let mut attrs = vec![
1710            keyvalue("service.name", "api"),
1711            keyvalue("service.namespace", "shop"),
1712        ];
1713        process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
1714        assert_eq!(attr_value(&attrs, "job").as_deref(), Some("shop/api"));
1715
1716        let mut attrs = vec![keyvalue("service.name", "api")];
1717        process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
1718        assert_eq!(attr_value(&attrs, "job").as_deref(), Some("api"));
1719
1720        let mut attrs = vec![
1721            keyvalue("service.name", "api"),
1722            keyvalue("service.namespace", ""),
1723        ];
1724        process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
1725        assert_eq!(attr_value(&attrs, "job").as_deref(), Some("api"));
1726
1727        // no service.name: no job is fabricated
1728        let mut attrs = vec![
1729            keyvalue("service.namespace", "shop"),
1730            keyvalue("service.instance.id", "inst-1"),
1731        ];
1732        process_resource_attrs(&mut attrs, &OtlpMetricCtx::default());
1733        assert_eq!(attr_value(&attrs, "job"), None);
1734        assert_eq!(attr_value(&attrs, "instance").as_deref(), Some("inst-1"));
1735    }
1736
1737    #[test]
1738    fn test_encode_gauge() {
1739        let mut tables = MultiTableData::default();
1740
1741        let data_points = vec![
1742            NumberDataPoint {
1743                attributes: vec![keyvalue("host", "testsevrer")],
1744                time_unix_nano: 100,
1745                value: Some(Value::AsInt(100)),
1746                ..Default::default()
1747            },
1748            NumberDataPoint {
1749                attributes: vec![keyvalue("host", "testserver")],
1750                time_unix_nano: 105,
1751                value: Some(Value::AsInt(105)),
1752                ..Default::default()
1753            },
1754        ];
1755        let gauge = Gauge { data_points };
1756        encode_gauge(
1757            &mut tables,
1758            "datamon",
1759            &gauge,
1760            Some(&vec![]),
1761            Some(&vec![keyvalue("scope", "otel")]),
1762            &OtlpMetricCtx::default(),
1763        )
1764        .unwrap();
1765
1766        let table = tables.get_or_default_table_data("datamon", 0, 0);
1767        assert_eq!(table.num_rows(), 2);
1768        assert_eq!(table.num_columns(), 4);
1769        assert_eq!(
1770            table
1771                .columns()
1772                .iter()
1773                .map(|c| &c.column_name)
1774                .collect::<Vec<&String>>(),
1775            vec![
1776                "otel_scope_scope",
1777                "host",
1778                greptime_timestamp(),
1779                greptime_value()
1780            ]
1781        );
1782    }
1783
1784    #[test]
1785    fn test_encode_sum() {
1786        let mut tables = MultiTableData::default();
1787
1788        let data_points = vec![
1789            NumberDataPoint {
1790                attributes: vec![keyvalue("host", "testserver")],
1791                time_unix_nano: 100,
1792                value: Some(Value::AsInt(100)),
1793                ..Default::default()
1794            },
1795            NumberDataPoint {
1796                attributes: vec![keyvalue("host", "testserver")],
1797                time_unix_nano: 105,
1798                value: Some(Value::AsInt(0)),
1799                ..Default::default()
1800            },
1801        ];
1802        let sum = Sum {
1803            data_points,
1804            ..Default::default()
1805        };
1806        encode_sum(
1807            &mut tables,
1808            "datamon",
1809            &sum,
1810            Some(&vec![]),
1811            Some(&vec![keyvalue("scope", "otel")]),
1812            &OtlpMetricCtx::default(),
1813        )
1814        .unwrap();
1815
1816        let table = tables.get_or_default_table_data("datamon", 0, 0);
1817        assert_eq!(table.num_rows(), 2);
1818        assert_eq!(table.num_columns(), 4);
1819        assert_eq!(
1820            table
1821                .columns()
1822                .iter()
1823                .map(|c| &c.column_name)
1824                .collect::<Vec<&String>>(),
1825            vec![
1826                "otel_scope_scope",
1827                "host",
1828                greptime_timestamp(),
1829                greptime_value()
1830            ]
1831        );
1832    }
1833
1834    #[test]
1835    fn test_encode_summary() {
1836        let mut tables = MultiTableData::default();
1837
1838        let data_points = vec![SummaryDataPoint {
1839            attributes: vec![keyvalue("host", "testserver")],
1840            time_unix_nano: 100,
1841            count: 25,
1842            sum: 5400.0,
1843            quantile_values: vec![
1844                ValueAtQuantile {
1845                    quantile: 0.90,
1846                    value: 1000.0,
1847                },
1848                ValueAtQuantile {
1849                    quantile: 0.95,
1850                    value: 3030.0,
1851                },
1852            ],
1853            ..Default::default()
1854        }];
1855        let summary = Summary { data_points };
1856        encode_summary(
1857            &mut tables,
1858            "datamon",
1859            &summary,
1860            Some(&vec![]),
1861            Some(&vec![keyvalue("scope", "otel")]),
1862            &OtlpMetricCtx::default(),
1863        )
1864        .unwrap();
1865
1866        let table = tables.get_or_default_table_data("datamon", 0, 0);
1867        assert_eq!(table.num_rows(), 2);
1868        assert_eq!(table.num_columns(), 5);
1869        assert_eq!(
1870            table
1871                .columns()
1872                .iter()
1873                .map(|c| &c.column_name)
1874                .collect::<Vec<&String>>(),
1875            vec![
1876                "otel_scope_scope",
1877                "host",
1878                greptime_timestamp(),
1879                "quantile",
1880                greptime_value()
1881            ]
1882        );
1883
1884        let table = tables.get_or_default_table_data("datamon_count", 0, 0);
1885        assert_eq!(table.num_rows(), 1);
1886        assert_eq!(table.num_columns(), 4);
1887        assert_eq!(
1888            table
1889                .columns()
1890                .iter()
1891                .map(|c| &c.column_name)
1892                .collect::<Vec<&String>>(),
1893            vec![
1894                "otel_scope_scope",
1895                "host",
1896                greptime_timestamp(),
1897                greptime_value()
1898            ]
1899        );
1900
1901        let table = tables.get_or_default_table_data("datamon_sum", 0, 0);
1902        assert_eq!(table.num_rows(), 1);
1903        assert_eq!(table.num_columns(), 4);
1904        assert_eq!(
1905            table
1906                .columns()
1907                .iter()
1908                .map(|c| &c.column_name)
1909                .collect::<Vec<&String>>(),
1910            vec![
1911                "otel_scope_scope",
1912                "host",
1913                greptime_timestamp(),
1914                greptime_value()
1915            ]
1916        );
1917    }
1918
1919    #[test]
1920    fn test_encode_legacy_summary_keeps_legacy_column_names() {
1921        set_default_prefix(Some("custom")).unwrap();
1922        let mut tables = MultiTableData::default();
1923        let summary = Summary {
1924            data_points: vec![SummaryDataPoint {
1925                attributes: vec![keyvalue("host", "testserver")],
1926                time_unix_nano: 100,
1927                count: 25,
1928                quantile_values: vec![ValueAtQuantile {
1929                    quantile: 0.90,
1930                    value: 1000.0,
1931                }],
1932                ..Default::default()
1933            }],
1934        };
1935
1936        encode_summary(
1937            &mut tables,
1938            "datamon",
1939            &summary,
1940            None,
1941            None,
1942            &OtlpMetricCtx {
1943                is_legacy: true,
1944                ..Default::default()
1945            },
1946        )
1947        .unwrap();
1948
1949        let table = tables.get_or_default_table_data("datamon", 0, 0);
1950        assert_eq!(
1951            table
1952                .columns()
1953                .iter()
1954                .map(|column| column.column_name.as_str())
1955                .collect::<Vec<_>>(),
1956            vec!["host", "custom_timestamp", "greptime_p90", GREPTIME_COUNT,]
1957        );
1958    }
1959
1960    #[test]
1961    fn test_encode_histogram() {
1962        let mut tables = MultiTableData::default();
1963        let mut outcome = MetricsIngestOutcome::default();
1964
1965        let data_points = vec![HistogramDataPoint {
1966            attributes: vec![keyvalue("host", "testserver")],
1967            time_unix_nano: 100,
1968            start_time_unix_nano: 23,
1969            count: 25,
1970            sum: Some(100.),
1971            max: Some(200.),
1972            min: Some(0.03),
1973            bucket_counts: vec![2, 4, 6, 9, 4],
1974            explicit_bounds: vec![0.1, 1., 10., 100.],
1975            ..Default::default()
1976        }];
1977
1978        let histogram = Histogram {
1979            data_points,
1980            aggregation_temporality: AggregationTemporality::Delta.into(),
1981        };
1982        encode_histogram(
1983            &mut tables,
1984            "histo",
1985            &histogram,
1986            Some(&vec![]),
1987            Some(&vec![keyvalue("scope", "otel")]),
1988            &OtlpMetricCtx::default(),
1989            &mut outcome,
1990        )
1991        .unwrap();
1992
1993        assert_eq!(3, tables.num_tables());
1994        assert_eq!(1, outcome.accepted_data_points);
1995
1996        // bucket table
1997        let bucket_table = tables.get_or_default_table_data("histo_bucket", 0, 0);
1998        assert_eq!(bucket_table.num_rows(), 5);
1999        assert_eq!(bucket_table.num_columns(), 6);
2000        assert_eq!(
2001            bucket_table
2002                .columns()
2003                .iter()
2004                .map(|c| &c.column_name)
2005                .collect::<Vec<&String>>(),
2006            vec![
2007                "otel_scope_scope",
2008                "host",
2009                greptime_timestamp(),
2010                OTLP_AGGREGATION_TEMPORALITY_LABEL,
2011                "le",
2012                greptime_value(),
2013            ]
2014        );
2015
2016        let sum_table = tables.get_or_default_table_data("histo_sum", 0, 0);
2017        assert_eq!(sum_table.num_rows(), 1);
2018        assert_eq!(sum_table.num_columns(), 5);
2019        assert_eq!(
2020            sum_table
2021                .columns()
2022                .iter()
2023                .map(|c| &c.column_name)
2024                .collect::<Vec<&String>>(),
2025            vec![
2026                "otel_scope_scope",
2027                "host",
2028                greptime_timestamp(),
2029                OTLP_AGGREGATION_TEMPORALITY_LABEL,
2030                greptime_value()
2031            ]
2032        );
2033
2034        let count_table = tables.get_or_default_table_data("histo_count", 0, 0);
2035        assert_eq!(count_table.num_rows(), 1);
2036        assert_eq!(count_table.num_columns(), 5);
2037        assert_eq!(
2038            count_table
2039                .columns()
2040                .iter()
2041                .map(|c| &c.column_name)
2042                .collect::<Vec<&String>>(),
2043            vec![
2044                "otel_scope_scope",
2045                "host",
2046                greptime_timestamp(),
2047                OTLP_AGGREGATION_TEMPORALITY_LABEL,
2048                greptime_value()
2049            ]
2050        );
2051    }
2052
2053    use std::collections::BTreeMap;
2054
2055    use table::requests::validate_semantic_option;
2056
2057    fn decode(index: &SemanticIndex) -> BTreeMap<String, BTreeMap<String, String>> {
2058        let nested: BTreeMap<String, BTreeMap<String, BTreeMap<String, String>>> =
2059            serde_json::from_str(&index.encode("public").expect("non-empty index")).unwrap();
2060        nested.into_values().next().unwrap()
2061    }
2062
2063    fn record(metric: &Metric, metric_type: MetricType, name: &str) -> SemanticIndex {
2064        let ctx = OtlpMetricCtx {
2065            metric_type,
2066            ..Default::default()
2067        };
2068        let mut index = SemanticIndex::default();
2069        record_metric_semantics(&mut index, metric, name, &ctx);
2070        index
2071    }
2072
2073    #[test]
2074    fn test_metric_type_constants_validate() {
2075        for value in [
2076            METRIC_TYPE_COUNTER,
2077            METRIC_TYPE_UPDOWN_COUNTER,
2078            METRIC_TYPE_GAUGE,
2079            METRIC_TYPE_HISTOGRAM,
2080            METRIC_TYPE_SUMMARY,
2081        ] {
2082            assert!(
2083                validate_semantic_option(SEMANTIC_METRIC_TYPE, value),
2084                "metric.type value `{value}` must be in the vocabulary domain"
2085            );
2086        }
2087        for value in ["delta", "cumulative"] {
2088            assert!(validate_semantic_option(SEMANTIC_METRIC_TEMPORALITY, value));
2089        }
2090    }
2091
2092    #[test]
2093    fn test_record_monotonic_sum() {
2094        let metric = Metric {
2095            name: "claude_code.cost.usage".to_string(),
2096            unit: "USD".to_string(),
2097            data: Some(metric::Data::Sum(Sum {
2098                aggregation_temporality: AggregationTemporality::Delta as i32,
2099                is_monotonic: true,
2100                ..Default::default()
2101            })),
2102            ..Default::default()
2103        };
2104        let index = record(
2105            &metric,
2106            MetricType::MonotonicSum,
2107            "claude_code_cost_usage_USD_total",
2108        );
2109        let decoded = decode(&index);
2110        let t = &decoded["claude_code_cost_usage_USD_total"];
2111
2112        assert_eq!(
2113            t.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
2114            Some("counter")
2115        );
2116        assert_eq!(
2117            t.get(SEMANTIC_METRIC_TEMPORALITY).map(String::as_str),
2118            Some("delta")
2119        );
2120        assert_eq!(t.get(SEMANTIC_METRIC_UNIT).map(String::as_str), Some("USD"));
2121        assert_eq!(
2122            t.get(SEMANTIC_METRIC_ORIGINAL_NAME).map(String::as_str),
2123            Some("claude_code.cost.usage")
2124        );
2125        assert_eq!(
2126            t.get(SEMANTIC_METRIC_METADATA_QUALITY).map(String::as_str),
2127            Some("declared")
2128        );
2129    }
2130
2131    #[test]
2132    fn test_record_non_monotonic_sum() {
2133        let metric = Metric {
2134            name: "queue_size".to_string(),
2135            data: Some(metric::Data::Sum(Sum {
2136                aggregation_temporality: AggregationTemporality::Cumulative as i32,
2137                is_monotonic: false,
2138                ..Default::default()
2139            })),
2140            ..Default::default()
2141        };
2142        let index = record(&metric, MetricType::NonMonotonicSum, "queue_size");
2143        let decoded = decode(&index);
2144        let t = &decoded["queue_size"];
2145        assert_eq!(
2146            t.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
2147            Some("updown_counter")
2148        );
2149        assert_eq!(
2150            t.get(SEMANTIC_METRIC_TEMPORALITY).map(String::as_str),
2151            Some("cumulative")
2152        );
2153        // Name unchanged by translation -> no original_name.
2154        assert_eq!(t.get(SEMANTIC_METRIC_ORIGINAL_NAME), None);
2155    }
2156
2157    #[test]
2158    fn test_record_gauge_has_no_temporality() {
2159        let metric = Metric {
2160            name: "temperature".to_string(),
2161            data: Some(metric::Data::Gauge(Gauge::default())),
2162            ..Default::default()
2163        };
2164        let index = record(&metric, MetricType::Gauge, "temperature");
2165        let decoded = decode(&index);
2166        let t = &decoded["temperature"];
2167        assert_eq!(
2168            t.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
2169            Some("gauge")
2170        );
2171        assert_eq!(t.get(SEMANTIC_METRIC_TEMPORALITY), None);
2172    }
2173
2174    #[test]
2175    fn test_record_histogram_fans_out_with_distinct_types() {
2176        let metric = Metric {
2177            name: "request.duration".to_string(),
2178            unit: "s".to_string(),
2179            data: Some(metric::Data::Histogram(Histogram {
2180                aggregation_temporality: AggregationTemporality::Cumulative as i32,
2181                ..Default::default()
2182            })),
2183            ..Default::default()
2184        };
2185        let index = record(&metric, MetricType::Histogram, "request_duration");
2186        let decoded = decode(&index);
2187
2188        let bucket = &decoded["request_duration_bucket"];
2189        assert_eq!(
2190            bucket.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
2191            Some("histogram")
2192        );
2193        assert_eq!(
2194            bucket.get(SEMANTIC_METRIC_UNIT).map(String::as_str),
2195            Some("s")
2196        );
2197
2198        for companion in ["request_duration_sum", "request_duration_count"] {
2199            let t = &decoded[companion];
2200            assert_eq!(
2201                t.get(SEMANTIC_METRIC_TYPE).map(String::as_str),
2202                Some("counter")
2203            );
2204            assert_eq!(
2205                t.get(SEMANTIC_METRIC_TEMPORALITY).map(String::as_str),
2206                Some("cumulative")
2207            );
2208        }
2209    }
2210
2211    #[test]
2212    fn test_record_summary_fans_out() {
2213        let metric = Metric {
2214            name: "rpc.latency".to_string(),
2215            data: Some(metric::Data::Summary(Summary::default())),
2216            ..Default::default()
2217        };
2218        let index = record(&metric, MetricType::Summary, "rpc_latency");
2219        let decoded = decode(&index);
2220
2221        assert_eq!(
2222            decoded["rpc_latency"]
2223                .get(SEMANTIC_METRIC_TYPE)
2224                .map(String::as_str),
2225            Some("summary")
2226        );
2227        // Summary has no temporality.
2228        assert_eq!(
2229            decoded["rpc_latency"].get(SEMANTIC_METRIC_TEMPORALITY),
2230            None
2231        );
2232        for companion in ["rpc_latency_count", "rpc_latency_sum"] {
2233            assert_eq!(
2234                decoded[companion]
2235                    .get(SEMANTIC_METRIC_TYPE)
2236                    .map(String::as_str),
2237                Some("counter")
2238            );
2239        }
2240    }
2241
2242    fn exponential_buckets(
2243        offset: i32,
2244        bucket_counts: Vec<u64>,
2245    ) -> exponential_histogram_data_point::Buckets {
2246        exponential_histogram_data_point::Buckets {
2247            offset,
2248            bucket_counts,
2249        }
2250    }
2251
2252    fn exponential_point() -> ExponentialHistogramDataPoint {
2253        ExponentialHistogramDataPoint {
2254            start_time_unix_nano: 1_000_000,
2255            time_unix_nano: 2_000_000,
2256            count: 28,
2257            sum: None,
2258            scale: 9,
2259            zero_count: 7,
2260            positive: Some(exponential_buckets(-3, vec![1, 2, 3, 4, 5, 6])),
2261            zero_threshold: 0.0,
2262            ..Default::default()
2263        }
2264    }
2265
2266    fn native_field(value: &ValueData, name: &str) -> Option<ValueData> {
2267        let ValueData::StructValue(value) = value else {
2268            panic!("expected native histogram Struct value");
2269        };
2270        let index = common_query::native_histogram::NATIVE_HISTOGRAM_FIELD_NAMES
2271            .iter()
2272            .position(|field| *field == name)
2273            .unwrap();
2274        value.items[index].value_data.clone()
2275    }
2276
2277    fn i32_list(value: Option<ValueData>) -> Vec<i32> {
2278        let Some(ValueData::ListValue(value)) = value else {
2279            panic!("expected i32 list");
2280        };
2281        value
2282            .items
2283            .into_iter()
2284            .map(|item| match item.value_data {
2285                Some(ValueData::I32Value(value)) => value,
2286                _ => panic!("expected i32 value"),
2287            })
2288            .collect()
2289    }
2290
2291    fn i64_list(value: Option<ValueData>) -> Vec<i64> {
2292        let Some(ValueData::ListValue(value)) = value else {
2293            panic!("expected i64 list");
2294        };
2295        value
2296            .items
2297            .into_iter()
2298            .map(|item| match item.value_data {
2299                Some(ValueData::I64Value(value)) => value,
2300                _ => panic!("expected i64 value"),
2301            })
2302            .collect()
2303    }
2304
2305    #[test]
2306    fn test_downscale_bucket_index_uses_signed_floor_division() {
2307        assert_eq!(downscale_bucket_index(-3, 1).unwrap(), -2);
2308        assert_eq!(downscale_bucket_index(-2, 1).unwrap(), -1);
2309        assert_eq!(downscale_bucket_index(-1, 1).unwrap(), -1);
2310        assert_eq!(downscale_bucket_index(0, 1).unwrap(), 0);
2311        assert_eq!(downscale_bucket_index(1, 1).unwrap(), 0);
2312        assert_eq!(downscale_bucket_index(2, 1).unwrap(), 1);
2313        assert_eq!(downscale_bucket_index(i32::MIN, 32).unwrap(), -1);
2314        assert_eq!(downscale_bucket_index(i32::MAX, 32).unwrap(), 0);
2315    }
2316
2317    #[test]
2318    fn test_convert_bucket_range_downscales_before_prometheus_shift() {
2319        let buckets = exponential_buckets(-3, vec![1, 2, 3, 4, 5, 6]);
2320        let (spans, deltas, total) = convert_bucket_range("positive", Some(&buckets), 1).unwrap();
2321
2322        assert_eq!(
2323            spans,
2324            vec![BucketSpan {
2325                offset: -1,
2326                length: 4
2327            }]
2328        );
2329        assert_eq!(deltas, vec![1, 4, 4, -3]);
2330        assert_eq!(total, 21);
2331    }
2332
2333    #[test]
2334    fn test_convert_bucket_range_compacts_zero_runs() {
2335        let large = (1u64 << 53) + 1;
2336        for (offset, counts, shift, expected_spans, expected_deltas, expected_total) in [
2337            (
2338                -3,
2339                vec![0, 3, 0, 0, 7, 9, 0],
2340                0,
2341                vec![(-1, 1), (2, 2)],
2342                vec![3, 4, 2],
2343                19,
2344            ),
2345            (
2346                -4,
2347                vec![0, 0, 2, 3, 0, 0, 0, 0, 7, 0],
2348                1,
2349                vec![(0, 1), (2, 1)],
2350                vec![5, 2],
2351                12,
2352            ),
2353            (0, vec![0, 0, 0], 0, vec![], vec![], 0),
2354            (0, vec![], 0, vec![], vec![], 0),
2355            (i32::MIN, vec![0, 1], 0, vec![(i32::MIN + 2, 1)], vec![1], 1),
2356            (i32::MAX - 1, vec![1], 0, vec![(i32::MAX, 1)], vec![1], 1),
2357            (
2358                0,
2359                vec![large, 0, large + 2],
2360                0,
2361                vec![(1, 1), (1, 1)],
2362                vec![large as i64, 2],
2363                large * 2 + 2,
2364            ),
2365        ] {
2366            let buckets = exponential_buckets(offset, counts);
2367            let (spans, deltas, total) =
2368                convert_bucket_range("positive", Some(&buckets), shift).unwrap();
2369            assert_eq!(
2370                spans
2371                    .iter()
2372                    .map(|span| (span.offset, span.length))
2373                    .collect::<Vec<_>>(),
2374                expected_spans,
2375                "{buckets:?}, shift={shift}"
2376            );
2377            assert_eq!(deltas, expected_deltas, "{buckets:?}, shift={shift}");
2378            assert_eq!(total, expected_total);
2379        }
2380
2381        // Compaction must not hide invalid source indexes, even for empty buckets.
2382        let buckets = exponential_buckets(i32::MAX, vec![0]);
2383        assert!(
2384            convert_bucket_range("positive", Some(&buckets), 0)
2385                .unwrap_err()
2386                .contains("shifted bucket index overflows")
2387        );
2388    }
2389
2390    #[test]
2391    fn test_exponential_histogram_compaction_preserves_integer_counts() {
2392        use common_query::native_histogram::{
2393            COUNT_I64_FIELD, NEGATIVE_BUCKETS_I64_FIELD, NEGATIVE_SPAN_LENGTHS_FIELD,
2394            NEGATIVE_SPAN_OFFSETS_FIELD, POSITIVE_BUCKETS_I64_FIELD, POSITIVE_SPAN_LENGTHS_FIELD,
2395            POSITIVE_SPAN_OFFSETS_FIELD,
2396        };
2397
2398        let large = (1u64 << 53) + 1;
2399        let buckets = exponential_buckets(-2, vec![0, large, 0, 3, 0]);
2400        let point = ExponentialHistogramDataPoint {
2401            count: 2 * (large + 3),
2402            positive: Some(buckets.clone()),
2403            negative: Some(buckets),
2404            ..Default::default()
2405        };
2406        let (value, _) = exponential_histogram_value(&point).unwrap();
2407        assert_eq!(
2408            native_field(&value, COUNT_I64_FIELD),
2409            Some(ValueData::I64Value(point.count as i64))
2410        );
2411        for (offsets, lengths, counts) in [
2412            (
2413                POSITIVE_SPAN_OFFSETS_FIELD,
2414                POSITIVE_SPAN_LENGTHS_FIELD,
2415                POSITIVE_BUCKETS_I64_FIELD,
2416            ),
2417            (
2418                NEGATIVE_SPAN_OFFSETS_FIELD,
2419                NEGATIVE_SPAN_LENGTHS_FIELD,
2420                NEGATIVE_BUCKETS_I64_FIELD,
2421            ),
2422        ] {
2423            assert_eq!(i32_list(native_field(&value, offsets)), vec![0, 1]);
2424            assert_eq!(i32_list(native_field(&value, lengths)), vec![1, 1]);
2425            assert_eq!(
2426                i64_list(native_field(&value, counts)),
2427                vec![large as i64, 3]
2428            );
2429        }
2430    }
2431
2432    #[test]
2433    fn test_exponential_histogram_value_uses_integer_family() {
2434        use common_query::native_histogram::{
2435            COUNT_F64_FIELD, COUNT_I64_FIELD, POSITIVE_BUCKETS_I64_FIELD,
2436            POSITIVE_SPAN_LENGTHS_FIELD, POSITIVE_SPAN_OFFSETS_FIELD, SCHEMA_FIELD, SUM_FIELD,
2437            ZERO_COUNT_I64_FIELD,
2438        };
2439
2440        let (value, timestamp_nanos) = exponential_histogram_value(&exponential_point()).unwrap();
2441
2442        assert_eq!(timestamp_nanos, 2_000_000);
2443        assert_eq!(
2444            native_field(&value, SCHEMA_FIELD),
2445            Some(ValueData::I32Value(8))
2446        );
2447        assert_eq!(
2448            native_field(&value, COUNT_I64_FIELD),
2449            Some(ValueData::I64Value(28))
2450        );
2451        assert_eq!(
2452            native_field(&value, ZERO_COUNT_I64_FIELD),
2453            Some(ValueData::I64Value(7))
2454        );
2455        assert_eq!(native_field(&value, COUNT_F64_FIELD), None);
2456        assert_eq!(
2457            i32_list(native_field(&value, POSITIVE_SPAN_OFFSETS_FIELD)),
2458            vec![-1]
2459        );
2460        assert_eq!(
2461            i32_list(native_field(&value, POSITIVE_SPAN_LENGTHS_FIELD)),
2462            vec![4]
2463        );
2464        assert_eq!(
2465            i64_list(native_field(&value, POSITIVE_BUCKETS_I64_FIELD)),
2466            vec![1, 5, 9, 6]
2467        );
2468        let Some(ValueData::F64Value(sum)) = native_field(&value, SUM_FIELD) else {
2469            panic!("expected histogram sum");
2470        };
2471        assert_eq!(sum.to_bits(), f64::NAN.to_bits());
2472        assert_ne!(sum.to_bits(), PROMETHEUS_STALE_NAN_BITS);
2473    }
2474
2475    #[test]
2476    fn test_no_recorded_value_ignores_other_value_fields() {
2477        use common_query::native_histogram::{
2478            COUNT_I64_FIELD, POSITIVE_BUCKETS_I64_FIELD, SCHEMA_FIELD, SUM_FIELD,
2479            ZERO_COUNT_I64_FIELD, ZERO_THRESHOLD_FIELD,
2480        };
2481
2482        let point = ExponentialHistogramDataPoint {
2483            start_time_unix_nano: 1_000_000,
2484            time_unix_nano: 2_000_000,
2485            count: u64::MAX,
2486            sum: Some(1.0),
2487            scale: i32::MIN,
2488            zero_count: u64::MAX,
2489            positive: Some(exponential_buckets(i32::MAX, vec![u64::MAX])),
2490            flags: DataPointFlags::NoRecordedValueMask as u32,
2491            zero_threshold: f64::NAN,
2492            ..Default::default()
2493        };
2494        let (value, _) = exponential_histogram_value(&point).unwrap();
2495
2496        assert_eq!(
2497            native_field(&value, SCHEMA_FIELD),
2498            Some(ValueData::I32Value(0))
2499        );
2500        assert_eq!(
2501            native_field(&value, ZERO_THRESHOLD_FIELD),
2502            Some(ValueData::F64Value(0.0))
2503        );
2504        assert_eq!(
2505            native_field(&value, COUNT_I64_FIELD),
2506            Some(ValueData::I64Value(0))
2507        );
2508        assert_eq!(
2509            native_field(&value, ZERO_COUNT_I64_FIELD),
2510            Some(ValueData::I64Value(0))
2511        );
2512        assert!(i64_list(native_field(&value, POSITIVE_BUCKETS_I64_FIELD)).is_empty());
2513        let Some(ValueData::F64Value(sum)) = native_field(&value, SUM_FIELD) else {
2514            panic!("expected histogram sum");
2515        };
2516        assert_eq!(sum.to_bits(), PROMETHEUS_STALE_NAN_BITS);
2517    }
2518
2519    #[test]
2520    fn test_non_flag_nan_sum_is_normalized() {
2521        use common_query::native_histogram::SUM_FIELD;
2522
2523        let point = ExponentialHistogramDataPoint {
2524            sum: Some(f64::from_bits(PROMETHEUS_STALE_NAN_BITS)),
2525            ..Default::default()
2526        };
2527        let (value, _) = exponential_histogram_value(&point).unwrap();
2528        let Some(ValueData::F64Value(sum)) = native_field(&value, SUM_FIELD) else {
2529            panic!("expected histogram sum");
2530        };
2531        assert_eq!(sum.to_bits(), f64::NAN.to_bits());
2532        assert_ne!(sum.to_bits(), PROMETHEUS_STALE_NAN_BITS);
2533    }
2534
2535    #[test]
2536    fn test_exponential_histogram_rejects_invalid_values() {
2537        let mut cases = Vec::new();
2538
2539        let mut point = exponential_point();
2540        point.scale = -5;
2541        cases.push((point, "scale -5 is unsupported"));
2542
2543        let mut point = exponential_point();
2544        point.zero_threshold = f64::INFINITY;
2545        cases.push((point, "must be finite and non-negative"));
2546
2547        let mut point = exponential_point();
2548        point.start_time_unix_nano = point.time_unix_nano + 1;
2549        cases.push((point, "start_time_unix_nano"));
2550
2551        let mut point = exponential_point();
2552        point.count = 27;
2553        cases.push((point, "declared count is 27"));
2554
2555        let point = ExponentialHistogramDataPoint {
2556            count: u64::MAX,
2557            zero_count: u64::MAX,
2558            ..Default::default()
2559        };
2560        cases.push((point, "count 18446744073709551615 overflows i64"));
2561
2562        let point = ExponentialHistogramDataPoint {
2563            count: 1,
2564            scale: 8,
2565            positive: Some(exponential_buckets(i32::MAX, vec![1])),
2566            ..Default::default()
2567        };
2568        cases.push((point, "shifted bucket index overflows i32"));
2569
2570        for (point, expected) in cases {
2571            let error = exponential_histogram_value(&point).unwrap_err();
2572            assert!(
2573                error.contains(expected),
2574                "expected {expected:?}, got {error}"
2575            );
2576        }
2577
2578        let buckets = exponential_buckets(0, vec![u64::MAX, 1]);
2579        let error = convert_bucket_range("positive", Some(&buckets), 1).unwrap_err();
2580        assert!(
2581            error.contains("merged bucket count overflows u64"),
2582            "{error}"
2583        );
2584
2585        let buckets = exponential_buckets(i32::MAX, vec![1, 1]);
2586        let error = convert_bucket_range("positive", Some(&buckets), 1).unwrap_err();
2587        assert!(error.contains("bucket index overflows i32"), "{error}");
2588    }
2589
2590    fn metrics_request(metrics: Vec<Metric>) -> ExportMetricsServiceRequest {
2591        ExportMetricsServiceRequest {
2592            resource_metrics: vec![ResourceMetrics {
2593                scope_metrics: vec![ScopeMetrics {
2594                    metrics,
2595                    ..Default::default()
2596                }],
2597                ..Default::default()
2598            }],
2599        }
2600    }
2601
2602    fn exponential_metric(
2603        name: impl Into<String>,
2604        data_points: Vec<ExponentialHistogramDataPoint>,
2605        temporality: AggregationTemporality,
2606    ) -> Metric {
2607        Metric {
2608            name: name.into(),
2609            data: Some(metric::Data::ExponentialHistogram(ExponentialHistogram {
2610                data_points,
2611                aggregation_temporality: temporality as i32,
2612            })),
2613            ..Default::default()
2614        }
2615    }
2616
2617    fn histogram_metric(name: impl Into<String>) -> Metric {
2618        Metric {
2619            name: name.into(),
2620            data: Some(metric::Data::Histogram(Histogram {
2621                data_points: vec![HistogramDataPoint {
2622                    start_time_unix_nano: 1_000_000,
2623                    time_unix_nano: 2_000_000,
2624                    count: 1,
2625                    sum: Some(1.0),
2626                    bucket_counts: vec![1],
2627                    ..Default::default()
2628                }],
2629                aggregation_temporality: AggregationTemporality::Cumulative as i32,
2630            })),
2631            ..Default::default()
2632        }
2633    }
2634
2635    #[test]
2636    fn test_metric_engine_path_keeps_nanosecond_precision() {
2637        let time_unix_nano = 1_704_067_200_123_456_789u64;
2638        let request = metrics_request(vec![Metric {
2639            name: "my_gauge".to_string(),
2640            data: Some(metric::Data::Gauge(Gauge {
2641                data_points: vec![NumberDataPoint {
2642                    time_unix_nano,
2643                    value: Some(Value::AsDouble(1.0)),
2644                    ..Default::default()
2645                }],
2646            })),
2647            ..Default::default()
2648        }]);
2649
2650        // The metric engine path keeps the full nanosecond precision here;
2651        // `Inserter::handle_metric_row_inserts` converts the timestamps to
2652        // the physical table's time index unit afterwards.
2653        let mut metric_ctx = OtlpMetricCtx {
2654            with_metric_engine: true,
2655            ..Default::default()
2656        };
2657        let MetricsConversion { requests, .. } =
2658            to_grpc_insert_requests(request, &mut metric_ctx).unwrap();
2659
2660        let rows = requests.inserts[0].rows.as_ref().unwrap();
2661        let ts_index = rows
2662            .schema
2663            .iter()
2664            .position(|column| column.column_name == greptime_timestamp())
2665            .unwrap();
2666        assert_eq!(
2667            rows.schema[ts_index].datatype,
2668            ColumnDataType::TimestampNanosecond as i32
2669        );
2670        assert!(matches!(
2671            rows.rows[0].values[ts_index].value_data,
2672            Some(ValueData::TimestampNanosecondValue(
2673                1_704_067_200_123_456_789
2674            ))
2675        ));
2676
2677        // The non-metric prometheus-compatible path stays millisecond so
2678        // auto-created mito tables keep the millisecond time index.
2679        let mut compat_ctx = OtlpMetricCtx::default();
2680        let request = metrics_request(vec![Metric {
2681            name: "my_gauge".to_string(),
2682            data: Some(metric::Data::Gauge(Gauge {
2683                data_points: vec![NumberDataPoint {
2684                    time_unix_nano,
2685                    value: Some(Value::AsDouble(1.0)),
2686                    ..Default::default()
2687                }],
2688            })),
2689            ..Default::default()
2690        }]);
2691        let MetricsConversion { requests, .. } =
2692            to_grpc_insert_requests(request, &mut compat_ctx).unwrap();
2693
2694        let rows = requests.inserts[0].rows.as_ref().unwrap();
2695        let ts_index = rows
2696            .schema
2697            .iter()
2698            .position(|column| column.column_name == greptime_timestamp())
2699            .unwrap();
2700        assert_eq!(
2701            rows.schema[ts_index].datatype,
2702            ColumnDataType::TimestampMillisecond as i32
2703        );
2704        assert!(matches!(
2705            rows.rows[0].values[ts_index].value_data,
2706            Some(ValueData::TimestampMillisecondValue(1_704_067_200_123))
2707        ));
2708    }
2709
2710    #[test]
2711    fn test_exponential_histogram_rejection_metrics() {
2712        // Other conversion tests update the same process-global counters.
2713        const ISOLATED_ENV: &str = "GREPTIME_TEST_OTLP_REJECTION_METRICS_ISOLATED";
2714        if std::env::var_os(ISOLATED_ENV).is_none() {
2715            let output = std::process::Command::new(std::env::current_exe().unwrap())
2716                .args([
2717                    "--exact",
2718                    "otlp::metrics::tests::test_exponential_histogram_rejection_metrics",
2719                ])
2720                .env(ISOLATED_ENV, "1")
2721                .output()
2722                .unwrap();
2723            let stdout = String::from_utf8_lossy(&output.stdout);
2724            assert!(
2725                output.status.success() && stdout.contains("1 passed"),
2726                "isolated rejection metrics test failed\nstdout:\n{stdout}\nstderr:\n{}",
2727                String::from_utf8_lossy(&output.stderr)
2728            );
2729            return;
2730        }
2731
2732        let mut invalid = exponential_point();
2733        invalid.scale = -5;
2734        for (temporality, points, reason, accepted, rejected) in [
2735            (
2736                AggregationTemporality::Delta,
2737                vec![exponential_point(); 2],
2738                "delta_temporality",
2739                0,
2740                2,
2741            ),
2742            (
2743                AggregationTemporality::Unspecified,
2744                vec![exponential_point(); 2],
2745                "unspecified_temporality",
2746                0,
2747                2,
2748            ),
2749            (
2750                AggregationTemporality::Cumulative,
2751                vec![exponential_point(), invalid],
2752                "invalid_data_point",
2753                1,
2754                1,
2755            ),
2756            (
2757                AggregationTemporality::Cumulative,
2758                vec![exponential_point()],
2759                "invalid_data_point",
2760                1,
2761                0,
2762            ),
2763        ] {
2764            let counter =
2765                OTLP_EXPONENTIAL_HISTOGRAM_REJECTED_DATA_POINTS.with_label_values(&[reason]);
2766            let before = counter.get();
2767            let mut request =
2768                metrics_request(vec![exponential_metric("latency", points, temporality)]);
2769            request.resource_metrics[0].resource = Some(Resource {
2770                attributes: vec![keyvalue("service.name", "api")],
2771                ..Default::default()
2772            });
2773            let mut ctx = descriptor_ctx();
2774            let conversion = to_grpc_insert_requests(request, &mut ctx).unwrap();
2775            assert_eq!(conversion.outcome.accepted_data_points, accepted);
2776            assert_eq!(conversion.outcome.rejected_data_points, rejected);
2777            assert_eq!(counter.get() - before, rejected as u64, "{reason}");
2778            assert_eq!(conversion.resource_info.is_some(), accepted > 0);
2779
2780            let before = counter.get();
2781            let empty = metrics_request(vec![exponential_metric("empty", vec![], temporality)]);
2782            assert_eq!(
2783                to_grpc_insert_requests(empty, &mut ctx)
2784                    .unwrap()
2785                    .outcome
2786                    .rejected_data_points,
2787                0
2788            );
2789            assert_eq!(counter.get(), before, "empty metric: {reason}");
2790        }
2791    }
2792
2793    #[test]
2794    fn test_exponential_histogram_default_context_and_partial_outcome() {
2795        let request = metrics_request(vec![
2796            Metric {
2797                name: "temperature".to_string(),
2798                data: Some(metric::Data::Gauge(Gauge {
2799                    data_points: vec![NumberDataPoint::default()],
2800                })),
2801                ..Default::default()
2802            },
2803            exponential_metric(
2804                "latency",
2805                vec![exponential_point()],
2806                AggregationTemporality::Cumulative,
2807            ),
2808            exponential_metric(
2809                "delta_latency",
2810                vec![exponential_point()],
2811                AggregationTemporality::Delta,
2812            ),
2813        ]);
2814        let MetricsConversion {
2815            requests,
2816            semantic_index,
2817            outcome,
2818            ..
2819        } = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
2820
2821        assert_eq!(outcome.accepted_data_points, 2);
2822        assert_eq!(outcome.rejected_data_points, 1);
2823        assert!(
2824            outcome
2825                .error_message
2826                .as_deref()
2827                .unwrap()
2828                .contains("only cumulative temporality is supported")
2829        );
2830        assert_eq!(requests.inserts.len(), 2);
2831        let semantics = decode(&semantic_index);
2832        assert!(semantics.contains_key("temperature"));
2833        assert!(semantics.contains_key("latency"));
2834        assert!(!semantics.contains_key("delta_latency"));
2835
2836        let empty = metrics_request(vec![exponential_metric(
2837            "empty",
2838            vec![],
2839            AggregationTemporality::Cumulative,
2840        )]);
2841        let outcome = to_grpc_insert_requests(empty, &mut OtlpMetricCtx::default())
2842            .unwrap()
2843            .outcome;
2844        assert_eq!(outcome.rejected_data_points, 0);
2845        assert_eq!(outcome.error_message, None);
2846    }
2847
2848    #[test]
2849    fn test_exponential_histogram_cannot_share_table_with_scalar_metric() {
2850        let request = metrics_request(vec![
2851            Metric {
2852                name: "latency".to_string(),
2853                data: Some(metric::Data::Gauge(Gauge {
2854                    data_points: vec![NumberDataPoint {
2855                        value: Some(number_data_point::Value::AsDouble(1.0)),
2856                        ..Default::default()
2857                    }],
2858                })),
2859                ..Default::default()
2860            },
2861            exponential_metric(
2862                "latency",
2863                vec![exponential_point()],
2864                AggregationTemporality::Cumulative,
2865            ),
2866        ]);
2867        let mut ctx = OtlpMetricCtx::default();
2868
2869        let error = to_grpc_insert_requests(request, &mut ctx).unwrap_err();
2870        assert!(
2871            error
2872                .to_string()
2873                .contains("cannot mix native histogram and float sample fields")
2874        );
2875    }
2876
2877    #[test]
2878    fn test_histogram_cannot_replace_exponential_histogram_table() {
2879        let request = metrics_request(vec![
2880            exponential_metric(
2881                "latency_bucket",
2882                vec![exponential_point()],
2883                AggregationTemporality::Cumulative,
2884            ),
2885            histogram_metric("latency"),
2886        ]);
2887        let mut ctx = OtlpMetricCtx::default();
2888
2889        let error = to_grpc_insert_requests(request, &mut ctx).unwrap_err();
2890        assert!(
2891            error
2892                .to_string()
2893                .contains("cannot mix native histogram and float sample fields")
2894        );
2895    }
2896
2897    #[test]
2898    fn test_histograms_with_same_name_across_resources_are_merged() {
2899        let metric = histogram_metric("latency");
2900        let request = ExportMetricsServiceRequest {
2901            resource_metrics: ["service-a", "service-b"]
2902                .into_iter()
2903                .map(|service| ResourceMetrics {
2904                    resource: Some(Resource {
2905                        attributes: vec![keyvalue("service.name", service)],
2906                        ..Default::default()
2907                    }),
2908                    scope_metrics: vec![ScopeMetrics {
2909                        metrics: vec![metric.clone()],
2910                        ..Default::default()
2911                    }],
2912                    ..Default::default()
2913                })
2914                .collect(),
2915        };
2916
2917        let MetricsConversion {
2918            requests,
2919            rows,
2920            outcome,
2921            ..
2922        } = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default()).unwrap();
2923
2924        assert_eq!(outcome.accepted_data_points, 2);
2925        assert_eq!(rows, 6);
2926        assert_eq!(requests.inserts.len(), 3);
2927        for request in requests.inserts {
2928            assert_eq!(
2929                request.rows.unwrap().rows.len(),
2930                2,
2931                "{}",
2932                request.table_name
2933            );
2934        }
2935    }
2936
2937    #[test]
2938    fn test_exponential_histogram_rejects_temporality_before_stale_point() {
2939        let stale = ExponentialHistogramDataPoint {
2940            flags: DataPointFlags::NoRecordedValueMask as u32,
2941            ..Default::default()
2942        };
2943        for temporality in [
2944            AggregationTemporality::Delta,
2945            AggregationTemporality::Unspecified,
2946        ] {
2947            let request = metrics_request(vec![exponential_metric(
2948                "latency",
2949                vec![stale.clone()],
2950                temporality,
2951            )]);
2952            let mut ctx = OtlpMetricCtx::default();
2953            let MetricsConversion {
2954                requests,
2955                rows,
2956                semantic_index,
2957                outcome,
2958                ..
2959            } = to_grpc_insert_requests(request, &mut ctx).unwrap();
2960
2961            assert_eq!(outcome.accepted_data_points, 0);
2962            assert_eq!(outcome.rejected_data_points, 1);
2963            assert_eq!(rows, 0);
2964            assert!(requests.inserts.is_empty());
2965            assert!(semantic_index.is_empty());
2966        }
2967    }
2968
2969    #[test]
2970    fn test_exponential_histogram_legacy_and_new_modes_share_struct() {
2971        use common_query::prelude::greptime_native_histogram;
2972
2973        let mut point = exponential_point();
2974        point.sum = Some(42.0);
2975        let request = metrics_request(vec![exponential_metric(
2976            "request.duration",
2977            vec![point],
2978            AggregationTemporality::Cumulative,
2979        )]);
2980        let mut new_ctx = OtlpMetricCtx::default();
2981        let new_requests = to_grpc_insert_requests(request.clone(), &mut new_ctx)
2982            .unwrap()
2983            .requests;
2984        let mut legacy_ctx = OtlpMetricCtx {
2985            is_legacy: true,
2986            ..Default::default()
2987        };
2988        let legacy_requests = to_grpc_insert_requests(request, &mut legacy_ctx)
2989            .unwrap()
2990            .requests;
2991
2992        let new_insert = &new_requests.inserts[0];
2993        let legacy_insert = &legacy_requests.inserts[0];
2994        assert_eq!(new_insert.table_name, "request_duration");
2995        assert_eq!(legacy_insert.table_name, "request_duration");
2996        let new_rows = new_insert.rows.as_ref().unwrap();
2997        let legacy_rows = legacy_insert.rows.as_ref().unwrap();
2998        let field = greptime_native_histogram();
2999        let new_histogram = new_rows.rows[0].values[new_rows
3000            .schema
3001            .iter()
3002            .position(|column| column.column_name == field)
3003            .unwrap()]
3004        .clone();
3005        let legacy_histogram = legacy_rows.rows[0].values[legacy_rows
3006            .schema
3007            .iter()
3008            .position(|column| column.column_name == field)
3009            .unwrap()]
3010        .clone();
3011        assert_eq!(new_histogram, legacy_histogram);
3012        assert!(matches!(
3013            new_rows.rows[0].values[new_rows
3014                .schema
3015                .iter()
3016                .position(|column| column.column_name == greptime_timestamp())
3017                .unwrap()]
3018            .value_data,
3019            Some(ValueData::TimestampMillisecondValue(2))
3020        ));
3021        assert!(matches!(
3022            legacy_rows.rows[0].values[legacy_rows
3023                .schema
3024                .iter()
3025                .position(|column| column.column_name == greptime_timestamp())
3026                .unwrap()]
3027            .value_data,
3028            Some(ValueData::TimestampNanosecondValue(2_000_000))
3029        ));
3030    }
3031
3032    #[test]
3033    fn test_rejection_message_is_bounded() {
3034        let request = metrics_request(vec![exponential_metric(
3035            "x".repeat(1_000),
3036            vec![ExponentialHistogramDataPoint::default()],
3037            AggregationTemporality::Delta,
3038        )]);
3039        let outcome = to_grpc_insert_requests(request, &mut OtlpMetricCtx::default())
3040            .unwrap()
3041            .outcome;
3042
3043        assert_eq!(outcome.rejected_data_points, 1);
3044        assert!(outcome.error_message.unwrap().len() <= MAX_REJECTION_MESSAGE_BYTES);
3045    }
3046
3047    #[test]
3048    fn test_rejection_message_reason_is_lazy_after_cap() {
3049        let mut outcome = MetricsIngestOutcome {
3050            error_message: Some("x".repeat(MAX_REJECTION_MESSAGE_BYTES)),
3051            ..Default::default()
3052        };
3053        let mut reason_built = false;
3054
3055        reject_data_points(&mut outcome, 1, || {
3056            reason_built = true;
3057            "unused".to_string()
3058        })
3059        .unwrap();
3060
3061        assert_eq!(outcome.rejected_data_points, 1);
3062        assert!(!reason_built);
3063    }
3064}