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::sync::Arc;
34
35use datafusion::arrow::array::{Float64Array, Float64Builder, TimestampMillisecondArray};
36use datafusion::arrow::datatypes::TimeUnit;
37use datafusion::common::{DataFusionError, Result as DfResult};
38use datafusion::logical_expr::{ScalarUDF, Volatility};
39use datafusion::physical_plan::ColumnarValue;
40use datafusion_expr::create_udf;
41use datatypes::arrow::array::{Array, Int64Array};
42use datatypes::arrow::datatypes::DataType;
43
44use crate::functions::{extract_array, extract_range_dict};
45use crate::range_array::{RangeArray, unpack};
46
47pub type Delta = ExtrapolatedRate<false, false>;
48pub type Rate = ExtrapolatedRate<true, true>;
49pub type Increase = ExtrapolatedRate<true, false>;
50
51/// Part of the `extrapolatedRate` in Promql,
52/// from <https://github.com/prometheus/prometheus/blob/v0.40.1/promql/functions.go#L66>
53#[derive(Debug)]
54pub struct ExtrapolatedRate<const IS_COUNTER: bool, const IS_RATE: bool> {
55    /// Range length in milliseconds.
56    range_length: i64,
57}
58
59impl<const IS_COUNTER: bool, const IS_RATE: bool> ExtrapolatedRate<IS_COUNTER, IS_RATE> {
60    /// Constructor. Other public usage should use [scalar_udf()](ExtrapolatedRate::scalar_udf()) instead.
61    fn new(range_length: i64) -> Self {
62        Self { range_length }
63    }
64
65    fn func_name() -> &'static str {
66        match (IS_COUNTER, IS_RATE) {
67            (true, true) => "prom_rate",
68            (true, false) => "prom_increase",
69            (false, false) => "prom_delta",
70            (false, true) => {
71                unreachable!("gauge rate is not supported by ExtrapolatedRate")
72            }
73        }
74    }
75
76    fn scalar_udf_with_name(name: &str) -> ScalarUDF {
77        let input_types = vec![
78            // timestamp range vector
79            RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
80            // value range vector
81            RangeArray::convert_data_type(DataType::Float64),
82            // timestamp vector
83            DataType::Timestamp(TimeUnit::Millisecond, None),
84            // range length
85            DataType::Int64,
86        ];
87
88        create_udf(
89            name,
90            input_types,
91            DataType::Float64,
92            Volatility::Volatile,
93            Arc::new(move |input: &_| Self::create_function(input)?.calc(input)) as _,
94        )
95    }
96
97    fn create_function(inputs: &[ColumnarValue]) -> DfResult<Self> {
98        if inputs.len() != 4 {
99            return Err(DataFusionError::Plan(
100                "ExtrapolatedRate function should have 4 inputs".to_string(),
101            ));
102        }
103
104        let range_length_array = extract_array(&inputs[3])?;
105        let range_length_array = range_length_array
106            .as_any()
107            .downcast_ref::<Int64Array>()
108            .ok_or_else(|| {
109                DataFusionError::Execution(format!(
110                    "{}: expect Int64 as range length type, found {}",
111                    Self::func_name(),
112                    range_length_array.data_type()
113                ))
114            })?;
115        if range_length_array.is_empty() || range_length_array.is_null(0) {
116            return Err(DataFusionError::Execution(format!(
117                "{}: range length must contain a non-null Int64 value",
118                Self::func_name()
119            )));
120        }
121        let range_length = range_length_array.value(0);
122
123        Ok(Self::new(range_length))
124    }
125
126    /// Input parameters:
127    /// * 0: timestamp range vector
128    /// * 1: value range vector
129    /// * 2: timestamp vector
130    /// * 3: range length. Range duration in milliseconds
131    fn calc(&self, input: &[ColumnarValue]) -> DfResult<ColumnarValue> {
132        if input.len() != 4 {
133            return Err(DataFusionError::Plan(
134                "ExtrapolatedRate function should have 4 inputs".to_string(),
135            ));
136        }
137
138        let ts_dict = extract_range_dict(
139            &input[0],
140            Self::func_name(),
141            "timestamp range vector",
142            &DataType::Timestamp(TimeUnit::Millisecond, None),
143        )?;
144        let value_dict = extract_range_dict(
145            &input[1],
146            Self::func_name(),
147            "value range vector",
148            &DataType::Float64,
149        )?;
150        let eval_ts_array = extract_eval_timestamps(&input[2], Self::func_name())?;
151
152        let keys = ts_dict.keys().values();
153        let num_windows = keys.len();
154        if value_dict.keys().len() != num_windows {
155            return Err(DataFusionError::Execution(format!(
156                "{}: timestamp and value ranges should have the same number of windows, found {} and {}",
157                Self::func_name(),
158                num_windows,
159                value_dict.keys().len()
160            )));
161        }
162        if value_dict.keys().values() != keys {
163            return Err(DataFusionError::Execution(format!(
164                "{}: timestamp and value ranges should have the same window layout",
165                Self::func_name()
166            )));
167        }
168        if eval_ts_array.len() != num_windows {
169            return Err(DataFusionError::Execution(format!(
170                "{}: evaluation timestamp vector should have the same number of rows as range inputs, found {} and {}",
171                Self::func_name(),
172                eval_ts_array.len(),
173                num_windows
174            )));
175        }
176
177        let all_timestamps = ts_dict
178            .values()
179            .as_any()
180            .downcast_ref::<TimestampMillisecondArray>()
181            .expect("validated by extract_range_dict")
182            .values();
183        let all_values = value_dict
184            .values()
185            .as_any()
186            .downcast_ref::<Float64Array>()
187            .expect("validated by extract_range_dict")
188            .values();
189        let eval_ts = eval_ts_array.values();
190
191        let mut result_builder = Float64Builder::with_capacity(num_windows);
192        let range_length = self.range_length;
193        let range_length_secs = range_length as f64 / 1000.0;
194
195        let mut counter_correction = 0.0;
196        let mut prev_offset = usize::MAX;
197        let mut prev_length = 0usize;
198
199        for index in 0..num_windows {
200            let (raw_offset, raw_length) = unpack(keys[index]);
201            let offset = raw_offset as usize;
202            let length = raw_length as usize;
203
204            if length < 2 {
205                result_builder.append_null();
206                prev_offset = usize::MAX;
207                continue;
208            }
209
210            let end = offset + length;
211            let first_value = all_values[offset];
212            let last_value = all_values[end - 1];
213
214            let result_value = if IS_COUNTER {
215                // Adjacent normalized windows usually slide forward by one sample. Reuse the
216                // previous window's accumulated reset correction and adjust only the dropped and
217                // newly added edges, falling back to a full scan when the layout changes.
218                if prev_offset != usize::MAX && offset == prev_offset + 1 && length == prev_length {
219                    if all_values[prev_offset + 1] < all_values[prev_offset] {
220                        counter_correction -= all_values[prev_offset];
221                    }
222                    if all_values[end - 1] < all_values[end - 2] {
223                        counter_correction += all_values[end - 2];
224                    }
225                } else {
226                    counter_correction = 0.0;
227                    for pair in all_values[offset..end].windows(2) {
228                        if pair[1] < pair[0] {
229                            counter_correction += pair[0];
230                        }
231                    }
232                }
233                last_value - first_value + counter_correction
234            } else {
235                last_value - first_value
236            };
237
238            prev_offset = offset;
239            prev_length = length;
240
241            let first_ts = all_timestamps[offset];
242            let last_ts = all_timestamps[end - 1];
243            let range_end = eval_ts[index];
244            let range_start = range_end - range_length;
245            let sampled_interval_ms = (last_ts - first_ts) as f64;
246            let average_interval_ms = sampled_interval_ms / (length - 1) as f64;
247            let mut duration_to_start_ms = (first_ts - range_start) as f64;
248            let duration_to_end_ms = (range_end - last_ts) as f64;
249
250            // Counters cannot be negative, so Prometheus allows the extrapolation window to snap
251            // back to the inferred zero point instead of extending into negative values.
252            if IS_COUNTER && result_value > 0.0 && first_value >= 0.0 {
253                let duration_to_zero = sampled_interval_ms * (first_value / result_value);
254                if duration_to_zero < duration_to_start_ms {
255                    duration_to_start_ms = duration_to_zero;
256                }
257            }
258
259            let extrapolation_threshold = average_interval_ms * 1.1;
260            let mut extrapolated_interval_ms = sampled_interval_ms;
261
262            // Mirror Prometheus extrapolation: extend to the real range boundary when a sample is
263            // close enough, otherwise add half an average sampling interval on that side.
264            if duration_to_start_ms < extrapolation_threshold {
265                extrapolated_interval_ms += duration_to_start_ms;
266            } else {
267                extrapolated_interval_ms += average_interval_ms / 2.0;
268            }
269            if duration_to_end_ms < extrapolation_threshold {
270                extrapolated_interval_ms += duration_to_end_ms;
271            } else {
272                extrapolated_interval_ms += average_interval_ms / 2.0;
273            }
274
275            let mut factor = extrapolated_interval_ms / sampled_interval_ms;
276
277            if IS_RATE {
278                factor /= range_length_secs;
279            }
280
281            result_builder.append_value(result_value * factor);
282        }
283
284        let result = ColumnarValue::Array(Arc::new(result_builder.finish()));
285        Ok(result)
286    }
287}
288
289fn extract_eval_timestamps(
290    columnar_value: &ColumnarValue,
291    func_name: &str,
292) -> DfResult<TimestampMillisecondArray> {
293    let array = extract_array(columnar_value)?;
294    let timestamps = array
295        .as_any()
296        .downcast_ref::<TimestampMillisecondArray>()
297        .ok_or_else(|| {
298            DataFusionError::Execution(format!(
299                "{func_name}: expect evaluation timestamp vector as Timestamp(Millisecond), found {}",
300                array.data_type()
301            ))
302        })?;
303    Ok(timestamps.clone())
304}
305
306// delta
307impl ExtrapolatedRate<false, false> {
308    pub const fn name() -> &'static str {
309        "prom_delta"
310    }
311
312    pub fn scalar_udf() -> ScalarUDF {
313        Self::scalar_udf_with_name(Self::name())
314    }
315}
316
317// rate
318impl ExtrapolatedRate<true, true> {
319    pub const fn name() -> &'static str {
320        "prom_rate"
321    }
322
323    pub fn scalar_udf() -> ScalarUDF {
324        Self::scalar_udf_with_name(Self::name())
325    }
326}
327
328// increase
329impl ExtrapolatedRate<true, false> {
330    pub const fn name() -> &'static str {
331        "prom_increase"
332    }
333
334    pub fn scalar_udf() -> ScalarUDF {
335        Self::scalar_udf_with_name(Self::name())
336    }
337}
338
339impl Display for ExtrapolatedRate<false, false> {
340    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
341        f.write_str("PromQL Delta Function")
342    }
343}
344
345impl Display for ExtrapolatedRate<true, true> {
346    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
347        f.write_str("PromQL Rate Function")
348    }
349}
350
351impl Display for ExtrapolatedRate<true, false> {
352    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
353        f.write_str("PromQL Increase Function")
354    }
355}
356
357#[cfg(test)]
358mod test {
359
360    use datafusion::arrow::array::ArrayRef;
361    use datafusion_common::ScalarValue;
362
363    use super::*;
364
365    /// Range length is fixed to 5
366    fn extrapolated_rate_runner<const IS_COUNTER: bool, const IS_RATE: bool>(
367        ts_range: RangeArray,
368        value_range: RangeArray,
369        timestamps: ArrayRef,
370        expected: Vec<f64>,
371    ) {
372        let input = vec![
373            ColumnarValue::Array(Arc::new(ts_range.into_dict())),
374            ColumnarValue::Array(Arc::new(value_range.into_dict())),
375            ColumnarValue::Array(timestamps),
376            ColumnarValue::Array(Arc::new(Int64Array::from(vec![5]))),
377        ];
378        let output = extract_array(
379            &ExtrapolatedRate::<IS_COUNTER, IS_RATE>::new(5)
380                .calc(&input)
381                .unwrap(),
382        )
383        .unwrap()
384        .as_any()
385        .downcast_ref::<Float64Array>()
386        .unwrap()
387        .values()
388        .to_vec();
389        assert_eq!(output, expected);
390    }
391
392    fn sample_range_inputs() -> (ColumnarValue, ColumnarValue, ColumnarValue) {
393        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
394            [1, 2, 3].into_iter().map(Some),
395        ));
396        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
397        let ranges = [(0, 2), (1, 2)];
398
399        let ts_range = RangeArray::from_ranges(ts_values, ranges).unwrap();
400        let value_range = RangeArray::from_ranges(value_values, ranges).unwrap();
401        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter(
402            [2, 3].into_iter().map(Some),
403        )) as _;
404
405        (
406            ColumnarValue::Array(Arc::new(ts_range.into_dict())),
407            ColumnarValue::Array(Arc::new(value_range.into_dict())),
408            ColumnarValue::Array(eval_ts),
409        )
410    }
411
412    #[test]
413    fn rate_rejects_wrong_input_arity() {
414        let err = ExtrapolatedRate::<true, true>::new(5)
415            .calc(&[])
416            .unwrap_err();
417
418        assert!(err.to_string().contains("should have 4 inputs"));
419    }
420
421    #[test]
422    fn rate_rejects_non_int64_range_length() {
423        let (ts_range, value_range, eval_ts) = sample_range_inputs();
424
425        let err = ExtrapolatedRate::<true, true>::create_function(&[
426            ts_range,
427            value_range,
428            eval_ts,
429            ColumnarValue::Scalar(ScalarValue::Float64(Some(5.0))),
430        ])
431        .unwrap_err();
432
433        assert!(err.to_string().contains("range length type"));
434    }
435
436    #[test]
437    fn rate_rejects_empty_range_length() {
438        let (ts_range, value_range, eval_ts) = sample_range_inputs();
439
440        let err = ExtrapolatedRate::<true, true>::create_function(&[
441            ts_range,
442            value_range,
443            eval_ts,
444            ColumnarValue::Array(Arc::new(Int64Array::from(Vec::<i64>::new()))),
445        ])
446        .unwrap_err();
447
448        assert!(err.to_string().contains("range length must contain"));
449    }
450
451    #[test]
452    fn rate_rejects_null_range_length() {
453        let (ts_range, value_range, eval_ts) = sample_range_inputs();
454
455        let err = ExtrapolatedRate::<true, true>::create_function(&[
456            ts_range,
457            value_range,
458            eval_ts,
459            ColumnarValue::Array(Arc::new(Int64Array::from(vec![None]))),
460        ])
461        .unwrap_err();
462
463        assert!(err.to_string().contains("range length must contain"));
464    }
465
466    #[test]
467    fn increase_abnormal_input() {
468        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
469            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
470        ));
471        let values_array = Arc::new(Float64Array::from_iter([
472            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
473        ]));
474        let ranges = [(0, 2), (0, 5), (1, 1), (3, 3), (8, 1), (9, 0)];
475        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
476        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
477        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
478            Some(2),
479            Some(5),
480            Some(2),
481            Some(6),
482            Some(9),
483            None,
484        ])) as _;
485        extrapolated_rate_runner::<true, false>(
486            ts_range,
487            value_range,
488            timestamps,
489            vec![2.0, 5.0, 0.0, 2.5, 0.0, 0.0],
490        );
491    }
492
493    #[test]
494    fn increase_normal_input() {
495        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
496            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
497        ));
498        let values_array = Arc::new(Float64Array::from_iter([
499            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
500        ]));
501        let ranges = [
502            (0, 2),
503            (1, 2),
504            (2, 2),
505            (3, 2),
506            (4, 2),
507            (5, 2),
508            (6, 2),
509            (7, 2),
510        ];
511        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
512        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
513        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
514            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
515        )) as _;
516        extrapolated_rate_runner::<true, false>(
517            ts_range,
518            value_range,
519            timestamps,
520            // `2.0` is because that `duration_to_zero` less than `extrapolation_threshold`
521            vec![2.0, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5],
522        );
523    }
524
525    #[test]
526    fn increase_short_input() {
527        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
528            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
529        ));
530        let values_array = Arc::new(Float64Array::from_iter([
531            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
532        ]));
533        let ranges = [
534            (0, 1),
535            (1, 0),
536            (2, 1),
537            (3, 0),
538            (4, 3),
539            (5, 1),
540            (6, 0),
541            (7, 2),
542        ];
543        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
544        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
545        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
546            Some(1),
547            None,
548            Some(3),
549            None,
550            Some(7),
551            Some(6),
552            None,
553            Some(9),
554        ])) as _;
555        extrapolated_rate_runner::<true, false>(
556            ts_range,
557            value_range,
558            timestamps,
559            vec![0.0, 0.0, 0.0, 0.0, 2.5, 0.0, 0.0, 1.5],
560        );
561    }
562
563    #[test]
564    fn increase_counter_reset() {
565        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
566            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
567        ));
568        // this series should be treated like [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]
569        let values_array = Arc::new(Float64Array::from_iter([
570            1.0, 2.0, 3.0, 4.0, 1.0, 2.0, 3.0, 4.0, 5.0,
571        ]));
572        let ranges = [
573            (0, 2),
574            (1, 2),
575            (2, 2),
576            (3, 2),
577            (4, 2),
578            (5, 2),
579            (6, 2),
580            (7, 2),
581        ];
582        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
583        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
584        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
585            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
586        )) as _;
587        extrapolated_rate_runner::<true, false>(
588            ts_range,
589            value_range,
590            timestamps,
591            // that two `2.0` is because `duration_to_start` are shrunk to
592            // `duration_to_zero`, and causes `duration_to_zero` less than
593            // `extrapolation_threshold`.
594            vec![2.0, 1.5, 1.5, 1.5, 2.0, 1.5, 1.5, 1.5],
595        );
596    }
597
598    #[test]
599    fn increase_counter_reset_wide_windows() {
600        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
601            [1, 2, 3, 4, 5, 6, 7].into_iter().map(Some),
602        ));
603        let values_array = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0, 1.0, 2.0, 1.0, 2.0]));
604        let ranges = [(0, 4), (1, 4), (2, 4), (3, 4)];
605        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
606        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
607        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
608            [4, 5, 6, 7].into_iter().map(Some),
609        )) as _;
610        extrapolated_rate_runner::<true, false>(
611            ts_range,
612            value_range,
613            timestamps,
614            vec![4.0, 3.5, 3.5, 4.0],
615        );
616    }
617
618    #[test]
619    fn rate_rejects_non_array_timestamp_ranges() {
620        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0]));
621        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
622        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter([Some(2)]));
623
624        let err = ExtrapolatedRate::<true, true>::new(5)
625            .calc(&[
626                ColumnarValue::Scalar(ScalarValue::Int64(Some(0))),
627                ColumnarValue::Array(Arc::new(value_range.into_dict())),
628                ColumnarValue::Array(eval_ts),
629                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
630            ])
631            .unwrap_err();
632
633        assert!(err.to_string().contains("timestamp range vector"));
634    }
635
636    #[test]
637    fn rate_rejects_non_timestamp_timestamp_range_values() {
638        let ts_values = Arc::new(Int64Array::from_iter([1, 2]));
639        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0]));
640        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2)]).unwrap();
641        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
642        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter([Some(2)]));
643
644        let err = ExtrapolatedRate::<true, true>::new(5)
645            .calc(&[
646                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
647                ColumnarValue::Array(Arc::new(value_range.into_dict())),
648                ColumnarValue::Array(eval_ts),
649                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
650            ])
651            .unwrap_err();
652
653        assert!(err.to_string().contains("values of type Timestamp"));
654    }
655
656    #[test]
657    fn rate_rejects_non_float_value_range_values() {
658        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
659            [1, 2].into_iter().map(Some),
660        ));
661        let value_values = Arc::new(Int64Array::from_iter([1, 2]));
662        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2)]).unwrap();
663        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
664        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter([Some(2)]));
665
666        let err = ExtrapolatedRate::<true, true>::new(5)
667            .calc(&[
668                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
669                ColumnarValue::Array(Arc::new(value_range.into_dict())),
670                ColumnarValue::Array(eval_ts),
671                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
672            ])
673            .unwrap_err();
674
675        assert!(
676            err.to_string()
677                .contains("value range vector values of type Float64")
678        );
679    }
680
681    #[test]
682    fn rate_rejects_mismatched_range_counts() {
683        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
684            [1, 2, 3].into_iter().map(Some),
685        ));
686        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
687        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2), (1, 2)]).unwrap();
688        let value_range = RangeArray::from_ranges(value_values, [(0, 2)]).unwrap();
689        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter(
690            [2, 3].into_iter().map(Some),
691        ));
692
693        let err = ExtrapolatedRate::<true, true>::new(5)
694            .calc(&[
695                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
696                ColumnarValue::Array(Arc::new(value_range.into_dict())),
697                ColumnarValue::Array(eval_ts),
698                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
699            ])
700            .unwrap_err();
701
702        assert!(err.to_string().contains("same number of windows"));
703    }
704
705    #[test]
706    fn rate_rejects_mismatched_range_layouts() {
707        let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
708            [1, 2, 3, 4].into_iter().map(Some),
709        ));
710        let value_values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0, 4.0]));
711        let ts_range = RangeArray::from_ranges(ts_values, [(0, 2), (1, 2)]).unwrap();
712        let value_range = RangeArray::from_ranges(value_values, [(0, 2), (2, 2)]).unwrap();
713        let eval_ts = Arc::new(TimestampMillisecondArray::from_iter(
714            [2, 4].into_iter().map(Some),
715        ));
716
717        let err = ExtrapolatedRate::<true, true>::new(5)
718            .calc(&[
719                ColumnarValue::Array(Arc::new(ts_range.into_dict())),
720                ColumnarValue::Array(Arc::new(value_range.into_dict())),
721                ColumnarValue::Array(eval_ts),
722                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
723            ])
724            .unwrap_err();
725
726        assert!(err.to_string().contains("same window layout"));
727    }
728
729    #[test]
730    fn rate_rejects_non_timestamp_eval_vector() {
731        let (ts_range, value_range, _) = sample_range_inputs();
732
733        let err = ExtrapolatedRate::<true, true>::new(5)
734            .calc(&[
735                ts_range,
736                value_range,
737                ColumnarValue::Array(Arc::new(Float64Array::from_iter([2.0, 3.0]))),
738                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
739            ])
740            .unwrap_err();
741
742        assert!(err.to_string().contains("evaluation timestamp vector"));
743    }
744
745    #[test]
746    fn rate_rejects_mismatched_eval_timestamp_rows() {
747        let (ts_range, value_range, _) = sample_range_inputs();
748
749        let err = ExtrapolatedRate::<true, true>::new(5)
750            .calc(&[
751                ts_range,
752                value_range,
753                ColumnarValue::Array(Arc::new(TimestampMillisecondArray::from_iter([Some(2)]))),
754                ColumnarValue::Scalar(ScalarValue::Int64(Some(5))),
755            ])
756            .unwrap_err();
757
758        assert!(err.to_string().contains("same number of rows"));
759    }
760
761    #[test]
762    fn rate_counter_reset() {
763        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
764            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
765        ));
766        // this series should be treated like [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]
767        let values_array = Arc::new(Float64Array::from_iter([
768            1.0, 2.0, 3.0, 4.0, 1.0, 2.0, 3.0, 4.0, 5.0,
769        ]));
770        let ranges = [
771            (0, 2),
772            (1, 2),
773            (2, 2),
774            (3, 2),
775            (4, 2),
776            (5, 2),
777            (6, 2),
778            (7, 2),
779        ];
780        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
781        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
782        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
783            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
784        )) as _;
785        extrapolated_rate_runner::<true, true>(
786            ts_range,
787            value_range,
788            timestamps,
789            vec![400.0, 300.0, 300.0, 300.0, 400.0, 300.0, 300.0, 300.0],
790        );
791    }
792
793    #[test]
794    fn rate_normal_input() {
795        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
796            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
797        ));
798        let values_array = Arc::new(Float64Array::from_iter([
799            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
800        ]));
801        let ranges = [
802            (0, 2),
803            (1, 2),
804            (2, 2),
805            (3, 2),
806            (4, 2),
807            (5, 2),
808            (6, 2),
809            (7, 2),
810        ];
811        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
812        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
813        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
814            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
815        )) as _;
816        extrapolated_rate_runner::<true, true>(
817            ts_range,
818            value_range,
819            timestamps,
820            vec![400.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0, 300.0],
821        );
822    }
823
824    #[test]
825    fn delta_counter_reset() {
826        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
827            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
828        ));
829        // this series should be treated like [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0]
830        let values_array = Arc::new(Float64Array::from_iter([
831            1.0, 2.0, 3.0, 4.0, 1.0, 2.0, 3.0, 4.0, 5.0,
832        ]));
833        let ranges = [
834            (0, 2),
835            (1, 2),
836            (2, 2),
837            (3, 2),
838            (4, 2),
839            (5, 2),
840            (6, 2),
841            (7, 2),
842        ];
843        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
844        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
845        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
846            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
847        )) as _;
848        extrapolated_rate_runner::<false, false>(
849            ts_range,
850            value_range,
851            timestamps,
852            // delta doesn't handle counter reset, thus there is a negative value
853            vec![1.5, 1.5, 1.5, -4.5, 1.5, 1.5, 1.5, 1.5],
854        );
855    }
856
857    #[test]
858    fn delta_normal_input() {
859        let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
860            [1, 2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
861        ));
862        let values_array = Arc::new(Float64Array::from_iter([
863            1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0, 9.0,
864        ]));
865        let ranges = [
866            (0, 2),
867            (1, 2),
868            (2, 2),
869            (3, 2),
870            (4, 2),
871            (5, 2),
872            (6, 2),
873            (7, 2),
874        ];
875        let ts_range = RangeArray::from_ranges(ts_array, ranges).unwrap();
876        let value_range = RangeArray::from_ranges(values_array, ranges).unwrap();
877        let timestamps = Arc::new(TimestampMillisecondArray::from_iter(
878            [2, 3, 4, 5, 6, 7, 8, 9].into_iter().map(Some),
879        )) as _;
880        extrapolated_rate_runner::<false, false>(
881            ts_range,
882            value_range,
883            timestamps,
884            vec![1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5, 1.5],
885        );
886    }
887}