Skip to content

Commit 2f0c913

Browse files
fix(query-engine): invalidate changed aggregation state
1 parent b25a655 commit 2f0c913

4 files changed

Lines changed: 322 additions & 28 deletions

File tree

‎asap-common/dependencies/rs/asap_types/src/aggregation_config.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use crate::utils::normalize_spatial_filter;
99
use promql_utilities::data_model::KeyByLabelNames;
1010
use promql_utilities::query_logics::enums::AggregationType;
1111

12-
#[derive(Debug, Clone, Serialize, Deserialize)]
12+
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1313
pub struct AggregationConfig {
1414
pub aggregation_id: u64,
1515
pub aggregation_type: AggregationType,

‎asap-common/dependencies/rs/asap_types/src/capability_matching.rs‎

Lines changed: 157 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,11 @@ use promql_utilities::query_logics::enums::AggregationType;
1818
/// Returns the aggregation types that can serve this statistic.
1919
pub fn compatible_agg_types(stat: Statistic) -> &'static [AggregationType] {
2020
match stat {
21-
Statistic::Sum => &[AggregationType::Sum, AggregationType::MultipleSum],
21+
Statistic::Sum => &[
22+
AggregationType::Sum,
23+
AggregationType::MultipleSum,
24+
AggregationType::CountMinSketch,
25+
],
2226
Statistic::Count => &[
2327
AggregationType::CountMinSketch,
2428
AggregationType::CountMinSketchWithHeap,
@@ -187,6 +191,20 @@ pub fn topk_weighting_compatible(
187191
}
188192
}
189193

194+
/// Plain Count-Min Sketches are value-weighted or event-weighted according to
195+
/// their subtype. A sketch with the other subtype cannot serve this statistic.
196+
fn plain_cms_sub_type_compatible(stat: Statistic, config: &AggregationConfig) -> bool {
197+
if config.aggregation_type != AggregationType::CountMinSketch {
198+
return true;
199+
}
200+
201+
match stat {
202+
Statistic::Sum => config.aggregation_sub_type.eq_ignore_ascii_case("sum"),
203+
Statistic::Count => config.aggregation_sub_type.eq_ignore_ascii_case("count"),
204+
_ => false,
205+
}
206+
}
207+
190208
/// Aggregation priority comparator: prefer larger `window_size_ms` (descending).
191209
/// This is a separate function so callers can swap the policy without touching matching logic.
192210
pub fn aggregation_priority(a: &AggregationConfig, b: &AggregationConfig) -> Ordering {
@@ -244,6 +262,7 @@ pub fn find_compatible_aggregation(
244262
&c.spatial_filter_normalized,
245263
&requirements.spatial_filter_normalized,
246264
)
265+
&& plain_cms_sub_type_compatible(stat, c)
247266
&& topk_weighting_compatible(stat, c, requirements.topk_count_events);
248267
if !ok {
249268
debug!(
@@ -466,6 +485,62 @@ mod tests {
466485
assert_eq!(result.unwrap().aggregation_id_for_value, 1);
467486
}
468487

488+
#[test]
489+
fn plain_cms_matching_respects_sum_and_count_subtypes() {
490+
let mut configs = HashMap::new();
491+
configs.insert(
492+
1,
493+
make_config(
494+
1,
495+
"cpu",
496+
"CountMinSketch",
497+
"sum",
498+
300_000,
499+
"tumbling",
500+
&[],
501+
"",
502+
),
503+
);
504+
configs.insert(
505+
2,
506+
make_config(
507+
2,
508+
"cpu",
509+
"CountMinSketch",
510+
"count",
511+
300_000,
512+
"tumbling",
513+
&[],
514+
"",
515+
),
516+
);
517+
configs.insert(
518+
9,
519+
make_config(
520+
9,
521+
"cpu",
522+
"DeltaSetAggregator",
523+
"",
524+
300_000,
525+
"tumbling",
526+
&[],
527+
"",
528+
),
529+
);
530+
531+
let sum =
532+
find_compatible_aggregation(&configs, &req("cpu", &[Statistic::Sum], 300_000, &[], ""))
533+
.expect("SUM should select the value-weighted sketch");
534+
assert_eq!(sum.aggregation_id_for_value, 1);
535+
536+
let count = find_compatible_aggregation(
537+
&configs,
538+
&req("cpu", &[Statistic::Count], 300_000, &[], ""),
539+
)
540+
.expect("COUNT should select the event-weighted sketch");
541+
assert_eq!(count.aggregation_id_for_value, 2);
542+
}
543+
469544
#[test]
470545
fn quantile_any_value_finds_kll() {
471546
let configs = single_config(make_config(
@@ -1076,7 +1151,16 @@ mod tests {
10761151

10771152
#[test]
10781153
fn multi_pop_rejects_tumbling_delta_set_that_cannot_partition_sliding_value_window() {
1079-
let mut value = make_config(10, "req", "CountMinSketch", "", 6_000, "sliding", &[], "");
1154+
let mut value = make_config(
1155+
10,
1156+
"req",
1157+
"CountMinSketch",
1158+
"count",
1159+
6_000,
1160+
"sliding",
1161+
&[],
1162+
"",
1163+
);
10801164
value.slide_interval_ms = 1_000;
10811165
let delta_keys = make_config(
10821166
11,
@@ -1102,7 +1186,16 @@ mod tests {
11021186

11031187
#[test]
11041188
fn multi_pop_accepts_tumbling_delta_set_that_partitions_sliding_value_grid() {
1105-
let mut value = make_config(10, "req", "CountMinSketch", "", 6_000, "sliding", &[], "");
1189+
let mut value = make_config(
1190+
10,
1191+
"req",
1192+
"CountMinSketch",
1193+
"count",
1194+
6_000,
1195+
"sliding",
1196+
&[],
1197+
"",
1198+
);
11061199
value.slide_interval_ms = 1_000;
11071200
let delta_keys = make_config(
11081201
11,
@@ -1127,7 +1220,16 @@ mod tests {
11271220

11281221
#[test]
11291222
fn multi_pop_rejects_tumbling_set_key_on_mismatched_nonzero_grid_step() {
1130-
let value = make_config(10, "req", "CountMinSketch", "", 5_000, "tumbling", &[], "");
1223+
let value = make_config(
1224+
10,
1225+
"req",
1226+
"CountMinSketch",
1227+
"count",
1228+
5_000,
1229+
"tumbling",
1230+
&[],
1231+
"",
1232+
);
11311233
let mut keys = make_config(11, "req", "SetAggregator", "", 5_000, "tumbling", &[], "");
11321234
keys.slide_interval_ms = 1_000;
11331235
let configs = HashMap::from([(10, value), (11, keys)]);
@@ -1141,7 +1243,16 @@ mod tests {
11411243

11421244
#[test]
11431245
fn tumbling_set_pairing_normalizes_zero_slide_to_window_size() {
1144-
let mut value = make_config(10, "req", "CountMinSketch", "", 5_000, "tumbling", &[], "");
1246+
let mut value = make_config(
1247+
10,
1248+
"req",
1249+
"CountMinSketch",
1250+
"count",
1251+
5_000,
1252+
"tumbling",
1253+
&[],
1254+
"",
1255+
);
11451256
let mut key = make_config(11, "req", "SetAggregator", "", 5_000, "tumbling", &[], "");
11461257
value.slide_interval_ms = 0;
11471258
key.slide_interval_ms = 0;
@@ -1152,7 +1263,16 @@ mod tests {
11521263

11531264
#[test]
11541265
fn set_pairing_rejects_each_grid_mismatch_dimension() {
1155-
let value = make_config(10, "req", "CountMinSketch", "", 5_000, "sliding", &[], "");
1266+
let value = make_config(
1267+
10,
1268+
"req",
1269+
"CountMinSketch",
1270+
"count",
1271+
5_000,
1272+
"sliding",
1273+
&[],
1274+
"",
1275+
);
11561276
let mut key = make_config(11, "req", "SetAggregator", "", 5_000, "sliding", &[], "");
11571277
key.slide_interval_ms = 1_000;
11581278
assert!(!key_agg_compatible_with_value(&value, &key));
@@ -1166,7 +1286,16 @@ mod tests {
11661286

11671287
#[test]
11681288
fn delta_set_pairing_truth_table_checks_both_divisors() {
1169-
let mut value = make_config(10, "req", "CountMinSketch", "", 6_000, "sliding", &[], "");
1289+
let mut value = make_config(
1290+
10,
1291+
"req",
1292+
"CountMinSketch",
1293+
"count",
1294+
6_000,
1295+
"sliding",
1296+
&[],
1297+
"",
1298+
);
11701299
value.slide_interval_ms = 2_000;
11711300
let key_valid = make_config(
11721301
11,
@@ -1212,7 +1341,16 @@ mod tests {
12121341

12131342
#[test]
12141343
fn delta_set_pairing_for_tumbling_values_does_not_apply_sliding_rules() {
1215-
let value = make_config(10, "req", "CountMinSketch", "", 6_000, "tumbling", &[], "");
1344+
let value = make_config(
1345+
10,
1346+
"req",
1347+
"CountMinSketch",
1348+
"count",
1349+
6_000,
1350+
"tumbling",
1351+
&[],
1352+
"",
1353+
);
12161354
let key = make_config(
12171355
11,
12181356
"req",
@@ -1228,7 +1366,16 @@ mod tests {
12281366

12291367
#[test]
12301368
fn matching_skips_incompatible_key_candidate_and_selects_compatible_one() {
1231-
let value = make_config(10, "req", "CountMinSketch", "", 6_000, "sliding", &[], "");
1369+
let value = make_config(
1370+
10,
1371+
"req",
1372+
"CountMinSketch",
1373+
"count",
1374+
6_000,
1375+
"sliding",
1376+
&[],
1377+
"",
1378+
);
12321379
let mut incompatible =
12331380
make_config(11, "req", "SetAggregator", "", 6_000, "sliding", &[], "");
12341381
incompatible.slide_interval_ms = 2_000;
@@ -1255,7 +1402,7 @@ mod tests {
12551402
2,
12561403
"cpu",
12571404
"CountMinSketch",
1258-
"",
1405+
"count",
12591406
300_000,
12601407
"tumbling",
12611408
&["job"],

‎asap-query-engine/src/precompute_engine/engine.rs‎

Lines changed: 54 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,8 @@ impl PrecomputeEngineHandle {
3838
&self,
3939
config: &StreamingConfig,
4040
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
41+
validate_cms_aggregation_configs(config)?;
42+
4143
let agg_configs_map: HashMap<u64, Arc<AggregationConfig>> = config
4244
.get_all_aggregation_configs()
4345
.iter()
@@ -155,7 +157,7 @@ impl PrecomputeEngine {
155157
/// Start the precompute engine. This spawns worker tasks and all registered
156158
/// ingest sources, then blocks until shutdown.
157159
pub async fn run(mut self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
158-
validate_startup_aggregation_configs(&self.streaming_config)?;
160+
validate_cms_aggregation_configs(&self.streaming_config)?;
159161

160162
let num_workers = self.config.num_workers;
161163

@@ -256,12 +258,13 @@ impl PrecomputeEngine {
256258
}
257259
}
258260

259-
/// Validate aggregation configs before starting any worker or ingest task.
261+
/// Validate plain Count-Min Sketch configs before starting workers or ingest
262+
/// tasks, and before accepting runtime replacements.
260263
///
261264
/// Plain Count-Min Sketch configs carry their SUM-versus-COUNT contract in
262265
/// `aggregation_sub_type`; allowing an invalid value to reach the lazy worker
263266
/// path would leave the engine running while silently losing that contract.
264-
fn validate_startup_aggregation_configs(
267+
fn validate_cms_aggregation_configs(
265268
streaming_config: &StreamingConfig,
266269
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
267270
for (&aggregation_id, config) in streaming_config.get_all_aggregation_configs() {
@@ -346,4 +349,52 @@ mod tests {
346349
assert!(err.to_string().contains("aggregation_id 1"));
347350
assert!(err.to_string().contains("sum") && err.to_string().contains("count"));
348351
}
352+
353+
#[tokio::test]
354+
async fn runtime_update_rejects_invalid_cms_subtype_without_replacing_config() {
355+
let mut parameters = HashMap::new();
356+
parameters.insert("depth".to_string(), json!(3_u64));
357+
parameters.insert("width".to_string(), json!(128_u64));
358+
let valid = AggregationConfig::new(
359+
1,
360+
AggregationType::CountMinSketch,
361+
"sum".to_string(),
362+
parameters,
363+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
364+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![
365+
"host".to_string()
366+
]),
367+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
368+
String::new(),
369+
1_000,
370+
1_000,
371+
WindowType::Tumbling,
372+
"requests_total".to_string(),
373+
"requests_total".to_string(),
374+
None,
375+
None,
376+
None,
377+
None,
378+
);
379+
let engine = PrecomputeEngine::new(
380+
PrecomputeEngineConfig::default(),
381+
Arc::new(StreamingConfig::new(HashMap::from([(1, valid.clone())]))),
382+
Arc::new(NoopOutputSink::new()),
383+
vec![],
384+
);
385+
let handle = engine.handle();
386+
387+
let mut invalid = valid;
388+
invalid.aggregation_sub_type.clear();
389+
let result = handle
390+
.update_streaming_config(&StreamingConfig::new(HashMap::from([(1, invalid)])))
391+
.await;
392+
393+
let err = result.expect_err("invalid runtime CMS subtype must be rejected");
394+
assert!(err.to_string().contains("aggregation_id 1"));
395+
assert!(err.to_string().contains("sum") && err.to_string().contains("count"));
396+
let current = handle.ingest_agg_configs.load();
397+
assert_eq!(current.len(), 1);
398+
assert_eq!(current[0].aggregation_sub_type, "sum");
399+
}
349400
}

0 commit comments

Comments
 (0)