1use datafusion::logical_expr::ScalarUDF;
19
20use crate::functions::edge_count::{self, EdgeKind};
21
22#[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, TinyPrng, assert_execution_error, build_test_range_arrays,
47 invoke_range_udf, 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), ];
76
77 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 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 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(2.0), Some(0.0)]);
155 }
156
157 #[test]
158 fn changes_range_array_seeded_differential() {
159 let mut prng = TinyPrng(0x2f6e_2b1d_834a_90c5);
160 let raw_values = (0..48)
161 .map(|_| match prng.next_index(12) {
162 0 => -0.0,
163 1 => 0.0,
164 2 => -2.0,
165 3 => -1.0,
166 4 => 1.0,
167 5 => 2.0,
168 6 => f64::INFINITY,
169 7 => f64::NEG_INFINITY,
170 8 | 9 => f64::NAN,
171 _ => STALE_NAN,
172 })
173 .collect::<Vec<_>>();
174 let values = raw_values.iter().copied().map(Some).collect();
175 let mut timestamp_ranges = Vec::new();
176 let mut value_ranges = Vec::new();
177 for _ in 0..32 {
178 let length = prng.next_index(13) as u32;
179 timestamp_ranges.push((prng.next_index(65 - length as usize) as u32, length));
180 value_ranges.push((prng.next_index(49 - length as usize) as u32, length));
181 }
182
183 test_util::run_oracle_ranges(
184 values,
185 raw_values,
186 timestamp_ranges,
187 value_ranges,
188 changes_oracle,
189 Changes::scalar_udf(),
190 );
191 }
192
193 #[test]
194 fn changes_range_array_mismatch_errors() {
195 let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
196 Some(0),
197 Some(1),
198 Some(2),
199 ]));
200 let values = Arc::new(Float64Array::from_iter([1.0, 2.0, 3.0]));
201 let error = invoke_range_udf(
202 Changes::scalar_udf(),
203 RangeArray::from_ranges(timestamps.clone(), [(0, 1), (1, 1)]).unwrap(),
204 RangeArray::from_ranges(values.clone(), [(0, 1)]).unwrap(),
205 )
206 .unwrap_err();
207 assert_execution_error(
208 error,
209 "RangeArray have different lengths in PromQL function prom_changes: array1=2, array2=1",
210 );
211
212 let error = invoke_range_udf(
213 Changes::scalar_udf(),
214 RangeArray::from_ranges(timestamps, [(0, 1), (1, 2)]).unwrap(),
215 RangeArray::from_ranges(values, [(0, 1), (1, 1)]).unwrap(),
216 )
217 .unwrap_err();
218 assert_execution_error(
219 error,
220 "RangeArray's element 1 have different lengths in PromQL function prom_changes: array1=2, array2=1",
221 );
222 }
223
224 #[test]
225 fn changes_range_array_boundaries_and_raw_nulls() {
226 let (timestamps, values) = build_test_range_arrays(
228 vec![0, 1, 2, 3],
229 vec![0.0, 1.0, 1.0, 1.0],
230 vec![(1, 3), (1, 3)],
231 );
232 simple_range_udf_runner(
233 Changes::scalar_udf(),
234 timestamps,
235 values,
236 vec![],
237 vec![Some(0.0), Some(0.0)],
238 );
239
240 let (timestamps, values) = build_test_range_arrays(
241 vec![0, 1, 2, 3],
242 vec![1.0, 2.0, 3.0, 4.0],
243 vec![(0, 0), (2, 0), (4, 0)],
244 );
245 simple_range_udf_runner(
246 Changes::scalar_udf(),
247 timestamps,
248 values,
249 vec![],
250 vec![None, None, None],
251 );
252
253 let (timestamps, values) = build_test_range_arrays(
254 vec![0, 1, 2, 3],
255 vec![1.0, 2.0, 3.0, 4.0],
256 vec![(0, 1), (2, 1), (3, 1)],
257 );
258 simple_range_udf_runner(
259 Changes::scalar_udf(),
260 timestamps,
261 values,
262 vec![],
263 vec![Some(0.0), Some(0.0), Some(0.0)],
264 );
265
266 let values = Arc::new(Float64Array::new(
267 vec![10.0, 7.0, 10.0].into(),
268 Some(NullBuffer::from(vec![true, false, true])),
269 ));
270 assert!(!values.is_valid(1));
271 assert_eq!(values.value(1), 7.0);
272 let timestamps = Arc::new(TimestampMillisecondArray::from_iter([
273 Some(0),
274 Some(1),
275 Some(2),
276 ]));
277 simple_range_udf_runner(
278 Changes::scalar_udf(),
279 RangeArray::from_ranges(timestamps, [(0, 3)]).unwrap(),
280 RangeArray::from_ranges(values, [(0, 3)]).unwrap(),
281 vec![],
282 vec![Some(2.0)],
283 );
284 }
285}