diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 76c3a88c..716c6ee4 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -221,13 +221,14 @@ pub use summary_maintenance_dag_export::{ }; pub use summary_maintenance_lifecycle::{ assemble_selected_dag_with_summary_maintenance_lifecycles, - global_selection_with_summary_maintenance_lifecycles, plan_summary_maintenance_lifecycles, - SummaryMaintenanceCapabilities, SummaryMaintenanceDeployment, - SummaryMaintenanceLifecycleAlternative, SummaryMaintenanceLifecycleAssemblyError, - SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs, - SummaryMaintenanceLifecyclePlan, SummaryMaintenanceLifecyclePlanError, - SummaryMaintenanceLifecycleRejection, SummaryMaintenanceLifecycleSelectionError, - WorkloadDemand, + enumerate_summary_maintenance_lifecycles, global_selection_with_summary_maintenance_lifecycles, + plan_summary_maintenance_lifecycles, SummaryMaintenanceCapabilities, + SummaryMaintenanceDeployment, SummaryMaintenanceLifecycleAlternative, + SummaryMaintenanceLifecycleAssemblyError, SummaryMaintenanceLifecycleCandidates, + SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleChoiceError, + SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecyclePlan, + SummaryMaintenanceLifecyclePlanError, SummaryMaintenanceLifecycleRejection, + SummaryMaintenanceLifecycleSelectionError, WorkloadDemand, }; pub use topk_reuse::TopKLimitReuseStrategy; diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index b0d9d79f..7e40645f 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -287,6 +287,169 @@ pub enum SummaryMaintenanceLifecycleSelectionError { SummaryMaintenance(#[from] SummaryMaintenanceLifecyclePlanError), } +/// Every lifecycle alternative for each unique summary state of one fixed +/// root, before any lifecycle is chosen. +/// +/// Planner selection ([`plan_summary_maintenance_lifecycles`]) and a +/// deployment's explicit choice ([`Self::select`]) both finish from this value, +/// so they produce the same [`SummaryMaintenanceLifecyclePlan`] shape. +pub struct SummaryMaintenanceLifecycleCandidates<'a> { + /// Unselected plan: deployments carry alternatives but no guarantee or + /// window framework. + plan: SummaryMaintenanceLifecyclePlan, + components: Vec, + arrival: DataArrival, + required_accuracy: Vec, + cost_model: &'a dyn CostModel, + comparison_target: Option<&'a QueryExpr>, +} + +/// Why an explicit per-state lifecycle choice cannot be bound. +#[derive(Debug, thiserror::Error, PartialEq)] +pub enum SummaryMaintenanceLifecycleChoiceError { + #[error("summary {0:?} is not a deployment of this root")] + UnknownSummary(PostAsapNodeId), + #[error("summary {0:?} is chosen more than once")] + DuplicateChoice(PostAsapNodeId), + #[error("summary {0:?} has no chosen lifecycle")] + MissingChoice(PostAsapNodeId), + #[error("chosen lifecycle is not an enumerated alternative of summary {0:?}")] + NotAnAlternative(PostAsapNodeId), + #[error("chosen lifecycle of summary {post_asap_node_id:?} is rejected: {rejection:?}")] + Rejected { + post_asap_node_id: PostAsapNodeId, + rejection: Option, + }, + #[error("summary states on one maintenance path have different evaluation schedules")] + IncompatibleEvaluationSchedules, + #[error("the cost model supplied no complete estimate for the chosen combination")] + NoCompleteEstimate, +} + +impl SummaryMaintenanceLifecycleCandidates<'_> { + /// One entry per unique reachable `SummaryAgg`, with every alternative + /// and its rejection; no lifecycle or window framework is selected. + pub fn deployments(&self) -> &[SummaryMaintenanceDeployment] { + &self.plan.deployments + } + + /// Guarantee that binding `lifecycle` would attach under this workload's + /// data arrival, so a caller can price an alternative before choosing it. + pub fn guarantee( + &self, + lifecycle: &SummaryMaintenanceLifecycle, + ) -> SummaryMaintenanceLifecycleGuarantee { + lifecycle_guarantee(lifecycle, self.arrival) + } + + fn context(&self) -> CompleteCostContext<'_> { + CompleteCostContext { + root: &self.plan.root, + components: &self.components, + cost_model: self.cost_model, + comparison_target: self.comparison_target, + horizon: self.plan.horizon, + expected_reads: self.plan.expected_reads, + required_accuracy: &self.required_accuracy, + } + } + + fn finish( + mut self, + estimate: Option, + ) -> SummaryMaintenanceLifecyclePlan { + if let Some(estimate) = estimate { + self.plan.summary_total_cost = Some(estimate.cost); + self.plan.selected_window_implementation_id = estimate.physical_plan_id; + self.plan.window_accuracy_guarantee = estimate.window_accuracy_guarantee; + } + self.plan + } + + /// Planner's choice: the cheapest complete combination of eligible + /// alternatives. + fn select_cheapest(mut self) -> SummaryMaintenanceLifecyclePlan { + let estimate = select_complete_lifecycle_combination( + &self.plan.root, + &mut self.plan.deployments, + &self.components, + self.arrival, + self.cost_model, + self.comparison_target, + self.plan.horizon, + self.plan.expected_reads, + &self.required_accuracy, + ); + self.finish(estimate) + } + + /// Bind one caller-chosen lifecycle per summary state. Each choice must be + /// an alternative Planner itself could select; the complete estimate is + /// then obtained exactly as for Planner selection, so window framework and + /// cost are the model's and unknown cost is never replaced by zero. + pub fn select( + mut self, + choices: &[(PostAsapNodeId, SummaryMaintenanceLifecycle)], + ) -> Result { + use SummaryMaintenanceLifecycleChoiceError as E; + let deployments = &self.plan.deployments; + let mut chosen: Vec> = + vec![None; deployments.len()]; + let context = self.context(); + for (id, lifecycle) in choices { + let index = deployments + .iter() + .position(|deployment| deployment.post_asap_node_id == *id) + .ok_or(E::UnknownSummary(*id))?; + if chosen[index].is_some() { + return Err(E::DuplicateChoice(*id)); + } + let alternative = deployments[index] + .alternatives + .iter() + .find(|alternative| alternative.summary_maintenance_lifecycle == *lifecycle) + .ok_or(E::NotAnAlternative(*id))?; + if !context.eligible(alternative) { + return Err(E::Rejected { + post_asap_node_id: *id, + rejection: alternative.rejection.clone(), + }); + } + chosen[index] = Some(alternative); + } + let selected = chosen + .into_iter() + .enumerate() + .map(|(index, alternative)| { + let alternative = + alternative.ok_or(E::MissingChoice(deployments[index].post_asap_node_id))?; + Ok(( + index, + lifecycle_guarantee(&alternative.summary_maintenance_lifecycle, self.arrival), + // Reached only for costed alternatives or when the + // complete hook is authoritative, matching Planner search. + alternative.total_cost.unwrap_or(Cost::ZERO), + )) + }) + .collect::, E>>()?; + if selected.is_empty() { + return Ok(self.finish(None)); + } + if !context.schedules_compatible(&selected) { + return Err(E::IncompatibleEvaluationSchedules); + } + let estimate = context + .estimate(deployments, &selected) + .ok_or(E::NoCompleteEstimate)?; + let guarantees = selected + .into_iter() + .map(|(index, guarantee, _)| (index, guarantee)) + .collect(); + apply_selection(&mut self.plan.deployments, guarantees, &estimate); + Ok(self.finish(Some(estimate))) + } +} + /// Workload-wide evidence derived specifically for summary-maintenance /// lifecycle enumeration and costing. /// @@ -333,7 +496,30 @@ pub fn plan_summary_maintenance_lifecycles( capabilities: SummaryMaintenanceLifecycleCapabilities, cost_model: &dyn CostModel, ) -> Result { - plan_summary_maintenance_lifecycles_with_profile( + Ok(enumerate_summary_maintenance_lifecycles( + root, + demand, + now_ms, + horizon, + capabilities, + cost_model, + )? + .select_cheapest()) +} + +/// Validate a materialized plan and enumerate lifecycle alternatives for each +/// unique summary state without choosing one. A deployment that prices the +/// alternatives itself binds its choice with +/// [`SummaryMaintenanceLifecycleCandidates::select`]. +pub fn enumerate_summary_maintenance_lifecycles<'a>( + root: Rc, + demand: WorkloadDemand<'_>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &'a dyn CostModel, +) -> Result, SummaryMaintenanceLifecyclePlanError> { + enumerate_with_profile( root, demand, now_ms, @@ -349,16 +535,16 @@ pub fn plan_summary_maintenance_lifecycles( /// eligibility and data-arrival facts; `profile` supplies effective uses after /// DAG path multiplicity has been propagated by `PlanSpace`. #[expect(clippy::too_many_arguments, reason = "internal bound planning context")] -fn plan_summary_maintenance_lifecycles_with_profile( +fn enumerate_with_profile<'a>( root: Rc, demand: WorkloadDemand<'_>, now_ms: u64, horizon: Option, capabilities: SummaryMaintenanceLifecycleCapabilities, - cost_model: &dyn CostModel, + cost_model: &'a dyn CostModel, profile: Option, - comparison_target: Option<&QueryExpr>, -) -> Result { + comparison_target: Option<&'a QueryExpr>, +) -> Result, SummaryMaintenanceLifecyclePlanError> { demand.workload.validate()?; if let Some(data) = demand.data_workload { data.validate()?; @@ -392,7 +578,7 @@ fn plan_summary_maintenance_lifecycles_with_profile( collect_summary_aggs(&root, &mut HashSet::new(), &mut summaries); let node_ids = compile_post_asap_dag_with_node_ids(&root)?.node_ids; let components = summary_state_components(&summaries); - let mut deployments: Vec = summaries + let deployments: Vec = summaries .into_iter() .map(|summary| { let alternatives = alternatives_for( @@ -413,37 +599,26 @@ fn plan_summary_maintenance_lifecycles_with_profile( } }) .collect(); - let complete_estimate = select_complete_lifecycle_combination( - &root, - &mut deployments, - &components, - facts.arrival, + let selected_raw_recompute = matches!(root.expr, SummaryExpr::KeepPreAsap(_)); + Ok(SummaryMaintenanceLifecycleCandidates { + plan: SummaryMaintenanceLifecyclePlan { + root, + deployments, + horizon, + evaluation_rate: facts.evaluation_rate, + update_rate: facts.update_rate, + expected_reads: facts.reads, + selected_raw_recompute, + selected_window_implementation_id: None, + summary_total_cost: None, + window_accuracy_guarantee: None, + raw_recompute_total_cost: None, + }, + components, + arrival: facts.arrival, + required_accuracy: facts.required_accuracy, cost_model, comparison_target, - horizon, - facts.reads, - &facts.required_accuracy, - ); - let summary_total_cost = complete_estimate.as_ref().map(|estimate| estimate.cost); - let selected_window_implementation_id = complete_estimate - .as_ref() - .and_then(|estimate| estimate.physical_plan_id.clone()); - let window_accuracy_guarantee = complete_estimate - .as_ref() - .and_then(|estimate| estimate.window_accuracy_guarantee.clone()); - let selected_raw_recompute = matches!(root.expr, SummaryExpr::KeepPreAsap(_)); - Ok(SummaryMaintenanceLifecyclePlan { - root, - deployments, - horizon, - evaluation_rate: facts.evaluation_rate, - update_rate: facts.update_rate, - expected_reads: facts.reads, - selected_raw_recompute, - selected_window_implementation_id, - summary_total_cost, - window_accuracy_guarantee, - raw_recompute_total_cost: None, }) } @@ -483,7 +658,7 @@ pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( continue; }; costs.finalize_target(&group.target); - let plan = plan_summary_maintenance_lifecycles_with_profile( + let plan = enumerate_with_profile( Rc::clone(summary), WorkloadDemand { workload, @@ -496,7 +671,8 @@ pub fn global_selection_with_summary_maintenance_lifecycles<'a, Id>( cost_model, Some(profiles.for_target(&group.target)), Some(&group.target), - )?; + )? + .select_cheapest(); let raw = plan .expected_reads .and_then(|reads| cost_model.raw_query_recompute_total_cost(&group.target, reads)); @@ -529,7 +705,7 @@ pub fn assemble_selected_dag_with_summary_maintenance_lifecycles( selection .assemble_selected_dag(target)? .map(|root| { - let mut plan = plan_summary_maintenance_lifecycles_with_profile( + let mut plan = enumerate_with_profile( root, demand, now_ms, @@ -538,7 +714,8 @@ pub fn assemble_selected_dag_with_summary_maintenance_lifecycles( cost_model, None, Some(target), - )?; + )? + .select_cheapest(); plan.raw_recompute_total_cost = plan .expected_reads .and_then(|reads| cost_model.raw_query_recompute_total_cost(target, reads)); @@ -1104,6 +1281,99 @@ fn summary_state_components(summaries: &[Rc]) -> Vec { .collect() } +/// Inputs shared by every complete lifecycle-combination evaluation of one +/// root, whether Planner searches combinations or a caller supplies one. +struct CompleteCostContext<'a> { + root: &'a SummaryNode, + components: &'a [usize], + cost_model: &'a dyn CostModel, + comparison_target: Option<&'a QueryExpr>, + horizon: Option, + expected_reads: Option, + required_accuracy: &'a [AccuracyTarget], +} + +impl CompleteCostContext<'_> { + /// Planner's own admission rule for one alternative. Uncosted alternatives + /// are admitted only when the complete-candidate hook is authoritative. + fn eligible(&self, alternative: &SummaryMaintenanceLifecycleAlternative) -> bool { + alternative.selectable() + || (self + .cost_model + .complete_summary_candidate_estimate_covers_lifecycle_costs() + && alternative.rejection + == Some(SummaryMaintenanceLifecycleRejection::MissingCostEvidence)) + } + + /// `selected` holds one entry per deployment, in deployment order. + fn schedules_compatible( + &self, + selected: &[(usize, SummaryMaintenanceLifecycleGuarantee, Cost)], + ) -> bool { + !selected.iter().enumerate().any(|(left, (_, a, _))| { + selected.iter().enumerate().any(|(right, (_, b, _))| { + self.components[left] == self.components[right] + && a.evaluation_schedule != b.evaluation_schedule + }) + }) + } + + fn estimate( + &self, + deployments: &[SummaryMaintenanceDeployment], + selected: &[(usize, SummaryMaintenanceLifecycleGuarantee, Cost)], + ) -> Option { + if !self.schedules_compatible(selected) { + return None; + } + let costed: Vec<_> = selected + .iter() + .map(|(index, guarantee, cost)| CostedSummaryDeployment { + summary: &deployments[*index].summary, + guarantee, + selected_cost: *cost, + }) + .collect(); + let estimate = self.cost_model.complete_summary_candidate_estimate( + self.root, + self.comparison_target, + &costed, + self.horizon, + self.expected_reads, + self.required_accuracy, + )?; + (estimate.window_frameworks.len() == deployments.len()).then_some(estimate) + } +} + +fn lifecycle_guarantee( + lifecycle: &SummaryMaintenanceLifecycle, + arrival: DataArrival, +) -> SummaryMaintenanceLifecycleGuarantee { + SummaryMaintenanceLifecycleGuarantee { + summary_maintenance_mode: maintenance_mode(lifecycle, arrival), + evaluation_schedule: evaluation_schedule(lifecycle, arrival), + summary_maintenance_lifecycle: lifecycle.clone(), + output_representation: OutputRepresentation::SummaryState, + } +} + +fn apply_selection( + deployments: &mut [SummaryMaintenanceDeployment], + guarantees: Vec<(usize, SummaryMaintenanceLifecycleGuarantee)>, + estimate: &CompleteSummaryCandidateEstimate, +) { + for (index, guarantee) in guarantees { + deployments[index].summary_maintenance_lifecycle_guarantee = Some(guarantee); + } + for (deployment, framework) in deployments + .iter_mut() + .zip(estimate.window_frameworks.iter().cloned()) + { + deployment.selected_window_framework = framework; + } +} + #[expect(clippy::too_many_arguments, reason = "complete combination context")] fn select_complete_lifecycle_combination( root: &SummaryNode, @@ -1120,82 +1390,48 @@ fn select_complete_lifecycle_combination( if deployments.is_empty() { return None; } + let context = CompleteCostContext { + root, + components, + cost_model, + comparison_target, + horizon, + expected_reads, + required_accuracy, + }; // The whole-candidate hook is intentionally arbitrary, so partial costs // cannot soundly prune the search. Bound exhaustive enumeration and fail // closed instead of allowing an adversarial DAG to consume exponential // planner time. - let complete_costing = cost_model.complete_summary_candidate_estimate_covers_lifecycle_costs(); - let eligible = |alternative: &SummaryMaintenanceLifecycleAlternative| { - alternative.selectable() - || (complete_costing - && alternative.rejection - == Some(SummaryMaintenanceLifecycleRejection::MissingCostEvidence)) - }; let combinations = deployments .iter() .try_fold(1_usize, |product, deployment| { let selectable = deployment .alternatives .iter() - .filter(|alternative| eligible(alternative)) + .filter(|alternative| context.eligible(alternative)) .count(); product.checked_mul(selectable) })?; if combinations == 0 || combinations > MAX_COMPLETE_LIFECYCLE_COMBINATIONS { return None; } - #[expect( - clippy::too_many_arguments, - reason = "recursive lifecycle-combination search state" - )] + type Best = Option<( + CompleteSummaryCandidateEstimate, + Vec<(usize, SummaryMaintenanceLifecycleGuarantee)>, + )>; fn visit( index: usize, - root: &SummaryNode, + context: &CompleteCostContext<'_>, deployments: &[SummaryMaintenanceDeployment], - components: &[usize], arrival: DataArrival, - cost_model: &dyn CostModel, - comparison_target: Option<&QueryExpr>, - horizon: Option, - expected_reads: Option, - required_accuracy: &[AccuracyTarget], - complete_costing: bool, selected: &mut Vec<(usize, SummaryMaintenanceLifecycleGuarantee, Cost)>, - best: &mut Option<( - CompleteSummaryCandidateEstimate, - Vec<(usize, SummaryMaintenanceLifecycleGuarantee)>, - )>, + best: &mut Best, ) { if index == deployments.len() { - if selected.iter().enumerate().any(|(left, (_, a, _))| { - selected.iter().enumerate().any(|(right, (_, b, _))| { - components[left] == components[right] - && a.evaluation_schedule != b.evaluation_schedule - }) - }) { - return; - } - let costed: Vec<_> = selected - .iter() - .map(|(index, guarantee, cost)| CostedSummaryDeployment { - summary: &deployments[*index].summary, - guarantee, - selected_cost: *cost, - }) - .collect(); - let Some(estimate) = cost_model.complete_summary_candidate_estimate( - root, - comparison_target, - &costed, - horizon, - expected_reads, - required_accuracy, - ) else { + let Some(estimate) = context.estimate(deployments, selected) else { return; }; - if estimate.window_frameworks.len() != deployments.len() { - return; - } if best .as_ref() .is_none_or(|(best_estimate, _)| estimate.cost.0 < best_estimate.cost.0) @@ -1213,40 +1449,14 @@ fn select_complete_lifecycle_combination( for alternative in deployments[index] .alternatives .iter() - .filter(|alternative| { - alternative.selectable() - || (complete_costing - && alternative.rejection - == Some(SummaryMaintenanceLifecycleRejection::MissingCostEvidence)) - }) + .filter(|alternative| context.eligible(alternative)) { - let lifecycle = alternative.summary_maintenance_lifecycle.clone(); - let guarantee = SummaryMaintenanceLifecycleGuarantee { - summary_maintenance_mode: maintenance_mode(&lifecycle, arrival), - evaluation_schedule: evaluation_schedule(&lifecycle, arrival), - summary_maintenance_lifecycle: lifecycle, - output_representation: OutputRepresentation::SummaryState, - }; selected.push(( index, - guarantee, + lifecycle_guarantee(&alternative.summary_maintenance_lifecycle, arrival), alternative.total_cost.unwrap_or(Cost::ZERO), )); - visit( - index + 1, - root, - deployments, - components, - arrival, - cost_model, - comparison_target, - horizon, - expected_reads, - required_accuracy, - complete_costing, - selected, - best, - ); + visit(index + 1, context, deployments, arrival, selected, best); selected.pop(); } } @@ -1254,29 +1464,14 @@ fn select_complete_lifecycle_combination( let mut best = None; visit( 0, - root, + &context, deployments, - components, arrival, - cost_model, - comparison_target, - horizon, - expected_reads, - required_accuracy, - complete_costing, &mut Vec::new(), &mut best, ); let (estimate, guarantees) = best?; - for (index, guarantee) in guarantees { - deployments[index].summary_maintenance_lifecycle_guarantee = Some(guarantee); - } - for (deployment, framework) in deployments - .iter_mut() - .zip(estimate.window_frameworks.iter().cloned()) - { - deployment.selected_window_framework = framework; - } + apply_selection(deployments, guarantees, &estimate); Some(estimate) } @@ -2402,4 +2597,252 @@ mod tests { assert_eq!(batch.evaluation_rate, None); assert_eq!(batch.one_shot_consumers, 1); } + + fn continuous_candidates<'a>( + workload: &QueryWorkload, + data: &DataWorkload, + model: &'a dyn CostModel, + ) -> SummaryMaintenanceLifecycleCandidates<'a> { + enumerate_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new_with_data(workload, data, &[0]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities { + supports_shared: false, + ..SummaryMaintenanceLifecycleCapabilities::ALL + }, + model, + ) + .unwrap() + } + + fn choose( + candidates: &SummaryMaintenanceLifecycleCandidates<'_>, + lifecycle: SummaryMaintenanceLifecycle, + ) -> Vec<(PostAsapNodeId, SummaryMaintenanceLifecycle)> { + candidates + .deployments() + .iter() + .map(|deployment| (deployment.post_asap_node_id, lifecycle.clone())) + .collect() + } + + // Enumeration reports all four lifecycle kinds with their rejections and + // selects nothing. + #[test] + fn enumeration_exposes_every_lifecycle_without_selecting() { + let data = continuous(1_000, 60_000); + let workload = workload(vec![], vec![repeating()], data.clone()); + let candidates = continuous_candidates(&workload, &data, &UnitCosts); + let [deployment] = candidates.deployments() else { + panic!("one summary state"); + }; + assert_eq!(deployment.summary_maintenance_lifecycle_guarantee, None); + assert_eq!(deployment.selected_window_framework, None); + let outcome: Vec<_> = deployment + .alternatives + .iter() + .map(|alternative| { + ( + &alternative.summary_maintenance_lifecycle, + alternative.rejection.clone(), + alternative.total_cost.is_some(), + ) + }) + .collect(); + assert!(matches!( + outcome.as_slice(), + [ + (SummaryMaintenanceLifecycle::Ephemeral, None, true), + ( + SummaryMaintenanceLifecycle::Prepared { .. }, + Some(SummaryMaintenanceLifecycleRejection::RequiresPredictableOneTimeQuery), + false + ), + ( + SummaryMaintenanceLifecycle::Shared { .. }, + Some(SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime), + false + ), + ( + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + None, + true + ), + ] + )); + let guarantee = candidates.guarantee(&SummaryMaintenanceLifecycle::ContinuouslyMaintained); + assert_eq!( + guarantee.summary_maintenance_mode, + SummaryMaintenanceMode::Incremental + ); + assert_eq!(guarantee.evaluation_schedule, EvaluationSchedule::PerUpdate); + } + + // Explicitly choosing Planner's own selection reproduces Planner's plan. + #[test] + fn explicit_choice_of_planner_selection_reproduces_planner_plan() { + let data = continuous(1_000, 60_000); + let workload = workload(vec![], vec![repeating()], data.clone()); + let planned = plan_summary_maintenance_lifecycles( + summary(), + WorkloadDemand::new_with_data(&workload, &data, &[0]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities { + supports_shared: false, + ..SummaryMaintenanceLifecycleCapabilities::ALL + }, + &UnitCosts, + ) + .unwrap(); + let candidates = continuous_candidates(&workload, &data, &UnitCosts); + let choice = choose( + &candidates, + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + ); + let chosen = candidates.select(&choice).unwrap(); + assert_eq!(format!("{chosen:?}"), format!("{planned:?}")); + } + + // A deployment may bind a legal alternative Planner's estimate does not + // prefer; the plan carries that alternative's guarantee and cost. + #[test] + fn explicit_choice_may_bind_a_costlier_legal_alternative() { + let data = continuous(1_000, 60_000); + let workload = workload(vec![], vec![repeating()], data.clone()); + let candidates = continuous_candidates(&workload, &data, &UnitCosts); + let ephemeral_cost = candidates.deployments()[0].alternatives[0].total_cost; + let choice = choose(&candidates, SummaryMaintenanceLifecycle::Ephemeral); + let plan = candidates.select(&choice).unwrap(); + assert_eq!( + selected_summary_maintenance_lifecycle(&plan.deployments[0]), + Some(&SummaryMaintenanceLifecycle::Ephemeral) + ); + assert_eq!(plan.summary_total_cost, ephemeral_cost); + } + + // Choices that Planner could not select, or that do not cover exactly the + // enumerated states, are refused rather than bound. + #[test] + fn explicit_choice_rejects_illegal_or_incomplete_choices() { + use SummaryMaintenanceLifecycleChoiceError as E; + let data = continuous(1_000, 60_000); + let workload = workload(vec![], vec![repeating()], data.clone()); + let select = |model: &dyn CostModel, choice: &dyn Fn(PostAsapNodeId) -> Vec<_>| { + let candidates = continuous_candidates(&workload, &data, model); + let id = candidates.deployments()[0].post_asap_node_id; + (id, candidates.select(&choice(id)).unwrap_err()) + }; + let shared = SummaryMaintenanceLifecycle::Shared { + retention: DurationMs(10_000), + }; + let (id, error) = select(&UnitCosts, &|id| vec![(id, shared.clone())]); + assert_eq!( + error, + E::Rejected { + post_asap_node_id: id, + rejection: Some(SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime), + } + ); + let continuous = SummaryMaintenanceLifecycle::ContinuouslyMaintained; + let (id, error) = select(&crate::cost_model::DefaultCostModel, &|id| { + vec![(id, SummaryMaintenanceLifecycle::Ephemeral)] + }); + assert_eq!( + error, + E::Rejected { + post_asap_node_id: id, + rejection: Some(SummaryMaintenanceLifecycleRejection::MissingCostEvidence), + } + ); + let (id, error) = select(&UnitCosts, &|id| { + vec![( + id, + SummaryMaintenanceLifecycle::Shared { + retention: DurationMs(1), + }, + )] + }); + assert_eq!(error, E::NotAnAlternative(id)); + let (id, error) = select(&UnitCosts, &|_| vec![]); + assert_eq!(error, E::MissingChoice(id)); + let (id, error) = select(&UnitCosts, &|id| { + vec![(id, continuous.clone()), (id, continuous.clone())] + }); + assert_eq!(error, E::DuplicateChoice(id)); + let (_, error) = select(&UnitCosts, &|_| { + vec![(PostAsapNodeId(u32::MAX), continuous.clone())] + }); + assert_eq!(error, E::UnknownSummary(PostAsapNodeId(u32::MAX))); + } + + // Nested states on one maintenance path must share an evaluation schedule. + #[test] + fn explicit_choice_rejects_incompatible_nested_schedules() { + let workload = workload(vec![], vec![repeating()], continuous(1_000, 20_000)); + let candidates = enumerate_summary_maintenance_lifecycles( + nested_summary(), + WorkloadDemand::new_with_data(&workload, &continuous(1_000, 20_000), &[0]), + 1_000, + Some(Horizon(10.0)), + SummaryMaintenanceLifecycleCapabilities::ALL, + &IncompatibleNestedCosts, + ) + .unwrap(); + let [outer, inner] = candidates.deployments() else { + panic!("two summary states"); + }; + let choice = vec![ + ( + outer.post_asap_node_id, + SummaryMaintenanceLifecycle::Ephemeral, + ), + ( + inner.post_asap_node_id, + SummaryMaintenanceLifecycle::ContinuouslyMaintained, + ), + ]; + assert_eq!( + candidates.select(&choice).unwrap_err(), + SummaryMaintenanceLifecycleChoiceError::IncompatibleEvaluationSchedules + ); + } + + // A multi-summary root yields one candidate entry per unique state, with + // a shared `Rc` state listed once. + #[test] + fn enumeration_lists_each_unique_summary_state_once() { + let shared = summary(); + let root = Rc::new(SummaryNode { + expr: SummaryExpr::SummaryMerge { + timing: asap_types::post_asap::ExecutionTiming::IngestionTime, + children: vec![Rc::clone(&shared), Rc::clone(&shared), summary()], + }, + schema: shared.schema.clone(), + guarantee: None, + }); + let workload = workload(vec![batch(Predictability::AdHoc)], vec![], at_rest()); + let candidates = enumerate_summary_maintenance_lifecycles( + root, + WorkloadDemand::new_with_data(&workload, &at_rest(), &[0]), + 1_000, + None, + SummaryMaintenanceLifecycleCapabilities::ALL, + &UnitCosts, + ) + .unwrap(); + let ids: HashSet<_> = candidates + .deployments() + .iter() + .map(|deployment| deployment.post_asap_node_id) + .collect(); + assert_eq!(candidates.deployments().len(), 2); + assert_eq!(ids.len(), 2); + assert!(candidates + .deployments() + .iter() + .any(|deployment| Rc::ptr_eq(&deployment.summary, &shared))); + } } diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index 7e7b3bd0..2aad9ca7 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -236,6 +236,21 @@ using workload demand, window/freshness requirements and supported physical implementations. Backend selection uses runtime feasibility and cost after physical compilation. The following example follows one candidate. +Candidate generation and selection are separate steps. For every unique summary +state, enumeration reports each lifecycle (ephemeral, prepared, shared, +continuously maintained) as legal, with a Planner cost or explicitly unknown +cost, or as rejected with a reason. Planner does not remove a legal alternative +because its own estimate prefers another. A deployment prices the legal +alternatives over the whole workload, counting shared state once, and binds one +lifecycle per state. Binding checks that the choice is legal and that states on +one maintenance path share an evaluation schedule. An alternative whose cost is +unknown can be bound only when the deployment's cost model is authoritative for +complete-candidate cost; unknown cost is never treated as zero. It then yields the same +lifecycle guarantee and window framework the physical compiler consumes when +Planner selects. Planner's own cheapest-alternative selection remains available +for callers without deployment pricing. The window framework is decided for the +complete combination, not for one alternative in isolation. + For the running example, assume it selects: ```text diff --git a/docs/develop_docs/library-api.md b/docs/develop_docs/library-api.md index 92288067..cb541b94 100644 --- a/docs/develop_docs/library-api.md +++ b/docs/develop_docs/library-api.md @@ -633,6 +633,8 @@ that prepared or retained shared state is supported. | `plan_summary_maintenance_lifecycles` | Assembled logical DAG root, `WorkloadDemand`, `now_ms`, optional horizon, runtime capabilities, cost model | `Result` for that fixed root; does not revisit all semantic candidates | | `global_selection_with_summary_maintenance_lifecycles` | `PlanSpace`, workload/root-entry associations, time, horizon, capabilities, cost model | Lifecycle-aware compatible selection/error, using eligible cost evidence | | `assemble_selected_dag_with_summary_maintenance_lifecycles` | Selection, target root and lifecycle context | Optional lifecycle plan/error; attaches state deployment decisions | +| `enumerate_summary_maintenance_lifecycles` | Same inputs as `plan_summary_maintenance_lifecycles` | `SummaryMaintenanceLifecycleCandidates`: per unique state, every alternative with its cost or rejection; nothing selected. `guarantee(&lifecycle)` gives the mode/schedule that alternative would carry | +| `SummaryMaintenanceLifecycleCandidates::select(choices)` | One `(PostAsapNodeId, SummaryMaintenanceLifecycle)` per state, copied from `deployments()` | The same `SummaryMaintenanceLifecyclePlan` Planner selection would produce for that combination, or `SummaryMaintenanceLifecycleChoiceError` when a choice is unknown, missing, duplicated, rejected, schedule-incompatible, or not completely estimable | Inspect `deployments`, their selected lifecycle/alternatives/rejections, `selected_raw_recompute`, and optional summary/raw costs. Success of a function @@ -644,6 +646,13 @@ lifecycle analysis after structural selection can evaluate the selected root, but does not make the earlier selection lifecycle-optimal. An application may consume ranked candidates and perform this comparison downstream instead. +A deployment that prices lifecycles itself calls +`enumerate_summary_maintenance_lifecycles`, prices the alternatives, and binds +its choice with `select`. A choice is accepted only if Planner could select it: +an alternative with `MissingCostEvidence` is accepted only when the cost model's +complete-candidate hook covers lifecycle costs. Window frameworks and totals come +from that hook, as in Planner selection. + ## Optional whole-plan selection and DAG assembly ### What does global selection mean?