1use 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#[derive(Debug)]
54pub struct ExtrapolatedRate<const IS_COUNTER: bool, const IS_RATE: bool> {
55 range_length: i64,
57}
58
59impl<const IS_COUNTER: bool, const IS_RATE: bool> ExtrapolatedRate<IS_COUNTER, IS_RATE> {
60 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 RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
80 RangeArray::convert_data_type(DataType::Float64),
82 DataType::Timestamp(TimeUnit::Millisecond, None),
84 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 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 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 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 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
306impl 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
317impl 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
328impl 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 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 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 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 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 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 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 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}