diff --git a/asap-planner-rs/src/config/input.rs b/asap-planner-rs/src/config/input.rs index 413e9f5c..66626c18 100644 --- a/asap-planner-rs/src/config/input.rs +++ b/asap-planner-rs/src/config/input.rs @@ -10,6 +10,7 @@ use tracing::warn; #[serde(deny_unknown_fields)] pub struct ControllerConfig { pub query_groups: Vec, + pub windowing: Option, pub sketch_parameters: Option, pub aggregate_cleanup: Option, /// Optional hint: per-metric label sets used as a fallback when Prometheus @@ -93,6 +94,50 @@ pub struct AggregateCleanupConfig { pub policy: Option, } +#[derive(Debug, Clone, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct WindowingConfig { + #[serde(rename = "type")] + pub window_type: WindowingType, + pub window_size_ms: u64, + pub slide_interval_ms: Option, +} + +impl WindowingConfig { + pub fn validate(&self) -> Result<(), String> { + if self.window_size_ms == 0 { + return Err("windowing.window_size_ms must be greater than 0".to_string()); + } + match self.window_type { + WindowingType::Tumbling if self.slide_interval_ms.is_some() => { + Err("windowing.slide_interval_ms is only valid for sliding windows".to_string()) + } + WindowingType::Sliding => match self.slide_interval_ms { + None => Err("windowing.slide_interval_ms is required for sliding windows".to_string()), + Some(0) => { + Err("windowing.slide_interval_ms must be greater than 0".to_string()) + } + Some(slide) if slide > self.window_size_ms => Err( + "windowing.slide_interval_ms must be <= windowing.window_size_ms".to_string(), + ), + Some(slide) if !self.window_size_ms.is_multiple_of(slide) => Err(format!( + "windowing.window_size_ms ({}) must be evenly divisible by windowing.slide_interval_ms ({slide})", + self.window_size_ms + )), + Some(_) => Ok(()), + }, + WindowingType::Tumbling => Ok(()), + } + } +} + +#[derive(Debug, Clone, Copy, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum WindowingType { + Tumbling, + Sliding, +} + #[derive(Debug, Clone, Deserialize, Default)] pub struct SketchParameterOverrides { #[serde(rename = "CountMinSketch")] @@ -142,6 +187,7 @@ pub struct HllParams { pub struct SQLControllerConfig { pub query_groups: Vec, pub tables: Vec, + pub windowing: Option, pub sketch_parameters: Option, pub aggregate_cleanup: Option, } diff --git a/asap-planner-rs/src/error.rs b/asap-planner-rs/src/error.rs index f257f41d..3fc5a4cf 100644 --- a/asap-planner-rs/src/error.rs +++ b/asap-planner-rs/src/error.rs @@ -1,5 +1,7 @@ use thiserror::Error; +use crate::planner::window::WindowingError; + #[derive(Debug, Error)] pub enum ControllerError { #[error("IO error: {0}")] @@ -12,6 +14,8 @@ pub enum ControllerError { DuplicateQuery(String), #[error("Planner error: {0}")] PlannerError(String), + #[error("Windowing error: {0}")] + Windowing(#[from] WindowingError), #[error("Unknown metric: {0}")] UnknownMetric(String), #[error("SQL parse error: {0}")] diff --git a/asap-planner-rs/src/lib.rs b/asap-planner-rs/src/lib.rs index 9f83c515..a59bda27 100644 --- a/asap-planner-rs/src/lib.rs +++ b/asap-planner-rs/src/lib.rs @@ -15,6 +15,7 @@ pub use asap_types::PromQLSchema; pub use config::input::ControllerConfig; pub use config::input::ElasticDSLControllerConfig; pub use config::input::SQLControllerConfig; +pub use config::input::{WindowingConfig, WindowingType}; pub use elastic_dsl::ElasticController; pub use elastic_dsl::ElasticIndexSchemaBuilder; pub use elastic_dsl::ElasticRuntimeOptions; diff --git a/asap-planner-rs/src/optimizer/pipeline.rs b/asap-planner-rs/src/optimizer/pipeline.rs index 627ea9aa..55f4d1b0 100644 --- a/asap-planner-rs/src/optimizer/pipeline.rs +++ b/asap-planner-rs/src/optimizer/pipeline.rs @@ -123,6 +123,7 @@ mod tests { ControllerConfig { query_groups, + windowing: None, sketch_parameters: None, aggregate_cleanup: None, metrics: None, diff --git a/asap-planner-rs/src/planner/promql.rs b/asap-planner-rs/src/planner/promql.rs index 7257e558..a4b2a7e6 100644 --- a/asap-planner-rs/src/planner/promql.rs +++ b/asap-planner-rs/src/planner/promql.rs @@ -7,7 +7,7 @@ use promql_utilities::query_logics::enums::{ }; use promql_utilities::query_logics::parsing::get_metric_and_spatial_filter; -use crate::config::input::SketchParameterOverrides; +use crate::config::input::{SketchParameterOverrides, WindowingConfig}; use crate::error::ControllerError; use crate::planner::agg_config::{build_agg_configs_for_statistics, IntermediateAggConfig}; use crate::planner::cleanup::get_cleanup_param; @@ -77,6 +77,7 @@ pub struct SingleQueryProcessor { range_duration_ms: u64, step_ms: u64, cleanup_policy: CleanupPolicy, + windowing: Option, } impl SingleQueryProcessor { @@ -91,6 +92,7 @@ impl SingleQueryProcessor { range_duration_ms: u64, step_ms: u64, cleanup_policy: CleanupPolicy, + windowing: Option, ) -> Self { Self { query, @@ -102,6 +104,7 @@ impl SingleQueryProcessor { range_duration_ms, step_ms, cleanup_policy, + windowing, } } @@ -154,6 +157,7 @@ impl SingleQueryProcessor { self.range_duration_ms, self.step_ms, self.cleanup_policy, + self.windowing.clone(), ) } @@ -252,8 +256,15 @@ impl SingleQueryProcessor { self.data_ingestion_interval_ms, self.step_ms, &mut window_cfg, + self.windowing.is_none(), ) .map_err(ControllerError::PlannerError)?; + crate::planner::window::apply_windowing_override( + &mut window_cfg, + requirements.data_range_ms, + self.step_ms, + self.windowing.as_ref(), + )?; let subpopulation_labels = requirements.grouping_labels; let rollup = all_labels.difference(&subpopulation_labels); diff --git a/asap-planner-rs/src/planner/sql.rs b/asap-planner-rs/src/planner/sql.rs index f7d13f2a..64e1efa1 100644 --- a/asap-planner-rs/src/planner/sql.rs +++ b/asap-planner-rs/src/planner/sql.rs @@ -10,7 +10,7 @@ use sql_utilities::ast_matching::SQLSchema; use sqlparser::dialect::ClickHouseDialect; use sqlparser::parser::Parser as SqlParser; -use crate::config::input::{SketchParameterOverrides, TableDefinition}; +use crate::config::input::{SketchParameterOverrides, TableDefinition, WindowingConfig}; use crate::error::ControllerError; use crate::planner::agg_config::{build_agg_configs_for_statistics, IntermediateAggConfig}; use crate::planner::cleanup::get_sql_cleanup_param; @@ -27,6 +27,7 @@ pub struct SQLSingleQueryProcessor { streaming_engine: StreamingEngine, sketch_parameters: Option, cleanup_policy: CleanupPolicy, + windowing: Option, } impl SQLSingleQueryProcessor { @@ -39,6 +40,7 @@ impl SQLSingleQueryProcessor { streaming_engine: StreamingEngine, sketch_parameters: Option, cleanup_policy: CleanupPolicy, + windowing: Option, ) -> Self { Self { query_string, @@ -48,6 +50,7 @@ impl SQLSingleQueryProcessor { streaming_engine, sketch_parameters, cleanup_policy, + windowing, } } @@ -100,11 +103,19 @@ impl SQLSingleQueryProcessor { let value_column = agg_info.get_value_column_name().to_string(); // Compute window - let window_cfg = compute_sql_window( + let mut window_cfg = compute_sql_window( &sql_query.query_data[0].time_info, self.data_ingestion_interval_ms, self.t_repeat_ms, )?; + let data_range_ms = + (sql_query.query_data[0].time_info.get_duration() * 1000.0).round() as u64; + crate::planner::window::apply_windowing_override( + &mut window_cfg, + data_range_ms, + 0, + self.windowing.as_ref(), + )?; // Get all metadata columns for the table let all_metadata = get_all_metadata_columns(&self.table_definitions, table_name)?; diff --git a/asap-planner-rs/src/planner/window.rs b/asap-planner-rs/src/planner/window.rs index 7cbc9be9..5c4c6ce6 100644 --- a/asap-planner-rs/src/planner/window.rs +++ b/asap-planner-rs/src/planner/window.rs @@ -1,4 +1,7 @@ use asap_types::enums::WindowType; +use std::fmt; + +use crate::config::input::{WindowingConfig, WindowingType}; pub fn get_effective_repeat(t_repeat_ms: u64, step_ms: u64) -> u64 { if step_ms > 0 { @@ -44,6 +47,7 @@ pub fn set_window_parameters( data_ingestion_interval_ms: u64, step_ms: u64, config: &mut IntermediateWindowConfig, + validate_step_alignment: bool, ) -> Result<(), String> { if t_repeat_ms < data_ingestion_interval_ms { return Err(format!( @@ -76,7 +80,7 @@ pub fn set_window_parameters( // Catch a window size incompatible with its own planning-time step_ms // here, at planning time, instead of provisioning a window that would // reject every range query using exactly the step_ms it was planned for. - if step_ms > 0 && !step_ms.is_multiple_of(window_size_ms) { + if validate_step_alignment && step_ms > 0 && !step_ms.is_multiple_of(window_size_ms) { return Err(format!( "step_ms ({step_ms}ms) must be a multiple of the computed window size ({window_size_ms}ms)" )); @@ -88,6 +92,114 @@ pub fn set_window_parameters( Ok(()) } +pub fn apply_windowing_override( + config: &mut IntermediateWindowConfig, + data_range_ms: u64, + step_ms: u64, + windowing: Option<&WindowingConfig>, +) -> Result<(), WindowingError> { + let Some(windowing) = windowing else { + return Ok(()); + }; + windowing + .validate() + .map_err(WindowingError::InvalidConfig)?; + + if !data_range_ms.is_multiple_of(windowing.window_size_ms) { + return Err(WindowingError::DataRangeNotDivisible { + data_range_ms, + window_size_ms: windowing.window_size_ms, + }); + } + + match windowing.window_type { + WindowingType::Tumbling => { + config.window_type = WindowType::Tumbling; + config.window_size_ms = windowing.window_size_ms; + config.slide_interval_ms = config.window_size_ms; + } + WindowingType::Sliding => { + let slide_interval_ms = match windowing.slide_interval_ms { + Some(slide_interval_ms) => slide_interval_ms, + None => { + return Err(WindowingError::InvalidConfig( + "windowing.slide_interval_ms is required for sliding windows".to_string(), + )); + } + }; + config.window_size_ms = windowing.window_size_ms; + if !config.window_size_ms.is_multiple_of(slide_interval_ms) { + return Err(WindowingError::WindowSizeNotDivisible { + window_size_ms: config.window_size_ms, + slide_interval_ms, + }); + } + config.window_type = WindowType::Sliding; + config.slide_interval_ms = slide_interval_ms; + } + } + + let grid_interval_ms = match config.window_type { + WindowType::Tumbling => config.window_size_ms, + WindowType::Sliding => config.slide_interval_ms, + }; + if step_ms > 0 && !step_ms.is_multiple_of(grid_interval_ms) { + return Err(WindowingError::StepNotDivisible { + step_ms, + grid_interval_ms, + }); + } + Ok(()) +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum WindowingError { + InvalidConfig(String), + WindowSizeNotDivisible { + window_size_ms: u64, + slide_interval_ms: u64, + }, + DataRangeNotDivisible { + data_range_ms: u64, + window_size_ms: u64, + }, + StepNotDivisible { + step_ms: u64, + grid_interval_ms: u64, + }, +} + +impl fmt::Display for WindowingError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidConfig(message) => f.write_str(message), + Self::WindowSizeNotDivisible { + window_size_ms, + slide_interval_ms, + } => write!( + f, + "windowing.window_size_ms ({window_size_ms}) must be evenly divisible by windowing.slide_interval_ms ({slide_interval_ms})" + ), + Self::DataRangeNotDivisible { + data_range_ms, + window_size_ms, + } => write!( + f, + "data_range_ms ({data_range_ms}) must be evenly divisible by window_size_ms ({window_size_ms})" + ), + Self::StepNotDivisible { + step_ms, + grid_interval_ms, + } => write!( + f, + "step_ms ({step_ms}) must be evenly divisible by final window grid interval ({grid_interval_ms})" + ), + } + } +} + +impl std::error::Error for WindowingError {} + /// A mutable window config holder used during planning #[derive(Debug, Clone, Default)] pub struct IntermediateWindowConfig { @@ -118,7 +230,7 @@ mod tests { #[test] fn set_window_parameters_temporal_shape() { let mut config = IntermediateWindowConfig::default(); - set_window_parameters(300_000, 60_000, 15_000, 0, &mut config).unwrap(); + set_window_parameters(300_000, 60_000, 15_000, 0, &mut config, true).unwrap(); assert_eq!(config.window_size_ms, 60_000); assert_eq!(config.slide_interval_ms, 60_000); assert_eq!(config.window_type, WindowType::Tumbling); @@ -129,7 +241,7 @@ mod tests { // data_range_ms == data_ingestion_interval_ms (spatial-only query), // t_repeat_ms also equal to the interval: unaffected by the relaxation. let mut config = IntermediateWindowConfig::default(); - set_window_parameters(15_000, 15_000, 15_000, 0, &mut config).unwrap(); + set_window_parameters(15_000, 15_000, 15_000, 0, &mut config, true).unwrap(); assert_eq!(config.window_size_ms, 15_000); } @@ -141,26 +253,26 @@ mod tests { // any cadence gives the latest available answer. window_size stays // exactly one interval regardless of t_repeat_ms. let mut config = IntermediateWindowConfig::default(); - set_window_parameters(15_000, 60_000, 15_000, 0, &mut config).unwrap(); + set_window_parameters(15_000, 60_000, 15_000, 0, &mut config, true).unwrap(); assert_eq!(config.window_size_ms, 15_000); } #[test] fn set_window_parameters_rejects_t_repeat_below_interval() { let mut config = IntermediateWindowConfig::default(); - assert!(set_window_parameters(300_000, 10_000, 15_000, 0, &mut config).is_err()); + assert!(set_window_parameters(300_000, 10_000, 15_000, 0, &mut config, true).is_err()); } #[test] fn set_window_parameters_rejects_data_range_below_t_repeat() { let mut config = IntermediateWindowConfig::default(); - assert!(set_window_parameters(30_000, 60_000, 15_000, 0, &mut config).is_err()); + assert!(set_window_parameters(30_000, 60_000, 15_000, 0, &mut config, true).is_err()); } #[test] fn set_window_parameters_rejects_step_below_interval() { let mut config = IntermediateWindowConfig::default(); - assert!(set_window_parameters(300_000, 60_000, 15_000, 10_000, &mut config).is_err()); + assert!(set_window_parameters(300_000, 60_000, 15_000, 10_000, &mut config, true).is_err()); } #[test] @@ -171,8 +283,32 @@ mod tests { // otherwise provision a window that query-engine's // validate_range_query_params rejects for exactly this step_ms. let mut config = IntermediateWindowConfig::default(); - let result = set_window_parameters(300_000, 40_000, 10_000, 100_000, &mut config); + let result = set_window_parameters(300_000, 40_000, 10_000, 100_000, &mut config, true); assert!(result.is_err()); assert!(result.unwrap_err().contains("must be a multiple of")); } + + #[test] + fn sliding_override_rejects_data_range_not_multiple_of_window_size() { + let mut config = IntermediateWindowConfig { + window_size_ms: 60_000, + slide_interval_ms: 60_000, + window_type: WindowType::Tumbling, + }; + let windowing = WindowingConfig { + window_type: WindowingType::Sliding, + window_size_ms: 60_000, + slide_interval_ms: Some(15_000), + }; + + let error = apply_windowing_override(&mut config, 90_000, 0, Some(&windowing)).unwrap_err(); + + assert_eq!( + error, + WindowingError::DataRangeNotDivisible { + data_range_ms: 90_000, + window_size_ms: 60_000, + } + ); + } } diff --git a/asap-planner-rs/src/promql/controller.rs b/asap-planner-rs/src/promql/controller.rs index 6e85f7f6..be08fd36 100644 --- a/asap-planner-rs/src/promql/controller.rs +++ b/asap-planner-rs/src/promql/controller.rs @@ -37,6 +37,11 @@ impl Controller { ) -> Result { let yaml_str = std::fs::read_to_string(path)?; let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?; + if let Some(windowing) = &config.windowing { + windowing + .validate() + .map_err(ControllerError::PlannerError)?; + } config.warn_default_slas(); let all_queries: Vec = config .query_groups @@ -79,6 +84,11 @@ impl Controller { ) -> Result { let yaml_str = std::fs::read_to_string(path)?; let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?; + if let Some(windowing) = &config.windowing { + windowing + .validate() + .map_err(ControllerError::PlannerError)?; + } config.warn_default_slas(); Ok(Self { config, @@ -94,6 +104,11 @@ impl Controller { opts: RuntimeOptions, ) -> Result { let config: ControllerConfig = serde_yaml::from_str(yaml)?; + if let Some(windowing) = &config.windowing { + windowing + .validate() + .map_err(ControllerError::PlannerError)?; + } config.warn_default_slas(); Ok(Self { config, diff --git a/asap-planner-rs/src/promql/generator.rs b/asap-planner-rs/src/promql/generator.rs index b211f665..297ce683 100644 --- a/asap-planner-rs/src/promql/generator.rs +++ b/asap-planner-rs/src/promql/generator.rs @@ -13,6 +13,7 @@ use crate::generator::{ }; use crate::planner::agg_config::IntermediateAggConfig; use crate::planner::promql::{BinaryArm, SingleQueryProcessor}; +use crate::planner::window::WindowingError; use crate::RuntimeOptions; /// `(query_string, Vec<(identifying_key, cleanup_param)>)` pairs produced by binary leaf decomposition. @@ -31,6 +32,12 @@ pub fn generate_plan( ) -> Result { let metric_schema = schema.clone(); + if let Some(windowing) = &controller_config.windowing { + windowing + .validate() + .map_err(|error| ControllerError::Windowing(WindowingError::InvalidConfig(error)))?; + } + // Determine cleanup policy let cleanup_policy = controller_config .aggregate_cleanup @@ -54,6 +61,7 @@ pub fn generate_plan( let mut query_keys_map: IndexMap)>> = IndexMap::new(); let mut punted_queries: Vec = Vec::new(); + let mut windowing_errors: Vec = Vec::new(); for qg in &controller_config.query_groups { for query_string in &qg.queries { @@ -67,6 +75,7 @@ pub fn generate_plan( qg.range_duration_ms.unwrap_or(opts.range_duration_ms), qg.step_ms.unwrap_or(opts.step_ms), cleanup_policy, + controller_config.windowing.clone(), ); let mut should_process = processor.is_supported(); @@ -97,20 +106,38 @@ pub fn generate_plan( "skipping query referencing unknown metric" ); } + Err(ControllerError::Windowing(error)) => { + windowing_errors.push(format!("query '{query_string}': {error}")); + } Err(e) => return Err(e), } - } else if let Some(arm_entries) = - collect_binary_leaf_entries(&processor, &mut dedup_map)? - { - // Binary arithmetic: register each leaf arm in dedup_map and query_keys_map - for (arm_query, keys_for_arm) in arm_entries { - // Use `entry` so a standalone query that duplicates an arm wins - query_keys_map.entry(arm_query).or_insert(keys_for_arm); + } else { + let mut pending_dedup_map = IndexMap::new(); + if let Some(arm_entries) = collect_binary_leaf_entries( + &processor, + &dedup_map, + &mut pending_dedup_map, + &mut windowing_errors, + query_string, + )? { + dedup_map.extend(pending_dedup_map); + // Binary arithmetic: register each leaf arm in dedup_map and query_keys_map + for (arm_query, keys_for_arm) in arm_entries { + // Use `entry` so a standalone query that duplicates an arm wins + query_keys_map.entry(arm_query).or_insert(keys_for_arm); + } } } } } + if !windowing_errors.is_empty() { + return Err(ControllerError::PlannerError(format!( + "sliding window validation failed:\n{}", + windowing_errors.join("\n") + ))); + } + // Assign sequential IDs (1-indexed, insertion order) let mut id_map: HashMap = HashMap::new(); for (idx, key) in dedup_map.keys().enumerate() { @@ -141,7 +168,10 @@ pub fn generate_plan( /// Returns `Err` only on internal planner errors. fn collect_binary_leaf_entries( processor: &SingleQueryProcessor, - dedup_map: &mut IndexMap, + dedup_map: &IndexMap, + pending_dedup_map: &mut IndexMap, + windowing_errors: &mut Vec, + query_context: &str, ) -> Result, ControllerError> { let arms = match processor.get_binary_arm_queries() { Some(arms) => arms, @@ -149,6 +179,8 @@ fn collect_binary_leaf_entries( }; let mut all_entries: LeafEntries = Vec::new(); + let mut found_windowing_error = false; + let mut found_unsupported_arm = false; for arm in [arms.0, arms.1] { match arm { @@ -161,24 +193,48 @@ fn collect_binary_leaf_entries( if arm_processor.is_supported() { // Leaf arm: gather its streaming aggregation configs. let (configs, cleanup_param) = - arm_processor.get_streaming_aggregation_configs()?; + match arm_processor.get_streaming_aggregation_configs() { + Ok(result) => result, + Err(ControllerError::Windowing(error)) => { + windowing_errors.push(format!( + "query '{query_context}' (leaf '{arm_query}'): {error}" + )); + found_windowing_error = true; + continue; + } + Err(error) => return Err(error), + }; let mut keys_for_arm = Vec::new(); for config in configs { let key = config.identifying_key(); keys_for_arm.push((key.clone(), cleanup_param)); - dedup_map.entry(key).or_insert(config); + if !dedup_map.contains_key(&key) { + pending_dedup_map.entry(key).or_insert(config); + } } all_entries.push((arm_query, keys_for_arm)); } else { // The arm might itself be a binary expression — recurse. - match collect_binary_leaf_entries(&arm_processor, dedup_map)? { + let error_count = windowing_errors.len(); + match collect_binary_leaf_entries( + &arm_processor, + dedup_map, + pending_dedup_map, + windowing_errors, + query_context, + )? { Some(sub_entries) => { all_entries.extend(sub_entries); } None => { + if windowing_errors.len() > error_count { + found_windowing_error = true; + continue; + } // Arm is neither a supported leaf nor a binary expression. // This entire query cannot be accelerated. - return Ok(None); + found_unsupported_arm = true; + continue; } } } @@ -186,7 +242,11 @@ fn collect_binary_leaf_entries( } } - Ok(Some(all_entries)) + if found_windowing_error || found_unsupported_arm { + Ok(None) + } else { + Ok(Some(all_entries)) + } } fn build_streaming_yaml( diff --git a/asap-planner-rs/src/query_log/converter.rs b/asap-planner-rs/src/query_log/converter.rs index 9a9d70e3..756dcd6b 100644 --- a/asap-planner-rs/src/query_log/converter.rs +++ b/asap-planner-rs/src/query_log/converter.rs @@ -37,6 +37,7 @@ pub fn to_controller_config( ControllerConfig { query_groups, + windowing: None, sketch_parameters: None, aggregate_cleanup: Some(AggregateCleanupConfig { policy: Some(CleanupPolicy::ReadBased), diff --git a/asap-planner-rs/src/sql/controller.rs b/asap-planner-rs/src/sql/controller.rs index a1ebeb47..d9da13b9 100644 --- a/asap-planner-rs/src/sql/controller.rs +++ b/asap-planner-rs/src/sql/controller.rs @@ -38,6 +38,11 @@ impl SQLController { ) -> Result { let yaml_str = std::fs::read_to_string(path)?; let mut config: SQLControllerConfig = serde_yaml::from_str(&yaml_str)?; + if let Some(windowing) = &config.windowing { + windowing + .validate() + .map_err(ControllerError::PlannerError)?; + } for table in &mut config.tables { if table.metadata_columns.is_empty() { debug!( @@ -67,6 +72,11 @@ impl SQLController { pub fn from_yaml(yaml: &str, opts: SQLRuntimeOptions) -> Result { let config: SQLControllerConfig = serde_yaml::from_str(yaml)?; + if let Some(windowing) = &config.windowing { + windowing + .validate() + .map_err(ControllerError::PlannerError)?; + } Ok(Self { config, options: opts, diff --git a/asap-planner-rs/src/sql/generator.rs b/asap-planner-rs/src/sql/generator.rs index 560f35f1..17c33aab 100644 --- a/asap-planner-rs/src/sql/generator.rs +++ b/asap-planner-rs/src/sql/generator.rs @@ -13,6 +13,7 @@ use crate::generator::{ }; use crate::planner::agg_config::IntermediateAggConfig; use crate::planner::sql::SQLSingleQueryProcessor; +use crate::planner::window::WindowingError; use crate::StreamingEngine; pub struct SQLRuntimeOptions { @@ -25,6 +26,12 @@ pub fn generate_sql_plan( config: &SQLControllerConfig, opts: &SQLRuntimeOptions, ) -> Result { + if let Some(windowing) = &config.windowing { + windowing + .validate() + .map_err(|error| ControllerError::Windowing(WindowingError::InvalidConfig(error)))?; + } + let eval_time: f64 = opts.query_evaluation_time.unwrap_or_else(|| { SystemTime::now() .duration_since(UNIX_EPOCH) @@ -74,6 +81,7 @@ pub fn generate_sql_plan( let mut dedup_map: IndexMap = IndexMap::new(); // query_string -> Vec<(key, cleanup_param)> let mut query_keys_map: IndexMap)>> = IndexMap::new(); + let mut windowing_errors: Vec = Vec::new(); for qg in &config.query_groups { for query_string in &qg.queries { @@ -85,10 +93,18 @@ pub fn generate_sql_plan( opts.streaming_engine, config.sketch_parameters.clone(), cleanup_policy, + config.windowing.clone(), ); let (configs, cleanup_param) = - processor.get_streaming_aggregation_configs(eval_time)?; + match processor.get_streaming_aggregation_configs(eval_time) { + Ok(result) => result, + Err(ControllerError::Windowing(error)) => { + windowing_errors.push(format!("query '{query_string}': {error}")); + continue; + } + Err(error) => return Err(error), + }; let mut keys_for_query = Vec::new(); for config_item in configs { @@ -100,6 +116,13 @@ pub fn generate_sql_plan( } } + if !windowing_errors.is_empty() { + return Err(ControllerError::PlannerError(format!( + "sliding window validation failed:\n{}", + windowing_errors.join("\n") + ))); + } + // Assign sequential IDs let mut id_map: HashMap = HashMap::new(); for (idx, key) in dedup_map.keys().enumerate() { diff --git a/asap-planner-rs/tests/integration.rs b/asap-planner-rs/tests/integration.rs index 18b0f098..cd22af7f 100644 --- a/asap-planner-rs/tests/integration.rs +++ b/asap-planner-rs/tests/integration.rs @@ -27,6 +27,273 @@ fn http_requests_schema() -> PromQLSchema { ) } +#[test] +fn config_file_sliding_window_override_generates_per_query_candidates() { + let controller = Controller::from_file_with_schema( + Path::new("tests/test_data/windowing/promql_sliding.yaml"), + http_requests_schema(), + arroyo_opts(), + ) + .unwrap(); + + let output = controller.generate().unwrap(); + let streaming: serde_yaml::Value = + serde_yaml::from_str(&output.to_streaming_yaml_string().unwrap()).unwrap(); + let aggregations = streaming["aggregations"].as_sequence().unwrap(); + + assert_eq!(output.inference_query_count(), 2); + assert_eq!(aggregations.len(), 1); + + let mut candidates: Vec<(u64, u64, String)> = aggregations + .iter() + .map(|aggregation| { + ( + aggregation["windowSizeMs"].as_u64().unwrap(), + aggregation["slideIntervalMs"].as_u64().unwrap(), + aggregation["windowType"].as_str().unwrap().to_string(), + ) + }) + .collect(); + candidates.sort(); + + assert_eq!(candidates, vec![(60_000, 15_000, "sliding".to_string())]); + assert!(candidates + .iter() + .all(|(_, _, window_type)| window_type != "tumbling")); +} + +#[test] +fn sliding_window_config_requires_a_slide_interval() { + let result = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 +query_groups: [] +"#, + http_requests_schema(), + arroyo_opts(), + ); + let error = match result { + Ok(_) => panic!("sliding config without a divisor should be rejected"), + Err(error) => error, + }; + + assert!(error + .to_string() + .contains("windowing.slide_interval_ms is required for sliding windows")); +} + +#[test] +fn sliding_window_validation_reports_all_invalid_queries() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +query_groups: + - id: 1 + queries: + - "rate(http_requests_total[65s])" + repetition_delay_ms: 65000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 + - id: 2 + queries: + - "rate(http_requests_total[55s])" + repetition_delay_ms: 55000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + http_requests_schema(), + arroyo_opts(), + ) + .unwrap(); + + let result = controller.generate(); + let error = match result { + Ok(_) => panic!("invalid sliding windows should abort plan generation"), + Err(error) => error, + }; + let message = error.to_string(); + + assert!(message.contains("rate(http_requests_total[65s])")); + assert!(message.contains("rate(http_requests_total[55s])")); +} + +#[test] +fn sliding_window_validation_rejects_data_range_not_multiple_of_window() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +query_groups: + - id: 1 + queries: + - "rate(http_requests_total[90s])" + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + http_requests_schema(), + arroyo_opts(), + ) + .unwrap(); + + let error = match controller.generate() { + Ok(_) => panic!("invalid sliding window should abort plan generation"), + Err(error) => error.to_string(), + }; + + assert!(error.contains("data_range_ms (90000)")); + assert!(error.contains("window_size_ms (60000)")); +} + +#[test] +fn sliding_window_config_rejects_zero_slide_interval() { + let result = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 0 +query_groups: [] +"#, + http_requests_schema(), + arroyo_opts(), + ); + let error = match result { + Ok(_) => panic!("sliding config with divisor 1 should be rejected"), + Err(error) => error, + }; + + assert!(error + .to_string() + .contains("windowing.slide_interval_ms must be greater than 0")); +} + +#[test] +fn tumbling_window_config_rejects_a_slide_interval() { + let result = Controller::from_yaml_with_schema( + r#" +windowing: + type: "tumbling" + window_size_ms: 60000 + slide_interval_ms: 15000 +query_groups: [] +"#, + http_requests_schema(), + arroyo_opts(), + ); + let error = match result { + Ok(_) => panic!("tumbling config with a divisor should be rejected"), + Err(error) => error, + }; + + assert!(error + .to_string() + .contains("windowing.slide_interval_ms is only valid for sliding windows")); +} + +#[test] +fn explicit_tumbling_window_override_keeps_tumbling_candidates() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "tumbling" + window_size_ms: 60000 +query_groups: + - id: 1 + queries: + - "rate(http_requests_total[2m])" + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + http_requests_schema(), + arroyo_opts(), + ) + .unwrap(); + + let output = controller.generate().unwrap(); + let streaming: serde_yaml::Value = + serde_yaml::from_str(&output.to_streaming_yaml_string().unwrap()).unwrap(); + let aggregation = &streaming["aggregations"][0]; + + assert_eq!(aggregation["windowType"].as_str(), Some("tumbling")); + assert_eq!(aggregation["windowSizeMs"].as_u64(), Some(60_000)); + assert_eq!( + aggregation["slideIntervalMs"].as_u64(), + aggregation["windowSizeMs"].as_u64() + ); + assert_ne!(aggregation["windowType"].as_str(), Some("sliding")); +} + +#[test] +fn explicit_sliding_window_revalidates_final_step_grid() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 30000 +query_groups: + - id: 1 + queries: + - "rate(http_requests_total[2m])" + repetition_delay_ms: 40000 + step_ms: 80000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + http_requests_schema(), + arroyo_opts(), + ) + .unwrap(); + + let error = match controller.generate() { + Ok(_) => panic!("step misaligned with the explicit sliding grid should fail"), + Err(error) => error.to_string(), + }; + assert!(error.contains("final window grid interval (30000)")); +} + +#[test] +fn explicit_tumbling_window_rejects_non_multiple_lookback() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "tumbling" + window_size_ms: 60000 +query_groups: + - id: 1 + queries: + - "rate(http_requests_total[90s])" + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + http_requests_schema(), + arroyo_opts(), + ) + .unwrap(); + + let error = match controller.generate() { + Ok(_) => panic!("tumbling lookback must be an exact window multiple"), + Err(error) => error.to_string(), + }; + assert!(error.contains("data_range_ms (90000)")); +} + /// Schema for binary arithmetic tests: errors_total and requests_total. fn binary_arithmetic_schema() -> PromQLSchema { PromQLSchema::new() @@ -573,6 +840,78 @@ fn binary_arithmetic_with_non_acceleratable_arm_produces_no_configs() { assert_eq!(out.inference_query_count(), 0); } +#[test] +fn binary_arithmetic_aggregates_windowing_errors_from_all_leaves() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +query_groups: + - id: 1 + queries: + - "rate(errors_total[65s]) / rate(requests_total[65s])" + repetition_delay_ms: 65000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 + - id: 2 + queries: + - "rate(errors_total[55s]) / rate(requests_total[55s])" + repetition_delay_ms: 55000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + binary_arithmetic_schema(), + arroyo_opts(), + ) + .unwrap(); + + let error = match controller.generate() { + Ok(_) => panic!("invalid binary sliding windows should abort plan generation"), + Err(error) => error.to_string(), + }; + + assert!(error.contains("rate(errors_total[65s])")); + assert!(error.contains("rate(requests_total[65s])")); + assert!(error.contains("rate(errors_total[55s])")); + assert!(error.contains("rate(requests_total[55s])")); + assert_eq!(error.matches("(leaf '").count(), 4); +} + +#[test] +fn binary_arithmetic_scans_past_unsupported_arms_for_windowing_errors() { + let controller = Controller::from_yaml_with_schema( + r#" +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +query_groups: + - id: 1 + queries: + - "abs(errors_total) / rate(requests_total[65s])" + - "rate(requests_total[65s]) / abs(errors_total)" + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +"#, + binary_arithmetic_schema(), + arroyo_opts(), + ) + .unwrap(); + + let error = match controller.generate() { + Ok(_) => panic!("invalid binary sliding windows should abort plan generation"), + Err(error) => error.to_string(), + }; + assert!(error.contains("abs(errors_total) / rate(requests_total[65s])")); + assert!(error.contains("rate(requests_total[65s]) / abs(errors_total)")); +} + #[test] fn temporal_overlapping_rate_increase_deduped() { // rate and increase produce identical MultipleIncrease configs → 1 streaming entry shared, diff --git a/asap-planner-rs/tests/sql_integration.rs b/asap-planner-rs/tests/sql_integration.rs index 28a61072..ae7e109f 100644 --- a/asap-planner-rs/tests/sql_integration.rs +++ b/asap-planner-rs/tests/sql_integration.rs @@ -1,4 +1,6 @@ use asap_planner::{ControllerError, SQLController, SQLRuntimeOptions, StreamingEngine}; +use std::io::Write; +use tempfile::NamedTempFile; // ── helpers ────────────────────────────────────────────────────────────────── @@ -11,6 +13,165 @@ fn sql_opts() -> SQLRuntimeOptions { } } +#[test] +fn config_file_sliding_window_override_generates_per_query_candidates() { + let controller = SQLController::from_file( + std::path::Path::new("tests/test_data/windowing/sql_sliding.yaml"), + sql_opts(), + ) + .unwrap(); + + let output = controller.generate().unwrap(); + let streaming: serde_yaml::Value = + serde_yaml::from_str(&output.to_streaming_yaml_string().unwrap()).unwrap(); + let aggregations = streaming["aggregations"].as_sequence().unwrap(); + + assert_eq!(output.inference_query_count(), 2); + assert_eq!(aggregations.len(), 1); + + let mut candidates: Vec<(u64, u64, String)> = aggregations + .iter() + .map(|aggregation| { + ( + aggregation["windowSizeMs"].as_u64().unwrap(), + aggregation["slideIntervalMs"].as_u64().unwrap(), + aggregation["windowType"].as_str().unwrap().to_string(), + ) + }) + .collect(); + candidates.sort(); + + assert_eq!(candidates, vec![(60_000, 15_000, "sliding".to_string())]); + assert!(candidates + .iter() + .all(|(_, _, window_type)| window_type != "tumbling")); +} + +#[test] +fn sliding_window_validation_reports_all_invalid_sql_queries() { + let controller = SQLController::from_file( + std::path::Path::new("tests/test_data/windowing/sql_sliding_invalid.yaml"), + sql_opts(), + ) + .unwrap(); + + let result = controller.generate(); + let error = match result { + Ok(_) => panic!("invalid sliding windows should abort plan generation"), + Err(error) => error, + }; + let message = error.to_string(); + + assert!(message.contains("DATEADD(s, -75, NOW())")); + assert!(message.contains("DATEADD(s, -90, NOW())")); +} + +#[test] +fn sliding_window_validation_rejects_sql_data_range_not_multiple_of_window() { + let query = "SELECT MIN(cpu_usage) FROM metrics_table WHERE time BETWEEN DATEADD(s, -90, NOW()) AND NOW() GROUP BY hostname"; + let config = format!( + r#" +windowing: + type: sliding + window_size_ms: 60000 + slide_interval_ms: 15000 +tables: + - name: metrics_table + time_column: time + value_columns: [cpu_usage] + metadata_columns: [hostname] +query_groups: + - id: 1 + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.95 + latency_sla: 100.0 + queries: + - "{query}" +aggregate_cleanup: + policy: no_cleanup +"# + ); + + let error = match SQLController::from_yaml(&config, sql_opts()) + .unwrap() + .generate() + { + Ok(_) => panic!("invalid sliding window should abort plan generation"), + Err(error) => error.to_string(), + }; + + assert!(error.contains("data_range_ms (90000)")); + assert!(error.contains("window_size_ms (60000)")); +} + +#[test] +fn explicit_tumbling_window_rejects_sql_non_multiple_lookback() { + let query = "SELECT MIN(cpu_usage) FROM metrics_table WHERE time BETWEEN DATEADD(s, -90, NOW()) AND NOW() GROUP BY hostname"; + let config = format!( + r#" +windowing: + type: tumbling + window_size_ms: 60000 +tables: + - name: metrics_table + time_column: time + value_columns: [cpu_usage] + metadata_columns: [hostname] +query_groups: + - id: 1 + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.95 + latency_sla: 100.0 + queries: + - "{query}" +aggregate_cleanup: + policy: read_based +"# + ); + + let error = match SQLController::from_yaml(&config, sql_opts()) + .unwrap() + .generate() + { + Ok(_) => panic!("tumbling lookback must be an exact window multiple"), + Err(error) => error.to_string(), + }; + assert!(error.contains("data_range_ms (90000)")); +} + +#[test] +fn discovery_validates_windowing_before_contacting_clickhouse() { + let mut file = NamedTempFile::new().unwrap(); + file.write_all( + br#" +windowing: + type: sliding + window_size_ms: 60000 +tables: + - name: metrics_table + time_column: time + value_columns: [cpu_usage] + metadata_columns: [] +query_groups: [] +"#, + ) + .unwrap(); + + let error = match SQLController::from_file_with_discovery( + file.path(), + "http://127.0.0.1:1", + "default", + sql_opts(), + ) { + Ok(_) => panic!("invalid windowing config should be rejected"), + Err(error) => error.to_string(), + }; + + assert!(error.contains("windowing.slide_interval_ms is required for sliding windows")); +} + /// Single-query config with a 3-column metadata schema. /// /// Schema: metrics_table diff --git a/asap-planner-rs/tests/test_data/windowing/promql_sliding.yaml b/asap-planner-rs/tests/test_data/windowing/promql_sliding.yaml new file mode 100644 index 00000000..ad1093cf --- /dev/null +++ b/asap-planner-rs/tests/test_data/windowing/promql_sliding.yaml @@ -0,0 +1,21 @@ +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +query_groups: + - id: 1 + queries: + - "rate(http_requests_total[1m])" + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 + - id: 2 + queries: + - "rate(http_requests_total[2m])" + repetition_delay_ms: 120000 + controller_options: + accuracy_sla: 0.99 + latency_sla: 1.0 +aggregate_cleanup: + policy: "read_based" diff --git a/asap-planner-rs/tests/test_data/windowing/sql_sliding.yaml b/asap-planner-rs/tests/test_data/windowing/sql_sliding.yaml new file mode 100644 index 00000000..acf8a142 --- /dev/null +++ b/asap-planner-rs/tests/test_data/windowing/sql_sliding.yaml @@ -0,0 +1,32 @@ +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +tables: + - name: metrics_table + time_column: time + value_columns: [cpu_usage] + metadata_columns: [hostname, datacenter, region] +query_groups: + - id: 1 + repetition_delay_ms: 60000 + controller_options: + accuracy_sla: 0.95 + latency_sla: 100.0 + queries: + - >- + SELECT MIN(cpu_usage) FROM metrics_table + WHERE time BETWEEN DATEADD(s, -60, NOW()) AND NOW() + GROUP BY datacenter + - id: 2 + repetition_delay_ms: 120000 + controller_options: + accuracy_sla: 0.95 + latency_sla: 100.0 + queries: + - >- + SELECT MIN(cpu_usage) FROM metrics_table + WHERE time BETWEEN DATEADD(s, -120, NOW()) AND NOW() + GROUP BY datacenter +aggregate_cleanup: + policy: "read_based" diff --git a/asap-planner-rs/tests/test_data/windowing/sql_sliding_invalid.yaml b/asap-planner-rs/tests/test_data/windowing/sql_sliding_invalid.yaml new file mode 100644 index 00000000..91af6af8 --- /dev/null +++ b/asap-planner-rs/tests/test_data/windowing/sql_sliding_invalid.yaml @@ -0,0 +1,32 @@ +windowing: + type: "sliding" + window_size_ms: 60000 + slide_interval_ms: 15000 +tables: + - name: metrics_table + time_column: time + value_columns: [cpu_usage] + metadata_columns: [hostname, datacenter, region] +query_groups: + - id: 1 + repetition_delay_ms: 75000 + controller_options: + accuracy_sla: 0.95 + latency_sla: 100.0 + queries: + - >- + SELECT MIN(cpu_usage) FROM metrics_table + WHERE time BETWEEN DATEADD(s, -75, NOW()) AND NOW() + GROUP BY datacenter + - id: 2 + repetition_delay_ms: 90000 + controller_options: + accuracy_sla: 0.95 + latency_sla: 100.0 + queries: + - >- + SELECT MIN(cpu_usage) FROM metrics_table + WHERE time BETWEEN DATEADD(s, -90, NOW()) AND NOW() + GROUP BY datacenter +aggregate_cleanup: + policy: "read_based" diff --git a/asap-query-engine/src/planner_client.rs b/asap-query-engine/src/planner_client.rs index b2b794d6..5d3f3e7c 100644 --- a/asap-query-engine/src/planner_client.rs +++ b/asap-query-engine/src/planner_client.rs @@ -113,6 +113,7 @@ mod tests { step_ms: None, range_duration_ms: None, }], + windowing: None, metrics: Some(vec![MetricDefinition { metric: "http_requests_total".to_string(), labels: vec!["method".to_string(), "status".to_string()], diff --git a/asap-tools/experiments/CONFIG_PARAMETERS_REFERENCE.md b/asap-tools/experiments/CONFIG_PARAMETERS_REFERENCE.md index fba9c6ac..469321b6 100644 --- a/asap-tools/experiments/CONFIG_PARAMETERS_REFERENCE.md +++ b/asap-tools/experiments/CONFIG_PARAMETERS_REFERENCE.md @@ -204,6 +204,29 @@ These parameters must be provided for all experiment scripts: - **Example**: `"10s"` - **Usage**: Controls pre-computed metric updates +## Planner Windowing Override + +#### `windowing` (object, optional) +- **Description**: Global manual override for the planner's window type +- **Default**: Omitted, which preserves tumbling-window planning +- **Choices**: `tumbling`, `sliding` +- **Scope**: Applies to PromQL and SQL planner inputs; ElasticDSL is unchanged + +#### `windowing.type` (string, required when `windowing` is present) +- **Description**: Selects the planner window type +- **Choices**: `tumbling`, `sliding` + +#### `windowing.window_size_ms` (int, required) +- **Description**: Explicit precompute window size in milliseconds +- **Validation**: Must be greater than zero; every supported query's lookback + must be an exact multiple + +#### `windowing.slide_interval_ms` (int, required for sliding) +- **Description**: Explicit sliding interval in milliseconds +- **Validation**: Must be positive, no greater than `window_size_ms`, and evenly + divide it +- **Constraint**: Omit this field for `tumbling` + --- ## Monitoring Configuration diff --git a/asap-tools/experiments/config/config.yaml b/asap-tools/experiments/config/config.yaml index 0c6d5ef6..6288ef73 100644 --- a/asap-tools/experiments/config/config.yaml +++ b/asap-tools/experiments/config/config.yaml @@ -99,6 +99,14 @@ query_engine: controller: punting: true # Enable query punting based on performance heuristics (should_be_performant check) +# Optional manual planner windowing override. Omit this block to preserve the +# default tumbling-window behavior. For tumbling windows, omit +# slide_interval_ms. For sliding windows, both sizes are explicit. +# windowing: +# type: "sliding" +# window_size_ms: 60000 +# slide_interval_ms: 15000 + # Aggregate cleanup configuration # Policy options: # - "circular_buffer": Keep N most recent aggregates (fixed-count cleanup) diff --git a/asap-tools/experiments/experiment_only_ingest_path.py b/asap-tools/experiments/experiment_only_ingest_path.py index a7c16dc7..d85765d8 100644 --- a/asap-tools/experiments/experiment_only_ingest_path.py +++ b/asap-tools/experiments/experiment_only_ingest_path.py @@ -190,6 +190,7 @@ def main(cfg: DictConfig): local_experiment_root_dir, cfg.aggregate_cleanup, cfg.get("sketch_parameters", None), + cfg.get("windowing", None), ) ) sync.rsync_controller_client_configs( diff --git a/asap-tools/experiments/experiment_run_clickhouse.py b/asap-tools/experiments/experiment_run_clickhouse.py index ac218244..d52c7e0c 100644 --- a/asap-tools/experiments/experiment_run_clickhouse.py +++ b/asap-tools/experiments/experiment_run_clickhouse.py @@ -377,7 +377,10 @@ def main(cfg: DictConfig) -> None: # Generate and rsync the planner input config to the node planner_input_yaml = config.generate_sql_planner_input( - ep.query_groups, dataset_cfg, cfg.get("sketch_parameters", None) + ep.query_groups, + dataset_cfg, + cfg.get("sketch_parameters", None), + cfg.get("windowing", None), ) local_planner_input = os.path.join( local_controller_dir, "planner_input.yaml" diff --git a/asap-tools/experiments/experiment_run_e2e.py b/asap-tools/experiments/experiment_run_e2e.py index ce9ccba3..f52c5f08 100644 --- a/asap-tools/experiments/experiment_run_e2e.py +++ b/asap-tools/experiments/experiment_run_e2e.py @@ -200,6 +200,7 @@ def main(cfg: DictConfig): local_experiment_root_dir, cfg.aggregate_cleanup, cfg.get("sketch_parameters", None), + cfg.get("windowing", None), ) ) sync.rsync_controller_client_configs( diff --git a/asap-tools/experiments/experiment_run_grafana_demo.py b/asap-tools/experiments/experiment_run_grafana_demo.py index d7e98b48..0ecd3ddc 100644 --- a/asap-tools/experiments/experiment_run_grafana_demo.py +++ b/asap-tools/experiments/experiment_run_grafana_demo.py @@ -161,6 +161,7 @@ def main(cfg: DictConfig): local_experiment_root_dir, cfg.aggregate_cleanup, cfg.get("sketch_parameters", None), + cfg.get("windowing", None), ) ) sync.rsync_controller_client_configs( diff --git a/asap-tools/experiments/experiment_utils/config.py b/asap-tools/experiments/experiment_utils/config.py index 348c3953..495b663c 100644 --- a/asap-tools/experiments/experiment_utils/config.py +++ b/asap-tools/experiments/experiment_utils/config.py @@ -342,6 +342,7 @@ def generate_controller_client_configs( local_experiment_dir: str, aggregate_cleanup: DictConfig = None, sketch_parameters: DictConfig = None, + windowing: DictConfig = None, ) -> Tuple[List[str], List[str]]: """Generate controller client configurations from experiment parameters.""" # experiment_params is already loaded by Hydra @@ -358,6 +359,10 @@ def generate_controller_client_configs( sketch_params_config = OmegaConf.to_container(sketch_parameters, resolve=True) experiment_config["sketch_parameters"] = sketch_params_config + # Add the optional planner windowing override if provided. + if windowing is not None: + experiment_config["windowing"] = OmegaConf.to_container(windowing, resolve=True) + output_dir = os.path.join(local_experiment_dir, "controller_client_configs") os.makedirs(output_dir, exist_ok=True) @@ -376,6 +381,7 @@ def generate_controller_client_configs( "query_groups", "sketch_parameters", "aggregate_cleanup", + "windowing", "metrics", "existing_streaming_config", } @@ -828,7 +834,10 @@ def generate_clickhouse_client_configs( def generate_sql_planner_input( - query_groups: Any, dataset_cfg: Any, sketch_parameters: Any = None + query_groups: Any, + dataset_cfg: Any, + sketch_parameters: Any = None, + windowing: Any = None, ) -> str: """Generate the YAML input file for asap-planner in SQL mode. @@ -836,6 +845,7 @@ def generate_sql_planner_input( ``SQLControllerConfig`` YAML that contains: - ``tables``: schema of the tables being queried - ``query_groups``: SQL queries with controller options + - ``windowing``: optional global tumbling/sliding window override - ``sketch_parameters``: optional per-sketch-type overrides (e.g. ``DatasketchesKLL.K``), matching ``ControllerConfig``'s PromQL-mode field of the same name (``SketchParameterOverrides`` in @@ -854,6 +864,9 @@ def generate_sql_planner_input( top-level ``sketch_parameters`` section (``CountMinSketch``, ``DatasketchesKLL``, etc.). When ``None``, the planner falls back to its own defaults. + windowing: Optional DictConfig/dict mirroring ``config.yaml``'s + top-level ``windowing`` section. When ``None``, the planner uses + its default tumbling-window behavior. Returns: YAML string ready to write to disk and pass to asap-planner. @@ -913,6 +926,10 @@ def generate_sql_planner_input( if isinstance(sketch_parameters, (DictConfig, ListConfig)): sketch_parameters = OmegaConf.to_container(sketch_parameters, resolve=True) planner_input["sketch_parameters"] = sketch_parameters + if windowing is not None: + if isinstance(windowing, (DictConfig, ListConfig)): + windowing = OmegaConf.to_container(windowing, resolve=True) + planner_input["windowing"] = windowing return yaml.dump(planner_input, default_flow_style=False, allow_unicode=True)