1use std::cmp::Ordering;
17use std::collections::{BTreeMap, HashMap};
18use std::hash::BuildHasher;
19use std::ops::Range;
20
21use arrow::array::{Array, ArrayRef, AsArray, StructArray};
22use arrow::buffer::BooleanBuffer;
23use arrow::compute::kernels::cmp::distinct;
24use arrow::datatypes::{Float64Type, TimestampMillisecondType};
25use arrow_schema::DataType;
26use axum::Json;
27use axum::http::HeaderValue;
28use axum::response::{IntoResponse, Response};
29use common_error::ext::ErrorExt;
30use common_error::status_code::StatusCode;
31use common_query::native_histogram::{
32 NativeHistogram, is_native_histogram_value_type, read_histogram,
33};
34use common_query::prometheus::{format_prometheus_float, is_prometheus_stale_nan};
35use common_query::promql_annotations::{
36 PromqlAnnotationCollector, get_promql_annotation_collector,
37};
38use common_query::{Output, OutputData};
39use common_recordbatch::{RecordBatch, RecordBatches, SendableRecordBatchStream};
40use datatypes::arrow_array::string_array_value_at_index;
41use datatypes::prelude::ConcreteDataType;
42use datatypes::schema::Schema;
43use futures::TryStreamExt;
44use indexmap::map::RawEntryApiV1;
45use indexmap::map::raw_entry_v1::RawEntryMut;
46use indexmap::{Equivalent, IndexMap};
47use itertools::Either;
48use promql_parser::label::METRIC_NAME;
49use promql_parser::parser::value::ValueType;
50use ryu::Buffer;
51use serde::{Deserialize, Serialize, Serializer};
52use serde_json::Value;
53use snafu::{OptionExt, ResultExt};
54
55use crate::error::{
56 ArrowSnafu, CollectRecordbatchSnafu, DataFusionSnafu, Result, UnexpectedResultSnafu,
57 status_code_to_http_status,
58};
59use crate::http::header::{GREPTIME_DB_HEADER_METRICS, collect_plan_metrics};
60use crate::http::prometheus::{
61 PromData, PromNativeHistogram, PromQueryResult, PromSeriesMatrix, PromSeriesVector,
62 PrometheusResponse,
63};
64
65#[derive(Default)]
66struct PromSeriesSamples {
67 values: Vec<(f64, PromSampleValue)>,
68 histograms: Vec<(f64, PromNativeHistogram)>,
69}
70
71#[derive(Debug, Clone, Deserialize, PartialEq)]
78#[serde(untagged)]
79pub enum PromSampleValue {
80 #[serde(skip_deserializing)]
81 Number(f64),
82 Text(String),
83}
84
85impl Serialize for PromSampleValue {
86 fn serialize<S: Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
87 match self {
88 Self::Number(value) if value.is_finite() => {
89 serializer.serialize_str(Buffer::new().format_finite(*value))
90 }
91 Self::Number(value) => serializer.collect_str(value),
92 Self::Text(value) => serializer.serialize_str(value),
93 }
94 }
95}
96
97impl PromSampleValue {
98 fn into_string(self) -> String {
99 match self {
100 Self::Number(value) => format_prometheus_sample_value(value),
101 Self::Text(value) => value,
102 }
103 }
104}
105
106fn prometheus_native_histogram(histogram: &NativeHistogram) -> Result<PromNativeHistogram> {
107 Ok(PromNativeHistogram {
108 count: format_prometheus_float(histogram.count),
109 sum: format_prometheus_float(histogram.sum),
110 buckets: histogram
111 .to_prometheus_buckets()
112 .context(UnexpectedResultSnafu {
113 reason: "native histogram cannot be converted to Prometheus buckets",
114 })?,
115 })
116}
117
118fn format_prometheus_sample_value(value: f64) -> String {
124 if value.is_finite() {
125 Buffer::new().format_finite(value).to_string()
126 } else {
127 value.to_string()
128 }
129}
130
131#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
132pub struct PrometheusJsonResponse {
133 pub status: String,
134 #[serde(skip_serializing_if = "PrometheusResponse::is_none")]
135 #[serde(default)]
136 pub data: PrometheusResponse,
137 #[serde(skip_serializing_if = "Option::is_none")]
138 pub error: Option<String>,
139 #[serde(skip_serializing_if = "Option::is_none")]
140 #[serde(rename = "errorType")]
141 pub error_type: Option<String>,
142 #[serde(skip_serializing_if = "Option::is_none")]
143 pub warnings: Option<Vec<String>>,
144 #[serde(skip_serializing_if = "Option::is_none")]
145 pub infos: Option<Vec<String>>,
146
147 #[serde(skip)]
148 pub status_code: Option<StatusCode>,
149 #[serde(skip)]
151 #[serde(default)]
152 pub resp_metrics: HashMap<String, Value>,
153}
154
155impl IntoResponse for PrometheusJsonResponse {
156 fn into_response(self) -> Response {
157 let metrics = if self.resp_metrics.is_empty() {
158 None
159 } else {
160 serde_json::to_string(&self.resp_metrics).ok()
161 };
162
163 let http_code = self.status_code.map(|c| status_code_to_http_status(&c));
164
165 let mut resp = Json(self).into_response();
166
167 if let Some(http_code) = http_code {
168 *resp.status_mut() = http_code;
169 }
170
171 if let Some(m) = metrics.and_then(|m| HeaderValue::from_str(&m).ok()) {
172 resp.headers_mut().insert(&GREPTIME_DB_HEADER_METRICS, m);
173 }
174
175 resp
176 }
177}
178
179impl PrometheusJsonResponse {
180 pub fn error<S1>(error_type: StatusCode, reason: S1) -> Self
181 where
182 S1: Into<String>,
183 {
184 PrometheusJsonResponse {
185 status: "error".to_string(),
186 data: PrometheusResponse::None,
187 error: Some(reason.into()),
188 error_type: Some(error_type.to_string()),
189 warnings: None,
190 infos: None,
191 resp_metrics: Default::default(),
192 status_code: Some(error_type),
193 }
194 }
195
196 pub fn success(data: PrometheusResponse) -> Self {
197 PrometheusJsonResponse {
198 status: "success".to_string(),
199 data,
200 error: None,
201 error_type: None,
202 warnings: None,
203 infos: None,
204 resp_metrics: Default::default(),
205 status_code: None,
206 }
207 }
208
209 fn append_promql_annotations(&mut self, collector: &PromqlAnnotationCollector) {
211 let mut warnings = self.warnings.take().unwrap_or_default();
212 let mut infos = self.infos.take().unwrap_or_default();
213 collector.append_to(&mut warnings, &mut infos);
214 self.warnings = (!warnings.is_empty()).then_some(warnings);
215 self.infos = (!infos.is_empty()).then_some(infos);
216 }
217
218 pub(crate) fn append_query_response(&mut self, mut other: Self) {
220 self.data.append(other.data);
221 merge_annotations(&mut self.warnings, other.warnings.take());
222 merge_annotations(&mut self.infos, other.infos.take());
223 }
224
225 pub async fn from_query_result(
227 result: Result<Output>,
228 metric_name: Option<String>,
229 result_type: ValueType,
230 query_id: Option<&str>,
231 ) -> Self {
232 let collector = query_id.and_then(get_promql_annotation_collector);
234 let response: Result<Self> = async {
235 let result = result?;
236 let mut resp =
237 match result.data {
238 OutputData::RecordBatches(batches) => Self::success(
239 Self::record_batches_to_data(batches, metric_name, result_type)?,
240 ),
241 OutputData::Stream(stream) => Self::success(
242 Self::consume_stream_to_data(stream, metric_name, result_type).await?,
243 ),
244 OutputData::AffectedRows(_) => Self::error(
245 StatusCode::Unexpected,
246 "expected data result, but got affected rows",
247 ),
248 };
249
250 if let Some(physical_plan) = result.meta.plan {
251 let mut result_map = HashMap::new();
252 let mut tmp = vec![&mut result_map];
253 collect_plan_metrics(&physical_plan, &mut tmp);
254
255 let re = result_map
256 .into_iter()
257 .map(|(k, v)| (k, Value::from(v)))
258 .collect();
259 resp.resp_metrics = re;
260 }
261
262 Ok(resp)
263 }
264 .await;
265
266 let result_type_string = result_type.to_string();
267
268 let mut response = match response {
269 Ok(resp) => resp,
270 Err(err) => {
271 if err.status_code() == StatusCode::TableNotFound
273 || err.status_code() == StatusCode::TableColumnNotFound
274 {
275 Self::success(PrometheusResponse::PromData(PromData {
276 result_type: result_type_string,
277 ..Default::default()
278 }))
279 } else {
280 Self::error(err.status_code(), err.output_msg())
281 }
282 }
283 };
284 if let Some(collector) = collector {
285 response.append_promql_annotations(&collector);
286 }
287 response
288 }
289
290 fn record_batches_to_data(
292 batches: RecordBatches,
293 metric_name: Option<String>,
294 result_type: ValueType,
295 ) -> Result<PrometheusResponse> {
296 if batches.iter().next().is_none() {
298 return Ok(empty_data(result_type));
299 }
300
301 let layout = ColumnLayout::infer(&batches.schema())?;
302
303 let mut buffer = IndexMap::<Vec<(String, String)>, PromSeriesSamples>::new();
306
307 for batch in batches.iter() {
308 merge_batch(&mut buffer, &layout, batch, metric_name.as_deref())?;
309 }
310
311 samples_to_data(buffer, result_type)
312 }
313
314 async fn consume_stream_to_data(
320 mut stream: SendableRecordBatchStream,
321 metric_name: Option<String>,
322 result_type: ValueType,
323 ) -> Result<PrometheusResponse> {
324 let Some(first_batch) = stream.try_next().await.context(CollectRecordbatchSnafu)? else {
327 return Ok(empty_data(result_type));
328 };
329
330 let layout = ColumnLayout::infer(&stream.schema())?;
331
332 let mut buffer = IndexMap::<Vec<(String, String)>, PromSeriesSamples>::new();
335
336 merge_batch(&mut buffer, &layout, &first_batch, metric_name.as_deref())?;
337 drop(first_batch);
339
340 while let Some(batch) = stream.try_next().await.context(CollectRecordbatchSnafu)? {
341 merge_batch(&mut buffer, &layout, &batch, metric_name.as_deref())?;
342 }
343
344 samples_to_data(buffer, result_type)
345 }
346}
347
348fn empty_data(result_type: ValueType) -> PrometheusResponse {
350 PrometheusResponse::PromData(PromData {
351 result_type: result_type.to_string(),
352 ..Default::default()
353 })
354}
355
356struct ColumnLayout {
358 timestamp_column_index: usize,
359 tag_column_indices: Vec<usize>,
360 tag_names: Vec<String>,
361 first_field_column_index: Option<usize>,
362 native_histogram_column_index: Option<usize>,
363}
364
365impl ColumnLayout {
366 fn infer(schema: &Schema) -> Result<Self> {
369 let mut timestamp_column_index = None;
370 let mut tag_column_indices = Vec::new();
371 let mut first_field_column_index = None;
372 let mut native_histogram_column_index = None;
373
374 for (i, column) in schema.column_schemas().iter().enumerate() {
375 match column.data_type {
376 ConcreteDataType::Timestamp(datatypes::types::TimestampType::Millisecond(_))
377 if timestamp_column_index.is_none() =>
378 {
379 timestamp_column_index = Some(i);
380 }
381 ConcreteDataType::Float32(_)
383 | ConcreteDataType::Float64(_)
384 | ConcreteDataType::Int8(_)
385 | ConcreteDataType::Int16(_)
386 | ConcreteDataType::Int32(_)
387 | ConcreteDataType::Int64(_)
388 | ConcreteDataType::UInt8(_)
389 | ConcreteDataType::UInt16(_)
390 | ConcreteDataType::UInt32(_)
391 | ConcreteDataType::UInt64(_)
392 if first_field_column_index.is_none() =>
393 {
394 first_field_column_index = Some(i);
395 }
396 _ if native_histogram_column_index.is_none()
397 && is_native_histogram_value_type(&column.data_type) =>
398 {
399 native_histogram_column_index = Some(i);
400 }
401 ConcreteDataType::String(_) => {
402 tag_column_indices.push(i);
403 }
404 _ => {}
405 }
406 }
407
408 let timestamp_column_index = timestamp_column_index.context(UnexpectedResultSnafu {
409 reason: "no timestamp column found".to_string(),
410 })?;
411 if first_field_column_index.is_none() && native_histogram_column_index.is_none() {
412 return UnexpectedResultSnafu {
413 reason: "no value column found".to_string(),
414 }
415 .fail();
416 }
417
418 let tag_names = tag_column_indices
419 .iter()
420 .map(|index| schema.column_name_by_index(*index).to_string())
421 .collect::<Vec<_>>();
422
423 Ok(Self {
424 timestamp_column_index,
425 tag_column_indices,
426 tag_names,
427 first_field_column_index,
428 native_histogram_column_index,
429 })
430 }
431}
432
433struct SeriesKeyLookup<'a>(&'a [(&'a str, &'a str)]);
440
441impl Equivalent<Vec<(String, String)>> for SeriesKeyLookup<'_> {
442 fn equivalent(&self, key: &Vec<(String, String)>) -> bool {
443 self.0.len() == key.len()
444 && self
445 .0
446 .iter()
447 .zip(key)
448 .all(|((name, value), (key_name, key_value))| {
449 *name == key_name.as_str() && *value == key_value.as_str()
450 })
451 }
452}
453
454fn merge_batch(
456 buffer: &mut IndexMap<Vec<(String, String)>, PromSeriesSamples>,
457 layout: &ColumnLayout,
458 batch: &RecordBatch,
459 metric_name: Option<&str>,
460) -> Result<()> {
461 let mut tags: Vec<(&str, &str)> = Vec::with_capacity(layout.tag_column_indices.len() + 1);
466
467 let tag_columns = layout
469 .tag_column_indices
470 .iter()
471 .map(|i| batch.column(*i))
472 .collect::<Vec<_>>();
473 let timestamp_column = batch
474 .column(layout.timestamp_column_index)
475 .as_primitive::<TimestampMillisecondType>();
476
477 let field_array = layout
478 .first_field_column_index
479 .map(|index| arrow::compute::cast(batch.column(index), &DataType::Float64))
480 .transpose()
481 .context(ArrowSnafu)?;
482 let field_column = field_array
483 .as_ref()
484 .map(|array| array.as_primitive::<Float64Type>());
485 let native_histogram_column = layout
486 .native_histogram_column_index
487 .map(|index| {
488 batch
489 .column(index)
490 .as_any()
491 .downcast_ref::<StructArray>()
492 .with_context(|| UnexpectedResultSnafu {
493 reason: "native histogram column is not a struct array",
494 })
495 })
496 .transpose()?;
497
498 let label_runs = if prefer_label_runs(&tag_columns, batch.num_rows()) {
502 Either::Left(label_runs(&tag_columns, batch.num_rows())?.into_iter())
503 } else {
504 Either::Right((0..batch.num_rows()).map(|row| row..row + 1))
505 };
506
507 for run in label_runs {
509 let mut run_entry_index = None;
510 for row_index in run {
511 let value = field_column.and_then(|field_column| {
512 if !field_column.is_valid(row_index) {
513 return None;
514 }
515 let value = field_column.value(row_index);
516 (!is_prometheus_stale_nan(value))
517 .then_some((timestamp_column.value(row_index), value))
518 });
519 let histogram = native_histogram_column
520 .and_then(|column| {
521 read_histogram(column, row_index)
522 .context(DataFusionSnafu)
523 .transpose()
524 })
525 .transpose()?
526 .filter(|histogram| !is_prometheus_stale_nan(histogram.sum))
527 .map(|histogram| {
528 prometheus_native_histogram(&histogram)
529 .map(|histogram| (timestamp_column.value(row_index), histogram))
530 })
531 .transpose()?;
532
533 if value.is_none() && histogram.is_none() {
534 continue;
535 }
536
537 let entry_index = match run_entry_index {
538 Some(index) => index,
539 None => {
540 tags.clear();
542 if let Some(metric_name) = metric_name {
543 tags.push((METRIC_NAME, metric_name));
544 }
545 for (tag_column, tag_name) in tag_columns.iter().zip(layout.tag_names.iter()) {
546 if let Some(tag_value) = string_array_value_at_index(tag_column, row_index)
547 {
548 tags.push((tag_name.as_str(), tag_value));
549 }
550 }
551
552 let hash = buffer.hasher().hash_one(&tags);
556 let entry = buffer
557 .raw_entry_mut_v1()
558 .from_key_hashed_nocheck(hash, &SeriesKeyLookup(tags.as_slice()));
559 let index = entry.index();
560 if let RawEntryMut::Vacant(entry) = entry {
561 let key = tags
564 .iter()
565 .map(|(name, value)| (name.to_string(), value.to_string()))
566 .collect::<Vec<_>>();
567 entry.insert_hashed_nocheck(hash, key, PromSeriesSamples::default());
568 }
569 run_entry_index = Some(index);
570 index
571 }
572 };
573 let samples = &mut buffer[entry_index];
574 if let Some((timestamp_millis, histogram)) = histogram {
575 samples
576 .histograms
577 .push((timestamp_millis as f64 / 1000.0, histogram));
578 } else if let Some((timestamp_millis, value)) = value {
579 samples.values.push((
580 timestamp_millis as f64 / 1000.0,
581 PromSampleValue::Number(value),
582 ));
583 }
584 }
585 }
586 Ok(())
587}
588
589fn samples_to_data(
592 buffer: IndexMap<Vec<(String, String)>, PromSeriesSamples>,
593 result_type: ValueType,
594) -> Result<PrometheusResponse> {
595 let mut result = match result_type {
597 ValueType::Vector => PromQueryResult::Vector(vec![]),
598 ValueType::Matrix => PromQueryResult::Matrix(vec![]),
599 ValueType::Scalar => PromQueryResult::Scalar(None),
600 ValueType::String => PromQueryResult::String(None),
601 };
602
603 buffer.into_iter().for_each(|(tags, mut samples)| {
605 let metric = tags.into_iter().collect::<BTreeMap<String, String>>();
606 match result {
607 PromQueryResult::Vector(ref mut v) => {
608 let histogram = samples.histograms.pop();
609 let value = if histogram.is_none() {
610 samples
611 .values
612 .pop()
613 .map(|(timestamp, value)| (timestamp, value.into_string()))
614 } else {
615 None
616 };
617 v.push(PromSeriesVector {
618 metric,
619 value,
620 histogram,
621 });
622 }
623 PromQueryResult::Matrix(ref mut v) => {
624 if !samples.values.is_sorted_by(|a, b| a.0 <= b.0) {
626 samples
627 .values
628 .sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(Ordering::Equal));
629 }
630 if !samples.histograms.is_sorted_by(|a, b| a.0 <= b.0) {
631 samples
632 .histograms
633 .sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(Ordering::Equal));
634 }
635
636 v.push(PromSeriesMatrix {
637 metric,
638 values: samples.values,
639 histograms: samples.histograms,
640 });
641 }
642 PromQueryResult::Scalar(ref mut v) => {
643 *v = samples
644 .values
645 .pop()
646 .map(|(timestamp, value)| (timestamp, value.into_string()));
647 }
648 PromQueryResult::String(ref mut _v) => {
649 }
651 }
652 });
653
654 if let PromQueryResult::Matrix(ref mut v) = result {
657 v.sort_by(|a, b| a.metric.cmp(&b.metric));
658 }
659
660 let result_type_string = result_type.to_string();
661 let data = PrometheusResponse::PromData(PromData {
662 result_type: result_type_string,
663 result,
664 });
665
666 Ok(data)
667}
668
669fn label_runs(columns: &[&ArrayRef], rows: usize) -> Result<Vec<Range<usize>>> {
678 let mut boundaries: Option<BooleanBuffer> = None;
679 if rows >= 2 {
680 for column in columns {
681 let changed = distinct(&column.slice(0, rows - 1), &column.slice(1, rows - 1))
682 .context(ArrowSnafu)?;
683 boundaries = Some(match boundaries {
684 Some(accumulated) => &accumulated | changed.values(),
685 None => changed.values().clone(),
686 });
687 }
688 }
689
690 let mut runs = Vec::new();
691 let mut start = 0;
692 for boundary in boundaries.iter().flat_map(BooleanBuffer::set_indices) {
693 runs.push(start..boundary + 1);
694 start = boundary + 1;
695 }
696 runs.push(start..rows);
697 Ok(runs)
698}
699
700fn prefer_label_runs(columns: &[&ArrayRef], rows: usize) -> bool {
709 if columns.is_empty() {
710 return true;
711 }
712 if rows < 2 {
713 return false;
714 }
715 let mut position = 0x9e37_79b9_u32;
716 let mut changes = 0;
717 for _ in 0..8 {
718 position ^= position << 13;
719 position ^= position >> 17;
720 position ^= position << 5;
721 let row = position as usize % (rows - 1);
722 if columns.iter().any(|column| {
723 string_array_value_at_index(column, row) != string_array_value_at_index(column, row + 1)
724 }) {
725 changes += 1;
726 if changes > 2 {
727 return false;
728 }
729 }
730 }
731 true
732}
733
734fn merge_annotations(target: &mut Option<Vec<String>>, source: Option<Vec<String>>) {
735 let Some(source) = source else {
736 return;
737 };
738 let target = target.get_or_insert_default();
739 target.extend(source);
740 target.sort();
741 target.dedup();
742}
743
744#[cfg(test)]
745mod tests {
746 use std::sync::Arc;
747
748 use arrow::array::StringViewArray;
749 use common_query::native_histogram::{
750 CUSTOM_BUCKETS_SCHEMA, CounterResetHint, NativeHistogram, Span, build_histogram_array,
751 native_histogram_value_type,
752 };
753 use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS;
754 use common_query::promql_annotations::{
755 PromqlAnnotationCollector, promql_annotation_collector,
756 };
757 use common_recordbatch::{RecordBatch, RecordBatches};
758 use datatypes::data_type::ConcreteDataType;
759 use datatypes::schema::{ColumnSchema, Schema, SchemaRef};
760 use datatypes::vectors::{
761 Float64Vector, StringVector, StructVector, TimestampMillisecondVector, VectorRef,
762 };
763
764 use super::*;
765
766 #[tokio::test]
767 async fn empty_stream_reads_as_empty_data() {
768 let schema = Arc::new(Schema::new(vec![
772 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
773 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false),
774 ]));
775 let stream = RecordBatches::try_new(schema, vec![]).unwrap().as_stream();
776
777 let response =
778 PrometheusJsonResponse::consume_stream_to_data(stream, None, ValueType::Matrix).await;
779
780 assert!(matches!(
781 response.unwrap(),
782 PrometheusResponse::PromData(PromData {
783 result: PromQueryResult::Matrix(series),
784 ..
785 }) if series.is_empty()
786 ));
787 }
788
789 #[tokio::test]
790 async fn stream_result_matches_record_batches_result() {
791 let schema = Arc::new(Schema::new(vec![
792 ColumnSchema::new(
793 "timestamp",
794 ConcreteDataType::timestamp_millisecond_datatype(),
795 false,
796 ),
797 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
798 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false),
799 ]));
800 let batch_rows: [Vec<(&str, i64, f64)>; 2] = [
804 vec![
805 ("a", 4_000, 0.0),
806 ("a", 5_000, 1.0),
807 ("a", 6_000, 2.0),
808 ("a", 7_000, 3.0),
809 ],
810 vec![
811 ("a", 0, 4.0),
812 ("a", 1_000, 5.0),
813 ("b", 2_000, 6.0),
814 ("b", 3_000, 7.0),
815 ],
816 ];
817 let batches: Vec<_> = batch_rows
818 .into_iter()
819 .map(|rows| {
820 RecordBatch::new(
821 schema.clone(),
822 vec![
823 Arc::new(TimestampMillisecondVector::from_values(
824 rows.iter().map(|(_, ts, _)| *ts),
825 )) as _,
826 Arc::new(StringVector::from(
827 rows.iter()
828 .map(|(host, _, _)| Some(*host))
829 .collect::<Vec<_>>(),
830 )) as _,
831 Arc::new(Float64Vector::from_values(rows.iter().map(|(_, _, v)| *v))) as _,
832 ],
833 )
834 .unwrap()
835 })
836 .collect();
837
838 let streamed = PrometheusJsonResponse::consume_stream_to_data(
839 RecordBatches::try_new(schema.clone(), batches.clone())
840 .unwrap()
841 .as_stream(),
842 Some("metric".to_string()),
843 ValueType::Matrix,
844 )
845 .await
846 .unwrap();
847 let eager = PrometheusJsonResponse::record_batches_to_data(
848 RecordBatches::try_new(schema, batches).unwrap(),
849 Some("metric".to_string()),
850 ValueType::Matrix,
851 )
852 .unwrap();
853
854 assert_eq!(
855 serde_json::to_value(&streamed).unwrap(),
856 serde_json::to_value(&eager).unwrap()
857 );
858 let PrometheusResponse::PromData(PromData {
859 result: PromQueryResult::Matrix(series),
860 ..
861 }) = streamed
862 else {
863 panic!("expected matrix response");
864 };
865 assert_eq!(series.len(), 2);
866 let series_a = series
870 .iter()
871 .find(|s| s.metric.get("host").map(String::as_str) == Some("a"))
872 .unwrap();
873 let samples_a: Vec<(i64, f64)> = series_a
874 .values
875 .iter()
876 .map(|(ts, value)| {
877 let PromSampleValue::Number(v) = value else {
878 panic!("expected numeric sample");
879 };
880 ((ts * 1000.0) as i64, *v)
881 })
882 .collect();
883 assert_eq!(
884 samples_a,
885 vec![
886 (0, 4.0),
887 (1_000, 5.0),
888 (4_000, 0.0),
889 (5_000, 1.0),
890 (6_000, 2.0),
891 (7_000, 3.0),
892 ]
893 );
894 }
895
896 #[tokio::test]
897 async fn stream_mixed_histogram_result_matches_record_batches_result() {
898 let schema = Arc::new(Schema::new(vec![
899 ColumnSchema::new(
900 "timestamp",
901 ConcreteDataType::timestamp_millisecond_datatype(),
902 false,
903 ),
904 ColumnSchema::new("kind", ConcreteDataType::string_datatype(), false),
905 ColumnSchema::new("float", ConcreteDataType::float64_datatype(), true),
906 ColumnSchema::new("histogram", native_histogram_value_type().clone(), true),
907 ]));
908 let first = RecordBatch::new(
912 schema.clone(),
913 vec![
914 Arc::new(TimestampMillisecondVector::from_values([1_000, 2_000])) as _,
915 Arc::new(StringVector::from(vec![Some("float"); 2])) as _,
916 Arc::new(Float64Vector::from(vec![Some(1.25), Some(0.5)])) as _,
917 histogram_vector(&[None, None]),
918 ],
919 )
920 .unwrap();
921 let second = RecordBatch::new(
922 schema.clone(),
923 vec![
924 Arc::new(TimestampMillisecondVector::from_values([3_000, 3_000])) as _,
925 Arc::new(StringVector::from(vec![Some("float"), Some("histogram")])) as _,
926 Arc::new(Float64Vector::from(vec![Some(2.0), None])) as _,
927 histogram_vector(&[None, Some(sample_histogram())]),
928 ],
929 )
930 .unwrap();
931
932 let streamed = PrometheusJsonResponse::consume_stream_to_data(
933 RecordBatches::try_new(schema.clone(), vec![first.clone(), second.clone()])
934 .unwrap()
935 .as_stream(),
936 Some("mixed_metric".to_string()),
937 ValueType::Vector,
938 )
939 .await
940 .unwrap();
941 let eager = PrometheusJsonResponse::record_batches_to_data(
942 RecordBatches::try_new(schema, vec![first, second]).unwrap(),
943 Some("mixed_metric".to_string()),
944 ValueType::Vector,
945 )
946 .unwrap();
947
948 assert_eq!(
949 serde_json::to_value(&streamed).unwrap(),
950 serde_json::to_value(&eager).unwrap()
951 );
952 let PrometheusResponse::PromData(PromData {
953 result: PromQueryResult::Vector(series),
954 ..
955 }) = streamed
956 else {
957 panic!("expected vector response");
958 };
959 assert_eq!(series.len(), 2);
960 let histogram = series
961 .iter()
962 .find(|series| series.metric["kind"] == "histogram")
963 .unwrap();
964 assert!(histogram.histogram.is_some());
965 assert!(histogram.value.is_none());
966 }
967
968 struct LazyDropCheckingStream {
975 schema: SchemaRef,
976 batch_index: usize,
977 last_batch: Option<std::sync::Weak<RecordBatch>>,
980 }
981
982 impl LazyDropCheckingStream {
983 fn new(schema: SchemaRef) -> Self {
984 Self {
985 schema,
986 batch_index: 0,
987 last_batch: None,
988 }
989 }
990
991 fn make_batch(&self) -> RecordBatch {
992 let base = self.batch_index as i64 * 1000;
993 RecordBatch::new(
994 self.schema.clone(),
995 vec![
996 Arc::new(TimestampMillisecondVector::from_values([base])) as _,
997 Arc::new(StringVector::from(vec![Some("host")])) as _,
998 Arc::new(Float64Vector::from_values([base as f64])) as _,
999 ],
1000 )
1001 .unwrap()
1002 }
1003 }
1004
1005 impl futures::Stream for LazyDropCheckingStream {
1006 type Item = common_recordbatch::error::Result<RecordBatch>;
1007
1008 fn poll_next(
1009 mut self: std::pin::Pin<&mut Self>,
1010 _cx: &mut std::task::Context<'_>,
1011 ) -> std::task::Poll<Option<Self::Item>> {
1012 if let Some(weak) = &self.last_batch {
1015 assert!(
1016 weak.upgrade().is_none(),
1017 "previous batch was not released before polling the next one"
1018 );
1019 }
1020
1021 if self.batch_index >= 3 {
1022 return std::task::Poll::Ready(None);
1023 }
1024
1025 let batch = Arc::new(self.make_batch());
1026 self.last_batch = Some(Arc::downgrade(&batch));
1027 self.batch_index += 1;
1028 std::task::Poll::Ready(Some(Ok((*batch).clone())))
1029 }
1030 }
1031
1032 impl common_recordbatch::RecordBatchStream for LazyDropCheckingStream {
1033 fn schema(&self) -> SchemaRef {
1034 self.schema.clone()
1035 }
1036
1037 fn output_ordering(&self) -> Option<&[common_recordbatch::OrderOption]> {
1038 None
1039 }
1040
1041 fn metrics(&self) -> Option<common_recordbatch::adapter::RecordBatchMetrics> {
1042 None
1043 }
1044 }
1045
1046 #[tokio::test]
1047 async fn stream_drops_each_batch_before_polling_next() {
1048 let schema = Arc::new(Schema::new(vec![
1049 ColumnSchema::new(
1050 "timestamp",
1051 ConcreteDataType::timestamp_millisecond_datatype(),
1052 false,
1053 ),
1054 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
1055 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false),
1056 ]));
1057 let stream = LazyDropCheckingStream::new(schema);
1058
1059 let response = PrometheusJsonResponse::consume_stream_to_data(
1060 Box::pin(stream),
1061 None,
1062 ValueType::Matrix,
1063 )
1064 .await;
1065
1066 assert!(matches!(
1067 response.unwrap(),
1068 PrometheusResponse::PromData(PromData {
1069 result: PromQueryResult::Matrix(series),
1070 ..
1071 }) if series.len() == 1 && series[0].values.len() == 3
1072 ));
1073 }
1074
1075 #[tokio::test]
1076 async fn query_response_preserves_and_merges_promql_annotations() {
1077 let query_id = "query_response_preserves_and_merges_promql_annotations";
1078 let left = promql_annotation_collector(query_id);
1079 left.record_warning("shared warning");
1080 left.record_info("left info");
1081 let mut response = PrometheusJsonResponse::from_query_result(
1082 Ok(Output::new_with_record_batches(RecordBatches::empty())),
1083 None,
1084 ValueType::Vector,
1085 Some(query_id),
1086 )
1087 .await;
1088
1089 let right = PromqlAnnotationCollector::default();
1090 right.record_warning("shared warning");
1091 right.record_info("right info");
1092 let mut other = PrometheusJsonResponse::success(PrometheusResponse::None);
1093 other.append_promql_annotations(&right);
1094 response.append_query_response(other);
1095
1096 assert_eq!(response.warnings, Some(vec!["shared warning".to_string()]));
1097 assert_eq!(
1098 response.infos,
1099 Some(vec!["left info".to_string(), "right info".to_string()])
1100 );
1101 let json = serde_json::to_value(response).unwrap();
1102 assert_eq!(json["warnings"], serde_json::json!(["shared warning"]));
1103 assert_eq!(
1104 json["infos"],
1105 serde_json::json!(["left info", "right info"])
1106 );
1107 }
1108
1109 fn sample_histogram() -> NativeHistogram {
1110 NativeHistogram {
1111 schema: 0,
1112 zero_threshold: 0.001,
1113 sum: 3.0,
1114 reset_hint: CounterResetHint::Unknown,
1115 start_timestamp: Some(0),
1116 custom_values: vec![],
1117 positive_spans: vec![Span {
1118 offset: 0,
1119 length: 1,
1120 }],
1121 negative_spans: vec![],
1122 count: 2.0,
1123 zero_count: 1.0,
1124 positive_buckets: vec![1.0],
1125 negative_buckets: vec![],
1126 }
1127 }
1128
1129 fn histogram_vector(values: &[Option<NativeHistogram>]) -> VectorRef {
1130 let histogram_array = build_histogram_array(values);
1131 let histogram_array = histogram_array
1132 .as_any()
1133 .downcast_ref::<StructArray>()
1134 .unwrap()
1135 .clone();
1136 let ConcreteDataType::Struct(histogram_type) = native_histogram_value_type().clone() else {
1137 unreachable!("native histogram type must be a struct")
1138 };
1139 Arc::new(StructVector::try_new(histogram_type, histogram_array).unwrap())
1140 }
1141
1142 #[test]
1143 fn format_prometheus_sample_value_uses_ryu_for_finite_values() {
1144 let values = [
1145 1.5,
1146 0.1,
1147 1.0,
1148 0.0,
1149 -0.0,
1150 100.0,
1151 1e-6,
1152 1e-7,
1153 1e21,
1154 1e30,
1155 f64::MAX,
1156 f64::MIN_POSITIVE,
1157 ];
1158
1159 for value in values {
1160 let output = format_prometheus_sample_value(value);
1161 assert_eq!(output.parse::<f64>().unwrap().to_bits(), value.to_bits());
1162 }
1163
1164 assert_eq!(format_prometheus_sample_value(1.0), "1.0");
1166 assert_eq!(format_prometheus_sample_value(-0.0), "-0.0");
1167 assert_eq!(format_prometheus_sample_value(100.0), "100.0");
1168 assert_eq!(format_prometheus_sample_value(1e-6), "1e-6");
1169 assert_eq!(format_prometheus_sample_value(1e-7), "1e-7");
1170 assert_eq!(format_prometheus_sample_value(1e21), "1e21");
1171
1172 for value in [
1175 f64::from_bits(0x42374876e8000400),
1176 f64::from_bits(0x3ff0000800000000),
1177 f64::from_bits(0x430a8e5672bc7312),
1178 ] {
1179 let ryu_output = format_prometheus_sample_value(value);
1180 let std_output = value.to_string();
1181 assert_ne!(ryu_output, std_output);
1182 assert_eq!(ryu_output.parse::<f64>().unwrap(), value);
1183 assert_eq!(std_output.parse::<f64>().unwrap(), value);
1184 }
1185 }
1186
1187 #[test]
1188 fn format_prometheus_sample_value_preserves_nonfinite_values() {
1189 assert_eq!(format_prometheus_sample_value(f64::NAN), "NaN");
1190 assert_eq!(format_prometheus_sample_value(f64::INFINITY), "inf");
1191 assert_eq!(format_prometheus_sample_value(f64::NEG_INFINITY), "-inf");
1192
1193 assert!(
1196 format_prometheus_sample_value(f64::NAN)
1197 .parse::<f64>()
1198 .unwrap()
1199 .is_nan()
1200 );
1201 }
1202
1203 #[test]
1204 fn sample_value_serialization_matches_eager_formatting() {
1205 let mut values = vec![
1206 0.0,
1207 -0.0,
1208 f64::MAX,
1209 f64::MIN,
1210 f64::MIN_POSITIVE,
1211 f64::from_bits(1),
1212 1e-7,
1213 1e21,
1214 f64::NAN,
1215 f64::INFINITY,
1216 f64::NEG_INFINITY,
1217 ];
1218 let mut bits = 0x1234_5678_9876_5432_u64;
1219 for _ in 0..1000 {
1220 bits ^= bits << 13;
1221 bits ^= bits >> 7;
1222 bits ^= bits << 17;
1223 values.push(f64::from_bits(bits));
1224 }
1225 for value in values {
1226 assert_eq!(
1227 serde_json::to_string(&PromSampleValue::Number(value)).unwrap(),
1228 serde_json::to_string(&format_prometheus_sample_value(value)).unwrap()
1229 );
1230 }
1231 for value in ["1.00", "+Inf", "-0", "NaN", "not-a-number", "", "\"\\\n"] {
1232 let json = serde_json::to_string(value).unwrap();
1233 let parsed: PromSampleValue = serde_json::from_str(&json).unwrap();
1234 assert_eq!(serde_json::to_string(&parsed).unwrap(), json);
1235 }
1236 for value in ["1", "null", "true", "[]", "{}"] {
1237 assert!(serde_json::from_str::<PromSampleValue>(value).is_err());
1238 }
1239 }
1240
1241 #[tokio::test]
1242 async fn matrix_response_body_matches_eagerly_formatted_json() {
1243 let schema = Arc::new(Schema::new(vec![
1244 ColumnSchema::new(
1245 "timestamp",
1246 ConcreteDataType::timestamp_millisecond_datatype(),
1247 false,
1248 ),
1249 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
1250 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false),
1251 ]));
1252 let batches = RecordBatches::try_new(
1253 schema.clone(),
1254 vec![
1255 RecordBatch::new(
1256 schema,
1257 vec![
1258 Arc::new(TimestampMillisecondVector::from_values([
1259 1000, 2000, 3000, 4000, 5000,
1260 ])) as _,
1261 Arc::new(StringVector::from(vec![Some("a"); 5])) as _,
1262 Arc::new(Float64Vector::from_values([
1263 -0.0,
1264 f64::NAN,
1265 1e-7,
1266 f64::INFINITY,
1267 f64::NEG_INFINITY,
1268 ])) as _,
1269 ],
1270 )
1271 .unwrap(),
1272 ],
1273 )
1274 .unwrap();
1275 let actual = PrometheusJsonResponse::from_query_result(
1276 Ok(Output::new_with_record_batches(batches)),
1277 None,
1278 ValueType::Matrix,
1279 None,
1280 )
1281 .await;
1282 let expected: PrometheusJsonResponse = serde_json::from_value(serde_json::json!({
1285 "status": "success",
1286 "data": {"resultType": "matrix", "result": [{
1287 "metric": {"host": "a"},
1288 "values": [[1.0, "-0.0"], [2.0, "NaN"], [3.0, "1e-7"], [4.0, "inf"], [5.0, "-inf"]]
1289 }]}
1290 }))
1291 .unwrap();
1292 assert_eq!(
1293 serde_json::to_string(&actual).unwrap(),
1294 serde_json::to_string(&expected).unwrap()
1295 );
1296 }
1297
1298 #[test]
1299 fn matrix_response_preserves_ordinary_nan_and_filters_stale_markers() {
1300 let schema = Arc::new(Schema::new(vec![
1301 ColumnSchema::new(
1302 "timestamp",
1303 ConcreteDataType::timestamp_millisecond_datatype(),
1304 false,
1305 ),
1306 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
1307 ]));
1308 let batch = RecordBatch::new(
1309 schema.clone(),
1310 vec![
1311 Arc::new(TimestampMillisecondVector::from_vec(vec![
1312 1_000, 2_000, 3_000, 4_000,
1313 ])) as _,
1314 Arc::new(Float64Vector::from(vec![
1315 Some(1.0),
1316 Some(f64::from_bits(0x7ff8_0000_0000_0000)),
1317 Some(f64::from_bits(0x7ff0_0000_0000_0002)),
1318 None,
1319 ])) as _,
1320 ],
1321 )
1322 .unwrap();
1323 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1324
1325 let response =
1326 PrometheusJsonResponse::record_batches_to_data(batches, None, ValueType::Matrix)
1327 .unwrap();
1328 let PrometheusResponse::PromData(data) = response else {
1329 panic!("expected Prometheus data response");
1330 };
1331 let PromQueryResult::Matrix(series) = data.result else {
1332 panic!("expected matrix result");
1333 };
1334
1335 assert_eq!(series.len(), 1);
1336 assert_eq!(
1337 serde_json::to_value(&series[0].values).unwrap(),
1338 serde_json::json!([[1.0, "1.0"], [2.0, "NaN"]])
1339 );
1340 }
1341
1342 #[test]
1343 fn record_batches_to_data_formats_values_with_ryu() {
1344 let schema = Arc::new(Schema::new(vec![
1345 ColumnSchema::new(
1346 "timestamp",
1347 ConcreteDataType::timestamp_millisecond_datatype(),
1348 false,
1349 ),
1350 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
1351 ]));
1352 let batch = RecordBatch::new(
1353 schema.clone(),
1354 vec![
1355 Arc::new(TimestampMillisecondVector::from_vec(vec![
1356 1_000, 2_000, 3_000, 4_000, 5_000, 6_000,
1357 ])) as _,
1358 Arc::new(Float64Vector::from(vec![
1359 Some(0.0),
1360 Some(-0.0),
1361 Some(1.25),
1362 Some(1e30),
1363 Some(1e-7),
1364 Some(f64::MAX),
1365 ])) as _,
1366 ],
1367 )
1368 .unwrap();
1369 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1370
1371 let response =
1372 PrometheusJsonResponse::record_batches_to_data(batches, None, ValueType::Matrix)
1373 .unwrap();
1374 let PrometheusResponse::PromData(PromData {
1375 result: PromQueryResult::Matrix(series),
1376 ..
1377 }) = response
1378 else {
1379 panic!("expected matrix response");
1380 };
1381
1382 assert_eq!(series.len(), 1);
1383 let input_values = [0.0, -0.0, 1.25, 1e30, 1e-7, f64::MAX];
1384 let expected = input_values
1387 .into_iter()
1388 .enumerate()
1389 .map(|(index, value)| {
1390 let expected_value = if value.is_finite() {
1391 let mut buffer = Buffer::new();
1392 buffer.format_finite(value).to_string()
1393 } else {
1394 value.to_string()
1395 };
1396 ((index + 1) as f64, expected_value)
1397 })
1398 .collect::<Vec<_>>();
1399 assert_eq!(
1400 serde_json::to_value(&series[0].values).unwrap(),
1401 serde_json::to_value(expected).unwrap()
1402 );
1403 }
1404
1405 #[test]
1406 fn record_batches_to_data_preserves_infinity_output() {
1407 let schema = Arc::new(Schema::new(vec![
1409 ColumnSchema::new(
1410 "timestamp",
1411 ConcreteDataType::timestamp_millisecond_datatype(),
1412 false,
1413 ),
1414 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
1415 ]));
1416 let batch = RecordBatch::new(
1417 schema.clone(),
1418 vec![
1419 Arc::new(TimestampMillisecondVector::from_vec(vec![
1420 1_000, 2_000, 3_000,
1421 ])) as _,
1422 Arc::new(Float64Vector::from(vec![
1423 Some(f64::INFINITY),
1424 Some(f64::NEG_INFINITY),
1425 Some(f64::NAN),
1426 ])) as _,
1427 ],
1428 )
1429 .unwrap();
1430 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1431
1432 let response =
1433 PrometheusJsonResponse::record_batches_to_data(batches, None, ValueType::Matrix)
1434 .unwrap();
1435 let PrometheusResponse::PromData(PromData {
1436 result: PromQueryResult::Matrix(series),
1437 ..
1438 }) = response
1439 else {
1440 panic!("expected matrix response");
1441 };
1442
1443 assert_eq!(series.len(), 1);
1444 assert_eq!(
1445 serde_json::to_value(&series[0].values).unwrap(),
1446 serde_json::json!([[1.0, "inf"], [2.0, "-inf"], [3.0, "NaN"]])
1447 );
1448 }
1449
1450 #[test]
1451 fn record_batches_to_data_groups_clustered_series() {
1452 let schema = Arc::new(Schema::new(vec![
1456 ColumnSchema::new(
1457 "timestamp",
1458 ConcreteDataType::timestamp_millisecond_datatype(),
1459 false,
1460 ),
1461 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
1462 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
1463 ]));
1464 let batch = RecordBatch::new(
1465 schema.clone(),
1466 vec![
1467 Arc::new(TimestampMillisecondVector::from_vec(vec![
1468 1_000, 2_000, 3_000, 4_000, 5_000, 6_000,
1469 ])) as _,
1470 Arc::new(StringVector::from(vec![
1471 Some("a"),
1472 Some("b"),
1473 Some("a"),
1474 Some("a"),
1475 Some("b"),
1476 Some("c"),
1477 ])) as _,
1478 Arc::new(Float64Vector::from(vec![
1479 Some(1.0),
1480 Some(2.0),
1481 Some(3.0),
1482 Some(4.0),
1483 Some(5.0),
1484 Some(6.0),
1485 ])) as _,
1486 ],
1487 )
1488 .unwrap();
1489 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1490
1491 let response = PrometheusJsonResponse::record_batches_to_data(
1492 batches,
1493 Some("metric".to_string()),
1494 ValueType::Vector,
1495 )
1496 .unwrap();
1497 let PrometheusResponse::PromData(PromData {
1498 result: PromQueryResult::Vector(series),
1499 ..
1500 }) = response
1501 else {
1502 panic!("expected vector response");
1503 };
1504
1505 assert_eq!(series.len(), 3);
1506 assert_eq!(
1508 series
1509 .iter()
1510 .map(|series| series.metric["host"].as_str())
1511 .collect::<Vec<_>>(),
1512 vec!["a", "b", "c"]
1513 );
1514 assert_eq!(series[0].value, Some((4.0, "4.0".to_string())));
1516 assert_eq!(series[1].value, Some((5.0, "5.0".to_string())));
1517 assert_eq!(series[2].value, Some((6.0, "6.0".to_string())));
1518 }
1519
1520 #[test]
1521 fn label_strategy_switches_preserve_all_rows_when_probes_miss_changes() {
1522 let schema = Arc::new(Schema::new(vec![
1523 ColumnSchema::new(
1524 "timestamp",
1525 ConcreteDataType::timestamp_millisecond_datatype(),
1526 false,
1527 ),
1528 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
1529 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false),
1530 ]));
1531 let mut batches = Vec::new();
1532 let mut expected = BTreeMap::<String, Vec<(f64, PromSampleValue)>>::new();
1533 for layout in 0..3 {
1534 let labels = (0..1024)
1535 .map(|row| {
1536 let alternating = match layout {
1537 0 => true,
1538 1 => row < 64,
1539 _ => row >= 64,
1540 };
1541 if alternating && row % 2 != 0 {
1542 "b"
1543 } else {
1544 "a"
1545 }
1546 })
1547 .collect::<Vec<_>>();
1548 for (row, label) in labels.iter().enumerate() {
1549 let value = (layout * 1024 + row) as f64;
1550 expected
1551 .entry((*label).to_string())
1552 .or_default()
1553 .push((value, PromSampleValue::Number(value)));
1554 }
1555 let batch = RecordBatch::new(
1556 schema.clone(),
1557 vec![
1558 Arc::new(TimestampMillisecondVector::from_values(
1559 (0..1024).map(|row| (layout * 1024 + row) as i64 * 1000),
1560 )) as _,
1561 Arc::new(StringVector::from(
1562 labels.into_iter().map(Some).collect::<Vec<_>>(),
1563 )) as _,
1564 Arc::new(Float64Vector::from(
1565 (0..1024)
1566 .map(|row| Some((layout * 1024 + row) as f64))
1567 .collect::<Vec<_>>(),
1568 )) as _,
1569 ],
1570 )
1571 .unwrap();
1572 assert_eq!(
1575 prefer_label_runs(&[batch.column(1)], batch.num_rows()),
1576 layout == 1
1577 );
1578 batches.push(batch);
1579 }
1580 let response = PrometheusJsonResponse::record_batches_to_data(
1581 RecordBatches::try_new(schema, batches).unwrap(),
1582 None,
1583 ValueType::Matrix,
1584 )
1585 .unwrap();
1586 let PrometheusResponse::PromData(PromData {
1587 result: PromQueryResult::Matrix(series),
1588 ..
1589 }) = response
1590 else {
1591 panic!("expected matrix response");
1592 };
1593 assert_eq!(series.len(), expected.len());
1594 for series in series {
1595 assert_eq!(series.metric.len(), 1);
1596 assert!(series.histograms.is_empty());
1597 assert_eq!(
1598 series.values,
1599 expected.remove(&series.metric["host"]).unwrap()
1600 );
1601 }
1602 assert!(expected.is_empty());
1603 }
1604
1605 #[test]
1606 fn label_runs_keep_null_and_empty_labels_apart_across_batches() {
1607 let schema = Arc::new(Schema::new(vec![
1608 ColumnSchema::new(
1609 "timestamp",
1610 ConcreteDataType::timestamp_millisecond_datatype(),
1611 false,
1612 ),
1613 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
1614 ColumnSchema::new("rack", ConcreteDataType::string_datatype(), true),
1615 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), false),
1616 ]));
1617 let mut batches = Vec::new();
1621 let mut expected = vec![Vec::new(), Vec::new()];
1622 for (batch_index, leading_null) in [true, false].into_iter().enumerate() {
1623 let mut racks = Vec::new();
1624 let mut values = Vec::new();
1625 for row in 0..1024 {
1626 let value = (batch_index * 1024 + row) as f64;
1627 let null_rack = (row < 512) == leading_null;
1628 racks.push((!null_rack).then_some(""));
1629 values.push(Some(value));
1630 expected[usize::from(!null_rack)].push((value, PromSampleValue::Number(value)));
1631 }
1632 batches.push(
1633 RecordBatch::new(
1634 schema.clone(),
1635 vec![
1636 Arc::new(TimestampMillisecondVector::from_values(
1637 values.iter().map(|value| value.unwrap() as i64 * 1000),
1638 )) as _,
1639 Arc::new(StringVector::from(vec![Some("a"); 1024])) as _,
1640 Arc::new(StringVector::from(racks)) as _,
1641 Arc::new(Float64Vector::from(values)) as _,
1642 ],
1643 )
1644 .unwrap(),
1645 );
1646 }
1647 for batch in &batches {
1648 assert!(prefer_label_runs(
1649 &[batch.column(1), batch.column(2)],
1650 batch.num_rows()
1651 ));
1652 }
1653 for series in &mut expected {
1654 series.sort_by(|left, right| left.0.total_cmp(&right.0));
1655 }
1656
1657 let response = PrometheusJsonResponse::record_batches_to_data(
1658 RecordBatches::try_new(schema, batches).unwrap(),
1659 None,
1660 ValueType::Matrix,
1661 )
1662 .unwrap();
1663 let PrometheusResponse::PromData(PromData {
1664 result: PromQueryResult::Matrix(series),
1665 ..
1666 }) = response
1667 else {
1668 panic!("expected matrix response");
1669 };
1670 assert_eq!(series.len(), 2);
1671 assert_eq!(
1672 series[0].metric,
1673 BTreeMap::from([("host".into(), "a".into())])
1674 );
1675 assert_eq!(series[0].values, expected[0]);
1676 assert_eq!(
1677 series[1].metric,
1678 BTreeMap::from([("host".into(), "a".into()), ("rack".into(), "".into())])
1679 );
1680 assert_eq!(series[1].values, expected[1]);
1681 }
1682
1683 #[test]
1684 fn matrix_response_is_independent_of_input_row_order() {
1685 let schema = Arc::new(Schema::new(vec![
1689 ColumnSchema::new(
1690 "timestamp",
1691 ConcreteDataType::timestamp_millisecond_datatype(),
1692 false,
1693 ),
1694 ColumnSchema::new("host", ConcreteDataType::string_datatype(), true),
1695 ColumnSchema::new("rack", ConcreteDataType::string_datatype(), true),
1696 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
1697 ColumnSchema::new("histogram", native_histogram_value_type().clone(), true),
1698 ]));
1699 let histogram = |sum: f64| NativeHistogram {
1700 sum,
1701 ..sample_histogram()
1702 };
1703 type Row = (
1705 i64,
1706 Option<&'static str>,
1707 Option<&'static str>,
1708 Option<f64>,
1709 Option<NativeHistogram>,
1710 );
1711 let rows: Vec<Row> = vec![
1712 (1_000, Some("a"), Some("r"), Some(1.0), None),
1713 (3_000, Some("a"), Some("r"), Some(3.0), None),
1714 (2_000, Some("a"), Some("r"), Some(2.0), None),
1715 (5_000, Some("a"), None, Some(5.0), None),
1716 (4_000, Some("a"), None, Some(4.0), None),
1717 (7_000, Some(""), None, Some(7.0), None),
1718 (8_000, None, None, Some(8.0), None),
1719 (6_000, None, None, Some(6.0), None),
1720 (2_000, Some("h"), None, None, Some(histogram(20.0))),
1721 (1_000, Some("h"), None, None, Some(histogram(10.0))),
1722 ];
1723 let matrix = |order: &[usize], splits: &[usize]| {
1724 let batch = RecordBatch::new(
1725 schema.clone(),
1726 vec![
1727 Arc::new(TimestampMillisecondVector::from_vec(
1728 order.iter().map(|&row| rows[row].0).collect(),
1729 )) as _,
1730 Arc::new(StringVector::from(
1731 order.iter().map(|&row| rows[row].1).collect::<Vec<_>>(),
1732 )) as _,
1733 Arc::new(StringVector::from(
1734 order.iter().map(|&row| rows[row].2).collect::<Vec<_>>(),
1735 )) as _,
1736 Arc::new(Float64Vector::from(
1737 order.iter().map(|&row| rows[row].3).collect::<Vec<_>>(),
1738 )) as _,
1739 histogram_vector(
1740 &order
1741 .iter()
1742 .map(|&row| rows[row].4.clone())
1743 .collect::<Vec<_>>(),
1744 ),
1745 ],
1746 )
1747 .unwrap();
1748 let mut batches = Vec::new();
1749 let mut start = 0;
1750 for &end in splits.iter().chain(std::iter::once(&order.len())) {
1751 batches.push(batch.slice(start, end - start).unwrap());
1752 start = end;
1753 }
1754 let response = PrometheusJsonResponse::record_batches_to_data(
1755 RecordBatches::try_new(schema.clone(), batches).unwrap(),
1756 Some("metric".to_string()),
1757 ValueType::Matrix,
1758 )
1759 .unwrap();
1760 let PrometheusResponse::PromData(PromData {
1761 result: PromQueryResult::Matrix(series),
1762 ..
1763 }) = response
1764 else {
1765 panic!("expected matrix response");
1766 };
1767 series
1768 };
1769
1770 let clustered = matrix(&[0, 2, 1, 4, 3, 5, 7, 6, 9, 8], &[5]);
1771 let interleaved = matrix(&[8, 5, 1, 6, 0, 3, 9, 2, 7, 4], &[3, 6]);
1772 assert_eq!(
1773 serde_json::to_value(&interleaved).unwrap(),
1774 serde_json::to_value(&clustered).unwrap()
1775 );
1776
1777 assert_eq!(
1779 serde_json::to_value(&clustered[..4]).unwrap(),
1780 serde_json::json!([
1781 {"metric": {"__name__": "metric"}, "values": [[6.0, "6.0"], [8.0, "8.0"]]},
1782 {"metric": {"__name__": "metric", "host": ""}, "values": [[7.0, "7.0"]]},
1783 {"metric": {"__name__": "metric", "host": "a"},
1784 "values": [[4.0, "4.0"], [5.0, "5.0"]]},
1785 {"metric": {"__name__": "metric", "host": "a", "rack": "r"},
1786 "values": [[1.0, "1.0"], [2.0, "2.0"], [3.0, "3.0"]]},
1787 ])
1788 );
1789 assert_eq!(clustered[4].metric["host"], "h");
1790 assert_eq!(
1791 clustered[4]
1792 .histograms
1793 .iter()
1794 .map(|(timestamp, histogram)| (*timestamp, histogram.sum.as_str()))
1795 .collect::<Vec<_>>(),
1796 vec![(1.0, "10"), (2.0, "20")]
1797 );
1798 }
1799
1800 #[test]
1801 fn record_batches_to_data_preserves_mixed_float_and_histogram_rows() {
1802 let schema = Arc::new(Schema::new(vec![
1803 ColumnSchema::new(
1804 "timestamp",
1805 ConcreteDataType::timestamp_millisecond_datatype(),
1806 false,
1807 ),
1808 ColumnSchema::new("kind", ConcreteDataType::string_datatype(), false),
1809 ColumnSchema::new("float", ConcreteDataType::float64_datatype(), true),
1810 ColumnSchema::new("histogram", native_histogram_value_type().clone(), true),
1811 ]));
1812 let batch = RecordBatch::new(
1813 schema.clone(),
1814 vec![
1815 Arc::new(TimestampMillisecondVector::from_values([1_000, 1_000])) as _,
1816 Arc::new(StringVector::from(vec![Some("float"), Some("histogram")])) as _,
1817 Arc::new(Float64Vector::from(vec![Some(1.25), None])) as _,
1818 histogram_vector(&[None, Some(sample_histogram())]),
1819 ],
1820 )
1821 .unwrap();
1822 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1823
1824 let response = PrometheusJsonResponse::record_batches_to_data(
1825 batches,
1826 Some("mixed_metric".to_string()),
1827 ValueType::Vector,
1828 )
1829 .unwrap();
1830 let PrometheusResponse::PromData(PromData {
1831 result: PromQueryResult::Vector(series),
1832 ..
1833 }) = response
1834 else {
1835 panic!("expected vector response");
1836 };
1837
1838 assert_eq!(series.len(), 2);
1839 let float = series
1840 .iter()
1841 .find(|series| series.metric["kind"] == "float")
1842 .unwrap();
1843 assert_eq!(float.value, Some((1.0, "1.25".to_string())));
1844 assert!(float.histogram.is_none());
1845
1846 let histogram = series
1847 .iter()
1848 .find(|series| series.metric["kind"] == "histogram")
1849 .unwrap();
1850 assert!(histogram.value.is_none());
1851 let (timestamp, histogram) = histogram.histogram.as_ref().unwrap();
1852 assert_eq!(*timestamp, 1.0);
1853 assert_eq!(histogram.count, "2");
1854 assert_eq!(histogram.sum, "3");
1855 }
1856
1857 #[test]
1858 fn label_replace_with_utf8view_labels_does_not_panic() {
1859 let schema = Arc::new(Schema::new(vec![
1864 ColumnSchema::new(
1865 "timestamp",
1866 ConcreteDataType::timestamp_millisecond_datatype(),
1867 false,
1868 ),
1869 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
1870 ColumnSchema::new("host_copy", ConcreteDataType::utf8_view_datatype(), false),
1871 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
1872 ]));
1873 let batch = RecordBatch::new(
1874 schema.clone(),
1875 vec![
1876 Arc::new(TimestampMillisecondVector::from_values([1_000])) as _,
1877 Arc::new(StringVector::from(vec![Some("server-01")])) as _,
1878 Arc::new(StringVector::from(StringViewArray::from(vec![Some(
1879 "server-01",
1880 )]))) as _,
1881 Arc::new(Float64Vector::from(vec![Some(1.0)])) as _,
1882 ],
1883 )
1884 .unwrap();
1885 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1886
1887 let response = PrometheusJsonResponse::record_batches_to_data(
1888 batches,
1889 Some("label_replace_repro".to_string()),
1890 ValueType::Vector,
1891 )
1892 .unwrap();
1893 let PrometheusResponse::PromData(PromData {
1894 result: PromQueryResult::Vector(series),
1895 ..
1896 }) = response
1897 else {
1898 panic!("expected vector response");
1899 };
1900
1901 assert_eq!(series.len(), 1);
1902 assert_eq!(series[0].metric["__name__"], "label_replace_repro");
1903 assert_eq!(series[0].metric["host"], "server-01");
1904 assert_eq!(series[0].metric["host_copy"], "server-01");
1905 assert_eq!(series[0].value, Some((1.0, "1.0".to_string())));
1906 }
1907
1908 #[test]
1909 fn matrix_response_preserves_ordinary_histogram_nan_and_filters_stale_marker() {
1910 let schema = Arc::new(Schema::new(vec![
1911 ColumnSchema::new(
1912 "timestamp",
1913 ConcreteDataType::timestamp_millisecond_datatype(),
1914 false,
1915 ),
1916 ColumnSchema::new("histogram", native_histogram_value_type().clone(), true),
1917 ]));
1918 let mut ordinary_nan = sample_histogram();
1919 ordinary_nan.sum = f64::NAN;
1920 let mut stale = sample_histogram();
1921 stale.sum = f64::from_bits(PROMETHEUS_STALE_NAN_BITS);
1922 let batch = RecordBatch::new(
1923 schema.clone(),
1924 vec![
1925 Arc::new(TimestampMillisecondVector::from_values([1_000, 2_000])) as _,
1926 histogram_vector(&[Some(ordinary_nan), Some(stale)]),
1927 ],
1928 )
1929 .unwrap();
1930 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
1931
1932 let response =
1933 PrometheusJsonResponse::record_batches_to_data(batches, None, ValueType::Matrix)
1934 .unwrap();
1935 let PrometheusResponse::PromData(PromData {
1936 result: PromQueryResult::Matrix(series),
1937 ..
1938 }) = response
1939 else {
1940 panic!("expected matrix response");
1941 };
1942
1943 assert_eq!(series.len(), 1);
1944 assert_eq!(series[0].histograms.len(), 1);
1945 assert_eq!(series[0].histograms[0].0, 1.0);
1946 assert_eq!(series[0].histograms[0].1.sum, "NaN");
1947 }
1948
1949 #[test]
1950 fn native_histogram_json_closes_custom_bucket_zero() {
1951 let mut histogram = sample_histogram();
1952 histogram.schema = CUSTOM_BUCKETS_SCHEMA;
1953 histogram.zero_threshold = 0.0;
1954 histogram.custom_values = vec![1.0];
1955 histogram.positive_spans = vec![Span {
1956 offset: 0,
1957 length: 1,
1958 }];
1959 histogram.count = 1.0;
1960 histogram.sum = 0.0;
1961 histogram.zero_count = 0.0;
1962
1963 let json = serde_json::to_value(prometheus_native_histogram(&histogram).unwrap()).unwrap();
1964 assert_eq!(json["buckets"], serde_json::json!([[3, "-Inf", "1", "1"]]));
1965 }
1966
1967 #[test]
1968 fn native_histogram_json_preserves_terminal_finite_bucket() {
1969 let mut histogram = sample_histogram();
1970 histogram.positive_spans = vec![Span {
1971 offset: 1024,
1972 length: 2,
1973 }];
1974 histogram.positive_buckets = vec![1.0, 1.0];
1975 histogram.zero_count = 0.0;
1976
1977 let json = serde_json::to_value(prometheus_native_histogram(&histogram).unwrap()).unwrap();
1978 assert_eq!(
1979 json["buckets"],
1980 serde_json::json!([
1981 [0, 2.0_f64.powi(1023).to_string(), f64::MAX.to_string(), "1"],
1982 [0, f64::MAX.to_string(), "+Inf", "1"]
1983 ])
1984 );
1985 }
1986}