1use std::collections::VecDeque;
16use std::sync::Arc;
17
18use common_macro::range_fn;
19use datafusion::arrow::array::{Float64Array, TimestampMillisecondArray};
20use datafusion::common::DataFusionError;
21use datafusion::logical_expr::{ScalarUDF, Volatility};
22use datafusion::physical_plan::ColumnarValue;
23use datatypes::arrow::array::Array;
24use datatypes::arrow::compute;
25use datatypes::arrow::datatypes::{DataType, TimeUnit};
26
27use crate::functions::{compensated_sum_inc, extract_array, extract_range_dict};
28use crate::range_array::{RangeArray, unpack};
29
30#[derive(Clone, Copy)]
31enum PresenceEvaluator {
32 Count,
33 Last,
34 Absent,
35 Present,
36}
37
38fn count_over_time_evaluator(
39 input: &[ColumnarValue],
40 name: &str,
41) -> Result<ColumnarValue, DataFusionError> {
42 evaluate_presence(input, name, PresenceEvaluator::Count)
43}
44
45fn last_over_time_evaluator(
46 input: &[ColumnarValue],
47 name: &str,
48) -> Result<ColumnarValue, DataFusionError> {
49 evaluate_presence(input, name, PresenceEvaluator::Last)
50}
51
52fn absent_over_time_evaluator(
53 input: &[ColumnarValue],
54 name: &str,
55) -> Result<ColumnarValue, DataFusionError> {
56 evaluate_presence(input, name, PresenceEvaluator::Absent)
57}
58
59fn present_over_time_evaluator(
60 input: &[ColumnarValue],
61 name: &str,
62) -> Result<ColumnarValue, DataFusionError> {
63 evaluate_presence(input, name, PresenceEvaluator::Present)
64}
65
66fn evaluate_presence(
67 input: &[ColumnarValue],
68 name: &str,
69 operation: PresenceEvaluator,
70) -> Result<ColumnarValue, DataFusionError> {
71 assert_eq!(input.len(), 2);
72
73 let timestamp_ranges = RangeArray::try_new(extract_array(&input[0])?.to_data().into())?;
74 let value_ranges = RangeArray::try_new(extract_array(&input[1])?.to_data().into())?;
75 let len = timestamp_ranges.len();
76 if len != value_ranges.len() {
77 return Err(DataFusionError::Execution(format!(
78 "RangeArray have different lengths in PromQL function {name}: array1={len}, array2={}",
79 value_ranges.len()
80 )));
81 }
82
83 if timestamp_ranges.is_empty() {
84 return Ok(ColumnarValue::Array(Arc::new(Float64Array::from_iter(
85 std::iter::empty::<Option<f64>>(),
86 ))));
87 }
88
89 assert!(
92 timestamp_ranges
93 .values()
94 .as_any()
95 .is::<TimestampMillisecondArray>()
96 );
97 let values = value_ranges
98 .values()
99 .as_any()
100 .downcast_ref::<Float64Array>()
101 .unwrap();
102 let evaluator: fn(&Float64Array, usize, usize) -> Option<f64> = if values.null_count() == 0 {
103 match operation {
104 PresenceEvaluator::Count => |_, _, length| (length != 0).then_some(length as f64),
105 PresenceEvaluator::Last => {
106 |values, offset, length| (length != 0).then(|| values.value(offset + length - 1))
107 }
108 PresenceEvaluator::Absent => |_, _, length| (length == 0).then_some(1.0),
109 PresenceEvaluator::Present => |_, _, length| (length != 0).then_some(1.0),
110 }
111 } else {
112 match operation {
113 PresenceEvaluator::Count => |values, offset, length| {
114 let count = (offset..offset + length)
115 .filter(|&index| values.is_valid(index))
116 .count();
117 (count != 0).then_some(count as f64)
118 },
119 PresenceEvaluator::Last => |values, offset, length| {
120 (offset..offset + length)
121 .rev()
122 .find(|&index| values.is_valid(index))
123 .map(|index| values.value(index))
124 },
125 PresenceEvaluator::Absent => |values, offset, length| {
126 (!(offset..offset + length).any(|index| values.is_valid(index))).then_some(1.0)
127 },
128 PresenceEvaluator::Present => |values, offset, length| {
129 (offset..offset + length)
130 .any(|index| values.is_valid(index))
131 .then_some(1.0)
132 },
133 }
134 };
135
136 let mut result = Vec::with_capacity(len);
137 for index in 0..len {
138 let (_, timestamp_length) = timestamp_ranges.get_offset_length(index).unwrap();
139 let (value_offset, value_length) = value_ranges.get_offset_length(index).unwrap();
140 if timestamp_length != value_length {
141 return Err(DataFusionError::Execution(format!(
142 "RangeArray's element {index} have different lengths in PromQL function {name}: array1={timestamp_length}, array2={value_length}"
143 )));
144 }
145 result.push(evaluator(values, value_offset, value_length));
146 }
147
148 Ok(ColumnarValue::Array(Arc::new(Float64Array::from_iter(
149 result,
150 ))))
151}
152
153#[range_fn(
155 name = AvgOverTime,
156 ret = Float64Array,
157 display_name = prom_avg_over_time
158)]
159pub fn avg_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
160 let sample_count = values.len() - values.null_count();
163 compute::sum(values).map(|result| result / sample_count as f64)
164}
165
166#[derive(Debug)]
168pub struct MinOverTime {}
169
170impl MinOverTime {
171 pub const fn name() -> &'static str {
172 "prom_min_over_time"
173 }
174
175 pub fn scalar_udf() -> ScalarUDF {
176 datafusion_expr::create_udf(
177 Self::name(),
178 Self::input_type(),
179 Self::return_type(),
180 Volatility::Volatile,
181 Arc::new(Self::calc) as _,
182 )
183 }
184
185 fn input_type() -> Vec<DataType> {
186 min_max_input_type()
187 }
188
189 fn return_type() -> DataType {
190 Float64Array::new_null(0).data_type().clone()
191 }
192
193 fn calc(input: &[ColumnarValue]) -> Result<ColumnarValue, DataFusionError> {
194 min_max_over_time_batch(input, true, Self::name())
195 }
196}
197
198#[derive(Debug)]
200pub struct MaxOverTime {}
201
202impl MaxOverTime {
203 pub const fn name() -> &'static str {
204 "prom_max_over_time"
205 }
206
207 pub fn scalar_udf() -> ScalarUDF {
208 datafusion_expr::create_udf(
209 Self::name(),
210 Self::input_type(),
211 Self::return_type(),
212 Volatility::Volatile,
213 Arc::new(Self::calc) as _,
214 )
215 }
216
217 fn input_type() -> Vec<DataType> {
218 min_max_input_type()
219 }
220
221 fn return_type() -> DataType {
222 Float64Array::new_null(0).data_type().clone()
223 }
224
225 fn calc(input: &[ColumnarValue]) -> Result<ColumnarValue, DataFusionError> {
226 min_max_over_time_batch(input, false, Self::name())
227 }
228}
229
230fn min_max_input_type() -> Vec<DataType> {
231 vec![
232 RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
233 RangeArray::convert_data_type(DataType::Float64),
234 ]
235}
236
237const MIN_SLIDING_WINDOWS: usize = 4;
239const MIN_SLIDING_WINDOW_LENGTH: u64 = 32;
242const MAX_SLIDING_STEP_FRACTION: u64 = 4;
244
245fn reuses_candidates(window_keys: &[i64]) -> bool {
258 if window_keys.len() < MIN_SLIDING_WINDOWS {
259 return false;
260 }
261
262 let windows = window_keys.len() as u64;
263 let mut total_length = 0u64;
264 let mut lowest_offset = u32::MAX;
265 let mut highest_offset = 0u32;
266 for &key in window_keys {
267 let (offset, length) = unpack(key);
268 if length == 0 {
272 continue;
273 }
274 total_length += u64::from(length);
275 lowest_offset = lowest_offset.min(offset);
276 highest_offset = highest_offset.max(offset);
277 }
278 let total_advance = u64::from(highest_offset.saturating_sub(lowest_offset));
280
281 total_length >= windows * MIN_SLIDING_WINDOW_LENGTH
286 && total_advance * MAX_SLIDING_STEP_FRACTION * windows <= total_length * (windows - 1)
287}
288
289fn is_better(value: f64, current: f64, is_min: bool) -> bool {
290 if is_min {
291 value < current
292 } else {
293 value > current
294 }
295}
296
297fn min_max_over_time_batch(
298 input: &[ColumnarValue],
299 is_min: bool,
300 func_name: &str,
301) -> Result<ColumnarValue, DataFusionError> {
302 if input.len() != 2 {
303 return Err(DataFusionError::Execution(format!(
304 "{func_name}: expected 2 inputs, found {}",
305 input.len()
306 )));
307 }
308
309 let timestamps = extract_range_dict(
310 &input[0],
311 func_name,
312 "timestamp range vector",
313 &DataType::Timestamp(TimeUnit::Millisecond, None),
314 )?;
315 let values = extract_range_dict(
316 &input[1],
317 func_name,
318 "value range vector",
319 &DataType::Float64,
320 )?;
321
322 let timestamp_keys = timestamps.keys().values();
323 let value_keys = values.keys().values();
324 if timestamp_keys.len() != value_keys.len() {
325 return Err(DataFusionError::Execution(format!(
326 "{func_name}: timestamp and value ranges should have the same number of windows, found {} and {}",
327 timestamp_keys.len(),
328 value_keys.len()
329 )));
330 }
331
332 let values = values
333 .values()
334 .as_any()
335 .downcast_ref::<Float64Array>()
336 .ok_or_else(|| {
337 DataFusionError::Execution(format!(
338 "{func_name}: expect value range vector values of type Float64"
339 ))
340 })?;
341 let mut sliding = reuses_candidates(value_keys).then(|| SlidingExtrema::new(is_min));
342 let mut result = Vec::with_capacity(value_keys.len());
343
344 for index in 0..value_keys.len() {
345 let (_, timestamp_length) = unpack(timestamp_keys[index]);
346 let (value_offset, value_length) = unpack(value_keys[index]);
347 if timestamp_length != value_length {
348 return Err(DataFusionError::Execution(format!(
349 "{func_name}: timestamp and value ranges have different lengths at window {index}: {timestamp_length} and {value_length}"
350 )));
351 }
352
353 let start = value_offset as usize;
354 let end = start + value_length as usize;
355 result.push(match sliding.as_mut() {
356 Some(sliding) => sliding.evaluate(values, start, end),
357 None => scan_extremum(values, start, end, is_min),
358 });
359 }
360
361 Ok(ColumnarValue::Array(Arc::new(Float64Array::from_iter(
362 result,
363 ))))
364}
365
366struct SlidingExtrema {
368 is_min: bool,
369 candidates: VecDeque<usize>,
374 latest_nan: Option<usize>,
376 previous_window: Option<(usize, usize)>,
377}
378
379impl SlidingExtrema {
380 fn new(is_min: bool) -> Self {
381 Self {
382 is_min,
383 candidates: VecDeque::new(),
384 latest_nan: None,
385 previous_window: None,
386 }
387 }
388
389 fn evaluate(&mut self, values: &Float64Array, start: usize, end: usize) -> Option<f64> {
390 let append_start = match self.previous_window {
391 Some((previous_start, previous_end))
392 if start >= previous_start && end >= previous_end =>
393 {
394 while self.candidates.front().is_some_and(|&index| index < start) {
395 self.candidates.pop_front();
396 }
397 if self.latest_nan.is_some_and(|index| index < start) {
398 self.latest_nan = None;
399 }
400 previous_end.max(start)
403 }
404 Some(_) | None => {
406 self.candidates.clear();
407 self.latest_nan = None;
408 start
409 }
410 };
411 self.append(values, append_start, end);
412 self.previous_window = Some((start, end));
413
414 self.candidates
415 .front()
416 .or(self.latest_nan.as_ref())
417 .map(|&index| values.value(index))
418 }
419
420 fn append(&mut self, values: &Float64Array, start: usize, end: usize) {
421 for index in start..end {
422 if values.is_null(index) {
423 continue;
424 }
425
426 let value = values.value(index);
427 if value.is_nan() {
428 self.latest_nan = Some(index);
430 continue;
431 }
432
433 while let Some(&tail_index) = self.candidates.back() {
434 if !is_better(value, values.value(tail_index), self.is_min) {
435 break;
436 }
437 self.candidates.pop_back();
438 }
439 self.candidates.push_back(index);
440 }
441 }
442}
443
444fn scan_extremum(values: &Float64Array, start: usize, end: usize, is_min: bool) -> Option<f64> {
446 let mut extremum: Option<f64> = None;
447 for index in start..end {
448 if values.is_null(index) {
449 continue;
450 }
451
452 let value = values.value(index);
453 let replace = match extremum {
454 None => true,
455 Some(current) => current.is_nan() || is_better(value, current, is_min),
457 };
458 if replace {
459 extremum = Some(value);
460 }
461 }
462 extremum
463}
464
465#[cfg(test)]
466fn min_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
467 let mut valid_values = values.iter().flatten();
468 let mut min = valid_values.next()?;
469 for value in valid_values {
470 if value < min || min.is_nan() {
471 min = value;
472 }
473 }
474 Some(min)
475}
476
477#[cfg(test)]
478fn max_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
479 let mut valid_values = values.iter().flatten();
480 let mut max = valid_values.next()?;
481 for value in valid_values {
482 if value > max || max.is_nan() {
483 max = value;
484 }
485 }
486 Some(max)
487}
488
489#[range_fn(
491 name = SumOverTime,
492 ret = Float64Array,
493 display_name = prom_sum_over_time
494)]
495pub fn sum_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
496 compute::sum(values)
497}
498
499#[range_fn(
501 name = CountOverTime,
502 ret = Float64Array,
503 display_name = prom_count_over_time,
504 evaluator = count_over_time_evaluator
505)]
506pub fn count_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
507 let count = values.iter().flatten().count();
508 (count != 0).then_some(count as f64)
509}
510
511#[range_fn(
513 name = LastOverTime,
514 ret = Float64Array,
515 display_name = prom_last_over_time,
516 evaluator = last_over_time_evaluator
517)]
518pub fn last_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
519 values.iter().flatten().last()
520}
521
522#[range_fn(
526 name = AbsentOverTime,
527 ret = Float64Array,
528 display_name = prom_absent_over_time,
529 evaluator = absent_over_time_evaluator
530)]
531pub fn absent_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
532 values.iter().flatten().next().is_none().then_some(1.0)
533}
534
535#[range_fn(
537 name = PresentOverTime,
538 ret = Float64Array,
539 display_name = prom_present_over_time,
540 evaluator = present_over_time_evaluator
541)]
542pub fn present_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
543 values.iter().flatten().next().is_some().then_some(1.0)
544}
545
546#[range_fn(
550 name = StdvarOverTime,
551 ret = Float64Array,
552 display_name = prom_stdvar_over_time
553)]
554pub fn stdvar_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
555 let mut count = 0;
556 let mut mean: f64 = 0.0;
557 let mut result: f64 = 0.0;
558 for value in values.iter().flatten() {
559 let new_count = count + 1;
560 let delta1 = value - mean;
561 let new_mean = delta1 / new_count as f64 + mean;
562 let delta2 = value - new_mean;
563 let new_result = result + delta1 * delta2;
564
565 count = new_count;
566 mean = new_mean;
567 result = new_result;
568 }
569 (count > 0).then(|| result / count as f64)
570}
571
572#[range_fn(
575 name = StddevOverTime,
576 ret = Float64Array,
577 display_name = prom_stddev_over_time
578)]
579pub fn stddev_over_time(_: &TimestampMillisecondArray, values: &Float64Array) -> Option<f64> {
580 let mut count = 0.0;
581 let mut mean = 0.0;
582 let mut comp_mean = 0.0;
583 let mut deviations_sum_sq = 0.0;
584 let mut comp_deviations_sum_sq = 0.0;
585 for current_value in values.iter().flatten() {
586 count += 1.0;
587 let delta = current_value - (mean + comp_mean);
588 let (new_mean, new_comp_mean) = compensated_sum_inc(delta / count, mean, comp_mean);
589 mean = new_mean;
590 comp_mean = new_comp_mean;
591 let (new_deviations_sum_sq, new_comp_deviations_sum_sq) = compensated_sum_inc(
592 delta * (current_value - (mean + comp_mean)),
593 deviations_sum_sq,
594 comp_deviations_sum_sq,
595 );
596 deviations_sum_sq = new_deviations_sum_sq;
597 comp_deviations_sum_sq = new_comp_deviations_sum_sq;
598 }
599 (count > 0.0).then(|| ((deviations_sum_sq + comp_deviations_sum_sq) / count).sqrt())
600}
601
602#[cfg(test)]
603mod test {
604 use datafusion::arrow::buffer::NullBuffer;
605
606 use super::*;
607 use crate::functions::test_util::{invoke_range_udf, simple_range_udf_runner};
608
609 fn assert_option_bits(actual: &[Option<f64>], expected: &[Option<f64>]) {
610 assert_eq!(actual.len(), expected.len());
611 for (actual, expected) in actual.iter().zip(expected) {
612 match (actual, expected) {
613 (Some(actual), Some(expected)) => assert_eq!(actual.to_bits(), expected.to_bits()),
614 (None, None) => {}
615 (actual, expected) => panic!("expected {expected:?}, got {actual:?}"),
616 }
617 }
618 }
619
620 fn old_min_max_over_windows(
621 values: &Float64Array,
622 ranges: &[(u32, u32)],
623 is_min: bool,
624 ) -> Vec<Option<f64>> {
625 ranges
626 .iter()
627 .map(|&(offset, length)| {
628 let values = values.slice(offset as usize, length as usize);
629 let values = values.as_any().downcast_ref::<Float64Array>().unwrap();
630 let timestamps = TimestampMillisecondArray::new_null(values.len());
631 if is_min {
632 min_over_time(×tamps, values)
633 } else {
634 max_over_time(×tamps, values)
635 }
636 })
637 .collect()
638 }
639
640 fn run_min_max_udf(
641 udf: ScalarUDF,
642 timestamps: RangeArray,
643 values: RangeArray,
644 ) -> Vec<Option<f64>> {
645 let result = invoke_range_udf(udf, timestamps, values).unwrap();
646 extract_array(&result)
647 .unwrap()
648 .as_any()
649 .downcast_ref::<Float64Array>()
650 .unwrap()
651 .iter()
652 .collect()
653 }
654
655 fn slice_range_array(range: RangeArray, offset: usize, length: usize) -> RangeArray {
656 RangeArray::try_new(range.into_dict().slice(offset, length).to_data().into()).unwrap()
657 }
658
659 fn sliced_float_values(values: Vec<Option<f64>>) -> Float64Array {
660 let mut backing = vec![Some(1234.0)];
661 backing.extend(values);
662 let length = backing.len() - 1;
663 let backing = Float64Array::from(backing);
664 backing
665 .slice(1, length)
666 .as_any()
667 .downcast_ref::<Float64Array>()
668 .unwrap()
669 .clone()
670 }
671
672 fn min_max_range_inputs(
673 timestamps: &TimestampMillisecondArray,
674 values: &Float64Array,
675 timestamp_ranges: &[(u32, u32)],
676 value_ranges: &[(u32, u32)],
677 ) -> (RangeArray, RangeArray) {
678 let timestamp_prefix = [(0, 1)];
679 let value_prefix = [(0, 1)];
680 let timestamps = RangeArray::from_ranges(
681 Arc::new(timestamps.clone()),
682 timestamp_prefix
683 .into_iter()
684 .chain(timestamp_ranges.iter().copied()),
685 )
686 .unwrap();
687 let values = RangeArray::from_ranges(
688 Arc::new(values.clone()),
689 value_prefix.into_iter().chain(value_ranges.iter().copied()),
690 )
691 .unwrap();
692 (
693 slice_range_array(timestamps, 1, timestamp_ranges.len()),
694 slice_range_array(values, 1, value_ranges.len()),
695 )
696 }
697
698 fn window_keys(ranges: &[(u32, u32)]) -> Vec<i64> {
699 let length = ranges
700 .iter()
701 .map(|&(offset, length)| (offset + length) as usize)
702 .max()
703 .unwrap_or_default();
704 RangeArray::from_ranges(
705 Arc::new(Float64Array::from(vec![0.0; length])),
706 ranges.iter().copied(),
707 )
708 .unwrap()
709 .into_dict()
710 .keys()
711 .values()
712 .to_vec()
713 }
714
715 fn assert_min_max_udfs_match_oracle(
716 timestamp_values: &TimestampMillisecondArray,
717 all_values: &Float64Array,
718 timestamp_ranges: &[(u32, u32)],
719 value_ranges: &[(u32, u32)],
720 ) {
721 let expected_min = old_min_max_over_windows(all_values, value_ranges, true);
722 let expected_max = old_min_max_over_windows(all_values, value_ranges, false);
723 let (timestamps, values) =
724 min_max_range_inputs(timestamp_values, all_values, timestamp_ranges, value_ranges);
725 assert_option_bits(
726 &run_min_max_udf(MinOverTime::scalar_udf(), timestamps, values),
727 &expected_min,
728 );
729 let (timestamps, values) =
730 min_max_range_inputs(timestamp_values, all_values, timestamp_ranges, value_ranges);
731 assert_option_bits(
732 &run_min_max_udf(MaxOverTime::scalar_udf(), timestamps, values),
733 &expected_max,
734 );
735 }
736
737 #[test]
738 fn min_max_over_time_batch_preserves_bits_across_irregular_windows() {
739 let first_nan = f64::from_bits(0x7ff8_0000_0000_00a1);
740 let second_nan = f64::from_bits(0x7ff8_0000_0000_00b2);
741 let third_nan = f64::from_bits(0x7ff8_0000_0000_00c3);
742 let values = sliced_float_values(vec![
743 None,
744 Some(first_nan),
745 Some(-0.0),
746 Some(0.0),
747 Some(2.0),
748 Some(2.0),
749 Some(f64::INFINITY),
750 Some(f64::NEG_INFINITY),
751 Some(second_nan),
752 None,
753 Some(3.0),
754 Some(third_nan),
755 ]);
756 let timestamps = TimestampMillisecondArray::from_iter((0..16).map(Some));
757 let value_ranges = [
758 (0, 0),
759 (0, 2),
760 (0, 4),
761 (1, 4),
762 (4, 4),
763 (4, 4),
764 (8, 2),
765 (5, 2),
766 (3, 3),
767 (11, 1),
768 ];
769 let timestamp_ranges = [
770 (8, 0),
771 (7, 2),
772 (6, 4),
773 (5, 4),
774 (4, 4),
775 (4, 4),
776 (3, 2),
777 (2, 2),
778 (1, 3),
779 (0, 1),
780 ];
781
782 assert_min_max_udfs_match_oracle(×tamps, &values, ×tamp_ranges, &value_ranges);
783 }
784
785 #[test]
786 fn min_max_over_time_batch_preserves_first_signed_zero_in_both_orders() {
787 let timestamps = TimestampMillisecondArray::from(vec![0, 1, 2, 3]);
788 for (values, expected) in [
789 (Float64Array::from(vec![-0.0, 0.0]), -0.0),
790 (Float64Array::from(vec![0.0, -0.0]), 0.0),
791 ] {
792 let (timestamp_ranges, value_ranges) =
793 min_max_range_inputs(×tamps, &values, &[(2, 2)], &[(0, 2)]);
794 assert_option_bits(
795 &run_min_max_udf(MinOverTime::scalar_udf(), timestamp_ranges, value_ranges),
796 &[Some(expected)],
797 );
798 let (timestamp_ranges, value_ranges) =
799 min_max_range_inputs(×tamps, &values, &[(2, 2)], &[(0, 2)]);
800 assert_option_bits(
801 &run_min_max_udf(MaxOverTime::scalar_udf(), timestamp_ranges, value_ranges),
802 &[Some(expected)],
803 );
804 }
805 }
806
807 #[test]
808 fn min_max_over_time_batch_rebuilds_after_right_bound_retreat() {
809 let timestamps = TimestampMillisecondArray::from(vec![0, 1, 2]);
810 let ranges = [(0, 3), (0, 2)];
811 assert_min_max_udfs_match_oracle(
812 ×tamps,
813 &Float64Array::from(vec![3.0, 2.0, 1.0]),
814 &ranges,
815 &ranges,
816 );
817 assert_min_max_udfs_match_oracle(
818 ×tamps,
819 &Float64Array::from(vec![1.0, 2.0, 3.0]),
820 &ranges,
821 &ranges,
822 );
823 }
824
825 #[test]
826 fn min_max_over_time_batch_returns_null_after_expiring_last_nan() {
827 let nan = f64::from_bits(0x7ff8_0000_0000_00d4);
828 let timestamps = TimestampMillisecondArray::from(vec![0, 1]);
829 let values = Float64Array::from(vec![Some(nan), None]);
830 let ranges = [(0, 1), (1, 1)];
831 let (timestamp_ranges, value_ranges) =
832 min_max_range_inputs(×tamps, &values, &ranges, &ranges);
833 assert_option_bits(
834 &run_min_max_udf(MinOverTime::scalar_udf(), timestamp_ranges, value_ranges),
835 &[Some(nan), None],
836 );
837 let (timestamp_ranges, value_ranges) =
838 min_max_range_inputs(×tamps, &values, &ranges, &ranges);
839 assert_option_bits(
840 &run_min_max_udf(MaxOverTime::scalar_udf(), timestamp_ranges, value_ranges),
841 &[Some(nan), None],
842 );
843 }
844
845 #[test]
846 fn min_max_over_time_batch_matches_scalar_oracle_for_random_windows() {
847 let values = sliced_float_values(
848 (0..64)
849 .map(|index| match index % 11 {
850 0 => None,
851 1 => Some(f64::from_bits(0x7ff8_0000_0000_0100 + index)),
852 2 => Some(-0.0),
853 3 => Some(0.0),
854 4 => Some(f64::INFINITY),
855 5 => Some(f64::NEG_INFINITY),
856 _ => Some((index as f64 * 17.0).sin()),
857 })
858 .collect(),
859 );
860 let timestamps = TimestampMillisecondArray::from_iter((0..128).map(Some));
861 let mut seed = 0x5eed_u64;
862 let mut value_ranges = Vec::new();
863 let mut timestamp_ranges = Vec::new();
864 for _ in 0..256 {
865 seed = seed.wrapping_mul(6364136223846793005).wrapping_add(1);
866 let start = (seed as usize) % values.len();
867 let length = ((seed >> 32) as usize) % (values.len() - start + 1);
868 value_ranges.push((start as u32, length as u32));
869 timestamp_ranges.push(((values.len() - start) as u32, length as u32));
870 }
871
872 assert_min_max_udfs_match_oracle(×tamps, &values, ×tamp_ranges, &value_ranges);
873 }
874
875 #[test]
876 fn min_max_over_time_batch_validates_inputs() {
877 let error = MinOverTime::calc(&[]).unwrap_err();
878 assert!(error.to_string().contains("expected 2 inputs, found 0"));
879
880 let timestamps =
881 RangeArray::from_ranges(Arc::new(TimestampMillisecondArray::from(vec![0])), [(0, 1)])
882 .unwrap();
883 let error = MaxOverTime::calc(&[
884 ColumnarValue::Array(Arc::new(timestamps.into_dict())),
885 ColumnarValue::Array(Arc::new(datatypes::arrow::array::Int64Array::from(vec![1]))),
886 ])
887 .unwrap_err();
888 assert!(
889 error
890 .to_string()
891 .contains("expect value range vector as DictionaryArray<Int64>")
892 );
893
894 let null_keys = datatypes::arrow::array::Int64Array::from_iter([Some(0), None]);
895 let null_key_dict = datatypes::arrow::array::DictionaryArray::<
896 datatypes::arrow::datatypes::Int64Type,
897 >::try_new(
898 null_keys,
899 Arc::new(TimestampMillisecondArray::from(vec![0, 1])),
900 )
901 .unwrap();
902 let error = MinOverTime::calc(&[
903 ColumnarValue::Array(Arc::new(null_key_dict)),
904 ColumnarValue::Array(Arc::new(datatypes::arrow::array::Int64Array::from(vec![1]))),
905 ])
906 .unwrap_err();
907 assert!(error.to_string().contains("Empty range is not expected"));
908 }
909
910 #[test]
911 fn min_max_over_time_batch_rejects_mismatched_windows() {
912 let timestamps = Arc::new(TimestampMillisecondArray::from(vec![0, 1, 2]));
913 let values = Arc::new(Float64Array::from(vec![1.0, 2.0, 3.0]));
914 let timestamp_ranges = RangeArray::from_ranges(timestamps.clone(), [(0, 1)]).unwrap();
915 let value_ranges = RangeArray::from_ranges(values.clone(), [(0, 1), (1, 1)]).unwrap();
916 let error = invoke_range_udf(MinOverTime::scalar_udf(), timestamp_ranges, value_ranges)
917 .unwrap_err();
918 assert!(error.to_string().contains("same number of windows"));
919
920 let timestamp_ranges = RangeArray::from_ranges(timestamps, [(0, 1)]).unwrap();
921 let value_ranges = RangeArray::from_ranges(values, [(1, 2)]).unwrap();
922 let error = invoke_range_udf(MaxOverTime::scalar_udf(), timestamp_ranges, value_ranges)
923 .unwrap_err();
924 assert!(error.to_string().contains("different lengths at window 0"));
925 }
926
927 #[test]
928 fn sliding_evaluator_is_selected_only_for_overlapping_window_batches() {
929 assert!(reuses_candidates(&window_keys(&[
931 (0, 32),
932 (8, 32),
933 (16, 32),
934 (24, 32)
935 ])));
936
937 assert!(!reuses_candidates(&window_keys(&[
939 (0, 32),
940 (8, 32),
941 (16, 32)
942 ])));
943 assert!(!reuses_candidates(&window_keys(&[
945 (0, 31),
946 (7, 31),
947 (14, 31),
948 (21, 31)
949 ])));
950 assert!(!reuses_candidates(&window_keys(&[
952 (0, 32),
953 (9, 32),
954 (18, 32),
955 (27, 32)
956 ])));
957 }
958
959 #[test]
960 fn sliding_evaluator_survives_an_unrepresentative_leading_window() {
961 let overlapping = (0..20u32).map(|index| (index * 8, 40));
966 for leading in [(0, 1), (0, 0)] {
967 let ranges = std::iter::once(leading)
968 .chain(overlapping.clone())
969 .collect::<Vec<_>>();
970 assert!(reuses_candidates(&window_keys(&ranges)));
971 }
972 }
973
974 #[test]
975 fn empty_windows_do_not_shorten_the_measured_advance() {
976 let disjoint = (0..4u32)
980 .map(|index| (index * 240, 240))
981 .collect::<Vec<_>>();
982 assert!(!reuses_candidates(&window_keys(&disjoint)));
983 let with_trailing_empty = [disjoint.as_slice(), &[(0, 0)]].concat();
984 assert!(!reuses_candidates(&window_keys(&with_trailing_empty)));
985 }
986
987 #[test]
988 fn min_max_over_time_batch_matches_oracle_on_overlapping_window_batches() {
989 let values = sliced_float_values(
990 (0..96)
991 .map(|index| match index % 13 {
992 0 => None,
993 1 => Some(f64::from_bits(0x7ff8_0000_0000_0200 + index as u64)),
994 2 => Some(-0.0),
995 3 => Some(0.0),
996 4 => Some(f64::INFINITY),
997 5 => Some(f64::NEG_INFINITY),
998 6 | 7 => Some(5.0),
999 _ => Some(((index * 37) % 23) as f64),
1000 })
1001 .collect(),
1002 );
1003 let timestamps = TimestampMillisecondArray::from_iter((0..96).map(Some));
1004 let ranges = [
1007 (0, 32),
1008 (4, 32),
1009 (8, 32),
1010 (12, 32),
1011 (16, 40),
1012 (20, 36),
1013 (20, 36),
1014 (64, 32),
1015 (64, 0),
1016 (60, 20),
1017 (0, 96),
1018 ];
1019 assert!(reuses_candidates(&window_keys(&ranges)));
1020
1021 assert_min_max_udfs_match_oracle(×tamps, &values, &ranges, &ranges);
1022 }
1023
1024 #[test]
1025 fn both_extrema_paths_match_the_scalar_oracle_on_every_four_sample_window() {
1026 let alphabet = [
1027 None,
1028 Some(-0.0),
1029 Some(0.0),
1030 Some(-3.0),
1031 Some(2.0),
1032 Some(f64::INFINITY),
1033 Some(f64::NEG_INFINITY),
1034 Some(f64::from_bits(0x7ff8_0000_0000_0042)),
1035 Some(f64::from_bits(0xfff8_0000_0000_0066)),
1036 ];
1037 let windows = (0..=4)
1040 .flat_map(|start| (0..=4 - start).map(move |length| (start, start + length)))
1041 .collect::<Vec<_>>();
1042
1043 for encoded in 0..alphabet.len().pow(4) {
1044 let mut remaining = encoded;
1045 let values = Float64Array::from(
1046 (0..4)
1047 .map(|_| {
1048 let value = alphabet[remaining % alphabet.len()];
1049 remaining /= alphabet.len();
1050 value
1051 })
1052 .collect::<Vec<_>>(),
1053 );
1054 let mut sliding_min = SlidingExtrema::new(true);
1055 let mut sliding_max = SlidingExtrema::new(false);
1056
1057 for &(start, end) in windows.iter().chain(windows.iter().rev()) {
1058 let window = values.slice(start, end - start);
1059 let timestamps = TimestampMillisecondArray::new_null(window.len());
1060
1061 let expected = min_over_time(×tamps, &window).map(f64::to_bits);
1062 assert_eq!(
1063 sliding_min.evaluate(&values, start, end).map(f64::to_bits),
1064 expected,
1065 "sliding min, values {encoded}, window {start}..{end}"
1066 );
1067 assert_eq!(
1068 scan_extremum(&values, start, end, true).map(f64::to_bits),
1069 expected,
1070 "scanned min, values {encoded}, window {start}..{end}"
1071 );
1072
1073 let expected = max_over_time(×tamps, &window).map(f64::to_bits);
1074 assert_eq!(
1075 sliding_max.evaluate(&values, start, end).map(f64::to_bits),
1076 expected,
1077 "sliding max, values {encoded}, window {start}..{end}"
1078 );
1079 assert_eq!(
1080 scan_extremum(&values, start, end, false).map(f64::to_bits),
1081 expected,
1082 "scanned max, values {encoded}, window {start}..{end}"
1083 );
1084 }
1085 }
1086 }
1087
1088 fn assert_over_time_value(actual: Option<f64>, expected: Option<f64>) {
1089 match (actual, expected) {
1090 (Some(actual), Some(expected)) if expected.is_nan() => assert!(actual.is_nan()),
1091 (Some(actual), Some(expected)) => assert_eq!(actual, expected),
1092 (None, None) => {}
1093 (actual, expected) => panic!("expected {expected:?}, got {actual:?}"),
1094 }
1095 }
1096
1097 fn assert_min_max(
1098 values: Vec<Option<f64>>,
1099 expected_min: Option<f64>,
1100 expected_max: Option<f64>,
1101 ) {
1102 let timestamps = TimestampMillisecondArray::from(vec![0; values.len()]);
1103 let values = Float64Array::from(values);
1104
1105 assert_over_time_value(min_over_time(×tamps, &values), expected_min);
1106 assert_over_time_value(max_over_time(×tamps, &values), expected_max);
1107 }
1108
1109 fn special_ranges() -> (RangeArray, RangeArray) {
1110 use datafusion::arrow::buffer::NullBuffer;
1111
1112 let timestamps = Arc::new(TimestampMillisecondArray::from_iter_values(0..10)).slice(1, 8);
1113 let values = Arc::new(Float64Array::new(
1114 vec![
1115 99.0,
1116 -0.0,
1117 0.0,
1118 f64::INFINITY,
1119 f64::NEG_INFINITY,
1120 f64::from_bits(0x7ff8_0000_0000_0042),
1121 f64::from_bits(0x7ff8_0000_0000_0066),
1122 3.0,
1123 -7.0,
1124 42.0,
1125 ]
1126 .into(),
1127 Some(NullBuffer::from(vec![
1128 true, true, true, true, true, false, true, true, true, true,
1129 ])),
1130 ))
1131 .slice(1, 8);
1132 let timestamp_ranges = [
1135 (0, 0),
1136 (1, 1),
1137 (2, 1),
1138 (3, 1),
1139 (4, 1),
1140 (5, 1),
1141 (6, 1),
1142 (0, 3),
1143 (3, 2),
1144 (4, 2),
1145 (6, 2),
1146 (5, 1),
1147 (2, 1),
1148 ];
1149 let value_ranges = [
1150 (7, 0),
1151 (0, 1),
1152 (1, 1),
1153 (2, 1),
1154 (3, 1),
1155 (4, 1),
1156 (5, 1),
1157 (0, 3),
1158 (3, 2),
1159 (4, 2),
1160 (1, 2),
1161 (5, 1),
1162 (1, 1),
1163 ];
1164
1165 let timestamp_ranges = RangeArray::from_ranges(
1166 Arc::new(timestamps),
1167 std::iter::once((0, 0)).chain(timestamp_ranges),
1168 )
1169 .unwrap()
1170 .into_dict()
1171 .slice(1, timestamp_ranges.len());
1172 let value_ranges = RangeArray::from_ranges(
1173 Arc::new(values),
1174 std::iter::once((0, 0)).chain(value_ranges),
1175 )
1176 .unwrap()
1177 .into_dict()
1178 .slice(1, value_ranges.len());
1179
1180 (
1181 RangeArray::try_new(timestamp_ranges).unwrap(),
1182 RangeArray::try_new(value_ranges).unwrap(),
1183 )
1184 }
1185
1186 fn assert_specialized_matches_oracle(
1187 udf: ScalarUDF,
1188 kernel: fn(&TimestampMillisecondArray, &Float64Array) -> Option<f64>,
1189 ) {
1190 use crate::functions::test_util::invoke_range_udf;
1191
1192 let (timestamps, values) = special_ranges();
1193 let expected = (0..timestamps.len())
1194 .map(|index| {
1195 let timestamps = timestamps.get(index).unwrap();
1196 let values = values.get(index).unwrap();
1197 kernel(
1198 timestamps
1199 .as_any()
1200 .downcast_ref::<TimestampMillisecondArray>()
1201 .unwrap(),
1202 values.as_any().downcast_ref::<Float64Array>().unwrap(),
1203 )
1204 })
1205 .collect::<Vec<_>>();
1206 let (timestamps, values) = special_ranges();
1207 let output = invoke_range_udf(udf, timestamps, values).unwrap();
1208 let output_array = extract_array(&output).unwrap();
1209 let output = output_array
1210 .as_any()
1211 .downcast_ref::<Float64Array>()
1212 .unwrap();
1213
1214 assert_eq!(output.len(), expected.len());
1215 for (index, expected) in expected.into_iter().enumerate() {
1216 assert_eq!(output.is_null(index), expected.is_none());
1217 if let Some(expected) = expected {
1218 assert_eq!(output.value(index).to_bits(), expected.to_bits());
1219 }
1220 }
1221 }
1222
1223 #[test]
1224 fn specialized_presence_range_udfs_match_slice_oracles() {
1225 assert_specialized_matches_oracle(CountOverTime::scalar_udf(), count_over_time);
1226 assert_specialized_matches_oracle(LastOverTime::scalar_udf(), last_over_time);
1227 assert_specialized_matches_oracle(AbsentOverTime::scalar_udf(), absent_over_time);
1228 assert_specialized_matches_oracle(PresentOverTime::scalar_udf(), present_over_time);
1229 }
1230
1231 #[test]
1232 fn specialized_presence_range_udfs_ignore_null_samples() {
1233 use datafusion::arrow::buffer::NullBuffer;
1234
1235 let make_ranges = || {
1236 let timestamps =
1237 Arc::new(TimestampMillisecondArray::from_iter_values(0..9)).slice(1, 7);
1238 let values = Arc::new(Float64Array::new(
1239 vec![99.0, 1.0, 0.0, 0.0, 4.0, f64::NAN, 0.0, -0.0, 42.0].into(),
1240 Some(NullBuffer::from(vec![
1241 true, true, false, false, true, true, true, true, true,
1242 ])),
1243 ))
1244 .slice(1, 7);
1245 let ranges = [(0, 4), (0, 3), (1, 2), (4, 1), (5, 2), (7, 0)];
1246
1247 (
1248 RangeArray::from_ranges(Arc::new(timestamps), ranges).unwrap(),
1249 RangeArray::from_ranges(Arc::new(values), ranges).unwrap(),
1250 )
1251 };
1252 let cases = [
1255 (
1256 CountOverTime::scalar_udf(),
1257 [Some(2.0), Some(1.0), None, Some(1.0), Some(2.0), None],
1258 ),
1259 (
1260 LastOverTime::scalar_udf(),
1261 [Some(4.0), Some(1.0), None, Some(f64::NAN), Some(-0.0), None],
1262 ),
1263 (
1264 AbsentOverTime::scalar_udf(),
1265 [None, None, Some(1.0), None, None, Some(1.0)],
1266 ),
1267 (
1268 PresentOverTime::scalar_udf(),
1269 [Some(1.0), Some(1.0), None, Some(1.0), Some(1.0), None],
1270 ),
1271 ];
1272
1273 for (udf, expected) in cases {
1274 let (timestamps, values) = make_ranges();
1275 let output =
1276 crate::functions::test_util::invoke_range_udf(udf, timestamps, values).unwrap();
1277 let output_array = extract_array(&output).unwrap();
1278 let output = output_array
1279 .as_any()
1280 .downcast_ref::<Float64Array>()
1281 .unwrap();
1282
1283 for (index, expected) in expected.into_iter().enumerate() {
1284 assert_eq!(output.is_valid(index), expected.is_some());
1285 if let Some(expected) = expected {
1286 assert_eq!(output.value(index).to_bits(), expected.to_bits());
1287 }
1288 }
1289 }
1290 }
1291
1292 #[test]
1293 fn specialized_presence_range_udfs_preserve_errors_and_metadata() {
1294 use datafusion::arrow::array::{DictionaryArray, Int64Array};
1295 use datafusion::arrow::datatypes::{Field, Int64Type};
1296 use datafusion_common::config::ConfigOptions;
1297 use datafusion_expr::ScalarFunctionArgs;
1298
1299 use crate::functions::test_util::{assert_execution_error, invoke_range_udf};
1300
1301 for (udf, name) in [
1302 (CountOverTime::scalar_udf(), "prom_count_over_time"),
1303 (LastOverTime::scalar_udf(), "prom_last_over_time"),
1304 (AbsentOverTime::scalar_udf(), "prom_absent_over_time"),
1305 (PresentOverTime::scalar_udf(), "prom_present_over_time"),
1306 ] {
1307 assert_eq!(udf.name(), name);
1308 assert_eq!(udf.signature().volatility, Volatility::Volatile);
1309 }
1310
1311 let timestamps = Arc::new(TimestampMillisecondArray::from_iter_values(0..3));
1312 let values = Arc::new(Float64Array::from_iter_values([1.0, 2.0, 3.0]));
1313 let error = invoke_range_udf(
1314 CountOverTime::scalar_udf(),
1315 RangeArray::from_ranges(timestamps.clone(), [(0, 1), (1, 1)]).unwrap(),
1316 RangeArray::from_ranges(values.clone(), [(0, 1)]).unwrap(),
1317 )
1318 .unwrap_err();
1319 assert_execution_error(
1320 error,
1321 "RangeArray have different lengths in PromQL function prom_count_over_time: array1=2, array2=1",
1322 );
1323
1324 let error = invoke_range_udf(
1325 CountOverTime::scalar_udf(),
1326 RangeArray::from_ranges(timestamps.clone(), [(0, 1), (1, 2)]).unwrap(),
1327 RangeArray::from_ranges(values.clone(), [(0, 1), (1, 1)]).unwrap(),
1328 )
1329 .unwrap_err();
1330 assert_execution_error(
1331 error,
1332 "RangeArray's element 1 have different lengths in PromQL function prom_count_over_time: array1=2, array2=1",
1333 );
1334
1335 let invoke_dict = |timestamps: DictionaryArray<Int64Type>,
1336 values: DictionaryArray<Int64Type>| {
1337 let args = vec![
1338 ColumnarValue::Array(Arc::new(timestamps)),
1339 ColumnarValue::Array(Arc::new(values)),
1340 ];
1341 CountOverTime::scalar_udf().invoke_with_args(ScalarFunctionArgs {
1342 arg_fields: args
1343 .iter()
1344 .enumerate()
1345 .map(|(index, value)| {
1346 Arc::new(Field::new(
1347 format!("c{index}"),
1348 value.data_type().clone(),
1349 true,
1350 ))
1351 })
1352 .collect(),
1353 args,
1354 number_rows: 1,
1355 return_field: Arc::new(Field::new("out", DataType::Float64, true)),
1356 config_options: Arc::new(ConfigOptions::default()),
1357 })
1358 };
1359 let empty = invoke_dict(
1360 DictionaryArray::new(
1361 Int64Array::from(Vec::<i64>::new()),
1362 Arc::new(Float64Array::from_iter_values([1.0])),
1363 ),
1364 DictionaryArray::new(
1365 Int64Array::from(Vec::<i64>::new()),
1366 Arc::new(TimestampMillisecondArray::from_iter_values([0])),
1367 ),
1368 )
1369 .unwrap();
1370 assert!(extract_array(&empty).unwrap().is_empty());
1371
1372 let null_keys = DictionaryArray::new(
1373 Int64Array::from(vec![None]),
1374 Arc::new(TimestampMillisecondArray::from_iter_values([0])),
1375 );
1376 let values_range = RangeArray::from_ranges(values.clone(), [(0, 1)]).unwrap();
1377 assert_eq!(
1378 invoke_dict(null_keys, values_range.into_dict())
1379 .unwrap_err()
1380 .to_string(),
1381 "External error: Empty range is not expected"
1382 );
1383
1384 let timestamps_range = RangeArray::from_ranges(timestamps, [(0, 1)]).unwrap();
1385 let invalid_values = unsafe { RangeArray::from_ranges_unchecked(values, [(2, 2)]) };
1386 assert_eq!(
1387 invoke_dict(timestamps_range.into_dict(), invalid_values.into_dict())
1388 .unwrap_err()
1389 .to_string(),
1390 "External error: Illegal range: offset 2, length 2, array len 3"
1391 );
1392
1393 let output = invoke_range_udf(
1394 SumOverTime::scalar_udf(),
1395 RangeArray::from_ranges(
1396 Arc::new(TimestampMillisecondArray::from_iter_values([0, 1])),
1397 [(0, 2)],
1398 )
1399 .unwrap(),
1400 RangeArray::from_ranges(
1401 Arc::new(Float64Array::from_iter_values([1.0, 2.0])),
1402 [(0, 2)],
1403 )
1404 .unwrap(),
1405 )
1406 .unwrap();
1407 assert_eq!(
1408 extract_array(&output)
1409 .unwrap()
1410 .as_any()
1411 .downcast_ref::<Float64Array>()
1412 .unwrap()
1413 .value(0),
1414 3.0
1415 );
1416 }
1417
1418 #[test]
1419 fn min_max_over_time_ignore_ordinary_nan_when_finite_values_exist() {
1420 let ordinary_nan = f64::from_bits(0x7ff8_0000_0000_0000);
1421
1422 assert_min_max(
1423 vec![Some(ordinary_nan), Some(3.0), Some(-2.0)],
1424 Some(-2.0),
1425 Some(3.0),
1426 );
1427 assert_min_max(
1428 vec![Some(3.0), Some(ordinary_nan), Some(-2.0)],
1429 Some(-2.0),
1430 Some(3.0),
1431 );
1432 assert_min_max(
1433 vec![Some(-2.0), Some(3.0), Some(ordinary_nan)],
1434 Some(-2.0),
1435 Some(3.0),
1436 );
1437 assert_min_max(
1438 vec![Some(ordinary_nan), Some(ordinary_nan)],
1439 Some(ordinary_nan),
1440 Some(ordinary_nan),
1441 );
1442 assert_min_max(
1443 vec![Some(3.0), Some(-2.0), Some(1.0)],
1444 Some(-2.0),
1445 Some(3.0),
1446 );
1447 assert_min_max(vec![], None, None);
1448 assert_min_max(vec![None, None], None, None);
1449 }
1450
1451 fn build_test_range_arrays() -> (RangeArray, RangeArray) {
1453 let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1454 [
1455 1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000, 17000, 200000, 500000,
1456 ]
1457 .into_iter()
1458 .map(Some),
1459 ));
1460 let ranges = [
1461 (0, 2),
1462 (0, 5),
1463 (1, 1), (2, 0), (2, 0), (3, 3),
1467 (4, 3),
1468 (5, 3),
1469 (8, 1), (9, 0), ];
1472
1473 let values_array = Arc::new(Float64Array::from_iter([
1474 12.345678, 87.654321, 31.415927, 27.182818, 70.710678, 41.421356, 57.735027, 69.314718,
1475 98.019802, 1.98019802, 61.803399,
1476 ]));
1477
1478 let ts_range_array = RangeArray::from_ranges(ts_array, ranges).unwrap();
1479 let value_range_array = RangeArray::from_ranges(values_array, ranges).unwrap();
1480
1481 (ts_range_array, value_range_array)
1482 }
1483
1484 #[test]
1485 fn calculate_avg_over_time() {
1486 let (ts_array, value_array) = build_test_range_arrays();
1487 simple_range_udf_runner(
1488 AvgOverTime::scalar_udf(),
1489 ts_array,
1490 value_array,
1491 vec![],
1492 vec![
1493 Some(49.9999995),
1494 Some(45.8618844),
1495 Some(87.654321),
1496 None,
1497 None,
1498 Some(46.438284),
1499 Some(56.62235366666667),
1500 Some(56.15703366666667),
1501 Some(98.019802),
1502 None,
1503 ],
1504 );
1505 }
1506
1507 #[test]
1508 fn calculate_min_over_time() {
1509 let (ts_array, value_array) = build_test_range_arrays();
1510 simple_range_udf_runner(
1511 MinOverTime::scalar_udf(),
1512 ts_array,
1513 value_array,
1514 vec![],
1515 vec![
1516 Some(12.345678),
1517 Some(12.345678),
1518 Some(87.654321),
1519 None,
1520 None,
1521 Some(27.182818),
1522 Some(41.421356),
1523 Some(41.421356),
1524 Some(98.019802),
1525 None,
1526 ],
1527 );
1528 }
1529
1530 #[test]
1531 fn calculate_max_over_time() {
1532 let (ts_array, value_array) = build_test_range_arrays();
1533 simple_range_udf_runner(
1534 MaxOverTime::scalar_udf(),
1535 ts_array,
1536 value_array,
1537 vec![],
1538 vec![
1539 Some(87.654321),
1540 Some(87.654321),
1541 Some(87.654321),
1542 None,
1543 None,
1544 Some(70.710678),
1545 Some(70.710678),
1546 Some(69.314718),
1547 Some(98.019802),
1548 None,
1549 ],
1550 );
1551 }
1552
1553 #[test]
1554 fn calculate_sum_over_time() {
1555 let (ts_array, value_array) = build_test_range_arrays();
1556 simple_range_udf_runner(
1557 SumOverTime::scalar_udf(),
1558 ts_array,
1559 value_array,
1560 vec![],
1561 vec![
1562 Some(99.999999),
1563 Some(229.309422),
1564 Some(87.654321),
1565 None,
1566 None,
1567 Some(139.314852),
1568 Some(169.867061),
1569 Some(168.471101),
1570 Some(98.019802),
1571 None,
1572 ],
1573 );
1574 }
1575
1576 #[test]
1577 fn calculate_count_over_time() {
1578 let (ts_array, value_array) = build_test_range_arrays();
1579 simple_range_udf_runner(
1580 CountOverTime::scalar_udf(),
1581 ts_array,
1582 value_array,
1583 vec![],
1584 vec![
1585 Some(2.0),
1586 Some(5.0),
1587 Some(1.0),
1588 None,
1589 None,
1590 Some(3.0),
1591 Some(3.0),
1592 Some(3.0),
1593 Some(1.0),
1594 None,
1595 ],
1596 );
1597 }
1598
1599 #[test]
1600 fn calculate_last_over_time() {
1601 let (ts_array, value_array) = build_test_range_arrays();
1602 simple_range_udf_runner(
1603 LastOverTime::scalar_udf(),
1604 ts_array,
1605 value_array,
1606 vec![],
1607 vec![
1608 Some(87.654321),
1609 Some(70.710678),
1610 Some(87.654321),
1611 None,
1612 None,
1613 Some(41.421356),
1614 Some(57.735027),
1615 Some(69.314718),
1616 Some(98.019802),
1617 None,
1618 ],
1619 );
1620 }
1621
1622 #[test]
1623 fn calculate_absent_over_time() {
1624 let (ts_array, value_array) = build_test_range_arrays();
1625 simple_range_udf_runner(
1626 AbsentOverTime::scalar_udf(),
1627 ts_array,
1628 value_array,
1629 vec![],
1630 vec![
1631 None,
1632 None,
1633 None,
1634 Some(1.0),
1635 Some(1.0),
1636 None,
1637 None,
1638 None,
1639 None,
1640 Some(1.0),
1641 ],
1642 );
1643 }
1644
1645 #[test]
1646 fn calculate_present_over_time() {
1647 let (ts_array, value_array) = build_test_range_arrays();
1648 simple_range_udf_runner(
1649 PresentOverTime::scalar_udf(),
1650 ts_array,
1651 value_array,
1652 vec![],
1653 vec![
1654 Some(1.0),
1655 Some(1.0),
1656 Some(1.0),
1657 None,
1658 None,
1659 Some(1.0),
1660 Some(1.0),
1661 Some(1.0),
1662 Some(1.0),
1663 None,
1664 ],
1665 );
1666 }
1667
1668 #[test]
1669 fn calculate_stdvar_over_time() {
1670 let (ts_array, value_array) = build_test_range_arrays();
1671 simple_range_udf_runner(
1672 StdvarOverTime::scalar_udf(),
1673 ts_array,
1674 value_array,
1675 vec![],
1676 vec![
1677 Some(1417.8479276253622),
1678 Some(808.999919713209),
1679 Some(0.0),
1680 None,
1681 None,
1682 Some(328.3638826418587),
1683 Some(143.5964181766362),
1684 Some(130.91830542386285),
1685 Some(0.0),
1686 None,
1687 ],
1688 );
1689
1690 let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1692 [1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000]
1693 .into_iter()
1694 .map(Some),
1695 ));
1696 let values_array = Arc::new(Float64Array::from_iter([
1697 1.5990505637277868,
1698 1.5990505637277868,
1699 1.5990505637277868,
1700 0.0,
1701 8.0,
1702 8.0,
1703 2.0,
1704 3.0,
1705 ]));
1706 let ranges = [(0, 3), (3, 5)];
1707 simple_range_udf_runner(
1708 StdvarOverTime::scalar_udf(),
1709 RangeArray::from_ranges(ts_array, ranges).unwrap(),
1710 RangeArray::from_ranges(values_array, ranges).unwrap(),
1711 vec![],
1712 vec![Some(0.0), Some(10.559999999999999)],
1713 );
1714 }
1715
1716 #[test]
1717 fn calculate_std_dev_over_time() {
1718 let (ts_array, value_array) = build_test_range_arrays();
1719 simple_range_udf_runner(
1720 StddevOverTime::scalar_udf(),
1721 ts_array,
1722 value_array,
1723 vec![],
1724 vec![
1725 Some(37.6543215),
1726 Some(28.442923895289123),
1727 Some(0.0),
1728 None,
1729 None,
1730 Some(18.12081352042062),
1731 Some(11.983172291869804),
1732 Some(11.441953741554055),
1733 Some(0.0),
1734 None,
1735 ],
1736 );
1737
1738 let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1740 [1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000]
1741 .into_iter()
1742 .map(Some),
1743 ));
1744 let values_array = Arc::new(Float64Array::from_iter([
1745 1.5990505637277868,
1746 1.5990505637277868,
1747 1.5990505637277868,
1748 0.0,
1749 8.0,
1750 8.0,
1751 2.0,
1752 3.0,
1753 ]));
1754 let ranges = [(0, 3), (3, 5)];
1755 simple_range_udf_runner(
1756 StddevOverTime::scalar_udf(),
1757 RangeArray::from_ranges(ts_array, ranges).unwrap(),
1758 RangeArray::from_ranges(values_array, ranges).unwrap(),
1759 vec![],
1760 vec![Some(0.0), Some(3.249615361854384)],
1761 );
1762 }
1763
1764 fn null_sample_range_arrays() -> (RangeArray, RangeArray) {
1766 let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([
1767 0i64, 1000, 2000, 3000,
1768 ]));
1769 let values_array = Arc::new(Float64Array::new(
1772 vec![2.0, 1000.0, -1000.0, 8.0].into(),
1773 Some(NullBuffer::from_iter([true, false, false, true])),
1774 ));
1775 let ranges = [(0, 4), (1, 2)];
1777
1778 (
1779 RangeArray::from_ranges(ts_array, ranges).unwrap(),
1780 RangeArray::from_ranges(values_array, ranges).unwrap(),
1781 )
1782 }
1783
1784 #[test]
1785 fn avg_over_time_divides_by_sample_count() {
1786 let (ts_array, value_array) = null_sample_range_arrays();
1787 simple_range_udf_runner(
1788 AvgOverTime::scalar_udf(),
1789 ts_array,
1790 value_array,
1791 vec![],
1792 vec![Some(5.0), None],
1793 );
1794 }
1795
1796 #[test]
1797 fn stdvar_and_stddev_over_time_skip_null_samples() {
1798 let (ts_array, value_array) = null_sample_range_arrays();
1799 simple_range_udf_runner(
1800 StdvarOverTime::scalar_udf(),
1801 ts_array,
1802 value_array,
1803 vec![],
1804 vec![Some(9.0), None],
1805 );
1806
1807 let (ts_array, value_array) = null_sample_range_arrays();
1808 simple_range_udf_runner(
1809 StddevOverTime::scalar_udf(),
1810 ts_array,
1811 value_array,
1812 vec![],
1813 vec![Some(3.0), None],
1814 );
1815 }
1816}