Skip to main content

promql/functions/
extrapolate_rate.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15// This file also contains some code from prometheus project.
16
17// Copyright 2015 The Prometheus Authors
18// Licensed under the Apache License, Version 2.0 (the "License");
19// you may not use this file except in compliance with the License.
20// You may obtain a copy of the License at
21//
22// http://www.apache.org/licenses/LICENSE-2.0
23//
24// Unless required by applicable law or agreed to in writing, software
25// distributed under the License is distributed on an "AS IS" BASIS,
26// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
27// See the License for the specific language governing permissions and
28// limitations under the License.
29
30//! Implementations of `rate`, `increase` and `delta` functions in PromQL.
31
32use std::fmt::Display;
33use std::ops::Range;
34use std::sync::Arc;
35
36use datafusion::arrow::array::{Float64Array, Float64Builder, TimestampMillisecondArray};
37use datafusion::arrow::datatypes::TimeUnit;
38use datafusion::common::{DataFusionError, Result as DfResult};
39use datafusion::logical_expr::{ScalarUDF, Volatility};
40use datafusion::physical_plan::ColumnarValue;
41use datafusion_expr::create_udf;
42use datatypes::arrow::array::{Array, Int64Array};
43use datatypes::arrow::datatypes::DataType;
44
45use crate::functions::{extract_array, extract_range_dict};
46use crate::range_array::{RangeArray, unpack};
47
48pub type Delta = ExtrapolatedRate<false, false>;
49pub type Rate = ExtrapolatedRate<true, true>;
50pub type Increase = ExtrapolatedRate<true, false>;
51
52/// Part of the `extrapolatedRate` in Promql,
53/// from <https://github.com/prometheus/prometheus/blob/v0.40.1/promql/functions.go#L66>
54#[derive(Debug)]
55pub struct ExtrapolatedRate<const IS_COUNTER: bool, const IS_RATE: bool> {
56    /// Range length in milliseconds.
57    range_length: i64,
58}
59
60impl<const IS_COUNTER: bool, const IS_RATE: bool> ExtrapolatedRate<IS_COUNTER, IS_RATE> {
61    /// Constructor. Other public usage should use [scalar_udf()](ExtrapolatedRate::scalar_udf()) instead.
62    fn new(range_length: i64) -> Self {
63        Self { range_length }
64    }
65
66    fn func_name() -> &'static str {
67        match (IS_COUNTER, IS_RATE) {
68            (true, true) => "prom_rate",
69            (true, false) => "prom_increase",
70            (false, false) => "prom_delta",
71            (false, true) => {
72                unreachable!("gauge rate is not supported by ExtrapolatedRate")
73            }
74        }
75    }
76
77    fn scalar_udf_with_name(name: &str) -> ScalarUDF {
78        let input_types = vec![
79            // timestamp range vector
80            RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
81            // value range vector
82            RangeArray::convert_data_type(DataType::Float64),
83            // timestamp vector
84            DataType::Timestamp(TimeUnit::Millisecond, None),
85            // range length
86            DataType::Int64,
87        ];
88
89        create_udf(
90            name,
91            input_types,
92            DataType::Float64,
93            Volatility::Volatile,
94            Arc::new(move |input: &_| Self::create_function(input)?.calc(input)) as _,
95        )
96    }
97
98    fn create_function(inputs: &[ColumnarValue]) -> DfResult<Self> {
99        if inputs.len() != 4 {
100            return Err(DataFusionError::Plan(
101                "ExtrapolatedRate function should have 4 inputs".to_string(),
102            ));
103        }
104
105        let range_length_array = extract_array(&inputs[3])?;
106        let range_length_array = range_length_array
107            .as_any()
108            .downcast_ref::<Int64Array>()
109            .ok_or_else(|| {
110                DataFusionError::Execution(format!(
111                    "{}: expect Int64 as range length type, found {}",
112                    Self::func_name(),
113                    range_length_array.data_type()
114                ))
115            })?;
116        if range_length_array.is_empty() || range_length_array.is_null(0) {
117            return Err(DataFusionError::Execution(format!(
118                "{}: range length must contain a non-null Int64 value",
119                Self::func_name()
120            )));
121        }
122        let range_length = range_length_array.value(0);
123
124        Ok(Self::new(range_length))
125    }
126
127    /// Input parameters:
128    /// * 0: timestamp range vector
129    /// * 1: value range vector
130    /// * 2: timestamp vector
131    /// * 3: range length. Range duration in milliseconds
132    fn calc(&self, input: &[ColumnarValue]) -> DfResult<ColumnarValue> {
133        if input.len() != 4 {
134            return Err(DataFusionError::Plan(
135                "ExtrapolatedRate function should have 4 inputs".to_string(),
136            ));
137        }
138
139        let ts_dict = extract_range_dict(
140            &input[0],
141            Self::func_name(),
142            "timestamp range vector",
143            &DataType::Timestamp(TimeUnit::Millisecond, None),
144        )?;
145        let value_dict = extract_range_dict(
146            &input[1],
147            Self::func_name(),
148            "value range vector",
149            &DataType::Float64,
150        )?;
151        let eval_ts_array = extract_eval_timestamps(&input[2], Self::func_name())?;
152
153        let keys = ts_dict.keys().values();
154        let num_windows = keys.len();
155        if value_dict.keys().len() != num_windows {
156            return Err(DataFusionError::Execution(format!(
157                "{}: timestamp and value ranges should have the same number of windows, found {} and {}",
158                Self::func_name(),
159                num_windows,
160                value_dict.keys().len()
161            )));
162        }
163        if value_dict.keys().values() != keys {
164            return Err(DataFusionError::Execution(format!(
165                "{}: timestamp and value ranges should have the same window layout",
166                Self::func_name()
167            )));
168        }
169        if eval_ts_array.len() != num_windows {
170            return Err(DataFusionError::Execution(format!(
171                "{}: evaluation timestamp vector should have the same number of rows as range inputs, found {} and {}",
172                Self::func_name(),
173                eval_ts_array.len(),
174                num_windows
175            )));
176        }
177
178        let all_timestamps = ts_dict
179            .values()
180            .as_any()
181            .downcast_ref::<TimestampMillisecondArray>()
182            .expect("validated by extract_range_dict")
183            .values();
184        let value_array = value_dict
185            .values()
186            .as_any()
187            .downcast_ref::<Float64Array>()
188            .expect("validated by extract_range_dict");
189        // A NULL field value means the series has no sample at that timestamp, so the padding
190        // under a null slot must not be read. Skip the per-window null scan when the whole
191        // backing array is null-free, which is the common case.
192        let has_nulls = value_array.null_count() > 0;
193        let all_values = value_array.values();
194        let eval_ts = eval_ts_array.values();
195
196        let mut result_builder = Float64Builder::with_capacity(num_windows);
197        let range_length = self.range_length;
198        let range_length_secs = range_length as f64 / 1000.0;
199
200        // Range windows normally overlap heavily, so scanning every one for resets costs far
201        // more than a single pass over the values. Index the reset positions once that is the
202        // cheaper side, and stop counting as soon as the requested pairs pass that budget,
203        // which heavy overlap does within the first few windows. A short lookback with a long
204        // step is the shape that never reaches it, and there the per-window scans do win.
205        let mut reset_index = if IS_COUNTER && !has_nulls {
206            let budget = all_values.len().saturating_sub(1);
207            let mut scanned_pairs = 0usize;
208            keys.iter()
209                .any(|&key| {
210                    scanned_pairs =
211                        scanned_pairs.saturating_add(unpack(key).1.saturating_sub(1) as usize);
212                    scanned_pairs > budget
213                })
214                .then(|| CounterResetIndex::new(all_values))
215        } else {
216            None
217        };
218
219        for index in 0..num_windows {
220            let (raw_offset, raw_length) = unpack(keys[index]);
221            let offset = raw_offset as usize;
222            let length = raw_length as usize;
223
224            let end = offset + length;
225            let (first_index, last_index, sample_count) = if has_nulls {
226                match valid_window_bounds(value_array, offset, length) {
227                    Some(bounds) => bounds,
228                    None => {
229                        result_builder.append_null();
230                        continue;
231                    }
232                }
233            } else {
234                (offset, end.saturating_sub(1), length)
235            };
236
237            if sample_count < 2 {
238                result_builder.append_null();
239                continue;
240            }
241
242            let first_value = all_values[first_index];
243            let last_value = all_values[last_index];
244
245            let mut result_value = last_value - first_value;
246            if IS_COUNTER {
247                result_value = if has_nulls {
248                    add_counter_resets_between_samples(
249                        result_value,
250                        value_array,
251                        first_index,
252                        last_index,
253                    )
254                } else {
255                    match &mut reset_index {
256                        Some(reset_index) => reset_index.add_resets(result_value, offset, end),
257                        None => add_counter_resets(result_value, &all_values[offset..end]),
258                    }
259                };
260            }
261
262            let first_ts = all_timestamps[first_index];
263            let last_ts = all_timestamps[last_index];
264            let range_end = eval_ts[index];
265            let range_start = range_end - range_length;
266            let sampled_interval_ms = (last_ts - first_ts) as f64;
267            let average_interval_ms = sampled_interval_ms / (sample_count - 1) as f64;
268            let mut duration_to_start_ms = (first_ts - range_start) as f64;
269            let mut duration_to_end_ms = (range_end - last_ts) as f64;
270            let extrapolation_threshold = average_interval_ms * 1.1;
271
272            // Mirror Prometheus extrapolation: extend to the real range boundary when a sample is
273            // close enough, otherwise only half an average sampling interval, which is the guess
274            // for where the series actually starts or ends.
275            if duration_to_start_ms >= extrapolation_threshold {
276                duration_to_start_ms = average_interval_ms / 2.0;
277            }
278            // Counters cannot be negative, so the extrapolation can snap back to the inferred
279            // zero point instead of extending into negative values. Prometheus applies this
280            // after the threshold clamp, so it can only shorten the leading extrapolation.
281            if IS_COUNTER && result_value > 0.0 && first_value >= 0.0 {
282                let duration_to_zero = sampled_interval_ms * (first_value / result_value);
283                if duration_to_zero < duration_to_start_ms {
284                    duration_to_start_ms = duration_to_zero;
285                }
286            }
287            if duration_to_end_ms >= extrapolation_threshold {
288                duration_to_end_ms = average_interval_ms / 2.0;
289            }
290
291            // Samples sharing one timestamp leave nothing to extrapolate over.
292            let mut factor = if sampled_interval_ms == 0.0 {
293                1.0
294            } else {
295                (sampled_interval_ms + duration_to_start_ms + duration_to_end_ms)
296                    / sampled_interval_ms
297            };
298
299            if IS_RATE {
300                factor /= range_length_secs;
301            }
302
303            result_builder.append_value(result_value * factor);
304        }
305
306        let result = ColumnarValue::Array(Arc::new(result_builder.finish()));
307        Ok(result)
308    }
309}
310
311/// Adds the value preceding every counter reset in `values` to `result`, in sample order.
312///
313/// Prometheus accumulates the resets into the running result rather than summing them on
314/// their own, and the two are not interchangeable in f64: a reset large enough to swallow a
315/// later one in an isolated sum still leaves it visible once the first difference is folded
316/// in first.
317fn add_counter_resets(result: f64, values: &[f64]) -> f64 {
318    values
319        .windows(2)
320        .filter(|pair| pair[1] < pair[0])
321        .fold(result, |result, pair| result + pair[0])
322}
323
324/// Positions of the counter resets in a value array, so that a window can accumulate the
325/// resets it contains instead of scanning all of its samples.
326struct CounterResetIndex<'a> {
327    values: &'a [f64],
328    /// Ascending indices `i` where `values[i] < values[i - 1]`.
329    positions: Vec<usize>,
330    /// Slice of `positions` covered by the last window.
331    active: Range<usize>,
332    /// That window, so the next one can tell whether it advanced.
333    previous: Range<usize>,
334    /// `positions[active.start]`: the reset a later `start` would drop. `usize::MAX` when the
335    /// active slice reaches the end of `positions`.
336    drops_at: usize,
337    /// `positions[active.end]`: the reset a later `end` would gain, saturated the same way.
338    gains_at: usize,
339}
340
341impl<'a> CounterResetIndex<'a> {
342    fn new(values: &'a [f64]) -> Self {
343        let positions: Vec<usize> = (1..values.len())
344            .filter(|&i| values[i] < values[i - 1])
345            .collect();
346        let first = positions.first().copied().unwrap_or(usize::MAX);
347        Self {
348            values,
349            positions,
350            active: 0..0,
351            previous: 0..0,
352            drops_at: first,
353            gains_at: first,
354        }
355    }
356
357    /// Same additions [`add_counter_resets`] performs over `values[start..end]`, in the same
358    /// order, reached through the index instead of by scanning the window.
359    #[inline]
360    fn add_resets(&mut self, result: f64, start: usize, end: usize) -> f64 {
361        // The active slice only stays put if the window advanced without reaching either of
362        // the resets that bound it.
363        if start < self.previous.start
364            || end < self.previous.end
365            || start >= self.drops_at
366            || end > self.gains_at
367        {
368            self.locate(start, end);
369        }
370        self.previous = start..end;
371
372        if self.active.start == self.active.end {
373            // A counter that has not reset inside this window, which is the normal case, would
374            // otherwise pay a range bounds check and an empty iterator for nothing.
375            return result;
376        }
377
378        let values = self.values;
379        self.positions[self.active.start..self.active.end]
380            .iter()
381            .fold(result, |result, &i| result + values[i - 1])
382    }
383
384    fn locate(&mut self, start: usize, end: usize) {
385        // Walk the bounds forward from the previous window and only search when they move
386        // back. On a series that resets often the searches cost more than the additions they
387        // locate, because they run deep and a window holds a handful of resets.
388        let (left, right) = if start < self.previous.start || end < self.previous.end {
389            (
390                self.positions.partition_point(|&i| i <= start),
391                self.positions.partition_point(|&i| i < end),
392            )
393        } else {
394            let mut left = self.active.start;
395            while left < self.positions.len() && self.positions[left] <= start {
396                left += 1;
397            }
398            let mut right = self.active.end.max(left);
399            while right < self.positions.len() && self.positions[right] < end {
400                right += 1;
401            }
402            (left, right)
403        };
404        self.active = left..right;
405        self.drops_at = self.positions.get(left).copied().unwrap_or(usize::MAX);
406        self.gains_at = self.positions.get(right).copied().unwrap_or(usize::MAX);
407    }
408}
409
410/// Same additions [`add_counter_resets`] performs, over the samples in `[first, last]` instead
411/// of over every slot.
412fn add_counter_resets_between_samples(
413    result: f64,
414    values: &Float64Array,
415    first: usize,
416    last: usize,
417) -> f64 {
418    let raw_values = values.values();
419    let mut result = result;
420    let mut previous = raw_values[first];
421    for index in first + 1..=last {
422        if values.is_null(index) {
423            continue;
424        }
425        let current = raw_values[index];
426        if current < previous {
427            result += previous;
428        }
429        previous = current;
430    }
431    result
432}
433
434/// Locates the samples inside `[offset, offset + length)`, returning the first and last
435/// non-null index together with the number of non-null slots. Returns `None` when the
436/// window holds no sample.
437fn valid_window_bounds(
438    values: &Float64Array,
439    offset: usize,
440    length: usize,
441) -> Option<(usize, usize, usize)> {
442    let mut first = None;
443    let mut last = 0;
444    let mut count = 0;
445    for index in offset..offset + length {
446        if values.is_null(index) {
447            continue;
448        }
449        first.get_or_insert(index);
450        last = index;
451        count += 1;
452    }
453    first.map(|first| (first, last, count))
454}
455
456fn extract_eval_timestamps(
457    columnar_value: &ColumnarValue,
458    func_name: &str,
459) -> DfResult<TimestampMillisecondArray> {
460    let array = extract_array(columnar_value)?;
461    let timestamps = array
462        .as_any()
463        .downcast_ref::<TimestampMillisecondArray>()
464        .ok_or_else(|| {
465            DataFusionError::Execution(format!(
466                "{func_name}: expect evaluation timestamp vector as Timestamp(Millisecond), found {}",
467                array.data_type()
468            ))
469        })?;
470    Ok(timestamps.clone())
471}
472
473// delta
474impl ExtrapolatedRate<false, false> {
475    pub const fn name() -> &'static str {
476        "prom_delta"
477    }
478
479    pub fn scalar_udf() -> ScalarUDF {
480        Self::scalar_udf_with_name(Self::name())
481    }
482}
483
484// rate
485impl ExtrapolatedRate<true, true> {
486    pub const fn name() -> &'static str {
487        "prom_rate"
488    }
489
490    pub fn scalar_udf() -> ScalarUDF {
491        Self::scalar_udf_with_name(Self::name())
492    }
493}
494
495// increase
496impl ExtrapolatedRate<true, false> {
497    pub const fn name() -> &'static str {
498        "prom_increase"
499    }
500
501    pub fn scalar_udf() -> ScalarUDF {
502        Self::scalar_udf_with_name(Self::name())
503    }
504}
505
506impl Display for ExtrapolatedRate<false, false> {
507    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
508        f.write_str("PromQL Delta Function")
509    }
510}
511
512impl Display for ExtrapolatedRate<true, true> {
513    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
514        f.write_str("PromQL Rate Function")
515    }
516}
517
518impl Display for ExtrapolatedRate<true, false> {
519    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
520        f.write_str("PromQL Increase Function")
521    }
522}
523
524#[cfg(test)]
525mod test {
526
527    use datafusion::arrow::array::ArrayRef;
528    use datafusion::arrow::buffer::NullBuffer;
529    use datafusion_common::ScalarValue;
530
531    use super::*;
532    use crate::functions::test_util::TinyPrng;
533
534    /// Range length is fixed to 5
535    fn extrapolated_rate_runner<const IS_COUNTER: bool, const IS_RATE: bool>(
536        ts_range: RangeArray,
537        value_range: RangeArray,
538        timestamps: ArrayRef,
539        expected: Vec<f64>,
540    ) {
541        let input = vec![
542            ColumnarValue::Array(Arc::new(ts_range.into_dict())),
543            ColumnarValue::Array(Arc::new(value_range.into_dict())),
544            ColumnarValue::Array(timestamps),
545            ColumnarValue::Array(Arc::new(Int64Array::from(vec![5]))),
546        ];
547        let output = extract_array(
548            &ExtrapolatedRate::<IS_COUNTER, IS_RATE>::new(5)
549                .calc(&input)
550                .unwrap(),
551        )
552        .unwrap()
553        .as_any()
554        .downcast_ref::<Float64Array>()
555        .unwrap()
556        .values()
557        .to_vec();
558        assert_eq!(output, expected);
559    }
560
561    fn sample_range_inputs() -> (ColumnarValue, ColumnarValue, ColumnarValue) {
562        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
563            [1, 2, 3].into_iter().map(Some),
564        ));
565        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
566        let ranges = [(0, 2), (1, 2)];
567
568        let ts_range = RangeArray::from_ranges(ts_values, ranges).unwrap();
569        let value_range = RangeArray::from_ranges(value_values, ranges).unwrap();
570        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter(
571            [2, 3].into_iter().map(Some),
572        )) as _;
573
574        (
575            ColumnarValue::Array(Arc::new(ts_range.into_dict())),
576            ColumnarValue::Array(Arc::new(value_range.into_dict())),
577            ColumnarValue::Array(eval_ts),
578        )
579    }
580
581    /// Evaluates `ranges` as one batch and asserts every window is bit-identical to evaluating
582    /// that window on its own, which always takes the direct per-window reduction.
583    fn assert_counter_windows_match_single(values: &[f64], ranges: &[(u32, u32)]) {
584        let timestamps = Arc::new(TimestampMillisecondArray::from_iter_values(
585            (0..values.len()).map(|i| i as i64 * 30_000 + (i % 5) as i64 * 1_000),
586        ));
587        let values = Arc::new(Float64Array::from(values.to_vec()));
588        let evaluate = |ranges: &[(u32, u32)]| {
589            let eval_ts = Arc::new(TimestampMillisecondArray::from_iter_values(
590                ranges.iter().map(|&(offset, length)| {
591                    timestamps.value((offset + length.saturating_sub(1)) as usize) + 5_000
592                }),
593            ));
594            let input = [
595                ColumnarValue::Array(Arc::new(
596                    RangeArray::from_ranges(timestamps.clone(), ranges.iter().copied())
597                        .unwrap()
598                        .into_dict(),
599                )),
600                ColumnarValue::Array(Arc::new(
601                    RangeArray::from_ranges(values.clone(), ranges.iter().copied())
602                        .unwrap()
603                        .into_dict(),
604                )),
605                ColumnarValue::Array(eval_ts),
606                ColumnarValue::Scalar(ScalarValue::Int64(Some(3_600_000))),
607            ];
608            extract_array(&Rate::new(3_600_000).calc(&input).unwrap())
609                .unwrap()
610                .as_any()
611                .downcast_ref::<Float64Array>()
612                .unwrap()
613                .iter()
614                .collect::<Vec<_>>()
615        };
616
617        for (range, batched) in ranges.iter().zip(evaluate(ranges)) {
618            let single = evaluate(std::slice::from_ref(range))[0];
619            match (batched, single) {
620                (None, None) => {}
621                (Some(batched), Some(single)) => assert!(
622                    batched.to_bits() == single.to_bits() || (batched.is_nan() && single.is_nan()),
623                    "range {range:?}: batched {batched} != single {single}"
624                ),
625                _ => panic!("range {range:?}: batched {batched:?} != single {single:?}"),
626            }
627        }
628    }
629
630    #[test]
631    fn counter_resets_accumulate_into_the_running_result() {
632        // Summed on their own, 1e16 and 1.0 round to 1e16, which then cancels against the
633        // first sample and reports no increase at all. Folding each reset into `last - first`
634        // as Prometheus does keeps the 1.0. Both paths detect the same two resets.
635        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
636            [1, 2, 3, 4].into_iter().map(Some),
637        ));
638        let values_array = Arc::new(Float64Array::from_iter([1e16, 1.0, 0.0, 1.0]));
639        let ranges = [(0, 4)];
640        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
641        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
642        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([Some(4)])) as _;
643
644        extrapolated_rate_runner::<true, false>(
645            ts_range,
646            value_range,
647            timestamps,
648            vec![1.1666666666666667],
649        );
650    }
651
652    #[test]
653    fn counter_correction_survives_huge_and_infinite_resets() {
654        // Both series reset twice inside the first window and once inside the second, but the
655        // first reset is large enough to swallow the second one when they are summed together.
656        for values in [
657            vec![1e16, 1.0, 0.0, 1.0, 2.0],
658            vec![f64::INFINITY, 1.0, 0.0, 1.0, 2.0],
659        ] {
660            assert_counter_windows_match_single(&values, &[(0, 4), (1, 4)]);
661        }
662    }
663
664    #[test]
665    fn counter_correction_matches_single_window_on_irregular_layouts() {
666        let mut values: Vec<f64> = (0..512).map(|i| (i % 37) as f64 * 0.25).collect();
667        values[20] = f64::NAN;
668        values[70] = f64::INFINITY;
669        values[140] = f64::NEG_INFINITY;
670        values[220] = 1e300;
671        values[221] = 1e-200;
672
673        // A stride of one walks both bounds across the reset positions, which sit every 37
674        // samples: `start` reaches one when that reset has to leave the window, `end` when the
675        // next one has to stay out of it, both while the cached bounds are live.
676        let mut ranges: Vec<(u32, u32)> = (0..390).map(|i| (i, 120)).collect();
677        // Empty, too-short, backward and disjoint windows all break a forward-only slide.
678        ranges.extend([(400, 0), (400, 1), (2, 20), (450, 30), (0, 120)]);
679        ranges.extend((0..390).rev().step_by(10).map(|i| (i, 120)));
680
681        assert_counter_windows_match_single(&values, &ranges);
682    }
683
684    /// Builds a value array whose null slots keep a distinguishable raw payload, so a
685    /// function that reads the padding instead of the samples produces a different result.
686    fn values_with_nulls(values: Vec<Option<f64>>, padding: f64) -> Arc<Float64Array> {
687        let raw = values
688            .iter()
689            .map(|value| value.unwrap_or(padding))
690            .collect::<Vec<_>>();
691        Arc::new(Float64Array::new(
692            raw.into(),
693            Some(NullBuffer::from_iter(
694                values.iter().map(|value| value.is_some()),
695            )),
696        ))
697    }
698
699    fn nullable_rate_runner<const IS_COUNTER: bool, const IS_RATE: bool>(
700        timestamps: Vec<i64>,
701        values: Arc<Float64Array>,
702        ranges: Vec<(u32, u32)>,
703        eval_timestamps: Vec<i64>,
704        range_length: i64,
705    ) -> Vec<Option<f64>> {
706        let ts_array = Arc::new(TimestampMillisecondArray::from_iter_values(timestamps));
707        let ts_range = RangeArray::from_ranges(ts_array, ranges.clone()).unwrap();
708        let value_range = RangeArray::from_ranges(values, ranges).unwrap();
709        let input = vec![
710            ColumnarValue::Array(Arc::new(ts_range.into_dict())),
711            ColumnarValue::Array(Arc::new(value_range.into_dict())),
712            ColumnarValue::Array(Arc::new(TimestampMillisecondArray::from_iter_values(
713                eval_timestamps,
714            ))),
715            ColumnarValue::Array(Arc::new(Int64Array::from(vec![range_length]))),
716        ];
717        let output = extract_array(
718            &ExtrapolatedRate::<IS_COUNTER, IS_RATE>::new(range_length)
719                .calc(&input)
720                .unwrap(),
721        )
722        .unwrap();
723        let output = output.as_any().downcast_ref::<Float64Array>().unwrap();
724        output.iter().collect()
725    }
726
727    #[test]
728    fn rate_uses_samples_not_null_padding() {
729        // Samples are 1.0@0 and 4.0@3000; the padding under the null slots would add two more.
730        let output = nullable_rate_runner::<true, true>(
731            vec![0, 1000, 2000, 3000],
732            values_with_nulls(vec![Some(1.0), None, None, Some(4.0)], 99.0),
733            vec![(0, 4)],
734            vec![3000],
735            4000,
736        );
737
738        assert_eq!(output, vec![Some(1.0)]);
739    }
740
741    #[test]
742    fn rate_returns_null_for_windows_without_enough_samples() {
743        let output = nullable_rate_runner::<true, true>(
744            vec![0, 1000, 2000],
745            values_with_nulls(vec![None, Some(2.0), None], 7.0),
746            vec![(0, 3), (0, 2), (2, 1)],
747            vec![2000, 2000, 2000],
748            4000,
749        );
750
751        assert_eq!(output, vec![None, None, None]);
752    }
753
754    #[test]
755    fn increase_corrects_counter_reset_between_samples() {
756        // The sample sequence is 5.0 -> 3.0, one reset. Reading the padding would see
757        // 5.0 -> 100.0 -> 3.0 and charge the correction against the wrong value.
758        let output = nullable_rate_runner::<true, false>(
759            vec![0, 1000, 2000],
760            values_with_nulls(vec![Some(5.0), None, Some(3.0)], 100.0),
761            vec![(0, 3)],
762            vec![2000],
763            2000,
764        );
765
766        assert_eq!(output, vec![Some(3.0)]);
767    }
768
769    #[test]
770    fn delta_extrapolates_from_sample_timestamps() {
771        // The window spans (-1000, 3000] but its samples only cover 1000..2000, so the
772        // extrapolation adds half an average interval on the leading side.
773        let output = nullable_rate_runner::<false, false>(
774            vec![0, 1000, 2000, 3000],
775            values_with_nulls(vec![None, Some(2.0), Some(5.0), None], 42.0),
776            vec![(0, 4)],
777            vec![3000],
778            4000,
779        );
780
781        assert_eq!(output, vec![Some(7.5)]);
782    }
783
784    /// Line-by-line port of Prometheus `extrapolatedRate` (promql/functions.go), float path
785    /// without start timestamps. Kept as a second implementation so that the order of the
786    /// threshold clamp and the counter zero-snap stays pinned to the upstream one.
787    fn prometheus_extrapolated_rate(
788        timestamps: &[i64],
789        values: &[f64],
790        eval_ts: i64,
791        range_ms: i64,
792        is_counter: bool,
793        is_rate: bool,
794    ) -> Option<f64> {
795        if values.len() < 2 {
796            return None;
797        }
798        let num_samples_minus_one = values.len() - 1;
799        let first_t = timestamps[0];
800        let last_t = timestamps[num_samples_minus_one];
801        let mut result = values[num_samples_minus_one] - values[0];
802        if is_counter {
803            for index in 1..values.len() {
804                if values[index] < values[index - 1] {
805                    result += values[index - 1];
806                }
807            }
808        }
809
810        let range_start = eval_ts - range_ms;
811        let mut duration_to_start = (first_t - range_start) as f64 / 1000.0;
812        let mut duration_to_end = (eval_ts - last_t) as f64 / 1000.0;
813        let sampled_interval = (last_t - first_t) as f64 / 1000.0;
814        let average_duration_between_samples = sampled_interval / num_samples_minus_one as f64;
815        let extrapolation_threshold = average_duration_between_samples * 1.1;
816
817        if duration_to_start >= extrapolation_threshold {
818            duration_to_start = average_duration_between_samples / 2.0;
819        }
820        if is_counter {
821            let mut duration_to_zero = duration_to_start;
822            if result > 0.0 && values[0] >= 0.0 {
823                duration_to_zero = sampled_interval * (values[0] / result);
824            }
825            if duration_to_zero < duration_to_start {
826                duration_to_start = duration_to_zero;
827            }
828        }
829        if duration_to_end >= extrapolation_threshold {
830            duration_to_end = average_duration_between_samples / 2.0;
831        }
832
833        let mut factor = 1.0;
834        if sampled_interval != 0.0 {
835            factor = (sampled_interval + duration_to_start + duration_to_end) / sampled_interval;
836        }
837        if is_rate {
838            factor /= range_ms as f64 / 1000.0;
839        }
840        Some(result * factor)
841    }
842
843    fn assert_matches_prometheus<const IS_COUNTER: bool, const IS_RATE: bool>(
844        timestamps: &[i64],
845        values: &[f64],
846        ranges: &[(u32, u32)],
847        eval_timestamps: &[i64],
848        range_ms: i64,
849    ) {
850        let actual = nullable_rate_runner::<IS_COUNTER, IS_RATE>(
851            timestamps.to_vec(),
852            Arc::new(Float64Array::from(values.to_vec())),
853            ranges.to_vec(),
854            eval_timestamps.to_vec(),
855            range_ms,
856        );
857
858        for (index, ((offset, length), eval_ts)) in ranges.iter().zip(eval_timestamps).enumerate() {
859            let window = *offset as usize..(*offset + *length) as usize;
860            let expected = prometheus_extrapolated_rate(
861                &timestamps[window.clone()],
862                &values[window],
863                *eval_ts,
864                range_ms,
865                IS_COUNTER,
866                IS_RATE,
867            );
868            match (actual[index], expected) {
869                (None, None) => {}
870                (Some(actual), Some(expected)) => assert!(
871                    (actual - expected).abs() <= expected.abs() * 1e-9,
872                    "window {index} {:?}: got {actual}, Prometheus gives {expected}",
873                    ranges[index]
874                ),
875                (actual, expected) => {
876                    panic!("window {index}: got {actual:?}, Prometheus gives {expected:?}")
877                }
878            }
879        }
880    }
881
882    #[test]
883    fn extrapolation_matches_prometheus_on_seeded_windows() {
884        let mut prng = TinyPrng(0x51ed_270b_8f26_1a37);
885        // Uneven spacing so the average interval, and with it the extrapolation threshold,
886        // differs from window to window.
887        let timestamps = (0..48)
888            .scan(0i64, |clock, _| {
889                *clock += 1_000 + prng.next_index(4) as i64 * 500;
890                Some(*clock)
891            })
892            .collect::<Vec<_>>();
893        // A counter that resets a few times, so the zero-snap branch is reached with both
894        // small and large leading values.
895        let values = (0..48)
896            .scan(0.0f64, |counter, _| {
897                *counter = match prng.next_index(8) {
898                    0 => 0.0,
899                    1 => *counter / 2.0,
900                    _ => *counter + prng.next_index(50) as f64,
901                };
902                Some(*counter)
903            })
904            .collect::<Vec<_>>();
905
906        let mut ranges = Vec::new();
907        let mut eval_timestamps = Vec::new();
908        for _ in 0..32 {
909            let length = 2 + prng.next_index(10) as u32;
910            let offset = prng.next_index(48 - length as usize) as u32;
911            ranges.push((offset, length));
912            // Land the range boundary at varying distances from the samples, so the clamp
913            // fires on neither, one, or both sides.
914            let last = timestamps[(offset + length - 1) as usize];
915            eval_timestamps.push(last + prng.next_index(5) as i64 * 500);
916        }
917
918        assert_matches_prometheus::<true, true>(
919            &timestamps,
920            &values,
921            &ranges,
922            &eval_timestamps,
923            20_000,
924        );
925        assert_matches_prometheus::<true, false>(
926            &timestamps,
927            &values,
928            &ranges,
929            &eval_timestamps,
930            20_000,
931        );
932        assert_matches_prometheus::<false, false>(
933            &timestamps,
934            &values,
935            &ranges,
936            &eval_timestamps,
937            20_000,
938        );
939    }
940
941    #[test]
942    fn rate_rejects_wrong_input_arity() {
943        let err = ExtrapolatedRate::<true, true>::new(5)
944            .calc(&[])
945            .unwrap_err();
946
947        assert!(err.to_string().contains("should have 4 inputs"));
948    }
949
950    #[test]
951    fn rate_rejects_non_int64_range_length() {
952        let (ts_range, value_range, eval_ts) = sample_range_inputs();
953
954        let err = ExtrapolatedRate::<true, true>::create_function(&[
955            ts_range,
956            value_range,
957            eval_ts,
958            ColumnarValue::Scalar(ScalarValue::Float64(Some(5.0))),
959        ])
960        .unwrap_err();
961
962        assert!(err.to_string().contains("range length type"));
963    }
964
965    #[test]
966    fn rate_rejects_empty_range_length() {
967        let (ts_range, value_range, eval_ts) = sample_range_inputs();
968
969        let err = ExtrapolatedRate::<true, true>::create_function(&[
970            ts_range,
971            value_range,
972            eval_ts,
973            ColumnarValue::Array(Arc::new(Int64Array::from(Vec::<i64>::new()))),
974        ])
975        .unwrap_err();
976
977        assert!(err.to_string().contains("range length must contain"));
978    }
979
980    #[test]
981    fn rate_rejects_null_range_length() {
982        let (ts_range, value_range, eval_ts) = sample_range_inputs();
983
984        let err = ExtrapolatedRate::<true, true>::create_function(&[
985            ts_range,
986            value_range,
987            eval_ts,
988            ColumnarValue::Array(Arc::new(Int64Array::from(vec![None]))),
989        ])
990        .unwrap_err();
991
992        assert!(err.to_string().contains("range length must contain"));
993    }
994
995    #[test]
996    fn increase_abnormal_input() {
997        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
998            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
999        ));
1000        let values_array = Arc::new(Float64Array::from_iter([
1001            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
1002        ]));
1003        let ranges = [(0, 2), (0, 5), (1, 1), (3, 3), (8, 1), (9, 0)];
1004        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1005        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1006        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
1007            Some(2),
1008            Some(5),
1009            Some(2),
1010            Some(6),
1011            Some(9),
1012            None,
1013        ])) as _;
1014        extrapolated_rate_runner::<true, false>(
1015            ts_range,
1016            value_range,
1017            timestamps,
1018            vec![1.5, 5.0, 0.0, 2.5, 0.0, 0.0],
1019        );
1020    }
1021
1022    #[test]
1023    fn increase_normal_input() {
1024        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1025            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1026        ));
1027        let values_array = Arc::new(Float64Array::from_iter([
1028            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
1029        ]));
1030        let ranges = [
1031            (0, 2),
1032            (1, 2),
1033            (2, 2),
1034            (3, 2),
1035            (4, 2),
1036            (5, 2),
1037            (6, 2),
1038            (7, 2),
1039        ];
1040        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1041        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1042        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1043            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1044        )) as _;
1045        extrapolated_rate_runner::<true, false>(
1046            ts_range,
1047            value_range,
1048            timestamps,
1049            vec![1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5],
1050        );
1051    }
1052
1053    #[test]
1054    fn increase_short_input() {
1055        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1056            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1057        ));
1058        let values_array = Arc::new(Float64Array::from_iter([
1059            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
1060        ]));
1061        let ranges = [
1062            (0, 1),
1063            (1, 0),
1064            (2, 1),
1065            (3, 0),
1066            (4, 3),
1067            (5, 1),
1068            (6, 0),
1069            (7, 2),
1070        ];
1071        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1072        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1073        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
1074            Some(1),
1075            None,
1076            Some(3),
1077            None,
1078            Some(7),
1079            Some(6),
1080            None,
1081            Some(9),
1082        ])) as _;
1083        extrapolated_rate_runner::<true, false>(
1084            ts_range,
1085            value_range,
1086            timestamps,
1087            vec![0.0, 0.0, 0.0, 0.0, 2.5, 0.0, 0.0, 1.5],
1088        );
1089    }
1090
1091    #[test]
1092    fn increase_counter_reset() {
1093        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1094            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1095        ));
1096        // this series should be treated like [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]
1097        let values_array = Arc::new(Float64Array::from_iter([
1098            1.0, 2.0, 3.0, 4.0, 1.0, 2.0, 3.0, 4.0, 5.0,
1099        ]));
1100        let ranges = [
1101            (0, 2),
1102            (1, 2),
1103            (2, 2),
1104            (3, 2),
1105            (4, 2),
1106            (5, 2),
1107            (6, 2),
1108            (7, 2),
1109        ];
1110        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1111        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1112        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1113            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1114        )) as _;
1115        extrapolated_rate_runner::<true, false>(
1116            ts_range,
1117            value_range,
1118            timestamps,
1119            // that two `2.0` is because `duration_to_start` are shrunk to
1120            // `duration_to_zero`, and causes `duration_to_zero` less than
1121            // `extrapolation_threshold`.
1122            vec![1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5],
1123        );
1124    }
1125
1126    #[test]
1127    fn increase_counter_reset_wide_windows() {
1128        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1129            [1, 2, 3, 4, 5, 6, 7].into_iter().map(Some),
1130        ));
1131        let values_array = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0, 1.0, 2.0, 1.0, 2.0]));
1132        let ranges = [(0, 4), (1, 4), (2, 4), (3, 4)];
1133        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1134        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1135        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1136            [4, 5, 6, 7].into_iter().map(Some),
1137        )) as _;
1138        extrapolated_rate_runner::<true, false>(
1139            ts_range,
1140            value_range,
1141            timestamps,
1142            vec![3.5, 3.5, 3.5, 3.5],
1143        );
1144    }
1145
1146    #[test]
1147    fn rate_rejects_non_array_timestamp_ranges() {
1148        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0]));
1149        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
1150        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter([Some(2)]));
1151
1152        let err = ExtrapolatedRate::<true, true>::new(5)
1153            .calc(&[
1154                ColumnarValue::Scalar(ScalarValue::Int64(Some(0))),
1155                ColumnarValue::Array(Arc::new(value_range.into_dict())),
1156                ColumnarValue::Array(eval_ts),
1157                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1158            ])
1159            .unwrap_err();
1160
1161        assert!(err.to_string().contains("timestamp range vector"));
1162    }
1163
1164    #[test]
1165    fn rate_rejects_non_timestamp_timestamp_range_values() {
1166        let ts_values = Arc::new(Int64Array::from_iter([1, 2]));
1167        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0]));
1168        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2)]).unwrap();
1169        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
1170        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter([Some(2)]));
1171
1172        let err = ExtrapolatedRate::<true, true>::new(5)
1173            .calc(&[
1174                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
1175                ColumnarValue::Array(Arc::new(value_range.into_dict())),
1176                ColumnarValue::Array(eval_ts),
1177                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1178            ])
1179            .unwrap_err();
1180
1181        assert!(err.to_string().contains("values of type Timestamp"));
1182    }
1183
1184    #[test]
1185    fn rate_rejects_non_float_value_range_values() {
1186        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
1187            [1, 2].into_iter().map(Some),
1188        ));
1189        let value_values = Arc::new(Int64Array::from_iter([1, 2]));
1190        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2)]).unwrap();
1191        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
1192        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter([Some(2)]));
1193
1194        let err = ExtrapolatedRate::<true, true>::new(5)
1195            .calc(&[
1196                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
1197                ColumnarValue::Array(Arc::new(value_range.into_dict())),
1198                ColumnarValue::Array(eval_ts),
1199                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1200            ])
1201            .unwrap_err();
1202
1203        assert!(
1204            err.to_string()
1205                .contains("value range vector values of type Float64")
1206        );
1207    }
1208
1209    #[test]
1210    fn rate_rejects_mismatched_range_counts() {
1211        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
1212            [1, 2, 3].into_iter().map(Some),
1213        ));
1214        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
1215        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2), (1, 2)]).unwrap();
1216        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
1217        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter(
1218            [2, 3].into_iter().map(Some),
1219        ));
1220
1221        let err = ExtrapolatedRate::<true, true>::new(5)
1222            .calc(&[
1223                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
1224                ColumnarValue::Array(Arc::new(value_range.into_dict())),
1225                ColumnarValue::Array(eval_ts),
1226                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1227            ])
1228            .unwrap_err();
1229
1230        assert!(err.to_string().contains("same number of windows"));
1231    }
1232
1233    #[test]
1234    fn rate_rejects_mismatched_range_layouts() {
1235        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
1236            [1, 2, 3, 4].into_iter().map(Some),
1237        ));
1238        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0, 4.0]));
1239        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2), (1, 2)]).unwrap();
1240        let value_range = RangeArray::from_ranges(value_values, [(0, 2), (2, 2)]).unwrap();
1241        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter(
1242            [2, 4].into_iter().map(Some),
1243        ));
1244
1245        let err = ExtrapolatedRate::<true, true>::new(5)
1246            .calc(&[
1247                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
1248                ColumnarValue::Array(Arc::new(value_range.into_dict())),
1249                ColumnarValue::Array(eval_ts),
1250                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1251            ])
1252            .unwrap_err();
1253
1254        assert!(err.to_string().contains("same window layout"));
1255    }
1256
1257    #[test]
1258    fn rate_rejects_non_timestamp_eval_vector() {
1259        let (ts_range, value_range, _) = sample_range_inputs();
1260
1261        let err = ExtrapolatedRate::<true, true>::new(5)
1262            .calc(&[
1263                ts_range,
1264                value_range,
1265                ColumnarValue::Array(Arc::new(Float64Array::from_iter([2.0, 3.0]))),
1266                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1267            ])
1268            .unwrap_err();
1269
1270        assert!(err.to_string().contains("evaluation timestamp vector"));
1271    }
1272
1273    #[test]
1274    fn rate_rejects_mismatched_eval_timestamp_rows() {
1275        let (ts_range, value_range, _) = sample_range_inputs();
1276
1277        let err = ExtrapolatedRate::<true, true>::new(5)
1278            .calc(&[
1279                ts_range,
1280                value_range,
1281                ColumnarValue::Array(Arc::new(TimestampMillisecondArray::from_iter([Some(2)]))),
1282                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
1283            ])
1284            .unwrap_err();
1285
1286        assert!(err.to_string().contains("same number of rows"));
1287    }
1288
1289    #[test]
1290    fn rate_counter_reset() {
1291        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1292            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1293        ));
1294        // this series should be treated like [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]
1295        let values_array = Arc::new(Float64Array::from_iter([
1296            1.0, 2.0, 3.0, 4.0, 1.0, 2.0, 3.0, 4.0, 5.0,
1297        ]));
1298        let ranges = [
1299            (0, 2),
1300            (1, 2),
1301            (2, 2),
1302            (3, 2),
1303            (4, 2),
1304            (5, 2),
1305            (6, 2),
1306            (7, 2),
1307        ];
1308        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1309        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1310        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1311            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1312        )) as _;
1313        extrapolated_rate_runner::<true, true>(
1314            ts_range,
1315            value_range,
1316            timestamps,
1317            vec![300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0],
1318        );
1319    }
1320
1321    #[test]
1322    fn rate_normal_input() {
1323        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1324            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1325        ));
1326        let values_array = Arc::new(Float64Array::from_iter([
1327            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
1328        ]));
1329        let ranges = [
1330            (0, 2),
1331            (1, 2),
1332            (2, 2),
1333            (3, 2),
1334            (4, 2),
1335            (5, 2),
1336            (6, 2),
1337            (7, 2),
1338        ];
1339        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1340        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1341        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1342            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1343        )) as _;
1344        extrapolated_rate_runner::<true, true>(
1345            ts_range,
1346            value_range,
1347            timestamps,
1348            vec![300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0],
1349        );
1350    }
1351
1352    #[test]
1353    fn delta_counter_reset() {
1354        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1355            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1356        ));
1357        // this series should be treated like [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]
1358        let values_array = Arc::new(Float64Array::from_iter([
1359            1.0, 2.0, 3.0, 4.0, 1.0, 2.0, 3.0, 4.0, 5.0,
1360        ]));
1361        let ranges = [
1362            (0, 2),
1363            (1, 2),
1364            (2, 2),
1365            (3, 2),
1366            (4, 2),
1367            (5, 2),
1368            (6, 2),
1369            (7, 2),
1370        ];
1371        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1372        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1373        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1374            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1375        )) as _;
1376        extrapolated_rate_runner::<false, false>(
1377            ts_range,
1378            value_range,
1379            timestamps,
1380            // delta doesn't handle counter reset, thus there is a negative value
1381            vec![1.5, 1.5, 1.5, -4.5, 1.5, 1.5, 1.5, 1.5],
1382        );
1383    }
1384
1385    #[test]
1386    fn delta_normal_input() {
1387        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
1388            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1389        ));
1390        let values_array = Arc::new(Float64Array::from_iter([
1391            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
1392        ]));
1393        let ranges = [
1394            (0, 2),
1395            (1, 2),
1396            (2, 2),
1397            (3, 2),
1398            (4, 2),
1399            (5, 2),
1400            (6, 2),
1401            (7, 2),
1402        ];
1403        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
1404        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
1405        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
1406            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
1407        )) as _;
1408        extrapolated_rate_runner::<false, false>(
1409            ts_range,
1410            value_range,
1411            timestamps,
1412            vec![1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5],
1413        );
1414    }
1415}