diff --git a/Cargo.lock b/Cargo.lock index 7191e6d9e..17da1fc90 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -364,6 +364,7 @@ version = "0.1.0" dependencies = [ "asap-ir", "asap-sketch", + "serde_json", "thiserror", ] diff --git a/crates/plan/Cargo.toml b/crates/plan/Cargo.toml index 079e9396b..cdbd0439a 100644 --- a/crates/plan/Cargo.toml +++ b/crates/plan/Cargo.toml @@ -10,3 +10,4 @@ edition = "2021" asap-ir = { path = "../ir" } asap-sketch = { path = "../sketch" } thiserror = "2" +serde_json = "1" diff --git a/crates/plan/src/bind.rs b/crates/plan/src/bind.rs index eb407c0a7..dd9215ce8 100644 --- a/crates/plan/src/bind.rs +++ b/crates/plan/src/bind.rs @@ -134,7 +134,7 @@ fn bind_summary_agg( let state_idx = summary_col_index(node, &out_schema, by); let col = summarised_column(intent, &child_schema); - let query = estimate.then(|| readout(intent, &col)); + let query = estimate.then(|| readout(intent, &col, cost_model)); let mut state_schema = lift(&out_schema); if let Some(field) = state_schema.fields.get_mut(state_idx) { @@ -211,12 +211,23 @@ fn summarised_column(intent: &AggIntent, child_schema: &Schema) -> ColumnRef { } /// The `SummaryEstimate` readout for a sketch-bound intent. -fn readout(intent: &AggIntent, col: &ColumnRef) -> SketchQuery { +fn readout(intent: &AggIntent, col: &ColumnRef, cost_model: &dyn CostModel) -> SketchQuery { match intent { AggIntent::Quantile { q, .. } => SketchQuery::Quantile { q: *q }, AggIntent::Cardinality { .. } => SketchQuery::Cardinality, AggIntent::TopK { k, .. } => SketchQuery::TopK { k: *k }, - AggIntent::Count { .. } => SketchQuery::PointCount { key: col.clone() }, + AggIntent::Count { .. } => SketchQuery::PointCount { + key: col.clone(), + value: None, + }, + // Core doesn't know the shape of a deployment-specific `Extension` + // intent, so it can't build its readout either — delegate to the + // same `CostModel` that decided (via `realize_extension`) this + // intent gets a sketch realization at all. See `readout_extension`'s + // doc for the invariant this depends on. + AggIntent::Extension { ext_kind, payload } => { + cost_model.readout_extension(ext_kind, payload, col) + } other => unreachable!("no sketch realization for {other:?} (boundary::implementation_for)"), } } @@ -386,6 +397,104 @@ mod tests { assert_eq!(params, &SummaryParams::DDSketch { alpha: 0.01 }); } + /// A deployment-supplied `CostModel` can realize an `AggIntent::Extension` + /// intent as a real sketch instead of the default `PassThrough` (issue + /// #150) — `implement_tree_with` must consult `realize_extension` for + /// the `Extension` arm, and `readout` must consult `readout_extension` + /// to build its `SketchQuery` without panicking. + struct FrequencyCostModel; + + impl CostModel for FrequencyCostModel { + fn rank_candidates( + &self, + _intent: &AggIntent, + candidates: &[SummaryKind], + ) -> Vec { + candidates.to_vec() + } + + fn realize_extension( + &self, + ext_kind: &str, + _payload: &serde_json::Value, + ) -> crate::boundary::Implementation { + if ext_kind == "frequency" { + crate::boundary::Implementation::Sketch { + kind: SummaryKind::CountSketch, + params: SummaryParams::CountSketch { + width: 256, + depth: 4, + }, + } + } else { + crate::boundary::Implementation::PassThrough + } + } + + fn readout_extension( + &self, + ext_kind: &str, + payload: &serde_json::Value, + _col: &ColumnRef, + ) -> SketchQuery { + assert_eq!(ext_kind, "frequency"); + let value = payload["item"].as_str().map(str::to_string); + SketchQuery::PointCount { + key: ColumnRef::Named("item".into()), + value, + } + } + } + + #[test] + fn extension_intent_stays_logical_by_default() { + // Without a CostModel overriding `realize_extension`, an + // `Extension` intent must stay `PassThrough` -- today's behavior, + // unchanged. + let intent = AggIntent::Extension { + ext_kind: "frequency".to_string(), + payload: serde_json::json!({ "item": "checkout" }), + }; + let q = agg(vec![], intent, metric_scan(&[])); + let root = implement_tree(&q).unwrap(); + assert!(matches!(root.expr, SummaryExpr::Logical(_))); + } + + #[test] + fn extension_intent_binds_via_custom_cost_model() { + let intent = AggIntent::Extension { + ext_kind: "frequency".to_string(), + payload: serde_json::json!({ "item": "checkout" }), + }; + let q = agg(vec![], intent, metric_scan(&[])); + let root = implement_tree_with(&q, &FrequencyCostModel).unwrap(); + + let SummaryExpr::SummaryEstimate { + sketch_input, + query, + } = &root.expr + else { + panic!("expected SummaryEstimate root, got {:?}", root.expr); + }; + assert!(matches!( + query, + SketchQuery::PointCount { key: ColumnRef::Named(k), value: Some(v) } + if k == "item" && v == "checkout" + )); + + let SummaryExpr::SummaryAgg { sketch, params, .. } = &sketch_input.expr else { + panic!("expected SummaryAgg, got {:?}", sketch_input.expr); + }; + assert_eq!(sketch, &SummaryKind::CountSketch); + assert_eq!( + params, + &SummaryParams::CountSketch { + width: 256, + depth: 4 + } + ); + } + #[test] fn exact_sum_binds_accumulator_without_estimate() { let q = agg(vec![2], AggIntent::Sum { col: None }, metric_scan(&["job"])); diff --git a/crates/plan/src/boundary.rs b/crates/plan/src/boundary.rs index a66a6cd74..c5020c16d 100644 --- a/crates/plan/src/boundary.rs +++ b/crates/plan/src/boundary.rs @@ -192,11 +192,13 @@ pub fn implementation_for_with(intent: &AggIntent, cost_model: &dyn CostModel) - AggIntent::Group | AggIntent::CountValues { .. } => Implementation::PassThrough, // ── Extension (deployment-model-specific, issue #131) — core has no - // realization opinion for a shape it doesn't know. The owning - // deployment model is expected to bind it via its own L4 rules - // before this generic boundary pass ever sees it; if one reaches - // here unbound, pass it through rather than guessing. - AggIntent::Extension { .. } => Implementation::PassThrough, + // realization opinion for a shape it doesn't know, so it defers + // entirely to the `CostModel` (issue #150): `realize_extension` + // defaults to `PassThrough`, preserving today's behavior for + // every deployment that doesn't override it. + AggIntent::Extension { ext_kind, payload } => { + cost_model.realize_extension(ext_kind, payload) + } } } diff --git a/crates/plan/src/cost_model.rs b/crates/plan/src/cost_model.rs index 54de1b3fc..ec9405e73 100644 --- a/crates/plan/src/cost_model.rs +++ b/crates/plan/src/cost_model.rs @@ -29,7 +29,10 @@ //! byte. use asap_ir::intent_algebra::agg_intent::AggIntent; -use asap_sketch::{SummaryKind, SummaryParams}; +use asap_ir::intent_algebra::expr_ir::ColumnRef; +use asap_sketch::{SketchQuery, SummaryKind, SummaryParams}; + +use crate::boundary::Implementation; /// Ranks the candidate summary families for one [`AggIntent`], best choice /// first. @@ -76,6 +79,39 @@ pub trait CostModel { ) -> SummaryParams { crate::boundary::default_size_params(kind, intent, eps, delta) } + + /// Realize an `AggIntent::Extension { ext_kind, payload }` — a + /// deployment-specific intent shape core has no realization opinion + /// for (issue #131). `boundary::implementation_for_with` consults this + /// for every `Extension` node instead of hardcoding `PassThrough` + /// (issue #150). Default: `PassThrough` — preserves today's behavior + /// for every deployment that doesn't override this, exactly like + /// `size_params`'s default-delegates pattern above. + fn realize_extension(&self, _ext_kind: &str, _payload: &serde_json::Value) -> Implementation { + Implementation::PassThrough + } + + /// Build the `SummaryEstimate` readout for an `Extension` intent this + /// same `CostModel` realized as `Implementation::Sketch` via + /// [`realize_extension`](Self::realize_extension). Only ever called + /// when `realize_extension` returned `Sketch` for the same + /// `(ext_kind, payload)` — `bind::readout` has no other way to build a + /// `SketchQuery` for a shape core doesn't know. A deployment that + /// overrides `realize_extension` to return `Sketch` for some + /// `ext_kind` MUST also override this for that same `ext_kind`, or + /// this default panics loudly (rather than silently misinterpreting + /// `payload`) the first time that intent is actually read out. + fn readout_extension( + &self, + ext_kind: &str, + _payload: &serde_json::Value, + _col: &ColumnRef, + ) -> SketchQuery { + unimplemented!( + "CostModel::realize_extension returned Sketch for ext_kind={ext_kind:?} but \ + readout_extension wasn't overridden to match" + ) + } } /// The default cost model: preserves [`summary_candidates`]'s built-in static diff --git a/crates/sketch/src/sketch.rs b/crates/sketch/src/sketch.rs index f1547438a..931e66421 100644 --- a/crates/sketch/src/sketch.rs +++ b/crates/sketch/src/sketch.rs @@ -105,8 +105,18 @@ pub enum SummaryParams { pub enum SketchQuery { /// Extract the value at quantile rank `q` ∈ (0, 1]. Quantile { q: f64 }, - /// Estimated count / frequency for a specific key. - PointCount { key: ColumnRef }, + /// Estimated count / frequency. `key` names which column is being + /// queried (`ColumnRef::SampleValue` for the bare bucket total, with + /// `value: None`); a `Named`/`Qualified` `key` paired with + /// `value: Some(v)` is a per-item point lookup (e.g. + /// `count(cms_metric{item="checkout"})` — `key` is `item`, `value` is + /// `"checkout"`). `value` is carried here rather than resolved by the + /// `SummaryExecutor` from a `Filter` predicate because `readout`'s + /// trait signature has no tree access — see `CostModel::readout_extension`. + PointCount { + key: ColumnRef, + value: Option, + }, /// Estimated number of distinct elements. Cardinality, /// Top-k most frequent (key, count) pairs.