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
22 changes: 13 additions & 9 deletions crates/asap-physical-operators/src/operators/summary/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,8 +99,14 @@ impl Operator {
) -> Result<Self, Error> {
crate::values::validate_family(&family)?;
validate_groups(&input, &groups)?;
if plain(&input, value)?.0 != &DataType::Float64 {
return Err(invalid("summary numeric update requires Float64"));
let dtype = plain(&input, value)?.0;
let typed_frequency = matches!(&family,
SummaryFamilyType::Sketch(kind, _) if kind.algorithm() == &planner_types::post_asap::SketchAlgorithm::UnivMon);
if *dtype != DataType::Float64
&& !(typed_frequency
&& matches!(dtype, DataType::Utf8 | DataType::Int64 | DataType::Bool))
{
return Err(invalid("summary update type is unsupported by its family"));
}
if let Some(time) = time {
if plain(&input, time)? != (&DataType::Timestamp, false) {
Expand Down Expand Up @@ -435,11 +441,10 @@ async fn build_summary(
states.get_mut(&key).expect("inserted group");
// SQL aggregates ignore NULL samples while retaining the group.
// A missing counter sample also contributes no observation.
let value = match row[value] {
Value::Float64(value) => value,
Value::Null => continue,
_ => return Err(invalid("summary update type")),
};
let value = &row[value];
if matches!(value, Value::Null) {
continue;
}
let timestamp = if let Some(time) = time {
let Value::Timestamp(time) = row[time] else {
return Err(invalid("summary time type"));
Expand All @@ -455,9 +460,8 @@ async fn build_summary(
));
}
updater
.validate_single_input(value)
.update_value(value, timestamp)
.map_err(Error::Operator)?;
updater.update_single(value, timestamp);
*previous = Some(timestamp);
memory.resize(updater.memory_usage_bytes() + *overhead)?;
}
Expand Down
20 changes: 20 additions & 0 deletions crates/asap-physical-operators/src/summary_kernels/factory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,21 @@ pub trait AccumulatorUpdater: Send {
}
}

/// Update from the typed native row. Numeric kernels retain their existing
/// Float64 contract; frequency kernels may accept nonnumeric identities.
fn update_value(
&mut self,
value: &crate::values::Value,
timestamp_ms: i64,
) -> Result<(), String> {
let crate::values::Value::Float64(value) = value else {
return Err("summary update requires Float64".into());
};
self.validate_single_input(*value)?;
self.update_single(*value, timestamp_ms);
Ok(())
}

/// Feed a single (value, timestamp_ms) pair — for SingleSubpopulation types.
fn update_single(&mut self, value: f64, timestamp_ms: i64);

Expand Down Expand Up @@ -688,6 +703,11 @@ impl AccumulatorUpdater for HllUpdater {
}

impl AccumulatorUpdater for UnivMonUpdater {
fn update_value(&mut self, value: &crate::values::Value, _: i64) -> Result<(), String> {
self.acc
.insert_value(value)
.map_err(|error| error.to_string())
}
fn is_keyed(&self) -> bool {
false
}
Expand Down
85 changes: 76 additions & 9 deletions crates/asap-physical-operators/src/summary_kernels/univmon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,25 @@ impl UnivMonAccumulator {
Ok(())
}

/// SQL identities are kept in their original type: converting Int64 to
/// Float64 would collapse neighboring keys above 2^53.
pub fn insert_value(&mut self, value: &crate::values::Value) -> Result<(), Error> {
use crate::values::Value;
match value {
Value::Null => return Ok(()),
Value::Float64(number) if number.is_finite() => return self.insert_sample(*number),
Value::Bool(_) | Value::Int64(_) | Value::Utf8(_) => {}
_ => return Err("unsupported UnivMon identity type or nonfinite sample".into()),
}
let key = value.key()?;
self.inner
.bucket_size
.checked_add(1)
.ok_or("UnivMon count overflow")?;
self.inner.insert(&DataInput::Bytes(&key), 1);
Ok(())
}

fn compatible(&self, other: &Self) -> bool {
(
self.inner.heap_size,
Expand Down Expand Up @@ -112,15 +131,25 @@ impl UnivMonAccumulator {

impl AggregateCore for UnivMonAccumulator {
fn approx_memory_bytes(&self) -> usize {
std::mem::size_of::<Self>().saturating_add(
self.inner.layer_size.saturating_mul(
self.inner
.sketch_row
.saturating_mul(self.inner.sketch_col)
.saturating_mul(16)
.saturating_add(self.inner.heap_size.saturating_mul(256)),
),
)
let key_bytes = (0..self.inner.layer_size)
.flat_map(|layer| self.inner.hh_layers[layer].heap())
.map(|item| match &item.key {
asap_sketchlib::HeapItem::String(key) => key.capacity(),
asap_sketchlib::HeapItem::Bytes(key) => key.capacity(),
_ => 0,
})
.fold(0usize, usize::saturating_add);
key_bytes
.saturating_add(std::mem::size_of::<Self>())
.saturating_add(
self.inner.layer_size.saturating_mul(
self.inner
.sketch_row
.saturating_mul(self.inner.sketch_col)
.saturating_mul(16)
.saturating_add(self.inner.heap_size.saturating_mul(256)),
),
)
}
fn clone_boxed_core(&self) -> Box<dyn AggregateCore> {
Box::new(self.clone())
Expand Down Expand Up @@ -218,4 +247,42 @@ mod tests {
assert!(UnivMonAccumulator::from_sketch(terminal).is_ok());
assert!(UnivMonAccumulator::from_sketch(UnivMon::init_univmon(4, 21, 16, 2)).is_err());
}
/// Typed keys remain distinct after merging and restoring persisted native state.
#[test]
fn typed_keys_merge_and_roundtrip() {
use crate::values::Value;
let mut left = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap();
let mut right = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap();
for _ in 0..2 {
left.insert_value(&Value::Utf8("192.0.2.1".into())).unwrap();
right
.insert_value(&Value::Utf8("192.0.2.2".into()))
.unwrap();
}
left.merge_in_place(&right).unwrap();
let bytes = left.sketch().serialize_to_bytes().unwrap();
let restored =
UnivMonAccumulator::from_sketch(UnivMon::deserialize_from_bytes(&bytes).unwrap())
.unwrap();
for (statistic, expected) in [
(SketchStatistic::Cardinality, 2.0),
(SketchStatistic::FrequencyL2, 8.0_f64.sqrt()),
(SketchStatistic::FrequencyEntropy, 1.0),
] {
assert!((restored.estimate(&statistic).unwrap() - expected).abs() < 0.01);
}
}

/// Retained variable-length identities contribute to the runtime memory reservation.
#[test]
fn memory_accounts_for_string_identities() {
let mut state = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap();
let empty = state.approx_memory_bytes();
state
.insert_value(&crate::values::Value::Utf8("x".repeat(4096).into()))
.unwrap();
assert!(state.approx_memory_bytes() >= empty + 4096);
state.clear();
assert_eq!(state.approx_memory_bytes(), empty);
}
}
73 changes: 73 additions & 0 deletions crates/asap-physical-operators/tests/univmon_execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -110,3 +110,76 @@ fn one_univmon_state_answers_distinct_l2_and_entropy() -> Result<(), Error> {
assert!(value(batches[2][0].rows().first().unwrap()) > 0.0);
Ok(())
}

/// SQL source IPs and integer identifiers enter frequency state without numeric coercion.
#[test]
fn typed_frequency_keys_preserve_identity() -> Result<(), Error> {
for (dtype, keys) in [
(
DataType::Utf8,
vec![
Value::Utf8("192.0.2.1".into()),
Value::Utf8("192.0.2.2".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)]),
] {
let input = Arc::new(planner_types::pre_asap::Schema::new(vec![Field::plain(
"src_ip", dtype, true,
)]));
let batch = Batch::try_new(
input.clone(),
vec![
vec![keys[0].clone()],
vec![keys[0].clone()],
vec![keys[1].clone()],
vec![keys[1].clone()],
vec![Value::Null],
],
)?;
let build = Operator::summary_build(input, family(), 0, None, vec![])?;
let output = build.output_schema();
let mut dag = PhysicalDAG::default();
dag.add(
0,
vec![],
Operator::source(batch.schema().clone(), vec![batch])?,
)?;
dag.add(1, vec![0], build)?;
for (id, statistic) in [
(2, SketchStatistic::Cardinality),
(3, SketchStatistic::FrequencyL2),
(4, SketchStatistic::FrequencyEntropy),
] {
dag.add(
id,
vec![1],
Operator::evaluation(output.clone(), 0, SummaryEvaluation::Sketch(statistic))?,
)?;
}
let context = RunContext::new(
Scope::Query {
evaluation_time_ms: 0,
revision: 1,
},
Limits::default(),
)?;
let outputs = block_on(futures::future::join_all(
dag.execute(&[2, 3, 4], context)?
.into_iter()
.map(|stream| stream.collect::<Vec<_>>()),
));
for (stream, expected) in outputs.into_iter().zip([2.0, 8.0_f64.sqrt(), 1.0]) {
let batches = stream.into_iter().collect::<Result<Vec<_>, _>>()?;
assert!((value(&batches[0].rows()[0]) - expected).abs() < 0.01);
}
}
Ok(())
}
12 changes: 12 additions & 0 deletions docs/develop_docs/operator-design-acceptance.md
Original file line number Diff line number Diff line change
Expand Up @@ -142,3 +142,15 @@ The public planner's `slow_cheap_raw_recompute_cannot_bypass_the_response_bound`
and `slow_raw_only_query_reports_no_latency_feasible_plan` reproduce both paths.
As with summary latency quotes, missing raw latency evidence remains unchecked;
this does not establish a latency guarantee for an unmeasured deployment.

## Planner-layering follow-up: typed frequency inputs

The native UnivMon build accepts Utf8, Int64 and Boolean frequency identities,
as well as Float64 samples. Typed keys prevent integer rounding and numeric
coercion from changing cardinality. NULL inputs are skipped as before; other
summary families keep their numeric input contracts. Variable-length heap keys
contribute to state memory accounting. `univmon_execution::typed_frequency_keys_preserve_identity`
executes one shared state with distinct/L2/entropy readers for each type,
including neighboring integers above 2^53. Kernel tests also merge and persist
string-key states and check memory after reset. SQL entropy/L2 idiom recognition
is not established by these native input tests.
Loading