Skip to content

Commit 3f4f83c

Browse files
fix(query-engine): MultipleSum weights samples by 1 for count sub_type (#717)
* fix(query-engine): MultipleSum weights samples by 1 for count sub_type MultipleSumAccumulatorUpdater always summed the raw sample value, so an ingested count precompute held Σvalue instead of Σ1 whenever aggregation_sub_type=="count". Mirror the CountMinSketch fix: extract a shared, strict sub_type-to-weighting helper and use it for both. Fixes #503 Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RzTeeyqi31cRJ1BkMFWWRx * fix(planner,query-engine): propagate MultipleSum's strict subtype contract #503 made MultipleSum reject an empty/invalid aggregation_sub_type, but two producers/validators of MultipleSum configs weren't updated to match (roborev review #197): - candidate_gen.rs derived "" for (Sum, MultipleSum) candidates instead of "sum", so every optimizer-selected MultipleSum candidate for a Sum statistic would fail to ingest. - engine startup validation only ran create_accumulator_updater over CountMinSketch/CountMinSketchWithHeap configs, so an invalid MultipleSum subtype would pass startup and fail lazily per-batch instead of failing fast. Deferred: capability_matching.rs's Sum/Count candidate lists aren't subtype-aware for MultipleSum either, but neither the streaming planner nor the optimizer ever generates a MultipleSum+"count" config today (Statistic::Count routes to CountMinSketch/CountMinSketchWithHeap only), so this is unreachable outside a hand-authored config -- same category as MultipleSubpopulation, which #503 already scoped out. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RzTeeyqi31cRJ1BkMFWWRx --------- Co-authored-by: Claude Sonnet 5 <noreply@anthropic.com>
1 parent 0b5309c commit 3f4f83c

4 files changed

Lines changed: 272 additions & 26 deletions

File tree

‎asap-planner-rs/src/optimizer/candidate_gen.rs‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -239,7 +239,7 @@ fn derive_sub_type(stat: Statistic, agg_type: AggregationType) -> String {
239239
(Statistic::Min, _) => "min",
240240
(Statistic::Max, _) => "max",
241241
(Statistic::Topk, _) => "topk",
242-
(Statistic::Sum, AggregationType::CountMinSketch) => "sum",
242+
(Statistic::Sum, AggregationType::CountMinSketch | AggregationType::MultipleSum) => "sum",
243243
(Statistic::Count, AggregationType::CountMinSketch) => "count",
244244
_ => "",
245245
}
@@ -363,6 +363,27 @@ mod tests {
363363
.any(|c| c.config.is_none() && c.query_method == QueryMethod::Exact));
364364
}
365365

366+
#[test]
367+
fn multiple_sum_candidates_get_a_non_empty_sub_type() {
368+
// MultipleSum's factory now rejects an empty aggregation_sub_type (#503) --
369+
// the optimizer must derive "sum" for it, same as it already does for
370+
// CountMinSketch, or every MultipleSum candidate fails to ingest.
371+
let aqe = make_aqe(Statistic::Sum, 300_000, 60_000);
372+
let candidates = enumerate_candidates(&aqe, 15_000);
373+
let multiple_sum_configs: Vec<_> = candidates
374+
.iter()
375+
.filter_map(|c| c.config.as_ref())
376+
.filter(|cfg| cfg.aggregation_type == AggregationType::MultipleSum)
377+
.collect();
378+
assert!(
379+
!multiple_sum_configs.is_empty(),
380+
"expected at least one MultipleSum candidate"
381+
);
382+
for cfg in multiple_sum_configs {
383+
assert_eq!(cfg.aggregation_sub_type, "sum");
384+
}
385+
}
386+
366387
#[test]
367388
fn stamps_dataset_label_group_count_on_every_candidate() {
368389
let aqe = make_aqe(Statistic::Sum, 300_000, 60_000);

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

Lines changed: 135 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -461,22 +461,18 @@ impl AccumulatorUpdater for DeltaSetAggregatorUpdater {
461461

462462
pub struct MultipleSumAccumulatorUpdater {
463463
acc: MultipleSumAccumulator,
464+
count_events: bool,
464465
}
465466

466467
impl MultipleSumAccumulatorUpdater {
467-
pub fn new() -> Self {
468+
pub fn new(count_events: bool) -> Self {
468469
Self {
469470
acc: MultipleSumAccumulator::new(),
471+
count_events,
470472
}
471473
}
472474
}
473475

474-
impl Default for MultipleSumAccumulatorUpdater {
475-
fn default() -> Self {
476-
Self::new()
477-
}
478-
}
479-
480476
impl AccumulatorUpdater for MultipleSumAccumulatorUpdater {
481477
fn update_single(&mut self, _value: f64, _timestamp_ms: i64) {
482478
debug_assert!(
@@ -486,7 +482,8 @@ impl AccumulatorUpdater for MultipleSumAccumulatorUpdater {
486482
}
487483

488484
fn update_keyed(&mut self, key: &KeyByLabelValues, value: f64, _timestamp_ms: i64) {
489-
self.acc.update(key.clone(), value);
485+
let weight = if self.count_events { 1.0 } else { value };
486+
self.acc.update(key.clone(), weight);
490487
}
491488

492489
impl_accumulator_methods!(acc);
@@ -845,20 +842,19 @@ fn cms_params(config: &AggregationConfig) -> Result<(usize, usize), String> {
845842
Ok((row_num, col_num))
846843
}
847844

848-
/// Resolve the weighting semantics for a plain Count-Min Sketch.
845+
/// Resolve the weighting semantics for aggregation types that dispatch SUM vs
846+
/// COUNT via `aggregation_sub_type` (plain CountMinSketch, MultipleSum).
849847
///
850-
/// Unlike the heap variant, plain CMS uses `aggregation_sub_type` to
851-
/// distinguish approximate SUM from approximate COUNT. Do not silently
852-
/// default malformed configs: the wrong weighting produces plausible but
853-
/// incorrect results.
854-
fn cms_count_events_for_sub_type(sub_type: &str) -> Result<bool, String> {
848+
/// Do not silently default malformed configs: the wrong weighting produces
849+
/// plausible but incorrect results.
850+
fn sum_or_count_events_for_sub_type(agg_type_name: &str, sub_type: &str) -> Result<bool, String> {
855851
if sub_type.eq_ignore_ascii_case("count") {
856852
Ok(true)
857853
} else if sub_type.eq_ignore_ascii_case("sum") {
858854
Ok(false)
859855
} else {
860856
Err(format!(
861-
"CountMinSketch requires aggregation_sub_type 'sum' or 'count', got '{sub_type}'"
857+
"{agg_type_name} requires aggregation_sub_type 'sum' or 'count', got '{sub_type}'"
862858
))
863859
}
864860
}
@@ -948,7 +944,7 @@ pub fn create_accumulator_updater(
948944
other => Err(format!("Unknown SingleSubpopulation sub_type '{other}'")),
949945
},
950946
AggregationType::MultipleSubpopulation => match sub_type {
951-
"Sum" | "sum" => Ok(Box::new(MultipleSumAccumulatorUpdater::new())),
947+
"Sum" | "sum" => Ok(Box::new(MultipleSumAccumulatorUpdater::new(false))),
952948
"Min" | "min" => Ok(Box::new(MultipleMinMaxAccumulatorUpdater::new(false))),
953949
"Max" | "max" => Ok(Box::new(MultipleMinMaxAccumulatorUpdater::new(true))),
954950
"Increase" | "increase" => Ok(Box::new(MultipleIncreaseAccumulatorUpdater::new())),
@@ -969,7 +965,10 @@ pub fn create_accumulator_updater(
969965
AggregationType::DatasketchesKLL => {
970966
Ok(Box::new(KllAccumulatorUpdater::new(kll_k_param(config)?)))
971967
}
972-
AggregationType::MultipleSum => Ok(Box::new(MultipleSumAccumulatorUpdater::new())),
968+
AggregationType::MultipleSum => {
969+
let count_events = sum_or_count_events_for_sub_type("MultipleSum", sub_type)?;
970+
Ok(Box::new(MultipleSumAccumulatorUpdater::new(count_events)))
971+
}
973972
AggregationType::MultipleIncrease => {
974973
Ok(Box::new(MultipleIncreaseAccumulatorUpdater::new()))
975974
}
@@ -983,7 +982,7 @@ pub fn create_accumulator_updater(
983982
AggregationType::Increase => Ok(Box::new(IncreaseAccumulatorUpdater::new())),
984983
AggregationType::CountMinSketch => {
985984
let (row_num, col_num) = cms_params(config)?;
986-
let count_events = cms_count_events_for_sub_type(sub_type)?;
985+
let count_events = sum_or_count_events_for_sub_type("CountMinSketch", sub_type)?;
987986
Ok(Box::new(CmsAccumulatorUpdater::new(
988987
row_num,
989988
col_num,
@@ -1017,6 +1016,7 @@ pub fn create_accumulator_updater(
10171016
#[cfg(test)]
10181017
mod tests {
10191018
use super::*;
1019+
use crate::data_model::MultipleSubpopulationAggregate;
10201020
use asap_types::enums::{AggregationType, WindowType};
10211021

10221022
#[test]
@@ -1066,7 +1066,7 @@ mod tests {
10661066

10671067
#[test]
10681068
fn test_multiple_sum_updater() {
1069-
let mut updater = MultipleSumAccumulatorUpdater::new();
1069+
let mut updater = MultipleSumAccumulatorUpdater::new(false);
10701070
assert!(updater.is_keyed());
10711071

10721072
let key_a = KeyByLabelValues::new_with_labels(vec!["a".to_string()]);
@@ -1137,7 +1137,7 @@ mod tests {
11371137
)));
11381138
assert!(config_is_keyed(&make_config(
11391139
AggregationType::MultipleSum,
1140-
""
1140+
"sum"
11411141
)));
11421142
assert!(config_is_keyed(&make_config(
11431143
AggregationType::MultipleIncrease,
@@ -1157,7 +1157,7 @@ mod tests {
11571157
for (agg_type, sub_type) in &[
11581158
(AggregationType::SingleSubpopulation, "Sum"),
11591159
(AggregationType::MultipleSubpopulation, "Sum"),
1160-
(AggregationType::MultipleSum, ""),
1160+
(AggregationType::MultipleSum, "sum"),
11611161
] {
11621162
let config = make_config(*agg_type, sub_type);
11631163
let updater = create_accumulator_updater(&config).unwrap();
@@ -1615,6 +1615,120 @@ mod tests {
16151615
assert!(err.contains(" count "));
16161616
}
16171617

1618+
fn multiple_sum_config(sub_type: &str) -> AggregationConfig {
1619+
AggregationConfig::new(
1620+
200,
1621+
AggregationType::MultipleSum,
1622+
sub_type.to_string(),
1623+
std::collections::HashMap::new(),
1624+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
1625+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
1626+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
1627+
String::new(),
1628+
1_000,
1629+
1_000,
1630+
WindowType::Tumbling,
1631+
"test_metric".to_string(),
1632+
"test_metric".to_string(),
1633+
None,
1634+
None,
1635+
None,
1636+
None,
1637+
)
1638+
}
1639+
1640+
#[test]
1641+
fn test_multiple_sum_count_subtype_uses_unit_weight() {
1642+
let config = multiple_sum_config("count");
1643+
let mut updater = create_accumulator_updater(&config).unwrap();
1644+
let key = KeyByLabelValues::new_with_labels(vec!["host-a".to_string()]);
1645+
1646+
for _ in 0..5 {
1647+
updater.update_keyed(&key, 1_000.0, 0);
1648+
}
1649+
1650+
let acc = updater.take_accumulator();
1651+
let multiple_sum = acc
1652+
.as_any()
1653+
.downcast_ref::<MultipleSumAccumulator>()
1654+
.expect("MultipleSum accumulator");
1655+
assert_eq!(
1656+
multiple_sum
1657+
.query(
1658+
promql_utilities::query_logics::enums::Statistic::Count,
1659+
&key,
1660+
None
1661+
)
1662+
.unwrap(),
1663+
5.0
1664+
);
1665+
}
1666+
1667+
#[test]
1668+
fn test_multiple_sum_sum_subtype_uses_sample_weight() {
1669+
let config = multiple_sum_config("sum");
1670+
let mut updater = create_accumulator_updater(&config).unwrap();
1671+
let key = KeyByLabelValues::new_with_labels(vec!["host-a".to_string()]);
1672+
1673+
for _ in 0..5 {
1674+
updater.update_keyed(&key, 10.0, 0);
1675+
}
1676+
1677+
let acc = updater.take_accumulator();
1678+
let multiple_sum = acc
1679+
.as_any()
1680+
.downcast_ref::<MultipleSumAccumulator>()
1681+
.expect("MultipleSum accumulator");
1682+
assert_eq!(
1683+
multiple_sum
1684+
.query(
1685+
promql_utilities::query_logics::enums::Statistic::Sum,
1686+
&key,
1687+
None
1688+
)
1689+
.unwrap(),
1690+
50.0
1691+
);
1692+
}
1693+
1694+
#[test]
1695+
fn test_multiple_sum_rejects_empty_subtype() {
1696+
let config = multiple_sum_config("");
1697+
let err = match create_accumulator_updater(&config) {
1698+
Ok(_) => panic!("empty MultipleSum subtype must fail"),
1699+
Err(err) => err,
1700+
};
1701+
assert!(err.contains("sum") && err.contains("count"));
1702+
}
1703+
1704+
#[test]
1705+
fn test_multiple_sum_rejects_unknown_subtype() {
1706+
let config = multiple_sum_config("frequency");
1707+
let err = match create_accumulator_updater(&config) {
1708+
Ok(_) => panic!("unknown MultipleSum subtype must fail"),
1709+
Err(err) => err,
1710+
};
1711+
assert!(err.contains("frequency"));
1712+
}
1713+
1714+
#[test]
1715+
fn test_multiple_sum_accepts_case_insensitive_subtype() {
1716+
for sub_type in ["COUNT", "SuM"] {
1717+
create_accumulator_updater(&multiple_sum_config(sub_type))
1718+
.unwrap_or_else(|err| panic!("subtype '{sub_type}' should be accepted: {err}"));
1719+
}
1720+
}
1721+
1722+
#[test]
1723+
fn test_multiple_sum_rejects_whitespace_padded_subtype() {
1724+
let config = multiple_sum_config(" count ");
1725+
let err = match create_accumulator_updater(&config) {
1726+
Ok(_) => panic!("whitespace-padded MultipleSum subtype must fail"),
1727+
Err(err) => err,
1728+
};
1729+
assert!(err.contains(" count "));
1730+
}
1731+
16181732
fn cms_heap_params_required() -> std::collections::HashMap<String, serde_json::Value> {
16191733
let mut p = std::collections::HashMap::new();
16201734
p.insert("depth".to_string(), serde_json::json!(3_u64));

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

Lines changed: 50 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -258,9 +258,10 @@ impl PrecomputeEngine {
258258

259259
/// Validate aggregation configs before starting any worker or ingest task.
260260
///
261-
/// Count-Min Sketch configs carry their semantic contract in subtype fields or
262-
/// parameters; allowing an invalid value to reach the lazy worker path would
263-
/// leave the engine running while silently losing that contract.
261+
/// Count-Min Sketch and MultipleSum configs carry their semantic contract in
262+
/// subtype fields or parameters; allowing an invalid value to reach the lazy
263+
/// worker path would leave the engine running while silently losing that
264+
/// contract.
264265
fn validate_startup_aggregation_configs(
265266
streaming_config: &StreamingConfig,
266267
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
@@ -269,7 +270,9 @@ fn validate_startup_aggregation_configs(
269270
for (&aggregation_id, config) in streaming_config.get_all_aggregation_configs() {
270271
if !matches!(
271272
config.aggregation_type,
272-
AggregationType::CountMinSketch | AggregationType::CountMinSketchWithHeap
273+
AggregationType::CountMinSketch
274+
| AggregationType::CountMinSketchWithHeap
275+
| AggregationType::MultipleSum
273276
) {
274277
continue;
275278
}
@@ -362,6 +365,49 @@ mod tests {
362365
assert!(err.to_string().contains("sum") && err.to_string().contains("count"));
363366
}
364367

368+
#[tokio::test]
369+
async fn run_rejects_invalid_multiple_sum_subtype_before_starting_workers() {
370+
let multiple_sum = AggregationConfig::new(
371+
1,
372+
AggregationType::MultipleSum,
373+
String::new(),
374+
HashMap::new(),
375+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
376+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![
377+
"host".to_string()
378+
]),
379+
promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]),
380+
String::new(),
381+
1_000,
382+
1_000,
383+
WindowType::Tumbling,
384+
"requests_total".to_string(),
385+
"requests_total".to_string(),
386+
None,
387+
None,
388+
None,
389+
None,
390+
);
391+
let engine = PrecomputeEngine::new(
392+
PrecomputeEngineConfig {
393+
num_workers: 1,
394+
late_data_policy: LateDataPolicy::Drop,
395+
..PrecomputeEngineConfig::default()
396+
},
397+
Arc::new(StreamingConfig::new(HashMap::from([(1, multiple_sum)]))),
398+
Arc::new(NoopOutputSink::new()),
399+
vec![Box::new(ShutdownSource)],
400+
);
401+
402+
let result = engine.run().await;
403+
let err = match result {
404+
Ok(()) => panic!("invalid MultipleSum subtype must fail before startup"),
405+
Err(err) => err,
406+
};
407+
assert!(err.to_string().contains("aggregation_id 1"));
408+
assert!(err.to_string().contains("sum") && err.to_string().contains("count"));
409+
}
410+
365411
#[tokio::test]
366412
async fn run_rejects_invalid_cms_with_heap_subtype_before_starting_workers() {
367413
let mut parameters = HashMap::new();

0 commit comments

Comments
 (0)