Skip to main content

promql/functions/
native_histogram.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//! Native histogram PromQL helpers.
16
17use std::hash::{Hash, Hasher};
18use std::mem::size_of;
19use std::sync::Arc;
20
21use common_query::native_histogram::*;
22use common_query::prometheus::format_prometheus_float;
23use common_query::promql_annotations::PromqlAnnotationCollector;
24use datafusion::arrow::array::{
25    Array, ArrayRef, BooleanArray, Float64Array, Float64Builder, Int64Array, StringBuilder,
26    StructArray, TimestampMillisecondArray, UInt64Array,
27};
28use datafusion::arrow::compute::filter;
29use datafusion::arrow::datatypes::{DataType, Field, TimeUnit};
30use datafusion::common::{DataFusionError, Result as DfResult};
31use datafusion::logical_expr::{Accumulator as DfAccumulator, AggregateUDF, ScalarUDF, Volatility};
32use datafusion::physical_plan::ColumnarValue;
33use datafusion_common::ScalarValue;
34use datafusion_expr::function::AccumulatorArgs;
35use datafusion_expr::{ScalarFunctionArgs, ScalarUDFImpl, Signature, create_udaf, create_udf};
36
37use crate::functions::{
38    AvgOverTime, Deriv, DoubleExponentialSmoothing, IDelta, Increase, LastOverTime, MaxOverTime,
39    MinOverTime, PredictLinear, QuantileOverTime, Rate, StddevOverTime, StdvarOverTime,
40    SumOverTime, extract_array, extract_range_dict,
41};
42use crate::range_array::{RangeArray, unpack};
43
44fn extract_histogram_array(value: &ColumnarValue, func_name: &str) -> DfResult<ArrayRef> {
45    let array = extract_array(value)?;
46    if array.data_type() != &native_histogram_arrow_type() {
47        return Err(DataFusionError::Execution(format!(
48            "{func_name}: expected native histogram struct, found {}",
49            array.data_type()
50        )));
51    }
52    Ok(array)
53}
54
55fn read_scalar_f64_arg(
56    value: &ColumnarValue,
57    row: usize,
58    len: usize,
59    func_name: &str,
60) -> DfResult<f64> {
61    match value {
62        ColumnarValue::Scalar(ScalarValue::Float64(value)) => Ok(value.unwrap_or(f64::NAN)),
63        ColumnarValue::Array(array) => {
64            let array = array
65                .as_any()
66                .downcast_ref::<Float64Array>()
67                .ok_or_else(|| {
68                    DataFusionError::Execution(format!(
69                        "{func_name}: expected Float64 argument, found {}",
70                        array.data_type()
71                    ))
72                })?;
73            if array.len() != len {
74                return Err(DataFusionError::Execution(format!(
75                    "{func_name}: Float64 argument length mismatch: {} vs {len}",
76                    array.len()
77                )));
78            }
79            Ok(if array.is_null(row) {
80                f64::NAN
81            } else {
82                array.value(row)
83            })
84        }
85        other => Err(DataFusionError::Execution(format!(
86            "{func_name}: expected Float64 argument, found {}",
87            other.data_type()
88        ))),
89    }
90}
91
92#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
93enum AnnotationReturn {
94    FloatNull,
95    BooleanTrue,
96    BooleanFalse,
97}
98
99#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
100enum AnnotationLevel {
101    Info,
102    Warning,
103}
104
105impl AnnotationReturn {
106    fn data_type(self) -> DataType {
107        match self {
108            Self::FloatNull => DataType::Float64,
109            Self::BooleanTrue | Self::BooleanFalse => DataType::Boolean,
110        }
111    }
112
113    fn scalar_value(self) -> ScalarValue {
114        match self {
115            Self::FloatNull => ScalarValue::Float64(None),
116            Self::BooleanTrue => ScalarValue::Boolean(Some(true)),
117            Self::BooleanFalse => ScalarValue::Boolean(Some(false)),
118        }
119    }
120}
121
122#[derive(Debug, Clone)]
123struct NativeHistogramAnnotationUdf {
124    name: &'static str,
125    signature: Signature,
126    return_kind: AnnotationReturn,
127    level: AnnotationLevel,
128    message: String,
129    collector: Option<PromqlAnnotationCollector>,
130}
131
132impl NativeHistogramAnnotationUdf {
133    fn new(
134        name: &'static str,
135        return_kind: AnnotationReturn,
136        level: AnnotationLevel,
137        message: String,
138        collector: Option<PromqlAnnotationCollector>,
139    ) -> Self {
140        Self {
141            name,
142            signature: Signature::variadic_any(Volatility::Volatile),
143            return_kind,
144            level,
145            message,
146            collector,
147        }
148    }
149}
150
151impl PartialEq for NativeHistogramAnnotationUdf {
152    fn eq(&self, other: &Self) -> bool {
153        self.name == other.name
154            && self.return_kind == other.return_kind
155            && self.level == other.level
156            && self.message == other.message
157    }
158}
159
160impl Eq for NativeHistogramAnnotationUdf {}
161
162impl Hash for NativeHistogramAnnotationUdf {
163    fn hash<H: Hasher>(&self, state: &mut H) {
164        self.name.hash(state);
165        self.return_kind.hash(state);
166        self.level.hash(state);
167        self.message.hash(state);
168    }
169}
170
171impl ScalarUDFImpl for NativeHistogramAnnotationUdf {
172    fn name(&self) -> &str {
173        self.name
174    }
175
176    fn signature(&self) -> &Signature {
177        &self.signature
178    }
179
180    fn return_type(&self, _arg_types: &[DataType]) -> DfResult<DataType> {
181        Ok(self.return_kind.data_type())
182    }
183
184    fn invoke_with_args(&self, args: ScalarFunctionArgs) -> DfResult<ColumnarValue> {
185        let has_dropped_sample = !args.args.is_empty()
186            && (0..args.number_rows).any(|row| {
187                args.args.iter().all(|arg| match arg {
188                    ColumnarValue::Array(array) => array.is_valid(row),
189                    ColumnarValue::Scalar(value) => !value.is_null(),
190                })
191            });
192        if has_dropped_sample
193            && let Some(collector) = args
194                .config_options
195                .extensions
196                .get::<PromqlAnnotationCollector>()
197                .cloned()
198                .or_else(|| self.collector.clone())
199        {
200            match self.level {
201                AnnotationLevel::Info => collector.record_info(self.message.clone()),
202                AnnotationLevel::Warning => collector.record_warning(self.message.clone()),
203            }
204        }
205        Ok(ColumnarValue::Scalar(self.return_kind.scalar_value()))
206    }
207}
208
209pub struct NativeHistogramDrop;
210
211impl NativeHistogramDrop {
212    const fn float_null_name() -> &'static str {
213        "prom_native_histogram_drop_float"
214    }
215
216    const fn bool_false_name() -> &'static str {
217        "prom_native_histogram_drop_bool"
218    }
219
220    const fn bool_true_name() -> &'static str {
221        "prom_native_histogram_keep_bool"
222    }
223
224    pub fn float_null_udf(
225        message: String,
226        collector: Option<PromqlAnnotationCollector>,
227    ) -> ScalarUDF {
228        ScalarUDF::new_from_impl(NativeHistogramAnnotationUdf::new(
229            Self::float_null_name(),
230            AnnotationReturn::FloatNull,
231            AnnotationLevel::Info,
232            message,
233            collector,
234        ))
235    }
236
237    pub fn bool_false_udf(
238        message: String,
239        collector: Option<PromqlAnnotationCollector>,
240    ) -> ScalarUDF {
241        ScalarUDF::new_from_impl(NativeHistogramAnnotationUdf::new(
242            Self::bool_false_name(),
243            AnnotationReturn::BooleanFalse,
244            AnnotationLevel::Info,
245            message,
246            collector,
247        ))
248    }
249
250    pub fn bool_true_udf(
251        message: String,
252        collector: Option<PromqlAnnotationCollector>,
253    ) -> ScalarUDF {
254        ScalarUDF::new_from_impl(NativeHistogramAnnotationUdf::new(
255            Self::bool_true_name(),
256            AnnotationReturn::BooleanTrue,
257            AnnotationLevel::Info,
258            message,
259            collector,
260        ))
261    }
262
263    pub fn warning_bool_false_udf(
264        message: String,
265        collector: Option<PromqlAnnotationCollector>,
266    ) -> ScalarUDF {
267        ScalarUDF::new_from_impl(NativeHistogramAnnotationUdf::new(
268            Self::bool_false_name(),
269            AnnotationReturn::BooleanFalse,
270            AnnotationLevel::Warning,
271            message,
272            collector,
273        ))
274    }
275}
276
277fn record_info(collector: &Option<PromqlAnnotationCollector>, message: impl Into<String>) {
278    if let Some(collector) = collector {
279        collector.record_info(message);
280    }
281}
282
283fn record_warning(collector: &Option<PromqlAnnotationCollector>, message: impl Into<String>) {
284    if let Some(collector) = collector {
285        collector.record_warning(message);
286    }
287}
288
289fn record_custom_reconciliation(
290    collector: &Option<PromqlAnnotationCollector>,
291    name: &'static str,
292    lhs: &NativeHistogram,
293    rhs: &NativeHistogram,
294) {
295    if lhs.needs_custom_reconciliation(rhs) {
296        record_info(
297            collector,
298            format!("{name}: reconciled native histograms with different custom buckets"),
299        );
300    }
301}
302
303fn record_counter_reset_contradiction(
304    collector: &Option<PromqlAnnotationCollector>,
305    name: &'static str,
306    lhs: &NativeHistogram,
307    rhs: &NativeHistogram,
308) {
309    if lhs.counter_reset_hints_contradict(rhs) {
310        record_counter_reset_contradiction_warning(collector, name);
311    }
312}
313
314fn record_counter_reset_contradiction_warning(
315    collector: &Option<PromqlAnnotationCollector>,
316    name: &'static str,
317) {
318    record_warning(
319        collector,
320        format!("{name}: native histogram counter reset hints contradict"),
321    );
322}
323
324fn scalar_histogram_udf<F>(
325    name: &'static str,
326    extra_input_types: Vec<DataType>,
327    calc: F,
328) -> ScalarUDF
329where
330    F: Fn(&NativeHistogram, &[ColumnarValue], usize, usize, &'static str) -> DfResult<f64>
331        + Send
332        + Sync
333        + 'static,
334{
335    let mut input_types = vec![native_histogram_arrow_type()];
336    input_types.extend(extra_input_types);
337    create_udf(
338        name,
339        input_types,
340        DataType::Float64,
341        Volatility::Volatile,
342        Arc::new(move |input: &[ColumnarValue]| {
343            if input.is_empty() {
344                return Err(DataFusionError::Plan(format!(
345                    "{name} requires a native histogram argument"
346                )));
347            }
348            let histograms = extract_histogram_array(&input[0], name)?;
349            let histograms = histograms
350                .as_any()
351                .downcast_ref::<StructArray>()
352                .expect("validated native histogram struct");
353            let mut result = Float64Builder::with_capacity(histograms.len());
354            for row in 0..histograms.len() {
355                match read_histogram(histograms, row)? {
356                    Some(histogram) => {
357                        result.append_value(calc(&histogram, input, row, histograms.len(), name)?)
358                    }
359                    None => result.append_null(),
360                }
361            }
362            Ok(ColumnarValue::Array(Arc::new(result.finish())))
363        }) as _,
364    )
365}
366
367fn histogram_pair_udf(
368    name: &'static str,
369    op: fn(&NativeHistogram, &NativeHistogram) -> Option<NativeHistogram>,
370) -> ScalarUDF {
371    histogram_pair_udf_with_collector(name, op, None)
372}
373
374fn histogram_pair_udf_with_collector(
375    name: &'static str,
376    op: fn(&NativeHistogram, &NativeHistogram) -> Option<NativeHistogram>,
377    collector: Option<PromqlAnnotationCollector>,
378) -> ScalarUDF {
379    create_udf(
380        name,
381        vec![native_histogram_arrow_type(), native_histogram_arrow_type()],
382        native_histogram_arrow_type(),
383        Volatility::Volatile,
384        Arc::new(move |input: &[ColumnarValue]| {
385            let lhs = extract_histogram_array(&input[0], name)?;
386            let rhs = extract_histogram_array(&input[1], name)?;
387            if lhs.len() != rhs.len() {
388                return Err(DataFusionError::Execution(format!(
389                    "{name}: native histogram argument length mismatch: {} vs {}",
390                    lhs.len(),
391                    rhs.len()
392                )));
393            }
394
395            let lhs = lhs
396                .as_any()
397                .downcast_ref::<StructArray>()
398                .expect("validated native histogram struct");
399            let rhs = rhs
400                .as_any()
401                .downcast_ref::<StructArray>()
402                .expect("validated native histogram struct");
403            let mut result = Vec::with_capacity(lhs.len());
404            for row in 0..lhs.len() {
405                result.push(
406                    match (read_histogram(lhs, row)?, read_histogram(rhs, row)?) {
407                        (Some(lhs), Some(rhs)) => {
408                            record_custom_reconciliation(&collector, name, &lhs, &rhs);
409                            record_counter_reset_contradiction(&collector, name, &lhs, &rhs);
410                            let result = op(&lhs, &rhs);
411                            if result.is_none() {
412                                record_warning(
413                                    &collector,
414                                    format!(
415                                    "{name}: dropped native histogram sample with incompatible schemas"
416                                    ),
417                                );
418                            }
419                            result
420                        }
421                        _ => None,
422                    },
423                );
424            }
425            Ok(ColumnarValue::Array(build_histogram_array(&result)))
426        }) as _,
427    )
428}
429
430fn histogram_transform_udf(
431    name: &'static str,
432    op: fn(NativeHistogram) -> NativeHistogram,
433) -> ScalarUDF {
434    create_udf(
435        name,
436        vec![native_histogram_arrow_type()],
437        native_histogram_arrow_type(),
438        Volatility::Volatile,
439        Arc::new(move |input: &[ColumnarValue]| {
440            let histograms = extract_histogram_array(&input[0], name)?;
441            let histograms = histograms
442                .as_any()
443                .downcast_ref::<StructArray>()
444                .expect("validated native histogram struct");
445            let mut result = Vec::with_capacity(histograms.len());
446            for row in 0..histograms.len() {
447                result.push(read_histogram(histograms, row)?.map(op));
448            }
449            Ok(ColumnarValue::Array(build_histogram_array(&result)))
450        }) as _,
451    )
452}
453
454fn histogram_string_udf(name: &'static str) -> ScalarUDF {
455    create_udf(
456        name,
457        vec![native_histogram_arrow_type()],
458        DataType::Utf8,
459        Volatility::Volatile,
460        Arc::new(move |input: &[ColumnarValue]| {
461            let histograms = extract_histogram_array(&input[0], name)?;
462            let histograms = histograms
463                .as_any()
464                .downcast_ref::<StructArray>()
465                .expect("validated native histogram struct");
466            let mut result = StringBuilder::with_capacity(histograms.len(), histograms.len() * 32);
467            for row in 0..histograms.len() {
468                match read_histogram(histograms, row)? {
469                    Some(histogram) => result.append_value(histogram.promql_string()),
470                    None => result.append_null(),
471                }
472            }
473            Ok(ColumnarValue::Array(Arc::new(result.finish())))
474        }) as _,
475    )
476}
477
478fn histogram_scalar_udf(
479    name: &'static str,
480    input_types: Vec<DataType>,
481    histogram_index: usize,
482    scalar_index: usize,
483    op: fn(NativeHistogram, f64) -> Option<NativeHistogram>,
484) -> ScalarUDF {
485    create_udf(
486        name,
487        input_types,
488        native_histogram_arrow_type(),
489        Volatility::Volatile,
490        Arc::new(move |input: &[ColumnarValue]| {
491            let histograms = extract_histogram_array(&input[histogram_index], name)?;
492            let histograms = histograms
493                .as_any()
494                .downcast_ref::<StructArray>()
495                .expect("validated native histogram struct");
496            let mut result = Vec::with_capacity(histograms.len());
497            for row in 0..histograms.len() {
498                result.push(match read_histogram(histograms, row)? {
499                    Some(histogram) => {
500                        let scalar =
501                            read_scalar_f64_arg(&input[scalar_index], row, histograms.len(), name)?;
502                        op(histogram, scalar)
503                    }
504                    None => None,
505                });
506            }
507            Ok(ColumnarValue::Array(build_histogram_array(&result)))
508        }) as _,
509    )
510}
511
512fn histogram_compare_udf(
513    name: &'static str,
514    op: fn(&NativeHistogram, &NativeHistogram) -> bool,
515) -> ScalarUDF {
516    create_udf(
517        name,
518        vec![native_histogram_arrow_type(), native_histogram_arrow_type()],
519        DataType::Boolean,
520        Volatility::Volatile,
521        Arc::new(move |input: &[ColumnarValue]| {
522            let lhs = extract_histogram_array(&input[0], name)?;
523            let rhs = extract_histogram_array(&input[1], name)?;
524            if lhs.len() != rhs.len() {
525                return Err(DataFusionError::Execution(format!(
526                    "{name}: native histogram argument length mismatch: {} vs {}",
527                    lhs.len(),
528                    rhs.len()
529                )));
530            }
531
532            let lhs = lhs
533                .as_any()
534                .downcast_ref::<StructArray>()
535                .expect("validated native histogram struct");
536            let rhs = rhs
537                .as_any()
538                .downcast_ref::<StructArray>()
539                .expect("validated native histogram struct");
540            let mut result = Vec::with_capacity(lhs.len());
541            for row in 0..lhs.len() {
542                result.push(
543                    match (read_histogram(lhs, row)?, read_histogram(rhs, row)?) {
544                        (Some(lhs), Some(rhs)) => Some(op(&lhs, &rhs)),
545                        _ => None,
546                    },
547                );
548            }
549            Ok(ColumnarValue::Array(Arc::new(BooleanArray::from(result))))
550        }) as _,
551    )
552}
553
554pub struct NativeHistogramAdd;
555
556impl NativeHistogramAdd {
557    pub const fn name() -> &'static str {
558        "prom_native_histogram_add"
559    }
560
561    pub fn scalar_udf() -> ScalarUDF {
562        histogram_pair_udf(Self::name(), NativeHistogram::add)
563    }
564
565    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
566        histogram_pair_udf_with_collector(Self::name(), NativeHistogram::add, collector)
567    }
568}
569
570pub struct NativeHistogramSub;
571
572impl NativeHistogramSub {
573    pub const fn name() -> &'static str {
574        "prom_native_histogram_sub"
575    }
576
577    pub fn scalar_udf() -> ScalarUDF {
578        histogram_pair_udf(Self::name(), NativeHistogram::sub)
579    }
580
581    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
582        histogram_pair_udf_with_collector(Self::name(), NativeHistogram::sub, collector)
583    }
584}
585
586pub struct NativeHistogramMulScalar;
587
588impl NativeHistogramMulScalar {
589    pub const fn name() -> &'static str {
590        "prom_native_histogram_mul_scalar"
591    }
592
593    pub fn scalar_udf() -> ScalarUDF {
594        histogram_scalar_udf(
595            Self::name(),
596            vec![native_histogram_arrow_type(), DataType::Float64],
597            0,
598            1,
599            |histogram, scalar| Some(histogram.scale(scalar)),
600        )
601    }
602}
603
604pub struct NativeHistogramScalarMul;
605
606impl NativeHistogramScalarMul {
607    pub const fn name() -> &'static str {
608        "prom_native_histogram_scalar_mul"
609    }
610
611    pub fn scalar_udf() -> ScalarUDF {
612        histogram_scalar_udf(
613            Self::name(),
614            vec![DataType::Float64, native_histogram_arrow_type()],
615            1,
616            0,
617            |histogram, scalar| Some(histogram.scale(scalar)),
618        )
619    }
620}
621
622pub struct NativeHistogramDivScalar;
623
624impl NativeHistogramDivScalar {
625    pub const fn name() -> &'static str {
626        "prom_native_histogram_div_scalar"
627    }
628
629    pub fn scalar_udf() -> ScalarUDF {
630        histogram_scalar_udf(
631            Self::name(),
632            vec![native_histogram_arrow_type(), DataType::Float64],
633            0,
634            1,
635            |histogram, scalar| Some(histogram.divide_by(scalar)),
636        )
637    }
638}
639
640pub struct NativeHistogramNeg;
641
642impl NativeHistogramNeg {
643    pub const fn name() -> &'static str {
644        "prom_native_histogram_neg"
645    }
646
647    pub fn scalar_udf() -> ScalarUDF {
648        histogram_transform_udf(Self::name(), NativeHistogram::negated)
649    }
650}
651
652pub struct NativeHistogramEq;
653
654impl NativeHistogramEq {
655    pub const fn name() -> &'static str {
656        "prom_native_histogram_eq"
657    }
658
659    pub fn scalar_udf() -> ScalarUDF {
660        histogram_compare_udf(Self::name(), NativeHistogram::promql_eq)
661    }
662}
663
664pub struct NativeHistogramNotEq;
665
666impl NativeHistogramNotEq {
667    pub const fn name() -> &'static str {
668        "prom_native_histogram_not_eq"
669    }
670
671    pub fn scalar_udf() -> ScalarUDF {
672        histogram_compare_udf(Self::name(), |lhs, rhs| !lhs.promql_eq(rhs))
673    }
674}
675
676pub struct NativeHistogramCount;
677
678impl NativeHistogramCount {
679    pub const fn name() -> &'static str {
680        "prom_native_histogram_count"
681    }
682
683    pub fn scalar_udf() -> ScalarUDF {
684        scalar_histogram_udf(Self::name(), vec![], |histogram, _, _, _, _| {
685            Ok(histogram.count)
686        })
687    }
688}
689
690pub struct NativeHistogramSum;
691
692impl NativeHistogramSum {
693    pub const fn name() -> &'static str {
694        "prom_native_histogram_sum"
695    }
696
697    pub fn scalar_udf() -> ScalarUDF {
698        scalar_histogram_udf(Self::name(), vec![], |histogram, _, _, _, _| {
699            Ok(histogram.sum)
700        })
701    }
702}
703
704pub struct NativeHistogramAvg;
705
706impl NativeHistogramAvg {
707    pub const fn name() -> &'static str {
708        "prom_native_histogram_avg"
709    }
710
711    pub fn scalar_udf() -> ScalarUDF {
712        scalar_histogram_udf(Self::name(), vec![], |histogram, _, _, _, _| {
713            Ok(histogram.sum / histogram.count)
714        })
715    }
716}
717
718pub struct NativeHistogramStddev;
719
720impl NativeHistogramStddev {
721    pub const fn name() -> &'static str {
722        "prom_native_histogram_stddev"
723    }
724
725    pub fn scalar_udf() -> ScalarUDF {
726        scalar_histogram_udf(Self::name(), vec![], |histogram, _, _, _, _| {
727            Ok(histogram.estimated_stddev())
728        })
729    }
730}
731
732pub struct NativeHistogramStdvar;
733
734impl NativeHistogramStdvar {
735    pub const fn name() -> &'static str {
736        "prom_native_histogram_stdvar"
737    }
738
739    pub fn scalar_udf() -> ScalarUDF {
740        scalar_histogram_udf(Self::name(), vec![], |histogram, _, _, _, _| {
741            Ok(histogram.estimated_stdvar())
742        })
743    }
744}
745
746/// Formats float samples as PromQL label values.
747pub struct PromqlFloatToString;
748
749impl PromqlFloatToString {
750    pub const fn name() -> &'static str {
751        "prom_float_to_string"
752    }
753
754    pub fn scalar_udf() -> ScalarUDF {
755        create_udf(
756            Self::name(),
757            vec![DataType::Float64],
758            DataType::Utf8,
759            Volatility::Volatile,
760            Arc::new(|input: &[ColumnarValue]| {
761                let values = extract_array(&input[0])?;
762                let values = values
763                    .as_any()
764                    .downcast_ref::<Float64Array>()
765                    .expect("validated Float64 input");
766                let mut result = StringBuilder::new();
767                for value in values.iter() {
768                    match value {
769                        Some(value) => result.append_value(format_prometheus_float(value)),
770                        None => result.append_null(),
771                    }
772                }
773                Ok(ColumnarValue::Array(Arc::new(result.finish())))
774            }),
775        )
776    }
777}
778
779pub struct NativeHistogramToString;
780
781impl NativeHistogramToString {
782    pub const fn name() -> &'static str {
783        "prom_native_histogram_to_string"
784    }
785
786    pub fn scalar_udf() -> ScalarUDF {
787        histogram_string_udf(Self::name())
788    }
789}
790
791pub struct NativeHistogramQuantile;
792
793impl NativeHistogramQuantile {
794    pub const fn name() -> &'static str {
795        "prom_native_histogram_quantile"
796    }
797
798    pub fn scalar_udf() -> ScalarUDF {
799        Self::scalar_udf_with_collector(None)
800    }
801
802    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
803        scalar_histogram_udf(
804            Self::name(),
805            vec![DataType::Float64],
806            move |histogram, input, row, len, name| {
807                let q = read_scalar_f64_arg(&input[1], row, len, name)?;
808                let (value, info) = histogram.quantile_with_info(q);
809                if let Some(info) = info {
810                    let message = match info {
811                        NativeHistogramQuantileInfo::NaNSkew => {
812                            "input to histogram_quantile has NaN observations, result is skewed higher"
813                        }
814                        NativeHistogramQuantileInfo::NaNResult => {
815                            "input to histogram_quantile has NaN observations, result is NaN"
816                        }
817                    };
818                    record_info(&collector, message);
819                }
820                Ok(value)
821            },
822        )
823    }
824}
825
826pub struct NativeHistogramFraction;
827
828impl NativeHistogramFraction {
829    pub const fn name() -> &'static str {
830        "prom_native_histogram_fraction"
831    }
832
833    pub fn scalar_udf() -> ScalarUDF {
834        Self::scalar_udf_with_collector(None)
835    }
836
837    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
838        scalar_histogram_udf(
839            Self::name(),
840            vec![DataType::Float64, DataType::Float64],
841            move |histogram, input, row, len, name| {
842                let lower = read_scalar_f64_arg(&input[1], row, len, name)?;
843                let upper = read_scalar_f64_arg(&input[2], row, len, name)?;
844                let (value, excluded_nans) = histogram.fraction_with_info(lower, upper);
845                if excluded_nans {
846                    record_info(
847                        &collector,
848                        "input to histogram_fraction has NaN observations, which are excluded from all fractions",
849                    );
850                }
851                Ok(value)
852            },
853        )
854    }
855}
856
857#[derive(Debug, Clone, Copy)]
858enum NativeHistogramAggregateKind {
859    Sum,
860    Avg,
861}
862
863impl NativeHistogramAggregateKind {
864    const fn name(self) -> &'static str {
865        match self {
866            Self::Sum => NativeHistogramAggSum::name(),
867            Self::Avg => NativeHistogramAggAvg::name(),
868        }
869    }
870
871    const fn needs_count(self) -> bool {
872        matches!(self, Self::Avg)
873    }
874}
875
876pub struct NativeHistogramAggSum;
877
878impl NativeHistogramAggSum {
879    pub const fn name() -> &'static str {
880        "prom_native_histogram_agg_sum"
881    }
882
883    pub fn aggregate_udf() -> AggregateUDF {
884        native_histogram_aggregate_udf(NativeHistogramAggregateKind::Sum, None)
885    }
886
887    pub fn aggregate_udf_with_collector(
888        collector: Option<PromqlAnnotationCollector>,
889    ) -> AggregateUDF {
890        native_histogram_aggregate_udf(NativeHistogramAggregateKind::Sum, collector)
891    }
892}
893
894pub struct NativeHistogramAggAvg;
895
896impl NativeHistogramAggAvg {
897    pub const fn name() -> &'static str {
898        "prom_native_histogram_agg_avg"
899    }
900
901    pub fn aggregate_udf() -> AggregateUDF {
902        native_histogram_aggregate_udf(NativeHistogramAggregateKind::Avg, None)
903    }
904
905    pub fn aggregate_udf_with_collector(
906        collector: Option<PromqlAnnotationCollector>,
907    ) -> AggregateUDF {
908        native_histogram_aggregate_udf(NativeHistogramAggregateKind::Avg, collector)
909    }
910}
911
912#[derive(Debug)]
913struct NativeHistogramAggregateAccumulator {
914    kind: NativeHistogramAggregateKind,
915    value: Option<NativeHistogram>,
916    count: u64,
917    dropped_incompatible: bool,
918    counter_reset_seen: bool,
919    not_counter_reset_seen: bool,
920    collector: Option<PromqlAnnotationCollector>,
921}
922
923impl NativeHistogramAggregateAccumulator {
924    fn new(
925        kind: NativeHistogramAggregateKind,
926        collector: Option<PromqlAnnotationCollector>,
927    ) -> Self {
928        Self {
929            kind,
930            value: None,
931            count: 0,
932            dropped_incompatible: false,
933            counter_reset_seen: false,
934            not_counter_reset_seen: false,
935            collector,
936        }
937    }
938
939    fn from_args(
940        kind: NativeHistogramAggregateKind,
941        collector: Option<PromqlAnnotationCollector>,
942        _args: AccumulatorArgs,
943    ) -> DfResult<Box<dyn DfAccumulator>> {
944        Ok(Box::new(Self::new(kind, collector)))
945    }
946
947    fn observe_reset_hints(&mut self, counter_reset_seen: bool, not_counter_reset_seen: bool) {
948        self.counter_reset_seen |= counter_reset_seen;
949        self.not_counter_reset_seen |= not_counter_reset_seen;
950        if self.counter_reset_seen && self.not_counter_reset_seen {
951            record_counter_reset_contradiction_warning(&self.collector, self.kind.name());
952        }
953    }
954
955    fn push_histogram(&mut self, histogram: NativeHistogram, count: u64) -> DfResult<()> {
956        if self.kind.needs_count() && count == 0 {
957            return Ok(());
958        }
959
960        self.observe_reset_hints(
961            histogram.reset_hint == COUNTER_RESET_HINT,
962            histogram.reset_hint == NOT_COUNTER_RESET_HINT,
963        );
964        if self.dropped_incompatible {
965            return Ok(());
966        }
967        let combined_count = if self.kind.needs_count() {
968            self.count.checked_add(count).ok_or_else(|| {
969                DataFusionError::Execution(format!(
970                    "{}: native histogram sample count overflow",
971                    self.kind.name()
972                ))
973            })?
974        } else {
975            self.count
976        };
977        let value = match self.value.take() {
978            Some(value) => {
979                record_custom_reconciliation(&self.collector, self.kind.name(), &value, &histogram);
980                let combined = match self.kind {
981                    NativeHistogramAggregateKind::Sum => value.add(&histogram),
982                    NativeHistogramAggregateKind::Avg => {
983                        weighted_histogram_mean(value, self.count, histogram, count, combined_count)
984                    }
985                };
986                match combined {
987                    Some(value) => Some(value),
988                    None => {
989                        self.record_incompatible();
990                        None
991                    }
992                }
993            }
994            None => Some(histogram),
995        };
996        if !self.dropped_incompatible {
997            self.value = value;
998            self.count = combined_count;
999        }
1000        Ok(())
1001    }
1002
1003    fn mark_incompatible(&mut self) {
1004        self.value = None;
1005        self.count = 0;
1006        self.dropped_incompatible = true;
1007    }
1008
1009    fn record_incompatible(&mut self) {
1010        self.mark_incompatible();
1011        record_warning(
1012            &self.collector,
1013            format!(
1014                "{}: dropped native histogram aggregate with incompatible schemas",
1015                self.kind.name()
1016            ),
1017        );
1018    }
1019}
1020
1021fn weighted_histogram_mean(
1022    left: NativeHistogram,
1023    left_count: u64,
1024    right: NativeHistogram,
1025    right_count: u64,
1026    total_count: u64,
1027) -> Option<NativeHistogram> {
1028    let total_count = total_count as f64;
1029    left.scale(left_count as f64 / total_count)
1030        .add(&right.scale(right_count as f64 / total_count))
1031}
1032
1033fn range_fold_histograms(
1034    samples: Vec<NativeHistogram>,
1035    kind: NativeHistogramAggregateKind,
1036    name: &'static str,
1037    collector: &Option<PromqlAnnotationCollector>,
1038) -> Option<NativeHistogram> {
1039    if samples
1040        .iter()
1041        .any(|histogram| histogram.reset_hint == COUNTER_RESET_HINT)
1042        && samples
1043            .iter()
1044            .any(|histogram| histogram.reset_hint == NOT_COUNTER_RESET_HINT)
1045    {
1046        record_counter_reset_contradiction_warning(collector, name);
1047    }
1048
1049    let mut value = None;
1050    let mut count = 0u64;
1051    for histogram in samples {
1052        value = match value {
1053            Some(value) => {
1054                record_custom_reconciliation(collector, name, &value, &histogram);
1055                let next_count = count.checked_add(1)?;
1056                let combined = match kind {
1057                    NativeHistogramAggregateKind::Sum => value.add(&histogram),
1058                    NativeHistogramAggregateKind::Avg => {
1059                        weighted_histogram_mean(value, count, histogram, 1, next_count)
1060                    }
1061                };
1062                match combined {
1063                    Some(value) => Some(value),
1064                    None => {
1065                        record_warning(
1066                            collector,
1067                            format!(
1068                                "{name}: dropped native histogram range with incompatible schemas"
1069                            ),
1070                        );
1071                        return None;
1072                    }
1073                }
1074            }
1075            None => Some(histogram),
1076        };
1077        count = count.checked_add(1)?;
1078    }
1079    value
1080}
1081
1082#[derive(Debug, Clone, Copy)]
1083enum NativeHistogramRangeHistogramKind {
1084    Sum,
1085    Avg,
1086    Last,
1087}
1088
1089#[derive(Debug, Clone, Copy)]
1090enum NativeHistogramRangeFloatKind {
1091    Absent,
1092    Count,
1093    Present,
1094    Changes,
1095    Resets,
1096}
1097
1098fn collect_window_histograms(
1099    histograms: &StructArray,
1100    offset: usize,
1101    length: usize,
1102) -> DfResult<Option<Vec<NativeHistogram>>> {
1103    let mut samples = Vec::with_capacity(length);
1104    for row in offset..offset + length {
1105        let Some(histogram) = read_histogram(histograms, row)? else {
1106            return Ok(None);
1107        };
1108        samples.push(histogram);
1109    }
1110    Ok(Some(samples))
1111}
1112
1113fn native_histogram_range_histogram(
1114    input: &[ColumnarValue],
1115    kind: NativeHistogramRangeHistogramKind,
1116    func_name: &'static str,
1117    collector: Option<PromqlAnnotationCollector>,
1118) -> DfResult<ColumnarValue> {
1119    if input.len() != 2 {
1120        return Err(DataFusionError::Plan(format!(
1121            "{func_name} function should have 2 inputs"
1122        )));
1123    }
1124
1125    let ts_range = extract_range_dict(
1126        &input[0],
1127        func_name,
1128        "timestamp range vector",
1129        &DataType::Timestamp(TimeUnit::Millisecond, None),
1130    )?;
1131    let value_range = extract_range_dict(
1132        &input[1],
1133        func_name,
1134        "value range vector",
1135        &native_histogram_arrow_type(),
1136    )?;
1137    if ts_range.keys().values() != value_range.keys().values() {
1138        return Err(DataFusionError::Execution(format!(
1139            "{func_name}: timestamp and value ranges should have the same window layout"
1140        )));
1141    }
1142
1143    let histograms = value_range
1144        .values()
1145        .as_any()
1146        .downcast_ref::<StructArray>()
1147        .expect("validated native histogram range");
1148    let mut result = Vec::with_capacity(value_range.keys().len());
1149    for key in value_range.keys().values() {
1150        let (offset, length) = unpack(*key);
1151        let offset = offset as usize;
1152        let length = length as usize;
1153        if length == 0 {
1154            result.push(None);
1155            continue;
1156        }
1157        if matches!(kind, NativeHistogramRangeHistogramKind::Last) {
1158            let histogram = if (offset..offset + length).any(|row| histograms.is_null(row)) {
1159                None
1160            } else {
1161                read_histogram(histograms, offset + length - 1)?
1162            };
1163            result.push(histogram);
1164            continue;
1165        }
1166        let Some(samples) = collect_window_histograms(histograms, offset, length)? else {
1167            result.push(None);
1168            continue;
1169        };
1170        let histogram = match kind {
1171            NativeHistogramRangeHistogramKind::Sum => range_fold_histograms(
1172                samples,
1173                NativeHistogramAggregateKind::Sum,
1174                func_name,
1175                &collector,
1176            ),
1177            NativeHistogramRangeHistogramKind::Avg => range_fold_histograms(
1178                samples,
1179                NativeHistogramAggregateKind::Avg,
1180                func_name,
1181                &collector,
1182            ),
1183            NativeHistogramRangeHistogramKind::Last => samples.last().cloned(),
1184        };
1185        result.push(histogram);
1186    }
1187
1188    Ok(ColumnarValue::Array(build_histogram_array(&result)))
1189}
1190
1191fn native_histogram_range_float(
1192    input: &[ColumnarValue],
1193    kind: NativeHistogramRangeFloatKind,
1194    func_name: &'static str,
1195) -> DfResult<ColumnarValue> {
1196    if input.len() != 2 {
1197        return Err(DataFusionError::Plan(format!(
1198            "{func_name} function should have 2 inputs"
1199        )));
1200    }
1201
1202    let ts_range = extract_range_dict(
1203        &input[0],
1204        func_name,
1205        "timestamp range vector",
1206        &DataType::Timestamp(TimeUnit::Millisecond, None),
1207    )?;
1208    let value_range = extract_range_dict(
1209        &input[1],
1210        func_name,
1211        "value range vector",
1212        &native_histogram_arrow_type(),
1213    )?;
1214    if ts_range.keys().values() != value_range.keys().values() {
1215        return Err(DataFusionError::Execution(format!(
1216            "{func_name}: timestamp and value ranges should have the same window layout"
1217        )));
1218    }
1219
1220    let timestamps = ts_range
1221        .values()
1222        .as_any()
1223        .downcast_ref::<TimestampMillisecondArray>()
1224        .expect("validated timestamp range")
1225        .values();
1226    let histograms = value_range
1227        .values()
1228        .as_any()
1229        .downcast_ref::<StructArray>()
1230        .expect("validated native histogram range");
1231    let mut result = Float64Builder::with_capacity(value_range.keys().len());
1232    for key in value_range.keys().values() {
1233        let (offset, length) = unpack(*key);
1234        let offset = offset as usize;
1235        let length = length as usize;
1236        if length == 0 {
1237            match kind {
1238                NativeHistogramRangeFloatKind::Absent => result.append_value(1.0),
1239                _ => result.append_null(),
1240            }
1241            continue;
1242        }
1243        if matches!(kind, NativeHistogramRangeFloatKind::Absent) {
1244            result.append_null();
1245            continue;
1246        }
1247        if matches!(
1248            kind,
1249            NativeHistogramRangeFloatKind::Count | NativeHistogramRangeFloatKind::Present
1250        ) {
1251            if (offset..offset + length).any(|row| histograms.is_null(row)) {
1252                result.append_null();
1253            } else if matches!(kind, NativeHistogramRangeFloatKind::Count) {
1254                result.append_value(length as f64);
1255            } else {
1256                result.append_value(1.0);
1257            }
1258            continue;
1259        }
1260        let Some(samples) = collect_window_histograms(histograms, offset, length)? else {
1261            result.append_null();
1262            continue;
1263        };
1264        let value = match kind {
1265            NativeHistogramRangeFloatKind::Absent => {
1266                result.append_null();
1267                continue;
1268            }
1269            NativeHistogramRangeFloatKind::Count => length as f64,
1270            NativeHistogramRangeFloatKind::Present => 1.0,
1271            NativeHistogramRangeFloatKind::Changes => samples
1272                .windows(2)
1273                .filter(|pair| !pair[0].promql_eq(&pair[1]))
1274                .count() as f64,
1275            NativeHistogramRangeFloatKind::Resets => samples
1276                .windows(2)
1277                .zip(timestamps[offset..offset + length].windows(2))
1278                .filter(|(pair, ts_pair)| {
1279                    (pair[0].reset_hint == GAUGE_RESET_HINT)
1280                        != (pair[1].reset_hint == GAUGE_RESET_HINT)
1281                        || pair[1].detect_counter_reset(&pair[0], ts_pair[0], ts_pair[1])
1282                })
1283                .count() as f64,
1284        };
1285        result.append_value(value);
1286    }
1287
1288    Ok(ColumnarValue::Array(Arc::new(result.finish())))
1289}
1290
1291fn create_native_range_histogram_udf(
1292    name: &'static str,
1293    kind: NativeHistogramRangeHistogramKind,
1294    collector: Option<PromqlAnnotationCollector>,
1295) -> ScalarUDF {
1296    create_udf(
1297        name,
1298        vec![
1299            RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
1300            RangeArray::convert_data_type(native_histogram_arrow_type()),
1301        ],
1302        native_histogram_arrow_type(),
1303        Volatility::Volatile,
1304        Arc::new(move |input: &[ColumnarValue]| {
1305            native_histogram_range_histogram(input, kind, name, collector.clone())
1306        }) as _,
1307    )
1308}
1309
1310fn create_native_range_float_udf(
1311    name: &'static str,
1312    kind: NativeHistogramRangeFloatKind,
1313) -> ScalarUDF {
1314    create_udf(
1315        name,
1316        vec![
1317            RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
1318            RangeArray::convert_data_type(native_histogram_arrow_type()),
1319        ],
1320        DataType::Float64,
1321        Volatility::Volatile,
1322        Arc::new(move |input: &[ColumnarValue]| native_histogram_range_float(input, kind, name))
1323            as _,
1324    )
1325}
1326
1327pub struct NativeHistogramSumOverTime;
1328pub struct NativeHistogramAvgOverTime;
1329pub struct NativeHistogramAbsentOverTime;
1330pub struct NativeHistogramCountOverTime;
1331pub struct NativeHistogramLastOverTime;
1332pub struct NativeHistogramPresentOverTime;
1333pub struct NativeHistogramChanges;
1334pub struct NativeHistogramResets;
1335
1336impl NativeHistogramSumOverTime {
1337    pub const fn name() -> &'static str {
1338        "prom_native_histogram_sum_over_time"
1339    }
1340
1341    pub fn scalar_udf() -> ScalarUDF {
1342        Self::scalar_udf_with_collector(None)
1343    }
1344
1345    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
1346        create_native_range_histogram_udf(
1347            Self::name(),
1348            NativeHistogramRangeHistogramKind::Sum,
1349            collector,
1350        )
1351    }
1352}
1353
1354impl NativeHistogramAvgOverTime {
1355    pub const fn name() -> &'static str {
1356        "prom_native_histogram_avg_over_time"
1357    }
1358
1359    pub fn scalar_udf() -> ScalarUDF {
1360        Self::scalar_udf_with_collector(None)
1361    }
1362
1363    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
1364        create_native_range_histogram_udf(
1365            Self::name(),
1366            NativeHistogramRangeHistogramKind::Avg,
1367            collector,
1368        )
1369    }
1370}
1371
1372impl NativeHistogramAbsentOverTime {
1373    pub const fn name() -> &'static str {
1374        "prom_native_histogram_absent_over_time"
1375    }
1376
1377    pub fn scalar_udf() -> ScalarUDF {
1378        create_native_range_float_udf(Self::name(), NativeHistogramRangeFloatKind::Absent)
1379    }
1380}
1381
1382impl NativeHistogramCountOverTime {
1383    pub const fn name() -> &'static str {
1384        "prom_native_histogram_count_over_time"
1385    }
1386
1387    pub fn scalar_udf() -> ScalarUDF {
1388        create_native_range_float_udf(Self::name(), NativeHistogramRangeFloatKind::Count)
1389    }
1390}
1391
1392impl NativeHistogramLastOverTime {
1393    pub const fn name() -> &'static str {
1394        "prom_native_histogram_last_over_time"
1395    }
1396
1397    pub fn scalar_udf() -> ScalarUDF {
1398        create_native_range_histogram_udf(
1399            Self::name(),
1400            NativeHistogramRangeHistogramKind::Last,
1401            None,
1402        )
1403    }
1404}
1405
1406impl NativeHistogramPresentOverTime {
1407    pub const fn name() -> &'static str {
1408        "prom_native_histogram_present_over_time"
1409    }
1410
1411    pub fn scalar_udf() -> ScalarUDF {
1412        create_native_range_float_udf(Self::name(), NativeHistogramRangeFloatKind::Present)
1413    }
1414}
1415
1416impl NativeHistogramChanges {
1417    pub const fn name() -> &'static str {
1418        "prom_native_histogram_changes"
1419    }
1420
1421    pub fn scalar_udf() -> ScalarUDF {
1422        create_native_range_float_udf(Self::name(), NativeHistogramRangeFloatKind::Changes)
1423    }
1424}
1425
1426impl NativeHistogramResets {
1427    pub const fn name() -> &'static str {
1428        "prom_native_histogram_resets"
1429    }
1430
1431    pub fn scalar_udf() -> ScalarUDF {
1432        create_native_range_float_udf(Self::name(), NativeHistogramRangeFloatKind::Resets)
1433    }
1434}
1435
1436/// Coordinated float/native-histogram range evaluation.
1437///
1438/// The function name is passed as the first scalar argument so these two UDF names are enough for
1439/// distributed plan decoding. The remaining leading arguments are timestamp, float, and histogram
1440/// ranges, followed by the ordinary function arguments.
1441pub struct MixedRange;
1442
1443impl MixedRange {
1444    const fn float_name() -> &'static str {
1445        "prom_mixed_range_float"
1446    }
1447
1448    const fn histogram_name() -> &'static str {
1449        "prom_mixed_range_histogram"
1450    }
1451
1452    pub fn float_udf(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
1453        ScalarUDF::new_from_impl(MixedRangeUdf::new(MixedRangeOutput::Float, collector))
1454    }
1455
1456    pub fn histogram_udf(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
1457        ScalarUDF::new_from_impl(MixedRangeUdf::new(MixedRangeOutput::Histogram, collector))
1458    }
1459}
1460
1461#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1462enum MixedRangeOutput {
1463    Float,
1464    Histogram,
1465}
1466
1467impl MixedRangeOutput {
1468    fn name(self) -> &'static str {
1469        match self {
1470            Self::Float => MixedRange::float_name(),
1471            Self::Histogram => MixedRange::histogram_name(),
1472        }
1473    }
1474
1475    fn data_type(self) -> DataType {
1476        match self {
1477            Self::Float => DataType::Float64,
1478            Self::Histogram => native_histogram_arrow_type(),
1479        }
1480    }
1481}
1482
1483#[derive(Debug, Clone)]
1484struct MixedRangeUdf {
1485    output: MixedRangeOutput,
1486    signature: Signature,
1487    collector: Option<PromqlAnnotationCollector>,
1488}
1489
1490impl MixedRangeUdf {
1491    fn new(output: MixedRangeOutput, collector: Option<PromqlAnnotationCollector>) -> Self {
1492        Self {
1493            output,
1494            signature: Signature::variadic_any(Volatility::Volatile),
1495            collector,
1496        }
1497    }
1498}
1499
1500impl PartialEq for MixedRangeUdf {
1501    fn eq(&self, other: &Self) -> bool {
1502        self.output == other.output
1503    }
1504}
1505
1506impl Eq for MixedRangeUdf {}
1507
1508impl Hash for MixedRangeUdf {
1509    fn hash<H: Hasher>(&self, state: &mut H) {
1510        self.output.hash(state);
1511    }
1512}
1513
1514impl ScalarUDFImpl for MixedRangeUdf {
1515    fn name(&self) -> &str {
1516        self.output.name()
1517    }
1518
1519    fn signature(&self) -> &Signature {
1520        &self.signature
1521    }
1522
1523    fn return_type(&self, _arg_types: &[DataType]) -> DfResult<DataType> {
1524        Ok(self.output.data_type())
1525    }
1526
1527    fn invoke_with_args(&self, args: ScalarFunctionArgs) -> DfResult<ColumnarValue> {
1528        let collector = args
1529            .config_options
1530            .extensions
1531            .get::<PromqlAnnotationCollector>()
1532            .cloned()
1533            .or_else(|| self.collector.clone());
1534        mixed_range(&args, self.output, collector)
1535    }
1536}
1537
1538#[derive(Debug, Clone, Copy)]
1539enum MixedRangeFunction {
1540    Rate,
1541    Increase,
1542    // Raw-delta modes sum floats while preserving mixed-range drop/warning semantics.
1543    RawDeltaRate,
1544    RawDeltaIncrease,
1545    Delta,
1546    IDelta,
1547    IRate,
1548    Changes,
1549    Resets,
1550    AvgOverTime,
1551    MinOverTime,
1552    MaxOverTime,
1553    SumOverTime,
1554    CountOverTime,
1555    LastOverTime,
1556    AbsentOverTime,
1557    PresentOverTime,
1558    StddevOverTime,
1559    StdvarOverTime,
1560    QuantileOverTime,
1561    Deriv,
1562    PredictLinear,
1563    DoubleExponentialSmoothing,
1564    HoltWinters,
1565}
1566
1567impl MixedRangeFunction {
1568    fn parse(value: &ColumnarValue) -> DfResult<Self> {
1569        let name = match value {
1570            ColumnarValue::Scalar(ScalarValue::Utf8(Some(name)))
1571            | ColumnarValue::Scalar(ScalarValue::LargeUtf8(Some(name)))
1572            | ColumnarValue::Scalar(ScalarValue::Utf8View(Some(name))) => name.as_str(),
1573            other => {
1574                return Err(DataFusionError::Execution(format!(
1575                    "mixed range function name must be a non-null string scalar, found {}",
1576                    other.data_type()
1577                )));
1578            }
1579        };
1580
1581        match name {
1582            "rate" => Ok(Self::Rate),
1583            "increase" => Ok(Self::Increase),
1584            "raw_delta_rate" => Ok(Self::RawDeltaRate),
1585            "raw_delta_increase" => Ok(Self::RawDeltaIncrease),
1586            "delta" => Ok(Self::Delta),
1587            "idelta" => Ok(Self::IDelta),
1588            "irate" => Ok(Self::IRate),
1589            "changes" => Ok(Self::Changes),
1590            "resets" => Ok(Self::Resets),
1591            "avg_over_time" => Ok(Self::AvgOverTime),
1592            "min_over_time" => Ok(Self::MinOverTime),
1593            "max_over_time" => Ok(Self::MaxOverTime),
1594            "sum_over_time" => Ok(Self::SumOverTime),
1595            "count_over_time" => Ok(Self::CountOverTime),
1596            "last_over_time" => Ok(Self::LastOverTime),
1597            "absent_over_time" => Ok(Self::AbsentOverTime),
1598            "present_over_time" => Ok(Self::PresentOverTime),
1599            "stddev_over_time" => Ok(Self::StddevOverTime),
1600            "stdvar_over_time" => Ok(Self::StdvarOverTime),
1601            "quantile_over_time" => Ok(Self::QuantileOverTime),
1602            "deriv" => Ok(Self::Deriv),
1603            "predict_linear" => Ok(Self::PredictLinear),
1604            "double_exponential_smoothing" => Ok(Self::DoubleExponentialSmoothing),
1605            "holt_winters" => Ok(Self::HoltWinters),
1606            _ => Err(DataFusionError::Execution(format!(
1607                "unsupported mixed range function: {name}"
1608            ))),
1609        }
1610    }
1611
1612    fn name(self) -> &'static str {
1613        match self {
1614            Self::Rate => "rate",
1615            Self::Increase => "increase",
1616            Self::RawDeltaRate => "rate",
1617            Self::RawDeltaIncrease => "increase",
1618            Self::Delta => "delta",
1619            Self::IDelta => "idelta",
1620            Self::IRate => "irate",
1621            Self::Changes => "changes",
1622            Self::Resets => "resets",
1623            Self::AvgOverTime => "avg_over_time",
1624            Self::MinOverTime => "min_over_time",
1625            Self::MaxOverTime => "max_over_time",
1626            Self::SumOverTime => "sum_over_time",
1627            Self::CountOverTime => "count_over_time",
1628            Self::LastOverTime => "last_over_time",
1629            Self::AbsentOverTime => "absent_over_time",
1630            Self::PresentOverTime => "present_over_time",
1631            Self::StddevOverTime => "stddev_over_time",
1632            Self::StdvarOverTime => "stdvar_over_time",
1633            Self::QuantileOverTime => "quantile_over_time",
1634            Self::Deriv => "deriv",
1635            Self::PredictLinear => "predict_linear",
1636            Self::DoubleExponentialSmoothing => "double_exponential_smoothing",
1637            Self::HoltWinters => "holt_winters",
1638        }
1639    }
1640
1641    fn policy(self) -> MixedRangePolicy {
1642        match self {
1643            Self::Rate
1644            | Self::Increase
1645            | Self::RawDeltaRate
1646            | Self::RawDeltaIncrease
1647            | Self::Delta
1648            | Self::AvgOverTime
1649            | Self::SumOverTime => MixedRangePolicy::DropMixed,
1650            Self::IDelta | Self::IRate => MixedRangePolicy::LastTwo,
1651            Self::LastOverTime => MixedRangePolicy::Last,
1652            Self::Changes
1653            | Self::Resets
1654            | Self::CountOverTime
1655            | Self::AbsentOverTime
1656            | Self::PresentOverTime => MixedRangePolicy::Combined,
1657            Self::MinOverTime
1658            | Self::MaxOverTime
1659            | Self::StddevOverTime
1660            | Self::StdvarOverTime
1661            | Self::QuantileOverTime
1662            | Self::Deriv
1663            | Self::PredictLinear
1664            | Self::DoubleExponentialSmoothing
1665            | Self::HoltWinters => MixedRangePolicy::FloatOnly,
1666        }
1667    }
1668
1669    fn float_udf(self) -> Option<ScalarUDF> {
1670        match self {
1671            Self::Rate => Some(Rate::scalar_udf()),
1672            Self::Increase => Some(Increase::scalar_udf()),
1673            Self::RawDeltaRate | Self::RawDeltaIncrease => Some(SumOverTime::scalar_udf()),
1674            Self::Delta => Some(crate::functions::Delta::scalar_udf()),
1675            Self::IDelta => Some(IDelta::<false>::scalar_udf()),
1676            Self::IRate => Some(IDelta::<true>::scalar_udf()),
1677            Self::AvgOverTime => Some(AvgOverTime::scalar_udf()),
1678            Self::MinOverTime => Some(MinOverTime::scalar_udf()),
1679            Self::MaxOverTime => Some(MaxOverTime::scalar_udf()),
1680            Self::SumOverTime => Some(SumOverTime::scalar_udf()),
1681            Self::LastOverTime => Some(LastOverTime::scalar_udf()),
1682            Self::StddevOverTime => Some(StddevOverTime::scalar_udf()),
1683            Self::StdvarOverTime => Some(StdvarOverTime::scalar_udf()),
1684            Self::QuantileOverTime => Some(QuantileOverTime::scalar_udf()),
1685            Self::Deriv => Some(Deriv::scalar_udf()),
1686            Self::PredictLinear => Some(PredictLinear::scalar_udf()),
1687            Self::DoubleExponentialSmoothing | Self::HoltWinters => {
1688                Some(DoubleExponentialSmoothing::scalar_udf())
1689            }
1690            Self::Changes
1691            | Self::Resets
1692            | Self::CountOverTime
1693            | Self::AbsentOverTime
1694            | Self::PresentOverTime => None,
1695        }
1696    }
1697
1698    fn histogram_udf(self, collector: Option<PromqlAnnotationCollector>) -> Option<ScalarUDF> {
1699        match self {
1700            Self::Rate => Some(NativeHistogramRate::scalar_udf_with_collector(collector)),
1701            Self::Increase => Some(NativeHistogramIncrease::scalar_udf_with_collector(
1702                collector,
1703            )),
1704            Self::Delta => Some(NativeHistogramDelta::scalar_udf_with_collector(collector)),
1705            Self::IDelta => Some(NativeHistogramIDelta::scalar_udf_with_collector(collector)),
1706            Self::IRate => Some(NativeHistogramIRate::scalar_udf_with_collector(collector)),
1707            Self::AvgOverTime => Some(NativeHistogramAvgOverTime::scalar_udf_with_collector(
1708                collector,
1709            )),
1710            Self::SumOverTime => Some(NativeHistogramSumOverTime::scalar_udf_with_collector(
1711                collector,
1712            )),
1713            Self::LastOverTime => Some(NativeHistogramLastOverTime::scalar_udf()),
1714            Self::RawDeltaRate | Self::RawDeltaIncrease => None,
1715            _ => None,
1716        }
1717    }
1718}
1719
1720#[derive(Debug, Clone, Copy)]
1721enum MixedRangePolicy {
1722    DropMixed,
1723    LastTwo,
1724    Last,
1725    Combined,
1726    FloatOnly,
1727}
1728
1729#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1730enum SampleLane {
1731    Float,
1732    Histogram,
1733}
1734
1735fn mixed_range(
1736    args: &ScalarFunctionArgs,
1737    output: MixedRangeOutput,
1738    collector: Option<PromqlAnnotationCollector>,
1739) -> DfResult<ColumnarValue> {
1740    if args.args.len() < 4 {
1741        return Err(DataFusionError::Plan(format!(
1742            "{} function should have at least 4 inputs",
1743            output.name()
1744        )));
1745    }
1746    let function = MixedRangeFunction::parse(&args.args[0])?;
1747    let name = function.name();
1748    let ts_range = extract_range_dict(
1749        &args.args[1],
1750        name,
1751        "timestamp range vector",
1752        &DataType::Timestamp(TimeUnit::Millisecond, None),
1753    )?;
1754    let float_range = extract_range_dict(
1755        &args.args[2],
1756        name,
1757        "float range vector",
1758        &DataType::Float64,
1759    )?;
1760    let histogram_range = extract_range_dict(
1761        &args.args[3],
1762        name,
1763        "native histogram range vector",
1764        &native_histogram_arrow_type(),
1765    )?;
1766    let keys = ts_range.keys().values();
1767    if float_range.keys().values() != keys || histogram_range.keys().values() != keys {
1768        return Err(DataFusionError::Execution(format!(
1769            "{name}: timestamp, float, and native histogram ranges should have the same layout"
1770        )));
1771    }
1772    if args.number_rows != keys.len() {
1773        return Err(DataFusionError::Execution(format!(
1774            "{name}: range inputs have {} windows but the batch has {} rows",
1775            keys.len(),
1776            args.number_rows
1777        )));
1778    }
1779
1780    let timestamps = ts_range
1781        .values()
1782        .as_any()
1783        .downcast_ref::<TimestampMillisecondArray>()
1784        .expect("validated timestamp range");
1785    let floats = float_range
1786        .values()
1787        .as_any()
1788        .downcast_ref::<Float64Array>()
1789        .expect("validated float range");
1790    let histograms = histogram_range
1791        .values()
1792        .as_any()
1793        .downcast_ref::<StructArray>()
1794        .expect("validated native histogram range");
1795    if floats.len() != timestamps.len() || histograms.len() != timestamps.len() {
1796        return Err(DataFusionError::Execution(format!(
1797            "{name}: timestamp, float, and native histogram values should be row-aligned"
1798        )));
1799    }
1800
1801    let bounds = checked_window_bounds(keys, timestamps.len(), name)?;
1802    let float_valid = (0..floats.len())
1803        .map(|row| floats.is_valid(row))
1804        .collect::<Vec<_>>();
1805    let histogram_valid = (0..histograms.len())
1806        .map(|row| histograms.is_valid(row))
1807        .collect::<Vec<_>>();
1808    let float_prefix = validity_prefix(&float_valid, name)?;
1809    let histogram_prefix = validity_prefix(&histogram_valid, name)?;
1810
1811    if matches!(function.policy(), MixedRangePolicy::Combined) {
1812        if output != MixedRangeOutput::Float {
1813            return Err(DataFusionError::Execution(format!(
1814                "{name} does not return native histograms"
1815            )));
1816        }
1817        return combined_range_float(
1818            function,
1819            &bounds,
1820            timestamps,
1821            floats,
1822            histograms,
1823            &float_valid,
1824            &histogram_valid,
1825            &float_prefix,
1826            &histogram_prefix,
1827        );
1828    }
1829
1830    let selections = select_mixed_range_lanes(
1831        function,
1832        &bounds,
1833        timestamps,
1834        &float_valid,
1835        &histogram_valid,
1836        &float_prefix,
1837        &histogram_prefix,
1838        &collector,
1839    );
1840    let (lane, values, valid, prefix) = match output {
1841        MixedRangeOutput::Float => (
1842            SampleLane::Float,
1843            floats as &dyn Array,
1844            float_valid.as_slice(),
1845            float_prefix.as_slice(),
1846        ),
1847        MixedRangeOutput::Histogram => (
1848            SampleLane::Histogram,
1849            histograms as &dyn Array,
1850            histogram_valid.as_slice(),
1851            histogram_prefix.as_slice(),
1852        ),
1853    };
1854    let input = compact_lane_input(
1855        timestamps,
1856        values,
1857        valid,
1858        prefix,
1859        &bounds,
1860        &selections,
1861        lane,
1862        name,
1863    )?;
1864
1865    let udf = match output {
1866        MixedRangeOutput::Float => function.float_udf(),
1867        MixedRangeOutput::Histogram => function.histogram_udf(collector),
1868    }
1869    .ok_or_else(|| {
1870        DataFusionError::Execution(format!(
1871            "{name} does not return {} values",
1872            match output {
1873                MixedRangeOutput::Float => "float",
1874                MixedRangeOutput::Histogram => "native histogram",
1875            }
1876        ))
1877    })?;
1878    let mut input = input;
1879    input.extend_from_slice(&args.args[4..]);
1880    invoke_range_udf(udf, input, args, bounds.len())
1881}
1882
1883/// Decodes packed range keys into validated half-open `(offset, end)` bounds.
1884fn checked_window_bounds(
1885    keys: &[i64],
1886    value_len: usize,
1887    name: &str,
1888) -> DfResult<Vec<(usize, usize)>> {
1889    keys.iter()
1890        .map(|key| {
1891            let (offset, length) = unpack(*key);
1892            let offset = offset as usize;
1893            let end = offset
1894                .checked_add(length as usize)
1895                .filter(|end| *end <= value_len)
1896                .ok_or_else(|| {
1897                    DataFusionError::Execution(format!(
1898                        "{name}: invalid range ({offset}, {length}) for {value_len} values"
1899                    ))
1900                })?;
1901            Ok((offset, end))
1902        })
1903        .collect()
1904}
1905
1906/// Builds prefix counts where `prefix[i]` is the number of valid samples in `valid[..i]`.
1907fn validity_prefix(valid: &[bool], name: &str) -> DfResult<Vec<usize>> {
1908    let capacity = valid.len().checked_add(1).ok_or_else(|| {
1909        DataFusionError::Execution(format!("{name}: sample validity length overflow"))
1910    })?;
1911    let mut prefix = Vec::with_capacity(capacity);
1912    prefix.push(0usize);
1913    for is_valid in valid {
1914        prefix.push(
1915            prefix
1916                .last()
1917                .copied()
1918                .unwrap()
1919                .checked_add(usize::from(*is_valid))
1920                .ok_or_else(|| {
1921                    DataFusionError::Execution(format!("{name}: sample count overflow"))
1922                })?,
1923        );
1924    }
1925    Ok(prefix)
1926}
1927
1928/// Returns the number of valid samples in the half-open window `[offset, end)`.
1929fn valid_count(prefix: &[usize], offset: usize, end: usize) -> usize {
1930    prefix[end]
1931        .checked_sub(prefix[offset])
1932        .expect("validity prefix is monotonic")
1933}
1934
1935/// Selects the sample lane for each window and records policy-required annotations.
1936#[allow(clippy::too_many_arguments)]
1937fn select_mixed_range_lanes(
1938    function: MixedRangeFunction,
1939    bounds: &[(usize, usize)],
1940    timestamps: &TimestampMillisecondArray,
1941    float_valid: &[bool],
1942    histogram_valid: &[bool],
1943    float_prefix: &[usize],
1944    histogram_prefix: &[usize],
1945    collector: &Option<PromqlAnnotationCollector>,
1946) -> Vec<Option<SampleLane>> {
1947    bounds
1948        .iter()
1949        .map(|(offset, end)| {
1950            let float_count = valid_count(float_prefix, *offset, *end);
1951            let histogram_count = valid_count(histogram_prefix, *offset, *end);
1952            match function.policy() {
1953                MixedRangePolicy::DropMixed => match (float_count > 0, histogram_count > 0) {
1954                    (true, true) => {
1955                        record_warning(
1956                            collector,
1957                            format!(
1958                                "{}: encountered a mix of float and native histogram samples",
1959                                function.name()
1960                            ),
1961                        );
1962                        None
1963                    }
1964                    (true, false) => Some(SampleLane::Float),
1965                    (false, true) => Some(SampleLane::Histogram),
1966                    (false, false) => None,
1967                },
1968                MixedRangePolicy::FloatOnly => {
1969                    if float_count > 0 && histogram_count > 0 {
1970                        record_info(
1971                            collector,
1972                            format!(
1973                                "{}: ignored native histogram samples",
1974                                function.name()
1975                            ),
1976                        );
1977                    }
1978                    (float_count > 0).then_some(SampleLane::Float)
1979                }
1980                MixedRangePolicy::Last => (*offset..*end).rev().find_map(|row| {
1981                    if histogram_valid[row] {
1982                        Some(SampleLane::Histogram)
1983                    } else if float_valid[row] {
1984                        Some(SampleLane::Float)
1985                    } else {
1986                        None
1987                    }
1988                }),
1989                MixedRangePolicy::LastTwo => {
1990                    let mut last_two = Vec::with_capacity(2);
1991                    for row in (*offset..*end).rev() {
1992                        // Prometheus keeps a float as the newest sample if both lanes have the
1993                        // same timestamp, while the histogram becomes the preceding sample.
1994                        if float_valid[row] {
1995                            last_two.push((SampleLane::Float, timestamps.value(row)));
1996                        }
1997                        if last_two.len() < 2 && histogram_valid[row] {
1998                            last_two.push((SampleLane::Histogram, timestamps.value(row)));
1999                        }
2000                        if last_two.len() == 2 {
2001                            break;
2002                        }
2003                    }
2004                    if last_two.len() < 2 || last_two[0].1 == last_two[1].1 {
2005                        None
2006                    } else if last_two[0].0 == last_two[1].0 {
2007                        Some(last_two[0].0)
2008                    } else {
2009                        record_warning(
2010                            collector,
2011                            format!(
2012                                "{}: encountered a mix of float and native histogram samples in the last two points",
2013                                function.name()
2014                            ),
2015                        );
2016                        None
2017                    }
2018                }
2019                MixedRangePolicy::Combined => unreachable!(),
2020            }
2021        })
2022        .collect()
2023}
2024
2025/// Filters null placeholders from one sample lane and remaps its selected windows.
2026/// Windows assigned to the other lane become empty.
2027#[allow(clippy::too_many_arguments)]
2028fn compact_lane_input(
2029    timestamps: &TimestampMillisecondArray,
2030    values: &dyn Array,
2031    valid: &[bool],
2032    prefix: &[usize],
2033    bounds: &[(usize, usize)],
2034    selections: &[Option<SampleLane>],
2035    lane: SampleLane,
2036    name: &str,
2037) -> DfResult<Vec<ColumnarValue>> {
2038    let mask = BooleanArray::from(valid.to_vec());
2039    let filtered_timestamps = filter(timestamps, &mask)?;
2040    let filtered_values = filter(values, &mask)?;
2041    let ranges = bounds
2042        .iter()
2043        .zip(selections)
2044        .map(|((offset, end), selected)| {
2045            let compact_offset = prefix[*offset];
2046            let compact_length = if *selected == Some(lane) {
2047                valid_count(prefix, *offset, *end)
2048            } else {
2049                0
2050            };
2051            Ok((
2052                u32::try_from(compact_offset).map_err(|_| {
2053                    DataFusionError::Execution(format!(
2054                        "{name}: compacted range offset exceeds u32"
2055                    ))
2056                })?,
2057                u32::try_from(compact_length).map_err(|_| {
2058                    DataFusionError::Execution(format!(
2059                        "{name}: compacted range length exceeds u32"
2060                    ))
2061                })?,
2062            ))
2063        })
2064        .collect::<DfResult<Vec<_>>>()?;
2065    let timestamp_range = RangeArray::from_ranges(filtered_timestamps, ranges.clone())
2066        .map_err(DataFusionError::from)?;
2067    let value_range =
2068        RangeArray::from_ranges(filtered_values, ranges).map_err(DataFusionError::from)?;
2069    Ok(vec![
2070        ColumnarValue::Array(Arc::new(timestamp_range.into_dict())),
2071        ColumnarValue::Array(Arc::new(value_range.into_dict())),
2072    ])
2073}
2074
2075fn invoke_range_udf(
2076    udf: ScalarUDF,
2077    input: Vec<ColumnarValue>,
2078    outer_args: &ScalarFunctionArgs,
2079    number_rows: usize,
2080) -> DfResult<ColumnarValue> {
2081    let arg_fields = input
2082        .iter()
2083        .enumerate()
2084        .map(|(index, value)| Arc::new(Field::new(format!("arg_{index}"), value.data_type(), true)))
2085        .collect();
2086    udf.invoke_with_args(ScalarFunctionArgs {
2087        args: input,
2088        arg_fields,
2089        number_rows,
2090        return_field: outer_args.return_field.clone(),
2091        config_options: outer_args.config_options.clone(),
2092    })
2093}
2094
2095enum MixedSample {
2096    Float(f64),
2097    Histogram(NativeHistogram),
2098}
2099
2100#[allow(clippy::too_many_arguments)]
2101fn combined_range_float(
2102    function: MixedRangeFunction,
2103    bounds: &[(usize, usize)],
2104    timestamps: &TimestampMillisecondArray,
2105    floats: &Float64Array,
2106    histograms: &StructArray,
2107    float_valid: &[bool],
2108    histogram_valid: &[bool],
2109    float_prefix: &[usize],
2110    histogram_prefix: &[usize],
2111) -> DfResult<ColumnarValue> {
2112    let mut result = Float64Builder::with_capacity(bounds.len());
2113    for (offset, end) in bounds {
2114        let sample_count = valid_count(float_prefix, *offset, *end)
2115            .checked_add(valid_count(histogram_prefix, *offset, *end))
2116            .ok_or_else(|| {
2117                DataFusionError::Execution(format!("{}: sample count overflow", function.name()))
2118            })?;
2119        match function {
2120            MixedRangeFunction::CountOverTime => {
2121                if sample_count == 0 {
2122                    result.append_null();
2123                } else {
2124                    result.append_value(sample_count as f64);
2125                }
2126            }
2127            MixedRangeFunction::AbsentOverTime => {
2128                if sample_count == 0 {
2129                    result.append_value(1.0);
2130                } else {
2131                    result.append_null();
2132                }
2133            }
2134            MixedRangeFunction::PresentOverTime => {
2135                if sample_count == 0 {
2136                    result.append_null();
2137                } else {
2138                    result.append_value(1.0);
2139                }
2140            }
2141            MixedRangeFunction::Changes | MixedRangeFunction::Resets => {
2142                if sample_count == 0 {
2143                    result.append_null();
2144                    continue;
2145                }
2146                let mut count = 0usize;
2147                let mut previous = None;
2148                for row in *offset..*end {
2149                    if float_valid[row] {
2150                        count += mixed_transition(
2151                            function,
2152                            &mut previous,
2153                            timestamps.value(row),
2154                            MixedSample::Float(floats.value(row)),
2155                        );
2156                    }
2157                    if histogram_valid[row] {
2158                        let histogram = read_histogram(histograms, row)?
2159                            .expect("validated native histogram sample");
2160                        count += mixed_transition(
2161                            function,
2162                            &mut previous,
2163                            timestamps.value(row),
2164                            MixedSample::Histogram(histogram),
2165                        );
2166                    }
2167                }
2168                result.append_value(count as f64);
2169            }
2170            _ => {
2171                return Err(DataFusionError::Internal(format!(
2172                    "{} does not support combined range evaluation",
2173                    function.name()
2174                )));
2175            }
2176        }
2177    }
2178    Ok(ColumnarValue::Array(Arc::new(result.finish())))
2179}
2180
2181fn mixed_transition(
2182    function: MixedRangeFunction,
2183    previous: &mut Option<(i64, MixedSample)>,
2184    timestamp: i64,
2185    current: MixedSample,
2186) -> usize {
2187    let changed = previous.as_ref().is_some_and(|(previous_ts, previous)| {
2188        match (function, previous, &current) {
2189            (
2190                MixedRangeFunction::Changes,
2191                MixedSample::Float(previous),
2192                MixedSample::Float(current),
2193            ) => current != previous && !(current.is_nan() && previous.is_nan()),
2194            (
2195                MixedRangeFunction::Changes,
2196                MixedSample::Histogram(previous),
2197                MixedSample::Histogram(current),
2198            ) => !current.promql_eq(previous),
2199            (MixedRangeFunction::Changes, _, _) => true,
2200            (
2201                MixedRangeFunction::Resets,
2202                MixedSample::Float(previous),
2203                MixedSample::Float(current),
2204            ) => current < previous,
2205            (
2206                MixedRangeFunction::Resets,
2207                MixedSample::Histogram(previous),
2208                MixedSample::Histogram(current),
2209            ) => {
2210                (previous.reset_hint == GAUGE_RESET_HINT)
2211                    != (current.reset_hint == GAUGE_RESET_HINT)
2212                    || current.detect_counter_reset(previous, *previous_ts, timestamp)
2213            }
2214            (MixedRangeFunction::Resets, _, _) => true,
2215            _ => unreachable!(),
2216        }
2217    });
2218    *previous = Some((timestamp, current));
2219    usize::from(changed)
2220}
2221
2222fn native_histogram_scalar(histogram: Option<NativeHistogram>) -> ScalarValue {
2223    let array = build_histogram_array(&[histogram]);
2224    let histogram = array
2225        .as_any()
2226        .downcast_ref::<StructArray>()
2227        .expect("native histogram array is a StructArray")
2228        .clone();
2229    ScalarValue::Struct(Arc::new(histogram))
2230}
2231
2232fn native_histogram_aggregate_udf(
2233    kind: NativeHistogramAggregateKind,
2234    collector: Option<PromqlAnnotationCollector>,
2235) -> AggregateUDF {
2236    let state_types = if kind.needs_count() {
2237        vec![
2238            native_histogram_arrow_type(),
2239            DataType::UInt64,
2240            DataType::Boolean,
2241            DataType::Boolean,
2242            DataType::Boolean,
2243        ]
2244    } else {
2245        vec![
2246            native_histogram_arrow_type(),
2247            DataType::Boolean,
2248            DataType::Boolean,
2249            DataType::Boolean,
2250        ]
2251    };
2252
2253    create_udaf(
2254        kind.name(),
2255        vec![native_histogram_arrow_type()],
2256        Arc::new(native_histogram_arrow_type()),
2257        Volatility::Volatile,
2258        Arc::new(move |args| {
2259            NativeHistogramAggregateAccumulator::from_args(kind, collector.clone(), args)
2260        }),
2261        Arc::new(state_types),
2262    )
2263}
2264
2265impl DfAccumulator for NativeHistogramAggregateAccumulator {
2266    fn update_batch(&mut self, values: &[ArrayRef]) -> DfResult<()> {
2267        let histograms = values
2268            .first()
2269            .and_then(|array| array.as_any().downcast_ref::<StructArray>())
2270            .ok_or_else(|| {
2271                DataFusionError::Execution(format!(
2272                    "{}: expected native histogram struct input",
2273                    self.kind.name()
2274                ))
2275            })?;
2276
2277        for row in 0..histograms.len() {
2278            let Some(histogram) = read_histogram(histograms, row)? else {
2279                continue;
2280            };
2281            self.push_histogram(histogram, 1)?;
2282        }
2283
2284        Ok(())
2285    }
2286
2287    fn evaluate(&mut self) -> DfResult<ScalarValue> {
2288        let histogram = match (self.kind, self.dropped_incompatible, self.value.clone()) {
2289            (_, true, _) => None,
2290            (_, false, None) => None,
2291            (NativeHistogramAggregateKind::Sum, false, value) => value,
2292            (NativeHistogramAggregateKind::Avg, false, Some(value)) if self.count > 0 => {
2293                Some(value)
2294            }
2295            (NativeHistogramAggregateKind::Avg, _, _) => None,
2296        };
2297
2298        Ok(native_histogram_scalar(histogram))
2299    }
2300
2301    fn size(&self) -> usize {
2302        size_of::<Self>()
2303            + self.value.as_ref().map_or(0, |histogram| {
2304                histogram.custom_values.capacity() * size_of::<f64>()
2305                    + histogram.positive_spans.capacity() * size_of::<Span>()
2306                    + histogram.negative_spans.capacity() * size_of::<Span>()
2307                    + histogram.positive_buckets.capacity() * size_of::<f64>()
2308                    + histogram.negative_buckets.capacity() * size_of::<f64>()
2309            })
2310    }
2311
2312    fn state(&mut self) -> DfResult<Vec<ScalarValue>> {
2313        let mut state = vec![native_histogram_scalar(self.value.clone())];
2314        if self.kind.needs_count() {
2315            state.push(ScalarValue::UInt64(Some(self.count)));
2316        }
2317        state.push(ScalarValue::Boolean(Some(self.dropped_incompatible)));
2318        state.push(ScalarValue::Boolean(Some(self.counter_reset_seen)));
2319        state.push(ScalarValue::Boolean(Some(self.not_counter_reset_seen)));
2320        Ok(state)
2321    }
2322
2323    fn merge_batch(&mut self, states: &[ArrayRef]) -> DfResult<()> {
2324        if states.is_empty() {
2325            return Ok(());
2326        }
2327
2328        let histograms = states[0]
2329            .as_any()
2330            .downcast_ref::<StructArray>()
2331            .ok_or_else(|| {
2332                DataFusionError::Execution(format!(
2333                    "{}: expected native histogram struct state",
2334                    self.kind.name()
2335                ))
2336            })?;
2337        let counts = if self.kind.needs_count() {
2338            Some(
2339                states
2340                    .get(1)
2341                    .and_then(|array| array.as_any().downcast_ref::<UInt64Array>())
2342                    .ok_or_else(|| {
2343                        DataFusionError::Execution(format!(
2344                            "{}: expected UInt64 count state",
2345                            self.kind.name()
2346                        ))
2347                    })?,
2348            )
2349        } else {
2350            None
2351        };
2352        let dropped_index = if self.kind.needs_count() { 2 } else { 1 };
2353        let dropped = states
2354            .get(dropped_index)
2355            .and_then(|array| array.as_any().downcast_ref::<BooleanArray>())
2356            .ok_or_else(|| {
2357                DataFusionError::Execution(format!(
2358                    "{}: expected Boolean dropped state",
2359                    self.kind.name()
2360                ))
2361            })?;
2362        let counter_reset_seen = states
2363            .get(dropped_index + 1)
2364            .and_then(|array| array.as_any().downcast_ref::<BooleanArray>())
2365            .ok_or_else(|| {
2366                DataFusionError::Execution(format!(
2367                    "{}: expected Boolean counter reset state",
2368                    self.kind.name()
2369                ))
2370            })?;
2371        let not_counter_reset_seen = states
2372            .get(dropped_index + 2)
2373            .and_then(|array| array.as_any().downcast_ref::<BooleanArray>())
2374            .ok_or_else(|| {
2375                DataFusionError::Execution(format!(
2376                    "{}: expected Boolean not-counter-reset state",
2377                    self.kind.name()
2378                ))
2379            })?;
2380
2381        for row in 0..histograms.len() {
2382            self.observe_reset_hints(
2383                counter_reset_seen.value(row),
2384                not_counter_reset_seen.value(row),
2385            );
2386            if dropped.value(row) {
2387                self.mark_incompatible();
2388            }
2389            if self.dropped_incompatible {
2390                continue;
2391            }
2392            let Some(histogram) = read_histogram(histograms, row)? else {
2393                continue;
2394            };
2395            let count = counts.map(|counts| counts.value(row)).unwrap_or(1);
2396            self.push_histogram(histogram, count)?;
2397        }
2398
2399        Ok(())
2400    }
2401}
2402
2403fn histogram_delta(
2404    samples: &[NativeHistogram],
2405    timestamps: &[i64],
2406    is_counter: bool,
2407) -> Option<NativeHistogram> {
2408    if samples.len() < 2 || samples.len() != timestamps.len() {
2409        return None;
2410    }
2411
2412    if !is_counter {
2413        return samples
2414            .last()?
2415            .sub(samples.first()?)
2416            .map(NativeHistogram::into_gauge);
2417    }
2418
2419    let first_reset = samples[1].detect_counter_reset(&samples[0], timestamps[0], timestamps[1]);
2420    let (initial, reset_scan_start) = if first_reset {
2421        // The first sample is irrelevant after an immediate reset. Adopt the
2422        // second sample's layout so an incompatible pre-reset layout is ignored.
2423        (samples[1].zero_like(), 2)
2424    } else {
2425        (samples[0].clone(), 1)
2426    };
2427    let mut result = samples.last()?.sub(&initial)?;
2428    for index in reset_scan_start..samples.len() {
2429        if samples[index].detect_counter_reset(
2430            &samples[index - 1],
2431            timestamps[index - 1],
2432            timestamps[index],
2433        ) {
2434            result = result.add(&samples[index - 1])?;
2435        }
2436    }
2437    Some(result.into_gauge())
2438}
2439
2440fn idelta_value(
2441    samples: &[NativeHistogram],
2442    is_rate: bool,
2443    previous_ts: i64,
2444    current_ts: i64,
2445    sampled_interval_secs: f64,
2446) -> Option<NativeHistogram> {
2447    if samples.len() < 2 {
2448        return None;
2449    }
2450    let previous = &samples[samples.len() - 2];
2451    let current = samples.last()?;
2452    let result = if is_rate && current.detect_counter_reset(previous, previous_ts, current_ts) {
2453        current.clone()
2454    } else {
2455        current.sub(previous)?
2456    };
2457    Some(
2458        if is_rate {
2459            result.scale(1.0 / sampled_interval_secs)
2460        } else {
2461            result
2462        }
2463        .into_gauge(),
2464    )
2465}
2466
2467fn native_extrapolated_rate<const IS_COUNTER: bool, const IS_RATE: bool>(
2468    input: &[ColumnarValue],
2469    range_length: i64,
2470    func_name: &'static str,
2471    collector: Option<PromqlAnnotationCollector>,
2472) -> DfResult<ColumnarValue> {
2473    if input.len() != 4 {
2474        return Err(DataFusionError::Plan(format!(
2475            "{func_name} function should have 4 inputs"
2476        )));
2477    }
2478
2479    let ts_dict = extract_range_dict(
2480        &input[0],
2481        func_name,
2482        "timestamp range vector",
2483        &DataType::Timestamp(TimeUnit::Millisecond, None),
2484    )?;
2485    let value_dict = extract_range_dict(
2486        &input[1],
2487        func_name,
2488        "value range vector",
2489        &native_histogram_arrow_type(),
2490    )?;
2491    let eval_ts = extract_array(&input[2])?;
2492    let eval_ts = eval_ts
2493        .as_any()
2494        .downcast_ref::<TimestampMillisecondArray>()
2495        .ok_or_else(|| {
2496            DataFusionError::Execution(format!(
2497                "{func_name}: expect evaluation timestamp vector as Timestamp(Millisecond), found {}",
2498                eval_ts.data_type()
2499            ))
2500        })?;
2501
2502    let keys = ts_dict.keys().values();
2503    if value_dict.keys().values() != keys || eval_ts.len() != keys.len() {
2504        return Err(DataFusionError::Execution(format!(
2505            "{func_name}: timestamp, value, and evaluation ranges should have the same layout"
2506        )));
2507    }
2508
2509    let all_timestamps = ts_dict
2510        .values()
2511        .as_any()
2512        .downcast_ref::<TimestampMillisecondArray>()
2513        .expect("validated timestamp range")
2514        .values();
2515    let all_histograms = value_dict
2516        .values()
2517        .as_any()
2518        .downcast_ref::<StructArray>()
2519        .expect("validated native histogram range");
2520    let range_length_secs = range_length as f64 / 1000.0;
2521    let mut result = Vec::with_capacity(keys.len());
2522
2523    for index in 0..keys.len() {
2524        let (raw_offset, raw_length) = unpack(keys[index]);
2525        let offset = raw_offset as usize;
2526        let length = raw_length as usize;
2527        if length == 0 {
2528            result.push(None);
2529            continue;
2530        }
2531
2532        let mut samples = Vec::with_capacity(length);
2533        let mut has_null = false;
2534        for row in offset..offset + length {
2535            let Some(histogram) = read_histogram(all_histograms, row)? else {
2536                has_null = true;
2537                break;
2538            };
2539            samples.push(histogram);
2540        }
2541        if has_null {
2542            result.push(None);
2543            continue;
2544        }
2545
2546        let first_ts = all_timestamps[offset];
2547        let last_ts = all_timestamps[offset + length - 1];
2548        let range_end = eval_ts.value(index);
2549        let range_start = range_end - range_length;
2550        let synthetic_zero_timestamp = IS_COUNTER
2551            .then_some(samples[0].start_timestamp)
2552            .flatten()
2553            .filter(|start| *start != 0 && range_start < *start && *start < first_ts);
2554        let synthetic_zero_start = synthetic_zero_timestamp.is_some();
2555        if length < 2 && !synthetic_zero_start {
2556            result.push(None);
2557            continue;
2558        }
2559
2560        let wrong_flavor = if IS_COUNTER {
2561            samples
2562                .iter()
2563                .any(|histogram| histogram.reset_hint == GAUGE_RESET_HINT)
2564        } else {
2565            samples[0].reset_hint != GAUGE_RESET_HINT
2566                || samples[samples.len() - 1].reset_hint != GAUGE_RESET_HINT
2567        };
2568        if wrong_flavor {
2569            let expected = if IS_COUNTER { "counter" } else { "gauge" };
2570            record_warning(
2571                &collector,
2572                format!("{func_name}: native histogram input should be a {expected} histogram"),
2573            );
2574        }
2575
2576        let timestamps = &all_timestamps[offset..offset + length];
2577        for pair in samples.windows(2) {
2578            record_custom_reconciliation(&collector, func_name, &pair[0], &pair[1]);
2579        }
2580        let mut histogram = if length == 1 {
2581            samples[0].clone()
2582        } else if let Some(histogram) = histogram_delta(&samples, timestamps, IS_COUNTER) {
2583            histogram
2584        } else {
2585            record_warning(
2586                &collector,
2587                format!("{func_name}: dropped native histogram range with incompatible schemas"),
2588            );
2589            result.push(None);
2590            continue;
2591        };
2592        if synthetic_zero_start && length > 1 {
2593            let Some(with_synthetic_zero) = histogram.add(&samples[0]) else {
2594                record_warning(
2595                    &collector,
2596                    format!(
2597                        "{func_name}: dropped native histogram range with incompatible schemas"
2598                    ),
2599                );
2600                result.push(None);
2601                continue;
2602            };
2603            histogram = with_synthetic_zero;
2604        }
2605
2606        let real_sampled_interval_ms = (last_ts - first_ts) as f64;
2607        let sampled_interval_ms = synthetic_zero_timestamp
2608            .map(|start| (last_ts - start) as f64)
2609            .unwrap_or(real_sampled_interval_ms);
2610        if sampled_interval_ms <= 0.0 {
2611            result.push(None);
2612            continue;
2613        }
2614        let average_interval_ms = if length > 1 {
2615            real_sampled_interval_ms / (length - 1) as f64
2616        } else {
2617            0.0
2618        };
2619        let mut duration_to_start_ms = if synthetic_zero_start {
2620            0.0
2621        } else {
2622            (first_ts - range_start) as f64
2623        };
2624        let duration_to_end_ms = (range_end - last_ts) as f64;
2625
2626        if IS_COUNTER && !synthetic_zero_start && histogram.count > 0.0 && samples[0].count >= 0.0 {
2627            let duration_to_zero = sampled_interval_ms * (samples[0].count / histogram.count);
2628            if duration_to_zero < duration_to_start_ms {
2629                duration_to_start_ms = duration_to_zero;
2630            }
2631        }
2632
2633        let extrapolation_threshold = average_interval_ms * 1.1;
2634        let mut extrapolated_interval_ms = sampled_interval_ms;
2635        if duration_to_start_ms < extrapolation_threshold {
2636            extrapolated_interval_ms += duration_to_start_ms;
2637        } else {
2638            extrapolated_interval_ms += average_interval_ms / 2.0;
2639        }
2640        if duration_to_end_ms < extrapolation_threshold {
2641            extrapolated_interval_ms += duration_to_end_ms;
2642        } else {
2643            extrapolated_interval_ms += average_interval_ms / 2.0;
2644        }
2645
2646        let mut factor = extrapolated_interval_ms / sampled_interval_ms;
2647        if IS_RATE {
2648            factor /= range_length_secs;
2649        }
2650        histogram = histogram.scale(factor).into_gauge();
2651        result.push(Some(histogram));
2652    }
2653
2654    Ok(ColumnarValue::Array(build_histogram_array(&result)))
2655}
2656
2657fn create_native_extrapolated_udf<const IS_COUNTER: bool, const IS_RATE: bool>(
2658    name: &'static str,
2659    collector: Option<PromqlAnnotationCollector>,
2660) -> ScalarUDF {
2661    let input_types = vec![
2662        RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
2663        RangeArray::convert_data_type(native_histogram_arrow_type()),
2664        DataType::Timestamp(TimeUnit::Millisecond, None),
2665        DataType::Int64,
2666    ];
2667    create_udf(
2668        name,
2669        input_types,
2670        native_histogram_arrow_type(),
2671        Volatility::Volatile,
2672        Arc::new(move |input: &[ColumnarValue]| {
2673            let range_length = extract_array(&input[3])?;
2674            let range_length = range_length
2675                .as_any()
2676                .downcast_ref::<Int64Array>()
2677                .ok_or_else(|| {
2678                    DataFusionError::Execution(format!(
2679                        "{name}: expect Int64 as range length type, found {}",
2680                        range_length.data_type()
2681                    ))
2682                })?;
2683            if range_length.is_empty() || range_length.is_null(0) {
2684                return Err(DataFusionError::Execution(format!(
2685                    "{name}: range length must contain a non-null Int64 value"
2686                )));
2687            }
2688            native_extrapolated_rate::<IS_COUNTER, IS_RATE>(
2689                input,
2690                range_length.value(0),
2691                name,
2692                collector.clone(),
2693            )
2694        }) as _,
2695    )
2696}
2697
2698pub struct NativeHistogramDelta;
2699pub struct NativeHistogramRate;
2700pub struct NativeHistogramIncrease;
2701
2702impl NativeHistogramDelta {
2703    pub const fn name() -> &'static str {
2704        "prom_native_histogram_delta"
2705    }
2706
2707    pub fn scalar_udf() -> ScalarUDF {
2708        Self::scalar_udf_with_collector(None)
2709    }
2710
2711    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
2712        create_native_extrapolated_udf::<false, false>(Self::name(), collector)
2713    }
2714}
2715
2716impl NativeHistogramRate {
2717    pub const fn name() -> &'static str {
2718        "prom_native_histogram_rate"
2719    }
2720
2721    pub fn scalar_udf() -> ScalarUDF {
2722        Self::scalar_udf_with_collector(None)
2723    }
2724
2725    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
2726        create_native_extrapolated_udf::<true, true>(Self::name(), collector)
2727    }
2728}
2729
2730impl NativeHistogramIncrease {
2731    pub const fn name() -> &'static str {
2732        "prom_native_histogram_increase"
2733    }
2734
2735    pub fn scalar_udf() -> ScalarUDF {
2736        Self::scalar_udf_with_collector(None)
2737    }
2738
2739    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
2740        create_native_extrapolated_udf::<true, false>(Self::name(), collector)
2741    }
2742}
2743
2744fn native_idelta<const IS_RATE: bool>(
2745    input: &[ColumnarValue],
2746    func_name: &'static str,
2747    collector: Option<PromqlAnnotationCollector>,
2748) -> DfResult<ColumnarValue> {
2749    if input.len() != 2 {
2750        return Err(DataFusionError::Plan(format!(
2751            "{func_name} function should have 2 inputs"
2752        )));
2753    }
2754
2755    let ts_range = extract_range_dict(
2756        &input[0],
2757        func_name,
2758        "timestamp range vector",
2759        &DataType::Timestamp(TimeUnit::Millisecond, None),
2760    )?;
2761    let value_range = extract_range_dict(
2762        &input[1],
2763        func_name,
2764        "value range vector",
2765        &native_histogram_arrow_type(),
2766    )?;
2767
2768    if ts_range.keys().values() != value_range.keys().values() {
2769        return Err(DataFusionError::Execution(format!(
2770            "{func_name}: timestamp and value ranges should have the same window layout"
2771        )));
2772    }
2773
2774    let ts_values = ts_range
2775        .values()
2776        .as_any()
2777        .downcast_ref::<TimestampMillisecondArray>()
2778        .expect("validated timestamp range")
2779        .values();
2780    let histograms = value_range
2781        .values()
2782        .as_any()
2783        .downcast_ref::<StructArray>()
2784        .expect("validated native histogram range");
2785    let mut result = Vec::with_capacity(ts_range.keys().len());
2786
2787    for key in ts_range.keys().values() {
2788        let (offset, length) = unpack(*key);
2789        let offset = offset as usize;
2790        let length = length as usize;
2791        if length < 2 {
2792            result.push(None);
2793            continue;
2794        }
2795
2796        let mut samples = Vec::with_capacity(2);
2797        let mut has_null = false;
2798        for row in offset + length - 2..offset + length {
2799            let Some(histogram) = read_histogram(histograms, row)? else {
2800                has_null = true;
2801                break;
2802            };
2803            samples.push(histogram);
2804        }
2805        if has_null {
2806            result.push(None);
2807            continue;
2808        }
2809
2810        let wrong_flavor = samples.iter().any(|histogram| {
2811            if IS_RATE {
2812                histogram.reset_hint == GAUGE_RESET_HINT
2813            } else {
2814                histogram.reset_hint != GAUGE_RESET_HINT
2815            }
2816        });
2817        if wrong_flavor {
2818            let expected = if IS_RATE { "counter" } else { "gauge" };
2819            record_warning(
2820                &collector,
2821                format!("{func_name}: native histogram input should be a {expected} histogram"),
2822            );
2823        }
2824
2825        let sampled_interval_secs =
2826            (ts_values[offset + length - 1] - ts_values[offset + length - 2]) as f64 / 1000.0;
2827        if sampled_interval_secs <= 0.0 {
2828            result.push(None);
2829            continue;
2830        }
2831        record_custom_reconciliation(&collector, func_name, &samples[0], &samples[1]);
2832        let value = idelta_value(
2833            &samples,
2834            IS_RATE,
2835            ts_values[offset + length - 2],
2836            ts_values[offset + length - 1],
2837            sampled_interval_secs,
2838        );
2839        if value.is_none() {
2840            record_warning(
2841                &collector,
2842                format!("{func_name}: dropped native histogram range with incompatible schemas"),
2843            );
2844        }
2845        result.push(value);
2846    }
2847
2848    Ok(ColumnarValue::Array(build_histogram_array(&result)))
2849}
2850
2851fn create_native_idelta_udf<const IS_RATE: bool>(
2852    name: &'static str,
2853    collector: Option<PromqlAnnotationCollector>,
2854) -> ScalarUDF {
2855    create_udf(
2856        name,
2857        vec![
2858            RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
2859            RangeArray::convert_data_type(native_histogram_arrow_type()),
2860        ],
2861        native_histogram_arrow_type(),
2862        Volatility::Volatile,
2863        Arc::new(move |input: &[ColumnarValue]| {
2864            native_idelta::<IS_RATE>(input, name, collector.clone())
2865        }) as _,
2866    )
2867}
2868
2869pub struct NativeHistogramIDelta;
2870pub struct NativeHistogramIRate;
2871
2872impl NativeHistogramIDelta {
2873    pub const fn name() -> &'static str {
2874        "prom_native_histogram_idelta"
2875    }
2876
2877    pub fn scalar_udf() -> ScalarUDF {
2878        Self::scalar_udf_with_collector(None)
2879    }
2880
2881    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
2882        create_native_idelta_udf::<false>(Self::name(), collector)
2883    }
2884}
2885
2886impl NativeHistogramIRate {
2887    pub const fn name() -> &'static str {
2888        "prom_native_histogram_irate"
2889    }
2890
2891    pub fn scalar_udf() -> ScalarUDF {
2892        Self::scalar_udf_with_collector(None)
2893    }
2894
2895    pub fn scalar_udf_with_collector(collector: Option<PromqlAnnotationCollector>) -> ScalarUDF {
2896        create_native_idelta_udf::<true>(Self::name(), collector)
2897    }
2898}
2899
2900#[cfg(test)]
2901mod tests {
2902    use datafusion::arrow::datatypes::Field;
2903    use datafusion_common::config::ConfigOptions;
2904    use datafusion_expr::ScalarFunctionArgs;
2905
2906    use super::*;
2907
2908    fn sample_histogram(count: f64, sum: f64, positive_buckets: Vec<f64>) -> NativeHistogram {
2909        NativeHistogram {
2910            schema: 0,
2911            zero_threshold: 0.0,
2912            sum,
2913            reset_hint: UNKNOWN_COUNTER_RESET_HINT,
2914            start_timestamp: None,
2915            custom_values: Vec::new(),
2916            positive_spans: vec![Span {
2917                offset: 0,
2918                length: positive_buckets.len() as i32,
2919            }],
2920            negative_spans: Vec::new(),
2921            count,
2922            zero_count: 0.0,
2923            positive_buckets,
2924            negative_buckets: Vec::new(),
2925        }
2926    }
2927
2928    fn run_scalar_udf(udf: ScalarUDF, input: Vec<ColumnarValue>) -> f64 {
2929        let result = run_udf(udf, input, DataType::Float64);
2930        extract_array(&result)
2931            .unwrap()
2932            .as_any()
2933            .downcast_ref::<Float64Array>()
2934            .unwrap()
2935            .value(0)
2936    }
2937
2938    fn run_udf(udf: ScalarUDF, input: Vec<ColumnarValue>, return_type: DataType) -> ColumnarValue {
2939        let arg_fields = input
2940            .iter()
2941            .enumerate()
2942            .map(|(idx, input)| Arc::new(Field::new(format!("arg_{idx}"), input.data_type(), true)))
2943            .collect();
2944        let args = ScalarFunctionArgs {
2945            args: input,
2946            arg_fields,
2947            number_rows: 1,
2948            return_field: Arc::new(Field::new("result", return_type, true)),
2949            config_options: Arc::new(ConfigOptions::default()),
2950        };
2951
2952        udf.invoke_with_args(args).unwrap()
2953    }
2954
2955    fn run_histogram_udf(udf: ScalarUDF, input: Vec<ColumnarValue>) -> NativeHistogram {
2956        let result = run_udf(udf, input, native_histogram_arrow_type());
2957        let array = extract_array(&result)
2958            .unwrap()
2959            .as_any()
2960            .downcast_ref::<StructArray>()
2961            .unwrap()
2962            .clone();
2963        read_histogram(&array, 0).unwrap().unwrap()
2964    }
2965
2966    fn evaluated_histogram(
2967        accumulator: &mut NativeHistogramAggregateAccumulator,
2968    ) -> NativeHistogram {
2969        let ScalarValue::Struct(array) = accumulator.evaluate().unwrap() else {
2970            panic!("native histogram accumulator returned a non-struct value");
2971        };
2972        read_histogram(&array, 0).unwrap().unwrap()
2973    }
2974
2975    fn histogram_range_input(values: Vec<Option<NativeHistogram>>) -> Vec<ColumnarValue> {
2976        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
2977            (0..values.len()).map(|idx| Some((idx as i64 + 1) * 1000)),
2978        ));
2979        let histograms = build_histogram_array(&values);
2980        let range = [(0, values.len() as u32)];
2981        let ts_range = RangeArray::from_ranges(timestamps, range).unwrap();
2982        let value_range = RangeArray::from_ranges(histograms, range).unwrap();
2983
2984        vec![
2985            ColumnarValue::Array(Arc::new(ts_range.into_dict())),
2986            ColumnarValue::Array(Arc::new(value_range.into_dict())),
2987        ]
2988    }
2989
2990    fn mixed_range_input(
2991        name: &str,
2992        floats: Vec<Option<f64>>,
2993        histograms: Vec<Option<NativeHistogram>>,
2994    ) -> Vec<ColumnarValue> {
2995        assert_eq!(floats.len(), histograms.len());
2996        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
2997            (0..floats.len()).map(|idx| Some((idx as i64 + 1) * 1000)),
2998        ));
2999        let floats = Arc::new(Float64Array::from(floats));
3000        let histograms = build_histogram_array(&histograms);
3001        let range = [(0, u32::try_from(floats.len()).unwrap())];
3002        vec![
3003            ColumnarValue::Scalar(ScalarValue::Utf8(Some(name.to_string()))),
3004            ColumnarValue::Array(Arc::new(
3005                RangeArray::from_ranges(timestamps, range)
3006                    .unwrap()
3007                    .into_dict(),
3008            )),
3009            ColumnarValue::Array(Arc::new(
3010                RangeArray::from_ranges(floats, range).unwrap().into_dict(),
3011            )),
3012            ColumnarValue::Array(Arc::new(
3013                RangeArray::from_ranges(histograms, range)
3014                    .unwrap()
3015                    .into_dict(),
3016            )),
3017        ]
3018    }
3019
3020    fn mixed_float_result(udf: ScalarUDF, input: Vec<ColumnarValue>) -> Option<f64> {
3021        let result = run_udf(udf, input, DataType::Float64);
3022        let result = extract_array(&result).unwrap();
3023        let result = result.as_any().downcast_ref::<Float64Array>().unwrap();
3024        result.is_valid(0).then(|| result.value(0))
3025    }
3026
3027    fn mixed_histogram_result(
3028        udf: ScalarUDF,
3029        input: Vec<ColumnarValue>,
3030    ) -> Option<NativeHistogram> {
3031        let result = run_udf(udf, input, native_histogram_arrow_type());
3032        let result = extract_array(&result).unwrap();
3033        let result = result.as_any().downcast_ref::<StructArray>().unwrap();
3034        read_histogram(result, 0).unwrap()
3035    }
3036
3037    fn run_histogram_range_udf(
3038        udf: ScalarUDF,
3039        histograms: Vec<NativeHistogram>,
3040    ) -> NativeHistogram {
3041        run_histogram_udf(
3042            udf,
3043            histogram_range_input(histograms.into_iter().map(Some).collect()),
3044        )
3045    }
3046
3047    fn run_float_range_udf(
3048        udf: ScalarUDF,
3049        histograms: Vec<Option<NativeHistogram>>,
3050    ) -> Option<f64> {
3051        let result = run_udf(udf, histogram_range_input(histograms), DataType::Float64);
3052        let result = extract_array(&result).unwrap();
3053        let result = result.as_any().downcast_ref::<Float64Array>().unwrap();
3054        (!result.is_null(0)).then(|| result.value(0))
3055    }
3056
3057    fn run_extrapolated_histogram_udf(
3058        udf: ScalarUDF,
3059        histograms: Vec<NativeHistogram>,
3060    ) -> NativeHistogram {
3061        let range_length = histograms.len() as i64 * 1000;
3062        let timestamps = (0..histograms.len())
3063            .map(|idx| (idx as i64 + 1) * 1000)
3064            .collect();
3065        extrapolated_histogram_result(udf, timestamps, histograms, range_length, range_length)
3066            .unwrap()
3067    }
3068
3069    fn extrapolated_histogram_result(
3070        udf: ScalarUDF,
3071        timestamps: Vec<i64>,
3072        histograms: Vec<NativeHistogram>,
3073        range_end: i64,
3074        range_length: i64,
3075    ) -> Option<NativeHistogram> {
3076        assert_eq!(timestamps.len(), histograms.len());
3077        let range = [(0, u32::try_from(histograms.len()).unwrap())];
3078        let timestamps = Arc::new(TimestampMillisecondArray::from(timestamps));
3079        let histograms =
3080            build_histogram_array(&histograms.into_iter().map(Some).collect::<Vec<_>>());
3081        let mut input = vec![
3082            ColumnarValue::Array(Arc::new(
3083                RangeArray::from_ranges(timestamps, range)
3084                    .unwrap()
3085                    .into_dict(),
3086            )),
3087            ColumnarValue::Array(Arc::new(
3088                RangeArray::from_ranges(histograms, range)
3089                    .unwrap()
3090                    .into_dict(),
3091            )),
3092        ];
3093        input.push(ColumnarValue::Array(Arc::new(
3094            TimestampMillisecondArray::from(vec![range_end]),
3095        )));
3096        input.push(ColumnarValue::Array(Arc::new(Int64Array::from(vec![
3097            range_length,
3098        ]))));
3099        let result = run_udf(udf, input, native_histogram_arrow_type());
3100        let result = extract_array(&result).unwrap();
3101        let result = result.as_any().downcast_ref::<StructArray>().unwrap();
3102        read_histogram(result, 0).unwrap()
3103    }
3104
3105    fn collected_warnings(collector: &PromqlAnnotationCollector) -> Vec<String> {
3106        let mut warnings = Vec::new();
3107        collector.append_to(&mut warnings, &mut Vec::new());
3108        warnings
3109    }
3110
3111    fn collected_infos(collector: &PromqlAnnotationCollector) -> Vec<String> {
3112        let mut infos = Vec::new();
3113        collector.append_to(&mut Vec::new(), &mut infos);
3114        infos
3115    }
3116
3117    #[test]
3118    fn quantile_and_fraction_report_nan_observations() {
3119        let histogram = sample_histogram(10.0, f64::NAN, vec![8.0]);
3120        let histogram_arg =
3121            || ColumnarValue::Array(build_histogram_array(&[Some(histogram.clone())]));
3122
3123        let quantile_collector = PromqlAnnotationCollector::default();
3124        let skewed = run_scalar_udf(
3125            NativeHistogramQuantile::scalar_udf_with_collector(Some(quantile_collector.clone())),
3126            vec![
3127                histogram_arg(),
3128                ColumnarValue::Scalar(ScalarValue::Float64(Some(0.5))),
3129            ],
3130        );
3131        assert!(skewed.is_finite());
3132        let nan = run_scalar_udf(
3133            NativeHistogramQuantile::scalar_udf_with_collector(Some(quantile_collector.clone())),
3134            vec![
3135                histogram_arg(),
3136                ColumnarValue::Scalar(ScalarValue::Float64(Some(0.9))),
3137            ],
3138        );
3139        assert!(nan.is_nan());
3140        let infos = collected_infos(&quantile_collector);
3141        assert!(
3142            infos
3143                .iter()
3144                .any(|info| info.ends_with("result is skewed higher"))
3145        );
3146        assert!(infos.iter().any(|info| info.ends_with("result is NaN")));
3147
3148        let fraction_collector = PromqlAnnotationCollector::default();
3149        assert_eq!(
3150            run_scalar_udf(
3151                NativeHistogramFraction::scalar_udf_with_collector(Some(
3152                    fraction_collector.clone(),
3153                )),
3154                vec![
3155                    histogram_arg(),
3156                    ColumnarValue::Scalar(ScalarValue::Float64(Some(f64::NEG_INFINITY))),
3157                    ColumnarValue::Scalar(ScalarValue::Float64(Some(f64::INFINITY))),
3158                ],
3159            ),
3160            0.8
3161        );
3162        assert_eq!(
3163            collected_infos(&fraction_collector),
3164            vec![
3165                "input to histogram_fraction has NaN observations, which are excluded from all fractions"
3166                    .to_string()
3167            ]
3168        );
3169    }
3170
3171    #[test]
3172    fn mixed_ranges_follow_prometheus_sample_type_semantics() {
3173        let first = sample_histogram(1.0, 1.0, vec![1.0]);
3174        let second = sample_histogram(3.0, 3.0, vec![3.0]);
3175        let collector = PromqlAnnotationCollector::default();
3176
3177        let mut rate = mixed_range_input(
3178            "rate",
3179            vec![Some(1.0), None, Some(3.0)],
3180            vec![None, Some(first.clone()), None],
3181        );
3182        rate.push(ColumnarValue::Array(Arc::new(
3183            TimestampMillisecondArray::from(vec![3000]),
3184        )));
3185        rate.push(ColumnarValue::Array(Arc::new(Int64Array::from(vec![3000]))));
3186        assert_eq!(
3187            mixed_float_result(MixedRange::float_udf(Some(collector.clone())), rate.clone()),
3188            None
3189        );
3190        assert_eq!(
3191            mixed_histogram_result(MixedRange::histogram_udf(Some(collector.clone())), rate),
3192            None
3193        );
3194        assert!(
3195            collected_warnings(&collector)
3196                .iter()
3197                .any(|warning| warning.contains("mix of float and native histogram"))
3198        );
3199
3200        let pure_rate = |floats, histograms| {
3201            let mut input = mixed_range_input("rate", floats, histograms);
3202            input.push(ColumnarValue::Array(Arc::new(
3203                TimestampMillisecondArray::from(vec![3000]),
3204            )));
3205            input.push(ColumnarValue::Array(Arc::new(Int64Array::from(vec![3000]))));
3206            input
3207        };
3208        let pure_float = pure_rate(
3209            vec![Some(1.0), Some(2.0), Some(3.0)],
3210            vec![None, None, None],
3211        );
3212        assert_eq!(
3213            mixed_float_result(MixedRange::float_udf(None), pure_float.clone()),
3214            Some(1.0)
3215        );
3216        assert_eq!(
3217            mixed_histogram_result(MixedRange::histogram_udf(None), pure_float),
3218            None
3219        );
3220        let pure_histogram = pure_rate(
3221            vec![None, None, None],
3222            vec![
3223                Some(sample_histogram(1.0, 1.0, vec![1.0])),
3224                Some(sample_histogram(2.0, 2.0, vec![2.0])),
3225                Some(sample_histogram(3.0, 3.0, vec![3.0])),
3226            ],
3227        );
3228        assert_eq!(
3229            mixed_float_result(MixedRange::float_udf(None), pure_histogram.clone()),
3230            None
3231        );
3232        assert_eq!(
3233            mixed_histogram_result(MixedRange::histogram_udf(None), pure_histogram)
3234                .unwrap()
3235                .count,
3236            1.0
3237        );
3238
3239        let idelta = mixed_range_input(
3240            "idelta",
3241            vec![Some(10.0), None, None],
3242            vec![None, Some(first.clone()), Some(second.clone())],
3243        );
3244        assert_eq!(
3245            mixed_float_result(MixedRange::float_udf(None), idelta.clone()),
3246            None
3247        );
3248        assert_eq!(
3249            mixed_histogram_result(MixedRange::histogram_udf(None), idelta)
3250                .unwrap()
3251                .count,
3252            2.0
3253        );
3254
3255        let alternating = || {
3256            mixed_range_input(
3257                "changes",
3258                vec![Some(1.0), None, None, Some(1.0)],
3259                vec![None, Some(first.clone()), Some(first.clone()), None],
3260            )
3261        };
3262        assert_eq!(
3263            mixed_float_result(MixedRange::float_udf(None), alternating()),
3264            Some(2.0)
3265        );
3266        let mut resets = alternating();
3267        resets[0] = ColumnarValue::Scalar(ScalarValue::Utf8(Some("resets".to_string())));
3268        assert_eq!(
3269            mixed_float_result(MixedRange::float_udf(None), resets),
3270            Some(2.0)
3271        );
3272
3273        for (name, expected) in [
3274            ("count_over_time", Some(4.0)),
3275            ("present_over_time", Some(1.0)),
3276            ("absent_over_time", None),
3277        ] {
3278            let mut input = alternating();
3279            input[0] = ColumnarValue::Scalar(ScalarValue::Utf8(Some(name.to_string())));
3280            assert_eq!(
3281                mixed_float_result(MixedRange::float_udf(None), input),
3282                expected,
3283                "{name}"
3284            );
3285        }
3286
3287        let last = mixed_range_input(
3288            "last_over_time",
3289            vec![None, Some(4.0)],
3290            vec![Some(first.clone()), None],
3291        );
3292        assert_eq!(
3293            mixed_float_result(MixedRange::float_udf(None), last.clone()),
3294            Some(4.0)
3295        );
3296        assert_eq!(
3297            mixed_histogram_result(MixedRange::histogram_udf(None), last),
3298            None
3299        );
3300
3301        let collector = PromqlAnnotationCollector::default();
3302        let histogram_only_min =
3303            mixed_range_input("min_over_time", vec![None], vec![Some(first.clone())]);
3304        assert_eq!(
3305            mixed_float_result(
3306                MixedRange::float_udf(Some(collector.clone())),
3307                histogram_only_min,
3308            ),
3309            None
3310        );
3311        assert!(collected_infos(&collector).is_empty());
3312
3313        let min = mixed_range_input(
3314            "min_over_time",
3315            vec![Some(3.0), None, Some(1.0)],
3316            vec![None, Some(first), None],
3317        );
3318        assert_eq!(
3319            mixed_float_result(MixedRange::float_udf(Some(collector.clone())), min),
3320            Some(1.0)
3321        );
3322        assert!(
3323            collected_infos(&collector)
3324                .iter()
3325                .any(|info| info.contains("ignored native histogram"))
3326        );
3327    }
3328
3329    #[test]
3330    fn count_sum_and_avg_read_struct() {
3331        let histograms = vec![Some(sample_histogram(6.0, 10.0, vec![2.0, 4.0]))];
3332        let array = build_histogram_array(&histograms);
3333        let input = vec![ColumnarValue::Array(array)];
3334
3335        let count = run_scalar_udf(NativeHistogramCount::scalar_udf(), input.clone());
3336        assert_eq!(count, 6.0);
3337
3338        let sum = run_scalar_udf(NativeHistogramSum::scalar_udf(), input.clone());
3339        assert_eq!(sum, 10.0);
3340
3341        let avg = run_scalar_udf(NativeHistogramAvg::scalar_udf(), input);
3342        assert_eq!(avg, 10.0 / 6.0);
3343    }
3344
3345    #[test]
3346    fn quantile_uses_bucket_bounds() {
3347        let histogram = sample_histogram(6.0, 10.0, vec![2.0, 4.0]);
3348        assert_eq!(histogram.quantile(0.0), 0.5);
3349        assert!(histogram.quantile(0.5) > 1.0);
3350        assert!(histogram.quantile(0.5) < 2.0);
3351    }
3352
3353    #[test]
3354    fn comparison_observes_explicit_sparse_zero_buckets() {
3355        let mut left = sample_histogram(1.0, 1.0, vec![1.0, 0.0]);
3356        left.reset_hint = COUNTER_RESET_HINT;
3357        left.start_timestamp = Some(1000);
3358        let mut right = sample_histogram(1.0, 1.0, vec![1.0]);
3359        right.reset_hint = NOT_COUNTER_RESET_HINT;
3360        right.start_timestamp = Some(2000);
3361
3362        let result = run_udf(
3363            NativeHistogramEq::scalar_udf(),
3364            vec![
3365                ColumnarValue::Array(build_histogram_array(&[Some(left)])),
3366                ColumnarValue::Array(build_histogram_array(&[Some(right)])),
3367            ],
3368            DataType::Boolean,
3369        );
3370        let values = extract_array(&result).unwrap();
3371        let values = values.as_any().downcast_ref::<BooleanArray>().unwrap();
3372        assert!(!values.value(0));
3373    }
3374
3375    #[test]
3376    fn unary_minus_returns_gauge_histogram() {
3377        let result = run_histogram_udf(
3378            NativeHistogramNeg::scalar_udf(),
3379            vec![ColumnarValue::Array(build_histogram_array(&[Some(
3380                sample_histogram(2.0, 3.0, vec![2.0]),
3381            )]))],
3382        );
3383
3384        assert_eq!(result.reset_hint, GAUGE_RESET_HINT);
3385        assert_eq!(result.count, -2.0);
3386        assert_eq!(result.sum, -3.0);
3387        assert_eq!(result.positive_buckets, vec![-2.0]);
3388    }
3389
3390    #[test]
3391    fn histogram_over_time_functions_preserve_reset_hints() {
3392        let mut first = sample_histogram(1.0, 1.0, vec![1.0]);
3393        first.reset_hint = COUNTER_RESET_HINT;
3394        let mut second = sample_histogram(2.0, 2.0, vec![2.0]);
3395        second.reset_hint = COUNTER_RESET_HINT;
3396
3397        let result = run_histogram_range_udf(
3398            NativeHistogramSumOverTime::scalar_udf(),
3399            vec![first.clone(), second.clone()],
3400        );
3401        assert_eq!(result.reset_hint, COUNTER_RESET_HINT);
3402        assert_eq!(result.count, 3.0);
3403        assert_eq!(result.sum, 3.0);
3404        assert_eq!(result.positive_buckets, vec![3.0]);
3405
3406        let result = run_histogram_range_udf(
3407            NativeHistogramAvgOverTime::scalar_udf(),
3408            vec![first.clone(), second.clone()],
3409        );
3410        assert_eq!(result.reset_hint, COUNTER_RESET_HINT);
3411        assert_eq!(result.count, 1.5);
3412        assert_eq!(result.sum, 1.5);
3413        assert_eq!(result.positive_buckets, vec![1.5]);
3414
3415        let result = run_histogram_range_udf(
3416            NativeHistogramLastOverTime::scalar_udf(),
3417            vec![first, second],
3418        );
3419        assert_eq!(result.reset_hint, COUNTER_RESET_HINT);
3420        assert_eq!(result.count, 2.0);
3421        assert_eq!(result.sum, 2.0);
3422        assert_eq!(result.positive_buckets, vec![2.0]);
3423    }
3424
3425    #[test]
3426    fn histogram_averages_avoid_sum_overflow() {
3427        let large = sample_histogram(1.0e308, 1.0e308, vec![1.0e308]);
3428
3429        let range_average = run_histogram_range_udf(
3430            NativeHistogramAvgOverTime::scalar_udf(),
3431            vec![large.clone(), large.clone()],
3432        );
3433        assert_eq!(range_average.count, 1.0e308);
3434        assert_eq!(range_average.sum, 1.0e308);
3435        assert_eq!(range_average.positive_buckets, vec![1.0e308]);
3436
3437        let mut aggregate =
3438            NativeHistogramAggregateAccumulator::new(NativeHistogramAggregateKind::Avg, None);
3439        aggregate.push_histogram(large.clone(), 1).unwrap();
3440        aggregate.push_histogram(large, 1).unwrap();
3441        let aggregate_average = evaluated_histogram(&mut aggregate);
3442        assert_eq!(aggregate_average.count, 1.0e308);
3443        assert_eq!(aggregate_average.sum, 1.0e308);
3444        assert_eq!(aggregate_average.positive_buckets, vec![1.0e308]);
3445    }
3446
3447    #[test]
3448    fn histogram_average_partial_states_are_weighted() {
3449        let mut first =
3450            NativeHistogramAggregateAccumulator::new(NativeHistogramAggregateKind::Avg, None);
3451        first
3452            .push_histogram(sample_histogram(1.0, 1.0, vec![1.0]), 1)
3453            .unwrap();
3454        first
3455            .push_histogram(sample_histogram(3.0, 3.0, vec![3.0]), 1)
3456            .unwrap();
3457        let first_state = first
3458            .state()
3459            .unwrap()
3460            .into_iter()
3461            .map(|value| value.to_array_of_size(1).unwrap())
3462            .collect::<Vec<_>>();
3463
3464        let mut second =
3465            NativeHistogramAggregateAccumulator::new(NativeHistogramAggregateKind::Avg, None);
3466        second
3467            .push_histogram(sample_histogram(8.0, 8.0, vec![8.0]), 1)
3468            .unwrap();
3469        let second_state = second
3470            .state()
3471            .unwrap()
3472            .into_iter()
3473            .map(|value| value.to_array_of_size(1).unwrap())
3474            .collect::<Vec<_>>();
3475
3476        let mut merged =
3477            NativeHistogramAggregateAccumulator::new(NativeHistogramAggregateKind::Avg, None);
3478        merged.merge_batch(&first_state).unwrap();
3479        merged.merge_batch(&second_state).unwrap();
3480        let average = evaluated_histogram(&mut merged);
3481        assert_eq!(average.count, 4.0);
3482        assert_eq!(average.sum, 4.0);
3483        assert_eq!(average.positive_buckets, vec![4.0]);
3484    }
3485
3486    #[test]
3487    fn histogram_average_rejects_sample_count_overflow() {
3488        let mut aggregate =
3489            NativeHistogramAggregateAccumulator::new(NativeHistogramAggregateKind::Avg, None);
3490        aggregate.value = Some(sample_histogram(1.0, 1.0, vec![1.0]));
3491        aggregate.count = u64::MAX;
3492
3493        let error = aggregate
3494            .push_histogram(sample_histogram(1.0, 1.0, vec![1.0]), 1)
3495            .unwrap_err();
3496        assert!(error.to_string().contains("sample count overflow"));
3497    }
3498
3499    #[test]
3500    fn presence_only_range_functions_preserve_null_semantics() {
3501        let first = sample_histogram(1.0, 1.0, vec![1.0]);
3502        let second = sample_histogram(2.0, 2.0, vec![2.0]);
3503        assert_eq!(
3504            run_float_range_udf(
3505                NativeHistogramCountOverTime::scalar_udf(),
3506                vec![Some(first.clone()), Some(second.clone())],
3507            ),
3508            Some(2.0)
3509        );
3510        assert_eq!(
3511            run_float_range_udf(
3512                NativeHistogramPresentOverTime::scalar_udf(),
3513                vec![Some(first.clone()), Some(second.clone())],
3514            ),
3515            Some(1.0)
3516        );
3517        assert_eq!(
3518            run_float_range_udf(
3519                NativeHistogramCountOverTime::scalar_udf(),
3520                vec![None, Some(second.clone())],
3521            ),
3522            None
3523        );
3524
3525        let result = run_udf(
3526            NativeHistogramLastOverTime::scalar_udf(),
3527            histogram_range_input(vec![None, Some(second)]),
3528            native_histogram_arrow_type(),
3529        );
3530        let result = extract_array(&result).unwrap();
3531        let result = result.as_any().downcast_ref::<StructArray>().unwrap();
3532        assert!(result.is_null(0));
3533    }
3534
3535    #[test]
3536    fn resets_counts_histogram_flavor_transitions() {
3537        let mut counter = sample_histogram(1.0, 1.0, vec![1.0]);
3538        counter.reset_hint = NOT_COUNTER_RESET_HINT;
3539        let mut gauge = sample_histogram(2.0, 2.0, vec![2.0]);
3540        gauge.reset_hint = GAUGE_RESET_HINT;
3541        assert_eq!(
3542            run_float_range_udf(
3543                NativeHistogramResets::scalar_udf(),
3544                vec![Some(counter), Some(gauge)],
3545            ),
3546            Some(1.0)
3547        );
3548
3549        let mut gauge = sample_histogram(1.0, 1.0, vec![1.0]);
3550        gauge.reset_hint = GAUGE_RESET_HINT;
3551        let mut counter = sample_histogram(2.0, 2.0, vec![2.0]);
3552        counter.reset_hint = NOT_COUNTER_RESET_HINT;
3553        assert_eq!(
3554            run_float_range_udf(
3555                NativeHistogramResets::scalar_udf(),
3556                vec![Some(gauge), Some(counter)],
3557            ),
3558            Some(1.0)
3559        );
3560    }
3561
3562    #[test]
3563    fn wrong_flavor_functions_record_warnings() {
3564        let mut gauge_first = sample_histogram(1.0, 1.0, vec![1.0]);
3565        gauge_first.reset_hint = GAUGE_RESET_HINT;
3566        let mut gauge_last = sample_histogram(2.0, 2.0, vec![2.0]);
3567        gauge_last.reset_hint = GAUGE_RESET_HINT;
3568        let mut counter_first = sample_histogram(1.0, 1.0, vec![1.0]);
3569        counter_first.reset_hint = NOT_COUNTER_RESET_HINT;
3570        let mut counter_last = sample_histogram(2.0, 2.0, vec![2.0]);
3571        counter_last.reset_hint = NOT_COUNTER_RESET_HINT;
3572        let collector = PromqlAnnotationCollector::default();
3573
3574        run_extrapolated_histogram_udf(
3575            NativeHistogramRate::scalar_udf_with_collector(Some(collector.clone())),
3576            vec![gauge_first.clone(), gauge_last.clone()],
3577        );
3578        run_extrapolated_histogram_udf(
3579            NativeHistogramDelta::scalar_udf_with_collector(Some(collector.clone())),
3580            vec![counter_first.clone(), counter_last.clone()],
3581        );
3582        run_histogram_range_udf(
3583            NativeHistogramIRate::scalar_udf_with_collector(Some(collector.clone())),
3584            vec![gauge_first, gauge_last],
3585        );
3586        run_histogram_range_udf(
3587            NativeHistogramIDelta::scalar_udf_with_collector(Some(collector.clone())),
3588            vec![counter_first, counter_last],
3589        );
3590
3591        let warnings = collected_warnings(&collector);
3592        for expected in [
3593            format!(
3594                "{}: native histogram input should be a counter histogram",
3595                NativeHistogramRate::name()
3596            ),
3597            format!(
3598                "{}: native histogram input should be a gauge histogram",
3599                NativeHistogramDelta::name()
3600            ),
3601            format!(
3602                "{}: native histogram input should be a counter histogram",
3603                NativeHistogramIRate::name()
3604            ),
3605            format!(
3606                "{}: native histogram input should be a gauge histogram",
3607                NativeHistogramIDelta::name()
3608            ),
3609        ] {
3610            assert!(warnings.contains(&expected), "missing warning: {expected}");
3611        }
3612    }
3613
3614    #[test]
3615    fn subtraction_records_reset_hint_contradictions_for_incompatible_histograms() {
3616        for incompatible in [false, true] {
3617            let mut left = sample_histogram(2.0, 2.0, vec![2.0]);
3618            left.reset_hint = COUNTER_RESET_HINT;
3619            let mut right = sample_histogram(1.0, 1.0, vec![1.0]);
3620            right.reset_hint = NOT_COUNTER_RESET_HINT;
3621            if incompatible {
3622                right.schema = CUSTOM_BUCKETS_SCHEMA;
3623                right.custom_values = vec![1.0];
3624            }
3625            let collector = PromqlAnnotationCollector::default();
3626
3627            run_udf(
3628                NativeHistogramSub::scalar_udf_with_collector(Some(collector.clone())),
3629                vec![
3630                    ColumnarValue::Array(build_histogram_array(&[Some(left)])),
3631                    ColumnarValue::Array(build_histogram_array(&[Some(right)])),
3632                ],
3633                native_histogram_arrow_type(),
3634            );
3635
3636            assert!(collected_warnings(&collector).contains(&format!(
3637                "{}: native histogram counter reset hints contradict",
3638                NativeHistogramSub::name()
3639            )));
3640        }
3641    }
3642
3643    #[test]
3644    fn counter_reset_hint_history_survives_folds_and_state_merges() {
3645        let mut reset = sample_histogram(1.0, 1.0, vec![1.0]);
3646        reset.reset_hint = COUNTER_RESET_HINT;
3647        let mut unknown = sample_histogram(2.0, 2.0, vec![2.0]);
3648        unknown.reset_hint = UNKNOWN_COUNTER_RESET_HINT;
3649        let mut not_reset = sample_histogram(3.0, 3.0, vec![3.0]);
3650        not_reset.reset_hint = NOT_COUNTER_RESET_HINT;
3651
3652        let range_collector = PromqlAnnotationCollector::default();
3653        assert!(
3654            range_fold_histograms(
3655                vec![reset.clone(), unknown.clone(), not_reset.clone()],
3656                NativeHistogramAggregateKind::Sum,
3657                NativeHistogramSumOverTime::name(),
3658                &Some(range_collector.clone()),
3659            )
3660            .is_some()
3661        );
3662        assert!(collected_warnings(&range_collector).contains(&format!(
3663            "{}: native histogram counter reset hints contradict",
3664            NativeHistogramSumOverTime::name()
3665        )));
3666
3667        for kind in [
3668            NativeHistogramAggregateKind::Sum,
3669            NativeHistogramAggregateKind::Avg,
3670        ] {
3671            let mut first_partial = NativeHistogramAggregateAccumulator::new(kind, None);
3672            first_partial.push_histogram(reset.clone(), 1).unwrap();
3673            first_partial.push_histogram(unknown.clone(), 1).unwrap();
3674            let first_states = first_partial
3675                .state()
3676                .unwrap()
3677                .into_iter()
3678                .map(|value| value.to_array_of_size(1).unwrap())
3679                .collect::<Vec<_>>();
3680
3681            let mut second_partial = NativeHistogramAggregateAccumulator::new(kind, None);
3682            second_partial.push_histogram(not_reset.clone(), 1).unwrap();
3683            let second_states = second_partial
3684                .state()
3685                .unwrap()
3686                .into_iter()
3687                .map(|value| value.to_array_of_size(1).unwrap())
3688                .collect::<Vec<_>>();
3689
3690            let collector = PromqlAnnotationCollector::default();
3691            let mut merged =
3692                NativeHistogramAggregateAccumulator::new(kind, Some(collector.clone()));
3693            merged.merge_batch(&first_states).unwrap();
3694            merged.merge_batch(&second_states).unwrap();
3695            assert!(collected_warnings(&collector).contains(&format!(
3696                "{}: native histogram counter reset hints contradict",
3697                kind.name()
3698            )));
3699        }
3700    }
3701
3702    #[test]
3703    fn incompatible_aggregates_do_not_hide_reset_hint_contradictions() {
3704        let mut reset = sample_histogram(1.0, 1.0, vec![1.0]);
3705        reset.reset_hint = COUNTER_RESET_HINT;
3706        let mut incompatible_reset = sample_histogram(2.0, 2.0, vec![2.0]);
3707        incompatible_reset.schema = CUSTOM_BUCKETS_SCHEMA;
3708        incompatible_reset.custom_values = vec![1.0];
3709        incompatible_reset.reset_hint = COUNTER_RESET_HINT;
3710        let mut not_reset = sample_histogram(3.0, 3.0, vec![3.0]);
3711        not_reset.reset_hint = NOT_COUNTER_RESET_HINT;
3712
3713        for kind in [
3714            NativeHistogramAggregateKind::Sum,
3715            NativeHistogramAggregateKind::Avg,
3716        ] {
3717            for histograms in [
3718                [reset.clone(), incompatible_reset.clone(), not_reset.clone()],
3719                [reset.clone(), not_reset.clone(), incompatible_reset.clone()],
3720            ] {
3721                let collector = PromqlAnnotationCollector::default();
3722                let mut aggregate =
3723                    NativeHistogramAggregateAccumulator::new(kind, Some(collector.clone()));
3724                for histogram in histograms {
3725                    aggregate.push_histogram(histogram, 1).unwrap();
3726                }
3727
3728                assert!(aggregate.dropped_incompatible);
3729                let warnings = collected_warnings(&collector);
3730                assert!(warnings.contains(&format!(
3731                    "{}: dropped native histogram aggregate with incompatible schemas",
3732                    kind.name()
3733                )));
3734                assert!(warnings.contains(&format!(
3735                    "{}: native histogram counter reset hints contradict",
3736                    kind.name()
3737                )));
3738            }
3739
3740            let mut dropped_partial = NativeHistogramAggregateAccumulator::new(kind, None);
3741            dropped_partial.push_histogram(reset.clone(), 1).unwrap();
3742            dropped_partial
3743                .push_histogram(incompatible_reset.clone(), 1)
3744                .unwrap();
3745            let dropped_states = dropped_partial
3746                .state()
3747                .unwrap()
3748                .into_iter()
3749                .map(|value| value.to_array_of_size(1).unwrap())
3750                .collect::<Vec<_>>();
3751
3752            let mut opposing_partial = NativeHistogramAggregateAccumulator::new(kind, None);
3753            opposing_partial
3754                .push_histogram(not_reset.clone(), 1)
3755                .unwrap();
3756            let opposing_states = opposing_partial
3757                .state()
3758                .unwrap()
3759                .into_iter()
3760                .map(|value| value.to_array_of_size(1).unwrap())
3761                .collect::<Vec<_>>();
3762
3763            for (first, second) in [
3764                (&dropped_states, &opposing_states),
3765                (&opposing_states, &dropped_states),
3766            ] {
3767                let collector = PromqlAnnotationCollector::default();
3768                let mut merged =
3769                    NativeHistogramAggregateAccumulator::new(kind, Some(collector.clone()));
3770                merged.merge_batch(first).unwrap();
3771                merged.merge_batch(second).unwrap();
3772
3773                assert!(merged.dropped_incompatible);
3774                assert!(collected_warnings(&collector).contains(&format!(
3775                    "{}: native histogram counter reset hints contradict",
3776                    kind.name()
3777                )));
3778            }
3779        }
3780
3781        for (kind, name) in [
3782            (
3783                NativeHistogramAggregateKind::Sum,
3784                NativeHistogramSumOverTime::name(),
3785            ),
3786            (
3787                NativeHistogramAggregateKind::Avg,
3788                NativeHistogramAvgOverTime::name(),
3789            ),
3790        ] {
3791            let collector = PromqlAnnotationCollector::default();
3792            assert!(
3793                range_fold_histograms(
3794                    vec![reset.clone(), incompatible_reset.clone(), not_reset.clone()],
3795                    kind,
3796                    name,
3797                    &Some(collector.clone()),
3798                )
3799                .is_none()
3800            );
3801            let warnings = collected_warnings(&collector);
3802            assert!(warnings.contains(&format!(
3803                "{name}: dropped native histogram range with incompatible schemas"
3804            )));
3805            assert!(warnings.contains(&format!(
3806                "{name}: native histogram counter reset hints contradict"
3807            )));
3808        }
3809    }
3810
3811    #[test]
3812    fn histogram_aggregate_accumulator_accounts_for_heap_allocations() {
3813        let histogram = sample_histogram(3.0, 3.0, vec![1.0, 2.0]);
3814        let heap_size = histogram.custom_values.capacity() * size_of::<f64>()
3815            + histogram.positive_spans.capacity() * size_of::<Span>()
3816            + histogram.negative_spans.capacity() * size_of::<Span>()
3817            + histogram.positive_buckets.capacity() * size_of::<f64>()
3818            + histogram.negative_buckets.capacity() * size_of::<f64>();
3819        let mut accumulator =
3820            NativeHistogramAggregateAccumulator::new(NativeHistogramAggregateKind::Sum, None);
3821        let empty_size = accumulator.size();
3822
3823        accumulator.push_histogram(histogram, 1).unwrap();
3824
3825        assert!(heap_size > 0);
3826        assert_eq!(accumulator.size(), empty_size + heap_size);
3827    }
3828
3829    #[test]
3830    fn absent_over_time_handles_histogram_ranges() {
3831        let values = vec![Some(sample_histogram(1.0, 1.0, vec![1.0]))];
3832        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([Some(1000)]));
3833        let histograms = build_histogram_array(&values);
3834        let ranges = [(0, 1), (0, 0)];
3835        let ts_range = RangeArray::from_ranges(timestamps, ranges).unwrap();
3836        let value_range = RangeArray::from_ranges(histograms, ranges).unwrap();
3837
3838        let result = run_udf(
3839            NativeHistogramAbsentOverTime::scalar_udf(),
3840            vec![
3841                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
3842                ColumnarValue::Array(Arc::new(value_range.into_dict())),
3843            ],
3844            DataType::Float64,
3845        );
3846        let result = extract_array(&result).unwrap();
3847        let result = result.as_any().downcast_ref::<Float64Array>().unwrap();
3848
3849        assert!(result.is_null(0));
3850        assert_eq!(result.value(1), 1.0);
3851    }
3852
3853    #[test]
3854    fn delta_requires_exact_layout() {
3855        let first = sample_histogram(2.0, 3.0, vec![1.0, 1.0]);
3856        let last = sample_histogram(5.0, 8.0, vec![2.0, 3.0]);
3857        let delta = histogram_delta(&[first, last], &[0, 1], false).unwrap();
3858        assert_eq!(delta.count, 3.0);
3859        assert_eq!(delta.sum, 5.0);
3860        assert_eq!(delta.reset_hint, GAUGE_RESET_HINT);
3861        assert_eq!(delta.positive_buckets, vec![1.0, 2.0]);
3862    }
3863
3864    #[test]
3865    fn reset_hint_shortcuts_detection() {
3866        let previous = sample_histogram(6.0, 10.0, vec![2.0, 4.0]);
3867
3868        let mut current = sample_histogram(7.0, 12.0, vec![3.0, 4.0]);
3869        current.reset_hint = COUNTER_RESET_HINT;
3870        assert!(current.detect_reset(&previous));
3871
3872        let mut current = sample_histogram(5.0, 8.0, vec![1.0, 4.0]);
3873        current.reset_hint = NOT_COUNTER_RESET_HINT;
3874        assert!(!current.detect_reset(&previous));
3875    }
3876
3877    #[test]
3878    fn start_timestamp_detects_counter_reset() {
3879        let first = sample_histogram(6.0, 10.0, vec![2.0, 4.0]);
3880        let mut last = sample_histogram(7.0, 12.0, vec![3.0, 4.0]);
3881        last.start_timestamp = Some(1500);
3882
3883        let delta = histogram_delta(&[first.clone(), last.clone()], &[1000, 2000], true).unwrap();
3884        assert_eq!(delta.count, 7.0);
3885        assert_eq!(delta.sum, 12.0);
3886
3887        let idelta = idelta_value(&[first, last], true, 1000, 2000, 1.0).unwrap();
3888        assert_eq!(idelta.count, 7.0);
3889        assert_eq!(idelta.sum, 12.0);
3890    }
3891
3892    #[test]
3893    fn extrapolated_rate_uses_start_timestamp_synthetic_zero() {
3894        let mut single = sample_histogram(1.0, 1.0, vec![1.0]);
3895        single.start_timestamp = Some(1_000);
3896
3897        let rate = extrapolated_histogram_result(
3898            NativeHistogramRate::scalar_udf(),
3899            vec![2_000],
3900            vec![single.clone()],
3901            3_000,
3902            3_000,
3903        )
3904        .unwrap();
3905        assert_eq!(rate.count, 1.0 / 3.0);
3906        assert_eq!(rate.sum, 1.0 / 3.0);
3907
3908        let increase = extrapolated_histogram_result(
3909            NativeHistogramIncrease::scalar_udf(),
3910            vec![2_000],
3911            vec![single],
3912            3_000,
3913            3_000,
3914        )
3915        .unwrap();
3916        assert_eq!(increase.count, 1.0);
3917        assert_eq!(increase.sum, 1.0);
3918
3919        let mut first = sample_histogram(2.0, 2.0, vec![2.0]);
3920        first.start_timestamp = Some(1_000);
3921        let last = sample_histogram(4.0, 4.0, vec![4.0]);
3922        let increase = extrapolated_histogram_result(
3923            NativeHistogramIncrease::scalar_udf(),
3924            vec![2_000, 3_000],
3925            vec![first, last],
3926            3_000,
3927            3_000,
3928        )
3929        .unwrap();
3930        assert_eq!(increase.count, 4.0);
3931        assert_eq!(increase.sum, 4.0);
3932    }
3933
3934    #[test]
3935    fn extrapolated_rate_requires_strictly_in_range_start_timestamp() {
3936        for (start_timestamp, range_end, range_length) in [
3937            (0, 3_000, 3_000),
3938            (500, 3_000, 2_000),
3939            (1_000, 3_000, 2_000),
3940            (2_000, 3_000, 3_000),
3941            (2_500, 3_000, 3_000),
3942        ] {
3943            let mut sample = sample_histogram(1.0, 1.0, vec![1.0]);
3944            sample.start_timestamp = Some(start_timestamp);
3945            assert!(
3946                extrapolated_histogram_result(
3947                    NativeHistogramRate::scalar_udf(),
3948                    vec![2_000],
3949                    vec![sample],
3950                    range_end,
3951                    range_length,
3952                )
3953                .is_none(),
3954                "start_timestamp={start_timestamp}"
3955            );
3956        }
3957
3958        let mut gauge = sample_histogram(1.0, 1.0, vec![1.0]);
3959        gauge.start_timestamp = Some(1_000);
3960        gauge.reset_hint = GAUGE_RESET_HINT;
3961        assert!(
3962            extrapolated_histogram_result(
3963                NativeHistogramDelta::scalar_udf(),
3964                vec![2_000],
3965                vec![gauge],
3966                3_000,
3967                3_000,
3968            )
3969            .is_none()
3970        );
3971    }
3972
3973    #[test]
3974    fn first_reset_ignores_incompatible_pre_reset_layout() {
3975        let first = sample_histogram(10.0, 10.0, vec![10.0]);
3976        let second = NativeHistogram {
3977            schema: CUSTOM_BUCKETS_SCHEMA,
3978            zero_threshold: 0.0,
3979            sum: 2.0,
3980            reset_hint: COUNTER_RESET_HINT,
3981            start_timestamp: None,
3982            custom_values: vec![1.0],
3983            positive_spans: vec![Span {
3984                offset: 0,
3985                length: 1,
3986            }],
3987            negative_spans: Vec::new(),
3988            count: 2.0,
3989            zero_count: 0.0,
3990            positive_buckets: vec![2.0],
3991            negative_buckets: Vec::new(),
3992        };
3993
3994        let delta = histogram_delta(&[first, second], &[1_000, 2_000], true).unwrap();
3995        assert_eq!(delta.schema, CUSTOM_BUCKETS_SCHEMA);
3996        assert_eq!(delta.count, 2.0);
3997        assert_eq!(delta.sum, 2.0);
3998        assert_eq!(delta.positive_buckets, vec![2.0]);
3999    }
4000
4001    #[test]
4002    fn counter_delta_handles_reset_segment_boundaries() {
4003        let first = sample_histogram(5.0, 5.0, vec![5.0]);
4004        let second = sample_histogram(7.0, 7.0, vec![7.0]);
4005        let mut reset = sample_histogram(2.0, 2.0, vec![2.0]);
4006        reset.reset_hint = COUNTER_RESET_HINT;
4007        let last = sample_histogram(4.0, 4.0, vec![4.0]);
4008
4009        let delta = histogram_delta(
4010            &[first, second, reset, last],
4011            &[1_000, 2_000, 3_000, 4_000],
4012            true,
4013        )
4014        .unwrap();
4015        assert_eq!(delta.count, 6.0);
4016        assert_eq!(delta.sum, 6.0);
4017        assert_eq!(delta.positive_buckets, vec![6.0]);
4018    }
4019}