diff --git a/crates/asap-aware-mapping/src/replacement.rs b/crates/asap-aware-mapping/src/replacement.rs index f8f55a1d5..fb6972510 100644 --- a/crates/asap-aware-mapping/src/replacement.rs +++ b/crates/asap-aware-mapping/src/replacement.rs @@ -560,6 +560,12 @@ pub enum ReplacementProvenance { /// [`Replacement::ExactComposition`] with /// [`OperationPlacement::Maintenance`] (issue #171). ValueOperationAtIngestionTime, + /// A finalized whole-query result over rows carrying the PromQL series + /// identity, which the logical root does not expose (see + /// [`ReplacementStrategy::propose_for_root`]). Default selection never + /// commits it, because its readout must be validated and priced by + /// deployment; otherwise it would silently replace the logical plan. + RootPhysicalRealization, } /// A candidate a strategy considered for a target but refused to propose on @@ -634,6 +640,15 @@ pub trait ReplacementStrategy { domain_error: None, } } + + /// Whole-query logical alternatives for a workload root under its + /// end-to-end `target`. These may need input rows the root does not expose + /// (for example, the PromQL series identity), so + /// [`search_workload_with_targets`] asks only workload roots, once each. + /// They decide what to compute, never placement. Default: none. + fn propose_for_root(&self, _root: &Rc, _target: &AccuracyTarget) -> Proposals { + Proposals::default() + } } // ── Realization: how one AggIntent may be realised ─────────────────────── @@ -1705,6 +1720,37 @@ impl ReplacementStrategy for SketchAlgorithmStrategy<'_> { fn propose(&self, target: &TargetSubDAG<'_>) -> Proposals { self.propose_with(target.root, None) } + + /// Heap realizations of an instant-vector ranking (current-series TopK). + /// They rank rows that carry the complete PromQL series identity, which + /// the logical root does not expose, so each is a finalized query result + /// for the identity-carrying root. Placement variants (for example, + /// fixed-window or query-time Rate aggregation) are not listed here: the + /// lifecycle assigns timing and the physical compiler reads it. + fn propose_for_root(&self, root: &Rc, target: &AccuracyTarget) -> Proposals { + let Ok(typed) = asap_types::pre_asap::schema::with_promql_series_identity(root) else { + return Proposals::default(); + }; + let typed = Rc::new(typed); + let mut proposals = self.current_series_topk_candidates(&typed, target); + for mut candidate in std::mem::take(&mut proposals.candidates) { + let Replacement::Summary(node) = candidate.replacement else { + continue; + }; + let Ok(node) = finalize_query_candidate(node, &typed) else { + continue; + }; + let duplicate = proposals.candidates.iter().any(|existing| { + matches!(&existing.replacement, Replacement::Summary(other) if *other == node) + }); + if !duplicate { + candidate.replacement = Replacement::Summary(node); + candidate.provenance = ReplacementProvenance::RootPhysicalRealization; + proposals.candidates.push(candidate); + } + } + proposals + } } /// A human-readable rationale for one candidate `Realization`, for @@ -5966,7 +6012,8 @@ fn is_cse_candidate(candidate: &ReplacementSubDAG) -> bool { } fn is_automatically_selectable(candidate: &ReplacementSubDAG, cost_model: &dyn CostModel) -> bool { - !candidate.has_missing_accuracy_evidence() + candidate.provenance != ReplacementProvenance::RootPhysicalRealization + && !candidate.has_missing_accuracy_evidence() && candidate.runtime_support_evidence(cost_model) != Some(false) } @@ -6468,6 +6515,28 @@ pub fn search_workload_with_targets<'s, Id>( .zip(targets) .filter_map(|((_, root), target)| target.map(|t| (Rc::as_ptr(root), t))) .collect(); + // Whole-root proposals join the root group before its target check. + for (index, (ptr, target)) in root_ptrs.iter().enumerate() { + if root_ptrs[..index].contains(&(*ptr, target.clone())) { + continue; + } + let group = space.groups.get_mut(ptr).expect("every root has a group"); + let root = Rc::clone(&group.target); + for strategy in strategies { + let name = strategy.name(); + let proposals = strategy.propose_for_root(&root, target); + for mut candidate in proposals.candidates { + candidate.strategy = name; + group.add_candidate(candidate); + } + group + .rejected + .extend(proposals.rejected.into_iter().map(|mut rejection| { + rejection.strategy = name; + rejection + })); + } + } let mut composition_targets: HashMap<_, Vec<_>> = HashMap::new(); for (ptr, target) in root_ptrs { composition_targets diff --git a/crates/asap-physical-operators/src/physical_planner/promql_rows.rs b/crates/asap-physical-operators/src/physical_planner/promql_rows.rs index cb6800963..8b887580c 100644 --- a/crates/asap-physical-operators/src/physical_planner/promql_rows.rs +++ b/crates/asap-physical-operators/src/physical_planner/promql_rows.rs @@ -1,7 +1,7 @@ //! A bounded PromQL source row carries the entire label set, not just labels //! mentioned by the query. The source adapter owns this lossless encoding. use super::*; -use planner_types::pre_asap::{Column, DataType, Source as LogicalSource}; +use planner_types::pre_asap::DataType; use std::rc::Rc; /// Not a legal PromQL label name, so it cannot shadow a user label. @@ -22,78 +22,10 @@ pub fn decode_series_identity(encoded: &str) -> Result, Ok(labels) } -/// Resolve the row representation before candidate search. `closed` describes -/// physical columns here: the final column contains every dynamic source label. -/// It does not assert that the query's projected labels are the full label set. -/// -/// This realization supports explicit `by` grouping and per-series computation. -/// Operators that rewrite or implicitly match dynamic label sets require their -/// own realization; they must not accidentally treat the opaque identity as a -/// user label or silently discard it. +/// Resolve the row representation before candidate search; see +/// [`planner_types::pre_asap::schema::with_promql_series_identity`]. pub fn with_series_identity(root: &QueryExpr) -> Result { - let mut root = root.clone(); - fn visit(node: &mut QueryExpr) -> Result<(), Error> { - use planner_types::pre_asap::Reduction; - match node { - QueryExpr::Scan { - source: LogicalSource::TimeSeries { .. }, - schema, - .. - } => { - if schema - .columns - .iter() - .any(|column| column.name == SERIES_IDENTITY_COLUMN) - { - return Err(invalid( - "source already contains a physical series identity", - )); - } - if schema.closed { - return Err(invalid( - "dynamic series identity requires an open PromQL source", - )); - } - schema - .columns - .push(Column::new(SERIES_IDENTITY_COLUMN, DataType::Utf8, false)); - schema.closed = true; - Ok(()) - } - QueryExpr::TimeRange { child, .. } | QueryExpr::Limit { child, .. } => { - visit(Rc::make_mut(child)) - } - QueryExpr::Aggregate { - child, reduction, .. - } => { - if matches!(reduction, Reduction::Reduce(keys) if keys.is_without()) { - return Err(invalid( - "dynamic without grouping requires label-set projection", - )); - } - visit(Rc::make_mut(child)) - } - QueryExpr::Sort { - child, - partition_by, - .. - } => { - if partition_by.is_without() { - return Err(invalid( - "dynamic without ranking requires label-set projection", - )); - } - visit(Rc::make_mut(child)) - } - _ => Err(invalid( - "operator has no dynamic series-identity realization", - )), - } - } - visit(&mut root)?; - root.output_schema() - .map_err(|error| invalid(error.to_string()))?; - Ok(root) + planner_types::pre_asap::schema::with_promql_series_identity(root).map_err(invalid) } /// Construct source rows only from full identities. The named label columns diff --git a/crates/asap-physical-operators/tests/planspace_series_identity_heap.rs b/crates/asap-physical-operators/tests/planspace_series_identity_heap.rs new file mode 100644 index 000000000..2f8d5249f --- /dev/null +++ b/crates/asap-physical-operators/tests/planspace_series_identity_heap.rs @@ -0,0 +1,250 @@ +//! Logical heap alternatives that need the PromQL series identity are part of +//! Planner's search space: `enumerate_candidate_dags_for_root` lists +//! current-series TopK heaps without a caller-side series-identity pass, cost +//! ranking, or workload Cartesian expansion. Placement variants are not listed. +use asap_aware_mapping::{ + accuracy::{AccuracyEvidenceProvider, DefaultAccuracyModel, PropagationStats}, + cost_model::DefaultCostModel, + replacement::{default_strategies_with_evidence, ReplacementProvenance}, + search_workload_with_targets, Proposals, ReplacementStrategy, ReplacementSubDAG, TargetSubDAG, +}; +use asap_physical_operators::physical_planner::promql_rows::{ + compile_current_series_readout, SERIES_IDENTITY_COLUMN, +}; +use planner_types::{ + post_asap::*, + pre_asap::QueryExpr, + types::AccuracyTarget, + workload::{ + AccuracyRequirement, BatchEntry, DataWorkload, DurationMs, Evidence as WorkloadEvidence, + PlanningWorkload, Predictability, Query, QueryLanguage, QueryRequirements, QueryWorkload, + TimeSelection, + }, +}; +use std::rc::Rc; + +struct Evidence; +impl AccuracyEvidenceProvider for Evidence { + fn topk_max_distinct_items(&self, _: &QueryExpr) -> Option { + Some(1000) + } + fn propagation_stats( + &self, + op: &CompositionOperator, + _: &SummaryFamilyType, + _: Option<&SketchQuery>, + ) -> PropagationStats { + if matches!(op, CompositionOperator::TopKSelection) { + PropagationStats { + topk_selected_lower_bound: Some(101.), + topk_excluded_upper_bound: Some(100.), + topk_interval_failure_probability: Some(0.001), + ..Default::default() + } + } else { + Default::default() + } + } +} + +/// Forwards everything except whole-root proposals: the pre-change search. +struct LogicalOnly(Box); +impl ReplacementStrategy for LogicalOnly { + fn name(&self) -> &'static str { + self.0.name() + } + fn matches(&self, target: &TargetSubDAG<'_>) -> bool { + self.0.matches(target) + } + fn replacements(&self, target: &TargetSubDAG<'_>) -> Vec { + self.0.replacements(target) + } + fn propose(&self, target: &TargetSubDAG<'_>) -> Proposals { + self.0.propose(target) + } +} + +fn lower(query: &str, accuracy: &AccuracyTarget) -> Rc { + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::PromQL, + query_batch: Some(vec![BatchEntry { + query: Query(query.into()), + requirements: QueryRequirements { + accuracy: AccuracyRequirement::Explicit(accuracy.clone()), + ..Default::default() + }, + predictability: Predictability::Unknown, + invocations: 1, + execute_at: None, + time_selection: TimeSelection::default(), + }]), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + data_ingestion_interval: WorkloadEvidence { + value: Some(DurationMs(1_000)), + ..Default::default() + }, + ..Default::default() + }), + }; + Rc::new( + asap_frontend_promql::lower_promql_workload(&workload, 0) + .unwrap() + .remove(0), + ) +} + +type Dag = Vec<(usize, Rc)>; + +/// Candidate DAGs for query 1 of a two-query workload, with and without +/// whole-root proposals. Query 0 is a bystander that must not multiply them. +fn inventories(query: &str, accuracy: AccuracyTarget) -> (Vec, Vec) { + let roots = vec![ + ( + 0, + lower("sum by(job)(m)", &AccuracyTarget::Exact), + Some(AccuracyTarget::Exact), + ), + (1, lower(query, &accuracy), Some(accuracy)), + ]; + let full = default_strategies_with_evidence(&DefaultCostModel, &Evidence); + let logical: Vec> = + default_strategies_with_evidence(&DefaultCostModel, &Evidence) + .into_iter() + .map(|strategy| Box::new(LogicalOnly(strategy)) as Box) + .collect(); + let enumerate = |strategies: &[Box]| { + search_workload_with_targets(roots.clone(), strategies, &DefaultAccuracyModel) + .enumerate_candidate_dags_for_root(&1, 65_536) + .unwrap() + .candidates + }; + (enumerate(&full), enumerate(&logical)) +} + +fn carries_identity(dag: &Dag) -> bool { + dag.iter().any(|(_, root)| { + compile_post_asap_dag(root) + .unwrap() + .nodes + .iter() + .any(|node| { + node.output_schema + .fields + .iter() + .any(|field| field.name == SERIES_IDENTITY_COLUMN) + }) + }) +} + +/// Shared acceptance checks; returns the added identity-carrying alternatives. +fn added_alternatives(query: &str, accuracy: AccuracyTarget) -> Vec> { + let (full, logical) = inventories(query, accuracy); + for (index, dag) in full.iter().enumerate() { + assert_eq!(dag.len(), 1, "one root per candidate, no workload product"); + assert!( + !full[..index].contains(dag), + "{query}: identical DAG listed twice" + ); + } + let (added, kept): (Vec<_>, Vec<_>) = full.into_iter().partition(carries_identity); + assert_eq!( + kept, logical, + "{query}: existing candidates must be unchanged" + ); + added.into_iter().map(|mut dag| dag.remove(0).1).collect() +} + +const CURRENT_SERIES_TOPK: &str = "topk by(job)(1, m)"; + +// Instant-vector TopK lists finalized current-series heap readouts. +#[test] +fn current_series_topk_lists_heap_readouts() { + let added = added_alternatives(CURRENT_SERIES_TOPK, AccuracyTarget::Epsilon(0.1)); + assert!(!added.is_empty()); + for root in added { + assert!(!matches!(root.expr, SummaryExpr::SummaryAgg { .. })); + assert!( + compile_current_series_readout(&root).is_ok(), + "unbindable alternative {root:?}" + ); + } +} + +// Rate queries gain no fixed-window or query-time placement variants. +#[test] +fn rate_placement_variants_are_not_listed() { + for (query, accuracy) in [ + ("sum by(job)(rate(m[1m]))", AccuracyTarget::Exact), + ("topk by(job)(2, rate(m[1m]))", AccuracyTarget::Epsilon(0.1)), + ("rate(m[1m])", AccuracyTarget::Exact), + ] { + assert!(added_alternatives(query, accuracy).is_empty(), "{query}"); + } +} + +// Queries without a current-series heap realization are unchanged. +#[test] +fn unrelated_queries_keep_their_inventory() { + for (query, accuracy) in [ + ("sum by(job)(m)", AccuracyTarget::Exact), + ( + "quantile_over_time(0.9, m[1m])", + AccuracyTarget::Epsilon(0.05), + ), + ("max_over_time(m[1m])", AccuracyTarget::Exact), + ] { + assert!(added_alternatives(query, accuracy).is_empty(), "{query}"); + } +} + +// Default cost-based selection keeps the logical plan; deployment prices heaps. +#[test] +fn global_selection_never_commits_a_series_identity_heap() { + let accuracy = AccuracyTarget::Epsilon(0.1); + let root = lower(CURRENT_SERIES_TOPK, &accuracy); + let strategies = default_strategies_with_evidence(&DefaultCostModel, &Evidence); + let space = search_workload_with_targets( + vec![(0, root, Some(accuracy))], + &strategies, + &DefaultAccuracyModel, + ); + let selected = space + .global_selection(&DefaultCostModel) + .assemble_selected_dag(&space.roots[0].1) + .unwrap() + .unwrap(); + assert!(!carries_identity(&vec![(0, selected)])); +} + +// A query repeated in the workload is proposed once, not once per copy. +#[test] +fn repeated_roots_do_not_duplicate_alternatives() { + let accuracy = AccuracyTarget::Epsilon(0.1); + let strategies = default_strategies_with_evidence(&DefaultCostModel, &Evidence); + let count = |copies: usize| { + let roots = (0..copies) + .map(|id| { + ( + id, + lower(CURRENT_SERIES_TOPK, &accuracy), + Some(accuracy.clone()), + ) + }) + .collect(); + let space = search_workload_with_targets(roots, &strategies, &DefaultAccuracyModel); + space + .candidates_for_target(&space.roots[0].1) + .unwrap() + .candidates + .iter() + .filter(|candidate| { + candidate.provenance == ReplacementProvenance::RootPhysicalRealization + }) + .count() + }; + assert!(count(1) > 0); + assert_eq!(count(2), count(1)); +} diff --git a/crates/types/src/pre_asap/schema.rs b/crates/types/src/pre_asap/schema.rs index fbcb9d674..938b9b7f6 100644 --- a/crates/types/src/pre_asap/schema.rs +++ b/crates/types/src/pre_asap/schema.rs @@ -158,6 +158,71 @@ pub struct Schema { /// label map. `$` cannot occur in a user PromQL label name. pub const PROMQL_SERIES_IDENTITY: &str = "$promql_series_identity"; +/// Resolve a PromQL root to rows carrying [`PROMQL_SERIES_IDENTITY`] before +/// candidate search. `closed` describes physical columns here: the final +/// column contains every dynamic source label. It does not assert that the +/// query's projected labels are the full label set. +/// +/// This realization supports explicit `by` grouping and per-series computation. +/// Operators that rewrite or implicitly match dynamic label sets require their +/// own realization; they must not accidentally treat the opaque identity as a +/// user label or silently discard it. +pub fn with_promql_series_identity(root: &super::QueryExpr) -> Result { + use super::{QueryExpr, Reduction, Source}; + use std::rc::Rc; + fn visit(node: &mut QueryExpr) -> Result<(), String> { + match node { + QueryExpr::Scan { + source: Source::TimeSeries { .. }, + schema, + .. + } => { + if schema + .columns + .iter() + .any(|column| column.name == PROMQL_SERIES_IDENTITY) + { + return Err("source already contains a physical series identity".into()); + } + if schema.closed { + return Err("dynamic series identity requires an open PromQL source".into()); + } + schema + .columns + .push(Column::new(PROMQL_SERIES_IDENTITY, DataType::Utf8, false)); + schema.closed = true; + Ok(()) + } + QueryExpr::TimeRange { child, .. } | QueryExpr::Limit { child, .. } => { + visit(Rc::make_mut(child)) + } + QueryExpr::Aggregate { + child, reduction, .. + } => { + if matches!(reduction, Reduction::Reduce(keys) if keys.is_without()) { + return Err("dynamic without grouping requires label-set projection".into()); + } + visit(Rc::make_mut(child)) + } + QueryExpr::Sort { + child, + partition_by, + .. + } => { + if partition_by.is_without() { + return Err("dynamic without ranking requires label-set projection".into()); + } + visit(Rc::make_mut(child)) + } + _ => Err("operator has no dynamic series-identity realization".into()), + } + } + let mut root = root.clone(); + visit(&mut root)?; + root.output_schema().map_err(|error| error.to_string())?; + Ok(root) +} + impl Schema { pub fn has_promql_series_identity(&self) -> bool { self.closed diff --git a/docs/design_docs/physical-planning-and-deployment.md b/docs/design_docs/physical-planning-and-deployment.md index fe71c8c3a..7e7b3bd03 100644 --- a/docs/design_docs/physical-planning-and-deployment.md +++ b/docs/design_docs/physical-planning-and-deployment.md @@ -79,6 +79,13 @@ sample values does not preserve instant-vector semantics. Replacement, rank decrease, expiry, grouping and the required approximation guarantee must be validated before admitting that physical candidate. +Planner's candidate space decides what to compute, not placement. For an +instant-vector PromQL TopK, Planner resolves rows that carry the complete series +identity and lists the current-series heap realizations per root with the other +candidates, unranked. Precompute or query-time placement of Rate and grouped Sum +is not a separate Planner candidate: the summary maintenance lifecycle assigns +each node's timing, and the physical compiler reads it. + This is the target ownership contract. A backend path that still reconstructs operators from logical candidates has not completed this integration. diff --git a/docs/develop_docs/library-api.md b/docs/develop_docs/library-api.md index bfeea0348..92288067d 100644 --- a/docs/develop_docs/library-api.md +++ b/docs/develop_docs/library-api.md @@ -267,6 +267,28 @@ they are not necessarily a globally sortable physical-cost scalar. Unavailable cost alternatives may remain for explanation. Inspect eligibility and evidence before physical selection; do not treat their presence as deployment permission. +### Enumerate candidate DAGs per root + +```text +PlanSpace::enumerate_candidate_dags_for_root(&self, id: &Id, expansion_limit: usize) + -> Result, RealizationError> +``` + +Returns every distinct finalized DAG for one root, unranked; other roots' +choices are not multiplied in. Exceeding `expansion_limit` is an error, never a +partial inventory. + +For PromQL roots that carry a target, `search_workload_with_targets` also asks +each strategy's `ReplacementStrategy::propose_for_root`. `SketchAlgorithmStrategy` +answers an instant-vector TopK with current-series heap realizations over rows +carrying the complete series identity (`$promql_series_identity`). They are +finalized, deduplicated, and marked `ReplacementProvenance::RootPhysicalRealization`. +Callers do not apply `with_series_identity` themselves. Compile each with +`promql_rows::compile_current_series_readout`; other queries keep their previous +inventory. `global_selection` never commits these candidates; the backend +compiles and prices them. PlanSpace lists no placement variants: node timing +comes from the summary maintenance lifecycle. + ## Choose strategies and models ### Strategy options