Skip to main content

promql/functions/
aggr_over_time.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    // The generic range-function wrapper downcasts both arrays for each window.
90    // Do it once after retaining its zero-window behavior above.
91    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/// The average value of all points in the specified interval.
154#[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    // `sum` already skips null slots and yields `None` for an all-null window, so only the
161    // divisor needs to count samples instead of slots.
162    let sample_count = values.len() - values.null_count();
163    compute::sum(values).map(|result| result / sample_count as f64)
164}
165
166/// The minimum value of all points in the specified interval.
167#[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/// The maximum value of all points in the specified interval.
199#[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
237/// Batches with fewer windows than this have too little repeated work to reclaim.
238const MIN_SLIDING_WINDOWS: usize = 4;
239/// Below this average window length the per-sample bookkeeping is comparable to the
240/// scan it would replace.
241const MIN_SLIDING_WINDOW_LENGTH: u64 = 32;
242/// Reuse is taken only when a step advances at most this fraction of the window.
243const MAX_SLIDING_STEP_FRACTION: u64 = 4;
244
245/// Whether reusing candidates across windows is expected to beat rescanning each one.
246///
247/// Reuse pays off in proportion to how much consecutive windows overlap, and loses to
248/// a plain scan on wide windows that barely overlap: maintaining the deque then costs
249/// more than the rescan it replaces.
250///
251/// The shape is read from the whole batch rather than from its leading windows.
252/// `RangeManipulate` holds the window duration and the evaluation step fixed, but the
253/// sample counts still vary: a series that starts inside the query range gets a first
254/// window covering roughly one step. Averages survive that; the first two windows do not.
255///
256/// A wrong answer costs time, not correctness — both evaluators return the same bits.
257fn 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        // `RangeManipulate` emits a window covering no sample as `(0, 0)`, including
269        // the trailing one a query gets when its last evaluation lands exactly one
270        // window past the last sample. Its offset says nothing about the batch.
271        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    // What reuse skips re-reading: the distance the left bound travels over the batch.
279    let total_advance = u64::from(highest_offset.saturating_sub(lowest_offset));
280
281    // `total_length / windows >= MIN_SLIDING_WINDOW_LENGTH` and
282    // `total_advance / (windows - 1) <= (total_length / windows) / MAX_SLIDING_STEP_FRACTION`,
283    // cross-multiplied to keep the averages exact. Empty windows stay in the window
284    // count, which only makes both conditions stricter.
285    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
366/// Reuses extrema candidates across windows whose bounds do not retreat.
367struct SlidingExtrema {
368    is_min: bool,
369    /// Indices of the samples that can still become the extremum, in arrival order.
370    /// A strictly better sample evicts the ones queued before it, so the front is the
371    /// extremum of the current window and tied values keep their arrival order — and
372    /// with it the sign of tied zeros.
373    candidates: VecDeque<usize>,
374    /// Backs the all-NaN window, which keeps the last NaN of the window.
375    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                // Samples between two disjoint windows belong to neither, and later
401                // windows only move right, so skipping them keeps the state exact.
402                previous_end.max(start)
403            }
404            // A retreating bound can bring back samples that are no longer tracked.
405            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                // Keep the latest payload while no non-NaN value is in the window.
429                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
444/// Folds one window on its own, the way the per-window scan used to.
445fn 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            // A NaN only holds the slot until any other sample arrives.
456            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/// The sum of all values in the specified interval.
490#[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/// The count of all values in the specified interval.
500#[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/// The most recent point value in specified interval.
512#[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/// absent_over_time returns an empty vector if the range vector passed to it has any
523/// elements (floats or native histograms) and a 1-element vector with the value 1 if
524/// the range vector passed to it has no elements.
525#[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/// the value 1 for any series in the specified interval.
536#[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/// the population standard variance of the values in the specified interval.
547/// DataFusion's implementation:
548/// <https://github.com/apache/arrow-datafusion/blob/292eb954fc0bad3a1febc597233ba26cb60bda3e/datafusion/physical-expr/src/aggregate/variance.rs#L224-#L241>
549#[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/// the population standard deviation of the values in the specified interval.
573/// Prometheus's implementation: <https://github.com/prometheus/prometheus/blob/f55ab2217984770aa1eecd0f2d5f54580029b1c0/promql/functions.go#L556-L569>
574#[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(&timestamps, values)
633                } else {
634                    max_over_time(&timestamps, 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(&timestamps, &values, &timestamp_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(&timestamps, &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(&timestamps, &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            &timestamps,
813            &Float64Array::from(vec![3.0, 2.0, 1.0]),
814            &ranges,
815            &ranges,
816        );
817        assert_min_max_udfs_match_oracle(
818            &timestamps,
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(&timestamps, &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(&timestamps, &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(&timestamps, &values, &timestamp_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        // Enough windows, long enough, advancing by at most a quarter of the window.
930        assert!(reuses_candidates(&window_keys(&[
931            (0, 32),
932            (8, 32),
933            (16, 32),
934            (24, 32)
935        ])));
936
937        // One window too few.
938        assert!(!reuses_candidates(&window_keys(&[
939            (0, 32),
940            (8, 32),
941            (16, 32)
942        ])));
943        // One sample per window too short.
944        assert!(!reuses_candidates(&window_keys(&[
945            (0, 31),
946            (7, 31),
947            (14, 31),
948            (21, 31)
949        ])));
950        // One sample per step too far apart.
951        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        // `RangeManipulate` gives a series that starts inside the query range a first
962        // window covering roughly one step, and emits a window covering no sample at
963        // all as (0, 0). Reading either one as the batch shape hides the overlap that
964        // the twenty windows behind it do have.
965        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        // A query whose last evaluation lands one window past the last sample ends on
977        // a (0, 0) window. Its zero offset must not read as a batch that never moved,
978        // which would put disjoint windows on the evaluator built for overlap.
979        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        // The first two windows open the sliding path; the rest then grow, repeat,
1005        // jump over a gap, collapse to empty, retreat, and rebuild from scratch.
1006        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(&timestamps, &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        // Every window of a four-sample array, walked forwards and then backwards, so
1038        // the evaluator sees growing, shrinking, disjoint, empty and retreating bounds.
1039        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(&timestamps, &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(&timestamps, &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(&timestamps, &values), expected_min);
1106        assert_over_time_value(max_over_time(&timestamps, &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        // Empty, singleton, overlapping, disjoint, repeated, and backward windows use
1133        // independent timestamp and value offsets.
1134        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        // The first window is [1, NULL, NULL, 4]; the old physical-length evaluator
1253        // incorrectly returned 4 for count_over_time.
1254        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    // build timestamp range and value range arrays for test
1452    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), // only 1 element
1464            (2, 0), // empty range
1465            (2, 0), // empty range
1466            (3, 3),
1467            (4, 3),
1468            (5, 3),
1469            (8, 1), // only 1 element
1470            (9, 0), // empty range
1471        ];
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        // add more assertions
1691        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        // add more assertions
1739        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    /// Timestamps and value ranges shared by the null-sample assertions below.
1765    fn null_sample_range_arrays() -> (RangeArray, RangeArray) {
1766        let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([
1767            0i64, 1000, 2000, 3000,
1768        ]));
1769        // Samples are 2.0@0 and 8.0@3000; the null slots keep a payload that would skew every
1770        // aggregate if it were read.
1771        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        // The second window holds no sample at all.
1776        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}