From dbbd254790bb591d0ab7fb03b11053f9fd0bf3f3 Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 21:22:50 +0000 Subject: [PATCH] refactor: own SDS stored-state decoding and dataset-bound semantic identity Backend now owns SDS frame decoding, delta reconstruction, sketch readouts and native batch frames, and the stored-definition semantic fragment. The logical dataset is recorded on the materialization and in `SummarySemantics::Planner`, so equal fragments over different datasets get different definition IDs. Planner stays deployment-agnostic. Co-Authored-By: Claude Opus 5.5 --- control_plane/src/main.rs | 2 +- control_plane/src/physical/compiler.rs | 27 +- crates/asap_types/src/aggregation_config.rs | 4 + crates/asap_types/src/precompute_plan.rs | 6 +- .../asap_types/src/precompute_plan/catalog.rs | 9 +- crates/asap_types/src/semantic_fragment.rs | 602 +++++++- crates/asap_types/src/summary_semantics.rs | 47 +- .../drivers/ingest/prometheus_remote_write.rs | 2 + data_plane/src/drivers/query/servers/http.rs | 1 + .../src/precompute_engine/output_sink.rs | 1 + data_plane/src/precompute_engine/revisions.rs | 6 +- .../asap_query_engine/summary_executor.rs | 15 +- .../src/storage_engines/sketch_db/data/mod.rs | 19 +- .../sketch_db/data/native_batch.rs | 368 +++++ .../sketch_db/lifecycle/eviction.rs | 1 + .../sketch_db/query/decoders.rs | 368 ++++- .../sketch_db/query/delta_apply.rs | 1283 ++++++++++++++++- .../storage_engines/sketch_db/query/mod.rs | 1 + .../sketch_db/query/sketch_readout.rs | 92 ++ .../tests/test_utilities/engine_factories.rs | 8 + 20 files changed, 2830 insertions(+), 32 deletions(-) create mode 100644 data_plane/src/storage_engines/sketch_db/data/native_batch.rs create mode 100644 data_plane/src/storage_engines/sketch_db/query/sketch_readout.rs diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index fae7acbb2..978910147 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -199,7 +199,7 @@ struct CompileAndPublishPhysicalPlanRequest { workload_cost_evidence: Option, queries: Vec, data_workload: planner_types::workload::DataWorkload, - dataset_identity: planner_types::post_asap::LogicalDatasetIdentity, + dataset_identity: asap_types::semantic_fragment::LogicalDatasetIdentity, #[serde(rename = "collector_ids")] target_collector_ids: Vec, capability_snapshot_id: String, diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs index d90199c2a..d688165bf 100644 --- a/control_plane/src/physical/compiler.rs +++ b/control_plane/src/physical/compiler.rs @@ -200,7 +200,7 @@ pub struct TopKMembershipEvidence { #[serde(deny_unknown_fields)] pub struct PhysicalDeploymentContext { /// Semantic dataset served by this deployment's input channel; never an endpoint. - pub dataset_identity: planner_types::post_asap::LogicalDatasetIdentity, + pub dataset_identity: asap_types::semantic_fragment::LogicalDatasetIdentity, pub target: PhysicalDeploymentTarget, #[serde(rename = "collector_ids")] pub target_collector_ids: Vec, @@ -1493,11 +1493,12 @@ impl DeploymentPlanCompiler { query_id: query.query_id.clone(), reason: "persisted semantic root is absent".into(), })?; + runtime_materialization.dataset_identity = + Some(environment.dataset_identity.clone()); runtime_materialization.semantic_fragment = Some( - asap_types::semantic_fragment::SemanticFragment::from_stored_output_in_dataset( + asap_types::semantic_fragment::SemanticFragment::from_stored_output( &compiled_dag.dag, semantic_root, - environment.dataset_identity.clone(), ) .map_err(|reason| CompileError::Query { query_id: query.query_id.clone(), @@ -4690,6 +4691,7 @@ pub(crate) mod tests { first.summary_catalog.definitions, relocated.summary_catalog.definitions ); + // An ingest binding forged to another dataset is rejected, as is a missing one. let mut forged = first.precompute_plan.clone(); forged.ingest.dataset_identity.as_mut().unwrap().namespace = "other-tenant".into(); assert!(forged @@ -4699,11 +4701,28 @@ pub(crate) mod tests { .contains("dataset")); forged.ingest.dataset_identity = None; assert!(forged.validate().is_err()); + // A materialization whose recorded dataset differs from the installed ingest binding is rejected. + let mut forged = first.precompute_plan.clone(); + let materialization = forged + .materializations + .iter_mut() + .find(|m| m.semantic_fragment.is_some()) + .expect("planner materialization"); + assert_eq!( + materialization.dataset_identity, + forged.ingest.dataset_identity + ); + materialization.dataset_identity.as_mut().unwrap().dataset = "other-metrics".into(); + assert!(forged + .validate() + .unwrap_err() + .to_string() + .contains("dataset")); } fn environment(now: u64) -> PhysicalDeploymentContext { PhysicalDeploymentContext { - dataset_identity: planner_types::post_asap::LogicalDatasetIdentity { + dataset_identity: asap_types::semantic_fragment::LogicalDatasetIdentity { namespace: "test".into(), dataset: "metrics".into(), }, diff --git a/crates/asap_types/src/aggregation_config.rs b/crates/asap_types/src/aggregation_config.rs index f256d74fd..0c525a8aa 100644 --- a/crates/asap_types/src/aggregation_config.rs +++ b/crates/asap_types/src/aggregation_config.rs @@ -98,6 +98,9 @@ pub struct PrecomputeMaterialization { /// Planner-selected dependency closure ending at the persisted output. #[serde(default, skip_serializing_if = "Option::is_none")] pub semantic_fragment: Option, + /// Logical dataset that the fragment's source names resolve in. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub dataset_identity: Option, pub aggregation_type: AggregationType, pub aggregation_sub_type: String, pub parameters: HashMap, @@ -279,6 +282,7 @@ impl PrecomputeMaterialization { Self { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, aggregation_type, aggregation_sub_type, parameters, diff --git a/crates/asap_types/src/precompute_plan.rs b/crates/asap_types/src/precompute_plan.rs index 05227b679..dd0486e22 100644 --- a/crates/asap_types/src/precompute_plan.rs +++ b/crates/asap_types/src/precompute_plan.rs @@ -139,7 +139,7 @@ pub enum TimestampUnit { #[serde(deny_unknown_fields)] pub struct IngestContract { #[serde(default, skip_serializing_if = "Option::is_none")] - pub dataset_identity: Option, + pub dataset_identity: Option, pub protocol: IngestProtocol, pub endpoint_path: String, pub timestamp_unit: TimestampUnit, @@ -366,7 +366,7 @@ impl PrecomputePlan { .collect(); let dataset_identity = materializations .iter() - .filter_map(|m| m.semantic_fragment.as_ref()?.dataset_identity.clone()) + .filter_map(|m| m.dataset_identity.clone()) .next(); let plan = Self { summary_catalog: None, @@ -484,7 +484,7 @@ impl PrecomputePlan { fragment .validate() .map_err(PrecomputePlanError::CatalogContract)?; - if fragment.dataset_identity != self.ingest.dataset_identity { + if materialization.dataset_identity != self.ingest.dataset_identity { return Err(PrecomputePlanError::CatalogContract( "semantic dataset differs from installed input binding".into(), )); diff --git a/crates/asap_types/src/precompute_plan/catalog.rs b/crates/asap_types/src/precompute_plan/catalog.rs index bad95e7e5..f2322cd93 100644 --- a/crates/asap_types/src/precompute_plan/catalog.rs +++ b/crates/asap_types/src/precompute_plan/catalog.rs @@ -70,10 +70,11 @@ impl PrecomputePlan { if stored_output.fingerprint() == config.policy_fingerprint()) { found = true; - let actual = match &self.ingest.dataset_identity { - Some(dataset) => crate::semantic_fragment::SemanticFragment::from_stored_output_in_dataset(&dag, *id, dataset.clone()), - None => crate::semantic_fragment::SemanticFragment::from_stored_output(&dag, *id), - }.map_err(invalid)?; + let actual = + crate::semantic_fragment::SemanticFragment::from_stored_output( + &dag, *id, + ) + .map_err(invalid)?; if &actual != expected { return Err(invalid( "semantic definition differs from Planner-selected producer", diff --git a/crates/asap_types/src/semantic_fragment.rs b/crates/asap_types/src/semantic_fragment.rs index 390354bd9..38d89be77 100644 --- a/crates/asap_types/src/semantic_fragment.rs +++ b/crates/asap_types/src/semantic_fragment.rs @@ -1,2 +1,600 @@ -//! Planner owns the versioned semantic description and its normalization. -pub use planner_types::post_asap::SummarySemanticFragment as SemanticFragment; +//! Persistable dependency closure of a Planner-selected stored output. Backend +//! owns this format; it canonicalizes Planner's typed operation vocabulary. +//! Node hashes are local semantic references, not executable or deployed IDs. +//! The logical dataset is recorded beside the fragment in `SummarySemantics`. +use planner_types::post_asap::{ + EdgeRole, ExecutableDag, ExecutableOperatorPayload, PostAsapNodeId, SummarySchema, +}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::collections::{BTreeMap, BTreeSet}; + +/// Stable identity of a logical input dataset, independent of its endpoint. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LogicalDatasetIdentity { + pub namespace: String, + pub dataset: String, +} + +impl LogicalDatasetIdentity { + pub fn validate(&self) -> Result<(), String> { + if self.namespace.trim().is_empty() || self.dataset.trim().is_empty() { + return Err("dataset namespace and identity must be nonempty".into()); + } + Ok(()) + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct SemanticFragment { + pub format_version: u32, + pub output: String, + pub nodes: BTreeMap, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct SemanticOperation { + // Wire ownership must be Send + Sync. These values are checked against the + // Planner types on export and on recovery; arbitrary JSON is not accepted. + pub operation: serde_json::Value, + pub output_schema: serde_json::Value, + pub inputs: Vec, + /// The direct input range is supplied by the stored record, not by a query lookback. + #[serde(default, skip_serializing_if = "std::ops::Not::not")] + pub record_range: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct SemanticInput { + pub role: EdgeRole, + pub node: String, +} + +pub(crate) fn canonical_bytes(value: &impl Serialize) -> Result, String> { + fn canonical(value: serde_json::Value) -> serde_json::Value { + match value { + serde_json::Value::Object(values) => serde_json::Value::Object( + values + .into_iter() + .map(|(k, v)| (k, canonical(v))) + .collect::>() + .into_iter() + .collect(), + ), + serde_json::Value::Array(values) => { + serde_json::Value::Array(values.into_iter().map(canonical).collect()) + } + value => value, + } + } + serde_json::to_vec(&canonical( + serde_json::to_value(value).map_err(|e| e.to_string())?, + )) + .map_err(|e| e.to_string()) +} +fn hash(value: &impl Serialize) -> Result { + Ok(format!("{:x}", Sha256::digest(canonical_bytes(value)?))) +} +fn role(role: EdgeRole) -> u8 { + match role { + EdgeRole::Input => 0, + EdgeRole::Left => 1, + EdgeRole::Right => 2, + } +} + +impl SemanticFragment { + pub fn from_stored_output(dag: &ExecutableDag, output: PostAsapNodeId) -> Result { + Self::export(dag, output, true) + } + + pub fn from_dag(dag: &ExecutableDag, output: PostAsapNodeId) -> Result { + Self::export(dag, output, false) + } + + fn export( + dag: &ExecutableDag, + output: PostAsapNodeId, + parameterize_range: bool, + ) -> Result { + let mut included = BTreeSet::new(); + let mut pending = vec![output]; + while let Some(id) = pending.pop() { + if included.insert(id) { + pending.extend( + dag.edges + .iter() + .filter(|e| e.consumer == id) + .map(|e| e.producer), + ); + } + } + let dag = ExecutableDag { + nodes: dag + .nodes + .iter() + .filter(|n| included.contains(&n.id)) + .cloned() + .collect(), + edges: dag + .edges + .iter() + .filter(|e| included.contains(&e.consumer)) + .cloned() + .collect(), + root: output, + }; + dag.validate().map_err(|e| e.to_string())?; + if dag.nodes.len() > 4096 { + return Err("semantic fragment exceeds node budget".into()); + } + // Open PromQL entities carry all labels. Nullable label columns demanded + // only by a downstream consumer do not change a per-entity scalar state. + let mut dag = dag; + let sample_only = matches!( + &dag.nodes + .iter() + .find(|n| n.id == output) + .ok_or("missing output")? + .payload, + ExecutableOperatorPayload::SummaryAgg { + reduction: planner_types::pre_asap::Reduction::PerEntity, + grouping: planner_types::post_asap::GroupingStrategy::PerSubpopulationInstance, + input: planner_types::post_asap::SummaryUpdate { + item: None, + weight: planner_types::post_asap::SummaryInputExpr::Column( + planner_types::pre_asap::ColumnRef::SampleValue + ), + .. + }, + .. + } + ); + if parameterize_range && sample_only { + let direct: BTreeSet<_> = dag + .edges + .iter() + .filter(|e| e.consumer == output) + .map(|e| e.producer) + .collect(); + let mut normalized = false; + for node in &mut dag.nodes { + if !direct.contains(&node.id) { + continue; + } + if let ExecutableOperatorPayload::Fallback { expression } = &mut node.payload { + let source = match expression { + planner_types::pre_asap::QueryExpr::TimeRange { child, .. } => { + std::rc::Rc::make_mut(child) + } + other => other, + }; + if let planner_types::pre_asap::QueryExpr::Scan { + source: planner_types::pre_asap::Source::TimeSeries { .. }, + predicates, + schema, + } = source + { + if !schema.closed + && predicates.is_empty() + && schema.unique_keys.is_empty() + && schema + .columns + .iter() + .take_while(|c| { + !(c.nullable + && c.dtype == planner_types::pre_asap::DataType::Utf8) + }) + .count() + + schema + .columns + .iter() + .rev() + .take_while(|c| { + c.nullable + && c.dtype == planner_types::pre_asap::DataType::Utf8 + }) + .count() + == schema.columns.len() + { + schema.columns.retain(|c| { + !(c.nullable && c.dtype == planner_types::pre_asap::DataType::Utf8) + }); + node.output_schema.fields.retain(|c| { + !(c.nullable + && c.dtype + == planner_types::post_asap::SummaryFamilyType::Plain( + planner_types::pre_asap::DataType::Utf8, + )) + }); + normalized = true; + } + } + } + } + if normalized { + dag.nodes + .iter_mut() + .find(|n| n.id == output) + .unwrap() + .output_schema + .fields + .retain(|c| { + !(c.nullable + && c.dtype + == planner_types::post_asap::SummaryFamilyType::Plain( + planner_types::pre_asap::DataType::Utf8, + )) + }); + } + } + let nodes: BTreeMap<_, _> = dag.nodes.iter().map(|n| (n.id, n)).collect(); + let mut hashes: BTreeMap = BTreeMap::new(); + let mut result = Self { + format_version: 1, + output: String::new(), + nodes: BTreeMap::new(), + }; + let mut stack = vec![(output, false)]; + while let Some((id, finish)) = stack.pop() { + if hashes.contains_key(&id) { + continue; + } + let node = nodes.get(&id).ok_or("missing semantic output")?; + let edges: Vec<_> = dag.edges.iter().filter(|e| e.consumer == id).collect(); + if !finish { + stack.push((id, true)); + for edge in &edges { + stack.push((edge.producer, false)); + } + continue; + } + let mut inputs = edges + .iter() + .map(|e| SemanticInput { + role: e.role, + node: hashes[&e.producer].clone(), + }) + .collect::>(); + inputs.sort_by(|a, b| (role(a.role), &a.node).cmp(&(role(b.role), &b.node))); + let mut payload = node.payload.clone(); + let mut record_range = false; + if parameterize_range + && matches!( + nodes[&output].payload, + ExecutableOperatorPayload::SummaryAgg { .. } + ) + && dag + .edges + .iter() + .any(|e| e.consumer == output && e.producer == id) + { + if let ExecutableOperatorPayload::Fallback { + expression: planner_types::pre_asap::QueryExpr::TimeRange { child, .. }, + } = &payload + { + payload = ExecutableOperatorPayload::Fallback { + expression: child.as_ref().clone(), + }; + record_range = true; + } + } + if let ExecutableOperatorPayload::RelationalJoin { pruning, .. } = &mut payload { + *pruning = None; + } + let operation = SemanticOperation { + record_range, + operation: serde_json::to_value(&payload).map_err(|e| e.to_string())?, + output_schema: serde_json::to_value(&node.output_schema) + .map_err(|e| e.to_string())?, + inputs, + }; + let key = hash(&operation)?; + result.nodes.insert(key.clone(), operation); + hashes.insert(id, key); + } + result.output = hashes + .remove(&output) + .ok_or("missing semantic output hash")?; + result.validate()?; + Ok(result) + } + + pub fn validate(&self) -> Result<(), String> { + if self.format_version != 1 + || self.nodes.is_empty() + || self.nodes.len() > 4096 + || canonical_bytes(self)?.len() > 4 * 1024 * 1024 + { + return Err("unsupported semantic fragment version or size".into()); + } + for (key, node) in &self.nodes { + let payload: ExecutableOperatorPayload = + serde_json::from_value(node.operation.clone()).map_err(|e| e.to_string())?; + if node.record_range { + let root = self + .nodes + .get(&self.output) + .ok_or("missing semantic root")?; + let root_payload: ExecutableOperatorPayload = + serde_json::from_value(root.operation.clone()).map_err(|e| e.to_string())?; + if !matches!(payload, ExecutableOperatorPayload::Fallback { .. }) + || !matches!(root_payload, ExecutableOperatorPayload::SummaryAgg { .. }) + || !root.inputs.iter().any(|input| &input.node == key) + { + return Err("record range must belong to a direct summary input".into()); + } + } + let _: SummarySchema = + serde_json::from_value(node.output_schema.clone()).map_err(|e| e.to_string())?; + if hash(node)? != *key + || node + .inputs + .iter() + .any(|i| !self.nodes.contains_key(&i.node)) + { + return Err("semantic fragment hash or dependency mismatch".into()); + } + if node + .inputs + .windows(2) + .any(|p| (role(p[0].role), &p[0].node) > (role(p[1].role), &p[1].node)) + { + return Err("noncanonical semantic input order".into()); + } + } + let mut seen = BTreeSet::new(); + let mut stack = vec![self.output.as_str()]; + while let Some(id) = stack.pop() { + let node = self.nodes.get(id).ok_or("missing semantic fragment root")?; + if seen.insert(id) { + stack.extend(node.inputs.iter().map(|i| i.node.as_str())); + } + } + if seen.len() != self.nodes.len() { + return Err("unrelated semantic fragment nodes".into()); + } + Ok(()) + } +} + +#[cfg(test)] +pub(crate) mod tests { + use super::*; + use planner_types::post_asap::{compile_executable_dag, SummaryExpr, SummaryNode}; + use planner_types::pre_asap::{Column, DataType, QueryExpr, Schema, Source}; + use std::rc::Rc; + + pub(crate) fn fixture(metric: &str) -> ExecutableDag { + let scan = QueryExpr::Scan { + source: Source::TimeSeries { + metric: metric.into(), + }, + predicates: vec![], + schema: Schema::new(vec![Column::new("value", DataType::Float64, false)]), + }; + let schema = SummarySchema { + fields: vec![planner_types::post_asap::SummaryField { + name: "value".into(), + dtype: planner_types::post_asap::SummaryFamilyType::Plain(DataType::Float64), + nullable: false, + }], + time_index: None, + }; + compile_executable_dag(&Rc::new(SummaryNode { + expr: SummaryExpr::KeepPreAsap(Rc::new(scan)), + schema, + guarantee: None, + })) + .unwrap() + } + + // Storage identity must ignore temporary identifiers and execution placement. + #[test] + fn identity_ignores_node_ids_and_phase() { + let dag = fixture("latency"); + let expected = SemanticFragment::from_dag(&dag, dag.root).unwrap(); + let mut other = dag.clone(); + other.root = PostAsapNodeId(71); + other.nodes[0].id = other.root; + other.nodes[0].output_state.timing = + planner_types::post_asap::ExecutionTiming::IngestionTime; + assert_eq!( + canonical_bytes(&expected).unwrap(), + canonical_bytes(&SemanticFragment::from_dag(&other, other.root).unwrap()).unwrap() + ); + } + + // Source identity and supported semantic format survive restart independently. + #[test] + fn semantics_roundtrip_and_reject_unknown_version() { + let a = fixture("latency"); + let b = fixture("bytes"); + let a = SemanticFragment::from_dag(&a, a.root).unwrap(); + let b = SemanticFragment::from_dag(&b, b.root).unwrap(); + assert_ne!(canonical_bytes(&a).unwrap(), canonical_bytes(&b).unwrap()); + let mut restored: SemanticFragment = + serde_json::from_slice(&canonical_bytes(&a).unwrap()).unwrap(); + restored.validate().unwrap(); + restored.format_version += 1; + assert!(restored.validate().is_err()); + } + // A transformed value cannot share state identity with its source column. + #[test] + fn value_expression_is_semantic_and_nonfinite_constants_are_rejected() { + use planner_types::pre_asap::{ProjectItem, ScalarValue}; + let original = fixture("latency"); + let expected = SemanticFragment::from_dag(&original, original.root).unwrap(); + let mut transformed = original.clone(); + let ExecutableOperatorPayload::Fallback { expression } = &mut transformed.nodes[0].payload + else { + unreachable!() + }; + *expression = QueryExpr::Project { + cols: vec![ProjectItem { + alias: Some("value".into()), + expr: QueryExpr::FunctionCall { + name: "ln".into(), + args: vec![QueryExpr::Column(0)], + }, + }], + qualifier: None, + child: Rc::new(expression.clone()), + }; + let logged = SemanticFragment::from_dag(&transformed, transformed.root).unwrap(); + assert_ne!( + canonical_bytes(&expected).unwrap(), + canonical_bytes(&logged).unwrap() + ); + let ExecutableOperatorPayload::Fallback { + expression: QueryExpr::Project { cols, .. }, + } = &mut transformed.nodes[0].payload + else { + unreachable!() + }; + cols[0].expr = QueryExpr::Literal(ScalarValue::Float64(f64::NAN)); + assert!(SemanticFragment::from_dag(&transformed, transformed.root).is_err()); + } + + // Changing a downstream consumer cannot change the persisted input definition. + #[test] + fn only_output_dependency_closure_is_exported() { + let mut dag = fixture("latency"); + let stored = dag.root; + let mut consumer = dag.nodes[0].clone(); + consumer.id = PostAsapNodeId(9); + consumer.payload = ExecutableOperatorPayload::Value { + operation: planner_types::post_asap::ValueOperation::Project { + cols: vec![], + qualifier: None, + }, + }; + dag.edges.push(planner_types::post_asap::ExecutableDagEdge { + producer: stored, + consumer: consumer.id, + role: EdgeRole::Input, + intermediate_schema: dag.nodes[0].output_schema.clone(), + data_state: dag.nodes[0].output_state, + grouping: planner_types::post_asap::GroupingEdgeCompatibility::NotApplicable, + window: planner_types::post_asap::WindowEdgeCompatibility::NotApplicable, + }); + dag.root = consumer.id; + dag.nodes.push(consumer); + let before = SemanticFragment::from_dag(&dag, stored).unwrap(); + dag.nodes.reverse(); + let after = SemanticFragment::from_dag(&dag, stored).unwrap(); + assert_eq!(before, after); + assert_eq!(before.nodes.len(), 1); + } + // Query lookback does not become the identity of each stored input pane. + #[test] + fn stored_input_range_is_parameterized_but_logical_range_is_preserved() { + use planner_types::post_asap::*; + use planner_types::pre_asap::{ColumnRef, Reduction}; + let make = |seconds| { + let mut dag = fixture("latency"); + let ExecutableOperatorPayload::Fallback { expression } = &mut dag.nodes[0].payload + else { + unreachable!() + }; + *expression = QueryExpr::TimeRange { + range: std::time::Duration::from_secs(seconds), + child: Rc::new(expression.clone()), + }; + let mut output = dag.nodes[0].clone(); + output.id = PostAsapNodeId(1); + let family = SummaryFamilyType::ExactAggregate(ExactKind::Sum, ExactParams::Sum); + output.payload = ExecutableOperatorPayload::SummaryAgg { + family: family.clone(), + input: SummaryUpdate { + item: None, + weight: SummaryInputExpr::Column(ColumnRef::SampleValue), + weight_domain: Default::default(), + }, + reduction: Reduction::PerEntity, + grouping: Default::default(), + }; + output.output_schema.fields[0].dtype = family; + output.output_state.primitive = DataPrimitive::SummaryState; + dag.edges.push(ExecutableDagEdge { + producer: dag.root, + consumer: output.id, + role: EdgeRole::Input, + intermediate_schema: dag.nodes[0].output_schema.clone(), + data_state: dag.nodes[0].output_state, + grouping: GroupingEdgeCompatibility::NotApplicable, + window: WindowEdgeCompatibility::NotApplicable, + }); + dag.root = output.id; + dag.nodes.push(output); + dag + }; + let one = make(60); + let five = make(300); + assert_ne!( + SemanticFragment::from_dag(&one, one.root).unwrap(), + SemanticFragment::from_dag(&five, five.root).unwrap() + ); + assert_eq!( + SemanticFragment::from_stored_output(&one, one.root).unwrap(), + SemanticFragment::from_stored_output(&five, five.root).unwrap() + ); + // Open entities retain their full label identity; consumer-demanded + // optional labels do not change the per-entity stored computation. + let mut open = one.clone(); + let ExecutableOperatorPayload::Fallback { + expression: QueryExpr::TimeRange { child, .. }, + } = &mut open.nodes[0].payload + else { + unreachable!() + }; + let QueryExpr::Scan { schema, .. } = Rc::make_mut(child) else { + unreachable!() + }; + schema.closed = false; + let expected = SemanticFragment::from_stored_output(&open, open.root).unwrap(); + let ExecutableOperatorPayload::Fallback { + expression: QueryExpr::TimeRange { child, .. }, + } = &mut open.nodes[0].payload + else { + unreachable!() + }; + let QueryExpr::Scan { schema, .. } = Rc::make_mut(child) else { + unreachable!() + }; + schema + .columns + .push(Column::new("job", DataType::Utf8, true)); + let label = SummaryField { + name: "job".into(), + dtype: SummaryFamilyType::Plain(DataType::Utf8), + nullable: true, + }; + for node in &mut open.nodes { + node.output_schema.fields.push(label.clone()); + } + open.edges[0].intermediate_schema.fields.push(label); + assert_eq!( + expected, + SemanticFragment::from_stored_output(&open, open.root).unwrap() + ); + let mut forged = SemanticFragment::from_stored_output(&one, one.root).unwrap(); + forged.nodes.values_mut().next().unwrap().operation = + serde_json::json!({"kind": "unknown"}); + assert!(forged.validate().is_err()); + } + // Persisted semantic format changes require an explicit migration/version review. + #[test] + fn semantic_format_v1_has_stable_wire_identity() { + let dag = fixture("latency"); + let exported = SemanticFragment::from_dag(&dag, dag.root).unwrap(); + assert_eq!( + exported.output, + "488a0550f37763997397403ae5ed3588dee5d2a09fa7840f4110b775095fe594" + ); + } +} diff --git a/crates/asap_types/src/summary_semantics.rs b/crates/asap_types/src/summary_semantics.rs index 757a7e15b..642bbba90 100644 --- a/crates/asap_types/src/summary_semantics.rs +++ b/crates/asap_types/src/summary_semantics.rs @@ -17,6 +17,9 @@ pub struct SummaryDefinition { #[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)] pub enum SummarySemantics { Planner { + /// Equal fragments over different datasets are different definitions. + #[serde(default, skip_serializing_if = "Option::is_none")] + dataset: Option, fragment: crate::semantic_fragment::SemanticFragment, }, /// Restricted raw-input adapter for explicit native summary configurations. @@ -76,6 +79,7 @@ impl SummaryDefinition { let fragment = config.semantic_fragment.as_ref(); if let Some(fragment) = fragment { self.semantics = SummarySemantics::Planner { + dataset: config.dataset_identity.clone(), fragment: fragment.clone(), }; } else if let SummarySemantics::Configured { @@ -98,10 +102,15 @@ impl SummaryDefinition { "unsupported semantic format version".into(), )); } - if let SummarySemantics::Planner { fragment } = &self.semantics { + if let SummarySemantics::Planner { dataset, fragment } = &self.semantics { fragment .validate() .map_err(SummaryCatalogError::Descriptor)?; + if let Some(dataset) = dataset { + dataset + .validate() + .map_err(SummaryCatalogError::Descriptor)?; + } } // Object keys are recursively sorted, independent of serde_json features. fn canonical(value: serde_json::Value) -> serde_json::Value { @@ -127,3 +136,39 @@ impl SummaryDefinition { Ok(crate::sds::SummaryDefinitionId::from_semantics(&bytes)) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::semantic_fragment::{LogicalDatasetIdentity, SemanticFragment}; + + fn definition(dataset: Option) -> SummaryDefinition { + let dag = crate::semantic_fragment::tests::fixture("m"); + SummaryDefinition { + semantic_format_version: 1, + semantics: SummarySemantics::Planner { + dataset, + fragment: SemanticFragment::from_dag(&dag, dag.root).unwrap(), + }, + } + } + + fn dataset(name: &str) -> Option { + Some(LogicalDatasetIdentity { + namespace: "ns".into(), + dataset: name.into(), + }) + } + + // Equal fragments over different datasets (or no dataset) get different definition IDs. + #[test] + fn dataset_identity_is_part_of_definition_id() { + let a = definition(dataset("a")).id().unwrap(); + let b = definition(dataset("b")).id().unwrap(); + let none = definition(None).id().unwrap(); + assert_eq!(a, definition(dataset("a")).id().unwrap()); + assert_ne!(a, b); + assert_ne!(none, a); + assert_ne!(none, b); + } +} diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 6f2c2ba2a..3448dbb38 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -1069,6 +1069,7 @@ mod tests { let aggregation = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), @@ -1199,6 +1200,7 @@ mod tests { PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index c48391baa..130c598bd 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -3248,6 +3248,7 @@ mod tests { let cfg = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), diff --git a/data_plane/src/precompute_engine/output_sink.rs b/data_plane/src/precompute_engine/output_sink.rs index eab97b320..fb33de7db 100644 --- a/data_plane/src/precompute_engine/output_sink.rs +++ b/data_plane/src/precompute_engine/output_sink.rs @@ -379,6 +379,7 @@ mod tests { PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), diff --git a/data_plane/src/precompute_engine/revisions.rs b/data_plane/src/precompute_engine/revisions.rs index 02a624622..05bbf8022 100644 --- a/data_plane/src/precompute_engine/revisions.rs +++ b/data_plane/src/precompute_engine/revisions.rs @@ -428,13 +428,11 @@ fn validate_records(output: u64, records: &[RevisionRecord]) -> Result<(), Revis Ok(()) } +use crate::storage_engines::sketch_db::data::native_batch as native; use crate::storage_engines::types::{ AggregateCore, InstalledPrecomputePlanHandle, RuntimePhysicalPlan, }; -use asap_physical_operators::{ - stored_state::native, - values::{Batch, Schema, Value}, -}; +use asap_physical_operators::values::{Batch, Schema, Value}; use planner_types::post_asap::{SummaryFamilyType, SummaryField, SummarySchema}; fn state_schema(family: SummaryFamilyType) -> Schema { diff --git a/data_plane/src/query_engines/asap_query_engine/summary_executor.rs b/data_plane/src/query_engines/asap_query_engine/summary_executor.rs index 0511e1aad..a7d72a44f 100644 --- a/data_plane/src/query_engines/asap_query_engine/summary_executor.rs +++ b/data_plane/src/query_engines/asap_query_engine/summary_executor.rs @@ -757,15 +757,16 @@ fn sketch_query_value( state: &SummaryState, query: &SketchQuery, ) -> Result { - asap_physical_operators::stored_state::readout::sketch_query_value(state, query).map_err( - |asap_physical_operators::stored_state::readout::Error::Unsupported(reason)| { - SummaryExecutorError::Unsupported(reason) - }, - ) + crate::storage_engines::sketch_db::query::sketch_readout::sketch_query_value(state, query) + .map_err( + |crate::storage_engines::sketch_db::query::sketch_readout::Error::Unsupported( + reason, + )| { SummaryExecutorError::Unsupported(reason) }, + ) } fn topk_ranked(state: &SummaryState, k: usize) -> Result, SummaryExecutorError> { - asap_physical_operators::stored_state::readout::topk_ranked(state, k).map_err( - |asap_physical_operators::stored_state::readout::Error::Unsupported(reason)| { + crate::storage_engines::sketch_db::query::sketch_readout::topk_ranked(state, k).map_err( + |crate::storage_engines::sketch_db::query::sketch_readout::Error::Unsupported(reason)| { SummaryExecutorError::Unsupported(reason) }, ) diff --git a/data_plane/src/storage_engines/sketch_db/data/mod.rs b/data_plane/src/storage_engines/sketch_db/data/mod.rs index 475fc6b14..7711f6ef3 100644 --- a/data_plane/src/storage_engines/sketch_db/data/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/data/mod.rs @@ -472,7 +472,24 @@ impl AccuracyBound { /// Per-sample sketch state. Stored as the payload column inside the /// per-sid `SidStoreData` columnar storage. -pub use asap_physical_operators::stored_state::{SketchEncoding, SketchSampleState}; +pub mod native_batch; + +#[derive(Debug, Clone)] +pub struct SketchSampleState { + pub bytes: Vec, + /// Wire-encoding hint from the OTLP DataPoint's `encoding` field. + pub encoding: SketchEncoding, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum SketchEncoding { + ProtoFull, + ProtoDelta, + MsgpackFull, + MsgpackDelta, + /// Versioned typed physical output; never a legacy sketch frame. + NativeBatchV1, +} /// One materialized series row returned by the query path. Resolved /// from the per-sid intern table at read time. diff --git a/data_plane/src/storage_engines/sketch_db/data/native_batch.rs b/data_plane/src/storage_engines/sketch_db/data/native_batch.rs new file mode 100644 index 000000000..308c63daf --- /dev/null +++ b/data_plane/src/storage_engines/sketch_db/data/native_batch.rs @@ -0,0 +1,368 @@ +//! Versioned physical output batches. Deployment identities and coverage remain +//! outside this payload and must be checked before decoding with the bound schema. +use asap_physical_operators::{ + summary_kernels::{ + datasketches_kll::DatasketchesKLLAccumulator, dd_sketch::DDSketchAccumulator, + exact::ExactAccumulator, hll_sketch::HllSketchAccumulator, + weighted_frequency::WeightedFrequency, SumAccumulator, + }, + values::{Batch, Schema, Value}, + AggregateCore, Error, +}; +use planner_types::post_asap::{SummaryFamilyType, SummarySchema}; +use serde::{Deserialize, Serialize}; +use std::sync::Arc; + +#[derive(Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +struct StoredBatch { + version: u32, + schema: SummarySchema, + rows: Vec>, +} + +#[derive(Serialize, Deserialize)] +enum Cell { + Plain(Value), + Summary { + family: SummaryFamilyType, + codec: StateCodec, + bytes: Vec, + }, +} + +// Codec identity is distinct from algorithm identity: old integer CMS/CS bytes +// must never be interpreted as Float64 weighted state with typed item tuples. +#[derive(Serialize, Deserialize)] +enum StateCodec { + WeightedFrequencyV1, + ExactAccumulatorV1, + SumAccumulatorV1, + KllMsgpackV1, + DdMsgpackV1, + HllMsgpackV1, +} +fn invalid(message: impl ToString) -> Error { + Error::Invalid(message.to_string()) +} +impl StateCodec { + fn for_state(state: &dyn AggregateCore) -> Result { + let state = state.as_any(); + if state.is::() { + Ok(Self::WeightedFrequencyV1) + } else if state.is::() { + Ok(Self::ExactAccumulatorV1) + } else if state.is::() { + Ok(Self::SumAccumulatorV1) + } else if state.is::() { + Ok(Self::KllMsgpackV1) + } else if state.is::() { + Ok(Self::DdMsgpackV1) + } else if state.is::() { + Ok(Self::HllMsgpackV1) + } else { + Err(invalid("physical summary has no persisted native codec")) + } + } + fn decode(&self, bytes: &[u8]) -> Result, Error> { + Ok(match self { + Self::WeightedFrequencyV1 => Arc::new(WeightedFrequency::from_bytes(bytes)?), + Self::ExactAccumulatorV1 => { + Arc::new(ExactAccumulator::deserialize_from_bytes(bytes).map_err(invalid)?) + } + Self::SumAccumulatorV1 => { + Arc::new(SumAccumulator::deserialize_from_bytes(bytes).map_err(invalid)?) + } + Self::KllMsgpackV1 => { + Arc::new(DatasketchesKLLAccumulator::from_msgpack_bytes(bytes).map_err(invalid)?) + } + Self::DdMsgpackV1 => { + Arc::new(DDSketchAccumulator::from_msgpack_bytes(bytes).map_err(invalid)?) + } + Self::HllMsgpackV1 => { + Arc::new(HllSketchAccumulator::from_msgpack_bytes(bytes).map_err(invalid)?) + } + }) + } +} + +/// Encode a validated physical output, preserving Float64 and typed identities. +/// This format is independent of the logical and physical plan wire formats. +pub fn encode_batch(batch: &Batch) -> Result, Error> { + let rows = batch + .rows() + .iter() + .map(|row| { + row.iter() + .map(|value| { + Ok(match value { + Value::Summary { family, state } => Cell::Summary { + family: family.clone(), + codec: StateCodec::for_state(state.as_ref())?, + bytes: state.serialize_to_bytes(), + }, + value => Cell::Plain(value.clone()), + }) + }) + .collect::, Error>>() + }) + .collect::, Error>>()?; + rmp_serde::to_vec_named(&StoredBatch { + version: 1, + schema: batch.schema().as_ref().clone(), + rows, + }) + .map_err(invalid) +} + +/// Decode only against the installed output contract. The caller supplies its +/// per-read payload limit; checking state parameters is part of Batch validation. +pub fn decode_batch(bytes: &[u8], expected: Schema, max_bytes: usize) -> Result { + if bytes.len() > max_bytes { + return Err(invalid("native output payload exceeds read budget")); + } + let stored: StoredBatch = rmp_serde::from_slice(bytes).map_err(invalid)?; + if stored.version != 1 { + return Err(invalid("unsupported native output format")); + } + if stored.schema != *expected { + return Err(invalid( + "native output schema differs from installed contract", + )); + } + let rows = stored + .rows + .into_iter() + .map(|row| { + row.into_iter() + .map(|cell| { + Ok(match cell { + Cell::Plain(value) => value, + Cell::Summary { + family, + codec, + bytes, + } => Value::Summary { + family, + state: codec.decode(&bytes)?, + }, + }) + }) + .collect::, Error>>() + }) + .collect::, Error>>()?; + Batch::try_new(expected, rows) +} + +#[cfg(test)] +mod tests { + use super::*; + use asap_physical_operators::summary_kernels::weighted_frequency::FrequencyAlgorithm; + use planner_types::{ + post_asap::{SketchAlgorithm, SketchKind, SketchParams, SummaryField}, + pre_asap::DataType, + }; + + fn weighted(algorithm: SketchAlgorithm) -> Batch { + let (native, params) = match algorithm { + SketchAlgorithm::CmsWithHeap => ( + FrequencyAlgorithm::Cms, + SketchParams::CmsWithHeap { + width: 64, + depth: 5, + heap_size: 8, + }, + ), + _ => ( + FrequencyAlgorithm::CountSketch, + SketchParams::CountSketchWithHeap { + width: 64, + depth: 5, + heap_size: 8, + }, + ), + }; + let family = + SummaryFamilyType::Sketch(SketchKind::new(algorithm, params), Default::default()); + let schema = Arc::new(SummarySchema { + fields: vec![ + SummaryField { + name: "group".into(), + dtype: SummaryFamilyType::Plain(DataType::Utf8), + nullable: false, + }, + SummaryField { + name: "state".into(), + dtype: family.clone(), + nullable: false, + }, + ], + time_index: None, + }); + let mut state = WeightedFrequency::new(native, 64, 5, 8).unwrap(); + state + .update(&[Value::Int64(7), Value::Utf8("service-a".into())], 0.125) + .unwrap(); + state + .update( + &[Value::Utf8("7".into()), Value::Utf8("service-b".into())], + 0.25, + ) + .unwrap(); + Batch::try_new( + schema, + vec![vec![ + Value::Utf8("job-a".into()), + Value::Summary { + family, + state: Arc::new(state), + }, + ]], + ) + .unwrap() + } + + // Fractional rates and distinct typed item tuples survive both heap codecs. + #[test] + fn weighted_outputs_roundtrip_without_integer_conversion() { + for algorithm in [ + SketchAlgorithm::CmsWithHeap, + SketchAlgorithm::CountSketchWithHeap, + ] { + let batch = weighted(algorithm); + let bytes = encode_batch(&batch).unwrap(); + let restored = decode_batch(&bytes, batch.schema().clone(), bytes.len()).unwrap(); + let scores = |b: &Batch| { + let Value::Summary { state, .. } = &b.rows()[0][1] else { + panic!() + }; + state + .as_any() + .downcast_ref::() + .unwrap() + .rows(8) + .iter() + .map(|row| row.iter().map(|v| v.key().unwrap()).collect::>()) + .collect::>() + }; + assert_eq!(scores(&batch), scores(&restored)); + assert!(decode_batch(&bytes, batch.schema().clone(), bytes.len() - 1).is_err()); + } + } + + // A storage tag cannot send weighted physical output through an integer heap decoder. + #[test] + fn legacy_sketch_reader_rejects_native_batch_frames() { + let sample = crate::storage_engines::sketch_db::data::SketchSampleState { + bytes: encode_batch(&weighted(SketchAlgorithm::CmsWithHeap)).unwrap(), + encoding: crate::storage_engines::sketch_db::data::SketchEncoding::NativeBatchV1, + }; + let result = crate::storage_engines::sketch_db::query::delta_apply::per_window_summary_states( + &[(60_000, &sample)], + crate::storage_engines::sketch_db::query::delta_apply::DeltaSketchKind::CmsWithHeap { + rows: 5, + cols: 64, + heap_size: 8, + }, + ); + assert!(result.is_err()); + } + + // Every admitted native summary codec survives the same typed boundary. + #[test] + fn native_summary_families_and_nonfinite_plain_values_roundtrip() { + use planner_types::post_asap::{ExactKind, ExactParams}; + let exact = SummaryFamilyType::ExactAggregate(ExactKind::Sum, ExactParams::Sum); + let sketch = |algorithm, params| { + SummaryFamilyType::Sketch(SketchKind::new(algorithm, params), Default::default()) + }; + let cases: Vec<(SummaryFamilyType, Arc)> = vec![ + ( + exact.clone(), + Arc::new(ExactAccumulator::new(exact.clone(), false).unwrap()), + ), + (exact, Arc::new(SumAccumulator::new())), + ( + sketch(SketchAlgorithm::Kll, SketchParams::Kll { k: 200 }), + Arc::new(DatasketchesKLLAccumulator::new(200)), + ), + ( + sketch( + SketchAlgorithm::DDSketch, + SketchParams::DDSketch { alpha: 0.01 }, + ), + Arc::new(DDSketchAccumulator::new(0.01)), + ), + ( + sketch(SketchAlgorithm::Hll, SketchParams::Hll { precision: 12 }), + Arc::new(HllSketchAccumulator::new( + asap_sketchlib::HllVariant::Regular, + 12, + )), + ), + ]; + for (family, state) in cases { + let schema = Arc::new(SummarySchema { + fields: vec![SummaryField { + name: "state".into(), + dtype: family.clone(), + nullable: false, + }], + time_index: None, + }); + let batch = + Batch::try_new(schema.clone(), vec![vec![Value::Summary { family, state }]]) + .unwrap(); + let bytes = encode_batch(&batch).unwrap(); + let restored = decode_batch(&bytes, schema, bytes.len()).unwrap(); + assert_eq!(encode_batch(&restored).unwrap(), bytes); + } + let schema = Arc::new(SummarySchema { + fields: vec![SummaryField { + name: "value".into(), + dtype: SummaryFamilyType::Plain(DataType::Float64), + nullable: false, + }], + time_index: None, + }); + let batch = Batch::try_new( + schema.clone(), + vec![ + vec![Value::Float64(f64::NAN)], + vec![Value::Float64(f64::INFINITY)], + ], + ) + .unwrap(); + let restored = decode_batch(&encode_batch(&batch).unwrap(), schema, usize::MAX).unwrap(); + assert!(matches!(restored.rows()[0][0], Value::Float64(v) if v.is_nan())); + assert!(matches!(restored.rows()[1][0], Value::Float64(v) if v == f64::INFINITY)); + } + + // Recovery validates the format, bound schema and actual sketch parameters. + #[test] + fn corrupt_or_relabelled_output_is_rejected() { + let batch = weighted(SketchAlgorithm::CmsWithHeap); + let bytes = encode_batch(&batch).unwrap(); + let wrong = weighted(SketchAlgorithm::CountSketchWithHeap); + assert!(decode_batch(&bytes, wrong.schema().clone(), usize::MAX).is_err()); + let mut stored: StoredBatch = rmp_serde::from_slice(&bytes).unwrap(); + stored.version = 2; + assert!(decode_batch( + &rmp_serde::to_vec_named(&stored).unwrap(), + batch.schema().clone(), + usize::MAX + ) + .is_err()); + stored.version = 1; + let Cell::Summary { bytes: payload, .. } = &mut stored.rows[0][1] else { + panic!() + }; + *payload = b"legacy integer heap".to_vec(); + assert!(decode_batch( + &rmp_serde::to_vec_named(&stored).unwrap(), + batch.schema().clone(), + usize::MAX + ) + .is_err()); + } +} diff --git a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs index 4788aaca0..b09fc4005 100644 --- a/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs +++ b/data_plane/src/storage_engines/sketch_db/lifecycle/eviction.rs @@ -204,6 +204,7 @@ mod tests { PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: AggregationType::Sum, aggregation_sub_type: String::new(), diff --git a/data_plane/src/storage_engines/sketch_db/query/decoders.rs b/data_plane/src/storage_engines/sketch_db/query/decoders.rs index d208cb57d..9e2d15d65 100644 --- a/data_plane/src/storage_engines/sketch_db/query/decoders.rs +++ b/data_plane/src/storage_engines/sketch_db/query/decoders.rs @@ -1,2 +1,366 @@ -//! Storage uses the Planner-owned state implementation. -pub use asap_physical_operators::stored_state::decoders::*; +//! SDS sketch frame decoding (modified-OTLP proto and msgpack). +use asap_sketchlib::CountMinSketch; +use asap_sketchlib::CountMinSketchDelta; +use asap_sketchlib::CountMinSketchWithHeap; +use asap_sketchlib::CountSketch; +use asap_sketchlib::CountSketchDelta; +use asap_sketchlib::CountSketchWithHeap; +use asap_sketchlib::CsHeapItem; +use asap_sketchlib::MessagePackCodec; + +use asap_physical_operators::summary_kernels::count_min_sketch_with_heap::CountMinSketchWithHeapAccumulator; + +/// Decode a `CountMinSketch` from the modified-OTLP wire bytes. +/// MSGPACK path round-trips `CountMinSketch::deserialize_msgpack`; +/// PROTO path decodes a `SketchEnvelope{count_min: CountMinState}` +/// (or bare `CountMinState`) and re-projects to a flat matrix. Mirrors +/// `precompute_operators::count_min_sketch::from_sketchlib_proto_bytes`. +pub fn decode_cms_from_proto(buffer: &[u8]) -> Result { + use asap_sketchlib::proto::sketchlib::{ + sketch_envelope, CountMinState, CounterType, SketchEnvelope, + }; + use prost::Message; + + let state = match SketchEnvelope::decode(buffer) { + Ok(env) => match env.sketch_state { + Some(sketch_envelope::SketchState::CountMin(st)) => st, + Some(_) => return Err("SketchEnvelope contains non-CountMin sketch".to_string()), + None => { + CountMinState::decode(buffer).map_err(|e| format!("decode CountMinState: {e}"))? + } + }, + Err(_) => { + CountMinState::decode(buffer).map_err(|e| format!("decode CountMinState: {e}"))? + } + }; + let rows = state.rows as usize; + let cols = state.cols as usize; + if rows == 0 || cols == 0 { + return Err(format!( + "CountMinState has zero dims (rows={rows}, cols={cols})" + )); + } + let expected_len = rows * cols; + let counter_type = CounterType::try_from(state.counter_type) + .map_err(|_| format!("CountMinState unknown counter_type {}", state.counter_type))?; + let flat: Vec = match counter_type { + CounterType::Int32 | CounterType::Int64 => { + if state.counts_int.len() != expected_len { + return Err(format!( + "CountMinState counts_int has {} entries, expected {}", + state.counts_int.len(), + expected_len + )); + } + state.counts_int.iter().map(|&v| v as f64).collect() + } + CounterType::Float64 => { + if state.counts_float.len() != expected_len { + return Err(format!( + "CountMinState counts_float has {} entries, expected {}", + state.counts_float.len(), + expected_len + )); + } + state.counts_float.clone() + } + other => { + return Err(format!( + "CountMinState counter_type {other:?} not yet supported in reducer" + )); + } + }; + let mut matrix = Vec::with_capacity(rows); + for r in 0..rows { + let start = r * cols; + matrix.push(flat[start..start + cols].to_vec()); + } + Ok(CountMinSketch::from_legacy_matrix(matrix, rows, cols)) +} + +/// Decode a `CountMinSketch` from msgpack bytes (sketch-core wire +/// format). Mirrors +/// `CountMinSketchAccumulator::from_msgpack_bytes`. +pub fn decode_cms_from_msgpack(buffer: &[u8]) -> Result { + CountMinSketch::from_msgpack(buffer) + .map_err(|e| format!("deserialize CountMinSketch msgpack: {e}")) +} + +/// Decode a `CountSketch` from the modified-OTLP proto wire bytes. +/// Mirrors +/// `precompute_operators::count_sketch::from_sketchlib_proto_bytes`. +pub fn decode_cs_from_proto(buffer: &[u8]) -> Result { + use asap_sketchlib::proto::sketchlib::{ + sketch_envelope, CountSketchState, CounterType, SketchEnvelope, + }; + use prost::Message; + + let state = match SketchEnvelope::decode(buffer) { + Ok(env) => match env.sketch_state { + Some(sketch_envelope::SketchState::CountSketch(st)) => st, + Some(_) => return Err("SketchEnvelope contains non-CountSketch sketch".to_string()), + None => CountSketchState::decode(buffer) + .map_err(|e| format!("decode CountSketchState: {e}"))?, + }, + Err(_) => { + CountSketchState::decode(buffer).map_err(|e| format!("decode CountSketchState: {e}"))? + } + }; + let rows = state.rows as usize; + let cols = state.cols as usize; + if rows == 0 || cols == 0 { + return Err(format!( + "CountSketchState has zero dims (rows={rows}, cols={cols})" + )); + } + let expected_len = rows * cols; + let counter_type = CounterType::try_from(state.counter_type).map_err(|_| { + format!( + "CountSketchState unknown counter_type {}", + state.counter_type + ) + })?; + let flat: Vec = match counter_type { + CounterType::Int32 | CounterType::Int64 => { + if state.counts_int.len() != expected_len { + return Err(format!( + "CountSketchState counts_int has {} entries, expected {}", + state.counts_int.len(), + expected_len + )); + } + state.counts_int.iter().map(|&v| v as f64).collect() + } + CounterType::Float64 => { + if state.counts_float.len() != expected_len { + return Err(format!( + "CountSketchState counts_float has {} entries, expected {}", + state.counts_float.len(), + expected_len + )); + } + state.counts_float.clone() + } + other => { + return Err(format!( + "CountSketchState counter_type {other:?} not yet supported in reducer" + )); + } + }; + let mut matrix = Vec::with_capacity(rows); + for r in 0..rows { + let start = r * cols; + matrix.push(flat[start..start + cols].to_vec()); + } + Ok(CountSketch::from_legacy_matrix(matrix, rows, cols)) +} + +/// Decode a `CountSketch` from msgpack bytes (sketch-core wire format). +pub fn decode_cs_from_msgpack(buffer: &[u8]) -> Result { + CountSketch::from_msgpack(buffer).map_err(|e| format!("deserialize CountSketch msgpack: {e}")) +} + +/// Decode a `CountMinSketchWithHeap` from msgpack bytes — the OTLP +/// `CountMinSketch` wire bytes when the gateway/precompute layer +/// marked the sid as CmsWithHeap (heap embedded in the +/// `CountMinSketchWithHeapSerialized` outer wrapper). Delegates to +/// `asap_sketchlib::CountMinSketchWithHeap::deserialize_msgpack`. +pub fn decode_cms_with_heap_from_msgpack(buffer: &[u8]) -> Result { + CountMinSketchWithHeap::from_msgpack(buffer) + .map_err(|e| format!("deserialize CountMinSketchWithHeap msgpack: {e}")) +} + +/// Decode a `CountSketchWithHeap` (median-estimator, Count Sketch family) +/// from msgpack bytes. Distinct wire type from `CountMinSketchWithHeap` +/// (min-estimator, Count-Min family) even though both are heap-bearing +/// frequency sketches — see `asap_sketchlib::CountSketchWithHeap`. +/// Delegates to `asap_sketchlib::CountSketchWithHeap::from_msgpack`. +pub fn decode_cs_with_heap_from_msgpack(buffer: &[u8]) -> Result { + CountSketchWithHeap::from_msgpack(buffer) + .map_err(|e| format!("deserialize CountSketchWithHeap msgpack: {e}")) +} + +// --------------------------------------------------------------------------- +// Delta decoders. Under the per-window-reset (PWR) contract +// (`asap-precompute-go/window.go`: a delta is that window's own state +// applied onto a freshly-reset per-series sketch), each stored *Delta +// frame reconstructs into the FULL window state when applied onto an +// EMPTY base of the frame's declared dimensions. The reducer's +// `FrequencyEstimate` / `FrequencyTopk` paths are per-window evaluations, +// so "empty + apply(this window's delta)" yields exactly the window's +// matrix/heap — no cross-window stitching needed (mirrors how the ingest +// accumulators reset_to_empty per window before applying). +// +// The proto path reuses the PUBLIC `asap_sketchlib::{CountSketch, +// CountMinSketch}::apply_delta`; the proto `*Delta` message is decoded via +// `asap_sketchlib::proto::sketchlib::{CountSketchDelta, CountMinDelta}`, +// exactly as `precompute_operators::{count_sketch, +// count_min_sketch}::apply_proto_delta_bytes` does. +// --------------------------------------------------------------------------- + +/// Decode a `CountMinSketch` PROTO_DELTA frame into a FULL sketch by +/// applying the sparse cell delta onto an empty base of the frame's +/// declared dimensions. Mirrors +/// `precompute_operators::count_min_sketch::apply_proto_delta_bytes`. +pub fn decode_cms_from_proto_delta(buffer: &[u8]) -> Result { + use asap_sketchlib::proto::sketchlib::CountMinDelta as PbDelta; + use prost::Message; + + let pb = PbDelta::decode(buffer).map_err(|e| format!("decode CountMinDelta: {e}"))?; + if pb.cell_rows.len() != pb.cell_cols.len() || pb.cell_rows.len() != pb.d_counts.len() { + return Err(format!( + "CountMinDelta packed-array length mismatch: cell_rows={}, cell_cols={}, d_counts={}", + pb.cell_rows.len(), + pb.cell_cols.len(), + pb.d_counts.len() + )); + } + let rows = pb.rows as usize; + let cols = pb.cols as usize; + if rows == 0 || cols == 0 { + return Err(format!( + "CountMinDelta has zero dims (rows={rows}, cols={cols})" + )); + } + let cells = pb + .cell_rows + .iter() + .zip(pb.cell_cols.iter()) + .zip(pb.d_counts.iter()) + .map(|((r, c), dc)| (*r, *c, *dc)) + .collect(); + // hh_keys is parsed off the wire by the precompute accumulator but + // intentionally dropped (the vendored Go proto bindings don't yet + // populate it); match that to keep behavior identical. + let delta = CountMinSketchDelta { + rows: pb.rows, + cols: pb.cols, + cells, + l1: pb.l1, + l2: pb.l2, + hh_keys: Vec::new(), + }; + let mut cms = CountMinSketch::from_legacy_matrix(vec![vec![0.0; cols]; rows], rows, cols); + cms.apply_delta(&delta) + .map_err(|e| format!("apply CountMinDelta onto empty base: {e}"))?; + Ok(cms) +} + +/// Decode a `CountSketch` PROTO_DELTA frame into a FULL sketch by applying +/// the sparse cell delta onto an empty base of the frame's declared +/// dimensions. Mirrors +/// `precompute_operators::count_sketch::apply_proto_delta_bytes`. +pub fn decode_cs_from_proto_delta(buffer: &[u8]) -> Result { + use asap_sketchlib::proto::sketchlib::CountSketchDelta as PbDelta; + use prost::Message; + + let pb = PbDelta::decode(buffer).map_err(|e| format!("decode CountSketchDelta: {e}"))?; + if pb.cell_rows.len() != pb.cell_cols.len() || pb.cell_rows.len() != pb.d_counts.len() { + return Err(format!( + "CountSketchDelta packed-array length mismatch: cell_rows={}, cell_cols={}, d_counts={}", + pb.cell_rows.len(), + pb.cell_cols.len(), + pb.d_counts.len() + )); + } + let rows = pb.rows as usize; + let cols = pb.cols as usize; + if rows == 0 || cols == 0 { + return Err(format!( + "CountSketchDelta has zero dims (rows={rows}, cols={cols})" + )); + } + let cells = pb + .cell_rows + .iter() + .zip(pb.cell_cols.iter()) + .zip(pb.d_counts.iter()) + .map(|((r, c), dc)| (*r, *c, *dc)) + .collect(); + let delta = CountSketchDelta { + rows: pb.rows, + cols: pb.cols, + cells, + l2: pb.l2, + hh_keys: Vec::new(), + }; + let mut cs = CountSketch::from_legacy_matrix(vec![vec![0.0; cols]; rows], rows, cols); + cs.apply_delta(&delta) + .map_err(|e| format!("apply CountSketchDelta onto empty base: {e}"))?; + Ok(cs) +} + +/// Decode a heap-bearing CountSketch MSGPACK_DELTA frame into a FULL +/// `CountMinSketchWithHeap` by applying the sparse matrix delta + full +/// heap onto an empty base of the frame's declared dimensions. This +/// REUSES the ingest-side delta-heap apply logic +/// (`CountMinSketchWithHeapAccumulator::from_msgpack_heap_delta_bytes` → +/// `apply_msgpack_heap_delta_bytes`), which decodes the frame generically +/// with `rmp_serde` — no `asap_sketchlib` delta API is added. +pub fn decode_cms_with_heap_from_msgpack_delta( + buffer: &[u8], +) -> Result { + let acc = CountMinSketchWithHeapAccumulator::from_msgpack_heap_delta_bytes(buffer) + .map_err(|e| format!("reconstruct CountMinSketchWithHeap from delta: {e}"))?; + Ok(acc.inner) +} + +/// Decode a heap-bearing CountSketch (median-estimator) MSGPACK_DELTA frame +/// into a FULL `asap_sketchlib::CountSketchWithHeap` by applying the sparse +/// matrix delta + full heap onto an empty base of the frame's declared +/// dimensions. Same DELTA-HEAP wire shape as the CmsWithHeap delta frame +/// (see `HeapDeltaWire`/`MatrixDeltaWire` in +/// `count_min_sketch_with_heap.rs`), decoded here directly +/// with `rmp_serde` since there is no CountSketchWithHeap ingest +/// accumulator to delegate to. No `asap_sketchlib` delta API needed — the +/// public `from_legacy_matrix` rebuilds both the matrix and heap. +pub fn decode_cs_with_heap_from_msgpack_delta( + buffer: &[u8], +) -> Result { + #[derive(serde::Deserialize)] + struct HeapDeltaWire { + is_delta: bool, + matrix_delta: MatrixDeltaWire, + topk_heap: Vec<(String, f64)>, + heap_size: u64, + } + #[derive(serde::Deserialize)] + struct MatrixDeltaWire { + rows: u32, + cols: u32, + cells: Vec<(u32, u32, i64)>, + } + + let wire: HeapDeltaWire = rmp_serde::from_slice(buffer) + .map_err(|e| format!("decode CountSketchWithHeap delta msgpack: {e}"))?; + if !wire.is_delta { + return Err("CountSketchWithHeap delta frame has is_delta=false".to_string()); + } + let rows = wire.matrix_delta.rows as usize; + let cols = wire.matrix_delta.cols as usize; + if rows == 0 || cols == 0 { + return Err(format!( + "CountSketchWithHeap delta frame has zero dims (rows={rows}, cols={cols})" + )); + } + let mut matrix = vec![vec![0.0; cols]; rows]; + for (r, c, dc) in &wire.matrix_delta.cells { + let (r, c) = (*r as usize, *c as usize); + if r >= rows || c >= cols { + continue; + } + matrix[r][c] += *dc as f64; + } + let heap: Vec = wire + .topk_heap + .into_iter() + .map(|(key, value)| CsHeapItem { key, value }) + .collect(); + Ok(CountSketchWithHeap::from_legacy_matrix( + matrix, + heap, + rows, + cols, + wire.heap_size as usize, + )) +} diff --git a/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs b/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs index 09869df2e..f2845e00a 100644 --- a/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs +++ b/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs @@ -1,10 +1,1287 @@ -//! Storage uses the Planner-owned state implementation. -pub use asap_physical_operators::stored_state::delta_apply::*; +//! SDS sketch state reconstruction from full and delta frames. +use asap_sketchlib::CountMinSketch; +use asap_sketchlib::CountMinSketchWithHeap; +use asap_sketchlib::CountSketch; +use asap_sketchlib::CountSketchWithHeap; +use asap_sketchlib::DdSketch; +use asap_sketchlib::HllSketch; +use asap_sketchlib::HllVariant; +use asap_sketchlib::KllSketch; +use asap_sketchlib::MessagePackCodec; + +use super::decoders::{ + decode_cms_from_msgpack, decode_cms_from_proto, decode_cms_from_proto_delta, + decode_cms_with_heap_from_msgpack, decode_cms_with_heap_from_msgpack_delta, + decode_cs_from_msgpack, decode_cs_from_proto, decode_cs_from_proto_delta, + decode_cs_with_heap_from_msgpack, decode_cs_with_heap_from_msgpack_delta, +}; +use crate::storage_engines::sketch_db::data::{SketchEncoding, SketchSampleState}; + +/// Which sketch family a candidate is, and the parameters needed to +/// *bootstrap an empty state* — required by the per-window-reset (PWR) +/// delta model where a window's FIRST frame is a delta-from-empty (no +/// carry-in Full). Most families' deltas embed their own params in the +/// wire fragment (decoded independently, then merged in — see +/// `SummaryState::apply_delta_bytes`); HLL register deltas and DD's +/// bucket-index deltas are applied onto a pre-sized structure instead, +/// so those two need the params known up front to allocate it. +#[derive(Debug, Clone, Copy)] +pub enum DeltaSketchKind { + UnivMon { + heap_size: u32, + sketch_rows: u32, + sketch_cols: u32, + layers: u8, + }, + DDSketch { + alpha: f64, + }, + Hll { + precision: u32, + }, + Kll { + k: u32, + }, + Cms { + rows: usize, + cols: usize, + }, + CountSketch { + rows: usize, + cols: usize, + }, + /// `CmsWithHeap` wraps `asap_sketchlib::CountMinSketchWithHeap` + /// (min-over-rows estimator) and `CountSketchWithHeap` wraps the + /// distinct `asap_sketchlib::CountSketchWithHeap` (median-of-signed-rows + /// estimator) -- different algorithms that happen to share a storage + /// shape. Kept as two variants (not one shared `Heap`) so + /// `merge_same_family` rejects merging one into the other the same + /// way it already rejects e.g. merging a `Cms` into a `Kll`; now the + /// type system enforces it too, since the two variants hold different + /// Rust types. + CmsWithHeap { + rows: usize, + cols: usize, + heap_size: usize, + }, + CountSketchWithHeap { + rows: usize, + cols: usize, + heap_size: usize, + }, +} + +impl DeltaSketchKind { + /// Construct an EMPTY state for this kind, used to seed a new window + /// when its first frame is a delta-from-empty (PWR). A delta applied + /// onto this empty base reconstructs exactly that window's state + /// (delta-from-empty ⊕ empty = window state). + fn bootstrap_empty(&self) -> SummaryState { + match self { + Self::UnivMon { + heap_size, + sketch_rows, + sketch_cols, + layers, + } => SummaryState::UnivMon( + asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator::new( + *heap_size as usize, + *sketch_rows as usize, + *sketch_cols as usize, + *layers as usize, + ) + .expect("validated UnivMon catalog dimensions"), + ), + DeltaSketchKind::DDSketch { alpha } => SummaryState::Dd(DdSketch::new(*alpha)), + DeltaSketchKind::Kll { k } => SummaryState::Kll(KllSketch::new(*k as u16)), + DeltaSketchKind::Hll { precision } => { + SummaryState::Hll(HllSketch::new(HllVariant::Regular, *precision)) + } + DeltaSketchKind::Cms { rows, cols } => { + SummaryState::Cms(CountMinSketch::new(*rows, *cols)) + } + DeltaSketchKind::CountSketch { rows, cols } => { + SummaryState::CountSketch(CountSketch::new(*rows, *cols)) + } + DeltaSketchKind::CmsWithHeap { + rows, + cols, + heap_size, + } => SummaryState::CmsWithHeap(CountMinSketchWithHeap::new(*rows, *cols, *heap_size)), + DeltaSketchKind::CountSketchWithHeap { + rows, + cols, + heap_size, + } => SummaryState::CountSketchWithHeap(CountSketchWithHeap::new( + *rows, *cols, *heap_size, + )), + } + } +} + +/// Try to decode a "full" sketch from the bytes (used by both +/// per-window and cumulative modes when the encoding is `*Full`). +fn decode_full( + kind: &DeltaSketchKind, + bytes: &[u8], + encoding: SketchEncoding, +) -> Result { + match (kind, encoding) { + ( + DeltaSketchKind::UnivMon { + heap_size, + sketch_rows, + sketch_cols, + layers, + }, + SketchEncoding::MsgpackFull, + ) => { + let state = + asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator::from_bytes( + bytes, + ) + .map_err(|e| e.to_string())?; + if state.dimensions() + != ( + *heap_size as usize, + *sketch_rows as usize, + *sketch_cols as usize, + *layers as usize, + ) + { + return Err("UnivMon payload dimensions differ from installed catalog".into()); + } + Ok(SummaryState::UnivMon(state)) + } + (DeltaSketchKind::DDSketch { .. }, SketchEncoding::ProtoFull) => { + let sk = dd_from_proto(bytes)?; + Ok(SummaryState::Dd(sk)) + } + (DeltaSketchKind::DDSketch { .. }, SketchEncoding::MsgpackFull) => { + let sk = DdSketch::from_msgpack(bytes) + .map_err(|e| format!("deserialize DDSketch msgpack: {e}"))?; + Ok(SummaryState::Dd(sk)) + } + (DeltaSketchKind::Hll { .. }, SketchEncoding::ProtoFull) => { + let sk = hll_from_proto(bytes)?; + Ok(SummaryState::Hll(sk)) + } + (DeltaSketchKind::Hll { .. }, SketchEncoding::MsgpackFull) => { + let sk = HllSketch::from_msgpack(bytes) + .map_err(|e| format!("deserialize HllSketch msgpack: {e}"))?; + Ok(SummaryState::Hll(sk)) + } + (DeltaSketchKind::Kll { .. }, SketchEncoding::ProtoFull) => { + let sk = kll_from_proto(bytes)?; + Ok(SummaryState::Kll(sk)) + } + (DeltaSketchKind::Kll { .. }, SketchEncoding::MsgpackFull) => { + let sk = KllSketch::from_msgpack(bytes) + .map_err(|e| format!("deserialize KllSketch msgpack: {e}"))?; + Ok(SummaryState::Kll(sk)) + } + (DeltaSketchKind::Cms { .. }, SketchEncoding::ProtoFull) => { + Ok(SummaryState::Cms(decode_cms_from_proto(bytes)?)) + } + (DeltaSketchKind::Cms { .. }, SketchEncoding::MsgpackFull) => { + Ok(SummaryState::Cms(decode_cms_from_msgpack(bytes)?)) + } + (DeltaSketchKind::CountSketch { .. }, SketchEncoding::ProtoFull) => { + Ok(SummaryState::CountSketch(decode_cs_from_proto(bytes)?)) + } + (DeltaSketchKind::CountSketch { .. }, SketchEncoding::MsgpackFull) => { + Ok(SummaryState::CountSketch(decode_cs_from_msgpack(bytes)?)) + } + // The heap-bearing wire format is msgpack-only in this + // deployment; `decode_cms_with_heap_from_msgpack` is the same + // "Full" decoder the reducer's existing per-frame dispatch falls + // through to for any non-MsgpackDelta encoding. + ( + DeltaSketchKind::CmsWithHeap { .. }, + SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull, + ) => Ok(SummaryState::CmsWithHeap( + decode_cms_with_heap_from_msgpack(bytes)?, + )), + ( + DeltaSketchKind::CountSketchWithHeap { .. }, + SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull, + ) => Ok(SummaryState::CountSketchWithHeap( + decode_cs_with_heap_from_msgpack(bytes)?, + )), + (_, e) => Err(format!("decode_full called with non-Full encoding {e:?}")), + } +} + +/// The reconstructed state one candidate sid contributes — either +/// folded across a window (or several) via delta application, or merged +/// in from another sid's own reconstruction. +pub enum SummaryState { + UnivMon(asap_physical_operators::summary_kernels::univmon::UnivMonAccumulator), + Dd(DdSketch), + Hll(HllSketch), + Kll(KllSketch), + Cms(CountMinSketch), + CountSketch(CountSketch), + /// See `DeltaSketchKind::CmsWithHeap`/`CountSketchWithHeap` for why + /// these are two variants holding two different sketchlib types. + CmsWithHeap(CountMinSketchWithHeap), + CountSketchWithHeap(CountSketchWithHeap), +} + +impl SummaryState { + /// Apply a delta-encoded payload from a window sample. For DD / KLL, + /// the delta is interpreted as a "mergeable fragment" decoded + /// through the same full-state decoder and merged into the + /// rolling state. For HLL, the wire delta is a sparse register + /// update applied via the sketch's `apply_delta`. + /// + /// On encoding mismatch (e.g. trying to apply an HllDelta to a + /// DDSketch rolling state) returns Err. + pub fn apply_delta_bytes( + &mut self, + bytes: &[u8], + encoding: SketchEncoding, + ) -> Result<(), String> { + if !matches!( + encoding, + SketchEncoding::ProtoDelta | SketchEncoding::MsgpackDelta + ) { + return Err(format!( + "apply_delta_bytes called with non-Delta encoding {encoding:?}" + )); + } + match self { + SummaryState::UnivMon(_) => Err("UnivMon requires full pane snapshots".into()), + SummaryState::Dd(sk) => { + match encoding { + // PROTO_DELTA: dispatch on the payload SHAPE, mirroring the + // supported DDSketch frame decoder, which tries the + // full-envelope decode first, then falls back to + // the bucket-delta proto. Two wire shapes can arrive on the + // ProtoDelta channel: + // + // 1. `SketchEnvelope{DdSketchState}` — a full-state + // fragment, mergeable via `DdSketch::merge`. (The edge + // sends this when `compute_delta_against` hits the + // empty-current / undecodable-prior fallback and ships + // a full snapshot tagged as a delta.) + // 2. `DDSketchDelta { buckets: [{index, d_count}] }` — a + // bucket-index delta proto, applied additively. This is + // the COMMON delta_transmission frame the edge emits + // under per-window-reset (`compute_delta(&empty)`). + // + // Before this fix the reducer decoded ONLY shape (1) via + // `decode_full`. A real shape-(2) frame failed with a wire- + // type mismatch on field 1 (delta field 1 = repeated + // submessage; state field 1 = `double alpha`) → the whole + // `quantile_over_time` returned `No result` for every + // delta_transmission DDSketch stream. We wrap the rolling + // `DdSketch` in a transient accumulator so the bucket-delta + // apply lands on `sk` in place. + SketchEncoding::ProtoDelta => { + // Shape (1): full envelope fragment → merge. Try this + // first (cheap decode attempt; a bucket-delta proto + // fails it on the field-1 wire-type mismatch). + if let Ok(SummaryState::Dd(other)) = decode_full( + &DeltaSketchKind::DDSketch { alpha: 0.0 }, + bytes, + SketchEncoding::ProtoFull, + ) { + sk.merge(&other) + .map_err(|e| format!("merge DDSketch delta envelope: {e}"))?; + return Ok(()); + } + // Shape (2): bucket-delta proto → additive apply via the + // SAME decoder the ingest delta path uses. + use asap_physical_operators::summary_kernels::dd_sketch::DDSketchAccumulator; + let mut acc = DDSketchAccumulator { + inner: std::mem::replace(sk, DdSketch::new(sk.alpha)), + sample_p: 1.0, + }; + let res = acc.apply_proto_delta_bytes(bytes); + *sk = acc.inner; + res.map_err(|e| format!("apply DDSketch proto bucket-delta: {e}"))?; + Ok(()) + } + // MSGPACK_DELTA: a serialized full-sketch fragment, mergeable + // via the full-state decoder. Kept for completeness — the + // edge wires PROTO_DELTA for DDSketch today. + SketchEncoding::MsgpackDelta => { + let other = match decode_full( + &DeltaSketchKind::DDSketch { alpha: 0.0 }, + bytes, + SketchEncoding::MsgpackFull, + ) { + Ok(SummaryState::Dd(s)) => s, + Ok(_) => { + return Err( + "decode_full(DDSketch) returned non-DDSketch state".to_string() + ) + } + Err(e) => return Err(e), + }; + sk.merge(&other) + .map_err(|e| format!("merge DDSketch delta: {e}"))?; + Ok(()) + } + _ => unreachable!(), + } + } + SummaryState::Hll(sk) => { + // HLL has a true sparse register delta in the proto + // wire format. Use the same path the precompute + // accumulator uses (`apply_proto_delta_bytes`-style). + if encoding == SketchEncoding::ProtoDelta { + apply_hll_proto_delta(sk, bytes) + } else { + // MsgpackDelta for HLL isn't a sparse encoding; + // it's a serialized HllSketch fragment, mergeable + // via `HllSketch::merge`. + let other = HllSketch::from_msgpack(bytes) + .map_err(|e| format!("deserialize HllSketch (delta-as-msgpack): {e}"))?; + sk.merge(&other) + .map_err(|e| format!("merge HLL delta: {e}"))?; + Ok(()) + } + } + SummaryState::Kll(sk) => { + let full_enc = match encoding { + SketchEncoding::ProtoDelta => SketchEncoding::ProtoFull, + SketchEncoding::MsgpackDelta => SketchEncoding::MsgpackFull, + _ => unreachable!(), + }; + let other = match decode_full(&DeltaSketchKind::Kll { k: 0 }, bytes, full_enc) { + Ok(SummaryState::Kll(s)) => s, + Ok(_) => return Err("decode_full(Kll) returned non-Kll state".to_string()), + Err(e) => return Err(e), + }; + sk.merge(&other) + .map_err(|e| format!("merge KLL delta: {e}"))?; + Ok(()) + } + // CMS/CountSketch/Heap have no true sparse in-place delta + // (unlike DD's bucket-index proto or HLL's register proto, + // above) — every delta frame already decodes into a + // complete, standalone state on its own (the PWR wire + // contract resets to empty at the source), so applying one + // is always "decode independently, then merge". + SummaryState::Cms(sk) => { + if encoding != SketchEncoding::ProtoDelta { + return Err( + "CountMin (heap-less) MSGPACK_DELTA is not a valid producer encoding \ + (msgpack-delta is the heap-bearing form)" + .to_string(), + ); + } + let other = decode_cms_from_proto_delta(bytes)?; + sk.merge(&other) + .map_err(|e| format!("merge CountMinSketch delta: {e}")) + } + SummaryState::CountSketch(sk) => { + if encoding != SketchEncoding::ProtoDelta { + return Err( + "CountSketch (heap-less) MSGPACK_DELTA is not a valid producer encoding \ + (msgpack-delta is the heap-bearing form)" + .to_string(), + ); + } + let other = decode_cs_from_proto_delta(bytes)?; + sk.merge(&other) + .map_err(|e| format!("merge CountSketch delta: {e}")) + } + SummaryState::CmsWithHeap(sk) => { + // Matches the existing per-frame reducer dispatch: only + // MsgpackDelta gets true delta treatment; ProtoDelta (not + // produced for this family in this deployment) falls + // through to the full-msgpack decoder, same as `decode_full`. + let other = if encoding == SketchEncoding::MsgpackDelta { + decode_cms_with_heap_from_msgpack_delta(bytes)? + } else { + decode_cms_with_heap_from_msgpack(bytes)? + }; + sk.merge(&other) + .map_err(|e| format!("merge CmsWithHeap delta: {e}")) + } + SummaryState::CountSketchWithHeap(sk) => { + let other = if encoding == SketchEncoding::MsgpackDelta { + decode_cs_with_heap_from_msgpack_delta(bytes)? + } else { + decode_cs_with_heap_from_msgpack(bytes)? + }; + sk.merge(&other) + .map_err(|e| format!("merge CountSketchWithHeap delta: {e}")) + } + } + } + + pub fn quantile(&self, q: f64) -> f64 { + match self { + SummaryState::Dd(sk) => sk.quantile(q).unwrap_or(0.0), + SummaryState::Kll(sk) => sk.quantile(q), + _ => 0.0, + } + } + + pub fn cardinality(&self) -> f64 { + match self { + SummaryState::Hll(sk) => sk.estimate(), + _ => 0.0, + } + } + + /// The bucket TOTAL — sum of row 0 of the underlying matrix. What a + /// bare `count_over_time`/`sum by (item) (rate(...))`-shaped query + /// (no specific item key) reads out. `0.0` for non-Frequency-family + /// states. + pub fn total(&self) -> f64 { + let matrix = match self { + SummaryState::Cms(c) => c.sketch(), + SummaryState::CountSketch(c) => c.sketch().clone(), + SummaryState::CmsWithHeap(h) => h.sketch_matrix(), + SummaryState::CountSketchWithHeap(h) => h.sketch_matrix(), + _ => return 0.0, + }; + matrix + .first() + .map(|row| row.iter().copied().sum::()) + .unwrap_or(0.0) + } + + /// Per-key point estimate — `count(metric{item="x"})`-shaped queries. + /// Unlike [`Self::topk_items`], no heap is needed: all four Frequency + /// variants (heap-bearing or not) already carry a keyed `estimate` + /// over their matrix. `None` for the quantile/cardinality states, + /// which have no item universe at all. + pub fn estimate(&self, key: &str) -> Option { + match self { + SummaryState::Cms(c) => Some(c.estimate(key)), + SummaryState::CountSketch(c) => Some(c.estimate(key)), + SummaryState::CmsWithHeap(h) => Some(h.estimate(key)), + SummaryState::CountSketchWithHeap(h) => Some(h.estimate(key)), + _ => None, + } + } + + /// Top-k `(key, value)` pairs from the heap, descending by value. + /// `None` for anything other than a heap-bearing state — the + /// heap-less Frequency states (`Cms`/`CountSketch`) carry no item + /// universe to enumerate, and the quantile/cardinality states have + /// no heap at all. + pub fn topk_items(&self) -> Option> { + match self { + SummaryState::CmsWithHeap(h) => Some( + h.topk_heap_items() + .into_iter() + .map(|item| (item.key, item.value)) + .collect(), + ), + SummaryState::CountSketchWithHeap(h) => Some( + h.topk_heap_items() + .into_iter() + .map(|item| (item.key, item.value)) + .collect(), + ), + _ => None, + } + } + + /// Merge `other` into `self` in place — both must be the same sketch + /// family. Used to combine several sids' reconstructed states + /// (`cumulative_summary_state`/`per_window_summary_states`) into one + /// cross-sid answer. `CmsWithHeap`/`CountSketchWithHeap` fall through + /// to the catch-all mismatch arm below like any other mixed pair — + /// and since the two variants now hold distinct sketchlib types + /// (`CountMinSketchWithHeap` vs `CountSketchWithHeap`), there is no + /// arm that could accidentally match them together — see their doc + /// on `DeltaSketchKind`. + pub fn merge_same_family(&mut self, other: &SummaryState) -> Result<(), String> { + match (self, other) { + (SummaryState::UnivMon(a), SummaryState::UnivMon(b)) => { + a.merge_in_place(b).map_err(|e| e.to_string()) + } + (SummaryState::Dd(a), SummaryState::Dd(b)) => { + a.merge(b).map_err(|e| format!("merge DDSketch: {e}")) + } + (SummaryState::Hll(a), SummaryState::Hll(b)) => { + a.merge(b).map_err(|e| format!("merge HLL: {e}")) + } + (SummaryState::Kll(a), SummaryState::Kll(b)) => { + a.merge(b).map_err(|e| format!("merge KLL: {e}")) + } + (SummaryState::Cms(a), SummaryState::Cms(b)) => { + a.merge(b).map_err(|e| format!("merge CountMinSketch: {e}")) + } + (SummaryState::CountSketch(a), SummaryState::CountSketch(b)) => { + a.merge(b).map_err(|e| format!("merge CountSketch: {e}")) + } + (SummaryState::CmsWithHeap(a), SummaryState::CmsWithHeap(b)) => { + a.merge(b).map_err(|e| format!("merge CmsWithHeap: {e}")) + } + (SummaryState::CountSketchWithHeap(a), SummaryState::CountSketchWithHeap(b)) => a + .merge(b) + .map_err(|e| format!("merge CountSketchWithHeap: {e}")), + (a, _) => Err(format!( + "SummaryState family mismatch in merge_same_family (self is {})", + a.family_name() + )), + } + } + + /// Diagnostic family name for error messages — not used for dispatch. + fn family_name(&self) -> &'static str { + match self { + SummaryState::UnivMon(_) => "UnivMon", + SummaryState::Dd(_) => "DDSketch", + SummaryState::Hll(_) => "Hll", + SummaryState::Kll(_) => "Kll", + SummaryState::Cms(_) => "Cms", + SummaryState::CountSketch(_) => "CountSketch", + SummaryState::CmsWithHeap(_) => "CmsWithHeap", + SummaryState::CountSketchWithHeap(_) => "CountSketchWithHeap", + } + } +} + +/// Fold every in-range window's frames for ONE series into a single +/// merged `SummaryState` (cumulative over `[t0, t1]`), returning `None` +/// if no Full frame ever landed (every sample was a leading delta). The +/// per-sid building block for a cross-sid answer: reconstruct each +/// candidate sid's state this way, then merge them (`merge_same_family`) +/// before reading out a quantile/cardinality over the combined data. +pub fn cumulative_summary_state( + samples: &[(i64, &SketchSampleState)], + kind: DeltaSketchKind, +) -> Result, String> { + let mut rolling: Option = None; + visit_window_summary_states(samples, kind, |_, state| { + if let Some(acc) = rolling.as_mut() { + acc.merge_same_family(&state)?; + } else { + rolling = Some(state); + } + Ok(()) + })?; + Ok(rolling) +} + +#[cfg(test)] +/// Walk a sorted-by-window-end slice of samples in time order and +/// produce ONE per-window scalar `(window_end_ms, scalar)`. +/// +/// ## Per-window-reset (PWR) delta model +/// +/// The edge emits frames grouped by window (all frames of one window +/// share the same `window_end` key; the key changes across windows). +/// The edge RESETS its snapshot base at each window boundary, so each +/// window's state is built *from empty*: +/// +/// * Within a window, frames accumulate to the window total. The first +/// frame may be a `Full` (window 1, or a periodic re-snapshot) or a +/// `Delta`-from-empty (windows 2+ under PWR); subsequent frames are +/// `Delta` INCREMENTS applied onto the window's running base. +/// * Across windows, the base MUST reset — a new `window_end` discards +/// the previous window's rolling state and starts from empty. Never +/// carry one window's state into the next (that would inflate via +/// cross-window accumulation). +/// +/// Concretely this fixes two bugs in the old "single rolling Option that +/// only ever resets on a Full" walk: +/// 1. A query range whose Full lives only in window 1 (or out of +/// range) left windows 2+ as deltas with `rolling=None`, all +/// skipped → empty result. +/// 2. A window 2+ delta applied onto window 1's leftover rolling state +/// → cross-window inflation. +/// +/// For a `Delta` that is the window's FIRST frame (the PWR delta-from- +/// empty case), we bootstrap an EMPTY rolling state of `kind` and apply +/// the delta onto it (delta-from-empty ⊕ empty = that window's state). +/// +/// The delta-OFF path (exactly one `Full` per window) still produces one +/// correct value per window: the window opens with a Full, has no +/// further frames, and emits that Full's scalar. +/// +/// `eval` reads a scalar from the rolling state (`quantile(q)` / +/// `cardinality()`). `skipped` counts frames that could not contribute +/// (a delta we genuinely couldn't bootstrap from — should be rare). +/// +/// Returns `Ok(per_window_samples, skipped)`. +pub fn per_window_evaluate( + samples: &[(i64, &SketchSampleState)], + kind: DeltaSketchKind, + eval: E, +) -> Result<(Vec<(i64, f64)>, usize), String> +where + E: Fn(&SummaryState) -> f64, +{ + let (states, skipped) = per_window_summary_states(samples, kind)?; + Ok(( + states.into_iter().map(|(w, rs)| (w, eval(&rs))).collect(), + skipped, + )) +} + +/// Walk a sorted-by-window-end slice of samples in time order and +/// reconstruct ONE sid's per-window `SummaryState` (same per-window-reset +/// walk as [`per_window_evaluate`], generalized to return the +/// reconstructed state itself instead of an already-evaluated scalar). +/// The per-sid building block for cross-sid per-window merging (unlike +/// [`cumulative_summary_state`], which folds a whole `[t0, t1]` range +/// into one answer, this keeps each window separate so a caller can +/// merge same-window states across several sids before evaluating -- +/// needed for a matrix/range-query answer, where each output point is +/// itself a cross-sid merge for that one window). +/// +/// Returns `Ok((per_window_states, skipped))`. +pub fn per_window_summary_states( + samples: &[(i64, &SketchSampleState)], + kind: DeltaSketchKind, +) -> Result<(Vec<(i64, SummaryState)>, usize), String> { + let mut out: Vec<(i64, SummaryState)> = Vec::new(); + let skipped = visit_window_summary_states(samples, kind, |end, state| { + out.push((end, state)); + Ok(()) + })?; + Ok((out, skipped)) +} + +// Both readout modes must reconstruct the same final pane population. The +// visitor lets cumulative merging stream panes without retaining every state. +fn visit_window_summary_states( + samples: &[(i64, &SketchSampleState)], + kind: DeltaSketchKind, + mut emit: impl FnMut(i64, SummaryState) -> Result<(), String>, +) -> Result { + let mut skipped = 0usize; + + // Rolling state for the CURRENT window only. Reset to None whenever + // `window_end` changes (a new window establishes its own base from + // empty). `cur_end` tracks which window `rolling` belongs to. + let mut rolling: Option = None; + let mut cur_end: Option = None; + + for (window_end, state) in samples { + // Window boundary: flush the previous window's final accumulated + // state, then reset the base so this window starts from empty. + if cur_end != Some(*window_end) { + if let (Some(prev_end), Some(rs)) = (cur_end, rolling.take()) { + emit(prev_end, rs)?; + } + cur_end = Some(*window_end); + } + + match state.encoding { + SketchEncoding::NativeBatchV1 => { + return Err("native physical outputs require the bound native batch decoder".into()) + } + SketchEncoding::ProtoFull | SketchEncoding::MsgpackFull => { + // A Full (re)sets this window's base. + rolling = Some(decode_full(&kind, &state.bytes, state.encoding)?); + } + SketchEncoding::ProtoDelta | SketchEncoding::MsgpackDelta => { + // Apply onto this window's running base. If this is the + // window's first frame (PWR delta-from-empty), bootstrap + // an empty base and apply onto it. + if rolling.is_none() { + rolling = Some(kind.bootstrap_empty()); + } + match rolling.as_mut() { + Some(rs) => rs.apply_delta_bytes(&state.bytes, state.encoding)?, + None => skipped += 1, + } + } + } + } + + // Flush the final window. + if let (Some(prev_end), Some(rs)) = (cur_end, rolling.take()) { + emit(prev_end, rs)?; + } + + Ok(skipped) +} + +// --------------------------------------------------------------------------- +// Proto-envelope decoders — P2-4: ONE decoder per family. +// +// These delegate to the precompute-side accumulators' +// `from_sketchlib_proto_bytes`, which are the single source of truth for +// the modified-OTLP proto wire format (envelope unwrapping, alpha/k/ +// precision validation, and — critically for HLL — SPARSE +// `registers_sparse` expansion). Folding the warm read path onto the +// same decoder the ingest path uses means the sparse-register fix (and +// any future format change) can never drift between the two copies again +// — the bug class P2-3 / P2-4 closed. We extract the accumulator's +// public `inner` sketch for the rolling-state merge. +// --------------------------------------------------------------------------- + +fn dd_from_proto(buffer: &[u8]) -> Result { + use asap_physical_operators::summary_kernels::dd_sketch::DDSketchAccumulator; + DDSketchAccumulator::from_sketchlib_proto_bytes(buffer) + .map(|acc| acc.inner) + .map_err(|e| e.to_string()) +} + +fn kll_from_proto(buffer: &[u8]) -> Result { + use asap_physical_operators::summary_kernels::datasketches_kll::DatasketchesKLLAccumulator; + DatasketchesKLLAccumulator::from_sketchlib_proto_bytes(buffer) + .map(|acc| acc.inner) + .map_err(|e| e.to_string()) +} + +fn hll_from_proto(buffer: &[u8]) -> Result { + use asap_physical_operators::summary_kernels::hll_sketch::HllSketchAccumulator; + HllSketchAccumulator::from_sketchlib_proto_bytes(buffer) + .map(|acc| acc.inner) + .map_err(|e| e.to_string()) +} + +/// Apply a proto-encoded `HllDelta` frame onto the HLL register vector — the +/// delta is a varint-packed (index_delta, value) blob; decode + apply +/// (register-wise max) via the shared sketch library so the unpacking stays a +/// single source of truth. +fn apply_hll_proto_delta(sk: &mut HllSketch, buffer: &[u8]) -> Result<(), String> { + sk.apply_delta_bytes(buffer) + .map_err(|e| format!("apply HLLDelta: {e}"))?; + Ok(()) +} #[cfg(test)] mod tests { + //! P2-3 / P2-4 regression tests for the consolidated single-decoder + //! path. These exercise the family proto decoders that now delegate + //! to the precompute accumulators (the single source of truth), so a + //! divergence between the warm read path and the ingest path — + //! notably the SPARSE-register HLL handling the deleted dead decoder + //! got wrong — fails the build. + use super::*; + use asap_sketchlib::HllVariant; + + fn encode_dd(sk: &DdSketch) -> Vec { + use asap_sketchlib::proto::sketchlib::{sketch_envelope, DdSketchState, SketchEnvelope}; + use prost::Message; + let state = DdSketchState { + alpha: sk.alpha, + store_counts: sk.store_counts.clone(), + store_offset: sk.store_offset, + ..Default::default() + }; + SketchEnvelope { + sketch_state: Some(sketch_envelope::SketchState::Ddsketch(state)), + ..Default::default() + } + .encode_to_vec() + } + + fn encode_kll(k: u16, items: &[f64]) -> Vec { + use asap_sketchlib::proto::sketchlib::{sketch_envelope, KllState, SketchEnvelope}; + use prost::Message; + let state = KllState { + k: k as u32, + items: items.to_vec(), + levels: vec![], + num_levels: 0, + ..Default::default() + }; + SketchEnvelope { + sketch_state: Some(sketch_envelope::SketchState::Kll(state)), + ..Default::default() + } + .encode_to_vec() + } + + fn encode_hll_dense(sk: &HllSketch) -> Vec { + use asap_sketchlib::proto::sketchlib::{ + sketch_envelope, HllVariant as ProtoVariant, HyperLogLogState, SketchEnvelope, + }; + use prost::Message; + let state = HyperLogLogState { + variant: ProtoVariant::Regular as i32, + precision: sk.precision, + registers: sk.registers.clone(), + hip_kxq0: sk.hip_kxq0, + hip_kxq1: sk.hip_kxq1, + hip_est: sk.hip_est, + registers_sparse: None, + }; + SketchEnvelope { + sketch_state: Some(sketch_envelope::SketchState::Hll(state)), + ..Default::default() + } + .encode_to_vec() + } + + /// Build a SPARSE HLL proto frame: dense `registers` left empty, + /// `registers_sparse.packed` = varint (index_delta, value) pairs. + /// This is exactly the wire form a low-cardinality producer emits + /// (sketchlib-go below its dense/sparse crossover) — the frame the + /// DELETED `HllSketch_from_sketchlib_proto_bytes` hard-rejected with + /// "registers has 0 bytes". + fn encode_hll_sparse(precision: u32, nonzero: &[(u64, u8)]) -> Vec { + use asap_sketchlib::proto::sketchlib::{ + sketch_envelope, HllSparseRegisters, HllVariant as ProtoVariant, HyperLogLogState, + SketchEnvelope, + }; + use prost::Message; + // Varint-pack (index_delta, value), ascending index order. + let mut packed: Vec = Vec::new(); + let mut prev: u64 = 0; + let put_uvarint = |buf: &mut Vec, mut v: u64| loop { + let b = (v & 0x7f) as u8; + v >>= 7; + if v != 0 { + buf.push(b | 0x80); + } else { + buf.push(b); + break; + } + }; + let mut sorted = nonzero.to_vec(); + sorted.sort_by_key(|(i, _)| *i); + for (idx, val) in &sorted { + put_uvarint(&mut packed, idx - prev); + put_uvarint(&mut packed, *val as u64); + prev = *idx; + } + let state = HyperLogLogState { + variant: ProtoVariant::Regular as i32, + precision, + registers: Vec::new(), // dense field empty → sparse path + hip_kxq0: 0.0, + hip_kxq1: 0.0, + hip_est: 0.0, + // `num_registers` is informational — the decoder expands + // against `expected_len` from precision, not this field. + registers_sparse: Some(HllSparseRegisters { + num_registers: 1u32 << precision, + packed, + }), + }; + SketchEnvelope { + sketch_state: Some(sketch_envelope::SketchState::Hll(state)), + ..Default::default() + } + .encode_to_vec() + } + + #[test] + fn hll_from_proto_accepts_sparse_frame() { + // The consolidated decoder must accept the sparse wire form (the + // deleted dead decoder rejected it). Build a sparse frame setting + // a handful of registers, decode it, and confirm those register + // slots came back set in the dense array. + let precision = 12u32; + let nonzero = [(3u64, 5u8), (100, 2), (4000, 7)]; + let bytes = encode_hll_sparse(precision, &nonzero); + let sk = hll_from_proto(&bytes).expect("sparse HLL frame must decode (P2-3 regression)"); + assert_eq!(sk.registers.len(), 1usize << precision); + for (idx, val) in nonzero { + assert_eq!( + sk.registers[idx as usize], val, + "sparse register {idx} expanded to wrong value" + ); + } + } + + #[test] + fn hll_from_proto_matches_accumulator_decoder() { + // P2-4: the warm read path and the ingest accumulator must decode + // the SAME bytes to the SAME sketch (one source of truth). + use asap_physical_operators::summary_kernels::hll_sketch::HllSketchAccumulator; + let mut sk = HllSketch::new(HllVariant::Regular, 12); + for i in 0..500u64 { + sk.update(format!("item-{i}").as_bytes()); + } + let bytes = encode_hll_dense(&sk); + let via_delta = hll_from_proto(&bytes).expect("delta_apply hll decode"); + let via_acc = HllSketchAccumulator::from_sketchlib_proto_bytes(&bytes) + .expect("accumulator hll decode") + .inner; + assert_eq!( + via_delta.registers, via_acc.registers, + "delta_apply and accumulator must produce identical HLL registers" + ); + assert!((via_delta.estimate() - via_acc.estimate()).abs() < 1e-9); + } + + #[test] + fn dd_from_proto_matches_accumulator_decoder() { + use asap_physical_operators::summary_kernels::dd_sketch::DDSketchAccumulator; + let mut sk = DdSketch::new(0.01); + for v in [1.0, 2.0, 5.0, 5.0, 9.0, 42.0] { + sk.update(v); + } + let bytes = encode_dd(&sk); + let via_delta = dd_from_proto(&bytes).expect("delta_apply dd decode"); + let via_acc = DDSketchAccumulator::from_sketchlib_proto_bytes(&bytes) + .expect("accumulator dd decode") + .inner; + // Same quantile answers from the same bytes through both paths. + assert_eq!(via_delta.quantile(0.5), via_acc.quantile(0.5)); + assert_eq!(via_delta.quantile(0.99), via_acc.quantile(0.99)); + } + + #[test] + fn kll_from_proto_matches_accumulator_decoder() { + use asap_physical_operators::summary_kernels::datasketches_kll::DatasketchesKLLAccumulator; + let items: Vec = (0..200).map(|i| i as f64).collect(); + let bytes = encode_kll(256, &items); + let via_delta = kll_from_proto(&bytes).expect("delta_apply kll decode"); + let via_acc = DatasketchesKLLAccumulator::from_sketchlib_proto_bytes(&bytes) + .expect("accumulator kll decode") + .inner; + assert_eq!(via_delta.quantile(0.5), via_acc.quantile(0.5)); + } + + // ----------------------------------------------------------------- + // Per-window-reset (PWR) delta-apply regression tests. + // + // The edge resets its snapshot base at every window boundary, so a + // window's first frame is either a Full (window 1 / re-snapshot) or + // a Delta-from-empty (windows 2+). The query-side walk must: + // * reset the rolling base when `window_end` changes, + // * bootstrap an empty base for a window's leading Delta, + // * emit ONE value per window (the window's final accumulated + // state), never per-frame and never cross-window-accumulated. + // ----------------------------------------------------------------- + + fn full(bytes: Vec) -> SketchSampleState { + SketchSampleState { + bytes, + encoding: SketchEncoding::ProtoFull, + } + } + fn delta(bytes: Vec) -> SketchSampleState { + SketchSampleState { + bytes, + encoding: SketchEncoding::ProtoDelta, + } + } + + fn dd_over(alpha: f64, vals: &[f64]) -> DdSketch { + let mut sk = DdSketch::new(alpha); + for &v in vals { + sk.update(v); + } + sk + } + + /// A full re-snapshot replaces its pane's earlier frames; cumulative + /// readout must merge the finalized panes without counting updates twice. + #[test] + fn cumulative_readout_counts_resnapshot_population_once() { + let first = full(encode_dd(&dd_over(0.01, &[1., 2.]))); + let updated = full(encode_dd(&dd_over(0.01, &[1., 2., 3.]))); + let next = delta(encode_dd(&dd_over(0.01, &[9.]))); + let samples = [(1000, &first), (1000, &updated), (2000, &next)]; + let state = cumulative_summary_state(&samples, DeltaSketchKind::DDSketch { alpha: 0.01 }) + .unwrap() + .unwrap(); + let SummaryState::Dd(state) = state else { + panic!("expected DDSketch state"); + }; + assert_eq!(state.store_counts.iter().sum::(), 4); + } + + /// PWR across 3 windows: window 1 is `[Full]`, windows 2 & 3 are + /// `[Delta-from-empty]` (NO Full carry-in). Each window must + /// reconstruct its OWN distribution's median — not empty (the old + /// "skip delta with no base" bug) and not cross-window-inflated. + #[test] + fn pwr_ddsketch_three_windows_delta_from_empty() { + let alpha = 0.01; + let w1 = dd_over(alpha, &[1.0, 2.0, 3.0, 4.0, 5.0]); + let w2 = dd_over(alpha, &[10.0, 20.0, 30.0, 40.0, 50.0]); + let w3 = dd_over(alpha, &[100.0, 200.0, 300.0, 400.0, 500.0]); + + // window 1 ships a Full; windows 2+ ship a delta-from-empty. + let s1 = full(encode_dd(&w1)); + let s2 = delta(encode_dd(&w2)); + let s3 = delta(encode_dd(&w3)); + let samples = vec![(1000_i64, &s1), (2000, &s2), (3000, &s3)]; + + let kind = DeltaSketchKind::DDSketch { alpha }; + let (out, skipped) = + per_window_evaluate(&samples, kind, |rs| rs.quantile(0.5)).expect("pwr eval"); + assert_eq!(skipped, 0, "PWR must not skip delta-from-empty frames"); + assert_eq!(out.len(), 3, "one value per window"); + + // Each window's median ≈ that window's own distribution median, + // independent of the others (no carry-in inflation). + let truth = [ + w1.quantile(0.5).unwrap(), + w2.quantile(0.5).unwrap(), + w3.quantile(0.5).unwrap(), + ]; + for (i, (w_end, est)) in out.iter().enumerate() { + assert_eq!(*w_end, (i as i64 + 1) * 1000); + let rel = (est - truth[i]).abs() / truth[i].max(1e-9); + assert!( + rel < 0.05, + "window {i}: est={est} truth={} rel={rel}", + truth[i] + ); + } + // Cross-window-inflation guard: window 2's median must NOT have + // absorbed window 1 (would pull it well below 30). + assert!( + out[1].1 > 20.0, + "window 2 median {} suggests cross-window accumulation", + out[1].1 + ); + } + + /// Sub-window producer: a SINGLE window carries multiple frames + /// `[Full, Delta, Delta]`, where each later delta is an increment + /// since the previous emit in that window. The walk must COLLAPSE + /// them to ONE value = the window's running total, not emit three. + #[test] + fn pwr_ddsketch_subwindow_frames_collapse_to_window_total() { + let alpha = 0.01; + // Three sub-window increments that together cover 1..=15. + let a = dd_over(alpha, &[1.0, 2.0, 3.0, 4.0, 5.0]); + let b = dd_over(alpha, &[6.0, 7.0, 8.0, 9.0, 10.0]); + let c = dd_over(alpha, &[11.0, 12.0, 13.0, 14.0, 15.0]); + let s_a = full(encode_dd(&a)); + let s_b = delta(encode_dd(&b)); + let s_c = delta(encode_dd(&c)); + // All three share the same window_end (one window, sub-window frames). + let samples = vec![(5000_i64, &s_a), (5000, &s_b), (5000, &s_c)]; + + let kind = DeltaSketchKind::DDSketch { alpha }; + let (out, skipped) = + per_window_evaluate(&samples, kind, |rs| rs.quantile(0.5)).expect("subwindow eval"); + assert_eq!(skipped, 0); + assert_eq!(out.len(), 1, "sub-window frames collapse to ONE value"); + assert_eq!(out[0].0, 5000); + + let truth = dd_over(alpha, &(1..=15).map(|v| v as f64).collect::>()) + .quantile(0.5) + .unwrap(); + let rel = (out[0].1 - truth).abs() / truth.max(1e-9); + assert!(rel < 0.05, "window total est={} truth={truth}", out[0].1); + } + + /// Same sub-window collapse, but the window's FIRST frame is a + /// Delta-from-empty (PWR window 2+ with sub-window frames): + /// `[Delta-from-empty, Delta, Delta]`. + #[test] + fn pwr_ddsketch_subwindow_first_frame_delta_from_empty() { + let alpha = 0.01; + let a = dd_over(alpha, &[1.0, 2.0, 3.0, 4.0, 5.0]); + let b = dd_over(alpha, &[6.0, 7.0, 8.0, 9.0, 10.0]); + let c = dd_over(alpha, &[11.0, 12.0, 13.0, 14.0, 15.0]); + let s_a = delta(encode_dd(&a)); // first frame is delta-from-empty + let s_b = delta(encode_dd(&b)); + let s_c = delta(encode_dd(&c)); + let samples = vec![(9000_i64, &s_a), (9000, &s_b), (9000, &s_c)]; + + let kind = DeltaSketchKind::DDSketch { alpha }; + let (out, skipped) = + per_window_evaluate(&samples, kind, |rs| rs.quantile(0.5)).expect("eval"); + assert_eq!(skipped, 0); + assert_eq!(out.len(), 1); + let truth = dd_over(alpha, &(1..=15).map(|v| v as f64).collect::>()) + .quantile(0.5) + .unwrap(); + let rel = (out[0].1 - truth).abs() / truth.max(1e-9); + assert!(rel < 0.05, "est={} truth={truth}", out[0].1); + } + + /// PWR for HLL across 3 windows, each a Delta-from-empty (sparse + /// register delta). Bootstrapping an EMPTY HLL of the right precision + /// is required (register deltas index into a pre-sized array). Each + /// window's cardinality must reflect its OWN item set. + #[test] + fn pwr_hll_three_windows_delta_from_empty() { + let precision = 12u32; + // Build per-window HLLs, then encode each as a register-delta + // against an EMPTY sketch (= that window's full register state, + // the PWR delta-from-empty wire form). + let empty = HllSketch::new(HllVariant::Regular, precision); + let mut frames = Vec::new(); + let truths = [200usize, 800, 1500]; + for (w, &n) in truths.iter().enumerate() { + let mut sk = HllSketch::new(HllVariant::Regular, precision); + let base = (w as u64) * 100_000; // disjoint item sets per window + for i in 0..n as u64 { + sk.update(format!("u-{}", base + i).as_bytes()); + } + let bytes = sk.compute_delta(&empty, 0); + frames.push((((w as u64) + 1) * 1000, delta(bytes))); + } + let samples: Vec<(i64, &SketchSampleState)> = + frames.iter().map(|(t, s)| (*t as i64, s)).collect(); + + let kind = DeltaSketchKind::Hll { precision }; + let (out, skipped) = + per_window_evaluate(&samples, kind, |rs| rs.cardinality()).expect("hll pwr eval"); + assert_eq!(skipped, 0, "HLL delta-from-empty must bootstrap, not skip"); + assert_eq!(out.len(), 3); + for (i, (_w_end, est)) in out.iter().enumerate() { + let n = truths[i] as f64; + let rel = (est - n).abs() / n; + assert!( + rel < 0.15, + "window {i}: HLL est={est} truth={n} rel={rel} (each window independent)" + ); + } + } + + /// `CmsWithHeap` (min-over-rows estimator, `CountMinSketchWithHeap`) + /// and `CountSketchWithHeap` (median-of-signed-rows estimator, the + /// distinct `CountSketchWithHeap` type) are different sketch + /// algorithms that merely happen to share a storage shape — merging + /// one into the other must be rejected as a family mismatch, the + /// same as merging a `Cms` into a `Kll` would be. Since the two + /// `SummaryState` variants now hold genuinely different Rust types, + /// this is also enforced at compile time — there is no arm in + /// `merge_same_family` that type-checks a mixed pair together. + #[test] + fn cms_with_heap_and_count_sketch_with_heap_are_not_the_same_family() { + use asap_sketchlib::{CountMinSketchWithHeap, CountSketchWithHeap, MessagePackCodec}; + + let mut cms_heap = CountMinSketchWithHeap::new(4, 256, 10); + cms_heap.update("a", 1.0); + let mut cs_heap = CountSketchWithHeap::new(4, 256, 10); + cs_heap.update("b", 1.0); + + let mut a = SummaryState::CmsWithHeap( + CountMinSketchWithHeap::from_msgpack(&cms_heap.to_msgpack().unwrap()).unwrap(), + ); + let b = SummaryState::CountSketchWithHeap( + CountSketchWithHeap::from_msgpack(&cs_heap.to_msgpack().unwrap()).unwrap(), + ); + + match a.merge_same_family(&b) { + Err(msg) => assert!( + msg.contains("family mismatch"), + "expected a family-mismatch error, got: {msg}" + ), + Ok(()) => panic!( + "CmsWithHeap must not merge with CountSketchWithHeap -- \ + different algorithms sharing only a storage shape" + ), + } + } + + fn encode_delta_heap( + rows: u32, + cols: u32, + cells: &[(u32, u32, i64)], + heap: &[(&str, f64)], + heap_size: u64, + ) -> Vec { + #[derive(serde::Serialize)] + struct W<'a>( + bool, + (u32, u32, &'a [(u32, u32, i64)]), + Vec<(String, f64)>, + u64, + ); + let heap_owned: Vec<(String, f64)> = + heap.iter().map(|(k, v)| (k.to_string(), *v)).collect(); + let w = W(true, (rows, cols, cells), heap_owned, heap_size); + rmp_serde::to_vec(&w).expect("encode delta-heap") + } + + /// `SummaryState::CountSketchWithHeap` must decode both FULL and + /// DELTA-HEAP msgpack frames through the genuine + /// `asap_sketchlib::CountSketchWithHeap` (median-of-signed-rows + /// estimator) rather than the CMS-family `CountMinSketchWithHeap` + /// (min-over-rows estimator) it used to alias — the bug this split + /// fixed. Built via real `update()` calls (not a hand-crafted matrix) + /// so the sign-hashed row semantics are genuinely exercised, then + /// checks both decode paths reproduce the same matrix and the same + /// `estimate()` as the in-memory sketch they were encoded from. + #[test] + fn count_sketch_with_heap_full_and_delta_decode_via_new_asap_sketchlib_type() { + use asap_sketchlib::{CountSketchWithHeap, MessagePackCodec}; + + let mut built = CountSketchWithHeap::new(4, 64, 10); + for _ in 0..50 { + built.update("k", 1.0); + } + let expected_matrix = built.sketch_matrix(); + let expected_estimate = built.estimate("k"); + + // FULL path. + let full_bytes = built.to_msgpack().expect("encode full CountSketchWithHeap"); + let full_state = decode_full( + &DeltaSketchKind::CountSketchWithHeap { + rows: 4, + cols: 64, + heap_size: 10, + }, + &full_bytes, + SketchEncoding::MsgpackFull, + ) + .expect("decode_full CountSketchWithHeap"); + match full_state { + SummaryState::CountSketchWithHeap(inner) => { + assert_eq!(inner.sketch_matrix(), expected_matrix); + assert_eq!(inner.estimate("k"), expected_estimate); + } + other => panic!( + "expected CountSketchWithHeap state, got {}", + other.family_name() + ), + } + + // DELTA-HEAP path: same cells + heap against an empty base (PWR + // contract), encoded the way the Go producer does. + let cells: Vec<(u32, u32, i64)> = expected_matrix + .iter() + .enumerate() + .flat_map(|(r, row)| { + row.iter().enumerate().filter_map(move |(c, v)| { + if *v != 0.0 { + Some((r as u32, c as u32, *v as i64)) + } else { + None + } + }) + }) + .collect(); + let heap_pairs: Vec<(String, f64)> = built + .topk_heap_items() + .into_iter() + .map(|item| (item.key, item.value)) + .collect(); + assert!(!heap_pairs.is_empty(), "expected \"k\" in the top-k heap"); + let heap_refs: Vec<(&str, f64)> = + heap_pairs.iter().map(|(k, v)| (k.as_str(), *v)).collect(); + let delta_bytes = encode_delta_heap(4, 64, &cells, &heap_refs, 10); + + let mut rolling = DeltaSketchKind::CountSketchWithHeap { + rows: 4, + cols: 64, + heap_size: 10, + } + .bootstrap_empty(); + rolling + .apply_delta_bytes(&delta_bytes, SketchEncoding::MsgpackDelta) + .expect("apply CountSketchWithHeap delta"); + match rolling { + SummaryState::CountSketchWithHeap(inner) => { + assert_eq!( + inner.sketch_matrix(), + expected_matrix, + "delta path must reconstruct the identical matrix" + ); + assert_eq!(inner.estimate("k"), expected_estimate); + } + other => panic!( + "expected CountSketchWithHeap state, got {}", + other.family_name() + ), + } + } +} + +#[cfg(test)] +mod native_frame_tests { use super::*; - use asap_physical_operators::stored_state::{SketchEncoding, SketchSampleState}; + use crate::storage_engines::sketch_db::data::{SketchEncoding, SketchSampleState}; #[test] fn native_batches_are_not_legacy_sketch_frames() { diff --git a/data_plane/src/storage_engines/sketch_db/query/mod.rs b/data_plane/src/storage_engines/sketch_db/query/mod.rs index 5b4f64493..460779214 100644 --- a/data_plane/src/storage_engines/sketch_db/query/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/query/mod.rs @@ -6,6 +6,7 @@ pub mod asap_tier_result; pub mod decoders; pub mod delta_apply; +pub mod sketch_readout; pub mod timeline; pub mod timeline_dispatch; pub mod window_merger; diff --git a/data_plane/src/storage_engines/sketch_db/query/sketch_readout.rs b/data_plane/src/storage_engines/sketch_db/query/sketch_readout.rs new file mode 100644 index 000000000..8b568be0b --- /dev/null +++ b/data_plane/src/storage_engines/sketch_db/query/sketch_readout.rs @@ -0,0 +1,92 @@ +//! Planner-declared readouts over reconstructed SDS summary states. +use super::delta_apply::SummaryState; +use planner_types::{post_asap::SketchQuery, pre_asap::ColumnRef}; +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("{0}")] + Unsupported(&'static str), +} +pub fn sketch_query_value(rs: &SummaryState, query: &SketchQuery) -> Result { + if let SummaryState::UnivMon(state) = rs { + use asap_physical_operators::AggregateCore; + let statistic = match query { + SketchQuery::Cardinality => asap_physical_operators::Statistic::Cardinality, + SketchQuery::FrequencyL2 => asap_physical_operators::Statistic::FrequencyL2, + SketchQuery::FrequencyEntropy => asap_physical_operators::Statistic::FrequencyEntropy, + SketchQuery::PointCount { + key: ColumnRef::SampleValue, + value: None, + } => asap_physical_operators::Statistic::Count, + _ => return Err(Error::Unsupported("unsupported UnivMon readout")), + }; + return state + .query_statistic(statistic, &None, &Default::default()) + .map_err(|_| Error::Unsupported("UnivMon readout failed")); + } + match query { + SketchQuery::FrequencyL2 | SketchQuery::FrequencyEntropy => Err(Error::Unsupported( + "frequency moment readout requires UnivMon", + )), + SketchQuery::Quantile { q } => match rs { + // Typed PromQL/continuous-percentile readout uses interpolation; + // portable DDS `quantile` deliberately retains lower-rank parity. + SummaryState::Dd(sketch) => sketch.quantile_interpolated(*q).ok_or(Error::Unsupported( + "DDS interpolated quantile is unavailable", + )), + _ => Ok(rs.quantile(*q)), + }, + SketchQuery::Cardinality => Ok(rs.cardinality()), + // `key: ColumnRef::SampleValue, value: None` means "no specific + // item" -- the bare bucket total. `key: Named(_), value: Some(v)` + // is a per-item point lookup (e.g. `count(cms_metric{item="x"})`) + // -- `value` is where the filter's actual value lives (see + // `planner_types::post_asap::SketchQuery::PointCount`'s doc for why `readout` + // can't resolve it itself). Any other combination (e.g. a `Named` + // key with no value, or `SampleValue` with a value) is a shape + // this executor doesn't expect to see and reports rather than + // silently misreading. + SketchQuery::PointCount { + key: ColumnRef::SampleValue, + value: None, + } => Ok(rs.total()), + SketchQuery::PointCount { + key: ColumnRef::Named(_) | ColumnRef::Qualified { .. }, + value: Some(v), + } => rs.estimate(v).ok_or(Error::Unsupported( + "PointCount by key requires a Frequency-family sketch (Cms/CountSketch/..WithHeap)", + )), + SketchQuery::PointCount { .. } => Err(Error::Unsupported( + "unrecognized PointCount shape (key/value combination not expected)", + )), + // Both readout callers branch on `TopK` before ever calling this + // function (see `readout_cumulative`/`readout_per_window`), so + // this arm is unreachable in practice; kept for match + // exhaustiveness (`SketchQuery` has no `#[non_exhaustive]`) and to + // fail loudly rather than panic if that invariant is ever broken. + SketchQuery::TopK { .. } => Err(Error::Unsupported( + "TopK must be read out via topk_ranked, not sketch_query_value", + )), + } +} + +/// Rank a merged `SummaryState`'s top-k heap items descending by value and +/// cap at the requested `k`. The sort is load-bearing, not defensive +/// polish: `SummaryState::topk_items` reads back a bounded min-heap's +/// backing array as-is (`HHHeap::heap()`, asap_sketchlib) -- it does NOT +/// actually guarantee order despite its own doc wording. Errors for a +/// heap-less family (`Dd`/`Hll`/`Kll`/`Cms`/`CountSketch` -- no item +/// universe to rank), not for an empty heap (a heap-bearing family that +/// simply never received any updates yields `Ok(vec![])`, not an error). +pub fn topk_ranked(rs: &SummaryState, k: usize) -> Result, Error> { + let mut items = rs.topk_items().ok_or(Error::Unsupported( + "TopK requires a heap-bearing family (CmsWithHeap/CountSketchWithHeap) -- \ + this state's family carries no item universe to rank", + ))?; + items.sort_by(|a, b| { + b.1.partial_cmp(&a.1) + .unwrap_or(std::cmp::Ordering::Equal) + .then_with(|| a.0.cmp(&b.0)) // deterministic tie-break for equal counts + }); + items.truncate(k); + Ok(items) +} diff --git a/data_plane/src/tests/test_utilities/engine_factories.rs b/data_plane/src/tests/test_utilities/engine_factories.rs index 62a7621ec..ca044d6f4 100644 --- a/data_plane/src/tests/test_utilities/engine_factories.rs +++ b/data_plane/src/tests/test_utilities/engine_factories.rs @@ -92,6 +92,7 @@ pub fn create_engine_single_pop_with_aggregated( let agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), @@ -185,6 +186,7 @@ pub fn create_engine_dual_input( let value_agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: value_agg_type, aggregation_sub_type: String::new(), @@ -217,6 +219,7 @@ pub fn create_engine_dual_input( let keys_agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: key_agg_type, aggregation_sub_type: String::new(), @@ -319,6 +322,7 @@ pub fn create_engine_two_metrics( let agg_config_a = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: aggregation_type_a, aggregation_sub_type: String::new(), @@ -350,6 +354,7 @@ pub fn create_engine_two_metrics( let agg_config_b = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: aggregation_type_b, aggregation_sub_type: String::new(), @@ -462,6 +467,7 @@ pub fn create_engine_three_metrics( let cfg = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type: agg_type, aggregation_sub_type: String::new(), @@ -551,6 +557,7 @@ pub fn create_engine_multi_timestamp( let agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(), @@ -632,6 +639,7 @@ pub fn create_engine_multi_timestamp_with_window( let agg_config = PrecomputeMaterialization { stored_output_id: None, semantic_fragment: None, + dataset_identity: None, population_key_encoding: Default::default(), aggregation_type, aggregation_sub_type: String::new(),