Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions crates/asap-aware-mapping/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -211,3 +211,7 @@ pub use rewrite::{AvgToSumOverCountStrategy, SemanticEquivalentRewriteStrategy};
pub use topk_reuse::TopKLimitReuseStrategy;

pub mod maintained_population;

/// Local candidate generation over the unified IR. No execution timing is
/// assigned: that is a Stage 2 materialization decision.
pub mod logical_candidates;
147 changes: 147 additions & 0 deletions crates/asap-aware-mapping/src/logical_candidates.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,147 @@
//! Pass 1 local alternatives over the unified logical IR.
//!
//! Alternatives are nominal realization descriptors attached to their original
//! target, not ranked plans or accuracy certificates. Workload composition and
//! physical planning consume this inventory later; empirical models belong to
//! selection. The legacy search API remains until planner cutover.
use std::collections::HashSet;
use std::rc::Rc;

use asap_types::ir::{NonASAPOp, OperatorNode, QueryRoot, SchemaDerivationError};
use asap_types::post_asap::{ExactKind, ExactParams, SketchKind};
use asap_types::pre_asap::AggIntent;
use asap_types::types::AccuracyTarget;
use thiserror::Error;

use crate::replacement::{
accuracy_budget, accuracy_target, default_size_params, summary_candidates, Realization,
};

/// All local realizations of one single-measure aggregate. The target retains
/// source, grouping, filters, input expressions and evaluation context.
#[derive(Debug, Clone)]
pub struct LocalLogicalTarget {
pub target: Rc<OperatorNode>,
pub alternatives: Vec<Realization>,
}

/// Compact Pass 1 inventory; roots and nested producer dependencies are retained.
#[derive(Debug, Clone)]
pub struct LocalLogicalCandidates<Id> {
pub roots: Vec<(Id, QueryRoot)>,
pub targets: Vec<LocalLogicalTarget>,
}

#[derive(Debug, Error)]
pub enum LogicalCandidateError {
#[error(transparent)]
Structure(#[from] SchemaDerivationError),
#[error("logical candidate input already has assigned execution timing")]
AssignedTiming,
#[error("approximate accuracy requires finite positive epsilon and delta in (0, 1)")]
InvalidAccuracy,
}

/// Enumerate exact and summary choices in stable catalog order, without ranking
/// or empirical assessment. Parameters are candidate dimensions, not a claim
/// that a deployment meets the request's accuracy requirement.
pub fn local_realizations_for_intent(
intent: &AggIntent,
) -> Result<Vec<Realization>, LogicalCandidateError> {
let mut choices = vec![Realization::PassThrough];
let exact = match intent {
AggIntent::Count { .. } => Some((ExactKind::Count, ExactParams::Count)),
AggIntent::Sum { .. } => Some((ExactKind::Sum, ExactParams::Sum)),
AggIntent::Min { .. } => Some((ExactKind::Min, ExactParams::Min)),
AggIntent::Max { .. } => Some((ExactKind::Max, ExactParams::Max)),
AggIntent::Rate => Some((ExactKind::Rate, ExactParams::Rate)),
AggIntent::IRate => Some((ExactKind::IRate, ExactParams::IRate)),
AggIntent::Increase => Some((ExactKind::Increase, ExactParams::Increase)),
_ => None,
};
if let Some((kind, params)) = exact {
choices.push(Realization::ExactAggregate { kind, params });
}
if let Some(target) = accuracy_target(intent) {
if *target != AccuracyTarget::Exact {
let (epsilon, delta) = accuracy_budget(target);
if !epsilon.is_finite()
|| epsilon <= 0.0
|| !delta.is_finite()
|| !(0.0..1.0).contains(&delta)
|| delta == 0.0
{
return Err(LogicalCandidateError::InvalidAccuracy);
}
for algorithm in summary_candidates(intent) {
choices.push(Realization::Sketch(SketchKind::new(
algorithm.clone(),
default_size_params(algorithm.clone(), intent, epsilon, delta),
)));
}
}
}
Ok(choices)
}

/// Discover single-measure targets, including operator plans read by scalar roots.
/// Multi-measure aggregates remain intact pending a semantics-preserving split.
pub fn enumerate_local_logical_candidates<Id>(
roots: Vec<(Id, QueryRoot)>,
) -> Result<LocalLogicalCandidates<Id>, LogicalCandidateError> {
let mut seen = HashSet::new();
let mut targets = Vec::new();
for (_, root) in &roots {
root.validate_structure()?;
let operators = match root {
QueryRoot::Operator(node) => vec![node],
QueryRoot::Scalar(expr) => expr.operator_refs(),
};
for root in operators {
for node in OperatorNode::reachable(root) {
if !seen.insert(Rc::as_ptr(&node)) {
continue;
}
if node.timing.is_some() {
return Err(LogicalCandidateError::AssignedTiming);
}
if let Some(NonASAPOp::Aggregate { measures, .. }) = node.non_asap() {
if let [intent] = measures.as_slice() {
targets.push(LocalLogicalTarget {
alternatives: local_realizations_for_intent(intent)?,
target: node,
});
}
}
}
}
}
Ok(LocalLogicalCandidates { roots, targets })
}

#[cfg(test)]
mod tests {
use super::*;
/// Approximate requests must retain the exact execution alternative too.
#[test]
fn approximate_count_keeps_exact_and_universal_choices() {
let choices = local_realizations_for_intent(&AggIntent::Count {
accuracy: AccuracyTarget::EpsilonDelta {
epsilon: 0.05,
delta: 0.01,
},
})
.unwrap();
assert!(choices
.iter()
.any(|choice| matches!(choice, Realization::PassThrough)));
assert!(choices.iter().any(|choice| matches!(
choice,
Realization::ExactAggregate {
kind: ExactKind::Count,
..
}
)));
assert!(choices.iter().any(|choice| matches!(choice, Realization::Sketch(kind) if *kind.algorithm() == asap_types::post_asap::SketchAlgorithm::UnivMon)));
}
}
243 changes: 243 additions & 0 deletions crates/asap-aware-mapping/tests/logical_candidates.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,243 @@
//! Frontend-to-Pass-1 acceptance: candidate discovery precedes empirical selection.
use asap_aware_mapping::{
logical_candidates::{
enumerate_local_logical_candidates, local_realizations_for_intent, LogicalCandidateError,
},
Realization,
};
use asap_types::{
ir::operator_properties::{Reduction, Source},
ir::{NonASAPOp, Operator, OperatorNode, QueryRoot, ScalarExpr},
post_asap::{ExactKind, SketchAlgorithm},
pre_asap::{AggIntent, DataType, Field, Schema},
types::AccuracyTarget,
};
use std::rc::Rc;

fn approximate() -> AccuracyTarget {
AccuracyTarget::EpsilonDelta {
epsilon: 0.05,
delta: 0.01,
}
}
fn algorithms(choices: &[Realization]) -> Vec<SketchAlgorithm> {
choices
.iter()
.filter_map(|choice| match choice {
Realization::Sketch(kind) => Some(kind.algorithm().clone()),
_ => None,
})
.collect()
}
fn aggregate(intent: AggIntent) -> Rc<OperatorNode> {
let child = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Scan {
source: Source::Table {
table_ref: "flows".into(),
},
predicates: vec![],
schema: Schema::lifted(vec![Field::plain("src_ip", DataType::Utf8, false)], None),
}))
.unwrap();
OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::Aggregate {
child,
reduction: Reduction::by(vec![]),
measures: vec![intent],
output_names: vec![],
filters: vec![],
having: None,
}))
.unwrap()
}

/// Example 2 preserves specialized distinct summaries and the universal alternative.
#[test]
fn cardinality_keeps_exact_specialized_and_universal_alternatives() {
let choices = local_realizations_for_intent(&AggIntent::Cardinality {
cols: vec![0],
accuracy: approximate(),
})
.unwrap();
assert!(matches!(choices[0], Realization::PassThrough));
assert_eq!(
algorithms(&choices),
vec![
SketchAlgorithm::Hll,
SketchAlgorithm::Theta,
SketchAlgorithm::Kmv,
SketchAlgorithm::UnivMon
]
);
let tuple = local_realizations_for_intent(&AggIntent::Cardinality {
cols: vec![0, 1],
accuracy: approximate(),
})
.unwrap();
assert!(!algorithms(&tuple).contains(&SketchAlgorithm::UnivMon));
}

/// Frequency moments retain exact execution and a universal sketch without certification.
#[test]
fn frequency_statistics_keep_universal_choices() {
for intent in [
AggIntent::FrequencyL2 {
col: Some(0),
accuracy: approximate(),
},
AggIntent::FrequencyEntropy {
col: Some(0),
accuracy: approximate(),
},
] {
let choices = local_realizations_for_intent(&intent).unwrap();
assert!(matches!(choices[0], Realization::PassThrough));
assert_eq!(algorithms(&choices), vec![SketchAlgorithm::UnivMon]);
}
}

/// An exact request cannot acquire an approximate sketch merely because one is available.
#[test]
fn exact_quantile_stays_exact_and_approximate_keeps_both_families() {
let choices = local_realizations_for_intent(&AggIntent::Quantile {
col: Some(0),
q: 0.99,
accuracy: approximate(),
})
.unwrap();
assert_eq!(
algorithms(&choices),
vec![SketchAlgorithm::Kll, SketchAlgorithm::DDSketch]
);
let exact = local_realizations_for_intent(&AggIntent::Quantile {
col: Some(0),
q: 0.99,
accuracy: AccuracyTarget::Exact,
})
.unwrap();
assert_eq!(exact, vec![Realization::PassThrough]);
}

/// Scalar roots expose their producer targets; repeated references retain one target identity.
#[test]
fn scalar_root_producers_are_discovered_once() {
let producer = aggregate(AggIntent::Cardinality {
cols: vec![0],
accuracy: approximate(),
});
let roots = vec![
(
"scalar",
QueryRoot::Scalar(ScalarExpr::ScalarSubquery(producer.clone())),
),
("relation", QueryRoot::Operator(producer.clone())),
];
let candidates = enumerate_local_logical_candidates(roots).unwrap();
assert_eq!(candidates.roots.len(), 2);
for (_, root) in &candidates.roots {
asap_types::ir::export::compile_logical_asap_query(root)
.unwrap()
.validate()
.unwrap();
}
assert_eq!(candidates.targets.len(), 1);
assert!(Rc::ptr_eq(&candidates.targets[0].target, &producer));
assert!(producer.timing.is_none());
assert!(producer.guarantee.is_none());
asap_types::ir::export::compile_logical_asap_dag(&producer)
.unwrap()
.validate()
.unwrap();
}

/// Example 1 rate lowering reaches the exact accumulator choice without a cost model.
#[test]
fn promql_lowering_reaches_local_candidates_without_execution_timing() {
use asap_types::workload::{
AccuracyRequirement, BatchEntry, PlanningWorkload, Query, QueryLanguage, QueryRequirements,
QueryWorkload,
};
let workload = PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: Some(vec![BatchEntry {
query: Query("sum by (job) (rate(http_requests_total[1m]))".into()),
requirements: QueryRequirements {
accuracy: AccuracyRequirement::Explicit(approximate()),
..Default::default()
},
predictability: Default::default(),
invocations: 1,
execute_at: None,
time_selection: Default::default(),
}]),
repeating_queries: None,
},
data_workload: Some(asap_types::workload::DataWorkload {
data_ingestion_interval: asap_types::workload::Evidence {
value: Some(asap_types::workload::DurationMs(1000)),
..Default::default()
},
..Default::default()
}),
};
let roots = asap_frontend_promql::lower_promql_query_workload(&workload, 0).unwrap();
let candidates =
enumerate_local_logical_candidates(roots.into_iter().enumerate().collect()).unwrap();
assert!(candidates
.targets
.iter()
.any(|target| target.alternatives.iter().any(|choice| matches!(
choice,
Realization::ExactAggregate {
kind: ExactKind::Rate,
..
}
))));
assert!(candidates
.targets
.iter()
.all(|target| target.target.timing.is_none()));
}

/// Physical annotations and invalid probability requirements fail at the stage boundary.
#[test]
fn assigned_timing_and_invalid_accuracy_are_rejected() {
let mut producer = (*aggregate(AggIntent::Count {
accuracy: approximate(),
}))
.clone();
producer.timing = Some(asap_types::post_asap::ExecutionTiming::QueryTime);
assert!(matches!(
enumerate_local_logical_candidates(vec![(0, QueryRoot::Operator(Rc::new(producer)))]),
Err(LogicalCandidateError::AssignedTiming)
));
for target in [
AccuracyTarget::Epsilon(f64::NAN),
AccuracyTarget::EpsilonDelta {
epsilon: 0.1,
delta: 0.0,
},
] {
assert!(matches!(
local_realizations_for_intent(&AggIntent::Count { accuracy: target }),
Err(LogicalCandidateError::InvalidAccuracy)
));
}
}

/// Local TopK keeps both declared heap substrates without choosing an implementation.
#[test]
fn topk_keeps_both_specialized_heap_choices() {
let choices = local_realizations_for_intent(&AggIntent::TopK {
k: 10,
accuracy: approximate(),
})
.unwrap();
assert!(matches!(choices[0], Realization::PassThrough));
assert_eq!(
algorithms(&choices),
vec![
SketchAlgorithm::CmsWithHeap,
SketchAlgorithm::CountSketchWithHeap
]
);
}
Loading