1use 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
63const 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
73const 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#[derive(Debug)]
110pub struct MetricsConversion {
111 pub requests: RowInsertRequests,
112 pub rows: usize,
114 pub semantic_index: SemanticIndex,
117 pub resource_info: Option<RowInsertRequests>,
120 pub outcome: MetricsIngestOutcome,
121}
122
123pub 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 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
249fn 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
283fn 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
299fn 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 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
352fn 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
362pub(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 let ServiceIdentity { job, instance } = service_identity(attrs);
407
408 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 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 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
718pub(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 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 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 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 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 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
1116fn 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
1153fn 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
1203fn 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}