1use 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
746pub 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
1436pub 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 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
1883fn 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
1906fn 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
1928fn 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#[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 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#[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, ¤t) {
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 (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}