Skip to main content

servers/http/result/
prometheus_resp.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! prom supply the prometheus HTTP API Server compliance
16use 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/// A sample value of the Prometheus HTTP API JSON format.
72///
73/// Samples read out of a query result are kept as `f64` and formatted while the
74/// response is serialized, which avoids one `String` per sample. Samples parsed
75/// from a JSON body keep their original spelling, so a response that is
76/// deserialized and serialized again is unchanged.
77#[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
118/// Formats a sample value for the Prometheus HTTP API.
119///
120/// Finite sample strings use ryu's shortest-roundtrip representation and may
121/// differ textually from previous Rust/Prometheus wire formatting; values parse
122/// identically. Non-finite values retain Rust's `f64::to_string()` output.
123fn 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    // placeholder for header value
150    #[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    /// Adds collected PromQL warnings and infos to the response.
210    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    /// Merges data and annotations from another expanded PromQL query response.
219    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    /// Convert from `Result<Output>`
226    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        // Hold the collector while a streaming result is consumed.
233        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                // Prometheus won't report error if querying nonexist label and metric
272                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    /// Convert [RecordBatches] to [PromData]
291    fn record_batches_to_data(
292        batches: RecordBatches,
293        metric_name: Option<String>,
294        result_type: ValueType,
295    ) -> Result<PrometheusResponse> {
296        // Return empty result if no batches
297        if batches.iter().next().is_none() {
298            return Ok(empty_data(result_type));
299        }
300
301        let layout = ColumnLayout::infer(&batches.schema())?;
302
303        // Preserves the order of output tags.
304        // Tag order matters, e.g., after sorc and sort_desc, the output order must be kept.
305        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    /// Consume a streaming query result into [PromData].
315    ///
316    /// Every batch is merged into the per-series buffer and dropped before the
317    /// next one is polled, so the whole query result and the assembled series
318    /// never coexist in memory.
319    async fn consume_stream_to_data(
320        mut stream: SendableRecordBatchStream,
321        metric_name: Option<String>,
322        result_type: ValueType,
323    ) -> Result<PrometheusResponse> {
324        // An empty stream carries no batch to infer the column roles from, and
325        // reads as an empty result, like empty [RecordBatches].
326        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        // Preserves the order of output tags.
333        // Tag order matters, e.g., after sorc and sort_desc, the output order must be kept.
334        let mut buffer = IndexMap::<Vec<(String, String)>, PromSeriesSamples>::new();
335
336        merge_batch(&mut buffer, &layout, &first_batch, metric_name.as_deref())?;
337        // Release the batch before polling the next one.
338        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
348/// The empty data of a Prometheus HTTP API response of `result_type`.
349fn empty_data(result_type: ValueType) -> PrometheusResponse {
350    PrometheusResponse::PromData(PromData {
351        result_type: result_type.to_string(),
352        ..Default::default()
353    })
354}
355
356/// The role of every column of a query result.
357struct 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    /// Infer the role of every column from `schema`.
367    // TODO(ruihang): wish there is a better way to do this.
368    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                // Treat all value types as field
382                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
433/// The lookup key of a series in the per-series sample buffer.
434///
435/// The blanket `Equivalent` implementation compares through `Borrow`, which
436/// `Vec<(String, String)>` cannot provide for labels borrowed from a batch, so
437/// this wrapper compares the labels itself. It is built for a single lookup and
438/// hashes like the owned key it looks up.
439struct 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
454/// Merge the rows of one batch into the per-series sample buffer.
455fn merge_batch(
456    buffer: &mut IndexMap<Vec<(String, String)>, PromSeriesSamples>,
457    layout: &ColumnLayout,
458    batch: &RecordBatch,
459    metric_name: Option<&str>,
460) -> Result<()> {
461    // Scratch buffer of the labels of the series being read, borrowed from
462    // `batch` and `layout` and refilled for every series that is not in
463    // `buffer` yet, so looking a series up allocates nothing and only a series
464    // that is new allocates its labels.
465    let mut tags: Vec<(&str, &str)> = Vec::with_capacity(layout.tag_column_indices.len() + 1);
466
467    // prepare things...
468    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    // Read the labels once per run of rows that share them instead of
499    // once per row. `label_runs` finds every boundary, so the probe only
500    // decides whether looking for runs is worth its cost.
501    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    // assemble rows
508    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                    // retrieve tags
541                    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                    // Borrowed and owned labels hash the same: both
553                    // hash their `str` contents, so this hash serves the
554                    // lookup and the insertion of a new series below.
555                    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                        // A series that is new to the buffer takes
562                        // ownership of its labels.
563                        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
589/// Assemble the per-series samples of `buffer` into the data of a Prometheus
590/// HTTP API response of `result_type`.
591fn samples_to_data(
592    buffer: IndexMap<Vec<(String, String)>, PromSeriesSamples>,
593    result_type: ValueType,
594) -> Result<PrometheusResponse> {
595    // initialize result to return
596    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    // accumulate data into result
604    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                // sort values by timestamp
625                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                // TODO(ruihang): Not supported yet
650            }
651        }
652    });
653
654    // sort matrix by metric
655    // see: https://prometheus.io/docs/prometheus/3.5/querying/api/#range-vectors
656    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
669/// Ranges of consecutive rows whose label values are all equal.
670///
671/// `arrow::compute::partition` computes the same ranges, but its contract takes
672/// lexicographically sorted columns, which query output is not: range queries
673/// run without the plan's output sort, and `sort`/`topk` order by value.
674/// `distinct` is element-wise, so it holds for any row order, and its null
675/// handling is the one a series key needs: a null label and an empty one are
676/// distinct, and two nulls are not.
677fn 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
700/// Decides whether to group the rows of a batch into runs that share their
701/// label values.
702///
703/// Sampling adjacent row pairs keeps the decision independent of the batch
704/// size. Probe positions come from a xorshift sequence instead of a fixed
705/// stride, which would alias with periodic series layouts. A wrong guess only
706/// costs time: `partition` still validates every boundary, and the row-by-row
707/// path builds the labels of every row.
708fn 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        // No timestamp column on purpose: if the empty-stream short-circuit ran
769        // after `ColumnLayout::infer`, inference would fail with
770        // "no timestamp column found" instead of returning empty data.
771        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        // Two batches, with series "a" split across them: batch 0 holds its
801        // later samples, batch 1 its earlier ones plus series "b", so streaming
802        // must merge samples of one series across a batch boundary.
803        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        // Pin the merged content of the cross-batch series, not just the
867        // two-path agreement: both paths share merge_batch, so a regression
868        // there could make them agree on the same wrong output.
869        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        // The float series splits across the batch boundary so streaming must
909        // merge its samples like the eager path; the histogram row only arrives
910        // in the second batch.
911        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    /// A stream that lazily creates one batch per poll and asserts the previous
969    /// batch has been released before yielding the next one.
970    ///
971    /// Unlike `RecordBatches::as_stream`, this stream does not retain yielded
972    /// batches, so it can verify the consumer drops each batch before polling
973    /// the next one.
974    struct LazyDropCheckingStream {
975        schema: SchemaRef,
976        batch_index: usize,
977        /// Weak reference to the Arc wrapping the previously yielded batch.
978        /// The consumer must have dropped its clone by the next poll.
979        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            // The consumer must have released the previous batch before polling
1013            // for the next one.
1014            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        // Representative integral values use ryu's explicit .0 form.
1165        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        // These known shortest-roundtrip tie cases have different text but
1173        // remain numerically equivalent to Rust's representation.
1174        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        // Parsing a NaN does not preserve its payload bits, so NaN is checked
1194        // by its required semantic spelling rather than by to_bits().
1195        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        // Deserializing the expectation yields `PromSampleValue::Text`, so this
1283        // compares the deferred numeric encoding against eagerly built strings.
1284        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        // Keep this expected-value generation independent from the production
1385        // formatter while still asserting the Arrow/batch-to-Prometheus path.
1386        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        // NaN and infinities use Rust's `f64::to_string()` output.
1408        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        // Rows are clustered by series (a, b, a, a, b, c) and `a` is revisited
1453        // after `b`. The result must keep the first-occurrence order and
1454        // accumulate values per series.
1455        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        // Output order is first-occurrence order: a, b, c.
1507        assert_eq!(
1508            series
1509                .iter()
1510                .map(|series| series.metric["host"].as_str())
1511                .collect::<Vec<_>>(),
1512            vec!["a", "b", "c"]
1513        );
1514        // Vector results keep the last sample of each series.
1515        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            // The middle batch alternates only over a prefix the probes miss,
1573            // so it takes the run path even though most rows are not runs.
1574            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        // Two batches of two runs each, with only `rack` changing: a null label
1618        // and an empty one must not share a run, and the run opening the second
1619        // batch continues the series that ended the first one.
1620        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        // Range queries run without the plan's output sort, so this function sees
1686        // series interleaved across batches with timestamps out of order. The
1687        // serialized matrix must be the same either way.
1688        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        // timestamp, host, rack, float value, histogram value
1704        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        // Pin the canonical arrangement itself, not only its stability.
1778        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        // A PromQL `label_replace` query produces its new label through DataFusion's
1860        // `regexp_replace`, whose output materializes as a `Utf8View` array even when
1861        // the source label is a plain `Utf8`. Serializing such labels must not assume
1862        // the column is a `StringArray`.
1863        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}