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
67 changes: 52 additions & 15 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4758,22 +4758,21 @@ fn collect_selected_materializations_with(
}
if let Some(readout) = readout {
let mut parameters = sketch_params_json(kind.params());
let mut item_label = None;
use planner_types::post_asap::SummaryInputExpr;
let item_label = match &input.item {
Some(SummaryInputExpr::Column(
planner_types::pre_asap::ColumnRef::Named(label),
)) => Some(label.clone()),
Some(SummaryInputExpr::Column(
planner_types::pre_asap::ColumnRef::Qualified { name, .. },
)) => Some(name.clone()),
// Legacy/direct TopK plans did not carry an item
// projection. Keep their generic item-key behavior;
// typed Planner plans name the inner aggregate
// identity explicitly through SummaryUpdate.item.
_ => None,
};
if matches!(readout, SketchQuery::TopK { .. }) {
use planner_types::post_asap::SummaryInputExpr;
item_label = match &input.item {
Some(SummaryInputExpr::Column(
planner_types::pre_asap::ColumnRef::Named(label),
)) => Some(label.clone()),
Some(SummaryInputExpr::Column(
planner_types::pre_asap::ColumnRef::Qualified { name, .. },
)) => Some(name.clone()),
// Legacy/direct TopK plans did not carry an item
// projection. Keep their generic item-key behavior;
// typed Planner plans name the inner aggregate
// identity explicitly through SummaryUpdate.item.
_ => None,
};
let mode = match &input.weight {
SummaryInputExpr::Constant(value) if *value == 1.0 => "count",
SummaryInputExpr::Column(
Expand Down Expand Up @@ -9174,6 +9173,44 @@ pub(crate) mod tests {
}
}

// Both plan transports retain nonempty typed DAGs and numeric binding keys.
#[test]
fn installed_dag_plan_round_trips_in_json_and_yaml() {
let plan = DeploymentPlanCompiler
.compile_promql(
request("roundtrip", "sum_over_time(m[1m])"),
environment(10_000),
)
.unwrap()
.precompute_plan;
assert!(!plan.executable_dags.is_empty());
let json = serde_json::to_value(&plan).unwrap();
let from_json: PrecomputePlan = serde_json::from_value(json.clone()).unwrap();
let from_yaml: PrecomputePlan =
serde_yaml::from_str(&serde_yaml::to_string(&plan).unwrap()).unwrap();
assert_eq!(serde_json::to_value(from_json).unwrap(), json);
assert_eq!(serde_json::to_value(from_yaml).unwrap(), json);
}

// Non-TopK keyed sketches retain the same item dimension as their Planner update.
#[test]
fn keyed_non_topk_materialization_preserves_item_label() {
let request = request("keyed", "quantile_over_time(0.90, m[1m])");
let mut root = request.queries[0].selected_plan_root.as_ref().clone();
let SummaryExpr::SummaryEstimate { summary_input, .. } = &mut root.expr else {
panic!("expected estimate");
};
let SummaryExpr::SummaryAgg { input, .. } = &mut Rc::make_mut(summary_input).expr else {
panic!("expected summary state");
};
input.item = Some(planner_types::post_asap::SummaryInputExpr::Column(
planner_types::pre_asap::ColumnRef::Named("customer".into()),
));
let root = Rc::new(root);
let selected = collect_selected_materializations(&root, false).unwrap();
assert_eq!(selected[0].item_label.as_deref(), Some("customer"));
}

#[test]
fn multiple_readouts_share_one_precompute_materialization() {
let mut compilation_request = request("q-p90", "quantile_over_time(0.90, m[1m])");
Expand Down
Loading
Loading