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
11 changes: 10 additions & 1 deletion crates/asap-aware-mapping/src/pass/major.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,13 +72,22 @@ impl OptimizationPass for MajorPass {
models.cost,
)
.map_err(OptimizeError::LifecycleSelection)?;
// Each root's lifecycle is planned against the entries that consume
// it — the same binding selection costed it with — not the whole
// workload, so one query's reads never amortize another's state.
let bindings = space
.workload_entries_by_target(demand.workload, &entry_indices)
.map_err(|error| OptimizeError::LifecycleSelection(error.into()))?;

let mut plans = Vec::with_capacity(space.roots.len());
for (entry_index, root) in &space.roots {
let plan = assemble_selected_dag_with_summary_maintenance_lifecycles(
&selection,
root,
demand,
WorkloadDemand {
entry_indices: &bindings[&Rc::as_ptr(root)],
..demand
},
lifecycle.now_ms,
lifecycle.horizon,
lifecycle.capabilities,
Expand Down
45 changes: 45 additions & 0 deletions crates/planner/tests/e2e_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -392,3 +392,48 @@ async fn lifecycle_decisions_ride_inside_each_plan() {
let _: &Rc<_> = &output.plans[0].plan.root;
assert_eq!(output.dags().len(), 1);
}

/// Each root's lifecycle is planned against the entries that read it: a
/// query polled every minute and an unrelated one polled every ten minutes
/// each see only their own reads over the hour, not the workload's 66.
#[tokio::test]
async fn each_plan_counts_only_its_own_entries_reads() {
let repeating = |query: &str, interval_ms: u32| RepeatingEntry {
query: Query(query.into()),
demand: RepeatedDemand::FixedInterval(RepetitionInterval(interval_ms)),
requirements: approximate(),
predictability: Predictability::Unknown,
time_selection: TimeSelection::default(),
};
let workload = PlanningWorkload {
query_workload: QueryWorkload {
language: QueryLanguage::PromQL,
query_batch: None,
repeating_queries: Some(vec![
repeating("count_over_time(up[5m])", 60_000),
repeating("sum_over_time(latency[5m])", 600_000),
]),
},
data_workload: Some(DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,
data_ingestion_interval: Evidence {
value: Some(DurationMs(15_000)),
..Default::default()
},
..Default::default()
}),
};
let input = UserInput::new(
&workload,
FrontendInput::Promql {
now_ms: NOW_MS,
histograms: None,
},
PlanningModels::builtin(),
lifecycle().with_horizon(Horizon(3_600.0)),
);

let output = e2e_plan(input).await.expect("workload plans");
let reads: Vec<_> = output.plans.iter().map(|p| p.plan.expected_reads).collect();
assert_eq!(reads, vec![Some(60.0), Some(6.0)]);
}
Loading