From c536785436940f9b1f3114c6a145b6f82f57bf6d Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Sun, 4 Oct 2026 11:17:27 -0400 Subject: [PATCH] fix(planner): apply shared cleanup thresholds --- .../rs/asap_types/src/streaming_config.rs | 61 ++++++- asap-planner-rs/src/optimizer/translator.rs | 165 ++++++++++++++++-- 2 files changed, 209 insertions(+), 17 deletions(-) diff --git a/asap-common/dependencies/rs/asap_types/src/streaming_config.rs b/asap-common/dependencies/rs/asap_types/src/streaming_config.rs index e73dbe37..df300e50 100644 --- a/asap-common/dependencies/rs/asap_types/src/streaming_config.rs +++ b/asap-common/dependencies/rs/asap_types/src/streaming_config.rs @@ -58,12 +58,17 @@ impl StreamingConfig { for aggregation in &query_config.aggregations { let aggregation_id = aggregation.aggregation_id; if let Some(num_aggregates) = aggregation.num_aggregates_to_retain { - // OLD: Keep last value only (for backwards compatibility) - retention_map.insert(aggregation_id, num_aggregates); + retention_map + .entry(aggregation_id) + .and_modify(|retain| *retain = (*retain).max(num_aggregates)) + .or_insert(num_aggregates); + } - // NEW: Sum up num_aggregates_to_retain across all queries - *read_count_threshold_map.entry(aggregation_id).or_insert(0) += - num_aggregates; + if let Some(read_count_threshold) = aggregation.read_count_threshold { + let threshold = read_count_threshold_map.entry(aggregation_id).or_insert(0); + *threshold = threshold + .checked_add(read_count_threshold) + .expect("read-count threshold overflowed"); } } } @@ -136,6 +141,10 @@ impl Default for StreamingConfig { mod tests { use super::*; use crate::aggregation_config::AggregationConfigError; + use crate::aggregation_reference::AggregationReference; + use crate::enums::{CleanupPolicy, QueryLanguage}; + use crate::inference_config::InferenceConfig; + use crate::query_config::QueryConfig; #[test] fn rejects_heap_config_with_invalid_sub_type() { @@ -365,4 +374,46 @@ aggregations: StreamingConfig::from_yaml_data(&yaml, None) .expect("MinMax config with 'MAX' subtype must be accepted"); } + + #[test] + fn shared_references_keep_the_largest_circular_buffer_retention() { + // A later short-range query must not reduce a shared aggregation's retention. + let yaml = minmax_yaml("MinMax", "max"); + for references in [[7, 2], [2, 7]] { + let mut inference = + InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::CircularBuffer); + inference.query_configs = references + .into_iter() + .map(|retain| { + QueryConfig::new(format!("max(metric[{retain}m])")) + .add_aggregation(AggregationReference::new(1, Some(retain))) + }) + .collect(); + + let config = StreamingConfig::from_yaml_data(&yaml, Some(&inference)) + .expect("shared aggregation references should load"); + + assert_eq!(config[1].num_aggregates_to_retain, Some(7)); + assert_eq!(config[1].read_count_threshold, None); + } + } + + #[test] + fn shared_references_sum_read_based_thresholds() { + // Each shared query consumes a read before the aggregate can be cleaned up. + let yaml = minmax_yaml("MinMax", "max"); + let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::ReadBased); + inference.query_configs = vec![ + QueryConfig::new("max(metric[10m])".into()) + .add_aggregation(AggregationReference::with_read_count_threshold(1, Some(7))), + QueryConfig::new("max(metric[1m])".into()) + .add_aggregation(AggregationReference::with_read_count_threshold(1, Some(2))), + ]; + + let config = StreamingConfig::from_yaml_data(&yaml, Some(&inference)) + .expect("shared aggregation references should load"); + + assert_eq!(config[1].num_aggregates_to_retain, None); + assert_eq!(config[1].read_count_threshold, Some(9)); + } } diff --git a/asap-planner-rs/src/optimizer/translator.rs b/asap-planner-rs/src/optimizer/translator.rs index 465e5d0b..8a307e61 100644 --- a/asap-planner-rs/src/optimizer/translator.rs +++ b/asap-planner-rs/src/optimizer/translator.rs @@ -2,6 +2,7 @@ use asap_types::aggregation_reference::AggregationReference; use asap_types::inference_config::InferenceConfig; use asap_types::query_config::QueryConfig; use asap_types::streaming_config::StreamingConfig; +use std::collections::HashMap; use super::solution::{OptimizerSolution, QueryMethod}; @@ -9,28 +10,38 @@ use super::solution::{OptimizerSolution, QueryMethod}; /// Arroyo and the query engine. /// pub fn translate(solution: &OptimizerSolution) -> (StreamingConfig, InferenceConfig) { - let streaming_config = build_streaming_config(solution); let inference_config = build_inference_config(solution); + let streaming_config = build_streaming_config(solution, &inference_config); (streaming_config, inference_config) } -fn build_streaming_config(solution: &OptimizerSolution) -> StreamingConfig { - // Deployed configs map directly to AggregationConfigs — the types are the same. - StreamingConfig::new(solution.deployed_configs().clone()) +fn build_streaming_config( + solution: &OptimizerSolution, + inference_config: &InferenceConfig, +) -> StreamingConfig { + let mut configs = solution.deployed_configs().clone(); + for (aggregation_id, threshold) in read_count_thresholds(inference_config) { + let config = configs + .get_mut(&aggregation_id) + .expect("every assigned aggregation must be deployed"); + config.num_aggregates_to_retain = None; + config.read_count_threshold = Some(threshold); + } + StreamingConfig::new(configs) } fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig { use asap_types::enums::{CleanupPolicy, QueryLanguage}; - let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::NoCleanup); + let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::ReadBased); for assignment in &solution.assignments { let aggregation_id = assignment.aggregation_id; - let retain = retention_count_for_assignment(&assignment.query_method); - let agg_ref = AggregationReference::new(aggregation_id, Some(retain)); + let cleanup_count = cleanup_count_for_assignment(&assignment.query_method); + let agg_ref = + AggregationReference::with_read_count_threshold(aggregation_id, Some(cleanup_count)); let key_ref = assignment.key_aggregation_id.map(|key_id| { - let key_retain = solution.deployed_configs()[&key_id].num_aggregates_to_retain; - AggregationReference::new(key_id, key_retain) + AggregationReference::with_read_count_threshold(key_id, Some(cleanup_count)) }); for query_string in &assignment.item.query_strings { @@ -47,9 +58,26 @@ fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig { inference } -/// For a Merge assignment, the number of retained windows to configure in the -/// inference config (num_aggregates_to_retain on the AggregationReference). -pub fn retention_count_for_assignment(query_method: &QueryMethod) -> u64 { +fn read_count_thresholds(inference_config: &InferenceConfig) -> HashMap { + let mut thresholds = HashMap::new(); + for query_config in &inference_config.query_configs { + for aggregation in &query_config.aggregations { + let Some(cleanup_count) = aggregation.read_count_threshold else { + continue; + }; + let threshold = thresholds + .entry(aggregation.aggregation_id) + .or_insert(0_u64); + *threshold = threshold + .checked_add(cleanup_count) + .expect("read-count threshold overflowed"); + } + } + thresholds +} + +/// Number of aggregate reads an assignment needs before its windows may be cleaned up. +pub fn cleanup_count_for_assignment(query_method: &QueryMethod) -> u64 { match query_method { QueryMethod::Direct => 1, QueryMethod::Merge { num_windows } => *num_windows, @@ -77,3 +105,116 @@ impl TranslationSummary { } } } + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use asap_types::aggregation_config::AggregationConfig; + use asap_types::enums::{CleanupPolicy, WindowType}; + use asap_types::query_requirements::QueryRequirements; + use promql_utilities::data_model::KeyByLabelNames; + use promql_utilities::query_logics::enums::{AggregationType, Statistic}; + + use super::super::solution::{AQEAssignment, OptimizerItem}; + use super::*; + + fn assignment(aggregation_id: u64, query_method: QueryMethod) -> AQEAssignment { + AQEAssignment { + item: OptimizerItem { + requirements: QueryRequirements { + metric: "metric".into(), + statistics: vec![Statistic::Sum], + data_range_ms: 60_000, + grouping_labels: KeyByLabelNames::empty(), + spatial_filter_normalized: String::new(), + topk_count_events: None, + topk_by_labels: None, + }, + query_strings: vec!["sum(metric)".into()], + query_frequency_hz: 1.0, + t_repeat_ms: 60_000, + accuracy_sla: 0.01, + latency_sla: 1.0, + }, + aggregation_id, + key_aggregation_id: None, + query_method, + estimated_query_cost_per_sec: 0.0, + } + } + + #[test] + fn translation_uses_read_based_cleanup_and_sums_shared_thresholds() { + let mut solution = OptimizerSolution::empty(); + let aggregation_id = solution.register_config(AggregationConfig::new( + 0, + AggregationType::Sum, + "sum".into(), + HashMap::new(), + KeyByLabelNames::empty(), + KeyByLabelNames::empty(), + KeyByLabelNames::empty(), + String::new(), + 60_000, + 60_000, + WindowType::Tumbling, + String::new(), + "metric".into(), + None, + None, + None, + None, + )); + solution.assignments = vec![ + assignment(aggregation_id, QueryMethod::Merge { num_windows: 3 }), + assignment(aggregation_id, QueryMethod::Direct), + ]; + + let (streaming, inference) = translate(&solution); + + assert_eq!(inference.cleanup_policy, CleanupPolicy::ReadBased); + assert_eq!(streaming[aggregation_id].read_count_threshold, Some(4)); + assert!(inference.query_configs.iter().all(|query_config| { + query_config.aggregations.iter().all(|aggregation| { + aggregation.num_aggregates_to_retain.is_none() + && aggregation.read_count_threshold.is_some() + }) + })); + } + + #[test] + fn translation_counts_each_query_string_in_a_shared_threshold() { + let mut solution = OptimizerSolution::empty(); + let aggregation_id = solution.register_config(AggregationConfig::new( + 0, + AggregationType::Sum, + "sum".into(), + HashMap::new(), + KeyByLabelNames::empty(), + KeyByLabelNames::empty(), + KeyByLabelNames::empty(), + String::new(), + 60_000, + 60_000, + WindowType::Tumbling, + String::new(), + "metric".into(), + None, + None, + None, + None, + )); + let mut shared_assignment = + assignment(aggregation_id, QueryMethod::Merge { num_windows: 3 }); + shared_assignment + .item + .query_strings + .push("sum by (job) (metric)".into()); + solution.assignments = vec![shared_assignment]; + + let (streaming, _) = translate(&solution); + + assert_eq!(streaming[aggregation_id].read_count_threshold, Some(6)); + } +}