Skip to content
Draft
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
2 changes: 2 additions & 0 deletions crates/asap-aware-mapping/src/query_physical_lowering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1203,6 +1203,8 @@ fn supports_hash_aggregate(
| AggIntent::PearsonCorr { .. }
| AggIntent::Group
| AggIntent::CountValues { .. }
| AggIntent::FrequencyL2 { .. }
| AggIntent::FrequencyEntropy { .. }
)
})
}
Expand Down
55 changes: 54 additions & 1 deletion crates/asap-physical-operators/src/operators/aggregate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,15 @@ impl Operator {
}
(DataType::Float64, false)
}
Reduction::FrequencyL2(i) | Reduction::FrequencyEntropy(i) => {
if !matches!(
plain(&input, *i)?.0,
DataType::Bool | DataType::Int64 | DataType::Float64 | DataType::Utf8
) {
return Err(invalid("frequency aggregate requires a Boolean, Int64, Float64 or Utf8 identity"));
}
(DataType::Float64, false)
}
Reduction::Min(i) | Reduction::Max(i) => {
let (t, nullable) = plain(&input, *i)?;
if !ordered(t) {
Expand Down Expand Up @@ -129,6 +138,10 @@ pub enum Reduction {
Avg(usize),
Min(usize),
Max(usize),
/// L2 norm of unit-update frequencies; NULL identities are skipped.
FrequencyL2(usize),
/// Shannon entropy in bits; NULL identities are skipped.
FrequencyEntropy(usize),
/// PromQL `quantile`: linear interpolation between closest ranks.
Quantile {
column: usize,
Expand Down Expand Up @@ -205,7 +218,7 @@ async fn reduce(
.map(|&i| rows[0][i].clone())
.collect::<Vec<_>>();
for measure in measures {
result.push(reduce_one(&rows, measure, input, &mut work).await?);
result.push(reduce_one(&rows, measure, input, &mut work, context).await?);
}
workspace.grow(row_bytes(&result))?;
output.push(result);
Expand Down Expand Up @@ -243,6 +256,7 @@ async fn reduce_one(
measure: &Reduction,
input: &SchemaRef,
work: &mut Cooperative,
context: &RunContext,
) -> Result<Value, Error> {
let column = match measure {
Reduction::Count => {
Expand All @@ -251,6 +265,45 @@ async fn reduce_one(
))
}
Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i,
Reduction::FrequencyL2(column) | Reduction::FrequencyEntropy(column) => {
let mut workspace = Workspace::new(context)?;
let mut counts = BTreeMap::<Vec<u8>, u64>::new();
let mut total = 0_u64;
for row in rows {
work.checkpoint().await?;
let value = &row[*column];
if matches!(value, Value::Null) {
continue;
}
if matches!(value, Value::Float64(v) if !v.is_finite()) {
return Err(invalid(
"frequency aggregate requires a finite floating identity",
));
}
let key = value.key()?;
if !counts.contains_key(&key) {
workspace.grow(
64 + std::mem::size_of::<Vec<u8>>()
+ key.len()
+ std::mem::size_of::<u64>(),
)?;
}
*counts.entry(key).or_default() += 1;
total += 1;
}
let mut result = 0.0_f64;
for count in counts.into_values() {
work.checkpoint().await?;
let count = count as f64;
if matches!(measure, Reduction::FrequencyL2(_)) {
result = result.hypot(count);
} else {
let probability = count / total as f64;
result -= probability * probability.log2();
}
}
return Ok(Value::Float64(result));
}
Reduction::Quantile { column, q } => {
let mut values = Vec::with_capacity(rows.len());
for row in rows {
Expand Down
6 changes: 6 additions & 0 deletions crates/asap-physical-operators/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -828,6 +828,12 @@ fn bind_operation(node: &PostAsapDAGNode, inputs: &[SchemaRef]) -> Result<Operat
AggIntent::Count { .. } => Reduction::Count,
AggIntent::Sum { col } => Reduction::Sum(column(*col)?),
AggIntent::Avg { col } => Reduction::Avg(column(*col)?),
AggIntent::FrequencyL2 { col, .. } => {
Reduction::FrequencyL2(column(*col)?)
}
AggIntent::FrequencyEntropy { col, .. } => {
Reduction::FrequencyEntropy(column(*col)?)
}
AggIntent::Min { col } => Reduction::Min(column(*col)?),
AggIntent::Max { col } => Reduction::Max(column(*col)?),
_ => {
Expand Down
43 changes: 43 additions & 0 deletions crates/asap-physical-operators/tests/blocking_resources.rs
Original file line number Diff line number Diff line change
Expand Up @@ -248,3 +248,46 @@ fn weighted_summary_build_yields_within_a_batch() {
drop(output);
assert_eq!(run.retained_bytes(), 0);
}

// Frequency dictionaries count against the budget and release reservations on failure.
#[test]
fn frequency_dictionary_enforces_memory_budget() {
use asap_physical_operators::operators::Reduction;
let input = schema(1);
let mut sources = PhysicalDAG::default();
sources
.add(
0,
vec![],
Operator::source(
input.clone(),
vec![Batch::try_new(
input.clone(),
(0..64).map(|i| vec![Value::Int64(i)]).collect(),
)
.unwrap()],
)
.unwrap(),
)
.unwrap();
for (reduction, succeeds) in [
(Reduction::Count, true),
(Reduction::FrequencyL2(0), false),
(Reduction::FrequencyEntropy(0), false),
] {
let run = context(12_000);
let inputs = sources.execute(&[0], run.clone()).unwrap();
let operator =
Operator::aggregate(input.clone(), vec![], vec![("result".into(), reduction)]).unwrap();
let mut output = operator.start(inputs, run.clone()).unwrap();
let result = block_on(output.next()).unwrap();
if succeeds {
assert!(result.is_ok());
} else {
assert!(matches!(result, Err(Error::MemoryLimit)));
}
drop(result);
drop(output);
assert_eq!(run.retained_bytes(), 0);
}
}
117 changes: 117 additions & 0 deletions crates/asap-physical-operators/tests/physical_semantics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -714,3 +714,120 @@ fn empty_exact_summary_extrema_agree_with_ordinary_aggregation() {
assert!(matches!(rows[0][0], Value::Null));
}
}

// Exact frequency intents bind to native reducers without a sketch or numeric key conversion.
#[test]
fn exact_frequency_intents_execute_typed_keys_and_empty_input() {
use asap_physical_operators::physical_planner::compile_node;
use planner_types::{
post_asap::ExecutionDataState,
pre_asap::{AggIntent, GroupKeys, Reduction as PlanReduction},
types::AccuracyTarget,
};
for (dtype, values) in [
(
DataType::Utf8,
vec![Value::Utf8("a".into()), Value::Utf8("b".into())],
),
(
DataType::Int64,
vec![
Value::Int64(9_007_199_254_740_992),
Value::Int64(9_007_199_254_740_993),
],
),
(DataType::Bool, vec![Value::Bool(false), Value::Bool(true)]),
(
DataType::Float64,
vec![Value::Float64(-0.0), Value::Float64(1.0)],
),
] {
let input = schema(&[("key", dtype, true)]);
for (measure, name, expected) in [
(
AggIntent::FrequencyL2 {
col: Some(0),
accuracy: AccuracyTarget::Exact,
},
"frequency_l2",
8.0_f64.sqrt(),
),
(
AggIntent::FrequencyEntropy {
col: Some(0),
accuracy: AccuracyTarget::Exact,
},
"frequency_entropy",
1.0,
),
] {
let node = PostAsapDAGNode {
id: PostAsapNodeId(1),
payload: PostAsapOperatorPayload::Relational {
operator: ValueOperation::Aggregate {
reduction: PlanReduction::Reduce(GroupKeys::none()),
measures: vec![measure],
output_names: vec![name.into()],
filters: vec![],
having: None,
},
},
output_state: ExecutionDataState::QUERY_ROWS,
output_schema: (*schema(&[(name, DataType::Float64, false)])).clone(),
guarantee: None,
};
let operator = compile_node(&node, std::slice::from_ref(&input))
.expect("exact frequency intent binds");
let rows = values
.iter()
.flat_map(|v| [vec![v.clone()], vec![v.clone()]])
.chain([vec![Value::Null]])
.collect();
let result = unary(input.clone(), vec![rows], operator.clone());
assert!(matches!(result[0][0], Value::Float64(v) if (v - expected).abs() < 1e-12));
for batches in [vec![], vec![vec![vec![Value::Null]]]] {
let result = unary(input.clone(), batches, operator.clone());
assert!(matches!(result[0][0], Value::Float64(0.0)));
}
}
}
}

// Each group gets its own frequency population, including one canonical signed-zero identity.
#[test]
fn exact_frequency_grouping_and_entropy_bits() {
let input = schema(&[
("group", DataType::Int64, false),
("key", DataType::Float64, true),
]);
let operator = Operator::aggregate(
input.clone(),
vec![0],
vec![
("l2".into(), Reduction::FrequencyL2(1)),
("entropy".into(), Reduction::FrequencyEntropy(1)),
],
)
.unwrap();
let rows = vec![
vec![Value::Int64(1), Value::Float64(-0.0)],
vec![Value::Int64(1), Value::Float64(0.0)],
vec![Value::Int64(1), Value::Float64(0.0)],
vec![Value::Int64(1), Value::Float64(1.0)],
vec![Value::Int64(2), Value::Float64(2.0)],
vec![Value::Int64(2), Value::Null],
vec![Value::Int64(3), Value::Null],
];
let result = unary(input.clone(), vec![rows], operator.clone());
assert_eq!(result.len(), 3);
let expected_entropy = -0.75_f64 * 0.75_f64.log2() - 0.25_f64 * 0.25_f64.log2();
for (row, l2, entropy) in [
(&result[0], 10.0_f64.sqrt(), expected_entropy),
(&result[1], 1.0, 0.0),
(&result[2], 0.0, 0.0),
] {
assert!(matches!(row[1], Value::Float64(v) if (v - l2).abs() < 1e-12));
assert!(matches!(row[2], Value::Float64(v) if (v - entropy).abs() < 1e-12));
}
assert!(unary(input, vec![], operator).is_empty());
}
47 changes: 47 additions & 0 deletions docs/develop_docs/planner-layering-status.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
# Planner-layering implementation status

Audience: planner developers. Audit baseline: PR #557 (`e0e1e2d7`), against
[the #509 proposal](../design_docs/proposals/planner-layering.md). The proposal
is a target contract, not a statement that its examples execute today.

| Proposal contract | Evidence at #557 | Remaining scope |
| --- | --- | --- |
| Language frontends and common logical IR | SQL/PromQL/MetricsQL lower to unified operators and scalars. | Example 2 SQL frequency L2 and entropy idioms are not recognized. Preserve alias lineage, filters, NULL groups, empty inputs, count overflow and entropy units when adding recognition. |
| Local exact and summary alternatives | `replacement::summary_candidates`, realization rules and candidate inventory exist; supplied accuracy models reach Pass 1. | Specialized entropy/norm families in Example 2 are illustrative, not registered families. UnivMon certifies only unit-update total count; L2, entropy and cardinality epsilon/delta bounds need verified evidence or a deployment model. |
| Summary-capability sharing | CSE interns structurally identical producers, including states with different readers. | It does not enumerate all partial sharing partitions or resize compatible states to the strictest consumer. Example 2's 37 candidates are not an acceptance result. |
| Window composition | Mergeable state IR/native merge exists; physical pane compatibility and reuse cost helpers exist. | Automatic logical sliding/tumbling/EH alternatives over differing windows, boundary coverage and error proofs are absent. A merge kernel alone does not implement Examples 1/3. |
| Physical materialization | Ephemeral/prepared/shared/continuously maintained lifecycle alternatives, costing, capabilities and latency checks exist. | Incremental query-time pane retention, historical backfill and the complete Example 4 matrix need executable implementations and explicit state/input contracts. |
| Whole-workload selection | One unified selected DAG; shared states are interned and costed across their consumers. | `replacement.rs` documents its selection as non-exhaustive over interacting choices. The proposal's cheapest complete candidate guarantee and 54/156 inventories need a complete workload search/selection path. |
| Deployment inputs and execution | `PlanningModels` bundles cost, accuracy, evidence and capabilities; native typed UnivMon supports one build with three readouts. | At #557 native exact frequency L2/entropy fallback is absent. End-to-end SQL Example 2 is not established by the native UnivMon fixture. |
| Subtract/delete, parallelism, partitioning and resource planning | Some runtime memory/cancellation limits and maintenance capability flags exist. | These remain proposal TODOs; capability flags do not supply missing IR operators or a physical resource search. |

## Follow-up sequence

1. **Exact frequency execution.** Bind the existing L2/entropy intents to native
reducers with typed identities, bits for entropy, NULL skipping, zero for an
empty population, grouping, memory accounting and cooperative cancellation.
This change supplies the fallback prerequisite; it does not certify UnivMon.
2. **SQL frequency recognition.** Add narrow, proven idiom recognition while
retaining the original exact relational computation. Entropy with `LN` is in
nats; the core intent is in bits. SQL `COUNT(*) GROUP BY nullable_key` counts
a NULL group, while the frequency intents skip NULL. Do not erase these
differences or SQL's NULL result for `SUM` over no groups.
3. **Summary sharing alternatives.** Enumerate compatible consumer partitions
and size shared states against all consumers using the supplied accuracy
model. Keep independent candidates. Tests must inspect the inventory and
selected producer count, rather than manually construct a shared state.
4. **Tumbling-window composition.** Start with aligned, fixed windows and
mergeable KLL, then add query-time retention and ingestion-time lifecycle
choices. Validate offsets, boundary alignment, retention and raw fallback.
5. **Sliding/EH composition.** Separate changes for overlapping active windows
and historical bucket coverage/error. Do not reuse a merge guarantee as a
boundary-error guarantee.
6. **Complete physical candidate selection.** Explore interacting workload
choices and lifecycle assignments, charge each shared producer once and
check all consumers' accuracy/latency/capabilities. A bounded exhaustive
implementation must fail explicitly on its budget instead of truncate.
7. **Proposal TODOs.** Design subtract/delete and parallelism/partitioning/
resource inputs before implementing their planning choices.

The examples' numerical candidate counts depend on their stated rule sets.
Tests should establish those rule sets explicitly before asserting the counts.
Loading