Skip to content
Merged
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
71 changes: 70 additions & 1 deletion crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<QueryExpr>, _target: &AccuracyTarget) -> Proposals {
Proposals::default()
}
}

// ── Realization: how one AggIntent may be realised ───────────────────────
Expand Down Expand Up @@ -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<QueryExpr>, 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
Expand Down Expand Up @@ -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)
}

Expand Down Expand Up @@ -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
Expand Down
76 changes: 4 additions & 72 deletions crates/asap-physical-operators/src/physical_planner/promql_rows.rs
Original file line number Diff line number Diff line change
@@ -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.
Expand All @@ -22,78 +22,10 @@ pub fn decode_series_identity(encoded: &str) -> Result<BTreeMap<String, String>,
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<QueryExpr, Error> {
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
Expand Down
Loading
Loading