Skip to content

Commit cf2af9c

Browse files
test(query-engine): cover SQL dual populations
1 parent 3f041e8 commit cf2af9c

3 files changed

Lines changed: 89 additions & 7 deletions

File tree

‎asap-query-engine/src/engines/simple_engine/mod.rs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -777,6 +777,9 @@ impl SimpleEngine {
777777
window_size_ms,
778778
slide_interval_ms,
779779
};
780+
// Range-step covers overlap heavily. Keep one sorted set so each
781+
// exact stored window is fetched once, in deterministic order, even
782+
// when multiple output timestamps require it.
780783
let mut windows = BTreeSet::new();
781784
for &output_timestamp in output_timestamps {
782785
let cover = plan_exact_cover(output_timestamp, lookback_ms, spec).map_err(|error| {

‎asap-query-engine/src/tests/native_range_query_tests.rs‎

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -257,11 +257,8 @@ mod tests {
257257
/// but the KEY aggregation is a Sliding window with
258258
/// `key_slide_interval_ms < key_window_size_ms` (#600). Real Sliding
259259
/// buckets are persisted on the slide_interval_ms grid, not the
260-
/// window_size_ms grid (`precompute_engine/window_manager.rs`), so the
261-
/// keys bucket span here is `key_slide_interval_ms`, not
262-
/// `key_window_size_ms` -- unlike the value side, which stays Tumbling
263-
/// (span == window) exactly as `create_range_engine_dual_input_with_windows`
264-
/// already does.
260+
/// window_size_ms grid (`precompute_engine/window_manager.rs`), and the
261+
/// value side uses the same Sliding grid for these tests.
265262
#[allow(clippy::too_many_arguments)]
266263
fn create_range_engine_dual_input_sliding_keys(
267264
metric: &str,
@@ -411,8 +408,8 @@ mod tests {
411408
// t=5000, whose keys lookback window is
412409
// [5000 - key_window_size_ms, 5000) = [3000,5000) -- lining up
413410
// exactly with the inserted bucket.
414-
// - value_data is also placed at timestamp=5000 (Tumbling, 1000ms
415-
// wide -> bucket [4000,5000)) purely so the CountMinSketch value
411+
// - value_data is also placed at timestamp=5000 (Sliding, 2000ms
412+
// wide -> bucket [3000,5000)) purely so the CountMinSketch value
416413
// side resolves at the same t=5000 step; it's unrelated to #600.
417414
let mut keys_add = SetAggregatorAccumulator::new();
418415
keys_add.add_key(KeyByLabelValues {

‎asap-query-engine/src/tests/query_equivalence_tests.rs‎

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -137,6 +137,88 @@ mod tests {
137137
assert_eq!(vector.values[0].value, 18.0);
138138
}
139139

140+
#[test]
141+
fn sql_executes_wider_sliding_dual_population_query() {
142+
use crate::data_model::{
143+
AggregationConfig, AggregationType, CleanupPolicy, KeyByLabelValues, PrecomputedOutput,
144+
};
145+
use crate::engines::query_result::QueryResult;
146+
use crate::precompute_operators::{CountMinSketchAccumulator, SetAggregatorAccumulator};
147+
use crate::stores::simple_map_store::SimpleMapStore;
148+
149+
let sql_query = "SELECT COUNT(value) FROM events WHERE time BETWEEN DATEADD(s, -10, NOW()) AND NOW() GROUP BY host";
150+
let (_, mut sql_config, mut streaming_config) = TestConfigBuilder::new("events")
151+
.with_grouping_labels(vec!["host"])
152+
.add_temporal_query(
153+
"count_over_time(events[10s])",
154+
sql_query,
155+
1,
156+
5_000,
157+
WindowType::Sliding,
158+
)
159+
.build_both();
160+
161+
let value_config = streaming_config
162+
.get_aggregation_config(1)
163+
.expect("value config")
164+
.clone();
165+
let mut value_config = AggregationConfig {
166+
aggregation_type: AggregationType::CountMinSketch,
167+
..value_config
168+
};
169+
value_config.slide_interval_ms = 1_000;
170+
let mut key_config = value_config.clone();
171+
key_config.aggregation_id = 2;
172+
key_config.aggregation_type = AggregationType::SetAggregator;
173+
streaming_config = Arc::new(crate::data_model::StreamingConfig {
174+
aggregation_configs: HashMap::from([(1, value_config), (2, key_config)]),
175+
});
176+
sql_config.query_configs[1] = sql_config.query_configs[1]
177+
.clone()
178+
.add_aggregation(crate::data_model::AggregationReference::new(2, None));
179+
180+
let store = Arc::new(SimpleMapStore::new(
181+
streaming_config.clone(),
182+
CleanupPolicy::NoCleanup,
183+
));
184+
for (start, end, count) in [(0, 5_000, 2.0), (5_000, 10_000, 3.0)] {
185+
let mut value = CountMinSketchAccumulator::new(4, 64);
186+
value.inner.update("host-a", count);
187+
let mut keys = SetAggregatorAccumulator::new();
188+
keys.add_key(KeyByLabelValues {
189+
labels: vec!["host-a".to_string()],
190+
});
191+
store
192+
.insert_precomputed_output(
193+
PrecomputedOutput::new(start, end, None, 1),
194+
Box::new(value),
195+
)
196+
.unwrap();
197+
store
198+
.insert_precomputed_output(
199+
PrecomputedOutput::new(start, end, None, 2),
200+
Box::new(keys),
201+
)
202+
.unwrap();
203+
}
204+
205+
let engine = SimpleEngine::new(
206+
store,
207+
sql_config,
208+
streaming_config,
209+
1_000,
210+
QueryLanguage::sql,
211+
);
212+
let (_, result) = engine
213+
.handle_query_sql(sql_query.to_string(), 10.0)
214+
.expect("SQL dual-population query should execute");
215+
let QueryResult::Vector(vector) = result else {
216+
panic!("expected an instant SQL vector");
217+
};
218+
assert_eq!(vector.values.len(), 1);
219+
assert_eq!(vector.values[0].value, 5.0);
220+
}
221+
140222
#[test]
141223
fn test_temporal_sum_equivalence() {
142224
let scrape_interval_ms = 1000;

0 commit comments

Comments
 (0)