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
266 changes: 262 additions & 4 deletions crates/asap-aware-mapping/src/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,18 @@
//! 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::collections::{BTreeMap, HashMap, 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::ir::summary_coverage::{CoverageRegion, SummaryCoverage};
use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, QueryRoot, SchemaDerivationError};
use asap_types::post_asap::{
EntityIdentity, ExactKind, ExactParams, FieldDataType, GroupingStrategy,
NonNegativeWeightProof, SketchAlgorithm, SketchKind, SketchStatistic, SummaryInputExpr,
SummaryUpdate, WeightDomain,
};
use asap_types::pre_asap::expr_ir::ColumnRef;
use asap_types::pre_asap::{AggIntent, Reduction, Schema};
use asap_types::types::AccuracyTarget;
use thiserror::Error;

Expand Down Expand Up @@ -40,6 +46,10 @@ pub enum LogicalCandidateError {
AssignedTiming,
#[error("approximate accuracy requires finite positive epsilon and delta in (0, 1)")]
InvalidAccuracy,
#[error("choice must name one listed alternative per target")]
InvalidChoice,
#[error("unsupported local realization: {0}")]
Unsupported(&'static str),
}

/// Enumerate exact and summary choices in stable catalog order, without ranking
Expand Down Expand Up @@ -119,6 +129,254 @@ pub fn enumerate_local_logical_candidates<Id>(
Ok(LocalLogicalCandidates { roots, targets })
}

/// Build one whole-workload candidate (#509 Stage 1): `choice[i]` indexes
/// `inventory.targets[i].alternatives`. Each chosen non-pass-through target is
/// replaced by `SummaryAgg` followed by `SummaryEstimate` (sketch) or
/// `FinalizeExactAccumulator` (exact accumulator). Untouched sub-DAGs keep
/// their identity, so sharing between roots is preserved.
pub fn compose_logical_candidate<Id: Clone>(
inventory: &LocalLogicalCandidates<Id>,
choice: &[usize],
) -> Result<Vec<(Id, QueryRoot)>, LogicalCandidateError> {
if choice.len() != inventory.targets.len() {
return Err(LogicalCandidateError::InvalidChoice);
}
let chosen = inventory
.targets
.iter()
.zip(choice)
.map(|(target, &index)| {
target
.alternatives
.get(index)
.map(|alternative| (Rc::as_ptr(&target.target), alternative))
.ok_or(LogicalCandidateError::InvalidChoice)
})
.collect::<Result<HashMap<_, _>, _>>()?;
let mut memo = HashMap::new();
inventory
.roots
.iter()
.map(|(id, root)| {
let root = match root {
QueryRoot::Operator(node) => {
QueryRoot::Operator(rewrite(node, &chosen, &mut memo)?)
}
QueryRoot::Scalar(expr) => {
for node in expr.operator_refs() {
rewrite(node, &chosen, &mut memo)?;
}
QueryRoot::Scalar(
expr.map_operator_refs(&mut |node| memo[&Rc::as_ptr(node)].clone()),
)
}
};
Ok((id.clone(), root))
})
.collect()
}

type Memo = HashMap<*const OperatorNode, Rc<OperatorNode>>;

fn rewrite(
node: &Rc<OperatorNode>,
chosen: &HashMap<*const OperatorNode, &Realization>,
memo: &mut Memo,
) -> Result<Rc<OperatorNode>, LogicalCandidateError> {
if let Some(done) = memo.get(&Rc::as_ptr(node)) {
return Ok(done.clone());
}
for child in node.children() {
rewrite(child, chosen, memo)?;
}
let changed = node
.children()
.iter()
.any(|child| !Rc::ptr_eq(child, &memo[&Rc::as_ptr(child)]));
let rebuilt = match chosen.get(&Rc::as_ptr(node)) {
Some(realization) if **realization != Realization::PassThrough => {
realize(node, realization, memo)?
}
_ if changed => Rc::new(node.map_children(|child| memo[&Rc::as_ptr(child)].clone())?),
_ => node.clone(),
};
memo.insert(Rc::as_ptr(node), rebuilt.clone());
Ok(rebuilt)
}

fn realize(
target: &OperatorNode,
realization: &Realization,
memo: &Memo,
) -> Result<Rc<OperatorNode>, LogicalCandidateError> {
let Some(NonASAPOp::Aggregate {
child,
reduction,
measures,
filters,
having,
..
}) = target.non_asap()
else {
return Err(LogicalCandidateError::Unsupported(
"target is not an aggregate",
));
};
let [intent] = measures.as_slice() else {
return Err(LogicalCandidateError::Unsupported(
"multi-measure aggregate",
));
};
if !filters.is_empty() || having.is_some() {
return Err(LogicalCandidateError::Unsupported(
"filtered or HAVING aggregate",
));
}
let child = memo[&Rc::as_ptr(child)].clone();
let (family, query) = match realization {
Realization::ExactAggregate { kind, params } => (
FieldDataType::ExactAggregate(kind.clone(), params.clone()),
None,
),
Realization::Sketch(kind) => (
FieldDataType::Sketch(kind.clone(), GroupingStrategy::default()),
Some(statistic(intent)?),
),
_ => return Err(LogicalCandidateError::Unsupported("summary family")),
};
let input = summary_update(intent, &family, reduction, &child.schema)?;
let state = OperatorNode::new(Operator::ASAP(ASAPOp::SummaryAgg {
child: child.clone(),
family,
input,
reduction: reduction.clone(),
grouping: GroupingStrategy::default(),
filter: None,
}))?;
// Whole-source coverage is declared, not proven: Pass 1 trusts that the
// state holds every observation of its source that reaches it (#570).
let coverage = SummaryCoverage {
source: single_source(&child)?,
regions: vec![CoverageRegion {
time_ms: None,
population: BTreeMap::new(),
}],
};
let state = Rc::new(state.with_coverage(coverage)?);
let evaluation = match query {
Some(query) => ASAPOp::SummaryEstimate {
summary_input: state,
query,
},
None => ASAPOp::FinalizeExactAccumulator { child: state },
};
Ok(OperatorNode::new_shared(Operator::ASAP(evaluation))?)
}

fn single_source(
node: &Rc<OperatorNode>,
) -> Result<asap_types::pre_asap::Source, LogicalCandidateError> {
let mut sources = Vec::new();
for node in OperatorNode::reachable(node) {
if let Some(NonASAPOp::Scan { source, .. }) = node.non_asap() {
if !sources.contains(source) {
sources.push(source.clone());
}
}
}
match <[_; 1]>::try_from(sources) {
Ok([source]) => Ok(source),
Err(_) => Err(LogicalCandidateError::Unsupported(
"summary coverage needs exactly one source",
)),
}
}

/// What each input row contributes, following the legacy realization rules:
/// heap sketches rank series identities, frequency sketches count values, and
/// every other family reads the measure's input column.
fn summary_update(
intent: &AggIntent,
family: &FieldDataType,
reduction: &Reduction,
child: &Schema,
) -> Result<SummaryUpdate, LogicalCandidateError> {
let algorithm = match family {
FieldDataType::Sketch(kind, _) => Some(kind.algorithm()),
_ => None,
};
let weight = crate::replacement::summarised_input(intent, child)
.map_err(|_| LogicalCandidateError::Unsupported("input column outside child schema"))?;
Ok(match (intent, algorithm) {
(AggIntent::TopK { .. }, Some(_)) => {
// SQL rows carry no implicit series identity to rank; closed
// PromQL rows carry it as a column.
if child.closed && !child.has_promql_series_identity() {
return Err(LogicalCandidateError::Unsupported(
"Top-K item identity for closed schemas",
));
}
if reduction.group_keys().is_some_and(|keys| keys.is_without()) {
return Err(LogicalCandidateError::Unsupported(
"Top-K partitions given by `without`",
));
}
let excluding = reduction
.group_keys()
.into_iter()
.flat_map(|keys| keys.iter())
.filter_map(|&index| child.fields.get(index))
.map(crate::replacement::column_ref)
.collect();
// Rows that carry the full series identity rank it as a column,
// the item form the runtime builds keyed summaries from.
let item = if child.has_promql_series_identity() {
SummaryInputExpr::Column(ColumnRef::Named(
asap_types::pre_asap::schema::PROMQL_SERIES_IDENTITY.into(),
))
} else {
SummaryInputExpr::EntityIdentity(EntityIdentity::PromqlLabelSet { excluding })
};
SummaryUpdate {
item: Some(item),
weight: SummaryInputExpr::Column(ColumnRef::SampleValue),
// Not proven non-negative; selection decides whether CMS is admissible.
weight_domain: WeightDomain::UnknownOrSigned,
}
}
(_, Some(SketchAlgorithm::UnivMon))
| (AggIntent::Count { .. }, Some(SketchAlgorithm::Cms | SketchAlgorithm::CountSketch)) => {
SummaryUpdate {
item: Some(weight),
weight: SummaryInputExpr::Constant(1.0),
weight_domain: WeightDomain::NonNegative {
proof: NonNegativeWeightProof::UnitCount,
},
}
}
_ => SummaryUpdate {
item: None,
weight,
weight_domain: WeightDomain::UnknownOrSigned,
},
})
}

fn statistic(intent: &AggIntent) -> Result<SketchStatistic, LogicalCandidateError> {
Ok(match intent {
AggIntent::Quantile { q, .. } => SketchStatistic::Quantile { q: *q },
AggIntent::Cardinality { .. } => SketchStatistic::Cardinality,
AggIntent::FrequencyL2 { .. } => SketchStatistic::FrequencyL2,
AggIntent::FrequencyEntropy { .. } => SketchStatistic::FrequencyEntropy,
AggIntent::TopK { k, .. } => SketchStatistic::TopK { k: *k },
AggIntent::Count { .. } => SketchStatistic::PointCount {
key: ColumnRef::SampleValue,
value: None,
},
_ => return Err(LogicalCandidateError::Unsupported("sketch evaluation")),
})
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
4 changes: 2 additions & 2 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3822,7 +3822,7 @@ fn summarised_column(intent: &AggIntent, child_schema: &Schema) -> ColumnRef {
}
}

fn column_ref(column: &Field) -> ColumnRef {
pub(crate) fn column_ref(column: &Field) -> ColumnRef {
match &column.table {
Some(t) => ColumnRef::Qualified {
table: t.clone(),
Expand All @@ -3840,7 +3840,7 @@ fn column_ref(column: &Field) -> ColumnRef {
/// A tuple leg outside the child schema is an error rather than
/// [`summarised_column`]'s sample-value fallback: a leg has no sample-value
/// reading, and silently dropping one would under-count.
fn summarised_input(
pub(crate) fn summarised_input(
intent: &AggIntent,
child_schema: &Schema,
) -> Result<SummaryInputExpr, RealizationError> {
Expand Down
38 changes: 38 additions & 0 deletions crates/asap-aware-mapping/tests/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -241,3 +241,41 @@ fn topk_keeps_both_specialized_heap_choices() {
]
);
}

/// A workload candidate replaces only chosen targets, declares whole-source
/// coverage on the summary, and keeps unchosen plans identical.
#[test]
fn composed_candidate_replaces_chosen_target_with_summary_evaluation() {
use asap_aware_mapping::logical_candidates::compose_logical_candidate;
use asap_types::ir::ASAPOp;
let producer = aggregate(AggIntent::Cardinality {
cols: vec![0],
accuracy: approximate(),
});
let inventory =
enumerate_local_logical_candidates(vec![(0, QueryRoot::Operator(producer.clone()))])
.unwrap();
let exact = compose_logical_candidate(&inventory, &[0]).unwrap();
assert!(matches!(&exact[0].1, QueryRoot::Operator(node) if Rc::ptr_eq(node, &producer)));

let hll = compose_logical_candidate(&inventory, &[1]).unwrap();
let QueryRoot::Operator(estimate) = &hll[0].1 else {
panic!("operator root expected")
};
let Some(ASAPOp::SummaryEstimate { summary_input, .. }) = estimate.asap() else {
panic!("summary evaluation expected")
};
let coverage = summary_input.coverage.as_ref().unwrap();
assert_eq!(
coverage.source,
Source::Table {
table_ref: "flows".into()
}
);
assert!(coverage.regions[0].time_ms.is_none() && coverage.regions[0].population.is_empty());
asap_types::ir::export::compile_logical_asap_workload(&[hll[0].1.clone()])
.unwrap()
.validate()
.unwrap();
assert!(compose_logical_candidate(&inventory, &[99]).is_err());
}
Loading