diff --git a/crates/asap-aware-mapping/src/accuracy/composition.rs b/crates/asap-aware-mapping/src/accuracy/composition.rs index a284ab32..2fa1a11b 100644 --- a/crates/asap-aware-mapping/src/accuracy/composition.rs +++ b/crates/asap-aware-mapping/src/accuracy/composition.rs @@ -276,10 +276,10 @@ impl DefaultAccuracyModel { if inputs.len() != 2 || inputs .iter() - .any(|input| input.metric != ErrorMetric::RelativeValue) + .any(|input| input.metric != ErrorMetric::RelativeValue && !input.is_exact()) { return Err(unsupported( - "division needs exactly two RelativeValue guarantees".into(), + "division needs two relative-value or exact guarantees".into(), )); } let Some(numerator) = inputs[0].bound.evaluate() else { diff --git a/crates/asap-aware-mapping/src/cost_model.rs b/crates/asap-aware-mapping/src/cost_model.rs index 366f0289..b2c42953 100644 --- a/crates/asap-aware-mapping/src/cost_model.rs +++ b/crates/asap-aware-mapping/src/cost_model.rs @@ -876,6 +876,13 @@ pub trait CostModel { self.raw_query_recompute_cost(target) .map(|per_read| Cost(per_read.0 * expected_reads)) } + /// Physical feasibility evidence for a complete summary candidate. + /// `None` defers admission to physical/deployment compilation; `Some(false)` + /// excludes the candidate without changing its computation or parameters. + fn summary_support_evidence(&self, _summary: &SummaryNode) -> Option { + None + } + /// Which mixed exact/summary execution shapes the downstream runtime /// advertises (issue #171). Gates candidate *generation* in /// [`crate::exact_composition::ExactCompositionStrategy`]: a shape the diff --git a/crates/asap-aware-mapping/src/maintained_population.rs b/crates/asap-aware-mapping/src/maintained_population.rs index 180fc224..504c8100 100644 --- a/crates/asap-aware-mapping/src/maintained_population.rs +++ b/crates/asap-aware-mapping/src/maintained_population.rs @@ -50,6 +50,7 @@ fn recognize(root: &QueryExpr) -> Option<(MaintainedPopulation, PopulationReadou AggIntent::Quantile { q, col, .. } if q.is_finite() => { (*col, PopulationReadout::Quantile { q: *q }) } + AggIntent::TopK { k, .. } => (None, PopulationReadout::TopK { k: *k }), AggIntent::Sum { col } => (*col, PopulationReadout::Sum), AggIntent::Count { .. } => (None, PopulationReadout::Count), AggIntent::Avg { col } => (*col, PopulationReadout::Average), @@ -142,8 +143,11 @@ fn recognize(root: &QueryExpr) -> Option<(MaintainedPopulation, PopulationReadou if value_column.is_some_and(|c| schema.columns.get(c).is_none_or(|c| c.name != "value")) { return None; } - // Open time-series schemas distinguish instant PromQL populations from table rows. - if metric.is_empty() || schema.closed || schema.time_index.is_none() { + // PromQL can retain open labels or resolve them into a complete identity column. + if metric.is_empty() + || (schema.closed && !schema.has_promql_series_identity()) + || schema.time_index.is_none() + { return None; } let label = |col: usize| -> Option { diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index ffad44bb..f8f55a1d 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -520,15 +520,15 @@ impl ReplacementSubDAG { ) } - /// Runtime support for this candidate. Summary implementations remain - /// unknown until backend binding; a pure logical rewrite needs no new - /// physical operator. `Some(false)` disproves mixed-operation support. + /// Physical feasibility evidence for this candidate. A pure logical + /// rewrite needs no new operator. Unknown support is checked during + /// physical/deployment compilation; explicit rejection prevents selection. pub fn runtime_support_evidence(&self, cost_model: &dyn CostModel) -> Option { match &self.replacement { Replacement::ExactComposition(composition) => { cost_model.value_operation_support_evidence(&composition.op, composition.placement) } - Replacement::Summary(_) => None, + Replacement::Summary(node) => cost_model.summary_support_evidence(node), Replacement::Rewrite(_) => Some(true), } } @@ -834,6 +834,11 @@ pub(crate) fn realizations_for_intent( )), ], AccuracyTarget::Exact => vec![exact_realization(intent)], + _ if matches!(intent, AggIntent::Count { .. }) => { + let mut candidates = sketch_realizations(intent, accuracy, cost_model); + candidates.push(exact_realization(intent)); + candidates + } _ => sketch_realizations(intent, accuracy, cost_model), }, @@ -1277,6 +1282,173 @@ impl<'a> SketchAlgorithmStrategy<'a> { } } + /// Preserve the canonical Sort/Limit representation while exploring heap + /// realizations of an instant-vector ranking under the caller's target. + /// The input must carry the complete dynamic series identity. This never + /// treats a range of historical samples as the instant vector. + pub fn current_series_topk_candidates( + &self, + root: &Rc, + accuracy: &AccuracyTarget, + ) -> Proposals { + let QueryExpr::Limit { + n, + offset: 0, + child, + } = root.as_ref() + else { + return Proposals::default(); + }; + let QueryExpr::Sort { + keys, + partition_by, + child, + } = child.as_ref() + else { + return Proposals::default(); + }; + let [key] = keys.as_slice() else { + return Proposals::default(); + }; + let QueryExpr::Column(value) = key.expr else { + return Proposals::default(); + }; + let Ok(schema) = child.output_schema() else { + return Proposals::default(); + }; + if key.ascending + || key.nulls_first + || partition_by.is_without() + || !schema.has_promql_series_identity() + || !schema + .columns + .get(value) + .is_some_and(|column| column.name == "value") + || !is_current_series_source(child) + { + return Proposals::default(); + } + let ranked = Rc::new(QueryExpr::Aggregate { + reduction: Reduction::Reduce(partition_by.clone()), + measures: vec![AggIntent::TopK { + k: *n, + accuracy: accuracy.clone(), + }], + output_names: vec![], + having: None, + child: Rc::clone(child), + }); + 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) -> Proposals { + fn place(node: &Rc) -> Option> { + 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) -> Proposals { + fn query_time(node: &Rc) -> Rc { + 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 } } @@ -1286,6 +1458,22 @@ impl<'a> SketchAlgorithmStrategy<'a> { /// differs — see [`realize_child_with`]). fn propose_with(&self, root: &Rc, intent_override: Option<&AggIntent>) -> Proposals { let mut proposals = Proposals::default(); + // A selected logical rewrite otherwise remains KeepPreAsap during DAG + // assembly. Also expose its concrete summary realization for selection. + if intent_override.is_none() { + if let Some(rewritten) = crate::rewrite::composed_aggregate_rewrite(root) { + if let Ok(node) = realize_child_with(&rewritten, self.planning_inputs, None) { + if !matches!(node.expr, SummaryExpr::KeepPreAsap(_)) { + proposals.candidates.push(ReplacementSubDAG { + replacement: Replacement::Summary(node), + strategy: "SketchAlgorithmStrategy", + provenance: ReplacementProvenance::SummaryRealization, + rationale: "realize a schema-preserving composition of temporal and grouped accumulators".into(), + }); + } + } + } + } if let Ok(Some(node)) = exact_topk_over_temporal_values(root, self.planning_inputs) { proposals.candidates.push(ReplacementSubDAG { replacement: Replacement::Summary(node), @@ -1525,7 +1713,7 @@ fn describe_realization(intent: &AggIntent, realization: &Realization) -> String match realization { Realization::Sketch(kind) => format!( "{} realizes as a {:?} sketch — one of summary_candidates' \ - alternatives for this intent (asap_aware_mapping::replacement::realizations_for_intent)", + candidates for this intent (asap_aware_mapping::replacement::realizations_for_intent)", describe_intent(intent), kind.algorithm() ), @@ -1638,11 +1826,7 @@ fn exact_topk_over_temporal_values( else { return Ok(None); }; - let [AggIntent::TopK { - k, - accuracy: AccuracyTarget::Exact, - }] = measures.as_slice() - else { + let [AggIntent::TopK { k, .. }] = measures.as_slice() else { return Ok(None); }; let QueryExpr::Aggregate { @@ -1665,7 +1849,7 @@ fn exact_topk_over_temporal_values( { return Ok(None); } - let values = finalize_exact_accumulator(values, child)?; + let values = finalize_query_candidate(values, child)?; let partition_by = reduction .group_keys() .ok_or(RealizationError::PhysicalRealization( @@ -1821,11 +2005,30 @@ fn realize_binary( return Ok(None); } let domains = lhs_domain.zip(rhs_domain).map(|(lhs, rhs)| [lhs, rhs]); + let has_mean = [lhs, rhs] + .iter() + .any(|expr| matches!(bindable_intent(expr), Some(AggIntent::Avg { .. }))); + if has_mean + && domains.as_ref().is_none_or(|domains| { + domains.iter().any(|domain| { + !(domain.lower.abs().max(domain.upper.abs()) * domain.max_samples as f64) + .is_finite() + }) + }) + { + return Ok(None); + } lhs_node = realize_ddsketch_quantile_operand(lhs, planning_inputs, &target)?; rhs_node = realize_ddsketch_quantile_operand(rhs, planning_inputs, &target)?; if let Some(domains) = domains.as_ref() { for (domain, node) in domains.iter().zip([&lhs_node, &rhs_node]) { if !ddsketch_quantile_alpha(node) + .or_else(|| { + node.guarantee + .as_ref() + .is_some_and(ResultGuarantee::is_exact) + .then_some(alpha) + }) .is_some_and(|alpha| domain.supports_ddsketch(alpha)) { return Ok(None); @@ -1900,16 +2103,20 @@ fn realize_binary( return Ok(None); } - lhs_node = finalize_exact_accumulator(lhs_node, lhs)?; - rhs_node = finalize_exact_accumulator(rhs_node, rhs)?; + lhs_node = finalize_query_candidate(lhs_node, lhs)?; + rhs_node = finalize_query_candidate(rhs_node, rhs)?; let has_ratio_domains = ratio_domains.is_some(); let guarantee = if matches!(op, BinaryOpKind::Arithmetic(ArithmeticOpKind::Div)) && direct_ddsketch_ratio && has_ratio_domains - && ddsketch_quantile_alpha(&lhs_node).is_some() - && ddsketch_quantile_alpha(&rhs_node).is_some() - { + && [&lhs_node, &rhs_node].iter().all(|node| { + ddsketch_quantile_alpha(node).is_some() + || node + .guarantee + .as_ref() + .is_some_and(ResultGuarantee::is_exact) + }) { [lhs_node.guarantee.clone(), rhs_node.guarantee.clone()] .into_iter() .collect::>>() @@ -1961,7 +2168,7 @@ fn realize_binary( /// Put an explicit read boundary between maintained exact state and a /// query-time value consumer. Approximate summaries must already carry a /// `SummaryEstimate`, so they deliberately do not pass this predicate. -fn finalize_exact_accumulator( +pub fn finalize_query_candidate( node: Rc, logical_output: &QueryExpr, ) -> Result, RealizationError> { @@ -2018,9 +2225,8 @@ fn is_promql_scalar(expr: &QueryExpr) -> bool { ) } -/// Both direct quantile operands inherit the workload target during PromQL -/// lowering. Reuse that one target for the ratio rather than interpreting it -/// as two independent error budgets. +/// Quantile operands inherit one workload target. A temporal mean is exact +/// on its checked finite domain and needs no approximation budget. fn shared_quantile_target(lhs: &QueryExpr, rhs: &QueryExpr) -> Option { let quantile_target = |expr: &QueryExpr| match bindable_intent(expr) { Some(AggIntent::Quantile { accuracy, q, .. }) @@ -2030,9 +2236,16 @@ fn shared_quantile_target(lhs: &QueryExpr, rhs: &QueryExpr) -> Option None, }; - let lhs = quantile_target(lhs)?; - let rhs = quantile_target(rhs)?; - (lhs == rhs).then_some(lhs) + match (quantile_target(lhs), quantile_target(rhs)) { + (Some(lhs), Some(rhs)) => (lhs == rhs).then_some(lhs), + (Some(target), None) if matches!(bindable_intent(rhs), Some(AggIntent::Avg { .. })) => { + Some(target) + } + (None, Some(target)) if matches!(bindable_intent(lhs), Some(AggIntent::Avg { .. })) => { + Some(target) + } + _ => None, + } } /// For `a / b`, two DDSketches with the same relative bound `alpha` produce @@ -2244,7 +2457,7 @@ pub(crate) fn construct_summary_with( child_target, allocation, )?; - if is_counter_weighted_topk(intent, child) { + if is_snapshot_weighted_topk(intent, child) { return finish_weighted_topk(candidate, expr, intent); } return Ok(candidate); @@ -2351,13 +2564,25 @@ fn finish_weighted_topk( Ok(result) } -fn is_counter_weighted_topk(intent: &AggIntent, child: &QueryExpr) -> bool { +fn is_current_series_source(child: &QueryExpr) -> bool { + let source = match child { + QueryExpr::TimeRange { child, .. } => child.as_ref(), + source => source, + }; + matches!(source, QueryExpr::Scan { + source: asap_types::pre_asap::Source::TimeSeries { .. }, schema, .. + } if schema.has_promql_series_identity()) +} + +fn is_snapshot_weighted_topk(intent: &AggIntent, child: &QueryExpr) -> bool { matches!(intent, AggIntent::TopK { .. }) - && matches!(child, + && (is_current_series_source(child) + || matches!(child, QueryExpr::Aggregate { measures, child, .. } - if matches!(measures.as_slice(), [AggIntent::Sum { .. }]) - && matches!(child.as_ref(), QueryExpr::Aggregate { measures, .. } - if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase]))) + if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase]) + || (matches!(measures.as_slice(), [AggIntent::Sum { .. }]) + && matches!(child.as_ref(), QueryExpr::Aggregate { measures, .. } + if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase]))))) } /// Translate an [`Realization`] into the `(family, needs a @@ -2413,6 +2638,8 @@ type PhysicalSummaryInputRule = fn( /// `construct_summary_agg`. const PHYSICAL_SUMMARY_INPUT_RULES: &[PhysicalSummaryInputRule] = &[ realize_value_frequency_summary_input, + realize_counter_value_summary_input, + realize_current_series_summary_input, realize_keyed_additive_summary_input, ]; @@ -2568,10 +2795,10 @@ fn construct_summary_agg( SummaryFamilyType::Sketch(kind, _) if matches!(kind.algorithm(), SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap) ); - let rate_weighted = matches!(node, QueryExpr::Aggregate { child, .. } - if is_counter_weighted_topk(intent, child)); + let snapshot_weighted = matches!(node, QueryExpr::Aggregate { child, .. } + if is_snapshot_weighted_topk(intent, child)); let mut family = family; - let score_population = if rate_weighted { + let score_population = if snapshot_weighted { let bound = planning_inputs.evidence.topk_max_distinct_items(node); if bound.is_some_and(|n| n == 0 || n > (1u64 << 53)) { return Err(RealizationError::PhysicalRealization( @@ -2595,7 +2822,7 @@ fn construct_summary_agg( } else { None }; - let physical_reduction = if rate_weighted { + let physical_reduction = if snapshot_weighted { let QueryExpr::Aggregate { child, .. } = node else { unreachable!() }; @@ -2646,7 +2873,7 @@ fn construct_summary_agg( let state_idx = summary_col_index(&out_schema, &by, per_series); let readout_schema = if keyed_heap - && matches!(node, QueryExpr::Aggregate { child, .. } if is_counter_weighted_topk(intent, child)) + && matches!(node, QueryExpr::Aggregate { child, .. } if is_snapshot_weighted_topk(intent, child)) { keyed_heap_readout_schema(&input, node)? } else { @@ -2655,7 +2882,7 @@ fn construct_summary_agg( let summary_input = input.input; let query = estimate.then(|| { - if rate_weighted { + if snapshot_weighted { if let SummaryFamilyType::Sketch(kind, _) = &family { let capacity = match kind.params() { SketchParams::CmsWithHeap { heap_size, .. } @@ -2674,7 +2901,7 @@ fn construct_summary_agg( if keyed_heap { let mut state = state_schema.fields[state_idx].clone(); state.dtype = family.clone(); - let mut fields = if rate_weighted { + let mut fields = if snapshot_weighted { readout_schema.fields[..reduction.group_keys().map_or(0, |keys| keys.len())].to_vec() } else { Vec::new() @@ -2696,16 +2923,33 @@ fn construct_summary_agg( let bound_child = realize_child_with( &input.child, planning_inputs, - if rate_weighted { + if snapshot_weighted { Some(&AccuracyTarget::Exact) } else { child_target }, )?; - let bound_child = if rate_weighted { + let bound_child = if snapshot_weighted && is_current_series_source(&input.child) { + // Explicit snapshot selection prevents historical observations from + // becoming repeated weights in an instant-vector heap. + let root = Rc::new(node.clone()); + let population = crate::maintained_population::MaintainedPopulationStrategy::new( + std::slice::from_ref(&root), + ) + .candidate(&root) + .ok_or(RealizationError::PhysicalRealization( + "snapshot ranking requires a supported current-series population", + ))?; + let SummaryExpr::ValueOperation { child, .. } = &population.expr else { + return Err(RealizationError::PhysicalRealization( + "missing population readout", + )); + }; + 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. - finalize_exact_accumulator(bound_child, &input.child)? + finalize_query_candidate(bound_child, &input.child)? } else { let child = finalize_exact_accumulator_at( bound_child, @@ -2729,7 +2973,7 @@ fn construct_summary_agg( planning_inputs.evidence.estimator_contract(node), local_target, ); - let membership_query = if rate_weighted { + let membership_query = if snapshot_weighted { Some(readout(intent, &summary_input, planning_inputs.cost)) } else { query.clone() @@ -2744,7 +2988,7 @@ fn construct_summary_agg( allocation, )?; - if rate_weighted { + if snapshot_weighted { use asap_types::post_asap::{BoundExpr, ProbabilityExpr}; let target = accuracy_target(intent).expect("TopK target"); guarantee = if let Some(mut score) = @@ -2968,6 +3212,19 @@ fn ranking_score_index( logical: &QueryExpr, values: &SummarySchema, ) -> Result { + if is_current_series_source(logical) { + return values + .fields + .iter() + .position(|field| { + field.name == "value" + && field.dtype + == SummaryFamilyType::Plain(asap_types::pre_asap::DataType::Float64) + }) + .ok_or(RealizationError::PhysicalRealization( + "snapshot ranking requires the sample value column", + )); + } let QueryExpr::Aggregate { reduction, measures, @@ -3015,6 +3272,124 @@ fn ranking_score_index( Ok(index) } +/// Rebuild a heap from this evaluation's finalized per-series counter values. +/// The rate window is preserved; raw counter samples never become CMS weights. +fn realize_counter_value_summary_input( + intent: &AggIntent, + family: &SummaryFamilyType, + output_reduction: &Reduction, + child: &Rc, +) -> PhysicalSummaryInputRuleResult { + if !matches!(intent, AggIntent::TopK { .. }) + || !matches!(family, SummaryFamilyType::Sketch(kind, _) if matches!(kind.algorithm(), SketchAlgorithm::CmsWithHeap | SketchAlgorithm::CountSketchWithHeap)) + || !matches!(child.as_ref(), QueryExpr::Aggregate { reduction: Reduction::PerEntity, measures, .. } if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase])) + { + return PhysicalSummaryInputRuleResult::NotApplicable; + } + let Ok(schema) = child.output_schema() else { + return PhysicalSummaryInputRuleResult::Unsupported( + "counter ranking needs a valid value schema", + ); + }; + if !schema.closed { + return PhysicalSummaryInputRuleResult::Unsupported( + "counter ranking needs the complete resolved series identity", + ); + } + let Reduction::Reduce(groups) = output_reduction else { + return PhysicalSummaryInputRuleResult::Unsupported( + "counter ranking requires explicit partitions", + ); + }; + if groups.is_without() { + return PhysicalSummaryInputRuleResult::Unsupported( + "counter ranking requires resolved partitions", + ); + } + // Retain the evaluation timestamp in each returned row. This sketch is a + // snapshot, not an additive history of successive rate evaluations. + let items = schema + .columns + .iter() + .enumerate() + .filter(|(index, column)| column.name != "value" && !groups.contains(index)) + .map(|(index, _)| schema_column_ref(child, index).map(SummaryInputExpr::Column)) + .collect::>>(); + let Some(items) = items.filter(|items| !items.is_empty()) else { + return PhysicalSummaryInputRuleResult::Unsupported("counter ranking has no item columns"); + }; + PhysicalSummaryInputRuleResult::Realized(PhysicalSummaryInput { + child: Rc::clone(child), + input: SummaryUpdate { + item: Some(SummaryInputExpr::Tuple(items)), + weight: SummaryInputExpr::Column(ColumnRef::SampleValue), + weight_domain: WeightDomain::NonNegative { + proof: NonNegativeWeightProof::ResetAwareCounterDerivative, + }, + }, + }) +} + +/// An instant-vector source has one current value per full series identity. +/// Rebuild the state for each evaluation; historical samples are not updates. +fn realize_current_series_summary_input( + intent: &AggIntent, + family: &SummaryFamilyType, + output_reduction: &Reduction, + child: &Rc, +) -> PhysicalSummaryInputRuleResult { + if !matches!(intent, AggIntent::TopK { .. }) || !is_current_series_source(child) { + return PhysicalSummaryInputRuleResult::NotApplicable; + } + let SummaryFamilyType::Sketch(kind, _) = family else { + return PhysicalSummaryInputRuleResult::NotApplicable; + }; + match kind.algorithm() { + SketchAlgorithm::CountSketchWithHeap => {} + SketchAlgorithm::CmsWithHeap => { + return PhysicalSummaryInputRuleResult::Unsupported( + "current sample values do not prove non-negative CMS weights", + ) + } + _ => return PhysicalSummaryInputRuleResult::NotApplicable, + } + let Reduction::Reduce(groups) = output_reduction else { + return PhysicalSummaryInputRuleResult::Unsupported( + "snapshot ranking requires explicit partitions", + ); + }; + if groups.is_without() { + return PhysicalSummaryInputRuleResult::Unsupported( + "snapshot ranking requires resolved partitions", + ); + } + let Ok(schema) = child.output_schema() else { + return PhysicalSummaryInputRuleResult::Unsupported( + "snapshot ranking requires a valid source schema", + ); + }; + let items = schema + .columns + .iter() + .enumerate() + .filter(|(index, column)| column.name != "value" && !groups.contains(index)) + .map(|(index, _)| schema_column_ref(child, index).map(SummaryInputExpr::Column)) + .collect::>>(); + let Some(items) = items.filter(|items| !items.is_empty()) else { + return PhysicalSummaryInputRuleResult::Unsupported( + "snapshot ranking has no item identity", + ); + }; + PhysicalSummaryInputRuleResult::Realized(PhysicalSummaryInput { + child: Rc::clone(child), + input: SummaryUpdate { + item: Some(SummaryInputExpr::Tuple(items)), + weight: SummaryInputExpr::Column(ColumnRef::SampleValue), + weight_domain: WeightDomain::UnknownOrSigned, + }, + }) +} + /// Realize the composite heavy-hitter realization for /// `TopK(Count GROUP BY key)`. The heap sketch consumes the raw keyed stream; /// it does not consume an independently materialized Count result. @@ -3665,6 +4040,256 @@ impl PlanSpace { } } +/// DAG candidates assembled from an unpriced search space. +/// This is an internal planning stage: callers must still validate lifecycle +/// requirements and compile supported physical operators before deployment. +/// The caller supplies a finite expansion budget; exceeding it is an error, +/// never a silently truncated inventory presented as exhaustive. +#[derive(Debug)] +pub struct CandidateDagInventory { + pub candidates: Vec)>>, + pub rejected_assemblies: Vec, +} + +type CandidateDagChoice<'a> = (Option<&'a ReplacementSubDAG>, Option>); + +impl PlanSpace { + pub fn enumerate_candidate_dags( + &self, + expansion_limit: usize, + ) -> Result, RealizationError> { + self.enumerate_candidate_roots(&self.roots, expansion_limit) + } + + /// Enumerate one workload root without expanding independent roots' choices. + /// Discovery and composition proofs still come from the shared workload + /// space. Deployment may price combinations lazily; this API does not rank + /// candidates or claim that independently cheapest roots minimize shared cost. + pub fn enumerate_candidate_dags_for_root( + &self, + id: &Id, + expansion_limit: usize, + ) -> Result, RealizationError> { + let roots = self + .roots + .iter() + .filter(|(candidate, _)| candidate == id) + .cloned() + .collect::>(); + if roots.len() != 1 { + return Err(RealizationError::PhysicalRealization( + "candidate enumeration requires one uniquely identified workload root", + )); + } + self.enumerate_candidate_roots(&roots, expansion_limit) + } + + fn enumerate_candidate_roots( + &self, + roots: &[(Id, Rc)], + expansion_limit: usize, + ) -> Result, RealizationError> { + let mut reachable = Vec::new(); + let mut nodes = HashMap::new(); + let mut counts = HashMap::new(); + for (_, root) in roots { + walk(root, &mut reachable, &mut nodes, &mut counts); + } + // Rewrites may introduce descendants absent from the original root. + let mut cursor = 0; + while cursor < reachable.len() { + let ptr = reachable[cursor]; + cursor += 1; + if let Some(group) = self.groups.get(&ptr) { + for candidate in &group.candidates { + if let Replacement::Rewrite(rewritten) = &candidate.replacement { + walk(rewritten, &mut reachable, &mut nodes, &mut counts); + } + } + } + } + let order = self + .order + .iter() + .copied() + .filter(|ptr| counts.contains_key(ptr)) + .collect::>(); + // Composition plans carry the proofs established during discovery. + // No cost ranking is consulted while expanding these choices. + let options: Vec>> = order + .iter() + .map(|ptr| { + let group = &self.groups[ptr]; + let mut choices = vec![(None, None)]; + for candidate in &group.candidates { + match &candidate.replacement { + Replacement::ExactComposition(operation) => { + for prepared in &self.composition_plans { + if prepared.target == *ptr + && prepared.operation.placement == operation.placement + && prepared.operation.op == operation.op + && Rc::ptr_eq( + &prepared.operation.child_target, + &operation.child_target, + ) + { + choices + .push((Some(candidate), Some(Rc::clone(&prepared.plan)))); + } + } + } + _ => choices.push((Some(candidate), None)), + } + } + choices + }) + .collect(); + let combinations = options + .iter() + .try_fold(1usize, |n, choices| n.checked_mul(choices.len())) + .filter(|n| *n <= expansion_limit) + .ok_or(RealizationError::PhysicalRealization( + "candidate expansion budget exceeded; no partial inventory returned", + ))?; + let mut inventory = CandidateDagInventory { + candidates: Vec::new(), + rejected_assemblies: Vec::new(), + }; + // Hash buckets avoid quadratic comparisons across a large workload + // inventory. Equality still decides deduplication, including collisions. + let mut seen = HashMap::>::new(); + for mut ordinal in 0..combinations { + let mut groups = HashMap::new(); + let mut assembled_nodes = HashMap::new(); + for (ptr, choices) in order.iter().zip(&options) { + let (chosen, prepared) = &choices[ordinal % choices.len()]; + ordinal /= choices.len(); + let group = &self.groups[ptr]; + if let Some(node) = prepared { + assembled_nodes.insert(*ptr, Rc::clone(node)); + } + groups.insert( + *ptr, + TargetSubDAGSelection { + target: &group.target, + consumer_count: group.consumer_count, + effective_consumer_count: group.consumer_count, + chosen: *chosen, + composition: None, + }, + ); + } + let assembly = GlobalSelection { + order: order.clone(), + groups, + assembled_nodes: RefCell::new(assembled_nodes), + }; + let roots = roots + .iter() + .map(|(id, root)| { + assembly + .assemble_target(root) + // Exposed query candidates return values. Internal assembly + // still retains accumulator states for sharing and storage. + .and_then(|node| finalize_query_candidate(node, root)) + .map(|node| (id.clone(), node)) + }) + .collect::, _>>(); + match roots { + Ok(roots) => { + let roots = asap_types::post_asap::share_common_summary_subtrees(roots); + use std::hash::{Hash, Hasher}; + let mut hash = std::collections::hash_map::DefaultHasher::new(); + let mut pending = roots + .iter() + .map(|(_, node)| node.as_ref()) + .collect::>(); + while let Some(node) = pending.pop() { + std::mem::discriminant(&node.expr).hash(&mut hash); + let raw = match &node.expr { + SummaryExpr::KeepPreAsap(raw) => Some(raw.as_ref()), + _ => None, + }; + let operation = match &node.expr { + SummaryExpr::ValueOperation { + timing, operation, .. + } => serde_json::json!((timing, operation)), + SummaryExpr::BinaryOp { + timing, operator, .. + } => serde_json::json!((timing, operator)), + SummaryExpr::SummaryMerge { timing, .. } => serde_json::json!(timing), + _ => serde_json::Value::Null, + }; + let mut value = + serde_json::to_value((&node.schema, &node.guarantee, raw, operation)) + .map_err(|_| { + RealizationError::PhysicalRealization( + "candidate identity serialization failed", + ) + })?; + fn normalize(value: &mut serde_json::Value) { + match value { + serde_json::Value::Number(number) + if number.as_f64() == Some(0.0) => + { + *value = serde_json::json!(0); + } + serde_json::Value::Array(values) => { + values.iter_mut().for_each(normalize) + } + serde_json::Value::Object(values) => { + values.values_mut().for_each(normalize) + } + _ => {} + } + } + normalize(&mut value); + value.sort_all_objects(); + value.to_string().hash(&mut hash); + match &node.expr { + SummaryExpr::KeepPreAsap(_) => {} + SummaryExpr::BinaryOp { lhs, rhs, .. } => { + pending.extend([lhs.as_ref(), rhs.as_ref()]) + } + SummaryExpr::RelationalJoin { left, right, .. } + | SummaryExpr::SummarySubtract { left, right } => { + pending.extend([left.as_ref(), right.as_ref()]) + } + SummaryExpr::ValueOperation { child, .. } + | SummaryExpr::SummaryAgg { child, .. } => pending.push(child.as_ref()), + SummaryExpr::SummaryJoin { outer, inner, .. } => { + pending.extend([outer.as_ref(), inner.as_ref()]) + } + SummaryExpr::SummaryDelete { summary_input, .. } + | SummaryExpr::SummaryEstimate { summary_input, .. } => { + pending.push(summary_input.as_ref()) + } + SummaryExpr::SummaryMerge { children, .. } => { + pending.extend(children.iter().map(|child| child.as_ref())) + } + } + } + let bucket = seen.entry(hash.finish()).or_default(); + if !bucket + .iter() + .any(|&index| inventory.candidates[index] == roots) + { + bucket.push(inventory.candidates.len()); + inventory.candidates.push(roots); + } + } + Err(error) => { + let reason = error.to_string(); + if !inventory.rejected_assemblies.contains(&reason) { + inventory.rejected_assemblies.push(reason); + } + } + } + } + Ok(inventory) + } +} + /// Lifecycle-aware whole-subplan costs keyed by target and candidate identity. #[derive(Default)] pub(crate) struct CandidateCostOverrides { @@ -4316,9 +4941,17 @@ fn rank_group<'a>( // purpose. `total_cmp` gives deterministic placement to a model's NaN // placeholders without dropping any candidate. ranked.sort_by(|a, b| { - cost_model - .estimate_cost(a, &target) - .total_cmp(&cost_model.estimate_cost(b, &target)) + match ( + cost_model.candidate_cost(a, &target), + cost_model.candidate_cost(b, &target), + ) { + (Some(a), Some(b)) => a.0.total_cmp(&b.0), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => cost_model + .estimate_cost(a, &target) + .total_cmp(&cost_model.estimate_cost(b, &target)), + } }); ranked } @@ -4553,12 +5186,32 @@ impl<'a> GlobalSelection<'a> { self.assemble_target(target).map(Some) } + /// Assemble a complete query result, including an exact-state readout when + /// needed. `assemble_selected_dag` also serves internal state frontiers; + /// callers exposing query results must use this boundary instead. + pub fn assemble_selected_query( + &self, + target: &Rc, + ) -> Result>, RealizationError> { + self.assemble_selected_dag(target)? + .map(|node| finalize_query_candidate(node, target)) + .transpose() + } + fn assemble_target(&self, target: &Rc) -> Result, RealizationError> { let ptr = Rc::as_ptr(target); if let Some(node) = self.assembled_nodes.borrow().get(&ptr) { return Ok(Rc::clone(node)); } - let node = if query_time_nested_sum(target) { + 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))))); + let node = if query_time_nested_sum(target) && !selected_composed_summary { self.assemble_residual(target)? } else { match self @@ -4609,8 +5262,8 @@ impl<'a> GlobalSelection<'a> { let Some(pred) = normalized_pred else { return keep_pre_asap(target); }; - let left = finalize_exact_accumulator(self.assemble_target(left)?, left)?; - let right = finalize_exact_accumulator(self.assemble_target(right)?, right)?; + let left = finalize_query_candidate(self.assemble_target(left)?, left)?; + let right = finalize_query_candidate(self.assemble_target(right)?, right)?; let guarantee = relational_join_guarantee(left.guarantee.as_ref(), right.guarantee.as_ref()); let node = Rc::new(SummaryNode { @@ -4681,7 +5334,7 @@ impl<'a> GlobalSelection<'a> { ), _ => return keep_pre_asap(target), }; - let child = finalize_exact_accumulator(self.assemble_target(child_target)?, child_target)?; + let child = finalize_query_candidate(self.assemble_target(child_target)?, child_target)?; let guarantee = child.guarantee.clone(); let node = Rc::new(SummaryNode { expr: SummaryExpr::ValueOperation { @@ -4904,7 +5557,7 @@ fn composition_options<'a>( None => child_group.candidates.iter().collect(), }; for child_candidate in child_candidates { - if child_candidate.has_missing_accuracy_evidence() { + if !is_automatically_selectable(child_candidate, cost_model) { continue; } let Replacement::Summary(summary) = &child_candidate.replacement else { @@ -5081,7 +5734,7 @@ impl PlanSpace { .candidates .iter() .filter(|candidate| !is_composition_candidate(candidate)) - .filter(|candidate| is_automatically_selectable(candidate)) + .filter(|candidate| is_automatically_selectable(candidate, cost_model)) .filter_map(|candidate| { costs .get(&group.target, candidate) @@ -5107,7 +5760,7 @@ impl PlanSpace { .filter(|candidate| { !is_cse_candidate(candidate) && !is_composition_candidate(candidate) - && is_automatically_selectable(candidate) + && is_automatically_selectable(candidate, cost_model) }) .filter_map(|candidate| { cost_model @@ -5164,7 +5817,7 @@ impl PlanSpace { .filter(|candidate| { !is_cse_candidate(candidate) && !is_composition_candidate(candidate) - && is_automatically_selectable(candidate) + && is_automatically_selectable(candidate, cost_model) }) .filter_map(|candidate| { cost_model @@ -5209,7 +5862,7 @@ impl PlanSpace { // children (see `multiplier`'s `_ => effective` arm). None => rank_group(group, cost_model).into_iter().find(|candidate| { !is_composition_candidate(candidate) - && is_automatically_selectable(candidate) + && is_automatically_selectable(candidate, cost_model) && (cost_model .candidate_cost( candidate, @@ -5226,7 +5879,7 @@ impl PlanSpace { .find(|candidate| { !is_cse_candidate(candidate) && !is_composition_candidate(candidate) - && is_automatically_selectable(candidate) + && is_automatically_selectable(candidate, cost_model) && (cost_model .candidate_cost(candidate, &effective_target) .is_some() @@ -5312,8 +5965,9 @@ fn is_cse_candidate(candidate: &ReplacementSubDAG) -> bool { ) } -fn is_automatically_selectable(candidate: &ReplacementSubDAG) -> bool { +fn is_automatically_selectable(candidate: &ReplacementSubDAG, cost_model: &dyn CostModel) -> bool { !candidate.has_missing_accuracy_evidence() + && candidate.runtime_support_evidence(cost_model) != Some(false) } /// How much one direct reference to `parent_ptr` actually costs, once @@ -6232,6 +6886,150 @@ mod tests { use asap_types::types::AccuracyTarget; use std::collections::HashMap; + // Every exposed query result has a readout; internal accumulator frontiers stay states. + #[test] + fn query_candidate_roots_do_not_leak_exact_accumulator_state() { + for query in [ + "sum by(job)(rate(m[1m]))", + "sum by(job)(m)", + "sum_over_time(m[1m])", + ] { + let root = Rc::new(lower_promql(query, AccuracyTarget::Exact)); + let space = search_workload(vec![(0usize, root.clone())]); + let inventory = space.enumerate_candidate_dags(4096).unwrap(); + assert!(!inventory.candidates.is_empty()); + let strategy = SketchAlgorithmStrategy::new(&DefaultCostModel); + for candidate in strategy.propose(&TargetSubDAG::new(&root)).candidates { + if let Replacement::Summary(node) = candidate.replacement { + let output = finalize_query_candidate(node, &root).unwrap(); + assert!( + output + .schema + .fields + .iter() + .all(|field| matches!(field.dtype, SummaryFamilyType::Plain(_))), + "direct candidate {query} leaks state" + ); + } + } + let selected = space + .global_selection(&DefaultCostModel) + .assemble_selected_query(&space.roots[0].1) + .unwrap() + .unwrap(); + for node in inventory + .candidates + .iter() + .map(|forest| &forest[0].1) + .chain(std::iter::once(&selected)) + { + assert!( + node.schema + .fields + .iter() + .all(|field| matches!(field.dtype, SummaryFamilyType::Plain(_))), + "{query}: query root leaks state: {:?}", + node.schema + ); + } + } + } + + #[test] + fn unpriced_inventory_retains_quantile_families_and_raw_execution() { + let query = Rc::new(agg(vec![2], default_quantile(0.9), metric_scan(&["job"]))); + let space = search_workload(vec![(0usize, query)]); + let inventory = space.enumerate_candidate_dags(4096).unwrap(); + let roots = inventory + .candidates + .iter() + .map(|forest| format!("{:?}", forest[0].1)) + .collect::>(); + assert!(roots.iter().any(|root| root.contains("Kll"))); + assert!(roots.iter().any(|root| root.contains("DDSketch"))); + assert!(inventory + .candidates + .iter() + .any(|forest| matches!(forest[0].1.expr, SummaryExpr::KeepPreAsap(_)))); + } + + // Independent roots must not require materializing their Cartesian product. + #[test] + fn root_inventory_preserves_choices_without_workload_cartesian_expansion() { + let roots = (0..24usize) + .map(|id| { + ( + id, + Rc::new(agg( + vec![2], + default_quantile((id + 1) as f64 / 25.0), + metric_scan(&["job"]), + )), + ) + }) + .collect(); + let space = search_workload(roots); + assert!(space.enumerate_candidate_dags(4096).is_err()); + for id in 0..24 { + let inventory = space.enumerate_candidate_dags_for_root(&id, 4096).unwrap(); + assert!(inventory + .candidates + .iter() + .all(|forest| forest.len() == 1 && forest[0].0 == id)); + let descriptions = inventory + .candidates + .iter() + .map(|forest| format!("{:?}", forest[0].1)) + .collect::>(); + assert!(descriptions.iter().any(|node| node.contains("Kll"))); + assert!(descriptions.iter().any(|node| node.contains("DDSketch"))); + assert!(inventory + .candidates + .iter() + .any(|forest| matches!(forest[0].1.expr, SummaryExpr::KeepPreAsap(_)))); + } + assert!(space.enumerate_candidate_dags_for_root(&24, 4096).is_err()); + assert!(space.enumerate_candidate_dags_for_root(&0, 0).is_err()); + } + + // Factoring changes enumeration, not the set of root computations. + #[test] + fn root_inventory_matches_projection_of_exhaustive_workload_inventory() { + let roots = (0..2usize) + .map(|id| { + ( + id, + Rc::new(agg( + vec![2], + default_quantile(0.5 + id as f64 * 0.4), + metric_scan(&["job"]), + )), + ) + }) + .collect(); + let space = search_workload(roots); + let full = space.enumerate_candidate_dags(4096).unwrap(); + for id in 0..2 { + let inventory = space.enumerate_candidate_dags_for_root(&id, 4096).unwrap(); + for forest in &full.candidates { + let node = &forest.iter().find(|(root, _)| *root == id).unwrap().1; + assert!(inventory.candidates.iter().any(|one| &one[0].1 == node)); + } + for one in &inventory.candidates { + assert!(full.candidates.iter().any(|forest| forest + .iter() + .any(|(root, node)| *root == id && node == &one[0].1))); + } + } + } + + #[test] + fn inventory_budget_never_returns_a_silent_partial_search() { + let query = Rc::new(agg(vec![2], default_quantile(0.9), metric_scan(&["job"]))); + let space = search_workload(vec![(0usize, query)]); + assert!(space.enumerate_candidate_dags(0).is_err()); + } + fn equi_pred(left: ColumnId, right: ColumnId) -> Predicate { Predicate(Rc::new(QueryExpr::Compare { left: Rc::new(QueryExpr::Column(left)), @@ -6265,6 +7063,25 @@ mod tests { ); } + // Approximate requests also admit exact temporal ranking candidates. + #[test] + fn approximate_temporal_topk_admits_exact_maintained_values() { + let root = Rc::new(lower_promql( + "topk by(job)(1,count_over_time(a[5m]))", + AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }, + )); + let planning_inputs = + CandidatePlanningInputs::with_default_accuracy(&crate::cost_model::DefaultCostModel); + let node = exact_topk_over_temporal_values(&root, planning_inputs) + .unwrap() + .expect("exact ranking is legal for an approximate request"); + assert!(node.guarantee.as_ref().unwrap().is_exact()); + asap_types::post_asap::compile_post_asap_dag(&node).unwrap(); + } + // Exact Top-K consumes the Planner's maintained temporal values. #[test] fn exact_temporal_topk_has_a_maintained_value_candidate() { @@ -6318,6 +7135,43 @@ mod tests { } } + // A bounded exact mean can share the relative division proof with a quantile. + #[test] + fn bounded_mean_quantile_ratio_is_certified() { + struct Domain; + impl AccuracyEvidenceProvider for Domain { + fn quantile_input_domain( + &self, + _: &QueryExpr, + ) -> Option { + Some(crate::accuracy::QuantileInputDomain { + lower: 1.0, + upper: 1000.0, + max_samples: 10000, + contract: "finite test population".into(), + }) + } + } + let target = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }; + let inputs = CandidatePlanningInputs { + evidence: &Domain, + ..CandidatePlanningInputs::with_default_accuracy(&DefaultCostModel) + }; + for query in [ + "avg_over_time(a[5m]) / quantile_over_time(0.5,a[5m])", + "quantile_over_time(0.5,a[5m]) / avg_over_time(a[5m])", + ] { + let root = Rc::new(lower_promql(query, target.clone())); + let node = realize_binary(&root, inputs, Some(&target)) + .unwrap() + .expect("bounded ratio candidate"); + assert!(DefaultAccuracyModel.satisfies(node.guarantee.as_ref().unwrap(), &target)); + } + } + // Missing domain proof permits an uncertified direct quantile ratio only. #[test] fn quantile_ratio_without_input_proof_has_no_root_guarantee() { @@ -6614,6 +7468,23 @@ mod tests { ); } + // Exact counting remains a legal candidate under an approximate target. + #[test] + fn approximate_count_includes_exact_accumulator_candidate() { + let intent = AggIntent::Count { + accuracy: eps(0.01), + }; + assert!(realizations_for_intent(&intent, &DefaultCostModel) + .iter() + .any(|candidate| matches!( + candidate, + Realization::ExactAggregate { + kind: ExactKind::Count, + .. + } + ))); + } + #[test] fn epsilon_delta_sizes_cms_depth() { let intent = AggIntent::Count { @@ -7399,7 +8270,7 @@ mod tests { assert_eq!(agg_group.consumer_count, 1); assert_eq!( agg_group.candidates.len(), - 5, + 6, "Hydra candidates with unknown evidence remain available: {:?}", agg_group.candidates ); @@ -8458,6 +9329,125 @@ mod tests { ); } + // A cheap but physically infeasible candidate must not be selected. + #[test] + fn explicit_summary_infeasibility_prevents_selection() { + struct Unsupported; + impl CostModel for Unsupported { + fn rank_candidates( + &self, + _: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + fn candidate_cost(&self, _: &ReplacementSubDAG, _: &TargetSubDAG<'_>) -> Option { + Some(Cost(1.0)) + } + fn summary_support_evidence(&self, _: &SummaryNode) -> Option { + Some(false) + } + } + let root = Rc::new(lower_promql("sum_over_time(a[1m])", AccuracyTarget::Exact)); + let space = search_workload(vec![("q", root)]); + let selected = space.global_selection(&Unsupported); + assert!(selected + .for_target(&space.roots[0].1) + .unwrap() + .chosen + .is_none()); + } + + // Composable temporal/grouped Sum must be executable as one producer. + #[test] + fn grouped_temporal_sum_has_one_summary_producer_candidate() { + let root = Rc::new(lower_promql( + "sum by(job)(sum_over_time(a[1m]))", + AccuracyTarget::Exact, + )); + let candidates = + SketchAlgorithmStrategy::default_cost_model().replacements(&TargetSubDAG::new(&root)); + assert!(candidates + .iter() + .any(|candidate| matches!(&candidate.replacement, + Replacement::Summary(node) if matches!(&node.expr, + SummaryExpr::SummaryAgg { reduction: Reduction::Reduce(_), child, .. } + if matches!(child.expr, SummaryExpr::KeepPreAsap(_)))))); + struct PreferComposed; + impl CostModel for PreferComposed { + fn rank_candidates( + &self, + _: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + fn candidate_cost( + &self, + candidate: &ReplacementSubDAG, + _: &TargetSubDAG<'_>, + ) -> Option { + Some(Cost( + if matches!(&candidate.replacement, + Replacement::Summary(node) if matches!(&node.expr, + SummaryExpr::SummaryAgg { reduction: Reduction::Reduce(_), child, .. } + if matches!(child.expr, SummaryExpr::KeepPreAsap(_)))) + { + 1.0 + } else { + 100.0 + }, + )) + } + } + let space = search_workload(vec![("q", root.clone())]); + let selected = space.global_selection(&PreferComposed); + let node = selected.assemble_target(&space.roots[0].1).unwrap(); + assert!(matches!(&node.expr, + SummaryExpr::SummaryAgg { reduction: Reduction::Reduce(_), child, .. } + if matches!(child.expr, SummaryExpr::KeepPreAsap(_)))); + } + + // Mixed candidate ranking must honor explicit costs, not legacy estimates. + #[test] + fn mixed_candidate_ranking_uses_explicit_candidate_costs() { + struct ExplicitCosts; + impl CostModel for ExplicitCosts { + fn rank_candidates( + &self, + _: &AggIntent, + candidates: &[SketchAlgorithm], + ) -> Vec { + candidates.to_vec() + } + fn candidate_cost( + &self, + candidate: &ReplacementSubDAG, + _: &TargetSubDAG<'_>, + ) -> Option { + Some(Cost( + if candidate.provenance == ReplacementProvenance::LogicalRewrite { + 1.0 + } else { + 100.0 + }, + )) + } + } + let root = Rc::new(lower_promql( + "sum by(job)(sum_over_time(a[1m]))", + AccuracyTarget::Exact, + )); + let space = search_workload(vec![("q", root)]); + let selection = space.global_selection(&ExplicitCosts); + let selected = selection + .for_target(&space.roots[0].1) + .unwrap() + .chosen + .unwrap(); + assert_eq!(selected.provenance, ReplacementProvenance::LogicalRewrite); + } + #[test] fn global_selection_compares_a_logical_rewrite_with_the_cse_choice() { struct PreferLogicalRewrite; diff --git a/crates/asap-aware-mapping/src/rewrite.rs b/crates/asap-aware-mapping/src/rewrite.rs index 6081105e..44b685f8 100644 --- a/crates/asap-aware-mapping/src/rewrite.rs +++ b/crates/asap-aware-mapping/src/rewrite.rs @@ -247,7 +247,7 @@ fn build_rewrite(root: &Rc) -> Option> { /// Compose adjacent per-entity and cross-entity accumulators when their /// algebra, rather than a query-language spelling, proves equivalence. -fn composed_aggregate_rewrite(root: &Rc) -> Option> { +pub(crate) fn composed_aggregate_rewrite(root: &Rc) -> Option> { let original_schema = root.output_schema().ok()?; let QueryExpr::Aggregate { reduction: outer_reduction @ Reduction::Reduce(_), diff --git a/crates/devtools/src/bin/dag_export.rs b/crates/devtools/src/bin/dag_export.rs index 0fb2fd8b..4ec9fb1a 100644 --- a/crates/devtools/src/bin/dag_export.rs +++ b/crates/devtools/src/bin/dag_export.rs @@ -2869,8 +2869,10 @@ mod tests { } } + // TopK queries of different sizes share one exact ranked prefix and differ + // only in their final Limit. #[test] - fn workload_node_ids_make_smaller_topk_reuse_explicit() { + fn workload_node_ids_share_ranked_prefix_across_topk_sizes() { let small = lower_promql( "topk(5, rate(http_requests_total[5m]))", AccuracyTarget::Epsilon(0.01), @@ -2914,18 +2916,14 @@ mod tests { .unwrap() .1; let q3_root = &q3.nodes[q3.root as usize]; - assert_eq!(q3_root.label, "Limit(5)"); - assert_eq!( - q3_root - .decision - .as_ref() - .map(|decision| decision.strategy.as_str()), - Some("TopKLimitReuseStrategy") - ); - let q3_large = &q3.nodes[q3_root.children[0] as usize]; - let q4_large = &q4.nodes[q4.root as usize]; - assert_eq!(q3_large.label, "Limit(10)"); - assert_eq!(q3_large.workload_node_id, q4_large.workload_node_id); + let q4_root = &q4.nodes[q4.root as usize]; + assert!(q3_root.label.contains("Limit { n: 5,")); + assert!(q4_root.label.contains("Limit { n: 10,")); + assert_ne!(q3_root.workload_node_id, q4_root.workload_node_id); + let q3_ranked = &q3.nodes[q3_root.children[0] as usize]; + let q4_ranked = &q4.nodes[q4_root.children[0] as usize]; + assert!(q3_ranked.label.contains("Sort")); + assert_eq!(q3_ranked.workload_node_id, q4_ranked.workload_node_id); } #[tokio::test] diff --git a/crates/frontend-promql/src/promql.rs b/crates/frontend-promql/src/promql.rs index f0cf907f..c3df6c67 100644 --- a/crates/frontend-promql/src/promql.rs +++ b/crates/frontend-promql/src/promql.rs @@ -625,7 +625,10 @@ fn build_over_subtree(outer: Outer, keys: Vec, child: Unresolved) -> && matches!(sum_child.as_ref(), Unresolved::Aggregate { measures, .. } if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase])) ); - if descending && weighted_counter_ranking { + let direct_counter_ranking = matches!(&child, Unresolved::Aggregate { + measures, reduction: Reduction::PerEntity, .. + } if matches!(measures.as_slice(), [AggIntent::Rate | AggIntent::Increase])); + if descending && (weighted_counter_ranking || direct_counter_ranking) { return Ok(outer_aggregate( keys, AggIntent::TopK { @@ -1487,10 +1490,27 @@ fn build(inner: Inner, keys: Vec, outer: Outer) -> Result }) } Outer::TopK { k, descending } => { + // Preserve the counter-value ranking intent. Physical candidates + // may rebuild a heap over finalized rates or use exact Sort/Limit; + // neither is allowed to sum raw counter samples as ranking weights. + if descending && matches!(inner.func, Some(InnerFunc::Rate | InnerFunc::Increase)) { + let intent = inner_intent(inner.func.as_ref().expect("counter function")); + let ranked = windowed_aggregate(inner, vec![], intent); + return Ok(Unresolved::Aggregate { + reduction: Reduction::Reduce(keys.into()), + measures: vec![AggIntent::TopK { + k: k as usize, + accuracy: current_accuracy(), + }], + output_names: vec![], + having: None, + child: Rc::new(ranked), + }); + } // Heavy-hitter only when ranking by an additive measure (`count` // or `sum`): that is a // first-class aggregate intent → `TopK`. Any other ranking (topk - // over avg/quantile/rate, a bare selector's raw value, all bottomk) + // over avg/quantile, a bare selector's raw value, all bottomk) // is a generic order-by-value + limit and stays as the `Sort + Limit` // operator pair. The descending-plus-measure rule is shared with the // canonicalize-pass promotion so the two cannot drift (issue #38). diff --git a/crates/integration-tests/tests/promql_to_post_asap.rs b/crates/integration-tests/tests/promql_to_post_asap.rs index f94a02af..539db7bf 100644 --- a/crates/integration-tests/tests/promql_to_post_asap.rs +++ b/crates/integration-tests/tests/promql_to_post_asap.rs @@ -831,7 +831,7 @@ fn execute_topk_reference(plan: &SummaryNode) -> Vec<(String, f64)> { } #[test] -fn planner_topk_reference_execution_matches_ground_truth() { +fn planner_heap_topk_reference_execution_matches_ground_truth() { // Pin numeric results independently of the emitted IR: swapping weights, // losing identity, changing the window, or dropping k changes the answer. for (query, expected) in [ @@ -865,8 +865,11 @@ fn planner_topk_reference_execution_matches_ground_truth() { &EqualSplitAllocator, &SeparatedTopK, ); - let candidates = strategy.replacements(&TargetSubDAG::new(&pre)); - assert!(!candidates.is_empty(), "no plan for {query}"); + // This reference executor consumes keyed heap updates. The inventory + // also contains maintained exact values followed by sort/limit; those + // have a different execution contract and must not enter this fixture. + let candidates: Vec<_> = strategy.replacements(&TargetSubDAG::new(&pre)).into_iter().filter(|candidate| matches!(&candidate.replacement, Replacement::Summary(plan) if matches!(plan.expr, SummaryExpr::SummaryEstimate { query: SketchQuery::TopK { .. }, .. }))).collect(); + assert!(!candidates.is_empty(), "no heap candidate for {query}"); for candidate in candidates { let Replacement::Summary(plan) = candidate.replacement else { panic!("expected summary plan for {query}") diff --git a/crates/types/src/post_asap/maintained_population.rs b/crates/types/src/post_asap/maintained_population.rs index 30dacf82..3939c7c1 100644 --- a/crates/types/src/post_asap/maintained_population.rs +++ b/crates/types/src/post_asap/maintained_population.rs @@ -63,7 +63,7 @@ impl CurrentSeriesInput { }; if self.metric.is_empty() || *metric != self.metric - || schema.closed + || (schema.closed && !schema.has_promql_series_identity()) || schema.time_index.is_none() { return false; diff --git a/crates/types/src/pre_asap/query_expr.rs b/crates/types/src/pre_asap/query_expr.rs index dcb73bee..dec0107a 100644 --- a/crates/types/src/pre_asap/query_expr.rs +++ b/crates/types/src/pre_asap/query_expr.rs @@ -65,6 +65,8 @@ pub enum QueryExprError { /// used by `Project`'s own `output_schema` arm instead). #[error("a scalar expression has no row schema of its own")] ScalarHasNoRowSchema, + #[error("invalid per-series sample column: {0}")] + InvalidSampleColumn(String), } // ── Leaf / supporting types ─────────────────────────────────────────────────── @@ -1451,12 +1453,23 @@ impl QueryExpr { /// one value per series, so every label column of `input` is preserved and only /// the sample value is replaced — kept named `value` so the PromQL sample-value /// convention (and any outer `SampleValue` reference) still resolves it by name. -fn per_series_reduction_schema(input: &Schema, agg: &AggIntent) -> Schema { - let value_idx = input - .column_id("value") - .or_else(|| (0..input.columns.len()).find(|&i| Some(i) != input.time_index)); +fn per_series_reduction_schema(input: &Schema, agg: &AggIntent) -> Result { + let vi = if let Some(index) = agg.input_cols().first() { + *index + } else { + super::column_resolution::resolve_column_ref(&ColumnRef::SampleValue, input) + .map_err(|error| QueryExprError::InvalidSampleColumn(error.to_string()))? + }; + if !matches!( + input.columns.get(vi).map(|column| &column.dtype), + Some(DataType::Float64 | DataType::Int64) + ) { + return Err(QueryExprError::InvalidSampleColumn(format!( + "column {vi} is not numeric" + ))); + } let mut columns = input.columns.clone(); - if let Some(vi) = value_idx { + { let mut out = agg.output_column(&columns[vi]); out.name = "value".into(); // A per-series range reduction produces a PromQL sample value, which is @@ -1466,14 +1479,14 @@ fn per_series_reduction_schema(input: &Schema, agg: &AggIntent) -> Schema { out.dtype = DataType::Float64; columns[vi] = out; } - Schema { + Ok(Schema { columns, time_index: input.time_index, unique_keys: input.unique_keys.clone(), // Per-series reduction is label-preserving: it inherits its input's // completeness (an open scan stays open; a closed one stays closed). closed: input.closed, - } + }) } /// The output schema of an `Aggregate { reduction, measures }` over `in_schema` — @@ -1500,7 +1513,7 @@ pub fn aggregate_output_schema( 1, "a per-entity reduction is single-aggregate" ); - return Ok(per_series_reduction_schema(in_schema, &measures[0])); + return per_series_reduction_schema(in_schema, &measures[0]); } Reduction::Reduce(by) => by, }; @@ -1562,6 +1575,11 @@ pub fn aggregate_output_schema( .and_then(|id| in_schema.columns.get(*id)) .unwrap_or(&probe); let mut out = intent.output_column(in_col); + // A global extremum emits NULL for an empty input, even if its input + // column is non-nullable. Grouped extrema only emit existing groups. + if by.is_empty() && matches!(intent, AggIntent::Min { .. } | AggIntent::Max { .. }) { + out.nullable = true; + } if let Some((arg, _)) = intent .arg_selector_columns(in_schema) .map_err(QueryExprError::InvalidScalarSignature)? @@ -2281,6 +2299,26 @@ mod tests { assert_eq!(back, s); } + // Nested temporal aggregation must replace the sample, never the grouping label. + #[test] + fn temporal_reduction_of_grouped_sum_preserves_job() { + let input = Schema::new(vec![ + col("job", DataType::Utf8, true), + col("sum", DataType::Float64, false), + ]); + for aggregate in [ + AggIntent::Avg { col: None }, + AggIntent::Avg { col: Some(1) }, + AggIntent::Rate, + ] { + let output = + aggregate_output_schema(&input, &Reduction::PerEntity, &[aggregate], &[]).unwrap(); + assert_eq!(output.columns[0], input.columns[0]); + assert_eq!(output.columns[1].name, "value"); + assert_eq!(output.columns[1].dtype, DataType::Float64); + } + } + #[test] fn per_series_rate_preserves_labels() { // A per-series range reduction (`rate`) is label-preserving: it produces diff --git a/crates/types/src/pre_asap/schema.rs b/crates/types/src/pre_asap/schema.rs index 81b4ae1f..fbcb9d67 100644 --- a/crates/types/src/pre_asap/schema.rs +++ b/crates/types/src/pre_asap/schema.rs @@ -154,7 +154,21 @@ pub struct Schema { pub closed: bool, } +/// Reserved physical row column carrying canonical JSON of a complete PromQL +/// label map. `$` cannot occur in a user PromQL label name. +pub const PROMQL_SERIES_IDENTITY: &str = "$promql_series_identity"; + impl Schema { + pub fn has_promql_series_identity(&self) -> bool { + self.closed + && self.columns.iter().any(|column| { + column.name == PROMQL_SERIES_IDENTITY + && column.dtype == DataType::Utf8 + && !column.nullable + && column.table.is_none() + }) + } + /// Construct a `Schema` from columns alone — no time index, no /// unique-key constraint. Used by `Scan` over a tabular source /// when the catalog supplies no primary-key metadata.