Skip to content
Closed
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
39 changes: 37 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
[workspace]
members = [
"crates/asap-physical-operators",
"crates/types",
"crates/sql-function-catalog",
"crates/asap-aware-mapping",
Expand Down
15 changes: 8 additions & 7 deletions crates/asap-aware-mapping/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -221,13 +221,14 @@ pub use summary_maintenance_dag_export::{
};
pub use summary_maintenance_lifecycle::{
assemble_selected_dag_with_summary_maintenance_lifecycles,
global_selection_with_summary_maintenance_lifecycles, plan_summary_maintenance_lifecycles,
SummaryMaintenanceCapabilities, SummaryMaintenanceDeployment,
SummaryMaintenanceLifecycleAlternative, SummaryMaintenanceLifecycleAssemblyError,
SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs,
SummaryMaintenanceLifecyclePlan, SummaryMaintenanceLifecyclePlanError,
SummaryMaintenanceLifecycleRejection, SummaryMaintenanceLifecycleSelectionError,
WorkloadDemand,
enumerate_summary_maintenance_lifecycles, global_selection_with_summary_maintenance_lifecycles,
plan_summary_maintenance_lifecycles, SummaryMaintenanceCapabilities,
SummaryMaintenanceDeployment, SummaryMaintenanceLifecycleAlternative,
SummaryMaintenanceLifecycleAssemblyError, SummaryMaintenanceLifecycleCandidates,
SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleChoiceError,
SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecyclePlan,
SummaryMaintenanceLifecyclePlanError, SummaryMaintenanceLifecycleRejection,
SummaryMaintenanceLifecycleSelectionError, WorkloadDemand,
};
pub use topk_reuse::TopKLimitReuseStrategy;

Expand Down
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
Loading
Loading