Skip to main content

promql/functions/
changes.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//! Implementation of [`changes`](https://prometheus.io/docs/prometheus/latest/querying/functions/#changes) in PromQL. Refer to the [original
16//! implementation](https://github.com/prometheus/prometheus/blob/main/promql/functions.go#L1023-L1040).
17
18use datafusion::logical_expr::ScalarUDF;
19
20use crate::functions::edge_count::{self, EdgeKind};
21
22/// used to count the number of value changes that occur within a specific time range
23#[derive(Debug)]
24pub struct Changes {}
25
26impl Changes {
27    pub const fn name() -> &'static str {
28        "prom_changes"
29    }
30
31    pub fn scalar_udf() -> ScalarUDF {
32        edge_count::scalar_udf(Self::name(), EdgeKind::Changes)
33    }
34}
35
36#[cfg(test)]
37mod test {
38    use std::sync::Arc;
39
40    use datafusion::arrow::array::{Float64Array, TimestampMillisecondArray};
41    use datafusion::arrow::buffer::NullBuffer;
42    use datatypes::arrow::array::Array;
43
44    use super::*;
45    use crate::functions::test_util::{
46        self, STALE_NAN, assert_execution_error, build_test_range_arrays, invoke_range_udf,
47        simple_range_udf_runner,
48    };
49    use crate::range_array::RangeArray;
50
51    fn changes_oracle(values: &[f64]) -> Option<f64> {
52        let (first, rest) = values.split_first()?;
53        let mut changes = 0;
54        let mut previous = first;
55        for current in rest {
56            if current != previous && !(current.is_nan() && previous.is_nan()) {
57                changes += 1;
58            }
59            previous = current;
60        }
61        Some(changes as f64)
62    }
63
64    #[test]
65    fn calculate_changes() {
66        let timestamps = vec![
67            1000i64, 3000, 5000, 7000, 9000, 11000, 13000, 15000, 17000, 200000, 500000,
68        ];
69        let ranges = vec![
70            (0, 1),
71            (0, 4),
72            (0, 6),
73            (0, 10),
74            (0, 0), // empty range
75        ];
76
77        // assertion 1
78        let values_1 = vec![1.0, 2.0, 3.0, 0.0, 1.0, 0.0, 0.0, 1.0, 2.0, 0.0];
79        let (ts_array_1, value_array_1) =
80            build_test_range_arrays(timestamps.clone(), values_1, ranges.clone());
81        simple_range_udf_runner(
82            Changes::scalar_udf(),
83            ts_array_1,
84            value_array_1,
85            vec![],
86            vec![Some(0.0), Some(3.0), Some(5.0), Some(8.0), None],
87        );
88
89        // assertion 2
90        let values_2 = vec![1.0, 2.0, 3.0, 4.0, 5.0, 1.0, 2.0, 3.0, 4.0, 5.0];
91        let (ts_array_2, value_array_2) =
92            build_test_range_arrays(timestamps.clone(), values_2, ranges.clone());
93        simple_range_udf_runner(
94            Changes::scalar_udf(),
95            ts_array_2,
96            value_array_2,
97            vec![],
98            vec![Some(0.0), Some(3.0), Some(5.0), Some(9.0), None],
99        );
100
101        // assertion 3
102        let values_3 = vec![0.0, 0.0, 0.0, 0.0, 0.0, 1.0, 1.0, 1.0, 1.0, 1.0];
103        let (ts_array_3, value_array_3) = build_test_range_arrays(timestamps, values_3, ranges);
104        simple_range_udf_runner(
105            Changes::scalar_udf(),
106            ts_array_3,
107            value_array_3,
108            vec![],
109            vec![Some(0.0), Some(0.0), Some(1.0), Some(1.0), None],
110        );
111    }
112
113    #[test]
114    fn changes_range_array_oracle_edge_cases() {
115        let values = vec![
116            Some(-0.0),
117            Some(0.0),
118            Some(f64::INFINITY),
119            Some(f64::INFINITY),
120            Some(f64::NEG_INFINITY),
121            Some(f64::NAN),
122            Some(STALE_NAN),
123            Some(f64::NAN),
124            Some(3.0),
125            Some(1.0),
126            Some(2.0),
127            Some(3.0),
128            Some(1.0),
129            Some(3.0),
130            Some(0.0),
131        ];
132        let raw_values = values.iter().map(|value| value.unwrap()).collect();
133        let expected = test_util::run_oracle_ranges(
134            values,
135            raw_values,
136            vec![(5, 0), (2, 1), (7, 5), (0, 5), (11, 5), (16, 4)],
137            vec![(0, 0), (0, 1), (4, 5), (0, 5), (8, 5), (1, 4)],
138            changes_oracle,
139            Changes::scalar_udf(),
140        );
141        assert_eq!(
142            expected,
143            vec![None, Some(0.0), Some(2.0), Some(2.0), Some(4.0), Some(2.0)]
144        );
145
146        let expected = test_util::run_oracle_ranges(
147            vec![Some(2.0), None, Some(2.0)],
148            vec![2.0, 0.0, 2.0],
149            vec![(9, 3), (1, 1)],
150            vec![(0, 3), (1, 1)],
151            changes_oracle,
152            Changes::scalar_udf(),
153        );
154        assert_eq!(expected, vec![Some(0.0), None]);
155    }
156
157    #[test]
158    fn changes_range_array_seeded_differential() {
159        test_util::run_seeded_differential(changes_oracle, Changes::scalar_udf(), false);
160    }
161
162    #[test]
163    fn changes_range_array_seeded_differential_with_nulls() {
164        test_util::run_seeded_differential(changes_oracle, Changes::scalar_udf(), true);
165    }
166
167    #[test]
168    fn changes_range_array_mismatch_errors() {
169        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
170            Some(0),
171            Some(1),
172            Some(2),
173        ]));
174        let values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
175        let error = invoke_range_udf(
176            Changes::scalar_udf(),
177            RangeArray::from_ranges(timestamps.clone(), [(0, 1), (1, 1)]).unwrap(),
178            RangeArray::from_ranges(values.clone(), [(0, 1)]).unwrap(),
179        )
180        .unwrap_err();
181        assert_execution_error(
182            error,
183            "RangeArray have different lengths in PromQL function prom_changes: array1=2, array2=1",
184        );
185
186        let error = invoke_range_udf(
187            Changes::scalar_udf(),
188            RangeArray::from_ranges(timestamps, [(0, 1), (1, 2)]).unwrap(),
189            RangeArray::from_ranges(values, [(0, 1), (1, 1)]).unwrap(),
190        )
191        .unwrap_err();
192        assert_execution_error(
193            error,
194            "RangeArray's element 1 have different lengths in PromQL function prom_changes: array1=2, array2=1",
195        );
196    }
197
198    #[test]
199    fn changes_range_array_boundaries_and_raw_nulls() {
200        // The 0.0 -> 1.0 edge enters the range; both identical windows contain no changes.
201        let (timestamps, values) = build_test_range_arrays(
202            vec![0, 1, 2, 3],
203            vec![0.0, 1.0, 1.0, 1.0],
204            vec![(1, 3), (1, 3)],
205        );
206        simple_range_udf_runner(
207            Changes::scalar_udf(),
208            timestamps,
209            values,
210            vec![],
211            vec![Some(0.0), Some(0.0)],
212        );
213
214        let (timestamps, values) = build_test_range_arrays(
215            vec![0, 1, 2, 3],
216            vec![1.0, 2.0, 3.0, 4.0],
217            vec![(0, 0), (2, 0), (4, 0)],
218        );
219        simple_range_udf_runner(
220            Changes::scalar_udf(),
221            timestamps,
222            values,
223            vec![],
224            vec![None, None, None],
225        );
226
227        let (timestamps, values) = build_test_range_arrays(
228            vec![0, 1, 2, 3],
229            vec![1.0, 2.0, 3.0, 4.0],
230            vec![(0, 1), (2, 1), (3, 1)],
231        );
232        simple_range_udf_runner(
233            Changes::scalar_udf(),
234            timestamps,
235            values,
236            vec![],
237            vec![Some(0.0), Some(0.0), Some(0.0)],
238        );
239
240        // The raw payload under the null slot would look like two changes; the sample sequence
241        // is 10 -> 10, which is none. The last window holds no sample at all.
242        let values = Arc::new(Float64Array::new(
243            vec![10.0, 7.0, 10.0].into(),
244            Some(NullBuffer::from(vec![true, false, true])),
245        ));
246        assert!(!values.is_valid(1));
247        assert_eq!(values.value(1), 7.0);
248        let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
249            Some(0),
250            Some(1),
251            Some(2),
252        ]));
253        simple_range_udf_runner(
254            Changes::scalar_udf(),
255            RangeArray::from_ranges(timestamps, [(0, 3), (1, 1)]).unwrap(),
256            RangeArray::from_ranges(values, [(0, 3), (1, 1)]).unwrap(),
257            vec![],
258            vec![Some(0.0), None],
259        );
260    }
261}