diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index 493042499..280ecb5c3 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -179,6 +179,7 @@ pub mod summary_maintenance_lifecycle; #[cfg(test)] mod test_support; pub mod topk_reuse; +pub mod window_composition; pub use accuracy::reconciliation::AccuracyReconciliationStrategy; pub use accuracy::{ diff --git a/crates/asap-aware-mapping/src/window_composition.rs b/crates/asap-aware-mapping/src/window_composition.rs new file mode 100644 index 000000000..cb3ef2bda --- /dev/null +++ b/crates/asap-aware-mapping/src/window_composition.rs @@ -0,0 +1,279 @@ +//! Exact window composition over mergeable summary states. +//! +//! Candidate generation is independent of costing. Relative pane selectors keep +//! the original source, predicates and anchor; stored outputs must still be +//! bound to the evaluation window and revision by the physical runtime. + +use std::{rc::Rc, time::Duration}; + +use asap_types::{ + ir::{operator_properties::TimeShift, ASAPOp, NonASAPOp, Operator, OperatorNode}, + workload::{QueryRecurrence, QueryWorkloadEntry, RepeatedDemand}, +}; + +#[derive(Debug, thiserror::Error)] +pub enum WindowCompositionError { + #[error("window composition exceeds candidate budget; no partial inventory returned")] + BudgetExceeded, + #[error("window composition: {0}")] + Schema(#[from] asap_types::ir::SchemaDerivationError), +} + +/// Enumerate retaining each state and composing it from disjoint cadence-sized +/// panes. A fixed recurrence without a phase remains executable by rebuilding +/// relative panes; it does not certify reuse of a catalog's pane layout. +/// Every combination of eligible states is returned, including the original. +pub fn enumerate_window_compositions( + root: &Rc, + entry: &QueryWorkloadEntry, + limit: usize, +) -> Result>, WindowCompositionError> { + let cadence = match &entry.recurrence { + QueryRecurrence::Repeated(RepeatedDemand::FixedInterval(interval)) + | QueryRecurrence::Repeated(RepeatedDemand::FixedIntervalAt { interval, .. }) => { + u64::from(interval.0) + } + QueryRecurrence::Repeated(RepeatedDemand::Scheduled(times)) => times + .windows(2) + .filter_map(|pair| pair[1].0.checked_sub(pair[0].0)) + .fold(0, gcd), + _ => 0, + }; + if limit == 0 { + return Err(WindowCompositionError::BudgetExceeded); + } + if cadence == 0 { + return Ok(vec![Rc::clone(root)]); + } + let mut sites = Vec::new(); + for node in OperatorNode::reachable(root) { + if let Some(composed) = compose(&node, cadence, limit)? { + sites.push((node, composed)); + } + } + let count = 1usize + .checked_shl(u32::try_from(sites.len()).unwrap_or(u32::MAX)) + .filter(|count| *count <= limit) + .ok_or(WindowCompositionError::BudgetExceeded)?; + let mut candidates = Vec::with_capacity(count); + for mask in 0..count { + fn rebuild( + node: &Rc, + sites: &[(Rc, Rc)], + mask: usize, + ) -> Result, WindowCompositionError> { + if let Some((index, (_, composed))) = sites + .iter() + .enumerate() + .find(|(_, (original, _))| Rc::ptr_eq(original, node)) + { + if mask & (1 << index) != 0 { + return Ok(Rc::clone(composed)); + } + } + let children = node.children(); + let rebuilt = children + .iter() + .map(|child| rebuild(child, sites, mask)) + .collect::, _>>()?; + if children.iter().zip(&rebuilt).all(|(a, b)| Rc::ptr_eq(a, b)) { + return Ok(Rc::clone(node)); + } + let replacements = children.iter().zip(&rebuilt); + let mut result = node.map_children(|child| { + // map_children also visits scalar plan references. + replacements + .clone() + .find(|(original, _)| Rc::ptr_eq(original, child)) + .map_or_else(|| Rc::clone(child), |(_, rebuilt)| Rc::clone(rebuilt)) + })?; + result.guarantee = node.guarantee.clone(); + Ok(Rc::new(result)) + } + candidates.push(rebuild(root, &sites, mask)?); + } + Ok(candidates) +} + +fn gcd(mut a: u64, mut b: u64) -> u64 { + while b != 0 { + (a, b) = (b, a % b); + } + a +} + +fn compose( + state: &Rc, + cadence: u64, + limit: usize, +) -> Result>, WindowCompositionError> { + let Operator::ASAP(ASAPOp::SummaryAgg { child, .. }) = &state.operator else { + return Ok(None); + }; + let original_input = child; + let Some(NonASAPOp::TimeRange { range, kind, child }) = child.non_asap() else { + return Ok(None); + }; + let Ok(lookback) = u64::try_from(range.as_millis()) else { + return Ok(None); + }; + if *range != Duration::from_millis(lookback) || lookback == 0 { + return Ok(None); + } + let width = gcd(lookback, cadence); + let count = lookback / width; + if count <= 1 { + return Ok(None); + } + if count > limit as u64 || lookback > i64::MAX as u64 { + return Err(WindowCompositionError::BudgetExceeded); + } + let (source, shift) = match child.non_asap() { + Some(NonASAPOp::Scan { .. }) => (child, TimeShift::default()), + Some(NonASAPOp::TimeShift { child, shift }) + if matches!(child.non_asap(), Some(NonASAPOp::Scan { .. })) => + { + (child, *shift) + } + _ => return Ok(None), + }; + if !matches!( + source.non_asap(), + Some(NonASAPOp::Scan { + source: asap_types::ir::operator_properties::Source::TimeSeries { .. }, + .. + }) + ) { + return Ok(None); + } + let mut panes = Vec::new(); + for pane in 0..count { + let Some(offset_ms) = shift.offset_ms.checked_add((pane * width) as i64) else { + return Ok(None); + }; + let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { + child: Rc::clone(source), + shift: TimeShift { offset_ms, ..shift }, + }))?; + let input = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: Duration::from_millis(width), + kind: *kind, + child: shifted, + }))?; + panes.push(Rc::new(state.map_children(|original| { + if Rc::ptr_eq(original, original_input) { + Rc::clone(&input) + } else { + Rc::clone(original) + } + })?)); + } + // Derivation checks identical state schemas; downstream physical compilation + // remains responsible for the concrete family's merge capability. + let Ok(mut merged) = + OperatorNode::new(Operator::ASAP(ASAPOp::SummaryMerge { children: panes })) + else { + return Ok(None); + }; + merged.schema = state.schema.clone(); + merged.guarantee = state.guarantee.clone(); + Ok(Some(Rc::new(merged))) +} + +#[cfg(test)] +mod tests { + use super::*; + use asap_types::{ + ir::operator_properties::{Reduction, Source}, + ir::TimeRangeKind, + post_asap::{GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate}, + pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, + workload::{Predictability, Query, QueryRequirements, RepetitionInterval, TimeSelection}, + }; + + fn fixture() -> (Rc, QueryWorkloadEntry) { + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::TimeSeries { + metric: "events".into(), + }, + predicates: vec![], + schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + })) + .unwrap(); + let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: Duration::from_secs(300), + kind: TimeRangeKind::Range, + child: scan, + })) + .unwrap(); + let state = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: range, + family: FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ), + input: SummaryUpdate::column(ColumnRef::SampleValue), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + filter: None, + })) + .unwrap(); + let entry = QueryWorkloadEntry { + query: Query("quantile_over_time(0.99, events[5m])".into()), + recurrence: QueryRecurrence::Repeated(RepeatedDemand::FixedInterval( + RepetitionInterval(60_000), + )), + requirements: QueryRequirements::default(), + predictability: Predictability::AdHoc, + time_selection: TimeSelection::default(), + }; + (state, entry) + } + + /// Automatic composition covers the original interval once, without gaps or overlap. + #[test] + fn five_minute_window_has_original_and_five_exact_panes() { + let (state, entry) = fixture(); + let variants = enumerate_window_compositions(&state, &entry, 16).unwrap(); + assert_eq!(variants.len(), 2); + assert!(Rc::ptr_eq(&variants[0], &state)); + let Operator::ASAP(ASAPOp::SummaryMerge { children }) = &variants[1].operator else { + panic!("expected merge") + }; + assert_eq!(children.len(), 5); + for (index, pane) in children.iter().enumerate() { + let Operator::ASAP(ASAPOp::SummaryAgg { child, .. }) = &pane.operator else { + panic!() + }; + let Some(NonASAPOp::TimeRange { range, child, .. }) = child.non_asap() else { + panic!() + }; + assert_eq!(*range, Duration::from_secs(60)); + let Some(NonASAPOp::TimeShift { shift, .. }) = child.non_asap() else { + panic!() + }; + assert_eq!(shift.offset_ms, index as i64 * 60_000); + } + variants[1].validate_structure().unwrap(); + } + + /// Unknown cadence retains the original; an exhausted budget never returns a partial search. + #[test] + fn unknown_cadence_and_budget_are_explicit() { + let (state, mut entry) = fixture(); + assert!(matches!( + enumerate_window_compositions(&state, &entry, 4), + Err(WindowCompositionError::BudgetExceeded) + )); + entry.recurrence = QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }; + assert_eq!( + enumerate_window_compositions(&state, &entry, 4) + .unwrap() + .len(), + 1 + ); + } +} diff --git a/crates/asap-physical-operators/src/physical_planner/candidates.rs b/crates/asap-physical-operators/src/physical_planner/candidates.rs index 53ac6d14e..135d26f33 100644 --- a/crates/asap-physical-operators/src/physical_planner/candidates.rs +++ b/crates/asap-physical-operators/src/physical_planner/candidates.rs @@ -206,6 +206,36 @@ pub fn compile_candidates( } } +/// One enumerated frontier with its lowering result, including any rejection. +pub struct MaterializationCandidate { + pub frontier: Vec, + pub realization: Result, +} + +/// Automatically enumerate every legal materialization frontier and lower +/// its producer/reader split. Compile once, preserve individual cut failures, +/// and fail before returning an inventory if the exhaustive budget is exceeded. +/// An evaluator must price each complete split, including retained outputs. +pub fn compile_materialization_candidates( + dag: &PostAsapDAG, + inputs: BTreeMap, + roots: &[NodeId], + max_candidates: usize, +) -> Result, Error> { + let compiled = compile(dag, inputs, roots)?; + let frontiers = enumerate_compiled_frontiers(&compiled, max_candidates)?; + Ok(frontiers + .into_iter() + .map(|frontier| { + let candidate = cut_candidate(&compiled, &frontier); + MaterializationCandidate { + frontier, + realization: candidate, + } + }) + .collect()) +} + /// Complete workload cost supplied by scoped optimizer/deployment evidence. /// The evaluator includes build/update work, retained state, shared producers /// and recurrent reads over the same horizon; these are not per-query timings. diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index fb7af68fd..67ae02afe 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -39,8 +39,9 @@ pub mod promql_values; mod candidates; pub use candidates::{ - compile_candidate, compile_candidates, cut_candidate, enumerate_frontiers, - frontier_from_timing, select_candidate, CandidateCost, CandidateSelection, PhysicalASAPDAG, + compile_candidate, compile_candidates, compile_materialization_candidates, cut_candidate, + enumerate_frontiers, frontier_from_timing, select_candidate, CandidateCost, CandidateSelection, + PhysicalASAPDAG, }; mod compiled; @@ -399,7 +400,11 @@ fn compile_internal( "per-entity summary requires a resolved raw time range", )); }; - let Some(NonASAPOp::Scan { schema, .. }) = child.non_asap() else { + let source = match child.non_asap() { + Some(NonASAPOp::TimeShift { child, .. }) => child, + _ => child, + }; + let Some(NonASAPOp::Scan { schema, .. }) = source.non_asap() else { return Err(invalid("per-entity summary requires a resolved source")); }; if !schema.closed || update.item.is_some() { diff --git a/crates/integration-tests/tests/automatic_window_composition.rs b/crates/integration-tests/tests/automatic_window_composition.rs new file mode 100644 index 000000000..d7fad8f1a --- /dev/null +++ b/crates/integration-tests/tests/automatic_window_composition.rs @@ -0,0 +1,187 @@ +//! Automatically generated panes execute directly and through an enumerated materialization frontier. +mod physical_common; +use asap_physical_operators::{ + physical_planner::{compile_materialization_candidates, InputContract}, + runtime::Scope, + values::{Batch, Value}, +}; +use asap_types::{ + ir::operator_properties::{Reduction, Source}, + ir::{ASAPOp, NonASAPOp, Operator, OperatorNode}, + post_asap::{ + GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SketchStatistic, SummaryUpdate, + }, + pre_asap::{ColumnRef, DataType, Field, FieldDataType, Schema}, +}; +use std::{collections::BTreeMap, sync::Arc}; + +/// Generate panes from one five-minute state and find its producer/reader split automatically. +#[test] +fn automatic_panes_and_materialization_preserve_quantile() { + let family = FieldDataType::Sketch( + SketchKind::new(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Default::default(), + ); + let scan = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan { + source: Source::TimeSeries { + metric: "events".into(), + }, + predicates: vec![], + schema: Schema::new(vec![Field::plain("value", DataType::Float64, false)]), + })) + .unwrap(); + let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { + range: std::time::Duration::from_secs(300), + kind: asap_types::ir::TimeRangeKind::Range, + child: scan, + })) + .unwrap(); + let state = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: range, + family, + input: SummaryUpdate::column(ColumnRef::SampleValue), + reduction: Reduction::by(vec![]), + grouping: GroupingStrategy::default(), + filter: None, + })) + .unwrap(); + let entry = asap_types::workload::QueryWorkloadEntry { + query: asap_types::workload::Query("quantile_over_time(0.99, events[5m])".into()), + recurrence: asap_types::workload::QueryRecurrence::Repeated( + asap_types::workload::RepeatedDemand::FixedInterval( + asap_types::workload::RepetitionInterval(60_000), + ), + ), + requirements: Default::default(), + predictability: asap_types::workload::Predictability::AdHoc, + time_selection: Default::default(), + }; + let variants = + asap_aware_mapping::window_composition::enumerate_window_compositions(&state, &entry, 64) + .unwrap(); + assert_eq!(variants.len(), 2); + let merged = variants[1].clone(); + let root = OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryEstimate { + summary_input: merged, + query: SketchStatistic::Quantile { q: 0.99 }, + })) + .unwrap(); + root.validate_structure().unwrap(); + let wire = physical_common::compile_post_asap_dag(&root).unwrap(); + wire.validate().unwrap(); + let mut inputs = BTreeMap::new(); + let mut batches = BTreeMap::new(); + for (pane, node) in wire + .nodes + .iter() + .filter(|node| { + matches!( + node.payload, + asap_types::ir::export::PostAsapOperatorPayload::Relational { + operator: asap_types::ir::export::NonASAPOpKind::TimeRange { .. } + } + ) + }) + .enumerate() + { + let schema = Arc::new(node.output_schema.clone()); + let id = u64::from(node.id.0); + inputs.insert(id, InputContract::bounded(schema.clone())); + batches.insert( + id, + Batch::try_new( + schema, + (0..20) + .map(|i| vec![Value::Float64((pane * 20 + i) as f64)]) + .collect(), + ) + .unwrap(), + ); + } + let choices = + compile_materialization_candidates(&wire, inputs.clone(), &[u64::from(wire.root.0)], 1024) + .unwrap(); + assert!( + compile_materialization_candidates(&wire, inputs, &[u64::from(wire.root.0)], 1).is_err() + ); + let plan = &choices + .iter() + .find(|candidate| candidate.frontier.is_empty()) + .unwrap() + .realization + .as_ref() + .unwrap() + .query; + let outputs = physical_common::execute( + plan, + batches.clone(), + Scope::Query { + evaluation_time_ms: 300_000, + revision: 1, + }, + ); + // Pattern B1 stores five built panes; its query reads the same merge DAG + // as Pattern B2, which builds every pane from raw rows on each read. + let frontier: Vec<_> = wire + .nodes + .iter() + .filter_map(|node| { + matches!( + node.payload, + asap_types::ir::export::PostAsapOperatorPayload::SummaryAgg { .. } + ) + .then_some(u64::from(node.id.0)) + }) + .collect(); + let materialized = choices + .iter() + .find(|candidate| { + let mut candidate = candidate.frontier.clone(); + let mut frontier = frontier.clone(); + candidate.sort(); + frontier.sort(); + candidate == frontier + }) + .unwrap() + .realization + .as_ref() + .unwrap(); + let precompute = materialized.precompute.as_ref().unwrap(); + let stored = physical_common::execute( + precompute, + batches, + Scope::Ingestion { + window_start_ms: 0, + window_end_ms: 300_000, + revision: 1, + }, + ); + let retained = precompute + .roots() + .iter() + .copied() + .zip(stored) + .map(|(id, mut batches)| { + assert_eq!(batches.len(), 1); + (id, batches.remove(0)) + }) + .collect(); + let retained_outputs = physical_common::execute( + &materialized.query, + retained, + Scope::Query { + evaluation_time_ms: 300_000, + revision: 1, + }, + ); + assert_eq!(retained_outputs[0][0].rows().len(), 1); + assert!( + matches!(retained_outputs[0][0].rows()[0].as_slice(), [Value::Float64(value)] if *value == 98.0) + ); + assert_eq!(outputs[0][0].rows().len(), 1); + assert!( + matches!(outputs[0][0].rows()[0].as_slice(), [Value::Float64(value)] if *value == 98.0), + "{:?}", + outputs[0][0].rows() + ); +} diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md index eb6008d40..2acee4fe1 100644 --- a/docs/develop_docs/planner-layering-status.md +++ b/docs/develop_docs/planner-layering-status.md @@ -9,7 +9,7 @@ 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. | -| Window composition | Mergeable state IR/native merge exists; physical pane compatibility and reuse cost helpers exist. | Automatic logical sliding/tumbling/EH alternatives over differing windows, boundary coverage and error proofs are absent. A merge kernel alone does not implement Examples 1/3. | +| 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. | | 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. | @@ -108,3 +108,21 @@ frequency rewrite, comparing filtering, natural-log units, empty-input NULL and unequal frequencies. Raw analytical cost lowering still excludes this unordered SUM window; supplying native execution does not provide missing cost evidence or extend the analytical adapter's existing ordered-window rule. + +## Automatic window/materialization acceptance + +`window_composition::enumerate_window_compositions` derives disjoint relative +panes from a summary's lookback and workload cadence. Five minutes every minute +produces the original state plus a five-pane merge. Non-divisible lookbacks use +the greatest common divisor, so coverage is exact. Source predicates, grouping, +summary parameters and selector offsets/anchors stay attached to the panes. +Unknown recurrence keeps the original; budget exhaustion returns an error. + +`compile_materialization_candidates` lowers once and enumerates every legal +producer/reader frontier, including rebuilding all panes at query time. It +retains individual failures and never returns a truncated inventory. +`integration-tests/tests/automatic_window_composition.rs` verifies generated +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.