diff --git a/crates/asap-physical-operators/src/operators/summary/mod.rs b/crates/asap-physical-operators/src/operators/summary/mod.rs index 5dab8109f..d14f147bb 100644 --- a/crates/asap-physical-operators/src/operators/summary/mod.rs +++ b/crates/asap-physical-operators/src/operators/summary/mod.rs @@ -99,8 +99,14 @@ impl Operator { ) -> Result { 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) { @@ -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")); @@ -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)?; } diff --git a/crates/asap-physical-operators/src/summary_kernels/factory.rs b/crates/asap-physical-operators/src/summary_kernels/factory.rs index be819f8a7..c63b49fc5 100644 --- a/crates/asap-physical-operators/src/summary_kernels/factory.rs +++ b/crates/asap-physical-operators/src/summary_kernels/factory.rs @@ -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); @@ -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 } diff --git a/crates/asap-physical-operators/src/summary_kernels/univmon.rs b/crates/asap-physical-operators/src/summary_kernels/univmon.rs index 23114ab1b..13294f2b4 100644 --- a/crates/asap-physical-operators/src/summary_kernels/univmon.rs +++ b/crates/asap-physical-operators/src/summary_kernels/univmon.rs @@ -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, @@ -112,15 +131,25 @@ impl UnivMonAccumulator { impl AggregateCore for UnivMonAccumulator { fn approx_memory_bytes(&self) -> usize { - std::mem::size_of::().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::()) + .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 { Box::new(self.clone()) @@ -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); + } } diff --git a/crates/asap-physical-operators/tests/univmon_execution.rs b/crates/asap-physical-operators/tests/univmon_execution.rs index 6ddc1f0e7..6ac6b0e7e 100644 --- a/crates/asap-physical-operators/tests/univmon_execution.rs +++ b/crates/asap-physical-operators/tests/univmon_execution.rs @@ -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::>()), + )); + for (stream, expected) in outputs.into_iter().zip([2.0, 8.0_f64.sqrt(), 1.0]) { + let batches = stream.into_iter().collect::, _>>()?; + assert!((value(&batches[0].rows()[0]) - expected).abs() < 0.01); + } + } + Ok(()) +} diff --git a/docs/develop_docs/operator-design-acceptance.md b/docs/develop_docs/operator-design-acceptance.md index c71c80051..2bee5e151 100644 --- a/docs/develop_docs/operator-design-acceptance.md +++ b/docs/develop_docs/operator-design-acceptance.md @@ -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.