From 0fda0f428f24430930ca3df44895b2b6ee876d57 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 23 Aug 2026 22:38:27 -0400 Subject: [PATCH 1/2] fix(query-engine): range queries expand self-keyed accumulators (top-k) without keys_query MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit execute_range_query_pipeline's single-population branch used each value group's own store-level group_key directly and never called get_keys() on the merged value accumulator, so self-keyed accumulators like CountMinSketchWithHeap (top-k) returned empty over a range instead of expanding into their top-k keys — the instant path already does this via collect_results_same_aggregation. Now every group tries merged.get_keys() first (mirroring the instant path exactly) and only falls back to the precomputed key list — the store's group_key for single-population metrics, or the separate keys aggregation's expansion for dual-population metrics — when get_keys() returns None. This also covers self-keyed accumulators stored under a non-None outer key (real Arroyo/worker.rs ingestion always wraps Some(key), even for empty grouping), not just the None-keyed case from the original report. Fixes #584. Co-Authored-By: Claude Sonnet 5 --- .../src/engines/simple_engine/mod.rs | 48 ++-- .../src/tests/native_range_query_tests.rs | 229 +++++++++++++++++- 2 files changed, 258 insertions(+), 19 deletions(-) diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index e7f610c..2430253 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -1525,13 +1525,16 @@ impl SimpleEngine { window_mode ); - // Resolve, for every value group, which output label-keys it serves: - // its own key for single-population metrics, or every key the merged - // keys aggregation expands it to for dual-population metrics. Mirrors - // collect_results_separate_keys exactly, including its error - // semantics — an unresolvable key set fails the whole range query - // (so callers fall back to Prometheus) instead of silently returning - // a partial result. See #582 review. + // Resolve, for every value group, a fallback key list — used only if + // the merged value accumulator itself doesn't self-key (see below). + // - dual-population metrics: every key the merged keys aggregation + // expands it to. Mirrors collect_results_separate_keys exactly, + // including its error semantics — an unresolvable key set fails + // the whole range query (so callers fall back to Prometheus) + // instead of silently returning a partial result. See #582 + // review. + // - single-population metrics: just its own store-level group key + // (0 or 1 keys). let groups: Vec<( &Vec, Vec, @@ -1542,22 +1545,20 @@ impl SimpleEngine { let timestamped_buckets = all_data .get(group_key) .ok_or_else(|| format!("No value for key: {:?}", group_key))?; - let expansion_keys = keys_precompute + let fallback_keys = keys_precompute .get_keys() .ok_or_else(|| "Keys required for separate aggregation".to_string())?; - Ok((timestamped_buckets, expansion_keys)) + Ok((timestamped_buckets, fallback_keys)) }) .collect::, String>>()?, None => all_data .iter() - .filter_map(|(group_key, buckets)| { - group_key.as_ref().map(|k| (buckets, vec![k.clone()])) - }) + .map(|(group_key, buckets)| (buckets, group_key.clone().into_iter().collect())) .collect(), }; // Process each value group independently - for (timestamped_buckets, expansion_keys) in groups { + for (timestamped_buckets, fallback_keys) in groups { // Build lookup: bucket_start_timestamp -> all buckets sharing // that start. A Sliding aggregation can legitimately return more // than one bucket per start timestamp (#567/#570) — every one of @@ -1568,9 +1569,9 @@ impl SimpleEngine { } debug!( - "Group with {} start-timestamps, expands to keys: {:?}", + "Group with {} start-timestamps, fallback keys: {:?}", bucket_map.len(), - expansion_keys + fallback_keys ); // Iterate by OUTPUT timestamp, not by bucket index @@ -1599,9 +1600,20 @@ impl SimpleEngine { match merger.get_merged() { Ok(merged) => { + // Mirrors collect_results_same_aggregation + // (instant path): a self-keyed accumulator's own + // get_keys() always takes priority when present + // (#584) — it can only be read after merging this + // window's buckets, since which keys are "in" e.g. + // a top-k heap depends on the window's data. + // Non-self-keyed accumulators always return None + // here, so this falls through to fallback_keys + // unchanged from before. + let resolved_keys = + merged.get_keys().unwrap_or_else(|| fallback_keys.clone()); // Query statistic and emit a sample at current_time // for every expanded key sharing this value group. - for key in &expansion_keys { + for key in &resolved_keys { match self.query_precompute_for_statistic( merged.as_ref(), &context.base.metadata.statistic_to_compute, @@ -1626,7 +1638,7 @@ impl SimpleEngine { Err(e) => { debug!( "Failed to get merged result at t={} for keys {:?}: {}", - current_time, expansion_keys, e + current_time, fallback_keys, e ); } } @@ -1634,7 +1646,7 @@ impl SimpleEngine { // No data at all for this window - skip sample debug!( "Keys {:?}: skipping sample at {} - no data in window [{}, {})", - expansion_keys, current_time, window_start, current_time + fallback_keys, current_time, window_start, current_time ); } diff --git a/asap-query-engine/src/tests/native_range_query_tests.rs b/asap-query-engine/src/tests/native_range_query_tests.rs index 891e3c7..6655e66 100644 --- a/asap-query-engine/src/tests/native_range_query_tests.rs +++ b/asap-query-engine/src/tests/native_range_query_tests.rs @@ -29,7 +29,9 @@ mod tests { use crate::engines::query_result::{QueryResult, RangeVectorElement}; use crate::engines::simple_engine::SimpleEngine; use crate::precompute_operators::sum_accumulator::SumAccumulator; - use crate::precompute_operators::{CountMinSketchAccumulator, DeltaSetAggregatorAccumulator}; + use crate::precompute_operators::{ + CountMinSketchAccumulator, CountMinSketchWithHeapAccumulator, DeltaSetAggregatorAccumulator, + }; use crate::stores::simple_map_store::SimpleMapStore; use crate::stores::Store; use crate::tests::test_utilities::engine_factories::create_engine_multi_timestamp_with_window; @@ -371,4 +373,229 @@ mod tests { timestamps ); } + + /// Single-population counterpart to `create_range_engine_dual_input`: one + /// `CountMinSketchWithHeap` (self-keyed, top-k) aggregation, no separate + /// keys aggregation, values stored with `group_key = None`. + fn create_range_engine_self_keyed( + metric: &str, + aggregated_label: &str, + data: Vec<(u64, CountMinSketchWithHeapAccumulator)>, + promql_query: &str, + ) -> SimpleEngine { + let mut aggregation_configs = HashMap::new(); + aggregation_configs.insert( + 1u64, + AggregationConfig { + aggregation_id: 1, + aggregation_type: AggregationType::CountMinSketchWithHeap, + aggregation_sub_type: String::new(), + parameters: HashMap::new(), + grouping_labels: KeyByLabelNames::empty(), + aggregated_labels: KeyByLabelNames::new(vec![aggregated_label.to_string()]), + rollup_labels: KeyByLabelNames::empty(), + original_yaml: String::new(), + window_size_ms: WINDOW_MS, + slide_interval_ms: WINDOW_MS, + window_type: WindowType::Tumbling, + spatial_filter: String::new(), + spatial_filter_normalized: String::new(), + metric: metric.to_string(), + num_aggregates_to_retain: None, + read_count_threshold: None, + table_name: None, + value_column: None, + }, + ); + + let streaming_config = Arc::new(StreamingConfig { + aggregation_configs, + }); + + let store = Arc::new(SimpleMapStore::new( + streaming_config.clone(), + CleanupPolicy::NoCleanup, + )); + + for (timestamp, acc) in data { + let output = PrecomputedOutput::new(timestamp - WINDOW_MS, timestamp, None, 1); + store + .insert_precomputed_output(output, Box::new(acc)) + .unwrap(); + } + + let promql_schema = PromQLSchema::new().add_metric( + metric.to_string(), + KeyByLabelNames::new(vec![aggregated_label.to_string()]), + ); + + let query_config = QueryConfig::new(promql_query.to_string()) + .add_aggregation(AggregationReference::new(1, None)); + + let inference_config = InferenceConfig { + schema: SchemaConfig::PromQL(promql_schema), + query_configs: vec![query_config], + cleanup_policy: CleanupPolicy::NoCleanup, + }; + + SimpleEngine::new( + store, + inference_config, + streaming_config, + WINDOW_MS, + QueryLanguage::promql, + ) + } + + #[tokio::test(flavor = "multi_thread")] + async fn range_query_self_keyed_topk_expands_without_keys_query() { + // #584: execute_range_query_pipeline's single-population branch (no + // separate keys_query) used each value group's own group_key + // directly and never called get_keys() on the value accumulator + // itself. A self-keyed accumulator like CountMinSketchWithHeap + // (top-k) is stored with group_key = None (single population, no + // keys_query) and only exposes its output keys via get_keys() on the + // merged accumulator -- exactly what the instant path + // (collect_results_same_aggregation) already does. Without that + // call, this None-keyed group is dropped and the range query returns + // empty instead of the expanded top-k keys. + let mut sketch = CountMinSketchWithHeapAccumulator::new(3, 1024, 32); + sketch.inner.update("host-a", 30.0); + sketch.inner.update("host-b", 20.0); + sketch.inner.update("host-c", 10.0); + + let query = "topk(5, transfer_events)"; + let engine = + create_range_engine_self_keyed("transfer_events", "srcip", vec![(1000, sketch)], query); + + let result = engine.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0); + let (_, qr) = result.expect("range query failed"); + let elements = matrix_values(qr); + + assert_eq!( + elements.len(), + 3, + "expected the self-keyed accumulator's top-k heap to expand into 3 \ + separate series (one per key) instead of being dropped as an empty result, got: {:?}", + elements + .iter() + .map(|e| e.labels.labels.clone()) + .collect::>() + ); + + let returned: std::collections::HashSet = elements + .iter() + .map(|e| e.labels.labels.last().cloned().unwrap_or_default()) + .collect(); + assert_eq!( + returned, + std::collections::HashSet::from([ + "host-a".to_string(), + "host-b".to_string(), + "host-c".to_string(), + ]), + "expected one expanded series per top-k key" + ); + + let host_a = elements + .iter() + .find(|e| e.labels.labels.last().map(String::as_str) == Some("host-a")) + .expect("host-a series should be present among the expanded top-k keys"); + assert_eq!(host_a.samples.len(), 1); + assert!( + (host_a.samples[0].value - 30.0).abs() < 1e-9, + "expected host-a's count-min-sketch estimate (30.0), got {}", + host_a.samples[0].value + ); + } + + #[tokio::test(flavor = "multi_thread")] + async fn range_query_self_keyed_topk_expands_with_non_none_outer_key() { + // Same bug as range_query_self_keyed_topk_expands_without_keys_query, + // but the value data is stored under a real, non-None outer + // group_key -- mirrors what real Arroyo/worker.rs ingestion actually + // writes for "empty grouping" self-keyed accumulators (deserialize_ + // from_json_arroyo always wraps Some(key), even for empty grouping; + // "None" is specific to this codebase's existing CountMinSketchWithHeap + // test convention, not a guarantee). A fix that only calls + // merged.get_keys() when the outer key is None (rather than always, + // like collect_results_same_aggregation does) would still miss this + // case. + let metric = "transfer_events"; + let aggregated_label = "srcip"; + let mut aggregation_configs = HashMap::new(); + aggregation_configs.insert( + 1u64, + AggregationConfig { + aggregation_id: 1, + aggregation_type: AggregationType::CountMinSketchWithHeap, + aggregation_sub_type: String::new(), + parameters: HashMap::new(), + grouping_labels: KeyByLabelNames::empty(), + aggregated_labels: KeyByLabelNames::new(vec![aggregated_label.to_string()]), + rollup_labels: KeyByLabelNames::empty(), + original_yaml: String::new(), + window_size_ms: WINDOW_MS, + slide_interval_ms: WINDOW_MS, + window_type: WindowType::Tumbling, + spatial_filter: String::new(), + spatial_filter_normalized: String::new(), + metric: metric.to_string(), + num_aggregates_to_retain: None, + read_count_threshold: None, + table_name: None, + value_column: None, + }, + ); + let streaming_config = Arc::new(StreamingConfig { + aggregation_configs, + }); + let store = Arc::new(SimpleMapStore::new( + streaming_config.clone(), + CleanupPolicy::NoCleanup, + )); + + let mut sketch = CountMinSketchWithHeapAccumulator::new(3, 1024, 32); + sketch.inner.update("host-a", 30.0); + sketch.inner.update("host-b", 20.0); + sketch.inner.update("host-c", 10.0); + + // Non-None outer key -- e.g. Some(KeyByLabelValues{labels: [""]}), + // exactly what deserialize_from_json_arroyo produces for "" grouping. + let non_none_key = Some(KeyByLabelValues::new_with_labels(vec![String::new()])); + let output = PrecomputedOutput::new(0, 1000, non_none_key, 1); + store + .insert_precomputed_output(output, Box::new(sketch)) + .unwrap(); + + let promql_schema = PromQLSchema::new().add_metric( + metric.to_string(), + KeyByLabelNames::new(vec![aggregated_label.to_string()]), + ); + let query = "topk(5, transfer_events)"; + let query_config = + QueryConfig::new(query.to_string()).add_aggregation(AggregationReference::new(1, None)); + let inference_config = InferenceConfig { + schema: SchemaConfig::PromQL(promql_schema), + query_configs: vec![query_config], + cleanup_policy: CleanupPolicy::NoCleanup, + }; + let engine = SimpleEngine::new( + store, + inference_config, + streaming_config, + WINDOW_MS, + QueryLanguage::promql, + ); + + let result = engine.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0); + let (_, qr) = result.expect("range query failed"); + let elements = matrix_values(qr); + assert_eq!( + elements.len(), + 3, + "expected top-k expansion even though the outer group_key is Some(..), not None, got: {:?}", + elements.iter().map(|e| e.labels.labels.clone()).collect::>() + ); + } } From 01bbb44f51b2fbb0e3fd2d391fef7b88c53038e9 Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 23 Aug 2026 22:53:41 -0400 Subject: [PATCH 2/2] fix(query-engine): don't let self-keyed value accumulators override keys_query in range queries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review on #587 caught a regression: the previous fix called merged.get_keys() unconditionally on every group's merged value accumulator, including dual-population groups (separate keys_query present). collect_all_results (instant path) branches globally on query shape instead — dual-population always goes through collect_results_separate_keys, which never consults the value accumulator's own get_keys() at all. A real, tested capability-matched config (sql.rs) pairs a self-keyed CountMinSketchWithHeap value aggregation with a separate DeltaSetAggregator keys aggregation; the previous fix let the heap's own (window-shifting) top-k keys silently override the keys aggregation's expansion for that config. Range queries now mirror collect_all_results exactly: dual-population groups always use the keys aggregation's expansion, full stop; only single-population groups let the value accumulator's own get_keys() take priority (falling back to the store-level group key), which is what #584 actually needed. Co-Authored-By: Claude Sonnet 5 --- .../src/engines/simple_engine/mod.rs | 56 ++++++++++++------- .../src/tests/native_range_query_tests.rs | 49 ++++++++++++++++ 2 files changed, 84 insertions(+), 21 deletions(-) diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index 2430253..20af125 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -1525,16 +1525,28 @@ impl SimpleEngine { window_mode ); - // Resolve, for every value group, a fallback key list — used only if - // the merged value accumulator itself doesn't self-key (see below). - // - dual-population metrics: every key the merged keys aggregation - // expands it to. Mirrors collect_results_separate_keys exactly, - // including its error semantics — an unresolvable key set fails - // the whole range query (so callers fall back to Prometheus) - // instead of silently returning a partial result. See #582 - // review. - // - single-population metrics: just its own store-level group key - // (0 or 1 keys). + // Whether the value accumulator's own get_keys() is even consulted + // depends on the query SHAPE (dual- vs single-population), not on a + // per-group fallback — mirrors collect_all_results exactly: + // - dual-population (separate keys_query present): always expand + // via the keys aggregation's get_keys(), full stop. The value + // accumulator's own get_keys() is never consulted, even if the + // value accumulator itself happens to be self-keyed (e.g. a + // CountMinSketchWithHeap value paired with a DeltaSetAggregator + // keys aggregation is a real capability-matched config, see + // sql.rs). Otherwise a self-keyed value accumulator's own + // (possibly different, window-to-window-shifting) keys would + // silently override the keys aggregation's expansion. See #587 + // review. + // - single-population (no separate keys_query): the value + // accumulator's own get_keys() takes priority whenever present + // (#584, self-keyed accumulators like top-k), falling back to + // the store-level group key otherwise. + let is_dual_population = merged_keys.is_some(); + + // Resolve, for every value group, a fallback key list — used + // directly for dual-population groups, or as a fallback for + // single-population groups whose value accumulator doesn't self-key. let groups: Vec<( &Vec, Vec, @@ -1600,17 +1612,19 @@ impl SimpleEngine { match merger.get_merged() { Ok(merged) => { - // Mirrors collect_results_same_aggregation - // (instant path): a self-keyed accumulator's own - // get_keys() always takes priority when present - // (#584) — it can only be read after merging this - // window's buckets, since which keys are "in" e.g. - // a top-k heap depends on the window's data. - // Non-self-keyed accumulators always return None - // here, so this falls through to fallback_keys - // unchanged from before. - let resolved_keys = - merged.get_keys().unwrap_or_else(|| fallback_keys.clone()); + // See is_dual_population above: dual-population + // always trusts the keys aggregation's expansion; + // only single-population lets the value + // accumulator's own get_keys() (read after + // merging this window, since e.g. a top-k heap's + // keys depend on the window's data) take + // priority, falling back to the store-level + // group key otherwise. + let resolved_keys = if is_dual_population { + fallback_keys.clone() + } else { + merged.get_keys().unwrap_or_else(|| fallback_keys.clone()) + }; // Query statistic and emit a sample at current_time // for every expanded key sharing this value group. for key in &resolved_keys { diff --git a/asap-query-engine/src/tests/native_range_query_tests.rs b/asap-query-engine/src/tests/native_range_query_tests.rs index 6655e66..556d356 100644 --- a/asap-query-engine/src/tests/native_range_query_tests.rs +++ b/asap-query-engine/src/tests/native_range_query_tests.rs @@ -200,6 +200,55 @@ mod tests { ); } + #[tokio::test(flavor = "multi_thread")] + async fn range_query_dual_population_self_keyed_value_still_uses_keys_query() { + // Real, tested capability-matching config (sql.rs): CountMinSketchWithHeap + // (self-keyed) as the VALUE aggregation, paired with a separate + // DeltaSetAggregator KEYS aggregation. collect_results_separate_keys + // (instant path) never consults the value accumulator's own + // get_keys() -- it always expands via the keys aggregation. A range + // pipeline that calls merged.get_keys() unconditionally on the merged + // VALUE accumulator would incorrectly let the heap's own internal + // top-k keys override the keys_query's expansion. + let mut heap = CountMinSketchWithHeapAccumulator::new(3, 1024, 32); + // Heap's own top-k keys are deliberately disjoint from the keys_query's + // key, so a wrong implementation is caught by key-set mismatch, not + // just a wrong value. + heap.inner.update("host-x;evt-x", 30.0); + heap.inner.update("host-y;evt-y", 20.0); + + let mut keys = DeltaSetAggregatorAccumulator::new(); + keys.add_key(KeyByLabelValues { + labels: vec!["host-a".to_string(), "evt-1".to_string()], + }); + + let engine = create_range_engine_dual_input( + "event_frequency", + AggregationType::CountMinSketchWithHeap, + AggregationType::DeltaSetAggregator, + vec![], + vec!["host", "event"], + vec![(1000, None, Box::new(heap) as Box)], + vec![(1000, None, Box::new(keys) as Box)], + "count(event_frequency) by (host, event)", + ); + + let query = "count(event_frequency) by (host, event)"; + let result = engine.handle_range_query_promql(query.to_string(), 1.0, 1.5, 1.0); + let (_, qr) = result.expect("range query failed"); + let elements = matrix_values(qr); + + let returned: std::collections::HashSet> = + elements.iter().map(|e| e.labels.labels.clone()).collect(); + assert_eq!( + returned, + std::collections::HashSet::from([vec!["host-a".to_string(), "evt-1".to_string()]]), + "expected the keys_query's key (host-a, evt-1), not the value accumulator's own \ + top-k keys (host-x/host-y), got: {:?}", + returned + ); + } + #[tokio::test(flavor = "multi_thread")] async fn range_query_sliding_window_merges_both_buckets() { // Same fixture as