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, 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), ];
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(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 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 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}