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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions crates/plan/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -10,3 +10,4 @@ edition = "2021"
asap-ir = { path = "../ir" }
asap-sketch = { path = "../sketch" }
thiserror = "2"
serde_json = "1"
115 changes: 112 additions & 3 deletions crates/plan/src/bind.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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)"),
}
}
Expand Down Expand Up @@ -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<SummaryKind> {
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"]));
Expand Down
12 changes: 7 additions & 5 deletions crates/plan/src/boundary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
}

Expand Down
38 changes: 37 additions & 1 deletion crates/plan/src/cost_model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
14 changes: 12 additions & 2 deletions crates/sketch/src/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>,
},
/// Estimated number of distinct elements.
Cardinality,
/// Top-k most frequent (key, count) pairs.
Expand Down
Loading