diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 5b0bcbb76..fbaf8875d 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -870,6 +870,54 @@ pub trait CostModel { false } + /// Price one complete multi-root workload assignment. Deployments shared + /// by `Rc` identity carry costs for the union of their consumers. The + /// additive default charges such a producer once, retaining each root's + /// remaining complete-DAG cost. Models with interactions between roots + /// override this hook; unknown evidence never becomes a zero estimate. + fn complete_workload_candidate_cost( + &self, + plans: &[crate::summary_maintenance_lifecycle::SummaryMaintenanceLifecyclePlan], + scalar_roots: &[(usize, asap_types::ir::ScalarExpr)], + ) -> Option { + // The additive legacy hooks price operator DAGs only. A workload model + // must supply scalar/subquery execution evidence rather than omit it. + if !scalar_roots.is_empty() { + return None; + } + let mut total = 0.0; + let mut seen = std::collections::HashSet::new(); + for plan in plans { + let root_cost = if plan.selected_raw_recompute { + plan.raw_recompute_total_cost? + } else { + plan.summary_total_cost? + }; + if !root_cost.0.is_finite() || root_cost.0 < 0.0 { + return None; + } + total += root_cost.0; + for deployment in &plan.deployments { + if seen.insert(std::rc::Rc::as_ptr(&deployment.summary)) { + continue; + } + let guarantee = deployment + .summary_maintenance_lifecycle_guarantee + .as_ref()?; + let cost = deployment + .alternatives + .iter() + .find(|alternative| { + alternative.summary_maintenance_lifecycle + == guarantee.summary_maintenance_lifecycle + })? + .total_cost?; + total -= cost.0; + } + } + (total.is_finite() && total >= 0.0).then_some(Cost(total)) + } + /// Cost of evaluating `target` directly from its logical/raw inputs once. /// When known, lifecycle-aware materialization compares this fallback with /// the aggregate cost of the selected summary deployments. diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 280ecb5c3..c8a528f36 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -50,6 +50,10 @@ //! values containing DAG roots and maintenance decisions; callers do not need //! to run ordinary selection/assembly first. //! +//! - Choose [`pass::CompletePass`] for bounded exhaustive workload search, +//! including exact pane composition, partial sharing and lifecycle assignments. +//! It requires complete workload costs and never falls back to a heuristic. +//! //! Models and evidence determine which choices the helpers can justify. //! Physical operator binding, placement, storage, deployment, and execution //! remain downstream responsibilities. Neither taking the first candidate nor diff --git a/crates/asap-aware-mapping/src/pass/complete.rs b/crates/asap-aware-mapping/src/pass/complete.rs new file mode 100644 index 000000000..fbae0c967 --- /dev/null +++ b/crates/asap-aware-mapping/src/pass/complete.rs @@ -0,0 +1,370 @@ +//! Exhaustive workload search. Unknown costs and budget exhaustion are explicit. + +use std::{collections::HashMap, rc::Rc}; + +use asap_types::ir::{cse::share_common_sub_dags, OperatorNode}; + +use super::{OptimizationInput, OptimizationPass, OptimizeError, PlanOutput, QueryLifecyclePlan}; +use crate::{ + replacement::{default_strategies_with_models, search_workload_with_targets}, + summary_maintenance_lifecycle::{ + enumerate_assembled_plans, SummaryMaintenanceLifecyclePlan, WorkloadDemand, + }, + window_composition::enumerate_window_compositions, +}; + +/// Strict complete-cost alternative to `MajorPass`. A distinct pass is needed +/// because the legacy default supports rank-only models, which cannot certify +/// the cheapest complete workload. No heuristic result is returned by this pass. +#[derive(Debug, Clone, Copy)] +pub struct CompletePass { + /// Maximum inventory/assignment count at each exhaustive stage. Exceeding + /// it aborts the search, even if a priced candidate was already discovered. + pub max_candidates: usize, +} + +impl Default for CompletePass { + fn default() -> Self { + Self { + max_candidates: 65_536, + } + } +} + +/// Reviewable inventory, including explanations for invalid/unpriced assignments. +#[derive(Debug)] +pub struct CompleteWorkloadInventory { + pub candidates: Vec, + pub rejected: Vec, +} + +fn failure(message: impl ToString) -> OptimizeError { + OptimizeError::CompleteSelection(message.to_string()) +} + +fn product(options: &[usize], limit: usize) -> Result { + if limit == 0 { + return Err(failure("complete search requires a positive budget")); + } + options + .iter() + .try_fold(1usize, |n, m| n.checked_mul(*m)) + .filter(|n| *n <= limit) + .ok_or_else(|| { + failure( + "complete workload search exceeds candidate budget; no partial inventory returned", + ) + }) +} + +impl CompletePass { + /// Enumerate logical choices, exact pane compositions, independent/partial + /// state-sharing partitions and all compatible lifecycle assignments. + pub fn enumerate( + &self, + input: OptimizationInput<'_>, + ) -> Result { + input.validate()?; + let models = input.models; + let workload = input.workload; + let entries: Vec<_> = workload.query_workload().entries().collect(); + let strategies = + default_strategies_with_models(models.cost, models.accuracy, models.evidence); + let roots = workload + .entries() + .zip(workload.operator_indices().iter().copied()) + .map(|((entry, root), id)| { + ( + id, + Rc::clone(root), + Some(entry.requirements.accuracy.target()), + ) + }) + .collect(); + let space = search_workload_with_targets(roots, &strategies, models.accuracy); + let logical = space + .enumerate_candidate_dags(self.max_candidates) + .map_err(failure)?; + let mut inventory = CompleteWorkloadInventory { + candidates: Vec::new(), + rejected: logical.rejected_assemblies, + }; + let mut assignments = 0usize; + let mut assemblies = 0usize; + for logical in logical.candidates { + let windows = logical + .iter() + .map(|(id, root)| { + enumerate_window_compositions(root, &entries[*id], self.max_candidates) + .map_err(failure) + }) + .collect::, _>>()?; + let count = product( + &windows.iter().map(Vec::len).collect::>(), + self.max_candidates, + )?; + for mut ordinal in 0..count { + let roots = windows + .iter() + .enumerate() + .map(|(index, choices)| { + let root = Rc::clone(&choices[ordinal % choices.len()]); + ordinal /= choices.len(); + (logical[index].0, root) + }) + .collect(); + for roots in sharing_partitions(roots, self.max_candidates)? { + assemblies += 1; + if assemblies > self.max_candidates { + return Err(failure("complete workload assembly budget exceeded; no partial inventory returned")); + } + let mut state_entries = HashMap::new(); + // Lifecycle enumeration also collects maintained populations; + // bind every reachable node so those states get exact demand. + for (id, root) in &roots { + for node in OperatorNode::reachable(root) { + let indices: &mut Vec = + state_entries.entry(Rc::as_ptr(&node)).or_default(); + if !indices.contains(id) { + indices.push(*id); + } + } + } + let mut choices = Vec::new(); + let mut invalid = false; + for (index, (id, root)) in roots.iter().enumerate() { + match enumerate_assembled_plans( + Rc::clone(root), + &space.roots[index].1, + WorkloadDemand { + workload: workload.query_workload(), + data_workload: workload.data_workload(), + entry_indices: &[*id], + }, + &state_entries, + input.lifecycle.now_ms, + input.lifecycle.horizon, + models.capabilities, + models.cost, + self.max_candidates, + ) { + Ok(plans) => choices.push(plans), + Err(reason) if reason.contains("budget") => { + return Err(failure(reason)) + } + Err(reason) => { + inventory.rejected.push(reason); + invalid = true; + break; + } + } + } + if invalid { + continue; + } + let count = product( + &choices.iter().map(Vec::len).collect::>(), + self.max_candidates, + )?; + assignments = assignments.checked_add(count).filter(|n| *n <= self.max_candidates) + .ok_or_else(|| failure("complete workload assignment budget exceeded; no partial inventory returned"))?; + for mut ordinal in 0..count { + let plans: Vec<_> = choices + .iter() + .map(|choices| { + let plan = choices[ordinal % choices.len()].clone(); + ordinal /= choices.len(); + plan + }) + .collect(); + if !compatible(&plans) + || plans.iter().any(|plan| plan.execution_timed_dag().is_err()) + { + inventory.rejected.push( + "shared state has inconsistent lifecycle/window assignment".into(), + ); + continue; + } + let Some(cost) = models + .cost + .complete_workload_candidate_cost(&plans, workload.scalar_roots()) + else { + inventory + .rejected + .push("complete workload cost is unknown".into()); + continue; + }; + if !cost.0.is_finite() || cost.0 < 0.0 { + inventory + .rejected + .push("complete workload cost is invalid".into()); + continue; + } + let plans = plans + .into_iter() + .zip(&roots) + .map(|(plan, (id, _))| QueryLifecyclePlan { + entry_index: *id, + plan, + }) + .collect(); + let mut output = PlanOutput::new(plans); + output.scalar_roots = workload.scalar_roots().to_vec(); + output.workload_total_cost = Some(cost); + inventory.candidates.push(output); + } + } + } + } + inventory.rejected.sort(); + inventory.rejected.dedup(); + Ok(inventory) + } +} + +impl OptimizationPass for CompletePass { + fn name(&self) -> &'static str { + "complete" + } + + fn optimize(&self, input: OptimizationInput<'_>) -> Result { + let inventory = self.enumerate(input)?; + inventory + .candidates + .into_iter() + .min_by(|a, b| { + a.workload_total_cost + .unwrap() + .0 + .total_cmp(&b.workload_total_cost.unwrap().0) + }) + .ok_or_else(|| { + failure(format!( + "no feasible priced complete workload: {}", + inventory.rejected.join("; ") + )) + }) + } +} + +fn compatible(plans: &[SummaryMaintenanceLifecyclePlan]) -> bool { + let mut assignments = HashMap::new(); + for plan in plans { + for deployment in &plan.deployments { + let assignment = ( + &deployment.summary_maintenance_lifecycle_guarantee, + &deployment.selected_window_framework, + ); + if assignments + .insert(Rc::as_ptr(&deployment.summary), assignment) + .is_some_and(|previous| previous != assignment) + { + return false; + } + } + } + true +} + +/// Enumerate every set partition for each structurally identical state class. +/// Parents are rebuilt per reader so sharing an estimate cannot accidentally +/// force its producer to be shared after choosing independent states. +type WorkloadRoots = Vec<(usize, Rc)>; + +fn sharing_partitions( + roots: WorkloadRoots, + limit: usize, +) -> Result, OptimizeError> { + let roots = share_common_sub_dags(roots); + let mut classes: Vec<(Rc, Vec)> = Vec::new(); + for (reader, (_, root)) in roots.iter().enumerate() { + for node in OperatorNode::reachable(root) { + if let Some((_, readers)) = classes + .iter_mut() + .find(|(state, _)| Rc::ptr_eq(state, &node)) + { + if !readers.contains(&reader) { + readers.push(reader); + } + } else { + classes.push((node, vec![reader])); + } + } + } + classes.retain(|(_, readers)| readers.len() > 1); + fn partitions(n: usize, limit: usize) -> Result>, OptimizeError> { + let mut result = vec![vec![0]]; + for _ in 1..n { + let mut next = Vec::new(); + for partition in result { + for group in 0..=partition.iter().max().unwrap() + 1 { + let mut choice = partition.clone(); + choice.push(group); + next.push(choice); + if next.len() > limit { + return Err(failure("state-sharing partition budget exceeded")); + } + } + } + result = next; + } + Ok(result) + } + let partitions = classes + .iter() + .map(|(_, readers)| partitions(readers.len(), limit)) + .collect::, _>>()?; + let count = product(&partitions.iter().map(Vec::len).collect::>(), limit)?; + let mut result = Vec::new(); + for mut ordinal in 0..count { + let selected: Vec<_> = partitions + .iter() + .map(|partitions| { + let choice = &partitions[ordinal % partitions.len()]; + ordinal /= partitions.len(); + choice + }) + .collect(); + let mut memo = HashMap::new(); + fn rebuild( + node: &Rc, + reader: usize, + classes: &[(Rc, Vec)], + selected: &[&Vec], + memo: &mut HashMap<(*const OperatorNode, usize), Rc>, + ) -> Rc { + let group = classes + .iter() + .enumerate() + .find_map(|(index, (state, readers))| { + Rc::ptr_eq(state, node).then(|| { + selected[index][readers.iter().position(|r| *r == reader).unwrap()] + }) + }); + let key = (Rc::as_ptr(node), group.map_or(reader, roots_namespace)); + if let Some(node) = memo.get(&key) { + return Rc::clone(node); + } + let mut rebuilt = (**node).clone(); + rebuilt.operator = node + .operator + .map_children(|child| rebuild(child, reader, classes, selected, memo)); + let rebuilt = Rc::new(rebuilt); + memo.insert(key, Rc::clone(&rebuilt)); + rebuilt + } + let rebuilt = roots + .iter() + .enumerate() + .map(|(reader, (id, node))| { + (*id, rebuild(node, reader, &classes, &selected, &mut memo)) + }) + .collect(); + result.push(rebuilt); + } + Ok(result) +} + +fn roots_namespace(group: usize) -> usize { + usize::MAX - group +} diff --git a/crates/asap-aware-mapping/src/pass/mod.rs b/crates/asap-aware-mapping/src/pass/mod.rs index 34e691cc0..7fabc73a6 100644 --- a/crates/asap-aware-mapping/src/pass/mod.rs +++ b/crates/asap-aware-mapping/src/pass/mod.rs @@ -12,6 +12,7 @@ //! validates the input once for every pass and checks the output contract that //! downstream consumers rely on. +mod complete; mod major; use std::collections::BTreeMap; @@ -32,6 +33,7 @@ use crate::summary_maintenance_lifecycle::{ SummaryMaintenanceLifecyclePlan, SummaryMaintenanceLifecycleSelectionError, }; +pub use complete::{CompletePass, CompleteWorkloadInventory}; pub use major::MajorPass; static DEFAULT_COST_MODEL: DefaultCostModel = DefaultCostModel; @@ -196,6 +198,9 @@ pub struct PlanOutput { pub plans: Vec, /// Exact scalar expressions, keyed by workload entry; embedded plan reads remain visible. pub scalar_roots: Vec<(usize, asap_types::ir::ScalarExpr)>, + /// A certified comparison cost for the complete selected workload when + /// supplied by a complete-cost pass. Rank-only passes leave it unknown. + pub workload_total_cost: Option, } impl PlanOutput { @@ -203,6 +208,7 @@ impl PlanOutput { Self { plans, scalar_roots: Vec::new(), + workload_total_cost: None, } } @@ -281,6 +287,8 @@ impl PlanOutput { #[derive(Debug, thiserror::Error)] #[non_exhaustive] pub enum OptimizeError { + #[error("complete workload selection: {0}")] + CompleteSelection(String), #[error("optimization input: {0}")] Input(#[from] OptimizationInputError), #[error("entry {entry_index}: {source}")] @@ -341,6 +349,12 @@ fn check_contract( ) -> Result<(), OptimizeError> { let violation = |detail: String| OptimizeError::ContractViolation { pass, detail }; + if output + .workload_total_cost + .is_some_and(|cost| !cost.0.is_finite() || cost.0 < 0.0) + { + return Err(violation("invalid complete workload cost".into())); + } if output.len() != expected_len { return Err(violation(format!( "{} plan(s) for {expected_len} workload entry/entries", @@ -378,13 +392,16 @@ impl PassRegistry { Self::default() } - /// Only [`MajorPass`], under the name `major`. + /// The rank-compatible `major` and strict complete-cost `complete` passes. pub fn with_builtin() -> Self { let mut registry = Self::new(); registry .register(Box::new(MajorPass)) .expect("empty registry cannot conflict"); registry + .register(Box::new(CompletePass::default())) + .expect("distinct built-in names"); + registry } /// Keyed by `pass.name()`. Registering a name twice is an error rather @@ -469,6 +486,9 @@ mod tests { registry.register(Box::new(Stub("alpha"))).unwrap(); assert!(registry.get("major").is_some()); assert!(registry.get("absent").is_none()); - assert_eq!(registry.names().collect::>(), vec!["alpha", "major"]); + assert_eq!( + registry.names().collect::>(), + vec!["alpha", "complete", "major"] + ); } } diff --git a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index de56c100d..9aebf8007 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -391,6 +391,7 @@ pub enum SummaryMaintenanceLifecycleSelectionError { /// 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. +#[derive(Clone)] pub struct SummaryMaintenanceLifecycleCandidates<'a> { /// Unselected plan: deployments carry alternatives but no guarantee or /// window framework. @@ -405,6 +406,8 @@ pub struct SummaryMaintenanceLifecycleCandidates<'a> { /// Why an explicit per-state lifecycle choice cannot be bound. #[derive(Debug, thiserror::Error, PartialEq)] pub enum SummaryMaintenanceLifecycleChoiceError { + #[error("lifecycle enumeration exceeds candidate budget; no partial inventory returned")] + BudgetExceeded, #[error("summary {0:?} is not a deployment of this root")] UnknownSummary(PostAsapNodeId), #[error("summary {0:?} is chosen more than once")] @@ -425,6 +428,58 @@ pub enum SummaryMaintenanceLifecycleChoiceError { } impl SummaryMaintenanceLifecycleCandidates<'_> { + /// Every fully costed, legal lifecycle assignment, rather than only the + /// per-root cheapest assignment. Whole-workload selection needs this + /// inventory because sharing can change which assignment wins. + pub fn enumerate_plans( + &self, + limit: usize, + ) -> Result, SummaryMaintenanceLifecycleChoiceError> { + use SummaryMaintenanceLifecycleChoiceError as E; + if limit == 0 { + return Err(E::BudgetExceeded); + } + let options: Vec> = self + .plan + .deployments + .iter() + .map(|deployment| { + deployment + .alternatives + .iter() + .filter(|alternative| self.context().eligible(alternative)) + .map(|alternative| { + ( + deployment.post_asap_node_id, + alternative.summary_maintenance_lifecycle.clone(), + ) + }) + .collect() + }) + .collect(); + let count = options + .iter() + .try_fold(1usize, |count, options| count.checked_mul(options.len())) + .filter(|count| *count <= limit) + .ok_or(E::BudgetExceeded)?; + let mut plans = Vec::new(); + for mut ordinal in 0..count { + let choices: Vec<_> = options + .iter() + .map(|options| { + let choice = options[ordinal % options.len()].clone(); + ordinal /= options.len(); + choice + }) + .collect(); + match self.clone().select(&choices) { + Ok(plan) => plans.push(plan), + Err(E::NoCompleteEstimate | E::IncompatibleEvaluationSchedules) => {} + Err(error) => return Err(error), + } + } + Ok(plans) + } /// One entry per unique retained state (see /// [`SummaryMaintenanceLifecyclePlan::deployments`]), with every /// alternative and its rejection; no lifecycle or window framework is @@ -636,7 +691,7 @@ pub fn enumerate_summary_maintenance_lifecycles<'a>( /// eligibility and data-arrival facts; `profile` supplies effective uses after /// DAG path multiplicity has been propagated by `CandidateLogicalASAPDAGs`. #[expect(clippy::too_many_arguments, reason = "internal bound planning context")] -fn enumerate_with_profile<'a>( +pub(crate) fn enumerate_with_profile<'a>( root: Rc, demand: WorkloadDemand<'_>, now_ms: u64, @@ -1089,6 +1144,75 @@ pub(crate) fn plan_assembled_dag( Ok(plan) } +/// Enumerate assignments without replacing a summary candidate by raw fallback. +/// Each shared state is priced against exactly the entries consuming that state. +#[expect(clippy::too_many_arguments, reason = "bound workload planning context")] +pub(crate) fn enumerate_assembled_plans( + root: Rc, + target: &Rc, + demand: WorkloadDemand<'_>, + state_entries: &HashMap<*const OperatorNode, Vec>, + now_ms: u64, + horizon: Option, + capabilities: SummaryMaintenanceLifecycleCapabilities, + cost_model: &dyn CostModel, + limit: usize, +) -> Result, String> { + if !root.contains_asap() { + return plan_assembled_dag( + Rc::clone(&root), + &root, + demand, + now_ms, + horizon, + capabilities, + cost_model, + ) + .map(|plan| vec![plan]) + .map_err(|error| error.to_string()); + } + let mut candidates = enumerate_with_profile( + root, + demand, + now_ms, + horizon, + capabilities, + cost_model, + None, + Some(target), + ) + .map_err(|error| error.to_string())?; + for deployment in &mut candidates.plan.deployments { + let entries = &state_entries[&Rc::as_ptr(&deployment.summary)]; + let state = enumerate_with_profile( + Rc::clone(&deployment.summary), + WorkloadDemand { + entry_indices: entries, + ..demand + }, + now_ms, + horizon, + capabilities, + cost_model, + None, + None, + ) + .map_err(|error| error.to_string())?; + let alternatives = state + .plan + .deployments + .iter() + .find(|state| Rc::ptr_eq(&state.summary, &deployment.summary)) + .expect("a collected state is its own deployment"); + deployment + .alternatives + .clone_from(&alternatives.alternatives); + } + candidates + .enumerate_plans(limit) + .map_err(|error| error.to_string()) +} + /// A quote is per execution, so every consumer's bound applies even when /// repeated reads make this alternative cheap over the planning horizon. fn raw_response_latency_violation( diff --git a/crates/integration-tests/tests/complete_workload_selection.rs b/crates/integration-tests/tests/complete_workload_selection.rs new file mode 100644 index 000000000..69533c4e0 --- /dev/null +++ b/crates/integration-tests/tests/complete_workload_selection.rs @@ -0,0 +1,303 @@ +//! Complete selection considers locally expensive choices and executes its winner. +mod physical_common; + +use asap_aware_mapping::{ + cost_model::{Cost, CostModel, DefaultCostModel}, + pass::{optimize, CompletePass, LifecycleInput, OptimizationInput, PlanningModels}, + CostRate, SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecyclePlan, +}; +use asap_physical_operators::values::Value; +use asap_types::{ + ir::{ + operator_properties::{Reduction, Source}, + NonASAPOp, Operator, OperatorNode, + }, + parsed_workload::ParsedWorkload, + pre_asap::{AggIntent, DataType, Field, Schema}, + workload::{ + BatchEntry, PlanningWorkload, Predictability, Query, QueryLanguage, QueryWorkload, + SqlDialect, + }, +}; +use std::rc::Rc; + +struct InteractingCosts; +impl CostModel for InteractingCosts { + fn rank_candidates( + &self, + intent: &AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + DefaultCostModel.rank_candidates(intent, candidates) + } + fn raw_query_recompute_cost(&self, _: &OperatorNode) -> Option { + Some(Cost(1.0)) + } + fn summary_maintenance_lifecycle_cost_inputs( + &self, + _: &OperatorNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + SummaryMaintenanceLifecycleCostInputs { + build_cost: Some(Cost(10.0)), + maintenance_cost_per_update: Some(Cost::ZERO), + summary_read_cost: Some(Cost::ZERO), + retention_cost_rate: Some(CostRate(0.0)), + retirement_cost: Some(Cost::ZERO), + } + } + fn complete_workload_candidate_cost( + &self, + plans: &[SummaryMaintenanceLifecyclePlan], + _: &[(usize, asap_types::ir::ScalarExpr)], + ) -> Option { + // A deployment's complete implementation has a discount only when + // both summary consumers use it. Per-site ranking cannot see this. + Some(Cost( + if plans.iter().all(|plan| !plan.selected_raw_recompute) { + 0.25 + } else { + 2.0 + }, + )) + } +} + +fn fixture() -> ParsedWorkload { + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::Table { + table_ref: "events".into(), + }, + predicates: vec![], + schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + })) + .unwrap(); + let roots = [ + AggIntent::Sum { col: Some(0) }, + AggIntent::Max { col: Some(0) }, + ] + .into_iter() + .map(|intent| { + OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Aggregate { + reduction: Reduction::by(vec![]), + measures: vec![intent], + output_names: vec![], + filters: vec![], + having: None, + child: Rc::clone(&scan), + })) + .unwrap() + }) + .collect(); + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::SQL(SqlDialect::DataFusionSQL), + repeating_queries: None, + query_batch: Some( + [ + "SELECT SUM(value) FROM events", + "SELECT MAX(value) FROM events", + ] + .into_iter() + .map(|query| BatchEntry { + query: Query(query.into()), + requirements: Default::default(), + predictability: Predictability::AdHoc, + invocations: 1, + execute_at: None, + time_selection: Default::default(), + }) + .collect(), + ), + }, + data_workload: None, + }; + ParsedWorkload::new(workload, roots).unwrap() +} + +/// Both summaries are locally more expensive than raw; their complete implementation wins. +#[test] +fn globally_cheapest_assignment_executes_both_queries() { + let workload = fixture(); + let input = OptimizationInput::new( + &workload, + PlanningModels::builtin().with_cost(&InteractingCosts), + LifecycleInput::new(0), + ); + let pass = CompletePass { + max_candidates: 4096, + }; + let inventory = pass.enumerate(input).unwrap(); + assert!(inventory.candidates.iter().any(|candidate| candidate + .plans + .iter() + .all(|plan| plan.plan.selected_raw_recompute))); + let output = optimize(&pass, input).unwrap(); + assert_eq!(output.workload_total_cost, Some(Cost(0.25))); + assert!(output + .plans + .iter() + .all(|plan| !plan.plan.selected_raw_recompute)); + for (plan, expected) in output.plans.iter().zip([3.0, 2.0]) { + let result = physical_common::execute_raw_rows( + &plan.plan.root, + vec![vec![Value::Float64(1.0)], vec![Value::Float64(2.0)]], + ); + assert_eq!(result.len(), 1); + assert!(matches!(result[0].as_slice(), [Value::Float64(value)] if *value == expected)); + } +} + +/// Exhaustion and unknown complete costs never yield a heuristic or partially searched winner. +#[test] +fn incomplete_search_is_an_explicit_error() { + let workload = fixture(); + let input = OptimizationInput::new( + &workload, + PlanningModels::builtin().with_cost(&InteractingCosts), + LifecycleInput::new(0), + ); + assert!(CompletePass { max_candidates: 1 } + .enumerate(input) + .unwrap_err() + .to_string() + .contains("budget")); + let unknown = + OptimizationInput::new(&workload, PlanningModels::builtin(), LifecycleInput::new(0)); + assert!(optimize(&CompletePass::default(), unknown) + .unwrap_err() + .to_string() + .contains("no feasible priced")); +} + +struct PartialSharingCosts; +impl CostModel for PartialSharingCosts { + fn rank_candidates( + &self, + intent: &AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + DefaultCostModel.rank_candidates(intent, candidates) + } + fn raw_query_recompute_cost(&self, _: &OperatorNode) -> Option { + Some(Cost(1000.0)) + } + fn summary_maintenance_lifecycle_cost_inputs( + &self, + node: &OperatorNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + InteractingCosts.summary_maintenance_lifecycle_cost_inputs(node) + } + fn complete_workload_candidate_cost( + &self, + plans: &[SummaryMaintenanceLifecyclePlan], + _: &[(usize, asap_types::ir::ScalarExpr)], + ) -> Option { + if plans.iter().any(|plan| plan.deployments.is_empty()) { + return Some(Cost(1000.0)); + } + let states: Vec<_> = plans + .iter() + .map(|plan| &plan.deployments[0].summary) + .collect(); + Some(Cost( + if Rc::ptr_eq(states[0], states[1]) && !Rc::ptr_eq(states[0], states[2]) { + 0.1 + } else { + 10.0 + }, + )) + } +} + +/// A partial sharing partition can beat both all-shared and all-independent choices. +#[test] +fn three_consumers_can_share_only_a_subset() { + let two = fixture(); + let root = Rc::clone(two.entries().next().unwrap().1); + let mut workload = two.planning_workload().clone(); + let entry = workload.query_workload.query_batch.as_ref().unwrap()[0].clone(); + workload.query_workload.query_batch = Some(vec![entry.clone(), entry.clone(), entry]); + let workload = ParsedWorkload::new(workload, vec![root.clone(), root.clone(), root]).unwrap(); + let input = OptimizationInput::new( + &workload, + PlanningModels::builtin().with_cost(&PartialSharingCosts), + LifecycleInput::new(0), + ); + let output = optimize( + &CompletePass { + max_candidates: 65_536, + }, + input, + ) + .unwrap(); + assert_eq!(output.workload_total_cost, Some(Cost(0.1))); + let states: Vec<_> = output + .plans + .iter() + .map(|plan| &plan.plan.deployments[0].summary) + .collect(); + assert!(Rc::ptr_eq(states[0], states[1])); + assert!(!Rc::ptr_eq(states[0], states[2])); +} + +struct AdditiveCosts; +impl CostModel for AdditiveCosts { + fn rank_candidates( + &self, + intent: &AggIntent, + candidates: &[asap_types::post_asap::SketchAlgorithm], + ) -> Vec { + DefaultCostModel.rank_candidates(intent, candidates) + } + fn raw_query_recompute_cost(&self, _: &OperatorNode) -> Option { + Some(Cost(1000.0)) + } + fn summary_maintenance_lifecycle_cost_inputs( + &self, + node: &OperatorNode, + ) -> SummaryMaintenanceLifecycleCostInputs { + InteractingCosts.summary_maintenance_lifecycle_cost_inputs(node) + } +} + +/// Union demand permits retention, and the additive workload model charges the shared producer once. +#[test] +fn shared_producer_cost_is_not_multiplied_by_consumer_count() { + let two = fixture(); + let root = Rc::clone(two.entries().next().unwrap().1); + let mut workload = two.planning_workload().clone(); + workload.data_workload = Some(asap_types::workload::DataWorkload { + arrival: asap_types::workload::DataArrival::AtRest, + ..Default::default() + }); + let entry = workload.query_workload.query_batch.as_ref().unwrap()[0].clone(); + workload.query_workload.query_batch = Some(vec![entry.clone(), entry]); + let workload = ParsedWorkload::new(workload, vec![root.clone(), root]).unwrap(); + let input = OptimizationInput::new( + &workload, + PlanningModels::builtin().with_cost(&AdditiveCosts), + LifecycleInput::new(0).with_horizon(asap_aware_mapping::Horizon(100.0)), + ); + let output = optimize( + &CompletePass { + max_candidates: 4096, + }, + input, + ) + .unwrap(); + assert_eq!(output.workload_total_cost, Some(Cost(10.0))); + let left = &output.plans[0].plan.deployments[0]; + let right = &output.plans[1].plan.deployments[0]; + assert!(Rc::ptr_eq(&left.summary, &right.summary)); + assert_eq!( + left.summary_maintenance_lifecycle_guarantee, + right.summary_maintenance_lifecycle_guarantee + ); + assert!(matches!( + left.summary_maintenance_lifecycle_guarantee + .as_ref() + .unwrap() + .summary_maintenance_lifecycle, + asap_types::post_asap::SummaryMaintenanceLifecycle::Shared { .. } + )); +} diff --git a/crates/planner/src/lib.rs b/crates/planner/src/lib.rs index 501bc9ead..0ceef759b 100644 --- a/crates/planner/src/lib.rs +++ b/crates/planner/src/lib.rs @@ -25,8 +25,9 @@ use asap_frontend_sql::{lower_sql_dialect, SqlCatalog, SqlError}; // configures the same models and reads the same output whether it goes through // `e2e_plan` or straight to `optimize`. pub use asap_aware_mapping::pass::{ - optimize, LifecycleInput, MajorPass, OptimizationInput, OptimizationPass, OptimizeError, - PassRegistry, PlanOutput, PlanningModels, QueryLifecyclePlan, + optimize, CompletePass, CompleteWorkloadInventory, LifecycleInput, MajorPass, + OptimizationInput, OptimizationPass, OptimizeError, PassRegistry, PlanOutput, PlanningModels, + QueryLifecyclePlan, }; // ── Input ──────────────────────────────────────────────────────────────── diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md index 2acee4fe1..02b61ef3c 100644 --- a/docs/develop_docs/planner-layering-status.md +++ b/docs/develop_docs/planner-layering-status.md @@ -8,10 +8,10 @@ is a target contract, not a statement that its examples execute today. | --- | --- | --- | | Language frontends and common logical IR | SQL/PromQL/MetricsQL lower to unified operators and scalars. | Floating-point SQL frequency L2 products and the normalized natural-log entropy idiom now have conservative logical rewrites. Integer L2 products remain unrecognized. Preserve alias lineage, filters, NULL groups, empty inputs, count overflow and entropy units when adding recognition. | | Local exact and summary alternatives | `replacement::summary_candidates`, realization rules and candidate inventory exist; supplied accuracy models reach Pass 1. | Specialized entropy/norm families in Example 2 are illustrative, not registered families. UnivMon certifies only unit-update total count; L2, entropy and cardinality epsilon/delta bounds need verified evidence or a deployment model. | -| Summary-capability sharing | CSE interns structurally identical producers, including states with different readers. | It does not enumerate all partial sharing partitions or resize compatible states to the strictest consumer. Example 2's 37 candidates are not an acceptance result. | +| Summary-capability sharing | CSE interns structurally identical producers, including states with different readers. | The strict complete pass enumerates partial sharing partitions for identical admitted producers. It does not resize different state parameters to the strictest consumer. Example 2's 37 candidates are not an acceptance result. | | Window composition | Mergeable state IR/native merge exists; physical pane compatibility and reuse cost helpers exist. | Cadence-based disjoint pane composition and automatic native materialization-frontier enumeration are available (acceptance below). Retained rotating panes, historical EH buckets and boundary-error certificates remain separate runtime work. | | Physical materialization | Ephemeral/prepared/shared/continuously maintained lifecycle alternatives, costing, capabilities and latency checks exist. | Incremental query-time pane retention, historical backfill and the complete Example 4 matrix need executable implementations and explicit state/input contracts. | -| Whole-workload selection | One unified selected DAG; shared states are interned and costed across their consumers. | `replacement.rs` documents its selection as non-exhaustive over interacting choices. The proposal's cheapest complete candidate guarantee and 54/156 inventories need a complete workload search/selection path. | +| Whole-workload selection | One unified selected DAG; shared states are interned and costed across their consumers. | `MajorPass` remains non-exhaustive and rank-compatible. `CompletePass` supplies an exhaustive complete-cost path (acceptance below); the illustrative 54/156 counts are not asserted without their rule sets. | | Deployment inputs and execution | `PlanningModels` bundles cost, accuracy, evidence and capabilities; native typed UnivMon supports one build with three readouts. | At #557 native exact distinct/L2/entropy fallback is absent; the follow-ups supply those native bindings. End-to-end SQL Example 2 is not established by the native UnivMon fixture. | | Subtract/delete, parallelism, partitioning and resource planning | Some runtime memory/cancellation limits and maintenance capability flags exist. | These remain proposal TODOs; capability flags do not supply missing IR operators or a physical resource search. | @@ -126,3 +126,35 @@ panes execute identically with and without retained outputs. The deployment still supplies each selector's exact raw window and binds retained outputs to that window/revision. This does not introduce a rotating pane cache or claim that cadence alone certifies compatibility with a catalog's pane origin. + +## Complete workload selection acceptance + +`pass::CompletePass` is registered as `complete` and re-exported by the planner +facade. Choose it with `UserInput::with_pass(&CompletePass::default())`, or call +`optimize` on an existing `ParsedWorkload`. It is opt-in because the default +`major` pass accepts rank-only models; a rank cannot certify a complete cost. + +The pass enumerates the registered logical inventory, cadence-derived pane +compositions, every partial partition of identical sharing-legal producers, +and every legal lifecycle assignment. It checks shared lifecycle/window +consistency and executable phase contracts before comparing complete workload +quotes. Each state is bound to the union of precisely its consumers; a shared +producer is charged once by the additive cost hook. Nonadditive deployments +override `CostModel::complete_workload_candidate_cost`, including interactions +between roots and their physical implementation choices. Scalar workloads need +a complete-workload override because the legacy per-root hooks do not price +scalar execution. Missing/invalid quotes are retained as rejection reasons. + +`CompletePass::enumerate` exposes priced full assignments and rejection reasons; +selection returns the cheapest quote in that inventory and records +`PlanOutput::workload_total_cost`. An exceeded logical, pane, partition or +lifecycle budget is an error even if a priced candidate was already found. +There is no heuristic fallback. This guarantee is over registered alternatives +with valid complete evidence; it does not invent missing sketch certificates, +physical implementations or additional parameter-sizing rules. + +`integration-tests/tests/complete_workload_selection.rs` executes a winner that +local ranking would miss, finds a winning partial partition, checks union-demand +retention and single producer charging, and rejects budget exhaustion/unknown +costs. Native `compile_materialization_candidates` supplies the executable +frontier inventory to deployment models that compare physical placements.