Skip to content
Open
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
193 changes: 81 additions & 112 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1356,114 +1356,6 @@ impl<'a> SketchAlgorithmStrategy<'a> {
self.propose_with(&ranked, None)
}

/// Fixed-window maintenance can finalize each series' counter state and
/// build a fresh heap or grouped Sum for that evaluation window. Deployment must provide
/// a complete, synchronized population and bind the matching window; this
/// candidate never incrementally adds one window's rates to another.
pub fn fixed_window_rate_candidates(&self, root: &Rc<QueryExpr>) -> Proposals {
fn place(node: &Rc<SummaryNode>) -> Option<Rc<SummaryNode>> {
let mut next = node.as_ref().clone();
match &mut next.expr {
SummaryExpr::ValueOperation {
child,
operation: ValueOperation::FinalizeExactAccumulator,
timing,
} if matches!(&child.expr, SummaryExpr::SummaryAgg {
family: SummaryFamilyType::ExactAggregate(ExactKind::Rate, _),
reduction: Reduction::PerEntity, child: source, ..
} if matches!(&source.expr, SummaryExpr::KeepPreAsap(source) if matches!(source.as_ref(), QueryExpr::TimeRange { .. }))) =>
{
*timing = ExecutionTiming::IngestionTime;
}
SummaryExpr::ValueOperation { child, .. }
| SummaryExpr::SummaryAgg { child, .. } => *child = place(child)?,
SummaryExpr::SummaryEstimate { summary_input, .. } => {
*summary_input = place(summary_input)?
}
_ => return None,
}
Some(Rc::new(next))
}
let mut proposals = self.propose_with(root, None);
proposals.candidates.retain_mut(|candidate| {
let Replacement::Summary(node) = &candidate.replacement else {
return false;
};
let Ok(dag) = asap_types::post_asap::compile_post_asap_dag(node) else {
return false;
};
if !dag.nodes.iter().any(|node| match &node.payload {
asap_types::post_asap::PostAsapOperatorPayload::SummaryAgg {
family: SummaryFamilyType::Sketch(kind, _),
..
} => matches!(
kind.algorithm(),
SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap
),
asap_types::post_asap::PostAsapOperatorPayload::SummaryAgg {
family: SummaryFamilyType::ExactAggregate(ExactKind::Sum, _),
..
} => true,
_ => false,
}) {
return false;
}
let Some(placed) = place(node) else {
return false;
};
if asap_types::post_asap::compile_post_asap_dag(&placed).is_err() {
return false;
}
let Ok(placed) = finalize_query_candidate(placed, root) else {
return false;
};
candidate.replacement = Replacement::Summary(placed);
candidate
.rationale
.push_str("; fixed-window precompute over complete per-series counter states");
true
});
proposals
}

/// Retain grouped Sum after a per-series Rate readout as a query-time
/// candidate alongside its complete-window maintenance placement.
pub fn query_time_rate_aggregation_candidates(&self, root: &Rc<QueryExpr>) -> Proposals {
fn query_time(node: &Rc<SummaryNode>) -> Rc<SummaryNode> {
let mut next = node.as_ref().clone();
match &mut next.expr {
SummaryExpr::ValueOperation {
child,
operation: ValueOperation::FinalizeExactAccumulator,
timing,
} if matches!(
&child.expr,
SummaryExpr::SummaryAgg {
family: SummaryFamilyType::ExactAggregate(ExactKind::Rate, _),
..
}
) =>
{
*timing = ExecutionTiming::QueryTime;
}
SummaryExpr::ValueOperation { child, .. }
| SummaryExpr::SummaryAgg { child, .. } => *child = query_time(child),
_ => {}
}
Rc::new(next)
}
let mut proposals = self.fixed_window_rate_candidates(root);
proposals.candidates.retain_mut(|candidate| {
let Replacement::Summary(node) = &candidate.replacement else { return false };
if !matches!(&node.expr, SummaryExpr::ValueOperation { child, operation: ValueOperation::FinalizeExactAccumulator, .. }
if matches!(&child.expr, SummaryExpr::SummaryAgg { family: SummaryFamilyType::ExactAggregate(ExactKind::Sum, _), .. })) { return false; }
candidate.replacement = Replacement::Summary(query_time(node));
candidate.rationale = "query-time grouped Sum over complete per-series Rate readouts".into();
true
});
proposals
}

pub(crate) fn from_planning_inputs(planning_inputs: CandidatePlanningInputs<'a>) -> Self {
Self { planning_inputs }
}
Expand Down Expand Up @@ -2764,7 +2656,9 @@ fn realize_physical_summary_input(
/// Emit `SummaryAgg` (recursively binding the child), plus the
/// `SummaryEstimate` readout when `estimate` is set.
// Retain the exact expression and schema while placing its value production
// on the update path. Read-time consumers keep their original shared nodes.
// on the update path. This is the initial layout for values feeding a summary;
// lifecycle timing is authoritative. Read-time consumers keep their original
// shared nodes.
fn maintenance_exact_values(node: Rc<SummaryNode>) -> Option<Rc<SummaryNode>> {
let expr = match &node.expr {
// These guards can fall back at read time, but cannot recover a parent
Expand Down Expand Up @@ -2993,8 +2887,9 @@ fn construct_summary_agg(
};
Rc::clone(child)
} else if snapshot_weighted {
// A fresh query-time summary consumes this evaluation's finalized rates.
// Moving rate snapshots must never accumulate across evaluations.
// Each evaluation's finalized rates feed a fresh summary; rate snapshots
// must never accumulate across evaluations. Query time is only the
// initial layout; a retained summary's lifecycle moves it to ingestion.
finalize_query_candidate(bound_child, &input.child)?
} else {
let child = finalize_exact_accumulator_at(
Expand Down Expand Up @@ -5249,14 +5144,17 @@ impl<'a> GlobalSelection<'a> {
if let Some(node) = self.assembled_nodes.borrow().get(&ptr) {
return Ok(Rc::clone(node));
}
// A selected summary that realizes its inner aggregate, instead of
// hiding it in `KeepPreAsap`, is kept; lifecycle assignment decides
// whether it runs in precompute or at query time.
let selected_composed_summary = self
.groups
.get(&ptr)
.and_then(|sel| sel.chosen)
.is_some_and(|candidate| matches!(&candidate.replacement,
Replacement::Summary(node) if matches!(&node.expr,
SummaryExpr::SummaryAgg { child, .. }
if matches!(&child.expr, SummaryExpr::KeepPreAsap(raw) if !contains_aggregate(raw)))));
if !matches!(&child.expr, SummaryExpr::KeepPreAsap(raw) if contains_aggregate(raw)))));
let node = if query_time_nested_sum(target) && !selected_composed_summary {
self.assemble_residual(target)?
} else {
Expand Down Expand Up @@ -6955,6 +6853,77 @@ mod tests {
use asap_types::types::AccuracyTarget;
use std::collections::HashMap;

// Candidate shape without execution timing: what is computed, not where.
fn timing_free_shape(node: &Rc<SummaryNode>) -> serde_json::Value {
fn strip(value: &mut serde_json::Value) {
match value {
serde_json::Value::Object(fields) => {
fields.remove("timing");
fields.values_mut().for_each(strip);
}
serde_json::Value::Array(values) => values.iter_mut().for_each(strip),
_ => {}
}
}
let mut shape =
serde_json::to_value(asap_types::post_asap::compile_post_asap_dag(node).unwrap())
.unwrap();
strip(&mut shape);
shape
}

// Rate inventories never offer two candidates that differ only in timing.
#[test]
fn rate_candidate_inventories_have_no_timing_only_duplicates() {
for (query, accuracy) in [
("sum by(job)(rate(m[1m]))", AccuracyTarget::Exact),
("topk by(job)(2, rate(m[1m]))", AccuracyTarget::Epsilon(0.1)),
] {
let root = Rc::new(lower_promql(query, accuracy));
let inventory = search_workload(vec![(0usize, root)])
.enumerate_candidate_dags(4096)
.unwrap();
let shapes = inventory
.candidates
.iter()
.map(|forest| timing_free_shape(&forest[0].1))
.collect::<Vec<_>>();
for (i, shape) in shapes.iter().enumerate() {
assert!(!shapes[..i].contains(shape), "{query}: duplicate {i}");
}
}
}

// Grouped Sum over Rate readouts stays a summary state in the inventory,
// so lifecycle assignment can place it in precompute or at query time.
#[test]
fn grouped_rate_sum_inventory_keeps_sum_state_for_lifecycle_placement() {
let root = Rc::new(lower_promql(
"sum by(job)(rate(m[1m]))",
AccuracyTarget::Exact,
));
let inventory = search_workload(vec![(0usize, root)])
.enumerate_candidate_dags(4096)
.unwrap();
let is_exact = |node: &SummaryNode, kind: ExactKind| {
matches!(&node.expr, SummaryExpr::SummaryAgg {
family: SummaryFamilyType::ExactAggregate(k, _), ..
} if *k == kind)
};
assert!(inventory.candidates.iter().any(|forest| {
let SummaryExpr::ValueOperation { child: sum, .. } = &forest[0].1.expr else {
return false;
};
let SummaryExpr::SummaryAgg { child: rate, .. } = &sum.expr else {
return false;
};
is_exact(sum, ExactKind::Sum)
&& matches!(&rate.expr, SummaryExpr::ValueOperation {
child, operation: ValueOperation::FinalizeExactAccumulator, ..
} if is_exact(child, ExactKind::Rate))
}));
}

// Every exposed query result has a readout; internal accumulator frontiers stay states.
#[test]
fn query_candidate_roots_do_not_leak_exact_accumulator_state() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -249,16 +249,13 @@ pub fn compile_rate_ranking(
Ok((source, program))
}

/// The selected logical placement requires fresh aggregate state per closed window.
/// Compile both physical graphs before deployment chooses storage or scheduling.
/// The input is the complete collection of per-series exact counter states.
/// Compile a lifecycle-timed DAG whose heap or grouped Sum over per-series
/// Rate readouts runs at ingestion time: fresh aggregate state per closed
/// window. The input is the complete collection of per-series counter states.
pub fn compile_fixed_window_rate_aggregation(
selected: &Rc<planner_types::post_asap::SummaryNode>,
dag: &planner_types::post_asap::PostAsapDag,
) -> Result<PhysicalCandidate, Error> {
use planner_types::post_asap::{
compile_post_asap_dag, ExactKind, ExecutionTiming, SketchAlgorithm,
};
let dag = compile_post_asap_dag(selected).map_err(|e| invalid(e.to_string()))?;
use planner_types::post_asap::{ExactKind, ExecutionTiming, SketchAlgorithm};
let sources = dag
.nodes
.iter()
Expand Down Expand Up @@ -310,7 +307,7 @@ pub fn compile_fixed_window_rate_aggregation(
));
}
compile_candidate(
&dag,
dag,
BTreeMap::from([(
u64::from(source.id.0),
InputContract::bounded(Arc::new(source.output_schema.clone())),
Expand Down
7 changes: 5 additions & 2 deletions crates/asap-physical-operators/tests/precompute_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ fn grouped_rate() -> PostAsapDag {
let space = grouped_rate_space();
let selected = space
.global_selection(&DefaultCostModel)
.assemble_selected_dag(&space.roots[0].1)
.assemble_selected_query(&space.roots[0].1)
.unwrap()
.unwrap();
compile_post_asap_dag(&selected).unwrap()
Expand Down Expand Up @@ -108,7 +108,10 @@ fn grouped_rate_can_be_materialized_before_or_after_grouped_sum() {
PostAsapOperatorPayload::Value {
operation: ValueOperation::FinalizeExactAccumulator
}
)
) && dag
.edges
.iter()
.any(|edge| edge.producer == state.id && edge.consumer == node.id)
})
.unwrap();
let input_schema = Arc::new(state.output_schema.clone());
Expand Down
Loading
Loading