Skip to content
Open
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
255 changes: 255 additions & 0 deletions control_plane/src/planner_selection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,30 @@ pub enum SelectionError {
NoLegalCandidate,
#[error("ASAPPlanner sketch strategy produced a logical rewrite instead of a summary")]
UnexpectedRewrite,
/// A root offered materially different legal alternatives and the cost
/// model priced none of them. Selection must not resolve that group from
/// candidate discovery order — registration order is not optimizer policy.
#[error(
"cost model reported no comparable cost for target {target_id}: \
{candidate_count} legal alternatives ({strategies}) are unpriced, so selection \
has no evidence to rank them"
)]
CostUnavailable {
target_id: String,
candidate_count: usize,
strategies: String,
/// Per-candidate diagnostic, mirroring the selection trace entries.
candidates: Vec<CostUnavailableCandidate>,
},
}

/// One unpriced alternative in a [`SelectionError::CostUnavailable`] group.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct CostUnavailableCandidate {
pub candidate_id: String,
pub strategy: String,
pub replacement_kind: String,
pub provenance: String,
}

/// Versioned diagnostic identity over existing canonical IR, never Rc or rank IDs.
Expand Down Expand Up @@ -210,6 +234,17 @@ pub fn archive_only(intent: &AggIntent) -> bool {
)
}

/// Whether this alternative keeps the target pre-ASAP — the raw/exact
/// fallback the planner emits when an intent has no summary realization.
/// Matched on the post-ASAP IR rather than on rationale text, so the check
/// survives rewording upstream.
fn keeps_pre_asap(candidate: &asap_aware_mapping::ReplacementSubDAG) -> bool {
let Replacement::Summary(node) = &candidate.replacement else {
return false;
};
matches!(node.expr, SummaryExpr::KeepPreAsap(_))
}

/// Preserve an unsupported subtree explicitly at the post-ASAP boundary.
pub fn keep_pre_asap(expr: &QueryExpr) -> Result<Rc<SummaryNode>, SelectionError> {
let schema = expr.output_schema()?;
Expand Down Expand Up @@ -427,6 +462,75 @@ fn select_workload_impl(
}).collect::<Vec<_>>();
*trace = serde_json::json!({ "schema_version": 1, "group_id_scope": "this_selection", "groups": groups });
}
// A group with two or more materially different legal alternatives and no
// comparable cost cannot be resolved on evidence. Selection would fall back
// to candidate discovery order, which makes strategy registration order an
// undeclared optimizer policy — so fail with the candidate identities and
// let the caller supply costs or pick a documented policy instead.
if let Some(unpriced) = space
.cost_sorted(cost_model)
.iter()
.find(|group| {
group.candidates.len() > 1
&& group
.candidates
.iter()
.all(|candidate| cost_model.candidate_cost(candidate, &TargetSubDAG::with_consumer_count(group.target, group.consumer_count)).is_none_or(|cost| !cost.0.is_finite()))
// The complaint is specifically a *silent raw fallback*: the
// discovery-order winner keeps the subtree pre-ASAP while a
// realizable alternative sits behind it, unranked. A group
// whose order-chosen candidate is already realizable is not
// resolved by registration order in any way a reader would
// call raw, so it keeps planning.
&& keeps_pre_asap(group.candidates[0])
&& group.candidates[1..]
.iter()
.any(|candidate| !keeps_pre_asap(candidate))
})
.map(|group| {
let candidates = group
.candidates
.iter()
.map(|candidate| CostUnavailableCandidate {
candidate_id: replacement_identity(
group.target,
&candidate.replacement,
&accuracy,
)
.unwrap_or_else(|| "unidentified".into()),
strategy: candidate.strategy.to_string(),
replacement_kind: match &candidate.replacement {
Replacement::Summary(_) => "summary",
Replacement::Rewrite(_) => "rewrite",
Replacement::ExactComposition(_) => "exact_composition",
}
.to_string(),
provenance: format!("{:?}", candidate.provenance),
})
.collect::<Vec<_>>();
SelectionError::CostUnavailable {
target_id: target_identity(group.target, &accuracy)
.unwrap_or_else(|| "unidentified".into()),
candidate_count: candidates.len(),
strategies: candidates
.iter()
.map(|candidate| candidate.strategy.as_str())
.collect::<Vec<_>>()
.join(", "),
candidates,
}
})
{
if let Some(trace) = trace.as_deref_mut() {
trace["unresolved_group"] = serde_json::json!({
"reason": "cost_unavailable",
"policy": "fail_loudly",
"detail": unpriced.to_string(),
});
}
return Err(unpriced);
}

let roots = space
.roots
.iter()
Expand Down Expand Up @@ -819,3 +923,154 @@ mod workload_tests {
);
}
}

#[cfg(test)]
mod cost_unavailable_selection {
use super::*;
use crate::physical::post_asap::cost_model::ControlPlaneCostModel;

fn select(
query: &str,
) -> Result<(Vec<(usize, Rc<SummaryNode>)>, serde_json::Value), SelectionError> {
let accuracy = AccuracyTarget::Epsilon(0.05);
let root =
crate::query_parser::parse_query_expr_canonical(query, accuracy.clone()).unwrap();
select_workload_with_accuracy_model_and_trace(
vec![(0, Rc::new(root))],
accuracy.clone(),
&ControlPlaneCostModel::new(accuracy),
&asap_aware_mapping::NoAccuracyEvidence,
&asap_aware_mapping::DefaultAccuracyModel,
)
}

/// `avg by (job) (data)` offers two materially different legal
/// alternatives for the same root: `SketchAlgorithmStrategy`'s `Avg`
/// pass-through (no summary realization exists, so it degrades to raw)
/// and `SemanticEquivalentRewriteStrategy`'s realizable `sum / count`
/// rewrite. Neither is priced.
///
/// Before this guard, selection preserved candidate discovery order, and
/// the deployment registers the sketch strategy first — so the raw
/// pass-through won without any cost evidence, making strategy
/// registration order an undeclared optimizer policy.
#[test]
fn unpriced_alternatives_fail_instead_of_resolving_on_discovery_order() {
let error = select("avg by (job) (data)").expect_err("must not silently select a root");
let SelectionError::CostUnavailable {
target_id,
candidate_count,
candidates,
..
} = &error
else {
panic!("expected a typed cost-unavailable error, got: {error}");
};
assert!(target_id.starts_with("asap-explain-v1:target:"), "{error}");
assert_eq!(*candidate_count, 2, "{error}");

// The diagnostic has to name both alternatives, so a reader can tell
// which inputs the cost model owes rather than only that ranking failed.
let strategies: Vec<&str> = candidates
.iter()
.map(|candidate| candidate.strategy.as_str())
.collect();
assert!(
strategies.contains(&"SketchAlgorithmStrategy")
&& strategies.contains(&"SemanticEquivalentRewriteStrategy"),
"{strategies:?}"
);
let kinds: Vec<&str> = candidates
.iter()
.map(|candidate| candidate.replacement_kind.as_str())
.collect();
assert!(
kinds.contains(&"summary") && kinds.contains(&"rewrite"),
"the group must be materially different, not two rankings of one shape: {kinds:?}"
);
assert!(
candidates.iter().all(|candidate| candidate
.candidate_id
.starts_with("asap-explain-v1:candidate:")),
"{candidates:?}"
);
}

/// A root whose alternatives the model does price still plans. The guard
/// must fire on missing evidence, not on every unpriced candidate.
#[test]
fn priced_roots_still_select() {
let (roots, trace) =
select("quantile_over_time(0.9, m[1m])").expect("a priced root still plans");
assert_eq!(roots.len(), 1);
assert!(
trace.get("unresolved_group").is_none(),
"a resolved selection must not carry an unresolved-group diagnostic: {trace}"
);
}
}

#[cfg(test)]
mod probe_721_scope {
use super::*;
use crate::physical::post_asap::cost_model::ControlPlaneCostModel;

#[test]
fn survey() {
let queries = [
"avg by (job) (data)",
"rate(asap_demo_counter_total[5s])",
"increase(asap_demo_counter_total[5s])",
"sum(sum_over_time(asap_demo_gauge[5s]))",
"quantile_over_time(0.5, asap_demo_latency_ms[5s])",
"topk(1, sum_over_time(asap_demo_gauge[5s]))",
"topk(1, count_over_time(asap_demo_gauge[5s]))",
"sum by (zone) (http_requests_total)",
"count_over_time(m[1m])",
];
let accuracy = AccuracyTarget::Epsilon(0.05);
for q in queries {
let Ok(root) = crate::query_parser::parse_query_expr_canonical(q, accuracy.clone())
else {
eprintln!("SURVEY {q} -> parse error");
continue;
};
let cost_model = ControlPlaneCostModel::new(accuracy.clone());
let strategies = replacement_strategies(
&cost_model,
&asap_aware_mapping::NoAccuracyEvidence,
&asap_aware_mapping::DefaultAccuracyModel,
);
let space = asap_aware_mapping::search_workload_with_targets(
vec![(0usize, Rc::new(root), Some(accuracy.clone()))],
&strategies,
&asap_aware_mapping::DefaultAccuracyModel,
);
let mut mixed_unpriced = 0;
let mut rank0_passthrough = 0;
for group in space.cost_sorted(&cost_model) {
if group.candidates.len() < 2 {
continue;
}
let target = TargetSubDAG::with_consumer_count(group.target, group.consumer_count);
let all_unpriced = group.candidates.iter().all(|c| {
cost_model
.candidate_cost(c, &target)
.is_none_or(|x| !x.0.is_finite())
});
let kinds: std::collections::HashSet<_> = group
.candidates
.iter()
.map(|c| std::mem::discriminant(&c.replacement))
.collect();
if all_unpriced && kinds.len() > 1 {
mixed_unpriced += 1;
if group.candidates[0].rationale.contains("pass-through") {
rank0_passthrough += 1;
}
}
}
eprintln!("SURVEY {q} -> mixed_unpriced_groups={mixed_unpriced} rank0_is_passthrough={rank0_passthrough}");
}
}
}
11 changes: 8 additions & 3 deletions data_plane/tests/support/current_series_process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,12 @@ async fn current_series_quantiles_topk_share_and_replace_values() {
queries.push(format!("topk({k}, a)"));
queries.push(format!("topk by (job) ({k}, a)"));
}
for operation in ["sum", "count", "avg"] {
// `avg` has no summary realization, so its root offers a raw pass-through
// and a realizable sum/count rewrite with no comparable cost between them.
// Planning now fails loudly on that rather than resolving it from strategy
// registration order (#721); the behaviour is asserted in
// `control_plane::planner_selection`'s cost-unavailable regression tests.
for operation in ["sum", "count"] {
queries.push(format!("{operation}(a)"));
queries.push(format!("{operation} by (job) (a)"));
}
Expand Down Expand Up @@ -252,7 +257,7 @@ async fn current_series_quantiles_topk_share_and_replace_values() {
"{body}"
);
}
for operation in ["sum", "count", "avg"] {
for operation in ["sum", "count"] {
for text in [
format!("{operation}(a)"),
format!("{operation} by (job) (a)"),
Expand Down Expand Up @@ -314,7 +319,7 @@ async fn current_series_quantiles_topk_share_and_replace_values() {
&body,
)
.await;
for operation in ["sum", "count", "avg"] {
for operation in ["sum", "count"] {
for text in [
format!("{operation}(a)"),
format!("{operation} by (job) (a)"),
Expand Down
11 changes: 4 additions & 7 deletions data_plane/tests/support/issue_701_702_process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,10 +52,12 @@ fn queries() -> Vec<(String, u64, u64)> {
));
queries.push((format!("quantile by(job)({q}, issue701_data)"), 1, 1));
}
for operation in ["sum", "count", "avg", "min", "max"] {
// `avg` is excluded: it plans to a raw pass-through against a realizable
// rewrite with no comparable cost, which now fails loudly (#721).
for operation in ["sum", "count", "min", "max"] {
queries.push((format!("{operation}_over_time(issue701_data[5m])"), 300, 30));
}
for operation in ["sum", "count", "avg"] {
for operation in ["sum", "count"] {
queries.push((format!("{operation}(issue701_data)"), 1, 1));
queries.push((format!("{operation} by(job)(issue701_data)"), 1, 1));
}
Expand All @@ -77,11 +79,6 @@ fn queries() -> Vec<(String, u64, u64)> {
300,
60,
));
queries.push((
"avg_over_time(issue701_data[5m]) / quantile_over_time(0.5, issue701_data[5m])".into(),
300,
30,
));
queries
}

Expand Down
Loading