1use 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#[derive(Debug)]
55pub struct ExtrapolatedRate<const IS_COUNTER: bool, const IS_RATE: bool> {
56 range_length: i64,
58}
59
60impl<const IS_COUNTER: bool, const IS_RATE: bool> ExtrapolatedRate<IS_COUNTER, IS_RATE> {
61 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 RangeArray::convert_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
81 RangeArray::convert_data_type(DataType::Float64),
83 DataType::Timestamp(TimeUnit::Millisecond, None),
85 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 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 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 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 if duration_to_start_ms >= extrapolation_threshold {
276 duration_to_start_ms = average_interval_ms / 2.0;
277 }
278 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 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
311fn 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
324struct CounterResetIndex<'a> {
327 values: &'a [f64],
328 positions: Vec<usize>,
330 active: Range<usize>,
332 previous: Range<usize>,
334 drops_at: usize,
337 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 #[inline]
360 fn add_resets(&mut self, result: f64, start: usize, end: usize) -> f64 {
361 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 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 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
410fn 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
434fn 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
473impl 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
484impl 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
495impl 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 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 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 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 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 let mut ranges: Vec<(u32, u32)> = (0..390).map(|i| (i, 120)).collect();
677 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 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 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 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 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 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 ×tamps[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 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 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 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 ×tamps,
920 &values,
921 &ranges,
922 &eval_timestamps,
923 20_000,
924 );
925 assert_matches_prometheus::<true, false>(
926 ×tamps,
927 &values,
928 &ranges,
929 &eval_timestamps,
930 20_000,
931 );
932 assert_matches_prometheus::<false, false>(
933 ×tamps,
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 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 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 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 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 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}