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
95 changes: 88 additions & 7 deletions crates/asap-aware-mapping/src/pass/major.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,14 +8,15 @@

use std::rc::Rc;

use asap_types::post_asap::{share_common_summary_subtrees, SummaryNode};
use asap_types::pre_asap::query_expr::QueryExpr;
use asap_types::types::AccuracyTarget;

use super::{OptimizationInput, OptimizationPass, OptimizeError, PlanOutput, QueryLifecyclePlan};
use crate::replacement::{default_strategies_with_evidence, search_workload_with_targets};
use crate::summary_maintenance_lifecycle::{
assemble_selected_dag_with_summary_maintenance_lifecycles,
global_selection_with_summary_maintenance_lifecycles, WorkloadDemand,
global_selection_with_summary_maintenance_lifecycles, plan_assembled_dag, shared_state_cost,
summary_states, WorkloadDemand,
};

/// The shipped algorithm. Unit struct: its strategy set is the crate default,
Expand Down Expand Up @@ -79,13 +80,93 @@ impl OptimizationPass for MajorPass {
.workload_entries_by_target(demand.workload, &entry_indices)
.map_err(|error| OptimizeError::LifecycleSelection(error.into()))?;

let mut plans = Vec::with_capacity(space.roots.len());
// Assemble every root, then intern structurally identical summary
// producers across them once, so two queries that selected the same
// `SummaryAgg` reach one `Rc` (consumers dedupe states by pointer).
let mut assembled = Vec::with_capacity(space.roots.len());
for (entry_index, root) in &space.roots {
let plan = assemble_selected_dag_with_summary_maintenance_lifecycles(
&selection,
let dag = selection
.assemble_selected_dag(root)
.map_err(|source| OptimizeError::LifecycleAssembly {
entry_index: *entry_index,
source: source.into(),
})?
.ok_or_else(|| self.missing_group(*entry_index))?;
assembled.push(dag);
}
let interned =
share_common_summary_subtrees(assembled.iter().cloned().enumerate().collect());
let states: Vec<_> = interned
.iter()
.map(|(_, dag)| summary_states(dag))
.collect();

// A state reached from several roots is planned once against all of
// their entries, in every plan that reaches it, so each plan picks
// the same lifecycle for it. When that union cannot be costed the
// roots keep their own, unshared DAG and entries.
let mut shared_entries: Vec<(Rc<SummaryNode>, Option<Vec<usize>>)> = Vec::new();
for (position, (entry_index, root)) in space.roots.iter().enumerate() {
for state in &states[position] {
if shared_entries.iter().any(|(s, _)| Rc::ptr_eq(s, state)) {
continue;
}
let readers: Vec<_> = (0..space.roots.len())
.filter(|&other| states[other].iter().any(|s| Rc::ptr_eq(s, state)))
.map(|other| &space.roots[other].1)
.collect();
if readers.iter().all(|reader| Rc::ptr_eq(reader, root)) {
continue;
}
let mut entries: Vec<usize> = readers
.iter()
.flat_map(|reader| bindings[&Rc::as_ptr(reader)].iter().copied())
.collect();
entries.sort_unstable();
entries.dedup();
let cost = shared_state_cost(
state,
WorkloadDemand {
entry_indices: &entries,
..demand
},
lifecycle.now_ms,
lifecycle.horizon,
lifecycle.capabilities,
models.cost,
)
.map_err(|source| OptimizeError::LifecycleAssembly {
entry_index: *entry_index,
source: source.into(),
})?;
shared_entries.push((Rc::clone(state), cost.map(|_| entries)));
}
}

let mut plans = Vec::with_capacity(space.roots.len());
for (position, (entry_index, root)) in space.roots.iter().enumerate() {
let mut entries = bindings[&Rc::as_ptr(root)].clone();
let mut dag = Rc::clone(&interned[position].1);
for (state, shared) in &shared_entries {
if !states[position].iter().any(|s| Rc::ptr_eq(s, state)) {
continue;
}
match shared {
Some(shared) => entries.extend(shared),
None => {
entries = bindings[&Rc::as_ptr(root)].clone();
dag = Rc::clone(&assembled[position]);
break;
}
}
}
entries.sort_unstable();
entries.dedup();
let plan = plan_assembled_dag(
dag,
root,
WorkloadDemand {
entry_indices: &bindings[&Rc::as_ptr(root)],
entry_indices: &entries,
..demand
},
lifecycle.now_ms,
Expand All @@ -99,7 +180,7 @@ impl OptimizationPass for MajorPass {
})?;
plans.push(QueryLifecyclePlan {
entry_index: *entry_index,
plan: plan.ok_or_else(|| self.missing_group(*entry_index))?,
plan,
});
}
Ok(PlanOutput::new(plans))
Expand Down
5 changes: 5 additions & 0 deletions crates/asap-aware-mapping/src/pass/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,11 @@ pub struct QueryLifecyclePlan {

/// One plan per workload entry, in `QueryWorkload::entries()` order;
/// [`check_contract`] enforces that.
///
/// Plans are not deduplicated across entries: a summary state that several
/// queries share appears in each of their plans as the same `Rc` (with the
/// same lifecycle), so a consumer that deploys or costs the workload must
/// dedupe deployments by `Rc::ptr_eq` on the summary node.
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct PlanOutput {
Expand Down
2 changes: 1 addition & 1 deletion crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4256,7 +4256,7 @@ impl<Id: Clone + PartialEq> PlanSpace<Id> {
}

/// Lifecycle-aware whole-subplan costs keyed by target and candidate identity.
#[derive(Default)]
#[derive(Default, Clone)]
pub(crate) struct CandidateCostOverrides {
costs: HashMap<(*const QueryExpr, *const ReplacementSubDAG), Cost>,
raw_costs: HashMap<*const QueryExpr, Cost>,
Expand Down
Loading
Loading